Files
go-sip/docs/contracts/dispatcher-agent.md
T

371 lines
21 KiB
Markdown
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# 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。
**SaaS↔Dispatcher 边界已修订为 MQ-only**,详见 [SaaS↔Dispatcher 契约](./saas-dispatcher.md)。D 有全局唯一身份和独立接收 Topic;这不改变内部 Unary 或 Agent→OSS 直传。**OSS 配置由 D 配置文件维护,Agent 向 D 领取临时上传 TOKEN,SaaS 不再提供 OSS 配置/TOKEN。** D 的签发职责保留;本文“当前实现”仍不能证明配置文件/TOKEN 全部约束已接通,尤其 D 本地对象验证不能代替 SaaS MQ verified。本轮不改变 SaaS 最终校验归属,未改 Proto/代码。
## 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 水位。`dispatcher_epoch` 是运行代次,不是全局唯一的 Dispatcher 逻辑 ID;后者的 MQ 关联及与内部会话的绑定待新契约冻结,当前 Proto 未因此自动增加字段。
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` 的调用封装。该方法属于 D→A 内部准入职责,不把它直接等同 SaaS 任务控制;SaaS 业务控制只能经 MQ 进入 D,旧 HTTP 控制入口应移除。
### 5.5 `GetExecutionPermit`
请求:`meta`、`ExecutionBinding`、`resource_reservation_id`、`expected_task_revision`、`admission_generation`、`config_sha256`。
响应:`OperationReceipt` + `ExecutionPermit`:
| 字段 | 含义 |
| --- | --- |
| `permit_id` | 当前实现为 `permit-<execution_id>` |
| `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-<fact_id>`,并在同一 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`。目标是 D 读取自身 OSS 配置文件并向 A 提供临时 TOKEN/受限上传信息;以下记录当前 handler,不代表配置入口、TOKEN 形态和新版业务会话约束已全部验收:
- 以 `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` 时才替换,不自动续期。目标流程同样由 Agent 显式向 D 重新领取 TOKEN;D 使用自身配置,不向 SaaS 申请 TOKEN,不改变原资产/会话。
新请求当前返回 `UPLOAD_STATE_REQUESTED`;持久层状态为 `granted`。`UPLOADING` 枚举存在,但当前 Dispatcher handler 不把 Agent 的 PUT 过程映射为该状态。
目标流程为 **D 依据自身配置文件向 A 提供临时上传 TOKEN,不向 SaaS 申请 OSS 配置/TOKEN**。D 侧配置缺失/无效时明确失败,不切换配置源;长期凭据不交给 A,不写入示例、日志或证据。TOKEN 的精确形态、SDK 能力及与 `UploadGrant` 的映射须核验冻结,不能只把现有字段改称 TOKEN 就宣称完成。
既有 SaaS 业务会话/资产登记仍经 MQ,D 关联原租户/执行/资产后才交付相应上传信息;这不是由 SaaS 签 TOKEN。业务响应可能晚于 RPC deadline,W02/W11 仍须冻结 pending、有界等待、原操作重取及最终结果,不无限阻塞 Unary、不因超时另造 upload ID。当前 D 的本地签发能力保留复用,禁止删除后改为等待 SaaS 下发配置。本段不新增 RPC/Proto 字段。
### 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` 且标记可重试。
以上是旧本地验证事实。**目标流程必须由 D 经 MQ 提交 complete,SaaS 独立验证对象后经本 D 专用 Topic 返回 verified/oss_id;D 校验原请求、租户、资产和会话并持久化后,才可记完成及写 recording.ready outbox。** D 本地 HEAD、PUT 2xx 或 broker confirm 不能替代 SaaS verified。等待中/超时/重复响应及 A 获取最终结果的 Unary 衔接由 W02/W11 冻结,尚未完成;不凭此说明宣称当前 handler 已符合目标。
## 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. D 配置文件→临时 TOKEN→A 直传的完整接线/约束,以及 SaaS 业务会话/complete/verified 的 MQ 协调与 R12/R13 有界衔接仍待核验;D 已有签发能力保留复用,但本地对象验证不能替代 SaaS verified。旧 `recording-uploads` HTTP 方案继续废弃,不开发 client。
4. SaaS AI 配置/授权的 MQ 请求响应与持久绑定尚未接通;当前快照/授权由启动输入提供。旧 AI version GET 已废弃,不能作为后续实现方向。
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`