From d0113bec47d1325b7386ba6cd9d7443d0c46193f Mon Sep 17 00:00:00 2001 From: Rogee Date: Thu, 1 Oct 2026 12:19:40 +0800 Subject: [PATCH] Trigger task configuration refresh on start and resume controls --- AGENTS.md | 2 +- cmd/sip-go-agent/dispatcher_command.go | 2 +- .../local/examples/mq-control-start.json | 1 + contracts/local/manifest.json | 4 +- contracts/local/mq.schema.json | 2 +- .../saas-dispatcher-implementation.md | 2 + .../saas-dispatcher-p08-acceptance.md | 2 + docs/thirds/saas-dispatcher.md | 4 +- internal/dispatcher/control.go | 27 +++-- internal/dispatcher/control_test.go | 104 ++++++++++++++++++ internal/dispatcher/discovery.go | 82 -------------- internal/dispatcher/discovery_test.go | 87 --------------- internal/dispatcher/runtime.go | 36 ++---- .../dispatcher/runtime_integration_test.go | 32 +++++- internal/dispatcher/sip_reload.go | 13 ++- internal/dispatcher/sip_runtime_test.go | 10 +- internal/store/control_flow.go | 34 ++++++ 17 files changed, 222 insertions(+), 222 deletions(-) create mode 100644 contracts/local/examples/mq-control-start.json delete mode 100644 internal/dispatcher/discovery.go delete mode 100644 internal/dispatcher/discovery_test.go diff --git a/AGENTS.md b/AGENTS.md index d25ef36..c797c12 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -87,7 +87,7 @@ - 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 应用收讫。失败/确认丢失与重启只重发同身份消息,不重复拨号或捏造结果。 -- `call.execute` 只带获批 `task_id/callee`,调用线路、主叫、AI 和时限由该任务快照固定;`task.control` 的 pause/resume/stop 无旧 CAS/版本字段,控制回 task/action/status 处理结果,stop 不可恢复。发现分页同快照先全部校验后提交,MQ 控制持久状态高于偶发 HTTP 状态;冲突关准入、resume 须重新核验。历史命令不因消息年龄过期,但实际呼出前仍检查任务/白名单/时段/授权/额度和有效通话上限;任务结束不删未确认结果、恢复、幂等或未知占用。 +- `call.execute` 只带获批 `task_id/callee`,调用线路、主叫、AI 和时限由该任务快照固定;`task.control` 的 start/pause/resume/stop 无旧 CAS/版本字段,控制回 task/action/status 处理结果,stop 不可恢复。任务启动/重启完整读取列表,运行中不定时轮询任务列表;新建任务由 start 读取配置和租户额度,暂停后修改的任务由 resume 重读;没有 edit 事件。SIP 修订变化时仍全量核验任务绑定。发现分页同快照先全部校验后提交,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。 - AI 使用任务内不可变授权快照:仅经获批准百炼/火山 ASR、OpenAI 兼容 LLM、火山 TTS 能表达的参数进入每通话实例;ASR-only 不启动 LLM/TTS,完整 AI 不借旧语音测试的 LLM/TTS。只有最终用户 ASR 文本的明确字面关键词可触发拒联/挂断;不由 SDK 默认值、环境、CLI、metadata 或宽松 Schema 改写业务参数,不因 SDK 重试产生第二次发起/收费或重播。日志只存脱敏版本/摘要/计数,不存密钥、prompt、完整对话或音频。 diff --git a/cmd/sip-go-agent/dispatcher_command.go b/cmd/sip-go-agent/dispatcher_command.go index 2c2e654..295de50 100644 --- a/cmd/sip-go-agent/dispatcher_command.go +++ b/cmd/sip-go-agent/dispatcher_command.go @@ -140,7 +140,7 @@ func runDispatcher(ctx context.Context, mode string) (result error) { Control: dispatcher.ControlController{ DispatcherID: settings.DispatcherID, Store: database, Client: reader, Agent: originator, VerifySIP: originator.VerifySIP, Now: time.Now, }, - PollInterval: time.Second, DiscoveryInterval: time.Minute, Logger: slog.Default(), + PollInterval: time.Second, SIPInterval: time.Minute, Logger: slog.Default(), } listener, err := net.Listen("tcp", settings.Listen) if err != nil { diff --git a/contracts/local/examples/mq-control-start.json b/contracts/local/examples/mq-control-start.json new file mode 100644 index 0000000..280fa9a --- /dev/null +++ b/contracts/local/examples/mq-control-start.json @@ -0,0 +1 @@ +{"event_id":"control-start-example","event_type":"task.control","dispatcher_id":"c046b893-8628-4589-ae50-619d049248a6","tenant_id":1001,"issued_at":"2026-09-21T01:00:00Z","payload":{"task_id":"task-asr","action":"start","reason":"created"}} diff --git a/contracts/local/manifest.json b/contracts/local/manifest.json index 487be3e..888fa43 100644 --- a/contracts/local/manifest.json +++ b/contracts/local/manifest.json @@ -4,8 +4,8 @@ "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": "19102a481aae33320aa3f6df619d3ea39af6d41609a34d56b21ca078e735e393" + "docs/thirds/saas-dispatcher.md": "a337aac3bc34aca1c9fd4a63a26938b1aa3885cf1c6228cf8af629c421f3899f" }, - "bundle_sha256": "4ff0afffa2c865050091c042d8f98bbe344ba9a4f4c3652e721a5077217722ed", + "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/contracts/local/mq.schema.json b/contracts/local/mq.schema.json index 0151c9a..1a47293 100644 --- a/contracts/local/mq.schema.json +++ b/contracts/local/mq.schema.json @@ -34,7 +34,7 @@ "properties": { "event_id": {"$ref": "#/$defs/event_id"}, "event_type": {"const": "task.control"}, "dispatcher_id": {"$ref": "#/$defs/dispatcher_id"}, "tenant_id": {"$ref": "#/$defs/tenant_id"}, "issued_at": {"$ref": "#/$defs/issued_at"}, "payload": {"type": "object", "additionalProperties": false, "required": ["task_id", "action", "reason"], "properties": { - "task_id": {"$ref": "#/$defs/task_id"}, "action": {"enum": ["pause", "resume", "stop"]}, "reason": {"type": "string", "minLength": 1}, + "task_id": {"$ref": "#/$defs/task_id"}, "action": {"enum": ["start", "pause", "resume", "stop"]}, "reason": {"type": "string", "minLength": 1}, "options": {"type": "object", "additionalProperties": false, "properties": {"active_call_policy": {"enum": ["drain", "hangup"]}}} }} } diff --git a/docs/evidence/saas-dispatcher-implementation.md b/docs/evidence/saas-dispatcher-implementation.md index 2fb6205..67f8f9a 100644 --- a/docs/evidence/saas-dispatcher-implementation.md +++ b/docs/evidence/saas-dispatcher-implementation.md @@ -1,5 +1,7 @@ # SaaS↔Dispatcher 项目内实施证据 +> 本文记录当时的实施过程。当前运行期间不再定时轮询任务列表,也不再有 `DiscoveryFollower`;新建任务用 `task.control` 的 `start`、暂停后修改用 `resume` 重读任务配置。现行规则以 [`docs/thirds/saas-dispatcher.md`](../thirds/saas-dispatcher.md) 为准。 + ## 改动前基线 - 基线提交:`f5c2d6a92036a579e0b070beb1d9381b0976f81e`;执行分支:`feat/saas-dispatcher-contract`;开始前已拉取并核对 `origin/main`,无待同步提交;Go 1.27.1。 diff --git a/docs/evidence/saas-dispatcher-p08-acceptance.md b/docs/evidence/saas-dispatcher-p08-acceptance.md index 27cf77d..f4dcdf3 100644 --- a/docs/evidence/saas-dispatcher-p08-acceptance.md +++ b/docs/evidence/saas-dispatcher-p08-acceptance.md @@ -2,6 +2,8 @@ > 本表保留 P08 当时的本地测试结果、来源路径与哈希;后续整理未改写验收事实。原计划现位于 `docs/archive/sources/plan-saas-dispatcher-v05-v0.1.md`,完成阶段计划位于 `docs/archive/plan-saas-dispatcher-completed.md`;**当前唯一规范**见 [`../thirds/saas-dispatcher.md`](../thirds/saas-dispatcher.md)。 +**当前补充(不改写 P08 历史结果):**现行运行期间不再定时轮询任务列表,删除了 `DiscoveryFollower` 及其两个测试;新建任务的 `start` 和暂停后修改的 `resume` 均重读任务配置。`TestStartReadsNewTaskWithoutAgentControl`、`TestResumeRefreshesModifiedTaskRevision`、`TestStartDoesNotReopenPausedOrStoppedTasks` 与隔离 RabbitMQ 的 `TestRuntimeIsolatedControlBacklogExecuteAndSharedResult` 已通过;本次 `PATH=/tmp/sip-go-agent-tools/bin:$PATH make check`(含 Proto 重新生成、隔离 RabbitMQ 测试)、`go test -race ./...`、`go vet ./...` 和构建通过,业务覆盖率 72.2%;含隔离 RabbitMQ 的 Dispatcher 测试覆盖率 75.0%。 + 本记录按 `docs/plan-saas-dispatcher-v05-v0.1.md` §3.1、§6 核对 P01–P08、K01–K16 和 A01–A12。**范围仅限单节点、单 Dispatcher、单 Agent、单 Cell、单租户的隔离 Mock。** 本机证书、RabbitMQ 容器和 OSS/AI 模拟服务不代表 SaaS 应用收讫、供应商验收或真实拨号。原计划是当时的来源事实,不是第二份运行合同;当时的状态入口 `docs/plan-saas-dispatcher.md` 已归档。 ## 可复现的本地门禁与来源 diff --git a/docs/thirds/saas-dispatcher.md b/docs/thirds/saas-dispatcher.md index 565c370..69499bb 100644 --- a/docs/thirds/saas-dispatcher.md +++ b/docs/thirds/saas-dispatcher.md @@ -14,7 +14,7 @@ | 任务列表 | `GET /internal/v1/dispatcher/tasks`,后续 `?after=` | 首次/重启完整读到 **带 cursor 且 tasks=[]** 的终止页;非空短页不能提前结束。非空页持久成功后才使用下一 cursor;控制队列积压处理前不开新任务准入。不使用旧 snapshot_id/watermark/mode=snapshot | [`page`](../../contracts/local/examples/task-discovery-page.json) · [`end`](../../contracts/local/examples/task-discovery-end.json) | | 租户额度 | `GET /internal/v1/dispatcher/tenant/{tenant_id}/quota` | `quota_revision` 为业务修订号;未知占用不得算成已释放 | [`quota`](../../contracts/local/examples/config-read-quota.json) | -HTTP 非 200、正文不合法、缺字段、归属冲突、缓存失效或读取失败时拒绝新的相关准入并记录脱敏错误;已持久接纳的执行继续使用原快照。错误体示例 [`error`](../../contracts/local/examples/config-read-error.json),不能当成功配置解析。任务仅有运行/暂停/终止;终止是不可逆 stop,同一 `task_id` 不再启用。HTTP 旧 running 不得覆盖已持久的 pause/stop;resume 经最新任务配置核验后只能解除人工暂停。 +HTTP 非 200、正文不合法、缺字段、归属冲突、缓存失效或读取失败时拒绝新的相关准入并记录脱敏错误;已持久接纳的执行继续使用原快照。错误体示例 [`error`](../../contracts/local/examples/config-read-error.json),不能当成功配置解析。任务仅有运行/暂停/终止;终止是不可逆 stop,同一 `task_id` 不再启用。任务配置只能在暂停或停止后修改;没有 `edit` 事件,也不定时刷新任务列表。启动/重启时完整读取任务列表;运行中新建任务由 `task.control` 的 `start` 触发读取,人工暂停后修改配置由 `resume` 触发重读(含租户配额)。HTTP 旧 running 不得覆盖已持久的 pause/stop;resume 经最新任务配置核验后只能解除人工暂停。SIP 修订变更时仍需完整读取任务列表以验证所有任务绑定。 ## MQ:固定 v1 传输,唯一事件清单 @@ -25,7 +25,7 @@ RabbitMQ 是 Topic,**SaaS 独占创建、绑定、退役 exchange/queue,D | 事件 | 方向和处理规则 | 正例 | | --- | --- | --- | | `sip.config` | SaaS→D;`revision` 触发全量重新读取和实际加载核验,不用通知正文代替全量 | [`notification`](../../contracts/local/examples/mq-sip-change.json) | -| `task.control` | SaaS→D;pause/resume/stop,无控制去重/CAS 字段;省略 `active_call_policy` 默认 **hangup**,显式仅 drain/hangup。成功将操作发给 Agent 后才回应用回执,回执不是已完成排空/挂断的证据;stop 同 ID 不可恢复 | [`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 重读最新配置和额度后解除人工暂停;pause/stop 发给 Agent,省略 `active_call_policy` 默认 **hangup**,显式仅 drain/hangup。配置读取、归属、SIP 或 AI 核验失败时不接纳;成功操作后才回应用回执,回执不是已完成排空/挂断的证据;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 ed3a5f9..b5bb107 100644 --- a/internal/dispatcher/control.go +++ b/internal/dispatcher/control.go @@ -76,28 +76,39 @@ func (c *ControlController) ProcessControl(ctx context.Context, body []byte) err return fmt.Errorf("control event %q issued_at is in the future; do not apply early", event.EventID) } policy := event.Payload.Options.ActiveCallPolicy - if event.Payload.Action == "resume" { + if event.Payload.Action == "start" || event.Payload.Action == "resume" { if policy != "" { - return errors.New("resume control must not carry an active-call policy") + return errors.New("start/resume control must not carry an active-call policy") } - // A stale HTTP running status cannot override persisted pause/stop. // Only a fresh approved task plus applied SIP snapshot may authorize - // the transition from a durable pause back to admission. + // new admission; the durable pause/stop barrier remains authoritative. snapshot, err := c.Client.ReadTask(ctx, event.Payload.TaskID, event.TenantID) if err != nil { - return fmt.Errorf("fresh resume task configuration: %w", err) + return fmt.Errorf("fresh %s task configuration: %w", event.Payload.Action, err) } if snapshot.Task.Status != "running" { - return errors.New("fresh resume task is not running") + return fmt.Errorf("fresh %s task is not running", event.Payload.Action) } if err := c.VerifySIP(ctx, snapshot.SIP); err != nil { - return fmt.Errorf("resume SIP revision not applied: %w", err) + return fmt.Errorf("%s SIP revision not applied: %w", event.Payload.Action, err) } if err := validateAISnapshot(snapshot); err != nil { return err } if err := c.Store.SaveSnapshot(snapshot); err != nil { - return fmt.Errorf("bind fresh resume task: %w", err) + return fmt.Errorf("bind fresh %s task: %w", event.Payload.Action, err) + } + if event.Payload.Action == "start" { + if err := c.Store.CompleteStart(snapshot, event.EventID); err != nil { + if errors.Is(err, store.ErrControlRejected) { + return c.Store.RejectControl(event.DispatcherID, event.TenantID, event.EventID) + } + return fmt.Errorf("start task %q: %w", event.Payload.TaskID, err) + } + return nil + } + if err := c.Store.ApplyDiscoveryPage(event.DispatcherID, []configread.DiscoveredTask{{TaskID: snapshot.Task.TaskID, TenantID: snapshot.Task.TenantID, TaskRevision: snapshot.Task.TaskRevision, Status: snapshot.Task.Status}}); err != nil { + return fmt.Errorf("bind fresh resume task revision: %w", err) } } else { if policy == "" { diff --git a/internal/dispatcher/control_test.go b/internal/dispatcher/control_test.go index 4c93b57..25dac66 100644 --- a/internal/dispatcher/control_test.go +++ b/internal/dispatcher/control_test.go @@ -213,6 +213,110 @@ func TestResumeRequiresFreshTaskAndAppliedSIP(t *testing.T) { } } +func TestStartReadsNewTaskWithoutAgentControl(t *testing.T) { + controller, agent, s := newControlFixture(t) + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "application/json") + switch r.URL.Path { + case "/internal/v1/dispatcher/sip": + raw, _ := json.Marshal(policySnapshot(t).SIP) + _, _ = w.Write(raw) + case "/internal/v1/dispatcher/task/task-new": + _, _ = w.Write([]byte(strings.ReplaceAll(string(configExample(t, "config-read-task-asr")), "task-asr", "task-new"))) + case "/internal/v1/dispatcher/ai-providers": + _, _ = w.Write(configExample(t, "config-read-providers")) + case "/internal/v1/dispatcher/tenant/1001/quota": + _, _ = w.Write(configExample(t, "config-read-quota")) + default: + t.Errorf("unexpected HTTP path %s", r.URL.Path) + w.WriteHeader(http.StatusNotFound) + } + })) + defer server.Close() + var err error + controller.Client, err = configread.NewClient(server.URL, controller.DispatcherID, "test-secret", server.Client()) + if err != nil { + t.Fatal(err) + } + body := strings.ReplaceAll(string(controlBody(t, "start-new", "start", "")), "task-asr", "task-new") + if err := controller.ProcessControl(context.Background(), []byte(body)); err != nil { + t.Fatal(err) + } + if len(agent.calls) != 0 { + t.Fatalf("start unexpectedly sent Agent control: %+v", agent.calls) + } + if admitted, err := s.CanAdmit(controller.DispatcherID, 1001, "task-new"); err != nil || !admitted { + t.Fatalf("new task not admitted after start: %v %v", admitted, err) + } + outbox, err := s.ListPendingOutbox(controller.DispatcherID) + if err != nil || len(outbox) != 1 || !strings.Contains(string(outbox[0].Body), `"status":"applied"`) { + t.Fatalf("start acknowledgment missing: %+v %v", outbox, err) + } +} + +func TestStartDoesNotReopenPausedOrStoppedTasks(t *testing.T) { + for _, action := range []string{"pause", "stop"} { + t.Run(action, func(t *testing.T) { + controller, agent, s := newControlFixture(t) + if err := controller.ProcessControl(context.Background(), controlBody(t, "preceding-"+action, action, "")); err != nil { + t.Fatal(err) + } + if err := controller.ProcessControl(context.Background(), controlBody(t, "late-start-"+action, "start", "")); err != nil { + t.Fatal(err) + } + if len(agent.calls) != 1 { + t.Fatalf("start unexpectedly reached Agent: %+v", agent.calls) + } + if admitted, err := s.CanAdmit(controller.DispatcherID, 1001, "task-asr"); err != nil || admitted { + t.Fatalf("start reopened %s task: %v %v", action, admitted, err) + } + outbox, err := s.ListPendingOutbox(controller.DispatcherID) + if err != nil || len(outbox) != 2 || !strings.Contains(string(outbox[1].Body), `"status":"rejected"`) { + t.Fatalf("start rejection acknowledgment missing: %+v %v", outbox, err) + } + }) + } +} + +func TestResumeRefreshesModifiedTaskRevision(t *testing.T) { + controller, agent, s := newControlFixture(t) + if err := controller.ProcessControl(context.Background(), controlBody(t, "pause-before-edit", "pause", "")); err != nil { + t.Fatal(err) + } + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "application/json") + switch r.URL.Path { + case "/internal/v1/dispatcher/sip": + raw, _ := json.Marshal(policySnapshot(t).SIP) + _, _ = w.Write(raw) + case "/internal/v1/dispatcher/task/task-asr": + _, _ = w.Write([]byte(strings.Replace(string(configExample(t, "config-read-task-asr")), `"task_revision":1`, `"task_revision":2`, 1))) + case "/internal/v1/dispatcher/ai-providers": + _, _ = w.Write(configExample(t, "config-read-providers")) + case "/internal/v1/dispatcher/tenant/1001/quota": + _, _ = w.Write(configExample(t, "config-read-quota")) + default: + t.Errorf("unexpected HTTP path %s", r.URL.Path) + w.WriteHeader(http.StatusNotFound) + } + })) + defer server.Close() + var err error + controller.Client, err = configread.NewClient(server.URL, controller.DispatcherID, "test-secret", server.Client()) + if err != nil { + t.Fatal(err) + } + if err := controller.ProcessControl(context.Background(), controlBody(t, "resume-after-edit", "resume", "")); err != nil { + t.Fatal(err) + } + if len(agent.calls) != 2 || agent.calls[1].Action != "resume" { + t.Fatalf("edited task was not resumed: %+v", agent.calls) + } + if admitted, err := s.CanAdmit(controller.DispatcherID, 1001, "task-asr"); err != nil || !admitted { + t.Fatalf("edited task not admitted: %v %v", admitted, err) + } +} + func TestRuntimeSurfacesControlFailureToCloseAdmission(t *testing.T) { controller, agent, s := newControlFixture(t) agent.err = errors.New("injected Agent control failure") diff --git a/internal/dispatcher/discovery.go b/internal/dispatcher/discovery.go deleted file mode 100644 index e0142ea..0000000 --- a/internal/dispatcher/discovery.go +++ /dev/null @@ -1,82 +0,0 @@ -package dispatcher - -import ( - "bytes" - "context" - "encoding/json" - "errors" - "fmt" - - "git.ipao.vip/rogee/go-sip/internal/configread" - "git.ipao.vip/rogee/go-sip/internal/store" -) - -// DiscoveryFollower keeps only an in-memory cursor. A restart always -// takes a complete snapshot and drains the control queue before admitting. -type DiscoveryFollower struct { - DispatcherID string - Client *configread.Client - Store *store.Store - ApprovedSIP configread.SIP - VerifySIP func(context.Context, configread.SIP) error - Cursor string -} - -func (f *DiscoveryFollower) Poll(ctx context.Context) error { - if f == nil || f.Client == nil || f.Store == nil || f.DispatcherID == "" || f.ApprovedSIP.DispatcherID != f.DispatcherID || f.ApprovedSIP.Revision <= 0 || f.VerifySIP == nil || f.Cursor == "" { - return errors.New("discovery follower requires a verified full snapshot and in-memory cursor") - } - approved, err := json.Marshal(f.ApprovedSIP) - if err != nil { - return f.fail(fmt.Errorf("encode approved SIP snapshot: %w", err)) - } - for { - page, err := f.Client.ReadTasks(ctx, f.Cursor) - if err != nil { - return f.fail(fmt.Errorf("read task discovery after cursor: %w", err)) - } - if len(page.Tasks) == 0 { - f.Cursor = page.Cursor - return nil - } - if err := f.Store.ApplyDiscoveryPage(f.DispatcherID, page.Tasks); err != nil { - return f.fail(fmt.Errorf("persist discovery page before cursor advance: %w", err)) - } - for _, task := range page.Tasks { - if task.Status == "stopped" { - continue - } - snapshot, err := f.Client.ReadTask(ctx, task.TaskID, task.TenantID) - if err != nil { - return f.fail(fmt.Errorf("load discovered task %q: %w", task.TaskID, err)) - } - if snapshot.Task.TaskRevision != task.TaskRevision || snapshot.Task.Status != task.Status { - return f.fail(fmt.Errorf("discovery and task configuration conflict for %q", task.TaskID)) - } - actual, err := json.Marshal(snapshot.SIP) - if err != nil { - return f.fail(fmt.Errorf("encode task %q SIP snapshot: %w", task.TaskID, err)) - } - if !bytes.Equal(approved, actual) { - return f.fail(fmt.Errorf("task %q SIP differs from the approved loaded revision", task.TaskID)) - } - if err := f.VerifySIP(ctx, snapshot.SIP); err != nil { - return f.fail(fmt.Errorf("task %q SIP is not loaded by Agent/Asterisk: %w", task.TaskID, err)) - } - if err := validateAISnapshot(snapshot); err != nil { - return f.fail(err) - } - if err := f.Store.SaveSnapshot(snapshot); err != nil { - return f.fail(fmt.Errorf("bind discovered task %q: %w", task.TaskID, err)) - } - } - f.Cursor = page.Cursor - } -} - -func (f *DiscoveryFollower) fail(cause error) error { - if err := f.Store.CloseAdmission(f.DispatcherID); err != nil { - return errors.Join(cause, fmt.Errorf("close admission after discovery failure: %w", err)) - } - return cause -} diff --git a/internal/dispatcher/discovery_test.go b/internal/dispatcher/discovery_test.go deleted file mode 100644 index 542a4ce..0000000 --- a/internal/dispatcher/discovery_test.go +++ /dev/null @@ -1,87 +0,0 @@ -package dispatcher - -import ( - "context" - "encoding/json" - "fmt" - "net/http" - "net/http/httptest" - "testing" - - "git.ipao.vip/rogee/go-sip/internal/configread" -) - -func newFollowerFixture(t *testing.T, revision int64) (*DiscoveryFollower, func() bool) { - t.Helper() - executor, _, _, s := newExecuteFixture(t) - approved := policySnapshot(t).SIP - sipBody, err := json.Marshal(approved) - if err != nil { - t.Fatal(err) - } - terminalSeen := false - server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { - w.Header().Set("Content-Type", "application/json") - var body []byte - switch r.URL.Path { - case "/internal/v1/dispatcher/tasks": - switch r.URL.Query().Get("after") { - case "initial": - body = []byte(fmt.Sprintf(`{"dispatcher_id":%q,"cursor":"middle","tasks":[{"task_id":"task-asr","tenant_id":1001,"status":"running","task_revision":%d}]}`, executor.DispatcherID, revision)) - case "middle": - terminalSeen = true - body = []byte(fmt.Sprintf(`{"dispatcher_id":%q,"cursor":"opaque-tail","tasks":[]}`, executor.DispatcherID)) - default: - t.Errorf("unexpected discovery cursor: %s", r.URL.RawQuery) - w.WriteHeader(http.StatusBadRequest) - return - } - case "/internal/v1/dispatcher/sip": - body = sipBody - case "/internal/v1/dispatcher/ai-providers": - body = configExample(t, "config-read-providers") - case "/internal/v1/dispatcher/task/task-asr": - body = configExample(t, "config-read-task-asr") - case "/internal/v1/dispatcher/tenant/1001/quota": - body = configExample(t, "config-read-quota") - default: - t.Errorf("unexpected config path: %s", r.URL.Path) - w.WriteHeader(http.StatusNotFound) - return - } - _, _ = w.Write(body) - })) - t.Cleanup(server.Close) - client, err := configread.NewClient(server.URL, executor.DispatcherID, "test-secret", server.Client()) - if err != nil { - t.Fatal(err) - } - follower := &DiscoveryFollower{DispatcherID: executor.DispatcherID, Client: client, Store: s, ApprovedSIP: approved, VerifySIP: func(context.Context, configread.SIP) error { return nil }, Cursor: "initial"} - return follower, func() bool { return terminalSeen } -} - -func TestDiscoveryFollowerPersistsPageBeforeAdvancingToTerminalEmpty(t *testing.T) { - follower, terminal := newFollowerFixture(t, 1) - if err := follower.Poll(context.Background()); err != nil { - t.Fatal(err) - } - if follower.Cursor != "opaque-tail" || !terminal() { - t.Fatalf("cursor did not advance through terminal empty page: %q %v", follower.Cursor, terminal()) - } - if admitted, err := follower.Store.CanAdmit(follower.DispatcherID, 1001, "task-asr"); err != nil || !admitted { - t.Fatalf("valid discovery closed admission: %v %v", admitted, err) - } -} - -func TestDiscoveryFollowerFailedPageKeepsCursorAndClosesAdmission(t *testing.T) { - follower, terminal := newFollowerFixture(t, 2) - if err := follower.Poll(context.Background()); err == nil { - t.Fatal("task revision mismatch was accepted") - } - if follower.Cursor != "initial" || terminal() { - t.Fatalf("failed page advanced cursor: %q %v", follower.Cursor, terminal()) - } - if admitted, err := follower.Store.CanAdmit(follower.DispatcherID, 1001, "task-asr"); err != nil || admitted { - t.Fatalf("failed page left task admission open: %v %v", admitted, err) - } -} diff --git a/internal/dispatcher/runtime.go b/internal/dispatcher/runtime.go index e50dcf9..89c7162 100644 --- a/internal/dispatcher/runtime.go +++ b/internal/dispatcher/runtime.go @@ -18,13 +18,13 @@ import ( // Runtime owns task/control consumers for one Dispatcher. SaaS creates // every queue/binding; this process only checks and consumes predeclared ones. type Runtime struct { - Broker *mq.Broker - Bootstrap Bootstrap - Execute ExecuteController - Control ControlController - PollInterval time.Duration - DiscoveryInterval time.Duration - Logger *slog.Logger + Broker *mq.Broker + Bootstrap Bootstrap + Execute ExecuteController + Control ControlController + PollInterval time.Duration + SIPInterval time.Duration + Logger *slog.Logger gate sync.RWMutex // SIP admission barrier versus each pre-dial instruction locksMu sync.Mutex @@ -35,7 +35,7 @@ type Runtime struct { // Serve closes admission on every shutdown/failure; it never clears durable // calls, results, task queues, or the SQLite file. 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.DiscoveryInterval <= 0 || r.Logger == nil || r.Bootstrap.DrainControls != nil { + 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.SIPInterval <= 0 || r.Logger == nil || r.Bootstrap.DrainControls != nil { return errors.New("runtime requires one Dispatcher, durable state, verified SIP, independent clocks, and configured polling; external control drain is forbidden") } if err := r.Execute.validate(); err != nil { @@ -66,9 +66,7 @@ func (r *Runtime) Serve(ctx context.Context) (result error) { } }() - var cursor string var sip configread.SIP - r.Bootstrap.Cursor = &cursor r.Bootstrap.SIP = &sip r.Bootstrap.DrainControls = func(ctx context.Context) error { _, err := r.Broker.DrainControlPredeclared(ctx, r.Broker.ControlQueue(), r.handleControl) @@ -93,7 +91,6 @@ func (r *Runtime) Serve(ctx context.Context) (result error) { } r.Logger.Warn("control and outbox stay active while newer SIP revision waits; task admission remains closed", "dispatcher_id", r.Bootstrap.DispatcherID, "error", err) } - follower := &DiscoveryFollower{DispatcherID: r.Bootstrap.DispatcherID, Client: r.Bootstrap.Client, Store: r.Bootstrap.Store, ApprovedSIP: sip, VerifySIP: r.Bootstrap.VerifySIP, Cursor: cursor} if err := r.syncTaskConsumers(ctx, consumers); err != nil { return err } @@ -102,8 +99,8 @@ func (r *Runtime) Serve(ctx context.Context) (result error) { } poll := time.NewTicker(r.PollInterval) defer poll.Stop() - discovery := time.NewTicker(r.DiscoveryInterval) - defer discovery.Stop() + sipTick := time.NewTicker(r.SIPInterval) + defer sipTick.Stop() for { select { case <-ctx.Done(): @@ -120,19 +117,10 @@ func (r *Runtime) Serve(ctx context.Context) (result error) { if err := r.syncTaskConsumers(ctx, consumers); err != nil { return err } - case <-discovery.C: - if err := r.refreshSIP(ctx, follower); err != nil { + case <-sipTick.C: + if err := r.refreshSIP(ctx, &sip); err != nil { return fmt.Errorf("refresh approved SIP: %w", err) } - _, pending, err := r.Bootstrap.Store.SIPState(r.Bootstrap.DispatcherID) - if err != nil { - return err - } - if pending == 0 { - if err := follower.Poll(ctx); err != nil { - return fmt.Errorf("poll assigned tasks: %w", err) - } - } if err := r.syncTaskConsumers(ctx, consumers); err != nil { return err } diff --git a/internal/dispatcher/runtime_integration_test.go b/internal/dispatcher/runtime_integration_test.go index 85ff918..670afd3 100644 --- a/internal/dispatcher/runtime_integration_test.go +++ b/internal/dispatcher/runtime_integration_test.go @@ -65,8 +65,9 @@ func TestRuntimeIsolatedControlBacklogExecuteAndSharedResult(t *testing.T) { controlRoute, _ := tenant.ControlRoute(id) taskRoute, _ := tenant.TaskRoute(id, "task-asr") resultRoute, _ := tenant.ResultRoute(id) + newTaskRoute, _ := tenant.TaskRoute(id, "task-new") shared := "agent-call.saas.events.v1" - for _, queue := range []string{controlRoute.Queue, taskRoute.Queue, shared} { + for _, queue := range []string{controlRoute.Queue, taskRoute.Queue, newTaskRoute.Queue, shared} { if _, err := admin.QueueDeclare(queue, true, false, false, false, nil); err != nil { t.Fatal(err) } @@ -75,7 +76,7 @@ func TestRuntimeIsolatedControlBacklogExecuteAndSharedResult(t *testing.T) { t.Fatal(err) } } - for _, route := range []tenant.Route{controlRoute, taskRoute, resultRoute} { + for _, route := range []tenant.Route{controlRoute, taskRoute, newTaskRoute, resultRoute} { queue := route.Queue if route.BindingKey == resultRoute.BindingKey { queue = shared @@ -99,6 +100,7 @@ func TestRuntimeIsolatedControlBacklogExecuteAndSharedResult(t *testing.T) { if err != nil { t.Fatal(err) } + var taskListReads atomic.Int32 server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { w.Header().Set("Content-Type", "application/json") var body []byte @@ -106,6 +108,7 @@ func TestRuntimeIsolatedControlBacklogExecuteAndSharedResult(t *testing.T) { case "/internal/v1/dispatcher/sip": body = sipJSON case "/internal/v1/dispatcher/tasks": + taskListReads.Add(1) if r.URL.Query().Get("after") == "" { body = configExample(t, "task-discovery-page") } else if r.URL.Query().Get("after") != "" { @@ -116,6 +119,8 @@ func TestRuntimeIsolatedControlBacklogExecuteAndSharedResult(t *testing.T) { } case "/internal/v1/dispatcher/task/task-asr": body = configExample(t, "config-read-task-asr") + case "/internal/v1/dispatcher/task/task-new": + body = []byte(strings.ReplaceAll(string(configExample(t, "config-read-task-asr")), "task-asr", "task-new")) case "/internal/v1/dispatcher/ai-providers": body = configExample(t, "config-read-providers") case "/internal/v1/dispatcher/tenant/1001/quota": @@ -161,7 +166,7 @@ func TestRuntimeIsolatedControlBacklogExecuteAndSharedResult(t *testing.T) { return monday(8, 59) }}, Control: ControlController{DispatcherID: id, Store: db, Client: client, Agent: agent, VerifySIP: verify, Now: func() time.Time { return monday(9, 30) }}, - PollInterval: 30 * time.Millisecond, DiscoveryInterval: 120 * time.Millisecond, Logger: slog.Default(), + PollInterval: 30 * time.Millisecond, SIPInterval: 120 * time.Millisecond, Logger: slog.Default(), } ctx, cancel := context.WithCancel(context.Background()) finished := make(chan struct{}) @@ -259,6 +264,27 @@ func TestRuntimeIsolatedControlBacklogExecuteAndSharedResult(t *testing.T) { t.Fatalf("redelivery reoriginated call %s", spec.EventID) case <-time.After(200 * time.Millisecond): } + publish(controlRoute, []byte(strings.ReplaceAll(string(controlBody(t, "new-task-start", "start", "")), "task-asr", "task-new"))) + startDeadline := time.After(5 * time.Second) + for { + newQueue, err := admin.QueueInspect(newTaskRoute.Queue) + if err != nil { + t.Fatal(err) + } + if newQueue.Consumers == 1 { + break + } + select { + case <-finished: + t.Fatalf("runtime stopped before new task started: %v", runtimeErr) + case <-startDeadline: + t.Fatal("start did not attach the new task queue") + case <-time.After(20 * time.Millisecond): + } + } + if got := taskListReads.Load(); got != 2 { + t.Fatalf("unexpected periodic task-list reads: %d", got) + } // A temporary task rule wait is retained durably, then its consumer // stops so further instructions remain in SaaS's task queue. Admission // resumes from the original identity when the configured window opens. diff --git a/internal/dispatcher/sip_reload.go b/internal/dispatcher/sip_reload.go index 9db75ca..3f9eb11 100644 --- a/internal/dispatcher/sip_reload.go +++ b/internal/dispatcher/sip_reload.go @@ -6,20 +6,22 @@ import ( "encoding/json" "errors" "fmt" + + "git.ipao.vip/rogee/go-sip/internal/configread" ) // refreshSIP keeps control/result delivery running while a new approved SIP // revision waits for old calls to drain and for Agent/Asterisk to load it. // No HTTP snapshot or notification by itself authorizes a real call. -func (r *Runtime) refreshSIP(ctx context.Context, follower *DiscoveryFollower) error { - if r == nil || follower == nil || r.Bootstrap.Client == nil || r.Bootstrap.Store == nil || r.Bootstrap.VerifySIP == nil || r.Logger == nil { +func (r *Runtime) refreshSIP(ctx context.Context, approvedSIP *configread.SIP) error { + if r == nil || approvedSIP == nil || r.Bootstrap.Client == nil || r.Bootstrap.Store == nil || r.Bootstrap.VerifySIP == nil || r.Logger == nil { return errors.New("SIP refresh requires durable state and applied-revision verifier") } current, err := r.Bootstrap.Client.ReadSIP(ctx) if err != nil { return r.closeSIPAdmission(fmt.Errorf("read approved SIP full snapshot: %w", err)) } - old := follower.ApprovedSIP + old := *approvedSIP if current.Revision < old.Revision { return r.closeSIPAdmission(fmt.Errorf("approved SIP revision regressed from %d to %d", old.Revision, current.Revision)) } @@ -68,7 +70,7 @@ func (r *Runtime) refreshSIP(ctx context.Context, follower *DiscoveryFollower) e // A full discovery snapshot is required after an approved SIP revision // change; a delta page cannot prove that all assigned tasks use the same // loaded revision. MQ control consumption remains active during this read. - tasks, cursor, err := r.Bootstrap.Client.ReadAllTasks(ctx) + tasks, _, err := r.Bootstrap.Client.ReadAllTasks(ctx) if err != nil { return r.closeSIPAdmission(fmt.Errorf("reload complete assigned task list for SIP: %w", err)) } @@ -113,8 +115,7 @@ func (r *Runtime) refreshSIP(ctx context.Context, follower *DiscoveryFollower) e if err := r.Bootstrap.Store.MarkReadyForSIP(r.Bootstrap.DispatcherID, current.Revision); err != nil { return r.closeSIPAdmission(fmt.Errorf("verify and commit reloaded SIP revision: %w", err)) } - follower.ApprovedSIP = current - follower.Cursor = cursor + *approvedSIP = current if r.Bootstrap.SIP != nil { *r.Bootstrap.SIP = current } diff --git a/internal/dispatcher/sip_runtime_test.go b/internal/dispatcher/sip_runtime_test.go index 8229f75..c53d639 100644 --- a/internal/dispatcher/sip_runtime_test.go +++ b/internal/dispatcher/sip_runtime_test.go @@ -61,7 +61,7 @@ func TestSIPNotificationPersistsBarrierAndWaitsForLoadedFullSnapshot(t *testing. return nil } runtime := &Runtime{Bootstrap: Bootstrap{DispatcherID: executor.DispatcherID, Client: client, Store: s, VerifySIP: verify}, Logger: slog.Default()} - follower := &DiscoveryFollower{DispatcherID: executor.DispatcherID, Client: client, Store: s, ApprovedSIP: approved, VerifySIP: verify, Cursor: "opaque-end-token"} + body := []byte(strings.Replace(string(configExample(t, "mq-sip-change")), `"revision":8`, `"revision":9`, 1)) if err := runtime.handleControl(context.Background(), "", body); err != nil { t.Fatal(err) @@ -72,20 +72,20 @@ func TestSIPNotificationPersistsBarrierAndWaitsForLoadedFullSnapshot(t *testing. if applied, pending, err := s.SIPState(executor.DispatcherID); err != nil || applied != 8 || pending != 9 { t.Fatalf("notification not durable before ACK: %d %d %v", applied, pending, err) } - if err := runtime.refreshSIP(context.Background(), follower); err != nil { + if err := runtime.refreshSIP(context.Background(), &approved); err != nil { t.Fatal(err) } if admitted, err := s.CanAdmit(executor.DispatcherID, 1001, "task-asr"); err != nil || admitted { t.Fatalf("unloaded SIP reopened admission: %v %v", admitted, err) } loaded.Store(true) - if err := runtime.refreshSIP(context.Background(), follower); err != nil { + if err := runtime.refreshSIP(context.Background(), &approved); err != nil { t.Fatal(err) } if admitted, err := s.CanAdmit(executor.DispatcherID, 1001, "task-asr"); err != nil || !admitted { t.Fatalf("verified SIP did not reopen: %v %v", admitted, err) } - if follower.ApprovedSIP.Revision != 9 || follower.Cursor != "opaque-end-token" { - t.Fatalf("reloaded SIP/cursor not bound: rev=%d cursor=%q", follower.ApprovedSIP.Revision, follower.Cursor) + if approved.Revision != 9 { + t.Fatalf("reloaded SIP not bound: rev=%d", approved.Revision) } } diff --git a/internal/store/control_flow.go b/internal/store/control_flow.go index 0414ca9..b7cd199 100644 --- a/internal/store/control_flow.go +++ b/internal/store/control_flow.go @@ -6,12 +6,46 @@ import ( "errors" "fmt" + "git.ipao.vip/rogee/go-sip/internal/configread" "git.ipao.vip/rogee/go-sip/internal/contract" "git.ipao.vip/rogee/go-sip/internal/tenant" ) var ErrControlRejected = errors.New("task control rejected by durable task state") +// CompleteStart binds a freshly verified task to its predeclared queue and +// records the control acknowledgment together. It never reopens a task +// that is currently paused or stopped. +func (s *Store) CompleteStart(snapshot configread.Snapshot, eventID string) error { + task := snapshot.Task + if task.Status != "running" { + return fmt.Errorf("%w: start requires a running task", ErrControlRejected) + } + tx, err := s.db.Begin() + if err != nil { + return err + } + defer tx.Rollback() + var state, status string + err = tx.QueryRow(`SELECT control_state,status FROM dispatcher_tasks WHERE dispatcher_id=? AND tenant_id=? AND task_id=?`, task.DispatcherID, task.TenantID, task.TaskID).Scan(&state, &status) + if err != nil && !errors.Is(err, sql.ErrNoRows) { + return fmt.Errorf("check start task state: %w", err) + } + if err == nil && (state != "" || status != "running") { + return fmt.Errorf("%w: start cannot reopen a paused or stopped task", ErrControlRejected) + } + if err := applyDiscoveredTasks(tx, task.DispatcherID, []configread.DiscoveredTask{{TenantID: task.TenantID, TaskID: task.TaskID, TaskRevision: task.TaskRevision, Status: task.Status}}); err != nil { + return err + } + if err := enqueueControlAck(tx, task.DispatcherID, task.TenantID, eventID, "applied"); err != nil { + return err + } + if err := tx.Commit(); err != nil { + return fmt.Errorf("commit started task and control acknowledgment: %w", err) + } + return nil +} + // PrepareControl closes only this task's admission before requesting the // Agent action. No successful control acknowledgment exists at this point. // Repeated commands are prepared and dispatched again; this is not control