diff --git a/docs/evidence/saas-dispatcher-implementation.md b/docs/evidence/saas-dispatcher-implementation.md index 52b2c80..9c725fe 100644 --- a/docs/evidence/saas-dispatcher-implementation.md +++ b/docs/evidence/saas-dispatcher-implementation.md @@ -117,6 +117,7 @@ - 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 消息与其他引用仍须继续清理,未接触现存业务数据或真实外部服务。 +- Agent 旧执行日志根因:新增「首次激活→同一路径重启」测试,复现现行入口曾在初次启动自动写出废弃的 `.executions` 文件,下一次启动又将它识别为旧未交付状态并拒绝服务。移除旧日志写入、回放状态和闲置执行配置,只保留当前 `.approved` 日志、会话代际和旧文件存在时拒绝启动的保护;回归确认首次启动与重启不会产生新旧执行日志。已有 `.executions` 一律保留并失败关闭,不自动清理或猜测其业务内容;当前测试只使用临时目录。 ## 验收台账 diff --git a/internal/rpc/approved_server_test.go b/internal/rpc/approved_server_test.go index 2e6e5f4..b06a995 100644 --- a/internal/rpc/approved_server_test.go +++ b/internal/rpc/approved_server_test.go @@ -104,6 +104,36 @@ func TestApprovedAgentServerRegistersCallsBeforeExecuteAckAndDrainsOnControl(t * } } +func TestApprovedAgentServerFreshStartAndRestartDoNotCreateLegacyExecutionState(t *testing.T) { + state := filepath.Join(t.TempDir(), "agent-session.json") + options := ServerOptions{ + Mode: "mock", StatePath: state, ApprovedDispatcherID: "dispatcher-1", + Status: &agentpb.AgentStatus{AgentId: "agent-1", CellId: "cell-1", BootId: "boot-1"}, + LoadedSIP: func(context.Context) (map[string]int64, error) { return map[string]int64{"trunk-mock": 8}, nil }, + } + worker := &ApprovedCallWorker{Lifecycle: context.Background(), Calls: &agent.TaskCalls{}, + Prepare: prepareWorkerRun(func(context.Context, ApprovedExecution) error { return nil }), + OnFailure: func(ApprovedExecution, error) error { return nil }, + } + server, err := NewApprovedAgentServer(options, worker) + if err != nil { + t.Fatal(err) + } + if _, err := server.ActivateAgent(context.Background(), &agentpb.ActivateAgentRequest{ + Meta: testMeta("activate-without-legacy", "", 0), + Binding: &agentpb.AgentBinding{AgentId: "agent-1", CellId: "cell-1", ExpectedBootId: "boot-1", DispatcherEpoch: "epoch-1", SessionGeneration: 1, DispatcherId: "dispatcher-1"}, + ActivationOperationId: "activate-without-legacy", + }); err != nil { + t.Fatal(err) + } + if _, err := os.Lstat(state + ".executions"); !errors.Is(err, os.ErrNotExist) { + t.Fatalf("current Agent wrote an obsolete execution journal: %v", err) + } + if _, err := NewApprovedAgentServer(options, worker); err != nil { + t.Fatalf("approved Agent could not restart with its current session: %v", err) + } +} + func TestApprovedAgentServerRefusesMissingOrCompetingAdapters(t *testing.T) { process := context.Background() calls := &agent.TaskCalls{} diff --git a/internal/rpc/execution_journal.go b/internal/rpc/execution_journal.go deleted file mode 100644 index 55ad9d4..0000000 --- a/internal/rpc/execution_journal.go +++ /dev/null @@ -1,136 +0,0 @@ -package rpc - -import ( - "bytes" - "encoding/json" - "errors" - "io" - "os" - - agentpb "git.ipao.vip/rogee/go-sip/gen/agent" - "google.golang.org/grpc/codes" - "google.golang.org/grpc/status" - "google.golang.org/protobuf/proto" -) - -type savedOperation struct { - Digest string `json:"digest"` - Receipt *agentpb.OperationReceipt `json:"receipt"` - Control *agentpb.ApplyTaskControlResponse `json:"control,omitempty"` -} -type savedExecution struct { - Binding *agentpb.ExecutionBinding `json:"binding"` - Digest string `json:"digest"` - State agentpb.ExecutionState `json:"state"` - Revision int64 `json:"revision"` - CallState string `json:"call_state"` - TerminalObservedAtUnixMs int64 `json:"terminal_observed_at_unix_ms,omitempty"` - ControlAction agentpb.ControlAction `json:"control_action"` -} -type executionJournal struct { - Version int `json:"version"` - Mode string `json:"mode"` - AgentID string `json:"agent_id"` - CellID string `json:"cell_id"` - Operations map[string]savedOperation `json:"operations"` - Executions map[string]savedExecution `json:"executions"` -} - -func (s *Server) loadExecutionJournal() error { - if s.executionPath == "" { - return nil - } - raw, err := os.ReadFile(s.executionPath) - if errors.Is(err, os.ErrNotExist) { - if _, sessionErr := os.Stat(s.sessions.statePath); sessionErr == nil { - return errors.New("execution journal missing beside existing session journal") - } else if !errors.Is(sessionErr, os.ErrNotExist) { - return sessionErr - } - // Establish the empty execution journal before any activation can be - // acknowledged; later absence must not silently erase execution history. - return s.persistExecutionJournalLocked() - } - if err != nil { - return err - } - decoder := json.NewDecoder(bytes.NewReader(raw)) - decoder.DisallowUnknownFields() - var journal executionJournal - if err := decoder.Decode(&journal); err != nil { - return err - } - if err := decoder.Decode(new(any)); err != io.EOF { - return errors.New("execution journal has trailing data") - } - if journal.Version != 1 || journal.Mode != s.mode || journal.AgentID != s.status.AgentId || journal.CellID != s.status.CellId || journal.Operations == nil || journal.Executions == nil { - return errors.New("execution journal version or identity is invalid") - } - for key, record := range journal.Operations { - if record.Receipt == nil || record.Receipt.Meta == nil || len(record.Digest) != 64 { - return errors.New("invalid persisted operation") - } - if record.Control != nil && !proto.Equal(record.Control.Receipt, record.Receipt) { - return errors.New("control receipt differs from operation receipt") - } - s.operations[key] = operationRecord{digest: record.Digest, receipt: record.Receipt, control: record.Control} - } - for id, record := range journal.Executions { - if record.Binding == nil || record.Binding.ExecutionId != id || record.Revision != record.Binding.TaskRevision { - return errors.New("invalid persisted execution binding") - } - if _, ok := agentpb.ExecutionState_name[int32(record.State)]; !ok { - return errors.New("invalid persisted execution state") - } - if _, ok := agentpb.ControlAction_name[int32(record.ControlAction)]; !ok { - return errors.New("invalid persisted control action") - } - mockTerminal := record.CallState == "mock_no_answer" || record.CallState == "mock_deadline_closed_without_dial" - if mockTerminal && (record.TerminalObservedAtUnixMs <= 0 || record.State != agentpb.ExecutionState_EXECUTION_STATE_TERMINAL) || - !mockTerminal && record.TerminalObservedAtUnixMs != 0 { - return errors.New("invalid durable Mock terminal observation") - } - execution := &executionRecord{binding: record.Binding, executeDigest: record.Digest, taskRevision: record.Revision, callState: record.CallState, - terminalObservedAtUnixMs: record.TerminalObservedAtUnixMs, controlAction: record.ControlAction, state: record.State} - // A new process cannot infer Asterisk's state from an old local snapshot. - // Never restore permits or clear unknown occupancy because a process restarted. - if execution.state != agentpb.ExecutionState_EXECUTION_STATE_TERMINAL { - execution.state = agentpb.ExecutionState_EXECUTION_STATE_UNKNOWN - execution.unknown = true - } - s.executions[id] = execution - } - return nil -} - -// Caller holds s.mu. Failure poisons further admission: in-memory state must -// never be acknowledged as durable after an unsuccessful journal write. -func (s *Server) persistExecutionJournalLocked() error { - if s.executionErr != nil { - return status.Errorf(codes.Internal, "execution journal unavailable: %v", s.executionErr) - } - if s.executionPath == "" { - return nil - } - journal := executionJournal{Version: 1, Mode: s.mode, AgentID: s.status.AgentId, CellID: s.status.CellId, Operations: make(map[string]savedOperation, len(s.operations)), Executions: make(map[string]savedExecution, len(s.executions))} - for key, record := range s.operations { - journal.Operations[key] = savedOperation{Digest: record.digest, Receipt: record.receipt, Control: record.control} - } - for id, record := range s.executions { - journal.Executions[id] = savedExecution{Binding: record.binding, Digest: record.executeDigest, State: record.state, Revision: record.taskRevision, - CallState: record.callState, TerminalObservedAtUnixMs: record.terminalObservedAtUnixMs, ControlAction: record.controlAction} - } - if err := writeRPCJournal(s.executionPath, journal); err != nil { - s.executionErr = err - return status.Errorf(codes.Internal, "persist execution journal: %v", err) - } - return nil -} -func (s *Server) executionJournalReady() error { - s.mu.Lock() - defer s.mu.Unlock() - if s.executionErr != nil { - return status.Errorf(codes.Internal, "execution journal unavailable: %v", s.executionErr) - } - return nil -} diff --git a/internal/rpc/server.go b/internal/rpc/server.go index 9d70d48..60995f3 100644 --- a/internal/rpc/server.go +++ b/internal/rpc/server.go @@ -13,8 +13,6 @@ import ( agentpb "git.ipao.vip/rogee/go-sip/gen/agent" "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/contract" "google.golang.org/grpc" "google.golang.org/grpc/codes" @@ -24,29 +22,20 @@ import ( "google.golang.org/protobuf/proto" ) -// ServerOptions contains deployment-safe identity, immutable contract -// artifacts, and mock policy inputs. Production credentials are supplied to -// grpc.Server via TLS credentials; no certificate or secret is stored here. +// ServerOptions contains deployment-bound identity and mock policy inputs. +// Production credentials are supplied to grpc.Server through TLS credentials. type ServerOptions struct { Mode string Status *agentpb.AgentStatus - UploadPolicy *agentpb.UploadPolicy - ConfigReferences []*agentpb.ConfigReference StaticArtifactRaw []byte StaticArtifactExpected contract.StaticArtifactExpectation - AISnapshotRaw []byte - AIAuthorizationRaw []byte Now func() time.Time RequirePeerCertificate bool ApprovedDispatcherID string PeerAgentIDs map[string]string PeerCertificateFingerprints map[string]struct{} StatePath string - CallLogger *calllog.Logger - // MockAuthorizedOriginate is explicitly supplied only in isolated mock mode. - // Real/mixed dialing is not authorized through the new local contract. - MockAuthorizedOriginate func(context.Context, *agentpb.ExecuteAuthorizedRequest) error - // LoadedSIP must inspect what the mock Agent actually loaded; nil fails closed. + // LoadedSIP reports the revision the mock Agent actually loaded; nil fails closed. LoadedSIP func(context.Context) (map[string]int64, error) // The mock may have issued a call even if its outcome is unknown. MockApprovedOriginate func(context.Context, ApprovedExecution) error @@ -54,29 +43,21 @@ type ServerOptions struct { ApprovedTaskCalls *agent.TaskCalls } -// Server is the local Unary gRPC state boundary. It owns session/fencing and -// durable-operation semantics for the RPC layer; Dispatcher SQLite remains the -// business source of truth for quotas and task state. +// Server owns the current Agent session, execution and control RPC boundary. +// Dispatcher SQLite remains authoritative for quotas and task state. type Server struct { agentpb.UnimplementedAgentControlServiceServer mode string now func() time.Time status *agentpb.AgentStatus - uploadPolicy *agentpb.UploadPolicy - configReferences []*agentpb.ConfigReference staticArtifact contract.StaticCellArtifact staticArtifactEnabled bool staticArtifactError error - aiSnapshot ai.Snapshot - aiAuthorizationRaw []byte - aiConfigError error requirePeerCertificate bool approvedDispatcherID string peerAgentIDs map[string]string peerCertificateFingerprints map[string]struct{} - callLogger *calllog.Logger - mockAuthorizedOriginate func(context.Context, *agentpb.ExecuteAuthorizedRequest) error loadedSIP func(context.Context) (map[string]int64, error) mockApprovedOriginate func(context.Context, ApprovedExecution) error approvedTaskCalls *agent.TaskCalls @@ -84,46 +65,6 @@ type Server struct { approvedJournalDir string approvedMu sync.Mutex - executionPath string - executionErr error - mu sync.Mutex - operations map[string]operationRecord - admissions map[string]admissionRecord - executions map[string]*executionRecord - facts map[string]string - uploads map[string]uploadRecord -} - -type operationRecord struct { - digest string - receipt *agentpb.OperationReceipt - control *agentpb.ApplyTaskControlResponse -} - -type admissionRecord struct { - state agentpb.AdmissionState - generation uint64 -} - -type executionRecord struct { - executeDigest string - binding *agentpb.ExecutionBinding - state agentpb.ExecutionState - taskRevision int64 - callState string - terminalObservedAtUnixMs int64 - controlAction agentpb.ControlAction - permit *agentpb.ExecutionPermit - unknown bool - phone calllog.Identity -} - -type uploadRecord struct { - binding *agentpb.ExecutionBinding - asset *agentpb.AssetDescriptor - state agentpb.UploadState - grant *agentpb.UploadGrant - completed bool } // NewServer constructs a handler suitable for registration with a gRPC server. @@ -143,84 +84,34 @@ func NewServer(options ServerOptions) *Server { if statusValue.AdmissionState == agentpb.AdmissionState_ADMISSION_STATE_UNSPECIFIED { statusValue.AdmissionState = agentpb.AdmissionState_ADMISSION_STATE_CLOSED } - uploadPolicy := &agentpb.UploadPolicy{} - if options.UploadPolicy != nil { - uploadPolicy = proto.Clone(options.UploadPolicy).(*agentpb.UploadPolicy) - } var staticArtifact contract.StaticCellArtifact var staticArtifactError error staticArtifactEnabled := len(options.StaticArtifactRaw) != 0 if staticArtifactEnabled { staticArtifact, staticArtifactError = contract.ValidateStaticArtifact(options.StaticArtifactRaw, options.StaticArtifactExpected) } - var aiSnapshot ai.Snapshot - var aiConfigError error - aiConfigured := len(options.AISnapshotRaw) != 0 || len(options.AIAuthorizationRaw) != 0 - if aiConfigured { - if len(options.AISnapshotRaw) == 0 || len(options.AIAuthorizationRaw) == 0 { - aiConfigError = errors.New("AI snapshot and authorization are required together") - } else { - aiSnapshot, aiConfigError = ai.Validate(options.AISnapshotRaw) - } - } server := &Server{ mode: mode, now: now, status: statusValue, - uploadPolicy: uploadPolicy, - configReferences: cloneConfigReferences(options.ConfigReferences), staticArtifact: staticArtifact, staticArtifactEnabled: staticArtifactEnabled, staticArtifactError: staticArtifactError, - aiSnapshot: aiSnapshot, - aiAuthorizationRaw: append([]byte(nil), options.AIAuthorizationRaw...), - aiConfigError: aiConfigError, requirePeerCertificate: options.RequirePeerCertificate, approvedDispatcherID: options.ApprovedDispatcherID, peerAgentIDs: cloneStringMap(options.PeerAgentIDs), peerCertificateFingerprints: cloneSet(options.PeerCertificateFingerprints), - callLogger: options.CallLogger, - mockAuthorizedOriginate: options.MockAuthorizedOriginate, loadedSIP: options.LoadedSIP, mockApprovedOriginate: options.MockApprovedOriginate, approvedTaskCalls: options.ApprovedTaskCalls, sessions: NewSessionRegistry(options.StatePath), - operations: make(map[string]operationRecord), - admissions: make(map[string]admissionRecord), - executions: make(map[string]*executionRecord), - facts: make(map[string]string), - uploads: make(map[string]uploadRecord), } if options.StatePath != "" { server.approvedJournalDir = options.StatePath + ".approved" - server.executionPath = options.StatePath + ".executions" - server.executionErr = server.loadExecutionJournal() } return server } -func (s *Server) validateAIExecution(binding *agentpb.ExecutionBinding, configSHA256 string) error { - if s.aiConfigError == nil && len(s.aiAuthorizationRaw) == 0 { - return nil - } - if s.aiConfigError != nil { - return status.Errorf(codes.FailedPrecondition, "AI authorization configuration is invalid: %v", s.aiConfigError) - } - if binding == nil { - return status.Error(codes.InvalidArgument, "execution binding is required") - } - if configSHA256 == "" || configSHA256 != s.aiSnapshot.Digest { - return status.Error(codes.FailedPrecondition, "AI config digest does not match the authorized snapshot") - } - if binding.AgentVersionId != s.aiSnapshot.AgentVersionID { - return status.Error(codes.FailedPrecondition, "AI agent version does not match the authorized snapshot") - } - if _, err := ai.ValidateAuthorization(s.aiAuthorizationRaw, s.aiSnapshot, binding.TenantId, binding.TenantKey, s.now()); err != nil { - return status.Errorf(codes.PermissionDenied, "AI authorization rejected: %v", err) - } - return nil -} - // SessionRegistry keeps the newest Dispatcher-approved binding for each Agent. // A newer generation fences all older requests; it does not release unknown // work from an older boot. @@ -444,9 +335,6 @@ func (s *Server) GetAgentStatus(ctx context.Context, req *agentpb.GetAgentStatus } func (s *Server) ActivateAgent(ctx context.Context, req *agentpb.ActivateAgentRequest) (*agentpb.ActivateAgentResponse, error) { - if err := s.executionJournalReady(); err != nil { - return nil, err - } if req == nil || req.Meta == nil || req.Binding == nil { return nil, status.Error(codes.InvalidArgument, "activation metadata and binding are required") } @@ -493,9 +381,6 @@ func (s *Server) ActivateAgent(ctx context.Context, req *agentpb.ActivateAgentRe } func (s *Server) authorize(ctx context.Context, meta *agentpb.RequestMeta) error { - if err := s.executionJournalReady(); err != nil { - return err - } if meta == nil { return status.Error(codes.InvalidArgument, "request metadata is required") } @@ -563,45 +448,6 @@ func (s *Server) responseMeta(meta *agentpb.RequestMeta) *agentpb.ResponseMeta { return &agentpb.ResponseMeta{ProtocolVersion: meta.ProtocolVersion, RequestId: meta.RequestId, TraceId: meta.TraceId, OperationId: meta.OperationId, ObservedAtUnixMs: s.now().UnixMilli(), DispatcherEpoch: meta.DispatcherEpoch, AgentId: meta.AgentId, CellId: meta.CellId, BootId: meta.BootId, SessionGeneration: meta.SessionGeneration} } -func (s *Server) receipt(meta *agentpb.RequestMeta, result agentpb.ResultCode, code agentpb.FailureCode, detail string, retryable bool) *agentpb.OperationReceipt { - receipt := &agentpb.OperationReceipt{Meta: s.responseMeta(meta), Result: result, AcceptedAtUnixMs: s.now().UnixMilli()} - if code != agentpb.FailureCode_FAILURE_CODE_UNSPECIFIED { - receipt.Failure = s.failure(code, detail, retryable) - } - return receipt -} - -func (s *Server) failure(code agentpb.FailureCode, detail string, retryable bool) *agentpb.Failure { - return &agentpb.Failure{Code: code, Detail: detail, Retryable: retryable} -} - -func (s *Server) operationKey(meta *agentpb.RequestMeta) string { - if meta == nil || meta.IdempotencyKey == "" { - return "" - } - return meta.AgentId + "\x00" + meta.OperationId + "\x00" + meta.IdempotencyKey -} - -func (s *Server) replayOperation(meta *agentpb.RequestMeta, digest string) (*agentpb.OperationReceipt, *agentpb.OperationReceipt) { - key := s.operationKey(meta) - if key == "" { - return nil, nil - } - s.mu.Lock() - defer s.mu.Unlock() - if s.executionErr != nil { - return nil, s.receipt(meta, agentpb.ResultCode_RESULT_CODE_UNKNOWN, agentpb.FailureCode_FAILURE_CODE_UNAVAILABLE, "execution journal unavailable", false) - } - record, ok := s.operations[key] - if !ok { - return nil, nil - } - if record.digest != digest { - return nil, s.receipt(meta, agentpb.ResultCode_RESULT_CODE_CONFLICT, agentpb.FailureCode_FAILURE_CODE_ABORTED, "idempotency key content conflict", false) - } - return proto.Clone(record.receipt).(*agentpb.OperationReceipt), nil -} - func messageDigest(message proto.Message) string { encoded, err := proto.Marshal(message) if err != nil { @@ -611,14 +457,6 @@ func messageDigest(message proto.Message) string { return hex.EncodeToString(digest[:]) } -func randomToken() string { - value := make([]byte, 16) - if _, err := cryptorand.Read(value); err != nil { - return "unavailable" - } - return hex.EncodeToString(value) -} - func cloneSession(value *agentpb.Session) *agentpb.Session { if value == nil { return nil @@ -626,16 +464,6 @@ func cloneSession(value *agentpb.Session) *agentpb.Session { return proto.Clone(value).(*agentpb.Session) } -func cloneConfigReferences(values []*agentpb.ConfigReference) []*agentpb.ConfigReference { - result := make([]*agentpb.ConfigReference, 0, len(values)) - for _, value := range values { - if value != nil { - result = append(result, proto.Clone(value).(*agentpb.ConfigReference)) - } - } - return result -} - func cloneStringMap(values map[string]string) map[string]string { if values == nil { return nil diff --git a/internal/rpc/server_test.go b/internal/rpc/server_test.go index 44d7a68..f2eaaf7 100644 --- a/internal/rpc/server_test.go +++ b/internal/rpc/server_test.go @@ -127,7 +127,7 @@ func TestGeneratedUnaryServiceWiring(t *testing.T) { func activatedServer(now time.Time, t *testing.T) *Server { t.Helper() - server := NewServer(ServerOptions{Now: func() time.Time { return now }, UploadPolicy: &agentpb.UploadPolicy{Enabled: true, MaxAssetBytes: 16 << 20}}) + server := NewServer(ServerOptions{Now: func() time.Time { return now }}) meta := testMeta("activate", "", 0) _, 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"}) require.NoError(t, err)