diff --git a/pkg/config/config.go b/pkg/config/config.go index 1ba3931b4..107157e3f 100644 --- a/pkg/config/config.go +++ b/pkg/config/config.go @@ -1440,6 +1440,17 @@ func (cli *CLI) S3Config() s3.Config { } } +// SetS3Config applies an s3.Config to the CLI's S3 fields — the inverse of +// S3Config, for processes (ingest workers) that receive the S3 destination over +// a handshake instead of from flags. +func (cli *CLI) SetS3Config(c s3.Config) { + cli.S3Endpoint = c.Endpoint + cli.S3Bucket = c.Bucket + cli.S3AccessKeyID = c.AccessKeyID + cli.S3SecretAccessKey = c.SecretAccessKey + cli.S3Region = c.Region +} + // DebugRecordingFile is the write target returned by DebugRecordingCreate: an // *os.File on local disk, or an S3 upload that commits on Close. Name() reports // the destination (path or object key) for logging. @@ -1458,7 +1469,11 @@ type DebugRecordingFile interface { func (cli *CLI) DebugRecordingCreate(ctx context.Context, fpath []string, contentType string, overwrite bool) (DebugRecordingFile, error) { if cli.S3Configured() { key := strings.Join(fpath, "/") - return s3.NewUploadWriter(ctx, s3.NewClient(cli.S3Config()), cli.S3Bucket, key, contentType) + // The recording outlives the ingest session's ctx: Close commits the upload + // during teardown, after that ctx is typically cancelled — a cancelled ctx + // here would abort the upload and lose the object. Callers bound the commit + // with their own finalize waits instead. + return s3.NewUploadWriter(context.WithoutCancel(ctx), s3.NewClient(cli.S3Config()), cli.S3Bucket, key, contentType) } return cli.DataFileCreate(fpath, overwrite) } diff --git a/pkg/media/ingest_supervisor.go b/pkg/media/ingest_supervisor.go index 5f09a9849..f268d08f2 100644 --- a/pkg/media/ingest_supervisor.go +++ b/pkg/media/ingest_supervisor.go @@ -278,6 +278,14 @@ func (mm *MediaManager) buildWorkerConfig(ctx context.Context, ms MediaSigner) ( } else if rec { cfg.Record = true cfg.DataDir = mm.cli.DataDir + // The worker writes the recording, so it needs main's S3 destination too — + // without it, DebugRecordingCreate inside the worker would silently fall + // back to local disk under DataDir. Only sent when recording, to keep the + // S3 secret out of handshakes that don't need it. + if mm.cli.S3Configured() { + s3cfg := mm.cli.S3Config() + cfg.S3 = &s3cfg + } } // Node transcode signer lets the worker complete to dual-codec itself. If it's // unavailable, the worker emits single-codec (the node doesn't re-transcode the diff --git a/pkg/media/ingest_worker.go b/pkg/media/ingest_worker.go index 8221738a0..7125a12e0 100644 --- a/pkg/media/ingest_worker.go +++ b/pkg/media/ingest_worker.go @@ -13,6 +13,7 @@ import ( "stream.place/streamplace/pkg/gstinit" "stream.place/streamplace/pkg/log" "stream.place/streamplace/pkg/muxl" + "stream.place/streamplace/pkg/s3" ) // manifestHolder holds the worker's current C2PA manifest. It starts as the @@ -90,14 +91,18 @@ type IngestWorkerConfig struct { Chunked bool `json:"chunked,omitempty"` // Record, when true, makes the worker write a debug recording of this session - // (the MKV/RTMP push body, or the WHIP session) under - // DataDir/debug-recordings//. main evaluates the per-stream DebugRecording - // setting (which needs the DB) and the worker carries it out — so debug - // recording keeps working on the isolated paths without main being in the data - // path, and a recording even survives a main restart. DataDir is set (only when - // Record) to the node data dir the worker writes recordings under. - Record bool `json:"record,omitempty"` - DataDir string `json:"data_dir,omitempty"` + // (the MKV/RTMP push body, or the WHIP session). main evaluates the per-stream + // DebugRecording setting (which needs the DB) and the worker carries it out — + // so debug recording keeps working on the isolated paths without main being in + // the data path, and a recording even survives a main restart. The recording + // streams to S3 under debug-recordings// when S3 is set (production), and + // falls back to DataDir/debug-recordings// on local disk otherwise (dev). + // DataDir and S3 are set (only when Record) from main's config; S3 carries the + // secret key, which is fine here — the handshake exists to carry key material + // off argv/env. + Record bool `json:"record,omitempty"` + DataDir string `json:"data_dir,omitempty"` + S3 *s3.Config `json:"s3,omitempty"` // Transport selects the worker's ingest source: "" / "mkv" reads MKV media // (stdin or InputFD); "whip" makes the worker own the WebRTC PeerConnection, @@ -112,6 +117,18 @@ type IngestWorkerConfig struct { // IngestTransportWHIP is the cfg.Transport value selecting the WHIP worker. const IngestTransportWHIP = "whip" +// workerCLI assembles the minimal config.CLI a worker runs with: the +// broadcaster identity plus the debug-recording destination (S3 when main +// handed its config over the handshake, else local disk under DataDir). Shared +// by the MKV and WHIP workers so both record to the same place main would. +func (cfg IngestWorkerConfig) workerCLI() *config.CLI { + cli := &config.CLI{BroadcasterHost: cfg.BroadcasterHost, DataDir: cfg.DataDir} + if cfg.S3 != nil { + cli.SetS3Config(*cfg.S3) + } + return cli +} + // WorkerInput reconstructs the raw media stream the gst pipeline reads from the // fd-passed push connection: prepend any bytes main already read past the headers // (Prebuf), then de-chunk if the push used chunked transfer-encoding. For stdin @@ -211,23 +228,22 @@ func RunMKVIngestWorker(ctx context.Context, cfg IngestWorkerConfig, stdin io.Re // Minimal manager: the broadcaster identity the transcode completion // (finishTranscodedSegment) stamps into the node-signed AAC track, plus the - // data dir for an optional debug recording. - mm := &MediaManager{cli: &config.CLI{BroadcasterHost: cfg.BroadcasterHost, DataDir: cfg.DataDir}} + // destination for an optional debug recording. + mm := &MediaManager{cli: cfg.workerCLI()} onSegment, flush := mm.workerSegmentSink(ctx, cfg, frames) - // Debug recording: tee the ingest media to a file before it reaches gst. main - // decided this (cfg.Record) and handed us DataDir; recording here keeps main + // Debug recording: tee the ingest media before it reaches gst. main decided + // this (cfg.Record) and handed us the destination; recording here keeps main // out of the data path and lets the recording survive a main restart. media := stdin if cfg.Record { log.Log(ctx, "recording ingest media to file", "streamer", cfg.StreamerDID) - pr, pw := io.Pipe() - media = io.TeeReader(stdin, pw) - go func() { - if derr := mm.dumpToFile(ctx, pr, cfg.StreamerDID, ".rtmp.mkv"); derr != nil { - log.Error(ctx, "ingest worker: dump recording to file", "error", derr) - } - }() + var finalize func() + media, finalize = mm.recordTee(ctx, stdin, cfg.StreamerDID, ".rtmp.mkv") + // Registered before the pipeline's SetState(Null) defer, so it runs after + // the pipeline stops reading — and before this worker process exits, which + // would otherwise strand an uncommitted S3 upload. + defer finalize() } signerElem, done, err := muxlSignSegmentElem(ctx, mm.cli, workerSignStream(cfg, getManifest), onSegment) diff --git a/pkg/media/ingest_worker_test.go b/pkg/media/ingest_worker_test.go index 2615a6315..74336d8f1 100644 --- a/pkg/media/ingest_worker_test.go +++ b/pkg/media/ingest_worker_test.go @@ -6,9 +6,12 @@ import ( "errors" "fmt" "io" + "net/http" + "net/http/httptest" "os" "path/filepath" "strings" + "sync" "testing" "time" @@ -19,6 +22,7 @@ import ( "stream.place/streamplace/pkg/gstinit" "stream.place/streamplace/pkg/ingestframe" "stream.place/streamplace/pkg/muxl" + "stream.place/streamplace/pkg/s3" ) // TestWorkerInputDeframes checks the body-deframing the worker applies to the @@ -219,6 +223,115 @@ func TestRunMKVIngestWorkerRecords(t *testing.T) { }, 10*time.Second, 25*time.Millisecond, "worker records the ingest media verbatim") } +// fakeS3Server is a minimal path-style S3 endpoint speaking just enough of the +// multipart-upload protocol for UploadWriter: initiate → upload parts → +// complete. Completed objects land in objects keyed by "/". +type fakeS3Server struct { + mu sync.Mutex + parts map[string][]byte // "#" → body + objects map[string][]byte // completed "/" → body +} + +func newFakeS3Server() *fakeS3Server { + return &fakeS3Server{parts: map[string][]byte{}, objects: map[string][]byte{}} +} + +func (f *fakeS3Server) ServeHTTP(w http.ResponseWriter, r *http.Request) { + f.mu.Lock() + defer f.mu.Unlock() + path := strings.TrimPrefix(r.URL.Path, "/") + q := r.URL.Query() + switch { + case r.Method == "POST" && q.Has("uploads"): + fmt.Fprintf(w, `test-upload`) + case r.Method == "PUT" && q.Has("partNumber"): + body, _ := io.ReadAll(r.Body) + f.parts[path+"#"+q.Get("partNumber")] = body + w.Header().Set("ETag", `"part-`+q.Get("partNumber")+`"`) + case r.Method == "POST" && q.Has("uploadId"): + var buf []byte + for i := 1; ; i++ { + part, ok := f.parts[fmt.Sprintf("%s#%d", path, i)] + if !ok { + break + } + buf = append(buf, part...) + } + f.objects[path] = buf + fmt.Fprintf(w, `%s`, path) + default: + w.WriteHeader(http.StatusBadRequest) + } +} + +// objectWithPrefix finds a completed object whose key starts with prefix and +// ends with suffix (the recording's timestamped filename isn't predictable). +func (f *fakeS3Server) objectWithPrefix(prefix, suffix string) ([]byte, bool) { + f.mu.Lock() + defer f.mu.Unlock() + for path, b := range f.objects { + if strings.HasPrefix(path, prefix) && strings.HasSuffix(path, suffix) { + return b, true + } + } + return nil, false +} + +// TestRunMKVIngestWorkerRecordsToS3 proves the debug recording streams to S3 +// when main hands its S3 config over the handshake (cfg.S3) — the production +// shape. Without that plumbing the worker's minimal CLI has no S3 fields and +// DebugRecordingCreate silently falls back to local disk, which is exactly the +// regression this guards against: recordings must land in the bucket, not under +// DataDir. +func TestRunMKVIngestWorkerRecordsToS3(t *testing.T) { + ctx := context.Background() + ms := newBareSegmentSigner(t) + + keyPEM, err := signers.MarshalES256KPrivateKeyPEM(ms.Signer) + require.NoError(t, err) + manifest, err := ms.buildManifest(ctx, time.Now().UnixMilli()) + require.NoError(t, err) + + fake := newFakeS3Server() + srv := httptest.NewServer(fake) + defer srv.Close() + + dataDir := t.TempDir() + cfg := IngestWorkerConfig{ + StreamerDID: ms.Streamer(), + KeyPEM: keyPEM, + CertPEM: ms.Cert, + Manifest: manifest, + BroadcasterHost: "test.example.com", + Record: true, + DataDir: dataDir, + S3: &s3.Config{ + Endpoint: srv.URL, + Bucket: "debug-bucket", + AccessKeyID: "test-access", + SecretAccessKey: "test-secret", + Region: "auto", + }, + } + + mkv := makeH264AACMKV(t, ctx, getFixture("5sec.mp4")) + + require.NoError(t, RunMKVIngestWorker(ctx, cfg, bytes.NewReader(mkv), ingestframe.NewWriter(io.Discard), func() []byte { return cfg.Manifest })) + + // The upload commits asynchronously (the dump goroutine's Close); the object + // must appear at debug-bucket/debug-recordings//.rtmp.mkv holding + // exactly the ingested media. + wantPrefix := "debug-bucket/debug-recordings/" + ms.Streamer() + "/" + require.Eventually(t, func() bool { + got, ok := fake.objectWithPrefix(wantPrefix, ".rtmp.mkv") + return ok && bytes.Equal(got, mkv) + }, 10*time.Second, 25*time.Millisecond, "worker streams the recording to the S3 bucket verbatim") + + // And nothing fell back to local disk. + matches, _ := filepath.Glob(filepath.Join(dataDir, "debug-recordings", "*", "*")) + require.Empty(t, matches, "recording must go to S3, not DataDir") +} + // TestRunMKVIngestWorkerSelfWatchdog proves the worker's OWN watchdog contains a // wedge. This is the only wedge containment on the detached/WHIP paths, where // main can't kill a detached worker — so the worker has to notice it's stuck and diff --git a/pkg/media/mkv_ingest.go b/pkg/media/mkv_ingest.go index 64182fcf4..2e9451a5d 100644 --- a/pkg/media/mkv_ingest.go +++ b/pkg/media/mkv_ingest.go @@ -22,14 +22,9 @@ func (mm *MediaManager) MKVIngest(ctx context.Context, input io.Reader, ms Media } if shouldRecord { log.Log(ctx, "recording RTMP stream to file", "streamer", ms.Streamer()) - pr, pw := io.Pipe() - input = io.TeeReader(input, pw) - go func() { - err := mm.dumpToFile(ctx, pr, ms.Streamer(), ".rtmp.mkv") - if err != nil { - log.Error(ctx, "error dumping to file", "error", err) - } - }() + var finalize func() + input, finalize = mm.recordTee(ctx, input, ms.Streamer(), ".rtmp.mkv") + defer finalize() } else { log.Log(ctx, "not recording RTMP stream to file", "streamer", ms.Streamer()) } @@ -127,6 +122,38 @@ func buildMKVIngestPipeline(ctx context.Context, input io.Reader, signerElem *gs return pipeline, nil } +// debugRecordingFlushTimeout bounds how long ingest teardown waits for a debug +// recording to finalize — for S3 the commit only happens at Close, so an +// unbounded wait could wedge teardown while an unwaited exit loses the object. +const debugRecordingFlushTimeout = 30 * time.Second + +// recordTee wires up a debug recording: everything read through the returned +// reader is teed into an asynchronous dumpToFile. The returned finalize ends +// the dump (closing the tee's pipe — the dump's io.Copy never sees EOF +// otherwise, since a TeeReader doesn't propagate one) and waits, bounded, for +// it to commit. Callers MUST finalize after ingest ends: on the S3 path the +// object only exists once Close commits the upload, so skipping it (e.g. a +// worker process exiting) silently loses the recording. +func (mm *MediaManager) recordTee(ctx context.Context, r io.Reader, user string, filesuffix string) (io.Reader, func()) { + pr, pw := io.Pipe() + done := make(chan struct{}) + go func() { + defer close(done) + if err := mm.dumpToFile(ctx, pr, user, filesuffix); err != nil { + log.Error(ctx, "error dumping to file", "error", err, "streamer", user) + } + }() + finalize := func() { + pw.Close() + select { + case <-done: + case <-time.After(debugRecordingFlushTimeout): + log.Error(ctx, "debug recording did not finalize in time", "streamer", user) + } + } + return io.TeeReader(r, pw), finalize +} + func (mm *MediaManager) dumpToFile(ctx context.Context, r io.Reader, user string, filesuffix string) error { now := aqtime.FromTime(time.Now()) filename := fmt.Sprintf("%s%s", now.FileSafeString(), filesuffix) diff --git a/pkg/media/whip_worker.go b/pkg/media/whip_worker.go index de598bcbe..61a98e080 100644 --- a/pkg/media/whip_worker.go +++ b/pkg/media/whip_worker.go @@ -7,7 +7,6 @@ import ( "os" "github.com/pion/webrtc/v4" - "stream.place/streamplace/pkg/config" "stream.place/streamplace/pkg/gstinit" "stream.place/streamplace/pkg/log" "stream.place/streamplace/pkg/rtcrec" @@ -59,11 +58,11 @@ func ServeWHIPIngestWorkerSocket(ctx context.Context, cfg IngestWorkerConfig) er return runErr } - mm := &MediaManager{cli: &config.CLI{BroadcasterHost: cfg.BroadcasterHost, DataDir: cfg.DataDir}} + mm := &MediaManager{cli: cfg.workerCLI()} // The worker owns the PeerConnection (its own UDP sockets), built with the // same codec/interceptor setup as the in-process server. Debug recording is - // decided by main (cfg.Record) and written by the worker under cfg.DataDir. + // decided by main (cfg.Record) and written by the worker (S3 or cfg.DataDir). api, webrtcConfig, err := newWebRTCAPI() if err != nil { return finish(fmt.Errorf("webrtc api: %w", err)) @@ -110,5 +109,10 @@ func ServeWHIPIngestWorkerSocket(ctx context.Context, cfg IngestWorkerConfig) er cancel() <-signerDone flush() + // The recording commits asynchronously after pc.Close (drain sleep + S3 + // commit); wait for it, or this process exits and the object never appears. + if rpc, ok := pc.(*rtcrec.RecordingPeerConnection); ok { + rpc.FinalizeRecording(ctx) + } return finish(streamErr) } diff --git a/pkg/rtcrec/recording_peerconnection.go b/pkg/rtcrec/recording_peerconnection.go index 29de5eaca..34a8347d2 100644 --- a/pkg/rtcrec/recording_peerconnection.go +++ b/pkg/rtcrec/recording_peerconnection.go @@ -3,6 +3,7 @@ package rtcrec import ( "context" "fmt" + "sync" "time" "github.com/pion/rtcp" @@ -13,10 +14,12 @@ import ( ) type RecordingPeerConnection struct { - enabled bool - pionpc *webrtc.PeerConnection - file config.DebugRecordingFile - stream *RecorderStream + enabled bool + pionpc *webrtc.PeerConnection + file config.DebugRecordingFile + stream *RecorderStream + closeOnce sync.Once + recDone chan struct{} // closed once the recording file/upload is committed } func NewRecordingPeerConnection(ctx context.Context, cli config.CLI, user string, pionpc *webrtc.PeerConnection, enabled bool) (PeerConnection, error) { @@ -43,6 +46,7 @@ func NewRecordingPeerConnection(ctx context.Context, cli config.CLI, user string file: f, stream: stream, enabled: enabled, + recDone: make(chan struct{}), }, nil } @@ -53,14 +57,43 @@ func (pc *RecordingPeerConnection) Do(f func()) { } func (pc *RecordingPeerConnection) Close() error { - pc.Do(func() { + pc.Do(pc.finishRecording) + return pc.pionpc.Close() +} + +// finishRecording drains stragglers, commits the recording (for S3, Close IS +// the commit), and signals recDone. Idempotent — Close on the disconnect path +// and FinalizeRecording at worker exit can both trigger it. +func (pc *RecordingPeerConnection) finishRecording() { + pc.closeOnce.Do(func() { // This is sloppy but there might be other goroutines still writing so let's chill for a sec time.Sleep(10 * time.Second) pc.file.Close() + close(pc.recDone) }) - return pc.pionpc.Close() } +// FinalizeRecording blocks until the debug recording is committed (bounded). +// Call it before process exit on paths like the WHIP ingest worker: Close only +// *starts* the drain+commit on a goroutine, and a process that exits first +// strands an uncommitted S3 upload — the object never appears. No-op when not +// recording. +func (pc *RecordingPeerConnection) FinalizeRecording(ctx context.Context) { + if !pc.enabled { + return + } + go pc.finishRecording() // in case nothing called Close (e.g. pipeline error) + select { + case <-pc.recDone: + case <-time.After(recordingFinalizeTimeout): + log.Error(ctx, "debug recording did not finalize in time", "file", pc.file.Name()) + } +} + +// recordingFinalizeTimeout bounds FinalizeRecording: the 10s straggler drain in +// finishRecording plus generous headroom for the S3 commit. +const recordingFinalizeTimeout = 40 * time.Second + func (pc *RecordingPeerConnection) CreateAnswer(options *webrtc.AnswerOptions) (webrtc.SessionDescription, error) { now := time.Now() ret, err := pc.pionpc.CreateAnswer(options) diff --git a/pkg/s3/s3.go b/pkg/s3/s3.go index 2351f97c1..6bad8e1a5 100644 --- a/pkg/s3/s3.go +++ b/pkg/s3/s3.go @@ -15,13 +15,15 @@ import ( "stream.place/streamplace/pkg/log" ) -// Config holds the configuration for an S3-compatible upload target. +// Config holds the configuration for an S3-compatible upload target. The json +// tags exist because it rides the ingest-worker startup handshake (a dedicated +// pipe fd, never argv/env — it carries the secret key). type Config struct { - Endpoint string - Bucket string - AccessKeyID string - SecretAccessKey string - Region string + Endpoint string `json:"endpoint,omitempty"` + Bucket string `json:"bucket,omitempty"` + AccessKeyID string `json:"access_key_id,omitempty"` + SecretAccessKey string `json:"secret_access_key,omitempty"` + Region string `json:"region,omitempty"` } // Recorder is an optional persistence hook for S3Uploader. RecordStart is