From 4228574a4134c08eb62e7d15bbaf9a04de391d00 Mon Sep 17 00:00:00 2001 From: Rogee Date: Wed, 30 Sep 2026 11:11:28 +0800 Subject: [PATCH] Remove obsolete V3 MQ broker after SQLite retirement --- .../saas-dispatcher-implementation.md | 3 +- internal/mq/amqp_v3.go | 416 ---------------- internal/mq/amqp_v3_integration_test.go | 449 ------------------ .../mq/amqp_v3_queue_full_integration_test.go | 129 ----- 4 files changed, 2 insertions(+), 995 deletions(-) delete mode 100644 internal/mq/amqp_v3.go delete mode 100644 internal/mq/amqp_v3_integration_test.go delete mode 100644 internal/mq/amqp_v3_queue_full_integration_test.go diff --git a/docs/evidence/saas-dispatcher-implementation.md b/docs/evidence/saas-dispatcher-implementation.md index b8e2322..fe972f2 100644 --- a/docs/evidence/saas-dispatcher-implementation.md +++ b/docs/evidence/saas-dispatcher-implementation.md @@ -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 接收与应用收讫仍未验证。 ## 验收台账 diff --git a/internal/mq/amqp_v3.go b/internal/mq/amqp_v3.go deleted file mode 100644 index bc92bd1..0000000 --- a/internal/mq/amqp_v3.go +++ /dev/null @@ -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) - } - } - } -} diff --git a/internal/mq/amqp_v3_integration_test.go b/internal/mq/amqp_v3_integration_test.go deleted file mode 100644 index 13b1951..0000000 --- a/internal/mq/amqp_v3_integration_test.go +++ /dev/null @@ -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") - } -} diff --git a/internal/mq/amqp_v3_queue_full_integration_test.go b/internal/mq/amqp_v3_queue_full_integration_test.go deleted file mode 100644 index e7a5ede..0000000 --- a/internal/mq/amqp_v3_queue_full_integration_test.go +++ /dev/null @@ -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) - } -}