From 8d7c3ab03f9a68c05e57c7fed3818e2cd9fd93e1 Mon Sep 17 00:00:00 2001 From: Rogee Date: Tue, 29 Sep 2026 23:18:49 +0800 Subject: [PATCH] Persist failed OSS recordings for bounded recovery --- .../saas-dispatcher-implementation.md | 3 +- internal/agent/recording_recovery.go | 206 ++++++++++++++ internal/agent/recording_recovery_test.go | 157 +++++++++++ internal/agent/recording_retry.go | 153 ++++++++++ internal/agent/recording_retry_test.go | 261 ++++++++++++++++++ internal/agent/recording_scan.go | 105 +++++++ internal/agent/recording_scan_test.go | 171 ++++++++++++ 7 files changed, 1055 insertions(+), 1 deletion(-) create mode 100644 internal/agent/recording_recovery.go create mode 100644 internal/agent/recording_recovery_test.go create mode 100644 internal/agent/recording_retry.go create mode 100644 internal/agent/recording_retry_test.go create mode 100644 internal/agent/recording_scan.go create mode 100644 internal/agent/recording_scan_test.go diff --git a/docs/evidence/saas-dispatcher-implementation.md b/docs/evidence/saas-dispatcher-implementation.md index e8bdbef..dda66f2 100644 --- a/docs/evidence/saas-dispatcher-implementation.md +++ b/docs/evidence/saas-dispatcher-implementation.md @@ -59,7 +59,8 @@ ## P06:录音/OSS/最终结果(进行中,未签收) - 已新增 Agent 内存录音单次 PUT 组件 `UploadClient.UploadBytes`:复用受限授权、大小和 SHA-256 核对,成功路径不写临时录音或通话结果文件。OSS 明确拒绝保留状态码且不自动重试;传输结果不明时返回专用错误、停止自动重试并隐藏带签名的 URL。空授权头拒绝而非使服务崩溃。这里只验证组件,尚未连接实际录音或最终结果。 -- 已验证:`go test ./internal/agent -count=1`、`go test -race ./internal/agent -count=1`、`go vet ./...`、`go build ./...`、`git diff --check`。失败双文件持久化、48 小时恢复、录音生成失败/无录音、Dispatcher 唯一结果 outbox、MQ/重启恢复和端到端回归均未完成,不能宣称 P06 通过。 +- Agent 失败恢复隔离组件:仅在 OSS PUT 明确失败后,把原 bucket/object_key 对应的录音和通话信息两文件写入私有目录并同步落盘;双文件缺失或损坏明确报错、不伪造结果。完成保存后固定 48 小时窗口,按 1 分钟递增至最长 1 小时重试;到期保留原文件。PUT 前持久写入 in-flight,结果不明或进程重启不会二次 PUT;确认上传后先持久记录成功,再经注入的 Mock 回报通话结果,回报失败/重启仅重发原结果。启动扫描识别遗漏文件并提供不暴露原路径的稳定摘要;真实 Dispatcher 授权及回报尚未接线。 +- 已验证:`go test ./... -count=1`、`go test -race ./internal/agent -count=1`、`go vet ./...`、`go build ./...`、`PATH=/tmp/sip-go-agent-tools/bin:$PATH bash scripts/check-current-contracts.sh`、`PATH=/tmp/sip-go-agent-tools/bin:$PATH bash scripts/check-proto.sh`、`git diff --check`。尚未完成实际录音到直传/失败恢复的接线、无录音与生成失败、Dispatcher 唯一最终结果 outbox、真实 Agent↔Dispatcher 结果交付及 MQ/端到端验收,不能宣称 P06 通过。 ## 验收台账 diff --git a/internal/agent/recording_recovery.go b/internal/agent/recording_recovery.go new file mode 100644 index 0000000..43d6b8d --- /dev/null +++ b/internal/agent/recording_recovery.go @@ -0,0 +1,206 @@ +package agent + +import ( + "bytes" + "crypto/sha256" + "encoding/hex" + "encoding/json" + "errors" + "fmt" + "io" + "os" + "path/filepath" + "strings" + "sync" + "time" +) + +const recordingRetryWindow = 48 * time.Hour + +var ErrRecordingRecoveryIncomplete = errors.New("recording recovery files are incomplete or inconsistent") + +// RecordingRecoveryEntry contains the original OSS target and the exact final +// result bytes. Temporary grants, credentials and signed URLs are never saved. +type RecordingRecoveryEntry struct { + CallID string `json:"call_id"` + SourceEventID string `json:"source_event_id"` + RecordingID string `json:"recording_id"` + UploadID string `json:"upload_id"` + Bucket string `json:"bucket"` + ObjectKey string `json:"object_key"` + SHA256 string `json:"sha256"` + SizeBytes int64 `json:"size_bytes"` + ResultPayload []byte `json:"result_payload"` + Uploaded *UploadResult `json:"uploaded,omitempty"` + SavedAt time.Time `json:"saved_at"` + NextAttemptAt time.Time `json:"next_attempt_at,omitempty"` + ExpiresAt time.Time `json:"expires_at"` + Attempts int `json:"attempts"` + State string `json:"state"` +} + +// RecordingRecovery is only used after a failed OSS PUT. A successful normal +// upload never creates either of its recovery files. +type RecordingRecovery struct { + Root string + Now func() time.Time + Upload UploadClient + + mu sync.Mutex + writeState func(string, any) error // deterministic failure injection in package tests +} + +func (r *RecordingRecovery) SaveFailure(audio []byte, entry RecordingRecoveryEntry, cause error) (RecordingRecoveryEntry, error) { + r.mu.Lock() + defer r.mu.Unlock() + var rejected *UploadHTTPError + if !errors.Is(cause, ErrUploadOutcomeUnknown) && (!errors.As(cause, &rejected) || rejected.StatusCode < 300) { + return RecordingRecoveryEntry{}, errors.New("only a failed OSS PUT may create recording recovery files") + } + if len(audio) == 0 || entry.CallID == "" || entry.SourceEventID == "" || entry.RecordingID == "" || entry.UploadID == "" || len(entry.ResultPayload) == 0 || !json.Valid(entry.ResultPayload) { + return RecordingRecoveryEntry{}, errors.New("complete failed recording identity, audio and result are required") + } + if entry.State != "" || entry.Attempts != 0 || !entry.SavedAt.IsZero() || !entry.NextAttemptAt.IsZero() || !entry.ExpiresAt.IsZero() { + return RecordingRecoveryEntry{}, errors.New("new failed recording already has recovery state") + } + dir, audioPath, infoPath, err := r.paths(entry.Bucket, entry.ObjectKey) + if err != nil { + return RecordingRecoveryEntry{}, err + } + sha := sha256.Sum256(audio) + digest := hex.EncodeToString(sha[:]) + if entry.SHA256 != "" && !strings.EqualFold(entry.SHA256, digest) { + return RecordingRecoveryEntry{}, ErrUploadChecksumMismatch + } + if entry.SizeBytes != 0 && entry.SizeBytes != int64(len(audio)) { + return RecordingRecoveryEntry{}, errors.New("failed recording size does not match the original grant") + } + if err := r.makePrivateDirs(entry.Bucket, entry.ObjectKey); err != nil { + return RecordingRecoveryEntry{}, err + } + file, err := os.OpenFile(audioPath, os.O_WRONLY|os.O_CREATE|os.O_EXCL, 0600) + if err != nil { + return RecordingRecoveryEntry{}, fmt.Errorf("save failed recording: %w", err) + } + written, writeErr := file.Write(audio) + if writeErr == nil && written != len(audio) { + writeErr = io.ErrShortWrite + } + syncErr := file.Sync() + closeErr := file.Close() + if err := firstError(writeErr, syncErr, closeErr); err != nil { + return RecordingRecoveryEntry{}, fmt.Errorf("persist failed recording audio: %w", err) + } + if err := syncDirectory(dir); err != nil { + return RecordingRecoveryEntry{}, fmt.Errorf("persist failed recording path: %w", err) + } + entry.SHA256, entry.SizeBytes = digest, int64(len(audio)) + entry.SavedAt = r.now().UTC() + entry.ExpiresAt = entry.SavedAt.Add(recordingRetryWindow) + entry.Attempts = 1 + if errors.Is(cause, ErrUploadOutcomeUnknown) { + entry.State = "outcome_unknown" + } else { + entry.State = "retry_pending" + entry.NextAttemptAt = entry.SavedAt.Add(time.Minute) + } + if err := r.persistState(infoPath, entry); err != nil { + return RecordingRecoveryEntry{}, fmt.Errorf("persist failed recording info: %w", err) + } + return entry, nil +} + +func (r *RecordingRecovery) Load(bucket, objectKey string) (RecordingRecoveryEntry, error) { + _, audioPath, infoPath, err := r.paths(bucket, objectKey) + if err != nil { + return RecordingRecoveryEntry{}, err + } + data, err := os.ReadFile(infoPath) + if errors.Is(err, os.ErrNotExist) { + if _, audioErr := os.Stat(audioPath); audioErr == nil { + return RecordingRecoveryEntry{}, ErrRecordingRecoveryIncomplete + } + } + if err != nil { + return RecordingRecoveryEntry{}, err + } + var entry RecordingRecoveryEntry + decoder := json.NewDecoder(bytes.NewReader(data)) + decoder.DisallowUnknownFields() + if err := decoder.Decode(&entry); err != nil { + return RecordingRecoveryEntry{}, fmt.Errorf("%w: decode info: %v", ErrRecordingRecoveryIncomplete, err) + } + var trailing any + if err := decoder.Decode(&trailing); !errors.Is(err, io.EOF) { + return RecordingRecoveryEntry{}, fmt.Errorf("%w: unexpected data after info object (%v)", ErrRecordingRecoveryIncomplete, err) + } + if entry.Bucket != bucket || entry.ObjectKey != objectKey || entry.CallID == "" || entry.SourceEventID == "" || entry.RecordingID == "" || entry.UploadID == "" || entry.SizeBytes <= 0 || entry.SHA256 == "" || !json.Valid(entry.ResultPayload) || entry.Attempts < 1 || entry.SavedAt.IsZero() || !entry.ExpiresAt.Equal(entry.SavedAt.Add(recordingRetryWindow)) { + return RecordingRecoveryEntry{}, ErrRecordingRecoveryIncomplete + } + switch entry.State { + case "retry_pending", "outcome_unknown", "put_in_flight", "uploaded_unreported", "delivered", "expired": + default: + return RecordingRecoveryEntry{}, ErrRecordingRecoveryIncomplete + } + if entry.State == "uploaded_unreported" || entry.State == "delivered" { + if entry.Uploaded == nil || entry.Uploaded.StatusCode < 200 || entry.Uploaded.StatusCode >= 300 || entry.Uploaded.SizeBytes != entry.SizeBytes || !strings.EqualFold(entry.Uploaded.SHA256, entry.SHA256) { + return RecordingRecoveryEntry{}, ErrRecordingRecoveryIncomplete + } + } else if entry.Uploaded != nil { + return RecordingRecoveryEntry{}, ErrRecordingRecoveryIncomplete + } + audio, err := os.ReadFile(audioPath) + if err != nil { + return RecordingRecoveryEntry{}, fmt.Errorf("%w: read original audio: %v", ErrRecordingRecoveryIncomplete, err) + } + sha := sha256.Sum256(audio) + if int64(len(audio)) != entry.SizeBytes || !strings.EqualFold(hex.EncodeToString(sha[:]), entry.SHA256) { + return RecordingRecoveryEntry{}, ErrRecordingRecoveryIncomplete + } + return entry, nil +} + +func (r *RecordingRecovery) paths(bucket, objectKey string) (dir, audioPath, infoPath string, err error) { + if strings.TrimSpace(r.Root) == "" || validateName(bucket) != nil || objectKey == "" || strings.ContainsAny(objectKey, "\\\x00") { + return "", "", "", errors.New("recovery root, bucket or object key is invalid") + } + for _, segment := range strings.Split(objectKey, "/") { + if segment == "" || segment == "." || segment == ".." { + return "", "", "", errors.New("recovery object key cannot escape its bucket") + } + } + dir = filepath.Join(r.Root, bucket, filepath.FromSlash(objectKey)) + return dir, filepath.Join(dir, "recording.wav"), filepath.Join(dir, "info.json"), nil +} + +func (r *RecordingRecovery) makePrivateDirs(bucket, objectKey string) error { + if err := os.MkdirAll(r.Root, 0700); err != nil { + return err + } + current := r.Root + if err := checkRecoveryDirectory(current); err != nil { + return fmt.Errorf("recovery root: %w", err) + } + for _, segment := range append([]string{bucket}, strings.Split(objectKey, "/")...) { + child := filepath.Join(current, segment) + if err := os.Mkdir(child, 0700); err == nil { + if err := syncDirectory(current); err != nil { + return err + } + } else if !errors.Is(err, os.ErrExist) { + return err + } + if err := checkRecoveryDirectory(child); err != nil { + return fmt.Errorf("recovery target directory: %w", err) + } + current = child + } + return nil +} + +func (r *RecordingRecovery) now() time.Time { + if r.Now != nil { + return r.Now() + } + return time.Now() +} diff --git a/internal/agent/recording_recovery_test.go b/internal/agent/recording_recovery_test.go new file mode 100644 index 0000000..a62bd49 --- /dev/null +++ b/internal/agent/recording_recovery_test.go @@ -0,0 +1,157 @@ +package agent + +import ( + "bytes" + "crypto/sha256" + "encoding/hex" + "encoding/json" + "errors" + "os" + "path/filepath" + "reflect" + "strings" + "testing" + "time" +) + +func privateRecoveryRoot(t *testing.T) string { + t.Helper() + root := filepath.Join(t.TempDir(), "private-recording-recovery") + if err := os.Mkdir(root, 0700); err != nil { + t.Fatal(err) + } + if err := os.Chmod(root, 0700); err != nil { + t.Fatal(err) + } + return root +} + +func failedRecordingFixture() RecordingRecoveryEntry { + return RecordingRecoveryEntry{ + CallID: "call-1", SourceEventID: "execute-1", RecordingID: "rec-1", UploadID: "upload-1", + Bucket: "mock-bucket", ObjectKey: "tenant/call/rec.wav", + ResultPayload: json.RawMessage(`{"call_id":"call-1","transcript":[]}`), + } +} + +func TestFailedRecordingPersistsOriginalTargetAndTwoPrivateFiles(t *testing.T) { + root := privateRecoveryRoot(t) + now := time.Date(2026, 10, 1, 10, 0, 0, 0, time.UTC) + recovery := RecordingRecovery{Root: root, Now: func() time.Time { return now }} + audio := []byte("RIFF original failed upload") + entry, err := recovery.SaveFailure(audio, failedRecordingFixture(), &UploadHTTPError{StatusCode: 503}) + if err != nil { + t.Fatal(err) + } + sha := sha256.Sum256(audio) + if entry.State != "retry_pending" || entry.Attempts != 1 || entry.SavedAt != now || entry.NextAttemptAt != now.Add(time.Minute) || entry.ExpiresAt != now.Add(48*time.Hour) || entry.SHA256 != hex.EncodeToString(sha[:]) || entry.SizeBytes != int64(len(audio)) { + t.Fatalf("retry identity or 48-hour clock was lost: %+v", entry) + } + dir := filepath.Join(root, "mock-bucket", "tenant", "call", "rec.wav") + files, err := os.ReadDir(dir) + if err != nil || len(files) != 2 || files[0].Name() != "info.json" || files[1].Name() != "recording.wav" { + t.Fatalf("failed upload must preserve exactly the audio and call info: files=%v err=%v", files, err) + } + for _, name := range []string{"info.json", "recording.wav"} { + file, err := os.Stat(filepath.Join(dir, name)) + if err != nil { + t.Fatalf("recovery file %q is missing: %v", name, err) + } + if file.Mode().Perm() != 0600 { + t.Fatalf("recovery file %q is not restricted: mode=%v", name, file.Mode()) + } + } + stored, err := os.ReadFile(filepath.Join(dir, "recording.wav")) + if err != nil || !bytes.Equal(stored, audio) { + t.Fatalf("original recording was lost or changed: %v", err) + } + metadata, err := os.ReadFile(filepath.Join(dir, "info.json")) + if err != nil || strings.Contains(string(metadata), "signature=") || strings.Contains(string(metadata), "mock-token") { + t.Fatalf("recovery info includes grant credentials or cannot be read: %v", err) + } + loaded, err := recovery.Load(entry.Bucket, entry.ObjectKey) + if err != nil || !reflect.DeepEqual(loaded, entry) { + t.Fatalf("persistent identity did not survive a fresh read: loaded=%+v err=%v", loaded, err) + } +} + +func TestFailedRecordingRequiresActualOSSFailure(t *testing.T) { + root := privateRecoveryRoot(t) + recovery := RecordingRecovery{Root: root} + _, err := recovery.SaveFailure([]byte("recording"), failedRecordingFixture(), nil) + if err == nil { + t.Fatal("normal upload path must not create recovery files") + } + files, err := os.ReadDir(root) + if err != nil || len(files) != 0 { + t.Fatalf("no OSS failure still created files: entries=%v err=%v", files, err) + } +} + +func TestFailedRecordingMetadataWriteFailureDoesNotFabricateResult(t *testing.T) { + root := privateRecoveryRoot(t) + recovery := RecordingRecovery{Root: root, writeState: func(string, any) error { return errors.New("disk full") }} + fixture := failedRecordingFixture() + _, err := recovery.SaveFailure([]byte("original recording"), fixture, &UploadHTTPError{StatusCode: 503}) + if err == nil || !strings.Contains(err.Error(), "disk full") { + t.Fatalf("failed metadata persistence was hidden: %v", err) + } + dir := filepath.Join(root, fixture.Bucket, filepath.FromSlash(fixture.ObjectKey)) + if _, err := os.Stat(filepath.Join(dir, "recording.wav")); err != nil { + t.Fatalf("original audio must remain for manual repair: %v", err) + } + if _, err := os.Stat(filepath.Join(dir, "info.json")); !errors.Is(err, os.ErrNotExist) { + t.Fatalf("invalid metadata appeared durable: %v", err) + } + restarted := RecordingRecovery{Root: root} + if _, err := restarted.Load(fixture.Bucket, fixture.ObjectKey); err == nil { + t.Fatal("orphan recording must not be treated as a recoverable final result") + } +} + +func TestFailedRecordingAmbiguousPUTRequiresManualRecovery(t *testing.T) { + recovery := RecordingRecovery{Root: privateRecoveryRoot(t)} + entry, err := recovery.SaveFailure([]byte("original recording"), failedRecordingFixture(), ErrUploadOutcomeUnknown) + if err != nil || entry.State != "outcome_unknown" || !entry.NextAttemptAt.IsZero() { + t.Fatalf("uncertain PUT must never become an automatic retry: entry=%+v err=%v", entry, err) + } +} + +func TestFailedRecordingRejectsPathTraversalBeforeWriting(t *testing.T) { + root := privateRecoveryRoot(t) + fixture := failedRecordingFixture() + fixture.ObjectKey = "../outside.wav" + recovery := RecordingRecovery{Root: root} + _, err := recovery.SaveFailure([]byte("recording"), fixture, &UploadHTTPError{StatusCode: 503}) + if err == nil { + t.Fatal("object key escaped the recovery root") + } + files, err := os.ReadDir(root) + if err != nil || len(files) != 0 { + t.Fatalf("invalid object key wrote files: entries=%v err=%v", files, err) + } +} + +func TestFailedRecordingRejectsUnknownRecoveryStateFields(t *testing.T) { + root := privateRecoveryRoot(t) + recovery := RecordingRecovery{Root: root} + entry, err := recovery.SaveFailure([]byte("RIFF original audio"), failedRecordingFixture(), &UploadHTTPError{StatusCode: 503}) + if err != nil { + t.Fatal(err) + } + infoPath := filepath.Join(root, entry.Bucket, filepath.FromSlash(entry.ObjectKey), "info.json") + original, err := os.ReadFile(infoPath) + if err != nil { + t.Fatal(err) + } + changed := bytes.Replace(original, []byte(`"state"`), []byte(`"unrecognized_recovery_state":true,"state"`), 1) + if bytes.Equal(original, changed) || !json.Valid(changed) { + t.Fatal("test fixture did not preserve valid JSON while adding an unknown field") + } + if err := os.WriteFile(infoPath, changed, 0600); err != nil { + t.Fatal(err) + } + if _, err := recovery.Load(entry.Bucket, entry.ObjectKey); !errors.Is(err, ErrRecordingRecoveryIncomplete) { + t.Fatalf("unknown state fields were silently accepted: %v", err) + } +} diff --git a/internal/agent/recording_retry.go b/internal/agent/recording_retry.go new file mode 100644 index 0000000..dbd1d48 --- /dev/null +++ b/internal/agent/recording_retry.go @@ -0,0 +1,153 @@ +package agent + +import ( + "context" + "errors" + "fmt" + "strings" + "time" + + agentpb "git.ipao.vip/rogee/go-sip/gen/agent" +) + +var ErrRecordingNotDue = errors.New("recording recovery retry is not due") +var ErrRecordingRetryExpired = errors.New("recording recovery retry window expired") + +// RecoveryTarget must come from the Dispatcher for the original upload. The +// fresh grant may change its token, but never its bucket, key, ID or checksum. +type RecoveryTarget struct { + Bucket string + Grant *agentpb.UploadGrant +} + +type RecoveryGrant func(context.Context, RecordingRecoveryEntry) (RecoveryTarget, error) +type RecordingResultReporter func(context.Context, RecordingRecoveryEntry) error + +// Retry runs one due recovery entry. The persisted in-flight barrier prevents +// a process restart from sending another PUT when its previous outcome is +// unknown. Once a PUT succeeds, only the original result is reported again. +func (r *RecordingRecovery) Retry(ctx context.Context, bucket, objectKey string, grant RecoveryGrant, report RecordingResultReporter) (RecordingRecoveryEntry, error) { + r.mu.Lock() + defer r.mu.Unlock() + entry, err := r.Load(bucket, objectKey) + if err != nil { + return RecordingRecoveryEntry{}, err + } + _, audioPath, infoPath, err := r.paths(bucket, objectKey) + if err != nil { + return entry, err + } + if err := ctx.Err(); err != nil { + return entry, err + } + switch entry.State { + case "delivered": + return entry, nil + case "expired": + return entry, ErrRecordingRetryExpired + case "outcome_unknown", "put_in_flight": + return entry, ErrUploadOutcomeUnknown + case "uploaded_unreported": + return r.reportUploaded(ctx, infoPath, entry, report) + case "retry_pending": + default: + return entry, ErrRecordingRecoveryIncomplete + } + now := r.now().UTC() + if !now.Before(entry.ExpiresAt) { + entry.State = "expired" + entry.NextAttemptAt = time.Time{} + if err := r.persistState(infoPath, entry); err != nil { + return entry, fmt.Errorf("persist expired recording recovery: %w", err) + } + return entry, ErrRecordingRetryExpired + } + if now.Before(entry.NextAttemptAt) { + return entry, ErrRecordingNotDue + } + if grant == nil || report == nil { + return entry, errors.New("Dispatcher recovery grant and final result reporter are required") + } + target, err := grant(ctx, entry) + if err != nil { + if ctxErr := ctx.Err(); ctxErr != nil { + return entry, ctxErr + } + entry.Attempts++ + entry.NextAttemptAt = r.now().UTC().Add(recordingRetryDelay(entry.Attempts)) + if saveErr := r.persistState(infoPath, entry); saveErr != nil { + return entry, fmt.Errorf("persist failed recovery grant backoff: %w", saveErr) + } + return entry, fmt.Errorf("request original recovery grant failed (%T)", err) + } + if err := r.validateTarget(ctx, entry, target); err != nil { + return entry, err + } + inFlight := entry + inFlight.State = "put_in_flight" + inFlight.NextAttemptAt = time.Time{} + if err := r.persistState(infoPath, inFlight); err != nil { + return entry, fmt.Errorf("persist PUT in-flight barrier: %w", err) + } + uploaded, err := r.Upload.UploadFile(ctx, target.Grant, audioPath) + if err != nil { + var rejected *UploadHTTPError + if errors.As(err, &rejected) { + entry.Attempts++ + entry.NextAttemptAt = r.now().UTC().Add(recordingRetryDelay(entry.Attempts)) + if saveErr := r.persistState(infoPath, entry); saveErr != nil { + return inFlight, fmt.Errorf("OSS rejected retry (status %d), recovery barrier retained: %w", rejected.StatusCode, saveErr) + } + return entry, err + } + return inFlight, err + } + inFlight.State = "uploaded_unreported" + inFlight.Uploaded = &uploaded + if err := r.persistState(infoPath, inFlight); err != nil { + return inFlight, fmt.Errorf("persist confirmed OSS upload before reporting result: %w", err) + } + return r.reportUploaded(ctx, infoPath, inFlight, report) +} + +func (r *RecordingRecovery) validateTarget(ctx context.Context, entry RecordingRecoveryEntry, target RecoveryTarget) error { + if target.Grant == nil || target.Bucket != entry.Bucket || target.Grant.UploadId != entry.UploadID || target.Grant.ObjectKey != entry.ObjectKey || target.Grant.MaxBytes < entry.SizeBytes || target.Grant.MaxBytes <= 0 || !strings.EqualFold(target.Grant.RequiredChecksumSha256, entry.SHA256) { + return errors.New("recovery grant target does not match the original OSS asset") + } + if _, err := r.Upload.validateGrant(ctx, target.Grant); err != nil { + return fmt.Errorf("recovery grant is invalid: %w", err) + } + return nil +} + +func (r *RecordingRecovery) reportUploaded(ctx context.Context, infoPath string, entry RecordingRecoveryEntry, report RecordingResultReporter) (RecordingRecoveryEntry, error) { + if report == nil { + return entry, errors.New("Dispatcher final result reporter is required") + } + if err := report(ctx, entry); err != nil { + return entry, fmt.Errorf("report confirmed upload result: %w", err) + } + entry.State = "delivered" + if err := r.persistState(infoPath, entry); err != nil { + return entry, fmt.Errorf("persist delivered result receipt: %w", err) + } + return entry, nil +} + +func (r *RecordingRecovery) persistState(infoPath string, entry RecordingRecoveryEntry) error { + if r.writeState != nil { + return r.writeState(infoPath, entry) + } + return writeJSONAtomic(infoPath, entry) +} + +func recordingRetryDelay(attempts int) time.Duration { + delay := time.Minute + for i := 1; i < attempts && delay < time.Hour; i++ { + delay *= 2 + if delay > time.Hour { + return time.Hour + } + } + return delay +} diff --git a/internal/agent/recording_retry_test.go b/internal/agent/recording_retry_test.go new file mode 100644 index 0000000..79ec2e3 --- /dev/null +++ b/internal/agent/recording_retry_test.go @@ -0,0 +1,261 @@ +package agent + +import ( + "context" + "errors" + "io" + "net/http" + "net/http/httptest" + "os" + "path/filepath" + "strings" + "sync/atomic" + "testing" + "time" + + agentpb "git.ipao.vip/rogee/go-sip/gen/agent" +) + +func recoveryTarget(now time.Time, entry RecordingRecoveryEntry, signedURL string) RecoveryTarget { + return RecoveryTarget{Bucket: entry.Bucket, Grant: &agentpb.UploadGrant{ + UploadId: entry.UploadID, ObjectKey: entry.ObjectKey, TargetUrl: signedURL, + ExpiresAtUnixMs: now.Add(15 * time.Minute).UnixMilli(), MaxBytes: entry.SizeBytes, + RequiredChecksumSha256: entry.SHA256, + }} +} + +func TestRecordingRecoveryBackoffAndResultRetryDoNotRepeatSuccessfulPUT(t *testing.T) { + root := privateRecoveryRoot(t) + now := time.Date(2026, 10, 1, 10, 0, 0, 0, time.UTC) + recovery := RecordingRecovery{Root: root, Now: func() time.Time { return now }, Upload: UploadClient{AllowInsecureHTTP: true, Now: func() time.Time { return now }}} + entry, err := recovery.SaveFailure([]byte("RIFF original audio"), failedRecordingFixture(), &UploadHTTPError{StatusCode: 503}) + if err != nil { + t.Fatal(err) + } + var puts atomic.Int32 + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodPut || r.URL.Path != "/original" { + http.Error(w, "changed target", http.StatusBadRequest) + return + } + _, _ = io.Copy(io.Discard, r.Body) + if puts.Add(1) == 1 { + w.WriteHeader(http.StatusServiceUnavailable) + } + })) + defer server.Close() + var grants atomic.Int32 + grant := func(_ context.Context, e RecordingRecoveryEntry) (RecoveryTarget, error) { + grants.Add(1) + return recoveryTarget(now, e, server.URL+"/original?signature=DO_NOT_LOG"), nil + } + var reports atomic.Int32 + report := func(_ context.Context, e RecordingRecoveryEntry) error { + if !strings.EqualFold(e.SHA256, entry.SHA256) || e.SourceEventID != entry.SourceEventID || string(e.ResultPayload) != string(entry.ResultPayload) || e.Uploaded == nil { + t.Errorf("final result lost its original identity or confirmed upload: %+v", e) + } + if reports.Add(1) == 1 { + return errors.New("Dispatcher result receipt unavailable") + } + return nil + } + if _, err := recovery.Retry(context.Background(), entry.Bucket, entry.ObjectKey, grant, report); !errors.Is(err, ErrRecordingNotDue) || puts.Load() != 0 { + t.Fatalf("first retry must wait at least one minute: puts=%d err=%v", puts.Load(), err) + } + now = now.Add(time.Minute) + if _, err := recovery.Retry(context.Background(), entry.Bucket, entry.ObjectKey, grant, report); err == nil || puts.Load() != 1 { + t.Fatalf("first due retry should preserve the original failed target: puts=%d err=%v", puts.Load(), err) + } + pending, err := recovery.Load(entry.Bucket, entry.ObjectKey) + if err != nil || pending.State != "retry_pending" || pending.Attempts != 2 || pending.NextAttemptAt != now.Add(2*time.Minute) { + t.Fatalf("failure did not persist its next bounded retry: entry=%+v err=%v", pending, err) + } + now = now.Add(2 * time.Minute) + if _, err := recovery.Retry(context.Background(), entry.Bucket, entry.ObjectKey, grant, report); err == nil || puts.Load() != 2 || reports.Load() != 1 { + t.Fatalf("successful retry must not be called successful if reporting failed: puts=%d reports=%d err=%v", puts.Load(), reports.Load(), err) + } + uploaded, err := recovery.Load(entry.Bucket, entry.ObjectKey) + if err != nil || uploaded.State != "uploaded_unreported" || uploaded.Uploaded == nil || uploaded.Uploaded.StatusCode != 200 { + t.Fatalf("successful PUT was not durably marked before reporting: entry=%+v err=%v", uploaded, err) + } + now = now.Add(time.Minute) + restarted := RecordingRecovery{Root: root, Now: func() time.Time { return now }, Upload: recovery.Upload} + neverGrant := func(context.Context, RecordingRecoveryEntry) (RecoveryTarget, error) { + t.Fatal("confirmed successful OSS upload requested another PUT grant") + return RecoveryTarget{}, errors.New("duplicate grant") + } + processed, err := restarted.RecoverDue(context.Background(), neverGrant, report) + finished, loadErr := restarted.Load(entry.Bucket, entry.ObjectKey) + if err != nil || loadErr != nil || processed != 1 || finished.State != "delivered" || puts.Load() != 2 || grants.Load() != 2 || reports.Load() != 2 { + t.Fatalf("restart must resend only original result: recovered=%d state=%s puts=%d grants=%d reports=%d err=%v load=%v", processed, finished.State, puts.Load(), grants.Load(), reports.Load(), err, loadErr) + } + if _, err := os.Stat(filepath.Join(root, "mock-bucket", "tenant", "call", "rec.wav", "recording.wav")); err != nil { + t.Fatalf("failed-path original audio must not be cleared automatically: %v", err) + } +} + +func TestRecordingRecoveryExpiresAfterExactly48HoursWithoutResult(t *testing.T) { + now := time.Date(2026, 10, 1, 10, 0, 0, 0, time.UTC) + recovery := RecordingRecovery{Root: privateRecoveryRoot(t), Now: func() time.Time { return now }} + entry, err := recovery.SaveFailure([]byte("RIFF original audio"), failedRecordingFixture(), &UploadHTTPError{StatusCode: 503}) + if err != nil { + t.Fatal(err) + } + now = now.Add(48 * time.Hour) + grantCalls, reportCalls := 0, 0 + _, err = recovery.Retry(context.Background(), entry.Bucket, entry.ObjectKey, + func(context.Context, RecordingRecoveryEntry) (RecoveryTarget, error) { + grantCalls++ + return RecoveryTarget{}, nil + }, + func(context.Context, RecordingRecoveryEntry) error { reportCalls++; return nil }) + if !errors.Is(err, ErrRecordingRetryExpired) || grantCalls != 0 || reportCalls != 0 { + t.Fatalf("expired recording was uploaded or given a fabricated result: grant=%d report=%d err=%v", grantCalls, reportCalls, err) + } + retained, err := recovery.Load(entry.Bucket, entry.ObjectKey) + if err != nil || retained.State != "expired" { + t.Fatalf("expired files must remain for manual handling: state=%s err=%v", retained.State, err) + } +} + +func TestRecordingRecoveryUnknownPUTNeverRetriesAfterRestart(t *testing.T) { + now := time.Date(2026, 10, 1, 10, 0, 0, 0, time.UTC) + root := privateRecoveryRoot(t) + recovery := RecordingRecovery{Root: root, Now: func() time.Time { return now }} + entry, err := recovery.SaveFailure([]byte("RIFF original audio"), failedRecordingFixture(), ErrUploadOutcomeUnknown) + if err != nil { + t.Fatal(err) + } + now = now.Add(time.Hour) + grants := 0 + restarted := RecordingRecovery{Root: root, Now: func() time.Time { return now }} + _, err = restarted.Retry(context.Background(), entry.Bucket, entry.ObjectKey, + func(context.Context, RecordingRecoveryEntry) (RecoveryTarget, error) { + grants++ + return RecoveryTarget{}, nil + }, + func(context.Context, RecordingRecoveryEntry) error { + t.Fatal("unknown PUT fabricated a result") + return nil + }) + if !errors.Is(err, ErrUploadOutcomeUnknown) || grants != 0 { + t.Fatalf("unknown PUT was retried: grants=%d err=%v", grants, err) + } +} + +func TestRecordingRecoveryRejectsChangedBucketBeforePUT(t *testing.T) { + now := time.Date(2026, 10, 1, 10, 0, 0, 0, time.UTC) + recovery := RecordingRecovery{Root: privateRecoveryRoot(t), Now: func() time.Time { return now }} + entry, err := recovery.SaveFailure([]byte("RIFF original audio"), failedRecordingFixture(), &UploadHTTPError{StatusCode: 503}) + if err != nil { + t.Fatal(err) + } + now = now.Add(time.Minute) + _, err = recovery.Retry(context.Background(), entry.Bucket, entry.ObjectKey, + func(context.Context, RecordingRecoveryEntry) (RecoveryTarget, error) { + return RecoveryTarget{Bucket: "changed-bucket"}, nil + }, + func(context.Context, RecordingRecoveryEntry) error { + t.Fatal("changed target fabricated a result") + return nil + }) + if err == nil || !strings.Contains(err.Error(), "target") { + t.Fatalf("changed OSS bucket was silently used: %v", err) + } + pending, err := recovery.Load(entry.Bucket, entry.ObjectKey) + if err != nil || pending.State != "retry_pending" { + t.Fatalf("original target was overwritten: state=%s err=%v", pending.State, err) + } +} + +func TestRecordingRecoveryPersistsInFlightBeforeUnknownPUT(t *testing.T) { + now := time.Date(2026, 10, 1, 10, 0, 0, 0, time.UTC) + root := privateRecoveryRoot(t) + var puts atomic.Int32 + recovery := RecordingRecovery{Root: root, Now: func() time.Time { return now }, Upload: UploadClient{Now: func() time.Time { return now }, HTTPClient: &http.Client{Transport: ambiguousBytesTransport{requests: &puts}}}} + entry, err := recovery.SaveFailure([]byte("RIFF original audio"), failedRecordingFixture(), &UploadHTTPError{StatusCode: 503}) + if err != nil { + t.Fatal(err) + } + now = now.Add(time.Minute) + grant := func(context.Context, RecordingRecoveryEntry) (RecoveryTarget, error) { + return recoveryTarget(now, entry, "https://oss.example.invalid/original?signature=DO_NOT_LOG"), nil + } + _, err = recovery.Retry(context.Background(), entry.Bucket, entry.ObjectKey, grant, func(context.Context, RecordingRecoveryEntry) error { + t.Fatal("unknown PUT fabricated a result") + return nil + }) + if !errors.Is(err, ErrUploadOutcomeUnknown) || puts.Load() != 1 { + t.Fatalf("transport outcome was hidden: puts=%d err=%v", puts.Load(), err) + } + loaded, err := recovery.Load(entry.Bucket, entry.ObjectKey) + if err != nil || loaded.State != "put_in_flight" { + t.Fatalf("PUT started without a durable in-flight barrier: state=%s err=%v", loaded.State, err) + } + now = now.Add(time.Minute) + restarted := RecordingRecovery{Root: root, Now: func() time.Time { return now }} + _, err = restarted.Retry(context.Background(), entry.Bucket, entry.ObjectKey, grant, func(context.Context, RecordingRecoveryEntry) error { t.Fatal("repeated unknown PUT"); return nil }) + if !errors.Is(err, ErrUploadOutcomeUnknown) || puts.Load() != 1 { + t.Fatalf("restart repeated an uncertain OSS PUT: puts=%d err=%v", puts.Load(), err) + } +} + +func TestRecordingRetryBackoffStopsAtOneHour(t *testing.T) { + for _, tc := range []struct { + attempts int + want time.Duration + }{{1, time.Minute}, {2, 2 * time.Minute}, {3, 4 * time.Minute}, {6, 32 * time.Minute}, {7, time.Hour}, {30, time.Hour}} { + if got := recordingRetryDelay(tc.attempts); got != tc.want { + t.Fatalf("attempt %d has retry delay %s, want %s", tc.attempts, got, tc.want) + } + } +} + +func TestRecordingRecoveryCannotRepeatSuccessfulPUTAfterMetadataDiskFault(t *testing.T) { + now := time.Date(2026, 10, 1, 10, 0, 0, 0, time.UTC) + root := privateRecoveryRoot(t) + var puts, reports atomic.Int32 + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + _, _ = io.Copy(io.Discard, r.Body) + puts.Add(1) + })) + defer server.Close() + recovery := RecordingRecovery{Root: root, Now: func() time.Time { return now }, Upload: UploadClient{AllowInsecureHTTP: true, Now: func() time.Time { return now }}} + entry, err := recovery.SaveFailure([]byte("RIFF original audio"), failedRecordingFixture(), &UploadHTTPError{StatusCode: 503}) + if err != nil { + t.Fatal(err) + } + recovery.writeState = func(path string, value any) error { + state, ok := value.(RecordingRecoveryEntry) + if ok && state.State == "uploaded_unreported" { + return errors.New("disk full after confirmed PUT") + } + return writeJSONAtomic(path, value) + } + now = now.Add(time.Minute) + _, err = recovery.Retry(context.Background(), entry.Bucket, entry.ObjectKey, + func(_ context.Context, e RecordingRecoveryEntry) (RecoveryTarget, error) { + return recoveryTarget(now, e, server.URL), nil + }, + func(context.Context, RecordingRecoveryEntry) error { reports.Add(1); return nil }) + if err == nil || !strings.Contains(err.Error(), "persist confirmed OSS upload") || puts.Load() != 1 || reports.Load() != 0 { + t.Fatalf("unrecorded PUT success fabricated a result or triggered another PUT: puts=%d reports=%d err=%v", puts.Load(), reports.Load(), err) + } + persisted, err := recovery.Load(entry.Bucket, entry.ObjectKey) + if err != nil || persisted.State != "put_in_flight" { + t.Fatalf("crash barrier was not retained: state=%s err=%v", persisted.State, err) + } + restarted := RecordingRecovery{Root: root, Now: func() time.Time { return now }} + _, err = restarted.Retry(context.Background(), entry.Bucket, entry.ObjectKey, + func(context.Context, RecordingRecoveryEntry) (RecoveryTarget, error) { + t.Fatal("duplicate grant after successful PUT") + return RecoveryTarget{}, nil + }, + func(context.Context, RecordingRecoveryEntry) error { + t.Fatal("unconfirmed upload fabricated a result") + return nil + }) + if !errors.Is(err, ErrUploadOutcomeUnknown) || puts.Load() != 1 { + t.Fatalf("restart retried a confirmed but unjournaled upload: puts=%d err=%v", puts.Load(), err) + } +} diff --git a/internal/agent/recording_scan.go b/internal/agent/recording_scan.go new file mode 100644 index 0000000..3e76b19 --- /dev/null +++ b/internal/agent/recording_scan.go @@ -0,0 +1,105 @@ +package agent + +import ( + "context" + "crypto/sha256" + "encoding/hex" + "errors" + "fmt" + "io/fs" + "os" + "path/filepath" + "strings" +) + +// RecoverDue discovers failure-only recording pairs after an Agent restart. +// Blocked/corrupt entries remain untouched and are returned as errors while +// other independent entries continue; none can produce a fabricated result. +func (r *RecordingRecovery) RecoverDue(ctx context.Context, grant RecoveryGrant, report RecordingResultReporter) (int, error) { + if strings.TrimSpace(r.Root) == "" { + return 0, errors.New("recording recovery root is required") + } + if err := ctx.Err(); err != nil { + return 0, err + } + if err := os.MkdirAll(r.Root, 0700); err != nil { + return 0, fmt.Errorf("prepare recording recovery root: %w", err) + } + if err := checkRecoveryDirectory(r.Root); err != nil { + return 0, err + } + processed := 0 + var blocked []error + addBlocked := func(dir string, cause error) { + rel, err := filepath.Rel(r.Root, dir) + if err != nil { + blocked = append(blocked, fmt.Errorf("recording recovery path unavailable: %w", cause)) + return + } + sum := sha256.Sum256([]byte(filepath.ToSlash(rel))) + blocked = append(blocked, fmt.Errorf("recording recovery %s: %w", hex.EncodeToString(sum[:6]), cause)) + } + walkErr := filepath.WalkDir(r.Root, func(path string, item fs.DirEntry, err error) error { + if err != nil { + return err + } + if err := ctx.Err(); err != nil { + return err + } + if item.IsDir() { + return checkRecoveryDirectory(path) + } + dir := filepath.Dir(path) + switch item.Name() { + case "recording.wav": + if _, err := os.Stat(filepath.Join(dir, "info.json")); err != nil { + addBlocked(dir, fmt.Errorf("%w: original audio has no call information", ErrRecordingRecoveryIncomplete)) + } + return nil + case "info.json": + default: + addBlocked(dir, fmt.Errorf("%w: unexpected recovery file", ErrRecordingRecoveryIncomplete)) + return nil + } + rel, err := filepath.Rel(r.Root, dir) + if err != nil { + addBlocked(dir, err) + return nil + } + parts := strings.Split(filepath.ToSlash(rel), "/") + if len(parts) < 2 { + addBlocked(dir, ErrRecordingRecoveryIncomplete) + return nil + } + bucket, objectKey := parts[0], strings.Join(parts[1:], "/") + before, err := r.Load(bucket, objectKey) + if err != nil { + addBlocked(dir, err) + return nil + } + after, err := r.Retry(ctx, bucket, objectKey, grant, report) + if errors.Is(err, ErrRecordingNotDue) { + return nil + } + if err != nil { + addBlocked(dir, err) + return nil + } + if before.State != "delivered" && after.State == "delivered" { + processed++ + } + return nil + }) + return processed, errors.Join(append(blocked, walkErr)...) +} + +func checkRecoveryDirectory(path string) error { + info, err := os.Stat(path) + if err != nil { + return err + } + if !info.IsDir() || info.Mode().Perm()&0077 != 0 { + return fmt.Errorf("recording recovery directory must be private (mode=%#o dir=%t)", info.Mode().Perm(), info.IsDir()) + } + return nil +} diff --git a/internal/agent/recording_scan_test.go b/internal/agent/recording_scan_test.go new file mode 100644 index 0000000..2cdb1a8 --- /dev/null +++ b/internal/agent/recording_scan_test.go @@ -0,0 +1,171 @@ +package agent + +import ( + "context" + "crypto/sha256" + "encoding/hex" + "errors" + "io" + "net/http" + "net/http/httptest" + "os" + "path/filepath" + "strings" + "sync" + "sync/atomic" + "testing" + "time" +) + +func TestRecordingRecoveryDiscoversAndResumesAfterProcessRestart(t *testing.T) { + root := privateRecoveryRoot(t) + now := time.Date(2026, 10, 1, 10, 0, 0, 0, time.UTC) + original := RecordingRecovery{Root: root, Now: func() time.Time { return now }} + entry, err := original.SaveFailure([]byte("RIFF original audio"), failedRecordingFixture(), &UploadHTTPError{StatusCode: 503}) + if err != nil { + t.Fatal(err) + } + now = now.Add(time.Minute) + var puts, reports atomic.Int32 + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + _, _ = io.Copy(io.Discard, r.Body) + puts.Add(1) + })) + defer server.Close() + restarted := RecordingRecovery{Root: root, Now: func() time.Time { return now }, Upload: UploadClient{AllowInsecureHTTP: true, Now: func() time.Time { return now }}} + grant := func(_ context.Context, e RecordingRecoveryEntry) (RecoveryTarget, error) { + return recoveryTarget(now, e, server.URL+"/original"), nil + } + report := func(_ context.Context, e RecordingRecoveryEntry) error { + if e.SourceEventID != entry.SourceEventID || e.Uploaded == nil { + t.Errorf("restarted reporter lost durable original upload identity: %+v", e) + } + reports.Add(1) + return nil + } + processed, err := restarted.RecoverDue(context.Background(), grant, report) + if err != nil || processed != 1 || puts.Load() != 1 || reports.Load() != 1 { + t.Fatalf("startup recovery did not restore the original result: processed=%d puts=%d reports=%d err=%v", processed, puts.Load(), reports.Load(), err) + } + processed, err = restarted.RecoverDue(context.Background(), grant, report) + if err != nil || processed != 0 || puts.Load() != 1 || reports.Load() != 1 { + t.Fatalf("delivered result was re-uploaded after second scan: processed=%d puts=%d reports=%d err=%v", processed, puts.Load(), reports.Load(), err) + } +} + +func TestRecordingRecoveryScannerReportsOrphanInsteadOfInventingResult(t *testing.T) { + root := privateRecoveryRoot(t) + broken := RecordingRecovery{Root: root, writeState: func(string, any) error { return errors.New("metadata disk full") }} + entry := failedRecordingFixture() + if _, err := broken.SaveFailure([]byte("RIFF original audio"), entry, &UploadHTTPError{StatusCode: 503}); err == nil { + t.Fatal("metadata disk fault was hidden") + } + restarted := RecordingRecovery{Root: root} + processed, err := restarted.RecoverDue(context.Background(), + func(context.Context, RecordingRecoveryEntry) (RecoveryTarget, error) { + t.Fatal("orphan must not PUT") + return RecoveryTarget{}, nil + }, + func(context.Context, RecordingRecoveryEntry) error { + t.Fatal("orphan must not report a result") + return nil + }) + if processed != 0 || !errors.Is(err, ErrRecordingRecoveryIncomplete) { + t.Fatalf("orphan was silently skipped: processed=%d err=%v", processed, err) + } + trace := sha256.Sum256([]byte(filepath.ToSlash(filepath.Join(entry.Bucket, entry.ObjectKey)))) + if !strings.Contains(err.Error(), hex.EncodeToString(trace[:6])) { + t.Fatalf("orphan error lacks a safe stable path reference: %v", err) + } + if _, err := os.Stat(filepath.Join(root, "mock-bucket", "tenant", "call", "rec.wav", "recording.wav")); err != nil { + t.Fatalf("original audio was removed while blocked: %v", err) + } +} + +func TestRecordingRecoveryConcurrentRetryOnlyPutsOnce(t *testing.T) { + root := privateRecoveryRoot(t) + now := time.Date(2026, 10, 1, 10, 0, 0, 0, time.UTC) + recovery := RecordingRecovery{Root: root, Now: func() time.Time { return now }, Upload: UploadClient{AllowInsecureHTTP: true, Now: func() time.Time { return now }}} + entry, err := recovery.SaveFailure([]byte("RIFF original audio"), failedRecordingFixture(), &UploadHTTPError{StatusCode: 503}) + if err != nil { + t.Fatal(err) + } + now = now.Add(time.Minute) + var puts, reports atomic.Int32 + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + _, _ = io.Copy(io.Discard, r.Body) + puts.Add(1) + time.Sleep(15 * time.Millisecond) + })) + defer server.Close() + grant := func(_ context.Context, e RecordingRecoveryEntry) (RecoveryTarget, error) { + return recoveryTarget(now, e, server.URL), nil + } + report := func(context.Context, RecordingRecoveryEntry) error { reports.Add(1); return nil } + var wg sync.WaitGroup + errs := make(chan error, 2) + for i := 0; i < 2; i++ { + wg.Add(1) + go func() { + defer wg.Done() + _, err := recovery.Retry(context.Background(), entry.Bucket, entry.ObjectKey, grant, report) + errs <- err + }() + } + wg.Wait() + close(errs) + for err := range errs { + if err != nil { + t.Fatalf("serialized retry failed: %v", err) + } + } + if puts.Load() != 1 || reports.Load() != 1 { + t.Fatalf("concurrent retry repeated a completed side effect: puts=%d reports=%d", puts.Load(), reports.Load()) + } +} + +func TestRecordingRecoveryRootMustBePrivate(t *testing.T) { + root := filepath.Join(t.TempDir(), "world-readable") + if err := os.Mkdir(root, 0700); err != nil { + t.Fatal(err) + } + if err := os.Chmod(root, 0755); err != nil { + t.Fatal(err) + } + recovery := RecordingRecovery{Root: root} + _, err := recovery.SaveFailure([]byte("RIFF original audio"), failedRecordingFixture(), &UploadHTTPError{StatusCode: 503}) + if err == nil { + t.Fatal("restricted recording was saved under a public recovery root") + } + files, err := os.ReadDir(root) + if err != nil || len(files) != 0 { + t.Fatalf("rejected root nevertheless contains recording files: files=%v err=%v", files, err) + } +} + +func TestRecordingRecoveryGrantFailureAlsoBacksOff(t *testing.T) { + now := time.Date(2026, 10, 1, 10, 0, 0, 0, time.UTC) + recovery := RecordingRecovery{Root: privateRecoveryRoot(t), Now: func() time.Time { return now }} + entry, err := recovery.SaveFailure([]byte("RIFF original audio"), failedRecordingFixture(), &UploadHTTPError{StatusCode: 503}) + if err != nil { + t.Fatal(err) + } + now = now.Add(time.Minute) + grantCalls := 0 + grant := func(context.Context, RecordingRecoveryEntry) (RecoveryTarget, error) { + grantCalls++ + return RecoveryTarget{}, errors.New("temporary Dispatcher outage") + } + _, err = recovery.Retry(context.Background(), entry.Bucket, entry.ObjectKey, grant, func(context.Context, RecordingRecoveryEntry) error { return nil }) + if err == nil || grantCalls != 1 { + t.Fatalf("failed grant request was hidden: calls=%d err=%v", grantCalls, err) + } + pending, err := recovery.Load(entry.Bucket, entry.ObjectKey) + if err != nil || pending.State != "retry_pending" || pending.Attempts != 2 || pending.NextAttemptAt != now.Add(2*time.Minute) { + t.Fatalf("grant failure caused an immediate unbounded retry: state=%+v err=%v", pending, err) + } + now = now.Add(time.Minute) + if _, err := recovery.Retry(context.Background(), entry.Bucket, entry.ObjectKey, grant, func(context.Context, RecordingRecoveryEntry) error { return nil }); !errors.Is(err, ErrRecordingNotDue) || grantCalls != 1 { + t.Fatalf("grant request retried before the next window: calls=%d err=%v", grantCalls, err) + } +}