Files
jaap-jan 9608d73747
ci / build and test (push) Failing after 2s
Come back from a sync position the server will not accept
"The server returned 400: The sync cursor is not valid for this vault. Resync
from the beginning." told the user exactly what to do and gave them no way to
do it. The cursor is the only thing a pull sends, so the refusal was permanent:
the next pass read the same stored cursor and was told the same thing, once a
minute, for ever. And because the pull runs first, the exception ended the pass
before it reached the outbox — so the vault stopped receiving other machines'
changes and stopped sending its own. A machine that met this went quietly
read-only until somebody deleted its cache.

The engine now does what the message asks. A pull refused with the
invalid-cursor problem code — the code, never the prose, which is free to
change — drops this vault's position, writes that down, and reads the log again
from the beginning. The restarted request carries no cursor, which is the one
position a server cannot reject, so the retry cannot loop; a refusal of that is
rethrown rather than retried, and a restart is allowed once per pull. The
position is saved before the replay starts, so a process that dies halfway
through begins the next one from the beginning too rather than meeting the same
refusal again.

The mirror is deliberately kept. Replaying rewrites every row the server still
has and applying a change is a blind overwrite, so the re-pull repairs the
mirror on its way past; clearing it first would claim more than the evidence
supports — the position was refused, not the contents — and would leave a
machine that lost its connection mid-replay with less than it started with.
That leaves one gap, named in the remarks rather than left to be discovered:
once tombstone collection exists, a replay stops carrying deletions older than
the retention window.

None of the causes are the user's doing — a rotated cursor signing key, a vault
served from a restored database, a cache copied between machines — so nothing
asks them to decide anything. The report carries ResyncedFromStart and the
status line says the position was not recognised and the vault was read again.
It is kept out of NeedsAttention, because nothing is outstanding, but the
background pass breaks its usual silence for it: a sync that pulled the whole
vault on a day nobody changed anything otherwise reads as a fault.

The fake server grew a switch that refuses cursors the way a rotated signing
key does, including ones it minted itself. Three cases: the vault is re-read
and the change on the far side of the refused position arrives; the edits
waiting in the outbox are still pushed in that same pass, which is the half
that made this worth recovering from rather than merely reporting; and a server
that refuses the beginning itself is surfaced instead of replayed against.

dotnet build is clean at zero warnings, dotnet format is clean, and the sync
and app suites pass — 109 and 101.
2026-07-31 11:32:14 +02:00

533 lines
19 KiB
C#

