diff --git a/pkg/api/api_internal.go b/pkg/api/api_internal.go index 84e203d13..f3f5a5295 100644 --- a/pkg/api/api_internal.go +++ b/pkg/api/api_internal.go @@ -262,7 +262,11 @@ func (a *StreamplaceAPI) InternalHandler(ctx context.Context) (http.Handler, err return } - err = a.MediaManager.MKVIngest(reqCtx, r, mediaSigner) + if a.CLI.IsolatedIngest { + err = a.MediaManager.MKVIngestIsolated(reqCtx, r, mediaSigner) + } else { + err = a.MediaManager.MKVIngest(reqCtx, r, mediaSigner) + } if err != nil { log.Log(reqCtx, "stream error", "error", err) diff --git a/pkg/cmd/streamplace.go b/pkg/cmd/streamplace.go index e4393ee91..99dae0456 100644 --- a/pkg/cmd/streamplace.go +++ b/pkg/cmd/streamplace.go @@ -4,9 +4,11 @@ import ( "bytes" "context" "crypto/rand" + "encoding/json" "errors" "flag" "fmt" + "io" "net/url" "os" "os/signal" @@ -30,6 +32,7 @@ import ( "stream.place/streamplace/pkg/bus" "stream.place/streamplace/pkg/director" "stream.place/streamplace/pkg/gstinit" + "stream.place/streamplace/pkg/ingestframe" "stream.place/streamplace/pkg/iroh/generated/iroh_streamplace" "stream.place/streamplace/pkg/localdb" "stream.place/streamplace/pkg/log" @@ -73,6 +76,7 @@ func start(build *config.BuildFlags, platformJobs []jobFunc) error { makeSelfTestCommand(build), makeVODTestCommand(build), makeStreamCommand(build), + makeIngestWorkerCommand(build), makeLiveCommand(build), makeWhepCommand(build), makeWhipCommand(build), @@ -834,6 +838,48 @@ func makeStreamCommand(build *config.BuildFlags) *urfavecli.Command { } } +// makeIngestWorkerCommand is the per-stream isolated ingest worker (Stage 1: +// MKV/RTMP push). The node spawns it; it is not meant for direct use. It reads +// the config handshake from fd 3, the MKV media from stdin, runs the mux + sign +// pipeline, and writes signed canonical .m4s frames to fd 4 — dedicated fds so +// stray stdout/stderr can't corrupt the frame stream. A clean run ends with an +// End frame; a fatal error emits an Error frame before exiting non-zero. +func makeIngestWorkerCommand(build *config.BuildFlags) *urfavecli.Command { + return &urfavecli.Command{ + Name: "ingest-worker", + Usage: "internal: per-stream isolated ingest worker (spawned by the node)", + Hidden: true, + Action: func(ctx context.Context, cmd *urfavecli.Command) error { + cfgFile := os.NewFile(3, "ingest-config") + if cfgFile == nil { + return fmt.Errorf("ingest-worker: missing config fd 3") + } + cfgBytes, err := io.ReadAll(cfgFile) + cfgFile.Close() + if err != nil { + return fmt.Errorf("ingest-worker: read config: %w", err) + } + var cfg media.IngestWorkerConfig + if err := json.Unmarshal(cfgBytes, &cfg); err != nil { + return fmt.Errorf("ingest-worker: parse config: %w", err) + } + + framesFile := os.NewFile(4, "ingest-frames") + if framesFile == nil { + return fmt.Errorf("ingest-worker: missing frames fd 4") + } + defer framesFile.Close() + frames := ingestframe.NewWriter(framesFile) + + if err := media.RunMKVIngestWorker(ctx, cfg, os.Stdin, frames); err != nil { + _ = frames.Error(err.Error()) + return err + } + return frames.End() + }, + } +} + func makeLiveCommand(build *config.BuildFlags) *urfavecli.Command { cli := config.CLI{Build: build} liveCmd := cli.NewCommand("live") diff --git a/pkg/config/config.go b/pkg/config/config.go index 6c50de9e2..1a0ffecce 100644 --- a/pkg/config/config.go +++ b/pkg/config/config.go @@ -75,6 +75,7 @@ type CLI struct { RTMPSAddonAddr string Secure bool NoMist bool + IsolatedIngest bool MistAdminPort int MistHTTPPort int MistRTMPPort int @@ -228,6 +229,13 @@ func (cli *CLI) NewCommand(name string) *urfavecli.Command { Destination: &cli.Secure, Sources: urfavecli.EnvVars("SP_SECURE"), }, + &urfavecli.BoolFlag{ + Name: "isolated-ingest", + Usage: "Run each MKV/RTMP-push ingest in an isolated worker subprocess (fault isolation)", + Value: false, + Destination: &cli.IsolatedIngest, + Sources: urfavecli.EnvVars("SP_ISOLATED_INGEST"), + }, &urfavecli.StringFlag{ Name: "tls-cert", Usage: fmt.Sprintf(`Path to TLS certificate (default: "%s")`, filepath.Join(SPDataDir, "tls", "tls.crt")), diff --git a/pkg/ingestframe/frame.go b/pkg/ingestframe/frame.go new file mode 100644 index 000000000..31f16d390 --- /dev/null +++ b/pkg/ingestframe/frame.go @@ -0,0 +1,147 @@ +// Package ingestframe defines the wire protocol a per-stream ingest worker uses +// to stream canonical MUXL fragments back to the main streamplace process. +// +// Each incoming live stream is handled by an isolated worker subprocess that +// owns the socket, muxes + transcodes the media, and signs each GoP. It emits +// the resulting signed canonical .m4s segments to the main process as a sequence +// of typed, length-prefixed frames. +// +// The framing is deliberately transport-agnostic: today it rides the worker's +// stdout pipe, but a detached / reattachable worker (the zero-downtime-upgrade +// path, where workers keep buffering signed segments across a main restart) can +// carry the identical frames over a unix socket. Nothing above this package +// cares which. +package ingestframe + +import ( + "encoding/binary" + "fmt" + "io" + "sync" +) + +// Type identifies a frame's payload. +type Type uint8 + +const ( + // Segment carries one signed canonical .m4s segment — the unit main's + // ValidateMP4 ingests. Payload: the bare canonical segment bytes. + Segment Type = 1 + // End signals the worker finished the stream cleanly (graceful EOS). No + // payload. Its ABSENCE before EOF is how main tells a crash from a clean end. + End Type = 2 + // Error carries a worker-side fatal error message (UTF-8). The worker emits + // it just before exiting so main can log a cause, not a bare "worker exited". + Error Type = 3 +) + +func (t Type) String() string { + switch t { + case Segment: + return "segment" + case End: + return "end" + case Error: + return "error" + default: + return fmt.Sprintf("unknown(%d)", uint8(t)) + } +} + +// magic prefixes every frame so a desynced/corrupt stream is caught immediately +// rather than mis-parsed as a length. +var magic = [4]byte{'S', 'P', 'F', '1'} + +// MaxPayload bounds a single frame so a corrupt or hostile length can't make the +// reader allocate unboundedly. Canonical GoP segments are well under this. +const MaxPayload = 64 << 20 // 64 MiB + +const headerSize = 4 + 1 + 4 // magic + type + uint32 length + +// Writer serializes frames to an underlying stream. Safe for concurrent use: a +// worker emits segments from more than one goroutine (the source signer and the +// transcoder's completion callback), and frames must never interleave. +type Writer struct { + mu sync.Mutex + w io.Writer +} + +// NewWriter wraps w. w is typically the worker's os.Stdout. +func NewWriter(w io.Writer) *Writer { + return &Writer{w: w} +} + +// WriteFrame writes one whole frame atomically with respect to other WriteFrame +// calls on the same Writer. +func (fw *Writer) WriteFrame(t Type, payload []byte) error { + if len(payload) > MaxPayload { + return fmt.Errorf("ingestframe: payload %d exceeds max %d", len(payload), MaxPayload) + } + var hdr [headerSize]byte + copy(hdr[0:4], magic[:]) + hdr[4] = byte(t) + binary.BigEndian.PutUint32(hdr[5:9], uint32(len(payload))) + + fw.mu.Lock() + defer fw.mu.Unlock() + if _, err := fw.w.Write(hdr[:]); err != nil { + return err + } + if len(payload) > 0 { + if _, err := fw.w.Write(payload); err != nil { + return err + } + } + return nil +} + +// Segment frames a signed canonical .m4s segment. +func (fw *Writer) Segment(seg []byte) error { return fw.WriteFrame(Segment, seg) } + +// End frames a clean end-of-stream marker. +func (fw *Writer) End() error { return fw.WriteFrame(End, nil) } + +// Error frames a fatal worker-side error message. +func (fw *Writer) Error(msg string) error { return fw.WriteFrame(Error, []byte(msg)) } + +// Reader decodes frames from an underlying stream. +type Reader struct { + r io.Reader +} + +// NewReader wraps r, typically the worker's stdout pipe. +func NewReader(r io.Reader) *Reader { + return &Reader{r: r} +} + +// ReadFrame decodes the next frame. It returns io.EOF only at a clean frame +// boundary (the stream ended between frames); a stream that dies mid-frame +// surfaces as io.ErrUnexpectedEOF, so an abrupt worker death is distinguishable +// from a clean close. +func (fr *Reader) ReadFrame() (Type, []byte, error) { + var hdr [headerSize]byte + if _, err := io.ReadFull(fr.r, hdr[:]); err != nil { + // io.EOF here = clean boundary. io.ReadFull maps a partial read to + // ErrUnexpectedEOF, which we keep: a torn header is an abrupt death. + return 0, nil, err + } + if [4]byte(hdr[0:4]) != magic { + return 0, nil, fmt.Errorf("ingestframe: bad magic %q (stream desynced)", hdr[0:4]) + } + t := Type(hdr[4]) + n := binary.BigEndian.Uint32(hdr[5:9]) + if n > MaxPayload { + return 0, nil, fmt.Errorf("ingestframe: frame length %d exceeds max %d", n, MaxPayload) + } + if n == 0 { + return t, nil, nil + } + payload := make([]byte, n) + if _, err := io.ReadFull(fr.r, payload); err != nil { + if err == io.EOF { + err = io.ErrUnexpectedEOF + } + return 0, nil, err + } + return t, payload, nil +} diff --git a/pkg/ingestframe/frame_test.go b/pkg/ingestframe/frame_test.go new file mode 100644 index 000000000..70881fcf6 --- /dev/null +++ b/pkg/ingestframe/frame_test.go @@ -0,0 +1,150 @@ +package ingestframe + +import ( + "bytes" + "encoding/binary" + "errors" + "fmt" + "io" + "sort" + "sync" + "testing" + + "github.com/stretchr/testify/require" +) + +// TestRoundTrip writes a mix of frame types/sizes and reads them back verbatim, +// then confirms a clean EOF at the boundary after the last frame. +func TestRoundTrip(t *testing.T) { + var buf bytes.Buffer + w := NewWriter(&buf) + + big := bytes.Repeat([]byte{0xAB}, 500_000) + require.NoError(t, w.Segment([]byte("seg-one"))) + require.NoError(t, w.Segment(nil)) // zero-length segment is legal + require.NoError(t, w.Segment(big)) + require.NoError(t, w.Error("something broke")) + require.NoError(t, w.End()) + + r := NewReader(&buf) + + assertFrame := func(wantT Type, wantPayload []byte) { + t.Helper() + gotT, got, err := r.ReadFrame() + require.NoError(t, err) + require.Equal(t, wantT, gotT) + require.Equal(t, wantPayload, got) + } + assertFrame(Segment, []byte("seg-one")) + assertFrame(Segment, nil) + assertFrame(Segment, big) + assertFrame(Error, []byte("something broke")) + assertFrame(End, nil) + + // Clean boundary after the last frame. + _, _, err := r.ReadFrame() + require.ErrorIs(t, err, io.EOF) +} + +// TestTruncatedFrameIsUnexpectedEOF is the crash-vs-clean-end distinction the +// supervisor relies on: a worker that dies mid-segment must NOT look like a +// graceful end. +func TestTruncatedFrameIsUnexpectedEOF(t *testing.T) { + var buf bytes.Buffer + require.NoError(t, NewWriter(&buf).Segment(bytes.Repeat([]byte{1}, 1000))) + + // Lop off the back half of the payload — an abrupt death mid-frame. + full := buf.Bytes() + torn := full[:len(full)-400] + + _, _, err := NewReader(bytes.NewReader(torn)).ReadFrame() + require.ErrorIs(t, err, io.ErrUnexpectedEOF) +} + +// TestTornHeaderIsUnexpectedEOF: dying partway through the header is also an +// abrupt death, not a clean boundary. +func TestTornHeaderIsUnexpectedEOF(t *testing.T) { + var buf bytes.Buffer + require.NoError(t, NewWriter(&buf).End()) + torn := buf.Bytes()[:headerSize-2] + + _, _, err := NewReader(bytes.NewReader(torn)).ReadFrame() + require.ErrorIs(t, err, io.ErrUnexpectedEOF) +} + +// TestBadMagicRejected: a desynced/corrupt stream is caught, not mis-parsed. +func TestBadMagicRejected(t *testing.T) { + junk := append([]byte("XXXX"), make([]byte, headerSize)...) + _, _, err := NewReader(bytes.NewReader(junk)).ReadFrame() + require.Error(t, err) + require.Contains(t, err.Error(), "bad magic") +} + +// TestOversizeLengthRejected: a hostile length can't trigger an unbounded alloc. +func TestOversizeLengthRejected(t *testing.T) { + var hdr [headerSize]byte + copy(hdr[0:4], magic[:]) + hdr[4] = byte(Segment) + binary.BigEndian.PutUint32(hdr[5:9], uint32(MaxPayload+1)) + + _, _, err := NewReader(bytes.NewReader(hdr[:])).ReadFrame() + require.Error(t, err) + require.Contains(t, err.Error(), "exceeds max") +} + +// TestWriteOversizeRejected: the writer refuses to emit an over-cap frame. +func TestWriteOversizeRejected(t *testing.T) { + err := NewWriter(io.Discard).Segment(make([]byte, MaxPayload+1)) + require.Error(t, err) + require.Contains(t, err.Error(), "exceeds max") +} + +// TestConcurrentWritesDoNotInterleave: the worker emits segments from multiple +// goroutines (source signer + transcoder completion). Frames must stay whole. +func TestConcurrentWritesDoNotInterleave(t *testing.T) { + var buf bytes.Buffer + w := NewWriter(&buf) + + const writers = 8 + const each = 50 + var wg sync.WaitGroup + for g := 0; g < writers; g++ { + wg.Add(1) + go func(g int) { + defer wg.Done() + for i := 0; i < each; i++ { + // Distinct, self-identifying payloads so interleaving is detectable. + payload := []byte(fmt.Sprintf("g%02d-i%02d-%s", g, i, bytes.Repeat([]byte("x"), i))) + require.NoError(t, w.Segment(payload)) + } + }(g) + } + wg.Wait() + + r := NewReader(&buf) + var got []string + for { + typ, payload, err := r.ReadFrame() + if errors.Is(err, io.EOF) { + break + } + require.NoError(t, err) + require.Equal(t, Segment, typ) + // Every payload must be one of the well-formed strings — a torn/interleaved + // frame would fail this prefix shape or the count. + require.Regexp(t, `^g\d\d-i\d\d-x*$`, string(payload)) + got = append(got, string(payload)) + } + require.Len(t, got, writers*each, "every frame arrives exactly once, intact") + + // And every expected payload is present exactly once. + want := make([]string, 0, writers*each) + for g := 0; g < writers; g++ { + for i := 0; i < each; i++ { + want = append(want, fmt.Sprintf("g%02d-i%02d-%s", g, i, bytes.Repeat([]byte("x"), i))) + } + } + sort.Strings(got) + sort.Strings(want) + require.Equal(t, want, got) +} diff --git a/pkg/media/ingest_subprocess_test.go b/pkg/media/ingest_subprocess_test.go new file mode 100644 index 000000000..c852325cc --- /dev/null +++ b/pkg/media/ingest_subprocess_test.go @@ -0,0 +1,126 @@ +package media + +import ( + "bytes" + "context" + "encoding/json" + "errors" + "io" + "os" + "os/exec" + "testing" + "time" + + "github.com/stretchr/testify/require" + "stream.place/streamplace/pkg/crypto/signers" + "stream.place/streamplace/pkg/ingestframe" + "stream.place/streamplace/pkg/muxl" +) + +// runIngestWorkerHelper is what the test binary becomes when re-exec'd with the +// `ingest-worker` arg (see TestMain). It mirrors makeIngestWorkerCommand exactly: +// config on fd 3, frames on fd 4, MKV on stdin; clean run ends with End, a fatal +// error with an Error frame and a non-zero exit. +func runIngestWorkerHelper() int { + cfgFile := os.NewFile(3, "ingest-config") + if cfgFile == nil { + return 2 + } + cfgBytes, err := io.ReadAll(cfgFile) + cfgFile.Close() + if err != nil { + return 2 + } + var cfg IngestWorkerConfig + if err := json.Unmarshal(cfgBytes, &cfg); err != nil { + return 2 + } + framesFile := os.NewFile(4, "ingest-frames") + if framesFile == nil { + return 2 + } + defer framesFile.Close() + frames := ingestframe.NewWriter(framesFile) + if err := RunMKVIngestWorker(context.Background(), cfg, os.Stdin, frames); err != nil { + _ = frames.Error(err.Error()) + return 1 + } + _ = frames.End() + return 0 +} + +// TestIngestWorkerSubprocess exercises the real process boundary: it spawns the +// worker as an actual subprocess (config over fd 3, MKV over stdin, frames over +// fd 4 — the exact wiring MKVIngestIsolated uses) and verifies the worker +// produces valid signed segments, a clean End frame, and a zero exit. This is +// the part the in-process worker test can't cover: fd passing, the framed wire +// protocol over a real pipe, and process lifecycle. +func TestIngestWorkerSubprocess(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) + cfgJSON, err := json.Marshal(IngestWorkerConfig{ + StreamerDID: ms.Streamer(), + KeyPEM: keyPEM, + CertPEM: ms.Cert, + Manifest: manifest, + }) + require.NoError(t, err) + + mkv := makeH264AACMKV(t, ctx, getFixture("5sec.mp4")) + + exe, err := os.Executable() + require.NoError(t, err) + cmd := exec.CommandContext(ctx, exe, "ingest-worker") + // Quiet gst in the child (it inherits the parent test's verbose leak-tracer env). + cmd.Env = append(os.Environ(), "GST_DEBUG=0", "GST_TRACERS=") + cmd.Stdin = bytes.NewReader(mkv) + cmd.Stderr = os.Stderr + + cfgR, cfgW, err := os.Pipe() + require.NoError(t, err) + framesR, framesW, err := os.Pipe() + require.NoError(t, err) + cmd.ExtraFiles = []*os.File{cfgR, framesW} // → child fd 3, fd 4 + + require.NoError(t, cmd.Start()) + cfgR.Close() + framesW.Close() + go func() { + _, _ = cfgW.Write(cfgJSON) + cfgW.Close() + }() + + fr := ingestframe.NewReader(framesR) + var segs int + var sawEnd bool + for { + typ, payload, rerr := fr.ReadFrame() + if errors.Is(rerr, io.EOF) { + break + } + require.NoError(t, rerr) + switch typ { + case ingestframe.Segment: + require.False(t, sawEnd, "no segments after End") + out, verr := muxl.RunMuxlVerify(ctx, bytes.NewReader(payload)) + require.NoError(t, verr, "segment %d verify", segs) + require.NotContains(t, out, `"validation_state":"Invalid"`, "segment %d must validate", segs) + segs++ + case ingestframe.End: + sawEnd = true + case ingestframe.Error: + t.Fatalf("worker emitted error frame: %s", payload) + } + } + framesR.Close() + + require.NoError(t, cmd.Wait(), "worker subprocess exits cleanly") + require.GreaterOrEqual(t, segs, 1, "worker subprocess emitted at least one signed segment") + require.True(t, sawEnd, "worker subprocess emitted a clean End frame") + t.Logf("worker subprocess emitted %d valid signed segments + clean End", segs) +} diff --git a/pkg/media/ingest_supervisor.go b/pkg/media/ingest_supervisor.go new file mode 100644 index 000000000..ce558d7e3 --- /dev/null +++ b/pkg/media/ingest_supervisor.go @@ -0,0 +1,209 @@ +package media + +import ( + "bufio" + "bytes" + "context" + "crypto/ecdsa" + "encoding/json" + "errors" + "fmt" + "io" + "os" + "os/exec" + "sync" + "time" + + "stream.place/streamplace/pkg/crypto/signers" + "stream.place/streamplace/pkg/ingestframe" + "stream.place/streamplace/pkg/log" +) + +// MKVIngestIsolated is the process-isolated counterpart to MKVIngest. Instead of +// running the demux + sign pipeline in this process — where a native gst fault, +// OOM, or deadlock would take the whole node down — it spawns a dedicated +// `ingest-worker` subprocess that owns the pipeline and streams signed canonical +// .m4s segments back. This process only reads frames and runs ValidateMP4 +// (memory-safe Go + wasm), so a single failing stream can at worst kill its own +// worker; the node survives. +// +// Per the locked design the worker signs everything, so main hands it the +// streamer key + cert + a once-built manifest over a dedicated config fd (kept +// off argv/env). See buildWorkerConfig for the interim key-custody note. +func (mm *MediaManager) MKVIngestIsolated(ctx context.Context, input io.Reader, ms MediaSigner) error { + cfg, err := mm.buildWorkerConfig(ctx, ms) + if err != nil { + return err + } + cfgJSON, err := json.Marshal(cfg) + if err != nil { + return fmt.Errorf("marshal worker config: %w", err) + } + + // Optional recording stays in main: tee the raw input before it reaches the + // worker, so the worker needs no data-dir access. + if shouldRecord, rerr := mm.shouldRecord(ctx, ms.Streamer()); rerr == nil && shouldRecord { + log.Log(ctx, "recording RTMP stream to file", "streamer", ms.Streamer()) + pr, pw := io.Pipe() + input = io.TeeReader(input, pw) + go func() { + if derr := mm.dumpToFile(ctx, pr, ms.Streamer(), ".rtmp.mkv"); derr != nil { + log.Error(ctx, "error dumping to file", "error", derr) + } + }() + } + + exe, err := os.Executable() + if err != nil { + return fmt.Errorf("locate self: %w", err) + } + ctx, cancel := context.WithCancel(ctx) + defer cancel() + + cmd := exec.CommandContext(ctx, exe, "ingest-worker") + + // Dedicated pipes: fd 3 carries the config in, fd 4 carries the frame stream + // out. Keeping frames off stdout means nothing the worker (or gst, or the + // self-test) writes to stdout/stderr can corrupt them; those stay plain logs. + cfgR, cfgW, err := os.Pipe() + if err != nil { + return fmt.Errorf("config pipe: %w", err) + } + defer cfgW.Close() + framesR, framesW, err := os.Pipe() + if err != nil { + cfgR.Close() + return fmt.Errorf("frames pipe: %w", err) + } + defer framesR.Close() + cmd.ExtraFiles = []*os.File{cfgR, framesW} // → child fd 3, fd 4 + + stdin, err := cmd.StdinPipe() + if err != nil { + cfgR.Close() + framesW.Close() + return err + } + stdout, err := cmd.StdoutPipe() + if err != nil { + cfgR.Close() + framesW.Close() + return err + } + stderr, err := cmd.StderrPipe() + if err != nil { + cfgR.Close() + framesW.Close() + return err + } + + if err := cmd.Start(); err != nil { + cfgR.Close() + framesW.Close() + return fmt.Errorf("start ingest worker: %w", err) + } + cfgR.Close() // the child holds its own copy now + framesW.Close() // ditto; the parent only reads framesR + + go func() { + _, _ = cfgW.Write(cfgJSON) + cfgW.Close() // EOF so the worker's config read completes + }() + + // Pump the media to the worker; closing stdin on input EOF is the worker's + // end-of-stream. Managed here (not cmd.Stdin) so a worker exit can't wedge + // cmd.Wait on a still-blocked body read. + go func() { + defer stdin.Close() + _, _ = io.Copy(stdin, input) + }() + + // Forward stdout + stderr to the node logger. Drain both fully before Wait. + var logsWG sync.WaitGroup + logsWG.Add(2) + go func() { defer logsWG.Done(); streamWorkerLogs(ctx, stdout, ms.Streamer()) }() + go func() { defer logsWG.Done(); streamWorkerLogs(ctx, stderr, ms.Streamer()) }() + + // Read signed-segment frames and feed each into the normal chokepoint. + sawEnd, readErr := mm.consumeWorkerFrames(ctx, framesR, ms.Streamer()) + logsWG.Wait() + werr := cmd.Wait() + + switch { + case readErr != nil: + return fmt.Errorf("ingest worker stream: %w", readErr) + case werr != nil && !sawEnd: + // A non-zero exit without a clean End frame means the worker died — that + // failure is contained to the subprocess; the node is unaffected. + return fmt.Errorf("ingest worker exited: %w", werr) + case werr != nil: + log.Warn(ctx, "ingest worker signalled clean end but exited nonzero", "streamer", ms.Streamer(), "error", werr) + } + return nil +} + +// 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, stdout io.Reader, streamer string) (sawEnd bool, _ error) { + fr := ingestframe.NewReader(stdout) + for { + typ, payload, err := fr.ReadFrame() + if err != nil { + if errors.Is(err, io.EOF) { + return sawEnd, nil + } + return sawEnd, err + } + switch typ { + case ingestframe.Segment: + if verr := mm.ValidateMP4(ctx, bytes.NewReader(payload), true); verr != nil { + log.Error(ctx, "ingest worker: validate segment failed", "streamer", streamer, "error", verr) + } + case ingestframe.End: + sawEnd = true + case ingestframe.Error: + log.Error(ctx, "ingest worker: reported error", "streamer", streamer, "error", string(payload)) + } + } +} + +// streamWorkerLogs forwards the worker's stderr lines into the node logger. +func streamWorkerLogs(ctx context.Context, stderr io.Reader, streamer string) { + scan := bufio.NewScanner(stderr) + scan.Buffer(make([]byte, 0, 64*1024), 1024*1024) + for scan.Scan() { + if line := scan.Text(); line != "" { + log.Log(ctx, "[ingest-worker] "+line, "streamer", streamer) + } + } +} + +// buildWorkerConfig extracts the handshake the worker needs to sign on main's +// behalf. INTERIM key custody: requires a software MediaSignerLocal — the +// MKV/RTMP push path always provides one; anything else errors so the caller can +// fall back to the in-process path. +func (mm *MediaManager) buildWorkerConfig(ctx context.Context, ms MediaSigner) (IngestWorkerConfig, error) { + local, ok := ms.(*MediaSignerLocal) + if !ok { + return IngestWorkerConfig{}, fmt.Errorf("isolated ingest requires a local signer, got %T", ms) + } + if _, ok := local.Signer.(*ecdsa.PrivateKey); !ok { + return IngestWorkerConfig{}, fmt.Errorf("isolated ingest requires a software signer, got %T", local.Signer) + } + keyPEM, err := signers.MarshalES256KPrivateKeyPEM(local.Signer) + if err != nil { + return IngestWorkerConfig{}, fmt.Errorf("marshal streamer key: %w", err) + } + manifest, err := local.buildManifest(ctx, time.Now().UnixMilli()) + if err != nil { + return IngestWorkerConfig{}, fmt.Errorf("build manifest: %w", err) + } + return IngestWorkerConfig{ + StreamerDID: ms.Streamer(), + KeyPEM: keyPEM, + CertPEM: local.Cert, + Manifest: manifest, + }, nil +} diff --git a/pkg/media/ingest_worker.go b/pkg/media/ingest_worker.go new file mode 100644 index 000000000..22d64f5dd --- /dev/null +++ b/pkg/media/ingest_worker.go @@ -0,0 +1,99 @@ +package media + +import ( + "context" + "fmt" + "io" + + "github.com/go-gst/go-gst/gst" + "stream.place/streamplace/pkg/config" + "stream.place/streamplace/pkg/gstinit" + "stream.place/streamplace/pkg/ingestframe" + "stream.place/streamplace/pkg/log" + "stream.place/streamplace/pkg/muxl" +) + +// IngestWorkerConfig is the startup handshake the main process hands an ingest +// worker over a dedicated pipe fd — kept off argv/env so key material never +// lands in a process listing. +// +// INTERIM key custody: per the locked design the worker signs everything, so it +// receives the streamer key directly. This is the deliberately-temporary +// approach; the detach/reattach work will revisit how a worker holds keys. +type IngestWorkerConfig struct { + StreamerDID string `json:"streamer_did"` + // KeyPEM is the streamer's ES256K signing key in PEM. The MKV/RTMP push path + // always yields a software key, which is what muxl-sign wants here; it is + // forwarded verbatim, no reconstruction. + KeyPEM []byte `json:"key_pem"` + CertPEM []byte `json:"cert_pem"` + // Manifest is the C2PA manifest JSON, built ONCE by main at stream start. + // muxl-sign stamps each segment's signing time into it as it signs. NOTE: + // static for the worker's lifetime — mid-stream manifest changes (e.g. a + // pre-live → live transition) don't yet cross the boundary; that needs a + // control channel and is tracked as future work. + Manifest []byte `json:"manifest"` +} + +// RunMKVIngestWorker is the body of the `ingest-worker` subcommand. It reads an +// MKV stream from stdin, runs the same demux + Opus re-encode + muxl-sign +// pipeline as the in-process MKVIngest, and emits each signed canonical .m4s +// segment to frames; the main process reads those frames and runs ValidateMP4 +// over each, exactly as if onSegment had called it directly. +// +// It returns when the stream ends cleanly (EOS) or the pipeline errors. The +// caller frames End or Error accordingly. All segment frames are guaranteed +// flushed before it returns, so a trailing End can never race ahead of the last +// Segment. +func RunMKVIngestWorker(ctx context.Context, cfg IngestWorkerConfig, stdin io.Reader, frames *ingestframe.Writer) error { + gstinit.InitGST() + ctx, cancel := context.WithCancel(ctx) + defer cancel() + + // The worker signs everything itself: forward the streamer key PEM + cert + + // prebuilt manifest straight to muxl-sign. No MediaSigner / model / DB needed. + signStream := func(ctx context.Context, input io.Reader, eventCh chan *muxl.MuxlEvent) error { + fetchManifest := func() ([]byte, error) { return cfg.Manifest, nil } + return muxl.RunMuxlSignSegment(ctx, input, muxl.SignerInput{ + CertPEM: cfg.CertPEM, + KeyPEM: cfg.KeyPEM, + TrackManifestFn: fetchManifest, + WrapperManifestFn: fetchManifest, + }, nil, nil, eventCh) + } + + onSegment := func(_ context.Context, segment []byte) error { + return frames.Segment(segment) + } + + signerElem, done, err := muxlSignSegmentElem(ctx, &config.CLI{}, signStream, onSegment) + if err != nil { + return fmt.Errorf("build signer element: %w", err) + } + pipeline, err := buildMKVIngestPipeline(ctx, stdin, signerElem) + if err != nil { + return fmt.Errorf("build pipeline: %w", err) + } + + busErr := make(chan error, 1) + go func() { + busErr <- HandleBusMessages(ctx, pipeline) + }() + + if err := pipeline.SetState(gst.StatePlaying); err != nil { + return fmt.Errorf("set playing: %w", err) + } + defer func() { + if err := pipeline.SetState(gst.StateNull); err != nil { + log.Error(ctx, "ingest worker: set null", "error", err) + } + }() + + // Wait for the pipeline to finish (EOS or error), then drain the signer: + // cancelling unblocks the signer's input pipe so it flushes the final GoP, + // and <-done guarantees every segment frame is written before we return. + pipeErr := <-busErr + cancel() + <-done + return pipeErr +} diff --git a/pkg/media/ingest_worker_test.go b/pkg/media/ingest_worker_test.go new file mode 100644 index 000000000..9eb8d6c16 --- /dev/null +++ b/pkg/media/ingest_worker_test.go @@ -0,0 +1,102 @@ +package media + +import ( + "bytes" + "context" + "errors" + "io" + "strings" + "testing" + "time" + + "github.com/go-gst/go-gst/gst" + "github.com/go-gst/go-gst/gst/app" + "github.com/stretchr/testify/require" + "stream.place/streamplace/pkg/crypto/signers" + "stream.place/streamplace/pkg/gstinit" + "stream.place/streamplace/pkg/ingestframe" + "stream.place/streamplace/pkg/muxl" +) + +// makeH264AACMKV builds a clean, single-track, streamable H264+AAC MKV from an +// H264+Opus MP4 fixture (video passed through, audio transcoded Opus→AAC). The +// repo's only AAC fixture (sample-stream.mkv) carries four audio tracks, which +// the single-audio ingest pipeline leaves three of unlinked — wedging +// matroskademux with no EOS. This produces exactly the 1-video-1-audio AAC MKV +// the RTMP push path actually delivers. +func makeH264AACMKV(t *testing.T, ctx context.Context, srcMP4 string) []byte { + t.Helper() + gstinit.InitGST() + desc := strings.Join([]string{ + "filesrc location=" + srcMP4 + " ! qtdemux name=d", + "d. ! queue ! h264parse ! matroskamux name=mux streamable=true ! appsink name=sink", + "d. ! queue ! opusdec ! audioconvert ! audioresample ! fdkaacenc ! aacparse ! mux.", + }, "\n") + pipeline, err := gst.NewPipelineFromString(desc) + require.NoError(t, err) + + sinkEle, err := pipeline.GetElementByName("sink") + require.NoError(t, err) + var buf bytes.Buffer + app.SinkFromElement(sinkEle).SetCallbacks(&app.SinkCallbacks{ + NewSampleFunc: WriterNewSample(ctx, &buf), + }) + + busErr := make(chan error, 1) + go func() { busErr <- HandleBusMessages(ctx, pipeline) }() + require.NoError(t, pipeline.SetState(gst.StatePlaying)) + defer func() { _ = pipeline.SetState(gst.StateNull) }() + require.NoError(t, <-busErr, "remux to H264+AAC MKV") + require.NotEmpty(t, buf.Bytes(), "remux produced an MKV") + return buf.Bytes() +} + +// TestRunMKVIngestWorkerProducesValidSignedFrames drives the isolated ingest +// worker's core directly (no subprocess): feed it an H264+AAC MKV, collect the +// framed output, and verify every emitted segment is a valid signed canonical +// .m4s. This is the contract the supervisor relies on — frames it can hand +// straight to ValidateMP4. (The real subprocess spawn + fault injection is +// Stage 3.) +func TestRunMKVIngestWorkerProducesValidSignedFrames(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) + + cfg := IngestWorkerConfig{ + StreamerDID: ms.Streamer(), + KeyPEM: keyPEM, + CertPEM: ms.Cert, + Manifest: manifest, + } + + mkv := makeH264AACMKV(t, ctx, getFixture("5sec.mp4")) + + // All frame writes complete before RunMKVIngestWorker returns (it waits on + // the signer drain), so reading the buffer single-threaded afterwards is safe. + var buf bytes.Buffer + frames := ingestframe.NewWriter(&buf) + require.NoError(t, RunMKVIngestWorker(ctx, cfg, bytes.NewReader(mkv), frames)) + + r := ingestframe.NewReader(&buf) + var segs int + 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") + 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) + segs++ + } + require.GreaterOrEqual(t, segs, 1, "worker emitted at least one signed segment") + t.Logf("worker emitted %d valid signed segments", segs) +} diff --git a/pkg/media/leak_test.go b/pkg/media/leak_test.go index 1d129a6c8..eb7d634d7 100644 --- a/pkg/media/leak_test.go +++ b/pkg/media/leak_test.go @@ -48,6 +48,12 @@ var LeakReportMutex sync.Mutex var LeakDoneCh = make(chan struct{}) func TestMain(m *testing.M) { + // When the parent test re-execs us as an isolated ingest worker (mirroring the + // `streamplace ingest-worker` subcommand), act as that worker and exit — before + // any leak-tracer setup, so the child stays a clean media pipeline. + if len(os.Args) > 1 && os.Args[1] == "ingest-worker" { + os.Exit(runIngestWorkerHelper()) + } if os.Getenv(IgnoreLeaks) != "" { gstinit.InitGST() os.Exit(m.Run()) diff --git a/pkg/media/mkv_ingest.go b/pkg/media/mkv_ingest.go index f4f15bccb..ca8bf5626 100644 --- a/pkg/media/mkv_ingest.go +++ b/pkg/media/mkv_ingest.go @@ -34,75 +34,76 @@ func (mm *MediaManager) MKVIngest(ctx context.Context, input io.Reader, ms Media } ctx, cancel := context.WithCancel(ctx) defer cancel() - pipelineSlice := []string{ - "appsrc name=streamsrc ! matroskademux name=demux", - "demux. ! queue ! h264parse name=parse", - "demux. ! queue ! fdkaacdec ! audioresample ! opusenc name=audioenc", - } - pipeline, err := gst.NewPipelineFromString(strings.Join(pipelineSlice, "\n")) - if err != nil { - return fmt.Errorf("error creating MKVIngest pipeline: %w", err) - } - - srcele, err := pipeline.GetElementByName("streamsrc") - if err != nil { - return err - } - // defer runtime.KeepAlive(srcele) - src := app.SrcFromElement(srcele) - src.SetCallbacks(&app.SourceCallbacks{ - NeedDataFunc: ReaderNeedDataIncremental(ctx, input), - }) - parseEle, err := pipeline.GetElementByName("parse") - if err != nil { - return err - } signer, err := mm.SegmentAndSignElem(ctx, ms) if err != nil { return err } - - err = pipeline.Add(signer) - if err != nil { - return err - } - err = parseEle.Link(signer) - if err != nil { - return err - } - audioenc, err := pipeline.GetElementByName("audioenc") - if err != nil { - return err - } - err = audioenc.Link(signer) + pipeline, err := buildMKVIngestPipeline(ctx, input, signer) if err != nil { return err } busErr := make(chan error) go func() { - err := HandleBusMessages(ctx, pipeline) - busErr <- err + busErr <- HandleBusMessages(ctx, pipeline) }() go mm.HandleKeyRevocation(ctx, ms, pipeline) - err = pipeline.SetState(gst.StatePlaying) - if err != nil { + if err := pipeline.SetState(gst.StatePlaying); err != nil { return err } - defer func() { - err := pipeline.SetState(gst.StateNull) - if err != nil { + if err := pipeline.SetState(gst.StateNull); err != nil { log.Error(ctx, "error setting pipeline to null state", "error", err) } }() - err = <-busErr + return <-busErr +} - return err +// buildMKVIngestPipeline builds the H264+AAC MKV demux graph (video → h264parse, +// audio → Opus re-encode) and links both branches into signerElem — the muxl +// signing bin that emits one bare canonical .m4s per GoP. Shared by the +// in-process MKVIngest and the isolated ingest worker, which differ only in +// where signerElem routes its segments (ValidateMP4 vs. a frame writer to the +// main process). +func buildMKVIngestPipeline(ctx context.Context, input io.Reader, signerElem *gst.Element) (*gst.Pipeline, error) { + pipelineSlice := []string{ + "appsrc name=streamsrc ! matroskademux name=demux", + "demux. ! queue ! h264parse name=parse", + "demux. ! queue ! fdkaacdec ! audioresample ! opusenc name=audioenc", + } + pipeline, err := gst.NewPipelineFromString(strings.Join(pipelineSlice, "\n")) + if err != nil { + return nil, fmt.Errorf("error creating MKVIngest pipeline: %w", err) + } + srcele, err := pipeline.GetElementByName("streamsrc") + if err != nil { + return nil, err + } + app.SrcFromElement(srcele).SetCallbacks(&app.SourceCallbacks{ + NeedDataFunc: ReaderNeedDataIncremental(ctx, input), + }) + parseEle, err := pipeline.GetElementByName("parse") + if err != nil { + return nil, err + } + if err := pipeline.Add(signerElem); err != nil { + return nil, err + } + if err := parseEle.Link(signerElem); err != nil { + return nil, err + } + audioenc, err := pipeline.GetElementByName("audioenc") + if err != nil { + return nil, err + } + if err := audioenc.Link(signerElem); err != nil { + return nil, err + } + return pipeline, nil } func (mm *MediaManager) dumpToFile(ctx context.Context, r io.Reader, user string, filesuffix string) error { diff --git a/pkg/media/muxl_segment.go b/pkg/media/muxl_segment.go index 4d6b94e2b..1de14d82f 100644 --- a/pkg/media/muxl_segment.go +++ b/pkg/media/muxl_segment.go @@ -14,6 +14,13 @@ import ( "stream.place/streamplace/pkg/muxl" ) +// SignSegmentStreamFunc drives muxl-sign's streaming per-segment signer over an +// fMP4 input, emitting one signed-segment event per GoP on eventCh. It is the +// only thing muxlSignSegmentElem needs from a signer, so the isolated ingest +// worker can supply a key-PEM-backed closure without a full MediaSigner (and +// without the model/DB a MediaSignerLocal carries). +type SignSegmentStreamFunc func(ctx context.Context, input io.Reader, eventCh chan *muxl.MuxlEvent) error + // MuxlSignSegmentElem builds the gstreamer bin that muxes the incoming // video+audio into a fragmented MP4 stream, then drives muxl-sign's streaming // per-segment signer over it. For each GoP it assembles the bare canonical @@ -23,6 +30,17 @@ import ( // produced here. Presentation headers are synthesized downstream (ValidateMP4 // / playback) only when needed. func MuxlSignSegmentElem(ctx context.Context, cli *config.CLI, ms MediaSigner, onSegment func(ctx context.Context, segment []byte) error) (*gst.Element, error) { + elem, _, err := muxlSignSegmentElem(ctx, cli, ms.SignSegmentStream, onSegment) + return elem, err +} + +// muxlSignSegmentElem is MuxlSignSegmentElem's core, parameterized by the raw +// sign-stream function and additionally returning a done channel that closes +// once every signed segment has been drained to onSegment (the signer goroutine +// has finished and the event loop has emptied). The isolated ingest worker waits +// on it to guarantee all segment frames are flushed before it signals a clean +// end-of-stream. +func muxlSignSegmentElem(ctx context.Context, cli *config.CLI, signStream SignSegmentStreamFunc, onSegment func(ctx context.Context, segment []byte) error) (*gst.Element, <-chan struct{}, error) { ctx = log.WithLogValues(ctx, "func", "MuxlSignSegmentElem") bin := gst.NewBin("muxl-segment-bin") elem, err := gst.NewElementWithProperties("mp4mux", map[string]any{ @@ -31,46 +49,46 @@ func MuxlSignSegmentElem(ctx context.Context, cli *config.CLI, ms MediaSigner, o "fragment-duration": 1, }) if err != nil { - return nil, err + return nil, nil, err } if err := bin.Add(elem); err != nil { - return nil, fmt.Errorf("failed to add mp4mux to bin: %w", err) + return nil, nil, fmt.Errorf("failed to add mp4mux to bin: %w", err) } videoPad := elem.GetRequestPad("video_%u") if videoPad == nil { - return nil, fmt.Errorf("failed to get video pad") + return nil, nil, fmt.Errorf("failed to get video pad") } videoGhost := gst.NewGhostPad("video_0", videoPad) if videoGhost == nil { - return nil, fmt.Errorf("failed to create video ghost pad") + return nil, nil, fmt.Errorf("failed to create video ghost pad") } audioPad := elem.GetRequestPad("audio_%u") if audioPad == nil { - return nil, fmt.Errorf("failed to get audio pad") + return nil, nil, fmt.Errorf("failed to get audio pad") } audioGhost := gst.NewGhostPad("audio_0", audioPad) if audioGhost == nil { - return nil, fmt.Errorf("failed to create audio ghost pad") + return nil, nil, fmt.Errorf("failed to create audio ghost pad") } if ok := bin.AddPad(videoGhost.Pad); !ok { - return nil, fmt.Errorf("failed to add video ghost pad to bin") + return nil, nil, fmt.Errorf("failed to add video ghost pad to bin") } if ok := bin.AddPad(audioGhost.Pad); !ok { - return nil, fmt.Errorf("failed to add audio ghost pad to bin") + return nil, nil, fmt.Errorf("failed to add audio ghost pad to bin") } appsink, err := gst.NewElementWithProperties("appsink", map[string]any{ "name": "muxl-appsink", }) if err != nil { - return nil, fmt.Errorf("failed to create appsink element: %w", err) + return nil, nil, fmt.Errorf("failed to create appsink element: %w", err) } if err := bin.Add(appsink); err != nil { - return nil, fmt.Errorf("failed to add appsink to bin: %w", err) + return nil, nil, fmt.Errorf("failed to add appsink to bin: %w", err) } if err := elem.Link(appsink); err != nil { - return nil, fmt.Errorf("failed to link mp4mux to appsink: %w", err) + return nil, nil, fmt.Errorf("failed to link mp4mux to appsink: %w", err) } r, w := io.Pipe() @@ -83,13 +101,15 @@ func MuxlSignSegmentElem(ctx context.Context, cli *config.CLI, ms MediaSigner, o // GoP's per-track signed canonical segments. eventCh := make(chan *muxl.MuxlEvent, 16) go func() { - err := ms.SignSegmentStream(ctx, r, eventCh) + err := signStream(ctx, r, eventCh) close(eventCh) if err != nil && ctx.Err() == nil { log.Error(ctx, "error running muxl sign-segment", "error", err) } }() + done := make(chan struct{}) go func() { + defer close(done) for ev := range eventCh { if ev.Type != "signed-segment" { continue @@ -107,7 +127,7 @@ func MuxlSignSegmentElem(ctx context.Context, cli *config.CLI, ms MediaSigner, o NewSampleFunc: WriterNewSample(ctx, w), }) - return bin.Element, nil + return bin.Element, done, nil } // concatTracksSorted joins the per-track canonical segment bytes for one GoP