Remove obsolete V3 MQ broker after SQLite retirement
This commit is contained in:
@@ -101,8 +101,9 @@
|
||||
- RPC 信任边界:录音服务原先复用旧 Dispatcher server 的可选 peer 校验函数;新 `peer_test.go` 先验证当前录音必须具备经 mTLS 验证的证书和预授权指纹,随后移入无可选关闭分支的独立校验函数。旧 Dispatcher 上传、拆分事件服务及只验证旧链路的测试已经删除;现行录音的双向 TLS 传输、共享最终结果及失败恢复测试保持通过。其余旧 Agent/Dispatcher 路径尚待清查,不将此批视为 P07 完成。
|
||||
- Dispatcher 旧任务消费分支:删除 `task_runtime_v3`、`task_queue_v3`、`task_control_v3` 及其专属测试和仅依赖这些旧测试夹具的旧 Mock 恢复测试;当前 `CurrentRuntime`、`CurrentBootstrap`、`CurrentDiscoveryFollower` 和根命令隔离测试仍独立通过。
|
||||
- Dispatcher 旧本地执行与控制分支:删除 `local_v01`、旧 Mock 呼叫/上传/恢复、旧路由与时间策略、旧 `Dispatcher`/lease 实现及依赖它们的测试;现行命令、任务消费、录音/结果链路的定向测试继续通过。
|
||||
- Agent 会话共用边界:保留当前根命令使用的 `AgentCoordinator` 注册、状态探测、激活、最新会话授权和过期拒绝;移除仅供旧执行链路使用的 `ExecuteRaw`、旧控制、查询补偿与旧 SIP applied-config 核验及其测试。当前 SIP 门禁仍由现行加载回报与隔离测试校验,不冒充 Asterisk 实际加载;旧 MQ/Store 和其余契约尚待清理,不能将此批视为 P07 完成。
|
||||
- Agent 会话共用边界:保留当前根命令使用的 `AgentCoordinator` 注册、状态探测、激活、最新会话授权和过期拒绝;移除仅供旧执行链路使用的 `ExecuteRaw`、旧控制、查询补偿与旧 SIP applied-config 核验及其测试。当前 SIP 门禁仍由现行加载回报与隔离测试校验,不冒充 Asterisk 实际加载;其他旧 MQ 实现、代次标记和契约尚待清理,不能将此批视为 P07 完成。
|
||||
- SQLite 代码归属:当前 Store 所需的 SQLite 驱动和 Dispatcher 身份/命令冲突错误已迁到当前入口;随后删除旧 Store 的源码、旧测试和旧迁移 SQL 文件。当前 Store 的独立建表、拒绝旧布局及外来 Dispatcher 抢占测试继续通过;这些操作仅移除仓库内过时源码,不读取、改写或清理现存 SQLite、spool 或 outbox 业务数据。
|
||||
- MQ 旧任务队列分支:旧 Store 的任务队列引用删除后,移除 `V3Broker` 和三份只验证旧拓扑的代码/测试;共享结果队列和无配置权限下的投递仍由当前 MQ 隔离测试验证。真实 SaaS/MQ 接收与应用收讫仍未验证。
|
||||
|
||||
## 验收台账
|
||||
|
||||
|
||||
@@ -1,416 +0,0 @@
|
||||
package mq
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"git.ipao.vip/rogee/go-sip/internal/tenant"
|
||||
amqp "github.com/rabbitmq/amqp091-go"
|
||||
)
|
||||
|
||||
const (
|
||||
CommandsExchangeV3 = "agent-call.dispatchers.v3"
|
||||
ResultsExchangeV3 = "agent-call.saas.v3"
|
||||
DeadLetterExchangeV3 = "agent-call.dead-letter.v3"
|
||||
MaxV3MessageBytes = 8 << 20
|
||||
)
|
||||
|
||||
type V3Broker struct {
|
||||
conn *amqp.Connection
|
||||
dispatcherID string
|
||||
prefetch int
|
||||
controlQueue string
|
||||
resultQueue string
|
||||
resultRoute string
|
||||
closed <-chan *amqp.Error
|
||||
mu sync.Mutex
|
||||
}
|
||||
|
||||
// OpenV3 verifies SaaS-provisioned topology passively. It never declares,
|
||||
// binds, or deletes exchanges or queues.
|
||||
func V3ControlQueueName(dispatcherID string) string {
|
||||
return "agent-call.d." + dispatcherID + ".control.v3"
|
||||
}
|
||||
|
||||
func OpenV3(url, dispatcherID string, prefetch int) (*V3Broker, error) {
|
||||
if strings.TrimSpace(url) == "" {
|
||||
return nil, errors.New("rabbitmq URL is required")
|
||||
}
|
||||
if prefetch <= 0 {
|
||||
return nil, errors.New("prefetch must be positive")
|
||||
}
|
||||
if err := tenant.ValidateDispatcherID(dispatcherID); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
conn, err := amqp.Dial(url)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("dial rabbitmq: %w", err)
|
||||
}
|
||||
channel, err := conn.Channel()
|
||||
if err != nil {
|
||||
_ = conn.Close()
|
||||
return nil, fmt.Errorf("open rabbitmq v3 channel: %w", err)
|
||||
}
|
||||
for _, exchange := range []string{CommandsExchangeV3, ResultsExchangeV3, DeadLetterExchangeV3} {
|
||||
if err := channel.ExchangeDeclarePassive(exchange, "topic", true, false, false, false, nil); err != nil {
|
||||
_ = channel.Close()
|
||||
_ = conn.Close()
|
||||
return nil, fmt.Errorf("required SaaS exchange %s unavailable: %w", exchange, err)
|
||||
}
|
||||
}
|
||||
controlQueue := V3ControlQueueName(dispatcherID)
|
||||
if _, err := channel.QueueDeclarePassive(controlQueue, true, false, false, false, nil); err != nil {
|
||||
_ = channel.Close()
|
||||
_ = conn.Close()
|
||||
return nil, fmt.Errorf("required SaaS control queue unavailable: %w", err)
|
||||
}
|
||||
resultQueue := "agent-call.saas.d." + dispatcherID + ".v3"
|
||||
if _, err := channel.QueueDeclarePassive(resultQueue, true, false, false, false, nil); err != nil {
|
||||
_ = channel.Close()
|
||||
_ = conn.Close()
|
||||
return nil, fmt.Errorf("required SaaS result queue unavailable: %w", err)
|
||||
}
|
||||
_ = channel.Close()
|
||||
return &V3Broker{
|
||||
conn: conn,
|
||||
dispatcherID: dispatcherID,
|
||||
prefetch: prefetch,
|
||||
controlQueue: controlQueue,
|
||||
resultQueue: resultQueue,
|
||||
resultRoute: "d." + dispatcherID + ".out",
|
||||
closed: conn.NotifyClose(make(chan *amqp.Error, 1)),
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (b *V3Broker) Close() error {
|
||||
b.mu.Lock()
|
||||
defer b.mu.Unlock()
|
||||
if b.conn == nil {
|
||||
return nil
|
||||
}
|
||||
err := b.conn.Close()
|
||||
b.conn = nil
|
||||
return err
|
||||
}
|
||||
|
||||
func (b *V3Broker) Done() <-chan *amqp.Error { return b.closed }
|
||||
|
||||
func (b *V3Broker) Publish(ctx context.Context, exchange, routingKey string, body []byte) error {
|
||||
if exchange != ResultsExchangeV3 || routingKey != b.resultRoute || len(body) == 0 || len(body) > MaxV3MessageBytes {
|
||||
return fmt.Errorf("v3 publish requires result exchange, Dispatcher route and 1..%d body bytes", MaxV3MessageBytes)
|
||||
}
|
||||
var identity struct {
|
||||
EventID string `json:"event_id"`
|
||||
MessageID string `json:"message_id"`
|
||||
}
|
||||
if err := json.Unmarshal(body, &identity); err != nil {
|
||||
return fmt.Errorf("decode outbound message identity: %w", err)
|
||||
}
|
||||
if (identity.EventID == "") == (identity.MessageID == "") {
|
||||
return errors.New("outbound message requires exactly one event_id or message_id")
|
||||
}
|
||||
messageID := identity.EventID + identity.MessageID
|
||||
if len(messageID) > 255 {
|
||||
return errors.New("outbound message identity exceeds AMQP limit")
|
||||
}
|
||||
ctx, cancel := context.WithTimeout(ctx, 10*time.Second)
|
||||
defer cancel()
|
||||
b.mu.Lock()
|
||||
defer b.mu.Unlock()
|
||||
if err := ctx.Err(); err != nil {
|
||||
return err
|
||||
}
|
||||
if b.conn == nil || b.conn.IsClosed() {
|
||||
return errors.New("rabbitmq v3 connection is closed")
|
||||
}
|
||||
channel, err := b.conn.Channel()
|
||||
if err != nil {
|
||||
return fmt.Errorf("open v3 publication channel: %w", err)
|
||||
}
|
||||
defer channel.Close()
|
||||
if _, err := channel.QueueDeclarePassive(b.resultQueue, true, false, false, false, nil); err != nil {
|
||||
return fmt.Errorf("required SaaS result queue unavailable: %w", err)
|
||||
}
|
||||
if err := channel.Confirm(false); err != nil {
|
||||
return fmt.Errorf("enable v3 publisher confirms: %w", err)
|
||||
}
|
||||
returned := channel.NotifyReturn(make(chan amqp.Return, 1))
|
||||
confirmation, err := channel.PublishWithDeferredConfirmWithContext(ctx, exchange, routingKey, true, false, amqp.Publishing{
|
||||
ContentType: "application/json", DeliveryMode: amqp.Persistent, Body: body, MessageId: messageID,
|
||||
})
|
||||
if err != nil {
|
||||
return fmt.Errorf("publish v3 result: %w", err)
|
||||
}
|
||||
if confirmation == nil {
|
||||
return errors.New("rabbitmq v3 publisher confirmation unavailable")
|
||||
}
|
||||
acked, err := confirmation.WaitContext(ctx)
|
||||
if err != nil {
|
||||
return fmt.Errorf("wait for v3 publisher confirmation: %w", err)
|
||||
}
|
||||
select {
|
||||
case result, ok := <-returned:
|
||||
if !ok {
|
||||
return errors.New("v3 publication channel closed before routing was established")
|
||||
}
|
||||
return fmt.Errorf("v3 publication returned: code=%d", result.ReplyCode)
|
||||
default:
|
||||
}
|
||||
if !acked {
|
||||
return errors.New("rabbitmq v3 publisher was negatively acknowledged")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
type V3Consumer struct {
|
||||
cancel context.CancelFunc
|
||||
done chan struct{}
|
||||
mu sync.Mutex
|
||||
err error
|
||||
}
|
||||
|
||||
func (c *V3Consumer) Wait(ctx context.Context) error {
|
||||
select {
|
||||
case <-c.done:
|
||||
c.mu.Lock()
|
||||
defer c.mu.Unlock()
|
||||
return c.err
|
||||
case <-ctx.Done():
|
||||
return ctx.Err()
|
||||
}
|
||||
}
|
||||
|
||||
func (c *V3Consumer) Stop(ctx context.Context) error {
|
||||
c.cancel()
|
||||
err := c.Wait(ctx)
|
||||
if errors.Is(err, context.Canceled) {
|
||||
return nil
|
||||
}
|
||||
return err
|
||||
}
|
||||
|
||||
func (b *V3Broker) ConsumePredeclared(ctx context.Context, queue string, handler MessageHandler) error {
|
||||
consumer, err := b.StartPredeclaredConsumer(ctx, queue, handler)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
err = consumer.Wait(context.Background())
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return ctx.Err()
|
||||
}
|
||||
|
||||
// StartPredeclaredConsumer subscribes only to an existing SaaS-owned queue.
|
||||
// Stopping it closes the channel so every unacknowledged delivery is requeued.
|
||||
func (b *V3Broker) StartPredeclaredConsumer(ctx context.Context, queue string, handler MessageHandler) (*V3Consumer, error) {
|
||||
if queue == "" || handler == nil {
|
||||
return nil, errors.New("predeclared queue and handler are required")
|
||||
}
|
||||
if len(queue) > 255 {
|
||||
return nil, errors.New("queue name exceeds AMQP limit")
|
||||
}
|
||||
if err := ctx.Err(); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
b.mu.Lock()
|
||||
conn, prefetch, dispatcherID := b.conn, b.prefetch, b.dispatcherID
|
||||
closed := conn == nil || conn.IsClosed()
|
||||
b.mu.Unlock()
|
||||
if closed {
|
||||
return nil, errors.New("rabbitmq v3 connection is closed")
|
||||
}
|
||||
channel, err := conn.Channel()
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("open v3 consumer channel: %w", err)
|
||||
}
|
||||
if _, err := channel.QueueDeclarePassive(queue, true, false, false, false, nil); err != nil {
|
||||
_ = channel.Close()
|
||||
return nil, fmt.Errorf("required SaaS-owned queue %s unavailable: %w", queue, err)
|
||||
}
|
||||
if err := channel.Qos(prefetch, 0, false); err != nil {
|
||||
_ = channel.Close()
|
||||
return nil, fmt.Errorf("set v3 consumer prefetch: %w", err)
|
||||
}
|
||||
consumerTag := fmt.Sprintf("sip-go-agent-v3-%d", consumerSequence.Add(1))
|
||||
deliveries, err := channel.Consume(queue, consumerTag, false, false, false, false, nil)
|
||||
if err != nil {
|
||||
_ = channel.Close()
|
||||
return nil, fmt.Errorf("consume SaaS-owned v3 queue: %w", err)
|
||||
}
|
||||
consumerCtx, cancel := context.WithCancel(ctx)
|
||||
consumer := &V3Consumer{cancel: cancel, done: make(chan struct{})}
|
||||
go func() {
|
||||
consumeErr := b.consumeV3Deliveries(consumerCtx, dispatcherID, deliveries, handler)
|
||||
if errors.Is(consumeErr, context.Canceled) && consumerCtx.Err() != nil {
|
||||
consumeErr = nil
|
||||
}
|
||||
if err := channel.Cancel(consumerTag, false); err != nil {
|
||||
consumeErr = errors.Join(consumeErr, fmt.Errorf("cancel v3 consumer: %w", err))
|
||||
}
|
||||
if err := channel.Close(); err != nil {
|
||||
consumeErr = errors.Join(consumeErr, fmt.Errorf("close v3 consumer channel: %w", err))
|
||||
}
|
||||
consumer.mu.Lock()
|
||||
consumer.err = consumeErr
|
||||
consumer.mu.Unlock()
|
||||
close(consumer.done)
|
||||
}()
|
||||
return consumer, nil
|
||||
}
|
||||
|
||||
func (b *V3Broker) DrainPredeclared(ctx context.Context, queue string) (int, error) {
|
||||
if queue == "" || len(queue) > 255 {
|
||||
return 0, errors.New("valid SaaS-owned queue name is required")
|
||||
}
|
||||
if err := ctx.Err(); err != nil {
|
||||
return 0, err
|
||||
}
|
||||
b.mu.Lock()
|
||||
conn := b.conn
|
||||
closed := conn == nil || conn.IsClosed()
|
||||
b.mu.Unlock()
|
||||
if closed {
|
||||
return 0, errors.New("rabbitmq v3 connection is closed")
|
||||
}
|
||||
channel, err := conn.Channel()
|
||||
if err != nil {
|
||||
return 0, fmt.Errorf("open v3 drain channel: %w", err)
|
||||
}
|
||||
defer channel.Close()
|
||||
if _, err := channel.QueueDeclarePassive(queue, true, false, false, false, nil); err != nil {
|
||||
return 0, fmt.Errorf("required SaaS-owned queue %s unavailable: %w", queue, err)
|
||||
}
|
||||
drained := 0
|
||||
for {
|
||||
if err := ctx.Err(); err != nil {
|
||||
return drained, err
|
||||
}
|
||||
delivery, ok, err := channel.Get(queue, false)
|
||||
if err != nil {
|
||||
return drained, fmt.Errorf("read SaaS-owned queue %s for stop drain: %w", queue, err)
|
||||
}
|
||||
if !ok {
|
||||
return drained, nil
|
||||
}
|
||||
if err := delivery.Ack(false); err != nil {
|
||||
return drained, fmt.Errorf("ack stopped task backlog: %w", err)
|
||||
}
|
||||
drained++
|
||||
}
|
||||
}
|
||||
|
||||
// DrainControlPredeclared processes the already queued controls before task admission.
|
||||
// Unlike stopped task backlogs, control deliveries must pass through the handler
|
||||
// before ACK; a transient failure is requeued and closes startup admission.
|
||||
func (b *V3Broker) DrainControlPredeclared(ctx context.Context, queue string, handler MessageHandler) (int, error) {
|
||||
if queue == "" || queue != b.controlQueue || handler == nil {
|
||||
return 0, errors.New("configured SaaS-owned control queue and handler are required")
|
||||
}
|
||||
if err := ctx.Err(); err != nil {
|
||||
return 0, err
|
||||
}
|
||||
b.mu.Lock()
|
||||
conn := b.conn
|
||||
closed := conn == nil || conn.IsClosed()
|
||||
b.mu.Unlock()
|
||||
if closed {
|
||||
return 0, errors.New("rabbitmq v3 connection is closed")
|
||||
}
|
||||
channel, err := conn.Channel()
|
||||
if err != nil {
|
||||
return 0, fmt.Errorf("open control backlog channel: %w", err)
|
||||
}
|
||||
defer channel.Close()
|
||||
if _, err := channel.QueueDeclarePassive(queue, true, false, false, false, nil); err != nil {
|
||||
return 0, fmt.Errorf("required SaaS-owned control queue %s unavailable: %w", queue, err)
|
||||
}
|
||||
processed := 0
|
||||
for {
|
||||
if err := ctx.Err(); err != nil {
|
||||
return processed, err
|
||||
}
|
||||
delivery, ok, err := channel.Get(queue, false)
|
||||
if err != nil {
|
||||
return processed, fmt.Errorf("read SaaS-owned control backlog %s: %w", queue, err)
|
||||
}
|
||||
if !ok {
|
||||
return processed, nil
|
||||
}
|
||||
if len(delivery.Body) == 0 || len(delivery.Body) > MaxV3MessageBytes || !json.Valid(delivery.Body) {
|
||||
if err := delivery.Reject(false); err != nil {
|
||||
return processed, fmt.Errorf("reject invalid control backlog message: %w", err)
|
||||
}
|
||||
slog.Warn("MQ control backlog message rejected", "dispatcher_id", b.dispatcherID, "delivery_tag", delivery.DeliveryTag, "reason", "invalid_size_or_json", "bytes", len(delivery.Body))
|
||||
processed++
|
||||
continue
|
||||
}
|
||||
if err := handler(ctx, delivery.RoutingKey, delivery.Body); err != nil {
|
||||
if IsPermanent(err) {
|
||||
if rejectErr := delivery.Reject(false); rejectErr != nil {
|
||||
return processed, fmt.Errorf("control handler error %v; reject: %w", err, rejectErr)
|
||||
}
|
||||
processed++
|
||||
continue
|
||||
}
|
||||
if nackErr := delivery.Nack(false, true); nackErr != nil {
|
||||
return processed, fmt.Errorf("control handler error %v; nack: %w", err, nackErr)
|
||||
}
|
||||
return processed, fmt.Errorf("control backlog handler failed, admission stays closed: %w", err)
|
||||
}
|
||||
if err := delivery.Ack(false); err != nil {
|
||||
return processed, fmt.Errorf("ack processed control backlog: %w", err)
|
||||
}
|
||||
processed++
|
||||
}
|
||||
}
|
||||
|
||||
func (b *V3Broker) ControlQueue() string { return b.controlQueue }
|
||||
|
||||
func (b *V3Broker) consumeV3Deliveries(ctx context.Context, dispatcherID string, deliveries <-chan amqp.Delivery, handler MessageHandler) error {
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return ctx.Err()
|
||||
case delivery, ok := <-deliveries:
|
||||
if !ok {
|
||||
return errors.New("rabbitmq v3 delivery channel closed")
|
||||
}
|
||||
if len(delivery.Body) == 0 || len(delivery.Body) > MaxV3MessageBytes {
|
||||
if err := delivery.Reject(false); err != nil {
|
||||
return fmt.Errorf("reject invalid v3 message size: %w", err)
|
||||
}
|
||||
slog.Warn("MQ v3 message rejected", "dispatcher_id", dispatcherID, "delivery_tag", delivery.DeliveryTag, "reason", "invalid_message_size", "bytes", len(delivery.Body))
|
||||
continue
|
||||
}
|
||||
if !json.Valid(delivery.Body) {
|
||||
if err := delivery.Reject(false); err != nil {
|
||||
return fmt.Errorf("reject invalid v3 JSON: %w", err)
|
||||
}
|
||||
slog.Warn("MQ v3 message rejected", "dispatcher_id", dispatcherID, "delivery_tag", delivery.DeliveryTag, "reason", "invalid_json")
|
||||
continue
|
||||
}
|
||||
if err := handler(ctx, delivery.RoutingKey, delivery.Body); err != nil {
|
||||
if IsPermanent(err) {
|
||||
if rejectErr := delivery.Reject(false); rejectErr != nil {
|
||||
return fmt.Errorf("permanent v3 handler error %v; reject: %w", err, rejectErr)
|
||||
}
|
||||
continue
|
||||
}
|
||||
if nackErr := delivery.Nack(false, true); nackErr != nil {
|
||||
return fmt.Errorf("v3 handler error %v; nack: %w", err, nackErr)
|
||||
}
|
||||
continue
|
||||
}
|
||||
if err := delivery.Ack(false); err != nil {
|
||||
return fmt.Errorf("ack v3 delivery: %w", err)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,449 +0,0 @@
|
||||
//go:build integration
|
||||
|
||||
package mq
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"os"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
amqp "github.com/rabbitmq/amqp091-go"
|
||||
)
|
||||
|
||||
func TestV3BrokerConsumesSaaSOwnedQueueAndPublishesConfirmedResult(t *testing.T) {
|
||||
url := os.Getenv("RABBITMQ_URL")
|
||||
if url == "" {
|
||||
t.Skip("RABBITMQ_URL not set")
|
||||
}
|
||||
provisionerURL := os.Getenv("RABBITMQ_PROVISIONER_URL")
|
||||
if provisionerURL == "" {
|
||||
provisionerURL = url
|
||||
}
|
||||
const dispatcherID = "550e8400-e29b-41d4-a716-446655440000"
|
||||
const taskID = "task-1"
|
||||
const routingKey = "d." + dispatcherID + ".task." + taskID + ".in"
|
||||
const taskQueue = "agent-call.d." + dispatcherID + ".task." + taskID + ".v3"
|
||||
const controlQueue = "agent-call.d." + dispatcherID + ".control.v3"
|
||||
const controlRoute = "d." + dispatcherID + ".control.in"
|
||||
const resultRoute = "d." + dispatcherID + ".out"
|
||||
const resultQueue = "agent-call.saas.d." + dispatcherID + ".v3"
|
||||
const deadRoute = "d." + dispatcherID + ".dead-letter"
|
||||
const deadQueue = "agent-call.d." + dispatcherID + ".dead-letter.v3"
|
||||
|
||||
mockConn, err := amqp.Dial(provisionerURL)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer mockConn.Close()
|
||||
mockChannel, err := mockConn.Channel()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer mockChannel.Close()
|
||||
for _, exchange := range []string{CommandsExchangeV3, ResultsExchangeV3, DeadLetterExchangeV3} {
|
||||
if err := mockChannel.ExchangeDeclare(exchange, "topic", true, false, false, false, nil); err != nil {
|
||||
t.Fatalf("Mock SaaS declare exchange %s: %v", exchange, err)
|
||||
}
|
||||
}
|
||||
if _, err := mockChannel.QueueDeclare(deadQueue, true, false, false, false, nil); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := mockChannel.QueueBind(deadQueue, deadRoute, DeadLetterExchangeV3, false, nil); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := mockChannel.QueueDeclare(taskQueue, true, false, false, false, amqp.Table{
|
||||
"x-dead-letter-exchange": DeadLetterExchangeV3,
|
||||
"x-dead-letter-routing-key": deadRoute,
|
||||
}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := mockChannel.QueueBind(taskQueue, routingKey, CommandsExchangeV3, false, nil); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := mockChannel.QueueDeclare(controlQueue, true, false, false, false, nil); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := mockChannel.QueueBind(controlQueue, controlRoute, CommandsExchangeV3, false, nil); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := mockChannel.QueueDeclare(resultQueue, true, false, false, false, nil); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := mockChannel.QueueBind(resultQueue, resultRoute, ResultsExchangeV3, false, nil); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for _, queue := range []string{taskQueue, controlQueue, resultQueue, deadQueue} {
|
||||
if _, err := mockChannel.QueuePurge(queue, false); err != nil {
|
||||
t.Fatalf("purge Mock-owned queue %s: %v", queue, err)
|
||||
}
|
||||
defer func(queue string) { _, _ = mockChannel.QueueDelete(queue, false, false, false) }(queue)
|
||||
}
|
||||
|
||||
if os.Getenv("RABBITMQ_EXPECT_NO_CONFIG") == "1" {
|
||||
forbiddenQueue := fmt.Sprintf("agent-call.d.%s.task.forbidden-%d.v3", dispatcherID, time.Now().UnixNano())
|
||||
dispatcherConn, err := amqp.Dial(url)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
dispatcherChannel, err := dispatcherConn.Channel()
|
||||
if err != nil {
|
||||
_ = dispatcherConn.Close()
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := dispatcherChannel.QueueDeclare(forbiddenQueue, true, false, false, false, nil); err == nil {
|
||||
_, _ = mockChannel.QueueDelete(forbiddenQueue, false, false, false)
|
||||
_ = dispatcherChannel.Close()
|
||||
_ = dispatcherConn.Close()
|
||||
t.Fatal("Dispatcher identity unexpectedly has queue configure permission")
|
||||
}
|
||||
_ = dispatcherChannel.Close()
|
||||
_ = dispatcherConn.Close()
|
||||
checkChannel, err := mockConn.Channel()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := checkChannel.QueueDeclarePassive(forbiddenQueue, true, false, false, false, nil); err == nil {
|
||||
_ = checkChannel.Close()
|
||||
_, _ = mockChannel.QueueDelete(forbiddenQueue, false, false, false)
|
||||
t.Fatal("unauthorized Dispatcher queue declaration created a queue")
|
||||
}
|
||||
_ = checkChannel.Close()
|
||||
}
|
||||
|
||||
broker, err := OpenV3(url, dispatcherID, DefaultPrefetch)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer broker.Close()
|
||||
if broker.ControlQueue() != controlQueue {
|
||||
t.Fatalf("wrong SaaS-owned control queue: %q", broker.ControlQueue())
|
||||
}
|
||||
oldStop := []byte(`{"action":"stop","issued_at":"2025-01-01T00:00:00Z"}`)
|
||||
publishControl := func(body []byte) {
|
||||
t.Helper()
|
||||
if err := mockChannel.PublishWithContext(context.Background(), CommandsExchangeV3, controlRoute, true, false, amqp.Publishing{
|
||||
ContentType: "application/json", DeliveryMode: amqp.Persistent, Body: body,
|
||||
}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
publishControl(oldStop)
|
||||
attempts := 0
|
||||
processed, err := broker.DrainControlPredeclared(context.Background(), controlQueue, func(_ context.Context, key string, body []byte) error {
|
||||
attempts++
|
||||
if key != controlRoute || string(body) != string(oldStop) {
|
||||
return fmt.Errorf("wrong historical control: key=%s body=%s", key, body)
|
||||
}
|
||||
return errors.New("mock SQLite write failure")
|
||||
})
|
||||
if err == nil || processed != 0 || attempts != 1 {
|
||||
t.Fatalf("transient control failure was ACKed or not surfaced: processed=%d attempts=%d err=%v", processed, attempts, err)
|
||||
}
|
||||
processed, err = broker.DrainControlPredeclared(context.Background(), controlQueue, func(_ context.Context, key string, body []byte) error {
|
||||
attempts++
|
||||
if key != controlRoute || string(body) != string(oldStop) {
|
||||
return fmt.Errorf("old stop not redelivered: key=%s body=%s", key, body)
|
||||
}
|
||||
return nil
|
||||
})
|
||||
if err != nil || processed != 1 || attempts != 2 {
|
||||
t.Fatalf("old stop did not recover under its original identity: processed=%d attempts=%d err=%v", processed, attempts, err)
|
||||
}
|
||||
publishControl([]byte(`{`))
|
||||
publishControl([]byte(`{"action":"invalid-control"}`))
|
||||
processed, err = broker.DrainControlPredeclared(context.Background(), controlQueue, func(_ context.Context, key string, body []byte) error {
|
||||
if key != controlRoute || string(body) != `{"action":"invalid-control"}` {
|
||||
return fmt.Errorf("invalid JSON delivered to control handler: key=%s body=%s", key, body)
|
||||
}
|
||||
return Permanent(errors.New("invalid control contract"))
|
||||
})
|
||||
if err != nil || processed != 2 {
|
||||
t.Fatalf("malformed and permanently invalid controls were not rejected: processed=%d err=%v", processed, err)
|
||||
}
|
||||
if state, err := mockChannel.QueueInspect(controlQueue); err != nil || state.Messages != 0 {
|
||||
t.Fatalf("control backlog did not drain: messages=%d err=%v", state.Messages, err)
|
||||
}
|
||||
missingQueue := "agent-call.d." + dispatcherID + ".task.missing.v3"
|
||||
missingCtx, missingCancel := context.WithTimeout(context.Background(), time.Second)
|
||||
if err := broker.ConsumePredeclared(missingCtx, missingQueue, func(context.Context, string, []byte) error { return nil }); err == nil {
|
||||
missingCancel()
|
||||
t.Fatal("missing task queue should fail passive declaration")
|
||||
}
|
||||
missingCancel()
|
||||
checkChannel, err := mockConn.Channel()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := checkChannel.QueueDeclarePassive(missingQueue, true, false, false, false, nil); err == nil {
|
||||
_ = checkChannel.Close()
|
||||
t.Fatal("Dispatcher created a missing SaaS-owned task queue")
|
||||
}
|
||||
_ = checkChannel.Close()
|
||||
|
||||
commandBody := []byte(`{"event_id":"command-event-1"}`)
|
||||
if err := mockChannel.PublishWithContext(context.Background(), CommandsExchangeV3, routingKey, true, false, amqp.Publishing{
|
||||
ContentType: "application/json", DeliveryMode: amqp.Persistent, Body: commandBody,
|
||||
}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
consumed := make(chan []byte, 1)
|
||||
consumeErr := make(chan error, 1)
|
||||
go func() {
|
||||
consumeErr <- broker.ConsumePredeclared(ctx, taskQueue, func(_ context.Context, key string, body []byte) error {
|
||||
if key != routingKey {
|
||||
return fmt.Errorf("routing key = %q, want %q", key, routingKey)
|
||||
}
|
||||
consumed <- append([]byte(nil), body...)
|
||||
return nil
|
||||
})
|
||||
}()
|
||||
select {
|
||||
case got := <-consumed:
|
||||
if string(got) != string(commandBody) {
|
||||
t.Fatalf("consumed body = %s, want %s", got, commandBody)
|
||||
}
|
||||
case err := <-consumeErr:
|
||||
cancel()
|
||||
t.Fatalf("consume task queue: %v", err)
|
||||
case <-time.After(5 * time.Second):
|
||||
cancel()
|
||||
t.Fatal("timed out waiting for command delivery")
|
||||
}
|
||||
cancel()
|
||||
select {
|
||||
case <-consumeErr:
|
||||
case <-time.After(5 * time.Second):
|
||||
t.Fatal("consumer did not stop after cancellation")
|
||||
}
|
||||
if leftover, ok, err := mockChannel.Get(taskQueue, false); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if ok {
|
||||
if string(leftover.Body) != string(commandBody) {
|
||||
t.Fatalf("unexpected command after first consumer stopped: %s", leftover.Body)
|
||||
}
|
||||
if err := leftover.Ack(false); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
|
||||
// A transient SQLite write error must NACK the same original delivery;
|
||||
// prefetch=1 and a later barrier prove the retry was ACKed only after
|
||||
// the handler recovered.
|
||||
transientBody := []byte(`{"event_id":"sqlite-write-failed-original"}`)
|
||||
if err := mockChannel.PublishWithContext(context.Background(), CommandsExchangeV3, routingKey, true, false, amqp.Publishing{
|
||||
ContentType: "application/json", DeliveryMode: amqp.Persistent, Body: transientBody,
|
||||
}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
retryCtx, retryCancel := context.WithCancel(context.Background())
|
||||
failed := make(chan struct{}, 1)
|
||||
retryEntered := make(chan struct{}, 1)
|
||||
releaseRetry := make(chan struct{})
|
||||
barrierSeen := make(chan struct{}, 1)
|
||||
retryBarrierBody := []byte(`{"event_id":"sqlite-recovery-barrier"}`)
|
||||
attempts = 0
|
||||
retryConsumer, err := broker.StartPredeclaredConsumer(retryCtx, taskQueue, func(_ context.Context, _ string, body []byte) error {
|
||||
if string(body) == string(retryBarrierBody) {
|
||||
barrierSeen <- struct{}{}
|
||||
return nil
|
||||
}
|
||||
if string(body) != string(transientBody) {
|
||||
return fmt.Errorf("unexpected redelivery: %s", body)
|
||||
}
|
||||
attempts++
|
||||
if attempts == 1 {
|
||||
failed <- struct{}{}
|
||||
return fmt.Errorf("injected SQLite write failure")
|
||||
}
|
||||
if attempts != 2 {
|
||||
return fmt.Errorf("original command delivered %d times", attempts)
|
||||
}
|
||||
retryEntered <- struct{}{}
|
||||
<-releaseRetry
|
||||
return nil
|
||||
})
|
||||
if err != nil {
|
||||
retryCancel()
|
||||
t.Fatal(err)
|
||||
}
|
||||
for _, signal := range []<-chan struct{}{failed, retryEntered} {
|
||||
select {
|
||||
case <-signal:
|
||||
case <-time.After(5 * time.Second):
|
||||
t.Fatal("transient SQLite write failure did not redeliver the original task")
|
||||
}
|
||||
}
|
||||
if err := mockChannel.PublishWithContext(context.Background(), CommandsExchangeV3, routingKey, true, false, amqp.Publishing{
|
||||
ContentType: "application/json", DeliveryMode: amqp.Persistent, Body: retryBarrierBody,
|
||||
}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
close(releaseRetry)
|
||||
select {
|
||||
case <-barrierSeen:
|
||||
case <-time.After(5 * time.Second):
|
||||
t.Fatal("recovered command was not ACKed before the barrier")
|
||||
}
|
||||
if err := retryConsumer.Stop(context.Background()); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
retryCancel()
|
||||
if leftover, ok, err := mockChannel.Get(taskQueue, false); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if ok {
|
||||
if string(leftover.Body) != string(retryBarrierBody) {
|
||||
t.Fatalf("original command was requeued after its ACK: %s", leftover.Body)
|
||||
}
|
||||
if err := leftover.Ack(false); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
if attempts != 2 {
|
||||
t.Fatalf("expected one failed and one successful delivery, got %d", attempts)
|
||||
}
|
||||
|
||||
pauseBody := []byte(`{"event_id":"pause-requeue"}`)
|
||||
if err := mockChannel.PublishWithContext(context.Background(), CommandsExchangeV3, routingKey, true, false, amqp.Publishing{
|
||||
ContentType: "application/json", DeliveryMode: amqp.Persistent, Body: pauseBody,
|
||||
}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for i := 1; i < 100; i++ {
|
||||
body := []byte(fmt.Sprintf(`{"event_id":"pause-backlog-%d"}`, i))
|
||||
if err := mockChannel.PublishWithContext(context.Background(), CommandsExchangeV3, routingKey, true, false, amqp.Publishing{
|
||||
ContentType: "application/json", DeliveryMode: amqp.Persistent, Body: body,
|
||||
}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
pauseCtx, pauseCancel := context.WithCancel(context.Background())
|
||||
delivered := make(chan struct{}, 1)
|
||||
consumer, err := broker.StartPredeclaredConsumer(pauseCtx, taskQueue, func(ctx context.Context, _ string, _ []byte) error {
|
||||
delivered <- struct{}{}
|
||||
<-ctx.Done()
|
||||
return ctx.Err()
|
||||
})
|
||||
if err != nil {
|
||||
pauseCancel()
|
||||
t.Fatal(err)
|
||||
}
|
||||
select {
|
||||
case <-delivered:
|
||||
case <-time.After(5 * time.Second):
|
||||
_ = consumer.Stop(context.Background())
|
||||
pauseCancel()
|
||||
t.Fatal("timed out waiting for pause-test delivery")
|
||||
}
|
||||
if err := consumer.Stop(context.Background()); err != nil {
|
||||
pauseCancel()
|
||||
t.Fatalf("stop paused consumer: %v", err)
|
||||
}
|
||||
pauseCancel()
|
||||
queued, err := mockChannel.QueueInspect(taskQueue)
|
||||
if err != nil || queued.Messages != 100 {
|
||||
t.Fatalf("pause did not preserve all 100 task commands: count=%d err=%v", queued.Messages, err)
|
||||
}
|
||||
// Prefetch=1 means seeing this last barrier proves all 100 preceding
|
||||
// deliveries were ACKed before the resumed consumer can be stopped.
|
||||
barrierBody := []byte(`{"event_id":"resume-barrier"}`)
|
||||
if err := mockChannel.PublishWithContext(context.Background(), CommandsExchangeV3, routingKey, true, false, amqp.Publishing{
|
||||
ContentType: "application/json", DeliveryMode: amqp.Persistent, Body: barrierBody,
|
||||
}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
resumeCtx, resumeCancel := context.WithCancel(context.Background())
|
||||
resumed := make(chan string, 100)
|
||||
barrier := make(chan struct{}, 1)
|
||||
resumedConsumer, err := broker.StartPredeclaredConsumer(resumeCtx, taskQueue, func(_ context.Context, _ string, body []byte) error {
|
||||
if string(body) == string(barrierBody) {
|
||||
barrier <- struct{}{}
|
||||
return nil
|
||||
}
|
||||
resumed <- string(body)
|
||||
return nil
|
||||
})
|
||||
if err != nil {
|
||||
resumeCancel()
|
||||
t.Fatal(err)
|
||||
}
|
||||
seen := make(map[string]bool, 100)
|
||||
for len(seen) < 100 {
|
||||
select {
|
||||
case body := <-resumed:
|
||||
if seen[body] {
|
||||
t.Fatalf("resumed task command delivered twice: %s", body)
|
||||
}
|
||||
seen[body] = true
|
||||
case <-time.After(15 * time.Second):
|
||||
t.Fatalf("resumed only %d of 100 original task commands", len(seen))
|
||||
}
|
||||
}
|
||||
if !seen[string(pauseBody)] {
|
||||
t.Fatal("paused in-flight task was not resumed from its original queue")
|
||||
}
|
||||
select {
|
||||
case <-barrier:
|
||||
case <-time.After(5 * time.Second):
|
||||
t.Fatal("resumed task ACK barrier was not reached")
|
||||
}
|
||||
if err := resumedConsumer.Stop(context.Background()); err != nil {
|
||||
resumeCancel()
|
||||
t.Fatal(err)
|
||||
}
|
||||
resumeCancel()
|
||||
// Cancellation can requeue the barrier itself; it is not one of the 100
|
||||
// task commands and cannot cause another task execution.
|
||||
if last, ok, err := mockChannel.Get(taskQueue, false); err != nil {
|
||||
t.Fatal(err)
|
||||
} else if ok {
|
||||
if string(last.Body) != string(barrierBody) {
|
||||
t.Fatalf("resumed task remained after barrier: %s", last.Body)
|
||||
}
|
||||
if err := last.Ack(false); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
|
||||
for i := 0; i < 100; i++ {
|
||||
body := []byte(fmt.Sprintf(`{"event_id":"stop-backlog-%d"}`, i))
|
||||
if err := mockChannel.PublishWithContext(context.Background(), CommandsExchangeV3, routingKey, true, false, amqp.Publishing{
|
||||
ContentType: "application/json", DeliveryMode: amqp.Persistent, Body: body,
|
||||
}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
drained, err := broker.DrainPredeclared(context.Background(), taskQueue)
|
||||
if err != nil || drained != 100 {
|
||||
t.Fatalf("stop drain count=%d err=%v, want 100", drained, err)
|
||||
}
|
||||
if _, ok, err := mockChannel.Get(taskQueue, true); err != nil || ok {
|
||||
t.Fatalf("stop drain left task messages: ok=%v err=%v", ok, err)
|
||||
}
|
||||
if _, ok, err := mockChannel.Get(resultQueue, true); err != nil || ok {
|
||||
t.Fatalf("stop drain emitted per-call results: ok=%v err=%v", ok, err)
|
||||
}
|
||||
|
||||
resultBody := []byte(`{"event_id":"result-event-1"}`)
|
||||
if err := broker.Publish(context.Background(), ResultsExchangeV3, resultRoute, resultBody); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
resultDeliveries, err := mockChannel.Consume(resultQueue, "", true, false, false, false, nil)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
select {
|
||||
case delivery := <-resultDeliveries:
|
||||
if string(delivery.Body) != string(resultBody) {
|
||||
t.Fatalf("published result = %s, want %s", delivery.Body, resultBody)
|
||||
}
|
||||
if delivery.DeliveryMode != amqp.Persistent {
|
||||
t.Fatalf("result delivery mode = %d, want persistent", delivery.DeliveryMode)
|
||||
}
|
||||
case <-time.After(5 * time.Second):
|
||||
t.Fatal("timed out waiting for confirmed result")
|
||||
}
|
||||
}
|
||||
@@ -1,129 +0,0 @@
|
||||
//go:build integration
|
||||
|
||||
package mq
|
||||
|
||||
import (
|
||||
"context"
|
||||
"os"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
amqp "github.com/rabbitmq/amqp091-go"
|
||||
)
|
||||
|
||||
// A full task queue is SaaS-owned: the publisher must observe the rejection
|
||||
// and retain its command, rather than expecting Dispatcher to create a queue
|
||||
// or invent the missing message after restart.
|
||||
func TestV3FullSaaSOwnedTaskQueueRejectsSecondCommand(t *testing.T) {
|
||||
url := os.Getenv("RABBITMQ_URL")
|
||||
if url == "" {
|
||||
t.Skip("RABBITMQ_URL not set")
|
||||
}
|
||||
provisionerURL := os.Getenv("RABBITMQ_PROVISIONER_URL")
|
||||
if provisionerURL == "" {
|
||||
provisionerURL = url
|
||||
}
|
||||
const dispatcherID = "550e8400-e29b-41d4-a716-446655440099"
|
||||
const taskQueue = "agent-call.d." + dispatcherID + ".task.queue-full-1.v3"
|
||||
const routingKey = "d." + dispatcherID + ".task.queue-full-1.in"
|
||||
provisioner, err := amqp.Dial(provisionerURL)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer provisioner.Close()
|
||||
channel, err := provisioner.Channel()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer channel.Close()
|
||||
for _, exchange := range []string{CommandsExchangeV3, ResultsExchangeV3, DeadLetterExchangeV3} {
|
||||
if err := channel.ExchangeDeclare(exchange, "topic", true, false, false, false, nil); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
if _, err := channel.QueueDeclare(taskQueue, true, false, false, false, amqp.Table{
|
||||
"x-max-length": int32(1), "x-overflow": "reject-publish",
|
||||
}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer channel.QueueDelete(taskQueue, false, false, false)
|
||||
if err := channel.QueueBind(taskQueue, routingKey, CommandsExchangeV3, false, nil); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
controlQueue := V3ControlQueueName(dispatcherID)
|
||||
resultQueue := "agent-call.saas.d." + dispatcherID + ".v3"
|
||||
for _, queue := range []string{controlQueue, resultQueue} {
|
||||
if _, err := channel.QueueDeclare(queue, true, false, false, false, nil); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer channel.QueueDelete(queue, false, false, false)
|
||||
}
|
||||
if err := channel.QueueBind(controlQueue, "d."+dispatcherID+".control.in", CommandsExchangeV3, false, nil); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := channel.QueueBind(resultQueue, "d."+dispatcherID+".out", ResultsExchangeV3, false, nil); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := channel.Confirm(false); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
confirms := channel.NotifyPublish(make(chan amqp.Confirmation, 2))
|
||||
returned := channel.NotifyReturn(make(chan amqp.Return, 2))
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second)
|
||||
defer cancel()
|
||||
publish := func(body string) bool {
|
||||
t.Helper()
|
||||
if err := channel.PublishWithContext(ctx, CommandsExchangeV3, routingKey, true, false, amqp.Publishing{
|
||||
ContentType: "application/json", DeliveryMode: amqp.Persistent, Body: []byte(body),
|
||||
}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
select {
|
||||
case confirmation := <-confirms:
|
||||
if !confirmation.Ack {
|
||||
return false
|
||||
}
|
||||
select {
|
||||
case <-returned:
|
||||
return false
|
||||
default:
|
||||
return true
|
||||
}
|
||||
case <-ctx.Done():
|
||||
t.Fatal("SaaS publisher received no RabbitMQ confirmation")
|
||||
return false
|
||||
}
|
||||
}
|
||||
const first = `{"event_id":"original-task-command"}`
|
||||
if !publish(first) || publish(`{"event_id":"rejected-task-command"}`) {
|
||||
t.Fatal("full SaaS-owned queue did not confirm the first command and reject the second")
|
||||
}
|
||||
queued, err := channel.QueueInspect(taskQueue)
|
||||
if err != nil || queued.Messages != 1 {
|
||||
t.Fatalf("queue-full publication invented or discarded the original command: count=%d err=%v", queued.Messages, err)
|
||||
}
|
||||
broker, err := OpenV3(url, dispatcherID, 1)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer broker.Close()
|
||||
got := make(chan string, 1)
|
||||
consumer, err := broker.StartPredeclaredConsumer(ctx, taskQueue, func(_ context.Context, _ string, body []byte) error {
|
||||
got <- string(body)
|
||||
return nil
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
select {
|
||||
case body := <-got:
|
||||
if body != first {
|
||||
t.Fatalf("Dispatcher consumed an unpublished command: %s", body)
|
||||
}
|
||||
case <-ctx.Done():
|
||||
t.Fatal("Dispatcher did not consume the SaaS-retained original command")
|
||||
}
|
||||
if err := consumer.Stop(context.Background()); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user