From 2a6f309d0dcae0acf719a37ce93b1d316bb20cc1 Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Sat, 6 Jun 2026 14:57:41 -0700 Subject: [PATCH] media: pass streamer DID on the ingest-worker argv MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Workers all showed up as a bare `libstreamplace ingest-worker` in a process listing, so an operator couldn't tell which user's session a given worker belonged to. Pass the streamer DID as a positional argv element at every spawn site (MKVIngestIsolated, SpawnIngestWorkerDetached) so `ps` reads `libstreamplace ingest-worker did:plc:...`. The DID is public, non-sensitive identity — unlike the signing key / manifest, which stay on the fd-3 config. The argv copy is purely for identification; the worker still reads the authoritative DID from fd 3. The subcommand logs it early ("ingest-worker starting") and threads it into the logger for log correlation too. Co-Authored-By: Claude Opus 4.8 --- pkg/cmd/streamplace.go | 15 ++++++++++++--- pkg/media/ingest_daemon.go | 4 +++- pkg/media/ingest_supervisor.go | 4 +++- 3 files changed, 18 insertions(+), 5 deletions(-) diff --git a/pkg/cmd/streamplace.go b/pkg/cmd/streamplace.go index 7f166c24..aad9e0e3 100644 --- a/pkg/cmd/streamplace.go +++ b/pkg/cmd/streamplace.go @@ -857,10 +857,19 @@ func makeStreamCommand(build *config.BuildFlags) *urfavecli.Command { // End frame; a fatal error emits an Error frame before exiting non-zero. func makeIngestWorkerCommand(build *config.BuildFlags) *urfavecli.Command { return &urfavecli.Command{ - Name: "ingest-worker", - Usage: "internal: per-stream isolated ingest worker (spawned by the node)", - Hidden: true, + Name: "ingest-worker", + Usage: "internal: per-stream isolated ingest worker (spawned by the node)", + ArgsUsage: "[streamer-did]", + Hidden: true, Action: func(ctx context.Context, cmd *urfavecli.Command) error { + // The streamer DID is passed on argv purely so the worker is + // identifiable in a process listing (ps); the authoritative copy + // still arrives in the fd-3 config. Thread it into the logger for + // log correlation. + if did := cmd.Args().First(); did != "" { + ctx = log.WithLogValues(ctx, "streamer", did) + log.Log(ctx, "ingest-worker starting") + } cfgFile := os.NewFile(3, "ingest-config") if cfgFile == nil { return fmt.Errorf("ingest-worker: missing config fd 3") diff --git a/pkg/media/ingest_daemon.go b/pkg/media/ingest_daemon.go index 45525381..9d9c068b 100644 --- a/pkg/media/ingest_daemon.go +++ b/pkg/media/ingest_daemon.go @@ -50,7 +50,9 @@ func SpawnIngestWorkerDetached(cfg IngestWorkerConfig, media *os.File) (*os.Proc defer cfgR.Close() defer cfgW.Close() - cmd := exec.Command(exe, "ingest-worker") + // The streamer DID rides argv (public, non-sensitive) so a worker is + // identifiable in a process listing; key material stays on fd 3. + cmd := exec.Command(exe, "ingest-worker", cfg.StreamerDID) setDetached(cmd) // own session, survives a main restart (Linux) // fd 3 = config; fd 4 = the fd-passed media connection (MKV/RTMP). WHIP owns // its own PeerConnection, so it passes no media fd. diff --git a/pkg/media/ingest_supervisor.go b/pkg/media/ingest_supervisor.go index 0dd018b1..cdc7ae70 100644 --- a/pkg/media/ingest_supervisor.go +++ b/pkg/media/ingest_supervisor.go @@ -67,7 +67,9 @@ func (mm *MediaManager) MKVIngestIsolated(ctx context.Context, input io.Reader, ctx, cancel := context.WithCancel(ctx) defer cancel() - cmd := exec.CommandContext(ctx, exe, "ingest-worker") + // The streamer DID rides argv (public, non-sensitive) so a worker is + // identifiable in a process listing; key material stays on fd 3. + cmd := exec.CommandContext(ctx, exe, "ingest-worker", cfg.StreamerDID) // Dedicated pipes: fd 3 carries the config in, fd 4 carries the frame stream // out. Keeping frames off stdout means nothing the worker (or gst, or the -- 2.51.2