From c94b5ac1d29a2a1bf961c08eeb58d3076409e71d Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Mon, 1 Jun 2026 18:37:24 -0700 Subject: [PATCH] media: buffered frame server for zero-downtime worker reconnect MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The transport core of the detach/reattach path. frameServer delivers a worker's signed segments to main over a reconnectable socket, buffering — bounded, drop-oldest, counted — whenever no client is attached. That's what makes a main restart lossless: main disconnects to upgrade, the worker keeps signing segments into the buffer, and the reconnecting main replays the buffer in order before going live. An outage longer than the buffer window drops the oldest tail loudly (droppedCount) rather than growing without bound; a write to a dead client auto-detaches and re-buffers. RunMKVIngestWorker now takes a FrameWriter interface (satisfied by both *ingestframe.Writer for the Stage-1 pipe and *frameServer for the socket path), so the worker body is transport-agnostic. Tests (deterministic, no gst/subprocess): buffer-while-disconnected then replay-in-order across a simulated restart with zero loss; drop-oldest beyond the bound with correct accounting; and the real accept loop flushing buffered frames on connect then streaming live. Co-Authored-By: Claude Opus 4.8 --- pkg/media/frame_server.go | 136 ++++++++++++++++++++++++++++ pkg/media/frame_server_test.go | 158 +++++++++++++++++++++++++++++++++ pkg/media/ingest_worker.go | 3 +- 3 files changed, 295 insertions(+), 2 deletions(-) create mode 100644 pkg/media/frame_server.go create mode 100644 pkg/media/frame_server_test.go diff --git a/pkg/media/frame_server.go b/pkg/media/frame_server.go new file mode 100644 index 000000000..8c55e71cc --- /dev/null +++ b/pkg/media/frame_server.go @@ -0,0 +1,136 @@ +package media + +import ( + "bytes" + "context" + "io" + "net" + "sync" + + "stream.place/streamplace/pkg/ingestframe" + "stream.place/streamplace/pkg/log" +) + +// 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 +// *ingestframe.Writer, so RunMKVIngestWorker is agnostic to which it gets. +type FrameWriter interface { + Segment(seg []byte) error + End() error + Error(msg string) error +} + +type bufferedFrame struct { + typ ingestframe.Type + payload []byte +} + +// frameServer delivers worker frames to the main process over a reconnectable +// transport (a per-session unix socket), buffering — bounded, drop-oldest — +// whenever no client is attached. This is the heart of the zero-downtime upgrade +// path: main disconnects for a restart, the worker keeps signing segments into +// the buffer, and the reconnecting main drains the buffer before going live, so +// segments produced during a brief restart are not lost. Drops are bounded and +// counted (a main outage longer than the buffer window loses the oldest tail, +// loudly, rather than growing without limit). +// +// Safe for concurrent push (the worker) vs attach/detach (the socket accept +// 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 +} + +// newFrameServer creates a server that buffers up to maxBuf frames while no +// client is attached. +func newFrameServer(maxBuf int) *frameServer { + return &frameServer{maxBuf: maxBuf} +} + +func (s *frameServer) push(typ ingestframe.Type, payload []byte) { + s.mu.Lock() + defer s.mu.Unlock() + if s.w != nil { + if err := s.w.WriteFrame(typ, payload); err == nil { + return + } + // Client gone; drop it and buffer this frame instead. + s.conn, s.w = nil, nil + } + s.pending = append(s.pending, bufferedFrame{typ, bytes.Clone(payload)}) + for len(s.pending) > s.maxBuf { + s.pending = s.pending[1:] + s.dropped++ + } +} + +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 } + +// dropped reports how many buffered frames were discarded because the buffer +// overflowed (main was disconnected longer than the buffer window). +func (s *frameServer) droppedCount() int { + s.mu.Lock() + defer s.mu.Unlock() + return s.dropped +} + +// attach binds a freshly-connected client, replaying buffered frames in order +// before going live. On a replay error the client is dropped and the buffer kept +// intact for the next reconnect. +func (s *frameServer) attach(conn net.Conn) { + s.mu.Lock() + defer s.mu.Unlock() + w := ingestframe.NewWriter(conn) + for _, f := range s.pending { + if err := w.WriteFrame(f.typ, f.payload); err != nil { + return + } + } + s.pending = nil + s.conn, s.w = conn, w +} + +// detachConn drops the named client if it's still the current one (a stale +// connection's teardown must not clobber a newer one that already reattached). +// The buffer and drop count are preserved. +func (s *frameServer) detachConn(conn net.Conn) { + s.mu.Lock() + defer s.mu.Unlock() + if s.conn == conn { + s.conn, s.w = nil, nil + } +} + +// 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 +// when ctx is cancelled or the listener is closed. +func serveFrameSocket(ctx context.Context, ln net.Listener, s *frameServer) { + go func() { + <-ctx.Done() + ln.Close() + }() + for { + conn, err := ln.Accept() + if err != nil { + return // listener closed + } + log.Log(ctx, "ingest worker: main attached to frame socket") + s.attach(conn) + go func(c net.Conn) { + // Main only reads frames; this drains anything it sends (nothing today) + // and unblocks when it disconnects, at which point we revert to buffering. + _, _ = io.Copy(io.Discard, c) + s.detachConn(c) + log.Log(ctx, "ingest worker: main detached from frame socket") + }(conn) + } +} diff --git a/pkg/media/frame_server_test.go b/pkg/media/frame_server_test.go new file mode 100644 index 000000000..68b4007fa --- /dev/null +++ b/pkg/media/frame_server_test.go @@ -0,0 +1,158 @@ +package media + +import ( + "context" + "fmt" + "net" + "path/filepath" + "testing" + + "github.com/stretchr/testify/require" + "stream.place/streamplace/pkg/ingestframe" +) + +// unixPair returns a connected (client, server) unix-socket pair. Unlike +// net.Pipe these are kernel-buffered, so small writes don't block on a reader — +// the frameServer can flush its buffer without a concurrent drainer. +func unixPair(t *testing.T) (client net.Conn, server net.Conn) { + t.Helper() + sock := filepath.Join(t.TempDir(), "p.sock") + ln, err := net.Listen("unix", sock) + require.NoError(t, err) + defer ln.Close() + type res struct { + c net.Conn + e error + } + ch := make(chan res, 1) + go func() { + c, e := ln.Accept() + ch <- res{c, e} + }() + client, err = net.Dial("unix", sock) + require.NoError(t, err) + r := <-ch + require.NoError(t, r.e) + t.Cleanup(func() { client.Close(); r.c.Close() }) + return client, r.c +} + +func seg(i int) []byte { return []byte(fmt.Sprintf("seg-%04d", i)) } + +func readSegs(t *testing.T, r *ingestframe.Reader, n int) [][]byte { + t.Helper() + out := make([][]byte, 0, n) + for i := 0; i < n; i++ { + typ, payload, err := r.ReadFrame() + require.NoError(t, err, "read frame %d", i) + require.Equal(t, ingestframe.Segment, typ) + out = append(out, payload) + } + return out +} + +// TestFrameServerBufferFlushReconnect is the zero-downtime guarantee at the +// transport level: segments produced while main is disconnected are buffered and +// replayed, in order, when main reconnects — nothing is lost across the gap. The +// attach/detach are driven explicitly so the assertion is deterministic. +func TestFrameServerBufferFlushReconnect(t *testing.T) { + srv := newFrameServer(1000) // ample buffer: no drops + + // Detached: the first 3 segments buffer. + for i := 0; i < 3; i++ { + require.NoError(t, srv.Segment(seg(i))) + } + + // Main connects: buffered 0,1,2 replayed, then 3,4 live. + clientA, serverA := unixPair(t) + srv.attach(serverA) + for i := 3; i < 5; i++ { + require.NoError(t, srv.Segment(seg(i))) + } + got := readSegs(t, ingestframe.NewReader(clientA), 5) + for i := 0; i < 5; i++ { + require.Equal(t, seg(i), got[i], "frame %d in order before the restart", i) + } + + // Main restarts: detach + drop the connection. Segments 5,6 produced while + // it's gone must buffer, not vanish. + srv.detachConn(serverA) + clientA.Close() + serverA.Close() + for i := 5; i < 7; i++ { + require.NoError(t, srv.Segment(seg(i))) + } + + // Main reconnects on a fresh connection: buffered 5,6 replayed, then 7 live. + clientB, serverB := unixPair(t) + srv.attach(serverB) + require.NoError(t, srv.Segment(seg(7))) + got = readSegs(t, ingestframe.NewReader(clientB), 3) + for i := 5; i < 8; i++ { + require.Equal(t, seg(i), got[i-5], "frame %d replayed/live after reconnect", i) + } + + require.Equal(t, 0, srv.droppedCount(), "ample buffer drops nothing across a brief restart") +} + +// TestFrameServerDropsOldestBeyondBound: a main outage longer than the buffer +// window drops the OLDEST frames (bounded memory), loudly via droppedCount — +// never grows without limit. +func TestFrameServerDropsOldestBeyondBound(t *testing.T) { + srv := newFrameServer(3) + for i := 0; i < 6; i++ { // 0,1,2 should be dropped; 3,4,5 retained + require.NoError(t, srv.Segment(seg(i))) + } + require.Equal(t, 3, srv.droppedCount(), "oldest 3 dropped") + + client, server := unixPair(t) + srv.attach(server) + got := readSegs(t, ingestframe.NewReader(client), 3) + for i := 3; i < 6; i++ { + require.Equal(t, seg(i), got[i-3], "newest 3 survive in order") + } +} + +// TestServeFrameSocketAttachAndFlush exercises the real accept loop: frames +// buffered before any client are flushed on connect, then live frames stream. +func TestServeFrameSocketAttachAndFlush(t *testing.T) { + sock := filepath.Join(t.TempDir(), "frames.sock") + ln, err := net.Listen("unix", sock) + require.NoError(t, err) + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + + srv := newFrameServer(1000) + for i := 0; i < 3; i++ { // buffered before anyone connects + require.NoError(t, srv.Segment(seg(i))) + } + go serveFrameSocket(ctx, ln, srv) + + client, err := net.Dial("unix", sock) + require.NoError(t, err) + defer client.Close() + r := ingestframe.NewReader(client) + + // The 3 buffered frames are replayed once the accept loop attaches us. + got := readSegs(t, r, 3) + for i := 0; i < 3; i++ { + require.Equal(t, seg(i), got[i]) + } + // Reading them confirms we're attached; subsequent pushes stream live. + for i := 3; i < 6; i++ { + require.NoError(t, srv.Segment(seg(i))) + } + got = readSegs(t, r, 3) + for i := 3; i < 6; i++ { + require.Equal(t, seg(i), got[i-3]) + } + + // End frame then a clean EOF on cancel. + require.NoError(t, srv.End()) + typ, _, err := r.ReadFrame() + require.NoError(t, err) + require.Equal(t, ingestframe.End, typ) + + cancel() + ln.Close() +} diff --git a/pkg/media/ingest_worker.go b/pkg/media/ingest_worker.go index 84128fdce..a0f1c53fc 100644 --- a/pkg/media/ingest_worker.go +++ b/pkg/media/ingest_worker.go @@ -8,7 +8,6 @@ import ( "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" ) @@ -53,7 +52,7 @@ type IngestWorkerConfig struct { // 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 { +func RunMKVIngestWorker(ctx context.Context, cfg IngestWorkerConfig, stdin io.Reader, frames FrameWriter) error { gstinit.InitGST() ctx, cancel := context.WithCancel(ctx) defer cancel() -- 2.51.2