From b6848e6a4f5be37cb11dfed12721aa6227e88763 Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Mon, 1 Jun 2026 19:10:29 -0700 Subject: [PATCH] =?UTF-8?q?media:=20detached=20worker=20lifecycle=20?= =?UTF-8?q?=E2=80=94=20spawn,=20discover,=20reconnecting=20consume?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The main-side half of zero-downtime. SpawnIngestWorkerDetached launches the worker in its own session (Setsid) so it outlives a main restart, fd-passing the authed ingest connection (fd 4) and config (fd 3, written synchronously so nothing must outlive a restarting main); it is not ctx-tied or waited on here. ConsumeWorkerSocket connects to the worker's frame socket and feeds segments to an injected handler (ValidateMP4 in prod), reconnecting across transient disconnects — the worker buffers + replays while disconnected, and ValidateMP4 dedup makes replayed overlap idempotent. DiscoverWorkerSockets lists running workers under the socket dir for a restarting main to reconnect to. consumeWorkerFrames now takes an injected onSegment so the same reader serves both the fd-4 pipe supervisor and the socket consumer. TestDetachedWorkerZeroDowntime runs it end to end: spawn a DETACHED worker (asserts it leads its own process group — the property that lets it survive a main restart) fd-passing real media, discover its socket the way a restarting main would, consume via the reconnecting consumer, and verify valid signed dual-codec segments through a clean End. Plus a DiscoverWorkerSockets unit test. Remaining: wire handleIncomingStream to hijack + fd-pass the real push connection into SpawnIngestWorkerDetached, and run discovery on main startup. Co-Authored-By: Claude Opus 4.8 --- pkg/media/ingest_daemon.go | 124 ++++++++++++++++++++++++++++++++ pkg/media/ingest_daemon_test.go | 100 ++++++++++++++++++++++++++ pkg/media/ingest_supervisor.go | 23 ++++-- 3 files changed, 242 insertions(+), 5 deletions(-) create mode 100644 pkg/media/ingest_daemon.go create mode 100644 pkg/media/ingest_daemon_test.go diff --git a/pkg/media/ingest_daemon.go b/pkg/media/ingest_daemon.go new file mode 100644 index 000000000..00ea09edd --- /dev/null +++ b/pkg/media/ingest_daemon.go @@ -0,0 +1,124 @@ +package media + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "net" + "os" + "os/exec" + "path/filepath" + "strings" + "syscall" + "time" + + "stream.place/streamplace/pkg/log" +) + +// ingestReconnectBackoff paces redials when a worker is up but main's connection +// dropped (a restart in progress). +const ingestReconnectBackoff = 250 * time.Millisecond + +// SpawnIngestWorkerDetached launches an ingest worker in its OWN session +// (Setsid) so it outlives a main restart, fd-passing the (already authed) ingest +// connection as the worker's media input (fd 4) and having it serve signed +// segments over cfg.SocketPath with buffered reconnect. The config rides a pipe +// on fd 3, written synchronously so nothing has to outlive a restarting main. +// +// The worker is deliberately NOT tied to main's context and NOT waited on here: +// it self-terminates after its stream drains (or its watchdog/orphan-grace +// fires), and is reaped by init once main exits. The returned process lets a +// caller that stays alive reap it; a restarting main just lets it go. +func SpawnIngestWorkerDetached(cfg IngestWorkerConfig, media *os.File) (*os.Process, error) { + if cfg.SocketPath == "" { + return nil, fmt.Errorf("SpawnIngestWorkerDetached: empty socket path") + } + exe, err := os.Executable() + if err != nil { + return nil, err + } + cfgJSON, err := json.Marshal(cfg) + if err != nil { + return nil, err + } + cfgR, cfgW, err := os.Pipe() + if err != nil { + return nil, err + } + defer cfgR.Close() + defer cfgW.Close() + + cmd := exec.Command(exe, "ingest-worker") + cmd.SysProcAttr = &syscall.SysProcAttr{Setsid: true} // detach into its own session + cmd.ExtraFiles = []*os.File{cfgR, media} // → child fd 3 (config), fd 4 (media) + cmd.Stderr = os.Stderr + if err := cmd.Start(); err != nil { + return nil, fmt.Errorf("spawn ingest worker: %w", err) + } + // Small, synchronous write: the bytes sit in the pipe buffer for the worker to + // read even after we close + return (and even if main then exits/restarts). + if _, err := cfgW.Write(cfgJSON); err != nil { + _ = cmd.Process.Kill() + return nil, fmt.Errorf("send worker config: %w", err) + } + return cmd.Process, nil +} + +// ConsumeWorkerSocket connects to a worker's frame socket and feeds its segments +// to onSegment, reconnecting across transient disconnects. A worker that +// outlived a main restart keeps buffering and replays on reconnect; ValidateMP4's +// dedup makes any replayed overlap idempotent. Returns nil on a clean End, or an +// error if the socket vanishes without one (worker crashed — contained). +func (mm *MediaManager) ConsumeWorkerSocket(ctx context.Context, socketPath, streamer string, onSegment func([]byte) error) error { + connectedOnce := false + for { + conn, err := net.Dial("unix", socketPath) + if err != nil { + if !connectedOnce { + select { // worker may still be coming up + case <-ctx.Done(): + return ctx.Err() + case <-time.After(ingestReconnectBackoff): + continue + } + } + return fmt.Errorf("ingest worker socket gone before End: %w", err) + } + connectedOnce = true + sawEnd, _ := mm.consumeWorkerFrames(ctx, conn, streamer, onSegment, nil) + conn.Close() + if sawEnd { + return nil + } + if ctx.Err() != nil { + return ctx.Err() + } + // Transient disconnect: reconnect and let the worker replay its buffer. + log.Log(ctx, "ingest worker connection dropped; reconnecting", "streamer", streamer) + select { + case <-ctx.Done(): + return ctx.Err() + case <-time.After(ingestReconnectBackoff): + } + } +} + +// DiscoverWorkerSockets lists the worker frame sockets under dir — the running +// workers a restarting main should reconnect to and resume consuming. +func DiscoverWorkerSockets(dir string) ([]string, error) { + entries, err := os.ReadDir(dir) + if err != nil { + if errors.Is(err, os.ErrNotExist) { + return nil, nil + } + return nil, err + } + var socks []string + for _, e := range entries { + if strings.HasSuffix(e.Name(), ".sock") { + socks = append(socks, filepath.Join(dir, e.Name())) + } + } + return socks, nil +} diff --git a/pkg/media/ingest_daemon_test.go b/pkg/media/ingest_daemon_test.go new file mode 100644 index 000000000..7d2ab55ce --- /dev/null +++ b/pkg/media/ingest_daemon_test.go @@ -0,0 +1,100 @@ +package media + +import ( + "bytes" + "context" + "os" + "path/filepath" + "syscall" + "testing" + "time" + + "github.com/stretchr/testify/require" + "stream.place/streamplace/pkg/crypto/signers" + "stream.place/streamplace/pkg/muxl" +) + +func TestDiscoverWorkerSockets(t *testing.T) { + dir := t.TempDir() + require.NoError(t, os.WriteFile(filepath.Join(dir, "notes.txt"), []byte("x"), 0o600)) // ignored + for _, n := range []string{"a.sock", "b.sock"} { + require.NoError(t, os.WriteFile(filepath.Join(dir, n), nil, 0o600)) + } + socks, err := DiscoverWorkerSockets(dir) + require.NoError(t, err) + require.Len(t, socks, 2) + + socks, err = DiscoverWorkerSockets(filepath.Join(dir, "nope")) // missing dir → empty, no error + require.NoError(t, err) + require.Empty(t, socks) +} + +// TestDetachedWorkerZeroDowntime exercises the main-side lifecycle end to end: +// spawn a worker DETACHED (its own session, so it survives a main restart), +// fd-passing its media; discover its socket the way a restarting main would; +// consume via the reconnecting consumer; and verify it served valid signed +// dual-codec segments through a clean End. The session check confirms the +// detachment that lets the worker outlive main. +func TestDetachedWorkerZeroDowntime(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) + + dir := t.TempDir() + sock := filepath.Join(dir, "stream.sock") + cfg := IngestWorkerConfig{ + StreamerDID: ms.Streamer(), + KeyPEM: keyPEM, + CertPEM: ms.Cert, + Manifest: manifest, + NodeCertPEM: ms.Cert, + NodeKeyPEM: keyPEM, + BroadcasterHost: "test.example.com", + SocketPath: sock, + InputFD: 4, + } + + mkv := makeH264AACMKV(t, ctx, getFixture("5sec.mp4")) + mediaR, mediaW, err := os.Pipe() + require.NoError(t, err) + + proc, err := SpawnIngestWorkerDetached(cfg, mediaR) + require.NoError(t, err) + mediaR.Close() // the worker holds its own dup + go func() { + _, _ = mediaW.Write(mkv) + mediaW.Close() + }() + + // A (re)starting main discovers the running worker by scanning the socket dir. + var socks []string + require.Eventually(t, func() bool { + socks, _ = DiscoverWorkerSockets(dir) + return len(socks) == 1 + }, 15*time.Second, 100*time.Millisecond, "worker socket appears for discovery") + + // The worker runs in its own session/process group (Setsid) — the property + // that lets it outlive a main restart. + wpgid, err := syscall.Getpgid(proc.Pid) + require.NoError(t, err) + require.Equal(t, proc.Pid, wpgid, "worker leads its own process group (Setsid detached)") + require.NotEqual(t, syscall.Getpgrp(), wpgid, "worker group differs from the test's") + + // Consume through the reconnecting consumer; verify dual-codec signed output. + var segs int + onSegment := func(s []byte) error { + out, verr := muxl.RunMuxlVerify(ctx, bytes.NewReader(s)) + require.NoError(t, verr) + require.NotContains(t, out, `"validation_state":"Invalid"`) + segs++ + return nil + } + require.NoError(t, (&MediaManager{}).ConsumeWorkerSocket(ctx, socks[0], ms.Streamer(), onSegment)) + require.GreaterOrEqual(t, segs, 1, "detached worker served signed segments") + + _, _ = proc.Wait() // reap the detached worker + t.Logf("detached worker (pgid=%d) served %d dual-codec segments via discovery+reconnect", wpgid, segs) +} diff --git a/pkg/media/ingest_supervisor.go b/pkg/media/ingest_supervisor.go index 52d365eb1..0dd018b1b 100644 --- a/pkg/media/ingest_supervisor.go +++ b/pkg/media/ingest_supervisor.go @@ -142,7 +142,7 @@ func (mm *MediaManager) MKVIngestIsolated(ctx context.Context, input io.Reader, defer watchdog.Stop() // Read signed-segment frames and feed each into the normal chokepoint. - sawEnd, readErr := mm.consumeWorkerFrames(ctx, framesR, ms.Streamer(), func() { + sawEnd, readErr := mm.consumeWorkerFrames(ctx, framesR, ms.Streamer(), mm.validateSegment(ctx), func() { watchdog.Reset(ingestWorkerWatchdog) }) logsWG.Wait() @@ -165,8 +165,8 @@ func (mm *MediaManager) MKVIngestIsolated(ctx context.Context, input io.Reader, // 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, onProgress func()) (sawEnd bool, _ error) { - fr := ingestframe.NewReader(stdout) +func (mm *MediaManager) consumeWorkerFrames(ctx context.Context, r io.Reader, streamer string, onSegment func([]byte) error, onProgress func()) (sawEnd bool, _ error) { + fr := ingestframe.NewReader(r) for { typ, payload, err := fr.ReadFrame() if err != nil { @@ -180,8 +180,12 @@ func (mm *MediaManager) consumeWorkerFrames(ctx context.Context, stdout io.Reade } 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) + if onSegment != nil { + if serr := onSegment(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.End: sawEnd = true @@ -191,6 +195,15 @@ func (mm *MediaManager) consumeWorkerFrames(ctx context.Context, stdout io.Reade } } +// 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) + } +} + // 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) -- 2.51.2