using System.Net.WebSockets;
using System.Security.Cryptography;
using System.Text.Json;
using System.Threading.Channels;
using DodoSSH.Contracts;
namespace DodoSSH.Client.Api;
///
/// A server's "pull now" notices, as everything above the transport needs them.
///
///
///
/// A queue to read from rather than an event to subscribe to, and that shape is the point: the one
/// consumer is a synchronisation loop that already waits on a timer, so it can wait on this the same
/// way and keep every continuation on the thread it started from. An event would deliver on whichever
/// thread the socket happened to complete on, which in a user interface is the difference between
/// working and an intermittent rendering fault nobody can reproduce.
///
///
/// Reading this is never how a change is applied. A notice says which vault moved and nothing
/// else; the answer to it is the ordinary delta pull. See ADR 0012.
///
///
public interface IVaultEventStream : IDisposable
{
/// Whether a socket is currently established.
///
/// For the interface to say whether it is live, not for a caller to branch on before reading:
/// synchronising is correct whether or not this is true, because the timer is the fallback.
///
bool IsConnected { get; }
///
/// Waits for the next notice.
///
///
/// Connects on the first call and reconnects for as long as it is read, so a caller neither starts
/// nor restarts anything. A server that cannot be reached is not an error here — it is a wait that
/// has not finished — because the caller's alternative is the timer it is already running.
///
ValueTask ReadAsync(CancellationToken cancellationToken);
///
/// Takes a notice if one is already waiting, without blocking.
///
/// Whether there was one.
///
/// How a caller coalesces a burst. Five people saving at once produces five notices whose answer
/// is a single synchronisation pass, so the loop reads one, waits a moment, and swallows the rest
/// rather than running the same pull five times.
///
bool TryRead(out VaultEvent notice);
}
///
/// A stream that never delivers anything.
///
///
/// For a server that does not advertise the events feature, and for tests. Deliberately waits
/// for ever rather than completing: a caller selecting between this and a timer must fall through to
/// the timer, and a read that returned immediately would spin that loop as fast as the machine allows.
///
public sealed class IdleVaultEventStream : IVaultEventStream
{
/// The one instance. It holds nothing.
public static IdleVaultEventStream Instance { get; } = new();
///
public bool IsConnected => false;
///
public async ValueTask ReadAsync(CancellationToken cancellationToken)
{
await Task.Delay(System.Threading.Timeout.Infinite, cancellationToken).ConfigureAwait(false);
// Unreachable: the delay above only ever ends by throwing.
return new VaultEvent(VaultEventKinds.Ping);
}
///
public bool TryRead(out VaultEvent notice)
{
notice = new VaultEvent(VaultEventKinds.Ping);
return false;
}
///
public void Dispose()
{
// Nothing is held.
}
}
/// Tuning for .
///
/// Every value here bounds a reconnection rather than a feature. With the socket permanently
/// unavailable the client synchronises on its timer, so the cost of getting these wrong is latency,
/// never correctness.
///
public sealed record VaultEventStreamOptions
{
/// The defaults.
public static VaultEventStreamOptions Default { get; } = new();
/// How long to wait before the first reconnection attempt.
public TimeSpan InitialBackoff { get; init; } = TimeSpan.FromSeconds(1);
/// The longest the backoff may grow to.
///
/// A minute, which is the polling interval: past that point reconnecting sooner buys nothing,
/// because the timer has already done the work the socket would have prompted.
///
public TimeSpan MaxBackoff { get; init; } = TimeSpan.FromMinutes(1);
///
/// How long a socket may be silent before it is presumed dead.
///
///
/// The server pings on an interval it states in its hello, so silence past a multiple of
/// that means the connection is gone rather than idle — which is otherwise indistinguishable, and
/// is exactly what a reverse proxy that quietly drops idle sockets produces. Used only until a
/// hello arrives; after that the server's own figure is trusted.
///
public TimeSpan InitialSilenceTimeout { get; init; } = TimeSpan.FromSeconds(90);
/// How many notices may be waiting before the oldest are dropped.
///
/// Small on purpose. A notice means "pull that vault", so a newer one subsumes the one it
/// displaces; a backlog would only make the loop pull repeatedly for work it has already done.
///
public int QueueDepth { get; init; } = 32;
}
///
/// Holds a socket to one server open, and hands over what it says.
///
///
///
/// The whole class is a reconnection policy. A dropped socket is the ordinary case — laptops sleep,
/// proxies time out, tokens expire, servers are redeployed — so nothing here treats a failure as
/// exceptional: it backs off and dials again, for as long as somebody is reading.
///
///
/// It is safe to have no server at all. Every failure path ends in "wait, then try again", and the
/// caller's synchronisation timer runs regardless, which is what makes it correct for this class to
/// stay silent about problems rather than surface them.
///
///
public sealed class VaultEventStream : IVaultEventStream, IAsyncDisposable
{
private readonly Uri endpoint;
private readonly IAccessTokenProvider tokens;
private readonly TimeProvider clock;
private readonly VaultEventStreamOptions options;
private readonly Func> connect;
private readonly Channel notices;
private readonly CancellationTokenSource closing = new();
private readonly Lock starting = new();
private Task? pump;
private bool disposed;
/// Creates a stream against one server.
/// The server's base URL, as an ordinary http or https address.
/// Supplies a bearer token, refreshing it when it is due.
/// Time source, for the backoff and the silence timeout.
/// Tuning, or null for the defaults.
public VaultEventStream(
Uri serverUrl,
IAccessTokenProvider tokens,
TimeProvider clock,
VaultEventStreamOptions? options = null)
: this(serverUrl, tokens, clock, DialAsync, options)
{
}
///
/// The connector is injected so the suite can drive this against a test host's in-memory socket.
/// Reconnection is the entire behaviour of this class, and testing it against a real network would
/// mean testing it against the one thing that cannot be made to fail on demand.
///
internal VaultEventStream(
Uri serverUrl,
IAccessTokenProvider tokens,
TimeProvider clock,
Func> connect,
VaultEventStreamOptions? options = null)
{
ArgumentNullException.ThrowIfNull(serverUrl);
ArgumentNullException.ThrowIfNull(tokens);
ArgumentNullException.ThrowIfNull(clock);
ArgumentNullException.ThrowIfNull(connect);
endpoint = EventsUrl(serverUrl);
this.tokens = tokens;
this.clock = clock;
this.connect = connect;
this.options = options ?? VaultEventStreamOptions.Default;
notices = Channel.CreateBounded(new BoundedChannelOptions(this.options.QueueDepth)
{
FullMode = BoundedChannelFullMode.DropOldest,
SingleReader = true,
SingleWriter = true,
});
}
///
public bool IsConnected { get; private set; }
///
public ValueTask ReadAsync(CancellationToken cancellationToken)
{
ObjectDisposedException.ThrowIf(disposed, this);
Start();
return notices.Reader.ReadAsync(cancellationToken);
}
///
///
/// Does not start the connection, unlike : a caller draining a burst has
/// already read one notice, and "is there another right now" is not a reason to dial a server.
///
public bool TryRead(out VaultEvent notice) => notices.Reader.TryRead(out notice!);
///
/// Ends the connection, without waiting for it to unwind.
///
///
/// What a shell calls when a connection is dropped, from a synchronous path that must not block —
/// IVaultServer is , and blocking on a socket teardown from the
/// user-interface thread is exactly the sync-over-async this repository bans. Cancelling is enough:
/// every loop reads the token, and the pump has nothing to flush.
///
public void Dispose()
{
if (disposed)
{
return;
}
disposed = true;
closing.Cancel();
// The source is deliberately left undisposed. The pump may still be inside a linked token
// source derived from this one, and disposing a parent out from under a live child is how a
// clean shutdown becomes an ObjectDisposedException on a background thread. It holds no timer
// and no handle once cancelled; DisposeAsync is the path that cleans it up properly.
}
/// Ends the connection and waits for it to unwind.
/// The deterministic form, for a caller that can await one — tests, mostly.
public async ValueTask DisposeAsync()
{
if (disposed)
{
return;
}
disposed = true;
await closing.CancelAsync().ConfigureAwait(false);
if (pump is not null)
{
try
{
await pump.ConfigureAwait(false);
}
catch (OperationCanceledException)
{
// The point of the cancel above.
}
}
closing.Dispose();
}
///
/// Turns a server's base URL into its event socket's.
///
///
/// The path is replaced rather than appended, matching every other call in this client: request
/// paths here are absolute — /api/v1/… — so a deployment behind a path prefix is already
/// unsupported, and pretending otherwise in this one place would be a difference nobody could act
/// on.
///
private static Uri EventsUrl(Uri serverUrl) =>
new UriBuilder(serverUrl)
{
Scheme = string.Equals(serverUrl.Scheme, Uri.UriSchemeHttps, StringComparison.OrdinalIgnoreCase)
? "wss"
: "ws",
Path = VaultEvents.Path,
Query = string.Empty,
Fragment = string.Empty,
}.Uri;
private static async Task DialAsync(Uri url, string token, CancellationToken cancellationToken)
{
var socket = new ClientWebSocket();
try
{
socket.Options.AddSubProtocol(VaultEvents.SubProtocol);
// A header rather than the Sec-WebSocket-Protocol smuggling ADR 0004 needs for the relay:
// this client is a native application and can set one, and the token here is the ordinary
// bearer credential rather than a ticket.
socket.Options.SetRequestHeader("Authorization", $"Bearer {token}");
await socket.ConnectAsync(url, cancellationToken).ConfigureAwait(false);
return socket;
}
catch
{
socket.Dispose();
throw;
}
}
private void Start()
{
if (pump is not null)
{
return;
}
lock (starting)
{
pump ??= Task.Run(() => RunAsync(closing.Token), closing.Token);
}
}
/// Connects, reads until it cannot, waits, and does it again.
private async Task RunAsync(CancellationToken cancellationToken)
{
var backoff = options.InitialBackoff;
while (!cancellationToken.IsCancellationRequested)
{
var outcome = await AttemptAsync(cancellationToken).ConfigureAwait(false);
// A socket that lived long enough to say hello proves the server is there and willing, so
// the next failure starts from the bottom again rather than inheriting the backoff that
// got us here. Without this a laptop that woke, connected, and then lost its network an
// hour later would wait a full minute before trying, having already proved it need not.
if (outcome == Outcome.Established)
{
backoff = options.InitialBackoff;
}
// The server said this token is spent, which the token provider can fix without waiting.
// Reconnecting at once is the whole reason that close code is distinct.
var wait = outcome == Outcome.TokenExpired ? TimeSpan.Zero : Jitter(backoff);
if (wait > TimeSpan.Zero)
{
try
{
await Task.Delay(wait, clock, cancellationToken).ConfigureAwait(false);
}
catch (OperationCanceledException)
{
return;
}
backoff = backoff < options.MaxBackoff
? Shorter(backoff * 2, options.MaxBackoff)
: options.MaxBackoff;
}
}
}
/// One connection, from dial to close.
private async Task AttemptAsync(CancellationToken cancellationToken)
{
WebSocket? socket = null;
try
{
var token = await tokens.GetAccessTokenAsync(cancellationToken).ConfigureAwait(false);
socket = await connect(endpoint, token, cancellationToken).ConfigureAwait(false);
IsConnected = true;
return await PumpAsync(socket, cancellationToken).ConfigureAwait(false);
}
catch (OperationCanceledException)
{
return Outcome.Cancelled;
}
catch (Exception exception) when (exception is not OutOfMemoryException)
{
// Every failure this can meet — no network, a refused upgrade, a server that has not been
// deployed with this feature, a token that cannot be refreshed — has the same remedy, and
// none of them is worth telling a user about. The synchronisation timer is still running.
return Outcome.Failed;
}
finally
{
IsConnected = false;
socket?.Dispose();
}
}
/// Reads frames until the socket ends or goes quiet.
private async Task PumpAsync(WebSocket socket, CancellationToken cancellationToken)
{
var buffer = new byte[8 * 1024];
var silence = options.InitialSilenceTimeout;
var established = false;
while (socket.State == WebSocketState.Open && !cancellationToken.IsCancellationRequested)
{
// Rebuilt per frame rather than reset, because a linked source cannot be un-cancelled and
// the deadline is what detects a socket that has silently gone away.
using var deadline = new CancellationTokenSource(silence, clock);
using var quiet = CancellationTokenSource.CreateLinkedTokenSource(
cancellationToken, deadline.Token);
WebSocketReceiveResult received;
try
{
received = await socket.ReceiveAsync(buffer, quiet.Token).ConfigureAwait(false);
}
catch (OperationCanceledException) when (!cancellationToken.IsCancellationRequested)
{
// Silent for longer than the server said it would be. The socket is gone in a way that
// only reconnecting can discover, which is what a proxy dropping an idle connection
// looks like from this end.
return Ended(established);
}
if (received.MessageType == WebSocketMessageType.Close)
{
return (int?)received.CloseStatus == VaultEvents.TokenExpiredCloseCode
? Outcome.TokenExpired
: Ended(established);
}
// Binary is reserved by ADR 0012 for shared-session data, and text that arrived in pieces
// is longer than anything this protocol defines. Skipped rather than fatal, so a newer
// server does not cost this client its push for the whole session.
if (received.MessageType != WebSocketMessageType.Text || !received.EndOfMessage)
{
continue;
}
if (Parse(buffer.AsSpan(0, received.Count)) is not { } frame)
{
continue;
}
established = true;
silence = await AbsorbAsync(socket, frame, silence, cancellationToken).ConfigureAwait(false);
}
return Ended(established);
}
/// Deals with one frame, and says how long the socket may now stay quiet.
private async Task AbsorbAsync(
WebSocket socket,
VaultEvent frame,
TimeSpan silence,
CancellationToken cancellationToken)
{
if (frame.HeartbeatSeconds is > 0 and var seconds)
{
// Three missed heartbeats. Two is within one paused thread of a false positive, and a
// false positive here costs a reconnection rather than anything a user sees.
silence = TimeSpan.FromSeconds(seconds * 3);
}
if (string.Equals(frame.Kind, VaultEventKinds.Ping, StringComparison.Ordinal))
{
await socket.SendAsync(
JsonSerializer.SerializeToUtf8Bytes(
new VaultEvent(VaultEventKinds.Pong), DodoSshJsonContext.Default.VaultEvent),
WebSocketMessageType.Text,
endOfMessage: true,
cancellationToken)
.ConfigureAwait(false);
return silence;
}
// Everything else, including a kind this build has never heard of, goes to the reader — which
// is what makes the frame table extensible. An unrecognised kind is one the caller ignores;
// refusing it here would be this class deciding what a newer server may say.
notices.Writer.TryWrite(frame);
return silence;
}
private static Outcome Ended(bool established) =>
established ? Outcome.Established : Outcome.Failed;
private static VaultEvent? Parse(ReadOnlySpan utf8)
{
try
{
return JsonSerializer.Deserialize(utf8, DodoSshJsonContext.Default.VaultEvent);
}
catch (JsonException)
{
return null;
}
}
///
/// Spreads reconnections out, so a server that restarts is not met by every client at once.
///
///
/// RandomNumberGenerator because System.Random is banned repo-wide. Nothing here is
/// security-relevant — the ban exists so that nothing key-, token- or nonce-adjacent can reach for
/// the weak one by habit, and paying a few microseconds to keep that rule absolute is the cheaper
/// side of the trade.
///
private static TimeSpan Jitter(TimeSpan delay)
{
var milliseconds = (int)Math.Clamp(delay.TotalMilliseconds, 1, int.MaxValue / 2);
return TimeSpan.FromMilliseconds(
milliseconds + RandomNumberGenerator.GetInt32(0, Math.Max(1, milliseconds / 2)));
}
private static TimeSpan Shorter(TimeSpan left, TimeSpan right) => left < right ? left : right;
/// How one connection attempt ended.
private enum Outcome
{
/// Never got as far as a frame. Back off.
Failed,
/// Ran, and then ended. Back off, but from the bottom.
Established,
/// The server closed it because the token expired. Reconnect at once with a new one.
TokenExpired,
/// The stream is being disposed.
Cancelled,
}
}