From efebb48cf21c6fa76877a0e7dcb592269ff27fee Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Thu, 1 Oct 2026 22:24:39 -0700 Subject: [PATCH] Send live segments without waiting for captions; record them synced Holding each segment for SP_CAPTIONS_MASTER_DELAY traded latency for caption timing: 7 s gave well-placed captions but 7 s of delay, while the 1.5 s default put most words in later segments. Live and recorded segments now get separate layouts. - Live: the signer no longer waits for speech recognition, only (at most 0.5 s) for the ingest caption tap. Late words keep playing in order from the next GoP, as before. - Recorded: an archive pass in the caption master lays each signed GoP out again once recognition covers it, bounded by SP_CAPTIONS_MASTER_DELAY (now 10 s, recording-only). When that layout differs from the live one it signs fresh text runs with the streamer's key (muxl SignTextRuns, same span and dc:date). The director swaps them into the completed segment before the S3 upload; audio, video and the node's AAC run stay byte-identical. S3 operations keep arrival order while copies wait in parallel. - Isolated workers sign their archive runs and send them to main as Captions frames before End, so resumed workers need no key in main. In-process, the archive rides the ingest context, so segments validated after the signer returns still find it. - Entries are keyed by the GoP's media-timeline start: after a stall, re-anchored GoPs can share a signed start millisecond. - A real-speech run (2 minutes, real time, bundled models) placed 54/54 phrases in the GoP where they were spoken, once three fixes it prompted were in: on-time cues no longer queue behind a late one in the archive; recognition coverage stops short of the window's last 2 s unless the speaker paused; and a cue timed at its predecessor's start follows it instead of erasing it. - Audio completion inserts its AAC run in ascending track-ID order instead of after any text runs; the VOD indexer rejects that order. - Every declared text track now appears in every GoP, so a recorded copy never drops one. --- go.mod | 2 +- go.sum | 2 + .../src/content/docs/features-dev/captions.md | 20 +- .../docs/guides/installing/captions.md | 33 +- pkg/captions/recognizer.go | 7 +- pkg/captions/recognizer_test.go | 8 +- pkg/config/config.go | 4 +- pkg/director/s3_upload.go | 32 +- pkg/director/s3_upload_test.go | 38 ++ pkg/director/stream_session.go | 1 + pkg/ingestframe/frame.go | 10 + pkg/media/captions_archive.go | 202 +++++++++ pkg/media/captions_archive_test.go | 31 ++ pkg/media/captions_master.go | 408 +++++++++++++----- pkg/media/captions_master_control_test.go | 2 +- pkg/media/captions_master_feed.go | 5 +- pkg/media/captions_master_socket.go | 14 +- pkg/media/captions_master_test.go | 202 ++++++--- pkg/media/captions_transcode_test.go | 6 +- pkg/media/frame_server.go | 12 +- pkg/media/ingest_daemon.go | 14 +- pkg/media/ingest_daemon_test.go | 2 +- pkg/media/ingest_supervisor.go | 34 +- pkg/media/ingest_worker.go | 22 +- pkg/media/ingest_worker_test.go | 19 +- pkg/media/media.go | 3 + pkg/media/media_signer.go | 8 + pkg/media/segmenter.go | 5 + pkg/media/transcode.go | 17 +- pkg/media/transcode_stream.go | 14 +- pkg/media/validate.go | 3 + pkg/media/whip_worker.go | 5 +- pkg/muxl/muxl.go | 11 + 33 files changed, 956 insertions(+), 240 deletions(-) create mode 100644 pkg/director/s3_upload_test.go create mode 100644 pkg/media/captions_archive.go create mode 100644 pkg/media/captions_archive_test.go diff --git a/go.mod b/go.mod index 8bda888de..280c517ed 100644 --- a/go.mod +++ b/go.mod @@ -73,7 +73,7 @@ require ( github.com/streamplace/atmoq/go v0.0.4-0.20260701223355-13757de4ae08 github.com/streamplace/atproto-oauth-golang v0.0.0-20260413212710-98956064d06c github.com/streamplace/glex v0.0.0-20260820164827-814f46540f22 - github.com/streamplace/muxl/go v0.3.6-0.20261001171246-1b33b3629baa + github.com/streamplace/muxl/go v0.3.6-0.20261002024311-09684073c273 github.com/streamplace/oatproxy v0.0.0-20260710202406-60d97b9d780b github.com/stretchr/testify v1.11.1 github.com/tdewolff/canvas v0.0.0-20250728095813-50d4cb1eee71 diff --git a/go.sum b/go.sum index 438cac16e..4f3730888 100644 --- a/go.sum +++ b/go.sum @@ -1399,6 +1399,8 @@ github.com/streamplace/indigo v0.0.0-20260218231908-939cdaf0c507 h1:e8M3qPLr37Nx github.com/streamplace/indigo v0.0.0-20260218231908-939cdaf0c507/go.mod h1:Pm2I1+iDXn/hLbF7XCg/DsZi6uDCiOo7hZGWprSM7k0= github.com/streamplace/muxl/go v0.3.6-0.20261001171246-1b33b3629baa h1:J+ii89pziuLu8XCQ3Nm6AXt7WsBqysGwxQGF1+V00lU= github.com/streamplace/muxl/go v0.3.6-0.20261001171246-1b33b3629baa/go.mod h1:aCyYTW3o6c1Kush9UJ/Yv6EYMUbj8l8GTD7cHKcSxw8= +github.com/streamplace/muxl/go v0.3.6-0.20261002024311-09684073c273 h1:Tja2cGAcEEI60ihDh3e4TLFntsNcPPAbztLS94rqrn0= +github.com/streamplace/muxl/go v0.3.6-0.20261002024311-09684073c273/go.mod h1:aCyYTW3o6c1Kush9UJ/Yv6EYMUbj8l8GTD7cHKcSxw8= github.com/streamplace/oatproxy v0.0.0-20260710202406-60d97b9d780b h1:eWbwCtBbMyrDTHLYIold07OR2hmvzXsbAUxi57ElMLk= github.com/streamplace/oatproxy v0.0.0-20260710202406-60d97b9d780b/go.mod h1:wpY+T/wE00jrUhgh2dKXbbE91D36u86KGlENK/hWFkE= github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= diff --git a/js/docs/src/content/docs/features-dev/captions.md b/js/docs/src/content/docs/features-dev/captions.md index fe31fb5cd..26c78a4b6 100644 --- a/js/docs/src/content/docs/features-dev/captions.md +++ b/js/docs/src/content/docs/features-dev/captions.md @@ -24,17 +24,29 @@ The MUXL text-track metadata convention is: an empty or unrecognized label as `ingest`; it is not a free-form display name. Tracks are declared when their first cues arrive and retain stable IDs and -immutable configuration during the ingest session. A track replaced by another -source continues with empty text segments rather than changing its metadata. +immutable configuration during the ingest session. A declared track appears in +every later GoP; a track replaced by another source continues with empty text +segments rather than changing its metadata. Cue pieces crossing GoP boundaries retain their identity and are clipped to the corresponding GoP ranges. A track's cues never overlap: a cue that arrives after its GoP was signed starts in the first unsigned GoP, after the previous cue has had its own duration, and keeps its whole duration; the next cue on the track ends it. Automatic captions stay up to 3 s past their last word until replaced, so the track reads continuously between recognized phrases. + +Live segments are never held for speech recognition, so their late words take +that path. The **recorded** copy of each origin segment (the S3 live recording +and the VODs finalized from it) is laid out a second time, once recognition +covers the GoP or `--captions-master-delay` passes. It is the live segment with +only its text runs replaced by runs signed with the streamer's key for the same +GoP span and `dc:date`, in ascending track-ID order. Audio, video, and the +node's transcoded audio run are byte-identical, so their signatures and the +transcode's source binding still verify. Isolated ingest workers send these +text runs to the main process as `ingestframe.Captions` frames before `End`. Origin text tracks start at reserved numeric ID 100, above node-added AV -renditions. Continuous audio completion decodes only AV tracks; late text -declarations remain in signed source bytes and completed archival segments. +renditions. Continuous audio completion decodes only AV tracks and inserts its +audio run before any text tracks; late text declarations remain in signed +source bytes and completed archival segments. Every policy stamps GoPs on the media clock, anchored to the first fragment's arrival, not signing time; drift over one second reanchors the next GoP. Pushed wall-clock cues use the inverse mapping. diff --git a/js/docs/src/content/docs/guides/installing/captions.md b/js/docs/src/content/docs/guides/installing/captions.md index f68c73f35..0217ee766 100644 --- a/js/docs/src/content/docs/guides/installing/captions.md +++ b/js/docs/src/content/docs/guides/installing/captions.md @@ -22,14 +22,14 @@ The equivalent CLI flags are listed below. | `SP_CAPTIONS` | `--captions` | `true` | Enables optional node-generated sidecars when the streamer's policy allows them. It does not disable an origin streamer's requested canonical automatic captions, or pass-through of existing captions. | | `SP_CAPTIONS_CPU_BUDGET` | `--captions-cpu-budget` | `0.5` | Fraction of logical CPUs available to the speech-recognition engine across streams. Must be greater than zero and at most one. The scheduler measures models and chooses the largest model that fits each stream's available budget. | | `SP_CAPTIONS_MODEL_DIR` | `--captions-model-dir` | Unset (no extra models) | Directory of additional `ggml-*.bin` Whisper models, offered alongside bundled models. Models are loaded and benchmarked; Silero-named files are not recognition models. | -| `SP_CAPTIONS_MASTER_DELAY` | `--captions-master-delay` | `1.5s` | Maximum time after a GoP closes that the origin may hold it waiting for captions before signing its canonical MUXL segment. Late final words move into the next unsigned GoP. | +| `SP_CAPTIONS_MASTER_DELAY` | `--captions-master-delay` | `10s` | Maximum time after a GoP closes that its **recorded** copy waits for speech recognition to cover it before the captions are laid out for the recording. Live segments are never held for captions. | For example: ```ini SP_CAPTIONS=true SP_CAPTIONS_CPU_BUDGET=0.5 -SP_CAPTIONS_MASTER_DELAY=1.5s +SP_CAPTIONS_MASTER_DELAY=10s # Optional extra models: # SP_CAPTIONS_MODEL_DIR=/var/lib/streamplace/whisper-models ``` @@ -87,17 +87,26 @@ two recognition passes agree on them, or once two seconds of audio follow them. With the bundled models a pass takes one to three seconds on a typical server, so captions usually trail speech by two to five seconds. -The signer waits until recognition covers the GoP end or until GoP closure -plus `SP_CAPTIONS_MASTER_DELAY`, whichever comes first. Buffering separates -that wait from ingest. A larger delay gives recognition more time to arrive -in the matching segment but increases origin latency **by up to that delay**; -a smaller delay releases media sooner but places more captions in a later -segment than the speech. Such late captions play in order from the first -unsigned segment, each for its full duration. This is not the total player -latency: encoding, segment length, network, and playback buffering also -contribute. See +Live segments are signed and sent as soon as the ingest tap has parsed them; +the origin never holds video or audio for speech recognition. A word recognized +after its segment was signed is shown live in the first unsigned segment, after +the previous caption has had its own reading time, so viewers see captions a few +seconds behind the speech. + +When the origin records the stream to S3, the recording gets captions in the +segment where they were spoken. A second, archival layout of each segment waits +until recognition covers it, or until GoP closure plus +`SP_CAPTIONS_MASTER_DELAY`, whichever comes first; with pushed captions, which +give no coverage signal, it waits the full delay. The recorded segment is the +live one with only its text runs replaced, re-signed with the streamer's key; its +audio and video bytes are unchanged. A larger delay therefore costs no live +latency, only memory for segments awaiting upload. Words that arrive after the +delay are recorded like late live captions. Isolated ingest workers send their +archival text to the main process before their stream ends. See [`captions_master.go`](https://github.com/streamplace/streamplace/blob/main/pkg/media/captions_master.go) -for the hold and late-cue behavior. +for both layouts and +[`captions_archive.go`](https://github.com/streamplace/streamplace/blob/main/pkg/media/captions_archive.go) +for the recorded copy. Sources wait for the first signed policy snapshot before recognition admission, so ingest-only/off never briefly starts automatic recognition. Lossless byte queues block at 32 MiB instead of dropping media; aborts discard retained bytes diff --git a/pkg/captions/recognizer.go b/pkg/captions/recognizer.go index d8fc6ae27..e7d99ef3b 100644 --- a/pkg/captions/recognizer.go +++ b/pkg/captions/recognizer.go @@ -415,8 +415,13 @@ func (r *Recognizer) maybePass(force bool) { r.publishInterim() } // Examining audio is not enough: the signer may consume only immutable - // finals, so never release it past a still-open cue or uncommitted word. + // finals, so never release it past a still-open cue or uncommitted word, + // nor, while the speaker is still talking, into the window's last + // recognizerSettled, where whisper may not have heard a word yet. finalized := winEnd + if !paused && !force { + finalized = settled + } if cue, ok := r.grouper.Current(r.prev); ok && cue.Start.Before(finalized) { finalized = cue.Start } diff --git a/pkg/captions/recognizer_test.go b/pkg/captions/recognizer_test.go index 7de3e7a3c..46a7360a2 100644 --- a/pkg/captions/recognizer_test.go +++ b/pkg/captions/recognizer_test.go @@ -515,7 +515,7 @@ func TestWrapLines(t *testing.T) { require.Equal(t, "", WrapLines(" ", 37, 2)) } -func TestRecognizerCoverageFollowsDecisionsAtMonotonicWindowEnd(t *testing.T) { +func TestRecognizerCoverageIsSettledAndMonotonic(t *testing.T) { model := &fakeModel{} var covered []time.Time start := time.UnixMilli(1000) @@ -531,7 +531,11 @@ func TestRecognizerCoverageFollowsDecisionsAtMonotonicWindowEnd(t *testing.T) { require.Empty(t, covered, "PCM receipt alone must not release a held GoP") r.Push(start.Add(time.Second), speech(time.Second)) r.settle() - require.Equal(t, []time.Time{start.Add(2 * time.Second)}, covered) + r.Push(start.Add(2*time.Second), speech(time.Second)) + r.settle() + // Mid-speech, a pass releases only audio that has recognizerSettled after + // it: whisper may not have heard a word at the window's end yet. + require.Equal(t, []time.Time{start, start.Add(time.Second)}, covered) // A pure-silence decision also covers its complete window. r.Push(start.Add(3*time.Second), silence(2*time.Second)) r.settle() diff --git a/pkg/config/config.go b/pkg/config/config.go index 7b2c45627..2f04bc739 100644 --- a/pkg/config/config.go +++ b/pkg/config/config.go @@ -1074,8 +1074,8 @@ func (cli *CLI) NewCommand(name string) *urfavecli.Command { }, &urfavecli.DurationFlag{ Name: "captions-master-delay", - Usage: "how long the origin may hold a segment waiting for its captions before mastering it into the canonical stream. Late words move to the next segment.", - Value: 1500 * time.Millisecond, + Usage: "how long the recorded (S3) copy of a segment waits for speech recognition to cover it before its captions are laid out. Live segments are never held: words recognized late are shown in the next live segment.", + Value: 10 * time.Second, Destination: &cli.CaptionsMasterDelay, Sources: urfavecli.EnvVars("SP_CAPTIONS_MASTER_DELAY"), }, diff --git a/pkg/director/s3_upload.go b/pkg/director/s3_upload.go index 4e158e0e3..1f5732154 100644 --- a/pkg/director/s3_upload.go +++ b/pkg/director/s3_upload.go @@ -103,11 +103,10 @@ func (ss *StreamSession) s3Upload(ctx context.Context, notif *media.NewSegmentNo if ss.s3Uploader == nil { return } - ss.Go(ctx, func() error { - // notif.Muxl is the bare canonical segment; it concatenates directly - // (the S3 uploader synthesizes one init and prepends it per object). - return ss.s3Uploader.AddSegment(ctx, notif.Muxl) - }) + // notif.ArchiveCopy is the bare canonical segment (with its captions laid + // out again for the recording, which can take seconds); it concatenates + // directly (the S3 uploader synthesizes one init and prepends it per object). + ss.s3InOrder(ctx, notif.ArchiveCopy, ss.s3Uploader.AddSegment) } // s3Cutover completes the current live-recording object so it's immediately @@ -119,11 +118,32 @@ func (ss *StreamSession) s3Cutover(ctx context.Context) { if ss.s3Uploader == nil { return } - ss.Go(ctx, func() error { + ss.s3InOrder(ctx, nil, func(ctx context.Context, _ []byte) error { return ss.s3Uploader.Cutover(ctx) }) } +// s3InOrder applies op to the recording after every earlier S3 operation of +// the session. prepare runs right away, so archive copies waiting on their +// captions overlap instead of adding up. Unlike a lane, nothing is skipped: a +// recording must not drop a slow segment. NewSegment, the only caller, runs on +// the director's dispatch goroutine, which orders s3Prev. +func (ss *StreamSession) s3InOrder(ctx context.Context, prepare func(context.Context) []byte, op func(context.Context, []byte) error) { + prev, done := ss.s3Prev, make(chan struct{}) + ss.s3Prev = done + ss.Go(ctx, func() error { + defer close(done) + var data []byte + if prepare != nil { + data = prepare(ctx) + } + if prev != nil { + <-prev + } + return op(ctx, data) + }) +} + func (ss *StreamSession) s3Close(ctx context.Context) { if ss.s3Uploader == nil { return diff --git a/pkg/director/s3_upload_test.go b/pkg/director/s3_upload_test.go new file mode 100644 index 000000000..5c9c60fd9 --- /dev/null +++ b/pkg/director/s3_upload_test.go @@ -0,0 +1,38 @@ +package director + +import ( + "context" + "sync" + "testing" + "time" + + "github.com/stretchr/testify/require" + "golang.org/x/sync/errgroup" +) + +// A recording applies uploads and cutovers in arrival order, even when an +// earlier segment's archive copy takes longer to prepare than later ones. +func TestS3OperationsKeepArrivalOrder(t *testing.T) { + ss := &StreamSession{started: make(chan struct{}), g: &errgroup.Group{}} + close(ss.started) + var mu sync.Mutex + var order []string + record := func(_ context.Context, data []byte) error { + mu.Lock() + defer mu.Unlock() + order = append(order, string(data)) + return nil + } + prepare := func(name string, delay time.Duration) func(context.Context) []byte { + return func(context.Context) []byte { + time.Sleep(delay) + return []byte(name) + } + } + ctx := context.Background() + ss.s3InOrder(ctx, prepare("first", 50*time.Millisecond), record) + ss.s3InOrder(ctx, nil, func(ctx context.Context, _ []byte) error { return record(ctx, []byte("cutover")) }) + ss.s3InOrder(ctx, prepare("third", 0), record) + require.NoError(t, ss.g.Wait()) + require.Equal(t, []string{"first", "cutover", "third"}, order) +} diff --git a/pkg/director/stream_session.go b/pkg/director/stream_session.go index a6f84baf7..09abfb15d 100644 --- a/pkg/director/stream_session.go +++ b/pkg/director/stream_session.go @@ -75,6 +75,7 @@ type StreamSession struct { lastLivestreamTime time.Time lastViewCountTime time.Time s3Uploader *s3.S3Uploader + s3Prev chan struct{} // closed when the latest S3 operation is done; see s3InOrder // localRole runs once this node takes up the ingest node's jobs for the // session (recording, multistream targets): on the first local segment, // whether that is the session's first segment or one that arrives after diff --git a/pkg/ingestframe/frame.go b/pkg/ingestframe/frame.go index c335185c5..68e0735d4 100644 --- a/pkg/ingestframe/frame.go +++ b/pkg/ingestframe/frame.go @@ -62,6 +62,11 @@ const ( // worker reconnecting — it has no model of its own to notice the change. // Payload: the manifest JSON. Manifest Type = 6 + // Captions carries one GoP's archival caption text runs (JSON) from the + // worker's caption master. Main swaps them into the recorded copy of that + // GoP, whose live copy went out before speech recognition settled. Every + // Captions frame precedes End. + Captions Type = 7 ) func (t Type) String() string { @@ -78,6 +83,8 @@ func (t Type) String() string { return "event" case Manifest: return "manifest" + case Captions: + return "captions" default: return fmt.Sprintf("unknown(%d)", uint8(t)) } @@ -144,6 +151,9 @@ func (fw *Writer) Event(payload []byte) error { return fw.WriteFrame(Event, payl // Manifest frames an updated C2PA manifest (main → worker). func (fw *Writer) Manifest(payload []byte) error { return fw.WriteFrame(Manifest, payload) } +// Captions frames one GoP's archival caption text runs (JSON payload). +func (fw *Writer) Captions(payload []byte) error { return fw.WriteFrame(Captions, payload) } + // Reader decodes frames from an underlying stream. The decoder buffers/reads // ahead, so a Reader OWNS its stream for the stream's lifetime — don't create a // second Reader on the same connection (it would lose the first's buffered diff --git a/pkg/media/captions_archive.go b/pkg/media/captions_archive.go new file mode 100644 index 000000000..2ab26b91e --- /dev/null +++ b/pkg/media/captions_archive.go @@ -0,0 +1,202 @@ +package media + +import ( + "context" + "fmt" + "strconv" + "sync" + "time" + + "stream.place/streamplace/pkg/log" + "stream.place/streamplace/pkg/muxl" +) + +// archiveText is one GoP's archival text runs, keyed by the start of the GoP's +// span on the media timeline (the TextRequest muxl made for it). Runs is nil +// when the live GoP's text already matches. +type archiveText struct { + StartMs uint64 `json:"startMs"` + Runs map[uint32][]byte `json:"runs,omitempty"` +} + +// captionArchiveSlack covers signing and delivery after the archive pass's +// hold, so a recording waits for a slow GoP rather than giving up early. +const captionArchiveSlack = 5 * time.Second + +// captionArchiveLimit bounds GoPs nothing claims, e.g. when the stream is not +// being recorded. +const captionArchiveLimit = 64 + +// captionArchive collects an ingest session's archival text runs until the +// recording claims them. Live segments go out as soon as they are signed; the +// recorded copy of each one waits here for captions laid out in the GoP where +// they were spoken. It rides the ingest context; see withCaptionArchive. +type captionArchive struct { + wait time.Duration + mu sync.Mutex + changed chan struct{} + texts map[uint64]archiveText + finished bool +} + +func newCaptionArchive(hold time.Duration) *captionArchive { + return &captionArchive{wait: hold + captionArchiveSlack, changed: make(chan struct{}), texts: make(map[uint64]archiveText)} +} + +func (a *captionArchive) put(t archiveText) error { + a.mu.Lock() + defer a.mu.Unlock() + a.texts[t.StartMs] = t + for len(a.texts) > captionArchiveLimit { + oldest := t.StartMs + for start := range a.texts { + oldest = min(oldest, start) + } + delete(a.texts, oldest) + } + a.signal() + return nil +} + +// finish records that the archive pass has handed over every GoP, so claims +// for anything else stop waiting. +func (a *captionArchive) finish() { + a.mu.Lock() + a.finished = true + a.signal() + a.mu.Unlock() +} + +func (a *captionArchive) signal() { close(a.changed); a.changed = make(chan struct{}) } + +// claim waits for the archive pass to report the GoP starting at startMs. It +// gives up once the pass has finished without it, or after the hold. +func (a *captionArchive) claim(ctx context.Context, startMs uint64) (archiveText, bool) { + timer := time.NewTimer(a.wait) + defer timer.Stop() + a.mu.Lock() + defer a.mu.Unlock() + for { + if t, ok := a.texts[startMs]; ok { + delete(a.texts, startMs) + return t, true + } + if a.finished { + return archiveText{}, false + } + changed := a.changed + a.mu.Unlock() + select { + case <-ctx.Done(): + a.mu.Lock() + return archiveText{}, false + case <-timer.C: + log.Warn(ctx, "recording the live captions: no archive pass for this GoP", "startMs", startMs) + a.mu.Lock() + return archiveText{}, false + case <-changed: + } + a.mu.Lock() + } +} + +// copy returns the bytes to record for the live GoP seg: seg with its text +// runs replaced by the archival ones. It falls back to seg if the archive pass +// never reports that GoP. +func (a *captionArchive) copy(ctx context.Context, seg []byte) []byte { + events, err := unwrapMuxlEvents(ctx, seg) + catalog, segment := catalogAndSegment(events) + start, ok := gopStartMs(catalog, segment) + if err != nil || !ok { + log.Warn(ctx, "recording the live captions: cannot place the segment's GoP", "error", err) + return seg + } + text, ok := a.claim(ctx, start) + if !ok || text.Runs == nil { + return seg + } + out, err := replaceTextRuns(catalog, segment.Tracks, text.Runs) + if err != nil { + log.Warn(ctx, "recording the live captions: could not replace them", "error", err) + return seg + } + return out +} + +// gopStartMs is the start of a GoP's span as muxl's streaming signer reports it +// to TextFn: the first decode time of the lowest video track, else of the +// lowest audio track, in whole milliseconds. +func gopStartMs(catalog *muxl.MuxlCatalog, segment *muxl.MuxlEvent) (uint64, bool) { + if catalog == nil || segment == nil { + return 0, false + } + var id, scale uint32 + pick := func(trackID, timescale uint32) { + if id == 0 || trackID < id { + id, scale = trackID, timescale + } + } + if catalog.Video != nil { + for _, video := range catalog.Video.Renditions { + pick(video.TrackID(), video.Timescale()) + } + } + if id == 0 && catalog.Audio != nil { + for _, audio := range catalog.Audio.Renditions { + pick(audio.TrackID(), audio.Timescale()) + } + } + ticks, ok := segment.FirstDecodeTimes[strconv.FormatUint(uint64(id), 10)] + if id == 0 || !ok { + return 0, false + } + if scale == 0 { + return ticks, true + } + ts := uint64(scale) + return ticks/ts*1000 + ticks%ts*1000/ts, true +} + +// replaceTextRuns swaps a GoP's text runs for archival ones. Every other run is +// kept byte for byte, and the result is in ascending track-ID order. +func replaceTextRuns(catalog *muxl.MuxlCatalog, tracks map[string][]byte, runs map[uint32][]byte) ([]byte, error) { + if catalog.Text != nil { + for _, text := range catalog.Text.Renditions { + delete(tracks, strconv.FormatUint(uint64(text.TrackID()), 10)) + } + } + for id, run := range runs { + key := strconv.FormatUint(uint64(id), 10) + if _, ok := tracks[key]; ok { + return nil, fmt.Errorf("archive text track %d collides with a media track", id) + } + tracks[key] = run + } + return concatTracksByID(tracks), nil +} + +type captionArchiveKey struct{} + +// withCaptionArchive marks ctx's ingest session as recording archival captions +// into archive. The context reaches both the session's signer and ValidateMP4 +// for its segments, so a segment carries its own session's archive even when +// it is validated after the signer has returned. +func withCaptionArchive(ctx context.Context, archive *captionArchive) context.Context { + return context.WithValue(ctx, captionArchiveKey{}, archive) +} + +func captionArchiveFrom(ctx context.Context) *captionArchive { + archive, _ := ctx.Value(captionArchiveKey{}).(*captionArchive) + return archive +} + +// ArchiveCopy returns the bytes to record for this segment. For an ingest +// session that masters captions that is Muxl with its text runs laid out again +// by the archive pass, which can take up to --captions-master-delay; otherwise +// it is Muxl itself. +func (n *NewSegmentNotification) ArchiveCopy(ctx context.Context) []byte { + if n.archive == nil { + return n.Muxl + } + return n.archive.copy(ctx, n.Muxl) +} diff --git a/pkg/media/captions_archive_test.go b/pkg/media/captions_archive_test.go new file mode 100644 index 000000000..1d312614d --- /dev/null +++ b/pkg/media/captions_archive_test.go @@ -0,0 +1,31 @@ +package media + +import ( + "context" + "testing" + "time" + + "github.com/stretchr/testify/require" +) + +func TestCaptionArchiveClaimWaitsForItsGoP(t *testing.T) { + ctx := context.Background() + archive := newCaptionArchive(time.Hour) + got := make(chan bool, 1) + go func() { + _, ok := archive.claim(ctx, 2000) + got <- ok + }() + require.NoError(t, archive.put(archiveText{StartMs: 0})) + select { + case <-got: + t.Fatal("another GoP's captions released the recording") + case <-time.After(50 * time.Millisecond): + } + require.NoError(t, archive.put(archiveText{StartMs: 2000})) + require.True(t, <-got) + + archive.finish() + _, ok := archive.claim(ctx, 4000) + require.False(t, ok, "a finished session never holds up the recording") +} diff --git a/pkg/media/captions_master.go b/pkg/media/captions_master.go index 57b0a56a6..28d1b75fa 100644 --- a/pkg/media/captions_master.go +++ b/pkg/media/captions_master.go @@ -5,6 +5,8 @@ import ( "encoding/json" "fmt" "io" + "reflect" + "sort" "strings" "sync" "time" @@ -13,6 +15,7 @@ import ( "stream.place/streamplace/pkg/captions" "stream.place/streamplace/pkg/config" + "stream.place/streamplace/pkg/log" "stream.place/streamplace/pkg/muxl" "stream.place/streamplace/pkg/placestream" "stream.place/streamplace/pkg/stt" @@ -26,6 +29,16 @@ const CaptionTrackIDBase uint32 = 100 // between agreed batches. const autoCaptionHold = 3 * time.Second +// liveTapWait bounds how long a live GoP waits for the ingest tap to parse it, +// so embedded captions land in their own GoP. Live media never waits for +// speech recognition; the archive pass places late words for recordings. +const liveTapWait = 500 * time.Millisecond + +// archiveTimingSlack absorbs whisper re-timing words between passes: a word a +// later pass commits can start slightly before audio an earlier pass already +// finalized, so the archive pass waits for coverage a little past its GoP. +const archiveTimingSlack = time.Second + type captionSession interface { captionPolicy() (captions.Policy, error) push(captions.Track, []captions.Cue) error @@ -43,6 +56,7 @@ type captionMaster struct { policyReady bool stopped bool ingestSeen bool + pushed bool recognitionUnavailable bool parsedUntil uint64 mediaFinished bool @@ -52,16 +66,45 @@ type captionMaster struct { changed chan struct{} closes map[uint64]time.Time gopTimes map[uint64]time.Time - signedUntil uint64 - consumed map[string]int64 - pending map[string]muxl.TextCue tracks map[string]muxl.TextTrack - next map[string]int64 // per track, the earliest start of its next cue nextID uint32 + live captionLayout + archive *captionArchivePass +} + +// captionLayout places final cues on one output's timeline. The live stream +// and the archive pass keep separate layouts over the same cues. +type captionLayout struct { + until uint64 // end of the last GoP laid out; written under captionMaster.mu + consumed map[string]int64 + pending map[string]muxl.TextCue + next map[string]int64 // per track, the earliest start of its next cue + // keepTimes starts every cue that arrives in time at its own time, even + // if that cuts short a cue still showing, as long as it doesn't erase + // it. Without it a cue also waits for the reading time of the cue before + // it, which keeps live captions readable but lets one late cue delay + // everything after it. + keepTimes bool +} + +func newCaptionLayout(keepTimes bool) captionLayout { + return captionLayout{consumed: make(map[string]int64), pending: make(map[string]muxl.TextCue), next: make(map[string]int64), keepTimes: keepTimes} +} + +// erases reports whether a cue starting at start would wipe out one of the +// track's pending cues, as when whisper times a cue to start with the one +// before it. +func (l *captionLayout) erases(trackID string, start int64) bool { + for key, shown := range l.pending { + if strings.HasPrefix(key, trackID+"/") && shown.Start >= uint64(start) { + return true + } + } + return false } func newCaptionMaster(ctx context.Context, streamer string, cli *config.CLI, engine stt.Engine) *captionMaster { - return &captionMaster{ctx: ctx, streamer: streamer, sessionID: uuid.NewString(), cli: cli, engine: engine, hub: captions.NewHub(0), current: captions.DefaultPolicy(), changed: make(chan struct{}), closes: make(map[uint64]time.Time), gopTimes: make(map[uint64]time.Time), consumed: make(map[string]int64), pending: make(map[string]muxl.TextCue), tracks: make(map[string]muxl.TextTrack), next: make(map[string]int64)} + return &captionMaster{ctx: ctx, streamer: streamer, sessionID: uuid.NewString(), cli: cli, engine: engine, hub: captions.NewHub(0), current: captions.DefaultPolicy(), changed: make(chan struct{}), closes: make(map[uint64]time.Time), gopTimes: make(map[uint64]time.Time), tracks: make(map[string]muxl.TextTrack), nextID: CaptionTrackIDBase, live: newCaptionLayout(false)} } func captionPolicyFromManifest(data []byte) captions.Policy { @@ -125,6 +168,36 @@ func (m *captionMaster) waitPolicy() (captions.Policy, error) { return policy, nil } func (m *captionMaster) signal() { close(m.changed); m.changed = make(chan struct{}) } + +// waitLocked waits, with m.mu held on entry and on return, until the master +// changes, deadline passes (a zero deadline never does), or ctx ends. +func (m *captionMaster) waitLocked(ctx context.Context, deadline time.Time) error { + changed := m.changed + m.mu.Unlock() + defer m.mu.Lock() + var timeout <-chan time.Time + if !deadline.IsZero() { + timer := time.NewTimer(time.Until(deadline)) + defer timer.Stop() + timeout = timer.C + } + select { + case <-ctx.Done(): + return ctx.Err() + case <-changed: + case <-timeout: + } + return nil +} + +// stop ends the session: no more media will be signed, and the archive pass +// drains what was. +func (m *captionMaster) stop() { + m.mu.Lock() + m.stopped = true + m.signal() + m.mu.Unlock() +} func (m *captionMaster) coverage(end time.Time) { m.mu.Lock() if end.After(m.covered) { @@ -168,6 +241,15 @@ func (m *captionMaster) segmentTime(start uint64) time.Time { delete(m.gopTimes, key) } } + // muxl asks for the time right after the GoP's text, which completes the + // GoP the archive pass will lay out again. + if a := m.archive; a != nil && a.signed != nil && a.signed.req.StartMs == start { + gop := *a.signed + gop.when = prediction + a.gops = append(a.gops, gop) + a.signed = nil + m.signal() + } return prediction } func (m *captionMaster) closeGopAt(end uint64, at time.Time) { @@ -182,7 +264,7 @@ func (m *captionMaster) closeGopAt(end uint64, at time.Time) { if end > m.parsedUntil { m.parsedUntil = end } - if end > m.signedUntil { + if end > m.live.until { if _, ok := m.closes[end]; !ok { m.closes[end] = at } @@ -199,6 +281,7 @@ func (m *captionMaster) push(track captions.Track, cues []captions.Cue) error { } arrival, origin := m.arrival, m.mediaOrigin m.ingestSeen = true + m.pushed = true m.signal() m.mu.Unlock() track.Origin = captions.OriginCanonical @@ -212,7 +295,9 @@ func (m *captionMaster) push(track captions.Track, cues []captions.Cue) error { } // text holds only the signer, never appsink. It consumes immutable final cues -// from a private hub. A late final is carried into the next unsigned GoP. +// from a private hub. Live media waits only for the ingest tap: a final that +// arrives after its GoP was signed is carried into the next unsigned GoP, and +// the archive pass puts it back in its own GoP for recordings. func (m *captionMaster) text(ctx context.Context, req muxl.TextRequest) (*muxl.TextAttachment, error) { m.mu.Lock() closeAt, ok := m.closes[req.EndMs] @@ -224,135 +309,234 @@ func (m *captionMaster) text(ctx context.Context, req muxl.TextRequest) (*muxl.T delete(m.closes, end) } } - delay := m.cli.CaptionsMasterDelay - deadline := closeAt.Add(delay) - for { - decision := captions.Decide(captions.Situation{Policy: m.current, Origin: true, IngestCaptions: m.ingestSeen}) - tapReady := m.mediaFinished || m.parsedUntil >= req.EndMs - recognitionReady := !decision.Recognize() || m.engine == nil || m.recognitionUnavailable || m.covered.UnixMilli() >= int64(req.EndMs) - if m.current.Canonical == captions.CanonicalOff || (tapReady && recognitionReady) || !time.Now().Before(deadline) { - break - } - changed := m.changed - m.mu.Unlock() - timer := time.NewTimer(time.Until(deadline)) - select { - case <-ctx.Done(): - timer.Stop() - return nil, ctx.Err() - case <-changed: - timer.Stop() - case <-timer.C: + deadline := closeAt.Add(liveTapWait) + for m.current.Canonical != captions.CanonicalOff && !m.mediaFinished && m.parsedUntil < req.EndMs && time.Now().Before(deadline) { + if err := m.waitLocked(ctx, deadline); err != nil { + m.mu.Unlock() + return nil, err } - m.mu.Lock() } - p, ingest := m.current, m.ingestSeen - until := m.signedUntil m.mu.Unlock() - attachment := &muxl.TextAttachment{} - if p.Canonical == captions.CanonicalOff { + attachment := m.layout(&m.live, req) + if m.archive != nil { m.mu.Lock() - m.signedUntil = req.EndMs - m.pending = make(map[string]muxl.TextCue) + m.archive.signed = &archiveGoP{req: req, closeAt: closeAt, live: attachment} m.mu.Unlock() - return attachment, nil - } - if m.nextID == 0 { - m.nextID = CaptionTrackIDBase - } - horizon := int64(req.StartMs) - captions.DefaultRetention.Milliseconds() - from, to := time.UnixMilli(horizon), time.UnixMilli(int64(req.EndMs)) - for key, end := range m.consumed { - if end <= horizon { - delete(m.consumed, key) - } } - items := make(map[uint32]*muxl.TextTrackAttachment) - for _, track := range m.hub.Tracks(m.streamer) { - if track.Source == captions.SourceAuto && !strings.HasPrefix(track.ID, "push-") && (p.Canonical != captions.CanonicalAuto || ingest) { - for key := range m.pending { - if strings.HasPrefix(key, track.ID+"/") { - delete(m.pending, key) - } + return attachment, nil +} + +// layout places the final cues overlapping req on l's timeline. Cues play in +// order: one that arrives after its GoP was laid out starts in the next GoP, +// after the previous cue has had its own reading time, and keeps its whole +// duration; it replaces whatever the track still shows. Every declared track +// appears in every GoP, so neither output drops a track mid-stream. +func (m *captionMaster) layout(l *captionLayout, req muxl.TextRequest) *muxl.TextAttachment { + m.mu.Lock() + p, ingest := m.current, m.ingestSeen + m.mu.Unlock() + cues := make(map[uint32][]muxl.TextCue) + if p.Canonical == captions.CanonicalOff { + l.pending = make(map[string]muxl.TextCue) + } else { + horizon := int64(req.StartMs) - captions.DefaultRetention.Milliseconds() + from, to := time.UnixMilli(horizon), time.UnixMilli(int64(req.EndMs)) + for key, end := range l.consumed { + if end <= horizon { + delete(l.consumed, key) } - continue - } - language := track.Language - if (language == "" || language == "und") && track.Source == captions.SourceIngest && len(p.Languages) > 0 { - language = p.Languages[0] } - if language == "" { - language = "und" - } - trackKey := string(track.Source) + "/" + strings.ToLower(language) - for _, cue := range m.hub.Cues(m.streamer, track.ID, from, to) { - key := track.ID + "/" + cue.ID - if _, ok := m.consumed[key]; ok { + for _, track := range m.hub.Tracks(m.streamer) { + if track.Source == captions.SourceAuto && !strings.HasPrefix(track.ID, "push-") && (p.Canonical != captions.CanonicalAuto || ingest) { + for key := range l.pending { + if strings.HasPrefix(key, track.ID+"/") { + delete(l.pending, key) + } + } continue } - m.consumed[key] = cue.End.UnixMilli() - start, end := cue.Start.UnixMilli(), cue.End.UnixMilli() - if end <= 0 { - continue + language := track.Language + if (language == "" || language == "und") && track.Source == captions.SourceIngest && len(p.Languages) > 0 { + language = p.Languages[0] } - start = max(start, 0) - // Live captions play in order on the canonical timeline. A late - // cue starts in the first unsigned GoP, after the previous cue - // has had its own reading time, and keeps its whole duration; it - // replaces whatever the track still shows. - duration := end - start - start = max(start, int64(until), m.next[track.ID]) - m.next[track.ID] = start + duration - end = start + duration - if track.Source == captions.SourceAuto { - end += autoCaptionHold.Milliseconds() + if language == "" { + language = "und" } - for k, shown := range m.pending { - if strings.HasPrefix(k, track.ID+"/") && shown.End > uint64(start) { - shown.End = uint64(start) - m.pending[k] = shown - if shown.End <= shown.Start { - delete(m.pending, k) + trackKey := string(track.Source) + "/" + strings.ToLower(language) + for _, cue := range m.hub.Cues(m.streamer, track.ID, from, to) { + key := track.ID + "/" + cue.ID + if _, ok := l.consumed[key]; ok { + continue + } + l.consumed[key] = cue.End.UnixMilli() + start, end := cue.Start.UnixMilli(), cue.End.UnixMilli() + if end <= 0 { + continue + } + start = max(start, 0) + duration := end - start + if !l.keepTimes || start < int64(l.until) || l.erases(track.ID, start) { + start = max(start, int64(l.until), l.next[track.ID]) + } + l.next[track.ID] = start + duration + end = start + duration + if track.Source == captions.SourceAuto { + end += autoCaptionHold.Milliseconds() + } + for k, shown := range l.pending { + if strings.HasPrefix(k, track.ID+"/") && shown.End > uint64(start) { + shown.End = uint64(start) + l.pending[k] = shown + if shown.End <= shown.Start { + delete(l.pending, k) + } } } + l.pending[key] = muxl.TextCue{Start: uint64(start), End: uint64(end), Text: cue.Text, ID: m.sessionID + "/" + key} + m.declareTrack(trackKey, language, track.Source) } - m.pending[key] = muxl.TextCue{Start: uint64(start), End: uint64(end), Text: cue.Text, ID: m.sessionID + "/" + key} - if _, ok := m.tracks[trackKey]; !ok { - m.tracks[trackKey] = muxl.TextTrack{TrackID: m.nextID, Language: language, Label: string(track.Source)} - m.nextID++ + m.mu.Lock() + config, ok := m.tracks[trackKey] + m.mu.Unlock() + if !ok { + continue + } + for key, cue := range l.pending { + if !strings.HasPrefix(key, track.ID+"/") { + continue + } + if cue.Start < req.EndMs && cue.End > req.StartMs { + clipped := cue + clipped.Start = max(clipped.Start, req.StartMs) + clipped.End = min(clipped.End, req.EndMs) + cues[config.TrackID] = append(cues[config.TrackID], clipped) + } + if cue.End <= req.EndMs { + delete(l.pending, key) + } } } - config, ok := m.tracks[trackKey] - if !ok { - continue + } + attachment := &muxl.TextAttachment{} + m.mu.Lock() + l.until = req.EndMs + for _, config := range m.tracks { + attachment.Tracks = append(attachment.Tracks, muxl.TextTrackAttachment{TextTrack: config, Cues: cues[config.TrackID]}) + } + m.mu.Unlock() + sort.Slice(attachment.Tracks, func(i, j int) bool { return attachment.Tracks[i].TrackID < attachment.Tracks[j].TrackID }) + for _, track := range attachment.Tracks { + sort.Slice(track.Cues, func(i, j int) bool { + a, b := track.Cues[i], track.Cues[j] + if a.Start != b.Start { + return a.Start < b.Start + } + if a.End != b.End { + return a.End < b.End + } + return a.ID < b.ID + }) + } + return attachment +} + +// declareTrack gives a source and language their track ID on first use. IDs +// and configurations never change during the session, in either output. +func (m *captionMaster) declareTrack(key, language string, source captions.Source) { + m.mu.Lock() + defer m.mu.Unlock() + if _, ok := m.tracks[key]; !ok { + m.tracks[key] = muxl.TextTrack{TrackID: m.nextID, Language: language, Label: string(source)} + m.nextID++ + } +} + +// captionArchivePass lays each live-signed GoP out again once recognition has +// covered it, or CaptionsMasterDelay after it closed, so a recording shows +// words in the GoP they were spoken in rather than where they arrived live. +type captionArchivePass struct { + layout captionLayout + signed *archiveGoP // laid out live, until segmentTime stamps it + gops []archiveGoP // stamped and awaiting the pass; guarded by captionMaster.mu + signer muxl.SignerInput + put func(archiveText) error + done chan struct{} +} + +type archiveGoP struct { + req muxl.TextRequest + closeAt time.Time + when time.Time // the GoP's signed start time + live *muxl.TextAttachment +} + +// archiveTo starts the archive pass. signer carries the cert, the key or Sign +// callback, and the manifest for its text runs; put receives every GoP the +// streaming signer signed, in order. +func (m *captionMaster) archiveTo(signer muxl.SignerInput, put func(archiveText) error) { + m.archive = &captionArchivePass{layout: newCaptionLayout(true), signer: signer, put: put, done: make(chan struct{})} + go m.runArchive() +} + +// awaitArchive waits until the archive pass has handed over every signed GoP. +// It drains once the session stops. +func (m *captionMaster) awaitArchive() { + if m.archive != nil { + <-m.archive.done + } +} + +func (m *captionMaster) runArchive() { + a := m.archive + defer close(a.done) + for { + m.mu.Lock() + for len(a.gops) == 0 && !m.stopped { + _ = m.waitLocked(context.Background(), time.Time{}) } - item := items[config.TrackID] - if item == nil { - item = &muxl.TextTrackAttachment{TextTrack: config} - items[config.TrackID] = item + if len(a.gops) == 0 { + m.mu.Unlock() + return } - for key, cue := range m.pending { - if !strings.HasPrefix(key, track.ID+"/") { - continue - } - if cue.Start < req.EndMs && cue.End > req.StartMs { - clipped := cue - clipped.Start = max(clipped.Start, req.StartMs) - clipped.End = min(clipped.End, req.EndMs) - item.Cues = append(item.Cues, clipped) + gop := a.gops[0] + a.gops = a.gops[1:] + deadline := gop.closeAt.Add(m.cli.CaptionsMasterDelay) + for !m.archiveReady(gop.req.EndMs) && time.Now().Before(deadline) { + if m.waitLocked(m.ctx, deadline) != nil { + break } - if cue.End <= req.EndMs { - delete(m.pending, key) + } + m.mu.Unlock() + text := archiveText{StartMs: gop.req.StartMs} + if attachment := m.layout(&a.layout, gop.req); !reflect.DeepEqual(attachment, gop.live) { + in := a.signer + in.SegmentTimeFn = func(uint64) time.Time { return gop.when } + runs, err := muxl.RunMuxlSignTextRuns(context.WithoutCancel(m.ctx), gop.req, attachment.Tracks, in) + if err != nil { + log.Warn(m.ctx, "sign archive captions; recording the live text", "error", err, "streamer", m.streamer) + } else { + text.Runs = runs } } + if err := a.put(text); err != nil { + log.Warn(m.ctx, "hand over archive captions", "error", err, "streamer", m.streamer) + } } - for _, item := range items { - attachment.Tracks = append(attachment.Tracks, *item) +} + +// archiveReady reports, with m.mu held, whether the archive pass has every +// caption for the GoP ending at end: recognition has covered it, nothing +// recognizes it, or the media ended. Pushed captions give no such signal, so +// a session that has had pushes waits out the hold. +func (m *captionMaster) archiveReady(end uint64) bool { + if m.mediaFinished || m.current.Canonical == captions.CanonicalOff { + return true } - m.mu.Lock() - m.signedUntil = req.EndMs - m.mu.Unlock() - return attachment, nil + if m.pushed || m.parsedUntil < end { + return false + } + decision := captions.Decide(captions.Situation{Policy: m.current, Origin: true, IngestCaptions: m.ingestSeen}) + return !decision.Recognize() || m.engine == nil || m.recognitionUnavailable || m.covered.UnixMilli() >= int64(end)+archiveTimingSlack.Milliseconds() } type captionManagerKey struct{} diff --git a/pkg/media/captions_master_control_test.go b/pkg/media/captions_master_control_test.go index 64203f855..749d4d567 100644 --- a/pkg/media/captions_master_control_test.go +++ b/pkg/media/captions_master_control_test.go @@ -24,7 +24,7 @@ func TestCaptionMasterWorkerControlAndReconnectIDs(t *testing.T) { path := filepath.Join(t.TempDir(), "ingest.sock") stop, err := master.servePush(path) require.NoError(t, err) - unregister := mm.registerWorkerCaptionMaster(path, master.streamer) + _, unregister := mm.registerWorkerCaptionMaster(ctx, path, master.streamer) policy, live := mm.OriginCaptionPolicy(master.streamer) require.True(t, live) require.Equal(t, captions.CanonicalIngest, policy.Canonical) diff --git a/pkg/media/captions_master_feed.go b/pkg/media/captions_master_feed.go index 0ec9e2022..a7ebef489 100644 --- a/pkg/media/captions_master_feed.go +++ b/pkg/media/captions_master_feed.go @@ -59,10 +59,7 @@ func (m *captionMaster) tee(input io.Reader) (io.Reader, func()) { } }() return media, func() { - m.mu.Lock() - m.stopped = true - m.signal() - m.mu.Unlock() + m.stop() _ = media.CloseWithError(context.Canceled) _ = audio.CloseWithError(context.Canceled) if closer, ok := input.(io.ReadCloser); ok { diff --git a/pkg/media/captions_master_socket.go b/pkg/media/captions_master_socket.go index 5c42d80f0..299044c2c 100644 --- a/pkg/media/captions_master_socket.go +++ b/pkg/media/captions_master_socket.go @@ -1,6 +1,7 @@ package media import ( + "context" "encoding/json" "fmt" "net" @@ -88,6 +89,15 @@ func (r *remoteCaptionMaster) push(track captions.Track, cues []captions.Cue) er _, err := r.call(captionControl{Track: track, Cues: cues}) return err } -func (mm *MediaManager) registerWorkerCaptionMaster(path, streamer string) func() { - return mm.registerCaptionMaster(streamer, &remoteCaptionMaster{path: path}) + +// registerWorkerCaptionMaster routes pushes to a worker's caption socket. When +// main records segments, the returned ctx also collects the worker's archival +// captions (Captions frames) for them; a ctx that already does is reused. +func (mm *MediaManager) registerWorkerCaptionMaster(ctx context.Context, path, streamer string) (context.Context, func()) { + unregister := mm.registerCaptionMaster(streamer, &remoteCaptionMaster{path: path}) + if captionArchiveFrom(ctx) != nil || mm.cli == nil || !mm.cli.S3Configured() { + return ctx, unregister + } + archive := newCaptionArchive(mm.cli.CaptionsMasterDelay) + return withCaptionArchive(ctx, archive), func() { unregister(); archive.finish() } } diff --git a/pkg/media/captions_master_test.go b/pkg/media/captions_master_test.go index b3af41eba..67c26ec6f 100644 --- a/pkg/media/captions_master_test.go +++ b/pkg/media/captions_master_test.go @@ -6,6 +6,7 @@ import ( "encoding/json" "io" "os" + "strconv" "sync" "testing" "time" @@ -13,6 +14,7 @@ import ( "github.com/stretchr/testify/require" "stream.place/streamplace/pkg/captions" "stream.place/streamplace/pkg/config" + "stream.place/streamplace/pkg/crypto/signers" "stream.place/streamplace/pkg/muxl" "stream.place/streamplace/pkg/stt" ) @@ -49,49 +51,56 @@ func masterTrack(source captions.Source) captions.Track { return captions.Track{ID: captions.TrackID(captions.OriginCanonical, source, "en-US"), Language: "en-US", Source: source, Origin: captions.OriginCanonical, Kind: captions.KindCaptions} } -func TestCaptionMasterHoldsUntilRecognitionCoversSpan(t *testing.T) { +func TestCaptionMasterLiveSkipsTheHoldTheArchiveWaitsOut(t *testing.T) { ctx := context.Background() engine := &captionTestEngine{enter: make(chan struct{}), release: make(chan struct{})} m := newCaptionMaster(ctx, "streamer", &config.CLI{CaptionsMasterDelay: time.Hour}, engine) - m.mediaFinished = true + m.setManifest(captionManifest("auto")) + m.clockAt(time.UnixMilli(0), time.Now()) + m.closeGopAt(3000, time.Now()) + signer := newBareSegmentSigner(t) + key, err := signers.MarshalES256KPrivateKeyPEM(signer.Signer) + require.NoError(t, err) + archived := make(chan archiveText, 1) + m.archiveTo(muxl.SignerInput{CertPEM: signer.Cert, KeyPEM: key, TrackManifest: signer.PrebuiltManifest}, func(text archiveText) error { + archived <- text + return nil + }) + defer m.awaitArchive() + defer m.stop() r, err := captions.NewRecognizer(ctx, captions.RecognizerOptions{Streamer: m.streamer, Origin: captions.OriginCanonical, Hub: m.hub, Engine: engine, OnCoverage: m.coverage, Step: time.Millisecond, MinWindow: time.Millisecond, SilenceFlush: 100 * time.Millisecond}) require.NoError(t, err) defer r.Close() + var release sync.Once + defer release.Do(func() { close(engine.release) }) pcm := make([]float32, stt.SampleRate*2) for i := range stt.SampleRate { pcm[i] = 0.2 } - r.Push(time.UnixMilli(0), pcm) + r.Push(time.UnixMilli(2000), pcm) <-engine.enter - result := make(chan *muxl.TextAttachment, 1) - started := make(chan struct{}) - go func() { - close(started) - attachment, err := m.text(ctx, muxl.TextRequest{StartMs: 0, EndMs: 1000}) - if err == nil { - result <- attachment - } - }() - <-started + live, err := m.text(ctx, muxl.TextRequest{StartMs: 2000, EndMs: 3000}) + require.NoError(t, err) + require.Empty(t, live.Tracks, "live media is signed without waiting for recognition") + m.segmentTime(2000) select { - case <-result: - t.Fatal("signed before recognition finished") - default: + case <-archived: + t.Fatal("archived before recognition covered the GoP") + case <-time.After(100 * time.Millisecond): } - close(engine.release) + release.Do(func() { close(engine.release) }) + var text archiveText select { - case attachment := <-result: - require.Len(t, attachment.Tracks, 1) - require.Equal(t, "en-US", attachment.Tracks[0].Language) - require.Equal(t, "auto", attachment.Tracks[0].Label) - require.Len(t, attachment.Tracks[0].Cues, 1) - require.Equal(t, uint64(250), attachment.Tracks[0].Cues[0].Start) - require.GreaterOrEqual(t, attachment.Tracks[0].Cues[0].End, uint64(750)) - require.LessOrEqual(t, attachment.Tracks[0].Cues[0].End, uint64(1000)) - require.Equal(t, "held words", attachment.Tracks[0].Cues[0].Text) + case text = <-archived: case <-time.After(5 * time.Second): - t.Fatal("coverage did not release GoP") + t.Fatal("coverage did not release the archive pass") } + require.Equal(t, uint64(2000), text.StartMs, "keyed by the GoP's media start") + cues, err := muxl.RunMuxlReadTextCues(ctx, bytes.NewReader(text.Runs[CaptionTrackIDBase]), CaptionTrackIDBase) + require.NoError(t, err) + require.Len(t, cues, 1) + require.Equal(t, "held words", cues[0].Text) + require.Equal(t, uint64(2250), cues[0].Start, "the archive places words in the GoP they were spoken in") } func TestCaptionMasterLaysOutLateCuesInOrder(t *testing.T) { @@ -126,6 +135,38 @@ func TestCaptionMasterLaysOutLateCuesInOrder(t *testing.T) { require.Equal(t, []muxl.TextCue{{Start: 2000, End: 3000, Text: "crossing", ID: id("crossing")}}, third.Tracks[0].Cues) } +func TestCaptionArchiveLayoutKeepsOnTimeCuesAfterALateOne(t *testing.T) { + m := newCaptionMaster(context.Background(), "streamer", &config.CLI{}, nil) + archive := newCaptionLayout(true) + // The first GoP was laid out before its speech was recognized. + m.layout(&archive, muxl.TextRequest{StartMs: 0, EndMs: 1000}) + track := masterTrack(captions.SourceHuman) + m.hub.Publish(m.streamer, track, captions.Cue{ID: "late", Start: time.UnixMilli(500), End: time.UnixMilli(1100), Text: "late", Final: true}) + m.hub.Publish(m.streamer, track, captions.Cue{ID: "on-time", Start: time.UnixMilli(1200), End: time.UnixMilli(1800), Text: "on time", Final: true}) + got := m.layout(&archive, muxl.TextRequest{StartMs: 1000, EndMs: 2000}) + id := func(cue string) string { return m.sessionID + "/" + track.ID + "/" + cue } + require.Equal(t, []muxl.TextCue{ + {Start: 1000, End: 1200, Text: "late", ID: id("late")}, + {Start: 1200, End: 1800, Text: "on time", ID: id("on-time")}, + }, got.Tracks[0].Cues, "a late cue cannot push the speech after it late") +} + +func TestCaptionArchiveLayoutNeverErasesACue(t *testing.T) { + m := newCaptionMaster(context.Background(), "streamer", &config.CLI{}, nil) + archive := newCaptionLayout(true) + track := masterTrack(captions.SourceHuman) + // Whisper sometimes times a cue to start with the one before it. + m.hub.Publish(m.streamer, track, captions.Cue{ID: "a", Start: time.UnixMilli(200), End: time.UnixMilli(500), Text: "a", Final: true}) + m.hub.Publish(m.streamer, track, captions.Cue{ID: "b", Start: time.UnixMilli(200), End: time.UnixMilli(400), Text: "b", Final: true}) + got := m.layout(&archive, muxl.TextRequest{StartMs: 0, EndMs: 1000}) + var texts []string + for _, cue := range got.Tracks[0].Cues { + require.Greater(t, cue.End, cue.Start) + texts = append(texts, cue.Text) + } + require.ElementsMatch(t, []string{"a", "b"}, texts) +} + func TestCaptionMasterPolicySwitchPreservesTruthfulTracks(t *testing.T) { m := newCaptionMaster(context.Background(), "streamer", &config.CLI{}, nil) m.mediaFinished = true @@ -141,13 +182,17 @@ func TestCaptionMasterPolicySwitchPreservesTruthfulTracks(t *testing.T) { require.NoError(t, m.push(ingest, []captions.Cue{{ID: "ingest", Start: m.arrival.Add(1200 * time.Millisecond), End: m.arrival.Add(1500 * time.Millisecond), Text: "supplied", Final: true}})) second, err := m.text(context.Background(), muxl.TextRequest{StartMs: 1000, EndMs: 2000}) require.NoError(t, err) - require.Len(t, second.Tracks, 1) - require.Equal(t, CaptionTrackIDBase+1, second.Tracks[0].TrackID) - require.Equal(t, "ingest", second.Tracks[0].Label) + require.Len(t, second.Tracks, 2) + require.Empty(t, second.Tracks[0].Cues, "the replaced automatic track continues, empty") + require.Equal(t, CaptionTrackIDBase+1, second.Tracks[1].TrackID) + require.Equal(t, "ingest", second.Tracks[1].Label) m.setManifest(captionManifest("off")) off, err := m.text(context.Background(), muxl.TextRequest{StartMs: 2000, EndMs: 3000}) require.NoError(t, err) - require.Empty(t, off.Tracks) + require.Len(t, off.Tracks, 2) + for _, track := range off.Tracks { + require.Empty(t, track.Cues) + } } type pushedFixtureManifester struct { @@ -247,14 +292,18 @@ func TestCaptionMasterOriginStreamingRecognizesDecodedAudio(t *testing.T) { fixture, err := os.ReadFile(getFixture("h264-opus-frag.mp4")) require.NoError(t, err) engine := &captionTestEngine{enter: make(chan struct{}), release: make(chan struct{})} - mm := NewOffline(&config.CLI{CaptionsMasterDelay: 500 * time.Millisecond}) + mm := NewOffline(&config.CLI{CaptionsMasterDelay: 5 * time.Second}) mm.STT = engine ms := newBareSegmentSigner(t) ms.PrebuiltManifest = captionManifest("auto") + archive := newCaptionArchive(mm.cli.CaptionsMasterDelay) events := make(chan *muxl.MuxlEvent, 16) done := make(chan error, 1) - go func() { done <- mm.SignOriginStream(ctx, ms, bytes.NewReader(fixture), events); close(events) }() - var archived bytes.Buffer + go func() { + done <- mm.SignOriginStream(withCaptionArchive(ctx, archive), ms, bytes.NewReader(fixture), events) + close(events) + }() + var live [][]byte first := true for event := range events { if event.Type != "signed-segment" { @@ -264,54 +313,101 @@ func TestCaptionMasterOriginStreamingRecognizesDecodedAudio(t *testing.T) { first = false close(engine.release) } - archived.Write(concatTracksByID(event.Tracks)) + live = append(live, concatTracksByID(event.Tracks)) } require.NoError(t, <-done) - tracks, err := muxl.RunMuxlTextTracks(ctx, bytes.NewReader(archived.Bytes())) + stream := bytes.Join(live, nil) + tracks, err := muxl.RunMuxlTextTracks(ctx, bytes.NewReader(stream)) require.NoError(t, err) require.Equal(t, []muxl.TextTrack{{TrackID: CaptionTrackIDBase, Language: "en-US", Label: "auto"}}, tracks) - cues, err := muxl.RunMuxlReadTextCues(ctx, bytes.NewReader(archived.Bytes()), CaptionTrackIDBase) + cues, err := muxl.RunMuxlReadTextCues(ctx, bytes.NewReader(stream), CaptionTrackIDBase) require.NoError(t, err) require.Len(t, cues, 1) require.Equal(t, "held words", cues[0].Text) require.GreaterOrEqual(t, cues[0].Start, uint64(1000), "late words cannot be written into the already signed first GoP") require.GreaterOrEqual(t, cues[0].End-cues[0].Start, uint64(500), "a late cue keeps its whole duration") - report, err := muxl.RunMuxlVerify(ctx, bytes.NewReader(archived.Bytes())) + requireValidMuxl(t, ctx, stream) + + recorded := recordSegments(t, ctx, archive, live) + cues, err = muxl.RunMuxlReadTextCues(ctx, bytes.NewReader(recorded), CaptionTrackIDBase) require.NoError(t, err) - var verified struct { - Segments []struct { - ValidationState string `json:"validation_state"` - } `json:"segments"` - } - require.NoError(t, json.Unmarshal([]byte(report), &verified)) - for _, track := range verified.Segments { - require.NotEqual(t, "Invalid", track.ValidationState) - } + require.Len(t, cues, 1) + require.Equal(t, "held words", cues[0].Text) + require.Equal(t, uint64(250), cues[0].Start, "the recording places words in the GoP they were spoken in") + requireValidMuxl(t, ctx, recorded) } -func TestCaptionMasterVoicedEOFFinalsReachLastSignedGoP(t *testing.T) { +func TestCaptionMasterRecordingKeepsVoicedEOFFinals(t *testing.T) { ctx, cancel := context.WithTimeout(context.Background(), 20*time.Second) defer cancel() fixture, err := os.ReadFile(getFixture("h264-opus-frag.mp4")) require.NoError(t, err) engine := &captionTestEngine{scripted: &stt.Result{Language: "en-US", Words: []stt.Word{{Text: "final voiced words", Start: 1500 * time.Millisecond, End: 1800 * time.Millisecond, Prob: 0.99}}}} - mm := NewOffline(&config.CLI{CaptionsMasterDelay: time.Second}) + mm := NewOffline(&config.CLI{CaptionsMasterDelay: 5 * time.Second}) mm.STT = engine ms := newBareSegmentSigner(t) ms.PrebuiltManifest = captionManifest("auto") + archive := newCaptionArchive(mm.cli.CaptionsMasterDelay) events := make(chan *muxl.MuxlEvent, 16) done := make(chan error, 1) - go func() { done <- mm.SignOriginStream(ctx, ms, bytes.NewReader(fixture), events); close(events) }() - var last []byte + go func() { + done <- mm.SignOriginStream(withCaptionArchive(ctx, archive), ms, bytes.NewReader(fixture), events) + close(events) + }() + var live [][]byte for event := range events { if event.Type == "signed-segment" { - last = concatTracksByID(event.Tracks) + live = append(live, concatTracksByID(event.Tracks)) } } require.NoError(t, <-done) - cues, err := muxl.RunMuxlReadTextCues(ctx, bytes.NewReader(last), CaptionTrackIDBase) + recorded := recordSegments(t, ctx, archive, live) + cues, err := muxl.RunMuxlReadTextCues(ctx, bytes.NewReader(recorded), CaptionTrackIDBase) require.NoError(t, err) require.Len(t, cues, 1) require.Equal(t, "final voiced words", cues[0].Text) require.Equal(t, uint64(1500), cues[0].Start) } + +// recordSegments records live segments the way the director does: each one's +// archive copy, once the session's archive pass has reported its GoP. +func recordSegments(t *testing.T, ctx context.Context, archive *captionArchive, live [][]byte) []byte { + t.Helper() + var recorded bytes.Buffer + for _, seg := range live { + copied := archive.copy(ctx, seg) + liveTracks, copiedTracks := segmentTracks(t, ctx, seg), segmentTracks(t, ctx, copied) + for id, run := range liveTracks { + if n, err := strconv.Atoi(id); err == nil && uint32(n) < CaptionTrackIDBase { + require.Equal(t, run, copiedTracks[id], "track %s is recorded byte for byte", id) + } + } + require.Equal(t, concatTracksByID(copiedTracks), copied, "recorded runs keep ascending track order") + recorded.Write(copied) + } + return recorded.Bytes() +} + +func segmentTracks(t *testing.T, ctx context.Context, seg []byte) map[string][]byte { + t.Helper() + events, err := unwrapMuxlEvents(ctx, seg) + require.NoError(t, err) + _, tracks := catalogAndTracks(events) + return tracks +} + +func requireValidMuxl(t *testing.T, ctx context.Context, data []byte) { + t.Helper() + report, err := muxl.RunMuxlVerify(ctx, bytes.NewReader(data)) + require.NoError(t, err) + var verified struct { + Segments []struct { + ValidationState string `json:"validation_state"` + } `json:"segments"` + } + require.NoError(t, json.Unmarshal([]byte(report), &verified)) + require.NotEmpty(t, verified.Segments) + for _, track := range verified.Segments { + require.NotEqual(t, "Invalid", track.ValidationState) + } +} diff --git a/pkg/media/captions_transcode_test.go b/pkg/media/captions_transcode_test.go index 5e269e139..688b0fa50 100644 --- a/pkg/media/captions_transcode_test.go +++ b/pkg/media/captions_transcode_test.go @@ -64,7 +64,11 @@ func TestCaptionMasterCanonicalNamespaceSurvivesAudioCompletion(t *testing.T) { require.GreaterOrEqual(t, len(sources), 3, "exercise another AV GoP after declaring text") } for index, segment := range completed { - require.True(t, bytes.HasPrefix(segment, sources[index]), "audio completion must retain every signed source byte") + completedTracks := segmentTracks(t, ctx, segment) + for id, run := range segmentTracks(t, ctx, sources[index]) { + require.Equal(t, run, completedTracks[id], "audio completion must retain every signed source run (track %s)", id) + } + require.Equal(t, concatTracksByID(completedTracks), segment, "the added audio run keeps ascending track order, ahead of text") events, err := unwrapMuxlEvents(ctx, segment) require.NoError(t, err) catalog, tracks := catalogAndTracks(events) diff --git a/pkg/media/frame_server.go b/pkg/media/frame_server.go index a60d7e387..f4b55f441 100644 --- a/pkg/media/frame_server.go +++ b/pkg/media/frame_server.go @@ -30,6 +30,8 @@ const workerDrainGrace = 60 * time.Second // *ingestframe.Writer, so RunMP4IngestWorker is agnostic to which it gets. type FrameWriter interface { Segment(seg []byte) error + // Captions sends one GoP's archival caption text runs; see archiveText. + Captions(payload []byte) error End() error Error(msg string) error } @@ -85,9 +87,13 @@ func (s *frameServer) push(typ ingestframe.Type, payload []byte) { } func (s *frameServer) Segment(seg []byte) error { s.push(ingestframe.Segment, seg); return nil } -func (s *frameServer) End() error { s.push(ingestframe.End, nil); return nil } -func (s *frameServer) Error(msg string) error { s.push(ingestframe.Error, []byte(msg)); return nil } -func (s *frameServer) Answer(sdp string) error { s.push(ingestframe.Answer, []byte(sdp)); return nil } +func (s *frameServer) Captions(payload []byte) error { + s.push(ingestframe.Captions, payload) + return nil +} +func (s *frameServer) End() error { s.push(ingestframe.End, nil); return nil } +func (s *frameServer) Error(msg string) error { s.push(ingestframe.Error, []byte(msg)); return nil } +func (s *frameServer) Answer(sdp string) error { s.push(ingestframe.Answer, []byte(sdp)); return nil } // dropped reports how many buffered frames were discarded because the buffer // overflowed (main was disconnected longer than the buffer window). diff --git a/pkg/media/ingest_daemon.go b/pkg/media/ingest_daemon.go index fb1e8900a..239696591 100644 --- a/pkg/media/ingest_daemon.go +++ b/pkg/media/ingest_daemon.go @@ -135,8 +135,8 @@ func removeWorkerFiles(socketPath string) { // manifestSource, when non-nil, is polled to refresh the worker's C2PA manifest // over the same socket (pushManifestUpdates) — so a pre-live → live transition // reaches a worker that has no model of its own. It's re-armed per connection. -func (mm *MediaManager) ConsumeWorkerSocket(ctx context.Context, socketPath, streamer string, onSegment func([]byte) error, manifestSource func() ([]byte, error)) error { - unregister := mm.registerWorkerCaptionMaster(socketPath, streamer) +func (mm *MediaManager) ConsumeWorkerSocket(ctx context.Context, socketPath, streamer string, onSegment func(context.Context, []byte) error, manifestSource func() ([]byte, error)) error { + ctx, unregister := mm.registerWorkerCaptionMaster(ctx, socketPath, streamer) defer unregister() connectedOnce := false giveUp := time.Now().Add(workerConnectGrace) @@ -303,7 +303,7 @@ func (mm *MediaManager) MP4IngestDetached(ctx context.Context, conn net.Conn, pr // it signs with a frozen one otherwise. Fixed start for stable change detection. start := time.Now().UnixMilli() manifestSource := func() ([]byte, error) { return mm.streamerManifest(ctx, ms.Streamer(), start) } - err = mm.ConsumeWorkerSocket(ctx, cfg.SocketPath, ms.Streamer(), mm.validateSegment(ctx), manifestSource) + err = mm.ConsumeWorkerSocket(ctx, cfg.SocketPath, ms.Streamer(), mm.validateSegment, manifestSource) recordWorkerExit("mp4", err, ctx.Err()) // Reap the worker unless we're deliberately leaving it running across a main // restart (ctx cancel). On a clean end OR a crash the worker has exited, so @@ -407,7 +407,7 @@ func (mm *MediaManager) WHIPIngestDetached(ctx context.Context, offerSDP string, return "", err } _ = conn.SetReadDeadline(time.Time{}) // clear; streaming has no deadline - unregister := mm.registerWorkerCaptionMaster(cfg.SocketPath, ms.Streamer()) + ctx, unregister := mm.registerWorkerCaptionMaster(ctx, cfg.SocketPath, ms.Streamer()) // Consume the signed segments in the background; the HTTP handler returns the // answer now and the WebRTC media establishes directly to the worker. @@ -428,13 +428,13 @@ func (mm *MediaManager) WHIPIngestDetached(ctx context.Context, offerSDP string, manifestSource := func() ([]byte, error) { return mm.streamerManifest(ctx, ms.Streamer(), start) } go pushManifestUpdates(wctx, conn, manifestSource) - sawEnd, _ := mm.consumeWorkerFrames(ctx, fr, ms.Streamer(), mm.validateSegment(ctx), nil) + sawEnd, _ := mm.consumeWorkerFrames(ctx, fr, ms.Streamer(), mm.validateSegment, nil) conn.Close() var exitErr error if !sawEnd && ctx.Err() == nil { // Connection dropped but the detached worker lives on — reconnect and // drain its buffer. Its terminal result is the worker's true outcome. - exitErr = mm.ConsumeWorkerSocket(ctx, cfg.SocketPath, ms.Streamer(), mm.validateSegment(ctx), manifestSource) + exitErr = mm.ConsumeWorkerSocket(ctx, cfg.SocketPath, ms.Streamer(), mm.validateSegment, manifestSource) } recordWorkerExit("whip", exitErr, ctx.Err()) go func() { _, _ = proc.Wait() }() @@ -483,7 +483,7 @@ func (mm *MediaManager) ResumeDetachedWorkers(ctx context.Context) { start := time.Now().UnixMilli() manifestSource = func() ([]byte, error) { return mm.streamerManifest(ctx, meta.StreamerDID, start) } } - cerr := mm.ConsumeWorkerSocket(wctx, sock, streamer, mm.validateSegment(ctx), manifestSource) + cerr := mm.ConsumeWorkerSocket(wctx, sock, streamer, mm.validateSegment, manifestSource) if cerr != nil { log.Error(ctx, "resumed ingest worker ended", "socket", sock, "error", cerr) } diff --git a/pkg/media/ingest_daemon_test.go b/pkg/media/ingest_daemon_test.go index d683ce0e4..04881c5fd 100644 --- a/pkg/media/ingest_daemon_test.go +++ b/pkg/media/ingest_daemon_test.go @@ -87,7 +87,7 @@ func TestDetachedWorkerZeroDowntime(t *testing.T) { // Consume through the reconnecting consumer; verify dual-codec signed output. var segs int - onSegment := func(s []byte) error { + onSegment := func(_ context.Context, s []byte) error { out, verr := muxl.RunMuxlVerify(ctx, bytes.NewReader(s)) require.NoError(t, verr) require.NotContains(t, out, `"validation_state":"Invalid"`) diff --git a/pkg/media/ingest_supervisor.go b/pkg/media/ingest_supervisor.go index d40634135..de6e8cece 100644 --- a/pkg/media/ingest_supervisor.go +++ b/pkg/media/ingest_supervisor.go @@ -50,7 +50,8 @@ func (mm *MediaManager) MP4IngestIsolated(ctx context.Context, input io.Reader, return err } cfg.CaptionSocketPath = filepath.Join(dir, uuid.NewString()+".sock") - defer mm.registerWorkerCaptionMaster(cfg.CaptionSocketPath, ms.Streamer())() + ctx, unregisterCaptions := mm.registerWorkerCaptionMaster(ctx, cfg.CaptionSocketPath, ms.Streamer()) + defer unregisterCaptions() defer os.Remove(cfg.CaptionSocketPath + ".captions") cfgJSON, err := json.Marshal(cfg) if err != nil { @@ -154,7 +155,7 @@ func (mm *MediaManager) MP4IngestIsolated(ctx context.Context, input io.Reader, }) // Read signed-segment frames and feed each into the normal chokepoint. - sawEnd, readErr := mm.consumeWorkerFrames(ctx, ingestframe.NewReader(framesR), ms.Streamer(), mm.validateSegment(ctx), func() { + sawEnd, readErr := mm.consumeWorkerFrames(ctx, ingestframe.NewReader(framesR), ms.Streamer(), mm.validateSegment, func() { watchdog.Reset(ingestWorkerWatchdog) }) logsWG.Wait() @@ -199,11 +200,12 @@ func recordWorkerExit(transport string, exitErr, ctxErr error) { } } -// consumeWorkerFrames reads framed segments from the worker and runs ValidateMP4 -// over each. It returns whether a clean End frame was seen and the terminal read -// error: nil on a clean close (End then EOF), or io.ErrUnexpectedEOF / a desync -// error when the worker died mid-frame. -func (mm *MediaManager) consumeWorkerFrames(ctx context.Context, fr *ingestframe.Reader, streamer string, onSegment func([]byte) error, onProgress func()) (sawEnd bool, _ error) { +// consumeWorkerFrames reads framed segments from the worker and runs onSegment +// over each with ctx, which carries the session's caption archive (if any) for +// the worker's Captions frames. It returns whether a clean End frame was seen +// and the terminal read error: nil on a clean close (End then EOF), or +// io.ErrUnexpectedEOF / a desync error when the worker died mid-frame. +func (mm *MediaManager) consumeWorkerFrames(ctx context.Context, fr *ingestframe.Reader, streamer string, onSegment func(context.Context, []byte) error, onProgress func()) (sawEnd bool, _ error) { for { typ, payload, err := fr.ReadFrame() if err != nil { @@ -218,12 +220,21 @@ func (mm *MediaManager) consumeWorkerFrames(ctx context.Context, fr *ingestframe switch typ { case ingestframe.Segment: if onSegment != nil { - if serr := onSegment(payload); serr != nil { + if serr := onSegment(ctx, payload); serr != nil { // Per-segment failures are logged, not fatal to the stream — a // bad GoP shouldn't tear down an otherwise-healthy ingest. log.Error(ctx, "ingest worker: segment handler failed", "streamer", streamer, "error", serr) } } + case ingestframe.Captions: + if archive := captionArchiveFrom(ctx); archive != nil { + var text archiveText + if uerr := json.Unmarshal(payload, &text); uerr != nil { + log.Error(ctx, "ingest worker: bad captions frame", "streamer", streamer, "error", uerr) + } else { + _ = archive.put(text) + } + } case ingestframe.End: sawEnd = true case ingestframe.Error: @@ -235,10 +246,8 @@ func (mm *MediaManager) consumeWorkerFrames(ctx context.Context, fr *ingestframe // validateSegment is the onSegment handler for ingested worker frames: it folds // each signed segment into the normal ValidateMP4 chokepoint (verify → archive → // live-HLS → notify). -func (mm *MediaManager) validateSegment(ctx context.Context) func([]byte) error { - return func(seg []byte) error { - return mm.ValidateMP4(ctx, bytes.NewReader(seg), true) - } +func (mm *MediaManager) validateSegment(ctx context.Context, seg []byte) error { + return mm.ValidateMP4(ctx, bytes.NewReader(seg), true) } // streamWorkerLogs forwards the worker's stderr lines into the node logger. @@ -280,6 +289,7 @@ func (mm *MediaManager) buildWorkerConfig(ctx context.Context, ms MediaSigner) ( BroadcasterHost: mm.cli.BroadcasterHost, CaptionEngineSocket: mm.CaptionEngineSocket, CaptionsMasterDelay: mm.cli.CaptionsMasterDelay, + ArchiveCaptions: mm.cli.S3Configured(), } // Debug recording: main owns the per-stream setting (it needs the DB); the // worker carries out the recording (it owns the data path). A lookup failure diff --git a/pkg/media/ingest_worker.go b/pkg/media/ingest_worker.go index b23cb9b91..d306360e6 100644 --- a/pkg/media/ingest_worker.go +++ b/pkg/media/ingest_worker.go @@ -3,6 +3,7 @@ package media import ( "bytes" "context" + "encoding/json" "fmt" "io" "net/http/httputil" @@ -62,6 +63,10 @@ type IngestWorkerConfig struct { Manifest []byte `json:"manifest"` CaptionEngineSocket string `json:"caption_engine_socket"` CaptionsMasterDelay time.Duration `json:"captions_master_delay"` + // ArchiveCaptions makes the worker lay its captions out a second time for + // recordings and send the re-signed text runs as Captions frames; main sets + // it when it records segments to S3. + ArchiveCaptions bool `json:"archive_captions,omitempty"` // Node transcode signer + broadcaster identity. When set, the worker completes // each single-codec source segment to dual-codec (Opus+AAC) itself — the @@ -153,7 +158,7 @@ func WorkerInput(cfg IngestWorkerConfig, raw io.Reader) io.Reader { // pushes mid-stream (e.g. pre-live → live) takes effect on the next GoP — the // same fresh-per-GoP shape as the in-process signer. Shared by the MP4 and WHIP // workers. -func workerSignStream(cfg IngestWorkerConfig, getManifest func() []byte) SignSegmentStreamFunc { +func workerSignStream(cfg IngestWorkerConfig, getManifest func() []byte, frames FrameWriter) SignSegmentStreamFunc { return func(ctx context.Context, input io.Reader, eventCh chan *muxl.MuxlEvent) error { cli := cfg.workerCLI() engine := stt.NewProxy(cfg.CaptionEngineSocket) @@ -169,6 +174,19 @@ func workerSignStream(cfg IngestWorkerConfig, getManifest func() []byte) SignSeg return fmt.Errorf("serve canonical caption pushes: %w", err) } defer stop() + if cfg.ArchiveCaptions { + // Every archive frame precedes End: the pass drains before the + // signer returns, and the worker frames End after that. + manifest := func() ([]byte, error) { return getManifest(), nil } + master.archiveTo(muxl.SignerInput{CertPEM: cfg.CertPEM, KeyPEM: cfg.KeyPEM, TrackManifestFn: manifest}, func(text archiveText) error { + payload, err := json.Marshal(text) + if err != nil { + return err + } + return frames.Captions(payload) + }) + defer master.awaitArchive() + } input, finish := master.tee(input) defer finish() fetchManifest := func() ([]byte, error) { data := getManifest(); master.setManifest(data); return data, nil } @@ -267,7 +285,7 @@ func RunMP4IngestWorker(ctx context.Context, cfg IngestWorkerConfig, stdin io.Re defer finalize() } - signerElem, done, err := muxlSignSegmentElem(ctx, mm.cli, workerSignStream(cfg, getManifest), onSegment) + signerElem, done, err := muxlSignSegmentElem(ctx, mm.cli, workerSignStream(cfg, getManifest, frames), onSegment) if err != nil { return fmt.Errorf("build signer element: %w", err) } diff --git a/pkg/media/ingest_worker_test.go b/pkg/media/ingest_worker_test.go index 011e6d9a2..f3acf9bab 100644 --- a/pkg/media/ingest_worker_test.go +++ b/pkg/media/ingest_worker_test.go @@ -3,6 +3,7 @@ package media import ( "bytes" "context" + "encoding/json" "errors" "fmt" "io" @@ -129,6 +130,7 @@ func TestRunMP4IngestWorkerProducesValidSignedFrames(t *testing.T) { NodeCertPEM: ms.Cert, NodeKeyPEM: keyPEM, BroadcasterHost: "test.example.com", + ArchiveCaptions: true, } mp4 := makeH264AACFMP4(t, ctx, getFixture("5sec.mp4")) @@ -141,18 +143,30 @@ func TestRunMP4IngestWorkerProducesValidSignedFrames(t *testing.T) { r := ingestframe.NewReader(&buf) var segs int + var segmentStarts, archiveStarts []uint64 for { typ, payload, err := r.ReadFrame() if errors.Is(err, io.EOF) { break } require.NoError(t, err) - require.Equal(t, ingestframe.Segment, typ, "worker emits only Segment frames; End is the subcommand's job") + if typ == ingestframe.Captions { + var text archiveText + require.NoError(t, json.Unmarshal(payload, &text)) + archiveStarts = append(archiveStarts, text.StartMs) + continue + } + require.Equal(t, ingestframe.Segment, typ, "worker emits Segment and Captions frames; End is the subcommand's job") require.NotEmpty(t, payload) out, err := muxl.RunMuxlVerify(ctx, bytes.NewReader(payload)) require.NoError(t, err, "segment %d verify", segs) require.NotContains(t, out, `"validation_state":"Invalid"`, "segment %d must validate", segs) + events, err := unwrapMuxlEvents(ctx, payload) + require.NoError(t, err) + start, ok := gopStartMs(catalogAndSegment(events)) + require.True(t, ok) + segmentStarts = append(segmentStarts, start) // With a node key the worker completes to dual-codec: every segment must // carry both the source Opus and a worker-transcoded AAC track. @@ -171,6 +185,9 @@ func TestRunMP4IngestWorkerProducesValidSignedFrames(t *testing.T) { segs++ } require.GreaterOrEqual(t, segs, 1, "worker emitted at least one signed dual-codec segment") + // Main finds each completed segment's archival captions by the GoP start it + // reads from that segment, so every GoP needs exactly one entry under it. + require.ElementsMatch(t, segmentStarts, archiveStarts, "one archival captions frame per GoP, keyed by its media start") t.Logf("worker emitted %d valid dual-codec segments", segs) } diff --git a/pkg/media/media.go b/pkg/media/media.go index 8b92082de..1a474cca7 100644 --- a/pkg/media/media.go +++ b/pkg/media/media.go @@ -131,6 +131,9 @@ type NewSegmentNotification struct { Muxl []byte Metadata *SegmentMetadata Local bool + // archive re-times this local segment's captions for its recorded copy; + // see ArchiveCopy. + archive *captionArchive } func RunSelfTest(ctx context.Context) error { diff --git a/pkg/media/media_signer.go b/pkg/media/media_signer.go index 7ba76ade2..c93593180 100644 --- a/pkg/media/media_signer.go +++ b/pkg/media/media_signer.go @@ -205,6 +205,14 @@ func (ms *MediaSignerLocal) SignSegmentStream(ctx context.Context, input io.Read in.Sign = muxl.SignerToCallback(ms.Signer, 32) span.SetAttributes(attribute.String("backend", "host-callback")) } + // A session whose segments are recorded masters a second, archival + // layout of its captions; see captionArchive. + if archive := captionArchiveFrom(ctx); archive != nil { + defer archive.finish() + manifest := func() ([]byte, error) { return ms.buildManifest(ctx, time.Now().UnixMilli()) } + master.archiveTo(muxl.SignerInput{CertPEM: in.CertPEM, KeyPEM: in.KeyPEM, Sign: in.Sign, TrackManifestFn: manifest}, archive.put) + defer master.awaitArchive() + } input, finish := master.tee(input) defer finish() diff --git a/pkg/media/segmenter.go b/pkg/media/segmenter.go index 3d8586d6b..149ab0c90 100644 --- a/pkg/media/segmenter.go +++ b/pkg/media/segmenter.go @@ -187,6 +187,11 @@ func (mm *MediaManager) SegmentAndSignElem(ctx context.Context, ms MediaSigner) // feeding the restarted timeline into the previous session's encoder. ctx = withIngestSession(ctx, mm.nextIngestSession()) ctx = withCaptionManager(ctx, mm) + if mm.cli.S3Configured() { + // Recorded segments get their captions laid out again once recognition + // settles; live viewers get them as soon as they are signed. + ctx = withCaptionArchive(ctx, newCaptionArchive(mm.cli.CaptionsMasterDelay)) + } // muxl path: stream the fMP4 through the per-segment signer. Each GoP // arrives as a bare canonical .m4s, which ValidateMP4 verifies, archives diff --git a/pkg/media/transcode.go b/pkg/media/transcode.go index ebd3ee104..849969049 100644 --- a/pkg/media/transcode.go +++ b/pkg/media/transcode.go @@ -308,8 +308,17 @@ func collectMuxlEvents(run func(chan *muxl.MuxlEvent) error) ([]*muxl.MuxlEvent, // bytes (from the first segment/signed-segment event) out of a muxl event // stream. func catalogAndTracks(events []*muxl.MuxlEvent) (*muxl.MuxlCatalog, map[string][]byte) { + cat, segment := catalogAndSegment(events) + if segment == nil { + return cat, nil + } + return cat, segment.Tracks +} + +// catalogAndSegment is catalogAndTracks with the whole first segment event. +func catalogAndSegment(events []*muxl.MuxlEvent) (*muxl.MuxlCatalog, *muxl.MuxlEvent) { var cat *muxl.MuxlCatalog - var tracks map[string][]byte + var segment *muxl.MuxlEvent for _, ev := range events { switch ev.Type { case "init": @@ -317,12 +326,12 @@ func catalogAndTracks(events []*muxl.MuxlEvent) (*muxl.MuxlCatalog, map[string][ cat = ev.Catalog } case "segment", "signed-segment": - if tracks == nil && len(ev.Tracks) > 0 { - tracks = ev.Tracks + if segment == nil && len(ev.Tracks) > 0 { + segment = ev } } } - return cat, tracks + return cat, segment } func maxU32(a, b uint32) uint32 { diff --git a/pkg/media/transcode_stream.go b/pkg/media/transcode_stream.go index 94f6891e5..330b59c9a 100644 --- a/pkg/media/transcode_stream.go +++ b/pkg/media/transcode_stream.go @@ -423,9 +423,9 @@ func (t *streamTranscoder) run(feedR *io.PipeReader) error { // the freshly-segmented transcoded audio (the audio track transTID within // transSeg) to a free track id so it won't collide with the source tracks, sign // it as a c2pa.transcoded derivative of the source segment's audio (node -// identity), and append it to the source segment. The relabel is a lossless -// re-container (no re-encode), so the continuous encoder's gaplessness is -// preserved. +// identity), and add it to the source segment's tracks. The relabel is a +// lossless re-container (no re-encode), so the continuous encoder's gaplessness +// is preserved. // // transSeg is the FULL emitted transcoded segment (video + transcoded audio), // not the audio track alone. The relabel goes through muxl's canonicalize, @@ -494,10 +494,10 @@ func (mm *MediaManager) finishTranscodedSegment(ctx context.Context, srcSeg, tra return nil, fmt.Errorf("sign transcoded track: %w", err) } - completed := make([]byte, 0, len(srcSeg)+len(signed)) - completed = append(completed, srcSeg...) - completed = append(completed, signed...) - return completed, nil + // The added track takes its place in ascending track-ID order, before any + // text tracks, as archives require. + tracks[strconv.FormatUint(uint64(freeTID), 10)] = signed + return concatTracksByID(tracks), nil } // audioTrackID returns the (single) audio rendition's track id in a catalog. diff --git a/pkg/media/validate.go b/pkg/media/validate.go index f51113a83..3dac29b69 100644 --- a/pkg/media/validate.go +++ b/pkg/media/validate.go @@ -84,6 +84,7 @@ type validatedSegment struct { repoDID string signingKeyDID string local bool + archive *captionArchive // the local ingest session's archival captions; see withCaptionArchive } // validateSource verifies + media-parses a bare canonical .m4s, resolves the @@ -184,6 +185,7 @@ func (mm *MediaManager) validateSource(ctx context.Context, buf []byte, local bo repoDID: repoDID, signingKeyDID: signingKeyDID, local: local, + archive: captionArchiveFrom(ctx), }, nil } @@ -282,6 +284,7 @@ func (mm *MediaManager) distributeSegment(ctx context.Context, vs *validatedSegm Muxl: seg, Metadata: meta, Local: vs.local, + archive: vs.archive, }) aqt := aqtime.FromTime(meta.StartTime.Time()) log.Log(ctx, "successfully ingested segment", "user", vs.repoDID, "signingKey", vs.signingKeyDID, "timestamp", aqt.FileSafeString(), "segmentID", vs.label) diff --git a/pkg/media/whip_worker.go b/pkg/media/whip_worker.go index 50c0465ca..b9ceb5a4d 100644 --- a/pkg/media/whip_worker.go +++ b/pkg/media/whip_worker.go @@ -83,8 +83,9 @@ func ServeWHIPIngestWorkerSocket(ctx context.Context, cfg IngestWorkerConfig) er wd := newWorkerWatchdog(ctx, ingestWorkerWatchdog, cancel, cfg.StreamerDID) defer wd.stop() - onSegment, flush := mm.workerSegmentSink(ctx, cfg, wd.wrap(srv)) - signerElem, signerDone, err := muxlSignSegmentElem(ctx, mm.cli, workerSignStream(cfg, manifest.get), onSegment) + frames := wd.wrap(srv) + onSegment, flush := mm.workerSegmentSink(ctx, cfg, frames) + signerElem, signerDone, err := muxlSignSegmentElem(ctx, mm.cli, workerSignStream(cfg, manifest.get, frames), onSegment) if err != nil { return finish(fmt.Errorf("build signer element: %w", err)) } diff --git a/pkg/muxl/muxl.go b/pkg/muxl/muxl.go index dc85d7bd6..1bc428cbe 100644 --- a/pkg/muxl/muxl.go +++ b/pkg/muxl/muxl.go @@ -337,6 +337,17 @@ func RunMuxlReadTextCues(ctx context.Context, input io.Reader, trackID uint32) ( return eng.ReadTextCues(ctx, input, trackID) } +// RunMuxlSignTextRuns mints and signs one standalone WebVTT run per track for +// the GoP span req, as the streaming signer would for that GoP, keyed by track +// ID. Exactly one of in.KeyPEM or in.Sign must be set. +func RunMuxlSignTextRuns(ctx context.Context, req TextRequest, tracks []TextTrackAttachment, in SignerInput) (map[uint32][]byte, error) { + eng, err := getEngine() + if err != nil { + return nil, err + } + return eng.SignTextRuns(ctx, req, tracks, in) +} + // FirstTFDT returns the baseMediaDecodeTime of the first tfdt box in a chunk // of ISO-BMFF boxes (a track's [c2pa uuid][muxl uuid][moof][mdat] segment, or // just the first bytes of one: a moof cut off by the end of the data is -- 2.51.2