From fabfbf71c64041c71e8e6953c91bde70ea0e3625 Mon Sep 17 00:00:00 2001 From: Rogee Date: Thu, 1 Oct 2026 18:08:09 +0800 Subject: [PATCH] Purge stopped task backlog after stopping its consumer --- AGENTS.md | 2 +- contracts/local/manifest.json | 2 +- docs/thirds/saas-dispatcher.md | 4 +- internal/dispatcher/control.go | 14 +++ internal/dispatcher/control_test.go | 49 ++++++++++- internal/dispatcher/runtime.go | 87 +++++++++++++------ .../dispatcher/runtime_integration_test.go | 10 +++ internal/mq/broker.go | 41 ++++----- internal/mq/broker_integration_test.go | 56 ++++++++++++ 9 files changed, 209 insertions(+), 56 deletions(-) diff --git a/AGENTS.md b/AGENTS.md index a2f07b9..176bff6 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -86,7 +86,7 @@ ## SaaS、Dispatcher 与 Agent 的现行边界 - SaaS→Dispatcher 的**五类只读配置**为 `GET /internal/v1/dispatcher/sip`、`/task/:task_id`、`/tasks`、`/tenant/:tenant_id/quota`、`/ai-providers`;路径前缀固定,均须校验 Dispatcher UUID/资源归属、数字 `tenant_id`、完整快照、来源、有效授权和版本。配置读失败、过期、矛盾或不确定时关新准入;没有旧 MQ 配置回退、通用业务 HTTP、ETag 兜底或偷偷启用旧执行字段。已接纳任务持久绑定原快照。 -- 呼叫、控制、必要回执和**每通话唯一最终结果**经固定 `v1` RabbitMQ Topic/队列,任务与控制队列由 SaaS 预建,Dispatcher 仅消费、不可自行建/删/绑定;结果进入指定共享 durable 队列。独立 D UUID 和接收队列不能广播后正文过滤;`tenant_key` 原值保留,任务只能由归属 D 执行。入站先校验和持久 inbox 后 ACK,状态/outbox 同事务;persistent、mandatory、无 return、publisher confirm 成功才记交付,confirm **不是** SaaS 应用收讫。失败/确认丢失与重启只重发同身份消息,不重复拨号或捏造结果。 +- 呼叫、控制、必要回执和**每通话唯一最终结果**经固定 `v1` RabbitMQ Topic/队列,任务与控制队列由 SaaS 预建,Dispatcher 不可自行建/删/绑定;stop 时先停该任务消费者并关闭持久准入,再 purge 仅该任务队列的待投递消息,失败不回成功;SaaS 停止继续投递已 stop 任务,后续误投递不执行。结果进入指定共享 durable 队列。独立 D UUID 和接收队列不能广播后正文过滤;`tenant_key` 原值保留,任务只能由归属 D 执行。入站先校验和持久 inbox 后 ACK,状态/outbox 同事务;persistent、mandatory、无 return、publisher confirm 成功才记交付,confirm **不是** SaaS 应用收讫。失败/确认丢失与重启只重发同身份消息,不重复拨号或捏造结果。 - `call.execute` 只带获批 `task_id/callee`,调用线路、主叫、AI 和时限由该任务快照固定;`task.control` 的 start/pause/resume/stop 无旧 CAS/版本字段,控制回 task/action/status 处理结果,stop 不可恢复。任务启动/重启完整读取列表,运行中不定时轮询任务列表;新建任务由 start 读取任务配置和租户额度,暂停后修改的任务由 resume 重读;全局 AI 服务商列表每次启动只读取一次,本进程全部任务复用,变更须重启后才生效;没有 edit 事件。启动时全量读取并核验 SIP,运行中只由 `sip.config` 通知触发全量 SIP 读取,无常规定时 SIP 轮询;任务 start/resume 使用已生效 SIP 快照,不自行拉取或向 Agent 核验 SIP。SIP 修订变化待旧呼叫结束且 Agent 已加载后,仅在本地重新绑定任务快照,不重读任务列表;待处理期间新准入关闭。发现分页同快照先全部校验后提交,MQ 控制持久状态高于偶发 HTTP 状态;冲突关准入、resume 须重新核验。历史命令不因消息年龄过期,但实际呼出前仍检查任务/白名单/时段/授权/额度和有效通话上限;任务结束不删未确认结果、恢复、幂等或未知占用。 - 独立 Dispatcher 的 SQLite 是任务、额度、inbox/outbox 的权威数据;Agent 无业务数据库,录音、执行与上传恢复只写受控私有文件。额度包含未知占用,新 boot/租约到期不得自动清除未知执行;不实现双活数据库、自动跨机热备、多 D 共享额度或第二租户公平。本轮不借本机 D1/D2 隔离夹具宣称多 D 运行。不得建立旧表/旧消息/旧 HTTP 执行兼容通道。 - Dispatcher↔Agent 复用 Unary gRPC 和受控 Endpoint;Agent 预绑定 D UUID 与服务端证书指纹,激活/会话代际、peer mTLS/SAN/SNI 和已签发期限须核对,新 boot 不清未知占用。Agent 不自行向 SaaS 取任务/AI/OSS 授权;Dispatcher 只用已经核验的 Agent `GetLoadedSIP` revision 开执行准入。本机 Mock 的加载回报不证明 Asterisk 已实际加载,SIP 配置的唯一编辑/审批面仍是 management。 diff --git a/contracts/local/manifest.json b/contracts/local/manifest.json index b2b5ef3..132bd5c 100644 --- a/contracts/local/manifest.json +++ b/contracts/local/manifest.json @@ -4,7 +4,7 @@ "sources": { "docs/archive/sources/v0.5-proposal.md": "612fdaee50aff6aa7fbef16c2d469d99857646c6d2235617d0e67f6098cd7ada", "docs/archive/sources/plan-saas-dispatcher-v05-v0.1.md": "666f39e56ea9f4b55661efcac82edd6f9729848e2d60e5f24cdf5aa3ac97ee87", - "docs/thirds/saas-dispatcher.md": "f74671b0a6e00f6dc6999dc32e71cf231ba83abe1b4e7988f7a8c496fc69fb11" + "docs/thirds/saas-dispatcher.md": "973ef69b1f71747d3e90f9798a7bf75d2fe2aac7a56d91dad97bbb8d4b9c566e" }, "bundle_sha256": "d89302b79228c884ce5e0c413bbe9c688c725724de73cb84892e1d38d638e2a5", "bundle_algorithm": "sha256 of sorted relative-path + space + sha256(file) + newline; only root-level JSON and examples/**/*.json, excluding manifest.json" diff --git a/docs/thirds/saas-dispatcher.md b/docs/thirds/saas-dispatcher.md index 062de1b..87820c1 100644 --- a/docs/thirds/saas-dispatcher.md +++ b/docs/thirds/saas-dispatcher.md @@ -18,14 +18,14 @@ HTTP 非 200、正文不合法、缺字段、归属冲突、缓存失效或读 ## MQ:固定 v1 传输,唯一事件清单 -RabbitMQ 是 Topic,**SaaS 独占创建、绑定、退役 exchange/queue,D 只消费/发布,无 configure 权限**。机器拓扑见 [`mq-topology.json`](../../contracts/local/mq-topology.json)。外呼维持每 D、每任务独立的 `d..task..in` 路由及 `agent-call.d..task..v1` 队列;每 D 单独控制路由 `d..control.in` 及队列 `agent-call.d..control.v1`。SaaS 将所有 D 的输出 `d..out` **精确绑定至同一个结果队列**(本地 Mock 名 `agent-call.saas.events.v1`,不是强制 SaaS 实际队列名)。入站 exchange `agent-call.dispatchers.v1`;出站 exchange `agent-call.saas.v1`;死信 exchange `agent-call.dead-letter.v1`。不广播后靠正文过滤,不使用独立通配段,不把被动声明/发布 confirm 冒充指定队列已收到:须以实际路由、mandatory return 与 confirm 联合验证。 +RabbitMQ 是 Topic,**SaaS 独占创建、绑定、退役 exchange/queue,D 只消费/发布,并仅在 stop 时对自己负责的任务队列执行 purge;不创建、删除或绑定队列,无 configure 权限**。机器拓扑见 [`mq-topology.json`](../../contracts/local/mq-topology.json)。外呼维持每 D、每任务独立的 `d..task..in` 路由及 `agent-call.d..task..v1` 队列;每 D 单独控制路由 `d..control.in` 及队列 `agent-call.d..control.v1`。SaaS 将所有 D 的输出 `d..out` **精确绑定至同一个结果队列**(本地 Mock 名 `agent-call.saas.events.v1`,不是强制 SaaS 实际队列名)。入站 exchange `agent-call.dispatchers.v1`;出站 exchange `agent-call.saas.v1`;死信 exchange `agent-call.dead-letter.v1`。不广播后靠正文过滤,不使用独立通配段,不把被动声明/发布 confirm 冒充指定队列已收到:须以实际路由、mandatory return 与 confirm 联合验证。 `event_id` 是入站/回执关联和内部防重复处理的消息身份;`dispatcher_id` 是 D 身份,`tenant_id` 是正整数。**不是所有事件共用一个统一必填信封**:`sip.config` 没有 tenant_id/issued_at;回执按其示例字段,最终结果使用 `issued_at` 而非 `occurred_at`。正文的 `schema_version` 已移除。只接受以下外发事件及必要入站事件;所有字段直接按 [`mq.schema.json`](../../contracts/local/mq.schema.json) 和正例校验: | 事件 | 方向和处理规则 | 正例 | | --- | --- | --- | | `sip.config` | SaaS→D;`revision` 触发全量重新读取和实际加载核验,不用通知正文代替全量 | [`notification`](../../contracts/local/examples/mq-sip-change.json) | -| `task.control` | SaaS→D;start/pause/resume/stop,无控制去重/CAS 字段。start 读取新任务配置及租户配额并持久绑定,无 Agent 控制动作;resume 重读最新任务配置和额度后解除人工暂停;均沿用启动时 AI 服务商列表及已生效 SIP 快照,不重复拉取全局服务商或 SIP,也不额外核验 SIP;pause/stop 发给 Agent,省略 `active_call_policy` 默认 **hangup**,显式仅 drain/hangup。配置读取、归属或 AI 核验失败时不接纳,SIP 待处理时新呼叫仍关闭;成功操作后才回应用回执,回执不是已完成排空/挂断的证据;stop 同 ID 不可恢复 | [`start`](../../contracts/local/examples/mq-control-start.json) · [`control`](../../contracts/local/examples/mq-control.json) · [`ack`](../../contracts/local/examples/mq-control-ack.json) | +| `task.control` | SaaS→D;start/pause/resume/stop,无控制去重/CAS 字段。start 读取新任务配置及租户配额并持久绑定,无 Agent 控制动作;resume 重读最新任务配置和额度后解除人工暂停;均沿用启动时 AI 服务商列表及已生效 SIP 快照,不重复拉取全局服务商或 SIP,也不额外核验 SIP;pause/stop 发给 Agent,省略 `active_call_policy` 默认 **hangup**,显式仅 drain/hangup。配置读取、归属或 AI 核验失败时不接纳,SIP 待处理时新呼叫仍关闭;stop 先停止该任务消费者(使在途未确认消息回队列),持久关闭任务准入并抑制本地待执行指令,再对该任务队列执行 purge 清除当时待投递消息;purge 失败不发送成功回执、保持准入关闭并重试原控制消息。SaaS 须停止继续投递已 stop 的任务;purge 不阻止之后新投递,这些消息不执行。成功操作后才回应用回执,回执不是已完成活跃通话排空/挂断、SaaS 已停止投递或未来消息不存在的证据;stop 同 ID 不可恢复 | [`start`](../../contracts/local/examples/mq-control-start.json) · [`control`](../../contracts/local/examples/mq-control.json) · [`ack`](../../contracts/local/examples/mq-control-ack.json) | | `call.execute` | SaaS→D 仅 `{task_id,callee}`;一次指令保留独立消息/执行身份,路由/主叫/AI/时限从绑定任务读取。`dispatched` 表示已发出呼叫指令,**不表示接通**;白名单/单号码格式不合规则回 `rejected,reason_code:null,reason_message`、不拨号不发最终结果、不暂停整任务 | [`execute`](../../contracts/local/examples/mq-execute.json) · [`dispatched`](../../contracts/local/examples/mq-execute-ack.json) · [`rejected`](../../contracts/local/examples/mq-execute-rejected.json) | | `call.execute.result` | D→SaaS;按 `task_id` + 原号码关联,每次呼叫仅一份最终结果;不新增外部 call_id/source_command_id;真正终结且录音成功上传、无录音或预期录音生成失败后才发送 | [`uploaded`](../../contracts/local/examples/mq-result-uploaded.json) · [`empty`](../../contracts/local/examples/mq-result-no-recording.json) | diff --git a/internal/dispatcher/control.go b/internal/dispatcher/control.go index 09bef10..af443f0 100644 --- a/internal/dispatcher/control.go +++ b/internal/dispatcher/control.go @@ -5,6 +5,7 @@ import ( "encoding/json" "errors" "fmt" + "log/slog" "time" "git.ipao.vip/rogee/go-sip/internal/configread" @@ -35,6 +36,7 @@ type ControlController struct { ApprovedSIP *configread.SIP ApprovedProviders *map[string]configread.Provider Now func() time.Time + PurgeTaskQueue func(context.Context, string) (int, error) } // ProcessControl applies each delivered control independently: it has no @@ -125,6 +127,18 @@ func (c *ControlController) ProcessControl(ctx context.Context, body []byte) err } return fmt.Errorf("prepare task control %q: %w", event.EventID, err) } + if event.Payload.Action == "stop" { + if c.PurgeTaskQueue == nil { + return errors.New("stop requires a task queue purger") + } + count, err := c.PurgeTaskQueue(ctx, event.Payload.TaskID) + if err != nil { + return fmt.Errorf("purge stopped task %q backlog: %w", event.Payload.TaskID, err) + } + // A purge covers Ready messages only. The durable stop barrier also + // rejects any later or already accepted deliveries. + slog.Info("stopped task backlog purged", "dispatcher_id", event.DispatcherID, "task_id", event.Payload.TaskID, "ready_messages", count) + } spec := ControlSpec{ DispatcherID: event.DispatcherID, TenantID: event.TenantID, TaskID: event.Payload.TaskID, Action: event.Payload.Action, diff --git a/internal/dispatcher/control_test.go b/internal/dispatcher/control_test.go index 81c635d..c0f1586 100644 --- a/internal/dispatcher/control_test.go +++ b/internal/dispatcher/control_test.go @@ -82,7 +82,7 @@ func newControlFixture(t *testing.T) (*ControlController, *fakeControlAgent, *st t.Fatal(err) } agent := &fakeControlAgent{} - controller := &ControlController{DispatcherID: id, Store: s, Client: client, Agent: agent, ApprovedSIP: &snapshot.SIP, ApprovedProviders: &providers, Now: func() time.Time { return monday(9, 30) }} + controller := &ControlController{DispatcherID: id, Store: s, Client: client, Agent: agent, ApprovedSIP: &snapshot.SIP, ApprovedProviders: &providers, Now: func() time.Time { return monday(9, 30) }, PurgeTaskQueue: func(context.Context, string) (int, error) { return 0, nil }} return controller, agent, s } @@ -162,6 +162,53 @@ func TestControlStopCannotResumeOrDispatchPending(t *testing.T) { } } +func TestControlStopPurgesAfterBarrierBeforeAgentAndAck(t *testing.T) { + controller, agent, s := newControlFixture(t) + controller.PurgeTaskQueue = func(_ context.Context, taskID string) (int, error) { + if taskID != "task-asr" { + t.Fatalf("purged wrong task: %s", taskID) + } + if admitted, err := s.CanAdmit(controller.DispatcherID, 1001, taskID); err != nil || admitted { + t.Fatalf("stop barrier missing at purge: %v %v", admitted, err) + } + if len(agent.calls) != 0 { + t.Fatal("Agent called before purge") + } + if pending, err := s.ListPendingOutbox(controller.DispatcherID); err != nil || len(pending) != 0 { + t.Fatalf("ack before purge: %v %v", pending, err) + } + return 42, nil + } + if err := controller.ProcessControl(context.Background(), controlBody(t, "stop-purge", "stop", "")); err != nil { + t.Fatal(err) + } + if len(agent.calls) != 1 { + t.Fatalf("Agent not called after purge: %v", agent.calls) + } +} + +func TestControlStopPurgeFailureDoesNotAckOrCallAgent(t *testing.T) { + controller, agent, s := newControlFixture(t) + controller.PurgeTaskQueue = func(context.Context, string) (int, error) { return 0, errors.New("rabbitmq purge failed") } + body := controlBody(t, "stop-fail", "stop", "") + if err := controller.ProcessControl(context.Background(), body); err == nil || !strings.Contains(err.Error(), "rabbitmq purge failed") { + t.Fatalf("purge failure hidden: %v", err) + } + if len(agent.calls) != 0 { + t.Fatalf("Agent called after purge failure: %v", agent.calls) + } + if pending, err := s.ListPendingOutbox(controller.DispatcherID); err != nil || len(pending) != 0 { + t.Fatalf("stop falsely acknowledged: %v %v", pending, err) + } + if admitted, err := s.CanAdmit(controller.DispatcherID, 1001, "task-asr"); err != nil || admitted { + t.Fatalf("failed purge reopened admission: %v %v", admitted, err) + } + controller.PurgeTaskQueue = func(context.Context, string) (int, error) { return 0, nil } + if err := controller.ProcessControl(context.Background(), body); err != nil { + t.Fatalf("redelivered stop did not recover: %v", err) + } +} + func TestStoppedTaskSilentlyAcknowledgesOldExecuteWithoutRecreatingInbox(t *testing.T) { control, _, s := newControlFixture(t) if err := control.ProcessControl(context.Background(), controlBody(t, "stop-before-execute", "stop", "")); err != nil { diff --git a/internal/dispatcher/runtime.go b/internal/dispatcher/runtime.go index d7cc762..9c2cec1 100644 --- a/internal/dispatcher/runtime.go +++ b/internal/dispatcher/runtime.go @@ -16,7 +16,8 @@ import ( ) // Runtime owns task/control consumers for one Dispatcher. SaaS creates -// every queue/binding; this process only checks and consumes predeclared ones. +// every queue/binding; this process checks and consumes predeclared ones and +// purges a stopped task's Ready backlog without changing queue topology. type Runtime struct { Broker *mq.Broker Bootstrap Bootstrap @@ -25,15 +26,17 @@ type Runtime struct { PollInterval time.Duration Logger *slog.Logger - gate sync.RWMutex // SIP admission barrier versus each pre-dial instruction - locksMu sync.Mutex - taskLocks map[string]*sync.Mutex - failures chan error - sipUpdates chan struct{} + gate sync.RWMutex // SIP admission barrier versus each pre-dial instruction + consumersMu sync.Mutex // stop consumer and purge cannot race consumer sync + taskConsumers map[string]*mq.Consumer + locksMu sync.Mutex + taskLocks map[string]*sync.Mutex + failures chan error + sipUpdates chan struct{} } // Serve closes admission on every shutdown/failure; it never clears durable -// calls, results, task queues, or the SQLite file. +// calls, results, or the SQLite file. Stop purges only the task's Ready backlog. func (r *Runtime) Serve(ctx context.Context) (result error) { if r == nil || r.Broker == nil || r.Bootstrap.Client == nil || r.Bootstrap.Store == nil || r.Bootstrap.VerifySIP == nil || r.Bootstrap.DispatcherID == "" || r.Execute.Store != r.Bootstrap.Store || r.Control.Store != r.Bootstrap.Store || r.Execute.DispatcherID != r.Bootstrap.DispatcherID || r.Control.DispatcherID != r.Bootstrap.DispatcherID || r.PollInterval <= 0 || r.Logger == nil || r.Bootstrap.DrainControls != nil { return errors.New("runtime requires one Dispatcher, durable state, verified SIP and configured pending-work polling; external control drain is forbidden") @@ -48,20 +51,24 @@ func (r *Runtime) Serve(ctx context.Context) (result error) { r.sipUpdates = make(chan struct{}, 1) r.taskLocks = make(map[string]*sync.Mutex) consumers := make(map[string]*mq.Consumer) + r.taskConsumers = consumers + r.Control.PurgeTaskQueue = r.Broker.PurgeTaskQueue var controlConsumer *mq.Consumer defer func() { stopCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second) defer cancel() - for queue, consumer := range consumers { - if err := consumer.Stop(stopCtx); err != nil { - result = errors.Join(result, fmt.Errorf("stop task consumer %q: %w", queue, err)) - } - } if controlConsumer != nil { if err := controlConsumer.Stop(stopCtx); err != nil { result = errors.Join(result, fmt.Errorf("stop control consumer: %w", err)) } } + r.consumersMu.Lock() + for queue, consumer := range consumers { + if err := consumer.Stop(stopCtx); err != nil { + result = errors.Join(result, fmt.Errorf("stop task consumer %q: %w", queue, err)) + } + } + r.consumersMu.Unlock() if err := r.Bootstrap.Store.CloseAdmission(r.Bootstrap.DispatcherID); err != nil { result = errors.Join(result, fmt.Errorf("close Dispatcher admission: %w", err)) } @@ -173,6 +180,7 @@ func (r *Runtime) handleControl(ctx context.Context, _ string, body []byte) erro EventType string `json:"event_type"` Payload struct { TaskID string `json:"task_id"` + Action string `json:"action"` Revision int64 `json:"revision"` } `json:"payload"` } @@ -183,6 +191,29 @@ func (r *Runtime) handleControl(ctx context.Context, _ string, body []byte) erro } switch event.EventType { case "task.control": + if event.Payload.Action == "stop" { + // Stop the consumer before purging: closing its channel requeues any + // unacknowledged deliveries so the purge covers them as Ready. + // Hold this lock through the control transition to prevent sync from + // restarting the consumer before the stop barrier is persisted. + r.consumersMu.Lock() + defer r.consumersMu.Unlock() + route, err := tenant.TaskRoute(r.Bootstrap.DispatcherID, event.Payload.TaskID) + if err != nil { + r.signalFailure(err) + return err + } + if consumer := r.taskConsumers[route.Queue]; consumer != nil { + stopCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + err := consumer.Stop(stopCtx) + cancel() + if err != nil { + r.signalFailure(err) + return fmt.Errorf("stop task consumer before purge %q: %w", route.Queue, err) + } + delete(r.taskConsumers, route.Queue) + } + } return r.withTask(event.Payload.TaskID, func() error { if err := r.Control.ProcessControl(ctx, body); err != nil { r.Logger.Error("task control failed", "dispatcher_id", r.Bootstrap.DispatcherID, "task_id", event.Payload.TaskID, "error", err) @@ -241,6 +272,8 @@ func (r *Runtime) processPending(ctx context.Context) error { } func (r *Runtime) syncTaskConsumers(ctx context.Context, consumers map[string]*mq.Consumer) error { + r.consumersMu.Lock() + defer r.consumersMu.Unlock() assigned, err := r.Bootstrap.Store.ListAssignedTasks(r.Bootstrap.DispatcherID) if err != nil { return err @@ -251,24 +284,22 @@ func (r *Runtime) syncTaskConsumers(ctx context.Context, consumers map[string]*m if err != nil { return err } - if task.ControlState == "paused" || task.ControlState == "pausing" || task.ControlState == "resuming" || task.Status == "paused" { + if task.ControlState == "paused" || task.ControlState == "pausing" || task.ControlState == "resuming" || task.Status == "paused" || task.ControlState == "stopping" || task.ControlState == "stopped" || task.Status == "stopped" { continue } - if task.ControlState != "stopped" && task.ControlState != "stopping" && task.Status != "stopped" { - admitted, err := r.Bootstrap.Store.CanAdmit(r.Bootstrap.DispatcherID, task.TenantID, task.TaskID) - if err != nil { - return err - } - if !admitted { - continue - } - waiting, err := r.Bootstrap.Store.PendingExecuteCount(r.Bootstrap.DispatcherID, task.TenantID, task.TaskID) - if err != nil { - return err - } - if waiting != 0 { - continue - } + admitted, err := r.Bootstrap.Store.CanAdmit(r.Bootstrap.DispatcherID, task.TenantID, task.TaskID) + if err != nil { + return err + } + if !admitted { + continue + } + waiting, err := r.Bootstrap.Store.PendingExecuteCount(r.Bootstrap.DispatcherID, task.TenantID, task.TaskID) + if err != nil { + return err + } + if waiting != 0 { + continue } wanted[route.Queue] = task } diff --git a/internal/dispatcher/runtime_integration_test.go b/internal/dispatcher/runtime_integration_test.go index 63affc4..4bea1d1 100644 --- a/internal/dispatcher/runtime_integration_test.go +++ b/internal/dispatcher/runtime_integration_test.go @@ -366,6 +366,12 @@ func TestRuntimeIsolatedControlBacklogExecuteAndSharedResult(t *testing.T) { t.Fatalf("pending SIP change originated call %s", spec.EventID) case <-time.After(120 * time.Millisecond): } + for _, eventID := range []string{"stop-backlog-1", "stop-backlog-2"} { + publish(taskRoute, executeBody(t, eventID, "15003164745")) + } + if state, err := admin.QueueInspect(taskRoute.Queue); err != nil || state.Messages < 3 { + t.Fatalf("expected isolated stop backlog: %+v %v", state, err) + } publish(controlRoute, controlBody(t, "stop-after-sip", "stop", "")) select { case spec := <-agent.controls: @@ -408,6 +414,10 @@ func TestRuntimeIsolatedControlBacklogExecuteAndSharedResult(t *testing.T) { if ack.Payload.Status != "applied" { t.Fatalf("stop control falsely acknowledged during SIP reload: %+v", ack) } + state, err := admin.QueueInspect(taskRoute.Queue) + if err != nil || state.Messages != 0 || state.Consumers != 0 { + t.Fatalf("stopped task backlog was not purged: %+v %v", state, err) + } break } } diff --git a/internal/mq/broker.go b/internal/mq/broker.go index 609dcdc..1187141 100644 --- a/internal/mq/broker.go +++ b/internal/mq/broker.go @@ -275,8 +275,12 @@ func (b *Broker) StartPredeclaredConsumer(ctx context.Context, queue string, han return consumer, nil } -func (b *Broker) DrainPredeclared(ctx context.Context, queue string) (int, error) { - if err := b.validateQueue(queue, false); err != nil { +// PurgeTaskQueue removes Ready messages from only this Dispatcher's task queue. +// The caller must first stop its consumer so unacknowledged deliveries are +// requeued before the purge; SaaS still owns queue creation and bindings. +func (b *Broker) PurgeTaskQueue(ctx context.Context, taskID string) (int, error) { + route, err := tenant.TaskRoute(b.dispatcherID, taskID) + if err != nil { return 0, err } if err := ctx.Err(); err != nil { @@ -291,34 +295,25 @@ func (b *Broker) DrainPredeclared(ctx context.Context, queue string) (int, error } channel, err := conn.Channel() if err != nil { - return 0, fmt.Errorf("open current drain channel: %w", err) + return 0, fmt.Errorf("open current purge 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) + if _, err := channel.QueueDeclarePassive(route.Queue, true, false, false, false, nil); err != nil { + return 0, fmt.Errorf("required SaaS-owned queue %s unavailable: %w", route.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++ + if err := ctx.Err(); err != nil { + return 0, err } + count, err := channel.QueuePurge(route.Queue, false) + if err != nil { + return 0, fmt.Errorf("purge SaaS-owned task queue %s: %w", route.Queue, err) + } + return count, nil } // 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. +// Control deliveries must pass through the handler before ACK; a transient +// failure is requeued and closes startup admission. 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") diff --git a/internal/mq/broker_integration_test.go b/internal/mq/broker_integration_test.go index 337542f..7bfa566 100644 --- a/internal/mq/broker_integration_test.go +++ b/internal/mq/broker_integration_test.go @@ -149,6 +149,62 @@ func TestBrokerSharedResultQueueAndNoConfigure(t *testing.T) { t.Fatalf("D2 stole D1 command: %+v %v", state, err) } + // Purge is permitted with read access and affects only the chosen task; + // it must not remove another Dispatcher's backlog or the control queue. + for i := 0; i < 3; i++ { + 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) + } + } + if err := admin.PublishWithContext(context.Background(), CommandsExchange, otherTask.BindingKey, true, false, amqp.Publishing{ContentType: "application/json", DeliveryMode: amqp.Persistent, Body: body}); err != nil { + t.Fatal(err) + } + count, err := brokers[0].PurgeTaskQueue(context.Background(), "task-asr") + if err != nil || count != 3 { + t.Fatalf("task-only purge: removed=%d err=%v", count, err) + } + if state, err := admin.QueueInspect(task.Queue); err != nil || state.Messages != 0 { + t.Fatalf("D1 queue not empty: %+v %v", state, err) + } + if state, err := admin.QueueInspect(otherTask.Queue); err != nil || state.Messages != 1 { + t.Fatalf("D2 queue was purged: %+v %v", state, err) + } + if _, err := brokers[0].PurgeTaskQueue(context.Background(), "../foreign"); err == nil { + t.Fatal("invalid task ID accepted for purge") + } + + // Stop closes the consumer channel before purge; an in-flight delivery + // must become Ready rather than surviving the purge as Unacked. + inFlight := make(chan struct{}, 1) + consumer, err := brokers[0].StartPredeclaredConsumer(context.Background(), task.Queue, func(ctx context.Context, _ string, _ []byte) error { + select { + case inFlight <- struct{}{}: + default: + } + <-ctx.Done() + return ctx.Err() + }) + if err != nil { + t.Fatal(err) + } + if err := admin.PublishWithContext(context.Background(), CommandsExchange, task.BindingKey, true, false, amqp.Publishing{ContentType: "application/json", DeliveryMode: amqp.Persistent, Body: currentMessage(t, "mq-execute")}); err != nil { + t.Fatal(err) + } + select { + case <-inFlight: + case <-time.After(3 * time.Second): + t.Fatal("task consumer did not receive in-flight delivery") + } + stopCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + if err := consumer.Stop(stopCtx); err != nil { + t.Fatal(err) + } + count, err = brokers[0].PurgeTaskQueue(context.Background(), "task-asr") + if err != nil || count != 1 { + t.Fatalf("in-flight delivery survived stop and purge: removed=%d err=%v", count, err) + } + resultBody := currentMessage(t, "mq-result-no-recording") for i, broker := range brokers { message := resultBody