Require MQ resume after discovery paused-to-running
This commit is contained in:
@@ -167,7 +167,7 @@
|
||||
- **本轮项目内目标(尚未替换真实外部运行)**:任务(含智能体)、SIP、任务发现和按 tenant_id 的租户额度走四条只读 HTTP;呼叫、控制及其必要回执/单份最终结果走 MQ,取消对外业务查询/补传和分散通话事件。不提供配置HTTP→MQ回退,也不恢复其它业务HTTP通道。SaaS须分发 management 已批准的唯一 SIP 版本,management 仍为唯一编辑/审批面。该本地语义按第三方契约及版本化 Schema/示例/hash 冻结,Mock C 是本地门禁,不等待外部签收;真实切换仍需另行授权。
|
||||
- 新HTTP配置字段阶段草案见 `docs/contracts/config-read-fields-v0.1-proposal.md` 及同目录 `config-read-v0.1.schema.json`/mock示例;截图只证实UI含义,英文响应键为项目自定义,绝非SaaS现网接口已确认字段。用户新增任务排除日期、线路时段等未见截图项按项目需求设计;SIP传输/鉴权/注册及额度未知不能猜默认值。Schema/Mock 校验可满足项目内 C,但不代表 SaaS/management 已发布或真实兼容;真实响应、审批来源、摘要和 Agent/Asterisk 实际加载仍须单独验证。当前唯一权威运行契约不因草案变化。
|
||||
- **每个Dispatcher必须有独立、全局唯一且不重复的ID及独立接收Topic/队列**;指定D的任务/现行MQ配置结果/上传结果不能由其它D抢收,也不能广播后仅靠正文过滤;新目标只读HTTP配置由该D UUID+SECRETKEY获取且须核对任务归属。身份与tenant/Agent/Cell ID、dispatcher_epoch分开;现行合同保留租户独立队列及原值tenant_key,完整新路由长度预算须重验。具体ID生成/持久化、Topic/绑定、消息字段/关联/错误/期限须随W01新版本冻结,不凭本文给旧严格Schema添加字段。
|
||||
- **本轮项目内任务队列所有权硬边界(见第三方契约 §1.2 与 `mq-topology-v0.1-proposal.json`):任务队列和 RabbitMQ 绑定只能由 SaaS 创建、维护、退役;Dispatcher 只消费,不能自行建队、绑定或删除;F03 v0.3 已本地验证,外部兼容性未验证。** SaaS 须先确认持久任务/控制队列与绑定就绪,再向归属 D+任务 ID 的队列发布 persistent 命令;D 离线时已存在队列可积压,缺队列而未入队的原消息由 SaaS 保留并按原身份重发。现行项目内任务发现从 `GET /internal/v1/dispatcher/tasks?after=0` 开始,以持久十进制事件游标逐页读取同一 `tasks[]`(每页最多 256、每 D 活跃任务最多 256),后续约每 30 秒从已提交游标续读;同任务更新可再次出现,`status=removed` 为墓碑,空页保持游标并表示追平。**无 410/游标过期重建、`changes` 兼容或配置回退;不得用旧 v0.2 checkpoint 解释事件 ID。** SaaS 必须保证分页完整有序、序号及撤销墓碑长期可读;累计墓碑不受活跃任务上限约束,容量和灾备连续性尚待外部签收。每页任务/归属/响应/游标同一 SQLite 事务提交;首次未追平、请求/响应/持久化错误均持久关新准入、停任务消费,但 MQ 即时控制和结果恢复照常处理,不自动拨号或重置游标。轮询周期不是端到端发现 SLA,也不能取代 MQ 控制。pause 保留积压,resume 重新核验后消费原积压且不延长期限;stop 持久生效后未接纳旧消息**静默消费/ACK,不拨号、不回逐条结果**,不删队列;控制本身及在途通话仍有处理/最终结果,本地错误/计数不静默。任务带 tenant_id/原值 tenant_key,同 D 各任务共用按 tenant_id 取得的租户额度;额度 0/缺失/过期拒新,降额不强挂或清未知占用,跨 D 仍需权威份额。已移除 `--tenant-key` 及自建租户队列,按 D+task ID 推导 SaaS 预建队列;本地 RabbitMQ 无 `configure` 权限测试通过。v0.1/v0.2 发现证据仅作历史,v0.3 本地证据见 `docs/evidence/f09-local-acceptance-v0.3.md`;不证明真实 SaaS/management、多 D 业务或生产可用。
|
||||
- **本轮项目内任务队列所有权硬边界(见第三方契约 §1.2 与 `mq-topology-v0.1-proposal.json`):任务队列和 RabbitMQ 绑定只能由 SaaS 创建、维护、退役;Dispatcher 只消费,不能自行建队、绑定或删除;F03 v0.3 已本地验证,外部兼容性未验证。** SaaS 须先确认持久任务/控制队列与绑定就绪,再向归属 D+任务 ID 的队列发布 persistent 命令;D 离线时已存在队列可积压,缺队列而未入队的原消息由 SaaS 保留并按原身份重发。现行项目内任务发现从 `GET /internal/v1/dispatcher/tasks?after=0` 开始,以持久十进制事件游标逐页读取同一 `tasks[]`(每页最多 256、每 D 活跃任务最多 256),后续约每 30 秒从已提交游标续读;同任务更新可再次出现,`status=removed` 为墓碑,空页保持游标并表示追平。**无 410/游标过期重建、`changes` 兼容或配置回退;不得用旧 v0.2 checkpoint 解释事件 ID。** SaaS 必须保证分页完整有序、序号及撤销墓碑长期可读;累计墓碑不受活跃任务上限约束,容量和灾备连续性尚待外部签收。每页任务/归属/响应/游标同一 SQLite 事务提交;首次未追平、请求/响应/持久化错误均持久关新准入、停任务消费,但 MQ 即时控制和结果恢复照常处理,不自动拨号或重置游标。轮询周期不是端到端发现 SLA,也不能取代 MQ 控制。发现页的较新 `running` 只更新任务状态和游标,不解除已持久的 paused;pause 保留积压,只有 MQ resume 重新读取单任务并确认 running 后才恢复原积压且不延长期限;stop 持久生效后未接纳旧消息**静默消费/ACK,不拨号、不回逐条结果**,不删队列;控制本身及在途通话仍有处理/最终结果,本地错误/计数不静默。任务带 tenant_id/原值 tenant_key,同 D 各任务共用按 tenant_id 取得的租户额度;额度 0/缺失/过期拒新,降额不强挂或清未知占用,跨 D 仍需权威份额。已移除 `--tenant-key` 及自建租户队列,按 D+task ID 推导 SaaS 预建队列;本地 RabbitMQ 无 `configure` 权限测试通过。v0.1/v0.2 发现证据仅作历史,v0.3 本地证据见 `docs/evidence/f09-local-acceptance-v0.3.md`;不证明真实 SaaS/management、多 D 业务或生产可用。
|
||||
- **OSS相关配置存于Dispatcher配置文件,Agent向Dispatcher领取临时上传TOKEN后直传OSS,不保存长期凭据;SaaS不再下发OSS配置/TOKEN。** D复用官方SDK提供受限TOKEN/目标信息,配置缺失/无效明确失败;不在样例、源码、日志或证据中保存实际密钥/完整TOKEN。过期只允许A显式向D重新申请,不自动续期或向SaaS申请TOKEN;精确配置格式/TOKEN形态/UploadGrant映射另行核验,不猜字段。
|
||||
- **本轮项目内目标:上传仅负责 Agent 直传及 D 可靠通知 MQ,SaaS 后续处理不属本项目职责;外部 v2 的现行上传事实不因本地目标自动改变。** D保留签发能力、不转发文件;R13持久保存事实与recording.uploaded outbox,消息为persistent、进入指定durable队列/绑定、mandatory无return且publisher confirm成功后才记交付完成。只写本地outbox不算入队;不申请SaaS会话、不等待verified/OSS ID、不新增VERIFYING、不伪造SaaS结果。
|
||||
- 新版recording.uploaded取代本项目recording.ready,字段为call_id/recording_id/upload_id/bucket/object_key/format/channels/sample_rate_hz/duration_ms/size_bytes/checksum_sha256;不含TOKEN/密钥/签名URL。MQ失败/确认丢失/重启只恢复原消息身份的交付,不重新PUT或新建资产。现行AI/控制等必要请求响应不受此收缩影响;新目标AI改用任务只读HTTP获取,控制仍走MQ。
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
{
|
||||
"manifest_version": "local-contract-manifest.v0.3",
|
||||
"hash_algorithm": "SHA-256",
|
||||
"source": {"path": "docs/thirds/v0.3.md", "sha256": "2fa8f5a10b1509044b04b4416a44496c5e9a2d9c741e5437862e2edcb92649ce"},
|
||||
"source": {"path": "docs/thirds/v0.3.md", "sha256": "9a22e9ab5da7d859b5fde4b77cb50ddc148e39567704609d633b91a529af564d"},
|
||||
"artifacts": [
|
||||
{"path": "docs/contracts/task-discovery-v0.3-proposal.schema.json", "sha256": "da8eda2e8f5416b2fb35f68e09a98f37d1417a878e9271b8e6d0d6824b94814c"},
|
||||
{"path": "docs/contracts/examples/task-discovery-empty-v0.3.json", "sha256": "4b1b03b75662d52b013cb080ceb13160f7811ab4396f07e908dcf684cf06f032"},
|
||||
|
||||
@@ -4,8 +4,8 @@
|
||||
|
||||
## 本地制品与合同
|
||||
|
||||
- `make release-check-local`:通过;从 Go 1.27.1 的独立项目构建一次性二进制及归档包,检查已存在 release/package 目录、压缩包或校验文件时均拒绝覆盖;校验二进制、归档、checksum/manifest。临时 release SHA-256:`243387b2441d0f88a156e955b821532c6f4246d98669ccbb9a3d91ddb4b2c0ea`。临时文件随测试删除,不冒充已发布制品。
|
||||
- `scripts/build-release.sh` 的 `contract_attestation` 和本地制品检查一致:其他项目内业务协议 `local-contract-manifest.v0.1`,**任务发现 `local-contract-manifest.v0.3`**,MQ 拓扑 `project-saas-dispatcher.v0.1`(v3 队列草案),以及当前 Proto manifest。v0.3 的版本化来源为 [`docs/thirds/v0.3.md`](../thirds/v0.3.md),来源 SHA-256 `2fa8f5a10b1509044b04b4416a44496c5e9a2d9c741e5437862e2edcb92649ce`;[`local-contract-manifest-v0.3.json`](../contracts/local-contract-manifest-v0.3.json) SHA-256 `436de2fe5314deb6a2b773f0217b72c61140db4195f68dcb2f9adb5290b93286`。旧 v0.2 Schema/清单保留追溯但不进入当前运行制品的发现 attestation。
|
||||
- `make release-check-local`:通过;从 Go 1.27.1 的独立项目构建一次性二进制及归档包,检查已存在 release/package 目录、压缩包或校验文件时均拒绝覆盖;校验二进制、归档、checksum/manifest。临时 release SHA-256:`0e250c476d95f6de40ee41932fc4cdb9c2d62324e90125bbbd6558b9b68c877d`。临时文件随测试删除,不冒充已发布制品。
|
||||
- `scripts/build-release.sh` 的 `contract_attestation` 和本地制品检查一致:其他项目内业务协议 `local-contract-manifest.v0.1`,**任务发现 `local-contract-manifest.v0.3`**,MQ 拓扑 `project-saas-dispatcher.v0.1`(v3 队列草案),以及当前 Proto manifest。v0.3 的版本化来源为 [`docs/thirds/v0.3.md`](../thirds/v0.3.md),来源 SHA-256 `9a22e9ab5da7d859b5fde4b77cb50ddc148e39567704609d633b91a529af564d`;[`local-contract-manifest-v0.3.json`](../contracts/local-contract-manifest-v0.3.json) SHA-256 `f32b114b998883901baddd7c4e0225a2643bf5ab68ce54ace5dfa12c543eea93`。旧 v0.2 Schema/清单保留追溯但不进入当前运行制品的发现 attestation。
|
||||
- `scripts/check-contracts.sh` 验证 v0.3 严格 Schema、8 个正例、2 个反例与 13 个来源/产物摘要;旧上游 `contracts/upstream/v1/` 没有更改。`./scripts/acceptance-local.sh`(含格式、`go test -race ./...`、`go vet ./...`、构建及受限 RabbitMQ)通过,具体矩阵与覆盖率见 [`f09-local-acceptance-v0.3.md`](f09-local-acceptance-v0.3.md)。
|
||||
- 制品已直接启动测试 `dispatcher` 的 `mixed` 和 `real` 两种模式:均返回隔离 Mock 限制错误,在创建 SQLite 文件前拒绝;旧 `--tenant-key`、`--consume` 开关不在制品中。隔离 RabbitMQ 测试账号无 `configure` 权限,Dispatcher 不建队或绑定,真实模式没有因 v0.3 合同放行。
|
||||
- 本地构建来源有未提交变更,制品仅标记 `local-development`、`production_approval=false`;Go/Proto/合同摘要可核对不等于发布、运行环境或真实配置已验收。
|
||||
|
||||
@@ -9,16 +9,16 @@
|
||||
| 同一合同及分页 | 首次与后续均读取 `tasks[]`,`after` 是规范十进制事件游标;首次多页、空页保持游标、同任务后续更新、`removed` 墓碑均通过。响应不含逐项 `event_id`;请求中的 `page_token`、旧 `changes`/410 路径不运行,也无 HTTP→MQ 回退。见 `internal/configread/discovery_v03_test.go`、`internal/dispatcher/task_runtime_v03_test.go`。 |
|
||||
| 并发翻页、页中断与重启 | Mock 在翻页时加入同任务更新事件,D 追至最终空页并应用较新状态;第二页 HTTP 失败时仅持久保存首个成功页,关新准入;关闭并重新打开 SQLite 后从已提交的游标续读,不从零重建,也不跳过页。见 `internal/dispatcher/task_runtime_v03_recovery_test.go`、`task_runtime_v03_test.go`。 |
|
||||
| 故障与事务 | 非规范/倒退游标、异常响应及 HTTP 410 均报错,不伪造空页;发现服务中断时停止旧任务消费者但保留 MQ 控制。SQLite 触发器模拟分配写失败后,任务、游标和响应证据全部回滚;未物理填满磁盘。v0.2 checkpoint 无法直接当事件 ID 使用,阻断静默切换。见 `internal/configread/discovery_v03_test.go`、`internal/store/local_v03_discovery_test.go`、`local_v03_failure_test.go`、`internal/dispatcher/task_runtime_v3_test.go`。 |
|
||||
| 控制与重复投递 | 发现页中的 `running` 不解除 MQ pause/stop 或已撤销墓碑;pause 保留积压、stop 对未接纳命令静默 ACK,在途终结仍出最终结果。重复执行恢复原回执,持久一次性拨号决定不产生第二次呼叫。见 `internal/store/local_v03_failure_test.go`、`execute_business_duplicate_test.go`、`local_origination_test.go`、`internal/dispatcher/task_control_v3_test.go`、`task_queue_v3_stop_active_test.go`。 |
|
||||
| 控制与重复投递 | 发现页先将任务从 `running` 改为 `paused`,后续更新为 `running`:SaaS 状态与游标已推进,但已持久的暂停不自动重开原队列;只有 MQ `resume` 读取单任务并确认新鲜 `running` 后准入与原队列才恢复。MQ pause/stop 及 removed 墓碑不因发现页运行状态而解除;pause 保留积压、stop 对未接纳命令静默 ACK,在途终结仍出最终结果。重复执行恢复原回执,持久一次性拨号决定不产生第二次呼叫。见 `internal/store/local_v03_discovery_test.go` 的 `TestEventPageRunningAfterDiscoveryPauseWaitsForMQResume`、`internal/dispatcher/task_discovery_resume_test.go` 的 `TestTaskDiscoveryRunningAfterPauseRequiresMQResumeToRestartQueue`、`TestEventRuntimeHTTPPausedToRunningNeedsMQResume`,以及 `internal/store/local_v03_failure_test.go`、`execute_business_duplicate_test.go`、`local_origination_test.go`、`internal/dispatcher/task_control_v3_test.go`、`task_queue_v3_stop_active_test.go`。 |
|
||||
| SaaS 建队与 MQ/最终结果 | 受限 Dispatcher RabbitMQ 用户的 `configure` 权限为空,只消费 SaaS Mock 预建任务/控制队列;队列缺失或满、确认丢失与重投不冒充已接纳。HTTP→SQLite→MQ→Mock 最终结果与 outbox 恢复集成通过。见 `internal/mq/amqp_v3_integration_test.go`、`amqp_v3_queue_full_integration_test.go`、`internal/dispatcher/local_v01_integration_test.go`。 |
|
||||
| 既有门禁 | 号码白名单、时段、配额、超时、录音结果与最后拨号判断继续使用 F04/F08 本地测试;V3 mixed/real 入口在资源连接前拒绝,Mock 不发真实呼叫。 |
|
||||
|
||||
## 执行记录与覆盖率
|
||||
|
||||
- `./scripts/acceptance-local.sh`:通过。依次核对 gofmt、Proto/本地 Schema 正反例和 SHA-256 清单、`go mod verify`、`go test -race ./... -count=1`、`go vet ./...`、构建及受限 RabbitMQ 集成测试。首次执行发现一个旧集成夹具仍使用 v0.2 全量响应;改为 v0.3 `after=0` → `after=1` 并在载入任务配置前追平后复测通过。
|
||||
- `scripts/check-contracts.sh`:v0.3 正例 8、反例 2、清单文件 13 均通过;旧 v0.2 清单仅作为历史,不参与运行。`docs/contracts/local-contract-manifest-v0.3.json` SHA-256 为 `436de2fe5314deb6a2b773f0217b72c61140db4195f68dcb2f9adb5290b93286`;`docs/thirds/v0.3.md` SHA-256 为 `2fa8f5a10b1509044b04b4416a44496c5e9a2d9c741e5437862e2edcb92649ce`。
|
||||
- `scripts/check-contracts.sh`:v0.3 正例 8、反例 2、清单文件 13 均通过;旧 v0.2 清单仅作为历史,不参与运行。`docs/contracts/local-contract-manifest-v0.3.json` SHA-256 为 `f32b114b998883901baddd7c4e0225a2643bf5ab68ce54ace5dfa12c543eea93`;`docs/thirds/v0.3.md` SHA-256 为 `9a22e9ab5da7d859b5fde4b77cb50ddc148e39567704609d633b91a529af564d`。
|
||||
- `go test -cover -count=1`:P1 业务包 `agent` 71.1%、`callwindow` 84.4%、`configread` 75.2%、`contract` 74.0%、`dispatcher` 70.5%、`rpc` 69.9%、`store` 67.4%、`tenant` 85.4%;RabbitMQ integration 下 `mq` 66.4%、`dispatcher` 70.6%。本轮 P1 业务包均 ≥65%。`cmd/sip-go-agent` 30.0%、历史 `internal/callruntime` 61.7% 单列技术债,不冒称全仓库达标;`mq` 未带 integration tag 的结果不用于验收。
|
||||
- 工具链为 Go 1.27.1;RabbitMQ 本地镜像 `rabbitmq:4.1-management-alpine`,受限账号无 `configure` 权限。`git diff --check` 通过;测试时工作树未提交,不把当前产物宣称为可复现生产制品。
|
||||
- 工具链为 Go 1.27.1;RabbitMQ 本地镜像 `rabbitmq:4.1-management-alpine`,受限账号无 `configure` 权限。`git diff --check` 通过;本地功能只提交到开发分支,未合并或推送 `main`,不把 Mock 产物宣称为生产制品。
|
||||
|
||||
## 未通过的外部门禁
|
||||
|
||||
|
||||
+1
-1
@@ -22,7 +22,7 @@
|
||||
- 每次 HTTP 200 都使用 `tasks[]` 和 `next_cursor`,不再区分 `snapshot`、`changes`。`tasks[]` 可以为空;空页表示本轮追平,游标不跳跃。非空页的游标必须前进。项目内 v0.3 已冻结为规范十进制游标、单页最多 256 项、仅响应级 `next_cursor`,不增加 `has_more` 或逐项 `event_id`;必须读到空页,不能凭“本页少于上限”推断结束。
|
||||
- SaaS 对该 D 的任务按**各自最新变更事件序号**升序分页;该序号必须持久、单调且不可复用。同一任务在两次查询之间多次更新,只需返回足以表达**最新状态**的记录,不要求保存每次中间状态。若分页期间再次更新,新的序号必须使该任务可再次被读到;绝不能因游标前进跳过尚未返回的任务。
|
||||
- 普通任务沿用已确认的身份、原值 `tenant_key`、归属和状态字段。`removed` 墓碑必须至少能唯一定位原 D、租户和任务;其余必填规则由项目内 v0.3 Schema 定义;`task_revision` 是任务状态修订,不能代替事件序号,也不从旧 `changes[].operation=removed` 猜字段。已存在任务的租户/任务归属与本地绑定冲突时拒绝处理,不静默迁移。
|
||||
- Dispatcher 对新 `task_id` 建立归属并核验 SaaS 预建队列;对已有任务只按最新 `status` 更新生命周期。`paused`/`stopped` 的本地持久屏障不能被较旧的 `running` 覆盖;`removed` 阻止新接纳,已有在途执行、最终结果和 outbox 按原身份收口,不由轮询触发重复拨号或擅自强挂。`status` 不变时可推进游标,但不重复执行控制动作。任务配置仍通过单独的只读任务接口核验,不让发现列表替代执行快照。
|
||||
- Dispatcher 对新 `task_id` 建立归属并核验 SaaS 预建队列;对已有任务按最新 `status` 更新 SaaS 状态并应用更严格的准入约束。发现页中的 `paused→running`(即使事件较新)只更新状态和游标,不解除已持久的暂停;须收到 MQ `resume`、重新读取单任务并确认 `running`,才恢复原队列。`stopped` 不可逆;`removed` 阻止新接纳,已有在途执行、最终结果和 outbox 按原身份收口,不由轮询触发重复拨号或擅自强挂。`status` 不变时可推进游标,但不重复执行控制动作。任务配置仍通过单独的只读任务接口核验,不让发现列表替代执行快照。
|
||||
|
||||
`tasks[]`、`next_cursor` 的项目内精确字段与正反例已按 [`v0.3 Schema`](contracts/task-discovery-v0.3-proposal.schema.json) 冻结并在 Mock 验证;**不是 SaaS 现网或对外正式合同**。逐项不带 `event_id`,Dispatcher 不能独立证明 SaaS 页内排序/没有漏项;完整有序返回、长期墓碑及授权绑定仍需 SaaS/业务另行签收。
|
||||
|
||||
|
||||
@@ -49,7 +49,7 @@ SaaS 先持久修改权威任务状态,再向独立 D 控制队列发布 `task
|
||||
1. **pause:**立即关闭新接纳并保存暂停屏障,停止从该任务队列继续取新执行;已投递但未接纳的有界消息退回原队列,不能 ACK 丢弃或搬进无界 SQLite 待拨队列。已有执行按 drain/hangup 处理。重启仍暂停。
|
||||
2. **resume:**重新 GET 单任务配置(不使用约 60 秒旧缓存),只有 SaaS 最新状态为 running、归属/授权/租户额度/时段有效、且本地未 stopped,才解除暂停并恢复消费原积压。无需 SaaS 重新发原消息;已过 `not_after` 的旧消息仍不可拨,暂停不能冻结或延长有效期。过期/窗口外非 stopped 命令给明确拒绝回执,不等待次日自动拨号。
|
||||
3. **stop:**持久化不可逆停止屏障,同任务 ID 不能 resume/重新启用;SaaS 停止新发布。D 不申请通话额度,继续小批量消费未接纳积压,仅计本地处置计数并 ACK,不建外呼结果/outbox,不拨号。已停止任务重投/重启、额度为 0、配置缓存失效时仍可排空;不得借此清掉别的任务、伪造已接纳执行终态或漏掉控制回执/既有通话结果。
|
||||
4. **优先级:**本地 stopped 永久高于任何 running;本地 paused 只能由有效 resume 解除。`tasks[]` 事件页及缓存更新可使状态更严格,不能清除已持久的暂停/停止屏障,且低于已知 `task_revision` 的状态不得覆盖新状态。冷启动恢复本地屏障并从已提交游标追平 SaaS 事件页后取更严格者;冲突/缺失关闭新准入,不丢已有执行事实。
|
||||
4. **优先级:**本地 stopped 永久高于任何 running;本地 paused 只能由有效 resume 解除。`tasks[]` 事件页及缓存更新可使状态更严格;即使较新的发现页把 `paused` 改为 `running`,也只更新 SaaS 状态和游标,不清除已持久的暂停/停止屏障。暂停只由上述 MQ `resume` 配合新鲜任务读取解除,且低于已知 `task_revision` 的状态不得覆盖新状态。冷启动恢复本地屏障并从已提交游标追平 SaaS 事件页后取更严格者;冲突/缺失关闭新准入,不丢已有执行事实。
|
||||
5. **无编号乱序:**控制重投可能重复返回回执,不保证按消息身份“只处理一次”。pause/stop 先关闭准入;与 SaaS 最新状态不符、不可核验或接收乱序时保持关闭并返回明确失败,不根据到达先后自动恢复。恢复必须重新发送有效 resume 并核对最新权威状态;不引入替代 command_id 或控制 CAS。SaaS 不可把发布成功视为控制已应用,丢失回执下状态不明不得主动扩量。
|
||||
|
||||
### 3.4 任务发现、发布与期限
|
||||
@@ -129,7 +129,7 @@ SaaS 清单须保留已 stopped 但未排空任务;撤销握手按第三方契
|
||||
| F07/F03 v0.3 | 项目内 `tasks[]` 事件游标页及来源/hash 冻结;Dispatcher 与 SQLite 按页原子提交,首次追平前和异常页拒绝新执行;同任务更新、撤销墓碑、断页重启从已提交游标续读、控制竞态与无重拨已完成 Mock。旧 `changes`/410 运行路径移除,SaaS 队列仍独占建队。 | 证据:[`f09-local-acceptance-v0.3.md`](evidence/f09-local-acceptance-v0.3.md)。外部 SaaS 有序返回、长期墓碑容量、旧状态受控切换和真实兼容仍未签收。 |
|
||||
| F04 | 本地隔离 Mock 已将任务×已选线路时段、任务排除日期、白名单/有效期、绑定快照及任务/AI 较小通话时限接入接纳与发起前门禁;执行身份只有一次持久 `issued`/`refused` 决定,队列消费后由 D 经版本化 Unary 向 Mock Agent 下发,Agent 不重算业务策略、不发送 SIP。V3 mixed/real 启动拒绝,原真实固定门禁未放宽。 | 证据:[`f04-local-dial-policy-v0.1.md`](evidence/f04-local-dial-policy-v0.1.md)。Proto/合同、本地 Mock、race、vet、build、RabbitMQ 无 `configure` 权限及 `scripts/acceptance-local.sh` 通过;外部合同/真实拨号未验证。F08 后续本地结果见 [`f08-local-final-result-v0.1.md`](evidence/f08-local-final-result-v0.1.md);F09 本地矩阵与覆盖率见 [`f09-local-acceptance-v0.2.md`](evidence/f09-local-acceptance-v0.2.md)。 |
|
||||
| F08 | 单份最终 `call.result`、确认终态后独立释放额度、uploaded/unavailable/not_created、本地 Agent→D 录音失败事实与原身份 outbox 已完成隔离 Mock 验证;V3 Mock 不并行写旧分散事件,V3 mixed/real 启动拒绝。证据:[`f08-local-final-result-v0.1.md`](evidence/f08-local-final-result-v0.1.md)。 | `scripts/acceptance-local.sh`(含 race、vet、构建及受限 RabbitMQ)通过;Mock 录音成功/失败由合成事实注入,默认 Mock 只模拟无应答,无实际媒体/OSS PUT,publisher confirm 不等于 SaaS 应用签收。F09 本地验证见 [`f09-local-acceptance-v0.2.md`](evidence/f09-local-acceptance-v0.2.md);F06 真实部署/切换仍未获授权。 |
|
||||
| F09 | v0.3 首次多页、空页、同任务重现及墓碑、翻页中更新、断页重启、坏游标/410 失败、SQLite 写失败同事务回滚、MQ 控制竞态、重复执行不重拨完成单节点/单租户本地 Mock 验证;旧 F09 v0.2 证据仅作历史。证据:[`f09-local-acceptance-v0.3.md`](evidence/f09-local-acceptance-v0.3.md)。 | `scripts/acceptance-local.sh`(格式/合同/Proto/hash、race、vet、构建、受限 RabbitMQ)通过;P1 业务包 66.4%–85.4%(MQ 含 integration)。CLI 30.0% 和历史 callruntime 61.7% 单列技术债,不声称全仓达标;SQLite 写失败用触发器注入,未物理填盘。外部有序返回/墓碑容量/真实验收未完成。 |
|
||||
| F09 | v0.3 首次多页、空页、同任务重现及墓碑、翻页中更新、断页重启、坏游标/410 失败、SQLite 写失败同事务回滚、MQ 控制竞态(含发现 `paused→running` 仅更新状态/游标、须 MQ `resume`+新鲜任务核验才恢复队列)、重复执行不重拨完成单节点/单租户本地 Mock 验证;旧 F09 v0.2 证据仅作历史。证据:[`f09-local-acceptance-v0.3.md`](evidence/f09-local-acceptance-v0.3.md)。 | `scripts/acceptance-local.sh`(格式/合同/Proto/hash、race、vet、构建、受限 RabbitMQ)通过;P1 业务包 66.4%–85.4%(MQ 含 integration)。CLI 30.0% 和历史 callruntime 61.7% 单列技术债,不声称全仓达标;SQLite 写失败用触发器注入,未物理填盘。外部有序返回/墓碑容量/真实验收未完成。 |
|
||||
| F06 | Go 1.27.1 一次性本地制品/校验和将 v0.1 其他业务+v0.3 任务发现+v3 MQ 拓扑及 Proto 版本/hash 写入清单;已有 release/package/归档拒绝覆盖,制品 mixed/real 在开 DB 前拒绝;旧 CLI 和分散事件不在新运行路径。证据:[`f06-local-release-gates-v0.3.md`](evidence/f06-local-release-gates-v0.3.md)。 | `make release-check-local` 与本地 acceptance 通过;来源 dirty、`production_approval=false`。外部 SaaS 有序返回/长期墓碑容量、v0.2 旧状态受控切换、真实部署/拨号与生产签收均**未完成**;本地阻断不代签真实切换。 |
|
||||
| F05 | 多D/第二租户不在本轮运行范围。 | 另获阶段与额度/资源授权。 |
|
||||
|
||||
|
||||
+4
-4
@@ -1,12 +1,12 @@
|
||||
# 第三方任务发现事件游标分页 v0.3(项目内提案)
|
||||
|
||||
> 仅替换 [`v0.2`](v0.2.md) §2.5 的任务发现目标;其他 SaaS 项目内业务接口仍按 v0.1。设计依据为 [`plan-0926.md`](../plan-0926.md)。本文件、Schema 和 Mock **未经 SaaS/业务签收,不是现网接口或生产合同**;新版未完成切换前 v0.2 仍是本地运行依据。
|
||||
> 仅替换 [`v0.2`](v0.2.md) §2.5 的任务发现目标;其他 SaaS 项目内业务接口仍按 v0.1。设计依据为 [`plan-0926.md`](../plan-0926.md)。本文件、Schema 和 Mock **未经 SaaS/业务签收,不是现网接口或生产合同**;当前本地仅运行 v0.3;v0.2 仅作历史,真实 SaaS 尚未切换。
|
||||
|
||||
## 2.5 唯一任务发现路径
|
||||
|
||||
`GET /internal/v1/dispatcher/tasks?after=<cursor>` 由指定 Dispatcher 使用既有 `X-DISPATCHER-id` / `X-DISPATCHER-SECRET-KEY` 读取自己的任务。首次无本协议游标时发送 `after=0`;之后只发送已和任务状态一起持久化的 `next_cursor`。`after` 是**该 D 范围内单调、不可复用的事件 ID**,不是任务 ID、页号或任务配置版本。规范形式是无前导零的十进制 uint64 字符串;只允许起点为 `0`。D 身份不可更换来绕过游标。
|
||||
|
||||
所有 HTTP 200 均采用同一形状:`schema_version=task-discovery.v0.3-proposal`、`dispatcher_id`、`tasks[]` 和 `next_cursor`。`tasks[]` 中每项为 `task_id`、`tenant_id`、原值 `tenant_key`、`status`、`task_revision`;**没有** `changes`、`operation`、`snapshot`、`mode`、单项 `event_id`、SaaS 队列名或任务页 token。同一任务有更晚的事件时可再次出现。SaaS 对每个任务只需返回**最新状态**,不要求回放所有中间状态;已有任务只有 `status` 变化才触发生命周期动作,其他字段仍须校验身份、版本及归属。任务具体执行配置仍由独立只读接口取得,不由发现列表取代。
|
||||
所有 HTTP 200 均采用同一形状:`schema_version=task-discovery.v0.3-proposal`、`dispatcher_id`、`tasks[]` 和 `next_cursor`。`tasks[]` 中每项为 `task_id`、`tenant_id`、原值 `tenant_key`、`status`、`task_revision`;**没有** `changes`、`operation`、`snapshot`、`mode`、单项 `event_id`、SaaS 队列名或任务页 token。同一任务有更晚的事件时可再次出现。SaaS 对每个任务只需返回**最新状态**,不要求回放所有中间状态;已有任务的 `status` 变化更新 SaaS 状态;新 `paused`、`stopped`、`finished`、`removed` 可关闭或保持准入。**后续 `running`(即使更新且已持久化)只更新 SaaS 状态与游标,不直接解除已持久的 `paused`;必须收到 MQ `resume` 并由任务只读接口新鲜确认 `running` 才能重开原队列。**其他字段仍须校验身份、版本及归属。任务具体执行配置仍由独立只读接口取得,不由发现列表取代。
|
||||
|
||||
SaaS 在响应中按其事件序号升序选取尚未返回的**最新任务记录**。每页至多 256 项,活动任务总量仍受单 D 256 上限约束;撤销墓碑的累计数量不在这个活动上限之内。非空页的 `next_cursor` 必须是**最后一个实际返回任务**的事件 ID,并严格大于请求的 `after`,绝不能前进到尚未返回的更新之后。空页表示本轮追平,`tasks=[]` 且 `next_cursor=after`;D 可在首次追平后开放经核验的新执行,其后约 30 秒再次查询。不得依靠“不足 256 项”断定追平:无论每页实际数量,D 都必须继续读取直到空页。若页内同一任务重复出现、任务身份发生冲突或游标不前进,D 拒绝整页而不是挑一条使用。
|
||||
|
||||
@@ -14,7 +14,7 @@ SaaS 在响应中按其事件序号升序选取尚未返回的**最新任务记
|
||||
|
||||
### 状态和撤销
|
||||
|
||||
`status` 只允许 `running`、`paused`、`stopped`、`finished`、`removed`。`removed` 是该 D 归属撤销的**任务墓碑**,必须保留原任务/租户身份;某任务没有出现在本页不表示撤销。SaaS 先持久化 stop/撤销及其必要回执、协调自己拥有的任务队列退役,再发出墓碑;D 收到后停止新接纳,不清理执行恢复、已入队结果、幂等或未知占用,也不自动强挂在途通话。paused 保留原积压;MQ 控制即时执行,不等轮询,也不允许较旧 `running` 解开已持久的 paused/stopped 屏障。
|
||||
`status` 只允许 `running`、`paused`、`stopped`、`finished`、`removed`。`removed` 是该 D 归属撤销的**任务墓碑**,必须保留原任务/租户身份;某任务没有出现在本页不表示撤销。SaaS 先持久化 stop/撤销及其必要回执、协调自己拥有的任务队列退役,再发出墓碑;D 收到后停止新接纳,不清理执行恢复、已入队结果、幂等或未知占用,也不自动强挂在途通话。paused 保留原积压;MQ 控制即时执行,不等轮询。发现页无论新旧 `running` 均不解除已持久的 paused/stopped 屏障;stopped 同任务不可逆,removed 不撤销旧执行事实。
|
||||
|
||||
### 一致性、错误和恢复
|
||||
|
||||
@@ -25,4 +25,4 @@ SaaS 在响应中按其事件序号升序选取尚未返回的**最新任务记
|
||||
|
||||
## 验收及容量阻断
|
||||
|
||||
本地 Mock 要覆盖初始多页、空页、同一任务再出现、撤销、跨页并发更新、重启原游标续读、重复页、事务回滚、暂停/停止竞态和 queue-only 消费;所有关键边界 fail-closed,不重拨。实时控制与最终结果的其他 v0.1 业务合同不变。活动任务 256 上限**不约束永久墓碑**;历史撤销增长、离线追赶时间、SaaS 数据保留及灾备后序号连续性未签收,不能声称容量有界或真实联调/生产可用。
|
||||
本地 Mock 要覆盖初始多页、空页、同一任务再出现、撤销、跨页并发更新、重启原游标续读、重复页、事务回滚、发现页 `paused→running` 后准入仍暂停直到 MQ `resume`+新鲜任务核验、暂停/停止竞态和 queue-only 消费;所有关键边界 fail-closed,不重拨。实时控制与最终结果的其他 v0.1 业务合同不变。活动任务 256 上限**不约束永久墓碑**;历史撤销增长、离线追赶时间、SaaS 数据保留及灾备后序号连续性未签收,不能声称容量有界或真实联调/生产可用。
|
||||
|
||||
@@ -0,0 +1,133 @@
|
||||
package dispatcher
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"git.ipao.vip/rogee/go-sip/internal/store"
|
||||
)
|
||||
|
||||
func TestTaskDiscoveryRunningAfterPauseRequiresMQResumeToRestartQueue(t *testing.T) {
|
||||
now := time.Date(2026, 9, 22, 10, 0, 0, 0, time.UTC)
|
||||
d, st, server := newLocalV01TestDispatcher(t, now)
|
||||
defer server.Close()
|
||||
defer st.Close()
|
||||
|
||||
for _, update := range []struct {
|
||||
after, next, status string
|
||||
revision int64
|
||||
}{
|
||||
{after: "1", next: "2", status: "paused", revision: 3},
|
||||
{after: "2", next: "3", status: "running", revision: 4},
|
||||
} {
|
||||
body, err := json.Marshal(map[string]any{
|
||||
"schema_version": "task-discovery.v0.3-proposal",
|
||||
"dispatcher_id": localTestDispatcherID,
|
||||
"next_cursor": update.next,
|
||||
"tasks": []map[string]any{{
|
||||
"task_id": localTestTaskID, "tenant_id": localTestTenantID, "tenant_key": localTestTenantKey,
|
||||
"status": update.status, "task_revision": update.revision,
|
||||
}},
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := st.ApplyLocalTaskDiscoveryPage(store.LocalTaskDiscoveryPage{
|
||||
DispatcherID: localTestDispatcherID, FromCursor: update.after, NextCursor: update.next,
|
||||
Tasks: []store.LocalDiscoveredTask{{TaskID: localTestTaskID, TenantID: localTestTenantID, TenantKey: localTestTenantKey,
|
||||
Status: update.status, TaskRevision: update.revision}}, Body: body, ObservedAt: now,
|
||||
}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
empty, err := json.Marshal(map[string]any{
|
||||
"schema_version": "task-discovery.v0.3-proposal", "dispatcher_id": localTestDispatcherID,
|
||||
"next_cursor": "3", "tasks": []any{},
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := st.ApplyLocalTaskDiscoveryPage(store.LocalTaskDiscoveryPage{
|
||||
DispatcherID: localTestDispatcherID, FromCursor: "3", NextCursor: "3", Body: empty, ObservedAt: now,
|
||||
}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
assignment, err := st.LocalTaskAssignment(localTestDispatcherID, localTestTaskID)
|
||||
if err != nil || assignment.Status != "running" || assignment.TaskRevision != 4 || assignment.AdmissionState != "paused" || taskAssignmentCanConsume(assignment) {
|
||||
t.Fatalf("new discovery status reopened paused queue without MQ resume: %+v err=%v", assignment, err)
|
||||
}
|
||||
cursor, exists, err := st.LocalTaskDiscoveryCursor(localTestDispatcherID)
|
||||
if err != nil || !exists || cursor != "3" {
|
||||
t.Fatalf("discovery update was not durably applied: cursor=%q exists=%v err=%v", cursor, exists, err)
|
||||
}
|
||||
reader := &fakeTaskControlStatusReader{status: localControlStatus("running", 4)}
|
||||
queues := &fakeTaskQueueController{}
|
||||
processor := newLocalTaskControlProcessor(d, reader, queues)
|
||||
if err := processor.Handle(context.Background(), localControlRoutingKey(), localTaskControlBody(t, now, "resume")); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
assignment, err = st.LocalTaskAssignment(localTestDispatcherID, localTestTaskID)
|
||||
if err != nil || assignment.Status != "running" || assignment.AdmissionState != "running" || !taskAssignmentCanConsume(assignment) {
|
||||
t.Fatalf("MQ resume failed to reopen the original task queue: %+v err=%v", assignment, err)
|
||||
}
|
||||
if reader.calls != 1 || len(queues.events) != 1 || queues.events[0] != "start" {
|
||||
t.Fatalf("resume did not verify fresh status and restart once: status_calls=%d queue_events=%v", reader.calls, queues.events)
|
||||
}
|
||||
assertTaskControlReceipt(t, st, "resume", "applied", "applied", "running")
|
||||
}
|
||||
|
||||
func TestEventRuntimeHTTPPausedToRunningNeedsMQResume(t *testing.T) {
|
||||
page := func(status string, revision int, cursor string) []byte {
|
||||
return []byte(fmt.Sprintf(`{"schema_version":"task-discovery.v0.3-proposal","dispatcher_id":%q,"tasks":[{"task_id":%q,"tenant_id":%q,"tenant_key":%q,"status":%q,"task_revision":%d}],"next_cursor":%q}`,
|
||||
localTestDispatcherID, localTestTaskID, localTestTenantID, localTestTenantKey, status, revision, cursor))
|
||||
}
|
||||
pages := map[string][]byte{
|
||||
"0": page("paused", 1, "1"),
|
||||
"1": page("running", 2, "2"),
|
||||
"2": []byte(fmt.Sprintf(`{"schema_version":"task-discovery.v0.3-proposal","dispatcher_id":%q,"tasks":[],"next_cursor":"2"}`, localTestDispatcherID)),
|
||||
}
|
||||
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
body, ok := pages[r.URL.Query().Get("after")]
|
||||
if r.URL.Path != "/internal/v1/dispatcher/tasks" || !ok {
|
||||
w.WriteHeader(http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
_, _ = w.Write(body)
|
||||
}))
|
||||
defer server.Close()
|
||||
runtime, st := newEventRuntime(t, server)
|
||||
broker := &fakeQueueBroker{}
|
||||
runtime.queues = newTaskQueueController(runtime.dispatcher, broker, nil)
|
||||
defer runtime.queues.StopAll(context.Background())
|
||||
if err := runtime.refreshDiscoveryState(context.Background()); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
assignment, err := st.LocalTaskAssignment(localTestDispatcherID, localTestTaskID)
|
||||
if err != nil || assignment.Status != "running" || assignment.TaskRevision != 2 || assignment.AdmissionState != "paused" {
|
||||
t.Fatalf("HTTP discovery did not retain paused admission: %+v err=%v", assignment, err)
|
||||
}
|
||||
if len(broker.started) != 0 {
|
||||
t.Fatalf("HTTP discovery started a paused task queue: %v", broker.started)
|
||||
}
|
||||
cursor, exists, err := st.LocalTaskDiscoveryCursor(localTestDispatcherID)
|
||||
if err != nil || !exists || cursor != "2" {
|
||||
t.Fatalf("HTTP discovery failed to commit status and cursor: cursor=%q exists=%v err=%v", cursor, exists, err)
|
||||
}
|
||||
reader := &fakeTaskControlStatusReader{status: localControlStatus("running", 2)}
|
||||
processor := newLocalTaskControlProcessor(runtime.dispatcher, reader, runtime.queues)
|
||||
if err := processor.Handle(context.Background(), localControlRoutingKey(), localTaskControlBody(t, time.Now().UTC(), "resume")); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
assignment, err = st.LocalTaskAssignment(localTestDispatcherID, localTestTaskID)
|
||||
if err != nil || assignment.AdmissionState != "running" || len(broker.started) != 1 || broker.started[0] != assignment.Queue.QueueName || reader.calls != 1 {
|
||||
t.Fatalf("MQ resume did not verify and restart the same queue: %+v started=%v calls=%d err=%v", assignment, broker.started, reader.calls, err)
|
||||
}
|
||||
assertTaskControlReceipt(t, st, "resume", "applied", "applied", "running")
|
||||
}
|
||||
@@ -143,6 +143,7 @@ func mergeLocalTaskState(current, incoming string) string {
|
||||
}
|
||||
switch incoming {
|
||||
case "running":
|
||||
// Discovery records the new SaaS status, but only MQ resume with a fresh task read can clear a persisted pause.
|
||||
return current
|
||||
case "paused", "stopped", "finished":
|
||||
return incoming
|
||||
|
||||
@@ -77,6 +77,32 @@ func TestEventPageAllowsSameTaskNewStatusAndRemovedTombstone(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestEventPageRunningAfterDiscoveryPauseWaitsForMQResume(t *testing.T) {
|
||||
st := openLocalDiscoveryTestStore(t)
|
||||
for _, page := range []LocalTaskDiscoveryPage{
|
||||
eventPage("0", "1", eventTask("task-a", "paused", 1)),
|
||||
eventPage("1", "1"),
|
||||
eventPage("1", "2", eventTask("task-a", "running", 2)),
|
||||
eventPage("2", "2"),
|
||||
} {
|
||||
if err := st.ApplyLocalTaskDiscoveryPage(page); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
assignment, err := st.LocalTaskAssignment("d-1", "task-a")
|
||||
if err != nil || assignment.Status != "running" || assignment.TaskRevision != 2 || assignment.AdmissionState != "paused" {
|
||||
t.Fatalf("discovery running incorrectly reopened admission: %+v err=%v", assignment, err)
|
||||
}
|
||||
cursor, exists, err := st.LocalTaskDiscoveryCursor("d-1")
|
||||
if err != nil || !exists || cursor != "2" {
|
||||
t.Fatalf("latest task status did not advance cursor: %q exists=%v err=%v", cursor, exists, err)
|
||||
}
|
||||
assignment, err = st.ResumeLocalTaskAdmission("d-1", "task-a", "tenant-a", "tenant-key-a", "running", 2)
|
||||
if err != nil || assignment.Status != "running" || assignment.AdmissionState != "running" {
|
||||
t.Fatalf("fresh MQ resume did not reopen admission: %+v err=%v", assignment, err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestEventPageCannotInterpretOldSnapshotCursor(t *testing.T) {
|
||||
st := openLocalDiscoveryTestStore(t)
|
||||
if _, err := st.db.Exec(`INSERT INTO local_v02_task_discovery_state(dispatcher_id,cursor,last_mode,ready,updated_at) VALUES('d-1','opaque-old','snapshot',1,'now')`); err != nil {
|
||||
|
||||
Reference in New Issue
Block a user