diff --git a/pkg/cmd/streamplace.go b/pkg/cmd/streamplace.go index 99dae045..715c1d79 100644 --- a/pkg/cmd/streamplace.go +++ b/pkg/cmd/streamplace.go @@ -864,6 +864,12 @@ func makeIngestWorkerCommand(build *config.BuildFlags) *urfavecli.Command { return fmt.Errorf("ingest-worker: parse config: %w", err) } + // Detach/reattach transport: serve frames over a unix socket with + // buffered reconnect (survives a main restart) instead of the fd-4 pipe. + if cfg.SocketPath != "" { + return media.ServeMKVIngestWorkerSocket(ctx, cfg, os.Stdin) + } + framesFile := os.NewFile(4, "ingest-frames") if framesFile == nil { return fmt.Errorf("ingest-worker: missing frames fd 4") diff --git a/pkg/media/frame_server.go b/pkg/media/frame_server.go index 8c55e71c..47585c77 100644 --- a/pkg/media/frame_server.go +++ b/pkg/media/frame_server.go @@ -3,14 +3,27 @@ package media import ( "bytes" "context" + "fmt" "io" "net" + "os" "sync" + "time" "stream.place/streamplace/pkg/ingestframe" "stream.place/streamplace/pkg/log" ) +// workerFrameBuffer bounds how many signed segments a worker holds while main is +// disconnected — ~10 min at one segment per ~1s GoP. Beyond this the oldest are +// dropped (a main outage longer than this loses the oldest tail, loudly). +const workerFrameBuffer = 600 + +// workerDrainGrace bounds how long a worker lingers after its stream ends waiting +// for main to drain the buffer. Generous enough for a main restart/upgrade; an +// orphaned worker (main never returns) exits after it. +const workerDrainGrace = 60 * time.Second + // FrameWriter is the worker's segment sink. Stage 1 uses a direct framed pipe // (*ingestframe.Writer); the zero-downtime path uses *frameServer, which buffers // across a disconnected main and replays on reconnect. Structurally satisfied by @@ -39,12 +52,13 @@ type bufferedFrame struct { // loop). A push to an attached-but-dead client fails the write, auto-detaches, // and re-buffers that frame, so a hard main disconnect degrades to buffering. type frameServer struct { - mu sync.Mutex - pending []bufferedFrame - conn net.Conn - w *ingestframe.Writer - maxBuf int - dropped int + mu sync.Mutex + pending []bufferedFrame + conn net.Conn + w *ingestframe.Writer + maxBuf int + dropped int + everAttached bool } // newFrameServer creates a server that buffers up to maxBuf frames while no @@ -96,6 +110,47 @@ func (s *frameServer) attach(conn net.Conn) { } s.pending = nil s.conn, s.w = conn, w + s.everAttached = true +} + +// waitDrained blocks until a client has attached and the buffer is fully written +// out to it (so main has the whole stream, including the trailing End), or until +// ctx is cancelled or grace elapses. The worker calls this after the stream ends +// so it doesn't exit — discarding the in-memory buffer — before a reconnecting +// main has drained it. grace bounds an orphaned worker whose main never returns. +func (s *frameServer) waitDrained(ctx context.Context, grace time.Duration) { + deadline := time.NewTimer(grace) + defer deadline.Stop() + tick := time.NewTicker(100 * time.Millisecond) + defer tick.Stop() + for { + s.mu.Lock() + drained := s.everAttached && len(s.pending) == 0 + s.mu.Unlock() + if drained { + return + } + select { + case <-ctx.Done(): + return + case <-deadline.C: + log.Warn(ctx, "ingest worker: main never drained the frame buffer; exiting", "pending", len(s.pending)) + return + case <-tick.C: + } + } +} + +// closeConn closes the current client connection, giving main a clean EOF after +// the trailing End — the signal that the worker is done and won't reconnect. +func (s *frameServer) closeConn() { + s.mu.Lock() + c := s.conn + s.conn, s.w = nil, nil + s.mu.Unlock() + if c != nil { + c.Close() + } } // detachConn drops the named client if it's still the current one (a stale @@ -109,6 +164,47 @@ func (s *frameServer) detachConn(conn net.Conn) { } } +// ServeMKVIngestWorkerSocket runs the ingest worker, delivering its signed +// segments to main over a per-session unix socket at cfg.SocketPath with +// buffered reconnect — the zero-downtime path. It listens, serves the frame +// stream (buffering across any main disconnect), runs the ingest, frames a +// trailing End/Error, then lingers until main has drained the buffer before +// removing the socket and returning. +func ServeMKVIngestWorkerSocket(ctx context.Context, cfg IngestWorkerConfig, stdin io.Reader) error { + if cfg.SocketPath == "" { + return fmt.Errorf("ServeMKVIngestWorkerSocket: empty socket path") + } + ctx, cancel := context.WithCancel(ctx) + defer cancel() + + _ = os.Remove(cfg.SocketPath) // clear any stale socket from a prior worker + ln, err := net.Listen("unix", cfg.SocketPath) + if err != nil { + return fmt.Errorf("listen %s: %w", cfg.SocketPath, err) + } + defer func() { + ln.Close() + _ = os.Remove(cfg.SocketPath) + }() + + srv := newFrameServer(workerFrameBuffer) + go serveFrameSocket(ctx, ln, srv) + + runErr := RunMKVIngestWorker(ctx, cfg, stdin, srv) + if runErr != nil { + _ = srv.Error(runErr.Error()) + } else { + _ = srv.End() + } + + // Don't exit (and drop the in-memory buffer) until main has the whole stream, + // including the trailing End — so a brief main restart loses nothing. + srv.waitDrained(ctx, workerDrainGrace) + // Close the connection so main reads End then a clean EOF (worker is done). + srv.closeConn() + return runErr +} + // serveFrameSocket accepts client connections on ln and attaches each to the // server, replacing any prior client (main reconnecting after a restart). Each // connection is watched for close so the server reverts to buffering. Returns diff --git a/pkg/media/frame_socket_e2e_test.go b/pkg/media/frame_socket_e2e_test.go new file mode 100644 index 00000000..33734803 --- /dev/null +++ b/pkg/media/frame_socket_e2e_test.go @@ -0,0 +1,91 @@ +package media + +import ( + "bytes" + "context" + "net" + "path/filepath" + "testing" + "time" + + "github.com/stretchr/testify/require" + "stream.place/streamplace/pkg/crypto/signers" + "stream.place/streamplace/pkg/ingestframe" + "stream.place/streamplace/pkg/muxl" +) + +// TestWorkerServesFramesOverSocket drives the zero-downtime transport end-to-end +// with a REAL ingest: ServeMKVIngestWorkerSocket runs the full mux+sign+transcode +// pipeline and serves the resulting signed dual-codec segments over a unix +// socket; a client connects and reads them through to a clean End. This proves +// the socket path carries real signed media (the frameServer reconnect tests +// cover the buffer-across-restart behavior deterministically). +func TestWorkerServesFramesOverSocket(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) + + sock := filepath.Join(t.TempDir(), "ingest.sock") + cfg := IngestWorkerConfig{ + StreamerDID: ms.Streamer(), + KeyPEM: keyPEM, + CertPEM: ms.Cert, + Manifest: manifest, + NodeCertPEM: ms.Cert, + NodeKeyPEM: keyPEM, + BroadcasterHost: "test.example.com", + SocketPath: sock, + } + + mkv := makeH264AACMKV(t, ctx, getFixture("5sec.mp4")) + + serveDone := make(chan error, 1) + go func() { serveDone <- ServeMKVIngestWorkerSocket(ctx, cfg, bytes.NewReader(mkv)) }() + + // Connect once the worker's listener is up (retry the dial briefly). + var conn net.Conn + for i := 0; i < 100; i++ { + if conn, err = net.Dial("unix", sock); err == nil { + break + } + time.Sleep(50 * time.Millisecond) + } + require.NoError(t, err, "connect to worker frame socket") + defer conn.Close() + + r := ingestframe.NewReader(conn) + var segs int + var sawEnd bool + for { + typ, payload, rerr := r.ReadFrame() + if rerr != nil { + break // EOF after the worker exits + } + 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 valid", segs) + segs++ + case ingestframe.End: + sawEnd = true + case ingestframe.Error: + t.Fatalf("worker error frame: %s", payload) + } + } + + require.GreaterOrEqual(t, segs, 1, "worker served at least one signed segment over the socket") + require.True(t, sawEnd, "worker served a clean End over the socket") + + select { + case serveErr := <-serveDone: + require.NoError(t, serveErr) + case <-time.After(30 * time.Second): + t.Fatal("ServeMKVIngestWorkerSocket did not return after the stream drained") + } + t.Logf("worker served %d signed segments + End over the socket", segs) +} diff --git a/pkg/media/ingest_worker.go b/pkg/media/ingest_worker.go index a0f1c53f..1b10a6f7 100644 --- a/pkg/media/ingest_worker.go +++ b/pkg/media/ingest_worker.go @@ -40,6 +40,11 @@ type IngestWorkerConfig struct { NodeCertPEM []byte `json:"node_cert_pem,omitempty"` NodeKeyPEM []byte `json:"node_key_pem,omitempty"` BroadcasterHost string `json:"broadcaster_host,omitempty"` + + // SocketPath, when set, switches the worker to the detach/reattach transport: + // it serves frames over this unix socket with buffered reconnect (survives a + // main restart) instead of the fd-4 pipe. Empty → fd-4 pipe (Stage 1). + SocketPath string `json:"socket_path,omitempty"` } // RunMKVIngestWorker is the body of the `ingest-worker` subcommand. It reads an