Trigger task configuration refresh on start and resume controls
This commit is contained in:
@@ -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、完整对话或音频。
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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"}}
|
||||
@@ -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"
|
||||
}
|
||||
|
||||
@@ -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"]}}}
|
||||
}}
|
||||
}
|
||||
|
||||
@@ -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。
|
||||
|
||||
@@ -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` 已归档。
|
||||
|
||||
## 可复现的本地门禁与来源
|
||||
|
||||
@@ -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` 不再启用。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) |
|
||||
|
||||
|
||||
@@ -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 == "" {
|
||||
|
||||
@@ -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")
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user