diff --git a/cmd/sip-go-agent/current_dispatcher_command.go b/cmd/sip-go-agent/current_dispatcher_command.go index ee33281..80da402 100644 --- a/cmd/sip-go-agent/current_dispatcher_command.go +++ b/cmd/sip-go-agent/current_dispatcher_command.go @@ -90,7 +90,7 @@ func runCurrentDispatcher(ctx context.Context, mode string) (result error) { if err := database.CloseAdmission(settings.DispatcherID); err != nil { return fmt.Errorf("close task admission before connecting external adapters: %w", err) } - broker, err := mq.OpenCurrent(settings.RabbitMQURL, settings.DispatcherID, 1) + broker, err := mq.Open(settings.RabbitMQURL, settings.DispatcherID, 1) if err != nil { return err } diff --git a/cmd/sip-go-agent/current_dispatcher_integration_test.go b/cmd/sip-go-agent/current_dispatcher_integration_test.go index 41fd1a5..6faa461 100644 --- a/cmd/sip-go-agent/current_dispatcher_integration_test.go +++ b/cmd/sip-go-agent/current_dispatcher_integration_test.go @@ -48,16 +48,16 @@ func TestCurrentDispatcherCommandStartsWithIsolatedMQHTTPAndAgent(t *testing.T) t.Fatal(err) } defer admin.Close() - for _, exchange := range []string{mq.CommandsExchangeCurrent, mq.ResultsExchangeCurrent, mq.DeadLetterExchangeCurrent} { + for _, exchange := range []string{mq.CommandsExchange, mq.ResultsExchange, mq.DeadLetterExchange} { if err := admin.ExchangeDeclare(exchange, "topic", true, false, false, false, nil); err != nil { t.Fatal(err) } } - control, err := tenant.CurrentControlRoute(dispatcherID) + control, err := tenant.ControlRoute(dispatcherID) if err != nil { t.Fatal(err) } - result, err := tenant.CurrentResultRoute(dispatcherID) + result, err := tenant.ResultRoute(dispatcherID) if err != nil { t.Fatal(err) } @@ -76,7 +76,7 @@ func TestCurrentDispatcherCommandStartsWithIsolatedMQHTTPAndAgent(t *testing.T) t.Fatal(err) } defer func() { _ = admin.QueueUnbind(result.Queue, result.BindingKey, result.Exchange, nil) }() - taskRoute, err := tenant.CurrentTaskRoute(dispatcherID, "task-asr") + taskRoute, err := tenant.TaskRoute(dispatcherID, "task-asr") if err != nil { t.Fatal(err) } diff --git a/docs/evidence/saas-dispatcher-implementation.md b/docs/evidence/saas-dispatcher-implementation.md index 349e410..e95c309 100644 --- a/docs/evidence/saas-dispatcher-implementation.md +++ b/docs/evidence/saas-dispatcher-implementation.md @@ -106,6 +106,7 @@ - MQ 旧任务队列分支:旧 Store 的任务队列引用删除后,移除 `V3Broker` 和三份只验证旧拓扑的代码/测试;共享结果队列和无配置权限下的投递仍由当前 MQ 隔离测试验证。真实 SaaS/MQ 接收与应用收讫仍未验证。 - MQ 旧声明拓扑分支:移除旧 `Broker`、旧租户 Topic 路由/队列声明及其测试;独立保留当前消费端的消息处理器、永久错误分类与 Dispatcher UUID v4 校验,补错误解包和规范身份的回归测试。当前 Broker 仅被动核对 SaaS 预建拓扑;隔离 MQ 测试通过不代表真实 SaaS 应用收讫。 - Store 名称收敛:移除现行 Store 的自有 `Current*` 类型、错误和 `OpenCurrent` 入口,保留单一 `Store`/`Open`;20 份 Go 源码与测试文件改为不带代次的路径,原现行任务/录音/结果测试仍执行。只重命名源码与调用,不修改现有 SQLite 表、记录或恢复数据。 +- MQ/租户路由名称收敛:现行 MQ 和租户路由的自有 `Current*` 代码标识改为唯一的 `Broker`/`Open`、`Route` 及控制/任务/结果路由入口;七份 Go 文件改为无代次路径。固定 Topic/队列的 `.v1` 名称和值保持原样,隔离 MQ 的精确路由和不声明拓扑检查继续执行。 ## 验收台账 diff --git a/internal/dispatcher/current_runtime.go b/internal/dispatcher/current_runtime.go index c541e1f..054d4d9 100644 --- a/internal/dispatcher/current_runtime.go +++ b/internal/dispatcher/current_runtime.go @@ -18,7 +18,7 @@ import ( // CurrentRuntime owns task/control consumers for one Dispatcher. SaaS creates // every queue/binding; this process only checks and consumes predeclared ones. type CurrentRuntime struct { - Broker *mq.CurrentBroker + Broker *mq.Broker Bootstrap CurrentBootstrap Execute CurrentExecuteController Control CurrentControlController @@ -46,8 +46,8 @@ func (r *CurrentRuntime) Serve(ctx context.Context) (result error) { } r.failures = make(chan error, 1) r.taskLocks = make(map[string]*sync.Mutex) - consumers := make(map[string]*mq.CurrentConsumer) - var controlConsumer *mq.CurrentConsumer + consumers := make(map[string]*mq.Consumer) + var controlConsumer *mq.Consumer defer func() { stopCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second) defer cancel() @@ -140,7 +140,7 @@ func (r *CurrentRuntime) Serve(ctx context.Context) (result error) { } } -func (r *CurrentRuntime) watchConsumer(ctx context.Context, queue string, consumer *mq.CurrentConsumer) { +func (r *CurrentRuntime) watchConsumer(ctx context.Context, queue string, consumer *mq.Consumer) { err := consumer.Wait(ctx) if err != nil && ctx.Err() == nil { r.Logger.Error("current MQ consumer failed", "dispatcher_id", r.Bootstrap.DispatcherID, "queue", queue, "error", err) @@ -243,14 +243,14 @@ func (r *CurrentRuntime) processPending(ctx context.Context) error { return nil } -func (r *CurrentRuntime) syncTaskConsumers(ctx context.Context, consumers map[string]*mq.CurrentConsumer) error { +func (r *CurrentRuntime) syncTaskConsumers(ctx context.Context, consumers map[string]*mq.Consumer) error { assigned, err := r.Bootstrap.Store.ListAssignedTasks(r.Bootstrap.DispatcherID) if err != nil { return err } wanted := make(map[string]store.AssignedTask, len(assigned)) for _, task := range assigned { - route, err := tenant.CurrentTaskRoute(r.Bootstrap.DispatcherID, task.TaskID) + route, err := tenant.TaskRoute(r.Bootstrap.DispatcherID, task.TaskID) if err != nil { return err } @@ -293,7 +293,7 @@ func (r *CurrentRuntime) syncTaskConsumers(ctx context.Context, consumers map[st } id, tenantID := task.TaskID, task.TenantID ready := make(chan struct{}) - var consumer *mq.CurrentConsumer + var consumer *mq.Consumer consumer, err = r.Broker.StartPredeclaredConsumer(ctx, queue, func(ctx context.Context, _ string, body []byte) error { <-ready // the subscription is assigned before its first delivery can pause itself if err := r.handleTask(ctx, id, body); err != nil { @@ -320,4 +320,4 @@ func (r *CurrentRuntime) syncTaskConsumers(ctx context.Context, consumers map[st // Compile-time interface checks: the same broker confirms bound persistent // results and consumes SaaS-owned queues without configure permissions. -var _ CurrentPublisher = (*mq.CurrentBroker)(nil) +var _ CurrentPublisher = (*mq.Broker)(nil) diff --git a/internal/dispatcher/current_runtime_integration_test.go b/internal/dispatcher/current_runtime_integration_test.go index 1f2035c..5ed7c94 100644 --- a/internal/dispatcher/current_runtime_integration_test.go +++ b/internal/dispatcher/current_runtime_integration_test.go @@ -57,14 +57,14 @@ func TestCurrentRuntimeIsolatedControlBacklogExecuteAndSharedResult(t *testing.T t.Fatal(err) } defer admin.Close() - for _, exchange := range []string{mq.CommandsExchangeCurrent, mq.ResultsExchangeCurrent, mq.DeadLetterExchangeCurrent} { + for _, exchange := range []string{mq.CommandsExchange, mq.ResultsExchange, mq.DeadLetterExchange} { if err := admin.ExchangeDeclare(exchange, "topic", true, false, false, false, nil); err != nil { t.Fatal(err) } } - controlRoute, _ := tenant.CurrentControlRoute(id) - taskRoute, _ := tenant.CurrentTaskRoute(id, "task-asr") - resultRoute, _ := tenant.CurrentResultRoute(id) + controlRoute, _ := tenant.ControlRoute(id) + taskRoute, _ := tenant.TaskRoute(id, "task-asr") + resultRoute, _ := tenant.ResultRoute(id) shared := "agent-call.saas.events.v1" for _, queue := range []string{controlRoute.Queue, taskRoute.Queue, shared} { if _, err := admin.QueueDeclare(queue, true, false, false, false, nil); err != nil { @@ -75,7 +75,7 @@ func TestCurrentRuntimeIsolatedControlBacklogExecuteAndSharedResult(t *testing.T t.Fatal(err) } } - for _, route := range []tenant.CurrentRoute{controlRoute, taskRoute, resultRoute} { + for _, route := range []tenant.Route{controlRoute, taskRoute, resultRoute} { queue := route.Queue if route.BindingKey == resultRoute.BindingKey { queue = shared @@ -85,7 +85,7 @@ func TestCurrentRuntimeIsolatedControlBacklogExecuteAndSharedResult(t *testing.T } defer func(q, key, exchange string) { _ = admin.QueueUnbind(q, key, exchange, nil) }(queue, route.BindingKey, route.Exchange) } - publish := func(route tenant.CurrentRoute, body []byte) { + publish := func(route tenant.Route, body []byte) { t.Helper() if err := admin.PublishWithContext(context.Background(), route.Exchange, route.BindingKey, true, false, amqp.Publishing{ContentType: "application/json", DeliveryMode: amqp.Persistent, Body: body}); err != nil { t.Fatal(err) @@ -137,7 +137,7 @@ func TestCurrentRuntimeIsolatedControlBacklogExecuteAndSharedResult(t *testing.T t.Fatal(err) } defer db.Close() - broker, err := mq.OpenCurrent(brokerURL, id, 1) + broker, err := mq.Open(brokerURL, id, 1) if err != nil { t.Fatal(err) } diff --git a/internal/mq/current.go b/internal/mq/broker.go similarity index 84% rename from internal/mq/current.go rename to internal/mq/broker.go index c48e84c..609dcdc 100644 --- a/internal/mq/current.go +++ b/internal/mq/broker.go @@ -16,13 +16,13 @@ import ( ) const ( - CommandsExchangeCurrent = "agent-call.dispatchers.v1" - ResultsExchangeCurrent = "agent-call.saas.v1" - DeadLetterExchangeCurrent = "agent-call.dead-letter.v1" - MaxCurrentMessageBytes = 8 << 20 + CommandsExchange = "agent-call.dispatchers.v1" + ResultsExchange = "agent-call.saas.v1" + DeadLetterExchange = "agent-call.dead-letter.v1" + MaxMessageBytes = 8 << 20 ) -type CurrentBroker struct { +type Broker struct { conn *amqp.Connection dispatcherID string prefetch int @@ -32,13 +32,14 @@ type CurrentBroker struct { mu sync.Mutex } -// OpenCurrent verifies SaaS-provisioned topology passively. It never declares, -// binds, or deletes exchanges or queues. -func CurrentControlQueueName(dispatcherID string) string { +// ControlQueueName returns this Dispatcher's SaaS-provisioned control queue. +func ControlQueueName(dispatcherID string) string { return "agent-call.d." + dispatcherID + ".control.v1" } -func OpenCurrent(url, dispatcherID string, prefetch int) (*CurrentBroker, error) { +// Open verifies SaaS-provisioned topology passively. It never declares, +// binds, or deletes exchanges or queues. +func Open(url, dispatcherID string, prefetch int) (*Broker, error) { if strings.TrimSpace(url) == "" { return nil, errors.New("rabbitmq URL is required") } @@ -57,14 +58,14 @@ func OpenCurrent(url, dispatcherID string, prefetch int) (*CurrentBroker, error) _ = conn.Close() return nil, fmt.Errorf("open rabbitmq current channel: %w", err) } - for _, exchange := range []string{CommandsExchangeCurrent, ResultsExchangeCurrent, DeadLetterExchangeCurrent} { + for _, exchange := range []string{CommandsExchange, ResultsExchange, DeadLetterExchange} { 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 := CurrentControlQueueName(dispatcherID) + controlQueue := ControlQueueName(dispatcherID) if _, err := channel.QueueDeclarePassive(controlQueue, true, false, false, false, nil); err != nil { _ = channel.Close() _ = conn.Close() @@ -73,7 +74,7 @@ func OpenCurrent(url, dispatcherID string, prefetch int) (*CurrentBroker, error) // The shared SaaS result queue is not readable by Dispatcher credentials. // mandatory + return + confirm and a SaaS-owned Mock consumer prove routing. _ = channel.Close() - return &CurrentBroker{ + return &Broker{ conn: conn, dispatcherID: dispatcherID, prefetch: prefetch, @@ -83,7 +84,7 @@ func OpenCurrent(url, dispatcherID string, prefetch int) (*CurrentBroker, error) }, nil } -func (b *CurrentBroker) Close() error { +func (b *Broker) Close() error { b.mu.Lock() defer b.mu.Unlock() if b.conn == nil { @@ -94,11 +95,11 @@ func (b *CurrentBroker) Close() error { return err } -func (b *CurrentBroker) Done() <-chan *amqp.Error { return b.closed } +func (b *Broker) Done() <-chan *amqp.Error { return b.closed } -func (b *CurrentBroker) Publish(ctx context.Context, exchange, routingKey string, body []byte) error { - if exchange != ResultsExchangeCurrent || routingKey != b.resultRoute || len(body) == 0 || len(body) > MaxCurrentMessageBytes { - return fmt.Errorf("current publish requires result exchange, Dispatcher route and 1..%d body bytes", MaxCurrentMessageBytes) +func (b *Broker) Publish(ctx context.Context, exchange, routingKey string, body []byte) error { + if exchange != ResultsExchange || routingKey != b.resultRoute || len(body) == 0 || len(body) > MaxMessageBytes { + return fmt.Errorf("publish requires result exchange, Dispatcher route and 1..%d body bytes", MaxMessageBytes) } if err := contract.ValidateCurrent("mq", body); err != nil { return fmt.Errorf("outbound MQ contract: %w", err) @@ -171,14 +172,14 @@ func (b *CurrentBroker) Publish(ctx context.Context, exchange, routingKey string return nil } -type CurrentConsumer struct { +type Consumer struct { cancel context.CancelFunc done chan struct{} mu sync.Mutex err error } -func (c *CurrentConsumer) Wait(ctx context.Context) error { +func (c *Consumer) Wait(ctx context.Context) error { select { case <-c.done: c.mu.Lock() @@ -191,11 +192,11 @@ func (c *CurrentConsumer) Wait(ctx context.Context) error { // RequestStop cancels consumption without waiting for the current handler. // Calling Stop from inside that handler would deadlock on its own completion. -func (c *CurrentConsumer) RequestStop() { +func (c *Consumer) RequestStop() { c.cancel() } -func (c *CurrentConsumer) Stop(ctx context.Context) error { +func (c *Consumer) Stop(ctx context.Context) error { c.cancel() err := c.Wait(ctx) if errors.Is(err, context.Canceled) { @@ -204,7 +205,7 @@ func (c *CurrentConsumer) Stop(ctx context.Context) error { return err } -func (b *CurrentBroker) ConsumePredeclared(ctx context.Context, queue string, handler MessageHandler) error { +func (b *Broker) ConsumePredeclared(ctx context.Context, queue string, handler MessageHandler) error { consumer, err := b.StartPredeclaredConsumer(ctx, queue, handler) if err != nil { return err @@ -218,7 +219,7 @@ func (b *CurrentBroker) ConsumePredeclared(ctx context.Context, queue string, ha // StartPredeclaredConsumer subscribes only to an existing SaaS-owned queue. // Stopping it closes the channel so every unacknowledged delivery is requeued. -func (b *CurrentBroker) StartPredeclaredConsumer(ctx context.Context, queue string, handler MessageHandler) (*CurrentConsumer, error) { +func (b *Broker) StartPredeclaredConsumer(ctx context.Context, queue string, handler MessageHandler) (*Consumer, error) { if handler == nil { return nil, errors.New("predeclared handler is required") } @@ -247,16 +248,16 @@ func (b *CurrentBroker) StartPredeclaredConsumer(ctx context.Context, queue stri _ = channel.Close() return nil, fmt.Errorf("set current consumer prefetch: %w", err) } - consumerTag := fmt.Sprintf("sip-go-agent-current-%d", consumerSequence.Add(1)) + consumerTag := fmt.Sprintf("sip-go-agent-%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 current queue: %w", err) } consumerCtx, cancel := context.WithCancel(ctx) - consumer := &CurrentConsumer{cancel: cancel, done: make(chan struct{})} + consumer := &Consumer{cancel: cancel, done: make(chan struct{})} go func() { - consumeErr := b.consumeCurrentDeliveries(consumerCtx, dispatcherID, queue, deliveries, handler) + consumeErr := b.consumeDeliveries(consumerCtx, dispatcherID, queue, deliveries, handler) if errors.Is(consumeErr, context.Canceled) && consumerCtx.Err() != nil { consumeErr = nil } @@ -274,7 +275,7 @@ func (b *CurrentBroker) StartPredeclaredConsumer(ctx context.Context, queue stri return consumer, nil } -func (b *CurrentBroker) DrainPredeclared(ctx context.Context, queue string) (int, error) { +func (b *Broker) DrainPredeclared(ctx context.Context, queue string) (int, error) { if err := b.validateQueue(queue, false); err != nil { return 0, err } @@ -318,7 +319,7 @@ func (b *CurrentBroker) DrainPredeclared(ctx context.Context, queue string) (int // 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 *CurrentBroker) DrainControlPredeclared(ctx context.Context, queue string, handler MessageHandler) (int, error) { +func (b *Broker) 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") } @@ -352,7 +353,7 @@ func (b *CurrentBroker) DrainControlPredeclared(ctx context.Context, queue strin if !ok { return processed, nil } - if len(delivery.Body) == 0 || len(delivery.Body) > MaxCurrentMessageBytes || !json.Valid(delivery.Body) { + if len(delivery.Body) == 0 || len(delivery.Body) > MaxMessageBytes || !json.Valid(delivery.Body) { if err := delivery.Reject(false); err != nil { return processed, fmt.Errorf("reject invalid control backlog message: %w", err) } @@ -388,12 +389,12 @@ func (b *CurrentBroker) DrainControlPredeclared(ctx context.Context, queue strin } } -func (b *CurrentBroker) ControlQueue() string { return b.controlQueue } +func (b *Broker) ControlQueue() string { return b.controlQueue } // validateInbound prevents a wrongly bound queue, foreign D message, old // envelope, or outbound acknowledgment from entering the command handler. -func (b *CurrentBroker) validateInbound(queue, routingKey string, body []byte) error { - if len(body) == 0 || len(body) > MaxCurrentMessageBytes { +func (b *Broker) validateInbound(queue, routingKey string, body []byte) error { + if len(body) == 0 || len(body) > MaxMessageBytes { return errors.New("invalid MQ message size") } if err := b.validateQueue(queue, true); err != nil { @@ -428,14 +429,14 @@ func (b *CurrentBroker) validateInbound(queue, routingKey string, body []byte) e } prefix := "agent-call.d." + b.dispatcherID + ".task." taskID := strings.TrimSuffix(strings.TrimPrefix(queue, prefix), ".v1") - route, err := tenant.CurrentTaskRoute(b.dispatcherID, taskID) + route, err := tenant.TaskRoute(b.dispatcherID, taskID) if err != nil || routingKey != route.BindingKey || identity.EventType != "call.execute" || identity.Payload.TaskID != taskID || identity.Payload.Callee == "" { return errors.New("unapproved task-queue event or routing key") } return nil } -func (b *CurrentBroker) validateQueue(queue string, allowControl bool) error { +func (b *Broker) validateQueue(queue string, allowControl bool) error { if allowControl && queue == b.controlQueue { return nil } @@ -444,14 +445,14 @@ func (b *CurrentBroker) validateQueue(queue string, allowControl bool) error { return errors.New("queue is not owned by this Dispatcher task") } taskID := strings.TrimSuffix(strings.TrimPrefix(queue, prefix), ".v1") - route, err := tenant.CurrentTaskRoute(b.dispatcherID, taskID) + route, err := tenant.TaskRoute(b.dispatcherID, taskID) if err != nil || route.Queue != queue { return errors.New("invalid assigned task queue") } return nil } -func (b *CurrentBroker) consumeCurrentDeliveries(ctx context.Context, dispatcherID, queue string, deliveries <-chan amqp.Delivery, handler MessageHandler) error { +func (b *Broker) consumeDeliveries(ctx context.Context, dispatcherID, queue string, deliveries <-chan amqp.Delivery, handler MessageHandler) error { for { select { case <-ctx.Done(): @@ -460,7 +461,7 @@ func (b *CurrentBroker) consumeCurrentDeliveries(ctx context.Context, dispatcher if !ok { return errors.New("rabbitmq current delivery channel closed") } - if len(delivery.Body) == 0 || len(delivery.Body) > MaxCurrentMessageBytes { + if len(delivery.Body) == 0 || len(delivery.Body) > MaxMessageBytes { if err := delivery.Reject(false); err != nil { return fmt.Errorf("reject invalid current message size: %w", err) } diff --git a/internal/mq/current_integration_test.go b/internal/mq/broker_integration_test.go similarity index 81% rename from internal/mq/current_integration_test.go rename to internal/mq/broker_integration_test.go index 0c7be66..84e0226 100644 --- a/internal/mq/current_integration_test.go +++ b/internal/mq/broker_integration_test.go @@ -14,13 +14,13 @@ import ( amqp "github.com/rabbitmq/amqp091-go" ) -func TestCurrentBrokerSharedResultQueueAndNoConfigure(t *testing.T) { +func TestBrokerSharedResultQueueAndNoConfigure(t *testing.T) { url := os.Getenv("RABBITMQ_URL") provisionerURL := os.Getenv("RABBITMQ_PROVISIONER_URL") if url == "" || provisionerURL == "" { t.Skip("isolated RabbitMQ Mock URLs not configured") } - ids := []string{currentTestDispatcherID, "550e8400-e29b-41d4-a716-446655440000"} + ids := []string{testDispatcherID, "550e8400-e29b-41d4-a716-446655440000"} conn, err := amqp.Dial(provisionerURL) if err != nil { t.Fatal(err) @@ -31,7 +31,7 @@ func TestCurrentBrokerSharedResultQueueAndNoConfigure(t *testing.T) { t.Fatal(err) } defer admin.Close() - for _, exchange := range []string{CommandsExchangeCurrent, ResultsExchangeCurrent, DeadLetterExchangeCurrent} { + for _, exchange := range []string{CommandsExchange, ResultsExchange, DeadLetterExchange} { if err := admin.ExchangeDeclare(exchange, "topic", true, false, false, false, nil); err != nil { t.Fatal(err) } @@ -43,10 +43,10 @@ func TestCurrentBrokerSharedResultQueueAndNoConfigure(t *testing.T) { defer func() { _, _ = admin.QueueDelete(shared, false, false, false) }() queues := []string{} for _, id := range ids { - control, _ := tenant.CurrentControlRoute(id) - task, _ := tenant.CurrentTaskRoute(id, "task-asr") - result, _ := tenant.CurrentResultRoute(id) - for _, route := range []tenant.CurrentRoute{control, task} { + control, _ := tenant.ControlRoute(id) + task, _ := tenant.TaskRoute(id, "task-asr") + result, _ := tenant.ResultRoute(id) + for _, route := range []tenant.Route{control, task} { if _, err := admin.QueueDeclare(route.Queue, true, false, false, false, nil); err != nil { t.Fatal(err) } @@ -58,7 +58,7 @@ func TestCurrentBrokerSharedResultQueueAndNoConfigure(t *testing.T) { if err := admin.QueueBind(shared, result.BindingKey, result.Exchange, false, nil); err != nil { t.Fatal(err) } - defer func(key string) { _ = admin.QueueUnbind(shared, key, ResultsExchangeCurrent, nil) }(result.BindingKey) + defer func(key string) { _ = admin.QueueUnbind(shared, key, ResultsExchange, nil) }(result.BindingKey) } defer func() { for _, name := range queues { @@ -88,19 +88,19 @@ func TestCurrentBrokerSharedResultQueueAndNoConfigure(t *testing.T) { _ = probe.Close() _ = probeConn.Close() - brokers := make([]*CurrentBroker, 0, 2) + brokers := make([]*Broker, 0, 2) for _, id := range ids { - broker, err := OpenCurrent(url, id, 1) + broker, err := Open(url, id, 1) if err != nil { t.Fatalf("Dispatcher without configure could not open %s: %v", id, err) } brokers = append(brokers, broker) defer broker.Close() } - control, _ := tenant.CurrentControlRoute(ids[0]) - task, _ := tenant.CurrentTaskRoute(ids[0], "task-asr") + control, _ := tenant.ControlRoute(ids[0]) + task, _ := tenant.TaskRoute(ids[0], "task-asr") controlBody := currentMessage(t, "mq-control") - if err := admin.PublishWithContext(context.Background(), CommandsExchangeCurrent, control.BindingKey, true, false, amqp.Publishing{ContentType: "application/json", DeliveryMode: amqp.Persistent, Body: controlBody}); err != nil { + if err := admin.PublishWithContext(context.Background(), CommandsExchange, control.BindingKey, true, false, amqp.Publishing{ContentType: "application/json", DeliveryMode: amqp.Persistent, Body: controlBody}); err != nil { t.Fatal(err) } processed, err := brokers[0].DrainControlPredeclared(context.Background(), control.Queue, func(_ context.Context, key string, body []byte) error { @@ -113,7 +113,7 @@ func TestCurrentBrokerSharedResultQueueAndNoConfigure(t *testing.T) { t.Fatalf("control backlog: processed=%d err=%v", processed, err) } body := currentMessage(t, "mq-execute") - if err := admin.PublishWithContext(context.Background(), CommandsExchangeCurrent, task.BindingKey, true, false, amqp.Publishing{ContentType: "application/json", DeliveryMode: amqp.Persistent, Body: body}); err != nil { + if err := admin.PublishWithContext(context.Background(), CommandsExchange, task.BindingKey, true, false, amqp.Publishing{ContentType: "application/json", DeliveryMode: amqp.Persistent, Body: body}); err != nil { t.Fatal(err) } ctx, cancel := context.WithCancel(context.Background()) @@ -144,7 +144,7 @@ func TestCurrentBrokerSharedResultQueueAndNoConfigure(t *testing.T) { case <-time.After(5 * time.Second): t.Fatal("consumer did not stop") } - otherTask, _ := tenant.CurrentTaskRoute(ids[1], "task-asr") + otherTask, _ := tenant.TaskRoute(ids[1], "task-asr") if state, err := admin.QueueInspect(otherTask.Queue); err != nil || state.Messages != 0 { t.Fatalf("D2 stole D1 command: %+v %v", state, err) } @@ -155,7 +155,7 @@ func TestCurrentBrokerSharedResultQueueAndNoConfigure(t *testing.T) { if i == 1 { message = []byte(strings.Replace(string(resultBody), ids[0], ids[1], 1)) } - route, _ := tenant.CurrentResultRoute(ids[i]) + route, _ := tenant.ResultRoute(ids[i]) if err := broker.Publish(context.Background(), route.Exchange, route.BindingKey, message); err != nil { t.Fatalf("result publish for D%d: %v", i+1, err) } @@ -177,16 +177,16 @@ func TestCurrentBrokerSharedResultQueueAndNoConfigure(t *testing.T) { } } for _, id := range ids { - route, _ := tenant.CurrentResultRoute(id) + route, _ := tenant.ResultRoute(id) if err := admin.QueueUnbind(shared, route.BindingKey, route.Exchange, nil); err != nil { t.Fatal(err) } } - resultRoute, _ := tenant.CurrentResultRoute(ids[0]) + resultRoute, _ := tenant.ResultRoute(ids[0]) if err := brokers[0].Publish(context.Background(), resultRoute.Exchange, resultRoute.BindingKey, resultBody); err == nil { t.Fatal("unroutable mandatory result falsely confirmed") } - controlRoute, err := tenant.CurrentControlRoute(ids[0]) + controlRoute, err := tenant.ControlRoute(ids[0]) if err != nil { t.Fatal(err) } diff --git a/internal/mq/current_test.go b/internal/mq/broker_test.go similarity index 71% rename from internal/mq/current_test.go rename to internal/mq/broker_test.go index fe750f0..85a98db 100644 --- a/internal/mq/current_test.go +++ b/internal/mq/broker_test.go @@ -8,7 +8,7 @@ import ( "testing" ) -const currentTestDispatcherID = "c046b893-8628-4589-ae50-619d049248a6" +const testDispatcherID = "c046b893-8628-4589-ae50-619d049248a6" func currentMessage(t *testing.T, name string) []byte { t.Helper() @@ -19,29 +19,29 @@ func currentMessage(t *testing.T, name string) []byte { return body } -func TestCurrentBrokerRejectsInvalidConfigBeforeConnecting(t *testing.T) { +func TestBrokerRejectsInvalidConfigBeforeConnecting(t *testing.T) { for _, tc := range []struct { url, id string prefetch int }{ - {"", currentTestDispatcherID, 1}, + {"", testDispatcherID, 1}, {"amqp://127.0.0.1:1", "not-a-uuid", 1}, - {"amqp://127.0.0.1:1", currentTestDispatcherID, 0}, + {"amqp://127.0.0.1:1", testDispatcherID, 0}, } { - if _, err := OpenCurrent(tc.url, tc.id, tc.prefetch); err == nil { + if _, err := Open(tc.url, tc.id, tc.prefetch); err == nil { t.Fatalf("accepted invalid broker settings: %+v", tc) } } } -func TestCurrentPublishRejectsLegacyAndWrongOwnerBeforeNetwork(t *testing.T) { - b := &CurrentBroker{dispatcherID: currentTestDispatcherID, resultRoute: "d." + currentTestDispatcherID + ".out"} +func TestPublishRejectsLegacyAndWrongOwnerBeforeNetwork(t *testing.T) { + b := &Broker{dispatcherID: testDispatcherID, resultRoute: "d." + testDispatcherID + ".out"} for _, tc := range []struct { name string body []byte }{ {"old version field", currentMessage(t, "invalid/mq-legacy-schema-version")}, - {"wrong owner", []byte(strings.Replace(string(currentMessage(t, "mq-result-no-recording")), currentTestDispatcherID, "c046b893-8628-4589-ae50-619d049248a7", 1))}, + {"wrong owner", []byte(strings.Replace(string(currentMessage(t, "mq-result-no-recording")), testDispatcherID, "c046b893-8628-4589-ae50-619d049248a7", 1))}, {"unapproved event", currentMessage(t, "mq-sip-change")}, {"empty event identity", []byte(strings.Replace(string(currentMessage(t, "mq-execute-ack")), `"event_id":"call-example"`, `"event_id":""`, 1))}, } { diff --git a/internal/mq/current_inbound_test.go b/internal/mq/inbound_test.go similarity index 72% rename from internal/mq/current_inbound_test.go rename to internal/mq/inbound_test.go index aaee1d0..f394b37 100644 --- a/internal/mq/current_inbound_test.go +++ b/internal/mq/inbound_test.go @@ -7,16 +7,16 @@ import ( "git.ipao.vip/rogee/go-sip/internal/tenant" ) -func TestCurrentBrokerInboundRequiresExactQueueRouteOwnerAndSchema(t *testing.T) { - control, err := tenant.CurrentControlRoute(currentTestDispatcherID) +func TestBrokerInboundRequiresExactQueueRouteOwnerAndSchema(t *testing.T) { + control, err := tenant.ControlRoute(testDispatcherID) if err != nil { t.Fatal(err) } - task, err := tenant.CurrentTaskRoute(currentTestDispatcherID, "task-asr") + task, err := tenant.TaskRoute(testDispatcherID, "task-asr") if err != nil { t.Fatal(err) } - broker := &CurrentBroker{dispatcherID: currentTestDispatcherID, controlQueue: control.Queue} + broker := &Broker{dispatcherID: testDispatcherID, controlQueue: control.Queue} for _, tc := range []struct { name, queue, route, fixture string wantError bool @@ -37,18 +37,18 @@ func TestCurrentBrokerInboundRequiresExactQueueRouteOwnerAndSchema(t *testing.T) } }) } - foreign := []byte(strings.Replace(string(currentMessage(t, "mq-execute")), currentTestDispatcherID, "c046b893-8628-4589-ae50-619d049248a7", 1)) + foreign := []byte(strings.Replace(string(currentMessage(t, "mq-execute")), testDispatcherID, "c046b893-8628-4589-ae50-619d049248a7", 1)) if err := broker.validateInbound(task.Queue, task.BindingKey, foreign); err == nil { t.Fatal("accepted wrong Dispatcher owner") } } -func TestCurrentBrokerRejectsControlQueueFromSilentTaskDrain(t *testing.T) { - control, err := tenant.CurrentControlRoute(currentTestDispatcherID) +func TestBrokerRejectsControlQueueFromSilentTaskDrain(t *testing.T) { + control, err := tenant.ControlRoute(testDispatcherID) if err != nil { t.Fatal(err) } - broker := &CurrentBroker{dispatcherID: currentTestDispatcherID, controlQueue: control.Queue} + broker := &Broker{dispatcherID: testDispatcherID, controlQueue: control.Queue} if err := broker.validateQueue(control.Queue, false); err == nil { t.Fatal("allowed task-drain to silently ACK controls") } diff --git a/internal/mq/current_pause_test.go b/internal/mq/pause_test.go similarity index 67% rename from internal/mq/current_pause_test.go rename to internal/mq/pause_test.go index 007784b..75444b5 100644 --- a/internal/mq/current_pause_test.go +++ b/internal/mq/pause_test.go @@ -9,9 +9,9 @@ import ( amqp "github.com/rabbitmq/amqp091-go" ) -func TestCurrentConsumerRequestStopDoesNotWaitForItself(t *testing.T) { +func TestConsumerRequestStopDoesNotWaitForItself(t *testing.T) { ctx, cancel := context.WithCancel(context.Background()) - consumer := &CurrentConsumer{cancel: cancel, done: make(chan struct{})} + consumer := &Consumer{cancel: cancel, done: make(chan struct{})} consumer.RequestStop() select { case <-ctx.Done(): @@ -27,18 +27,18 @@ func (*currentRejectRecorder) Ack(uint64, bool) error { return nil } func (*currentRejectRecorder) Nack(uint64, bool, bool) error { return nil } func (a *currentRejectRecorder) Reject(uint64, bool) error { a.rejected++; return nil } -func TestCurrentMalformedControlStopsConsumerAfterReject(t *testing.T) { +func TestMalformedControlStopsConsumerAfterReject(t *testing.T) { id := "c046b893-8628-4589-ae50-619d049248a6" - route, err := tenant.CurrentControlRoute(id) + route, err := tenant.ControlRoute(id) if err != nil { t.Fatal(err) } - broker := &CurrentBroker{dispatcherID: id, controlQueue: route.Queue} + broker := &Broker{dispatcherID: id, controlQueue: route.Queue} ack := ¤tRejectRecorder{} deliveries := make(chan amqp.Delivery, 1) deliveries <- amqp.Delivery{Body: []byte("{"), RoutingKey: route.BindingKey, Acknowledger: ack} close(deliveries) - err = broker.consumeCurrentDeliveries(context.Background(), id, route.Queue, deliveries, func(context.Context, string, []byte) error { t.Fatal("invalid control reached handler"); return nil }) + err = broker.consumeDeliveries(context.Background(), id, route.Queue, deliveries, func(context.Context, string, []byte) error { t.Fatal("invalid control reached handler"); return nil }) if err == nil || !strings.Contains(err.Error(), "invalid control") || ack.rejected != 1 { t.Fatalf("malformed control was silently discarded before admission: rejected=%d err=%v", ack.rejected, err) } diff --git a/internal/store/calls.go b/internal/store/calls.go index cb356fb..1953f89 100644 --- a/internal/store/calls.go +++ b/internal/store/calls.go @@ -357,7 +357,7 @@ func insertExecuteAck(tx *sql.Tx, dispatcherID, eventID string, tenantID int64, if err := contract.ValidateCurrent("mq", body); err != nil { return fmt.Errorf("call acknowledgment violates current MQ contract: %w", err) } - route, err := tenant.CurrentResultRoute(dispatcherID) + route, err := tenant.ResultRoute(dispatcherID) if err != nil { return err } diff --git a/internal/store/control_flow.go b/internal/store/control_flow.go index 4c14b85..0414ca9 100644 --- a/internal/store/control_flow.go +++ b/internal/store/control_flow.go @@ -134,7 +134,7 @@ func enqueueControlAck(tx *sql.Tx, dispatcherID string, tenantID int64, eventID, if eventID == "" || len(eventID) > 255 || tenantID <= 0 || (status != "applied" && status != "rejected") { return errors.New("invalid task control acknowledgment identity or status") } - route, err := tenant.CurrentResultRoute(dispatcherID) + route, err := tenant.ResultRoute(dispatcherID) if err != nil { return err } diff --git a/internal/store/result.go b/internal/store/result.go index 1b41464..7ec68e6 100644 --- a/internal/store/result.go +++ b/internal/store/result.go @@ -115,7 +115,7 @@ func (s *Store) recordCallResult(dispatcherID, sourceEventID string, payload []b } else if proof != nil || grantErr == nil { return OutboxEvent{}, false, ErrUploadUnverified } - route, err := tenant.CurrentResultRoute(dispatcherID) + route, err := tenant.ResultRoute(dispatcherID) if err != nil { return OutboxEvent{}, false, err } diff --git a/internal/store/result_test.go b/internal/store/result_test.go index aea1d5c..121776a 100644 --- a/internal/store/result_test.go +++ b/internal/store/result_test.go @@ -65,7 +65,7 @@ func TestFinalResultNeedsConfirmedEndAndRemainsExactlyOne(t *testing.T) { if err != nil || !created || event.EventType != "call.execute.result" || event.EventID == cmd.EventID { t.Fatalf("exact final result not persisted: event=%+v created=%t err=%v", event, created, err) } - route, err := tenant.CurrentResultRoute(cmd.DispatcherID) + route, err := tenant.ResultRoute(cmd.DispatcherID) if err != nil || event.RoutingKey != route.BindingKey || contract.ValidateCurrent("mq", event.Body) != nil { t.Fatalf("final result was not Schema-valid or routed to shared SaaS queue: route=%+v err=%v", route, err) } diff --git a/internal/tenant/current.go b/internal/tenant/current.go deleted file mode 100644 index 23431c3..0000000 --- a/internal/tenant/current.go +++ /dev/null @@ -1,49 +0,0 @@ -package tenant - -import ( - "fmt" - "regexp" -) - -// CurrentRoute describes a SaaS-owned exact Topic binding. Dispatcher only -// consumes/publishes; it never declares, binds, or deletes these resources. -type CurrentRoute struct { - Exchange string - BindingKey string - Queue string -} - -var currentTaskID = regexp.MustCompile(`^[A-Za-z0-9_-]{1,128}$`) - -func CurrentControlRoute(dispatcherID string) (CurrentRoute, error) { - if err := ValidateDispatcherID(dispatcherID); err != nil { - return CurrentRoute{}, err - } - return currentRoute("agent-call.dispatchers.v1", "d."+dispatcherID+".control.in", "agent-call.d."+dispatcherID+".control.v1") -} - -func CurrentTaskRoute(dispatcherID, taskID string) (CurrentRoute, error) { - if err := ValidateDispatcherID(dispatcherID); err != nil { - return CurrentRoute{}, err - } - if !currentTaskID.MatchString(taskID) { - return CurrentRoute{}, fmt.Errorf("invalid task ID for exact MQ route") - } - return currentRoute("agent-call.dispatchers.v1", "d."+dispatcherID+".task."+taskID+".in", "agent-call.d."+dispatcherID+".task."+taskID+".v1") -} - -func CurrentResultRoute(dispatcherID string) (CurrentRoute, error) { - if err := ValidateDispatcherID(dispatcherID); err != nil { - return CurrentRoute{}, err - } - // The shared queue name is only a project-local Mock fixture. SaaS owns - // the actual receiving queue and exact bindings for each Dispatcher. - return currentRoute("agent-call.saas.v1", "d."+dispatcherID+".out", "agent-call.saas.events.v1") -} - -func currentRoute(exchange, key, queue string) (CurrentRoute, error) { - if len(key) > 255 || len(queue) > 255 { - return CurrentRoute{}, fmt.Errorf("AMQP route exceeds 255 bytes") - } - return CurrentRoute{Exchange: exchange, BindingKey: key, Queue: queue}, nil -} diff --git a/internal/tenant/routing.go b/internal/tenant/routing.go new file mode 100644 index 0000000..3451b14 --- /dev/null +++ b/internal/tenant/routing.go @@ -0,0 +1,49 @@ +package tenant + +import ( + "fmt" + "regexp" +) + +// Route describes a SaaS-owned exact Topic binding. Dispatcher only +// consumes/publishes; it never declares, binds, or deletes these resources. +type Route struct { + Exchange string + BindingKey string + Queue string +} + +var taskIDPattern = regexp.MustCompile(`^[A-Za-z0-9_-]{1,128}$`) + +func ControlRoute(dispatcherID string) (Route, error) { + if err := ValidateDispatcherID(dispatcherID); err != nil { + return Route{}, err + } + return exactRoute("agent-call.dispatchers.v1", "d."+dispatcherID+".control.in", "agent-call.d."+dispatcherID+".control.v1") +} + +func TaskRoute(dispatcherID, taskID string) (Route, error) { + if err := ValidateDispatcherID(dispatcherID); err != nil { + return Route{}, err + } + if !taskIDPattern.MatchString(taskID) { + return Route{}, fmt.Errorf("invalid task ID for exact MQ route") + } + return exactRoute("agent-call.dispatchers.v1", "d."+dispatcherID+".task."+taskID+".in", "agent-call.d."+dispatcherID+".task."+taskID+".v1") +} + +func ResultRoute(dispatcherID string) (Route, error) { + if err := ValidateDispatcherID(dispatcherID); err != nil { + return Route{}, err + } + // The shared queue name is only a project-local Mock fixture. SaaS owns + // the actual receiving queue and exact bindings for each Dispatcher. + return exactRoute("agent-call.saas.v1", "d."+dispatcherID+".out", "agent-call.saas.events.v1") +} + +func exactRoute(exchange, key, queue string) (Route, error) { + if len(key) > 255 || len(queue) > 255 { + return Route{}, fmt.Errorf("AMQP route exceeds 255 bytes") + } + return Route{Exchange: exchange, BindingKey: key, Queue: queue}, nil +} diff --git a/internal/tenant/current_test.go b/internal/tenant/routing_test.go similarity index 76% rename from internal/tenant/current_test.go rename to internal/tenant/routing_test.go index c4320da..5df1215 100644 --- a/internal/tenant/current_test.go +++ b/internal/tenant/routing_test.go @@ -5,32 +5,32 @@ import ( "testing" ) -func TestCurrentTaskAndControlRoutesAreExactAndPartitioned(t *testing.T) { +func TestTaskAndControlRoutesAreExactAndPartitioned(t *testing.T) { d1 := "c046b893-8628-4589-ae50-619d049248a6" d2 := "c046b893-8628-4589-ae50-619d049248a7" - control, err := CurrentControlRoute(d1) + control, err := ControlRoute(d1) if err != nil { t.Fatal(err) } if control.Exchange != "agent-call.dispatchers.v1" || control.BindingKey != "d."+d1+".control.in" || control.Queue != "agent-call.d."+d1+".control.v1" { t.Fatalf("unapproved control route: %+v", control) } - task, err := CurrentTaskRoute(d1, "task-1") + task, err := TaskRoute(d1, "task-1") if err != nil { t.Fatal(err) } if task.Exchange != control.Exchange || task.BindingKey != "d."+d1+".task.task-1.in" || task.Queue != "agent-call.d."+d1+".task.task-1.v1" { t.Fatalf("unapproved task route: %+v", task) } - other, err := CurrentTaskRoute(d2, "task-1") + other, err := TaskRoute(d2, "task-1") if err != nil || task.Queue == other.Queue || task.BindingKey == other.BindingKey { t.Fatalf("Dispatcher queue collision: %+v, %+v, %v", task, other, err) } - shared1, err := CurrentResultRoute(d1) + shared1, err := ResultRoute(d1) if err != nil { t.Fatal(err) } - shared2, err := CurrentResultRoute(d2) + shared2, err := ResultRoute(d2) if err != nil { t.Fatal(err) } @@ -39,19 +39,19 @@ func TestCurrentTaskAndControlRoutesAreExactAndPartitioned(t *testing.T) { } } -func TestCurrentTaskRoutingRejectsWildcardAndOverflow(t *testing.T) { +func TestTaskRoutingRejectsWildcardAndOverflow(t *testing.T) { d := "c046b893-8628-4589-ae50-619d049248a6" for _, taskID := range []string{"", "task.*", "task.#", "task.with.dot", "task/with/slash", strings.Repeat("a", 129)} { - if _, err := CurrentTaskRoute(d, taskID); err == nil { + if _, err := TaskRoute(d, taskID); err == nil { t.Fatalf("accepted unsafe task identifier %q", taskID) } } for _, dID := range []string{"not-a-uuid", "C046B893-8628-4589-AE50-619D049248A6", ""} { - if _, err := CurrentControlRoute(dID); err == nil { + if _, err := ControlRoute(dID); err == nil { t.Fatalf("accepted unsafe D identity %q", dID) } } - long, err := CurrentTaskRoute(d, strings.Repeat("a", 128)) + long, err := TaskRoute(d, strings.Repeat("a", 128)) if err != nil || len(long.BindingKey) > 255 || len(long.Queue) > 255 { t.Fatalf("valid task exceeded AMQP route limits: %+v %v", long, err) }