From 793b6448e28acbaec58a3c760e695d8c3a2d1861 Mon Sep 17 00:00:00 2001 From: Rogee Date: Wed, 30 Sep 2026 13:21:37 +0800 Subject: [PATCH] refactor(rpc): remove retired Agent business handlers --- .../saas-dispatcher-implementation.md | 1 + internal/rpc/ai_authorization_test.go | 122 ---- internal/rpc/approved_execution_test.go | 31 + internal/rpc/authorized_execution.go | 123 ---- internal/rpc/authorized_execution_test.go | 201 ------ internal/rpc/calllog_test.go | 100 --- internal/rpc/control_policy_test.go | 32 - .../rpc/control_policy_validation_test.go | 48 -- .../rpc/execution_journal_failure_test.go | 42 -- internal/rpc/execution_journal_mode_test.go | 25 - internal/rpc/execution_journal_test.go | 74 --- internal/rpc/server.go | 589 ------------------ internal/rpc/server_test.go | 140 +---- internal/rpc/service_test.go | 13 + internal/rpc/upload_test.go | 71 --- 15 files changed, 47 insertions(+), 1565 deletions(-) delete mode 100644 internal/rpc/ai_authorization_test.go delete mode 100644 internal/rpc/authorized_execution.go delete mode 100644 internal/rpc/authorized_execution_test.go delete mode 100644 internal/rpc/calllog_test.go delete mode 100644 internal/rpc/control_policy_test.go delete mode 100644 internal/rpc/control_policy_validation_test.go delete mode 100644 internal/rpc/execution_journal_failure_test.go delete mode 100644 internal/rpc/execution_journal_mode_test.go delete mode 100644 internal/rpc/execution_journal_test.go delete mode 100644 internal/rpc/upload_test.go diff --git a/docs/evidence/saas-dispatcher-implementation.md b/docs/evidence/saas-dispatcher-implementation.md index b8ba754..52b2c80 100644 --- a/docs/evidence/saas-dispatcher-implementation.md +++ b/docs/evidence/saas-dispatcher-implementation.md @@ -116,6 +116,7 @@ - Agent 旧执行状态屏障:批准的 Mock Agent 启动时若发现旧执行日志,立即拒绝启动并明确要求人工处置;隔离测试先复现了原先允许启动的问题,再证明旧日志字节不变、新会话状态未写入。此屏障不迁移、不清理存量文件;没有检查或处置真实环境中的旧记录。 - Agent Proto 唯一服务面:新增服务方法集合测试先确认旧 RPC 仍可被服务描述符发现,随后将 `AgentControlService` 收敛为状态/激活、获批执行与控制、加载版本及录音/结果八个方法;重新生成 Go 类型、核验七项来源哈希,并同步改写当前 Proto 错误与恢复说明。旧消息定义、旧服务端实现及专属测试尚待清理,不能把本批当作 P07 完成或真实 Agent 联调。 - Agent 旧失败事实分支:旧 Mock 上传失败事实通过已退役的 `ReportExecutionEvent` RPC 回报,现删除该代码及专属测试。现行录音失败、未知 PUT、重启恢复和最终结果仍由 `recording_delivery*` 隔离测试覆盖;不复活额外通话事件或把未知上传当作成功。 +- Agent 旧业务 RPC 处理器:先以服务结构测试复现旧 `GetBootstrap` 等十个方法仍存在,再删除旧执行、许可、控制、查询、事件和上传处理器及仅依赖旧服务面的专属测试;原混合测试保留会话代际、证书指纹、状态和真实 gRPC 激活。另将旧执行日志写入失败保护迁至现行获批执行测试:日志目录不可写时两次相同请求均不得接受或发起呼叫,修复目录后只发起一次,结果未知时拒绝重拨。`go test ./...`、`go test -race ./...`、`go vet ./...`、`go build ./...`、Proto 来源/hash、当前合同及临时 RabbitMQ/HTTPS/双向 TLS 隔离链路均通过;旧执行日志结构、未使用的 Proto 消息与其他引用仍须继续清理,未接触现存业务数据或真实外部服务。 ## 验收台账 diff --git a/internal/rpc/ai_authorization_test.go b/internal/rpc/ai_authorization_test.go deleted file mode 100644 index b880b5f..0000000 --- a/internal/rpc/ai_authorization_test.go +++ /dev/null @@ -1,122 +0,0 @@ -package rpc - -import ( - "context" - "testing" - "time" - - "git.ipao.vip/rogee/go-sip/contracts" - agentpb "git.ipao.vip/rogee/go-sip/gen/agent" - "git.ipao.vip/rogee/go-sip/internal/ai" - "git.ipao.vip/rogee/go-sip/internal/contract" - "git.ipao.vip/rogee/go-sip/internal/testfixture" - "google.golang.org/grpc/codes" - "google.golang.org/grpc/status" -) - -func TestExecutionPermitEnforcesAIAuthorization(t *testing.T) { - snapshotRaw, err := contracts.Read("examples/agent-version-asr-only.json") - if err != nil { - t.Fatal(err) - } - snapshot, err := ai.ValidateForMode(snapshotRaw, ai.ModeASROnly) - if err != nil { - t.Fatal(err) - } - authorizationRaw, err := contracts.Files.ReadFile("local/v0.3/examples/ai-authorization-v0.2.json") - if err != nil { - t.Fatal(err) - } - server := NewServer(ServerOptions{ - Now: func() time.Time { return time.Date(2026, 9, 18, 0, 0, 30, 0, time.UTC) }, - AISnapshotRaw: snapshotRaw, - AIAuthorizationRaw: authorizationRaw, - }) - activateTestServer(t, server) - - raw, err := testfixture.Execute() - if err != nil { - t.Fatal(err) - } - _, payload, err := contract.DecodeExecute(raw) - if err != nil { - t.Fatal(err) - } - binding := &agentpb.ExecutionBinding{ - TenantId: "tenant-1", - TenantKey: "tenant-demo-key", - ExecutionId: payload.ExecutionID, - TaskId: payload.TaskID, - TaskItemId: payload.TaskItemID, - TaskRevision: payload.TaskRevision, - AgentVersionId: snapshot.AgentVersionID, - RoutePolicyId: payload.RoutePolicyID, - CallerProfileId: payload.CallerProfileID, - } - response, err := server.GetExecutionPermit(context.Background(), &agentpb.GetExecutionPermitRequest{ - Meta: testMeta("permit-ai", "permit-ai-key", 1), - Binding: binding, - ResourceReservationId: "reservation-ai", - ExpectedTaskRevision: payload.TaskRevision, - ConfigSha256: snapshot.Digest, - }) - if err != nil { - t.Fatal(err) - } - if response.Permit == nil || response.Receipt == nil || response.Receipt.Result != agentpb.ResultCode_RESULT_CODE_APPLIED { - t.Fatalf("unexpected authorized permit response: %+v", response) - } - - badDigest := &agentpb.GetExecutionPermitRequest{ - Meta: testMeta("permit-ai-bad-digest", "permit-ai-bad-digest-key", 1), - Binding: binding, - ResourceReservationId: "reservation-ai-2", - ExpectedTaskRevision: payload.TaskRevision, - ConfigSha256: "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa", - } - if _, err := server.GetExecutionPermit(context.Background(), badDigest); status.Code(err) != codes.FailedPrecondition { - t.Fatalf("bad digest error=%v, code=%s", err, status.Code(err)) - } -} - -func TestExecutionPermitRejectsRevokedAIAuthorization(t *testing.T) { - snapshotRaw, err := contracts.Read("examples/agent-version-asr-only.json") - if err != nil { - t.Fatal(err) - } - snapshot, err := ai.ValidateForMode(snapshotRaw, ai.ModeASROnly) - if err != nil { - t.Fatal(err) - } - authorizationRaw, err := contracts.Files.ReadFile("local/v0.3/examples/invalid-ai-authorization-revoked-v0.2.json") - if err != nil { - t.Fatal(err) - } - server := NewServer(ServerOptions{ - Now: func() time.Time { return time.Date(2026, 9, 18, 1, 0, 30, 0, time.UTC) }, - AISnapshotRaw: snapshotRaw, - AIAuthorizationRaw: authorizationRaw, - }) - activateTestServer(t, server) - response, err := server.GetExecutionPermit(context.Background(), &agentpb.GetExecutionPermitRequest{ - Meta: testMeta("permit-revoked", "permit-revoked-key", 1), - Binding: &agentpb.ExecutionBinding{TenantId: "tenant-1", TenantKey: "tenant-demo-key", ExecutionId: "execution-revoked", AgentVersionId: snapshot.AgentVersionID}, - ResourceReservationId: "reservation-revoked", - ConfigSha256: snapshot.Digest, - }) - if err == nil || status.Code(err) != codes.PermissionDenied || response != nil { - t.Fatalf("revoked authorization response=%+v err=%v code=%s", response, err, status.Code(err)) - } -} - -func activateTestServer(t *testing.T, server *Server) { - t.Helper() - _, err := server.ActivateAgent(context.Background(), &agentpb.ActivateAgentRequest{ - Meta: testMeta("activate-ai", "", 0), - Binding: &agentpb.AgentBinding{AgentId: "agent-1", CellId: "cell-1", ExpectedBootId: "boot-1", DispatcherEpoch: "epoch-1", SessionGeneration: 1}, - ActivationOperationId: "activate-ai", - }) - if err != nil { - t.Fatal(err) - } -} diff --git a/internal/rpc/approved_execution_test.go b/internal/rpc/approved_execution_test.go index 7c05abc..fa2d2d9 100644 --- a/internal/rpc/approved_execution_test.go +++ b/internal/rpc/approved_execution_test.go @@ -108,6 +108,37 @@ func TestExecuteApprovedBindsConfigAndPreventsRedialAcrossRestart(t *testing.T) } } +func TestExecuteApprovedJournalWriteFailureCannotReplayFromMemory(t *testing.T) { + now := time.Date(2026, 9, 20, 10, 0, 0, 0, time.UTC) + req := approvedTestRequest(t, now) + state := filepath.Join(t.TempDir(), "agent-session.json") + attempts := 0 + s := activatedApprovedServer(t, now, state, req.DispatcherId, 1, func(context.Context, ApprovedExecution) error { + attempts++ + return errors.New("mock originated; outcome unknown") + }) + blocker := state + ".approved" + if err := os.WriteFile(blocker, []byte("not a directory"), 0600); err != nil { + t.Fatal(err) + } + for i := 0; i < 2; i++ { + response, err := s.ExecuteApproved(context.Background(), req) + if err == nil || response != nil || attempts != 0 { + t.Fatalf("failed journal write acknowledged or originated call: attempt=%d response=%v err=%v", i, response, err) + } + } + if err := os.Remove(blocker); err != nil { + t.Fatal(err) + } + response, err := s.ExecuteApproved(context.Background(), req) + if err == nil || response != nil || attempts != 1 { + t.Fatalf("recovered journal must originate only once: attempts=%d response=%v err=%v", attempts, response, err) + } + if _, err := s.ExecuteApproved(context.Background(), req); status.Code(err) != codes.FailedPrecondition || attempts != 1 { + t.Fatalf("unknown outcome must not redial: attempts=%d err=%v", attempts, err) + } +} + func TestExecuteApprovedRejectsMismatchedIdentityAndCredentialsBeforeAttempt(t *testing.T) { now := time.Date(2026, 9, 20, 10, 0, 0, 0, time.UTC) req := approvedTestRequest(t, now) diff --git a/internal/rpc/authorized_execution.go b/internal/rpc/authorized_execution.go deleted file mode 100644 index e349fd3..0000000 --- a/internal/rpc/authorized_execution.go +++ /dev/null @@ -1,123 +0,0 @@ -package rpc - -import ( - "context" - "encoding/hex" - "log/slog" - - agentpb "git.ipao.vip/rogee/go-sip/gen/agent" - "google.golang.org/grpc/codes" - "google.golang.org/grpc/status" - "google.golang.org/protobuf/proto" -) - -const authorizedOriginationVersion = "agent-authorized-origination.v0.1" - -// ExecuteAuthorized only runs an explicit mock adapter. D has already made and -// durably claimed the final dial decision. A enforces the D-issued exclusive -// deadline and session identity, but never recalculates business policy. -func (s *Server) ExecuteAuthorized(ctx context.Context, req *agentpb.ExecuteAuthorizedRequest) (*agentpb.ExecuteAuthorizedResponse, error) { - if req == nil || req.Meta == nil || req.Binding == nil { - return nil, status.Error(codes.InvalidArgument, "request metadata and execution binding are required") - } - if s.mode != "mock" || s.mockAuthorizedOriginate == nil { - return nil, status.Error(codes.FailedPrecondition, "authorized origination requires an explicit mock adapter") - } - if err := s.authorize(ctx, req.Meta); err != nil { - return nil, err - } - if err := requireIdempotency(req.Meta); err != nil { - return nil, err - } - binding := req.Binding - if req.SchemaVersion != authorizedOriginationVersion || binding.ExecutionId == "" || binding.TenantId == "" || binding.TenantKey == "" || - binding.TaskId == "" || binding.TaskItemId == "" || req.SelectedTrunkId == "" || req.CallerId == "" || req.Callee == "" || - req.RingTimeoutMs <= 0 || req.MaxCallDurationMs <= 0 || req.DialBeforeUnixMs <= 0 || !validLowerSHA256(req.BoundSnapshotSha256) { - return nil, status.Error(codes.InvalidArgument, "invalid authorized origination version, binding or dial decision") - } - digest := messageDigest(req) - key := s.operationKey(req.Meta) - s.mu.Lock() - if s.executionErr != nil { - s.mu.Unlock() - return nil, status.Errorf(codes.Internal, "execution journal unavailable: %v", s.executionErr) - } - if previous, exists := s.operations[key]; exists { - if previous.digest != digest { - s.mu.Unlock() - return &agentpb.ExecuteAuthorizedResponse{Receipt: s.receipt(req.Meta, agentpb.ResultCode_RESULT_CODE_CONFLICT, agentpb.FailureCode_FAILURE_CODE_ABORTED, "idempotency key content conflict", false)}, nil - } - entry := s.executions[binding.ExecutionId] - if entry == nil { - s.mu.Unlock() - return nil, status.Error(codes.Internal, "persisted operation lacks execution state") - } - response := &agentpb.ExecuteAuthorizedResponse{Receipt: proto.Clone(previous.receipt).(*agentpb.OperationReceipt), State: entry.state} - s.mu.Unlock() - return response, nil - } - if _, exists := s.executions[binding.ExecutionId]; exists { - s.mu.Unlock() - return &agentpb.ExecuteAuthorizedResponse{Receipt: s.receipt(req.Meta, agentpb.ResultCode_RESULT_CODE_CONFLICT, agentpb.FailureCode_FAILURE_CODE_ABORTED, "execution already exists", false)}, nil - } - if s.now().UnixMilli() >= req.DialBeforeUnixMs { - s.mu.Unlock() - return nil, status.Error(codes.FailedPrecondition, "Dispatcher-issued dial deadline expired") - } - // Persist an unknown one-shot attempt BEFORE contacting the mock adapter. - // After a crash or lost response, replay/query cannot cause a second dial. - s.executions[binding.ExecutionId] = &executionRecord{ - executeDigest: digest, binding: proto.Clone(binding).(*agentpb.ExecutionBinding), - state: agentpb.ExecutionState_EXECUTION_STATE_UNKNOWN, taskRevision: binding.TaskRevision, - unknown: true, callState: "mock_originating_unknown", - } - pending := s.receipt(req.Meta, agentpb.ResultCode_RESULT_CODE_UNKNOWN, agentpb.FailureCode_FAILURE_CODE_UNAVAILABLE, "mock originate pending; do not retry", false) - s.operations[key] = operationRecord{digest: digest, receipt: pending} - if err := s.persistExecutionJournalLocked(); err != nil { - s.mu.Unlock() - return nil, err - } - s.mu.Unlock() - - // A newly closed deadline between durable claim and adapter invocation must - // not dial. This is a technical check of D's instruction, not a new policy. - expired := s.now().UnixMilli() >= req.DialBeforeUnixMs - var dialErr error - if !expired { - dialErr = s.mockAuthorizedOriginate(ctx, proto.Clone(req).(*agentpb.ExecuteAuthorizedRequest)) - } - s.mu.Lock() - defer s.mu.Unlock() - record := s.executions[binding.ExecutionId] - var result agentpb.ResultCode - var failure agentpb.FailureCode - var detail string - switch { - case expired: - record.state, record.unknown, record.callState = agentpb.ExecutionState_EXECUTION_STATE_TERMINAL, false, "mock_deadline_closed_without_dial" - record.terminalObservedAtUnixMs = s.now().UnixMilli() - result, failure, detail = agentpb.ResultCode_RESULT_CODE_REJECTED, agentpb.FailureCode_FAILURE_CODE_FAILED_PRECONDITION, "dial deadline expired before mock adapter" - case dialErr != nil: - slog.Error("mock originate failed; execution remains unknown and cannot be retried", "execution_id", binding.ExecutionId, "error", dialErr) - record.state, record.unknown, record.callState = agentpb.ExecutionState_EXECUTION_STATE_UNKNOWN, true, "mock_originating_unknown" - result, failure, detail = agentpb.ResultCode_RESULT_CODE_UNKNOWN, agentpb.FailureCode_FAILURE_CODE_UNAVAILABLE, "mock originate outcome unknown; reconcile instead of retry" - default: - record.state, record.unknown, record.callState = agentpb.ExecutionState_EXECUTION_STATE_TERMINAL, false, "mock_no_answer" - record.terminalObservedAtUnixMs = s.now().UnixMilli() - result, failure, detail = agentpb.ResultCode_RESULT_CODE_APPLIED, agentpb.FailureCode_FAILURE_CODE_UNSPECIFIED, "" - } - receipt := s.receipt(req.Meta, result, failure, detail, false) - s.operations[key] = operationRecord{digest: digest, receipt: receipt} - if err := s.persistExecutionJournalLocked(); err != nil { - return nil, err - } - return &agentpb.ExecuteAuthorizedResponse{Receipt: receipt, State: record.state}, nil -} - -func validLowerSHA256(value string) bool { - if len(value) != 64 { - return false - } - decoded, err := hex.DecodeString(value) - return err == nil && hex.EncodeToString(decoded) == value -} diff --git a/internal/rpc/authorized_execution_test.go b/internal/rpc/authorized_execution_test.go deleted file mode 100644 index ba2f32b..0000000 --- a/internal/rpc/authorized_execution_test.go +++ /dev/null @@ -1,201 +0,0 @@ -package rpc - -import ( - "context" - "errors" - "path/filepath" - "strings" - "testing" - "time" - - agentpb "git.ipao.vip/rogee/go-sip/gen/agent" - "google.golang.org/grpc/codes" - "google.golang.org/grpc/status" - "google.golang.org/protobuf/proto" -) - -func authorizedMockRequest(now time.Time) *agentpb.ExecuteAuthorizedRequest { - return &agentpb.ExecuteAuthorizedRequest{ - SchemaVersion: "agent-authorized-origination.v0.1", - Meta: testMeta("authorized-1", "authorized-key-1", 1), - Binding: &agentpb.ExecutionBinding{ - TenantId: "tenant-1", TenantKey: "tenant-original", ExecutionId: "execution-authorized-1", - TaskId: "task-1", TaskItemId: "command-1", TaskRevision: 2, - AgentVersionId: "agent-version-1", RoutePolicyId: "route-1", CallerProfileId: "caller-1", - }, - SelectedTrunkId: "trunk-1", - CallerId: "BD93205882", - Callee: "15003164745", - RingTimeoutMs: 12000, - MaxCallDurationMs: 30000, - DialBeforeUnixMs: now.Add(time.Minute).UnixMilli(), - BoundSnapshotSha256: strings.Repeat("a", 64), - } -} - -func TestExecuteAuthorizedMockOneShotAndIdentity(t *testing.T) { - now := time.Date(2026, 9, 21, 1, 30, 0, 0, time.UTC) - attempts := 0 - server := NewServer(ServerOptions{Mode: "mock", Now: func() time.Time { return now }, MockAuthorizedOriginate: func(_ context.Context, req *agentpb.ExecuteAuthorizedRequest) error { - attempts++ - if req.SelectedTrunkId != "trunk-1" || req.CallerId != "BD93205882" || req.MaxCallDurationMs != 30000 { - t.Fatalf("Agent received modified Dispatcher decision: %+v", req) - } - return nil - }}) - activateTestServer(t, server) - req := authorizedMockRequest(now) - first, err := server.ExecuteAuthorized(context.Background(), req) - if err != nil || first.GetReceipt().GetResult() != agentpb.ResultCode_RESULT_CODE_APPLIED || first.GetState() != agentpb.ExecutionState_EXECUTION_STATE_TERMINAL || attempts != 1 { - t.Fatalf("first mock execution=%+v attempts=%d err=%v", first, attempts, err) - } - replay, err := server.ExecuteAuthorized(context.Background(), req) - if err != nil || replay.GetReceipt().GetMeta().GetOperationId() != first.GetReceipt().GetMeta().GetOperationId() || attempts != 1 { - t.Fatalf("replay caused duplicate mock dial: %+v attempts=%d err=%v", replay, attempts, err) - } - changed := proto.Clone(req).(*agentpb.ExecuteAuthorizedRequest) - changed.Callee = "15830461047" - conflict, err := server.ExecuteAuthorized(context.Background(), changed) - if err != nil || conflict.GetReceipt().GetResult() != agentpb.ResultCode_RESULT_CODE_CONFLICT || attempts != 1 { - t.Fatalf("same key/different body: %+v attempts=%d err=%v", conflict, attempts, err) - } - newKey := proto.Clone(req).(*agentpb.ExecuteAuthorizedRequest) - newKey.Meta = testMeta("authorized-2", "authorized-key-2", 1) - conflict, err = server.ExecuteAuthorized(context.Background(), newKey) - if err != nil || conflict.GetReceipt().GetResult() != agentpb.ResultCode_RESULT_CODE_CONFLICT || attempts != 1 { - t.Fatalf("same execution/new operation: %+v attempts=%d err=%v", conflict, attempts, err) - } -} - -func TestExecuteAuthorizedRejectsExpiredAndInvalidMock(t *testing.T) { - now := time.Date(2026, 9, 21, 1, 30, 0, 0, time.UTC) - attempts := 0 - server := NewServer(ServerOptions{Mode: "mock", Now: func() time.Time { return now }, MockAuthorizedOriginate: func(context.Context, *agentpb.ExecuteAuthorizedRequest) error { - attempts++ - return nil - }}) - activateTestServer(t, server) - for _, tc := range []struct { - name string - change func(*agentpb.ExecuteAuthorizedRequest) - }{ - {"expired at Agent", func(r *agentpb.ExecuteAuthorizedRequest) { r.DialBeforeUnixMs = now.UnixMilli() }}, - {"missing selected trunk", func(r *agentpb.ExecuteAuthorizedRequest) { r.SelectedTrunkId = "" }}, - {"invalid snapshot digest", func(r *agentpb.ExecuteAuthorizedRequest) { r.BoundSnapshotSha256 = "bad" }}, - {"unknown contract version", func(r *agentpb.ExecuteAuthorizedRequest) { r.SchemaVersion = "other" }}, - } { - t.Run(tc.name, func(t *testing.T) { - req := authorizedMockRequest(now) - tc.change(req) - response, err := server.ExecuteAuthorized(context.Background(), req) - if response != nil || status.Code(err) != codes.InvalidArgument && status.Code(err) != codes.FailedPrecondition { - t.Fatalf("invalid decision response=%+v err=%v", response, err) - } - if attempts != 0 { - t.Fatalf("invalid decision reached mock dial %d times", attempts) - } - }) - } -} - -func TestExecuteAuthorizedMockFailureIsUnknownAndNeverRetried(t *testing.T) { - now := time.Date(2026, 9, 21, 1, 30, 0, 0, time.UTC) - attempts := 0 - server := NewServer(ServerOptions{Mode: "mock", Now: func() time.Time { return now }, MockAuthorizedOriginate: func(context.Context, *agentpb.ExecuteAuthorizedRequest) error { - attempts++ - return errors.New("mock transport failed") - }}) - activateTestServer(t, server) - req := authorizedMockRequest(now) - first, err := server.ExecuteAuthorized(context.Background(), req) - if err != nil || first.GetReceipt().GetResult() != agentpb.ResultCode_RESULT_CODE_UNKNOWN || attempts != 1 { - t.Fatalf("uncertain mock execution=%+v attempts=%d err=%v", first, attempts, err) - } - replay, err := server.ExecuteAuthorized(context.Background(), req) - if err != nil || replay.GetReceipt().GetResult() != agentpb.ResultCode_RESULT_CODE_UNKNOWN || attempts != 1 { - t.Fatalf("uncertain execution redialed: %+v attempts=%d err=%v", replay, attempts, err) - } -} - -func TestExecuteAuthorizedJournalPreventsRedialAfterAgentRestart(t *testing.T) { - for _, tc := range []struct { - name string - fail bool - }{ - {"mock completed", false}, - {"mock outcome unknown", true}, - } { - t.Run(tc.name, func(t *testing.T) { - now := time.Date(2026, 9, 21, 1, 30, 0, 0, time.UTC) - path := filepath.Join(t.TempDir(), "agent-session.json") - attempts := 0 - start := func(generation uint64, bootID string) *Server { - server := NewServer(ServerOptions{Mode: "mock", Now: func() time.Time { return now }, StatePath: path, - Status: &agentpb.AgentStatus{AgentId: "agent-1", CellId: "cell-1", BootId: bootID}, - MockAuthorizedOriginate: func(context.Context, *agentpb.ExecuteAuthorizedRequest) error { - attempts++ - if tc.fail { - return errors.New("mock adapter failed") - } - return nil - }, - }) - meta := testMeta("activate", "", 0) - meta.BootId = bootID - if _, err := server.ActivateAgent(context.Background(), &agentpb.ActivateAgentRequest{ - Meta: meta, Binding: &agentpb.AgentBinding{AgentId: "agent-1", CellId: "cell-1", ExpectedBootId: bootID, - DispatcherEpoch: "epoch-1", SessionGeneration: generation}, ActivationOperationId: "activate", - }); err != nil { - t.Fatal(err) - } - return server - } - first := start(1, "boot-1") - req := authorizedMockRequest(now) - response, err := first.ExecuteAuthorized(context.Background(), req) - if err != nil || response == nil || attempts != 1 { - t.Fatalf("first mock call: response=%+v attempts=%d err=%v", response, attempts, err) - } - terminalObservedAt := now.UnixMilli() - now = now.Add(time.Hour) - recovered := start(2, "boot-2") - replay := proto.Clone(req).(*agentpb.ExecuteAuthorizedRequest) - replay.Meta = testMeta("authorized-new-boot", "authorized-key-new-boot", 2) - replay.Meta.BootId = "boot-2" - retry, err := recovered.ExecuteAuthorized(context.Background(), replay) - if err != nil || retry.GetReceipt().GetResult() == agentpb.ResultCode_RESULT_CODE_APPLIED || attempts != 1 { - t.Fatalf("restart redialed execution: response=%+v attempts=%d err=%v", retry, attempts, err) - } - queryMeta := testMeta("query-recovered", "query-recovered-key", 2) - queryMeta.BootId = "boot-2" - query, err := recovered.QueryExecution(context.Background(), &agentpb.QueryExecutionRequest{Meta: queryMeta, Binding: req.Binding}) - if err != nil || query.Snapshot == nil || query.Snapshot.Unknown != tc.fail { - t.Fatalf("recovered state failed to retain uncertainty: %+v err=%v", query, err) - } - if !tc.fail && query.Snapshot.ObservedAtUnixMs != terminalObservedAt { - t.Fatalf("terminal observation drifted after restart: got=%d want=%d", query.Snapshot.ObservedAtUnixMs, terminalObservedAt) - } - }) - } -} - -func TestExecuteAuthorizedCannotDialOutsideMockOrWithoutAdapter(t *testing.T) { - now := time.Date(2026, 9, 21, 1, 30, 0, 0, time.UTC) - for _, tc := range []struct { - name string - mode string - }{ - {"mixed disabled", "mixed"}, - {"real disabled", "real"}, - {"mock adapter missing", "mock"}, - } { - t.Run(tc.name, func(t *testing.T) { - server := NewServer(ServerOptions{Mode: tc.mode, Now: func() time.Time { return now }}) - activateTestServer(t, server) - response, err := server.ExecuteAuthorized(context.Background(), authorizedMockRequest(now)) - if response != nil || status.Code(err) != codes.FailedPrecondition { - t.Fatalf("unavailable mode/adapter response=%+v err=%v", response, err) - } - }) - } -} diff --git a/internal/rpc/calllog_test.go b/internal/rpc/calllog_test.go deleted file mode 100644 index d438b17..0000000 --- a/internal/rpc/calllog_test.go +++ /dev/null @@ -1,100 +0,0 @@ -package rpc - -import ( - "context" - "encoding/json" - "os" - "path/filepath" - "strings" - "testing" - "time" - - agentpb "git.ipao.vip/rogee/go-sip/gen/agent" - "git.ipao.vip/rogee/go-sip/internal/calllog" - "git.ipao.vip/rogee/go-sip/internal/contract" - "git.ipao.vip/rogee/go-sip/internal/testfixture" -) - -func TestServerWritesPhoneCorrelatedCallBusinessLog(t *testing.T) { - now := time.Date(2026, 9, 18, 1, 0, 0, 0, time.UTC) - logger, err := calllog.New(filepath.Join(t.TempDir(), "business", "calls.jsonl"), []byte("0123456789abcdef"), func() time.Time { return now }) - if err != nil { - t.Fatal(err) - } - server := NewServer(ServerOptions{Now: func() time.Time { return now }, CallLogger: logger}) - meta := testMeta("activate-log", "", 0) - if _, err := server.ActivateAgent(context.Background(), &agentpb.ActivateAgentRequest{ - Meta: meta, - Binding: &agentpb.AgentBinding{AgentId: "agent-1", CellId: "cell-1", ExpectedBootId: "boot-1", DispatcherEpoch: "epoch-1", SessionGeneration: 1}, - ActivationOperationId: "activate-log", - }); err != nil { - t.Fatal(err) - } - raw, err := testfixture.Execute() - if err != nil { - t.Fatal(err) - } - envelope, payload, err := contract.DecodeExecute(raw) - if err != nil { - t.Fatal(err) - } - binding := &agentpb.ExecutionBinding{ - TenantId: envelope.TenantID, TenantKey: envelope.TenantKey, ExecutionId: payload.ExecutionID, - TaskId: payload.TaskID, TaskItemId: payload.TaskItemID, TaskRevision: payload.TaskRevision, - AgentVersionId: payload.AgentVersionID, RoutePolicyId: payload.RoutePolicyID, CallerProfileId: payload.CallerProfileID, - } - if _, err := server.Execute(context.Background(), &agentpb.ExecuteRequest{ - Meta: testMeta("execute-log", "execute-log-key", 1), Binding: binding, CallExecuteJson: raw, - ConfigSha256: "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa", - }); err != nil { - t.Fatal(err) - } - statusPayload := []byte(`{"call_id":"call-1","execution_id":"execution-1","call_state":"answered","attempt_id":"attempt-1","attempt_state":"active","route_policy_id":"route-sip-first","caller_profile_id":"caller-sip-first","trunk_id":"provider-primary","cell_id":"cell-1","sip_stage":"invite","sip_status_code":183}`) - if _, err := server.ReportExecutionEvent(context.Background(), &agentpb.ReportExecutionEventRequest{ - Meta: testMeta("fact-status", "fact-status-key", 1), Fact: &agentpb.ExecutionFact{ - FactId: "fact-status", ContentSha256: "digest-status", Binding: binding, - Kind: agentpb.FactKind_FACT_KIND_CALL_STATUS, PayloadJson: statusPayload, - }, - }); err != nil { - t.Fatal(err) - } - finishedPayload := []byte(`{"call_id":"call-1","execution_id":"execution-1","outcome":"no_answer","duration_ms":3000,"reason_code":"provider_480","asset_state":"failed","attempt_summary":[{"attempt_id":"attempt-1","state":"ended","trunk_id":"provider-primary","cell_id":"cell-1","reason_code":"provider_480"}],"recording_id":"recording-1","recording_state":"failed"}`) - if _, err := server.ReportExecutionEvent(context.Background(), &agentpb.ReportExecutionEventRequest{ - Meta: testMeta("fact-finished", "fact-finished-key", 1), Fact: &agentpb.ExecutionFact{ - FactId: "fact-finished", ContentSha256: "digest-finished", Binding: binding, - Kind: agentpb.FactKind_FACT_KIND_CALL_FINISHED, PayloadJson: finishedPayload, - }, - }); err != nil { - t.Fatal(err) - } - - data, err := os.ReadFile(logger.Path()) - if err != nil { - t.Fatal(err) - } - text := string(data) - if strings.Contains(text, payload.Callee) { - t.Fatalf("business log contains original phone: %s", text) - } - lines := strings.Split(strings.TrimSpace(text), "\n") - if len(lines) != 4 { - t.Fatalf("got %d business log lines, want prepared/status/finished/attempt: %s", len(lines), text) - } - var records []map[string]any - for _, line := range lines { - var record map[string]any - if err := json.Unmarshal([]byte(line), &record); err != nil { - t.Fatal(err) - } - records = append(records, record) - } - if records[0]["event_type"] != "execution.prepared" || records[1]["sip_status_code"] != float64(183) { - t.Fatalf("prepared or SIP status record missing: %#v", records) - } - if records[2]["recording_id"] != "recording-1" || records[2]["result"] != "no_answer" || records[3]["trunk_id"] != "provider-primary" { - t.Fatalf("finished or attempt record missing: %#v", records) - } - if records[0]["phone_ref"] != records[1]["phone_ref"] || records[1]["phone_ref"] != records[2]["phone_ref"] { - t.Fatalf("phone correlation changed across events: %#v", records) - } -} diff --git a/internal/rpc/control_policy_test.go b/internal/rpc/control_policy_test.go deleted file mode 100644 index 82bd12f..0000000 --- a/internal/rpc/control_policy_test.go +++ /dev/null @@ -1,32 +0,0 @@ -package rpc - -import ( - "context" - "testing" - "time" - - agentpb "git.ipao.vip/rogee/go-sip/gen/agent" - "git.ipao.vip/rogee/go-sip/internal/contract" - "git.ipao.vip/rogee/go-sip/internal/testfixture" -) - -func TestStopDrainPreservesCallAndCannotResume(t *testing.T) { - s := activatedServer(time.Date(2026, 9, 18, 1, 0, 0, 0, time.UTC), t) - raw, err := testfixture.Execute() - require.NoError(t, err) - envelope, payload, err := contract.DecodeExecute(raw) - require.NoError(t, err) - binding := &agentpb.ExecutionBinding{TenantId: envelope.TenantID, TenantKey: envelope.TenantKey, ExecutionId: payload.ExecutionID, TaskId: payload.TaskID, TaskItemId: payload.TaskItemID, TaskRevision: payload.TaskRevision, AgentVersionId: payload.AgentVersionID} - _, err = s.Execute(context.Background(), &agentpb.ExecuteRequest{Meta: testMeta("execute", "execute-key", 1), Binding: binding, CallExecuteJson: raw}) - require.NoError(t, err) - s.executions[binding.ExecutionId].state = agentpb.ExecutionState_EXECUTION_STATE_OBSERVED - stopped, err := s.ApplyTaskControl(context.Background(), &agentpb.ApplyTaskControlRequest{Meta: testMeta("stop", "stop-key", 1), Binding: binding, ExpectedTaskRevision: binding.TaskRevision, Action: agentpb.ControlAction_CONTROL_ACTION_STOP, ActiveCallPolicy: agentpb.ActiveCallPolicy_ACTIVE_CALL_POLICY_DRAIN}) - require.NoError(t, err) - require.Equal(t, agentpb.ResultCode_RESULT_CODE_APPLIED, stopped.Receipt.Result) - require.Equal(t, agentpb.ExecutionState_EXECUTION_STATE_OBSERVED, stopped.State) - binding.TaskRevision = stopped.AppliedTaskRevision - resumed, err := s.ApplyTaskControl(context.Background(), &agentpb.ApplyTaskControlRequest{Meta: testMeta("resume", "resume-key", 1), Binding: binding, ExpectedTaskRevision: binding.TaskRevision, Action: agentpb.ControlAction_CONTROL_ACTION_RESUME, ActiveCallPolicy: agentpb.ActiveCallPolicy_ACTIVE_CALL_POLICY_DRAIN}) - require.NoError(t, err) - require.Equal(t, agentpb.ResultCode_RESULT_CODE_REJECTED, resumed.Receipt.Result) - require.Equal(t, stopped.AppliedTaskRevision, resumed.AppliedTaskRevision) -} diff --git a/internal/rpc/control_policy_validation_test.go b/internal/rpc/control_policy_validation_test.go deleted file mode 100644 index 65ab331..0000000 --- a/internal/rpc/control_policy_validation_test.go +++ /dev/null @@ -1,48 +0,0 @@ -package rpc - -import ( - "context" - "testing" - "time" - - agentpb "git.ipao.vip/rogee/go-sip/gen/agent" -) - -func TestControlPolicyDoesNotInventMediaCompletion(t *testing.T) { - for _, tc := range []struct { - name, mode string - action agentpb.ControlAction - policy agentpb.ActiveCallPolicy - wantError, terminal bool - }{ - {name: "missing policy", mode: "mock", action: agentpb.ControlAction_CONTROL_ACTION_STOP, wantError: true}, - {name: "invalid policy", mode: "mock", action: agentpb.ControlAction_CONTROL_ACTION_STOP, policy: agentpb.ActiveCallPolicy(99), wantError: true}, - {name: "mixed hangup requires media adapter", mode: "mixed", action: agentpb.ControlAction_CONTROL_ACTION_STOP, policy: agentpb.ActiveCallPolicy_ACTIVE_CALL_POLICY_HANGUP, wantError: true}, - {name: "real pause hangup requires media adapter", mode: "real", action: agentpb.ControlAction_CONTROL_ACTION_PAUSE, policy: agentpb.ActiveCallPolicy_ACTIVE_CALL_POLICY_HANGUP, wantError: true}, - {name: "mock pause hangup", mode: "mock", action: agentpb.ControlAction_CONTROL_ACTION_PAUSE, policy: agentpb.ActiveCallPolicy_ACTIVE_CALL_POLICY_HANGUP, terminal: true}, - {name: "mock pause drain", mode: "mock", action: agentpb.ControlAction_CONTROL_ACTION_PAUSE, policy: agentpb.ActiveCallPolicy_ACTIVE_CALL_POLICY_DRAIN}, - } { - t.Run(tc.name, func(t *testing.T) { - s := activatedServer(time.Date(2026, 9, 18, 1, 0, 0, 0, time.UTC), t) - // Unit-only adapter mode selection: no provider or network call is made. - s.mode = tc.mode - binding := &agentpb.ExecutionBinding{ExecutionId: "execution-policy", TenantId: "tenant-1", TenantKey: "tenant-key", TaskId: "task-1", TaskItemId: "item-1", TaskRevision: 1} - s.executions[binding.ExecutionId] = &executionRecord{binding: binding, taskRevision: 1, state: agentpb.ExecutionState_EXECUTION_STATE_OBSERVED} - response, err := s.ApplyTaskControl(context.Background(), &agentpb.ApplyTaskControlRequest{Meta: testMeta("control", "control-key", 1), Binding: binding, ExpectedTaskRevision: 1, Action: tc.action, ActiveCallPolicy: tc.policy}) - if tc.wantError { - if err == nil || response != nil { - t.Fatal("unsupported policy reported applied") - } - if s.executions[binding.ExecutionId].taskRevision != 1 || s.executions[binding.ExecutionId].state != agentpb.ExecutionState_EXECUTION_STATE_OBSERVED { - t.Fatal("failed policy changed execution") - } - return - } - require.NoError(t, err) - require.Equal(t, agentpb.ResultCode_RESULT_CODE_APPLIED, response.Receipt.Result) - if (response.State == agentpb.ExecutionState_EXECUTION_STATE_TERMINAL) != tc.terminal { - t.Fatal("incorrect call termination") - } - }) - } -} diff --git a/internal/rpc/execution_journal_failure_test.go b/internal/rpc/execution_journal_failure_test.go deleted file mode 100644 index 7546ffc..0000000 --- a/internal/rpc/execution_journal_failure_test.go +++ /dev/null @@ -1,42 +0,0 @@ -package rpc - -import ( - "context" - "os" - "path/filepath" - "testing" - "time" - - agentpb "git.ipao.vip/rogee/go-sip/gen/agent" - "git.ipao.vip/rogee/go-sip/internal/contract" - "git.ipao.vip/rogee/go-sip/internal/testfixture" -) - -func TestExecutionJournalFailureCannotReplayMemoryAsSuccess(t *testing.T) { - s := activatedServer(time.Date(2026, 9, 18, 1, 0, 0, 0, time.UTC), t) - blocker := filepath.Join(t.TempDir(), "not-a-directory") - require.NoError(t, os.WriteFile(blocker, []byte("block"), 0600)) - s.executionPath = filepath.Join(blocker, "journal") - raw, err := testfixture.Execute() - require.NoError(t, err) - envelope, payload, err := contract.DecodeExecute(raw) - require.NoError(t, err) - req := &agentpb.ExecuteRequest{Meta: testMeta("execute", "execute-key", 1), Binding: &agentpb.ExecutionBinding{TenantId: envelope.TenantID, TenantKey: envelope.TenantKey, ExecutionId: payload.ExecutionID, TaskId: payload.TaskID, TaskItemId: payload.TaskItemID, TaskRevision: payload.TaskRevision, AgentVersionId: payload.AgentVersionID}, CallExecuteJson: raw} - response, err := s.Execute(context.Background(), req) - if err == nil || response != nil { - t.Fatal("failed write acknowledged") - } - response, err = s.Execute(context.Background(), req) - if err == nil || response != nil { - t.Fatal("memory replay bypassed failed journal") - } -} - -func TestMissingExecutionJournalWithExistingSessionFailsClosed(t *testing.T) { - path := filepath.Join(t.TempDir(), "session.json") - require.NoError(t, os.WriteFile(path, []byte(`{"generations":{"agent-1":1}}`), 0600)) - s := NewServer(ServerOptions{StatePath: path, Status: &agentpb.AgentStatus{AgentId: "agent-1", CellId: "cell-1", BootId: "boot-1"}}) - if s.executionJournalReady() == nil { - t.Fatal("missing execution history silently reset") - } -} diff --git a/internal/rpc/execution_journal_mode_test.go b/internal/rpc/execution_journal_mode_test.go deleted file mode 100644 index 13cd5d9..0000000 --- a/internal/rpc/execution_journal_mode_test.go +++ /dev/null @@ -1,25 +0,0 @@ -package rpc - -import ( - "path/filepath" - "testing" - - agentpb "git.ipao.vip/rogee/go-sip/gen/agent" -) - -func TestExecutionJournalCannotCrossAdapterModes(t *testing.T) { - for _, mode := range []string{"mixed", "real"} { - t.Run(mode, func(t *testing.T) { - path := filepath.Join(t.TempDir(), "session.json") - status := &agentpb.AgentStatus{AgentId: "agent-1", CellId: "cell-1", BootId: "boot-1"} - original := NewServer(ServerOptions{Mode: "mock", StatePath: path, Status: status}) - require.NoError(t, original.executionJournalReady()) - recovered := NewServer(ServerOptions{Mode: mode, StatePath: path, Status: status}) - if recovered.executionJournalReady() == nil { - t.Fatal("mock execution state reused by another adapter mode") - } - same := NewServer(ServerOptions{Mode: "mock", StatePath: path, Status: status}) - require.NoError(t, same.executionJournalReady()) - }) - } -} diff --git a/internal/rpc/execution_journal_test.go b/internal/rpc/execution_journal_test.go deleted file mode 100644 index 8467968..0000000 --- a/internal/rpc/execution_journal_test.go +++ /dev/null @@ -1,74 +0,0 @@ -package rpc - -import ( - "context" - "os" - "path/filepath" - "testing" - "time" - - agentpb "git.ipao.vip/rogee/go-sip/gen/agent" - "git.ipao.vip/rogee/go-sip/internal/contract" - "git.ipao.vip/rogee/go-sip/internal/testfixture" - "google.golang.org/protobuf/proto" -) - -func TestControlReceiptSurvivesRestartAndNewSession(t *testing.T) { - now := time.Date(2026, 9, 18, 1, 0, 0, 0, time.UTC) - path := filepath.Join(t.TempDir(), "session.json") - start := func(generation uint64) *Server { - bootID := "boot-1" - if generation > 1 { - bootID = "boot-2" - } - s := NewServer(ServerOptions{Mode: "mock", Now: func() time.Time { return now }, StatePath: path, Status: &agentpb.AgentStatus{AgentId: "agent-1", CellId: "cell-1", BootId: bootID}}) - activationMeta := testMeta("activate", "", 0) - activationMeta.BootId = bootID - _, err := s.ActivateAgent(context.Background(), &agentpb.ActivateAgentRequest{Meta: activationMeta, Binding: &agentpb.AgentBinding{AgentId: "agent-1", CellId: "cell-1", ExpectedBootId: bootID, DispatcherEpoch: "epoch-1", SessionGeneration: generation}, ActivationOperationId: "activate"}) - require.NoError(t, err) - return s - } - s := start(1) - raw, err := testfixture.Execute() - require.NoError(t, err) - envelope, payload, err := contract.DecodeExecute(raw) - require.NoError(t, err) - binding := &agentpb.ExecutionBinding{TenantId: envelope.TenantID, TenantKey: envelope.TenantKey, ExecutionId: payload.ExecutionID, TaskId: payload.TaskID, TaskItemId: payload.TaskItemID, TaskRevision: payload.TaskRevision, AgentVersionId: payload.AgentVersionID} - _, err = s.Execute(context.Background(), &agentpb.ExecuteRequest{Meta: testMeta("execute", "execute-key", 1), Binding: binding, CallExecuteJson: raw}) - require.NoError(t, err) - request := &agentpb.ApplyTaskControlRequest{Meta: testMeta("pause", "pause-key", 1), Binding: binding, Action: agentpb.ControlAction_CONTROL_ACTION_PAUSE, ActiveCallPolicy: agentpb.ActiveCallPolicy_ACTIVE_CALL_POLICY_DRAIN, ExpectedTaskRevision: binding.TaskRevision} - applied, err := s.ApplyTaskControl(context.Background(), request) - require.NoError(t, err) - require.Equal(t, agentpb.ResultCode_RESULT_CODE_APPLIED, applied.Receipt.Result) - recovered := start(2) - request.Meta.SessionGeneration = 2 - request.Meta.BootId = "boot-2" - freshMeta := func(operation, key string) *agentpb.RequestMeta { - meta := testMeta(operation, key, 2) - meta.BootId = "boot-2" - return meta - } - replayed, err := recovered.ApplyTaskControl(context.Background(), request) - require.NoError(t, err) - if !proto.Equal(applied, replayed) { - t.Fatal("lost original control receipt after restart") - } - snapshot, err := recovered.QueryExecution(context.Background(), &agentpb.QueryExecutionRequest{Meta: freshMeta("query", "query-key"), Binding: binding}) - require.NoError(t, err) - if snapshot.Snapshot == nil || !snapshot.Snapshot.Unknown || snapshot.Snapshot.Binding.TaskRevision != applied.AppliedTaskRevision { - t.Fatal("recovered active work must remain unknown with applied control revision") - } - retry, err := recovered.Execute(context.Background(), &agentpb.ExecuteRequest{Meta: freshMeta("new-execute", "new-key"), Binding: binding, CallExecuteJson: raw}) - require.NoError(t, err) - if retry.Receipt.Result == agentpb.ResultCode_RESULT_CODE_ACCEPTED { - t.Fatal("restarted execution was prepared again") - } -} - -func TestCorruptExecutionJournalPreventsActivation(t *testing.T) { - path := filepath.Join(t.TempDir(), "session.json") - require.NoError(t, os.WriteFile(path+".executions", []byte(`{"version":999}`), 0600)) - s := NewServer(ServerOptions{Mode: "mock", StatePath: path, Status: &agentpb.AgentStatus{AgentId: "agent-1", CellId: "cell-1", BootId: "boot-1"}}) - _, err := s.ActivateAgent(context.Background(), &agentpb.ActivateAgentRequest{Meta: testMeta("activate", "", 0), Binding: &agentpb.AgentBinding{AgentId: "agent-1", CellId: "cell-1", ExpectedBootId: "boot-1", DispatcherEpoch: "epoch-1", SessionGeneration: 1}, ActivationOperationId: "activate"}) - require.Error(t, err) -} diff --git a/internal/rpc/server.go b/internal/rpc/server.go index 42e89d8..9d70d48 100644 --- a/internal/rpc/server.go +++ b/internal/rpc/server.go @@ -7,7 +7,6 @@ import ( "encoding/hex" "encoding/json" "errors" - "fmt" "os" "sync" "time" @@ -16,7 +15,6 @@ import ( "git.ipao.vip/rogee/go-sip/internal/agent" "git.ipao.vip/rogee/go-sip/internal/ai" "git.ipao.vip/rogee/go-sip/internal/calllog" - "git.ipao.vip/rogee/go-sip/internal/callwindow" "git.ipao.vip/rogee/go-sip/internal/contract" "google.golang.org/grpc" "google.golang.org/grpc/codes" @@ -494,593 +492,6 @@ func (s *Server) ActivateAgent(ctx context.Context, req *agentpb.ActivateAgentRe return &agentpb.ActivateAgentResponse{Meta: s.responseMeta(req.Meta), State: state, Session: session}, nil } -func (s *Server) GetBootstrap(ctx context.Context, req *agentpb.GetBootstrapRequest) (*agentpb.GetBootstrapResponse, error) { - if req == nil || req.Meta == nil { - return nil, status.Error(codes.InvalidArgument, "request metadata is required") - } - if err := s.checkPeer(ctx, req.Meta.AgentId); err != nil { - return nil, err - } - if err := s.sessions.Authorize(req.Meta, s.now()); err != nil { - return nil, err - } - return &agentpb.GetBootstrapResponse{ - Meta: s.responseMeta(req.Meta), - State: agentpb.ActivationState_ACTIVATION_STATE_ACTIVE, - RuntimeConfigs: cloneConfigReferences(s.configReferences), - UploadPolicy: proto.Clone(s.uploadPolicy).(*agentpb.UploadPolicy), - }, nil -} - -func (s *Server) SetAdmissionState(ctx context.Context, req *agentpb.SetAdmissionStateRequest) (*agentpb.SetAdmissionStateResponse, error) { - if req == nil || req.Meta == nil || req.Target == nil { - return nil, status.Error(codes.InvalidArgument, "request metadata and target are required") - } - if err := s.authorize(ctx, req.Meta); err != nil { - return nil, err - } - if err := requireIdempotency(req.Meta); err != nil { - return nil, err - } - if req.State == agentpb.AdmissionState_ADMISSION_STATE_UNSPECIFIED { - return nil, status.Error(codes.InvalidArgument, "admission state is required") - } - key := req.Target.AgentId - s.mu.Lock() - defer s.mu.Unlock() - current := s.admissions[key] - if current.generation != req.ExpectedAdmissionGeneration { - return &agentpb.SetAdmissionStateResponse{Receipt: s.receipt(req.Meta, agentpb.ResultCode_RESULT_CODE_CONFLICT, agentpb.FailureCode_FAILURE_CODE_ABORTED, "admission generation conflict", false)}, nil - } - current.generation++ - current.state = req.State - s.admissions[key] = current - return &agentpb.SetAdmissionStateResponse{Receipt: s.receipt(req.Meta, agentpb.ResultCode_RESULT_CODE_APPLIED, agentpb.FailureCode_FAILURE_CODE_UNSPECIFIED, "", false), AppliedAdmissionGeneration: current.generation}, nil -} - -func (s *Server) Execute(ctx context.Context, req *agentpb.ExecuteRequest) (*agentpb.ExecuteResponse, error) { - if req == nil || req.Meta == nil || req.Binding == nil { - return nil, status.Error(codes.InvalidArgument, "request metadata and execution binding are required") - } - if err := s.authorize(ctx, req.Meta); err != nil { - return nil, err - } - if err := requireIdempotency(req.Meta); err != nil { - return nil, err - } - if req.Binding.ExecutionId == "" || req.Binding.TaskId == "" || req.Binding.TaskItemId == "" || req.Binding.TenantKey == "" || len(req.CallExecuteJson) == 0 { - return nil, status.Error(codes.InvalidArgument, "execution binding and command bytes are required") - } - envelope, payload, err := contract.DecodeExecute(req.CallExecuteJson) - if err != nil { - return nil, status.Errorf(codes.InvalidArgument, "call.execute contract: %v", err) - } - if envelope.TenantKey != req.Binding.TenantKey || payload.ExecutionID != req.Binding.ExecutionId || payload.TaskID != req.Binding.TaskId || payload.TaskItemID != req.Binding.TaskItemId || payload.AgentVersionID != req.Binding.AgentVersionId { - return nil, status.Error(codes.Aborted, "execution binding does not match command") - } - if err := s.validateAIExecution(req.Binding, req.ConfigSha256); err != nil { - return nil, err - } - phoneIdentity := calllog.Identity{} - if s.callLogger != nil { - phoneIdentity, err = s.callLogger.Identity(payload.Callee) - if err != nil { - return nil, status.Errorf(codes.InvalidArgument, "callee cannot be logged: %v", err) - } - } - digest := messageDigest(req) - if receipt, conflict := s.replayOperation(req.Meta, digest); receipt != nil || conflict != nil { - if conflict != nil { - return &agentpb.ExecuteResponse{Receipt: conflict}, nil - } - return &agentpb.ExecuteResponse{Receipt: receipt, State: agentpb.ExecutionState_EXECUTION_STATE_PREPARED}, nil - } - if s.mode != "mock" { - if err := callwindow.Check(s.now()); err != nil { - return nil, status.Error(codes.FailedPrecondition, err.Error()) - } - } - s.mu.Lock() - defer s.mu.Unlock() - if s.executionErr != nil { - return nil, status.Errorf(codes.Internal, "execution journal unavailable: %v", s.executionErr) - } - // Recheck after acquiring the execution lock: concurrent deliveries may - // have passed the earlier read-only replay check together. - if previous, ok := s.operations[s.operationKey(req.Meta)]; ok { - if previous.digest != digest { - return &agentpb.ExecuteResponse{Receipt: s.receipt(req.Meta, agentpb.ResultCode_RESULT_CODE_CONFLICT, agentpb.FailureCode_FAILURE_CODE_ABORTED, "idempotency key content conflict", false)}, nil - } - return &agentpb.ExecuteResponse{Receipt: proto.Clone(previous.receipt).(*agentpb.OperationReceipt), State: agentpb.ExecutionState_EXECUTION_STATE_PREPARED}, nil - } - var priorPermit *agentpb.ExecutionPermit - if previous := s.executions[req.Binding.ExecutionId]; previous != nil { - if previous.unknown { - return &agentpb.ExecuteResponse{Receipt: s.receipt(req.Meta, agentpb.ResultCode_RESULT_CODE_UNKNOWN, agentpb.FailureCode_FAILURE_CODE_FAILED_PRECONDITION, "execution requires reconciliation", false), State: agentpb.ExecutionState_EXECUTION_STATE_UNKNOWN}, nil - } - if previous.controlAction == agentpb.ControlAction_CONTROL_ACTION_PAUSE || previous.controlAction == agentpb.ControlAction_CONTROL_ACTION_STOP || previous.state == agentpb.ExecutionState_EXECUTION_STATE_TERMINAL { - return &agentpb.ExecuteResponse{Receipt: s.receipt(req.Meta, agentpb.ResultCode_RESULT_CODE_REJECTED, agentpb.FailureCode_FAILURE_CODE_FAILED_PRECONDITION, "execution control blocks preparation", false)}, nil - } - if previous.executeDigest != "" { - return &agentpb.ExecuteResponse{Receipt: s.receipt(req.Meta, agentpb.ResultCode_RESULT_CODE_CONFLICT, agentpb.FailureCode_FAILURE_CODE_ABORTED, "execution already prepared under another operation", false)}, nil - } - priorPermit = previous.permit - } - if s.callLogger != nil { - if err := s.callLogger.Append(calllog.Event{ - EventID: "execution:" + payload.ExecutionID + ":prepared", EventType: "execution.prepared", Phone: payload.Callee, - TenantID: envelope.TenantID, TraceID: envelope.TraceID, ExecutionID: payload.ExecutionID, - TaskID: payload.TaskID, TaskItemID: payload.TaskItemID, TaskRevision: payload.TaskRevision, - AgentID: req.Meta.AgentId, CellID: req.Meta.CellId, RoutePolicyID: payload.RoutePolicyID, - CallerProfileID: payload.CallerProfileID, Status: "accepted", Result: "prepared", ReasonCode: "accepted", - }); err != nil { - return nil, status.Errorf(codes.Internal, "write call business log: %v", err) - } - } - state := agentpb.ExecutionState_EXECUTION_STATE_PREPARED - if req.PermitId != "" { - state = agentpb.ExecutionState_EXECUTION_STATE_PERMIT_GRANTED - } - s.executions[req.Binding.ExecutionId] = &executionRecord{executeDigest: digest, binding: proto.Clone(req.Binding).(*agentpb.ExecutionBinding), state: state, taskRevision: req.Binding.TaskRevision, callState: "prepared", phone: phoneIdentity, permit: priorPermit} - receipt := s.receipt(req.Meta, agentpb.ResultCode_RESULT_CODE_ACCEPTED, agentpb.FailureCode_FAILURE_CODE_UNSPECIFIED, "", false) - s.operations[s.operationKey(req.Meta)] = operationRecord{digest: digest, receipt: proto.Clone(receipt).(*agentpb.OperationReceipt)} - if err := s.persistExecutionJournalLocked(); err != nil { - return nil, err - } - return &agentpb.ExecuteResponse{Receipt: receipt, State: state}, nil -} - -func (s *Server) GetExecutionPermit(ctx context.Context, req *agentpb.GetExecutionPermitRequest) (*agentpb.GetExecutionPermitResponse, error) { - if req == nil || req.Meta == nil || req.Binding == nil { - return nil, status.Error(codes.InvalidArgument, "request metadata and execution binding are required") - } - if err := s.authorize(ctx, req.Meta); err != nil { - return nil, err - } - if err := requireIdempotency(req.Meta); err != nil { - return nil, err - } - if req.Binding.ExecutionId == "" || req.ResourceReservationId == "" { - return nil, status.Error(codes.InvalidArgument, "execution and reservation are required") - } - if err := s.validateAIExecution(req.Binding, req.ConfigSha256); err != nil { - return nil, err - } - if s.mode != "mock" { - if err := callwindow.Check(s.now()); err != nil { - return nil, status.Error(codes.FailedPrecondition, err.Error()) - } - } - digest := messageDigest(req) - s.mu.Lock() - defer s.mu.Unlock() - if s.executionErr != nil { - return nil, status.Errorf(codes.Internal, "execution journal unavailable: %v", s.executionErr) - } - execution := s.executions[req.Binding.ExecutionId] - // Control and permission decisions share one lock: neither a fresh grant - // nor a replayed grant may cross an already-applied pause/stop barrier. - if execution != nil && (execution.unknown || execution.controlAction == agentpb.ControlAction_CONTROL_ACTION_PAUSE || execution.controlAction == agentpb.ControlAction_CONTROL_ACTION_STOP || execution.state == agentpb.ExecutionState_EXECUTION_STATE_TERMINAL) { - return &agentpb.GetExecutionPermitResponse{Receipt: s.receipt(req.Meta, agentpb.ResultCode_RESULT_CODE_REJECTED, agentpb.FailureCode_FAILURE_CODE_FAILED_PRECONDITION, "execution control blocks permission", false)}, nil - } - if previous, ok := s.operations[s.operationKey(req.Meta)]; ok { - if previous.digest != digest { - return &agentpb.GetExecutionPermitResponse{Receipt: s.receipt(req.Meta, agentpb.ResultCode_RESULT_CODE_CONFLICT, agentpb.FailureCode_FAILURE_CODE_ABORTED, "idempotency key content conflict", false)}, nil - } - var permit *agentpb.ExecutionPermit - if execution != nil && execution.permit != nil { - permit = proto.Clone(execution.permit).(*agentpb.ExecutionPermit) - } - return &agentpb.GetExecutionPermitResponse{Receipt: proto.Clone(previous.receipt).(*agentpb.OperationReceipt), Permit: permit}, nil - } - permitID := fmt.Sprintf("permit-%s", req.Binding.ExecutionId) - fencingToken := randomToken() - permit := &agentpb.ExecutionPermit{PermitId: permitID, ResourceReservationId: req.ResourceReservationId, IssuedAtUnixMs: s.now().UnixMilli(), ExpiresAtUnixMs: s.now().Add(time.Second).UnixMilli(), DispatcherEpoch: req.Meta.DispatcherEpoch, SessionGeneration: req.Meta.SessionGeneration, FencingToken: fencingToken, ConfigSha256: req.ConfigSha256} - if execution == nil { - execution = &executionRecord{binding: proto.Clone(req.Binding).(*agentpb.ExecutionBinding), taskRevision: req.Binding.TaskRevision, callState: "prepared"} - s.executions[req.Binding.ExecutionId] = execution - } - if execution.permit != nil && execution.permit.ResourceReservationId != req.ResourceReservationId { - return &agentpb.GetExecutionPermitResponse{Receipt: s.receipt(req.Meta, agentpb.ResultCode_RESULT_CODE_CONFLICT, agentpb.FailureCode_FAILURE_CODE_ABORTED, "execution already has a different permit", false)}, nil - } - execution.permit = proto.Clone(permit).(*agentpb.ExecutionPermit) - execution.state = agentpb.ExecutionState_EXECUTION_STATE_PERMIT_GRANTED - receipt := s.receipt(req.Meta, agentpb.ResultCode_RESULT_CODE_APPLIED, agentpb.FailureCode_FAILURE_CODE_UNSPECIFIED, "", false) - s.operations[s.operationKey(req.Meta)] = operationRecord{digest: digest, receipt: proto.Clone(receipt).(*agentpb.OperationReceipt)} - if err := s.persistExecutionJournalLocked(); err != nil { - return nil, err - } - return &agentpb.GetExecutionPermitResponse{Receipt: receipt, Permit: permit}, nil -} - -func (s *Server) ApplyTaskControl(ctx context.Context, req *agentpb.ApplyTaskControlRequest) (*agentpb.ApplyTaskControlResponse, error) { - if req == nil || req.Meta == nil || req.Binding == nil { - return nil, status.Error(codes.InvalidArgument, "request metadata and execution binding are required") - } - if err := s.authorize(ctx, req.Meta); err != nil { - return nil, err - } - if err := requireIdempotency(req.Meta); err != nil { - return nil, err - } - if req.Action != agentpb.ControlAction_CONTROL_ACTION_PAUSE && req.Action != agentpb.ControlAction_CONTROL_ACTION_RESUME && req.Action != agentpb.ControlAction_CONTROL_ACTION_STOP { - return nil, status.Error(codes.InvalidArgument, "a supported control action is required") - } - if req.ActiveCallPolicy != agentpb.ActiveCallPolicy_ACTIVE_CALL_POLICY_DRAIN && req.ActiveCallPolicy != agentpb.ActiveCallPolicy_ACTIVE_CALL_POLICY_HANGUP { - return nil, status.Error(codes.InvalidArgument, "an explicit drain or hangup policy is required") - } - // Authorization above checks the live session. Recovery may use a new - // session or trace, but must retain the original operation and business body. - identity := proto.Clone(req).(*agentpb.ApplyTaskControlRequest) - identity.Meta = &agentpb.RequestMeta{ - ProtocolVersion: req.Meta.ProtocolVersion, - AgentId: req.Meta.AgentId, CellId: req.Meta.CellId, - OperationId: req.Meta.OperationId, IdempotencyKey: req.Meta.IdempotencyKey, - } - encoded, err := (proto.MarshalOptions{Deterministic: true}).Marshal(identity) - if err != nil { - return nil, status.Errorf(codes.InvalidArgument, "encode control request: %v", err) - } - hash := sha256.Sum256(encoded) - digest := hex.EncodeToString(hash[:]) - key := s.operationKey(req.Meta) - s.mu.Lock() - defer s.mu.Unlock() - if s.executionErr != nil { - return nil, status.Errorf(codes.Internal, "execution journal unavailable: %v", s.executionErr) - } - if previous, ok := s.operations[key]; ok { - if previous.digest != digest || previous.control == nil { - return &agentpb.ApplyTaskControlResponse{Receipt: s.receipt(req.Meta, agentpb.ResultCode_RESULT_CODE_CONFLICT, agentpb.FailureCode_FAILURE_CODE_ABORTED, "idempotency key content conflict", false)}, nil - } - return proto.Clone(previous.control).(*agentpb.ApplyTaskControlResponse), nil - } - save := func(response *agentpb.ApplyTaskControlResponse) (*agentpb.ApplyTaskControlResponse, error) { - s.operations[key] = operationRecord{digest: digest, receipt: proto.Clone(response.Receipt).(*agentpb.OperationReceipt), control: proto.Clone(response).(*agentpb.ApplyTaskControlResponse)} - if err := s.persistExecutionJournalLocked(); err != nil { - return nil, err - } - return response, nil - } - execution := s.executions[req.Binding.ExecutionId] - if execution == nil { - return save(&agentpb.ApplyTaskControlResponse{Receipt: s.receipt(req.Meta, agentpb.ResultCode_RESULT_CODE_REJECTED, agentpb.FailureCode_FAILURE_CODE_NOT_FOUND, "execution not found", false)}) - } - if !proto.Equal(execution.binding, req.Binding) { - return save(&agentpb.ApplyTaskControlResponse{Receipt: s.receipt(req.Meta, agentpb.ResultCode_RESULT_CODE_CONFLICT, agentpb.FailureCode_FAILURE_CODE_ABORTED, "execution control binding mismatch", false)}) - } - if execution.taskRevision != req.ExpectedTaskRevision { - return save(&agentpb.ApplyTaskControlResponse{Receipt: s.receipt(req.Meta, agentpb.ResultCode_RESULT_CODE_CONFLICT, agentpb.FailureCode_FAILURE_CODE_ABORTED, "task revision conflict", false), AppliedTaskRevision: execution.taskRevision, State: execution.state}) - } - if (execution.state == agentpb.ExecutionState_EXECUTION_STATE_TERMINAL || execution.controlAction == agentpb.ControlAction_CONTROL_ACTION_STOP) && req.Action != agentpb.ControlAction_CONTROL_ACTION_STOP { - return save(&agentpb.ApplyTaskControlResponse{Receipt: s.receipt(req.Meta, agentpb.ResultCode_RESULT_CODE_REJECTED, agentpb.FailureCode_FAILURE_CODE_FAILED_PRECONDITION, "stopped execution cannot resume", false), AppliedTaskRevision: execution.taskRevision, State: execution.state}) - } - if req.Action != agentpb.ControlAction_CONTROL_ACTION_RESUME && req.ActiveCallPolicy == agentpb.ActiveCallPolicy_ACTIVE_CALL_POLICY_HANGUP && s.mode != "mock" { - return nil, status.Error(codes.FailedPrecondition, "media hangup adapter is not configured; control was not applied") - } - execution.taskRevision++ - execution.binding.TaskRevision = execution.taskRevision - execution.controlAction = req.Action - execution.permit = nil - if req.Action == agentpb.ControlAction_CONTROL_ACTION_STOP { - if req.ActiveCallPolicy == agentpb.ActiveCallPolicy_ACTIVE_CALL_POLICY_DRAIN { - execution.callState = "draining" - } else { - execution.state = agentpb.ExecutionState_EXECUTION_STATE_TERMINAL - execution.callState = "stopped" - } - } else if req.Action == agentpb.ControlAction_CONTROL_ACTION_PAUSE { - execution.callState = "paused" - if req.ActiveCallPolicy == agentpb.ActiveCallPolicy_ACTIVE_CALL_POLICY_HANGUP { - execution.state = agentpb.ExecutionState_EXECUTION_STATE_TERMINAL - } - } else { - execution.callState = "resumed" - } - return save(&agentpb.ApplyTaskControlResponse{Receipt: s.receipt(req.Meta, agentpb.ResultCode_RESULT_CODE_APPLIED, agentpb.FailureCode_FAILURE_CODE_UNSPECIFIED, "", false), AppliedTaskRevision: execution.taskRevision, State: execution.state}) -} - -func (s *Server) QueryExecution(ctx context.Context, req *agentpb.QueryExecutionRequest) (*agentpb.QueryExecutionResponse, error) { - if req == nil || req.Meta == nil || req.Binding == nil { - return nil, status.Error(codes.InvalidArgument, "request metadata and execution binding are required") - } - if err := s.authorize(ctx, req.Meta); err != nil { - return nil, err - } - s.mu.Lock() - execution := s.executions[req.Binding.ExecutionId] - if execution == nil { - s.mu.Unlock() - return &agentpb.QueryExecutionResponse{Meta: s.responseMeta(req.Meta), Failure: s.failure(agentpb.FailureCode_FAILURE_CODE_NOT_FOUND, "execution not found", false)}, nil - } - observedAt := s.now().UnixMilli() - if execution.callState == "mock_no_answer" || execution.callState == "mock_deadline_closed_without_dial" { - if execution.terminalObservedAtUnixMs <= 0 { - s.mu.Unlock() - return nil, status.Error(codes.Internal, "missing durable Mock terminal observation time") - } - observedAt = execution.terminalObservedAtUnixMs - } - snapshot := &agentpb.ExecutionSnapshot{Binding: proto.Clone(execution.binding).(*agentpb.ExecutionBinding), State: execution.state, CallState: execution.callState, AttemptId: execution.binding.AttemptId, ObservedAtUnixMs: observedAt, Unknown: execution.unknown} - s.mu.Unlock() - return &agentpb.QueryExecutionResponse{Meta: s.responseMeta(req.Meta), Snapshot: snapshot}, nil -} - -func (s *Server) ReportExecutionEvent(ctx context.Context, req *agentpb.ReportExecutionEventRequest) (*agentpb.ReportExecutionEventResponse, error) { - if req == nil || req.Meta == nil || req.Fact == nil { - return nil, status.Error(codes.InvalidArgument, "request metadata and fact are required") - } - if err := s.authorize(ctx, req.Meta); err != nil { - return nil, err - } - if err := requireIdempotency(req.Meta); err != nil { - return nil, err - } - if req.Fact.FactId == "" || req.Fact.ContentSha256 == "" { - return nil, status.Error(codes.InvalidArgument, "fact ID and content digest are required") - } - s.mu.Lock() - if previous, ok := s.facts[req.Fact.FactId]; ok { - s.mu.Unlock() - if previous != req.Fact.ContentSha256 { - return &agentpb.ReportExecutionEventResponse{Receipt: s.receipt(req.Meta, agentpb.ResultCode_RESULT_CODE_CONFLICT, agentpb.FailureCode_FAILURE_CODE_ABORTED, "fact digest conflict", false)}, nil - } - return &agentpb.ReportExecutionEventResponse{Receipt: s.receipt(req.Meta, agentpb.ResultCode_RESULT_CODE_ACCEPTED, agentpb.FailureCode_FAILURE_CODE_UNSPECIFIED, "duplicate fact", false)}, nil - } - execution := executionForFact(s.executions, req.Fact) - if s.callLogger != nil && execution != nil && execution.phone.Ref != "" { - for _, event := range callLogEvents(req.Fact, execution, req.Meta, s.now()) { - if err := s.callLogger.Append(event); err != nil { - s.mu.Unlock() - return nil, status.Errorf(codes.Internal, "write call business log: %v", err) - } - } - } - s.facts[req.Fact.FactId] = req.Fact.ContentSha256 - s.mu.Unlock() - return &agentpb.ReportExecutionEventResponse{Receipt: s.receipt(req.Meta, agentpb.ResultCode_RESULT_CODE_ACCEPTED, agentpb.FailureCode_FAILURE_CODE_UNSPECIFIED, "", false)}, nil -} - -type callFactPayload struct { - CallID string `json:"call_id"` - ExecutionID string `json:"execution_id"` - TaskID string `json:"task_id"` - TaskItemID string `json:"task_item_id"` - CallState string `json:"call_state"` - AttemptID string `json:"attempt_id"` - AttemptState string `json:"attempt_state"` - RoutePolicyID string `json:"route_policy_id"` - CallerProfileID string `json:"caller_profile_id"` - TrunkID string `json:"trunk_id"` - CellID string `json:"cell_id"` - ReasonCode string `json:"reason_code"` - Outcome string `json:"outcome"` - Result string `json:"result"` - Status string `json:"status"` - DurationMS int64 `json:"duration_ms"` - AssetState string `json:"asset_state"` - RecordingID string `json:"recording_id"` - RecordingState string `json:"recording_state"` - SizeBytes int64 `json:"size_bytes"` - ChecksumSHA256 string `json:"checksum_sha256"` - SHA256 string `json:"sha256"` - SIPStage string `json:"sip_stage"` - SIPStatusCode int `json:"sip_status_code"` - StatusCode int `json:"status_code"` - SIPReason string `json:"sip_reason"` - Stage string `json:"stage"` - AttemptSummary []struct { - AttemptID string `json:"attempt_id"` - State string `json:"state"` - TrunkID string `json:"trunk_id"` - CellID string `json:"cell_id"` - ReasonCode string `json:"reason_code"` - } `json:"attempt_summary"` -} - -func executionForFact(executions map[string]*executionRecord, fact *agentpb.ExecutionFact) *executionRecord { - if fact == nil || fact.Binding == nil || fact.Binding.ExecutionId == "" { - return nil - } - return executions[fact.Binding.ExecutionId] -} - -func callLogEvents(fact *agentpb.ExecutionFact, execution *executionRecord, meta *agentpb.RequestMeta, now time.Time) []calllog.Event { - if fact == nil || execution == nil { - return nil - } - payload := callFactPayload{} - _ = json.Unmarshal(fact.PayloadJson, &payload) - binding := execution.binding - if fact.Binding != nil { - binding = fact.Binding - } - event := calllog.Event{ - EventID: fact.FactId, EventType: factKindEventType(fact.Kind), OccurredAt: now, - PhoneRef: execution.phone.Ref, PhoneMask: execution.phone.Mask, - AttemptID: payload.AttemptID, CallID: payload.CallID, Status: payload.Status, - Result: payload.Result, ReasonCode: payload.ReasonCode, DurationMS: payload.DurationMS, - SIPStage: payload.SIPStage, SIPStatusCode: payload.SIPStatusCode, SIPReason: payload.SIPReason, - RecordingID: payload.RecordingID, RecordingState: payload.RecordingState, - RecordingSize: payload.SizeBytes, RecordingSHA256: firstNonEmpty(payload.ChecksumSHA256, payload.SHA256), - } - if meta != nil { - event.TraceID = meta.TraceId - event.AgentID = meta.AgentId - event.CellID = meta.CellId - } - if event.SIPStatusCode == 0 { - event.SIPStatusCode = payload.StatusCode - } - if event.RecordingState == "" { - event.RecordingState = payload.AssetState - } - if event.RecordingState == "" && payload.Stage != "" { - event.RecordingState = payload.Stage - } - if event.Result == "" { - event.Result = payload.Outcome - } - if event.Status == "" { - event.Status = payload.CallState - } - if event.Status == "" { - event.Status = payload.Outcome - } - if binding != nil { - event.TenantID = binding.TenantId - event.ExecutionID = binding.ExecutionId - event.TaskID = binding.TaskId - event.TaskItemID = binding.TaskItemId - event.TaskRevision = binding.TaskRevision - event.CallID = firstNonEmpty(event.CallID, binding.CallId) - event.AttemptID = firstNonEmpty(event.AttemptID, binding.AttemptId) - event.RoutePolicyID = firstNonEmpty(payload.RoutePolicyID, binding.RoutePolicyId) - event.CallerProfileID = firstNonEmpty(payload.CallerProfileID, binding.CallerProfileId) - } - if event.CallID == "" { - event.CallID = payload.CallID - } - if event.ExecutionID == "" { - event.ExecutionID = payload.ExecutionID - } - if event.TaskID == "" { - event.TaskID = payload.TaskID - } - if event.TaskItemID == "" { - event.TaskItemID = payload.TaskItemID - } - event.TrunkID = payload.TrunkID - if payload.CellID != "" { - event.CellID = payload.CellID - } - result := []calllog.Event{event} - for index, attempt := range payload.AttemptSummary { - if attempt.AttemptID == "" { - continue - } - result = append(result, calllog.Event{ - EventID: fact.FactId + ":attempt:" + attempt.AttemptID, EventType: "call.attempt", OccurredAt: now, - PhoneRef: execution.phone.Ref, PhoneMask: execution.phone.Mask, TenantID: event.TenantID, - ExecutionID: event.ExecutionID, TaskID: event.TaskID, TaskItemID: event.TaskItemID, - TaskRevision: event.TaskRevision, AttemptID: attempt.AttemptID, CallID: event.CallID, - RoutePolicyID: event.RoutePolicyID, CallerProfileID: event.CallerProfileID, - TrunkID: attempt.TrunkID, CellID: attempt.CellID, AttemptState: attempt.State, - Status: attempt.State, ReasonCode: attempt.ReasonCode, Result: fmt.Sprintf("attempt_%d", index+1), - }) - } - return result -} - -func factKindEventType(kind agentpb.FactKind) string { - switch kind { - case agentpb.FactKind_FACT_KIND_EXECUTION_ACCEPTED: - return "execution.accepted" - case agentpb.FactKind_FACT_KIND_CALL_STATUS: - return "call.status" - case agentpb.FactKind_FACT_KIND_CALL_FINISHED: - return "call.finished" - case agentpb.FactKind_FACT_KIND_TRANSCRIPT_UPDATED: - return "transcript.updated" - case agentpb.FactKind_FACT_KIND_TRANSCRIPT_FAILED: - return "transcript.failed" - case agentpb.FactKind_FACT_KIND_CONTACT_OPT_OUT: - return "contact.opt_out" - case agentpb.FactKind_FACT_KIND_RECORDING_PROGRESS: - return "recording.progress" - default: - return "execution.fact" - } -} - -func firstNonEmpty(values ...string) string { - for _, value := range values { - if value != "" { - return value - } - } - return "" -} - -func (s *Server) RequestUpload(ctx context.Context, req *agentpb.RequestUploadRequest) (*agentpb.RequestUploadResponse, error) { - if s.mode != "mock" { - return nil, status.Error(codes.Unimplemented, "mock upload path is disabled outside mock mode") - } - if req == nil || req.Meta == nil || req.Binding == nil || req.Asset == nil { - return nil, status.Error(codes.InvalidArgument, "request metadata, binding and asset are required") - } - if err := s.authorize(ctx, req.Meta); err != nil { - return nil, err - } - if err := requireIdempotency(req.Meta); err != nil { - return nil, err - } - if req.UploadId == "" { - return nil, status.Error(codes.InvalidArgument, "upload ID is required") - } - if !s.uploadPolicy.Enabled { - return &agentpb.RequestUploadResponse{Receipt: s.receipt(req.Meta, agentpb.ResultCode_RESULT_CODE_REJECTED, agentpb.FailureCode_FAILURE_CODE_FAILED_PRECONDITION, "uploads are disabled", false), State: agentpb.UploadState_UPLOAD_STATE_FAILED}, nil - } - s.mu.Lock() - defer s.mu.Unlock() - if existing, ok := s.uploads[req.UploadId]; ok { - if !proto.Equal(existing.binding, req.Binding) || !proto.Equal(existing.asset, req.Asset) { - return &agentpb.RequestUploadResponse{Receipt: s.receipt(req.Meta, agentpb.ResultCode_RESULT_CODE_CONFLICT, agentpb.FailureCode_FAILURE_CODE_ABORTED, "upload ID is bound to a different execution or asset", false), State: agentpb.UploadState_UPLOAD_STATE_FAILED}, nil - } - return &agentpb.RequestUploadResponse{Receipt: s.receipt(req.Meta, agentpb.ResultCode_RESULT_CODE_ACCEPTED, agentpb.FailureCode_FAILURE_CODE_UNSPECIFIED, "duplicate upload request", false), Grant: proto.Clone(existing.grant).(*agentpb.UploadGrant), State: existing.state}, nil - } - grant := &agentpb.UploadGrant{UploadId: req.UploadId, TargetUrl: "https://oss.mock.invalid/upload/" + req.UploadId, ExpiresAtUnixMs: s.now().Add(5 * time.Minute).UnixMilli(), ObjectKey: req.Asset.AssetId, Bucket: "mock-bucket", RequiredChecksumSha256: req.Asset.ChecksumSha256, MaxBytes: s.uploadPolicy.MaxAssetBytes} - if grant.MaxBytes == 0 { - grant.MaxBytes = req.Asset.SizeBytes - } - s.uploads[req.UploadId] = uploadRecord{binding: proto.Clone(req.Binding).(*agentpb.ExecutionBinding), asset: proto.Clone(req.Asset).(*agentpb.AssetDescriptor), state: agentpb.UploadState_UPLOAD_STATE_REQUESTED, grant: grant} - return &agentpb.RequestUploadResponse{Receipt: s.receipt(req.Meta, agentpb.ResultCode_RESULT_CODE_ACCEPTED, agentpb.FailureCode_FAILURE_CODE_UNSPECIFIED, "", false), Grant: proto.Clone(grant).(*agentpb.UploadGrant), State: agentpb.UploadState_UPLOAD_STATE_REQUESTED}, nil -} - -func (s *Server) CompleteUpload(ctx context.Context, req *agentpb.CompleteUploadRequest) (*agentpb.CompleteUploadResponse, error) { - if s.mode != "mock" { - return nil, status.Error(codes.Unimplemented, "mock upload path is disabled outside mock mode") - } - if req == nil || req.Meta == nil || req.Binding == nil || req.Asset == nil { - return nil, status.Error(codes.InvalidArgument, "request metadata, binding and asset are required") - } - if err := s.authorize(ctx, req.Meta); err != nil { - return nil, err - } - if err := requireIdempotency(req.Meta); err != nil { - return nil, err - } - s.mu.Lock() - defer s.mu.Unlock() - upload, ok := s.uploads[req.UploadId] - if !ok { - return &agentpb.CompleteUploadResponse{Receipt: s.receipt(req.Meta, agentpb.ResultCode_RESULT_CODE_REJECTED, agentpb.FailureCode_FAILURE_CODE_NOT_FOUND, "upload not found", false), State: agentpb.UploadState_UPLOAD_STATE_FAILED}, nil - } - if !proto.Equal(upload.binding, req.Binding) || !proto.Equal(upload.asset, req.Asset) { - return &agentpb.CompleteUploadResponse{Receipt: s.receipt(req.Meta, agentpb.ResultCode_RESULT_CODE_CONFLICT, agentpb.FailureCode_FAILURE_CODE_ABORTED, "upload completion binding does not match request", false), State: agentpb.UploadState_UPLOAD_STATE_FAILED}, nil - } - if upload.asset.ChecksumSha256 != req.UploadedChecksumSha256 || upload.asset.SizeBytes != req.UploadedSizeBytes || req.UploadedSizeBytes > upload.grant.MaxBytes { - return &agentpb.CompleteUploadResponse{Receipt: s.receipt(req.Meta, agentpb.ResultCode_RESULT_CODE_REJECTED, agentpb.FailureCode_FAILURE_CODE_INVALID_ARGUMENT, "uploaded asset does not match grant", false), State: agentpb.UploadState_UPLOAD_STATE_FAILED}, nil - } - if upload.completed { - return &agentpb.CompleteUploadResponse{Receipt: s.receipt(req.Meta, agentpb.ResultCode_RESULT_CODE_ACCEPTED, agentpb.FailureCode_FAILURE_CODE_UNSPECIFIED, "duplicate mock upload completion", false), State: agentpb.UploadState_UPLOAD_STATE_COMPLETED}, nil - } - if upload.grant.ExpiresAtUnixMs <= s.now().UnixMilli() { - return &agentpb.CompleteUploadResponse{Receipt: s.receipt(req.Meta, agentpb.ResultCode_RESULT_CODE_REJECTED, agentpb.FailureCode_FAILURE_CODE_FAILED_PRECONDITION, "upload grant has expired", true), State: agentpb.UploadState_UPLOAD_STATE_FAILED}, nil - } - upload.completed = true - upload.state = agentpb.UploadState_UPLOAD_STATE_COMPLETED - s.uploads[req.UploadId] = upload - return &agentpb.CompleteUploadResponse{Receipt: s.receipt(req.Meta, agentpb.ResultCode_RESULT_CODE_ACCEPTED, agentpb.FailureCode_FAILURE_CODE_UNSPECIFIED, "mock upload notification completion accepted", false), State: upload.state}, nil -} - -func requireIdempotency(meta *agentpb.RequestMeta) error { - if meta == nil || meta.OperationId == "" || meta.IdempotencyKey == "" { - return status.Error(codes.InvalidArgument, "operation and idempotency key are required") - } - return nil -} - func (s *Server) authorize(ctx context.Context, meta *agentpb.RequestMeta) error { if err := s.executionJournalReady(); err != nil { return err diff --git a/internal/rpc/server_test.go b/internal/rpc/server_test.go index 2caa8fa..44d7a68 100644 --- a/internal/rpc/server_test.go +++ b/internal/rpc/server_test.go @@ -9,8 +9,6 @@ import ( "time" agentpb "git.ipao.vip/rogee/go-sip/gen/agent" - "git.ipao.vip/rogee/go-sip/internal/contract" - "git.ipao.vip/rogee/go-sip/internal/testfixture" "google.golang.org/grpc" "google.golang.org/grpc/codes" "google.golang.org/grpc/credentials" @@ -18,7 +16,6 @@ import ( "google.golang.org/grpc/peer" "google.golang.org/grpc/status" "google.golang.org/grpc/test/bufconn" - "google.golang.org/protobuf/proto" "reflect" ) @@ -102,146 +99,13 @@ func TestSessionGenerationFencesOlderRequests(t *testing.T) { }) require.NoError(t, err) - _, err = server.GetBootstrap(context.Background(), &agentpb.GetBootstrapRequest{Meta: testMeta("old", "read-old", 1)}) + err = server.authorize(context.Background(), testMeta("old", "read-old", 1)) require.Error(t, err) require.Equal(t, codes.Aborted, status.Code(err)) fresh := testMeta("fresh", "read-fresh", 2) fresh.BootId = "boot-2" fresh.DispatcherEpoch = "epoch-2" - _, err = server.GetBootstrap(context.Background(), &agentpb.GetBootstrapRequest{Meta: fresh}) - require.NoError(t, err) -} - -func TestExecuteIdempotencyAndBinding(t *testing.T) { - now := time.Date(2026, 9, 18, 1, 0, 0, 0, time.UTC) - server := activatedServer(now, t) - raw, err := testfixture.Execute() - require.NoError(t, err) - envelope, payload, err := contract.DecodeExecute(raw) - require.NoError(t, err) - meta := testMeta("execute-1", "execute-key", 1) - req := &agentpb.ExecuteRequest{Meta: meta, Binding: &agentpb.ExecutionBinding{TenantId: envelope.TenantID, TenantKey: envelope.TenantKey, ExecutionId: payload.ExecutionID, TaskId: payload.TaskID, TaskItemId: payload.TaskItemID, TaskRevision: payload.TaskRevision, AgentVersionId: payload.AgentVersionID}, CallExecuteJson: raw, ConfigSha256: "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"} - first, err := server.Execute(context.Background(), req) - require.NoError(t, err) - require.Equal(t, agentpb.ResultCode_RESULT_CODE_ACCEPTED, first.Receipt.Result) - replay, err := server.Execute(context.Background(), req) - require.NoError(t, err) - require.Equal(t, first.Receipt.Meta.OperationId, replay.Receipt.Meta.OperationId) - conflictReq := proto.Clone(req).(*agentpb.ExecuteRequest) - conflictReq.ConfigSha256 = "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb" - conflict, err := server.Execute(context.Background(), conflictReq) - require.NoError(t, err) - require.Equal(t, agentpb.ResultCode_RESULT_CODE_CONFLICT, conflict.Receipt.Result) -} - -func TestRealModeRejectsExecutionOutsideCallWindow(t *testing.T) { - server := NewServer(ServerOptions{ - Mode: "mixed", - Now: func() time.Time { return time.Date(2026, 9, 18, 12, 0, 0, 0, time.UTC) }, - }) - activateTestServer(t, server) - response, err := server.GetExecutionPermit(context.Background(), &agentpb.GetExecutionPermitRequest{ - Meta: testMeta("permit-window", "permit-window-key", 1), - Binding: &agentpb.ExecutionBinding{ExecutionId: "execution-window"}, - ResourceReservationId: "reservation-window", - }) - if response != nil || status.Code(err) != codes.FailedPrecondition { - t.Fatalf("response=%+v err=%v code=%s", response, err, status.Code(err)) - } -} - -func TestRealModeRejectsMockUploadDataPlane(t *testing.T) { - server := NewServer(ServerOptions{Mode: "real"}) - response, err := server.RequestUpload(context.Background(), &agentpb.RequestUploadRequest{}) - if response != nil || status.Code(err) != codes.Unimplemented { - t.Fatalf("response=%+v err=%v code=%s", response, err, status.Code(err)) - } -} - -func TestAdmissionAndControlCAS(t *testing.T) { - server := activatedServer(time.Date(2026, 9, 18, 1, 0, 0, 0, time.UTC), t) - meta := testMeta("admission-1", "admission-key", 1) - admission, err := server.SetAdmissionState(context.Background(), &agentpb.SetAdmissionStateRequest{Meta: meta, Target: &agentpb.AgentBinding{AgentId: "agent-1", CellId: "cell-1"}, State: agentpb.AdmissionState_ADMISSION_STATE_OPEN}) - require.NoError(t, err) - require.Equal(t, uint64(1), admission.AppliedAdmissionGeneration) - conflict, err := server.SetAdmissionState(context.Background(), &agentpb.SetAdmissionStateRequest{Meta: testMeta("admission-2", "admission-key-2", 1), Target: &agentpb.AgentBinding{AgentId: "agent-1", CellId: "cell-1"}, State: agentpb.AdmissionState_ADMISSION_STATE_CLOSED}) - require.NoError(t, err) - require.Equal(t, agentpb.ResultCode_RESULT_CODE_CONFLICT, conflict.Receipt.Result) - - raw, err := testfixture.Execute() - require.NoError(t, err) - envelope, payload, err := contract.DecodeExecute(raw) - require.NoError(t, err) - executeMeta := testMeta("execute-control", "execute-control-key", 1) - _, err = server.Execute(context.Background(), &agentpb.ExecuteRequest{Meta: executeMeta, Binding: &agentpb.ExecutionBinding{TenantId: envelope.TenantID, TenantKey: envelope.TenantKey, ExecutionId: payload.ExecutionID, TaskId: payload.TaskID, TaskItemId: payload.TaskItemID, TaskRevision: payload.TaskRevision, AgentVersionId: payload.AgentVersionID}, CallExecuteJson: raw, ConfigSha256: "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"}) - require.NoError(t, err) - controlMeta := testMeta("control-1", "control-key", 1) - binding := &agentpb.ExecutionBinding{TenantId: envelope.TenantID, TenantKey: envelope.TenantKey, ExecutionId: payload.ExecutionID, TaskId: payload.TaskID, TaskItemId: payload.TaskItemID, TaskRevision: payload.TaskRevision, AgentVersionId: payload.AgentVersionID} - pauseRequest := &agentpb.ApplyTaskControlRequest{Meta: controlMeta, Binding: proto.Clone(binding).(*agentpb.ExecutionBinding), Action: agentpb.ControlAction_CONTROL_ACTION_PAUSE, ActiveCallPolicy: agentpb.ActiveCallPolicy_ACTIVE_CALL_POLICY_DRAIN, ExpectedTaskRevision: payload.TaskRevision} - foreign := proto.Clone(pauseRequest).(*agentpb.ApplyTaskControlRequest) - foreign.Meta = testMeta("foreign-control", "foreign-control-key", 1) - foreign.Binding.TenantKey = "different-tenant" - foreignResponse, err := server.ApplyTaskControl(context.Background(), foreign) - require.NoError(t, err) - require.Equal(t, agentpb.ResultCode_RESULT_CODE_CONFLICT, foreignResponse.Receipt.Result) - paused, err := server.ApplyTaskControl(context.Background(), pauseRequest) - require.NoError(t, err) - require.Equal(t, agentpb.ResultCode_RESULT_CODE_APPLIED, paused.Receipt.Result) - duplicate, err := server.ApplyTaskControl(context.Background(), pauseRequest) - require.NoError(t, err) - require.Equal(t, agentpb.ResultCode_RESULT_CODE_APPLIED, duplicate.Receipt.Result) - require.Equal(t, paused.AppliedTaskRevision, duplicate.AppliedTaskRevision) - blockedPermit, err := server.GetExecutionPermit(context.Background(), &agentpb.GetExecutionPermitRequest{Meta: testMeta("permit-paused", "permit-paused-key", 1), Binding: &agentpb.ExecutionBinding{ExecutionId: payload.ExecutionID, TaskRevision: paused.AppliedTaskRevision}, ResourceReservationId: "paused-reservation"}) - require.NoError(t, err) - require.Equal(t, agentpb.ResultCode_RESULT_CODE_REJECTED, blockedPermit.Receipt.Result) - if blockedPermit.Permit != nil { - t.Fatal("paused execution received a permit") - } - resetAttempt, err := server.Execute(context.Background(), &agentpb.ExecuteRequest{Meta: testMeta("execute-after-pause", "execute-after-pause-key", 1), Binding: &agentpb.ExecutionBinding{TenantId: envelope.TenantID, TenantKey: envelope.TenantKey, ExecutionId: payload.ExecutionID, TaskId: payload.TaskID, TaskItemId: payload.TaskItemID, TaskRevision: payload.TaskRevision, AgentVersionId: payload.AgentVersionID}, CallExecuteJson: raw, ConfigSha256: "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"}) - require.NoError(t, err) - require.Equal(t, agentpb.ResultCode_RESULT_CODE_REJECTED, resetAttempt.Receipt.Result) - retraced := proto.Clone(pauseRequest).(*agentpb.ApplyTaskControlRequest) - retraced.Meta.RequestId = "recovered-control-request" - retraced.Meta.TraceId = "recovered-control-trace" - replayed, err := server.ApplyTaskControl(context.Background(), retraced) - require.NoError(t, err) - if !proto.Equal(paused, replayed) { - t.Fatal("transport metadata changed a control's business identity") - } - changed := proto.Clone(pauseRequest).(*agentpb.ApplyTaskControlRequest) - changed.Reason = "different control content" - conflicted, err := server.ApplyTaskControl(context.Background(), changed) - require.NoError(t, err) - require.Equal(t, agentpb.ResultCode_RESULT_CODE_CONFLICT, conflicted.Receipt.Result) - binding.TaskRevision = paused.AppliedTaskRevision - stopped, err := server.ApplyTaskControl(context.Background(), &agentpb.ApplyTaskControlRequest{Meta: testMeta("control-2", "control-key-2", 1), Binding: proto.Clone(binding).(*agentpb.ExecutionBinding), Action: agentpb.ControlAction_CONTROL_ACTION_STOP, ActiveCallPolicy: agentpb.ActiveCallPolicy_ACTIVE_CALL_POLICY_HANGUP, ExpectedTaskRevision: paused.AppliedTaskRevision}) - require.NoError(t, err) - require.Equal(t, agentpb.ExecutionState_EXECUTION_STATE_TERMINAL, stopped.State) - original, err := server.ApplyTaskControl(context.Background(), pauseRequest) - require.NoError(t, err) - require.Equal(t, paused.AppliedTaskRevision, original.AppliedTaskRevision) - snapshot, err := server.QueryExecution(context.Background(), &agentpb.QueryExecutionRequest{Meta: testMeta("query-after-stop", "", 1), Binding: pauseRequest.Binding}) - require.NoError(t, err) - require.Equal(t, stopped.AppliedTaskRevision, snapshot.Snapshot.Binding.TaskRevision) - binding.TaskRevision = stopped.AppliedTaskRevision - resumed, err := server.ApplyTaskControl(context.Background(), &agentpb.ApplyTaskControlRequest{Meta: testMeta("control-3", "control-key-3", 1), Binding: proto.Clone(binding).(*agentpb.ExecutionBinding), Action: agentpb.ControlAction_CONTROL_ACTION_RESUME, ActiveCallPolicy: agentpb.ActiveCallPolicy_ACTIVE_CALL_POLICY_DRAIN, ExpectedTaskRevision: stopped.AppliedTaskRevision}) - require.NoError(t, err) - require.Equal(t, agentpb.ResultCode_RESULT_CODE_REJECTED, resumed.Receipt.Result) -} - -func TestFactDeduplication(t *testing.T) { - server := activatedServer(time.Unix(100, 0), t) - fact := &agentpb.ExecutionFact{FactId: "fact-1", ContentSha256: "digest-a", Binding: &agentpb.ExecutionBinding{ExecutionId: "execution-1"}, Kind: agentpb.FactKind_FACT_KIND_CALL_STATUS} - first, err := server.ReportExecutionEvent(context.Background(), &agentpb.ReportExecutionEventRequest{Meta: testMeta("fact-1", "fact-key-1", 1), Fact: fact}) - require.NoError(t, err) - require.Equal(t, agentpb.ResultCode_RESULT_CODE_ACCEPTED, first.Receipt.Result) - replay, err := server.ReportExecutionEvent(context.Background(), &agentpb.ReportExecutionEventRequest{Meta: testMeta("fact-2", "fact-key-2", 1), Fact: fact}) - require.NoError(t, err) - require.Equal(t, agentpb.ResultCode_RESULT_CODE_ACCEPTED, replay.Receipt.Result) - fact.ContentSha256 = "digest-b" - conflict, err := server.ReportExecutionEvent(context.Background(), &agentpb.ReportExecutionEventRequest{Meta: testMeta("fact-3", "fact-key-3", 1), Fact: fact}) - require.NoError(t, err) - require.Equal(t, agentpb.ResultCode_RESULT_CODE_CONFLICT, conflict.Receipt.Result) + require.NoError(t, server.authorize(context.Background(), fresh)) } func TestGeneratedUnaryServiceWiring(t *testing.T) { diff --git a/internal/rpc/service_test.go b/internal/rpc/service_test.go index 3d58bb5..085f35b 100644 --- a/internal/rpc/service_test.go +++ b/internal/rpc/service_test.go @@ -21,3 +21,16 @@ func TestAgentControlServiceOnlyExposesApprovedMethods(t *testing.T) { t.Fatalf("Agent gRPC method surface = %v, want %v", got, want) } } + +func TestAgentServerHasNoRetiredBusinessMethods(t *testing.T) { + server := reflect.TypeOf(&Server{}) + for _, name := range []string{ + "GetBootstrap", "SetAdmissionState", "Execute", "ExecuteAuthorized", + "GetExecutionPermit", "ApplyTaskControl", "QueryExecution", + "ReportExecutionEvent", "RequestUpload", "CompleteUpload", + } { + if _, exists := server.MethodByName(name); exists { + t.Fatalf("retired Agent RPC method %s is still implemented", name) + } + } +} diff --git a/internal/rpc/upload_test.go b/internal/rpc/upload_test.go deleted file mode 100644 index fba59a8..0000000 --- a/internal/rpc/upload_test.go +++ /dev/null @@ -1,71 +0,0 @@ -package rpc - -import ( - "context" - "strings" - "testing" - "time" - - agentpb "git.ipao.vip/rogee/go-sip/gen/agent" - "google.golang.org/grpc/codes" - "google.golang.org/grpc/status" - "google.golang.org/protobuf/proto" -) - -func TestUploadBindingAndExpiryGuards(t *testing.T) { - now := time.Unix(100, 0) - server := NewServer(ServerOptions{ - Now: func() time.Time { return now }, - UploadPolicy: &agentpb.UploadPolicy{Enabled: true, MaxAssetBytes: 1024}, - }) - _, err := server.ActivateAgent(context.Background(), &agentpb.ActivateAgentRequest{ - Meta: testMeta("activate-upload", "", 0), - Binding: &agentpb.AgentBinding{AgentId: "agent-1", CellId: "cell-1", ExpectedBootId: "boot-1", DispatcherEpoch: "epoch-1", SessionGeneration: 1}, - ActivationOperationId: "activate-upload", - SessionExpiresAtUnixMs: now.Add(time.Hour).UnixMilli(), - }) - if err != nil { - t.Fatal(err) - } - binding := &agentpb.ExecutionBinding{TenantId: "tenant-1", TenantKey: "tenant-demo-key", ExecutionId: "execution-upload", TaskId: "task-upload", TaskItemId: "item-upload", TaskRevision: 1} - asset := &agentpb.AssetDescriptor{Kind: agentpb.AssetKind_ASSET_KIND_RECORDING, AssetId: "recording-upload", ExecutionId: binding.ExecutionId, Format: "wav", SizeBytes: 4, ChecksumSha256: strings.Repeat("a", 64), Channels: 1, SampleRateHz: 16000, DurationMs: 1} - request := &agentpb.RequestUploadRequest{Meta: testMeta("upload-request", "upload-request-key", 1), Binding: binding, Asset: asset, UploadId: "upload-binding"} - created, err := server.RequestUpload(context.Background(), request) - if err != nil { - t.Fatal(err) - } - if created.Grant == nil || created.Grant.ExpiresAtUnixMs <= now.UnixMilli() { - t.Fatalf("invalid grant: %+v", created.Grant) - } - - mismatched := proto.Clone(request).(*agentpb.RequestUploadRequest) - mismatched.Meta = testMeta("upload-conflict", "upload-conflict-key", 1) - mismatched.Binding = proto.Clone(binding).(*agentpb.ExecutionBinding) - mismatched.Binding.ExecutionId = "other-execution" - conflict, err := server.RequestUpload(context.Background(), mismatched) - if err != nil { - t.Fatal(err) - } - if conflict.Receipt == nil || conflict.Receipt.Result != agentpb.ResultCode_RESULT_CODE_CONFLICT { - t.Fatalf("unexpected upload conflict: %+v", conflict) - } - - now = now.Add(6 * time.Minute) - completed, err := server.CompleteUpload(context.Background(), &agentpb.CompleteUploadRequest{ - Meta: testMeta("upload-complete-expired", "upload-complete-expired-key", 1), - Binding: binding, - Asset: asset, - UploadId: request.UploadId, - UploadedSizeBytes: asset.SizeBytes, - UploadedChecksumSha256: asset.ChecksumSha256, - }) - if err != nil { - t.Fatal(err) - } - if completed.Receipt == nil || completed.Receipt.Result != agentpb.ResultCode_RESULT_CODE_REJECTED || completed.Receipt.Failure == nil || completed.Receipt.Failure.Code != agentpb.FailureCode_FAILURE_CODE_FAILED_PRECONDITION { - t.Fatalf("unexpected expired completion: %+v", completed) - } - if status.Code(err) != codes.OK { - t.Fatalf("unexpected status code: %s", status.Code(err)) - } -}