diff --git a/pkg/media/ingest_daemon.go b/pkg/media/ingest_daemon.go index bf5a0edf9..804400b49 100644 --- a/pkg/media/ingest_daemon.go +++ b/pkg/media/ingest_daemon.go @@ -15,12 +15,19 @@ import ( "github.com/google/uuid" "stream.place/streamplace/pkg/ingestframe" "stream.place/streamplace/pkg/log" + "stream.place/streamplace/pkg/spmetrics" ) // ingestReconnectBackoff paces redials when a worker is up but main's connection // dropped (a restart in progress). const ingestReconnectBackoff = 250 * time.Millisecond +// workerConnectGrace bounds the INITIAL connect to a worker's socket. A freshly +// spawned worker takes a moment to listen; a socket that never accepts within +// this window is a dead/never-started worker — or a stale socket a SIGKILLed +// worker left behind — so we give up and unlink it rather than spin forever. +const workerConnectGrace = 15 * time.Second + // 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 @@ -80,10 +87,18 @@ func SpawnIngestWorkerDetached(cfg IngestWorkerConfig, media *os.File) (*os.Proc // 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 + giveUp := time.Now().Add(workerConnectGrace) for { conn, err := net.Dial("unix", socketPath) if err != nil { if !connectedOnce { + if time.Now().After(giveUp) { + // Never came up within the grace: a dead/never-started worker, or a + // stale socket a SIGKILLed worker left behind. Unlink it so it can't + // linger or re-spin a consumer on the next restart. + _ = os.Remove(socketPath) + return fmt.Errorf("ingest worker never came up at %s: %w", socketPath, err) + } select { // worker may still be coming up case <-ctx.Done(): return ctx.Err() @@ -91,6 +106,10 @@ func (mm *MediaManager) ConsumeWorkerSocket(ctx context.Context, socketPath, str continue } } + // Connected before, now the socket's gone: the worker exited. A clean + // exit already unlinked its socket; a SIGKILL leaves it behind, so remove + // the leftover (no-op if already gone). + _ = os.Remove(socketPath) return fmt.Errorf("ingest worker socket gone before End: %w", err) } connectedOnce = true @@ -164,13 +183,17 @@ func (mm *MediaManager) MKVIngestDetached(ctx context.Context, conn net.Conn, pr if err != nil { return fmt.Errorf("spawn detached worker: %w", err) } + spmetrics.IngestWorkerStarts.WithLabelValues("mkv").Inc() err = mm.ConsumeWorkerSocket(ctx, cfg.SocketPath, ms.Streamer(), mm.validateSegment(ctx)) - if err == nil { - go func() { _, _ = proc.Wait() }() // clean end: reap the exiting worker + recordWorkerExit("mkv", err, ctx.Err()) + // Reap the worker unless we're deliberately leaving it running across a main + // restart (ctx cancel). On a clean end OR a crash the worker has exited, so + // Wait() clears the zombie; only on main shutdown do we let it stay detached + // for a restarting main to rediscover via DiscoverWorkerSockets. + if ctx.Err() == nil { + go func() { _, _ = proc.Wait() }() } - // On ctx cancel (main shutting down) we deliberately leave the detached worker - // running; a restarting main reconnects via discovery. return err } @@ -241,6 +264,7 @@ func (mm *MediaManager) WHIPIngestDetached(ctx context.Context, offerSDP string, if err != nil { return "", fmt.Errorf("spawn detached whip worker: %w", err) } + spmetrics.IngestWorkerStarts.WithLabelValues("whip").Inc() // Connect + read the SDP answer (the worker's first frame), bounded so a // wedged setup can't hang the WHIP client. @@ -271,11 +295,13 @@ func (mm *MediaManager) WHIPIngestDetached(ctx context.Context, offerSDP string, go func() { sawEnd, _ := mm.consumeWorkerFrames(ctx, fr, ms.Streamer(), mm.validateSegment(ctx), nil) conn.Close() + var exitErr error if !sawEnd && ctx.Err() == nil { // Connection dropped but the detached worker lives on — reconnect and - // drain its buffer. - _ = mm.ConsumeWorkerSocket(ctx, cfg.SocketPath, ms.Streamer(), mm.validateSegment(ctx)) + // drain its buffer. Its terminal result is the worker's true outcome. + exitErr = mm.ConsumeWorkerSocket(ctx, cfg.SocketPath, ms.Streamer(), mm.validateSegment(ctx)) } + recordWorkerExit("whip", exitErr, ctx.Err()) go func() { _, _ = proc.Wait() }() }() return answer, nil @@ -299,9 +325,11 @@ func (mm *MediaManager) ResumeDetachedWorkers(ctx context.Context) { sock := sock log.Log(ctx, "resuming detached ingest worker", "socket", sock) go func() { - if cerr := mm.ConsumeWorkerSocket(ctx, sock, "resumed", mm.validateSegment(ctx)); cerr != nil { + cerr := mm.ConsumeWorkerSocket(ctx, sock, "resumed", mm.validateSegment(ctx)) + if cerr != nil { log.Error(ctx, "resumed ingest worker ended", "socket", sock, "error", cerr) } + recordWorkerExit("resumed", cerr, ctx.Err()) }() } } diff --git a/pkg/media/ingest_supervisor.go b/pkg/media/ingest_supervisor.go index 47fa4bc49..42454874c 100644 --- a/pkg/media/ingest_supervisor.go +++ b/pkg/media/ingest_supervisor.go @@ -17,6 +17,7 @@ import ( "stream.place/streamplace/pkg/crypto/signers" "stream.place/streamplace/pkg/ingestframe" "stream.place/streamplace/pkg/log" + "stream.place/streamplace/pkg/spmetrics" ) // ingestWorkerWatchdog bounds how long an isolated worker may go without @@ -103,6 +104,7 @@ func (mm *MediaManager) MKVIngestIsolated(ctx context.Context, input io.Reader, } cfgR.Close() // the child holds its own copy now framesW.Close() // ditto; the parent only reads framesR + spmetrics.IngestWorkerStarts.WithLabelValues("mkv-fd").Inc() go func() { _, _ = cfgW.Write(cfgJSON) @@ -140,6 +142,17 @@ func (mm *MediaManager) MKVIngestIsolated(ctx context.Context, input io.Reader, logsWG.Wait() werr := cmd.Wait() + // Classify the exit for the lifecycle metrics: a crash is a read error or a + // nonzero exit without a clean End; a clean End (even with a nonzero exit) is + // clean. ctx cancel (main shutdown SIGKILLing the worker) isn't counted. + var exitErr error + if readErr != nil { + exitErr = readErr + } else if werr != nil && !sawEnd { + exitErr = werr + } + recordWorkerExit("mkv-fd", exitErr, ctx.Err()) + switch { case readErr != nil: return fmt.Errorf("ingest worker stream: %w", readErr) @@ -153,6 +166,21 @@ func (mm *MediaManager) MKVIngestIsolated(ctx context.Context, input io.Reader, return nil } +// recordWorkerExit folds a worker's terminal state into the lifecycle metrics: +// a clean end vs a crash (any non-nil error). A ctx cancel (main shutting down, +// the worker forcibly stopped or deliberately left detached) is not a fault we +// count. +func recordWorkerExit(transport string, exitErr, ctxErr error) { + switch { + case ctxErr != nil: + // main shutdown — not a stream-level fault + case exitErr == nil: + spmetrics.IngestWorkerExits.WithLabelValues(transport, "clean").Inc() + default: + spmetrics.IngestWorkerExits.WithLabelValues(transport, "crash").Inc() + } +} + // 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 diff --git a/pkg/spmetrics/spmetrics.go b/pkg/spmetrics/spmetrics.go index 1113bc27f..9c25f3512 100644 --- a/pkg/spmetrics/spmetrics.go +++ b/pkg/spmetrics/spmetrics.go @@ -121,6 +121,25 @@ var FirehoseEventsDedupedTotal = promauto.NewCounterVec(prometheus.CounterOpts{ Help: "firehose events dropped as cross-relay duplicates, by kind", }, []string{"kind"}) +// --- isolated ingest workers ------------------------------------------------ + +// IngestWorkerStarts counts isolated ingest worker subprocesses spawned, by +// transport ("mkv-fd" = fd-4 fallback, "mkv" = detached, "whip" = detached WHIP). +var IngestWorkerStarts = promauto.NewCounterVec(prometheus.CounterOpts{ + Name: "streamplace_ingest_worker_starts_total", + Help: "isolated ingest worker subprocesses spawned, by transport", +}, []string{"transport"}) + +// IngestWorkerExits counts isolated ingest worker exits by transport and outcome +// ("clean" | "crash"). A rising crash rate is the signal that a stream is +// repeatedly faulting — the contained fault the node now survives but which would +// otherwise be invisible. A worker left running across a main shutdown is not an +// exit and is not counted; "resumed" is a worker reattached after a restart. +var IngestWorkerExits = promauto.NewCounterVec(prometheus.CounterOpts{ + Name: "streamplace_ingest_worker_exits_total", + Help: "isolated ingest worker exits, by transport and outcome (clean|crash)", +}, []string{"transport", "outcome"}) + // --- VOD processing --------------------------------------------------------- // VODProcessAttemptsTotal increments once per task dequeued for VOD