diff --git a/contracts/v1_bundle_test.go b/contracts/v1_bundle_test.go index 46b83ac..bb710a0 100644 --- a/contracts/v1_bundle_test.go +++ b/contracts/v1_bundle_test.go @@ -11,8 +11,6 @@ import ( "testing" "github.com/santhosh-tekuri/jsonschema/v6" - - "git.ipao.vip/rogee/go-sip/internal/tenant" ) const v1Dir = "upstream/v1" @@ -171,51 +169,6 @@ func TestV1DispatcherConfigDoesNotAcceptSecretsOrTTLOverride(t *testing.T) { } } -func TestV1TopologyMatchesDispatcherRoutes(t *testing.T) { - var topology struct { - Exchanges map[string]struct { - Type string `json:"type"` - Durable bool `json:"durable"` - } `json:"exchanges"` - Inbox string `json:"inbox_queue"` - Dead string `json:"dead_letter_queue"` - Inbound string `json:"inbound_key"` - Outbound string `json:"outbound_key"` - Durable bool `json:"business_queues_durable"` - Persistent bool `json:"message_persistent"` - Mandatory bool `json:"publish_mandatory"` - Confirm bool `json:"publisher_confirms"` - TenantBytes int `json:"tenant_key_max_utf8_bytes"` - TokenSeconds int `json:"upload_token_seconds"` - } - if err := json.Unmarshal(readV1(t, filepath.Join(v1Dir, "mq-topology.json")), &topology); err != nil { - t.Fatal(err) - } - if !topology.Durable || !topology.Persistent || !topology.Mandatory || !topology.Confirm || topology.TenantBytes != 196 || topology.TokenSeconds != 900 { - t.Fatal("topology weakens the approved queue/token boundary") - } - if len(topology.Exchanges) != 3 { - t.Fatal("unexpected exchange topology") - } - for _, name := range []string{"agent-call.dispatchers.v2", "agent-call.saas.v2", "agent-call.dead-letter.v2"} { - x, ok := topology.Exchanges[name] - if !ok || x.Type != "topic" || !x.Durable { - t.Fatalf("wrong exchange contract: %s", name) - } - } - id := "c046b893-8628-4589-ae50-619d049248a6" - for _, key := range []string{"tenant-a", "租户.甲", strings.Repeat("a", 196)} { - route, err := tenant.NewDispatcherRoute(id, key) - if err != nil { - t.Fatal(err) - } - r := strings.NewReplacer("", id, "", key) - if route.InboxQueue != r.Replace(topology.Inbox) || route.DeadLetterQueue != r.Replace(topology.Dead) || route.InboundKey != r.Replace(topology.Inbound) || route.OutboundKey != r.Replace(topology.Outbound) { - t.Fatal("implementation differs from the published route templates") - } - } -} - func TestV1ManifestPinsEveryBundleFile(t *testing.T) { var manifest struct { Version string `json:"version"` diff --git a/docs/evidence/saas-dispatcher-implementation.md b/docs/evidence/saas-dispatcher-implementation.md index fe972f2..d777ad7 100644 --- a/docs/evidence/saas-dispatcher-implementation.md +++ b/docs/evidence/saas-dispatcher-implementation.md @@ -104,6 +104,7 @@ - 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 接收与应用收讫仍未验证。 +- MQ 旧声明拓扑分支:移除旧 `Broker`、旧租户 Topic 路由/队列声明及其测试;独立保留当前消费端的消息处理器、永久错误分类与 Dispatcher UUID v4 校验,补错误解包和规范身份的回归测试。当前 Broker 仅被动核对 SaaS 预建拓扑;隔离 MQ 测试通过不代表真实 SaaS 应用收讫。 ## 验收台账 diff --git a/internal/mq/amqp.go b/internal/mq/amqp.go deleted file mode 100644 index a35472f..0000000 --- a/internal/mq/amqp.go +++ /dev/null @@ -1,294 +0,0 @@ -// Package mq contains the RabbitMQ adapter. Business state remains in the -// Dispatcher store; this package only declares topology and transports bytes. -package mq - -import ( - "context" - "encoding/json" - "errors" - "fmt" - "log/slog" - "sync" - "sync/atomic" - "time" - - "git.ipao.vip/rogee/go-sip/internal/tenant" - amqp "github.com/rabbitmq/amqp091-go" -) - -var consumerSequence atomic.Uint64 - -const ( - DefaultExchange = "agent-call.dispatchers.v2" - EventExchange = "agent-call.saas.v2" - DeadLetterExchange = "agent-call.dead-letter.v2" - SaaSQueue = "agent-call.saas.events.v2" - DefaultPrefetch = 1 - MaxMessageBytes = 256 << 10 -) - -type Publisher interface { - Publish(context.Context, string, string, []byte) error -} - -type permanentError struct{ err error } - -func (e permanentError) Error() string { return e.err.Error() } -func (e permanentError) Unwrap() error { return e.err } - -func Permanent(err error) error { - if err == nil { - return nil - } - return permanentError{err: err} -} - -func IsPermanent(err error) bool { - var target permanentError - return errors.As(err, &target) -} - -type Broker struct { - conn *amqp.Connection - channel *amqp.Channel - dispatcherID string - prefetch int - outbound map[string]bool - inboxes map[string]bool - closed <-chan *amqp.Error - mu sync.Mutex -} - -func Open(url, dispatcherID string) (*Broker, error) { - return OpenWithPrefetch(url, dispatcherID, DefaultPrefetch) -} - -func OpenWithPrefetch(url, dispatcherID string, prefetch int) (*Broker, error) { - if 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 channel: %w", err) - } - owner := "agent-call.d." + dispatcherID + ".owner.v2" - if _, err := channel.QueueDeclare(owner, false, false, true, false, nil); err != nil { - _ = conn.Close() - return nil, fmt.Errorf("claim dispatcher identity %s: %w", dispatcherID, err) - } - for _, exchange := range []string{DefaultExchange, EventExchange, DeadLetterExchange} { - if err := channel.ExchangeDeclare(exchange, "topic", true, false, false, false, nil); err != nil { - _ = conn.Close() - return nil, fmt.Errorf("declare topic exchange %s: %w", exchange, err) - } - } - return &Broker{conn: conn, channel: channel, dispatcherID: dispatcherID, prefetch: prefetch, - outbound: make(map[string]bool), inboxes: make(map[string]bool), closed: conn.NotifyClose(make(chan *amqp.Error, 1))}, nil -} - -func (b *Broker) Close() error { - b.mu.Lock() - defer b.mu.Unlock() - if b.conn == nil { - return nil - } - err := b.conn.Close() - b.channel, b.conn = nil, nil - return err -} - -// Done reports loss of the connection that owns this Dispatcher identity. -// The caller must stop admission when that ownership connection is lost. -func (b *Broker) Done() <-chan *amqp.Error { return b.closed } - -func (b *Broker) Publish(ctx context.Context, exchange, routingKey string, body []byte) error { - if exchange != EventExchange || routingKey == "" || len(body) == 0 || len(body) > MaxMessageBytes { - return errors.New("publish requires SaaS exchange, declared route and 1..262144 body bytes") - } - 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 connection is closed") - } - if !b.outbound[routingKey] { - return errors.New("outbound route does not belong to a declared Dispatcher/tenant") - } - // A separate channel isolates late returns/confirms after an ambiguous send. - // Connections are reused; a timed-out publication can never acknowledge the next one. - channel, err := b.conn.Channel() - if err != nil { - return fmt.Errorf("open publication channel: %w", err) - } - defer channel.Close() - if _, err := channel.QueueDeclarePassive(SaaSQueue, true, false, false, false, nil); err != nil { - return fmt.Errorf("required SaaS queue unavailable: %w", err) - } - if err := channel.Confirm(false); err != nil { - return fmt.Errorf("enable 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: %w", err) - } - if confirmation == nil { - return errors.New("rabbitmq publisher confirmation unavailable") - } - acked, err := confirmation.WaitContext(ctx) - if err != nil { - return fmt.Errorf("wait publisher confirmation: %w", err) - } - // RabbitMQ sends basic.return before its confirm; the SDK dispatches it - // before completing the deferred confirmation. A positive confirm alone is insufficient. - select { - case result, ok := <-returned: - if !ok { - return errors.New("publication channel closed before routing was established") - } - return fmt.Errorf("publication returned: code=%d", result.ReplyCode) - default: - } - if !acked { - return errors.New("rabbitmq publisher was negatively acknowledged") - } - return nil -} - -func (b *Broker) DeclareTenantQueue(tenantKey string) (string, error) { - route, err := tenant.NewDispatcherRoute(b.dispatcherID, tenantKey) - if err != nil { - return "", err - } - b.mu.Lock() - defer b.mu.Unlock() - if b.channel == nil { - return "", errors.New("rabbitmq channel is closed") - } - queue, routingKey, deadLetterQueue := route.InboxQueue, route.InboundKey, route.DeadLetterQueue - queueArgs := amqp.Table{ - "x-dead-letter-exchange": DeadLetterExchange, - "x-dead-letter-routing-key": routingKey, - } - if _, err := b.channel.QueueDeclare(queue, true, false, false, false, queueArgs); err != nil { - return "", fmt.Errorf("declare tenant queue: %w", err) - } - if _, err := b.channel.QueueDeclare(deadLetterQueue, true, false, false, false, nil); err != nil { - return "", fmt.Errorf("declare tenant dead-letter queue: %w", err) - } - if err := b.channel.QueueBind(queue, routingKey, DefaultExchange, false, nil); err != nil { - return "", fmt.Errorf("bind tenant queue: %w", err) - } - if err := b.channel.QueueBind(deadLetterQueue, routingKey, DeadLetterExchange, false, nil); err != nil { - return "", fmt.Errorf("bind tenant dead-letter queue: %w", err) - } - if _, err := b.channel.QueueDeclare(SaaSQueue, true, false, false, false, nil); err != nil { - return "", fmt.Errorf("declare SaaS queue: %w", err) - } - if err := b.channel.QueueBind(SaaSQueue, route.OutboundKey, EventExchange, false, nil); err != nil { - return "", fmt.Errorf("bind SaaS queue: %w", err) - } - b.outbound[route.OutboundKey], b.inboxes[queue] = true, true - return queue, nil -} - -type MessageHandler func(context.Context, string, []byte) error - -// Consume ACKs only after the handler returns nil. A transient handler error -// requeues; malformed or unauthorized messages can be rejected by the caller -// with Permanent, which RabbitMQ dead-letters through the tenant queue policy. -func (b *Broker) Consume(ctx context.Context, queue string, handler MessageHandler) error { - if queue == "" || handler == nil { - return errors.New("queue and handler are required") - } - b.mu.Lock() - if b.channel == nil { - b.mu.Unlock() - return errors.New("rabbitmq channel is closed") - } - if !b.inboxes[queue] { - b.mu.Unlock() - return errors.New("queue does not belong to a declared Dispatcher/tenant") - } - prefetch := b.prefetch - if err := b.channel.Qos(prefetch, 0, false); err != nil { - b.mu.Unlock() - return fmt.Errorf("set tenant prefetch: %w", err) - } - consumerTag := fmt.Sprintf("sip-go-agent-%d", consumerSequence.Add(1)) - deliveries, err := b.channel.Consume(queue, consumerTag, false, false, false, false, nil) - b.mu.Unlock() - if err != nil { - return fmt.Errorf("consume tenant queue: %w", err) - } - defer func() { - b.mu.Lock() - if b.channel != nil { - _ = b.channel.Cancel(consumerTag, false) - } - b.mu.Unlock() - }() - for { - select { - case <-ctx.Done(): - return ctx.Err() - case d, ok := <-deliveries: - if !ok { - return errors.New("rabbitmq delivery channel closed") - } - if len(d.Body) == 0 || len(d.Body) > MaxMessageBytes { - if err := d.Reject(false); err != nil { - return fmt.Errorf("reject invalid message size: %w", err) - } - slog.Warn("MQ message rejected", "dispatcher_id", b.dispatcherID, "delivery_tag", d.DeliveryTag, "reason", "invalid_message_size", "bytes", len(d.Body)) - continue - } - if err := handler(ctx, d.RoutingKey, d.Body); err != nil { - if IsPermanent(err) { - if rejectErr := d.Reject(false); rejectErr != nil { - return fmt.Errorf("permanent handler error %v; reject: %w", err, rejectErr) - } - continue - } - if nackErr := d.Nack(false, true); nackErr != nil { - return fmt.Errorf("handler error %v; nack: %w", err, nackErr) - } - continue - } - if err := d.Ack(false); err != nil { - return fmt.Errorf("ack delivery: %w", err) - } - } - } -} diff --git a/internal/mq/amqp_test.go b/internal/mq/amqp_test.go deleted file mode 100644 index 4b43762..0000000 --- a/internal/mq/amqp_test.go +++ /dev/null @@ -1,13 +0,0 @@ -package mq - -import ( - "strings" - "testing" -) - -func TestOpenWithPrefetchRejectsZero(t *testing.T) { - _, err := OpenWithPrefetch("amqp://unused", "", 0) - if err == nil || !strings.Contains(err.Error(), "prefetch") { - t.Fatalf("error=%v, want prefetch validation", err) - } -} diff --git a/internal/mq/disconnect_integration_test.go b/internal/mq/disconnect_integration_test.go deleted file mode 100644 index bb2e38d..0000000 --- a/internal/mq/disconnect_integration_test.go +++ /dev/null @@ -1,44 +0,0 @@ -package mq - -import ( - "context" - "strings" - "testing" - "time" - - "git.ipao.vip/rogee/go-sip/internal/tenant" - "github.com/google/uuid" -) - -func TestV2LocalBrokerDisconnectStopsPublisher(t *testing.T) { - address := localBrokerURL(t) - broker, err := Open(address, uuid.NewString()) - if err != nil { - t.Fatal(err) - } - defer broker.Close() - tenantKey := "disconnect-" + uuid.NewString() - if _, err := broker.DeclareTenantQueue(tenantKey); err != nil { - t.Fatal(err) - } - route, err := tenant.NewDispatcherRoute(broker.dispatcherID, tenantKey) - if err != nil { - t.Fatal(err) - } - // Closing the underlying AMQP connection models a broker disconnect while - // leaving the Broker object in the state observed by the runtime watcher. - if err := broker.conn.Close(); err != nil { - t.Fatal(err) - } - select { - case <-broker.Done(): - case <-time.After(3 * time.Second): - t.Fatal("broker disconnect was not surfaced") - } - ctx, cancel := context.WithTimeout(context.Background(), time.Second) - defer cancel() - err = broker.Publish(ctx, EventExchange, route.OutboundKey, []byte(`{"event_id":"after-disconnect"}`)) - if err == nil || !strings.Contains(err.Error(), "closed") { - t.Fatalf("publisher accepted a disconnected broker: %v", err) - } -} diff --git a/internal/mq/handler.go b/internal/mq/handler.go new file mode 100644 index 0000000..a395e8f --- /dev/null +++ b/internal/mq/handler.go @@ -0,0 +1,29 @@ +package mq + +import ( + "context" + "errors" + "sync/atomic" +) + +var consumerSequence atomic.Uint64 + +type MessageHandler func(context.Context, string, []byte) error + +type permanentError struct{ err error } + +func (e permanentError) Error() string { return e.err.Error() } +func (e permanentError) Unwrap() error { return e.err } + +// Permanent marks an inbound handler failure that must not be requeued. +func Permanent(err error) error { + if err == nil { + return nil + } + return permanentError{err: err} +} + +func IsPermanent(err error) bool { + var target permanentError + return errors.As(err, &target) +} diff --git a/internal/mq/handler_test.go b/internal/mq/handler_test.go new file mode 100644 index 0000000..0d9ec97 --- /dev/null +++ b/internal/mq/handler_test.go @@ -0,0 +1,24 @@ +package mq + +import ( + "errors" + "fmt" + "testing" +) + +func TestPermanentHandlerErrorClassification(t *testing.T) { + cause := errors.New("invalid inbound command") + if got := Permanent(nil); got != nil { + t.Fatalf("nil permanent error = %v, want nil", got) + } + permanent := Permanent(cause) + if !errors.Is(permanent, cause) { + t.Fatal("permanent error lost its cause") + } + if !IsPermanent(permanent) || !IsPermanent(fmt.Errorf("handler: %w", permanent)) { + t.Fatal("wrapped permanent handler error must reject delivery") + } + if IsPermanent(cause) || IsPermanent(nil) { + t.Fatal("transient or absent handler error must not reject delivery") + } +} diff --git a/internal/mq/integration_test.go b/internal/mq/integration_test.go deleted file mode 100644 index fe36d6b..0000000 --- a/internal/mq/integration_test.go +++ /dev/null @@ -1,239 +0,0 @@ -//go:build integration - -package mq - -import ( - "context" - "errors" - "fmt" - "os" - "strings" - "testing" - "time" - - "git.ipao.vip/rogee/go-sip/internal/tenant" - "github.com/google/uuid" - amqp "github.com/rabbitmq/amqp091-go" -) - -func TestLocalRabbitMQConfirmAckAndDeadLetter(t *testing.T) { - url := os.Getenv("RABBITMQ_URL") - if url == "" { - t.Skip("RABBITMQ_URL is not configured") - } - broker, err := OpenWithPrefetch(url, uuid.NewString(), 1) - if err != nil { - t.Fatal(err) - } - defer broker.Close() - - tenantKey := fmt.Sprintf("integration-%d", time.Now().UnixNano()) - queue, err := broker.DeclareTenantQueue(tenantKey) - if err != nil { - t.Fatal(err) - } - route, err := tenant.NewDispatcherRoute(broker.dispatcherID, tenantKey) - if err != nil { - t.Fatal(err) - } - routingKey := route.InboundKey - deadLetterQueue := route.DeadLetterQueue - body := []byte(`{"command_id":"integration-command"}`) - - ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second) - defer cancel() - publishCommand := func(ctx context.Context) error { - channel, err := broker.conn.Channel() - if err != nil { - return err - } - defer channel.Close() - if err := channel.Confirm(false); err != nil { - return err - } - returned := channel.NotifyReturn(make(chan amqp.Return, 1)) - confirmation, err := channel.PublishWithDeferredConfirmWithContext(ctx, DefaultExchange, routingKey, true, false, amqp.Publishing{ - DeliveryMode: amqp.Persistent, - Body: body, - }) - if err != nil { - return err - } - acked, err := confirmation.WaitContext(ctx) - if err != nil { - return err - } - select { - case result := <-returned: - return fmt.Errorf("command publication returned: code=%d", result.ReplyCode) - default: - } - if !acked { - return errors.New("command publication was negatively acknowledged") - } - return nil - } - acked := make(chan struct{}, 1) - consumeDone := make(chan error, 1) - go func() { - consumeDone <- broker.Consume(ctx, queue, func(_ context.Context, gotRoutingKey string, gotBody []byte) error { - if gotRoutingKey != routingKey || string(gotBody) != string(body) { - return fmt.Errorf("delivery mismatch: key=%q body=%q", gotRoutingKey, gotBody) - } - acked <- struct{}{} - return nil - }) - }() - if err := publishCommand(ctx); err != nil { - t.Fatal(err) - } - select { - case <-acked: - case <-ctx.Done(): - t.Fatal(ctx.Err()) - } - cancel() - select { - case <-consumeDone: - case <-time.After(3 * time.Second): - t.Fatal("consumer did not stop") - } - - permanentCtx, permanentCancel := context.WithTimeout(context.Background(), 15*time.Second) - defer permanentCancel() - permanentSeen := make(chan struct{}, 1) - permanentDone := make(chan error, 1) - go func() { - permanentDone <- broker.Consume(permanentCtx, queue, func(_ context.Context, _ string, _ []byte) error { - permanentSeen <- struct{}{} - return Permanent(errors.New("synthetic permanent command error")) - }) - }() - if err := publishCommand(permanentCtx); err != nil { - t.Fatal(err) - } - select { - case <-permanentSeen: - case <-permanentCtx.Done(): - t.Fatal(permanentCtx.Err()) - } - permanentCancel() - select { - case <-permanentDone: - case <-time.After(3 * time.Second): - t.Fatal("permanent consumer did not stop") - } - - dlqChannel, err := broker.conn.Channel() - if err != nil { - t.Fatal(err) - } - defer dlqChannel.Close() - if _, err := dlqChannel.QueueDeclarePassive(deadLetterQueue, true, false, false, false, nil); err != nil { - t.Fatal(err) - } - dlqDeliveries, err := dlqChannel.Consume(deadLetterQueue, "sip-go-agent-dlq", false, false, false, false, nil) - if err != nil { - t.Fatal(err) - } - deadLettered := make(chan struct{}, 1) - dlqDone := make(chan error, 1) - go func() { - for delivery := range dlqDeliveries { - if string(delivery.Body) != string(body) { - dlqDone <- fmt.Errorf("dead-letter body mismatch: %q", delivery.Body) - return - } - if err := delivery.Ack(false); err != nil { - dlqDone <- err - return - } - deadLettered <- struct{}{} - dlqDone <- nil - return - } - dlqDone <- errors.New("dead-letter delivery channel closed") - }() - select { - case <-deadLettered: - _ = dlqChannel.Cancel("sip-go-agent-dlq", false) - case err := <-dlqDone: - t.Fatal(err) - case <-time.After(15 * time.Second): - t.Fatal("dead-letter message did not arrive") - } - select { - case err := <-dlqDone: - if err != nil { - t.Fatal(err) - } - case <-time.After(3 * time.Second): - t.Fatal("dead-letter consumer did not stop") - } -} - -func TestLocalRabbitMQV2PublishesOnlyConfirmedRoutedMessages(t *testing.T) { - url := os.Getenv("RABBITMQ_URL") - if url == "" { - t.Skip("RABBITMQ_URL is not configured") - } - broker, err := OpenWithPrefetch(url, uuid.NewString(), 1) - if err != nil { - t.Fatal(err) - } - defer broker.Close() - if _, err := broker.DeclareTenantQueue("publish-" + uuid.NewString()); err != nil { - t.Fatal(err) - } - var outbound string - for route := range broker.outbound { - outbound = route - } - ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second) - defer cancel() - - reject := func(exchange, route string, body []byte, reason string) { - t.Helper() - if err := broker.Publish(ctx, exchange, route, body); err == nil || !strings.Contains(err.Error(), reason) { - t.Fatalf("publish exchange=%q route=%q body-size=%d: got %v, want %q", exchange, route, len(body), err, reason) - } - } - reject(DefaultExchange, outbound, []byte(`{"event_id":"wrong-exchange"}`), "SaaS exchange") - reject(EventExchange, "", []byte(`{"event_id":"empty-route"}`), "declared route") - reject(EventExchange, outbound, nil, "body bytes") - reject(EventExchange, outbound, make([]byte, MaxMessageBytes+1), "body bytes") - reject(EventExchange, outbound, []byte(`{`), "decode outbound message identity") - reject(EventExchange, outbound, []byte(`{"event_id":"a","message_id":"b"}`), "exactly one") - reject(EventExchange, outbound, []byte(`{}`), "exactly one") - reject(EventExchange, outbound, []byte(fmt.Sprintf(`{"event_id":"%s"}`, strings.Repeat("x", 256))), "identity exceeds") - reject(EventExchange, outbound+".other", []byte(`{"event_id":"unknown-route"}`), "does not belong") - - channel, err := broker.conn.Channel() - if err != nil { - t.Fatal(err) - } - defer channel.Close() - for _, idField := range []string{"event_id", "message_id"} { - id := uuid.NewString() - body := []byte(fmt.Sprintf(`{"%s":"%s"}`, idField, id)) - if err := broker.Publish(ctx, EventExchange, outbound, body); err != nil { - t.Fatalf("confirmed %s publication: %v", idField, err) - } - delivery, ok, err := channel.Get(SaaSQueue, true) - if err != nil || !ok || string(delivery.Body) != string(body) || delivery.MessageId != id || delivery.DeliveryMode != amqp.Persistent { - t.Fatalf("confirmed message lost identity or durability: delivery=%+v ok=%t err=%v", delivery, ok, err) - } - } - - if err := channel.QueueUnbind(SaaSQueue, outbound, EventExchange, nil); err != nil { - t.Fatal(err) - } - reject(EventExchange, outbound, []byte(`{"event_id":"unroutable"}`), "publication returned") - if err := channel.QueueBind(SaaSQueue, outbound, EventExchange, false, nil); err != nil { - t.Fatal(err) - } - if err := broker.Close(); err != nil { - t.Fatal(err) - } - reject(EventExchange, outbound, []byte(`{"event_id":"after-close"}`), "closed") -} diff --git a/internal/mq/v2_integration_test.go b/internal/mq/v2_integration_test.go deleted file mode 100644 index b9068d3..0000000 --- a/internal/mq/v2_integration_test.go +++ /dev/null @@ -1,172 +0,0 @@ -package mq - -import ( - "context" - "net/url" - "os" - "strings" - "testing" - "time" - - "git.ipao.vip/rogee/go-sip/internal/tenant" - "github.com/google/uuid" - amqp "github.com/rabbitmq/amqp091-go" -) - -// This opt-in is deliberately separate from any external broker setting. -func localBrokerURL(t *testing.T) string { - t.Helper() - value := os.Getenv("GO_SIP_LOCAL_MQ_URL") - if value == "" { - t.Skip("isolated local RabbitMQ not configured") - } - parsed, err := url.Parse(value) - if err != nil || (parsed.Hostname() != "127.0.0.1" && parsed.Hostname() != "::1" && parsed.Hostname() != "localhost") { - t.Fatal("MQ integration requires a loopback broker") - } - return value -} - -func TestV2LocalBrokerIdentityIsolationAndReliableRouting(t *testing.T) { - address := localBrokerURL(t) - id1, id2 := uuid.NewString(), uuid.NewString() - first, err := Open(address, id1) - if err != nil { - t.Fatal(err) - } - defer first.Close() - second, err := Open(address, id2) - if err != nil { - t.Fatal(err) - } - defer second.Close() - if duplicate, err := Open(address, id1); err == nil { - duplicate.Close() - t.Fatal("duplicate live identity accepted") - } - key := "租户." + uuid.NewString() - q1, err := first.DeclareTenantQueue(key) - if err != nil { - t.Fatal(err) - } - q2, err := second.DeclareTenantQueue(key) - if err != nil { - t.Fatal(err) - } - if q1 == q2 { - t.Fatal("different Dispatchers share a queue") - } - r1, _ := tenant.NewDispatcherRoute(id1, key) - r2, _ := tenant.NewDispatcherRoute(id2, key) - connection, err := amqp.Dial(address) - if err != nil { - t.Fatal(err) - } - defer connection.Close() - channel, err := connection.Channel() - if err != nil { - t.Fatal(err) - } - defer channel.Close() - defer channel.QueueDelete(q1, false, false, false) - defer channel.QueueDelete(q2, false, false, false) - defer channel.QueueDelete(r1.DeadLetterQueue, false, false, false) - defer channel.QueueDelete(r2.DeadLetterQueue, false, false, false) - defer channel.QueueUnbind(SaaSQueue, r1.OutboundKey, EventExchange, nil) - defer channel.QueueUnbind(SaaSQueue, r2.OutboundKey, EventExchange, nil) - if err := channel.Confirm(false); err != nil { - t.Fatal(err) - } - ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) - defer cancel() - confirmation, err := channel.PublishWithDeferredConfirmWithContext(ctx, DefaultExchange, r1.InboundKey, true, false, amqp.Publishing{DeliveryMode: amqp.Persistent, Body: []byte(`{"target":"first"}`)}) - if err != nil { - t.Fatal(err) - } - if ack, err := confirmation.WaitContext(ctx); err != nil || !ack { - t.Fatalf("command publication: ack=%v error=%v", ack, err) - } - if _, ok, err := channel.Get(q2, true); err != nil || ok { - t.Fatalf("second Dispatcher received first's message: ok=%v error=%v", ok, err) - } - message, ok, err := channel.Get(q1, true) - if err != nil || !ok || string(message.Body) != `{"target":"first"}` { - t.Fatalf("first did not receive its own message: ok=%v error=%v", ok, err) - } - if err := first.Consume(ctx, q2, func(context.Context, string, []byte) error { return nil }); err == nil { - t.Fatal("foreign queue consumption accepted") - } - if err := first.Publish(ctx, EventExchange, r2.OutboundKey, []byte(`{}`)); err == nil { - t.Fatal("foreign publication accepted") - } - if err := first.Publish(ctx, EventExchange, r1.OutboundKey, []byte(`{"event_id":"uploaded"}`)); err != nil { - t.Fatal(err) - } - message, ok, err = channel.Get(SaaSQueue, true) - if err != nil || !ok || message.RoutingKey != r1.OutboundKey || message.DeliveryMode != amqp.Persistent { - t.Fatalf("expected persistent target-queue delivery: ok=%v error=%v", ok, err) - } - if err := channel.QueueUnbind(SaaSQueue, r1.OutboundKey, EventExchange, nil); err != nil { - t.Fatal(err) - } - if err := first.Publish(ctx, EventExchange, r1.OutboundKey, []byte(`{"event_id":"unroutable"}`)); err == nil || !strings.Contains(err.Error(), "returned") { - t.Fatalf("unroutable confirm treated as delivery: %v", err) - } - if err := channel.QueueBind(SaaSQueue, r1.OutboundKey, EventExchange, false, nil); err != nil { - t.Fatal(err) - } - if err := first.Publish(ctx, EventExchange, r1.OutboundKey, []byte(`{"event_id":"recovered"}`)); err != nil { - t.Fatalf("return contaminated later publication: %v", err) - } - message, ok, err = channel.Get(SaaSQueue, true) - if err != nil || !ok || string(message.Body) != `{"event_id":"recovered"}` || message.MessageId != "recovered" { - t.Fatalf("recovery delivery missing: ok=%v error=%v", ok, err) - } - consumeCtx, stop := context.WithCancel(ctx) - defer stop() - called := make(chan struct{}, 1) - finished := make(chan error, 1) - go func() { - finished <- first.Consume(consumeCtx, q1, func(context.Context, string, []byte) error { called <- struct{}{}; return nil }) - }() - confirmation, err = channel.PublishWithDeferredConfirmWithContext(ctx, DefaultExchange, r1.InboundKey, true, false, amqp.Publishing{DeliveryMode: amqp.Persistent, Body: []byte(strings.Repeat("x", MaxMessageBytes+1))}) - if err != nil { - t.Fatal(err) - } - if ack, err := confirmation.WaitContext(ctx); err != nil || !ack { - t.Fatalf("large message publication: %v", err) - } - deadline := time.NewTimer(2 * time.Second) - defer deadline.Stop() - ticker := time.NewTicker(10 * time.Millisecond) - defer ticker.Stop() - for { - select { - case <-called: - t.Fatal("oversized message reached the business handler") - case <-deadline.C: - t.Fatal("oversized message was not durably dead-lettered") - case <-ticker.C: - message, ok, err := channel.Get(r1.DeadLetterQueue, true) - if err != nil { - t.Fatal(err) - } - if ok { - if len(message.Body) != MaxMessageBytes+1 { - t.Fatal("wrong dead-lettered message") - } - stop() - <-finished - return - } - } - } -} - -func TestOpenRejectsNonCanonicalIdentityBeforeDial(t *testing.T) { - for _, id := range []string{"", DefaultExchange, "dispatcher", "00000000-0000-0000-0000-000000000000"} { - if _, err := Open("amqp://unused.invalid", id); err == nil || !strings.Contains(err.Error(), "dispatcher_id") { - t.Fatalf("identity %q not rejected before dial: %v", id, err) - } - } -} diff --git a/internal/tenant/dispatcher.go b/internal/tenant/dispatcher.go deleted file mode 100644 index f15438f..0000000 --- a/internal/tenant/dispatcher.go +++ /dev/null @@ -1,58 +0,0 @@ -package tenant - -import ( - "fmt" - "strings" - - "github.com/google/uuid" - - "git.ipao.vip/rogee/go-sip/internal/contract" -) - -// DispatcherRoute contains exact v2 bindings. Tenant keys are never rewritten. -type DispatcherRoute struct { - InboxQueue string - DeadLetterQueue string - InboundKey string - OutboundKey string -} - -// ValidateDispatcherID requires the stable, deployment-assigned UUID v4 form. -func ValidateDispatcherID(id string) error { - parsed, err := uuid.Parse(id) - if err != nil || parsed.Version() != 4 || parsed.Variant() != uuid.RFC4122 || parsed.String() != id { - return fmt.Errorf("dispatcher_id must be a canonical lowercase UUID v4") - } - return nil -} - -// NewDispatcherRoute rejects Topic wildcard words rather than changing a tenant -// identity. The longest queue name, including its dead-letter suffix, determines -// the 196-byte tenant budget for a 36-byte Dispatcher ID. -func NewDispatcherRoute(dispatcherID, tenantKey string) (DispatcherRoute, error) { - if err := ValidateDispatcherID(dispatcherID); err != nil { - return DispatcherRoute{}, err - } - if err := contract.ValidateTenantKey(tenantKey); err != nil { - return DispatcherRoute{}, err - } - for _, word := range strings.Split(tenantKey, ".") { - if word == "*" || word == "#" { - return DispatcherRoute{}, fmt.Errorf("tenant_key contains a Topic wildcard word; preserve the source task without publishing") - } - } - queuePrefix := "agent-call.d." + dispatcherID + ".t." + tenantKey - keyPrefix := "d." + dispatcherID + ".t." + tenantKey - route := DispatcherRoute{ - InboxQueue: queuePrefix + ".v2", - DeadLetterQueue: queuePrefix + ".dlq.v2", - InboundKey: keyPrefix + ".in", - OutboundKey: keyPrefix + ".out", - } - for _, name := range []string{route.InboxQueue, route.DeadLetterQueue, route.InboundKey, route.OutboundKey} { - if len(name) > 255 { - return DispatcherRoute{}, fmt.Errorf("Dispatcher route exceeds AMQP's 255-byte limit; tenant_key maximum is 196 UTF-8 bytes") - } - } - return route, nil -} diff --git a/internal/tenant/dispatcher_test.go b/internal/tenant/dispatcher_test.go deleted file mode 100644 index c12f08d..0000000 --- a/internal/tenant/dispatcher_test.go +++ /dev/null @@ -1,100 +0,0 @@ -package tenant - -import ( - "strings" - "testing" -) - -const ( - testDispatcherA = "c046b893-8628-4589-ae50-619d049248a6" - testDispatcherB = "bd72ec77-7296-4da6-8742-732bdb3dbf97" -) - -func TestDispatcherRouteExactNames(t *testing.T) { - route, err := NewDispatcherRoute(testDispatcherA, "tenant-a") - if err != nil { - t.Fatal(err) - } - want := DispatcherRoute{ - InboxQueue: "agent-call.d." + testDispatcherA + ".t.tenant-a.v2", - DeadLetterQueue: "agent-call.d." + testDispatcherA + ".t.tenant-a.dlq.v2", - InboundKey: "d." + testDispatcherA + ".t.tenant-a.in", - OutboundKey: "d." + testDispatcherA + ".t.tenant-a.out", - } - if route != want { - t.Fatalf("route = %#v, want %#v", route, want) - } - other, err := NewDispatcherRoute(testDispatcherB, "tenant-a") - if err != nil { - t.Fatal(err) - } - if route.InboxQueue == other.InboxQueue || route.InboundKey == other.InboundKey || route.OutboundKey == other.OutboundKey || route.DeadLetterQueue == other.DeadLetterQueue { - t.Fatal("different Dispatcher identities share a route") - } -} - -func TestDispatcherRoutePreservesSafeTenantKeys(t *testing.T) { - for _, key := range []string{"tenant-a", "租户甲", "a.b", ".a..b.", "a*b", "a#b", "a*", "#a", "A B/甲"} { - t.Run(key, func(t *testing.T) { - route, err := NewDispatcherRoute(testDispatcherA, key) - if err != nil { - t.Fatal(err) - } - if route.InboundKey != "d."+testDispatcherA+".t."+key+".in" || route.InboxQueue != "agent-call.d."+testDispatcherA+".t."+key+".v2" { - t.Fatalf("tenant key changed: %#v", route) - } - }) - } -} - -func TestDispatcherRouteRejectsWildcardWords(t *testing.T) { - for _, key := range []string{"*", "#", "a.*", "#.a", "a.#.b", ".*.", "a..#", "*.#"} { - t.Run(key, func(t *testing.T) { - if _, err := NewDispatcherRoute(testDispatcherA, key); err == nil { - t.Fatalf("unsafe Topic binding accepted: %q", key) - } - }) - } -} - -func TestDispatcherRouteByteBudget(t *testing.T) { - for _, key := range []string{strings.Repeat("a", 196), strings.Repeat("甲", 65) + "a"} { - route, err := NewDispatcherRoute(testDispatcherA, key) - if err != nil { - t.Fatal(err) - } - if len(route.DeadLetterQueue) != 255 { - t.Fatalf("dead-letter queue = %d bytes, want 255", len(route.DeadLetterQueue)) - } - for _, name := range []string{route.InboxQueue, route.DeadLetterQueue, route.InboundKey, route.OutboundKey} { - if len(name) > 255 { - t.Fatalf("AMQP short-string budget exceeded: %d", len(name)) - } - } - } - for _, key := range []string{"", string([]byte{0xff}), strings.Repeat("a", 197), strings.Repeat("甲", 66), strings.Repeat("a", 224)} { - if _, err := NewDispatcherRoute(testDispatcherA, key); err == nil { - t.Fatalf("invalid or over-budget tenant key accepted: %q", key) - } - } -} - -func TestDispatcherIdentityRequiresCanonicalV4(t *testing.T) { - if err := ValidateDispatcherID(testDispatcherA); err != nil { - t.Fatal(err) - } - for _, id := range []string{ - "", "dispatcher-a", strings.ToUpper(testDispatcherA), - "00000000-0000-0000-0000-000000000000", - "c046b893-8628-1589-ae50-619d049248a6", - "c046b893-8628-4589-7e50-619d049248a6", - "c046b89386284589ae50619d049248a6", "urn:uuid:" + testDispatcherA, - } { - if err := ValidateDispatcherID(id); err == nil { - t.Fatalf("invalid Dispatcher identity accepted: %q", id) - } - if _, err := NewDispatcherRoute(id, "tenant-a"); err == nil { - t.Fatalf("route accepted invalid identity: %q", id) - } - } -} diff --git a/internal/tenant/identity.go b/internal/tenant/identity.go new file mode 100644 index 0000000..a76786b --- /dev/null +++ b/internal/tenant/identity.go @@ -0,0 +1,16 @@ +package tenant + +import ( + "fmt" + + "github.com/google/uuid" +) + +// ValidateDispatcherID requires a canonical lowercase UUID v4. +func ValidateDispatcherID(id string) error { + parsed, err := uuid.Parse(id) + if err != nil || parsed.Version() != 4 || parsed.Variant() != uuid.RFC4122 || parsed.String() != id { + return fmt.Errorf("dispatcher_id must be a canonical lowercase UUID v4") + } + return nil +} diff --git a/internal/tenant/identity_test.go b/internal/tenant/identity_test.go new file mode 100644 index 0000000..743e3e9 --- /dev/null +++ b/internal/tenant/identity_test.go @@ -0,0 +1,24 @@ +package tenant + +import ( + "strings" + "testing" +) + +func TestValidateDispatcherIDCanonicalV4(t *testing.T) { + const approved = "c046b893-8628-4589-ae50-619d049248a6" + if err := ValidateDispatcherID(approved); err != nil { + t.Fatal(err) + } + for _, id := range []string{ + "", "dispatcher-a", strings.ToUpper(approved), + "00000000-0000-0000-0000-000000000000", + "c046b893-8628-1589-ae50-619d049248a6", + "c046b893-8628-4589-7e50-619d049248a6", + "c046b89386284589ae50619d049248a6", "urn:uuid:" + approved, + } { + if err := ValidateDispatcherID(id); err == nil { + t.Fatalf("invalid Dispatcher identity accepted: %q", id) + } + } +} diff --git a/internal/tenant/routing.go b/internal/tenant/routing.go deleted file mode 100644 index 4a862a5..0000000 --- a/internal/tenant/routing.go +++ /dev/null @@ -1,58 +0,0 @@ -package tenant - -import ( - "fmt" - - "git.ipao.vip/rogee/go-sip/internal/contract" -) - -const ( - CommandExchange = "agent-call.commands.v1" - EventExchange = "agent-call.events.v1" - CommandQueuePrefix = "agent-call.executor." - CommandQueueSuffix = ".v1" - CommandRoutingPrefix = "agent-call.tenant." - CommandRoutingSuffix = ".call.execute" - DeadLetterQueueSuffix = ".dlq.v1" - maxAMQPNameBytes = 255 -) - -func CommandQueue(tenantKey string) (string, error) { - if err := contract.ValidateTenantKey(tenantKey); err != nil { - return "", err - } - return boundedQueueName(tenantKey, CommandQueueSuffix) -} - -func DeadLetterQueue(tenantKey string) (string, error) { - if err := contract.ValidateTenantKey(tenantKey); err != nil { - return "", err - } - return boundedQueueName(tenantKey, DeadLetterQueueSuffix) -} - -func boundedQueueName(tenantKey, suffix string) (string, error) { - queue := CommandQueuePrefix + tenantKey + suffix - if len([]byte(queue)) > maxAMQPNameBytes { - return "", fmt.Errorf("AMQP queue name exceeds %d UTF-8 bytes", maxAMQPNameBytes) - } - return queue, nil -} - -func CommandRoutingKey(tenantKey string) (string, error) { - if err := contract.ValidateTenantKey(tenantKey); err != nil { - return "", err - } - return CommandRoutingPrefix + tenantKey + CommandRoutingSuffix, nil -} - -func VerifyCommandRouting(tenantKey, routingKey string) error { - expected, err := CommandRoutingKey(tenantKey) - if err != nil { - return err - } - if expected != routingKey { - return fmt.Errorf("routing key does not match tenant_key: expected %q got %q", expected, routingKey) - } - return nil -} diff --git a/internal/tenant/routing_test.go b/internal/tenant/routing_test.go deleted file mode 100644 index 12d2e7b..0000000 --- a/internal/tenant/routing_test.go +++ /dev/null @@ -1,31 +0,0 @@ -package tenant - -import "testing" - -func TestCommandRoutingPreservesTenantKey(t *testing.T) { - got, err := CommandRoutingKey("租户-A") - if err != nil { - t.Fatal(err) - } - want := "agent-call.tenant.租户-A.call.execute" - if got != want { - t.Fatalf("routing key = %q, want %q", got, want) - } - if err := VerifyCommandRouting("租户-A", got); err != nil { - t.Fatal(err) - } - queue, err := CommandQueue("租户-A") - if err != nil { - t.Fatal(err) - } - if queue != "agent-call.executor.租户-A.v1" { - t.Fatalf("queue = %q", queue) - } - deadLetterQueue, err := DeadLetterQueue("租户-A") - if err != nil { - t.Fatal(err) - } - if deadLetterQueue != "agent-call.executor.租户-A.dlq.v1" { - t.Fatalf("dead-letter queue = %q", deadLetterQueue) - } -}