Skip to content

Commit aacf19f

Browse files
authored
Merge pull request #16 from Resgrid/develop
RR1-T102 Fixes and hardening and systemd support for cli
2 parents 4edaa95 + c006395 commit aacf19f

19 files changed

Lines changed: 856 additions & 23 deletions

‎Resgrid.Audio.Relay.Console/Program.cs‎

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,7 @@
99
using System.IO;
1010
using System.Linq;
1111
using System.Reflection;
12+
using System.Runtime.InteropServices;
1213
using System.Text.Json;
1314
using System.Threading;
1415
using System.Threading.Tasks;
@@ -85,6 +86,16 @@ private static async Task<int> RunAsync()
8586
};
8687
Cli.CancelKeyPress += cancelHandler;
8788

89+
// Also shut down gracefully on SIGTERM — the default signal `systemctl stop` / `docker stop`
90+
// send. Mirrors the Ctrl+C / SIGINT path so in-flight work (recordings, LiveKit) unwinds
91+
// cleanly instead of the process being killed. Disposed before the CTS (reverse using order),
92+
// so the handler can never fire against a disposed token source.
93+
using var sigTermRegistration = PosixSignalRegistration.Create(PosixSignal.SIGTERM, context =>
94+
{
95+
context.Cancel = true;
96+
cancellationTokenSource.Cancel();
97+
});
98+
8899
try
89100
{
90101
// The engine owns all mode wiring; the console just builds the service and runs it.
Lines changed: 139 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,139 @@
1+
using System;
2+
using System.Threading;
3+
using System.Threading.Tasks;
4+
using FluentAssertions;
5+
using NUnit.Framework;
6+
using Resgrid.Relay.Engine;
7+
using Resgrid.Relay.Engine.Configuration;
8+
using Resgrid.Relay.Engine.Services;
9+
10+
namespace Resgrid.Audio.Voice.Tests
11+
{
12+
/// <summary>
13+
/// Behavioral coverage for the LiveKit-mode resilience loop in
14+
/// <see cref="RelayServiceBase"/>: exponential back-off retries, the consecutive-failure
15+
/// circuit breaker, and graceful stop while retrying. Uses a tiny test subclass with a
16+
/// scripted <c>ExecuteAsync</c> and near-zero back-off so the tests run fast.
17+
/// </summary>
18+
[TestFixture]
19+
public class RelayResilienceTests
20+
{
21+
private static RelayHostOptions OptionsWith(ResilienceOptions resilience)
22+
{
23+
return new RelayHostOptions
24+
{
25+
Resilience = resilience,
26+
// No Sentry DSN ⇒ telemetry is the no-op singleton (no SDK init).
27+
Telemetry = new RelayTelemetryOptions()
28+
};
29+
}
30+
31+
private static ResilienceOptions FastResilience(int maxFailures, double healthySeconds = 1000) => new ResilienceOptions
32+
{
33+
Enabled = true,
34+
MaxConsecutiveFailures = maxFailures,
35+
InitialBackoffSeconds = 0.01,
36+
MaxBackoffSeconds = 0.05,
37+
HealthyRunSeconds = healthySeconds
38+
};
39+
40+
[Test]
41+
public void ComputeBackoff_GrowsExponentially_AndIsCappedWithJitter()
42+
{
43+
var r = new ResilienceOptions
44+
{
45+
InitialBackoffSeconds = 2,
46+
MaxBackoffSeconds = 60,
47+
MaxConsecutiveFailures = 10,
48+
HealthyRunSeconds = 30
49+
};
50+
51+
// failure 1 ≈ 2s ±20% ⇒ [1.6, 2.4]; failure 3 ≈ 8s ±20% ⇒ [6.4, 9.6];
52+
// failure 10 caps at 60s ±20% ⇒ [48, 72]; floor is always >= 0.5s.
53+
RelayServiceBase.ComputeBackoff(r, 1).TotalSeconds.Should().BeInRange(1.6, 2.4);
54+
RelayServiceBase.ComputeBackoff(r, 3).TotalSeconds.Should().BeInRange(6.4, 9.6);
55+
RelayServiceBase.ComputeBackoff(r, 10).TotalSeconds.Should().BeInRange(48, 72);
56+
RelayServiceBase.ComputeBackoff(r, 1).TotalSeconds.Should().BeGreaterThanOrEqualTo(0.5);
57+
}
58+
59+
[Test]
60+
public async Task AlwaysFailing_LiveKitMode_TripsBreaker_AndFaults()
61+
{
62+
var svc = new ScriptedRelayService(OptionsWith(FastResilience(maxFailures: 3)));
63+
64+
// ExecuteAsync throws immediately every time; the breaker opens on the 3rd failure.
65+
svc.ExecuteBehavior = (_, __) => throw new InvalidOperationException("boom");
66+
67+
await svc.StartAsync(CancellationToken.None);
68+
69+
svc.State.Should().Be(RelayServiceState.Faulted);
70+
svc.Attempts.Should().Be(3, "the breaker opens once MaxConsecutiveFailures quick failures pile up");
71+
svc.Status.LiveKit.Should().Be(ConnectionState.Degraded);
72+
}
73+
74+
[Test]
75+
public async Task TransientFailures_ThenGracefulStop_EndsStopped_NotFaulted()
76+
{
77+
var svc = new ScriptedRelayService(OptionsWith(FastResilience(maxFailures: 5)));
78+
79+
// First two attempts fail (transient); the third blocks until cancelled, then
80+
// returns by honoring the token ⇒ graceful Stopped, never reaching the breaker.
81+
svc.ExecuteBehavior = async (s, token) =>
82+
{
83+
if (s.Attempts < 3)
84+
throw new InvalidOperationException("transient");
85+
await Task.Delay(Timeout.Infinite, token).ConfigureAwait(false);
86+
};
87+
88+
var run = svc.StartAsync(CancellationToken.None);
89+
90+
// Wait until the run is parked in the long-lived (3rd) attempt, then stop it.
91+
var spun = SpinWait.SpinUntil(() => svc.Attempts >= 3 && svc.State == RelayServiceState.Running, TimeSpan.FromSeconds(5));
92+
spun.Should().BeTrue("the service should reach its healthy long-lived run after the transient failures");
93+
94+
await svc.StopAsync();
95+
await run;
96+
97+
svc.State.Should().Be(RelayServiceState.Stopped);
98+
svc.Attempts.Should().Be(3);
99+
}
100+
101+
[Test]
102+
public async Task ResilienceDisabled_RunsExecuteOnce_AndPropagatesFault()
103+
{
104+
var resilience = FastResilience(maxFailures: 3);
105+
resilience.Enabled = false;
106+
var svc = new ScriptedRelayService(OptionsWith(resilience));
107+
108+
svc.ExecuteBehavior = (_, __) => throw new InvalidOperationException("boom");
109+
110+
await svc.StartAsync(CancellationToken.None);
111+
112+
svc.State.Should().Be(RelayServiceState.Faulted);
113+
svc.Attempts.Should().Be(1, "with resilience disabled ExecuteAsync runs exactly once");
114+
}
115+
116+
/// <summary>Minimal LiveKit-mode service whose run is scripted by the test.</summary>
117+
private sealed class ScriptedRelayService : RelayServiceBase
118+
{
119+
private int _attempts;
120+
121+
public ScriptedRelayService(RelayHostOptions options)
122+
: base("test", options, null)
123+
{
124+
}
125+
126+
protected override bool IsLiveKitMode => true;
127+
128+
public int Attempts => Volatile.Read(ref _attempts);
129+
130+
public Func<ScriptedRelayService, CancellationToken, Task> ExecuteBehavior { get; set; }
131+
132+
protected override async Task ExecuteAsync(CancellationToken token)
133+
{
134+
Interlocked.Increment(ref _attempts);
135+
await ExecuteBehavior(this, token).ConfigureAwait(false);
136+
}
137+
}
138+
}
139+
}

‎Resgrid.Audio.Voice.Tests/Resgrid.Audio.Voice.Tests.csproj‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -16,5 +16,6 @@
1616

1717
<ItemGroup>
1818
<ProjectReference Include="..\Resgrid.Audio.Voice\Resgrid.Audio.Voice.csproj" />
19+
<ProjectReference Include="..\Resgrid.Relay.Engine\Resgrid.Relay.Engine.csproj" />
1920
</ItemGroup>
2021
</Project>

‎Resgrid.Relay.Engine/Configuration/RelayHostOptions.cs‎

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,24 @@ public sealed class RelayHostOptions
1717
public RecorderModeOptions Recorder { get; set; } = new RecorderModeOptions();
1818
public DispatchVoiceOptions DispatchVoice { get; set; } = new DispatchVoiceOptions();
1919
public TtsSettings Tts { get; set; } = new TtsSettings();
20+
21+
// ─── Fault tolerance for the LiveKit voice modes (retry + circuit-breaker) ───
22+
public ResilienceOptions Resilience { get; set; } = new ResilienceOptions();
23+
}
24+
25+
/// <summary>
26+
/// Retry/back-off and circuit-breaker tuning for the long-lived LiveKit voice modes
27+
/// (radio / record / dispatch). A run that faults is restarted with exponential
28+
/// back-off; the breaker opens (faulting the service) once failures pile up without a
29+
/// healthy run in between.
30+
/// </summary>
31+
public sealed class ResilienceOptions
32+
{
33+
public bool Enabled { get; set; } = true;
34+
public int MaxConsecutiveFailures { get; set; } = 5; // circuit-breaker trips after this many consecutive quick failures
35+
public double InitialBackoffSeconds { get; set; } = 2;
36+
public double MaxBackoffSeconds { get; set; } = 60;
37+
public double HealthyRunSeconds { get; set; } = 30; // a run lasting >= this resets the consecutive-failure counter
2038
}
2139

