Apply task configuration edits from control events
This commit is contained in:
@@ -88,7 +88,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 不可自行建/删/绑定;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 和时限由该任务快照固定;主叫由绑定线路的 SIP 快照唯一 `caller_id` 决定,不再由任务挑选主叫档案;`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 须重新核验。历史命令不因消息年龄过期,但实际呼出前仍检查任务归属/事件号码格式/时段/授权/额度和有效通话上限;任务结束不删未确认结果、恢复、幂等或未知占用。
|
||||
- `call.execute` 只带获批 `task_id/callee`,调用线路、AI 和时限由该任务快照固定;主叫由绑定线路的 SIP 快照唯一 `caller_id` 决定,不再由任务挑选主叫档案;`task.control` 的 start/pause/resume/stop/edit 无旧 CAS/版本字段,控制回 task/action/status 处理结果,stop 不可恢复。任务启动/重启完整读取一次列表,运行中不定时轮询任务列表或任务配置;新建任务由 start 读取任务配置和租户额度,resume 重读并解除人工暂停,edit 只重读、校验并持久更新已有任务的配置和额度,不改变人工暂停/停止状态、不向 Agent 发动作,也不影响已发出的通话快照;旧修订号及停止后 edit 拒绝。全局 AI 服务商列表每次启动只读取一次,本进程全部任务复用,变更须重启后才生效。启动时全量读取并核验 SIP,运行中只由 `sip.config` 通知触发全量 SIP 读取,无常规定时 SIP 轮询;任务 start/resume/edit 使用已生效 SIP 快照,不自行拉取或向 Agent 核验 SIP。SIP 修订变化待旧呼叫结束且 Agent 已加载后,仅在本地重新绑定任务快照,不重读任务列表;待处理期间新准入关闭。发现分页同快照先全部校验后提交,MQ 控制持久状态高于偶发 HTTP 状态;冲突关准入、resume 须重新核验。历史命令不因消息年龄过期,但实际呼出前仍检查任务归属/事件号码格式/时段/授权/额度和有效通话上限;任务结束不删未确认结果、恢复、幂等或未知占用。
|
||||
- 独立 Dispatcher 的 SQLite 是任务、额度、inbox/outbox 的权威数据;Agent 无业务数据库,录音、执行与上传恢复只写受控私有文件。额度包含未知占用,新 boot/租约到期不得自动清除未知执行;明确 Agent 在拨号前拒绝时须持久拒绝并释放,已接纳但证实 ARI originate 从未提交的失败由原 Agent 经正式终结回执释放,RPC 交付未知时仍占用。2026-10-04 修复版 `d2c4a99` 已安装并通过非呼出时段主机诊断,新隔离测试租户额度 2/revision 2 已启用并由 Dispatcher 只读拉取;原快照保留;该条旧执行仅按前述经使用者授权的人工核实记录解除 Dispatcher 占用,绝非新 boot 或超时自动清除未知执行。仅 Agent 用机器可读标记证实本次未发起的拒绝可自动释放;旧执行的同错误码不算证明。不实现双活数据库、自动跨机热备、多 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-only Agent 在配置原子写入、PJSIP reload 和运行态 endpoint/AOR/UDP transport 一致后持久标记 revision,每次加载查询重新核验。只支持明确的 UDP、IP/none 鉴权、无需 REGISTER 的 IPv4/PCMA 线路;未知字段和其他传输/鉴权/注册方式拒绝,不热更静态 transport。每条 SIP trunk 只允许一个原值 `caller_id`,Agent 把该值同时配置为 Asterisk endpoint 的 `from_user` 和 `callerid`,并在 `GetLoadedSIP` 前核对实际运行值;旧任务 `caller_profile_id` 与线路 `caller_profiles` 必须直接拒绝,不保留兼容路径。SIP 配置的唯一编辑/审批面仍是 management。
|
||||
- 非生产真实 Agent 的 Asterisk `res_hep`/`res_hep_pjsip` 仅镜像至本机 UDP:在 ARI Dial 前读取本次 SIP Call-ID 并绑定执行,按 Call-ID 和原始 INVITE transaction 只采集真实最终响应的状态码、状态行与完整 `raw`;丢失/无法关联以 `sip_capture_error` 显式报告,不从 ARI、号码或挂断原因推测。镜像缺失不阻止已证实终结的额度释放;未经确认的占用仍保留。配置及回滚见 [`deploys/cell/README.md`](deploys/cell/README.md),不得将完整 SIP 报文写入普通日志、提交或聊天;本地测试通过不等于测试机已部署 HEP 或真实接通。
|
||||
|
||||
@@ -0,0 +1 @@
|
||||
{"event_id":"control-edit-invalid","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":"edit","reason":"configuration changed","options":{"active_call_policy":"hangup"}}}
|
||||
@@ -0,0 +1 @@
|
||||
{"event_id":"control-edit-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":"edit","reason":"configuration changed"}}
|
||||
@@ -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": "f0dc3c419a3f2de2ad08e00a65e6259f7810035fe8c59feffcfcc024ed4a6a53"
|
||||
"docs/thirds/saas-dispatcher.md": "478cf69df793e6bacb6d19dcb50c1a9dc56d4cd9cae59ea5744ce94ee9e9d08c"
|
||||
},
|
||||
"bundle_sha256": "29e590fbca2a389e8c003ef475398f16968c69fa4bdc24d9bc35ea07e6934bd1",
|
||||
"bundle_sha256": "d25884427faaa4e096cfb2eb769d47216fef25f8e0359fb352a0c4934c8ed07c",
|
||||
"bundle_algorithm": "sha256 of sorted relative-path + space + sha256(file) + newline; only root-level JSON and examples/**/*.json, excluding manifest.json"
|
||||
}
|
||||
|
||||
@@ -34,9 +34,9 @@
|
||||
"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": ["start", "pause", "resume", "stop"]}, "reason": {"type": "string", "minLength": 1},
|
||||
"task_id": {"$ref": "#/$defs/task_id"}, "action": {"enum": ["start", "pause", "resume", "stop", "edit"]}, "reason": {"type": "string", "minLength": 1},
|
||||
"options": {"type": "object", "additionalProperties": false, "properties": {"active_call_policy": {"enum": ["drain", "hangup"]}}}
|
||||
}}
|
||||
}, "allOf": [{"if": {"properties": {"action": {"const": "edit"}}}, "then": {"not": {"required": ["options"]}}}] }
|
||||
}
|
||||
},
|
||||
"control_ack": {
|
||||
|
||||
@@ -14,7 +14,7 @@
|
||||
| 任务列表 | `GET /internal/v1/dispatcher/tasks`,后续 `?after=<cursor>` | 首次/重启完整读到 **带 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` 不再启用。任务配置只能在暂停或停止后修改;没有 `edit` 事件,也不定时刷新任务列表。启动/重启时完整读取任务列表及独立的全量 SIP;运行中新建任务由 `task.control` 的 `start` 触发读取,人工暂停后修改配置由 `resume` 触发重读任务和租户配额,两者沿用启动时的 AI 服务商列表和已生效的 SIP 快照,不另行读取或向 Agent 核验全量 SIP。HTTP 旧 running 不得覆盖已持久的 pause/stop;resume 经最新任务配置核验后只能解除人工暂停。SIP 无常规定时拉取;`sip.config` 持久关准入并触发 SIP 全量读取,等待已接纳呼叫结束及 Agent 加载核验后,将新 SIP 绑定到本地任务快照(不重读任务列表),然后才可恢复准入。通知待处理时可重试,读取或校验失败不回退旧 SIP,也不开准入。
|
||||
HTTP 非 200、正文不合法、缺字段、归属冲突、缓存失效或读取失败时拒绝新的相关准入并记录脱敏错误;已持久接纳的执行继续使用原快照。错误体示例 [`error`](../../contracts/local/examples/config-read-error.json),不能当成功配置解析。任务仅有运行/暂停/终止;终止是不可逆 stop,同一 `task_id` 不再启用。任务配置可在运行或暂停时通过 `task.control.edit` 更新:D 按事件重读该任务及租户额度,校验归属、修订号与 AI 绑定后持久切换未来呼叫所用快照;已发出的通话始终保留原快照,`edit` 不改变暂停/停止状态,也不下发 Agent 控制。旧修订号不得覆盖新修订号,停止的任务不能被 `edit` 恢复。启动/重启时完整读取一次任务列表及独立的全量 SIP;运行中不定时刷新任务列表或任务详情,新增任务由 `start` 触发读取,`resume` 仍重读最新任务和租户额度后解除人工暂停。start/resume/edit 均沿用启动时的 AI 服务商列表和已生效 SIP 快照,不另行读取或向 Agent 核验全量 SIP。HTTP 旧 running 不得覆盖已持久的 pause/stop;resume 经最新任务配置核验后只能解除人工暂停。SIP 无常规定时拉取;`sip.config` 持久关准入并触发 SIP 全量读取,等待已接纳呼叫结束及 Agent 加载核验后,将新 SIP 绑定到本地任务快照(不重读任务列表),然后才可恢复准入。通知待处理时可重试,读取或校验失败不回退旧 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;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) |
|
||||
| `task.control` | SaaS→D;start/pause/resume/stop/edit,无控制去重/CAS 字段。start 读取新任务配置及租户配额并持久绑定,无 Agent 控制动作;resume 重读最新任务配置和额度后解除人工暂停;edit 重读并持久更新已有任务的配置和额度,不改变控制状态,不向 Agent 发动作,不接收 `options`,已发出通话沿用原快照;三者均沿用启动时 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) · [`edit`](../../contracts/local/examples/mq-control-edit.json) · [`control`](../../contracts/local/examples/mq-control.json) · [`ack`](../../contracts/local/examples/mq-control-ack.json) |
|
||||
| `call.execute` | SaaS→D 仅 `{task_id,callee}`;一次指令保留独立消息/执行身份,路由/AI/时限从绑定任务读取,主叫由绑定 SIP 线路快照的唯一 `caller_id` 决定。`callee` 原值由通过归属及任务校验的 SaaS 事件确定,不从任务快照或本地号码列表另选;`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) |
|
||||
|
||||
|
||||
@@ -79,6 +79,26 @@ 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 == "edit" {
|
||||
if policy != "" {
|
||||
return errors.New("edit control must not carry an active-call policy")
|
||||
}
|
||||
// Refresh only the bound task configuration. Already issued calls
|
||||
// retain their original immutable execution snapshot.
|
||||
snapshot, err := c.Client.ReadTask(ctx, event.Payload.TaskID, event.TenantID, *c.ApprovedSIP, *c.ApprovedProviders)
|
||||
if err != nil {
|
||||
return fmt.Errorf("fresh edit task configuration: %w", err)
|
||||
}
|
||||
if err := validateAISnapshot(snapshot); err != nil {
|
||||
return err
|
||||
}
|
||||
status, err := c.Store.CompleteEdit(snapshot, event.EventID)
|
||||
if err != nil {
|
||||
return fmt.Errorf("commit edit task configuration: %w", err)
|
||||
}
|
||||
slog.Info("task edit event processed", "task_id", event.Payload.TaskID, "task_revision", snapshot.Task.TaskRevision, "status", status)
|
||||
return nil
|
||||
}
|
||||
if event.Payload.Action == "start" || event.Payload.Action == "resume" {
|
||||
if policy != "" {
|
||||
return errors.New("start/resume control must not carry an active-call policy")
|
||||
|
||||
@@ -4,6 +4,7 @@ import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
@@ -253,6 +254,58 @@ func TestResumeUsesApprovedSIPButPendingRevisionStillClosesAdmission(t *testing.
|
||||
}
|
||||
}
|
||||
|
||||
func TestEditRefreshesOnlyFutureCallsAndPreservesPause(t *testing.T) {
|
||||
controller, agent, s := newControlFixture(t)
|
||||
old, err := s.ReadSnapshot(controller.DispatcherID, 1001, "task-asr")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
revision := 2
|
||||
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/task/task-asr":
|
||||
_, _ = w.Write([]byte(strings.Replace(string(configExample(t, "config-read-task-asr")), `"task_revision":1`, fmt.Sprintf(`"task_revision":%d`, revision), 1)))
|
||||
case "/internal/v1/dispatcher/tenant/1001/quota":
|
||||
_, _ = w.Write(configExample(t, "config-read-quota"))
|
||||
default:
|
||||
http.NotFound(w, r)
|
||||
}
|
||||
}))
|
||||
defer server.Close()
|
||||
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, "edit-1", "edit", "")); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
fresh, err := s.ReadSnapshot(controller.DispatcherID, 1001, "task-asr")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if old.Task.TaskRevision != 1 || fresh.Task.TaskRevision != 2 || len(agent.calls) != 0 {
|
||||
t.Fatalf("edit did not preserve old snapshot and refresh future calls: old=%d new=%d calls=%d err=%v", old.Task.TaskRevision, fresh.Task.TaskRevision, len(agent.calls), err)
|
||||
}
|
||||
if err := controller.ProcessControl(context.Background(), controlBody(t, "pause-1", "pause", "drain")); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
revision = 3
|
||||
if err := controller.ProcessControl(context.Background(), controlBody(t, "edit-paused", "edit", "")); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
fresh, err = s.ReadSnapshot(controller.DispatcherID, 1001, "task-asr")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if fresh.Task.TaskRevision != 3 || len(agent.calls) != 1 {
|
||||
t.Fatalf("edit disturbed paused task: revision=%d calls=%d err=%v", fresh.Task.TaskRevision, len(agent.calls), err)
|
||||
}
|
||||
if admitted, err := s.CanAdmit(controller.DispatcherID, 1001, "task-asr"); err != nil || admitted {
|
||||
t.Fatalf("edit reopened paused task: %v %v", admitted, err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestStartReadsNewTaskWithoutAgentControl(t *testing.T) {
|
||||
controller, agent, s := newControlFixture(t)
|
||||
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
|
||||
@@ -79,6 +79,31 @@ func executeBody(t *testing.T, eventID, callee string) []byte {
|
||||
return []byte(body)
|
||||
}
|
||||
|
||||
func TestEditLeavesInFlightCallBoundToOldSnapshot(t *testing.T) {
|
||||
controller, originator, _, s := newExecuteFixture(t)
|
||||
if err := controller.ProcessExecute(context.Background(), executeBody(t, "before-edit", "15003164745")); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(originator.calls) != 1 || originator.calls[0].Snapshot.Task.TaskRevision != 1 {
|
||||
t.Fatalf("first call did not receive original task snapshot: %+v", originator.calls)
|
||||
}
|
||||
updated, err := s.ReadSnapshot(controller.DispatcherID, 1001, "task-asr")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
updated.Task.Raw = []byte(strings.Replace(string(updated.Task.Raw), `"task_revision":1`, `"task_revision":2`, 1))
|
||||
updated.Task.TaskRevision = 2
|
||||
if status, err := s.CompleteEdit(updated, "edit-between-calls"); err != nil || status != "applied" {
|
||||
t.Fatalf("edit status=%q: %v", status, err)
|
||||
}
|
||||
if err := controller.ProcessExecute(context.Background(), executeBody(t, "after-edit", "15830461047")); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(originator.calls) != 2 || originator.calls[0].Snapshot.Task.TaskRevision != 1 || originator.calls[1].Snapshot.Task.TaskRevision != 2 {
|
||||
t.Fatalf("edit changed an issued call or failed to update future admission: calls=%+v", originator.calls)
|
||||
}
|
||||
}
|
||||
|
||||
func TestExecuteRejectsInvalidCalleeWithoutStoppingTask(t *testing.T) {
|
||||
controller, originator, publisher, s := newExecuteFixture(t)
|
||||
if err := controller.ProcessExecute(context.Background(), executeBody(t, "bad-1", "not-a-phone")); err != nil {
|
||||
|
||||
@@ -46,6 +46,51 @@ func (s *Store) CompleteStart(snapshot configread.Snapshot, eventID string) erro
|
||||
return nil
|
||||
}
|
||||
|
||||
// CompleteEdit advances only a previously assigned task's configuration
|
||||
// revision and its acknowledgment. Pause/stop barriers and running calls stay
|
||||
// untouched; they are never inferred from a configuration edit.
|
||||
func (s *Store) CompleteEdit(snapshot configread.Snapshot, eventID string) (string, error) {
|
||||
task := snapshot.Task
|
||||
tx, err := s.db.Begin()
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
defer tx.Rollback()
|
||||
var state, status string
|
||||
var revision int64
|
||||
var present int
|
||||
err = tx.QueryRow(`SELECT control_state,status,task_revision,present FROM dispatcher_tasks
|
||||
WHERE dispatcher_id=? AND tenant_id=? AND task_id=?`, task.DispatcherID, task.TenantID, task.TaskID).Scan(&state, &status, &revision, &present)
|
||||
if err != nil && !errors.Is(err, sql.ErrNoRows) {
|
||||
return "", fmt.Errorf("read edit task state: %w", err)
|
||||
}
|
||||
if errors.Is(err, sql.ErrNoRows) || present != 1 || state == "stopped" || state == "stopping" || status == "stopped" || task.Status == "stopped" {
|
||||
if err := enqueueControlAck(tx, task.DispatcherID, task.TenantID, eventID, "rejected"); err != nil {
|
||||
return "", err
|
||||
}
|
||||
if err := tx.Commit(); err != nil {
|
||||
return "", fmt.Errorf("commit rejected edit acknowledgment: %w", err)
|
||||
}
|
||||
return "rejected", nil
|
||||
}
|
||||
if task.TaskRevision < revision {
|
||||
return "", fmt.Errorf("edit task %q revision regressed from %d to %d", task.TaskID, revision, task.TaskRevision)
|
||||
}
|
||||
if err := saveSnapshot(tx, snapshot); err != nil {
|
||||
return "", fmt.Errorf("bind edited task snapshot: %w", err)
|
||||
}
|
||||
if _, err := tx.Exec(`UPDATE dispatcher_tasks SET task_revision=? WHERE dispatcher_id=? AND tenant_id=? AND task_id=?`, task.TaskRevision, task.DispatcherID, task.TenantID, task.TaskID); err != nil {
|
||||
return "", fmt.Errorf("persist edited task revision: %w", 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 edited task and acknowledgment: %w", err)
|
||||
}
|
||||
return "applied", 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
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
package store
|
||||
|
||||
import (
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
@@ -46,4 +47,63 @@ func TestStopSuppressesUnstartedCallsButPreservesUnknown(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestEditCommitsSnapshotRevisionAndAckTogether(t *testing.T) {
|
||||
s := preparedCurrentCallStore(t)
|
||||
updated, err := s.ReadSnapshot(currentDispatcherID, 1001, "task-asr")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
updated.Task.Raw = []byte(strings.Replace(string(updated.Task.Raw), `"task_revision":1`, `"task_revision":2`, 1))
|
||||
updated.Task.TaskRevision = 2
|
||||
if status, err := s.CompleteEdit(updated, "edit-success"); err != nil || status != "applied" {
|
||||
t.Fatalf("edit status=%q: %v", status, err)
|
||||
}
|
||||
persisted, err := s.ReadSnapshot(currentDispatcherID, 1001, "task-asr")
|
||||
if err != nil || persisted.Task.TaskRevision != 2 {
|
||||
t.Fatalf("edit not durably bound: revision=%d err=%v", persisted.Task.TaskRevision, err)
|
||||
}
|
||||
var revision int64
|
||||
if err := s.db.QueryRow(`SELECT task_revision FROM dispatcher_tasks WHERE task_id='task-asr'`).Scan(&revision); err != nil || revision != 2 {
|
||||
t.Fatalf("task and snapshot revisions differ: revision=%d err=%v", revision, err)
|
||||
}
|
||||
stale := currentStoreSnapshot(t)
|
||||
if _, err := s.CompleteEdit(stale, "edit-stale"); err == nil {
|
||||
t.Fatal("stale edit downgraded the task")
|
||||
}
|
||||
conflicting := updated
|
||||
conflicting.Task.Raw = []byte(strings.Replace(string(updated.Task.Raw), `"status":"running"`, `"status":"paused"`, 1))
|
||||
if _, err := s.CompleteEdit(conflicting, "edit-conflicting"); err == nil {
|
||||
t.Fatal("same revision changed task contents")
|
||||
}
|
||||
persisted, err = s.ReadSnapshot(currentDispatcherID, 1001, "task-asr")
|
||||
if err != nil || persisted.Task.TaskRevision != 2 || string(persisted.Task.Raw) != string(updated.Task.Raw) {
|
||||
t.Fatalf("failed edit changed bound snapshot: revision=%d err=%v", persisted.Task.TaskRevision, err)
|
||||
}
|
||||
if err := s.ApplyControl(currentDispatcherID, 1001, "task-asr", "stop"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
updated.Task.Raw = []byte(strings.Replace(string(updated.Task.Raw), `"task_revision":2`, `"task_revision":3`, 1))
|
||||
updated.Task.TaskRevision = 3
|
||||
if status, err := s.CompleteEdit(updated, "edit-stopped"); err != nil || status != "rejected" {
|
||||
t.Fatalf("stopped edit status=%q: %v", status, err)
|
||||
}
|
||||
outbox, err := s.ListPendingOutbox(currentDispatcherID)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
foundRejected := false
|
||||
for _, item := range outbox {
|
||||
if strings.Contains(string(item.Body), `"event_id":"edit-stopped"`) && strings.Contains(string(item.Body), `"status":"rejected"`) {
|
||||
foundRejected = true
|
||||
}
|
||||
}
|
||||
if !foundRejected {
|
||||
t.Fatal("edit on stopped task lacked durable rejection receipt")
|
||||
}
|
||||
persisted, err = s.ReadSnapshot(currentDispatcherID, 1001, "task-asr")
|
||||
if err != nil || persisted.Task.TaskRevision != 2 {
|
||||
t.Fatalf("stopped edit modified task: revision=%d err=%v", persisted.Task.TaskRevision, err)
|
||||
}
|
||||
}
|
||||
|
||||
func currentMondayUTC() time.Time { return time.Date(2026, 9, 21, 1, 30, 0, 0, time.UTC) }
|
||||
|
||||
+14
-7
@@ -369,6 +369,18 @@ func (s *Store) CanAdmit(dispatcherID string, tenantID int64, taskID string) (bo
|
||||
// SaveSnapshot refuses an immutable task revision with different content. Only
|
||||
// the providers actually referenced by this task are included in its digest.
|
||||
func (s *Store) SaveSnapshot(snapshot configread.Snapshot) error {
|
||||
tx, err := s.db.Begin()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer tx.Rollback()
|
||||
if err := saveSnapshot(tx, snapshot); err != nil {
|
||||
return err
|
||||
}
|
||||
return tx.Commit()
|
||||
}
|
||||
|
||||
func saveSnapshot(tx *sql.Tx, snapshot configread.Snapshot) error {
|
||||
task := snapshot.Task
|
||||
if task.DispatcherID == "" || task.TenantID <= 0 || task.TaskID == "" || task.TaskRevision <= 0 || snapshot.SIP.DispatcherID != task.DispatcherID || snapshot.Quota.DispatcherID != task.DispatcherID || snapshot.Quota.TenantID != task.TenantID || snapshot.SIP.Revision <= 0 || snapshot.Quota.QuotaRevision <= 0 || !json.Valid(task.Raw) {
|
||||
return errors.New("invalid current task snapshot identity, revision, or body")
|
||||
@@ -408,11 +420,6 @@ func (s *Store) SaveSnapshot(snapshot configread.Snapshot) error {
|
||||
}
|
||||
sum := sha256.Sum256(binding)
|
||||
digest := hex.EncodeToString(sum[:])
|
||||
tx, err := s.db.Begin()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer tx.Rollback()
|
||||
if err := validateGlobalRevisions(tx, snapshot); err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -433,7 +440,7 @@ func (s *Store) SaveSnapshot(snapshot configread.Snapshot) error {
|
||||
if _, err := tx.Exec(`UPDATE dispatcher_configs SET sip_revision=?,quota_revision=?,snapshot_json=? WHERE dispatcher_id=? AND tenant_id=? AND task_id=?`, snapshot.SIP.Revision, snapshot.Quota.QuotaRevision, full, task.DispatcherID, task.TenantID, task.TaskID); err != nil {
|
||||
return fmt.Errorf("refresh global SIP/quota for task %q: %w", task.TaskID, err)
|
||||
}
|
||||
return tx.Commit()
|
||||
return nil
|
||||
}
|
||||
}
|
||||
if _, err = tx.Exec(`INSERT INTO dispatcher_configs(dispatcher_id,tenant_id,task_id,task_revision,sip_revision,quota_revision,content_sha256,snapshot_json)
|
||||
@@ -442,7 +449,7 @@ func (s *Store) SaveSnapshot(snapshot configread.Snapshot) error {
|
||||
content_sha256=excluded.content_sha256,snapshot_json=excluded.snapshot_json`, task.DispatcherID, task.TenantID, task.TaskID, task.TaskRevision, snapshot.SIP.Revision, snapshot.Quota.QuotaRevision, digest, full); err != nil {
|
||||
return fmt.Errorf("persist task snapshot %q: %w", task.TaskID, err)
|
||||
}
|
||||
return tx.Commit()
|
||||
return nil
|
||||
}
|
||||
|
||||
// validateGlobalRevisions prevents two task bindings from silently
|
||||
|
||||
Reference in New Issue
Block a user