diff --git a/docs/evidence/saas-dispatcher-implementation.md b/docs/evidence/saas-dispatcher-implementation.md index a00abce..ed32652 100644 --- a/docs/evidence/saas-dispatcher-implementation.md +++ b/docs/evidence/saas-dispatcher-implementation.md @@ -64,7 +64,8 @@ - Dispatcher 的无录音最终结果隔离组件:`CurrentStore.RecordCallResult` 仅在确认通话结束后,按持久任务快照校验任务、被叫、主叫和已选线路,并以源执行事件固定生成唯一最终结果身份;消息通过严格 MQ Schema 校验后与 outbox 在同一事务写入。同内容重投/重启只恢复原消息,冲突结果和 SQLite 写入失败均不会产生第二份结果。这里只验证隔离组件,Agent 实际回报尚未连通。 - Dispatcher 的原始 OSS 目标及已上传结果隔离组件:新 SQLite 布局把一次通话的 upload_id、bucket、object_key、录音格式/时长/大小和 SHA-256 唯一绑定到已保留的执行;不保存临时 URL 或 TOKEN。旧布局拒绝启动并原样保留待交付 outbox,不自动迁移或清理。已签发录音目标不能通过空录音结果绕过上传;已有空录音最终结果不能再签发录音目标。`RecordUploadedCallResult` 仅接受与持久绑定完全一致的录音事实及 Agent 所报告的成功 PUT 状态,录音确认与唯一最终结果 outbox 同事务提交;丢失回报或 MQ 投递时重用原消息,已确认后拒绝再次签发 PUT 授权。Mock 证明的是本地状态约束,不是独立 OSS 校验或真实 Agent 身份验证。 - Dispatcher 的终结与外呼回执顺序竞争隔离修复:原流程在 Agent 接受执行的 RPC 返回后才写入外呼回执,快速结束或 RPC 超时可能先到;现以一次 SQLite 事务在确认通话已结束后补齐原回执并释放占用,未知执行仅在确认结束后释放。迟到的执行响应、超时和重复结束不会产生第二份回执;注入 outbox 写入失败保留原占用。并发竞争及结束后立即生成唯一最终结果有单元测试;真实 Agent 结束事实的鉴权与传输尚未接线。 -- 已验证:`go test ./... -count=1`、`go test -race ./internal/agent -count=1`、`go test -race ./internal/store -count=1`、`go vet ./...`、`go build ./...`、`PATH=/tmp/sip-go-agent-tools/bin:$PATH bash scripts/check-current-contracts.sh`、`PATH=/tmp/sip-go-agent-tools/bin:$PATH bash scripts/check-proto.sh`、`git diff --check`。尚未完成实际录音到直传/失败恢复的接线、无录音与生成失败的真实调用、Dispatcher 实际签发授权与 Agent↔Dispatcher 上传事实/最终结果交付及 MQ/端到端验收,不能宣称 P06 通过。 +- Dispatcher 录音事实 Unary RPC 隔离服务:`RequestRecordingUpload`、`ReportCallEnded`、`ReportCallResult` 均要求已配置本 D、核验 mTLS 指纹及当前 Agent 会话、数字租户和已保留的执行;复用官方 SDK 仅对原始录音签发固定 15 分钟授权,显式重申请仍用相同 bucket/object_key。结束事实可先于外呼响应而持久化原回执;录音结果核对 D 已存目标和 Agent 报告的成功 PUT,再与唯一结果 outbox 同事务提交。Mock 覆盖会话/租户拒绝、同资产重申请、上传前结果拒绝、坏 JSON、已上传与无录音结果及重复回报。当前是服务方法的本地隔离测试,尚无真实 gRPC 握手、Agent 调用或实际 OSS PUT。 +- 已验证:`go test ./... -count=1`、`go test -race ./internal/agent ./internal/rpc ./internal/store -count=1`(分批执行)、`go vet ./...`、`go build ./...`、`PATH=/tmp/sip-go-agent-tools/bin:$PATH bash scripts/check-current-contracts.sh`、`PATH=/tmp/sip-go-agent-tools/bin:$PATH bash scripts/check-proto.sh`、`git diff --check`。尚未完成实际录音到直传/失败恢复的接线、无录音与生成失败的真实调用、主入口 Agent↔Dispatcher 受控会话/完整 gRPC 传输及上传事实/最终结果交付、MQ/端到端验收,不能宣称 P06 通过。 ## 验收台账 diff --git a/gen/agent/agent.pb.go b/gen/agent/agent.pb.go index 69a5db9..8b64500 100644 --- a/gen/agent/agent.pb.go +++ b/gen/agent/agent.pb.go @@ -4293,6 +4293,452 @@ func (x *CompleteUploadResponse) GetState() UploadState { return UploadState_UPLOAD_STATE_UNSPECIFIED } +// These facts refer only to the Dispatcher-approved execution identified by +// source_event_id; no audio bytes or temporary credential is persisted in MQ. +type RequestRecordingUploadRequest struct { + state protoimpl.MessageState `protogen:"open.v1"` + Meta *RequestMeta `protobuf:"bytes,1,opt,name=meta,proto3" json:"meta,omitempty"` + DispatcherId string `protobuf:"bytes,2,opt,name=dispatcher_id,json=dispatcherId,proto3" json:"dispatcher_id,omitempty"` + TenantId int64 `protobuf:"varint,3,opt,name=tenant_id,json=tenantId,proto3" json:"tenant_id,omitempty"` + SourceEventId string `protobuf:"bytes,4,opt,name=source_event_id,json=sourceEventId,proto3" json:"source_event_id,omitempty"` + UploadId string `protobuf:"bytes,5,opt,name=upload_id,json=uploadId,proto3" json:"upload_id,omitempty"` + Asset *AssetDescriptor `protobuf:"bytes,6,opt,name=asset,proto3" json:"asset,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *RequestRecordingUploadRequest) Reset() { + *x = RequestRecordingUploadRequest{} + mi := &file_agent_agent_proto_msgTypes[47] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *RequestRecordingUploadRequest) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*RequestRecordingUploadRequest) ProtoMessage() {} + +func (x *RequestRecordingUploadRequest) ProtoReflect() protoreflect.Message { + mi := &file_agent_agent_proto_msgTypes[47] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use RequestRecordingUploadRequest.ProtoReflect.Descriptor instead. +func (*RequestRecordingUploadRequest) Descriptor() ([]byte, []int) { + return file_agent_agent_proto_rawDescGZIP(), []int{47} +} + +func (x *RequestRecordingUploadRequest) GetMeta() *RequestMeta { + if x != nil { + return x.Meta + } + return nil +} + +func (x *RequestRecordingUploadRequest) GetDispatcherId() string { + if x != nil { + return x.DispatcherId + } + return "" +} + +func (x *RequestRecordingUploadRequest) GetTenantId() int64 { + if x != nil { + return x.TenantId + } + return 0 +} + +func (x *RequestRecordingUploadRequest) GetSourceEventId() string { + if x != nil { + return x.SourceEventId + } + return "" +} + +func (x *RequestRecordingUploadRequest) GetUploadId() string { + if x != nil { + return x.UploadId + } + return "" +} + +func (x *RequestRecordingUploadRequest) GetAsset() *AssetDescriptor { + if x != nil { + return x.Asset + } + return nil +} + +type RequestRecordingUploadResponse struct { + state protoimpl.MessageState `protogen:"open.v1"` + Grant *UploadGrant `protobuf:"bytes,1,opt,name=grant,proto3" json:"grant,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *RequestRecordingUploadResponse) Reset() { + *x = RequestRecordingUploadResponse{} + mi := &file_agent_agent_proto_msgTypes[48] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *RequestRecordingUploadResponse) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*RequestRecordingUploadResponse) ProtoMessage() {} + +func (x *RequestRecordingUploadResponse) ProtoReflect() protoreflect.Message { + mi := &file_agent_agent_proto_msgTypes[48] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use RequestRecordingUploadResponse.ProtoReflect.Descriptor instead. +func (*RequestRecordingUploadResponse) Descriptor() ([]byte, []int) { + return file_agent_agent_proto_rawDescGZIP(), []int{48} +} + +func (x *RequestRecordingUploadResponse) GetGrant() *UploadGrant { + if x != nil { + return x.Grant + } + return nil +} + +type ReportCallEndedRequest struct { + state protoimpl.MessageState `protogen:"open.v1"` + Meta *RequestMeta `protobuf:"bytes,1,opt,name=meta,proto3" json:"meta,omitempty"` + DispatcherId string `protobuf:"bytes,2,opt,name=dispatcher_id,json=dispatcherId,proto3" json:"dispatcher_id,omitempty"` + TenantId int64 `protobuf:"varint,3,opt,name=tenant_id,json=tenantId,proto3" json:"tenant_id,omitempty"` + SourceEventId string `protobuf:"bytes,4,opt,name=source_event_id,json=sourceEventId,proto3" json:"source_event_id,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *ReportCallEndedRequest) Reset() { + *x = ReportCallEndedRequest{} + mi := &file_agent_agent_proto_msgTypes[49] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *ReportCallEndedRequest) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*ReportCallEndedRequest) ProtoMessage() {} + +func (x *ReportCallEndedRequest) ProtoReflect() protoreflect.Message { + mi := &file_agent_agent_proto_msgTypes[49] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use ReportCallEndedRequest.ProtoReflect.Descriptor instead. +func (*ReportCallEndedRequest) Descriptor() ([]byte, []int) { + return file_agent_agent_proto_rawDescGZIP(), []int{49} +} + +func (x *ReportCallEndedRequest) GetMeta() *RequestMeta { + if x != nil { + return x.Meta + } + return nil +} + +func (x *ReportCallEndedRequest) GetDispatcherId() string { + if x != nil { + return x.DispatcherId + } + return "" +} + +func (x *ReportCallEndedRequest) GetTenantId() int64 { + if x != nil { + return x.TenantId + } + return 0 +} + +func (x *ReportCallEndedRequest) GetSourceEventId() string { + if x != nil { + return x.SourceEventId + } + return "" +} + +type ReportCallEndedResponse struct { + state protoimpl.MessageState `protogen:"open.v1"` + Receipt *OperationReceipt `protobuf:"bytes,1,opt,name=receipt,proto3" json:"receipt,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *ReportCallEndedResponse) Reset() { + *x = ReportCallEndedResponse{} + mi := &file_agent_agent_proto_msgTypes[50] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *ReportCallEndedResponse) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*ReportCallEndedResponse) ProtoMessage() {} + +func (x *ReportCallEndedResponse) ProtoReflect() protoreflect.Message { + mi := &file_agent_agent_proto_msgTypes[50] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use ReportCallEndedResponse.ProtoReflect.Descriptor instead. +func (*ReportCallEndedResponse) Descriptor() ([]byte, []int) { + return file_agent_agent_proto_rawDescGZIP(), []int{50} +} + +func (x *ReportCallEndedResponse) GetReceipt() *OperationReceipt { + if x != nil { + return x.Receipt + } + return nil +} + +type UploadObservation struct { + state protoimpl.MessageState `protogen:"open.v1"` + UploadId string `protobuf:"bytes,1,opt,name=upload_id,json=uploadId,proto3" json:"upload_id,omitempty"` + RecordingId string `protobuf:"bytes,2,opt,name=recording_id,json=recordingId,proto3" json:"recording_id,omitempty"` + PutStatusCode int32 `protobuf:"varint,3,opt,name=put_status_code,json=putStatusCode,proto3" json:"put_status_code,omitempty"` + SizeBytes int64 `protobuf:"varint,4,opt,name=size_bytes,json=sizeBytes,proto3" json:"size_bytes,omitempty"` + ChecksumSha256 string `protobuf:"bytes,5,opt,name=checksum_sha256,json=checksumSha256,proto3" json:"checksum_sha256,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *UploadObservation) Reset() { + *x = UploadObservation{} + mi := &file_agent_agent_proto_msgTypes[51] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *UploadObservation) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*UploadObservation) ProtoMessage() {} + +func (x *UploadObservation) ProtoReflect() protoreflect.Message { + mi := &file_agent_agent_proto_msgTypes[51] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use UploadObservation.ProtoReflect.Descriptor instead. +func (*UploadObservation) Descriptor() ([]byte, []int) { + return file_agent_agent_proto_rawDescGZIP(), []int{51} +} + +func (x *UploadObservation) GetUploadId() string { + if x != nil { + return x.UploadId + } + return "" +} + +func (x *UploadObservation) GetRecordingId() string { + if x != nil { + return x.RecordingId + } + return "" +} + +func (x *UploadObservation) GetPutStatusCode() int32 { + if x != nil { + return x.PutStatusCode + } + return 0 +} + +func (x *UploadObservation) GetSizeBytes() int64 { + if x != nil { + return x.SizeBytes + } + return 0 +} + +func (x *UploadObservation) GetChecksumSha256() string { + if x != nil { + return x.ChecksumSha256 + } + return "" +} + +type ReportCallResultRequest struct { + state protoimpl.MessageState `protogen:"open.v1"` + Meta *RequestMeta `protobuf:"bytes,1,opt,name=meta,proto3" json:"meta,omitempty"` + DispatcherId string `protobuf:"bytes,2,opt,name=dispatcher_id,json=dispatcherId,proto3" json:"dispatcher_id,omitempty"` + TenantId int64 `protobuf:"varint,3,opt,name=tenant_id,json=tenantId,proto3" json:"tenant_id,omitempty"` + SourceEventId string `protobuf:"bytes,4,opt,name=source_event_id,json=sourceEventId,proto3" json:"source_event_id,omitempty"` + ResultPayloadJson []byte `protobuf:"bytes,5,opt,name=result_payload_json,json=resultPayloadJson,proto3" json:"result_payload_json,omitempty"` + Upload *UploadObservation `protobuf:"bytes,6,opt,name=upload,proto3" json:"upload,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *ReportCallResultRequest) Reset() { + *x = ReportCallResultRequest{} + mi := &file_agent_agent_proto_msgTypes[52] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *ReportCallResultRequest) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*ReportCallResultRequest) ProtoMessage() {} + +func (x *ReportCallResultRequest) ProtoReflect() protoreflect.Message { + mi := &file_agent_agent_proto_msgTypes[52] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use ReportCallResultRequest.ProtoReflect.Descriptor instead. +func (*ReportCallResultRequest) Descriptor() ([]byte, []int) { + return file_agent_agent_proto_rawDescGZIP(), []int{52} +} + +func (x *ReportCallResultRequest) GetMeta() *RequestMeta { + if x != nil { + return x.Meta + } + return nil +} + +func (x *ReportCallResultRequest) GetDispatcherId() string { + if x != nil { + return x.DispatcherId + } + return "" +} + +func (x *ReportCallResultRequest) GetTenantId() int64 { + if x != nil { + return x.TenantId + } + return 0 +} + +func (x *ReportCallResultRequest) GetSourceEventId() string { + if x != nil { + return x.SourceEventId + } + return "" +} + +func (x *ReportCallResultRequest) GetResultPayloadJson() []byte { + if x != nil { + return x.ResultPayloadJson + } + return nil +} + +func (x *ReportCallResultRequest) GetUpload() *UploadObservation { + if x != nil { + return x.Upload + } + return nil +} + +type ReportCallResultResponse struct { + state protoimpl.MessageState `protogen:"open.v1"` + Receipt *OperationReceipt `protobuf:"bytes,1,opt,name=receipt,proto3" json:"receipt,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *ReportCallResultResponse) Reset() { + *x = ReportCallResultResponse{} + mi := &file_agent_agent_proto_msgTypes[53] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *ReportCallResultResponse) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*ReportCallResultResponse) ProtoMessage() {} + +func (x *ReportCallResultResponse) ProtoReflect() protoreflect.Message { + mi := &file_agent_agent_proto_msgTypes[53] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use ReportCallResultResponse.ProtoReflect.Descriptor instead. +func (*ReportCallResultResponse) Descriptor() ([]byte, []int) { + return file_agent_agent_proto_rawDescGZIP(), []int{53} +} + +func (x *ReportCallResultResponse) GetReceipt() *OperationReceipt { + if x != nil { + return x.Receipt + } + return nil +} + var File_agent_agent_proto protoreflect.FileDescriptor const file_agent_agent_proto_rawDesc = "" + @@ -4620,7 +5066,39 @@ const file_agent_agent_proto_rawDesc = "" + "\x18uploaded_checksum_sha256\x18\x06 \x01(\tR\x16uploadedChecksumSha256\"\x83\x01\n" + "\x16CompleteUploadResponse\x121\n" + "\areceipt\x18\x01 \x01(\v2\x17.agent.OperationReceiptR\areceipt\x12(\n" + - "\x05state\x18\x02 \x01(\x0e2\x12.agent.UploadStateR\x05stateJ\x04\b\x03\x10\x04R\x06oss_id*\xa9\x01\n" + + "\x05state\x18\x02 \x01(\x0e2\x12.agent.UploadStateR\x05stateJ\x04\b\x03\x10\x04R\x06oss_id\"\xfc\x01\n" + + "\x1dRequestRecordingUploadRequest\x12&\n" + + "\x04meta\x18\x01 \x01(\v2\x12.agent.RequestMetaR\x04meta\x12#\n" + + "\rdispatcher_id\x18\x02 \x01(\tR\fdispatcherId\x12\x1b\n" + + "\ttenant_id\x18\x03 \x01(\x03R\btenantId\x12&\n" + + "\x0fsource_event_id\x18\x04 \x01(\tR\rsourceEventId\x12\x1b\n" + + "\tupload_id\x18\x05 \x01(\tR\buploadId\x12,\n" + + "\x05asset\x18\x06 \x01(\v2\x16.agent.AssetDescriptorR\x05asset\"J\n" + + "\x1eRequestRecordingUploadResponse\x12(\n" + + "\x05grant\x18\x01 \x01(\v2\x12.agent.UploadGrantR\x05grant\"\xaa\x01\n" + + "\x16ReportCallEndedRequest\x12&\n" + + "\x04meta\x18\x01 \x01(\v2\x12.agent.RequestMetaR\x04meta\x12#\n" + + "\rdispatcher_id\x18\x02 \x01(\tR\fdispatcherId\x12\x1b\n" + + "\ttenant_id\x18\x03 \x01(\x03R\btenantId\x12&\n" + + "\x0fsource_event_id\x18\x04 \x01(\tR\rsourceEventId\"L\n" + + "\x17ReportCallEndedResponse\x121\n" + + "\areceipt\x18\x01 \x01(\v2\x17.agent.OperationReceiptR\areceipt\"\xc3\x01\n" + + "\x11UploadObservation\x12\x1b\n" + + "\tupload_id\x18\x01 \x01(\tR\buploadId\x12!\n" + + "\frecording_id\x18\x02 \x01(\tR\vrecordingId\x12&\n" + + "\x0fput_status_code\x18\x03 \x01(\x05R\rputStatusCode\x12\x1d\n" + + "\n" + + "size_bytes\x18\x04 \x01(\x03R\tsizeBytes\x12'\n" + + "\x0fchecksum_sha256\x18\x05 \x01(\tR\x0echecksumSha256\"\x8d\x02\n" + + "\x17ReportCallResultRequest\x12&\n" + + "\x04meta\x18\x01 \x01(\v2\x12.agent.RequestMetaR\x04meta\x12#\n" + + "\rdispatcher_id\x18\x02 \x01(\tR\fdispatcherId\x12\x1b\n" + + "\ttenant_id\x18\x03 \x01(\x03R\btenantId\x12&\n" + + "\x0fsource_event_id\x18\x04 \x01(\tR\rsourceEventId\x12.\n" + + "\x13result_payload_json\x18\x05 \x01(\fR\x11resultPayloadJson\x120\n" + + "\x06upload\x18\x06 \x01(\v2\x18.agent.UploadObservationR\x06upload\"M\n" + + "\x18ReportCallResultResponse\x121\n" + + "\areceipt\x18\x01 \x01(\v2\x17.agent.OperationReceiptR\areceipt*\xa9\x01\n" + "\n" + "ResultCode\x12\x1b\n" + "\x17RESULT_CODE_UNSPECIFIED\x10\x00\x12\x18\n" + @@ -4688,7 +5166,7 @@ const file_agent_agent_proto_rawDesc = "" + "\x1cFACT_KIND_TRANSCRIPT_UPDATED\x10\x04\x12\x1f\n" + "\x1bFACT_KIND_TRANSCRIPT_FAILED\x10\x05\x12\x1d\n" + "\x19FACT_KIND_CONTACT_OPT_OUT\x10\x06\x12 \n" + - "\x1cFACT_KIND_RECORDING_PROGRESS\x10\a2\xf9\b\n" + + "\x1cFACT_KIND_RECORDING_PROGRESS\x10\a2\x87\v\n" + "\x13AgentControlService\x12M\n" + "\x0eGetAgentStatus\x12\x1c.agent.GetAgentStatusRequest\x1a\x1d.agent.GetAgentStatusResponse\x12J\n" + "\rActivateAgent\x12\x1b.agent.ActivateAgentRequest\x1a\x1c.agent.ActivateAgentResponse\x12G\n" + @@ -4703,7 +5181,10 @@ const file_agent_agent_proto_rawDesc = "" + "\x0eQueryExecution\x12\x1c.agent.QueryExecutionRequest\x1a\x1d.agent.QueryExecutionResponse\x12_\n" + "\x14ReportExecutionEvent\x12\".agent.ReportExecutionEventRequest\x1a#.agent.ReportExecutionEventResponse\x12J\n" + "\rRequestUpload\x12\x1b.agent.RequestUploadRequest\x1a\x1c.agent.RequestUploadResponse\x12M\n" + - "\x0eCompleteUpload\x12\x1c.agent.CompleteUploadRequest\x1a\x1d.agent.CompleteUploadResponseB-Z+git.ipao.vip/rogee/go-sip/gen/agent;agentpbb\x06proto3" + "\x0eCompleteUpload\x12\x1c.agent.CompleteUploadRequest\x1a\x1d.agent.CompleteUploadResponse\x12e\n" + + "\x16RequestRecordingUpload\x12$.agent.RequestRecordingUploadRequest\x1a%.agent.RequestRecordingUploadResponse\x12P\n" + + "\x0fReportCallEnded\x12\x1d.agent.ReportCallEndedRequest\x1a\x1e.agent.ReportCallEndedResponse\x12S\n" + + "\x10ReportCallResult\x12\x1e.agent.ReportCallResultRequest\x1a\x1f.agent.ReportCallResultResponseB-Z+git.ipao.vip/rogee/go-sip/gen/agent;agentpbb\x06proto3" var ( file_agent_agent_proto_rawDescOnce sync.Once @@ -4718,177 +5199,198 @@ func file_agent_agent_proto_rawDescGZIP() []byte { } var file_agent_agent_proto_enumTypes = make([]protoimpl.EnumInfo, 10) -var file_agent_agent_proto_msgTypes = make([]protoimpl.MessageInfo, 48) +var file_agent_agent_proto_msgTypes = make([]protoimpl.MessageInfo, 55) var file_agent_agent_proto_goTypes = []any{ - (ResultCode)(0), // 0: agent.ResultCode - (FailureCode)(0), // 1: agent.FailureCode - (ActivationState)(0), // 2: agent.ActivationState - (AdmissionState)(0), // 3: agent.AdmissionState - (ControlAction)(0), // 4: agent.ControlAction - (ActiveCallPolicy)(0), // 5: agent.ActiveCallPolicy - (ExecutionState)(0), // 6: agent.ExecutionState - (AssetKind)(0), // 7: agent.AssetKind - (UploadState)(0), // 8: agent.UploadState - (FactKind)(0), // 9: agent.FactKind - (*RequestMeta)(nil), // 10: agent.RequestMeta - (*ResponseMeta)(nil), // 11: agent.ResponseMeta - (*Failure)(nil), // 12: agent.Failure - (*OperationReceipt)(nil), // 13: agent.OperationReceipt - (*AgentBinding)(nil), // 14: agent.AgentBinding - (*Capability)(nil), // 15: agent.Capability - (*ResourceSample)(nil), // 16: agent.ResourceSample - (*AppliedConfig)(nil), // 17: agent.AppliedConfig - (*AgentStatus)(nil), // 18: agent.AgentStatus - (*Session)(nil), // 19: agent.Session - (*ConfigReference)(nil), // 20: agent.ConfigReference - (*UploadPolicy)(nil), // 21: agent.UploadPolicy - (*ExecutionBinding)(nil), // 22: agent.ExecutionBinding - (*AssetDescriptor)(nil), // 23: agent.AssetDescriptor - (*Header)(nil), // 24: agent.Header - (*GetAgentStatusRequest)(nil), // 25: agent.GetAgentStatusRequest - (*GetAgentStatusResponse)(nil), // 26: agent.GetAgentStatusResponse - (*ActivateAgentRequest)(nil), // 27: agent.ActivateAgentRequest - (*ActivateAgentResponse)(nil), // 28: agent.ActivateAgentResponse - (*GetBootstrapRequest)(nil), // 29: agent.GetBootstrapRequest - (*GetBootstrapResponse)(nil), // 30: agent.GetBootstrapResponse - (*SetAdmissionStateRequest)(nil), // 31: agent.SetAdmissionStateRequest - (*SetAdmissionStateResponse)(nil), // 32: agent.SetAdmissionStateResponse - (*ExecuteRequest)(nil), // 33: agent.ExecuteRequest - (*ExecuteResponse)(nil), // 34: agent.ExecuteResponse - (*ExecuteAuthorizedRequest)(nil), // 35: agent.ExecuteAuthorizedRequest - (*ExecuteAuthorizedResponse)(nil), // 36: agent.ExecuteAuthorizedResponse - (*ExecuteApprovedRequest)(nil), // 37: agent.ExecuteApprovedRequest - (*ExecuteApprovedResponse)(nil), // 38: agent.ExecuteApprovedResponse - (*GetLoadedSIPRequest)(nil), // 39: agent.GetLoadedSIPRequest - (*GetLoadedSIPResponse)(nil), // 40: agent.GetLoadedSIPResponse - (*GetExecutionPermitRequest)(nil), // 41: agent.GetExecutionPermitRequest - (*ExecutionPermit)(nil), // 42: agent.ExecutionPermit - (*GetExecutionPermitResponse)(nil), // 43: agent.GetExecutionPermitResponse - (*ApplyTaskControlRequest)(nil), // 44: agent.ApplyTaskControlRequest - (*ApplyTaskControlResponse)(nil), // 45: agent.ApplyTaskControlResponse - (*QueryExecutionRequest)(nil), // 46: agent.QueryExecutionRequest - (*ExecutionSnapshot)(nil), // 47: agent.ExecutionSnapshot - (*QueryExecutionResponse)(nil), // 48: agent.QueryExecutionResponse - (*ExecutionFact)(nil), // 49: agent.ExecutionFact - (*ReportExecutionEventRequest)(nil), // 50: agent.ReportExecutionEventRequest - (*ReportExecutionEventResponse)(nil), // 51: agent.ReportExecutionEventResponse - (*RequestUploadRequest)(nil), // 52: agent.RequestUploadRequest - (*UploadGrant)(nil), // 53: agent.UploadGrant - (*RequestUploadResponse)(nil), // 54: agent.RequestUploadResponse - (*CompleteUploadRequest)(nil), // 55: agent.CompleteUploadRequest - (*CompleteUploadResponse)(nil), // 56: agent.CompleteUploadResponse - nil, // 57: agent.GetLoadedSIPResponse.TrunkRevisionEntry + (ResultCode)(0), // 0: agent.ResultCode + (FailureCode)(0), // 1: agent.FailureCode + (ActivationState)(0), // 2: agent.ActivationState + (AdmissionState)(0), // 3: agent.AdmissionState + (ControlAction)(0), // 4: agent.ControlAction + (ActiveCallPolicy)(0), // 5: agent.ActiveCallPolicy + (ExecutionState)(0), // 6: agent.ExecutionState + (AssetKind)(0), // 7: agent.AssetKind + (UploadState)(0), // 8: agent.UploadState + (FactKind)(0), // 9: agent.FactKind + (*RequestMeta)(nil), // 10: agent.RequestMeta + (*ResponseMeta)(nil), // 11: agent.ResponseMeta + (*Failure)(nil), // 12: agent.Failure + (*OperationReceipt)(nil), // 13: agent.OperationReceipt + (*AgentBinding)(nil), // 14: agent.AgentBinding + (*Capability)(nil), // 15: agent.Capability + (*ResourceSample)(nil), // 16: agent.ResourceSample + (*AppliedConfig)(nil), // 17: agent.AppliedConfig + (*AgentStatus)(nil), // 18: agent.AgentStatus + (*Session)(nil), // 19: agent.Session + (*ConfigReference)(nil), // 20: agent.ConfigReference + (*UploadPolicy)(nil), // 21: agent.UploadPolicy + (*ExecutionBinding)(nil), // 22: agent.ExecutionBinding + (*AssetDescriptor)(nil), // 23: agent.AssetDescriptor + (*Header)(nil), // 24: agent.Header + (*GetAgentStatusRequest)(nil), // 25: agent.GetAgentStatusRequest + (*GetAgentStatusResponse)(nil), // 26: agent.GetAgentStatusResponse + (*ActivateAgentRequest)(nil), // 27: agent.ActivateAgentRequest + (*ActivateAgentResponse)(nil), // 28: agent.ActivateAgentResponse + (*GetBootstrapRequest)(nil), // 29: agent.GetBootstrapRequest + (*GetBootstrapResponse)(nil), // 30: agent.GetBootstrapResponse + (*SetAdmissionStateRequest)(nil), // 31: agent.SetAdmissionStateRequest + (*SetAdmissionStateResponse)(nil), // 32: agent.SetAdmissionStateResponse + (*ExecuteRequest)(nil), // 33: agent.ExecuteRequest + (*ExecuteResponse)(nil), // 34: agent.ExecuteResponse + (*ExecuteAuthorizedRequest)(nil), // 35: agent.ExecuteAuthorizedRequest + (*ExecuteAuthorizedResponse)(nil), // 36: agent.ExecuteAuthorizedResponse + (*ExecuteApprovedRequest)(nil), // 37: agent.ExecuteApprovedRequest + (*ExecuteApprovedResponse)(nil), // 38: agent.ExecuteApprovedResponse + (*GetLoadedSIPRequest)(nil), // 39: agent.GetLoadedSIPRequest + (*GetLoadedSIPResponse)(nil), // 40: agent.GetLoadedSIPResponse + (*GetExecutionPermitRequest)(nil), // 41: agent.GetExecutionPermitRequest + (*ExecutionPermit)(nil), // 42: agent.ExecutionPermit + (*GetExecutionPermitResponse)(nil), // 43: agent.GetExecutionPermitResponse + (*ApplyTaskControlRequest)(nil), // 44: agent.ApplyTaskControlRequest + (*ApplyTaskControlResponse)(nil), // 45: agent.ApplyTaskControlResponse + (*QueryExecutionRequest)(nil), // 46: agent.QueryExecutionRequest + (*ExecutionSnapshot)(nil), // 47: agent.ExecutionSnapshot + (*QueryExecutionResponse)(nil), // 48: agent.QueryExecutionResponse + (*ExecutionFact)(nil), // 49: agent.ExecutionFact + (*ReportExecutionEventRequest)(nil), // 50: agent.ReportExecutionEventRequest + (*ReportExecutionEventResponse)(nil), // 51: agent.ReportExecutionEventResponse + (*RequestUploadRequest)(nil), // 52: agent.RequestUploadRequest + (*UploadGrant)(nil), // 53: agent.UploadGrant + (*RequestUploadResponse)(nil), // 54: agent.RequestUploadResponse + (*CompleteUploadRequest)(nil), // 55: agent.CompleteUploadRequest + (*CompleteUploadResponse)(nil), // 56: agent.CompleteUploadResponse + (*RequestRecordingUploadRequest)(nil), // 57: agent.RequestRecordingUploadRequest + (*RequestRecordingUploadResponse)(nil), // 58: agent.RequestRecordingUploadResponse + (*ReportCallEndedRequest)(nil), // 59: agent.ReportCallEndedRequest + (*ReportCallEndedResponse)(nil), // 60: agent.ReportCallEndedResponse + (*UploadObservation)(nil), // 61: agent.UploadObservation + (*ReportCallResultRequest)(nil), // 62: agent.ReportCallResultRequest + (*ReportCallResultResponse)(nil), // 63: agent.ReportCallResultResponse + nil, // 64: agent.GetLoadedSIPResponse.TrunkRevisionEntry } var file_agent_agent_proto_depIdxs = []int32{ - 1, // 0: agent.Failure.code:type_name -> agent.FailureCode - 11, // 1: agent.OperationReceipt.meta:type_name -> agent.ResponseMeta - 0, // 2: agent.OperationReceipt.result:type_name -> agent.ResultCode - 12, // 3: agent.OperationReceipt.failure:type_name -> agent.Failure - 3, // 4: agent.AgentStatus.admission_state:type_name -> agent.AdmissionState - 15, // 5: agent.AgentStatus.capabilities:type_name -> agent.Capability - 16, // 6: agent.AgentStatus.resources:type_name -> agent.ResourceSample - 17, // 7: agent.AgentStatus.applied_configs:type_name -> agent.AppliedConfig - 7, // 8: agent.AssetDescriptor.kind:type_name -> agent.AssetKind - 10, // 9: agent.GetAgentStatusRequest.meta:type_name -> agent.RequestMeta - 14, // 10: agent.GetAgentStatusRequest.target:type_name -> agent.AgentBinding - 11, // 11: agent.GetAgentStatusResponse.meta:type_name -> agent.ResponseMeta - 18, // 12: agent.GetAgentStatusResponse.status:type_name -> agent.AgentStatus - 12, // 13: agent.GetAgentStatusResponse.failure:type_name -> agent.Failure - 10, // 14: agent.ActivateAgentRequest.meta:type_name -> agent.RequestMeta - 14, // 15: agent.ActivateAgentRequest.binding:type_name -> agent.AgentBinding - 11, // 16: agent.ActivateAgentResponse.meta:type_name -> agent.ResponseMeta - 2, // 17: agent.ActivateAgentResponse.state:type_name -> agent.ActivationState - 19, // 18: agent.ActivateAgentResponse.session:type_name -> agent.Session - 12, // 19: agent.ActivateAgentResponse.failure:type_name -> agent.Failure - 10, // 20: agent.GetBootstrapRequest.meta:type_name -> agent.RequestMeta - 11, // 21: agent.GetBootstrapResponse.meta:type_name -> agent.ResponseMeta - 2, // 22: agent.GetBootstrapResponse.state:type_name -> agent.ActivationState - 20, // 23: agent.GetBootstrapResponse.runtime_configs:type_name -> agent.ConfigReference - 21, // 24: agent.GetBootstrapResponse.upload_policy:type_name -> agent.UploadPolicy - 12, // 25: agent.GetBootstrapResponse.failure:type_name -> agent.Failure - 10, // 26: agent.SetAdmissionStateRequest.meta:type_name -> agent.RequestMeta - 14, // 27: agent.SetAdmissionStateRequest.target:type_name -> agent.AgentBinding - 3, // 28: agent.SetAdmissionStateRequest.state:type_name -> agent.AdmissionState - 13, // 29: agent.SetAdmissionStateResponse.receipt:type_name -> agent.OperationReceipt - 10, // 30: agent.ExecuteRequest.meta:type_name -> agent.RequestMeta - 22, // 31: agent.ExecuteRequest.binding:type_name -> agent.ExecutionBinding - 13, // 32: agent.ExecuteResponse.receipt:type_name -> agent.OperationReceipt - 6, // 33: agent.ExecuteResponse.state:type_name -> agent.ExecutionState - 10, // 34: agent.ExecuteAuthorizedRequest.meta:type_name -> agent.RequestMeta - 22, // 35: agent.ExecuteAuthorizedRequest.binding:type_name -> agent.ExecutionBinding - 13, // 36: agent.ExecuteAuthorizedResponse.receipt:type_name -> agent.OperationReceipt - 6, // 37: agent.ExecuteAuthorizedResponse.state:type_name -> agent.ExecutionState - 10, // 38: agent.ExecuteApprovedRequest.meta:type_name -> agent.RequestMeta - 10, // 39: agent.GetLoadedSIPRequest.meta:type_name -> agent.RequestMeta - 57, // 40: agent.GetLoadedSIPResponse.trunk_revision:type_name -> agent.GetLoadedSIPResponse.TrunkRevisionEntry - 10, // 41: agent.GetExecutionPermitRequest.meta:type_name -> agent.RequestMeta - 22, // 42: agent.GetExecutionPermitRequest.binding:type_name -> agent.ExecutionBinding - 13, // 43: agent.GetExecutionPermitResponse.receipt:type_name -> agent.OperationReceipt - 42, // 44: agent.GetExecutionPermitResponse.permit:type_name -> agent.ExecutionPermit - 10, // 45: agent.ApplyTaskControlRequest.meta:type_name -> agent.RequestMeta - 22, // 46: agent.ApplyTaskControlRequest.binding:type_name -> agent.ExecutionBinding - 4, // 47: agent.ApplyTaskControlRequest.action:type_name -> agent.ControlAction - 5, // 48: agent.ApplyTaskControlRequest.active_call_policy:type_name -> agent.ActiveCallPolicy - 13, // 49: agent.ApplyTaskControlResponse.receipt:type_name -> agent.OperationReceipt - 6, // 50: agent.ApplyTaskControlResponse.state:type_name -> agent.ExecutionState - 10, // 51: agent.QueryExecutionRequest.meta:type_name -> agent.RequestMeta - 22, // 52: agent.QueryExecutionRequest.binding:type_name -> agent.ExecutionBinding - 22, // 53: agent.ExecutionSnapshot.binding:type_name -> agent.ExecutionBinding - 6, // 54: agent.ExecutionSnapshot.state:type_name -> agent.ExecutionState - 23, // 55: agent.ExecutionSnapshot.assets:type_name -> agent.AssetDescriptor - 11, // 56: agent.QueryExecutionResponse.meta:type_name -> agent.ResponseMeta - 47, // 57: agent.QueryExecutionResponse.snapshot:type_name -> agent.ExecutionSnapshot - 12, // 58: agent.QueryExecutionResponse.failure:type_name -> agent.Failure - 22, // 59: agent.ExecutionFact.binding:type_name -> agent.ExecutionBinding - 9, // 60: agent.ExecutionFact.kind:type_name -> agent.FactKind - 10, // 61: agent.ReportExecutionEventRequest.meta:type_name -> agent.RequestMeta - 49, // 62: agent.ReportExecutionEventRequest.fact:type_name -> agent.ExecutionFact - 13, // 63: agent.ReportExecutionEventResponse.receipt:type_name -> agent.OperationReceipt - 10, // 64: agent.RequestUploadRequest.meta:type_name -> agent.RequestMeta - 22, // 65: agent.RequestUploadRequest.binding:type_name -> agent.ExecutionBinding - 23, // 66: agent.RequestUploadRequest.asset:type_name -> agent.AssetDescriptor - 24, // 67: agent.UploadGrant.headers:type_name -> agent.Header - 13, // 68: agent.RequestUploadResponse.receipt:type_name -> agent.OperationReceipt - 53, // 69: agent.RequestUploadResponse.grant:type_name -> agent.UploadGrant - 8, // 70: agent.RequestUploadResponse.state:type_name -> agent.UploadState - 10, // 71: agent.CompleteUploadRequest.meta:type_name -> agent.RequestMeta - 22, // 72: agent.CompleteUploadRequest.binding:type_name -> agent.ExecutionBinding - 23, // 73: agent.CompleteUploadRequest.asset:type_name -> agent.AssetDescriptor - 13, // 74: agent.CompleteUploadResponse.receipt:type_name -> agent.OperationReceipt - 8, // 75: agent.CompleteUploadResponse.state:type_name -> agent.UploadState - 25, // 76: agent.AgentControlService.GetAgentStatus:input_type -> agent.GetAgentStatusRequest - 27, // 77: agent.AgentControlService.ActivateAgent:input_type -> agent.ActivateAgentRequest - 29, // 78: agent.AgentControlService.GetBootstrap:input_type -> agent.GetBootstrapRequest - 31, // 79: agent.AgentControlService.SetAdmissionState:input_type -> agent.SetAdmissionStateRequest - 33, // 80: agent.AgentControlService.Execute:input_type -> agent.ExecuteRequest - 35, // 81: agent.AgentControlService.ExecuteAuthorized:input_type -> agent.ExecuteAuthorizedRequest - 37, // 82: agent.AgentControlService.ExecuteApproved:input_type -> agent.ExecuteApprovedRequest - 39, // 83: agent.AgentControlService.GetLoadedSIP:input_type -> agent.GetLoadedSIPRequest - 41, // 84: agent.AgentControlService.GetExecutionPermit:input_type -> agent.GetExecutionPermitRequest - 44, // 85: agent.AgentControlService.ApplyTaskControl:input_type -> agent.ApplyTaskControlRequest - 46, // 86: agent.AgentControlService.QueryExecution:input_type -> agent.QueryExecutionRequest - 50, // 87: agent.AgentControlService.ReportExecutionEvent:input_type -> agent.ReportExecutionEventRequest - 52, // 88: agent.AgentControlService.RequestUpload:input_type -> agent.RequestUploadRequest - 55, // 89: agent.AgentControlService.CompleteUpload:input_type -> agent.CompleteUploadRequest - 26, // 90: agent.AgentControlService.GetAgentStatus:output_type -> agent.GetAgentStatusResponse - 28, // 91: agent.AgentControlService.ActivateAgent:output_type -> agent.ActivateAgentResponse - 30, // 92: agent.AgentControlService.GetBootstrap:output_type -> agent.GetBootstrapResponse - 32, // 93: agent.AgentControlService.SetAdmissionState:output_type -> agent.SetAdmissionStateResponse - 34, // 94: agent.AgentControlService.Execute:output_type -> agent.ExecuteResponse - 36, // 95: agent.AgentControlService.ExecuteAuthorized:output_type -> agent.ExecuteAuthorizedResponse - 38, // 96: agent.AgentControlService.ExecuteApproved:output_type -> agent.ExecuteApprovedResponse - 40, // 97: agent.AgentControlService.GetLoadedSIP:output_type -> agent.GetLoadedSIPResponse - 43, // 98: agent.AgentControlService.GetExecutionPermit:output_type -> agent.GetExecutionPermitResponse - 45, // 99: agent.AgentControlService.ApplyTaskControl:output_type -> agent.ApplyTaskControlResponse - 48, // 100: agent.AgentControlService.QueryExecution:output_type -> agent.QueryExecutionResponse - 51, // 101: agent.AgentControlService.ReportExecutionEvent:output_type -> agent.ReportExecutionEventResponse - 54, // 102: agent.AgentControlService.RequestUpload:output_type -> agent.RequestUploadResponse - 56, // 103: agent.AgentControlService.CompleteUpload:output_type -> agent.CompleteUploadResponse - 90, // [90:104] is the sub-list for method output_type - 76, // [76:90] is the sub-list for method input_type - 76, // [76:76] is the sub-list for extension type_name - 76, // [76:76] is the sub-list for extension extendee - 0, // [0:76] is the sub-list for field type_name + 1, // 0: agent.Failure.code:type_name -> agent.FailureCode + 11, // 1: agent.OperationReceipt.meta:type_name -> agent.ResponseMeta + 0, // 2: agent.OperationReceipt.result:type_name -> agent.ResultCode + 12, // 3: agent.OperationReceipt.failure:type_name -> agent.Failure + 3, // 4: agent.AgentStatus.admission_state:type_name -> agent.AdmissionState + 15, // 5: agent.AgentStatus.capabilities:type_name -> agent.Capability + 16, // 6: agent.AgentStatus.resources:type_name -> agent.ResourceSample + 17, // 7: agent.AgentStatus.applied_configs:type_name -> agent.AppliedConfig + 7, // 8: agent.AssetDescriptor.kind:type_name -> agent.AssetKind + 10, // 9: agent.GetAgentStatusRequest.meta:type_name -> agent.RequestMeta + 14, // 10: agent.GetAgentStatusRequest.target:type_name -> agent.AgentBinding + 11, // 11: agent.GetAgentStatusResponse.meta:type_name -> agent.ResponseMeta + 18, // 12: agent.GetAgentStatusResponse.status:type_name -> agent.AgentStatus + 12, // 13: agent.GetAgentStatusResponse.failure:type_name -> agent.Failure + 10, // 14: agent.ActivateAgentRequest.meta:type_name -> agent.RequestMeta + 14, // 15: agent.ActivateAgentRequest.binding:type_name -> agent.AgentBinding + 11, // 16: agent.ActivateAgentResponse.meta:type_name -> agent.ResponseMeta + 2, // 17: agent.ActivateAgentResponse.state:type_name -> agent.ActivationState + 19, // 18: agent.ActivateAgentResponse.session:type_name -> agent.Session + 12, // 19: agent.ActivateAgentResponse.failure:type_name -> agent.Failure + 10, // 20: agent.GetBootstrapRequest.meta:type_name -> agent.RequestMeta + 11, // 21: agent.GetBootstrapResponse.meta:type_name -> agent.ResponseMeta + 2, // 22: agent.GetBootstrapResponse.state:type_name -> agent.ActivationState + 20, // 23: agent.GetBootstrapResponse.runtime_configs:type_name -> agent.ConfigReference + 21, // 24: agent.GetBootstrapResponse.upload_policy:type_name -> agent.UploadPolicy + 12, // 25: agent.GetBootstrapResponse.failure:type_name -> agent.Failure + 10, // 26: agent.SetAdmissionStateRequest.meta:type_name -> agent.RequestMeta + 14, // 27: agent.SetAdmissionStateRequest.target:type_name -> agent.AgentBinding + 3, // 28: agent.SetAdmissionStateRequest.state:type_name -> agent.AdmissionState + 13, // 29: agent.SetAdmissionStateResponse.receipt:type_name -> agent.OperationReceipt + 10, // 30: agent.ExecuteRequest.meta:type_name -> agent.RequestMeta + 22, // 31: agent.ExecuteRequest.binding:type_name -> agent.ExecutionBinding + 13, // 32: agent.ExecuteResponse.receipt:type_name -> agent.OperationReceipt + 6, // 33: agent.ExecuteResponse.state:type_name -> agent.ExecutionState + 10, // 34: agent.ExecuteAuthorizedRequest.meta:type_name -> agent.RequestMeta + 22, // 35: agent.ExecuteAuthorizedRequest.binding:type_name -> agent.ExecutionBinding + 13, // 36: agent.ExecuteAuthorizedResponse.receipt:type_name -> agent.OperationReceipt + 6, // 37: agent.ExecuteAuthorizedResponse.state:type_name -> agent.ExecutionState + 10, // 38: agent.ExecuteApprovedRequest.meta:type_name -> agent.RequestMeta + 10, // 39: agent.GetLoadedSIPRequest.meta:type_name -> agent.RequestMeta + 64, // 40: agent.GetLoadedSIPResponse.trunk_revision:type_name -> agent.GetLoadedSIPResponse.TrunkRevisionEntry + 10, // 41: agent.GetExecutionPermitRequest.meta:type_name -> agent.RequestMeta + 22, // 42: agent.GetExecutionPermitRequest.binding:type_name -> agent.ExecutionBinding + 13, // 43: agent.GetExecutionPermitResponse.receipt:type_name -> agent.OperationReceipt + 42, // 44: agent.GetExecutionPermitResponse.permit:type_name -> agent.ExecutionPermit + 10, // 45: agent.ApplyTaskControlRequest.meta:type_name -> agent.RequestMeta + 22, // 46: agent.ApplyTaskControlRequest.binding:type_name -> agent.ExecutionBinding + 4, // 47: agent.ApplyTaskControlRequest.action:type_name -> agent.ControlAction + 5, // 48: agent.ApplyTaskControlRequest.active_call_policy:type_name -> agent.ActiveCallPolicy + 13, // 49: agent.ApplyTaskControlResponse.receipt:type_name -> agent.OperationReceipt + 6, // 50: agent.ApplyTaskControlResponse.state:type_name -> agent.ExecutionState + 10, // 51: agent.QueryExecutionRequest.meta:type_name -> agent.RequestMeta + 22, // 52: agent.QueryExecutionRequest.binding:type_name -> agent.ExecutionBinding + 22, // 53: agent.ExecutionSnapshot.binding:type_name -> agent.ExecutionBinding + 6, // 54: agent.ExecutionSnapshot.state:type_name -> agent.ExecutionState + 23, // 55: agent.ExecutionSnapshot.assets:type_name -> agent.AssetDescriptor + 11, // 56: agent.QueryExecutionResponse.meta:type_name -> agent.ResponseMeta + 47, // 57: agent.QueryExecutionResponse.snapshot:type_name -> agent.ExecutionSnapshot + 12, // 58: agent.QueryExecutionResponse.failure:type_name -> agent.Failure + 22, // 59: agent.ExecutionFact.binding:type_name -> agent.ExecutionBinding + 9, // 60: agent.ExecutionFact.kind:type_name -> agent.FactKind + 10, // 61: agent.ReportExecutionEventRequest.meta:type_name -> agent.RequestMeta + 49, // 62: agent.ReportExecutionEventRequest.fact:type_name -> agent.ExecutionFact + 13, // 63: agent.ReportExecutionEventResponse.receipt:type_name -> agent.OperationReceipt + 10, // 64: agent.RequestUploadRequest.meta:type_name -> agent.RequestMeta + 22, // 65: agent.RequestUploadRequest.binding:type_name -> agent.ExecutionBinding + 23, // 66: agent.RequestUploadRequest.asset:type_name -> agent.AssetDescriptor + 24, // 67: agent.UploadGrant.headers:type_name -> agent.Header + 13, // 68: agent.RequestUploadResponse.receipt:type_name -> agent.OperationReceipt + 53, // 69: agent.RequestUploadResponse.grant:type_name -> agent.UploadGrant + 8, // 70: agent.RequestUploadResponse.state:type_name -> agent.UploadState + 10, // 71: agent.CompleteUploadRequest.meta:type_name -> agent.RequestMeta + 22, // 72: agent.CompleteUploadRequest.binding:type_name -> agent.ExecutionBinding + 23, // 73: agent.CompleteUploadRequest.asset:type_name -> agent.AssetDescriptor + 13, // 74: agent.CompleteUploadResponse.receipt:type_name -> agent.OperationReceipt + 8, // 75: agent.CompleteUploadResponse.state:type_name -> agent.UploadState + 10, // 76: agent.RequestRecordingUploadRequest.meta:type_name -> agent.RequestMeta + 23, // 77: agent.RequestRecordingUploadRequest.asset:type_name -> agent.AssetDescriptor + 53, // 78: agent.RequestRecordingUploadResponse.grant:type_name -> agent.UploadGrant + 10, // 79: agent.ReportCallEndedRequest.meta:type_name -> agent.RequestMeta + 13, // 80: agent.ReportCallEndedResponse.receipt:type_name -> agent.OperationReceipt + 10, // 81: agent.ReportCallResultRequest.meta:type_name -> agent.RequestMeta + 61, // 82: agent.ReportCallResultRequest.upload:type_name -> agent.UploadObservation + 13, // 83: agent.ReportCallResultResponse.receipt:type_name -> agent.OperationReceipt + 25, // 84: agent.AgentControlService.GetAgentStatus:input_type -> agent.GetAgentStatusRequest + 27, // 85: agent.AgentControlService.ActivateAgent:input_type -> agent.ActivateAgentRequest + 29, // 86: agent.AgentControlService.GetBootstrap:input_type -> agent.GetBootstrapRequest + 31, // 87: agent.AgentControlService.SetAdmissionState:input_type -> agent.SetAdmissionStateRequest + 33, // 88: agent.AgentControlService.Execute:input_type -> agent.ExecuteRequest + 35, // 89: agent.AgentControlService.ExecuteAuthorized:input_type -> agent.ExecuteAuthorizedRequest + 37, // 90: agent.AgentControlService.ExecuteApproved:input_type -> agent.ExecuteApprovedRequest + 39, // 91: agent.AgentControlService.GetLoadedSIP:input_type -> agent.GetLoadedSIPRequest + 41, // 92: agent.AgentControlService.GetExecutionPermit:input_type -> agent.GetExecutionPermitRequest + 44, // 93: agent.AgentControlService.ApplyTaskControl:input_type -> agent.ApplyTaskControlRequest + 46, // 94: agent.AgentControlService.QueryExecution:input_type -> agent.QueryExecutionRequest + 50, // 95: agent.AgentControlService.ReportExecutionEvent:input_type -> agent.ReportExecutionEventRequest + 52, // 96: agent.AgentControlService.RequestUpload:input_type -> agent.RequestUploadRequest + 55, // 97: agent.AgentControlService.CompleteUpload:input_type -> agent.CompleteUploadRequest + 57, // 98: agent.AgentControlService.RequestRecordingUpload:input_type -> agent.RequestRecordingUploadRequest + 59, // 99: agent.AgentControlService.ReportCallEnded:input_type -> agent.ReportCallEndedRequest + 62, // 100: agent.AgentControlService.ReportCallResult:input_type -> agent.ReportCallResultRequest + 26, // 101: agent.AgentControlService.GetAgentStatus:output_type -> agent.GetAgentStatusResponse + 28, // 102: agent.AgentControlService.ActivateAgent:output_type -> agent.ActivateAgentResponse + 30, // 103: agent.AgentControlService.GetBootstrap:output_type -> agent.GetBootstrapResponse + 32, // 104: agent.AgentControlService.SetAdmissionState:output_type -> agent.SetAdmissionStateResponse + 34, // 105: agent.AgentControlService.Execute:output_type -> agent.ExecuteResponse + 36, // 106: agent.AgentControlService.ExecuteAuthorized:output_type -> agent.ExecuteAuthorizedResponse + 38, // 107: agent.AgentControlService.ExecuteApproved:output_type -> agent.ExecuteApprovedResponse + 40, // 108: agent.AgentControlService.GetLoadedSIP:output_type -> agent.GetLoadedSIPResponse + 43, // 109: agent.AgentControlService.GetExecutionPermit:output_type -> agent.GetExecutionPermitResponse + 45, // 110: agent.AgentControlService.ApplyTaskControl:output_type -> agent.ApplyTaskControlResponse + 48, // 111: agent.AgentControlService.QueryExecution:output_type -> agent.QueryExecutionResponse + 51, // 112: agent.AgentControlService.ReportExecutionEvent:output_type -> agent.ReportExecutionEventResponse + 54, // 113: agent.AgentControlService.RequestUpload:output_type -> agent.RequestUploadResponse + 56, // 114: agent.AgentControlService.CompleteUpload:output_type -> agent.CompleteUploadResponse + 58, // 115: agent.AgentControlService.RequestRecordingUpload:output_type -> agent.RequestRecordingUploadResponse + 60, // 116: agent.AgentControlService.ReportCallEnded:output_type -> agent.ReportCallEndedResponse + 63, // 117: agent.AgentControlService.ReportCallResult:output_type -> agent.ReportCallResultResponse + 101, // [101:118] is the sub-list for method output_type + 84, // [84:101] is the sub-list for method input_type + 84, // [84:84] is the sub-list for extension type_name + 84, // [84:84] is the sub-list for extension extendee + 0, // [0:84] is the sub-list for field type_name } func init() { file_agent_agent_proto_init() } @@ -4902,7 +5404,7 @@ func file_agent_agent_proto_init() { GoPackagePath: reflect.TypeOf(x{}).PkgPath(), RawDescriptor: unsafe.Slice(unsafe.StringData(file_agent_agent_proto_rawDesc), len(file_agent_agent_proto_rawDesc)), NumEnums: 10, - NumMessages: 48, + NumMessages: 55, NumExtensions: 0, NumServices: 1, }, diff --git a/gen/agent/agent_grpc.pb.go b/gen/agent/agent_grpc.pb.go index c33646d..23b9c5d 100644 --- a/gen/agent/agent_grpc.pb.go +++ b/gen/agent/agent_grpc.pb.go @@ -19,20 +19,23 @@ import ( const _ = grpc.SupportPackageIsVersion9 const ( - AgentControlService_GetAgentStatus_FullMethodName = "/agent.AgentControlService/GetAgentStatus" - AgentControlService_ActivateAgent_FullMethodName = "/agent.AgentControlService/ActivateAgent" - AgentControlService_GetBootstrap_FullMethodName = "/agent.AgentControlService/GetBootstrap" - AgentControlService_SetAdmissionState_FullMethodName = "/agent.AgentControlService/SetAdmissionState" - AgentControlService_Execute_FullMethodName = "/agent.AgentControlService/Execute" - AgentControlService_ExecuteAuthorized_FullMethodName = "/agent.AgentControlService/ExecuteAuthorized" - AgentControlService_ExecuteApproved_FullMethodName = "/agent.AgentControlService/ExecuteApproved" - AgentControlService_GetLoadedSIP_FullMethodName = "/agent.AgentControlService/GetLoadedSIP" - AgentControlService_GetExecutionPermit_FullMethodName = "/agent.AgentControlService/GetExecutionPermit" - AgentControlService_ApplyTaskControl_FullMethodName = "/agent.AgentControlService/ApplyTaskControl" - AgentControlService_QueryExecution_FullMethodName = "/agent.AgentControlService/QueryExecution" - AgentControlService_ReportExecutionEvent_FullMethodName = "/agent.AgentControlService/ReportExecutionEvent" - AgentControlService_RequestUpload_FullMethodName = "/agent.AgentControlService/RequestUpload" - AgentControlService_CompleteUpload_FullMethodName = "/agent.AgentControlService/CompleteUpload" + AgentControlService_GetAgentStatus_FullMethodName = "/agent.AgentControlService/GetAgentStatus" + AgentControlService_ActivateAgent_FullMethodName = "/agent.AgentControlService/ActivateAgent" + AgentControlService_GetBootstrap_FullMethodName = "/agent.AgentControlService/GetBootstrap" + AgentControlService_SetAdmissionState_FullMethodName = "/agent.AgentControlService/SetAdmissionState" + AgentControlService_Execute_FullMethodName = "/agent.AgentControlService/Execute" + AgentControlService_ExecuteAuthorized_FullMethodName = "/agent.AgentControlService/ExecuteAuthorized" + AgentControlService_ExecuteApproved_FullMethodName = "/agent.AgentControlService/ExecuteApproved" + AgentControlService_GetLoadedSIP_FullMethodName = "/agent.AgentControlService/GetLoadedSIP" + AgentControlService_GetExecutionPermit_FullMethodName = "/agent.AgentControlService/GetExecutionPermit" + AgentControlService_ApplyTaskControl_FullMethodName = "/agent.AgentControlService/ApplyTaskControl" + AgentControlService_QueryExecution_FullMethodName = "/agent.AgentControlService/QueryExecution" + AgentControlService_ReportExecutionEvent_FullMethodName = "/agent.AgentControlService/ReportExecutionEvent" + AgentControlService_RequestUpload_FullMethodName = "/agent.AgentControlService/RequestUpload" + AgentControlService_CompleteUpload_FullMethodName = "/agent.AgentControlService/CompleteUpload" + AgentControlService_RequestRecordingUpload_FullMethodName = "/agent.AgentControlService/RequestRecordingUpload" + AgentControlService_ReportCallEnded_FullMethodName = "/agent.AgentControlService/ReportCallEnded" + AgentControlService_ReportCallResult_FullMethodName = "/agent.AgentControlService/ReportCallResult" ) // AgentControlServiceClient is the client API for AgentControlService service. @@ -58,6 +61,10 @@ type AgentControlServiceClient interface { ReportExecutionEvent(ctx context.Context, in *ReportExecutionEventRequest, opts ...grpc.CallOption) (*ReportExecutionEventResponse, error) RequestUpload(ctx context.Context, in *RequestUploadRequest, opts ...grpc.CallOption) (*RequestUploadResponse, error) CompleteUpload(ctx context.Context, in *CompleteUploadRequest, opts ...grpc.CallOption) (*CompleteUploadResponse, error) + // Agent requests a fresh bounded token explicitly; the Dispatcher owns the original target. + RequestRecordingUpload(ctx context.Context, in *RequestRecordingUploadRequest, opts ...grpc.CallOption) (*RequestRecordingUploadResponse, error) + ReportCallEnded(ctx context.Context, in *ReportCallEndedRequest, opts ...grpc.CallOption) (*ReportCallEndedResponse, error) + ReportCallResult(ctx context.Context, in *ReportCallResultRequest, opts ...grpc.CallOption) (*ReportCallResultResponse, error) } type agentControlServiceClient struct { @@ -208,6 +215,36 @@ func (c *agentControlServiceClient) CompleteUpload(ctx context.Context, in *Comp return out, nil } +func (c *agentControlServiceClient) RequestRecordingUpload(ctx context.Context, in *RequestRecordingUploadRequest, opts ...grpc.CallOption) (*RequestRecordingUploadResponse, error) { + cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) + out := new(RequestRecordingUploadResponse) + err := c.cc.Invoke(ctx, AgentControlService_RequestRecordingUpload_FullMethodName, in, out, cOpts...) + if err != nil { + return nil, err + } + return out, nil +} + +func (c *agentControlServiceClient) ReportCallEnded(ctx context.Context, in *ReportCallEndedRequest, opts ...grpc.CallOption) (*ReportCallEndedResponse, error) { + cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) + out := new(ReportCallEndedResponse) + err := c.cc.Invoke(ctx, AgentControlService_ReportCallEnded_FullMethodName, in, out, cOpts...) + if err != nil { + return nil, err + } + return out, nil +} + +func (c *agentControlServiceClient) ReportCallResult(ctx context.Context, in *ReportCallResultRequest, opts ...grpc.CallOption) (*ReportCallResultResponse, error) { + cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) + out := new(ReportCallResultResponse) + err := c.cc.Invoke(ctx, AgentControlService_ReportCallResult_FullMethodName, in, out, cOpts...) + if err != nil { + return nil, err + } + return out, nil +} + // AgentControlServiceServer is the server API for AgentControlService service. // All implementations must embed UnimplementedAgentControlServiceServer // for forward compatibility. @@ -231,6 +268,10 @@ type AgentControlServiceServer interface { ReportExecutionEvent(context.Context, *ReportExecutionEventRequest) (*ReportExecutionEventResponse, error) RequestUpload(context.Context, *RequestUploadRequest) (*RequestUploadResponse, error) CompleteUpload(context.Context, *CompleteUploadRequest) (*CompleteUploadResponse, error) + // Agent requests a fresh bounded token explicitly; the Dispatcher owns the original target. + RequestRecordingUpload(context.Context, *RequestRecordingUploadRequest) (*RequestRecordingUploadResponse, error) + ReportCallEnded(context.Context, *ReportCallEndedRequest) (*ReportCallEndedResponse, error) + ReportCallResult(context.Context, *ReportCallResultRequest) (*ReportCallResultResponse, error) mustEmbedUnimplementedAgentControlServiceServer() } @@ -283,6 +324,15 @@ func (UnimplementedAgentControlServiceServer) RequestUpload(context.Context, *Re func (UnimplementedAgentControlServiceServer) CompleteUpload(context.Context, *CompleteUploadRequest) (*CompleteUploadResponse, error) { return nil, status.Errorf(codes.Unimplemented, "method CompleteUpload not implemented") } +func (UnimplementedAgentControlServiceServer) RequestRecordingUpload(context.Context, *RequestRecordingUploadRequest) (*RequestRecordingUploadResponse, error) { + return nil, status.Errorf(codes.Unimplemented, "method RequestRecordingUpload not implemented") +} +func (UnimplementedAgentControlServiceServer) ReportCallEnded(context.Context, *ReportCallEndedRequest) (*ReportCallEndedResponse, error) { + return nil, status.Errorf(codes.Unimplemented, "method ReportCallEnded not implemented") +} +func (UnimplementedAgentControlServiceServer) ReportCallResult(context.Context, *ReportCallResultRequest) (*ReportCallResultResponse, error) { + return nil, status.Errorf(codes.Unimplemented, "method ReportCallResult not implemented") +} func (UnimplementedAgentControlServiceServer) mustEmbedUnimplementedAgentControlServiceServer() {} func (UnimplementedAgentControlServiceServer) testEmbeddedByValue() {} @@ -556,6 +606,60 @@ func _AgentControlService_CompleteUpload_Handler(srv interface{}, ctx context.Co return interceptor(ctx, in, info, handler) } +func _AgentControlService_RequestRecordingUpload_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { + in := new(RequestRecordingUploadRequest) + if err := dec(in); err != nil { + return nil, err + } + if interceptor == nil { + return srv.(AgentControlServiceServer).RequestRecordingUpload(ctx, in) + } + info := &grpc.UnaryServerInfo{ + Server: srv, + FullMethod: AgentControlService_RequestRecordingUpload_FullMethodName, + } + handler := func(ctx context.Context, req interface{}) (interface{}, error) { + return srv.(AgentControlServiceServer).RequestRecordingUpload(ctx, req.(*RequestRecordingUploadRequest)) + } + return interceptor(ctx, in, info, handler) +} + +func _AgentControlService_ReportCallEnded_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { + in := new(ReportCallEndedRequest) + if err := dec(in); err != nil { + return nil, err + } + if interceptor == nil { + return srv.(AgentControlServiceServer).ReportCallEnded(ctx, in) + } + info := &grpc.UnaryServerInfo{ + Server: srv, + FullMethod: AgentControlService_ReportCallEnded_FullMethodName, + } + handler := func(ctx context.Context, req interface{}) (interface{}, error) { + return srv.(AgentControlServiceServer).ReportCallEnded(ctx, req.(*ReportCallEndedRequest)) + } + return interceptor(ctx, in, info, handler) +} + +func _AgentControlService_ReportCallResult_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { + in := new(ReportCallResultRequest) + if err := dec(in); err != nil { + return nil, err + } + if interceptor == nil { + return srv.(AgentControlServiceServer).ReportCallResult(ctx, in) + } + info := &grpc.UnaryServerInfo{ + Server: srv, + FullMethod: AgentControlService_ReportCallResult_FullMethodName, + } + handler := func(ctx context.Context, req interface{}) (interface{}, error) { + return srv.(AgentControlServiceServer).ReportCallResult(ctx, req.(*ReportCallResultRequest)) + } + return interceptor(ctx, in, info, handler) +} + // AgentControlService_ServiceDesc is the grpc.ServiceDesc for AgentControlService service. // It's only intended for direct use with grpc.RegisterService, // and not to be introspected or modified (even as a copy) @@ -619,6 +723,18 @@ var AgentControlService_ServiceDesc = grpc.ServiceDesc{ MethodName: "CompleteUpload", Handler: _AgentControlService_CompleteUpload_Handler, }, + { + MethodName: "RequestRecordingUpload", + Handler: _AgentControlService_RequestRecordingUpload_Handler, + }, + { + MethodName: "ReportCallEnded", + Handler: _AgentControlService_ReportCallEnded_Handler, + }, + { + MethodName: "ReportCallResult", + Handler: _AgentControlService_ReportCallResult_Handler, + }, }, Streams: []grpc.StreamDesc{}, Metadata: "agent/agent.proto", diff --git a/internal/rpc/recording_server.go b/internal/rpc/recording_server.go new file mode 100644 index 0000000..663f196 --- /dev/null +++ b/internal/rpc/recording_server.go @@ -0,0 +1,159 @@ +package rpc + +import ( + "context" + "crypto/sha256" + "errors" + "fmt" + "log" + "path" + "strconv" + "time" + + agentpb "git.ipao.vip/rogee/go-sip/gen/agent" + "git.ipao.vip/rogee/go-sip/internal/oss" + "git.ipao.vip/rogee/go-sip/internal/store" + + "google.golang.org/grpc/codes" + "google.golang.org/grpc/status" +) + +// RecordingServer is the Dispatcher-only endpoint for the approved recording +// lifecycle. It does not accept the old upload binding or provide a fallback +// when the active Agent session, mTLS peer or original execution is unknown. +type RecordingServer struct { + agentpb.UnimplementedAgentControlServiceServer + Store *store.CurrentStore + OSS *oss.Client + DispatcherID string + TrustedFingerprints map[string]struct{} + AuthorizeSession func(*agentpb.RequestMeta) error + Now func() time.Time +} + +func (s *RecordingServer) clock() time.Time { + if s.Now != nil { + return s.Now().UTC() + } + return time.Now().UTC() +} + +func (s *RecordingServer) authorize(ctx context.Context, meta *agentpb.RequestMeta, dispatcherID string, tenantID int64, sourceEventID string) error { + if s == nil || s.Store == nil || s.OSS == nil || s.DispatcherID == "" || len(s.TrustedFingerprints) == 0 || s.AuthorizeSession == nil { + return status.Error(codes.FailedPrecondition, "Dispatcher recording authority is not configured") + } + if dispatcherID == "" || sourceEventID == "" || tenantID <= 0 { + return status.Error(codes.InvalidArgument, "recording fact requires original Dispatcher, tenant and call identity") + } + if dispatcherID != s.DispatcherID { + return status.Error(codes.PermissionDenied, "recording fact targets another Dispatcher") + } + if err := validateDispatcherPeer(ctx, meta, true, s.TrustedFingerprints, nil); err != nil { + return err + } + if err := s.AuthorizeSession(meta); err != nil { + log.Printf("recording RPC denied dispatcher=%s event=%s reason=inactive_agent_session", dispatcherID, sourceEventID) + return status.Error(codes.PermissionDenied, "Agent session is not active for this Dispatcher") + } + if err := s.Store.RequireReservedCall(dispatcherID, sourceEventID, tenantID); err != nil { + log.Printf("recording RPC denied dispatcher=%s event=%s reason=unbound_tenant_execution", dispatcherID, sourceEventID) + return status.Error(codes.FailedPrecondition, "recording source is not the approved tenant execution") + } + return nil +} + +func recordingObjectKey(prefix, dispatcherID string, tenantID int64, sourceEventID, recordingID string) string { + identity := sha256.Sum256([]byte(dispatcherID + "\x00" + strconv.FormatInt(tenantID, 10) + "\x00" + sourceEventID + "\x00" + recordingID)) + return path.Join(prefix, "recordings", fmt.Sprintf("%x.wav", identity)) +} + +func (s *RecordingServer) RequestRecordingUpload(ctx context.Context, req *agentpb.RequestRecordingUploadRequest) (*agentpb.RequestRecordingUploadResponse, error) { + if err := s.authorize(ctx, req.GetMeta(), req.GetDispatcherId(), req.GetTenantId(), req.GetSourceEventId()); err != nil { + return nil, err + } + asset := req.GetAsset() + if asset == nil || req.GetUploadId() == "" || asset.GetKind() != agentpb.AssetKind_ASSET_KIND_RECORDING || + asset.GetExecutionId() != req.GetSourceEventId() || asset.GetCallId() != req.GetSourceEventId() || asset.GetAssetId() == "" || + asset.GetFormat() != "wav" || asset.GetChannels() < 1 || asset.GetChannels() > 2 || asset.GetSampleRateHz() != 16000 || asset.GetDurationMs() < 0 || + asset.GetSizeBytes() <= 0 || asset.GetChecksumSha256() == "" { + return nil, status.Error(codes.InvalidArgument, "recording asset does not match the approved call and media profile") + } + config := s.OSS.Config() + objectKey := recordingObjectKey(config.KeyPrefix, req.GetDispatcherId(), req.GetTenantId(), req.GetSourceEventId(), asset.GetAssetId()) + grant, err := s.OSS.Grant(ctx, req.GetUploadId(), objectKey, asset.GetChecksumSha256(), asset.GetSizeBytes(), s.clock()) + if err != nil { + log.Printf("recording grant failed dispatcher=%s event=%s stage=presign", req.GetDispatcherId(), req.GetSourceEventId()) + return nil, status.Error(codes.FailedPrecondition, "bounded OSS recording grant unavailable") + } + if grant.GetBucket() != config.Bucket || grant.GetObjectKey() != objectKey || grant.GetUploadId() != req.GetUploadId() || grant.GetRequiredChecksumSha256() != asset.GetChecksumSha256() || grant.GetMaxBytes() != asset.GetSizeBytes() { + log.Printf("recording grant failed dispatcher=%s event=%s stage=signed_target_mismatch", req.GetDispatcherId(), req.GetSourceEventId()) + return nil, status.Error(codes.Internal, "OSS grant differs from approved original asset") + } + binding := store.CurrentRecordingGrant{ + DispatcherID: req.GetDispatcherId(), SourceEventID: req.GetSourceEventId(), UploadID: req.GetUploadId(), RecordingID: asset.GetAssetId(), + Bucket: grant.GetBucket(), ObjectKey: grant.GetObjectKey(), ChecksumSHA256: asset.GetChecksumSha256(), SizeBytes: asset.GetSizeBytes(), + Format: asset.GetFormat(), Channels: int(asset.GetChannels()), SampleRateHz: int(asset.GetSampleRateHz()), DurationMS: asset.GetDurationMs(), + } + _, created, err := s.Store.BindRecordingUpload(binding) + if err != nil { + log.Printf("recording grant failed dispatcher=%s event=%s stage=bind error_class=%T", req.GetDispatcherId(), req.GetSourceEventId(), err) + switch { + case errors.Is(err, store.ErrCurrentUploadConflict), errors.Is(err, store.ErrCurrentResultConflict): + return nil, status.Error(codes.AlreadyExists, "original recording target cannot be changed") + case errors.Is(err, store.ErrCurrentUploadAlreadyConfirmed): + return nil, status.Error(codes.FailedPrecondition, "recording is already confirmed; another PUT is forbidden") + default: + return nil, status.Error(codes.Internal, "original recording target could not be persisted") + } + } + log.Printf("recording grant bound dispatcher=%s event=%s new=%t", req.GetDispatcherId(), req.GetSourceEventId(), created) + return &agentpb.RequestRecordingUploadResponse{Grant: grant}, nil +} + +func (s *RecordingServer) ReportCallEnded(ctx context.Context, req *agentpb.ReportCallEndedRequest) (*agentpb.ReportCallEndedResponse, error) { + if err := s.authorize(ctx, req.GetMeta(), req.GetDispatcherId(), req.GetTenantId(), req.GetSourceEventId()); err != nil { + return nil, err + } + if err := s.Store.FinishExecute(req.GetDispatcherId(), req.GetSourceEventId()); err != nil { + log.Printf("call end not committed dispatcher=%s event=%s stage=ack_and_release error_class=%T", req.GetDispatcherId(), req.GetSourceEventId(), err) + return nil, status.Error(codes.FailedPrecondition, "confirmed call end could not be persisted with its acknowledgment") + } + log.Printf("call end committed dispatcher=%s event=%s", req.GetDispatcherId(), req.GetSourceEventId()) + return &agentpb.ReportCallEndedResponse{Receipt: &agentpb.OperationReceipt{FactId: req.GetSourceEventId(), Result: agentpb.ResultCode_RESULT_CODE_APPLIED, AcceptedAtUnixMs: s.clock().UnixMilli()}}, nil +} + +func (s *RecordingServer) ReportCallResult(ctx context.Context, req *agentpb.ReportCallResultRequest) (*agentpb.ReportCallResultResponse, error) { + if err := s.authorize(ctx, req.GetMeta(), req.GetDispatcherId(), req.GetTenantId(), req.GetSourceEventId()); err != nil { + return nil, err + } + if len(req.GetResultPayloadJson()) == 0 { + return nil, status.Error(codes.InvalidArgument, "final result payload is required") + } + var event store.CurrentOutboxEvent + var created bool + var err error + if observation := req.GetUpload(); observation != nil { + event, created, err = s.Store.RecordUploadedCallResult(req.GetDispatcherId(), req.GetSourceEventId(), req.GetResultPayloadJson(), store.CurrentUploadProof{ + UploadID: observation.GetUploadId(), RecordingID: observation.GetRecordingId(), StatusCode: int(observation.GetPutStatusCode()), + SizeBytes: observation.GetSizeBytes(), SHA256: observation.GetChecksumSha256(), + }) + } else { + event, created, err = s.Store.RecordCallResult(req.GetDispatcherId(), req.GetSourceEventId(), req.GetResultPayloadJson()) + } + if err != nil { + log.Printf("final result not committed dispatcher=%s event=%s stage=outbox error_class=%T", req.GetDispatcherId(), req.GetSourceEventId(), err) + switch { + case errors.Is(err, store.ErrCurrentResultInvalid): + return nil, status.Error(codes.InvalidArgument, "final result violates the approved MQ contract") + case errors.Is(err, store.ErrCurrentUploadUnverified), errors.Is(err, store.ErrCurrentEndUnconfirmed): + return nil, status.Error(codes.FailedPrecondition, "final result requires confirmed call end and original upload outcome") + case errors.Is(err, store.ErrCurrentResultConflict): + return nil, status.Error(codes.AlreadyExists, "call already has a different final result") + default: + return nil, status.Error(codes.Internal, "final result could not be persisted") + } + } + checksum := sha256.Sum256(event.Body) + log.Printf("final result committed dispatcher=%s event=%s outbox=%s new=%t", req.GetDispatcherId(), req.GetSourceEventId(), event.EventID, created) + return &agentpb.ReportCallResultResponse{Receipt: &agentpb.OperationReceipt{FactId: event.EventID, Result: agentpb.ResultCode_RESULT_CODE_ACCEPTED, ContentSha256: fmt.Sprintf("%x", checksum), AcceptedAtUnixMs: s.clock().UnixMilli()}}, nil +} diff --git a/internal/rpc/recording_server_flow_test.go b/internal/rpc/recording_server_flow_test.go new file mode 100644 index 0000000..463c453 --- /dev/null +++ b/internal/rpc/recording_server_flow_test.go @@ -0,0 +1,270 @@ +package rpc + +import ( + "context" + "crypto/tls" + "crypto/x509" + "encoding/json" + "errors" + "os" + "path/filepath" + "strings" + "testing" + "time" + + agentpb "git.ipao.vip/rogee/go-sip/gen/agent" + "git.ipao.vip/rogee/go-sip/internal/configread" + "git.ipao.vip/rogee/go-sip/internal/oss" + "git.ipao.vip/rogee/go-sip/internal/store" + + "google.golang.org/grpc/codes" + "google.golang.org/grpc/credentials" + "google.golang.org/grpc/peer" + "google.golang.org/grpc/status" + "google.golang.org/protobuf/proto" +) + +const recordingDispatcherID = "c046b893-8628-4589-ae50-619d049248a6" + +func recordingRPCFixture(t *testing.T) (*RecordingServer, *store.CurrentStore, context.Context, *agentpb.RequestRecordingUploadRequest, configread.CurrentSnapshot) { + t.Helper() + database, err := store.OpenCurrent(filepath.Join(t.TempDir(), "recordings.db")) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = database.Close() }) + readExample := func(name string, dst any) { + t.Helper() + raw, err := os.ReadFile(filepath.Join("..", "..", "contracts", "local", "examples", name+".json")) + if err != nil { + t.Fatal(err) + } + if err := json.Unmarshal(raw, dst); err != nil { + t.Fatal(err) + } + } + var snapshot configread.CurrentSnapshot + readExample("config-read-task-asr", &snapshot.Task) + readExample("config-read-sip", &snapshot.SIP) + readExample("config-read-quota", &snapshot.Quota) + var providers struct { + Providers []configread.CurrentProvider `json:"providers"` + } + readExample("config-read-providers", &providers) + snapshot.Providers = make(map[string]configread.CurrentProvider) + for _, provider := range providers.Providers { + snapshot.Providers[provider.ProviderRef] = provider + } + snapshot.SIP.Trunks = []byte(strings.Replace(string(snapshot.SIP.Trunks), `"max_concurrent_calls":null`, `"max_concurrent_calls":2`, 1)) + if err := database.ApplyDiscoverySnapshot(recordingDispatcherID, []configread.CurrentDiscoveredTask{{TaskID: snapshot.Task.TaskID, TenantID: snapshot.Task.TenantID, TaskRevision: snapshot.Task.TaskRevision, Status: "running"}}); err != nil { + t.Fatal(err) + } + if err := database.SaveSnapshot(snapshot); err != nil { + t.Fatal(err) + } + if err := database.MarkReadyForSIP(recordingDispatcherID, 8); err != nil { + t.Fatal(err) + } + command := store.CurrentExecuteCommand{DispatcherID: recordingDispatcherID, EventID: "call-recording-1", TenantID: 1001, TaskID: "task-asr", Callee: "15003164745", IssuedAt: "2026-09-21T01:30:00Z"} + if _, _, err := database.RecordExecute(command); err != nil { + t.Fatal(err) + } + now := time.Date(2026, 9, 21, 1, 30, 0, 0, time.UTC) + if err := database.ReserveExecute(recordingDispatcherID, command.EventID, store.CurrentCallReservation{ + TrunkID: "trunk-mock", SIPRevision: 8, CallerID: "BD00000000", DialedCallee: command.Callee, Deadline: now.Add(2 * time.Minute), + }, now); err != nil { + t.Fatal(err) + } + client, err := oss.NewClient(oss.Config{ + Endpoint: "https://oss.example.invalid", Region: "cn-test", Bucket: "mock-bucket", KeyPrefix: "approved", + AccessKeyID: "local-test-key", AccessKeySecret: "local-test-secret", GrantTTL: 15 * time.Minute, MaxAssetBytes: 1024, + }) + if err != nil { + t.Fatal(err) + } + leaf := &x509.Certificate{Raw: []byte("isolated-test-agent-peer")} + fingerprint := CertificateFingerprint(leaf) + ctx := peer.NewContext(context.Background(), &peer.Peer{AuthInfo: credentials.TLSInfo{State: tls.ConnectionState{VerifiedChains: [][]*x509.Certificate{{leaf}}}}}) + server := &RecordingServer{ + Store: database, OSS: client, DispatcherID: recordingDispatcherID, + TrustedFingerprints: map[string]struct{}{fingerprint: {}}, Now: func() time.Time { return now }, + AuthorizeSession: func(meta *agentpb.RequestMeta) error { + if meta.AgentId != "agent-mock" || meta.CellId != "cell-mock" || meta.BootId != "boot-mock" || meta.DispatcherEpoch != "epoch-7" || meta.SessionGeneration != 1 { + return errors.New("Agent boot/session is not activated") + } + return nil + }, + } + request := &agentpb.RequestRecordingUploadRequest{ + Meta: &agentpb.RequestMeta{AgentId: "agent-mock", CellId: "cell-mock", BootId: "boot-mock", DispatcherEpoch: "epoch-7", SessionGeneration: 1, OperationId: "request-1", IdempotencyKey: "request-1"}, + DispatcherId: recordingDispatcherID, TenantId: 1001, SourceEventId: command.EventID, UploadId: "upload-recording-1", + Asset: &agentpb.AssetDescriptor{Kind: agentpb.AssetKind_ASSET_KIND_RECORDING, ExecutionId: command.EventID, CallId: command.EventID, + AssetId: "recording-1", Format: "wav", Channels: 1, SampleRateHz: 16000, DurationMs: 3000, SizeBytes: 128, ChecksumSha256: strings.Repeat("a", 64)}, + } + return server, database, ctx, request, snapshot +} + +func recordingResultPayload(t *testing.T, task configread.CurrentSnapshot, grant *agentpb.UploadGrant) []byte { + t.Helper() + result := map[string]any{ + "task_id": "task-asr", "caller_profile_id": task.Task.CallerProfileID, + "callee": "15003164745", "trunk_id": "trunk-mock", + "started_at": "2026-09-21T01:30:01Z", "ended_at": "2026-09-21T01:30:06Z", "duration_ms": 5000, + "outcome": "no_answer", "reason_code": 486, "reason_message": "busy", + "transcript": []any{}, "opt_out": false, "recording": map[string]any{}, + } + if grant != nil { + result["outcome"] = "answered" + result["reason_code"] = 200 + result["reason_message"] = "completed" + result["recording"] = map[string]any{ + "status": "uploaded", "bucket": grant.GetBucket(), "object_key": grant.GetObjectKey(), + "format": "wav", "channels": 1, "sample_rate_hz": 16000, "duration_ms": 3000, + "size_bytes": 128, "checksum_sha256": strings.Repeat("a", 64), + } + } + raw, err := json.Marshal(result) + if err != nil { + t.Fatal(err) + } + return raw +} + +func TestRecordingServerAuthenticatesOriginalGrantAndExplicitReissue(t *testing.T) { + server, database, ctx, request, _ := recordingRPCFixture(t) + if _, err := server.RequestRecordingUpload(context.Background(), request); status.Code(err) != codes.Unauthenticated { + t.Fatalf("upload token issued without verified mTLS peer: %v", err) + } + if _, err := database.LoadRecordingUpload(request.DispatcherId, request.SourceEventId); !errors.Is(err, store.ErrCurrentUploadNotFound) { + t.Fatalf("unauthenticated Agent persisted a grant: %v", err) + } + wrongDispatcher := proto.Clone(request).(*agentpb.RequestRecordingUploadRequest) + wrongDispatcher.DispatcherId = "another-dispatcher" + if _, err := server.RequestRecordingUpload(ctx, wrongDispatcher); status.Code(err) != codes.PermissionDenied { + t.Fatalf("another Dispatcher obtained this D's token: %v", err) + } + wrongTenant := proto.Clone(request).(*agentpb.RequestRecordingUploadRequest) + wrongTenant.TenantId++ + if _, err := server.RequestRecordingUpload(ctx, wrongTenant); status.Code(err) != codes.FailedPrecondition { + t.Fatalf("foreign tenant obtained an OSS grant: %v", err) + } + wrongSession := proto.Clone(request).(*agentpb.RequestRecordingUploadRequest) + wrongSession.Meta.BootId = "not-active-boot" + if _, err := server.RequestRecordingUpload(ctx, wrongSession); status.Code(err) != codes.PermissionDenied { + t.Fatalf("changed Agent boot obtained an OSS token: %v", err) + } + response, err := server.RequestRecordingUpload(ctx, request) + grant := response.GetGrant() + if err != nil || grant.GetBucket() != "mock-bucket" || grant.GetObjectKey() == "" || grant.GetUploadId() != request.UploadId || grant.GetExpiresAtUnixMs() != server.Now().Add(15*time.Minute).UnixMilli() { + t.Fatalf("approved bound recording grant absent or changed: bucket=%q key=%q err=%v", grant.GetBucket(), grant.GetObjectKey(), err) + } + stored, err := database.LoadRecordingUpload(request.DispatcherId, request.SourceEventId) + if err != nil || stored.Bucket != grant.GetBucket() || stored.ObjectKey != grant.GetObjectKey() || stored.UploadID != request.UploadId || stored.ChecksumSHA256 != request.Asset.ChecksumSha256 { + t.Fatalf("Dispatcher did not persist original OSS asset: stored=%+v err=%v", stored, err) + } + reissuedResponse, err := server.RequestRecordingUpload(ctx, request) + reissued := reissuedResponse.GetGrant() + if err != nil || reissued.GetBucket() != grant.GetBucket() || reissued.GetObjectKey() != grant.GetObjectKey() || reissued.GetUploadId() != grant.GetUploadId() { + t.Fatalf("explicit token request changed original target: err=%v", err) + } + changedAsset := proto.Clone(request).(*agentpb.RequestRecordingUploadRequest) + changedAsset.Asset.AssetId = "another-recording" + if _, err := server.RequestRecordingUpload(ctx, changedAsset); status.Code(err) != codes.AlreadyExists { + t.Fatalf("new object replaced approved recording: %v", err) + } +} + +func TestRecordingServerConfirmedEndAndUploadedResultOutbox(t *testing.T) { + server, database, ctx, request, snapshot := recordingRPCFixture(t) + response, err := server.RequestRecordingUpload(ctx, request) + if err != nil { + t.Fatal(err) + } + grant := response.GetGrant() + result := &agentpb.ReportCallResultRequest{Meta: request.Meta, DispatcherId: request.DispatcherId, TenantId: request.TenantId, SourceEventId: request.SourceEventId, + ResultPayloadJson: recordingResultPayload(t, snapshot, grant), + Upload: &agentpb.UploadObservation{UploadId: request.UploadId, RecordingId: request.Asset.AssetId, PutStatusCode: 200, SizeBytes: 128, ChecksumSha256: request.Asset.ChecksumSha256}} + if _, err := server.ReportCallResult(ctx, result); status.Code(err) != codes.FailedPrecondition { + t.Fatalf("result escaped before confirmed end: %v", err) + } + end := &agentpb.ReportCallEndedRequest{Meta: request.Meta, DispatcherId: request.DispatcherId, TenantId: request.TenantId, SourceEventId: request.SourceEventId} + if _, err := server.ReportCallEnded(ctx, end); err != nil { + t.Fatalf("early Agent end did not persist original acknowledgment: %v", err) + } + if err := database.MarkExecuteDispatched(request.DispatcherId, request.SourceEventId); err != nil { + t.Fatalf("late originate RPC response broke ended call: %v", err) + } + if _, err := server.ReportCallEnded(ctx, end); err != nil { + t.Fatalf("repeated Agent end fact lost idempotency: %v", err) + } + withoutUpload := proto.Clone(result).(*agentpb.ReportCallResultRequest) + withoutUpload.Upload = nil + withoutUpload.ResultPayloadJson = recordingResultPayload(t, snapshot, nil) + if _, err := server.ReportCallResult(ctx, withoutUpload); status.Code(err) != codes.FailedPrecondition { + t.Fatalf("OSS failure was relabeled as no recording: %v", err) + } + badProof := proto.Clone(result).(*agentpb.ReportCallResultRequest) + badProof.Upload.PutStatusCode = 503 + if _, err := server.ReportCallResult(ctx, badProof); status.Code(err) != codes.FailedPrecondition { + t.Fatalf("failed PUT was reported as uploaded: %v", err) + } + first, err := server.ReportCallResult(ctx, result) + if err != nil || first.GetReceipt().GetFactId() == "" { + t.Fatalf("confirmed original upload could not report one result: fact=%q err=%v", first.GetReceipt().GetFactId(), err) + } + stored, err := database.LoadRecordingUpload(request.DispatcherId, request.SourceEventId) + if err != nil || stored.ConfirmedAt == "" { + t.Fatalf("uploaded fact was not committed with result outbox: state=%+v err=%v", stored, err) + } + duplicate, err := server.ReportCallResult(ctx, result) + if err != nil || duplicate.GetReceipt().GetFactId() != first.GetReceipt().GetFactId() { + t.Fatalf("duplicate result created another event: original=%q duplicate=%q err=%v", first.GetReceipt().GetFactId(), duplicate.GetReceipt().GetFactId(), err) + } + outbox, err := database.ListPendingOutbox(request.DispatcherId) + if err != nil || len(outbox) != 2 || outbox[0].EventType != "call.execute" || outbox[1].EventType != "call.execute.result" { + t.Fatalf("Agent end and upload did not produce exactly one ACK+result: events=%+v err=%v", outbox, err) + } + if _, err := server.RequestRecordingUpload(ctx, request); status.Code(err) != codes.FailedPrecondition { + t.Fatalf("confirmed uploaded asset got another PUT token: %v", err) + } +} + +func TestRecordingServerRejectsMalformedResultBeforeOutbox(t *testing.T) { + server, database, ctx, request, _ := recordingRPCFixture(t) + end := &agentpb.ReportCallEndedRequest{Meta: request.Meta, DispatcherId: request.DispatcherId, TenantId: request.TenantId, SourceEventId: request.SourceEventId} + if _, err := server.ReportCallEnded(ctx, end); err != nil { + t.Fatal(err) + } + bad := &agentpb.ReportCallResultRequest{Meta: request.Meta, DispatcherId: request.DispatcherId, TenantId: request.TenantId, SourceEventId: request.SourceEventId, ResultPayloadJson: []byte(`{"task_id":`)} + if _, err := server.ReportCallResult(ctx, bad); status.Code(err) != codes.InvalidArgument { + t.Fatalf("malformed result was not reported as bad Agent input: %v", err) + } + pending, err := database.ListPendingOutbox(request.DispatcherId) + if err != nil || len(pending) != 1 || pending[0].EventType != "call.execute" { + t.Fatalf("malformed payload escaped into final result: events=%+v err=%v", pending, err) + } +} + +func TestRecordingServerNoRecordingResultCannotGainLaterGrant(t *testing.T) { + server, database, ctx, request, snapshot := recordingRPCFixture(t) + end := &agentpb.ReportCallEndedRequest{Meta: request.Meta, DispatcherId: request.DispatcherId, TenantId: request.TenantId, SourceEventId: request.SourceEventId} + if _, err := server.ReportCallEnded(ctx, end); err != nil { + t.Fatal(err) + } + result := &agentpb.ReportCallResultRequest{Meta: request.Meta, DispatcherId: request.DispatcherId, TenantId: request.TenantId, SourceEventId: request.SourceEventId, + ResultPayloadJson: recordingResultPayload(t, snapshot, nil)} + first, err := server.ReportCallResult(ctx, result) + if err != nil || first.GetReceipt().GetFactId() == "" { + t.Fatalf("no-recording confirmed end was not reported: fact=%q err=%v", first.GetReceipt().GetFactId(), err) + } + if repeated, err := server.ReportCallResult(ctx, result); err != nil || repeated.GetReceipt().GetFactId() != first.GetReceipt().GetFactId() { + t.Fatalf("no-recording result replay changed identity: err=%v", err) + } + if _, err := server.RequestRecordingUpload(ctx, request); status.Code(err) != codes.AlreadyExists { + t.Fatalf("final no-recording result gained a later object: %v", err) + } + outbox, err := database.ListPendingOutbox(request.DispatcherId) + if err != nil || len(outbox) != 2 { + t.Fatalf("no-recording end has multiple or missing results: events=%+v err=%v", outbox, err) + } +} diff --git a/internal/rpc/recording_server_test.go b/internal/rpc/recording_server_test.go new file mode 100644 index 0000000..63a7c92 --- /dev/null +++ b/internal/rpc/recording_server_test.go @@ -0,0 +1,36 @@ +package rpc + +import ( + "context" + "testing" + + agentpb "git.ipao.vip/rogee/go-sip/gen/agent" + + "google.golang.org/grpc/codes" + "google.golang.org/grpc/status" +) + +func TestRecordingServerRequiresConfiguredAuthority(t *testing.T) { + server := &RecordingServer{} + for _, call := range []struct { + name string + call func() error + }{ + {"grant", func() error { + _, err := server.RequestRecordingUpload(context.Background(), &agentpb.RequestRecordingUploadRequest{}) + return err + }}, + {"end", func() error { + _, err := server.ReportCallEnded(context.Background(), &agentpb.ReportCallEndedRequest{}) + return err + }}, + {"result", func() error { + _, err := server.ReportCallResult(context.Background(), &agentpb.ReportCallResultRequest{}) + return err + }}, + } { + if err := call.call(); status.Code(err) != codes.FailedPrecondition { + t.Fatalf("%s without original D authority was not rejected: %v", call.name, err) + } + } +} diff --git a/internal/store/current_result.go b/internal/store/current_result.go index 939caa8..76c12c3 100644 --- a/internal/store/current_result.go +++ b/internal/store/current_result.go @@ -16,6 +16,8 @@ import ( var ErrCurrentUploadUnverified = errors.New("recording outcome is not verified against its Dispatcher-approved upload") var ErrCurrentResultConflict = errors.New("call already has a different final result") +var ErrCurrentEndUnconfirmed = errors.New("final result requires confirmed call end and its frozen snapshot") +var ErrCurrentResultInvalid = errors.New("final result input violates the approved contract") func currentResultEventID(dispatcherID, sourceEventID string) string { identity := sha256.Sum256([]byte("call.execute.result\x00" + dispatcherID + "\x00" + sourceEventID)) @@ -37,7 +39,7 @@ func (s *CurrentStore) RecordUploadedCallResult(dispatcherID, sourceEventID stri func (s *CurrentStore) recordCallResult(dispatcherID, sourceEventID string, payload []byte, proof *CurrentUploadProof) (CurrentOutboxEvent, bool, error) { if dispatcherID == "" || sourceEventID == "" || len(sourceEventID) > 255 || len(payload) == 0 { - return CurrentOutboxEvent{}, false, errors.New("final result requires a durable call identity and payload") + return CurrentOutboxEvent{}, false, fmt.Errorf("%w: durable call identity and payload are required", ErrCurrentResultInvalid) } var result struct { TaskID string `json:"task_id"` @@ -49,15 +51,15 @@ func (s *CurrentStore) recordCallResult(dispatcherID, sourceEventID string, payl Recording json.RawMessage `json:"recording"` } if err := json.Unmarshal(payload, &result); err != nil { - return CurrentOutboxEvent{}, false, fmt.Errorf("decode final result: %w", err) + return CurrentOutboxEvent{}, false, fmt.Errorf("%w: decode JSON: %v", ErrCurrentResultInvalid, err) } start, err := time.Parse(time.RFC3339Nano, result.StartedAt) if err != nil { - return CurrentOutboxEvent{}, false, errors.New("final result has invalid start time") + return CurrentOutboxEvent{}, false, fmt.Errorf("%w: invalid start time", ErrCurrentResultInvalid) } end, err := time.Parse(time.RFC3339Nano, result.EndedAt) if err != nil || end.Before(start) { - return CurrentOutboxEvent{}, false, errors.New("final result end precedes start or is invalid") + return CurrentOutboxEvent{}, false, fmt.Errorf("%w: end precedes start or is invalid", ErrCurrentResultInvalid) } tx, err := s.db.Begin() if err != nil { @@ -75,7 +77,7 @@ func (s *CurrentStore) recordCallResult(dispatcherID, sourceEventID string, payl return CurrentOutboxEvent{}, false, fmt.Errorf("load completed call identity: %w", err) } if status != "finished" || trunkID == "" || len(snapshotJSON) == 0 { - return CurrentOutboxEvent{}, false, errors.New("final result requires confirmed call end and its frozen snapshot") + return CurrentOutboxEvent{}, false, ErrCurrentEndUnconfirmed } var snapshot struct { Task configread.CurrentTask `json:"task"` @@ -98,7 +100,7 @@ func (s *CurrentStore) recordCallResult(dispatcherID, sourceEventID string, payl ChecksumSHA256 string `json:"checksum_sha256"` } if err := json.Unmarshal(result.Recording, &recording); err != nil { - return CurrentOutboxEvent{}, false, errors.New("final result recording fact is invalid") + return CurrentOutboxEvent{}, false, fmt.Errorf("%w: recording fact is invalid", ErrCurrentResultInvalid) } grant, grantErr := loadRecordingUpload(tx.QueryRow(`SELECT dispatcher_id,source_event_id,upload_id,recording_id,bucket,object_key,checksum_sha256,size_bytes,format,channels,sample_rate_hz,duration_ms,COALESCE(confirmed_at,'') FROM dispatcher_recordings WHERE dispatcher_id=? AND source_event_id=?`, dispatcherID, sourceEventID)) if grantErr != nil && !errors.Is(grantErr, ErrCurrentUploadNotFound) { @@ -130,7 +132,7 @@ func (s *CurrentStore) recordCallResult(dispatcherID, sourceEventID string, payl return CurrentOutboxEvent{}, false, fmt.Errorf("encode final result: %w", err) } if err := contract.ValidateCurrent("mq", body); err != nil { - return CurrentOutboxEvent{}, false, fmt.Errorf("final result violates current MQ contract: %w", err) + return CurrentOutboxEvent{}, false, fmt.Errorf("%w: MQ contract: %v", ErrCurrentResultInvalid, err) } inserted, err := tx.Exec(`INSERT INTO dispatcher_outbox(dispatcher_id,event_id,event_type,routing_key,body) VALUES(?,?,?,?,?) ON CONFLICT(dispatcher_id,event_id) DO NOTHING`, dispatcherID, eventID, "call.execute.result", route.BindingKey, body) if err != nil { diff --git a/internal/store/current_upload.go b/internal/store/current_upload.go index 1f0ee2f..d95130d 100644 --- a/internal/store/current_upload.go +++ b/internal/store/current_upload.go @@ -41,6 +41,25 @@ type CurrentRecordingGrant struct { ConfirmedAt string } +// RequireReservedCall ties Agent-reported facts to the numeric tenant and the +// durable, already reserved execution. Callers must also authenticate the +// active Agent session before accepting any report or signing an OSS token. +func (s *CurrentStore) RequireReservedCall(dispatcherID, sourceEventID string, tenantID int64) error { + if dispatcherID == "" || sourceEventID == "" || tenantID <= 0 { + return errors.New("recording source lacks a Dispatcher execution identity") + } + var storedTenant int64 + var status, trunkID string + var snapshot []byte + if err := s.db.QueryRow(`SELECT tenant_id,status,COALESCE(selected_trunk_id,''),snapshot_json FROM dispatcher_inbox WHERE dispatcher_id=? AND event_id=?`, dispatcherID, sourceEventID).Scan(&storedTenant, &status, &trunkID, &snapshot); err != nil { + return fmt.Errorf("load reserved call identity: %w", err) + } + if storedTenant != tenantID || trunkID == "" || len(snapshot) == 0 || (status != "dispatching" && status != "dispatched" && status != "unknown" && status != "finished") { + return errors.New("recording source is not the approved tenant execution") + } + return nil +} + func (s *CurrentStore) BindRecordingUpload(proposed CurrentRecordingGrant) (CurrentRecordingGrant, bool, error) { if proposed.DispatcherID == "" || proposed.SourceEventID == "" || proposed.UploadID == "" || proposed.RecordingID == "" || proposed.Bucket == "" || proposed.ObjectKey == "" || proposed.SizeBytes <= 0 || proposed.Format != "wav" || (proposed.Channels != 1 && proposed.Channels != 2) || proposed.SampleRateHz != 16000 || proposed.DurationMS < 0 || proposed.ConfirmedAt != "" { return CurrentRecordingGrant{}, false, errors.New("recording grant lacks a bounded original asset and execution") @@ -65,7 +84,7 @@ func (s *CurrentStore) BindRecordingUpload(proposed CurrentRecordingGrant) (Curr if err != nil { return CurrentRecordingGrant{}, false, fmt.Errorf("load approved call for recording: %w", err) } - if status != "dispatched" && status != "unknown" && status != "finished" || trunkID == "" || len(snapshot) == 0 { + if status != "dispatching" && status != "dispatched" && status != "unknown" && status != "finished" || trunkID == "" || len(snapshot) == 0 { return CurrentRecordingGrant{}, false, errors.New("recording grant requires an already reserved approved call") } inserted, err := tx.Exec(`INSERT INTO dispatcher_recordings(dispatcher_id,source_event_id,upload_id,recording_id,bucket,object_key,checksum_sha256,size_bytes,format,channels,sample_rate_hz,duration_ms) diff --git a/internal/store/current_upload_test.go b/internal/store/current_upload_test.go index f58d631..7285542 100644 --- a/internal/store/current_upload_test.go +++ b/internal/store/current_upload_test.go @@ -7,6 +7,7 @@ import ( "reflect" "strings" "testing" + "time" ) func currentUploadBinding(cmd CurrentExecuteCommand) CurrentRecordingGrant { @@ -19,6 +20,27 @@ func currentUploadBinding(cmd CurrentExecuteCommand) CurrentRecordingGrant { } } +func TestCurrentUploadBoundWhileOriginateAcknowledgmentIsInFlight(t *testing.T) { + s := preparedCurrentCallStore(t) + cmd := currentCall("grant-during-dispatch") + if _, _, err := s.RecordExecute(cmd); err != nil { + t.Fatal(err) + } + if err := s.ReserveExecute(cmd.DispatcherID, cmd.EventID, currentReservation(), time.Date(2026, 9, 21, 1, 30, 0, 0, time.UTC)); err != nil { + t.Fatal(err) + } + binding := currentUploadBinding(cmd) + if _, created, err := s.BindRecordingUpload(binding); err != nil || !created { + t.Fatalf("approved Agent could not request a recording token before originate RPC returned: created=%t err=%v", created, err) + } + if err := s.MarkExecuteDispatched(cmd.DispatcherID, cmd.EventID); err != nil { + t.Fatalf("late originate acknowledgment conflicted with original upload: %v", err) + } + if _, created, err := s.BindRecordingUpload(binding); err != nil || created { + t.Fatalf("late acknowledgment replaced the bound asset: created=%t err=%v", created, err) + } +} + func TestCurrentUploadBindingPersistsOnlyOneOriginalAssetAcrossRestart(t *testing.T) { s := preparedCurrentCallStore(t) cmd := currentResultCall(t, s, "grant-stable-1", true) @@ -94,6 +116,35 @@ func TestCurrentUploadBindingRejectsChangedTargetAndUnreservedCalls(t *testing.T } } +func TestCurrentRecordingSourceMustMatchReservedCallAndNumericTenant(t *testing.T) { + s := preparedCurrentCallStore(t) + cmd := currentResultCall(t, s, "grant-auth-source", true) + if err := s.RequireReservedCall(cmd.DispatcherID, cmd.EventID, cmd.TenantID); err != nil { + t.Fatalf("approved call identity rejected: %v", err) + } + for _, check := range []struct { + name string + dispatcherID string + sourceEventID string + tenantID int64 + }{ + {"foreign dispatcher", "other-dispatcher", cmd.EventID, cmd.TenantID}, + {"foreign tenant", cmd.DispatcherID, cmd.EventID, cmd.TenantID + 1}, + {"unknown call", cmd.DispatcherID, "unknown-execution", cmd.TenantID}, + } { + if err := s.RequireReservedCall(check.dispatcherID, check.sourceEventID, check.tenantID); err == nil { + t.Fatalf("%s was accepted as original approved call", check.name) + } + } + pending := currentCall("grant-not-authorized") + if _, _, err := s.RecordExecute(pending); err != nil { + t.Fatal(err) + } + if err := s.RequireReservedCall(pending.DispatcherID, pending.EventID, pending.TenantID); err == nil { + t.Fatal("unreserved instruction was accepted as an approved recording source") + } +} + func TestCurrentUploadBindingRejectsReusedUploadIDOrOSSObject(t *testing.T) { s := preparedCurrentCallStore(t) first := currentUploadBinding(currentResultCall(t, s, "grant-unique-1", true)) diff --git a/proto/agent/agent.proto b/proto/agent/agent.proto index acf2560..f574076 100644 --- a/proto/agent/agent.proto +++ b/proto/agent/agent.proto @@ -23,6 +23,10 @@ service AgentControlService { rpc ReportExecutionEvent(ReportExecutionEventRequest) returns (ReportExecutionEventResponse); rpc RequestUpload(RequestUploadRequest) returns (RequestUploadResponse); rpc CompleteUpload(CompleteUploadRequest) returns (CompleteUploadResponse); + // Agent requests a fresh bounded token explicitly; the Dispatcher owns the original target. + rpc RequestRecordingUpload(RequestRecordingUploadRequest) returns (RequestRecordingUploadResponse); + rpc ReportCallEnded(ReportCallEndedRequest) returns (ReportCallEndedResponse); + rpc ReportCallResult(ReportCallResultRequest) returns (ReportCallResultResponse); } enum ResultCode { @@ -519,3 +523,50 @@ message CompleteUploadResponse { reserved 3; reserved "oss_id"; } + +// These facts refer only to the Dispatcher-approved execution identified by +// source_event_id; no audio bytes or temporary credential is persisted in MQ. +message RequestRecordingUploadRequest { + RequestMeta meta = 1; + string dispatcher_id = 2; + int64 tenant_id = 3; + string source_event_id = 4; + string upload_id = 5; + AssetDescriptor asset = 6; +} + +message RequestRecordingUploadResponse { + UploadGrant grant = 1; +} + +message ReportCallEndedRequest { + RequestMeta meta = 1; + string dispatcher_id = 2; + int64 tenant_id = 3; + string source_event_id = 4; +} + +message ReportCallEndedResponse { + OperationReceipt receipt = 1; +} + +message UploadObservation { + string upload_id = 1; + string recording_id = 2; + int32 put_status_code = 3; + int64 size_bytes = 4; + string checksum_sha256 = 5; +} + +message ReportCallResultRequest { + RequestMeta meta = 1; + string dispatcher_id = 2; + int64 tenant_id = 3; + string source_event_id = 4; + bytes result_payload_json = 5; + UploadObservation upload = 6; +} + +message ReportCallResultResponse { + OperationReceipt receipt = 1; +} diff --git a/proto/manifest.json b/proto/manifest.json index cef59b5..4865690 100644 --- a/proto/manifest.json +++ b/proto/manifest.json @@ -34,18 +34,18 @@ }, { "path": "proto/agent/agent.proto", - "bytes": 13281, - "sha256": "dc086c1963fa898d39f1f9e10de85abf3c4449834a84f0a87ec7b5e7aa281f91" + "bytes": 14721, + "sha256": "25f888bd3b5097a2d637bf4d1b4c97498102a5f7c67919d41d28e65413a296f7" }, { "path": "gen/agent/agent.pb.go", - "bytes": 163322, - "sha256": "a9da5dc255d2c5d3af736fcdd31909ebd59409f07bb4cef4de578d7baaf40563" + "bytes": 180766, + "sha256": "1f43d3f7e039f9cfd1fda248e71491eda00f25de6d4a06091a9a3328302b0ccd" }, { "path": "gen/agent/agent_grpc.pb.go", - "bytes": 28480, - "sha256": "1944231ebc402160813430db0a21f8ec9003ff8f58db0dc9d2b7013d7285da49" + "bytes": 34240, + "sha256": "4d9ef438fc4a32508ed055bc0cae6f6559b05530835c19851ebf2951680bc08d" } ] }