Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
49 changes: 35 additions & 14 deletions src/OpenIPC.Viewer.Video/Pipeline/AutoReconnectingVideoSession.cs
Original file line number Diff line number Diff line change
Expand Up @@ -14,13 +14,18 @@ namespace OpenIPC.Viewer.Video.Pipeline;
// attempts with no frame the wrapper surfaces Offline (the interactive error
// cell) while still probing at the 30s cadence. A watchdog forces a reconnect
// when a "Playing" stream stops delivering frames — FFmpeg can sit on a dead
// RTSP socket without erroring. Auth errors (401, Unauthorized, EACCES) abort
// RTSP socket without erroring — or when an attempt never gets past
// "Connecting" (issue #77: tiles stuck until an app restart). Auth errors (401, Unauthorized, EACCES) abort
// permanently — we never retry against a wrong password (would lock the camera
// out / DDoS it). Phase 12.3.
internal sealed class AutoReconnectingVideoSession : IVideoSession
{
// No decoded frame for this long while Playing → treat the stream as hung.
private const long FrameTimeoutTicks = 5 * TimeSpan.TicksPerSecond;
private const long FrameTimeoutMs = 5_000;
// Stuck in Connecting this long → abandon the attempt. Above the inner
// session's own open + probe budgets, so it only fires when something it
// can't interrupt (e.g. HW decoder setup) wedges the connect.
private const long ConnectTimeoutMs = 45_000;
// Consecutive failed attempts (no successful frame) before going Offline.
private const int ColdFailures = 5;

Expand All @@ -36,11 +41,13 @@ internal sealed class AutoReconnectingVideoSession : IVideoSession
private SessionState _state = SessionState.Idle;
private string? _lastError;

// Watchdog shared state. _lastActivityTicks is the UTC tick of the last
// frame (or the moment Playing was reached); _watching gates the watchdog
// so it only fires while the inner session believes it is Playing.
private long _lastActivityTicks;
// Watchdog shared state. _lastActivityMs is the Environment.TickCount64 of
// the last frame (or the moment Playing/Connecting was reached) — monotonic,
// so an NTP clock step can't blind the watchdog. _watching gates the stall
// check to Playing, _connecting the connect-timeout check to Connecting.
private long _lastActivityMs;
private volatile bool _watching;
private volatile bool _connecting;
// Transition-only logging — avoids a log line per retry attempt.
private SessionState _lastLoggedState = SessionState.Idle;

Expand Down Expand Up @@ -137,7 +144,7 @@ private async Task LoopAsync(CancellationToken ct)

using var framesSub = inner.Frames.Subscribe(f =>
{
Volatile.Write(ref _lastActivityTicks, DateTime.UtcNow.Ticks);
Volatile.Write(ref _lastActivityMs, Environment.TickCount64);
// A real decoded frame means the connection is healthy — reset the
// backoff so a later blip starts again from 1s, not the capped 30s.
if (!sawFrame) { sawFrame = true; attempt = 0; }
Expand All @@ -149,13 +156,21 @@ private async Task LoopAsync(CancellationToken ct)
{
if (s == SessionState.Playing)
{
Volatile.Write(ref _lastActivityTicks, DateTime.UtcNow.Ticks);
Volatile.Write(ref _lastActivityMs, Environment.TickCount64);
_connecting = false;
_watching = true;
LogTransition(SessionState.Playing, null);
}
else if (s == SessionState.Connecting)
{
Volatile.Write(ref _lastActivityMs, Environment.TickCount64);
_watching = false;
_connecting = true;
}
else
{
_watching = false;
_connecting = false;
}
SetState(s);
if (s is SessionState.Failed or SessionState.Idle)
Expand All @@ -172,6 +187,7 @@ private async Task LoopAsync(CancellationToken ct)
{
var error = await failed.Task.ConfigureAwait(false);
_watching = false;
_connecting = false;

if (IsAuthFailure(error))
{
Expand Down Expand Up @@ -220,9 +236,10 @@ private async Task LoopAsync(CancellationToken ct)
}

// Background frame watchdog. While the inner session reports Playing, a gap
// of FrameTimeoutTicks with no decoded frame means the stream is hung — push
// the same failure path as an explicit disconnect. Disposing the returned
// handle cancels the loop.
// of FrameTimeoutMs with no decoded frame means the stream is hung; while it
// reports Connecting, ConnectTimeoutMs without reaching Playing means the
// connect is wedged. Either pushes the same failure path as an explicit
// disconnect. Disposing the returned handle cancels the loop.
private IDisposable StartWatchdog(TaskCompletionSource<string?> failed, CancellationToken ct)
{
var cts = CancellationTokenSource.CreateLinkedTokenSource(ct);
Expand All @@ -233,13 +250,17 @@ private IDisposable StartWatchdog(TaskCompletionSource<string?> failed, Cancella
while (!cts.IsCancellationRequested)
{
await Task.Delay(1000, cts.Token).ConfigureAwait(false);
if (!_watching) continue;
var idle = DateTime.UtcNow.Ticks - Volatile.Read(ref _lastActivityTicks);
if (idle > FrameTimeoutTicks)
var idle = Environment.TickCount64 - Volatile.Read(ref _lastActivityMs);
if (_watching && idle > FrameTimeoutMs)
{
failed.TrySetResult("Stream stalled (no frames for 5s)");
return;
}
if (_connecting && idle > ConnectTimeoutMs)
{
failed.TrySetResult("Connect timed out (no stream after 45s)");
return;
Comment on lines +259 to +262

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Action required

1. Wedged connections leave decoder threads alive 🐞 Bug ☼ Reliability

The new Connecting watchdog sends a timed-out attempt through inner.DisposeAsync(), which returns
after a two-second join even if the decoder thread is still blocked in native setup. When that setup
cannot be interrupted, the wrapper starts another attempt while the old thread still holds native
resources and is no longer tracked for disposal.
Agent Prompt
## Issue description
The Connecting watchdog retries after disposing an inner session, but disposal can return while its native decoder thread remains blocked. Repeated retries can leave untracked threads and native resources alive.
## Fix Focus Areas
- src/OpenIPC.Viewer.Video/Pipeline/AutoReconnectingVideoSession.cs[259-262]
- src/OpenIPC.Viewer.Video/Pipeline/AutoReconnectingVideoSession.cs[198-204]
- src/OpenIPC.Viewer.Video/Pipeline/FfmpegVideoSession.cs[167-184]
## Recommended Fix
Do not treat a timed-out inner session as fully disposed when its decoder thread has not exited. Retain and track unfinished sessions for eventual cleanup, and bound or defer retries so an uninterruptible setup cannot accumulate decoder threads.

ⓘ Copy this prompt and use it to remediate the issue with your preferred AI generation tools

}
}
}
catch (Exception) { /* cancelled on dispose */ }
Expand Down
61 changes: 58 additions & 3 deletions src/OpenIPC.Viewer.Video/Pipeline/FfmpegVideoSession.cs
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,22 @@ internal sealed class FfmpegVideoSession : IVideoSession
// so the single-camera page (EnableAudio=true) starts with audio ready.
private volatile bool _audioEnabled;

// Wall-clock budgets for each blocking libavformat call, enforced by the
// interrupt callback. The socket timeout alone doesn't bound a camera that
// accepts the TCP connection and then never answers (or trickles bytes), so
// without these avformat_open_input could block forever and the tile sat on
// "Connecting…" until an app restart (issue #77).
private const long ConnectBudgetMs = 15_000;
private const long ReadBudgetMs = 10_000;
private const long CloseBudgetMs = 1_000;
// Environment.TickCount64 deadline for the call in flight — monotonic, so an
// NTP clock step during a long run can't stretch or skip it.
private long _ioDeadlineMs = long.MaxValue;
// Set by DisposeAsync: aborts whatever FFmpeg call is blocking right now.
private volatile bool _abortIo;
// Rooted so the unmanaged function pointer stays valid for the demuxer's life.
private AVIOInterruptCB_callback? _interruptDelegate;

private int _framesDecoded;
private DateTime _lastFpsTick;
private int _framesSinceFpsTick;
Expand Down Expand Up @@ -149,6 +165,7 @@ public void Resume()

public async ValueTask DisposeAsync()
{
_abortIo = true;
_cts?.Cancel();
// Unblock the decode gate so a paused thread observes cancellation and
// exits instead of parking forever.
Expand Down Expand Up @@ -185,13 +202,26 @@ private unsafe void Run()
{
BuildOpts(&opts);
fmtCtx = ffmpeg.avformat_alloc_context();
// Wire the interrupt callback before open so the budget and Dispose
// cover the connect/handshake too, not just the read loop.
_interruptDelegate = OnInterrupt;
fmtCtx->interrupt_callback = new AVIOInterruptCB
{
callback = new AVIOInterruptCB_callback_func
{
Pointer = Marshal.GetFunctionPointerForDelegate(_interruptDelegate),
},
opaque = null,
};
var url = BuildUrlWithCredentials(_options.RtspUri, _options.Credentials);

ArmIoDeadline(ConnectBudgetMs);
var ret = ffmpeg.avformat_open_input(&fmtCtx, url, null, &opts);
FfmpegError.ThrowIfError(ret, "avformat_open_input");
ThrowIfIoError(ret, "avformat_open_input", ConnectBudgetMs);

ArmIoDeadline(ConnectBudgetMs);
ret = ffmpeg.avformat_find_stream_info(fmtCtx, null);
FfmpegError.ThrowIfError(ret, "avformat_find_stream_info");
ThrowIfIoError(ret, "avformat_find_stream_info", ConnectBudgetMs);

for (var i = 0; i < (int)fmtCtx->nb_streams; i++)
{
Expand Down Expand Up @@ -275,6 +305,7 @@ private unsafe void Run()
audioProbedNoTrack = false; // a later re-enable should retry setup
}

ArmIoDeadline(ReadBudgetMs);
ret = ffmpeg.av_read_frame(fmtCtx, packet);
if (ret < 0)
Comment on lines +308 to 310

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Remediation recommended

2. Read timeouts lose their failure reason 🐞 Bug ◔ Observability

Run arms a deadline before av_read_frame, but its negative-result branch logs the returned error
and exits without checking whether that deadline expired or setting LastError. If the read is
interrupted by the new deadline before the frame watchdog signals failure, the wrapper retries with
a null error, so the tile and reconnect log omit the timeout cause.
Agent Prompt
## Issue description
An interrupted `av_read_frame` exits as an ordinary idle transition without recording that the new read deadline expired.
## Fix Focus Areas
- src/OpenIPC.Viewer.Video/Pipeline/FfmpegVideoSession.cs[308-320]
- src/OpenIPC.Viewer.Video/Pipeline/FfmpegVideoSession.cs[416-432]
## Recommended Fix
On a negative read result, distinguish cancellation and EOF from an expired read deadline. Report an expired deadline as a timeout through the session's failure path so `LastError` reaches the reconnect wrapper and UI.

ⓘ Copy this prompt and use it to remediate the issue with your preferred AI generation tools

{
Expand Down Expand Up @@ -372,13 +403,34 @@ private unsafe void Run()
if (packet != null) { var p = packet; ffmpeg.av_packet_free(&p); }
if (codecCtx != null) { var p = codecCtx; ffmpeg.avcodec_free_context(&p); }
if (hwDeviceCtx != null) { var p = hwDeviceCtx; ffmpeg.av_buffer_unref(&p); }
// Close sends RTSP TEARDOWN — give it a short budget so a dead camera
// can't wedge the thread on the way out either.
ArmIoDeadline(CloseBudgetMs);
if (fmtCtx != null) ffmpeg.avformat_close_input(&fmtCtx);
if (opts != null) ffmpeg.av_dict_free(&opts);
}

SetState(SessionState.Idle);
}

// The RTSP demuxer reports an interrupted read as "Invalid data found", so
// key off the expired budget rather than the error code and name it as the
// timeout it is — the error banner and logs then say why.
private void ThrowIfIoError(int ret, string operation, long budgetMs)
{
if (ret < 0 && !_abortIo && Environment.TickCount64 > Volatile.Read(ref _ioDeadlineMs))
throw new TimeoutException($"{operation} timed out after {budgetMs / 1000}s");
FfmpegError.ThrowIfError(ret, operation);
}

private void ArmIoDeadline(long budgetMs) =>
Volatile.Write(ref _ioDeadlineMs, Environment.TickCount64 + budgetMs);

// Returns 1 to make libavformat abort the blocking call in flight: on
// Dispose, or once the current call has overrun its budget.
private unsafe int OnInterrupt(void* opaque) =>
_abortIo || Environment.TickCount64 > Volatile.Read(ref _ioDeadlineMs) ? 1 : 0;

private unsafe bool TryEnableHw(AVCodecContext* ctx, HwAccelHint hint, AVBufferRef** outDeviceCtx)
{
var (deviceType, hwPixFmt) = HwAccelSelector.MapToFfmpeg(hint);
Expand Down Expand Up @@ -534,7 +586,10 @@ private unsafe void BuildOpts(AVDictionary** opts)
_ => "tcp",
};
ffmpeg.av_dict_set(opts, "rtsp_transport", transport, 0);
ffmpeg.av_dict_set(opts, "stimeout", "5000000", 0); // 5s socket timeout (µs)
// 5s socket I/O timeout (µs). "stimeout" was renamed to "timeout" and is
// silently ignored by the n7.1 build we bundle, so live view had no socket
// timeout at all.
ffmpeg.av_dict_set(opts, "timeout", "5000000", 0);
ffmpeg.av_dict_set(opts, "max_delay", "200000", 0); // 200ms reorder window
ffmpeg.av_dict_set(opts, "buffer_size", "1048576", 0);
ffmpeg.av_dict_set(opts, "reorder_queue_size", "0", 0);
Expand Down
Loading