fix(rpc): stop creating retired execution journal on Agent startup
This commit is contained in:
@@ -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` 一律保留并失败关闭,不自动清理或猜测其业务内容;当前测试只使用临时目录。
|
||||
|
||||
## 验收台账
|
||||
|
||||
|
||||
@@ -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{}
|
||||
|
||||
@@ -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
|
||||
}
|
||||
+5
-177
@@ -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
|
||||
|
||||
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user