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);
}