diff --git a/docs/evidence/saas-dispatcher-implementation.md b/docs/evidence/saas-dispatcher-implementation.md index 6ea5936..5f177c4 100644 --- a/docs/evidence/saas-dispatcher-implementation.md +++ b/docs/evidence/saas-dispatcher-implementation.md @@ -120,6 +120,7 @@ - Agent 旧执行日志根因:新增「首次激活→同一路径重启」测试,复现现行入口曾在初次启动自动写出废弃的 `.executions` 文件,下一次启动又将它识别为旧未交付状态并拒绝服务。移除旧日志写入、回放状态和闲置执行配置,只保留当前 `.approved` 日志、会话代际和旧文件存在时拒绝启动的保护;回归确认首次启动与重启不会产生新旧执行日志。已有 `.executions` 一律保留并失败关闭,不自动清理或猜测其业务内容;当前测试只使用临时目录。 - Agent 旧静态制品入口:现行命令从未提供 `StaticArtifactRaw/Expected`,旧激活分支及只服务于旧 `static-cell-artifact-v0.2` Schema 的手写解析器无法证明 Asterisk 实际加载。结构测试先复现残留,再移除旧 RPC 参数、解析器与专属测试;保留当前隔离 Mock 的 `LoadedSIP` revision 回报及 Dispatcher SIP 版本准入校验。此变更不等于真实 Agent/Asterisk 已加载或管理平台已审批,历史 Schema/来源事实另行辨析。 - Agent 旧 spool 准入边界:隔离测试先复现现行 Agent 对同一恢复根目录中旧 `.uploads`、`.upload-locks` 与逐执行 `state.json` 均会照常启动;现于创建会话和媒体状态前只读检查这些遗留标记,发现时明确拒绝启动并保留原文件。测试核实没有写新会话、没有修改标记;仅对配置的恢复根目录生效,不替代现存数据的人工核查或处置。 +- Agent 旧 spool 代码:旧执行状态、实时事件、上传尝试、失败事实及上传锁只由旧模块彼此调用,没有现行命令或录音恢复调用;删除专属实现与测试。录音恢复仍复用的原子写、目录同步、文件名校验等小函数移至 `internal/agent/file_state.go`,通过现行录音重试/未知结果测试确认恢复能力保持。删除源码不删除任何磁盘 spool;旧文件在当前恢复根目录触发上述拒绝启动,不能视作已经补传或处置。 ## 验收台账 diff --git a/internal/agent/events.go b/internal/agent/events.go deleted file mode 100644 index 713c925..0000000 --- a/internal/agent/events.go +++ /dev/null @@ -1,55 +0,0 @@ -package agent - -import ( - "encoding/json" - "errors" - "fmt" - "time" - - "git.ipao.vip/rogee/go-sip/internal/contract" -) - -// EventWriter keeps Agent-produced realtime facts in the approved event -// vocabulary before they are handed to Dispatcher/MQ. Transcript text is -// written to the Agent spool as an archive as well as returned for realtime -// publication; the archive is not a substitute for transcript.updated. -type EventWriter struct { - DispatcherID string - TenantID string - TenantKey string - TraceID string -} - -func (w EventWriter) TranscriptUpdated(now time.Time, eventID, callID, turnID, segmentID, role, text string, revision int64, final bool, startMS, endMS int64) ([]byte, error) { - if eventID == "" || callID == "" || turnID == "" || segmentID == "" || role == "" { - return nil, errors.New("transcript event identity is required") - } - if revision < 1 || startMS < 0 || endMS < startMS { - return nil, errors.New("transcript timing or revision is invalid") - } - return (contract.EventBuilder{ - DispatcherID: w.DispatcherID, TenantID: w.TenantID, TenantKey: w.TenantKey, TraceID: w.TraceID, - EventType: "transcript.updated", Aggregate: "transcript_segment", AggregateID: segmentID, Version: revision, - Payload: map[string]any{ - "call_id": callID, "turn_id": turnID, "segment_id": segmentID, - "role": role, "revision": revision, "text": text, "is_final": final, - "start_ms": startMS, "end_ms": endMS, "playback_state": "not_applicable", - }, - }).Marshal(now, eventID) -} - -func (s *Spool) AppendApprovedEvent(executionID string, event []byte) error { - var envelope struct { - EventType string `json:"event_type"` - } - if err := json.Unmarshal(event, &envelope); err != nil { - return fmt.Errorf("decode event envelope: %w", err) - } - if envelope.EventType != "transcript.updated" { - return fmt.Errorf("Agent transcript archive accepts transcript.updated only, got %q", envelope.EventType) - } - if err := contract.ValidateEvent(event); err != nil { - return err - } - return s.AppendTranscript(executionID, event) -} diff --git a/internal/agent/events_test.go b/internal/agent/events_test.go deleted file mode 100644 index 91ac5b5..0000000 --- a/internal/agent/events_test.go +++ /dev/null @@ -1,50 +0,0 @@ -package agent - -import ( - "strings" - "testing" - "time" - - "git.ipao.vip/rogee/go-sip/internal/contract" -) - -func TestEventWriterBuildsApprovedRealtimeTranscript(t *testing.T) { - writer := EventWriter{DispatcherID: "11111111-1111-4111-8111-111111111111", TenantID: "tenant-1", TenantKey: "tenant-demo-key", TraceID: "trace-1"} - event, err := writer.TranscriptUpdated(time.Unix(100, 0), "event-1", "call-1", "turn-1", "segment-1", "customer", "您好", 1, true, 0, 600) - if err != nil { - t.Fatal(err) - } - if err := contract.ValidateEvent(event); err != nil { - t.Fatal(err) - } - if string(event) == "" || strings.Contains(string(event), "call.transcript") { - t.Fatal("invalid transcript alias or empty event") - } - - spool, err := NewSpool(t.TempDir(), time.Now) - if err != nil { - t.Fatal(err) - } - if _, err := spool.Start("execution-1", 1, "session-1"); err != nil { - t.Fatal(err) - } - if err := spool.AppendApprovedEvent("execution-1", event); err != nil { - t.Fatal(err) - } -} - -func TestEventWriterRejectsLegacyRecordingEventAndWrongArchiveEvent(t *testing.T) { - spool, err := NewSpool(t.TempDir(), time.Now) - if err != nil { - t.Fatal(err) - } - if _, err := spool.Start("execution-1", 1, "session-1"); err != nil { - t.Fatal(err) - } - if err := spool.AppendApprovedEvent("execution-1", []byte(`{"event_type":"call.transcript"}`)); err == nil { - t.Fatal("expected invalid realtime event name to be rejected") - } - if err := spool.AppendApprovedEvent("execution-1", []byte(`{"event_type":"recording.ready"}`)); err == nil { - t.Fatal("legacy recording.ready must not enter the Agent archive") - } -} diff --git a/internal/agent/file_state.go b/internal/agent/file_state.go new file mode 100644 index 0000000..374375b --- /dev/null +++ b/internal/agent/file_state.go @@ -0,0 +1,65 @@ +package agent + +import ( + "encoding/json" + "fmt" + "os" + "path/filepath" + "strings" +) + +func writeJSONAtomic(path string, value any) error { + data, err := json.MarshalIndent(value, "", " ") + if err != nil { + return err + } + file, err := os.CreateTemp(filepath.Dir(path), ".state-*.tmp") + if err != nil { + return err + } + tmp := file.Name() + defer os.Remove(tmp) + if _, err := file.Write(append(data, '\n')); err != nil { + _ = file.Close() + return err + } + if err := file.Sync(); err != nil { + _ = file.Close() + _ = os.Remove(tmp) + return err + } + if err := file.Close(); err != nil { + _ = os.Remove(tmp) + return err + } + if err := os.Rename(tmp, path); err != nil { + _ = os.Remove(tmp) + return err + } + return syncDirectory(filepath.Dir(path)) +} + +func syncDirectory(path string) error { + dir, err := os.Open(path) + if err != nil { + return err + } + syncErr := dir.Sync() + return firstError(syncErr, dir.Close()) +} + +func validateName(name string) error { + if name == "" || name == "." || name == ".." || strings.ContainsAny(name, `/\\`) || strings.Contains(name, "..") || strings.TrimSpace(name) != name { + return fmt.Errorf("unsafe file name %q", name) + } + return nil +} + +func firstError(errs ...error) error { + for _, err := range errs { + if err != nil { + return err + } + } + return nil +} diff --git a/internal/agent/spool.go b/internal/agent/spool.go deleted file mode 100644 index f71093c..0000000 --- a/internal/agent/spool.go +++ /dev/null @@ -1,276 +0,0 @@ -// Package agent owns the Agent's file-backed execution and asset recovery -// state. It deliberately has no business database and never proxies audio to -// Dispatcher. -package agent - -import ( - "crypto/sha256" - "encoding/hex" - "encoding/json" - "errors" - "fmt" - "io" - "os" - "path/filepath" - "strings" - "sync" - "time" -) - -type State struct { - SchemaVersion string `json:"schema_version"` - ExecutionID string `json:"execution_id"` - TaskRevision int64 `json:"task_revision"` - SessionID string `json:"session_id,omitempty"` - Status string `json:"status"` - Unknown bool `json:"unknown"` - Reason string `json:"reason,omitempty"` - UpdatedAt time.Time `json:"updated_at"` -} - -type Spool struct { - root string - now func() time.Time - mu sync.Mutex -} - -func NewSpool(root string, now func() time.Time) (*Spool, error) { - if strings.TrimSpace(root) == "" { - return nil, errors.New("spool root is required") - } - if now == nil { - now = time.Now - } - if err := os.MkdirAll(root, 0o700); err != nil { - return nil, fmt.Errorf("create spool root: %w", err) - } - return &Spool{root: root, now: now}, nil -} - -func (s *Spool) Root() string { return s.root } - -func (s *Spool) Start(executionID string, revision int64, sessionID string) (State, error) { - if err := validateName(executionID); err != nil { - return State{}, err - } - if revision < 1 { - return State{}, errors.New("task revision must be positive") - } - state := State{SchemaVersion: "1", ExecutionID: executionID, TaskRevision: revision, SessionID: sessionID, Status: "reserved", UpdatedAt: s.now().UTC()} - s.mu.Lock() - defer s.mu.Unlock() - path := s.statePath(executionID) - if _, err := os.Stat(path); err == nil { - return State{}, fmt.Errorf("execution state already exists: %s", executionID) - } else if !errors.Is(err, os.ErrNotExist) { - return State{}, err - } - for _, dir := range []string{s.executionDir(executionID), filepath.Join(s.executionDir(executionID), "transcript"), filepath.Join(s.executionDir(executionID), "assets")} { - if err := os.MkdirAll(dir, 0o700); err != nil { - return State{}, err - } - } - if err := writeJSONAtomic(path, state); err != nil { - return State{}, err - } - return state, nil -} - -func (s *Spool) Load(executionID string) (State, error) { - if err := validateName(executionID); err != nil { - return State{}, err - } - data, err := os.ReadFile(s.statePath(executionID)) - if err != nil { - return State{}, err - } - var state State - if err := json.Unmarshal(data, &state); err != nil { - return State{}, fmt.Errorf("decode execution state: %w", err) - } - return state, nil -} - -func (s *Spool) Update(executionID, status, reason string) (State, error) { - if err := validateName(executionID); err != nil { - return State{}, err - } - if status == "" { - return State{}, errors.New("state status is required") - } - s.mu.Lock() - defer s.mu.Unlock() - data, err := os.ReadFile(s.statePath(executionID)) - if err != nil { - return State{}, err - } - var state State - if err := json.Unmarshal(data, &state); err != nil { - return State{}, fmt.Errorf("decode execution state: %w", err) - } - state.Status, state.Reason, state.UpdatedAt = status, reason, s.now().UTC() - state.Unknown = status == "unknown" - if err := writeJSONAtomic(s.statePath(executionID), state); err != nil { - return State{}, err - } - return state, nil -} - -// MarkUnknownOnBoot converts all in-flight local state to explicit unknown. -// It never deletes or releases a remote reservation; Dispatcher reconciliation -// must decide whether a recovered execution may proceed. -func (s *Spool) MarkUnknownOnBoot() (RecoveryReport, error) { - s.mu.Lock() - defer s.mu.Unlock() - entries, err := os.ReadDir(s.root) - if err != nil { - return RecoveryReport{}, err - } - var report RecoveryReport - for _, entry := range entries { - if !entry.IsDir() || strings.HasPrefix(entry.Name(), ".") { - continue - } - executionID := entry.Name() - path := s.statePath(executionID) - data, err := os.ReadFile(path) - if errors.Is(err, os.ErrNotExist) { - continue - } - if err != nil { - return report, err - } - var state State - if err := json.Unmarshal(data, &state); err != nil { - quarantine := path + ".corrupt-" + s.now().UTC().Format("20060102T150405.000000000Z") - if renameErr := os.Rename(path, quarantine); renameErr != nil { - return report, fmt.Errorf("quarantine corrupt state: %w (decode: %v)", renameErr, err) - } - report.Quarantined = append(report.Quarantined, executionID) - continue - } - if state.Status != "running" && state.Status != "reserved" && state.Status != "starting" && state.Status != "draining" { - continue - } - state.Status, state.Unknown, state.Reason, state.UpdatedAt = "unknown", true, "agent_boot_recovery", s.now().UTC() - if err := writeJSONAtomic(path, state); err != nil { - return report, err - } - report.Unknown = append(report.Unknown, executionID) - } - return report, nil -} - -type RecoveryReport struct { - Unknown []string - Quarantined []string -} - -func (s *Spool) AppendTranscript(executionID string, event []byte) error { - if err := validateName(executionID); err != nil { - return err - } - if len(event) == 0 { - return errors.New("transcript event is empty") - } - if !json.Valid(event) { - return errors.New("transcript event must be valid JSON") - } - path := filepath.Join(s.executionDir(executionID), "transcript", "events.jsonl") - file, err := os.OpenFile(path, os.O_WRONLY|os.O_CREATE|os.O_APPEND, 0o600) - if err != nil { - return err - } - defer file.Close() - if _, err := file.Write(append(event, '\n')); err != nil { - return err - } - return file.Sync() -} - -func (s *Spool) WriteAsset(executionID, assetID string, r io.Reader) (string, int64, string, error) { - if err := validateName(executionID); err != nil { - return "", 0, "", err - } - if err := validateName(assetID); err != nil { - return "", 0, "", err - } - if r == nil { - return "", 0, "", errors.New("asset reader is required") - } - dir := filepath.Join(s.executionDir(executionID), "assets") - if err := os.MkdirAll(dir, 0o700); err != nil { - return "", 0, "", err - } - part := filepath.Join(dir, assetID+".part") - final := filepath.Join(dir, assetID) - file, err := os.OpenFile(part, os.O_WRONLY|os.O_CREATE|os.O_TRUNC, 0o600) - if err != nil { - return "", 0, "", err - } - hash := sha256.New() - n, copyErr := io.Copy(io.MultiWriter(file, hash), r) - syncErr := file.Sync() - closeErr := file.Close() - if copyErr != nil || syncErr != nil || closeErr != nil { - _ = os.Remove(part) - return "", n, "", firstError(copyErr, syncErr, closeErr) - } - if err := os.Rename(part, final); err != nil { - _ = os.Remove(part) - return "", n, "", err - } - return final, n, hex.EncodeToString(hash.Sum(nil)), nil -} - -func (s *Spool) executionDir(executionID string) string { return filepath.Join(s.root, executionID) } -func (s *Spool) statePath(executionID string) string { - return filepath.Join(s.executionDir(executionID), "state.json") -} - -func writeJSONAtomic(path string, value any) error { - data, err := json.MarshalIndent(value, "", " ") - if err != nil { - return err - } - file, err := os.CreateTemp(filepath.Dir(path), ".state-*.tmp") - if err != nil { - return err - } - tmp := file.Name() - defer os.Remove(tmp) - if _, err := file.Write(append(data, '\n')); err != nil { - _ = file.Close() - return err - } - if err := file.Sync(); err != nil { - _ = file.Close() - _ = os.Remove(tmp) - return err - } - if err := file.Close(); err != nil { - _ = os.Remove(tmp) - return err - } - if err := os.Rename(tmp, path); err != nil { - _ = os.Remove(tmp) - return err - } - return syncDirectory(filepath.Dir(path)) -} - -func validateName(name string) error { - if name == "" || name == "." || name == ".." || strings.ContainsAny(name, `/\\`) || strings.Contains(name, "..") || strings.TrimSpace(name) != name { - return fmt.Errorf("unsafe file name %q", name) - } - return nil -} - -func firstError(errs ...error) error { - for _, err := range errs { - if err != nil { - return err - } - } - return nil -} diff --git a/internal/agent/spool_test.go b/internal/agent/spool_test.go deleted file mode 100644 index 7e9b489..0000000 --- a/internal/agent/spool_test.go +++ /dev/null @@ -1,84 +0,0 @@ -package agent - -import ( - "bytes" - "os" - "path/filepath" - "testing" - "time" -) - -func testSpool(t *testing.T) *Spool { - t.Helper() - now := time.Date(2026, 9, 18, 0, 0, 0, 0, time.UTC) - s, err := NewSpool(t.TempDir(), func() time.Time { return now }) - if err != nil { - t.Fatal(err) - } - return s -} - -func TestSpoolAtomicStateAndBootUnknown(t *testing.T) { - s := testSpool(t) - if _, err := s.Start("exec-1", 1, "session-1"); err != nil { - t.Fatal(err) - } - if _, err := s.Update("exec-1", "running", ""); err != nil { - t.Fatal(err) - } - report, err := s.MarkUnknownOnBoot() - if err != nil { - t.Fatal(err) - } - if len(report.Unknown) != 1 || report.Unknown[0] != "exec-1" { - t.Fatalf("recovery report = %+v", report) - } - state, err := s.Load("exec-1") - if err != nil { - t.Fatal(err) - } - if state.Status != "unknown" || !state.Unknown { - t.Fatalf("state = %+v", state) - } -} - -func TestSpoolQuarantinesCorruptStateAndNeverDeletesIt(t *testing.T) { - s := testSpool(t) - if err := os.MkdirAll(filepath.Join(s.Root(), "broken"), 0o700); err != nil { - t.Fatal(err) - } - if err := os.WriteFile(filepath.Join(s.Root(), "broken", "state.json"), []byte("{"), 0o600); err != nil { - t.Fatal(err) - } - report, err := s.MarkUnknownOnBoot() - if err != nil { - t.Fatal(err) - } - if len(report.Quarantined) != 1 || len(report.Unknown) != 0 { - t.Fatalf("recovery report = %+v", report) - } - matches, err := filepath.Glob(filepath.Join(s.Root(), "broken", "state.json.corrupt-*")) - if err != nil || len(matches) != 1 { - t.Fatalf("quarantine files = %v, err=%v", matches, err) - } -} - -func TestSpoolTranscriptAndAssetAreDurableFiles(t *testing.T) { - s := testSpool(t) - if _, err := s.Start("exec-2", 1, "session-2"); err != nil { - t.Fatal(err) - } - if err := s.AppendTranscript("exec-2", []byte(`{"text":"hello","final":true}`)); err != nil { - t.Fatal(err) - } - path, n, hash, err := s.WriteAsset("exec-2", "recording.pcm", bytes.NewReader([]byte("pcm"))) - if err != nil { - t.Fatal(err) - } - if n != 3 || hash == "" || path == "" { - t.Fatalf("asset result path=%q bytes=%d hash=%q", path, n, hash) - } - if _, err := os.Stat(filepath.Join(s.Root(), "exec-2", "assets", "recording.pcm.part")); !os.IsNotExist(err) { - t.Fatalf("temporary asset still exists: %v", err) - } -} diff --git a/internal/agent/upload_failure.go b/internal/agent/upload_failure.go deleted file mode 100644 index 53dfbb3..0000000 --- a/internal/agent/upload_failure.go +++ /dev/null @@ -1,109 +0,0 @@ -package agent - -import ( - "crypto/sha256" - "encoding/hex" - "errors" - "fmt" - "os" - "path/filepath" - - agentpb "git.ipao.vip/rogee/go-sip/gen/agent" - "git.ipao.vip/rogee/go-sip/internal/contract" - "github.com/google/uuid" - "google.golang.org/protobuf/proto" -) - -func validateUploadFailureFact(record UploadAttempt, fact *agentpb.ExecutionFact) error { - if record.State != "attempted" || record.Binding == nil || record.Asset == nil || record.Result.SizeBytes != 0 || - fact == nil || fact.Binding == nil || fact.Kind != agentpb.FactKind_FACT_KIND_RECORDING_PROGRESS || - !proto.Equal(record.Binding, fact.Binding) || fact.SourceBootId == "" || fact.SourceSequence == 0 || fact.ObservedAtUnixMs <= 0 { - return errors.New("Mock failure requires a matching uncompleted upload attempt") - } - parsed, err := uuid.Parse(fact.FactId) - if err != nil || parsed.Version() != 4 { - return errors.New("Mock failure fact ID must be UUID v4") - } - failure, err := contract.DecodeLocalMockRecordingFailure(fact.PayloadJson) - if err != nil { - return err - } - if failure.UploadID != record.UploadID || failure.RecordingID != record.Asset.AssetId { - return errors.New("Mock failure payload does not match claimed upload") - } - sum := sha256.Sum256(fact.PayloadJson) - if hex.EncodeToString(sum[:]) != fact.ContentSha256 { - return errors.New("Mock failure payload checksum mismatch") - } - return nil -} - -// RecordUploadFailureFact persists a single terminal failure identity before -// any network report. Unknown PUT results never enter this path and cannot be -// replayed as a second upload after restart. -func (s *Spool) RecordUploadFailureFact(id string, fact *agentpb.ExecutionFact) error { - s.mu.Lock() - defer s.mu.Unlock() - record, err := s.LoadUploadAttempt(id) - if err != nil { - return err - } - if err := validateUploadFailureFact(record, fact); err != nil { - return err - } - if record.FailureFact != nil { - if !proto.Equal(record.FailureFact, fact) { - return errors.New("Mock upload failure fact identity changed") - } - return nil - } - record.FailureFact = proto.Clone(fact).(*agentpb.ExecutionFact) - return writeJSONAtomic(filepath.Join(s.uploadAttemptDir(id), "state.json"), record) -} - -func (s *Spool) CompleteUploadFailureReport(id, factID string) error { - s.mu.Lock() - defer s.mu.Unlock() - record, err := s.LoadUploadAttempt(id) - if err != nil { - return err - } - if record.FailureFact == nil || record.FailureFact.FactId != factID || record.State != "attempted" { - return errors.New("failure report acknowledgement does not match upload") - } - if record.FailureDelivered { - return nil - } - record.FailureDelivered = true - return writeJSONAtomic(filepath.Join(s.uploadAttemptDir(id), "state.json"), record) -} - -// PendingUploadFailures retries only the fact notification. It never requests -// another token, reads the source recording or repeats an uncertain PUT. -func (s *Spool) PendingUploadFailures() ([]UploadAttempt, error) { - entries, err := os.ReadDir(filepath.Join(s.root, ".uploads")) - if errors.Is(err, os.ErrNotExist) { - return nil, nil - } - if err != nil { - return nil, err - } - var pending []UploadAttempt - for _, entry := range entries { - if !entry.IsDir() { - return nil, fmt.Errorf("unexpected upload journal entry %q", entry.Name()) - } - record, err := s.LoadUploadAttempt(entry.Name()) - if err != nil { - return nil, err - } - if record.FailureFact == nil || record.FailureDelivered { - continue - } - if err := validateUploadFailureFact(record, record.FailureFact); err != nil { - return nil, fmt.Errorf("invalid persisted Mock failure for %s: %w", record.UploadID, err) - } - pending = append(pending, record) - } - return pending, nil -} diff --git a/internal/agent/upload_failure_state_test.go b/internal/agent/upload_failure_state_test.go deleted file mode 100644 index 2285d12..0000000 --- a/internal/agent/upload_failure_state_test.go +++ /dev/null @@ -1,88 +0,0 @@ -package agent - -import ( - "crypto/sha256" - "encoding/hex" - "errors" - "os" - "path/filepath" - "testing" - "time" - - agentpb "git.ipao.vip/rogee/go-sip/gen/agent" - "github.com/google/uuid" - "google.golang.org/protobuf/proto" -) - -func TestMockRecordingFailureFactSurvivesRestartWithoutAnotherPUT(t *testing.T) { - root := t.TempDir() - spool, err := NewSpool(root, nil) - if err != nil { - t.Fatal(err) - } - binding := &agentpb.ExecutionBinding{TenantId: "tenant-a", TenantKey: "tenant-a", ExecutionId: "execution-a"} - asset := &agentpb.AssetDescriptor{AssetId: "recording-a"} - requestID := uuid.NewString() - if err := spool.ClaimUpload(UploadAttempt{ - RequestID: requestID, UploadID: "upload-a", Identity: "identity-a", State: "attempted", - Binding: binding, Asset: asset, - }); err != nil { - t.Fatal(err) - } - payload := []byte(`{"schema_version":"local-mock-recording-failure.v0.1","upload_id":"upload-a","recording_id":"recording-a","error_code":"upload_failed"}`) - sum := sha256.Sum256(payload) - fact := &agentpb.ExecutionFact{ - FactId: uuid.NewString(), Binding: binding, Kind: agentpb.FactKind_FACT_KIND_RECORDING_PROGRESS, - PayloadJson: payload, ContentSha256: hex.EncodeToString(sum[:]), ObservedAtUnixMs: time.Now().UnixMilli(), - SourceBootId: uuid.NewString(), SourceSequence: 1, - } - if err := spool.RecordUploadFailureFact("upload-a", fact); err != nil { - t.Fatal(err) - } - restarted, err := NewSpool(root, nil) - if err != nil { - t.Fatal(err) - } - pending, err := restarted.PendingUploadFailures() - if err != nil || len(pending) != 1 || pending[0].FailureFact.GetFactId() != fact.FactId { - t.Fatalf("failure fact lost across restart: pending=%+v err=%v", pending, err) - } - if err := restarted.RecordUploadFailureFact("upload-a", fact); err != nil { - t.Fatalf("same fact must be idempotent: %v", err) - } - conflict := proto.Clone(fact).(*agentpb.ExecutionFact) - conflict.FactId = uuid.NewString() - if err := restarted.RecordUploadFailureFact("upload-a", conflict); err == nil { - t.Fatal("a second failure identity replaced the original") - } - if err := restarted.ReserveUploadRetry("upload-a", uuid.NewString()); err == nil { - t.Fatal("final failure still permitted an explicit new PUT") - } - if err := restarted.RecordUploadResult("upload-a", UploadResult{SizeBytes: 4, SHA256: "checksum", StatusCode: 200}); err == nil { - t.Fatal("failed upload was replaced by a success") - } - if err := restarted.CompleteUploadFailureReport("upload-a", fact.FactId); err != nil { - t.Fatal(err) - } - pending, err = restarted.PendingUploadFailures() - if err != nil || len(pending) != 0 { - t.Fatalf("acknowledged fact was redelivered: pending=%d err=%v", len(pending), err) - } - if err := restarted.CompleteUploadFailureReport("upload-a", fact.FactId); err != nil { - t.Fatalf("same acknowledgement must be idempotent: %v", err) - } - if err := restarted.CompleteUploadFailureReport("upload-a", uuid.NewString()); err == nil { - t.Fatal("unrelated acknowledgement closed failure fact") - } - persisted, err := restarted.LoadUploadAttempt("upload-a") - if err != nil || persisted.State != "attempted" || !persisted.FailureDelivered { - t.Fatalf("failed PUT attempt regressed: %+v err=%v", persisted, err) - } - data, err := os.ReadFile(filepath.Join(root, ".uploads", "upload-a", "state.json")) - if err != nil { - t.Fatal(err) - } - if len(data) == 0 || errors.Is(err, os.ErrNotExist) { - t.Fatal("failure fact was not durably journaled") - } -} diff --git a/internal/agent/upload_lock.go b/internal/agent/upload_lock.go deleted file mode 100644 index 75d4b52..0000000 --- a/internal/agent/upload_lock.go +++ /dev/null @@ -1,47 +0,0 @@ -package agent - -import ( - "errors" - "os" - "path/filepath" - - "golang.org/x/sys/unix" -) - -var ErrUploadBusy = errors.New("upload is already being handled by another process") - -type UploadLock struct{ file *os.File } - -// LockUpload uses an OS lock released on process death. The stable lock inode -// is never removed, so independent processes cannot lock different inodes. -func (s *Spool) LockUpload(id string) (*UploadLock, error) { - if err := validateName(id); err != nil { - return nil, err - } - root := filepath.Join(s.root, ".upload-locks") - if err := os.MkdirAll(root, 0700); err != nil { - return nil, err - } - file, err := os.OpenFile(filepath.Join(root, id), os.O_CREATE|os.O_RDWR, 0600) - if err != nil { - return nil, err - } - if err := unix.Flock(int(file.Fd()), unix.LOCK_EX|unix.LOCK_NB); err != nil { - closeErr := file.Close() - if errors.Is(err, unix.EWOULDBLOCK) { - return nil, errors.Join(ErrUploadBusy, closeErr) - } - return nil, errors.Join(err, closeErr) - } - return &UploadLock{file: file}, nil -} - -func (l *UploadLock) Close() error { - if l == nil || l.file == nil { - return nil - } - unlockErr := unix.Flock(int(l.file.Fd()), unix.LOCK_UN) - closeErr := l.file.Close() - l.file = nil - return errors.Join(unlockErr, closeErr) -} diff --git a/internal/agent/upload_lock_test.go b/internal/agent/upload_lock_test.go deleted file mode 100644 index 5338bba..0000000 --- a/internal/agent/upload_lock_test.go +++ /dev/null @@ -1,35 +0,0 @@ -package agent - -import ( - "errors" - "testing" -) - -func TestUploadLockSerializesIndependentSpools(t *testing.T) { - root := t.TempDir() - first, err := NewSpool(root, nil) - if err != nil { - t.Fatal(err) - } - second, err := NewSpool(root, nil) - if err != nil { - t.Fatal(err) - } - lock, err := first.LockUpload("upload-a") - if err != nil { - t.Fatal(err) - } - if _, err := second.LockUpload("upload-a"); !errors.Is(err, ErrUploadBusy) { - t.Fatalf("parallel writer accepted: %v", err) - } - if err := lock.Close(); err != nil { - t.Fatal(err) - } - next, err := second.LockUpload("upload-a") - if err != nil { - t.Fatal(err) - } - if err := next.Close(); err != nil { - t.Fatal(err) - } -} diff --git a/internal/agent/upload_retry_test.go b/internal/agent/upload_retry_test.go deleted file mode 100644 index dc6b847..0000000 --- a/internal/agent/upload_retry_test.go +++ /dev/null @@ -1,30 +0,0 @@ -package agent - -import "testing" - -func TestExplicitRetryReservesEachRequestOnlyOnce(t *testing.T) { - spool, err := NewSpool(t.TempDir(), nil) - if err != nil { - t.Fatal(err) - } - record := UploadAttempt{UploadID: "upload-a", RequestID: "first-request", Identity: "identity-a", State: "attempted"} - if err := spool.ClaimUpload(record); err != nil { - t.Fatal(err) - } - for _, request := range []string{"second-request", "third-request"} { - if err := spool.ReserveUploadRetry(record.UploadID, request); err != nil { - t.Fatal(err) - } - } - for _, request := range []string{"first-request", "second-request", "third-request"} { - if err := spool.ReserveUploadRetry(record.UploadID, request); err == nil { - t.Fatalf("request %s reused", request) - } - } - if err := spool.RecordUploadResult(record.UploadID, UploadResult{SizeBytes: 4, SHA256: "checksum", StatusCode: 200}); err != nil { - t.Fatal(err) - } - if err := spool.ReserveUploadRetry(record.UploadID, "fourth-request"); err == nil { - t.Fatal("successful PUT was authorized again") - } -} diff --git a/internal/agent/upload_state.go b/internal/agent/upload_state.go deleted file mode 100644 index b776f92..0000000 --- a/internal/agent/upload_state.go +++ /dev/null @@ -1,197 +0,0 @@ -package agent - -import ( - "encoding/json" - "errors" - "fmt" - "os" - "path/filepath" - - agentpb "git.ipao.vip/rogee/go-sip/gen/agent" -) - -// UploadAttempt records no signed URLs or credentials. An attempted PUT whose -// result is unknown must never be repeated by restart recovery. -type UploadAttempt struct { - RequestID string `json:"request_id"` - UploadID string `json:"upload_id"` - Identity string `json:"identity"` - State string `json:"state"` - ObjectKey string `json:"object_key"` - Result UploadResult `json:"result"` - Binding *agentpb.ExecutionBinding `json:"binding"` - Asset *agentpb.AssetDescriptor `json:"asset"` - FailureFact *agentpb.ExecutionFact `json:"failure_fact,omitempty"` - FailureDelivered bool `json:"failure_delivered,omitempty"` -} - -func (s *Spool) uploadAttemptDir(id string) string { return filepath.Join(s.root, ".uploads", id) } - -func (s *Spool) ClaimUpload(record UploadAttempt) error { - if err := validateName(record.UploadID); err != nil { - return err - } - if record.RequestID != "" { - if err := validateName(record.RequestID); err != nil { - return err - } - } - if record.Identity == "" || record.State != "attempted" || record.FailureFact != nil || record.FailureDelivered { - return errors.New("upload attempt must start without a failure fact") - } - root := filepath.Join(s.root, ".uploads") - if err := os.MkdirAll(root, 0700); err != nil { - return err - } - // Exclusive directory creation arbitrates across processes, not just goroutines. - dir := s.uploadAttemptDir(record.UploadID) - if err := os.Mkdir(dir, 0700); err != nil { - return err - } - if err := syncDirectory(s.root); err != nil { - return err - } - if err := syncDirectory(root); err != nil { - return err - } - if record.RequestID != "" { - requests := filepath.Join(dir, "requests") - if err := os.Mkdir(requests, 0700); err != nil { - return err - } - if err := os.Mkdir(filepath.Join(requests, record.RequestID), 0700); err != nil { - return err - } - if err := syncDirectory(requests); err != nil { - return err - } - } - return writeJSONAtomic(filepath.Join(dir, "state.json"), record) -} - -// ReserveUploadRetry is used only for an explicit new request, while holding -// LockUpload. A consumed request identity is never made reusable after a crash. -func (s *Spool) ReserveUploadRetry(id, requestID string) error { - if err := validateName(requestID); err != nil { - return err - } - record, err := s.LoadUploadAttempt(id) - if err != nil { - return err - } - if record.State != "attempted" || record.FailureFact != nil || record.RequestID == "" || requestID == record.RequestID { - return errors.New("only an unsuccessful, non-terminal attempt can use an explicit new request") - } - requests := filepath.Join(s.uploadAttemptDir(id), "requests") - if err := os.Mkdir(filepath.Join(requests, requestID), 0700); err != nil { - return err - } - if err := syncDirectory(requests); err != nil { - return err - } - record.RequestID = requestID - record.Result = UploadResult{} - return writeJSONAtomic(filepath.Join(s.uploadAttemptDir(id), "state.json"), record) -} - -func (s *Spool) LoadUploadAttempt(id string) (UploadAttempt, error) { - if err := validateName(id); err != nil { - return UploadAttempt{}, err - } - dir := s.uploadAttemptDir(id) - data, err := os.ReadFile(filepath.Join(dir, "state.json")) - if errors.Is(err, os.ErrNotExist) { - if _, statErr := os.Stat(dir); statErr == nil { - return UploadAttempt{}, errors.New("upload attempt exists without durable state; PUT outcome is unknown") - } - } - if err != nil { - return UploadAttempt{}, err - } - var record UploadAttempt - if err := json.Unmarshal(data, &record); err != nil { - return record, err - } - if record.UploadID != id || record.Identity == "" { - return record, errors.New("invalid persisted upload identity") - } - switch record.State { - case "attempted", "uploaded", "completed": - default: - return record, fmt.Errorf("invalid persisted upload state %q", record.State) - } - if record.FailureDelivered && record.FailureFact == nil || record.FailureFact != nil && record.State != "attempted" { - return record, errors.New("persisted upload failure contradicts attempt state") - } - return record, nil -} - -func (s *Spool) RecordUploadResult(id string, result UploadResult) error { - s.mu.Lock() - defer s.mu.Unlock() - record, err := s.LoadUploadAttempt(id) - if err != nil { - return err - } - if record.State != "attempted" || record.FailureFact != nil { - return errors.New("only an attempted upload without a terminal failure may record its PUT result") - } - if result.SizeBytes <= 0 || result.SHA256 == "" || result.StatusCode < 200 || result.StatusCode >= 300 { - return errors.New("successful upload result is required") - } - record.State, record.Result = "uploaded", result - return writeJSONAtomic(filepath.Join(s.uploadAttemptDir(id), "state.json"), record) -} - -func (s *Spool) CompleteUploadNotification(id string) error { - s.mu.Lock() - defer s.mu.Unlock() - record, err := s.LoadUploadAttempt(id) - if err != nil { - return err - } - if record.State != "uploaded" && record.State != "completed" { - return errors.New("upload result must precede notification completion") - } - record.State = "completed" - return writeJSONAtomic(filepath.Join(s.uploadAttemptDir(id), "state.json"), record) -} - -// PendingUploadNotifications enumerates only successful PUTs. Unknown attempts -// remain reserved and are never returned as work to retry. -func (s *Spool) PendingUploadNotifications() ([]UploadAttempt, error) { - entries, err := os.ReadDir(filepath.Join(s.root, ".uploads")) - if errors.Is(err, os.ErrNotExist) { - return nil, nil - } - if err != nil { - return nil, err - } - var pending []UploadAttempt - for _, entry := range entries { - if !entry.IsDir() { - return nil, fmt.Errorf("unexpected upload journal entry %q", entry.Name()) - } - record, err := s.LoadUploadAttempt(entry.Name()) - if err != nil { - return nil, err - } - if record.State != "uploaded" { - continue - } - if record.Binding == nil || record.Asset == nil { - return nil, fmt.Errorf("upload %q lacks notification metadata", record.UploadID) - } - pending = append(pending, record) - } - return pending, nil -} - -func syncDirectory(path string) error { - dir, err := os.Open(path) - if err != nil { - return err - } - syncErr := dir.Sync() - return firstError(syncErr, dir.Close()) -} diff --git a/internal/agent/upload_state_test.go b/internal/agent/upload_state_test.go deleted file mode 100644 index 502acc8..0000000 --- a/internal/agent/upload_state_test.go +++ /dev/null @@ -1,50 +0,0 @@ -package agent - -import ( - "errors" - "os" - "testing" -) - -func TestUploadAttemptSurvivesRestartWithoutAnotherPUT(t *testing.T) { - root := t.TempDir() - spool, err := NewSpool(root, nil) - if err != nil { - t.Fatal(err) - } - record := UploadAttempt{UploadID: "upload-a", Identity: "identity-a", State: "attempted"} - if err := spool.ClaimUpload(record); err != nil { - t.Fatal(err) - } - restarted, err := NewSpool(root, nil) - if err != nil { - t.Fatal(err) - } - if err := restarted.ClaimUpload(record); !errors.Is(err, os.ErrExist) { - t.Fatalf("duplicate PUT was not prevented: %v", err) - } - recovered, err := restarted.LoadUploadAttempt("upload-a") - if err != nil { - t.Fatal(err) - } - if recovered.State != "attempted" { - t.Fatal("unknown PUT attempt lost") - } - result := UploadResult{SizeBytes: 10, SHA256: "checksum", StatusCode: 200} - if err := restarted.RecordUploadResult("upload-a", result); err != nil { - t.Fatal(err) - } - recovered, err = spool.LoadUploadAttempt("upload-a") - if err != nil { - t.Fatal(err) - } - if recovered.State != "uploaded" || recovered.Result != result { - t.Fatal("notification recovery lost original upload result") - } - if err := spool.CompleteUploadNotification("upload-a"); err != nil { - t.Fatal(err) - } - if err := spool.RecordUploadResult("upload-a", result); err == nil { - t.Fatal("completed upload regressed") - } -}