diff --git a/pkg/cmd/streamplace.go b/pkg/cmd/streamplace.go index 96ad214d..1781a95f 100644 --- a/pkg/cmd/streamplace.go +++ b/pkg/cmd/streamplace.go @@ -207,6 +207,11 @@ func runMain(ctx context.Context, build *config.BuildFlags, platformJobs []jobFu } } + // Captured before any listener can accept a stream: spool session dirs + // named (or modified) after this instant belong to live sessions of THIS + // process, and the salvage pass must not touch them. + bootTime := time.Now() + group, ctx := TimeoutGroupWithContext(ctx) out := carstore.SQLiteStore{} @@ -574,12 +579,13 @@ func runMain(ctx context.Context, build *config.BuildFlags, platformJobs []jobFu // Upload any live-rec spool data a prior run left behind (crash, or a // stream that ended while its bucket was failing) — the recording - // completes late instead of being lost. Startup-only by design: spool - // dirs present now are orphans; new sessions create fresh ones. + // completes late instead of being lost. Only spools from BEFORE bootTime + // are touched: sessions of this process create their dirs after it, so + // salvage can never race a live uploader. if cli.S3Configured() && cli.LiveRecSpoolMaxMB > 0 { group.Go(func() error { if err := sps3.SalvageSpools(ctx, cli.S3Config(), state, - director.LiveRecSpoolRoot(cli.DataDir), director.LiveRecKeyPrefix, sps3.DefaultCutoverEvery); err != nil { + director.LiveRecSpoolRoot(cli.DataDir), director.LiveRecKeyPrefix, sps3.DefaultCutoverEvery, bootTime); err != nil { log.Error(ctx, "salvaging live-rec spools", "error", err) } return nil diff --git a/pkg/s3/s3.go b/pkg/s3/s3.go index 1fa8fdb5..176b2ac8 100644 --- a/pkg/s3/s3.go +++ b/pkg/s3/s3.go @@ -271,6 +271,21 @@ type objectWriter struct { u *S3Uploader current *activeUpload objSeq int // disambiguates keys when two objects roll over within one second + + // keySuffix is appended to every object key before ".m4s". The salvage + // path sets it to the spool session name so salvaged keys can never + // collide with each other (objSeq resets per salvaged spool and mtimes + // have second granularity) or with keys the crashed run already wrote. + keySuffix string + + // strictRecorder makes Recorder failures upload failures. The spool loop + // and salvage set this: their segments survive on disk, so failing and + // retrying beats completing an object that no s3_segments row points at — + // such an object is invisible to finalize, and once the spool is acked + // the bytes are unrecoverable. The memory loop stays lenient: without a + // spool the segments are gone either way, and an untracked object in the + // bucket at least preserves the bytes. + strictRecorder bool } // start opens a new multipart object. started is the moment the object's @@ -280,7 +295,7 @@ type objectWriter struct { // still sort correctly in the assembled VOD. func (w *objectWriter) start(ctx context.Context, uri string, started time.Time) error { w.objSeq++ - key := fmt.Sprintf("%s%s-%d.m4s", w.u.keyPrefix, started.UTC().Format("2006-01-02T15-04-05"), w.objSeq) + key := fmt.Sprintf("%s%s-%d%s.m4s", w.u.keyPrefix, started.UTC().Format("2006-01-02T15-04-05"), w.objSeq, w.keySuffix) resp, err := w.u.client.CreateMultipartUpload(ctx, &s3.CreateMultipartUploadInput{ Bucket: aws.String(w.u.bucket), @@ -301,6 +316,14 @@ func (w *objectWriter) start(ctx context.Context, uri string, started time.Time) id, recErr := w.u.recorder.RecordStart(ctx, w.u.userDID, w.u.bucket, key, uri, started) if recErr != nil { log.Error(ctx, "recording S3 upload start", "key", key, "error", recErr) + if w.strictRecorder { + // An object without a row is invisible to finalize; fail now + // (the segments are still spooled) rather than upload into a + // black hole. + w.abortMultipart(ctx) + w.current = nil + return fmt.Errorf("recording S3 upload start for %s: %w", key, recErr) + } } w.current.recordID = id } @@ -385,6 +408,14 @@ func (w *objectWriter) complete(ctx context.Context) error { if w.u.recorder != nil && w.current.recordID != "" { if recErr := w.u.recorder.RecordComplete(ctx, w.current.recordID, int32(len(w.current.parts)), w.current.totalSize); recErr != nil { log.Error(ctx, "recording S3 upload completion", "key", w.current.key, "error", recErr) + if w.strictRecorder { + // The object completed in S3 but finalize will never see it + // (its row stays incomplete). Surface the failure so the spool + // segments are retried into a fresh, properly-recorded object; + // the completed-but-unrecorded one is orphaned storage, not + // lost footage. + return fmt.Errorf("recording S3 upload completion for %s: %w", w.current.key, recErr) + } } } w.current = nil @@ -533,7 +564,7 @@ var ( // crashes, costing only lag. func (u *S3Uploader) uploadLoopSpool(ctx context.Context) { ctx = log.WithLogValues(ctx, "func", "s3.uploadLoopSpool") - w := &objectWriter{u: u} + w := &objectWriter{u: u, strictRecorder: true} var nextSeq int64 = 1 // next spool seq to consume var objFirstSeq int64 // first seq of the in-flight object (rewind point) var lastConsumed int64 // last seq appended into the in-flight object @@ -576,8 +607,16 @@ func (u *S3Uploader) uploadLoopSpool(ctx context.Context) { // times (spool mtimes), so objects assembled during catch-up still map to // real stream time and sort honestly in finalize. process := func() { - if !retryAt.IsZero() && time.Now().Before(retryAt) { - return + if !retryAt.IsZero() { + if time.Now().Before(retryAt) { + return + } + // The retry is due: clear it so the wake timer stops re-arming + // with an in-the-past deadline (a ~100ms busy-wake otherwise). + // backoff itself only resets when an object completes, so a + // still-failing bucket keeps escalating rather than restarting + // at spoolRetryMin every part success. + retryAt = time.Time{} } for { seq, ok := u.spool.NextFrom(nextSeq) @@ -651,11 +690,14 @@ func (u *S3Uploader) uploadLoopSpool(ctx context.Context) { } if cmd.cutover { // Close out the current object so it's immediately finalize-able - // (e.g. the livestream just ended). During backoff there may be - // nothing in flight; the pending segments complete on retry (a - // late cutover can then merge across the cut — finalize orders by - // started-at, so the recording stays correct, just chunked - // differently). + // (e.g. the livestream just ended). A cutover overrides any + // backoff: the stream is ending, so make one immediate attempt + // at the backlog — if the bucket is still down it fails fast and + // re-arms the backoff, and the tail completes on a later retry + // or next-startup salvage (a late tail can merge across the cut; + // finalize orders by started-at, so the recording stays correct, + // just chunked differently). + retryAt = time.Time{} process() if err := completeAndAck(); err != nil { fail(fmt.Errorf("error completing upload on cutover: %w", err)) diff --git a/pkg/s3/salvage.go b/pkg/s3/salvage.go index 00c8d86b..816c2651 100644 --- a/pkg/s3/salvage.go +++ b/pkg/s3/salvage.go @@ -5,6 +5,7 @@ import ( "fmt" "os" "path/filepath" + "strconv" "time" "stream.place/streamplace/pkg/log" @@ -21,20 +22,22 @@ import ( // objects by started-at, so salvaged objects sort into the right place no // matter how late they upload. Drained spools are deleted. // -// Run this once at startup, before stream sessions pile up. It's safe against -// concurrent new sessions because every session creates a fresh uniquely-named -// spool dir: anything present at startup is by construction orphaned. (A spool +// Run this once at startup. startedBefore must be a timestamp from before any +// of this process's listeners could accept a stream: only session dirs created +// before it are touched (session dirs are named by their creation UnixNano), +// so salvage can never open a spool a live uploader of this process owns — +// even if a streamer reconnects while the scan is still running. (A spool // abandoned by a failing Close mid-run waits for the next restart; that needs // a failure at the exact end of a stream, and the data just sits on disk in // the meantime.) // // Failures leave the affected spool in place for the next attempt and move on // to the next one; the returned error is only ever a scan-level failure. -func SalvageSpools(ctx context.Context, cfg Config, recorder Recorder, root string, keyPrefixFor func(did string) string, cutoverEvery time.Duration) error { - return salvageSpools(ctx, NewClient(cfg), cfg.Bucket, recorder, root, keyPrefixFor, cutoverEvery) +func SalvageSpools(ctx context.Context, cfg Config, recorder Recorder, root string, keyPrefixFor func(did string) string, cutoverEvery time.Duration, startedBefore time.Time) error { + return salvageSpools(ctx, NewClient(cfg), cfg.Bucket, recorder, root, keyPrefixFor, cutoverEvery, startedBefore) } -func salvageSpools(ctx context.Context, client uploadAPI, bucket string, recorder Recorder, root string, keyPrefixFor func(did string) string, cutoverEvery time.Duration) error { +func salvageSpools(ctx context.Context, client uploadAPI, bucket string, recorder Recorder, root string, keyPrefixFor func(did string) string, cutoverEvery time.Duration, startedBefore time.Time) error { ctx = log.WithLogValues(ctx, "func", "s3.SalvageSpools") didDirs, err := os.ReadDir(root) if os.IsNotExist(err) { @@ -57,8 +60,12 @@ func salvageSpools(ctx context.Context, client uploadAPI, bucket string, recorde if !sess.IsDir() { continue } + if !sessionStartedBefore(sess, startedBefore) { + // A live session of this process — never touch it. + continue + } dir := filepath.Join(root, did, sess.Name()) - if err := salvageOne(ctx, client, bucket, recorder, dir, did, keyPrefixFor(did), cutoverEvery); err != nil { + if err := salvageOne(ctx, client, bucket, recorder, dir, did, keyPrefixFor(did), sess.Name(), cutoverEvery); err != nil { log.Error(ctx, "salvaging live-rec spool failed; leaving it for the next attempt", "dir", dir, "error", err) } } @@ -68,11 +75,27 @@ func salvageSpools(ctx context.Context, client uploadAPI, bucket string, recorde return nil } -// salvageOne drains a single leftover spool into live-rec objects. Any error +// sessionStartedBefore reports whether a spool session dir predates cutoff. +// Session dirs are named by their creation UnixNano (see the director); a +// non-numeric name (foreign layout) falls back to the dir mtime. +func sessionStartedBefore(sess os.DirEntry, cutoff time.Time) bool { + if nanos, err := strconv.ParseInt(sess.Name(), 10, 64); err == nil { + return nanos < cutoff.UnixNano() + } + info, err := sess.Info() + if err != nil { + return false // can't tell: leave it alone + } + return info.ModTime().Before(cutoff) +} + +// salvageOne drains a single leftover spool into live-rec objects. session +// (the spool dir name) is baked into every object key so salvaged keys can't +// collide across spools or with keys the crashed run already wrote. Any error // leaves the remaining segments on disk (already-completed objects stay // completed — their segments were acked, so a re-run picks up exactly where // this one failed). -func salvageOne(ctx context.Context, client uploadAPI, bucket string, recorder Recorder, dir, did, keyPrefix string, cutoverEvery time.Duration) error { +func salvageOne(ctx context.Context, client uploadAPI, bucket string, recorder Recorder, dir, did, keyPrefix, session string, cutoverEvery time.Duration) error { spool, err := OpenSpool(dir, 0) if err != nil { return err @@ -94,7 +117,7 @@ func salvageOne(ctx context.Context, client uploadAPI, bucket string, recorder R recorder: recorder, spool: spool, } - w := &objectWriter{u: u} + w := &objectWriter{u: u, strictRecorder: true, keySuffix: "-salvaged-" + session} var nextSeq, lastConsumed int64 = 1, 0 var salvaged int64 for { diff --git a/pkg/s3/salvage_test.go b/pkg/s3/salvage_test.go index a03576b8..9578bb92 100644 --- a/pkg/s3/salvage_test.go +++ b/pkg/s3/salvage_test.go @@ -2,6 +2,7 @@ package s3 import ( "context" + "fmt" "os" "path/filepath" "testing" @@ -37,7 +38,7 @@ func TestSalvageSpools(t *testing.T) { healthy := &fakeUploadAPI{} rec2 := &fakeRecorder{} require.NoError(t, salvageSpools(ctx, healthy, "bucket", rec2, root, - func(d string) string { return "live-rec/" + d + "/" }, time.Hour)) + func(d string) string { return "live-rec/" + d + "/" }, time.Hour, time.Now())) healthy.mu.Lock() defer healthy.mu.Unlock() @@ -56,6 +57,35 @@ func TestSalvageSpools(t *testing.T) { require.True(t, os.IsNotExist(err), "empty did dir must be removed") } +// TestSalvageSpoolsSkipsLiveSessions proves the boot-time cutoff: a session +// dir whose UnixNano name is at/after startedBefore belongs to a live session +// of this process and must not be touched — that's what makes salvage safe to +// run concurrently with new streams at startup. +func TestSalvageSpoolsSkipsLiveSessions(t *testing.T) { + root := filepath.Join(t.TempDir(), "live-rec-spool") + did := "did:plc:live" + cutoff := time.Now() + liveDir := filepath.Join(root, did, fmt.Sprintf("%d", time.Now().Add(time.Second).UnixNano())) + + spool, err := OpenSpool(liveDir, 0) + require.NoError(t, err) + ctx := context.Background() + _, err = spool.Append(ctx, make([]byte, 1024), "at://A") + require.NoError(t, err) + + healthy := &fakeUploadAPI{} + rec := &fakeRecorder{} + require.NoError(t, salvageSpools(ctx, healthy, "bucket", rec, root, + func(d string) string { return "live-rec/" + d + "/" }, time.Hour, cutoff)) + + healthy.mu.Lock() + defer healthy.mu.Unlock() + require.Equal(t, 0, healthy.creates, "a live session's spool must not be salvaged") + s2, err := OpenSpool(liveDir, 0) + require.NoError(t, err) + require.Equal(t, 1, s2.Len(), "the live session's segments must be untouched") +} + // TestSalvageSpoolsPartialFailure proves a mid-salvage failure keeps the // un-uploaded remainder on disk for the next attempt, without re-uploading // what already completed. @@ -76,7 +106,7 @@ func TestSalvageSpoolsPartialFailure(t *testing.T) { flaky := &fakeUploadAPI{completesBeforeFail: 1, failCompletes: 1 << 30} rec := &fakeRecorder{} require.NoError(t, salvageSpools(ctx, flaky, "bucket", rec, root, - func(d string) string { return "live-rec/" + d + "/" }, time.Hour)) + func(d string) string { return "live-rec/" + d + "/" }, time.Hour, time.Now())) // Object A's segment was acked away; B's remains for the next run. s2, err := OpenSpool(dir, 0) @@ -88,7 +118,7 @@ func TestSalvageSpoolsPartialFailure(t *testing.T) { healthy := &fakeUploadAPI{} rec3 := &fakeRecorder{} require.NoError(t, salvageSpools(ctx, healthy, "bucket", rec3, root, - func(d string) string { return "live-rec/" + d + "/" }, time.Hour)) + func(d string) string { return "live-rec/" + d + "/" }, time.Hour, time.Now())) healthy.mu.Lock() defer healthy.mu.Unlock() require.Equal(t, 1, healthy.completes) diff --git a/pkg/s3/spool.go b/pkg/s3/spool.go index 1a5095ee..74bff8f9 100644 --- a/pkg/s3/spool.go +++ b/pkg/s3/spool.go @@ -73,6 +73,13 @@ func OpenSpool(dir string, maxBytes int64) (*Spool, error) { } for _, e := range entries { name := e.Name() + if strings.HasSuffix(name, ".tmp") { + // A segment write torn by a crash; Append never acknowledged it, + // so it was never owed durability. Remove rather than salvage + // possibly-partial MUXL bytes into a recording. + _ = os.Remove(filepath.Join(dir, name)) + continue + } if !strings.HasSuffix(name, ".seg") { continue } @@ -114,9 +121,30 @@ func (s *Spool) Append(ctx context.Context, data []byte, uri string) (int64, err s.mu.Lock() defer s.mu.Unlock() seq := s.nextSeq - if err := os.WriteFile(filepath.Join(s.dir, spoolSegFile(seq)), data, 0o644); err != nil { + // Write-tmp + fsync + rename: after a host crash (not just a process + // crash) the segment file is either complete or absent — never torn bytes + // that salvage would splice into a recording. OpenSpool sweeps orphaned + // .tmp files. + final := filepath.Join(s.dir, spoolSegFile(seq)) + tmp := final + ".tmp" + f, err := os.OpenFile(tmp, os.O_CREATE|os.O_TRUNC|os.O_WRONLY, 0o644) + if err != nil { + return 0, fmt.Errorf("creating spool segment %d: %w", seq, err) + } + if _, err := f.Write(data); err != nil { + _ = f.Close() return 0, fmt.Errorf("writing spool segment %d: %w", seq, err) } + if err := f.Sync(); err != nil { + _ = f.Close() + return 0, fmt.Errorf("syncing spool segment %d: %w", seq, err) + } + if err := f.Close(); err != nil { + return 0, fmt.Errorf("closing spool segment %d: %w", seq, err) + } + if err := os.Rename(tmp, final); err != nil { + return 0, fmt.Errorf("renaming spool segment %d: %w", seq, err) + } s.nextSeq++ s.seqs = append(s.seqs, seq) s.sizes[seq] = int64(len(data)) @@ -129,6 +157,7 @@ func (s *Spool) Append(ctx context.Context, data []byte, uri string) (int64, err f, err := os.OpenFile(filepath.Join(s.dir, spoolMetaFile), os.O_APPEND|os.O_CREATE|os.O_WRONLY, 0o644) if err == nil { _, _ = f.Write(append(line, '\n')) + _ = f.Sync() // URI changes are rare; losing one misattributes segments _ = f.Close() } else { log.Error(ctx, "writing spool meta", "dir", s.dir, "error", err)