package rpc import ( "context" cryptorand "crypto/rand" "crypto/sha256" "encoding/hex" "encoding/json" "errors" "fmt" "os" "sync" "time" agentv1 "git.ipao.vip/rogee/go-sip/gen/agent/v1" "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" "google.golang.org/grpc/credentials" "google.golang.org/grpc/peer" "google.golang.org/grpc/status" "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. type ServerOptions struct { Mode string Status *agentv1.AgentStatus UploadPolicy *agentv1.UploadPolicy ConfigReferences []*agentv1.ConfigReference StaticArtifactRaw []byte StaticArtifactExpected contract.StaticArtifactExpectation AISnapshotRaw []byte AIAuthorizationRaw []byte AIEgressPoolID string Now func() time.Time RequirePeerCertificate bool PeerAgentIDs map[string]string PeerCertificateFingerprints map[string]struct{} StatePath string CallLogger *calllog.Logger } // 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. type Server struct { agentv1.UnimplementedAgentControlServiceServer mode string now func() time.Time status *agentv1.AgentStatus uploadPolicy *agentv1.UploadPolicy configReferences []*agentv1.ConfigReference staticArtifact contract.StaticCellArtifact staticArtifactEnabled bool staticArtifactError error aiSnapshot ai.Snapshot aiAuthorizationRaw []byte aiEgressPoolID string aiConfigError error requirePeerCertificate bool peerAgentIDs map[string]string peerCertificateFingerprints map[string]struct{} callLogger *calllog.Logger sessions *SessionRegistry 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 *agentv1.OperationReceipt control *agentv1.ApplyTaskControlResponse } type admissionRecord struct { state agentv1.AdmissionState generation uint64 } type executionRecord struct { executeDigest string binding *agentv1.ExecutionBinding state agentv1.ExecutionState taskRevision int64 callState string controlAction agentv1.ControlAction permit *agentv1.ExecutionPermit unknown bool phone calllog.Identity } type uploadRecord struct { binding *agentv1.ExecutionBinding asset *agentv1.AssetDescriptor state agentv1.UploadState grant *agentv1.UploadGrant completed bool } // NewServer constructs a handler suitable for registration with a gRPC server. func NewServer(options ServerOptions) *Server { mode := options.Mode if mode == "" { mode = "mock" } now := options.Now if now == nil { now = time.Now } statusValue := &agentv1.AgentStatus{} if options.Status != nil { statusValue = proto.Clone(options.Status).(*agentv1.AgentStatus) } if statusValue.AdmissionState == agentv1.AdmissionState_ADMISSION_STATE_UNSPECIFIED { statusValue.AdmissionState = agentv1.AdmissionState_ADMISSION_STATE_CLOSED } uploadPolicy := &agentv1.UploadPolicy{} if options.UploadPolicy != nil { uploadPolicy = proto.Clone(options.UploadPolicy).(*agentv1.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 || options.AIEgressPoolID == "" { aiConfigError = errors.New("AI snapshot, authorization and egress pool 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...), aiEgressPoolID: options.AIEgressPoolID, aiConfigError: aiConfigError, requirePeerCertificate: options.RequirePeerCertificate, peerAgentIDs: cloneStringMap(options.PeerAgentIDs), peerCertificateFingerprints: cloneSet(options.PeerCertificateFingerprints), callLogger: options.CallLogger, 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.executionPath = options.StatePath + ".executions" server.executionErr = server.loadExecutionJournal() } return server } func (s *Server) validateAIExecution(binding *agentv1.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.aiEgressPoolID, 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. type SessionRegistry struct { mu sync.Mutex sessions map[string]sessionRecord generations map[string]uint64 statePath string loadErr error } type sessionRecord struct { binding *agentv1.AgentBinding activationOperationID string digest string session *agentv1.Session } func NewSessionRegistry(statePath string) *SessionRegistry { registry := &SessionRegistry{sessions: make(map[string]sessionRecord), generations: make(map[string]uint64), statePath: statePath} if statePath != "" { registry.loadErr = registry.load() } return registry } type sessionJournal struct { Generations map[string]uint64 `json:"generations"` } func (r *SessionRegistry) load() error { data, err := os.ReadFile(r.statePath) if errors.Is(err, os.ErrNotExist) { return nil } if err != nil { return err } var journal sessionJournal if err := json.Unmarshal(data, &journal); err != nil { return err } for agentID, generation := range journal.Generations { if agentID != "" && generation > 0 { r.generations[agentID] = generation } } return nil } func (r *SessionRegistry) persistLocked() error { if r.statePath == "" { return nil } return writeRPCJournal(r.statePath, sessionJournal{Generations: r.generations}) } func (r *SessionRegistry) Activate(binding *agentv1.AgentBinding, activationOperationID, digest string, now time.Time) (*agentv1.Session, bool, error) { if binding == nil || binding.AgentId == "" || binding.CellId == "" || binding.DispatcherEpoch == "" || activationOperationID == "" { return nil, false, status.Error(codes.InvalidArgument, "agent, cell, epoch and activation operation are required") } r.mu.Lock() defer r.mu.Unlock() if r.loadErr != nil { return nil, false, status.Errorf(codes.Internal, "load session journal: %v", r.loadErr) } if existing, ok := r.sessions[binding.AgentId]; ok { if existing.activationOperationID == activationOperationID && existing.digest == digest { return cloneSession(existing.session), true, nil } if binding.SessionGeneration == 0 { binding = proto.Clone(binding).(*agentv1.AgentBinding) binding.SessionGeneration = existing.binding.SessionGeneration + 1 } if binding.SessionGeneration <= existing.binding.SessionGeneration { return nil, false, status.Error(codes.Aborted, "session generation is fenced") } } if binding.SessionGeneration == 0 { binding = proto.Clone(binding).(*agentv1.AgentBinding) binding.SessionGeneration = 1 } if previous, ok := r.generations[binding.AgentId]; ok && binding.SessionGeneration <= previous { return nil, false, status.Error(codes.Aborted, "persisted session generation is fenced") } credential := make([]byte, 32) if _, err := cryptorand.Read(credential); err != nil { return nil, false, status.Errorf(codes.Internal, "create session credential: %v", err) } session := &agentv1.Session{ DispatcherEpoch: binding.DispatcherEpoch, SessionGeneration: binding.SessionGeneration, ExpiresAtUnixMs: now.Add(10 * time.Minute).UnixMilli(), SessionCredential: credential, } previous, hadPrevious := r.sessions[binding.AgentId] r.sessions[binding.AgentId] = sessionRecord{ binding: proto.Clone(binding).(*agentv1.AgentBinding), activationOperationID: activationOperationID, digest: digest, session: cloneSession(session), } r.generations[binding.AgentId] = binding.SessionGeneration if err := r.persistLocked(); err != nil { if hadPrevious { r.sessions[binding.AgentId] = previous } else { delete(r.sessions, binding.AgentId) } return nil, false, status.Errorf(codes.Internal, "persist session journal: %v", err) } return session, false, nil } func (r *SessionRegistry) Authorize(meta *agentv1.RequestMeta, now time.Time) error { if meta == nil || meta.AgentId == "" || meta.CellId == "" || meta.BootId == "" || meta.DispatcherEpoch == "" || meta.SessionGeneration == 0 { return status.Error(codes.InvalidArgument, "complete session metadata is required") } r.mu.Lock() defer r.mu.Unlock() existing, ok := r.sessions[meta.AgentId] if !ok { return status.Error(codes.Unauthenticated, "agent session is not active") } if existing.binding.CellId != meta.CellId || existing.binding.ExpectedBootId != meta.BootId || existing.binding.DispatcherEpoch != meta.DispatcherEpoch || existing.binding.SessionGeneration != meta.SessionGeneration { return status.Error(codes.Aborted, "agent session is fenced") } if existing.session.ExpiresAtUnixMs <= now.UnixMilli() { return status.Error(codes.Unauthenticated, "agent session expired") } return nil } func (s *Server) GetAgentStatus(ctx context.Context, req *agentv1.GetAgentStatusRequest) (*agentv1.GetAgentStatusResponse, error) { if req == nil || req.Meta == nil || req.Meta.AgentId == "" || req.Meta.CellId == "" { return nil, status.Error(codes.InvalidArgument, "status metadata with agent and cell is required") } if err := s.checkConfiguredIdentity(req.Meta.AgentId, req.Meta.CellId); err != nil { return nil, err } if req.Target != nil { if req.Target.AgentId != "" && req.Target.AgentId != req.Meta.AgentId { return nil, status.Error(codes.PermissionDenied, "target agent does not match authenticated agent") } if req.Target.CellId != "" && req.Target.CellId != req.Meta.CellId { return nil, status.Error(codes.PermissionDenied, "target Cell does not match authenticated Cell") } } if err := s.checkPeer(ctx, req.Meta.AgentId); err != nil { return nil, err } preActivation := req.Meta.BootId == "" && req.Meta.DispatcherEpoch == "" && req.Meta.SessionGeneration == 0 if preActivation { if req.Target != nil && (req.Target.ExpectedBootId != "" || req.Target.DispatcherEpoch != "" || req.Target.SessionGeneration != 0) { return nil, status.Error(codes.InvalidArgument, "pre-activation status cannot include session binding") } } else if err := s.sessions.Authorize(req.Meta, s.now()); err != nil { return nil, err } result := proto.Clone(s.status).(*agentv1.AgentStatus) result.SessionActive = !preActivation result.MtlsAuthenticated = s.peerIsAuthenticated(ctx) return &agentv1.GetAgentStatusResponse{Meta: s.responseMeta(req.Meta), Status: result}, nil } func (s *Server) ActivateAgent(ctx context.Context, req *agentv1.ActivateAgentRequest) (*agentv1.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") } if req.Meta.AgentId != req.Binding.AgentId || req.Meta.CellId != req.Binding.CellId || req.Meta.BootId == "" || req.Meta.DispatcherEpoch == "" { return nil, status.Error(codes.InvalidArgument, "activation identity is inconsistent") } if err := s.checkConfiguredIdentity(req.Binding.AgentId, req.Binding.CellId); err != nil { return nil, err } if err := s.checkPeer(ctx, req.Binding.AgentId); err != nil { return nil, err } if s.staticArtifactEnabled { if s.staticArtifactError != nil { return nil, status.Errorf(codes.FailedPrecondition, "static Cell artifact is invalid: %v", s.staticArtifactError) } if req.Binding.CellId != s.staticArtifact.CellID { return nil, status.Error(codes.FailedPrecondition, "activation Cell does not match static artifact") } } binding := proto.Clone(req.Binding).(*agentv1.AgentBinding) if binding.ExpectedBootId == "" { binding.ExpectedBootId = req.Meta.BootId } if binding.ExpectedBootId != req.Meta.BootId { return nil, status.Error(codes.Aborted, "activation boot identity is fenced") } if req.ActivationOperationId == "" { return nil, status.Error(codes.InvalidArgument, "activation operation is required") } digest := messageDigest(req) session, replay, err := s.sessions.Activate(binding, req.ActivationOperationId, digest, s.now()) if err != nil { return nil, err } state := agentv1.ActivationState_ACTIVATION_STATE_ACTIVE if replay { state = agentv1.ActivationState_ACTIVATION_STATE_ACTIVE } return &agentv1.ActivateAgentResponse{Meta: s.responseMeta(req.Meta), State: state, Session: session}, nil } func (s *Server) GetBootstrap(ctx context.Context, req *agentv1.GetBootstrapRequest) (*agentv1.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 &agentv1.GetBootstrapResponse{ Meta: s.responseMeta(req.Meta), State: agentv1.ActivationState_ACTIVATION_STATE_ACTIVE, RuntimeConfigs: cloneConfigReferences(s.configReferences), UploadPolicy: proto.Clone(s.uploadPolicy).(*agentv1.UploadPolicy), }, nil } func (s *Server) SetAdmissionState(ctx context.Context, req *agentv1.SetAdmissionStateRequest) (*agentv1.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 == agentv1.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 &agentv1.SetAdmissionStateResponse{Receipt: s.receipt(req.Meta, agentv1.ResultCode_RESULT_CODE_CONFLICT, agentv1.FailureCode_FAILURE_CODE_ABORTED, "admission generation conflict", false)}, nil } current.generation++ current.state = req.State s.admissions[key] = current return &agentv1.SetAdmissionStateResponse{Receipt: s.receipt(req.Meta, agentv1.ResultCode_RESULT_CODE_APPLIED, agentv1.FailureCode_FAILURE_CODE_UNSPECIFIED, "", false), AppliedAdmissionGeneration: current.generation}, nil } func (s *Server) Execute(ctx context.Context, req *agentv1.ExecuteRequest) (*agentv1.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 &agentv1.ExecuteResponse{Receipt: conflict}, nil } return &agentv1.ExecuteResponse{Receipt: receipt, State: agentv1.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 &agentv1.ExecuteResponse{Receipt: s.receipt(req.Meta, agentv1.ResultCode_RESULT_CODE_CONFLICT, agentv1.FailureCode_FAILURE_CODE_ABORTED, "idempotency key content conflict", false)}, nil } return &agentv1.ExecuteResponse{Receipt: proto.Clone(previous.receipt).(*agentv1.OperationReceipt), State: agentv1.ExecutionState_EXECUTION_STATE_PREPARED}, nil } var priorPermit *agentv1.ExecutionPermit if previous := s.executions[req.Binding.ExecutionId]; previous != nil { if previous.unknown { return &agentv1.ExecuteResponse{Receipt: s.receipt(req.Meta, agentv1.ResultCode_RESULT_CODE_UNKNOWN, agentv1.FailureCode_FAILURE_CODE_FAILED_PRECONDITION, "execution requires reconciliation", false), State: agentv1.ExecutionState_EXECUTION_STATE_UNKNOWN}, nil } if previous.controlAction == agentv1.ControlAction_CONTROL_ACTION_PAUSE || previous.controlAction == agentv1.ControlAction_CONTROL_ACTION_STOP || previous.state == agentv1.ExecutionState_EXECUTION_STATE_TERMINAL { return &agentv1.ExecuteResponse{Receipt: s.receipt(req.Meta, agentv1.ResultCode_RESULT_CODE_REJECTED, agentv1.FailureCode_FAILURE_CODE_FAILED_PRECONDITION, "execution control blocks preparation", false)}, nil } if previous.executeDigest != "" { return &agentv1.ExecuteResponse{Receipt: s.receipt(req.Meta, agentv1.ResultCode_RESULT_CODE_CONFLICT, agentv1.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 := agentv1.ExecutionState_EXECUTION_STATE_PREPARED if req.PermitId != "" { state = agentv1.ExecutionState_EXECUTION_STATE_PERMIT_GRANTED } s.executions[req.Binding.ExecutionId] = &executionRecord{executeDigest: digest, binding: proto.Clone(req.Binding).(*agentv1.ExecutionBinding), state: state, taskRevision: req.Binding.TaskRevision, callState: "prepared", phone: phoneIdentity, permit: priorPermit} receipt := s.receipt(req.Meta, agentv1.ResultCode_RESULT_CODE_ACCEPTED, agentv1.FailureCode_FAILURE_CODE_UNSPECIFIED, "", false) s.operations[s.operationKey(req.Meta)] = operationRecord{digest: digest, receipt: proto.Clone(receipt).(*agentv1.OperationReceipt)} if err := s.persistExecutionJournalLocked(); err != nil { return nil, err } return &agentv1.ExecuteResponse{Receipt: receipt, State: state}, nil } func (s *Server) GetExecutionPermit(ctx context.Context, req *agentv1.GetExecutionPermitRequest) (*agentv1.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 == agentv1.ControlAction_CONTROL_ACTION_PAUSE || execution.controlAction == agentv1.ControlAction_CONTROL_ACTION_STOP || execution.state == agentv1.ExecutionState_EXECUTION_STATE_TERMINAL) { return &agentv1.GetExecutionPermitResponse{Receipt: s.receipt(req.Meta, agentv1.ResultCode_RESULT_CODE_REJECTED, agentv1.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 &agentv1.GetExecutionPermitResponse{Receipt: s.receipt(req.Meta, agentv1.ResultCode_RESULT_CODE_CONFLICT, agentv1.FailureCode_FAILURE_CODE_ABORTED, "idempotency key content conflict", false)}, nil } var permit *agentv1.ExecutionPermit if execution != nil && execution.permit != nil { permit = proto.Clone(execution.permit).(*agentv1.ExecutionPermit) } return &agentv1.GetExecutionPermitResponse{Receipt: proto.Clone(previous.receipt).(*agentv1.OperationReceipt), Permit: permit}, nil } permitID := fmt.Sprintf("permit-%s", req.Binding.ExecutionId) fencingToken := randomToken() permit := &agentv1.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).(*agentv1.ExecutionBinding), taskRevision: req.Binding.TaskRevision, callState: "prepared"} s.executions[req.Binding.ExecutionId] = execution } if execution.permit != nil && execution.permit.ResourceReservationId != req.ResourceReservationId { return &agentv1.GetExecutionPermitResponse{Receipt: s.receipt(req.Meta, agentv1.ResultCode_RESULT_CODE_CONFLICT, agentv1.FailureCode_FAILURE_CODE_ABORTED, "execution already has a different permit", false)}, nil } execution.permit = proto.Clone(permit).(*agentv1.ExecutionPermit) execution.state = agentv1.ExecutionState_EXECUTION_STATE_PERMIT_GRANTED receipt := s.receipt(req.Meta, agentv1.ResultCode_RESULT_CODE_APPLIED, agentv1.FailureCode_FAILURE_CODE_UNSPECIFIED, "", false) s.operations[s.operationKey(req.Meta)] = operationRecord{digest: digest, receipt: proto.Clone(receipt).(*agentv1.OperationReceipt)} if err := s.persistExecutionJournalLocked(); err != nil { return nil, err } return &agentv1.GetExecutionPermitResponse{Receipt: receipt, Permit: permit}, nil } func (s *Server) ApplyTaskControl(ctx context.Context, req *agentv1.ApplyTaskControlRequest) (*agentv1.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 != agentv1.ControlAction_CONTROL_ACTION_PAUSE && req.Action != agentv1.ControlAction_CONTROL_ACTION_RESUME && req.Action != agentv1.ControlAction_CONTROL_ACTION_STOP { return nil, status.Error(codes.InvalidArgument, "a supported control action is required") } if req.ActiveCallPolicy != agentv1.ActiveCallPolicy_ACTIVE_CALL_POLICY_DRAIN && req.ActiveCallPolicy != agentv1.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).(*agentv1.ApplyTaskControlRequest) identity.Meta = &agentv1.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 &agentv1.ApplyTaskControlResponse{Receipt: s.receipt(req.Meta, agentv1.ResultCode_RESULT_CODE_CONFLICT, agentv1.FailureCode_FAILURE_CODE_ABORTED, "idempotency key content conflict", false)}, nil } return proto.Clone(previous.control).(*agentv1.ApplyTaskControlResponse), nil } save := func(response *agentv1.ApplyTaskControlResponse) (*agentv1.ApplyTaskControlResponse, error) { s.operations[key] = operationRecord{digest: digest, receipt: proto.Clone(response.Receipt).(*agentv1.OperationReceipt), control: proto.Clone(response).(*agentv1.ApplyTaskControlResponse)} if err := s.persistExecutionJournalLocked(); err != nil { return nil, err } return response, nil } execution := s.executions[req.Binding.ExecutionId] if execution == nil { return save(&agentv1.ApplyTaskControlResponse{Receipt: s.receipt(req.Meta, agentv1.ResultCode_RESULT_CODE_REJECTED, agentv1.FailureCode_FAILURE_CODE_NOT_FOUND, "execution not found", false)}) } if !proto.Equal(execution.binding, req.Binding) { return save(&agentv1.ApplyTaskControlResponse{Receipt: s.receipt(req.Meta, agentv1.ResultCode_RESULT_CODE_CONFLICT, agentv1.FailureCode_FAILURE_CODE_ABORTED, "execution control binding mismatch", false)}) } if execution.taskRevision != req.ExpectedTaskRevision { return save(&agentv1.ApplyTaskControlResponse{Receipt: s.receipt(req.Meta, agentv1.ResultCode_RESULT_CODE_CONFLICT, agentv1.FailureCode_FAILURE_CODE_ABORTED, "task revision conflict", false), AppliedTaskRevision: execution.taskRevision, State: execution.state}) } if (execution.state == agentv1.ExecutionState_EXECUTION_STATE_TERMINAL || execution.controlAction == agentv1.ControlAction_CONTROL_ACTION_STOP) && req.Action != agentv1.ControlAction_CONTROL_ACTION_STOP { return save(&agentv1.ApplyTaskControlResponse{Receipt: s.receipt(req.Meta, agentv1.ResultCode_RESULT_CODE_REJECTED, agentv1.FailureCode_FAILURE_CODE_FAILED_PRECONDITION, "stopped execution cannot resume", false), AppliedTaskRevision: execution.taskRevision, State: execution.state}) } if req.Action != agentv1.ControlAction_CONTROL_ACTION_RESUME && req.ActiveCallPolicy == agentv1.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 == agentv1.ControlAction_CONTROL_ACTION_STOP { if req.ActiveCallPolicy == agentv1.ActiveCallPolicy_ACTIVE_CALL_POLICY_DRAIN { execution.callState = "draining" } else { execution.state = agentv1.ExecutionState_EXECUTION_STATE_TERMINAL execution.callState = "stopped" } } else if req.Action == agentv1.ControlAction_CONTROL_ACTION_PAUSE { execution.callState = "paused" if req.ActiveCallPolicy == agentv1.ActiveCallPolicy_ACTIVE_CALL_POLICY_HANGUP { execution.state = agentv1.ExecutionState_EXECUTION_STATE_TERMINAL } } else { execution.callState = "resumed" } return save(&agentv1.ApplyTaskControlResponse{Receipt: s.receipt(req.Meta, agentv1.ResultCode_RESULT_CODE_APPLIED, agentv1.FailureCode_FAILURE_CODE_UNSPECIFIED, "", false), AppliedTaskRevision: execution.taskRevision, State: execution.state}) } func (s *Server) QueryExecution(ctx context.Context, req *agentv1.QueryExecutionRequest) (*agentv1.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 &agentv1.QueryExecutionResponse{Meta: s.responseMeta(req.Meta), Failure: s.failure(agentv1.FailureCode_FAILURE_CODE_NOT_FOUND, "execution not found", false)}, nil } snapshot := &agentv1.ExecutionSnapshot{Binding: proto.Clone(execution.binding).(*agentv1.ExecutionBinding), State: execution.state, CallState: execution.callState, AttemptId: execution.binding.AttemptId, ObservedAtUnixMs: s.now().UnixMilli(), Unknown: execution.unknown} s.mu.Unlock() return &agentv1.QueryExecutionResponse{Meta: s.responseMeta(req.Meta), Snapshot: snapshot}, nil } func (s *Server) ReportExecutionEvent(ctx context.Context, req *agentv1.ReportExecutionEventRequest) (*agentv1.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 &agentv1.ReportExecutionEventResponse{Receipt: s.receipt(req.Meta, agentv1.ResultCode_RESULT_CODE_CONFLICT, agentv1.FailureCode_FAILURE_CODE_ABORTED, "fact digest conflict", false)}, nil } return &agentv1.ReportExecutionEventResponse{Receipt: s.receipt(req.Meta, agentv1.ResultCode_RESULT_CODE_ACCEPTED, agentv1.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 &agentv1.ReportExecutionEventResponse{Receipt: s.receipt(req.Meta, agentv1.ResultCode_RESULT_CODE_ACCEPTED, agentv1.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 *agentv1.ExecutionFact) *executionRecord { if fact == nil || fact.Binding == nil || fact.Binding.ExecutionId == "" { return nil } return executions[fact.Binding.ExecutionId] } func callLogEvents(fact *agentv1.ExecutionFact, execution *executionRecord, meta *agentv1.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 agentv1.FactKind) string { switch kind { case agentv1.FactKind_FACT_KIND_EXECUTION_ACCEPTED: return "execution.accepted" case agentv1.FactKind_FACT_KIND_CALL_STATUS: return "call.status" case agentv1.FactKind_FACT_KIND_CALL_FINISHED: return "call.finished" case agentv1.FactKind_FACT_KIND_TRANSCRIPT_UPDATED: return "transcript.updated" case agentv1.FactKind_FACT_KIND_TRANSCRIPT_FAILED: return "transcript.failed" case agentv1.FactKind_FACT_KIND_CONTACT_OPT_OUT: return "contact.opt_out" case agentv1.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 *agentv1.RequestUploadRequest) (*agentv1.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 &agentv1.RequestUploadResponse{Receipt: s.receipt(req.Meta, agentv1.ResultCode_RESULT_CODE_REJECTED, agentv1.FailureCode_FAILURE_CODE_FAILED_PRECONDITION, "uploads are disabled", false), State: agentv1.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 &agentv1.RequestUploadResponse{Receipt: s.receipt(req.Meta, agentv1.ResultCode_RESULT_CODE_CONFLICT, agentv1.FailureCode_FAILURE_CODE_ABORTED, "upload ID is bound to a different execution or asset", false), State: agentv1.UploadState_UPLOAD_STATE_FAILED}, nil } return &agentv1.RequestUploadResponse{Receipt: s.receipt(req.Meta, agentv1.ResultCode_RESULT_CODE_ACCEPTED, agentv1.FailureCode_FAILURE_CODE_UNSPECIFIED, "duplicate upload request", false), Grant: proto.Clone(existing.grant).(*agentv1.UploadGrant), State: existing.state}, nil } grant := &agentv1.UploadGrant{UploadId: req.UploadId, TargetUrl: "https://oss.mock.invalid/upload/" + req.UploadId, ExpiresAtUnixMs: s.now().Add(5 * time.Minute).UnixMilli(), ObjectKey: req.Asset.AssetId, 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).(*agentv1.ExecutionBinding), asset: proto.Clone(req.Asset).(*agentv1.AssetDescriptor), state: agentv1.UploadState_UPLOAD_STATE_REQUESTED, grant: grant} return &agentv1.RequestUploadResponse{Receipt: s.receipt(req.Meta, agentv1.ResultCode_RESULT_CODE_ACCEPTED, agentv1.FailureCode_FAILURE_CODE_UNSPECIFIED, "", false), Grant: proto.Clone(grant).(*agentv1.UploadGrant), State: agentv1.UploadState_UPLOAD_STATE_REQUESTED}, nil } func (s *Server) CompleteUpload(ctx context.Context, req *agentv1.CompleteUploadRequest) (*agentv1.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 &agentv1.CompleteUploadResponse{Receipt: s.receipt(req.Meta, agentv1.ResultCode_RESULT_CODE_REJECTED, agentv1.FailureCode_FAILURE_CODE_NOT_FOUND, "upload not found", false), State: agentv1.UploadState_UPLOAD_STATE_FAILED}, nil } if !proto.Equal(upload.binding, req.Binding) || !proto.Equal(upload.asset, req.Asset) { return &agentv1.CompleteUploadResponse{Receipt: s.receipt(req.Meta, agentv1.ResultCode_RESULT_CODE_CONFLICT, agentv1.FailureCode_FAILURE_CODE_ABORTED, "upload completion binding does not match request", false), State: agentv1.UploadState_UPLOAD_STATE_FAILED}, nil } if upload.asset.ChecksumSha256 != req.UploadedChecksumSha256 || upload.asset.SizeBytes != req.UploadedSizeBytes || req.UploadedSizeBytes > upload.grant.MaxBytes { return &agentv1.CompleteUploadResponse{Receipt: s.receipt(req.Meta, agentv1.ResultCode_RESULT_CODE_REJECTED, agentv1.FailureCode_FAILURE_CODE_INVALID_ARGUMENT, "uploaded asset does not match grant", false), State: agentv1.UploadState_UPLOAD_STATE_FAILED}, nil } if upload.completed { return &agentv1.CompleteUploadResponse{Receipt: s.receipt(req.Meta, agentv1.ResultCode_RESULT_CODE_ACCEPTED, agentv1.FailureCode_FAILURE_CODE_UNSPECIFIED, "duplicate mock upload completion", false), State: agentv1.UploadState_UPLOAD_STATE_COMPLETED}, nil } if upload.grant.ExpiresAtUnixMs <= s.now().UnixMilli() { return &agentv1.CompleteUploadResponse{Receipt: s.receipt(req.Meta, agentv1.ResultCode_RESULT_CODE_REJECTED, agentv1.FailureCode_FAILURE_CODE_FAILED_PRECONDITION, "upload grant has expired", true), State: agentv1.UploadState_UPLOAD_STATE_FAILED}, nil } upload.completed = true upload.state = agentv1.UploadState_UPLOAD_STATE_COMPLETED s.uploads[req.UploadId] = upload return &agentv1.CompleteUploadResponse{Receipt: s.receipt(req.Meta, agentv1.ResultCode_RESULT_CODE_ACCEPTED, agentv1.FailureCode_FAILURE_CODE_UNSPECIFIED, "mock upload notification completion accepted", false), State: upload.state}, nil } func requireIdempotency(meta *agentv1.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 *agentv1.RequestMeta) error { if err := s.executionJournalReady(); err != nil { return err } if meta == nil { return status.Error(codes.InvalidArgument, "request metadata is required") } if err := s.checkPeer(ctx, meta.AgentId); err != nil { return err } return s.sessions.Authorize(meta, s.now()) } func (s *Server) checkConfiguredIdentity(agentID, cellID string) error { if s.status.AgentId != "" && s.status.AgentId != agentID { return status.Error(codes.PermissionDenied, "request Agent identity is not bound to this endpoint") } if s.status.CellId != "" && s.status.CellId != cellID { return status.Error(codes.PermissionDenied, "request Cell identity is not bound to this endpoint") } return nil } func (s *Server) checkPeer(ctx context.Context, agentID string) error { if !s.requirePeerCertificate && len(s.peerAgentIDs) == 0 && len(s.peerCertificateFingerprints) == 0 { return nil } p, ok := peer.FromContext(ctx) if !ok { return status.Error(codes.Unauthenticated, "mTLS peer is missing") } tlsInfo, ok := p.AuthInfo.(credentials.TLSInfo) if !ok || len(tlsInfo.State.VerifiedChains) == 0 || len(tlsInfo.State.VerifiedChains[0]) == 0 { return status.Error(codes.Unauthenticated, "verified mTLS peer is required") } fingerprint := CertificateFingerprint(tlsInfo.State.VerifiedChains[0][0]) if len(s.peerAgentIDs) == 0 && len(s.peerCertificateFingerprints) == 0 { // The shared Agent certificate authenticates the certificate group. The // Dispatcher-approved session binding still authorizes the individual // agent/cell/boot tuple; no self-reported identity is trusted here. return nil } if len(s.peerCertificateFingerprints) > 0 { if _, allowed := s.peerCertificateFingerprints[fingerprint]; !allowed { return status.Error(codes.PermissionDenied, "mTLS certificate is not in the endpoint allowlist") } } if len(s.peerAgentIDs) > 0 { if expected := s.peerAgentIDs[fingerprint]; expected == "" || expected != agentID { return status.Error(codes.PermissionDenied, "mTLS certificate is not bound to this agent") } } return nil } func (s *Server) peerIsAuthenticated(ctx context.Context) bool { p, ok := peer.FromContext(ctx) if !ok { return false } _, ok = p.AuthInfo.(credentials.TLSInfo) return ok } func (s *Server) responseMeta(meta *agentv1.RequestMeta) *agentv1.ResponseMeta { if meta == nil { return nil } return &agentv1.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 *agentv1.RequestMeta, result agentv1.ResultCode, code agentv1.FailureCode, detail string, retryable bool) *agentv1.OperationReceipt { receipt := &agentv1.OperationReceipt{Meta: s.responseMeta(meta), Result: result, AcceptedAtUnixMs: s.now().UnixMilli()} if code != agentv1.FailureCode_FAILURE_CODE_UNSPECIFIED { receipt.Failure = s.failure(code, detail, retryable) } return receipt } func (s *Server) failure(code agentv1.FailureCode, detail string, retryable bool) *agentv1.Failure { return &agentv1.Failure{Code: code, Detail: detail, Retryable: retryable} } func (s *Server) operationKey(meta *agentv1.RequestMeta) string { if meta == nil || meta.IdempotencyKey == "" { return "" } return meta.AgentId + "\x00" + meta.OperationId + "\x00" + meta.IdempotencyKey } func (s *Server) replayOperation(meta *agentv1.RequestMeta, digest string) (*agentv1.OperationReceipt, *agentv1.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, agentv1.ResultCode_RESULT_CODE_UNKNOWN, agentv1.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, agentv1.ResultCode_RESULT_CODE_CONFLICT, agentv1.FailureCode_FAILURE_CODE_ABORTED, "idempotency key content conflict", false) } return proto.Clone(record.receipt).(*agentv1.OperationReceipt), nil } func messageDigest(message proto.Message) string { encoded, err := proto.Marshal(message) if err != nil { return "marshal-error" } digest := sha256.Sum256(encoded) 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 *agentv1.Session) *agentv1.Session { if value == nil { return nil } return proto.Clone(value).(*agentv1.Session) } func cloneConfigReferences(values []*agentv1.ConfigReference) []*agentv1.ConfigReference { result := make([]*agentv1.ConfigReference, 0, len(values)) for _, value := range values { if value != nil { result = append(result, proto.Clone(value).(*agentv1.ConfigReference)) } } return result } func cloneStringMap(values map[string]string) map[string]string { if values == nil { return nil } result := make(map[string]string, len(values)) for key, value := range values { result[key] = value } return result } func cloneSet(values map[string]struct{}) map[string]struct{} { if values == nil { return nil } result := make(map[string]struct{}, len(values)) for key := range values { result[key] = struct{}{} } return result } var _ agentv1.AgentControlServiceServer = (*Server)(nil) var _ = grpc.SupportPackageIsVersion9