Persist failed OSS recordings for bounded recovery
This commit is contained in:
@@ -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 通过。
|
||||
|
||||
## 验收台账
|
||||
|
||||
|
||||
@@ -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()
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user