diff --git a/internal/media/egress.go b/internal/media/egress.go index f2513f5..a641ba9 100644 --- a/internal/media/egress.go +++ b/internal/media/egress.go @@ -8,6 +8,7 @@ import ( "os" "strconv" "strings" + "sync" "time" "inno-live-server/internal/config" @@ -34,6 +35,33 @@ const ( egressStableFrames = 300 ) +// EgressPhase는 RTMP egress의 수명 단계다. +type EgressPhase string + +const ( + // EgressPhaseIdle: 첫 스폰 전, 프레임 수집·대기 중. + EgressPhaseIdle EgressPhase = "idle" + // EgressPhaseStreaming: FFmpeg 가동, 프레임 송출 중. + EgressPhaseStreaming EgressPhase = "streaming" + // EgressPhaseReconnecting: 송출 실패 후 백오프 대기 중. + EgressPhaseReconnecting EgressPhase = "reconnecting" + // EgressPhaseStopped: Run 종료(세션 teardown 또는 명시적 중지). + EgressPhaseStopped EgressPhase = "stopped" +) + +// EgressStatus는 Status()가 반환하는 egress 상태 스냅샷이다. 세션 계층이 +// StreamState 응답을 합성할 때 조회하는 pull 모델의 소스이며, egress 내부 +// 고루틴이 콜백으로 세션 락을 잡는 역방향 결합을 만들지 않기 위한 구조다. +type EgressStatus struct { + Phase EgressPhase + TargetURL string // 마스킹된 출력 URL + StartedAt *time.Time + StoppedAt *time.Time + UpdatedAt time.Time + LastError *string + ReconnectAttempts int +} + // RTMPEgress pushes processed (blurred) frames to an RTMP endpoint through a // dedicated FFmpeg child, following the FFmpegTranscoder process pattern. // @@ -60,6 +88,9 @@ type RTMPEgress struct { // resolution-derived video bitrate. bitrateOverride string input chan frame + + statusMu sync.Mutex + status EgressStatus } func NewRTMPEgress(path string, logger *slog.Logger, registry *metrics.Registry, options TranscoderOptions, outputURL string, audio *AudioPipe, latencyLog bool, audioOffset time.Duration, videoBitrate string) *RTMPEgress { @@ -78,6 +109,56 @@ func NewRTMPEgress(path string, logger *slog.Logger, registry *metrics.Registry, audioOffset: audioOffset, bitrateOverride: videoBitrate, input: make(chan frame, egressQueueSize), + status: EgressStatus{ + Phase: EgressPhaseIdle, + TargetURL: maskStreamKey(outputURL), + UpdatedAt: time.Now().UTC(), + }, + } +} + +// Status는 egress 상태 스냅샷을 반환한다. 어느 고루틴에서든 호출 가능하다. +func (e *RTMPEgress) Status() EgressStatus { + e.statusMu.Lock() + defer e.statusMu.Unlock() + return e.status +} + +// setStreaming은 스폰 성공 후 송출 단계 진입을 기록한다. StartedAt은 최초 +// 진입 시각을 보존한다(재연결·해상도 재기동으로 갱신하지 않는다). +func (e *RTMPEgress) setStreaming() { + e.statusMu.Lock() + defer e.statusMu.Unlock() + now := time.Now().UTC() + e.status.Phase = EgressPhaseStreaming + e.status.UpdatedAt = now + if e.status.StartedAt == nil { + e.status.StartedAt = &now + } +} + +// noteReconnect는 송출 실패로 백오프 대기에 들어감을 기록한다. +func (e *RTMPEgress) noteReconnect(cause error) { + e.statusMu.Lock() + defer e.statusMu.Unlock() + e.status.Phase = EgressPhaseReconnecting + e.status.UpdatedAt = time.Now().UTC() + e.status.ReconnectAttempts++ + if cause != nil { + message := cause.Error() + e.status.LastError = &message + } +} + +// setStopped는 Run 종료를 기록한다. +func (e *RTMPEgress) setStopped() { + e.statusMu.Lock() + defer e.statusMu.Unlock() + now := time.Now().UTC() + e.status.Phase = EgressPhaseStopped + e.status.UpdatedAt = now + if e.status.StoppedAt == nil { + e.status.StoppedAt = &now } } @@ -118,52 +199,77 @@ func (e *RTMPEgress) Enqueue(item frame) { // spawns the FFmpeg child, feeds it frames, and reconnects with exponential // backoff when the child dies or the RTMP write path fails. Frames arriving // while disconnected are dropped, never buffered. +// +// egress가 세션 수명으로 살게 되면서(#84) 트랙 교체로 프레임 해상도가 +// 바뀌어도 이 고루틴이 계속 담당한다: 바깥 루프가 해상도 한 세대를, +// 안쪽 루프가 그 해상도에서의 스폰·재연결을 관리한다. 해상도가 바뀐 +// 프레임을 만나면 백오프 없이 바깥 루프로 나가 새 해상도로 재측정한다. func (e *RTMPEgress) Run(ctx context.Context) { - startup, ok := e.collectStartupFrames(ctx) - if !ok { - return - } - width, height := startup[0].width, startup[0].height - fps := measureFPS(startup) - e.logger.Info("starting RTMP egress", - "url", maskStreamKey(e.outputURL), "width", width, "height", height, "fps", fps, - "video_bitrate", e.videoBitrateFor(width, height)) - + defer e.setStopped() backoff := egressBackoffMin - pending := startup + var seed []frame for ctx.Err() == nil { - process, err := e.start(ctx, width, height, fps) - if err != nil { + startup, ok := e.collectStartupFrames(ctx, seed) + if !ok { + return + } + seed = nil + width, height := startup[0].width, startup[0].height + fps := measureFPS(startup) + e.logger.Info("starting RTMP egress", + "url", maskStreamKey(e.outputURL), "width", width, "height", height, "fps", fps, + "video_bitrate", e.videoBitrateFor(width, height)) + + pending := startup + for ctx.Err() == nil { + process, err := e.start(ctx, width, height, fps) + if err != nil { + if ctx.Err() != nil { + return + } + e.logger.Error("start RTMP egress FFmpeg failed", "error", err) + e.noteReconnect(err) + e.waitBackoff(ctx, &backoff) + continue + } + e.setStreaming() + written, mismatch, err := e.writeFrames(ctx, process, pending, width, height) + pending = nil + process.close() + // Give the child EOF on pipe:3 so it exits, and free the Ogg stream so + // the next spawn attaches a fresh one. + if e.audio != nil { + e.audio.Detach(e.audioWriteEnd) + } if ctx.Err() != nil { return } - e.logger.Error("start RTMP egress FFmpeg failed", "error", err) + if mismatch != nil { + // 트랙 교체로 해상도가 바뀌었다. 송출 실패가 아니므로 백오프와 + // 재연결 집계 없이, 이 프레임을 시드로 즉시 재측정에 들어간다. + e.logger.Info("video resolution changed; restarting RTMP egress", + "old_width", width, "old_height", height, + "new_width", mismatch.width, "new_height", mismatch.height) + seed = []frame{*mismatch} + break + } + if written >= egressStableFrames { + backoff = egressBackoffMin + } + e.metrics.IncEgressReconnect() + e.noteReconnect(err) + e.logger.Warn("RTMP egress disconnected; reconnecting", + "error", err, "frames_written", written, "backoff", backoff) e.waitBackoff(ctx, &backoff) - continue - } - written, err := e.writeFrames(ctx, process, pending) - pending = nil - process.close() - // Give the child EOF on pipe:3 so it exits, and free the Ogg stream so - // the next spawn attaches a fresh one. - if e.audio != nil { - e.audio.Detach(e.audioWriteEnd) - } - if ctx.Err() != nil { - return - } - if written >= egressStableFrames { - backoff = egressBackoffMin } - e.metrics.IncEgressReconnect() - e.logger.Warn("RTMP egress disconnected; reconnecting", - "error", err, "frames_written", written, "backoff", backoff) - e.waitBackoff(ctx, &backoff) } } -func (e *RTMPEgress) collectStartupFrames(ctx context.Context) ([]frame, bool) { +// collectStartupFrames는 fps·해상도 측정에 쓸 프레임을 모은다. seed는 해상도 +// 재기동 시 이월되는 새 해상도의 첫 프레임으로, 수집 목표 수에 포함된다. +func (e *RTMPEgress) collectStartupFrames(ctx context.Context, seed []frame) ([]frame, bool) { startup := make([]frame, 0, egressMeasureFrames) + startup = append(startup, seed...) for len(startup) < egressMeasureFrames { select { case <-ctx.Done(): @@ -175,37 +281,54 @@ func (e *RTMPEgress) collectStartupFrames(ctx context.Context) ([]frame, bool) { return startup, true } -func (e *RTMPEgress) writeFrames(ctx context.Context, process *ffmpegProcess, pending []frame) (int, error) { +// writeFrames는 프레임을 FFmpeg에 공급한다. 스폰 해상도와 다른 프레임을 +// 만나면 그 프레임을 반환하고 즉시 중단한다 — 해상도가 다른 데이터를 밀면 +// 인코더가 어차피 죽고, 죽은 뒤에는 재기동 기준 해상도를 알 수 없기 때문에 +// 죽기 전에 감지해 재측정 시드로 넘기는 것이다. +func (e *RTMPEgress) writeFrames(ctx context.Context, process *ffmpegProcess, pending []frame, width, height uint16) (int, *frame, error) { written := 0 - write := func(item frame) error { + write := func(item frame) (*frame, error) { + if resolutionChanged(item, width, height) { + return &item, nil + } if !e.validFrame(item) { e.metrics.IncEgressFrameDropped() - return nil + return nil, nil } if _, err := process.stdin.Write(item.data); err != nil { - return err + return nil, err } written++ e.latency.observe(item.ingestAt) - return nil + return nil, nil } for _, item := range pending { - if err := write(item); err != nil { - return written, err + if mismatch, err := write(item); mismatch != nil || err != nil { + return written, mismatch, err } } for { select { case <-ctx.Done(): - return written, ctx.Err() + return written, nil, ctx.Err() case item := <-e.input: - if err := write(item); err != nil { - return written, err + if mismatch, err := write(item); mismatch != nil || err != nil { + return written, mismatch, err } } } } +// resolutionChanged는 프레임 해상도가 현재 스폰 해상도와 다른지 판정한다. +// 해상도 정보가 없는 프레임(0값)은 비교에서 제외해 기존 검증(validFrame) +// 경로에 맡긴다 — 0×0을 기준 삼아 재기동을 반복하는 것을 막는 가드다. +func resolutionChanged(item frame, width, height uint16) bool { + if item.width == 0 || item.height == 0 { + return false + } + return item.width != width || item.height != height +} + func (e *RTMPEgress) validFrame(item frame) bool { if e.wireFormat == config.WireFormatRaw { return len(item.data) == rawFrameSize(item.width, item.height) diff --git a/internal/media/egress_integration_test.go b/internal/media/egress_integration_test.go index 0b9ba1e..f2a601c 100644 --- a/internal/media/egress_integration_test.go +++ b/internal/media/egress_integration_test.go @@ -11,13 +11,23 @@ package media import ( + "bytes" "context" + "image" + "image/color" + "image/jpeg" + "io" + "log/slog" + "os" "os/exec" "strconv" "strings" "sync" "testing" "time" + + "inno-live-server/internal/config" + "inno-live-server/internal/metrics" ) func requireTool(t *testing.T, name string) { @@ -108,6 +118,121 @@ feed: strings.TrimSpace(ainfo), dur) } +// TestEgressIntegrationResolutionChange: 트랙 교체를 모사해 프레임 해상도를 +// 640x360 → 1280x720으로 바꾸면, egress가 백오프·재연결 집계 없이 FFmpeg를 +// 재기동하고 최종 출력이 새 해상도로 muxing되는지 실 ffmpeg로 검증한다(#84). +// +// 고정 시간 급전은 race 계측 환경에서 프레임 공급이 재측정 수(30)에 못 미쳐 +// 재스폰 전에 끝나버릴 수 있으므로, "starting RTMP egress" 로그 횟수를 관측해 +// 각 세대의 스폰이 실제로 일어난 뒤에만 다음 단계로 넘어간다. JPEG 인코딩도 +// 급전 루프 밖에서 미리 해둔다(루프 안 동기 인코딩은 race 계측 시 33ms 틱을 +// 넘겨 공급 부족을 일으킨 원인이었다). +func TestEgressIntegrationResolutionChange(t *testing.T) { + requireTool(t, "ffmpeg") + requireTool(t, "ffprobe") + + logs := &syncLogBuffer{} + logger := slog.New(slog.NewTextHandler(io.MultiWriter(os.Stderr, logs), &slog.HandlerOptions{Level: slog.LevelDebug})) + out := t.TempDir() + "/egress-resolution.flv" + egress := NewRTMPEgress("ffmpeg", logger, metrics.New(), TranscoderOptions{WireFormat: config.WireFormatJPEG}, out, nil, false, 0, "") + ctx, cancel := context.WithCancel(context.Background()) + var done sync.WaitGroup + done.Add(1) + go func() { defer done.Done(); egress.Run(ctx) }() + + const fps = 30 + smallFrame := sizedJPEG(t, 0, 640, 360) + bigFrame := sizedJPEG(t, 1, 1280, 720) + spawnCount := func() int { return strings.Count(logs.String(), "starting RTMP egress") } + + timestamp := uint32(90000) + step := uint32(videoClockRate / fps) + // feedUntil은 조건이 참이 될 때까지 프레임을 공급하고, 조건 달성 후에도 + // extra 프레임을 더 밀어 새로 뜬 FFmpeg가 유효한 출력을 쓰게 한다. + feedUntil := func(data []byte, width, height uint16, condition func() bool, extra int, label string) { + t.Helper() + ticker := time.NewTicker(time.Second / fps) + defer ticker.Stop() + deadline := time.After(60 * time.Second) + remaining := -1 + for { + select { + case <-deadline: + t.Fatalf("%s: condition not reached within 60s (spawns=%d)", label, spawnCount()) + case <-ticker.C: + egress.Enqueue(frame{data: data, timestamp: timestamp, width: width, height: height}) + timestamp += step + if remaining < 0 && condition() { + remaining = extra + continue + } + if remaining > 0 { + remaining-- + } + if remaining == 0 { + return + } + } + } + } + + // 1세대: 첫 스폰이 관측될 때까지 640x360 공급. + feedUntil(smallFrame, 640, 360, func() bool { return spawnCount() >= 1 }, 5, "first spawn") + // 2세대: 해상도 불일치 → 재측정 → 두 번째 스폰이 관측될 때까지 1280x720 + // 공급하고, FLV에 실제 프레임이 muxing되도록 1초 분량을 더 밀어준다. + feedUntil(bigFrame, 1280, 720, func() bool { return spawnCount() >= 2 }, fps, "respawn after resolution change") + cancel() + done.Wait() + + vinfo := ffprobeField(t, out, "v", "stream=width,height,codec_name") + if !strings.Contains(vinfo, "width=1280") || !strings.Contains(vinfo, "height=720") { + t.Errorf("final FLV is not 1280x720 (egress did not respawn for the new resolution):\n%s", vinfo) + } else if !strings.Contains(vinfo, "codec_name=h264") { + t.Errorf("video codec not h264:\n%s", vinfo) + } else { + t.Logf("resolution change verified: %s", strings.TrimSpace(strings.ReplaceAll(vinfo, "\n", " "))) + } + // 해상도 재기동은 송출 실패가 아니다: 재연결로 집계되면 안 된다. + if status := egress.Status(); status.ReconnectAttempts != 0 { + t.Errorf("resolution restart was counted as reconnect: attempts=%d", status.ReconnectAttempts) + } +} + +// syncLogBuffer는 여러 고루틴이 쓰는 slog 출력을 race 없이 모으는 버퍼다. +type syncLogBuffer struct { + mu sync.Mutex + buf bytes.Buffer +} + +func (b *syncLogBuffer) Write(p []byte) (int, error) { + b.mu.Lock() + defer b.mu.Unlock() + return b.buf.Write(p) +} + +func (b *syncLogBuffer) String() string { + b.mu.Lock() + defer b.mu.Unlock() + return b.buf.String() +} + +// sizedJPEG는 지정 해상도의 합성 JPEG 프레임을 만든다. harnessJPEG와 달리 +// 전역 해상도(EGRESS_WIDTH/HEIGHT)에 묶이지 않아 해상도 전환 시나리오에 쓴다. +func sizedJPEG(t *testing.T, index, width, height int) []byte { + t.Helper() + canvas := image.NewRGBA(image.Rect(0, 0, width, height)) + for y := 0; y < height; y++ { + for x := 0; x < width; x++ { + canvas.Set(x, y, color.RGBA{R: uint8(x / 3), G: uint8(y / 2), B: uint8(index * 2), A: 255}) + } + } + var encoded bytes.Buffer + if err := jpeg.Encode(&encoded, canvas, &jpeg.Options{Quality: 80}); err != nil { + t.Fatalf("encode sized frame: %v", err) + } + return encoded.Bytes() +} + func parseDuration(t *testing.T, ffprobeOut string) float64 { t.Helper() for _, line := range strings.Split(ffprobeOut, "\n") { diff --git a/internal/media/egress_test.go b/internal/media/egress_test.go index 492de91..dff2b8c 100644 --- a/internal/media/egress_test.go +++ b/internal/media/egress_test.go @@ -2,10 +2,12 @@ package media import ( "bytes" + "context" "io" "log/slog" "strings" "testing" + "time" "inno-live-server/internal/config" "inno-live-server/internal/metrics" @@ -183,6 +185,136 @@ func TestHandleStderrLineMasksKeyInErrors(t *testing.T) { } } +func TestEgressStatusTransitions(t *testing.T) { + e := newTestEgress(config.WireFormatJPEG, "rtmp://a.rtmp.youtube.com/live2/secretkey") + + initial := e.Status() + if initial.Phase != EgressPhaseIdle { + t.Fatalf("initial phase = %q, want %q", initial.Phase, EgressPhaseIdle) + } + if initial.TargetURL != "rtmp://a.rtmp.youtube.com/live2/****" { + t.Fatalf("initial TargetURL = %q, want masked URL", initial.TargetURL) + } + if initial.StartedAt != nil || initial.StoppedAt != nil { + t.Fatal("initial StartedAt/StoppedAt must be nil") + } + + e.setStreaming() + streaming := e.Status() + if streaming.Phase != EgressPhaseStreaming || streaming.StartedAt == nil { + t.Fatalf("after setStreaming: phase=%q started=%v", streaming.Phase, streaming.StartedAt) + } + firstStart := *streaming.StartedAt + + e.noteReconnect(io.ErrUnexpectedEOF) + reconnecting := e.Status() + if reconnecting.Phase != EgressPhaseReconnecting { + t.Fatalf("after noteReconnect: phase = %q", reconnecting.Phase) + } + if reconnecting.ReconnectAttempts != 1 { + t.Fatalf("ReconnectAttempts = %d, want 1", reconnecting.ReconnectAttempts) + } + if reconnecting.LastError == nil || *reconnecting.LastError != io.ErrUnexpectedEOF.Error() { + t.Fatalf("LastError = %v, want %q", reconnecting.LastError, io.ErrUnexpectedEOF) + } + + // 재연결 후 송출 재개: StartedAt은 최초 시각을 보존해야 한다. + e.setStreaming() + resumed := e.Status() + if resumed.StartedAt == nil || !resumed.StartedAt.Equal(firstStart) { + t.Fatalf("StartedAt changed across reconnect: %v -> %v", firstStart, resumed.StartedAt) + } + + e.setStopped() + stopped := e.Status() + if stopped.Phase != EgressPhaseStopped || stopped.StoppedAt == nil { + t.Fatalf("after setStopped: phase=%q stopped=%v", stopped.Phase, stopped.StoppedAt) + } + firstStop := *stopped.StoppedAt + e.setStopped() + if again := e.Status(); !again.StoppedAt.Equal(firstStop) { + t.Fatal("StoppedAt must not change on repeated setStopped") + } +} + +func TestCollectStartupFramesWithSeed(t *testing.T) { + e := newTestEgress(config.WireFormatJPEG, "out.flv") + seed := frame{timestamp: 999, width: 640, height: 360} + go func() { + for i := 0; i < egressMeasureFrames-1; i++ { + e.input <- frame{timestamp: uint32(i), width: 640, height: 360} + } + }() + startup, ok := e.collectStartupFrames(t.Context(), []frame{seed}) + if !ok { + t.Fatal("collectStartupFrames returned ok=false") + } + if len(startup) != egressMeasureFrames { + t.Fatalf("collected %d frames, want %d", len(startup), egressMeasureFrames) + } + if startup[0].timestamp != seed.timestamp { + t.Fatalf("startup[0].timestamp = %d, want seed %d", startup[0].timestamp, seed.timestamp) + } +} + +func TestCollectStartupFramesCancelled(t *testing.T) { + e := newTestEgress(config.WireFormatJPEG, "out.flv") + ctx, cancel := context.WithCancel(context.Background()) + cancel() + if _, ok := e.collectStartupFrames(ctx, nil); ok { + t.Fatal("collectStartupFrames must return ok=false on cancelled context") + } +} + +func TestResolutionChanged(t *testing.T) { + tests := []struct { + name string + item frame + width, height uint16 + want bool + }{ + {"same resolution", frame{width: 1280, height: 720}, 1280, 720, false}, + {"width changed", frame{width: 1920, height: 720}, 1280, 720, true}, + {"height changed", frame{width: 1280, height: 1080}, 1280, 720, true}, + {"both changed", frame{width: 1920, height: 1080}, 1280, 720, true}, + {"zero width skips check", frame{width: 0, height: 720}, 1280, 720, false}, + {"zero height skips check", frame{width: 1280, height: 0}, 1280, 720, false}, + } + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + if got := resolutionChanged(tc.item, tc.width, tc.height); got != tc.want { + t.Fatalf("resolutionChanged = %v, want %v", got, tc.want) + } + }) + } +} + +// TestRunStopsOnCancelBeforeSpawn: 프레임이 측정 수에 못 미치면 Run은 FFmpeg를 +// 스폰하지 않고 대기하다가, 컨텍스트 취소 시 stopped 상태로 종료해야 한다. +// 트랙이 한 번도 도착하지 않은 세션의 teardown 경로다. +func TestRunStopsOnCancelBeforeSpawn(t *testing.T) { + e := newTestEgress(config.WireFormatJPEG, "out.flv") + ctx, cancel := context.WithCancel(context.Background()) + done := make(chan struct{}) + go func() { + defer close(done) + e.Run(ctx) + }() + // 측정 수 미만의 프레임만 공급한다: 스폰 없이 수집 단계에 머문다. + for i := 0; i < 3; i++ { + e.Enqueue(frame{timestamp: uint32(i), width: 640, height: 360}) + } + cancel() + select { + case <-done: + case <-time.After(2 * time.Second): + t.Fatal("Run did not return after context cancellation") + } + if status := e.Status(); status.Phase != EgressPhaseStopped { + t.Fatalf("phase after Run = %q, want %q", status.Phase, EgressPhaseStopped) + } +} + // egressDroppedCount extracts the innolive_egress_frames_dropped_total counter. func egressDroppedCount(r *metrics.Registry) int { out := prometheusDump(r) diff --git a/internal/session/manager.go b/internal/session/manager.go index af01c65..333d9d3 100644 --- a/internal/session/manager.go +++ b/internal/session/manager.go @@ -101,7 +101,13 @@ type Session struct { processedTrackID string audioTrackID string audioPipe *media.AudioPipe - processor *media.Processor + // egress는 세션 수명으로 관리한다(#84). 트랙 수명(trackCtx)에 묶으면 + // 카메라 전환(트랙 교체) 순간 egress가 함께 종료되어 RTMP 연결이 끊기고 + // FFmpeg 재기동·백오프 동안 송출이 단절된다. egressCancel은 세션 종료와 + // 별개로 egress만 멈추는 경로(#83의 명시적 stream/stop)를 위해 분리해 둔다. + egress *media.RTMPEgress + egressCancel context.CancelFunc + processor *media.Processor ignoredTracks int offerReceivedAt time.Time answerCreatedAt time.Time @@ -257,6 +263,19 @@ func (m *Manager) CreateForUser(userID uuid.UUID, metadata map[string]string) (* s.audioPipe = media.NewAudioPipe(m.logger.With("session_id", id), m.metrics, 2) go s.audioPipe.Run(ctx) } + // egress도 트랙 도착 전에 세션 수명으로 미리 만든다(#84). Run은 첫 + // 프레임 수집(collectStartupFrames) 전에는 FFmpeg를 스폰하지 않으므로 + // 트랙 없는 세션에서는 고루틴 하나가 대기할 뿐 자원을 쓰지 않는다. + if m.cfg.YoutubeStreamKey != "" { + egressCtx, egressCancel := context.WithCancel(ctx) + s.egress = media.NewRTMPEgress(m.cfg.FFmpegPath, m.logger.With("session_id", id), m.metrics, media.TranscoderOptions{ + Gate: m.spawnGate, + WireFormat: m.cfg.AIWireFormat, + }, youtubeIngestURL+m.cfg.YoutubeStreamKey, s.audioPipe, m.cfg.EgressLatencyLog, m.cfg.EgressAudioOffset, m.cfg.EgressVideoBitrate) + s.egressCancel = egressCancel + go s.egress.Run(egressCtx) + m.logger.Info("YouTube RTMP egress enabled", "session_id", id, "url", youtubeIngestURL+"****") + } m.installHandlers(ctx, s) m.mu.Lock() m.sessions[id] = s @@ -521,16 +540,10 @@ func (m *Manager) installHandlers(ctx context.Context, s *Session) { } s.mu.Lock() s.processor = processor + // 세션 수명 egress를 재사용한다(#84): 트랙이 교체되어 이 파이프라인이 + // 새로 만들어져도 같은 egress에 Enqueue가 이어져 RTMP 연결이 유지된다. + egress := s.egress s.mu.Unlock() - var egress *media.RTMPEgress - if m.cfg.YoutubeStreamKey != "" { - egress = media.NewRTMPEgress(m.cfg.FFmpegPath, m.logger.With("session_id", s.ID), m.metrics, media.TranscoderOptions{ - Gate: m.spawnGate, - WireFormat: m.cfg.AIWireFormat, - }, youtubeIngestURL+m.cfg.YoutubeStreamKey, s.audioPipe, m.cfg.EgressLatencyLog, m.cfg.EgressAudioOffset, m.cfg.EgressVideoBitrate) - go egress.Run(trackCtx) - m.logger.Info("YouTube RTMP egress enabled", "session_id", s.ID, "url", youtubeIngestURL+"****") - } m.logger.Info("received WebRTC video track", "session_id", s.ID, "track_id", track.ID(), "codec", track.Codec().MimeType, "mode", m.cfg.PrivacyMode) trackID := track.ID() go func() { @@ -635,9 +648,31 @@ func (s *Session) Response() Response { if s.processor != nil { response.Media.AIFallbackActive = s.processor.FallbackActive() } + if s.egress != nil { + response.Stream = streamStateFromEgress(s.egress.Status(), s.rawTrackID != "") + } return response } +// streamStateFromEgress는 egress 상태 스냅샷을 API 응답 계약(StreamState)으로 +// 옮긴다. StopReason은 명시적 중지 개념이 생기는 #83에서 채운다. +func streamStateFromEgress(status media.EgressStatus, publisherActive bool) StreamState { + state := StreamState{ + Status: string(status.Phase), + StartedAt: status.StartedAt, + StoppedAt: status.StoppedAt, + UpdatedAt: status.UpdatedAt, + PublisherActive: publisherActive, + LastError: status.LastError, + ReconnectAttempts: status.ReconnectAttempts, + } + if status.TargetURL != "" { + target := status.TargetURL + state.TargetURL = &target + } + return state +} + func (s *Session) close(reason string, logger *slog.Logger) { s.mu.Lock() if s.closed { @@ -650,6 +685,11 @@ func (s *Session) close(reason string, logger *slog.Logger) { s.disconnectTimer.Stop() s.disconnectTimer = nil } + // 세션 ctx 취소로도 전파되지만, egress 종료가 세션 teardown의 일부임을 + // 명시하기 위해 전용 cancel을 직접 호출한다. + if s.egressCancel != nil { + s.egressCancel() + } s.cancel() s.mu.Unlock() if err := s.PC.Close(); err != nil { diff --git a/internal/session/manager_test.go b/internal/session/manager_test.go index 52f9665..55b839c 100644 --- a/internal/session/manager_test.go +++ b/internal/session/manager_test.go @@ -4,11 +4,13 @@ import ( "errors" "io" "log/slog" + "strings" "sync" "testing" "time" "inno-live-server/internal/config" + "inno-live-server/internal/media" "inno-live-server/internal/metrics" ) @@ -144,3 +146,79 @@ func TestReapsUnnegotiatedSessions(t *testing.T) { } t.Fatal("unnegotiated session was not reaped within the timeout window") } + +// TestCreateWithYoutubeKeyOwnsSessionScopedEgress: 전역 키가 설정되면 egress는 +// 세션 생성 시점에 세션 수명으로 만들어지고(#84 — 트랙 수명이 아니라), +// 응답 StreamState가 egress 상태를 반영하며, 세션 삭제 시 함께 종료돼야 한다. +func TestCreateWithYoutubeKeyOwnsSessionScopedEgress(t *testing.T) { + cfg := config.Config{ + PrivacyMode: config.PrivacyModeBypass, + FFmpegPath: "ffmpeg", + UDPPortMin: 42000, + UDPPortMax: 42100, + FrameQueueSize: 2, + YoutubeStreamKey: "test-stream-key", + } + manager, err := NewManager(cfg, slog.New(slog.NewTextHandler(io.Discard, nil)), metrics.New(), nil, nil) + if err != nil { + t.Fatal(err) + } + t.Cleanup(manager.CloseAll) + + created, _, err := manager.Create(nil) + if err != nil { + t.Fatal(err) + } + if created.egress == nil { + t.Fatal("session with YoutubeStreamKey must own an egress at creation") + } + if created.egressCancel == nil { + t.Fatal("session egress must have a dedicated cancel") + } + + stream := created.Response().Stream + if stream.Status != string(media.EgressPhaseIdle) { + t.Fatalf("initial Stream.Status = %q, want %q", stream.Status, media.EgressPhaseIdle) + } + if stream.TargetURL == nil { + t.Fatal("Stream.TargetURL must be populated when egress exists") + } + if strings.Contains(*stream.TargetURL, "test-stream-key") { + t.Fatalf("Stream.TargetURL leaks the stream key: %q", *stream.TargetURL) + } + if !strings.HasSuffix(*stream.TargetURL, "/****") { + t.Fatalf("Stream.TargetURL is not masked: %q", *stream.TargetURL) + } + + if err := manager.Delete(created.ID, "test"); err != nil { + t.Fatal(err) + } + deadline := time.Now().Add(2 * time.Second) + for time.Now().Before(deadline) { + if created.egress.Status().Phase == media.EgressPhaseStopped { + return + } + time.Sleep(10 * time.Millisecond) + } + t.Fatal("egress did not stop after session delete") +} + +// TestCreateWithoutYoutubeKeyKeepsIdleStream: 키가 없으면 egress를 만들지 않고 +// StreamState는 종전 기본값(idle)을 유지한다 — 기존 동작 보존 검증. +func TestCreateWithoutYoutubeKeyKeepsIdleStream(t *testing.T) { + manager := newTestManager(t, 0) + created, _, err := manager.Create(nil) + if err != nil { + t.Fatal(err) + } + if created.egress != nil || created.egressCancel != nil { + t.Fatal("session without YoutubeStreamKey must not own an egress") + } + stream := created.Response().Stream + if stream.Status != "idle" { + t.Fatalf("Stream.Status = %q, want %q", stream.Status, "idle") + } + if stream.TargetURL != nil { + t.Fatalf("Stream.TargetURL = %v, want nil", *stream.TargetURL) + } +}