using System.Globalization;
using System.Net;
using DodoSSH.Client.Api;
using DodoSSH.Contracts;
namespace DodoSSH.Client.Sync.Tests;
/// <summary>
/// An in-memory vault server with the real sync semantics.
/// </summary>
/// <remarks>
/// <para>
/// A faithful reimplementation of <c>DodoSSH.Api.Features.Sync.SyncService</c>'s decision table: the
/// version check, the tombstone-beats-late-upsert rule, idempotent deletes, operation receipts, the
/// change log, cursors that are opaque to the client, the pull filter, and the per-type rules about which
/// plaintext columns an item may carry. It is not a stub that returns canned answers — if it were, none of
/// the conflict tests would mean anything, because the interesting behaviour is exactly the server's
/// refusal to apply a stale write.
/// </para>
/// <para>
/// The duplication against the real service is deliberate and is the point of the exercise: two
/// independent expressions of the same rules, and <c>SyncEndpointTests</c> checks the other one against
/// real Postgres. A shared implementation would let a misreading of the protocol pass on both sides.
/// </para>
/// <para>
/// Rows are keyed on the entity type as well as the id, as the server's separate tables are and as the
/// client's cache is. Keying on the id alone would work for every test that uses one item type and would
/// silently make a host and a key with the same id the same row.
/// </para>
/// </remarks>
internal sealed class FakeVaultServer : ISyncApi
{
/// <summary>The item types this fake knows, mirroring the server's own registry.</summary>
private static readonly SyncEntityType[] Supported =
[
SyncEntityType.Host,
SyncEntityType.SshKey,
SyncEntityType.Credential,
SyncEntityType.KnownHostKey,
];
private readonly Dictionary<(SyncEntityType Type, Guid EntityId), Row> rows = [];
private readonly List<LogEntry> log = [];
private readonly Dictionary<Guid, Receipt> receipts = [];
internal FakeVaultServer(Guid vaultId, uint keyGeneration = 1)
{
VaultId = vaultId;
KeyGeneration = keyGeneration;
}
internal Guid VaultId { get; }
internal uint KeyGeneration { get; set; }
/// <summary>The server's clock, so a test can create skew deliberately.</summary>
internal DateTimeOffset Now { get; set; } = DateTimeOffset.FromUnixTimeSeconds(1_750_000_000);
/// <summary>Pull pages are capped here, as the real server clamps a client's requested limit.</summary>
internal int MaxPullLimit { get; set; } = 500;
/// <summary>Forces the next push to answer <see cref="SyncOperationStatus.Forbidden"/>.</summary>
internal bool DenyWrites { get; set; }
/// <summary>
/// Refuses every cursor a client sends, as a server does whose signing key has been rotated.
/// </summary>
/// <remarks>
/// A cursor this fake issued itself is refused just as readily, which is the whole point: the position
/// is not wrong, the deployment's ability to verify it is gone. A pull carrying no cursor is still
/// served, because "from the beginning" is what the real server tells a client to fall back to and is
/// the only position it cannot reject.
/// </remarks>
internal bool RefuseCursors { get; set; }
/// <summary>
/// Refuses a pull that carries no cursor as well, which no real server does.
/// </summary>
/// <remarks>
/// Here so a test can prove the client gives up on such a server rather than replaying the log against
/// it. Starting over is the only position a client may ask for, so being refused it has no next step.
/// </remarks>
internal bool RefuseEvenTheBeginning { get; set; }
/// <summary>Pushes received, so a test can prove a retry did or did not happen.</summary>
internal int PushCount { get; private set; }
/// <summary>The entity-type filter of the last pull, so a test can assert what was asked for.</summary>
internal IReadOnlyList<SyncEntityType>? LastPullTypes { get; private set; }
/// <summary>
/// Runs just before a push is applied, so a test can land another client's write in the window
/// between one client's pull and its push. That window is the whole subject of the cursor-gap test.
/// </summary>
internal Action? OnPush { get; set; }
internal int RowCount => rows.Count(entry => !entry.Value.IsDeleted);
/// <inheritdoc />
public Task<SyncPullResponse> SyncPullAsync(
Guid vaultId,
SyncPullRequest request,
CancellationToken cancellationToken)
{
var after = DecodeCursor(request.Cursor);
var limit = Math.Clamp(request.Limit ?? MaxPullLimit, 1, MaxPullLimit);
LastPullTypes = request.EntityTypes;
// Empty or absent means every type, as the contract says.
var wanted = request.EntityTypes is { Count: > 0 } types ? types : null;
var page = log
.Where(entry => entry.Sequence > after)
.Where(entry => wanted is null || wanted.Contains(entry.EntityType))
.Take(limit + 1)
.ToList();
var hasMore = page.Count > limit;
if (hasMore)
{
page.RemoveAt(page.Count - 1);
}
// When nothing came back the cursor must not move, or a write landing between this read and the
// next would be skipped for ever.
var next = page.Count > 0 ? page[^1].Sequence : after;
return Task.FromResult(new SyncPullResponse(
[.. page.Select(entry => Hydrate(entry))],
EncodeCursor(next),
hasMore,
Now,
KeyGeneration));
}
/// <inheritdoc />
public Task<SyncPushResponse> SyncPushAsync(
Guid vaultId,
SyncPushRequest request,
CancellationToken cancellationToken)
{
PushCount++;
var interleaved = OnPush;
OnPush = null;
interleaved?.Invoke();
var results = new List<SyncPushResult>(request.Operations.Count);
foreach (var operation in request.Operations)
{
results.Add(Apply(operation));
}
return Task.FromResult(new SyncPushResponse(results, EncodeCursor(Head)));
}
/// <summary>Applies a change as if another client had made it.</summary>
internal int ExternalUpsert(
Guid entityId,
EncryptedPayload payload,
SyncPlaintextFields? fields,
SyncEntityType entityType = SyncEntityType.Host)
{
var result = Apply(new SyncPushOperation(
Guid.CreateVersion7(),
entityType,
entityId,
SyncOperation.Upsert,
rows.TryGetValue((entityType, entityId), out var existing) && !existing.IsDeleted
? existing.Version
: null,
payload,
fields));
if (result.Status != SyncOperationStatus.Applied)
{
throw new InvalidOperationException(
$"The external write was not applied: {result.Status} — {result.Detail}.");
}
return result.Version!.Value;
}
/// <summary>Deletes as if another client had done it.</summary>
internal void ExternalDelete(Guid entityId, SyncEntityType entityType = SyncEntityType.Host)
{
var existing = rows[(entityType, entityId)];
var result = Apply(new SyncPushOperation(
Guid.CreateVersion7(),
entityType,
entityId,
SyncOperation.Delete,
existing.Version,
null,
null));
if (result.Status != SyncOperationStatus.Applied)
{
throw new InvalidOperationException($"The external delete was not applied: {result.Status}.");
}
}
internal Row? Find(Guid entityId, SyncEntityType entityType = SyncEntityType.Host) =>
rows.TryGetValue((entityType, entityId), out var row) ? row : null;
private long Head => log.Count == 0 ? 0 : log[^1].Sequence;
// ---- The decision table ----
private SyncPushResult Apply(SyncPushOperation operation)
{
if (!Supported.Contains(operation.EntityType))
{
return Invalid(operation, $"Entity type {operation.EntityType} is not yet supported.");
}
if (receipts.TryGetValue(operation.OperationId, out var receipt))
{
return new SyncPushResult(
operation.OperationId,
SyncOperationStatus.Duplicate,
receipt.Version,
receipt.Sequence,
null,
null);
}
if (DenyWrites)
{
return new SyncPushResult(
operation.OperationId, SyncOperationStatus.Forbidden, null, null, null, null);
}
rows.TryGetValue((operation.EntityType, operation.EntityId), out var existing);
return operation.Operation == SyncOperation.Delete
? ApplyDelete(operation, existing)
: ApplyUpsert(operation, existing);
}
private SyncPushResult ApplyUpsert(SyncPushOperation operation, Row? existing)
{
if (operation.Payload is null)
{
return Invalid(operation, "An upsert requires a payload.");
}
if (operation.Payload.WrappedDataKey.Length == 0 || operation.Payload.DataKeyId == Guid.Empty)
{
return Invalid(operation, "A payload requires its data key.");
}
var fields = operation.PlaintextFields ?? new SyncPlaintextFields();
if (!ValidateFields(operation.EntityType, fields, out var fieldError))
{
return Invalid(operation, fieldError);
}
if (existing is null || existing.IsDeleted)
{
return Create(operation, existing, fields);
}
if (operation.ExpectedVersion != existing.Version)
{
return Conflict(operation, existing);
}
var updated = existing with
{
Version = existing.Version + 1,
Payload = operation.Payload,
Fields = fields,
IsDeleted = false,
};
return Commit(operation, updated, SyncOperation.Upsert);
}
/// <summary>The per-type rules about which plaintext columns an item may carry.</summary>
/// <remarks>
/// One method per type, as the server has one class per type, because the differences are the interesting
/// part. Everything except a host is stricter rather than merely different: the relay concession belongs
/// to hosts alone, so anything else arriving with an address is a client bug and is refused with a reason
/// instead of being quietly dropped.
/// </remarks>
private static bool ValidateFields(
SyncEntityType entityType,
SyncPlaintextFields fields,
out string error) => entityType switch
{
SyncEntityType.SshKey => ValidateKeyFields(fields, out error),
SyncEntityType.Credential => ValidateCredentialFields(fields, out error),
SyncEntityType.KnownHostKey => ValidateKnownHostFields(fields, out error),
_ => ValidateHostFields(fields, out error),
};
private static bool ValidateKeyFields(SyncPlaintextFields fields, out string error)
{
error = string.Empty;
if (fields.RelayEnabled || fields.Hostname is not null || fields.Port is not null)
{
error = "An SSH key has no relay target; relay fields may only be set on a host.";
return false;
}
return true;
}
private static bool ValidateCredentialFields(SyncPlaintextFields fields, out string error)
{
error = string.Empty;
if (fields.RelayEnabled || fields.Hostname is not null || fields.Port is not null)
{
error = "A credential has no relay target; relay fields may only be set on a host.";
return false;
}
if (fields.PublicKeyFingerprint is not null)
{
error = "A credential has no public key.";
return false;
}
return true;
}
/// <remarks>
/// The type that does hold an address, and holds it inside the ciphertext. A pin arriving with one in the
/// clear would be the server being handed the list of endpoints a user reaches.
/// </remarks>
private static bool ValidateKnownHostFields(SyncPlaintextFields fields, out string error)
{
error = string.Empty;
if (fields.RelayEnabled || fields.Hostname is not null || fields.Port is not null)
{
error = "A known host key is not something the server dials; its address stays encrypted.";
return false;
}
if (fields.PublicKeyFingerprint is not null)
{
error = "A known host key's fingerprint stays inside its payload.";
return false;
}
return true;
}
private static bool ValidateHostFields(SyncPlaintextFields fields, out string error)
{
error = string.Empty;
if (!fields.RelayEnabled && (fields.Hostname is not null || fields.Port is not null))
{
error = "An address may only be supplied when relay is enabled.";
return false;
}
if (fields.RelayEnabled && (string.IsNullOrWhiteSpace(fields.Hostname) || fields.Port is null))
{
error = "Relay-enabled hosts require both a hostname and a port.";
return false;
}
return true;
}
private SyncPushResult Create(SyncPushOperation operation, Row? existing, SyncPlaintextFields fields)
{
// A tombstone beats a late upsert. The client is told so it can resurrect the item deliberately
// under a new id rather than silently undoing someone else's delete.
if (existing?.IsDeleted == true)
{
return Conflict(operation, existing);
}
if (operation.ExpectedVersion is not null)
{
// The client believes it is updating something that does not exist here.
return Conflict(operation, existing: null);
}
var created = new Row(
operation.EntityType, operation.EntityId, 1, 0, operation.Payload!, fields, false);
return Commit(operation, created, SyncOperation.Upsert);
}
private SyncPushResult ApplyDelete(SyncPushOperation operation, Row? existing)
{
if (existing is null)
{
return Invalid(operation, "Cannot delete an item that does not exist.");
}
if (existing.IsDeleted)
{
// Idempotent: a client retrying a delete it is unsure about should not have to tell these
// two situations apart.
return new SyncPushResult(
operation.OperationId,
SyncOperationStatus.Applied,
existing.Version,
existing.ChangeSequence,
null,
null);
}
if (operation.ExpectedVersion is not null && operation.ExpectedVersion != existing.Version)
{
return Conflict(operation, existing);
}
var tombstone = existing with
{
Version = existing.Version + 1,
IsDeleted = true,
// The address goes with the item, or the server stays able to resolve a host the user
// believes they deleted.
Fields = new SyncPlaintextFields(),
};
return Commit(operation, tombstone, SyncOperation.Delete);
}
private SyncPushResult Commit(SyncPushOperation operation, Row row, SyncOperation change)
{
var sequence = Head + 1;
log.Add(new LogEntry(sequence, row.EntityType, row.EntityId, change, row.Version, Now));
rows[(row.EntityType, row.EntityId)] = row with { ChangeSequence = sequence };
receipts[operation.OperationId] = new Receipt(row.Version, sequence);
return new SyncPushResult(
operation.OperationId, SyncOperationStatus.Applied, row.Version, sequence, null, null);
}
private SyncPushResult Conflict(SyncPushOperation operation, Row? existing) =>
new(
operation.OperationId,
SyncOperationStatus.Conflict,
existing?.Version,
existing?.ChangeSequence,
existing is null ? null : ToChange(existing),
null);
private static SyncPushResult Invalid(SyncPushOperation operation, string detail) =>
new(operation.OperationId, SyncOperationStatus.Invalid, null, null, null, detail);
private SyncChange Hydrate(LogEntry entry)
{
var row = rows[(entry.EntityType, entry.EntityId)];
return ToChange(row, entry.Sequence, entry.Revision, entry.OccurredAt);
}
private SyncChange ToChange(Row row, long? sequence = null, int? version = null, DateTimeOffset? at = null) =>
new(
row.EntityType,
row.EntityId,
row.IsDeleted ? SyncOperation.Delete : SyncOperation.Upsert,
version ?? row.Version,
sequence ?? row.ChangeSequence,
// A delete carries no payload: there is nothing left to decrypt, and shipping the pre-delete
// ciphertext would undermine the point of the tombstone.
row.IsDeleted ? null : row.Payload,
row.IsDeleted ? null : row.Fields,
at ?? Now);
// ---- Cursors ----
/// <remarks>
/// Prefixed and non-numeric so a client that tried to compute one would produce something this
/// rejects. The real server HMAC-tags them; the property that matters to the client is only that it
/// must round-trip what it is given.
/// </remarks>
private static string EncodeCursor(long sequence) =>
"fake-v1:" + sequence.ToString(CultureInfo.InvariantCulture);
private long DecodeCursor(string? cursor)
{
if (RefuseCursors && (RefuseEvenTheBeginning || !string.IsNullOrEmpty(cursor)))
{
// What the real endpoint answers: 400, the invalid-cursor code, and the sentence telling the
// client to start over. See DodoSSH.Api.Features.Sync.SyncEndpoints.
throw new DodoSshApiException(
HttpStatusCode.BadRequest,
ProblemCodes.InvalidCursor,
"The server returned 400: The sync cursor is not valid for this vault. "
+ "Resync from the beginning.");
}
if (string.IsNullOrEmpty(cursor))
{
return 0;
}
if (!cursor.StartsWith("fake-v1:", StringComparison.Ordinal)
|| !long.TryParse(cursor.AsSpan(8), CultureInfo.InvariantCulture, out var sequence))
{
throw new InvalidOperationException($"A client sent a cursor it should not have: '{cursor}'.");
}
return sequence;
}
internal sealed record Row(
SyncEntityType EntityType,
Guid EntityId,
int Version,
long ChangeSequence,
EncryptedPayload Payload,
SyncPlaintextFields Fields,
bool IsDeleted);
private sealed record LogEntry(
long Sequence,
SyncEntityType EntityType,
Guid EntityId,
SyncOperation Operation,
int Revision,
DateTimeOffset OccurredAt);
private sealed record Receipt(int Version, long Sequence);
}