refactor(agent): retire uncalled execution and upload spool

This commit is contained in:
2026-09-30 13:53:06 +08:00
parent 25bf73cc51
commit dc46a92fe4
13 changed files with 66 additions and 1021 deletions
@@ -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;旧文件在当前恢复根目录触发上述拒绝启动,不能视作已经补传或处置。
## 验收台账
-55
View File
@@ -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)
}
-50
View File
@@ -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")
}
}
+65
View File
@@ -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
}
-276
View File
@@ -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
}
-84
View File
@@ -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)
}
}
-109
View File
@@ -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
}
@@ -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")
}
}
-47
View File
@@ -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)
}
-35
View File
@@ -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)
}
}
-30
View File
@@ -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")
}
}
-197
View File
@@ -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())
}
-50
View File
@@ -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")
}
}