refactor(rpc): remove retired Agent business handlers

This commit is contained in:
2026-09-30 13:21:37 +08:00
parent 842735a17b
commit 793b6448e2
15 changed files with 47 additions and 1565 deletions
@@ -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 消息与其他引用仍须继续清理,未接触现存业务数据或真实外部服务。
## 验收台账
-122
View File
@@ -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)
}
}
+31
View File
@@ -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)
-123
View File
@@ -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
}
-201
View File
@@ -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)
}
})
}
}
-100
View File
@@ -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)
}
}
-32
View File
@@ -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)
}
@@ -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")
}
})
}
}
@@ -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")
}
}
@@ -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())
})
}
}
-74
View File
@@ -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)
}
-589
View File
@@ -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
+2 -138
View File
@@ -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) {
+13
View File
@@ -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)
}
}
}
-71
View File
@@ -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))
}
}