From 0a3e0aeec6f896d31c7dc7aed5fbfbc0d8f3e89a Mon Sep 17 00:00:00 2001 From: Rogee Date: Sun, 4 Oct 2026 22:22:53 +0800 Subject: [PATCH] persist pending Agent results and bound upload recovery --- internal/agent/recording_delivery.go | 107 ++++++--- internal/agent/recording_delivery_test.go | 100 ++++++++ internal/agent/recording_recovery.go | 2 + internal/agent/recording_scan.go | 6 + internal/agent/recording_scan_test.go | 17 ++ internal/agent/result_journal.go | 269 ++++++++++++++++++++++ internal/agent/result_journal_test.go | 138 +++++++++++ 7 files changed, 612 insertions(+), 27 deletions(-) create mode 100644 internal/agent/result_journal.go create mode 100644 internal/agent/result_journal_test.go diff --git a/internal/agent/recording_delivery.go b/internal/agent/recording_delivery.go index 23c1e4c..99cf4ae 100644 --- a/internal/agent/recording_delivery.go +++ b/internal/agent/recording_delivery.go @@ -36,9 +36,11 @@ type CompletedRecording struct { type RecordingDelivery struct { Call RecordingClient Recovery *RecordingRecovery + Journal *ResultJournal mu sync.Mutex ended bool + endJournaled bool attempted bool pending []byte pendingProof *agentpb.UploadObservation @@ -50,45 +52,40 @@ func (d *RecordingDelivery) Complete(ctx context.Context, recording CompletedRec if d.attempted { return ErrRecordingRetryManaged } + // Preserve a valid no-recording result before sending a call-end report + // whose response could be lost. Invalid payloads are still reported as + // ended, then rejected by the original validation below. + if !d.ended && d.Journal != nil && (!recording.Expected || recording.CaptureError != nil || len(recording.WAV) == 0) { + if payload, err := noUploadPayload(recording); err == nil { + if err := d.Journal.SaveEndPending(d.resultEntry(payload)); err != nil { + return err + } + d.endJournaled = true + } + } // Release the confirmed call independently of OSS and final-result delivery. if !d.ended { if err := d.Call.ReportEnded(ctx); err != nil { return err } d.ended = true + if d.endJournaled { + if err := d.Journal.MarkEnded(d.Call.SourceEventID); err != nil { + return err + } + } } - result, err := emptyRecordingResult(recording.ResultPayload) - if err != nil { - return err - } - if !recording.Expected { - if len(recording.WAV) != 0 || recording.CaptureError != nil || recording.RecordingID != "" || recording.UploadID != "" || recording.DurationMS != 0 { - return errors.New("call without a recording cannot contain audio or upload identity") - } - return d.report(ctx, recording.ResultPayload, nil) - } - if recording.CaptureError != nil || len(recording.WAV) == 0 { - if len(recording.WAV) != 0 { - return errors.New("recording has both an error and apparently uploadable audio") - } - var reason string - if err := json.Unmarshal(result["reason_message"], &reason); err != nil { - return errors.New("failed recording needs the original call reason") - } - if reason != "" { - reason += "; " - } - reason += "recording generation failed" - result["reason_message"], err = json.Marshal(reason) - if err != nil { - return err - } - payload, err := json.Marshal(result) + if !recording.Expected || recording.CaptureError != nil || len(recording.WAV) == 0 { + payload, err := noUploadPayload(recording) if err != nil { return err } return d.report(ctx, payload, nil) } + result, err := emptyRecordingResult(recording.ResultPayload) + if err != nil { + return err + } if recording.RecordingID == "" || recording.UploadID == "" || !validMonoWAV(recording.WAV, recording.DurationMS) { return errors.New("recording identity or 16-kHz mono WAV is invalid") } @@ -126,6 +123,13 @@ func (d *RecordingDelivery) Complete(ctx context.Context, recording CompletedRec return err } + // The result is durable before PUT. A crash before its success is + // observed leaves this entry blocked, never presumed uploaded. + if d.Journal != nil { + if err := d.Journal.SavePrepared(d.resultEntry(payload)); err != nil { + return err + } + } // A returned OSS failure needs the original exact result and bytes; no // failure in this path is allowed to fabricate an uploaded result. d.attempted = true @@ -135,6 +139,7 @@ func (d *RecordingDelivery) Complete(ctx context.Context, recording CompletedRec if errors.As(err, &rejected) || errors.Is(err, ErrUploadOutcomeUnknown) { entry := RecordingRecoveryEntry{ CallID: d.Call.SourceEventID, SourceEventID: d.Call.SourceEventID, + DispatcherID: d.Call.DispatcherID, TenantID: d.Call.TenantID, RecordingID: recording.RecordingID, UploadID: recording.UploadID, Bucket: grant.Bucket, ObjectKey: grant.ObjectKey, SHA256: digest, SizeBytes: int64(len(recording.WAV)), ResultPayload: payload, @@ -149,6 +154,11 @@ func (d *RecordingDelivery) Complete(ctx context.Context, recording CompletedRec UploadId: recording.UploadID, RecordingId: recording.RecordingID, PutStatusCode: int32(uploaded.StatusCode), SizeBytes: uploaded.SizeBytes, ChecksumSha256: uploaded.SHA256, } + if d.Journal != nil { + if err := d.Journal.MarkUploaded(d.Call.SourceEventID, proof); err != nil { + return fmt.Errorf("persist confirmed upload before reporting result: %w", err) + } + } return d.report(ctx, payload, proof) } @@ -163,13 +173,56 @@ func (d *RecordingDelivery) RetryResult(ctx context.Context) error { return d.report(ctx, d.pending, d.pendingProof) } +func noUploadPayload(recording CompletedRecording) ([]byte, error) { + result, err := emptyRecordingResult(recording.ResultPayload) + if err != nil { + return nil, err + } + if !recording.Expected { + if len(recording.WAV) != 0 || recording.CaptureError != nil || recording.RecordingID != "" || recording.UploadID != "" || recording.DurationMS != 0 { + return nil, errors.New("call without a recording cannot contain audio or upload identity") + } + return recording.ResultPayload, nil + } + if len(recording.WAV) != 0 { + return nil, errors.New("recording has both an error and apparently uploadable audio") + } + var reason string + if err := json.Unmarshal(result["reason_message"], &reason); err != nil { + return nil, errors.New("failed recording needs the original call reason") + } + if reason != "" { + reason += "; " + } + reason += "recording generation failed" + result["reason_message"], err = json.Marshal(reason) + if err != nil { + return nil, err + } + return json.Marshal(result) +} + +func (d *RecordingDelivery) resultEntry(payload []byte) ResultEntry { + return ResultEntry{SourceEventID: d.Call.SourceEventID, DispatcherID: d.Call.DispatcherID, TenantID: d.Call.TenantID, Payload: bytes.Clone(payload)} +} + func (d *RecordingDelivery) report(ctx context.Context, payload []byte, proof *agentpb.UploadObservation) error { + if d.Journal != nil && !d.attempted && !d.endJournaled { + if err := d.Journal.SaveReady(d.resultEntry(payload)); err != nil { + return err + } + } d.attempted = true d.pending = bytes.Clone(payload) d.pendingProof = proof if _, err := d.Call.ReportFinal(ctx, d.pending, d.pendingProof); err != nil { return err } + if d.Journal != nil { + if err := d.Journal.Ack(d.Call.SourceEventID); err != nil { + return fmt.Errorf("clear accepted final result journal: %w", err) + } + } d.pending = nil d.pendingProof = nil return nil diff --git a/internal/agent/recording_delivery_test.go b/internal/agent/recording_delivery_test.go index fa48e63..a0dfa89 100644 --- a/internal/agent/recording_delivery_test.go +++ b/internal/agent/recording_delivery_test.go @@ -11,6 +11,7 @@ import ( "net/http" "net/http/httptest" "os" + "path/filepath" "reflect" "sync/atomic" "testing" @@ -86,6 +87,86 @@ func TestRecordingDeliveryNoRecordingReportsOnlyAfterConfirmedEnd(t *testing.T) } } +func TestRecordingDeliveryConfirmedHangupSurvivesLostEndReport(t *testing.T) { + root := t.TempDir() + if err := os.Chmod(root, 0700); err != nil { + t.Fatal(err) + } + stub := &recordingDeliveryRPC{endError: errors.New("end reply lost")} + delivery := testRecordingDelivery(stub) + delivery.Journal = &ResultJournal{Root: root} + payload := []byte(`{"recording":{}}`) + if err := delivery.Complete(context.Background(), CompletedRecording{ResultPayload: payload}); err == nil { + t.Fatal("Dispatcher did not confirm the local hangup") + } + stub.endError = nil + if err := (&ResultJournal{Root: root}).RecoverWithEnd(context.Background(), func(ctx context.Context, entry ResultEntry) error { + call := delivery.Call + return call.ReportEnded(ctx) + }, func(ctx context.Context, entry ResultEntry) error { + call := delivery.Call + _, err := call.ReportFinal(ctx, entry.Payload, entry.Upload) + return err + }); err != nil { + t.Fatal(err) + } + if !reflect.DeepEqual(stub.calls, []string{"end", "end", "result"}) { + t.Fatalf("local hangup was not recovered under the original identity: %v", stub.calls) + } +} + +func TestRecordingFailurePreservesResultAcrossLostEndReport(t *testing.T) { + root := t.TempDir() + if err := os.Chmod(root, 0700); err != nil { + t.Fatal(err) + } + stub := &recordingDeliveryRPC{endError: errors.New("end reply lost")} + delivery := testRecordingDelivery(stub) + delivery.Journal = &ResultJournal{Root: root} + completed := CompletedRecording{Expected: true, CaptureError: errors.New("recording generation error"), ResultPayload: []byte(`{"recording":{},"reason_message":"media failed"}`)} + if err := delivery.Complete(context.Background(), completed); err == nil { + t.Fatal("lost end reply must block result delivery") + } + stub.endError = nil + if err := (&ResultJournal{Root: root}).RecoverWithEnd(context.Background(), func(ctx context.Context, entry ResultEntry) error { + return delivery.Call.ReportEnded(ctx) + }, func(ctx context.Context, entry ResultEntry) error { + _, err := delivery.Call.ReportFinal(ctx, entry.Payload, nil) + return err + }); err != nil { + t.Fatal(err) + } + if !reflect.DeepEqual(stub.calls, []string{"end", "end", "result"}) || !bytes.Contains(stub.resultBody, []byte("recording generation failed")) || stub.resultProof != nil { + t.Fatalf("recording failure was lost or reported as uploaded: calls=%v body=%s", stub.calls, stub.resultBody) + } +} + +func TestRecordingDeliveryUnconfirmedResultSurvivesRestart(t *testing.T) { + root := t.TempDir() + if err := os.Chmod(root, 0700); err != nil { + t.Fatal(err) + } + stub := &recordingDeliveryRPC{resultError: errors.New("reply lost")} + delivery := testRecordingDelivery(stub) + delivery.Journal = &ResultJournal{Root: root} + payload := []byte(`{"recording":{}}`) + if err := delivery.Complete(context.Background(), CompletedRecording{ResultPayload: payload}); err == nil { + t.Fatal("Dispatcher confirmation was not observed") + } + stub.resultError = nil + if err := (&ResultJournal{Root: root}).Recover(context.Background(), func(ctx context.Context, entry ResultEntry) error { + call := delivery.Call + call.SourceEventID = entry.SourceEventID + _, err := call.ReportFinal(ctx, entry.Payload, entry.Upload) + return err + }); err != nil { + t.Fatal(err) + } + if !reflect.DeepEqual(stub.calls, []string{"end", "result", "result"}) || !bytes.Equal(stub.resultBody, payload) { + t.Fatalf("restart changed or failed to resend the original result: %v", stub.calls) + } +} + func TestRecordingDeliveryCannotReturnResultWhenEndNotConfirmed(t *testing.T) { stub := &recordingDeliveryRPC{endError: errors.New("injected Dispatcher persistence failure")} delivery := testRecordingDelivery(stub) @@ -168,6 +249,25 @@ func TestRecordingDeliveryDirectPUTThenOneResultWithoutBusinessFiles(t *testing. } } +func TestRecordingDeliverySuccessfulPUTLeavesNoResultOrAudioFiles(t *testing.T) { + delivery, _, completed, puts, root := testDirectRecordingDelivery(t, http.StatusCreated) + delivery.Journal = &ResultJournal{Root: root} + if err := delivery.Complete(context.Background(), completed); err != nil { + t.Fatal(err) + } + if puts.Load() != 1 { + t.Fatalf("successful delivery repeated PUT %d times", puts.Load()) + } + entries, err := os.ReadDir(filepath.Join(root, ".results")) + if err != nil || len(entries) != 0 { + t.Fatalf("confirmed result was not removed: entries=%v err=%v", entries, err) + } + rootEntries, err := os.ReadDir(root) + if err != nil || len(rootEntries) != 1 || rootEntries[0].Name() != ".results" { + t.Fatalf("normal success retained audio or business files: entries=%v err=%v", rootEntries, err) + } +} + func TestRecordingDeliveryFailedPUTPersistsOriginalAndNeverRetriesImplicitly(t *testing.T) { delivery, stub, completed, puts, _ := testDirectRecordingDelivery(t, http.StatusServiceUnavailable) var rejected *UploadHTTPError diff --git a/internal/agent/recording_recovery.go b/internal/agent/recording_recovery.go index 43d6b8d..e62ffdf 100644 --- a/internal/agent/recording_recovery.go +++ b/internal/agent/recording_recovery.go @@ -24,6 +24,8 @@ var ErrRecordingRecoveryIncomplete = errors.New("recording recovery files are in type RecordingRecoveryEntry struct { CallID string `json:"call_id"` SourceEventID string `json:"source_event_id"` + DispatcherID string `json:"dispatcher_id,omitempty"` + TenantID int64 `json:"tenant_id,omitempty"` RecordingID string `json:"recording_id"` UploadID string `json:"upload_id"` Bucket string `json:"bucket"` diff --git a/internal/agent/recording_scan.go b/internal/agent/recording_scan.go index 3e76b19..78dc12b 100644 --- a/internal/agent/recording_scan.go +++ b/internal/agent/recording_scan.go @@ -47,6 +47,12 @@ func (r *RecordingRecovery) RecoverDue(ctx context.Context, grant RecoveryGrant, return err } if item.IsDir() { + if path == filepath.Join(r.Root, ".results") { + if err := checkRecoveryDirectory(path); err != nil { + return err + } + return filepath.SkipDir // result-only files have their own scanner + } return checkRecoveryDirectory(path) } dir := filepath.Dir(path) diff --git a/internal/agent/recording_scan_test.go b/internal/agent/recording_scan_test.go index 2cdb1a8..d7cb255 100644 --- a/internal/agent/recording_scan_test.go +++ b/internal/agent/recording_scan_test.go @@ -53,6 +53,23 @@ func TestRecordingRecoveryDiscoversAndResumesAfterProcessRestart(t *testing.T) { } } +func TestRecordingRecoveryScannerLeavesResultJournalToItsOwnRecovery(t *testing.T) { + root := privateRecoveryRoot(t) + journal := ResultJournal{Root: root} + if err := journal.SaveReady(ResultEntry{SourceEventID: "original", DispatcherID: "11111111-1111-4111-8111-111111111111", TenantID: 1, Payload: []byte(`{"recording":{}}`)}); err != nil { + t.Fatal(err) + } + processed, err := (&RecordingRecovery{Root: root}).RecoverDue(context.Background(), + func(context.Context, RecordingRecoveryEntry) (RecoveryTarget, error) { + t.Fatal("no failed PUT exists") + return RecoveryTarget{}, nil + }, + func(context.Context, RecordingRecoveryEntry) error { t.Fatal("no failed PUT exists"); return nil }) + if processed != 0 || err != nil { + t.Fatalf("result-only journal is not a failed PUT: processed=%d err=%v", processed, err) + } +} + func TestRecordingRecoveryScannerReportsOrphanInsteadOfInventingResult(t *testing.T) { root := privateRecoveryRoot(t) broken := RecordingRecovery{Root: root, writeState: func(string, any) error { return errors.New("metadata disk full") }} diff --git a/internal/agent/result_journal.go b/internal/agent/result_journal.go new file mode 100644 index 0000000..1cc8a7a --- /dev/null +++ b/internal/agent/result_journal.go @@ -0,0 +1,269 @@ +package agent + +import ( + "bytes" + "context" + "crypto/sha256" + "encoding/hex" + "encoding/json" + "errors" + "fmt" + "io" + "log" + "os" + "path/filepath" + "strings" + + agentpb "git.ipao.vip/rogee/go-sip/gen/agent" + "google.golang.org/protobuf/proto" +) + +// ResultEntry never includes audio, OSS credentials or a signed URL. Prepared +// results must not be reported: a PUT could have been interrupted or rejected. +type ResultEntry struct { + SourceEventID string `json:"source_event_id"` + DispatcherID string `json:"dispatcher_id"` + TenantID int64 `json:"tenant_id"` + Payload []byte `json:"payload"` + Upload *agentpb.UploadObservation `json:"upload,omitempty"` + State string `json:"state"` +} + +type ResultJournal struct{ Root string } + +func (j ResultJournal) path(id string) (string, error) { + if j.Root == "" || id == "" || strings.ContainsAny(id, "/\\\x00") { + return "", errors.New("private result journal and original call identity are required") + } + sum := sha256.Sum256([]byte(id)) + return filepath.Join(j.Root, ".results", hex.EncodeToString(sum[:])+".json"), nil +} + +func (j ResultJournal) prepare() (string, error) { + if j.Root == "" { + return "", errors.New("private result journal root is required") + } + if err := checkRecoveryDirectory(j.Root); err != nil { + return "", err + } + dir := filepath.Join(j.Root, ".results") + if err := os.Mkdir(dir, 0700); err != nil && !errors.Is(err, os.ErrExist) { + return "", err + } + if err := checkRecoveryDirectory(dir); err != nil { + return "", err + } + return dir, nil +} + +func (j ResultJournal) save(entry ResultEntry, state string) error { + if entry.DispatcherID == "" || entry.TenantID <= 0 || !json.Valid(entry.Payload) || len(entry.Payload) == 0 || entry.Upload != nil || entry.State != "" { + return errors.New("original result identity or payload is invalid") + } + path, err := j.path(entry.SourceEventID) + if err != nil { + return err + } + dir, err := j.prepare() + if err != nil { + return err + } + entry.State = state + data, err := json.Marshal(entry) + if err != nil { + return err + } + file, err := os.OpenFile(path, os.O_WRONLY|os.O_CREATE|os.O_EXCL, 0600) + if err != nil { + return fmt.Errorf("save original result: %w", err) + } + written, writeErr := file.Write(data) + if writeErr == nil && written != len(data) { + writeErr = io.ErrShortWrite + } + syncErr := file.Sync() + closeErr := file.Close() + if err := firstError(writeErr, syncErr, closeErr); err != nil { + return fmt.Errorf("persist original result: %w", err) + } + return syncDirectory(dir) +} + +func (j ResultJournal) SaveReady(entry ResultEntry) error { return j.save(entry, "ready") } +func (j ResultJournal) SavePrepared(entry ResultEntry) error { return j.save(entry, "prepared") } +func (j ResultJournal) SaveEndPending(entry ResultEntry) error { return j.save(entry, "end_pending") } + +func (j ResultJournal) load(id string) (ResultEntry, string, error) { + path, err := j.path(id) + if err != nil { + return ResultEntry{}, "", err + } + info, err := os.Lstat(path) + if err != nil { + return ResultEntry{}, "", err + } + if !info.Mode().IsRegular() || info.Mode().Perm()&0077 != 0 || info.Size() > 1<<20 { + return ResultEntry{}, "", errors.New("original result file is not private or exceeds the supported size") + } + data, err := os.ReadFile(path) + if err != nil { + return ResultEntry{}, "", err + } + var entry ResultEntry + decoder := json.NewDecoder(bytes.NewReader(data)) + decoder.DisallowUnknownFields() + if err := decoder.Decode(&entry); err != nil { + return ResultEntry{}, "", err + } + var trailing any + if err := decoder.Decode(&trailing); err == nil { + return ResultEntry{}, "", errors.New("original result has trailing data") + } else if !errors.Is(err, io.EOF) { + return ResultEntry{}, "", err + } + if entry.SourceEventID != id || entry.DispatcherID == "" || entry.TenantID <= 0 || !json.Valid(entry.Payload) || len(entry.Payload) == 0 || (entry.State != "prepared" && entry.State != "ready" && entry.State != "end_pending") || (entry.State != "ready" && entry.Upload != nil) { + return ResultEntry{}, "", errors.New("original result journal identity or state is invalid") + } + return entry, path, nil +} + +func (j ResultJournal) MarkEnded(id string) error { + entry, path, err := j.load(id) + if err != nil { + return err + } + if entry.State != "end_pending" { + return errors.New("original end is not awaiting confirmation") + } + entry.State = "ready" + return writeJSONAtomic(path, entry) +} + +func (j ResultJournal) MarkUploaded(id string, observation *agentpb.UploadObservation) error { + entry, path, err := j.load(id) + if err != nil { + return err + } + if entry.State != "prepared" || observation == nil || observation.GetUploadId() == "" || observation.GetRecordingId() == "" || observation.GetPutStatusCode() < 200 || observation.GetPutStatusCode() >= 300 || observation.GetSizeBytes() <= 0 || observation.GetChecksumSha256() == "" { + return errors.New("confirmed original upload observation is required") + } + entry.State = "ready" + entry.Upload = proto.Clone(observation).(*agentpb.UploadObservation) + return writeJSONAtomic(path, entry) +} + +func (j ResultJournal) Ack(id string) error { + entry, path, err := j.load(id) + if err != nil { + return err + } + if entry.State != "ready" { + return errors.New("cannot remove a result whose upload outcome is unknown") + } + if err := os.Remove(path); err != nil { + return err + } + return syncDirectory(filepath.Dir(path)) +} + +// DiscardPrepared is only called after the durable audio+result recovery entry +// has itself reported the same final result. Old entries have no journal. +func (j ResultJournal) DiscardPrepared(id string) error { + entry, path, err := j.load(id) + if err != nil { + if errors.Is(err, os.ErrNotExist) { + return nil + } + return err + } + if entry.State != "prepared" { + return errors.New("cannot discard a confirmed result") + } + if err := os.Remove(path); err != nil { + return err + } + return syncDirectory(filepath.Dir(path)) +} + +// Recover only re-reports ready results under their original identity. A +// prepared upload is left untouched for inspection, never retried or assumed +// uploaded. Independent entries continue when one is blocked. +func (j ResultJournal) Recover(ctx context.Context, report func(context.Context, ResultEntry) error) error { + return j.RecoverWithEnd(ctx, nil, report) +} + +func (j ResultJournal) RecoverWithEnd(ctx context.Context, confirmEnd, report func(context.Context, ResultEntry) error) error { + dir, err := j.prepare() + if err != nil { + return err + } + items, err := os.ReadDir(dir) + if err != nil { + return err + } + var blocked []error + for _, item := range items { + if err := ctx.Err(); err != nil { + return errors.Join(append(blocked, err)...) + } + path := filepath.Join(dir, item.Name()) + if item.IsDir() || !strings.HasSuffix(item.Name(), ".json") || len(item.Name()) != 69 { + blocked = append(blocked, errors.New("unexpected private result journal entry")) + continue + } + info, err := os.Lstat(path) + if err != nil || !info.Mode().IsRegular() || info.Mode().Perm()&0077 != 0 || info.Size() > 1<<20 { + blocked = append(blocked, errors.New("private result file is invalid")) + continue + } + data, err := os.ReadFile(path) + if err != nil { + blocked = append(blocked, err) + continue + } + var id struct { + SourceEventID string `json:"source_event_id"` + } + if err := json.Unmarshal(data, &id); err != nil { + blocked = append(blocked, errors.New("invalid private result journal")) + continue + } + entry, originalPath, err := j.load(id.SourceEventID) + if err != nil || path != originalPath { + blocked = append(blocked, errors.New("result journal identity mismatch")) + continue + } + if entry.State == "end_pending" { + if confirmEnd == nil { + blocked = append(blocked, errors.New("original call end requires Dispatcher confirmation")) + continue + } + if err := confirmEnd(ctx, entry); err != nil { + blocked = append(blocked, fmt.Errorf("confirm original call end: %w", err)) + continue + } + if err := j.MarkEnded(entry.SourceEventID); err != nil { + blocked = append(blocked, fmt.Errorf("persist original call end: %w", err)) + continue + } + entry.State = "ready" + } + if entry.State != "ready" { + log.Printf("Agent recovery blocked stage=original_upload_outcome event_digest=%s", item.Name()[:12]) + blocked = append(blocked, errors.New("original upload outcome requires inspection")) + continue + } + if report == nil { + blocked = append(blocked, errors.New("result reporter is required")) + continue + } + if err := report(ctx, entry); err != nil { + blocked = append(blocked, fmt.Errorf("report original result: %w", err)) + continue + } + if err := j.Ack(entry.SourceEventID); err != nil { + blocked = append(blocked, fmt.Errorf("confirm original result: %w", err)) + } + } + return errors.Join(blocked...) +} diff --git a/internal/agent/result_journal_test.go b/internal/agent/result_journal_test.go new file mode 100644 index 0000000..5ca137a --- /dev/null +++ b/internal/agent/result_journal_test.go @@ -0,0 +1,138 @@ +package agent + +import ( + "context" + "errors" + "os" + "path/filepath" + "testing" + + agentpb "git.ipao.vip/rogee/go-sip/gen/agent" +) + +func TestResultJournalResumesOnlyConfirmedResults(t *testing.T) { + root := t.TempDir() + if err := os.Chmod(root, 0700); err != nil { + t.Fatal(err) + } + j := ResultJournal{Root: root} + entry := ResultEntry{SourceEventID: "call-1", DispatcherID: "11111111-1111-4111-8111-111111111111", TenantID: 1, Payload: []byte(`{"event_id":"call-1","recording":{}}`)} + if err := j.SaveReady(entry); err != nil { + t.Fatal(err) + } + attempts := 0 + restarted := ResultJournal{Root: root} + if err := restarted.Recover(context.Background(), func(_ context.Context, got ResultEntry) error { + attempts++ + if got.SourceEventID != entry.SourceEventID || string(got.Payload) != string(entry.Payload) || got.Upload != nil { + t.Fatalf("changed original result: %+v", got) + } + if attempts == 1 { + return errors.New("Dispatcher response unknown") + } + return nil + }); err == nil { + t.Fatal("unconfirmed result must remain for a later retry") + } + if err := restarted.Recover(context.Background(), func(ctx context.Context, got ResultEntry) error { + attempts++ + return nil + }); err != nil { + t.Fatal(err) + } + if attempts != 2 { + t.Fatalf("expected one retry, got %d attempts", attempts) + } + if err := restarted.Recover(context.Background(), func(context.Context, ResultEntry) error { t.Fatal("confirmed result repeated"); return nil }); err != nil { + t.Fatal(err) + } +} + +func TestResultJournalNeverReportsPreparedUpload(t *testing.T) { + root := t.TempDir() + if err := os.Chmod(root, 0700); err != nil { + t.Fatal(err) + } + j := ResultJournal{Root: root} + entry := ResultEntry{SourceEventID: "call-2", DispatcherID: "11111111-1111-4111-8111-111111111111", TenantID: 1, Payload: []byte(`{"event_id":"call-2","recording":{"status":"uploaded"}}`)} + if err := j.SavePrepared(entry); err != nil { + t.Fatal(err) + } + if err := (&ResultJournal{Root: root}).Recover(context.Background(), func(context.Context, ResultEntry) error { + t.Fatal("a PUT with unknown outcome cannot generate a final result") + return nil + }); err == nil { + t.Fatal("prepared result must remain visibly blocked") + } + observation := &agentpb.UploadObservation{UploadId: "upload-2", RecordingId: "recording-2", PutStatusCode: 200, SizeBytes: 10, ChecksumSha256: "digest"} + if err := j.MarkUploaded("call-2", observation); err != nil { + t.Fatal(err) + } + if err := j.DiscardPrepared("call-2"); err == nil { + t.Fatal("a ready result must not be discarded as an incomplete upload") + } + if err := (&ResultJournal{Root: root}).Recover(context.Background(), func(_ context.Context, got ResultEntry) error { + if got.Upload == nil || got.Upload.GetUploadId() != observation.GetUploadId() { + t.Fatal("uploaded proof missing after restart") + } + return nil + }); err != nil { + t.Fatal(err) + } + if _, err := os.ReadDir(filepath.Join(root, ".results")); err != nil { + t.Fatal(err) + } +} + +func TestResultJournalRejectsTrailingCorruptData(t *testing.T) { + root := t.TempDir() + if err := os.Chmod(root, 0700); err != nil { + t.Fatal(err) + } + j := ResultJournal{Root: root} + entry := ResultEntry{SourceEventID: "call-corrupt", DispatcherID: "11111111-1111-4111-8111-111111111111", TenantID: 1, Payload: []byte(`{"recording":{}}`)} + if err := j.SaveReady(entry); err != nil { + t.Fatal(err) + } + path, err := j.path(entry.SourceEventID) + if err != nil { + t.Fatal(err) + } + file, err := os.OpenFile(path, os.O_APPEND|os.O_WRONLY, 0600) + if err != nil { + t.Fatal(err) + } + if _, err := file.WriteString("\n{}"); err != nil { + t.Fatal(err) + } + if err := file.Close(); err != nil { + t.Fatal(err) + } + if err := j.Recover(context.Background(), func(context.Context, ResultEntry) error { + t.Fatal("corrupt result was reported") + return nil + }); err == nil { + t.Fatal("corrupt result must remain visibly blocked") + } +} + +func TestResultJournalDiscardPreparedAfterRecovery(t *testing.T) { + root := t.TempDir() + if err := os.Chmod(root, 0700); err != nil { + t.Fatal(err) + } + j := ResultJournal{Root: root} + entry := ResultEntry{SourceEventID: "call-3", DispatcherID: "11111111-1111-4111-8111-111111111111", TenantID: 1, Payload: []byte(`{"recording":{"status":"uploaded"}}`)} + if err := j.SavePrepared(entry); err != nil { + t.Fatal(err) + } + if err := j.DiscardPrepared("call-3"); err != nil { + t.Fatal(err) + } + if err := j.DiscardPrepared("call-3"); err != nil { + t.Fatal(err) + } + if items, err := os.ReadDir(filepath.Join(root, ".results")); err != nil || len(items) != 0 { + t.Fatalf("confirmed recovered result still has a prepared journal: %v %v", items, err) + } +}