diff --git a/pkg/cmd/streamplace.go b/pkg/cmd/streamplace.go index 7f166c249..aad9e0e39 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 455253814..9d9c068b8 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 0dd018b1b..cdc7ae70f 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