From cff0be1cce912ee92a9e9871677415eafddc0b96 Mon Sep 17 00:00:00 2001 From: Rogee Date: Mon, 21 Sep 2026 10:04:07 +0800 Subject: [PATCH] docs: add SaaS dispatcher agent contracts --- docs/contracts/dispatcher-agent.md | 362 +++++++++++++++++++++++++++++ docs/contracts/saas-dispatcher.md | 344 +++++++++++++++++++++++++++ 2 files changed, 706 insertions(+) create mode 100644 docs/contracts/dispatcher-agent.md create mode 100644 docs/contracts/saas-dispatcher.md diff --git a/docs/contracts/dispatcher-agent.md b/docs/contracts/dispatcher-agent.md new file mode 100644 index 0000000..ad5cdb8 --- /dev/null +++ b/docs/contracts/dispatcher-agent.md @@ -0,0 +1,362 @@ +# Dispatcher ↔ Agent 对接契约(实现事实) + +## 1. 适用范围与实现边界 + +当前边界是项目自有的 Unary gRPC `agent.v1.AgentControlService`。消息只携带控制数据、执行绑定、事实和上传元信息,不经过 Dispatcher 传输音频字节。 + +| 方向/角色 | 当前实现 | +| --- | --- | +| Dispatcher → Agent | `internal/rpc.Server` 提供 Agent 侧服务;`internal/dispatcher.AgentCoordinator` 通过配置的 endpoint 调用 | +| Agent → Dispatcher | `internal/rpc.DispatcherServer` 提供 Dispatcher 侧同名服务,仅接收事件和上传 RPC | +| 传输 | gRPC Unary + mTLS;`internal/rpc.Client` 不做业务自动重试 | +| 业务权威 | Dispatcher SQLite 管理任务、配额、事实、outbox;Agent 只维护本地会话、执行/文件恢复事实 | +| 数据面 | Agent 按 Dispatcher 下发的 grant 直接 PUT 到 OSS;Dispatcher 不接收或转发录音内容 | + +精确字段号和枚举以 `proto/agent/v1/agent.proto` 为唯一源;本文不另造 protobuf。 + +## 2. Service 方法与方向 + +| RPC | 方向 | 当前代码状态 | 接收端 | +| --- | --- | --- | --- | +| `GetAgentStatus` | Dispatcher → Agent | 已由 `AgentCoordinator.Probe` 调用 | Agent `rpc.Server` | +| `ActivateAgent` | Dispatcher → Agent | 已由 `AgentCoordinator.Activate` 调用 | Agent `rpc.Server` | +| `GetBootstrap` | Dispatcher → Agent | Agent 侧已实现;当前启动绑定流程未调用 | Agent `rpc.Server` | +| `SetAdmissionState` | Dispatcher → Agent | Agent 侧已实现;当前 `AgentCoordinator` 没有调用封装 | Agent `rpc.Server` | +| `Execute` | Dispatcher → Agent | 已调用;当前 RPC handler 只准备并记录执行状态 | Agent `rpc.Server` | +| `GetExecutionPermit` | Dispatcher → Agent | 已由 `ExecuteRaw` 调用 | Agent `rpc.Server` | +| `ApplyTaskControl` | Dispatcher → Agent | 已由 `AgentCoordinator.Control` 调用 | Agent `rpc.Server` | +| `QueryExecution` | Dispatcher → Agent | 用于超时/响应丢失后的对账 | Agent `rpc.Server` | +| `ReportExecutionEvent` | Agent → Dispatcher | 已接收、去重并生成 MQ outbox | Dispatcher `rpc.DispatcherEventServer` | +| `RequestUpload` | Agent → Dispatcher | 已接收并签发 OSS grant | Dispatcher `rpc.DispatcherUploadServer` | +| `CompleteUpload` | Agent → Dispatcher | 已验证 OSS 对象并生成 `recording.ready` outbox | Dispatcher `rpc.DispatcherUploadServer` | + +`DispatcherServer` 对外只实现 `ReportExecutionEvent`、`RequestUpload`、`CompleteUpload`;其它 RPC 在 Dispatcher listener 上返回 `UNIMPLEMENTED`。Agent `rpc.Server` 虽实现完整 generated service,但其 upload handler 在非 `mock` 模式明确返回 `UNIMPLEMENTED`。 + +## 3. 连接、认证与会话 + +### 3.1 连接 + +1. Dispatcher 从受控 Agent endpoint 文件读取 `agent_id`、`cell_id`、地址和 `server_name`。 +2. `rpc.DialFromFiles` 使用 CA、客户端证书、私钥和 server name 建立 TLS gRPC 连接。 +3. Dispatcher 对每个 endpoint 先 `GetAgentStatus`,再以返回的 `boot_id` 调 `ActivateAgent`。 +4. Dispatcher 生成本次 `dispatcher_epoch`;Agent 用 `session_generation` 持久化 fencing 水位。 +5. Agent 新会话会 fence 旧的 `agent_id + cell_id + boot_id + epoch + generation` 组合;旧请求返回 `ABORTED`,不会自动释放未知执行。 + +### 3.2 mTLS 与身份 + +- Agent listener 要求 verified peer certificate;可按证书 fingerprint 和允许的 Agent ID 限制。 +- Dispatcher listener 同样要求 verified mTLS peer、fingerprint allowlist 和可选的 `AllowedAgentIDs`。 +- 共用证书只证明证书组;`agent_id`、`cell_id`、`boot_id` 必须经过 Dispatcher 激活绑定,不能信任 Agent 自报 endpoint。 +- gRPC RPC 失败不等于业务失败。尤其 `Execute`、permit 和上传完成超时后,调用方必须查询原执行/上传状态,不能换 execution ID、attempt ID 或 upload ID。 + +## 4. 公共消息结构 + +### 4.1 `RequestMeta` + +| 字段 | 类型 | 用途 | +| --- | --- | --- | +| `protocol_version` | string | 当前实现发送 `agent.v1` | +| `request_id` | string | 单次 RPC 请求关联 | +| `trace_id` | string | 跨模块追踪 | +| `operation_id` | string | 业务操作标识;写操作必填 | +| `deadline_unix_ms` | int64 | 协议字段;当前 handler 未单独执行该值的截止检查 | +| `dispatcher_epoch` | string | 当前 Dispatcher 会话代次 | +| `agent_id` / `cell_id` | string | endpoint 绑定身份 | +| `boot_id` | string | Agent 进程启动身份 | +| `session_generation` | uint64 | 会话 fencing 代次 | +| `idempotency_key` | string | 同操作重报必须复用;写操作必填 | + +### 4.2 `ResponseMeta`、`Failure`、`OperationReceipt` + +`ResponseMeta` 回显协议、请求、trace、operation、Dispatcher epoch、Agent/Cell/boot/generation,并增加 `observed_at_unix_ms`。 + +`Failure`: + +| 字段 | 类型 | +| --- | --- | +| `code` | `FailureCode` | +| `retryable` | bool | +| `detail` | string | +| `field` | string | + +`OperationReceipt`:`meta`、`result`、可选 `failure`、`fact_id`、`content_sha256`、`accepted_at_unix_ms`。 + +`ResultCode` 为 `ACCEPTED`、`APPLIED`、`REJECTED`、`UNKNOWN`、`CONFLICT`;`ACCEPTED` 只表示接收/持久记录,不能直接解释为已拨号或已挂断。 + +### 4.3 `ExecutionBinding` + +| 字段 | +| --- | +| `tenant_id`, `tenant_key` | +| `execution_id`, `task_id`, `task_item_id` | +| `task_revision` | +| `call_id`, `attempt_id` | +| `agent_version_id` | +| `route_policy_id`, `caller_profile_id` | + +Dispatcher 从已校验的 `call.execute` 构造 binding;Agent 回报不能改写租户、任务或资产归属。 + +### 4.4 `AssetDescriptor`、`ExecutionFact`、`UploadGrant` + +`AssetDescriptor`: + +| 字段 | 类型/约束 | +| --- | --- | +| `kind` | `RECORDING` 或 `TRANSCRIPT` | +| `asset_id`, `call_id`, `execution_id` | string | +| `format` | string | +| `size_bytes`, `duration_ms` | int64 | +| `checksum_sha256` | string | +| `channels`, `sample_rate_hz` | int32 | + +`ExecutionFact`: + +| 字段 | 类型/约束 | +| --- | --- | +| `fact_id` | 稳定事实 ID,必填 | +| `content_sha256` | 事实内容摘要,必填 | +| `binding` | `ExecutionBinding` | +| `kind` | `FactKind` | +| `observed_at_unix_ms` | Agent 事实发生时间,必填 | +| `source_boot_id` | 事实来源 boot,必填 | +| `source_sequence` | uint64,诊断/排序关联 | +| `payload_json` | JSON object 字节,必填 | + +`UploadGrant`:`upload_id`、`target_url`、`headers[]`、`expires_at_unix_ms`、`object_key`、`required_checksum_sha256`、`max_bytes`。grant 不包含长期 OSS 密钥。 + +## 5. Dispatcher → Agent RPC + +### 5.1 `GetAgentStatus` + +请求:`RequestMeta` + 可选 `AgentBinding target`。激活前仅发送 protocol/request/trace/operation/agent/cell,不能带 session binding。 + +响应:`ResponseMeta` + `AgentStatus`: + +- `agent_id`、`cell_id`、`boot_id`; +- `software_version`、`protocol_version`、`asterisk_version`; +- `admission_state`; +- `capabilities[]`; +- `resources`(CPU、内存、FD、spool、媒体端口及 `sample_fresh`); +- `applied_configs[]`; +- `mtls_authenticated`、`session_active`、`status_reason`。 + +当前 `Probe` 至少校验返回的 Agent/Cell 与目标一致且 `boot_id` 非空。 + +### 5.2 `ActivateAgent` + +请求:`RequestMeta`、`AgentBinding`、`activation_operation_id`、可选 `session_nonce`、`session_expires_at_unix_ms`。 + +`AgentBinding` 由 Dispatcher 提供:`agent_id`、`cell_id`、`expected_boot_id`、`dispatcher_epoch`、`session_generation`、`endpoint_id`。 + +响应:`ResponseMeta`、`state=ACTIVE`、`Session`: + +- `dispatcher_epoch`; +- `session_generation`; +- `expires_at_unix_ms`; +- `session_credential`(bytes,敏感数据)。 + +当前 Agent 实现实际将 session 有效期固定为 `now + 10 分钟`;`session_nonce` 和调用方提供的 `session_expires_at_unix_ms` 当前没有参与校验。后续请求由 mTLS + 会话绑定字段授权,当前 handler 未在每个请求中再次传输或校验 `session_credential`。 + +### 5.3 `GetBootstrap` + +请求:`RequestMeta`、`agent_id`、`cell_id`、`boot_id`、`session_generation`。当前 Agent handler 以 `RequestMeta` 会话授权为准。 + +响应:`ResponseMeta`、`state`、`runtime_configs[]`、`UploadPolicy`: + +- `ConfigReference`:`kind`、`version`、`sha256`、`source`; +- `UploadPolicy`:`enabled`、`max_asset_bytes`、`min_retention_ms`、`allowed_hosts[]`。 + +当前 Dispatcher 启动绑定流程只执行 status + activate,没有调用 bootstrap。 + +### 5.4 `SetAdmissionState` + +请求字段:`meta`、`target`、`state`(`OPEN/CLOSED/DRAINING/QUARANTINED`)、`barrier_id`、`expected_admission_generation`、`reason`。 + +响应:`OperationReceipt` + `applied_admission_generation`。Agent 侧按目标 Agent 做 generation CAS;代次不匹配返回 `CONFLICT`。当前代码没有 `AgentCoordinator` 的调用封装,也没有把此状态操作接到 SaaS HTTP 控制入口。 + +### 5.5 `GetExecutionPermit` + +请求:`meta`、`ExecutionBinding`、`resource_reservation_id`、`expected_task_revision`、`admission_generation`、`config_sha256`。 + +响应:`OperationReceipt` + `ExecutionPermit`: + +| 字段 | 含义 | +| --- | --- | +| `permit_id` | 当前实现为 `permit-` | +| `resource_reservation_id` | 对应 Dispatcher reservation | +| `issued_at_unix_ms` / `expires_at_unix_ms` | 许可时间窗 | +| `dispatcher_epoch` / `session_generation` | fencing 绑定 | +| `fencing_token` | 随机 fencing token | +| `config_sha256` | 执行配置摘要 | + +当前 `rpc.Server` 的 permit 默认有效期为 1 秒;Dispatcher 在收到传输错误时不重新申请另一个 reservation,而是先 `QueryExecution` 对账。 + +### 5.6 `Execute` + +请求字段:`meta`、`ExecutionBinding`、`call_execute_json`、`config_sha256`、`admission_generation`、`resource_reservation_id`、`permit_id`。 + +`call_execute_json` 必须是原始、已通过 `mq.schema.json` 的 `call.execute` body。Agent 当前会再次解码并校验:租户、execution、task、task item、AI version 必须与 binding 一致,并在已提供本地 AI 快照/授权时校验摘要、版本、租户、模式、有效期和 egress pool。 + +成功响应:`OperationReceipt` + `state`。当前 `rpc.Server.Execute` 的可证明副作用是: + +- 写入执行准备日志; +- 写入内存 execution record; +- 同 idempotency key 同内容返回原 receipt;不同内容返回 `CONFLICT`; +- 在非 mock 模式检查 Asia/Shanghai 外呼时间窗口。 + +当前该 RPC handler 本身不证明已调用 ARI/RTP 或已经产生 SIP 外呼。`AgentCoordinator.ExecuteRaw` 当前发送 binding、原始 command JSON、config digest 和 permit ID,不填充 `admission_generation` 与 `resource_reservation_id`。 + +### 5.7 `ApplyTaskControl` + +请求:`meta`、`ExecutionBinding`、`action`(`PAUSE/RESUME/STOP`)、`active_call_policy`(`DRAIN/HANGUP`)、`expected_task_revision`、`reason`。 + +响应:`OperationReceipt`、`applied_task_revision`、`state`。 + +当前 Agent handler: + +- 找不到 execution 返回 `NOT_FOUND`; +- task revision 不一致返回 `CONFLICT`; +- terminal execution 不允许 resume; +- `STOP` 将 call state 置为 `stopped`/terminal; +- `PAUSE`/`RESUME` 更新内存 call state。 + +当前 handler 没有依据 `active_call_policy` 实施实际 drain/hangup,也没有把 `reason` 写入业务事实;真正通话屏障仍需 ARI/执行器事实补齐。 + +### 5.8 `QueryExecution` + +请求:`meta` + `ExecutionBinding`。 + +成功响应包含 `ExecutionSnapshot`:binding、execution state、`call_state`、`attempt_id`、`reason_code`、`observed_at_unix_ms`、`unknown`、`assets[]`。 + +当前 handler 查询 Agent 内存 execution record;未找到返回 `NOT_FOUND`,当前实现不会填充资产列表。Dispatcher 只在原 RPC 结果未知或回包不完整时用原 binding 对账,不能据此自动重拨。 + +## 6. Agent → Dispatcher RPC + +### 6.1 `ReportExecutionEvent` + +请求:`RequestMeta` + `ExecutionFact`。 + +Dispatcher 侧额外要求: + +- verified mTLS peer;可选允许的 Agent ID; +- `meta.agent_id/cell_id/boot_id` 非空; +- operation/idempotency key 非空; +- fact ID、内容摘要、source boot、观测时间和 payload 非空; +- binding 至少包含 `tenant_id`、`tenant_key`、`execution_id`; +- `payload_json` 必须解码为 JSON object。 + +`FactKind` 到 SaaS MQ 事件的当前映射: + +| FactKind | Dispatcher 输出 | aggregate | +| --- | --- | --- | +| `EXECUTION_ACCEPTED` | `command.result` | `command` | +| `CALL_STATUS` | `call.status` | `call` | +| `CALL_FINISHED` | `call.finished` | `call` | +| `TRANSCRIPT_UPDATED` | `transcript.updated` | `transcript_segment` | +| `TRANSCRIPT_FAILED` | 尝试生成 `transcript.failed`,但当前 `aggregate_type=transcript` 不通过 `mq.schema.json` | 当前路径阻塞 | +| `CONTACT_OPT_OUT` | `contact.opt_out` | `call` | +| `RECORDING_PROGRESS` | 不发布 MQ 事件,只保存 fact | `execution_fact` | + +当前 Dispatcher 为有事件的 fact 生成 `event_id = execution-fact-`,并在同一 SQLite 事务内。注意:`TRANSCRIPT_FAILED` 的当前代码映射使用 `aggregate_type=transcript`,而 `mq.schema.json` 只允许 `transcript_segment`;因此该 fact 当前会在事件 Schema 校验阶段失败,不能按成功回报处理。 + +具体流程为: + +1. 以 `fact_id + content_sha256` 去重; +2. 保存 binding、payload、source boot、source sequence、观测时间; +3. 分配 aggregate version; +4. 写入权威事件 outbox。 + +重复 fact 同摘要返回 `ACCEPTED`;同 fact ID 不同摘要返回 `CONFLICT`。Agent 不能通过 payload 或请求字段指定 `aggregate_version`。 + +### 6.2 `RequestUpload` + +请求:`meta`、`ExecutionBinding`、`AssetDescriptor`、`upload_id`。 + +当前 Dispatcher 验证:Agent/Cell/operation/idempotency 元数据、完整 execution/tenant binding、合法 `tenant_key`、asset ID、upload ID、正数文件大小和 SHA-256;超过 OSS 最大文件大小返回 `RESOURCE_EXHAUSTED`。 + +成功响应:`OperationReceipt`、`UploadGrant`、`state`。当前生产 Dispatcher handler: + +- 以 `tenant_key + "\\0" + execution_id + "\\0" + asset_id` 的 SHA-256 hex 生成 object key,并加配置的 key prefix; +- 将 binding、asset、grant、object key、state 持久到 Dispatcher SQLite; +- grant 的过期时间由 OSS client 配置提供; +- 同 upload ID 同 binding/asset 返回原 grant;绑定不同返回冲突; +- 过期 grant 只有在显式再次 `RequestUpload` 时才替换,不自动续期。 + +新请求当前返回 `UPLOAD_STATE_REQUESTED`;持久层状态为 `granted`。`UPLOADING` 枚举存在,但当前 Dispatcher handler 不把 Agent 的 PUT 过程映射为该状态。 + +### 6.3 Agent → OSS 直接上传 + +Agent 获得 grant 后使用 `internal/agent.UploadClient.UploadFile`: + +- 只允许 HTTPS;除非显式配置,否则不允许 HTTP; +- 可限制目标 host;禁止 grant 注入 `Host` 和 `Content-Length`; +- 预先读取文件计算 SHA-256,再按 `Content-Length` PUT; +- 禁止重定向;2xx 才算 PUT 成功; +- 返回 HTTP status、文件大小、SHA-256、ETag; +- 文件内容不经过 Dispatcher,Agent 不把源文件删除或移动。 + +PUT 成功不等于 OSS verified,也不等于 SaaS 已应用。 + +### 6.4 `CompleteUpload` + +请求:`meta`、原 `ExecutionBinding`、原 `AssetDescriptor`、`upload_id`、`uploaded_size_bytes`、`uploaded_checksum_sha256`。 + +Dispatcher 当前执行: + +1. 校验 upload ID 对应的 binding/asset 完全一致; +2. 校验上传大小、SHA-256 和 grant `max_bytes`; +3. 校验 grant 未过期; +4. 调用 OSS client 独立验证 object; +5. 只允许 `AssetKind=RECORDING` 进入 verified recording 路径; +6. 在一个 SQLite 事务内更新 upload 为 completed、保存 `oss_id` 并写入 `recording.ready` outbox。 + +响应为 `OperationReceipt`、`state=COMPLETED`、`oss_id`。同 upload ID 同 OSS ID 的重复 complete 返回 `ACCEPTED`;未知 upload 返回 `NOT_FOUND`;对象校验失败返回 `FAILED_PRECONDITION` 且标记可重试。 + +## 7. 错误、幂等与未知结果 + +### 7.1 gRPC FailureCode + +`INVALID_ARGUMENT`、`UNAUTHENTICATED`、`PERMISSION_DENIED`、`FAILED_PRECONDITION`、`ABORTED`、`RESOURCE_EXHAUSTED`、`UNAVAILABLE`、`DEADLINE_EXCEEDED`、`NOT_FOUND`、`ALREADY_EXISTS`。 + +调用方处理原则: + +- 参数、Schema、绑定错误:拒绝,不改写成另一任务; +- CAS、session、fact、upload 绑定冲突:查询原对象,不换 ID 绕过; +- `UNAVAILABLE`/`DEADLINE_EXCEEDED`:结果可能已发生,先查询/对账; +- `Execute` 结果未知:Dispatcher 标记 reservation/task 为 unknown,禁止自动重新 originate; +- 上传回包丢失:复用原 `upload_id` 和原资产摘要; +- Agent 新 boot:保留旧未知占用,先恢复和对账,不因新 boot 的空状态释放配额。 + +### 7.2 RPC 幂等键 + +Agent 侧写操作的内存 operation key 为: + +`agent_id + "\\0" + operation_id + "\\0" + idempotency_key` + +同 key 同 protobuf 内容返回原 receipt;同 key 不同内容返回 `CONFLICT`。Dispatcher 事件侧使用 `fact_id + content_sha256`,上传侧使用 `upload_id + binding + asset` 绑定。 + +## 8. 当前未实现或未接通的部分 + +1. `GetBootstrap`、`SetAdmissionState` 虽有 handler,但当前 Dispatcher 启动流程没有调用完整 bootstrap/admission 编排。 +2. `Execute` 的当前 RPC 实现只证明准备/幂等/配置校验,不证明 ARI/RTP/SIP 已由该 RPC 直接完成。 +3. Dispatcher → SaaS 的 `recording-uploads` HTTP handshake 未在当前 Go client 中接入;当前完成路径是 Dispatcher 直接验证 OSS。 +4. Dispatcher → SaaS 的 AI version GET 未接入;当前 AI snapshot/authorization 由启动输入提供并在 Agent 侧校验。 +5. 不能把 generated service 中的全量方法数当作每个 listener 都可调用;实际 listener 能力以第 2 节和 `DispatcherServer` 代码为准。 + +## 9. 依据文件 + +- `proto/agent/v1/agent.proto` +- `internal/dispatcher/agent.go` +- `internal/rpc/client.go` +- `internal/rpc/server.go` +- `internal/rpc/dispatcher_server.go` +- `internal/rpc/dispatcher_events.go` +- `internal/rpc/dispatcher_upload.go` +- `internal/agent/upload.go` +- `internal/store/facts.go` +- `internal/store/uploads.go` +- `internal/ai/snapshot.go` +- `internal/ai/authorization.go` +- `cmd/sip-go-agent/main.go` +- `contracts/upstream/2026-09-19-p1-v1/event-payloads.schema.json` +- `contracts/upstream/2026-09-19-p1-v1/ai-authorization.schema.json` diff --git a/docs/contracts/saas-dispatcher.md b/docs/contracts/saas-dispatcher.md new file mode 100644 index 0000000..70780b4 --- /dev/null +++ b/docs/contracts/saas-dispatcher.md @@ -0,0 +1,344 @@ +# SaaS ↔ Dispatcher 对接契约(实现事实) + +## 1. 适用范围与事实等级 + +本文只描述当前 Go 实现能够证明的字段、方向、接口和状态;不把上游 OpenAPI 中尚未接入的 HTTP 客户端写成“已实现”。 + +| 内容 | 当前状态 | 事实来源 | +| --- | --- | --- | +| SaaS → Dispatcher:`call.execute` RabbitMQ 命令 | 已实现 | `internal/dispatcher/consumer.go`、`internal/store/store.go` | +| Dispatcher → SaaS:RabbitMQ 业务事件及 outbox | 已实现 | `internal/contract/contract.go`、`internal/store/facts.go`、`internal/dispatcher/dispatcher.go` | +| SaaS → Dispatcher:控制、命令查询、命令补传 HTTP | 已实现,但只有当前代码列出的行为 | `internal/control/http.go` | +| SaaS → Dispatcher:通话查询、通话补传 HTTP | 路由存在,当前固定返回 `404` | `internal/control/http.go` | +| Dispatcher → SaaS:录音 upload-session/complete HTTP | 上游契约已定义,当前代码没有 SaaS HTTP client | `contracts/upstream/2026-09-19-p1-v1/saas.openapi.yaml` | +| Dispatcher → SaaS:AI version GET | 上游契约已定义,当前代码没有 SaaS HTTP client;当前 Agent 启动时读取本地快照 | `contracts/upstream/2026-09-19-p1-v1/ai-config.openapi.yaml`、`internal/ai/*.go` | + +契约包固定来源为 `contracts.SourceCommit = 2026-09-19-p1-v1`。MQ 与事件正文必须以该目录中的 JSON Schema 为准,不维护第二套手写 Schema。 + +## 2. 通信拓扑与方向 + +| 方向 | 接口 | 传输 | 交付语义 | +| --- | --- | --- | --- | +| SaaS → Dispatcher | `call.execute` | RabbitMQ command exchange | Dispatcher 完成 Schema、租户路由、截止时间和 SQLite 持久化后才 ACK | +| Dispatcher → SaaS | `command.result` 等事件 | RabbitMQ event exchange + SQLite outbox | publisher confirm 只代表 broker 收到;不代表 SaaS 业务事务已应用 | +| SaaS → Dispatcher | 控制/查询/补传 | 内部 HTTP | 控制和补传的 `202` 只代表已接受/持久化,不代表 Agent 已执行 | +| Dispatcher → SaaS | 录音会话申请/完成 | 上游定义的内部 HTTP | 当前项目尚未接入;不能用本地 OSS 结果代替 SaaS `verified` | +| Dispatcher → SaaS | `GET /internal/v1/ai/agent-versions/{agent_version_id}` | 上游定义的内部 HTTP | 当前项目尚未接入;不能读取 `latest` 或用默认配置替代 | + +`tenant_key` 原值复制到消息体、队列、binding 和 routing key;必须是有效 UTF-8,最大 `224` 字节,不清洗、编码或截断。 + +## 3. RabbitMQ 命令:SaaS → Dispatcher + +### 3.1 拓扑 + +| 元素 | 值 | +| --- | --- | +| command exchange | `agent-call.commands.v1`,durable `direct` | +| tenant queue | `agent-call.executor.{tenant_key}.v1`,durable、单一精确 binding | +| routing key | `agent-call.tenant.{tenant_key}.call.execute` | +| event exchange | `agent-call.events.v1`,durable `topic` | +| Dispatcher dead-letter exchange | `agent-call.dead-letter.v1`,durable `topic` | +| 默认 prefetch | `1` | + +Dispatcher 为每个消费租户声明 command queue 和 `.dlq.v1` 队列。SaaS 结果队列由 SaaS 管理,拓扑基线为 `agent-call.saas.events.v1`,binding `agent-call.#`。 + +### 3.2 `call.execute` 外壳 + +消息为 JSON,`additionalProperties: false`: + +| 字段 | 类型 | 必填/约束 | +| --- | --- | --- | +| `schema_version` | string | 必须为 `1.0` | +| `command_type` | string | 必须为 `call.execute` | +| `command_id` | string | 1–128 字节;不可含空白、`/`、`\\` | +| `tenant_id` | string | 同上 | +| `tenant_key` | string | 非空;有效 UTF-8;实现额外限制 224 字节 | +| `trace_id` | string | 同 `id` 约束 | +| `issued_at` | RFC3339 时间 | 必填 | +| `not_after` | RFC3339 时间 | 必填;Dispatcher 接收时已到期则拒绝 | +| `payload` | object | 必须符合 `executePayload` | + +### 3.3 `payload` 数据结构 + +| 字段 | 类型 | 约束 | +| --- | --- | --- | +| `execution_id` | string | 必填 ID | +| `task_id` | string | 必填 ID | +| `task_item_id` | string | 必填 ID | +| `task_revision` | integer | `>= 1` | +| `callee` | string | 1–256 字符;保留业务原始被叫号码 | +| `route_policy_id` | string | 必填 ID | +| `caller_profile_id` | string | 必填 ID | +| `agent_version_id` | string | 必填 ID | +| `variables` | object | 必填;当前 Schema 允许任意附加属性,业务白名单仍由上游约束 | +| `ring_timeout_ms` | integer | `>= 1` | +| `max_call_duration_ms` | integer | `>= 1` | + +Go 实现对应 `internal/contract.ExecutePayload`;`payload` 以 `json.RawMessage` 保留原始 JSON,Dispatcher 不把变量转成另一套协议。 + +### 3.4 接收、幂等与 ACK + +1. `ConsumeTenant` 先校验 `tenant_key`,声明租户队列,再以 `prefetch=1` 消费。 +2. `AcceptCommand` 调用 `Store.IngestCommand`:校验源 Schema、routing key、`not_after`,计算原始 body 的 SHA-256。 +3. 以 `command_id` 查 inbox:同 ID 同 body 为重复;同 ID 不同 body 为冲突并拒绝。 +4. 新命令在一个 SQLite 事务内写入 inbox、task 和 `command.result(accepted)` outbox。 +5. handler 成功后才 ACK。Schema、JSON、租户路由等永久错误 `Reject(false)`,经死信队列处理;其它错误 `Nack(requeue=true)`。 +6. `command_id`、`execution_id`、`task_id` 等业务标识不因重投而更换;`not_after` 不因重投延期。 + +当前 `command.result` 初始 payload 为: + +```json +{ + "command_id": "", + "command_type": "call.execute", + "execution_id": "", + "status": "accepted", + "reason_code": "accepted", + "requested_task_revision": 1 +} +``` + +## 4. RabbitMQ 事件:Dispatcher → SaaS + +### 4.1 通用外壳 + +事件 routing key 固定为 `agent-call.{event_type}`。外壳字段全部必填: + +```json +{ + "schema_version": "1.0", + "event_id": "", + "event_type": "", + "tenant_id": "", + "tenant_key": "<原值>", + "trace_id": "", + "occurred_at": "2026-09-19T00:00:00Z", + "aggregate_type": "", + "aggregate_id": "", + "aggregate_version": 1, + "payload": {} +} +``` + +`event_type` 的源 Schema 枚举为: + +`command.result`、`call.status`、`transcript.updated`、`call.finished`、`recording.ready`、`recording.failed`、`transcript.failed`、`contact.opt_out`。 + +事件 payload 必须同时通过 `mq.schema.json` 和 `event-payloads.schema.json`。`EventBuilder` 不接受未知字段;`aggregate_version` 由 Dispatcher SQLite 按 `aggregate_type + aggregate_id` 递增,Agent 不能指定它。 + +### 4.2 当前代码实际生成的事件 + +| 事件 | 事实来源 | 当前代码行为 | +| --- | --- | --- | +| `command.result` | 命令持久受理、执行接受事实 | 命令接收路径直接生成;Agent `EXECUTION_ACCEPTED` fact 也映射到该类型 | +| `call.status` | Agent `CALL_STATUS` fact | Dispatcher 校验 fact 后生成 | +| `call.finished` | Agent `CALL_FINISHED` fact | Dispatcher 校验 fact 后生成 | +| `transcript.updated` | Agent `TRANSCRIPT_UPDATED` fact | 保留实时文字事件名,不使用 `call.transcript` | +| `transcript.failed` | Agent `TRANSCRIPT_FAILED` fact | 当前代码尝试映射,但使用 `aggregate_type=transcript`;`mq.schema.json` 只允许 `transcript_segment`,因此该路径当前无法通过双 Schema 校验 | +| `contact.opt_out` | Agent `CONTACT_OPT_OUT` fact | Dispatcher 校验 fact 后生成 | +| `recording.ready` | 已验证 OSS 上传完成 | `CompleteUpload` 将完成状态和 outbox 放在同一 SQLite 事务 | +| `recording.failed` | 上游 Schema | 当前 Go 代码没有对应 `FactKind` 或事件生成路径 | + +`RECORDING_PROGRESS` fact 只写入 Dispatcher 事实表,不生成 SaaS MQ 事件。`TRANSCRIPT_FAILED` 虽在两个事件 payload Schema 中声明,但当前 `internal/rpc/dispatcher_events.go` 将其聚合类型写成 `transcript`,随后被 `mq.schema.json` 拒绝;不能把该路径记为已交付。 + +### 4.3 已冻结的专属 payload 关键字段 + +完整约束在 `event-payloads.schema.json`;下面列出对接方必须使用的字段,不是新的 Schema: + +| 事件 | 必填字段 | +| --- | --- | +| `command.result` | `command_id`, `command_type`, `status`, `reason_code` | +| `call.status` | `call_id`, `execution_id`, `call_state`, `call_version`, `attempt_id`, `attempt_state` | +| `transcript.updated` | `call_id`, `turn_id`, `segment_id`, `role`, `revision`, `text`, `is_final`, `start_ms`, `end_ms`, `playback_state` | +| `call.finished` | `call_id`, `execution_id`, `call_version`, `outcome`, `started_at`, `ended_at`, `duration_ms`, `reason_code` | +| `recording.ready` | `call_id`, `recording_id`, `oss_id`, `format`, `channels`, `sample_rate_hz`, `duration_ms`, `size_bytes`, `checksum_sha256` | +| `recording.failed` | `call_id`, `recording_id`, `stage`, `reason_code`, `retryable` | +| `transcript.failed` | `call_id`, `reason_code`, `retryable` | +| `contact.opt_out` | `call_id`, `task_id`, `task_item_id`, `requested_at` | + +`recording.ready` 只能表示对象已验证;不能用 PUT 成功、ETag 或本地路径替代 `oss_id` 和 SHA-256。 + +### 4.4 Outbox 交付 + +`Store.RecordExecutionFact`、命令接收和上传完成均可在状态事务中写 outbox。`Dispatcher.FlushOutbox`: + +- claim `pending/retry` 记录并标记 `dispatching`; +- 通过 RabbitMQ persistent JSON message 发布并等待 publisher confirm; +- 成功标记 `published`;失败标记 `retry`; +- 重启时将 `dispatching` 恢复为 `retry`。 + +重复 fact 使用相同 `fact_id + content_sha256` 时不创建第二条事件;同 fact ID 不同摘要返回冲突。SaaS 必须以 `event_id` 做 inbox 幂等,broker confirm 不能当作 SaaS 应用收讫。 + +## 5. Dispatcher 提供的 HTTP 接口:SaaS → Dispatcher + +### 5.1 公共请求头 + +当前 `internal/control.Handler` 要求: + +- `Authorization: Bearer `;配置了 token 时必须精确匹配; +- `X-Tenant-ID`; +- `X-Request-ID`。 + +当前 Handler 未强制检查权限 scope。源 `executor.openapi.yaml` 另外要求 `Idempotency-Key`;控制接口当前代码只把该值传入存储层,没有把缺失值直接拒绝,这属于实现与源 OpenAPI 的已知差异。 + +错误响应为 JSON: + +```json +{ + "type": "about:blank", + "title": "", + "status": 400, + "code": "", + "detail": "", + "request_id": "", + "retryable": false +} +``` + +### 5.2 控制任务 + +**方向:SaaS → Dispatcher** + +`POST /internal/v1/outbound/tasks/{task_id}/controls` + +请求体: + +| 字段 | 类型/约束 | +| --- | --- | +| `command_id` | 必填 ID;当前代码用于响应,不单独生成控制 outbox | +| `action` | `pause`、`resume`、`stop` | +| `expected_task_revision` | integer,`>=1`;CAS 版本 | +| `active_call_policy` | 可选:`drain` 或 `hangup` | +| `reason` | 非空,最大 512 字符 | + +当前成功响应为 `202`: + +```json +{ + "command_id": "", + "tenant_id": "", + "tenant_key": "", + "task_id": "", + "status": "accepted", + "requested_task_revision": 1, + "accepted_at": "2026-09-19T00:00:00Z" +} +``` + +实现先按 `X-Tenant-ID + task_id` 查询任务,再以任务的 `execution_id` 调用 `ApplyControlDetailed`。版本冲突为 `409`;任务不存在为 `404`;输入错误为 `400`。`202` 不表示 Agent 已应用控制。 + +### 5.3 查询命令 + +**方向:SaaS → Dispatcher** + +`GET /internal/v1/outbound/commands/{command_id}` + +当前成功响应字段: + +```json +{ + "command_id": "", + "command_type": "call.execute", + "tenant_id": "", + "tenant_key": "", + "task_id": "", + "execution_id": "", + "call_id": null, + "status": "persisted", + "reason_code": null, + "wait_reason_code": null, + "accepted_at": "", + "waiting_since": null, + "admission_deadline": "", + "requested_task_revision": 1, + "applied_task_revision": null, + "task_state": null, + "aggregate_version": 1, + "updated_at": "" +} +``` + +查询按租户隔离;不存在返回 `404`。当前实现不返回 call snapshot、控制应用版本或完整状态机。 + +### 5.4 按 source command 补传 + +**方向:SaaS → Dispatcher** + +`POST /internal/v1/outbound/commands/{source_command_id}/replays` + +请求体: + +```json +{"command_id":"","reason":"<1..512 chars>"} +``` + +必须提供 `Idempotency-Key`。成功响应为 `202`: + +```json +{"command_id":"","status":"accepted","snapshot_cutoff":""} +``` + +当前实现读取原 inbox body,保留原 `command_id` 和原消息内容,按原租户 routing key 写入 outbox;`replay-` 只是 outbox event ID,不是新的业务 command ID。不存在 source command 返回 `404`。原消息的 `not_after` 不会被改写,重发后仍可能因过期被拒绝。 + +### 5.5 当前不可用的通话接口 + +以下路由已匹配,但当前固定返回 `404`: + +- `GET /internal/v1/outbound/calls/{call_id}`; +- `POST /internal/v1/outbound/calls/{call_id}/replays`。 + +不得依据上游 OpenAPI 的 `Call` 结构宣称当前代码已经提供通话查询或通话补传。 + +## 6. Dispatcher → SaaS 的已声明、未接入 HTTP + +### 6.1 录音 upload session / complete + +权威源为 `contracts/upstream/2026-09-19-p1-v1/saas.openapi.yaml`。 + +| 方向 | 方法 | 路径 | +| --- | --- | --- | +| Dispatcher → SaaS | `POST` | `/internal/v1/outbound/recording-uploads` | +| Dispatcher → SaaS | `POST` | `/internal/v1/outbound/recording-uploads/{upload_id}/complete` | + +公共要求:Bearer、`X-Tenant-ID`、`X-Request-ID`、`Idempotency-Key`。 + +申请请求的实际源字段为 `recording_id`、`call_id`、`content_type=audio/wav`、`size_bytes`、`checksum_algorithm=SHA-256`、`checksum`、`channels=1`、`sample_rate_hz`、`duration_ms`。返回 `upload_id`、`recording_id`、`expires_at`、`upload_method=PUT`、`upload_url`、`required_headers`、`constraints` 和可空 `oss_id`。 + +complete 请求为 `recording_id`、`size_bytes`、`checksum_algorithm=SHA-256`、`checksum`、可空 `etag`;成功返回 `status=verified`、`oss_id`、`verified_at`。 + +当前 Go 代码没有调用这两个 SaaS HTTP 路径。当前 Dispatcher upload handler 直接使用本地 `internal/oss` client 签发和验证 OSS grant;这条路径属于 Dispatcher ↔ Agent 的当前实现,不能写成 SaaS 已联调。 + +### 6.2 AI immutable version GET + +权威源为 `contracts/upstream/2026-09-19-p1-v1/ai-config.openapi.yaml`: + +`GET /internal/v1/ai/agent-versions/{agent_version_id}` + +请求要求租户和 request header,返回不可变版本的 `tenant_id`、`agent_version_id`、`status`、`immutable=true`、`content_sha256` 和 `config`。`config` 顶层由 `agent_version_id`、`immutable=true`、`asr`、`conversation` 以及 full-AI 模式所需的 `llm`、`prompt`、`tts` 组成;ASR-only 不得带后三者。 + +当前实现只在 `internal/rpc.ServerOptions` 接收本地 `AISnapshotRaw`/`AIAuthorizationRaw`,由 `internal/ai/snapshot.go` 和 `internal/ai/authorization.go` 校验版本、摘要、租户、模式、有效期和 egress;没有 SaaS GET client,也没有“断 SaaS 后使用任意默认配置”的回退。 + +## 7. 禁止误读 + +1. RabbitMQ publisher confirm ≠ SaaS 已应用。 +2. HTTP `202 accepted` ≠ Agent 已执行或控制已生效。 +3. 当前本地 OSS grant/verify ≠ SaaS recording session/verified 已完成。 +4. Schema 支持 `recording.failed`、通话查询或 AI GET ≠ 当前 Go 代码已经生成或提供这些接口。 +5. 任何重试都必须保留原 `command_id`、`execution_id`、`event_id` 或 `upload_id` 的业务语义;未知执行不能换 ID 重拨。 + +## 8. 依据文件 + +- `internal/contract/contract.go` +- `internal/mq/amqp.go` +- `internal/tenant/routing.go` +- `internal/dispatcher/consumer.go` +- `internal/dispatcher/dispatcher.go` +- `internal/store/store.go` +- `internal/store/facts.go` +- `internal/control/http.go` +- `contracts/upstream/2026-09-19-p1-v1/mq.schema.json` +- `contracts/upstream/2026-09-19-p1-v1/event-payloads.schema.json` +- `contracts/upstream/2026-09-19-p1-v1/mq-topology.md` +- `contracts/upstream/2026-09-19-p1-v1/executor.openapi.yaml` +- `contracts/upstream/2026-09-19-p1-v1/saas.openapi.yaml` +- `contracts/upstream/2026-09-19-p1-v1/ai-config.openapi.yaml`