using System.Globalization; using DodoSSH.Client.Api; using DodoSSH.Contracts; namespace DodoSSH.Client.Sync.Tests; /// /// An in-memory vault server with the real sync semantics. /// /// /// /// A faithful reimplementation of DodoSSH.Api.Features.Sync.SyncService'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. /// /// /// The duplication against the real service is deliberate and is the point of the exercise: two /// independent expressions of the same rules, and SyncEndpointTests checks the other one against /// real Postgres. A shared implementation would let a misreading of the protocol pass on both sides. /// /// /// 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. /// /// internal sealed class FakeVaultServer : ISyncApi { /// The item types this fake knows, mirroring the server's own registry. 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 log = []; private readonly Dictionary receipts = []; internal FakeVaultServer(Guid vaultId, uint keyGeneration = 1) { VaultId = vaultId; KeyGeneration = keyGeneration; } internal Guid VaultId { get; } internal uint KeyGeneration { get; set; } /// The server's clock, so a test can create skew deliberately. internal DateTimeOffset Now { get; set; } = DateTimeOffset.FromUnixTimeSeconds(1_750_000_000); /// Pull pages are capped here, as the real server clamps a client's requested limit. internal int MaxPullLimit { get; set; } = 500; /// Forces the next push to answer . internal bool DenyWrites { get; set; } /// Pushes received, so a test can prove a retry did or did not happen. internal int PushCount { get; private set; } /// The entity-type filter of the last pull, so a test can assert what was asked for. internal IReadOnlyList? LastPullTypes { get; private set; } /// /// 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. /// internal Action? OnPush { get; set; } internal int RowCount => rows.Count(entry => !entry.Value.IsDeleted); /// public Task 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)); } /// public Task SyncPushAsync( Guid vaultId, SyncPushRequest request, CancellationToken cancellationToken) { PushCount++; var interleaved = OnPush; OnPush = null; interleaved?.Invoke(); var results = new List(request.Operations.Count); foreach (var operation in request.Operations) { results.Add(Apply(operation)); } return Task.FromResult(new SyncPushResponse(results, EncodeCursor(Head))); } /// Applies a change as if another client had made it. 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; } /// Deletes as if another client had done it. 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); } /// The per-type rules about which plaintext columns an item may carry. /// /// 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. /// 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; } /// /// 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. /// 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 ---- /// /// 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. /// private static string EncodeCursor(long sequence) => "fake-v1:" + sequence.ToString(CultureInfo.InvariantCulture); private static long DecodeCursor(string? cursor) { 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); }