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()