2240
public sealed class RelayTelemetryOptions
Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,4 @@
1+
using System.Runtime.CompilerServices;
2+
3+
// Exposes engine internals (e.g. RelayServiceBase.ComputeBackoff) to the resilience tests.
4+
[assembly: InternalsVisibleTo("Resgrid.Audio.Voice.Tests")]

‎Resgrid.Relay.Engine/Services/DispatchRelayService.cs‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,8 @@ public DispatchRelayService(RelayHostOptions options, ILogger logger)
1717
{
1818
}
1919

20+
protected override bool IsLiveKitMode => true;
21+
2022
protected override async Task ExecuteAsync(CancellationToken token)
2123
{
2224
MutableStatus.LiveKit = ConnectionState.Connecting;

‎Resgrid.Relay.Engine/Services/RadioRelayService.cs‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,8 @@ public RadioRelayService(RelayHostOptions options, ILogger logger)
1818
{
1919
}
2020

21+
protected override bool IsLiveKitMode => true;
22+
2123
protected override async Task ExecuteAsync(CancellationToken token)
2224
{
2325
MutableStatus.LiveKit = ConnectionState.Connecting;

‎Resgrid.Relay.Engine/Services/RecordRelayService.cs‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,8 @@ public RecordRelayService(RelayHostOptions options, ILogger logger)
1717
{
1818
}
1919

20+
protected override bool IsLiveKitMode => true;
21+
2022
protected override async Task ExecuteAsync(CancellationToken token)
2123
{
2224
MutableStatus.LiveKit = ConnectionState.Connecting;

‎Resgrid.Relay.Engine/Services/RelayServiceBase.cs‎

Lines changed: 83 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@
22
using System.Threading;
33
using System.Threading.Tasks;
44
using Resgrid.Relay.Engine.Configuration;
5+
using Resgrid.Relay.Engine.Telemetry;
56
using Serilog;
67

78
namespace Resgrid.Relay.Engine.Services
@@ -15,6 +16,7 @@ namespace Resgrid.Relay.Engine.Services
1516
public abstract class RelayServiceBase : IRelayService
1617
{
1718
private readonly object _sync = new object();
19+
private readonly IRelayModeTelemetry _telemetry;
1820
private CancellationTokenSource _cts;
1921
private RelayServiceState _state = RelayServiceState.Stopped;
2022

@@ -23,8 +25,18 @@ protected RelayServiceBase(string mode, RelayHostOptions options, ILogger logger
2325
Mode = mode ?? throw new ArgumentNullException(nameof(mode));
2426
Options = options ?? throw new ArgumentNullException(nameof(options));
2527
Logger = logger;
28+
// Only the LiveKit voice modes wrap their run in the retry/circuit-breaker
29+
// loop and report to Sentry; everything else uses the cheap no-op telemetry.
30+
_telemetry = IsLiveKitMode ? RelayModeTelemetry.Create(options.Telemetry, logger) : NullRelayModeTelemetry.Instance;
2631
}
2732

33+
/// <summary>
34+
/// Whether this mode is a long-lived LiveKit voice mode (radio / record / dispatch)
35+
/// that should be wrapped in the resilience (retry + circuit-breaker) loop and
36+
/// reported to Sentry. Non-LiveKit modes (e.g. SMTP) leave this false.
37+
/// </summary>
38+
protected virtual bool IsLiveKitMode => false;
39+
2840
public string Mode { get; }
2941
public RelayServiceState State => _state;
3042
public event EventHandler<RelayStateChangedEventArgs> StateChanged;
@@ -78,21 +90,24 @@ public async Task StartAsync(CancellationToken token)
7890
// Starting during the startup window (it sets Stopping + cancels _cts under _sync).
7991
// Atomic, so a stop that wins the race keeps Stopping/Stopped instead of reverting.
8092
TryTransition(RelayServiceState.Starting, RelayServiceState.Running);
81-
await ExecuteAsync(_cts.Token).ConfigureAwait(false);
93+
await RunWithResilienceAsync(_cts.Token).ConfigureAwait(false);
94+
_telemetry.ModeStopped(Mode);
8295
TransitionTo(RelayServiceState.Stopped);
8396
}
8497
catch (OperationCanceledException) when (_cts.IsCancellationRequested)
8598
{
8699
// Graceful stop: only when OUR shutdown token was actually requested. A cancellation
87100
// from elsewhere (e.g. a dependency/HttpClient timeout) is a real fault and flows to
88101
// the catch below so Program surfaces a failure exit code.
102+
_telemetry.ModeStopped(Mode);
89103
TransitionTo(RelayServiceState.Stopped);
90104
}
91105
catch (Exception ex)
92106
{
93107
// IRelayService contract: StartAsync returns on fault — surface it via
94108
// State/StateChanged rather than throwing back to the caller.
95109
Logger?.Error(ex, "Relay mode '{Mode}' faulted", Mode);
110+
_telemetry.ModeFaulted(Mode, ex);
96111
TransitionTo(RelayServiceState.Faulted, ex.Message);
97112
}
98113
}
@@ -124,7 +139,73 @@ public virtual async ValueTask DisposeAsync()
124139
try { _cts?.Cancel(); }
125140
catch (ObjectDisposedException) { }
126141
_cts?.Dispose();
127-
await Task.CompletedTask.ConfigureAwait(false);
142+
if (_telemetry != null)
143+
await _telemetry.DisposeAsync().ConfigureAwait(false);
144+
}
145+
146+
/// <summary>
147+
/// Runs <see cref="ExecuteAsync"/>, wrapping the LiveKit voice modes in a
148+
/// retry-with-back-off loop guarded by a simple circuit breaker. A run that faults
149+
/// is restarted after an exponential, jittered delay; the breaker opens (rethrowing
150+
/// so <see cref="StartAsync"/> faults the service) once <see cref="ResilienceOptions.MaxConsecutiveFailures"/>
151+
/// failures occur without a healthy run resetting the counter. Non-LiveKit modes,
152+
/// and any mode with resilience disabled, run <see cref="ExecuteAsync"/> once directly.
153+
/// </summary>
154+
private async Task RunWithResilienceAsync(CancellationToken token)
155+
{
156+
var r = Options.Resilience ?? new ResilienceOptions();
157+
if (!IsLiveKitMode || !r.Enabled)
158+
{
159+
await ExecuteAsync(token).ConfigureAwait(false);
160+
return;
161+
}
162+
163+
_telemetry.ModeStarting(Mode);
164+
var consecutiveFailures = 0;
165+
while (true)
166+
{
167+
token.ThrowIfCancellationRequested();
168+
var startTs = System.Diagnostics.Stopwatch.GetTimestamp();
169+
try
170+
{
171+
await ExecuteAsync(token).ConfigureAwait(false);
172+
return;
173+
}
174+
catch (OperationCanceledException) when (token.IsCancellationRequested)
175+
{
176+
throw;
177+
}
178+
catch (Exception ex)
179+
{
180+
var ran = System.Diagnostics.Stopwatch.GetElapsedTime(startTs).TotalSeconds;
181+
if (ran >= r.HealthyRunSeconds)
182+
consecutiveFailures = 0; // it had been healthy ⇒ fresh transient failure
183+
consecutiveFailures++;
184+
MutableStatus.LiveKit = ConnectionState.Degraded;
185+
if (consecutiveFailures >= r.MaxConsecutiveFailures)
186+
{
187+
// Circuit open ⇒ stop retrying and let the outer catch fault the service.
188+
Logger?.Warning(ex, "Relay mode '{Mode}' circuit breaker open after {N} consecutive failures; faulting", Mode, consecutiveFailures);
189+
throw;
190+
}
191+
var delay = ComputeBackoff(r, consecutiveFailures);
192+
_telemetry.ModeRetrying(Mode, ex, consecutiveFailures, delay);
193+
Logger?.Warning(ex, "Relay mode '{Mode}' failed (failure {N}/{Max}); reconnecting in {Delay:n1}s", Mode, consecutiveFailures, r.MaxConsecutiveFailures, delay.TotalSeconds);
194+
await Task.Delay(delay, token).ConfigureAwait(false);
195+
}
196+
}
197+
}
198+
199+
/// <summary>
200+
/// Exponential back-off (doubling from <see cref="ResilienceOptions.InitialBackoffSeconds"/>,
201+
/// capped at <see cref="ResilienceOptions.MaxBackoffSeconds"/>) with ±20% jitter, floored at
202+
/// 0.5s so retries never busy-spin.
203+
/// </summary>
204+
internal static TimeSpan ComputeBackoff(ResilienceOptions r, int failures)
205+
{
206+
var baseSecs = Math.Min(r.InitialBackoffSeconds * Math.Pow(2, failures - 1), r.MaxBackoffSeconds);
207+
var jitter = baseSecs * 0.2 * (Random.Shared.NextDouble() * 2 - 1);
208+
return TimeSpan.FromSeconds(Math.Max(0.5, baseSecs + jitter));
128209
}
129210

130211
private void TransitionTo(RelayServiceState next, string error = null)
Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,19 @@
1+
using System;
2+
using System.Threading.Tasks;
3+
4+
namespace Resgrid.Relay.Engine.Telemetry
5+
{
6+
/// <summary>
7+
/// Lightweight observability hooks for the long-lived LiveKit voice modes
8+
/// (radio / record / dispatch). Implementations forward lifecycle and
9+
/// fault/retry events to a backend (e.g. Sentry) — or do nothing when
10+
/// monitoring is disabled.
11+
/// </summary>
12+
public interface IRelayModeTelemetry : IAsyncDisposable
13+
{
14+
void ModeStarting(string mode);
15+
void ModeRetrying(string mode, Exception ex, int attempt, TimeSpan nextDelay);
16+
void ModeFaulted(string mode, Exception ex);
17+
void ModeStopped(string mode);
18+
}
19+
}

0 commit comments

Comments
 (0)