Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion .github/workflows/dotnet.yml
Original file line number Diff line number Diff line change
Expand Up @@ -64,7 +64,8 @@ jobs:
FullyQualifiedName~KurrentDB.Client.Tests.PositionTests.|
FullyQualifiedName~KurrentDB.Client.Tests.NodePreferenceComparerTests.|
FullyQualifiedName~KurrentDB.Client.Tests.NodeSelectorTests.|
FullyQualifiedName~KurrentDB.Client.Tests.StreamPositionTests.
FullyQualifiedName~KurrentDB.Client.Tests.StreamPositionTests.|
FullyQualifiedName~KurrentDB.Client.Tests.Compatibility.
steps:
- name: Checkout
uses: actions/checkout@v5
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,228 @@
using System.Text.Json;
using KurrentDB.Client;

namespace TrogonEventStore.Client.Compatibility;

internal sealed class CompatibilityApplication(CompatibilityOptions options) {
static readonly JsonSerializerOptions JsonOptions = new(JsonSerializerDefaults.Web);
enum AppendRpc { Standard, Batch }

public async Task ExecuteAsync(CompatibilityCommand command, CancellationToken cancellationToken) {
switch (command) {
case CompatibilityCommand.Write write:
await Write(write.Stream, AppendRpc.Standard, "write", cancellationToken);
break;
case CompatibilityCommand.BatchWrite write:
await Write(write.Stream, AppendRpc.Batch, "batch-write", cancellationToken);
break;
case CompatibilityCommand.Read read:
await Read(read.Stream, cancellationToken);
break;
case CompatibilityCommand.Subscribe subscribe:
await Subscribe(subscribe.Stream, cancellationToken);
break;
case CompatibilityCommand.CreatePersistentSubscription create:
await CreatePersistentSubscription(create.Stream, create.Group, cancellationToken);
break;
case CompatibilityCommand.ConsumePersistentSubscription consume:
await ConsumePersistentSubscription(consume.Stream, consume.Group, cancellationToken);
break;
}
}

async Task Write(
StreamName stream,
AppendRpc appendRpc,
string command,
CancellationToken cancellationToken
) {
var settings = CreateClientSettings();
await using var client = new KurrentDBClient(settings);
var userCredentials = appendRpc switch {
AppendRpc.Standard => settings.DefaultCredentials ?? throw new ArgumentException(
$"Environment variable {CompatibilityContract.ServerUriName} must include credentials for write."
),
AppendRpc.Batch => null,
_ => throw new ArgumentOutOfRangeException(nameof(appendRpc), appendRpc, null)
};
var payload = new CompatibilityEvent(CompatibilityContract.Producer, options.RunId);
var eventData = new EventData(
Uuid.NewUuid(),
CompatibilityContract.EventType,
JsonSerializer.SerializeToUtf8Bytes(payload, JsonOptions),
"{}"u8.ToArray()
);

await client.AppendToStreamAsync(
stream.Value,
StreamState.NoStream,
[eventData],
userCredentials: userCredentials,
cancellationToken: cancellationToken
);

WriteResult(command, stream, null, payload, eventData.EventId.ToString());
}

async Task Read(StreamName stream, CancellationToken cancellationToken) {
await using var client = CreateStreamsClient();
var result = client.ReadStreamAsync(
Direction.Forwards,
stream.Value,
StreamPosition.Start,
cancellationToken: cancellationToken
);

await foreach (var resolvedEvent in result.WithCancellation(cancellationToken)) {
var payload = ReadPayload(resolvedEvent);
if (payload is null)
continue;

WriteResult("read", stream, null, payload, resolvedEvent.Event.EventId.ToString());
return;
}

throw MissingEvent(stream);
}

async Task Subscribe(StreamName stream, CancellationToken cancellationToken) {
await using var client = CreateStreamsClient();
await using var subscription = client.SubscribeToStream(
stream.Value,
FromStream.Start,
cancellationToken: cancellationToken
);

await foreach (var message in subscription.Messages.WithCancellation(cancellationToken)) {
if (message is StreamMessage.SubscriptionConfirmation) {
await SignalReady(cancellationToken);
continue;
}

if (message is not StreamMessage.Event(var resolvedEvent))
continue;

var payload = ReadPayload(resolvedEvent);
if (payload is null)
continue;

WriteResult("subscribe", stream, null, payload, resolvedEvent.Event.EventId.ToString());
return;
}

throw MissingEvent(stream);
}

async Task CreatePersistentSubscription(
StreamName stream,
GroupName group,
CancellationToken cancellationToken
) {
await using var client = CreatePersistentSubscriptionsClient();
await client.CreateToStreamAsync(
stream.Value,
group.Value,
new(startFrom: StreamPosition.Start),
cancellationToken: cancellationToken
);

WriteResult(
"create-persistent-subscription",
stream,
group,
new(CompatibilityContract.Producer, options.RunId),
null
);
}

async Task ConsumePersistentSubscription(
StreamName stream,
GroupName group,
CancellationToken cancellationToken
) {
await using var client = CreatePersistentSubscriptionsClient();
await using var subscription = client.SubscribeToStream(
stream.Value,
group.Value,
cancellationToken: cancellationToken
);

await foreach (var message in subscription.Messages.WithCancellation(cancellationToken)) {
if (message is PersistentSubscriptionMessage.SubscriptionConfirmation) {
await SignalReady(cancellationToken);
continue;
}

if (message is not PersistentSubscriptionMessage.Event(var resolvedEvent, _))
continue;

await subscription.Ack(resolvedEvent);
var payload = ReadPayload(resolvedEvent);
if (payload is null)
continue;

WriteResult(
"consume-persistent-subscription",
stream,
group,
payload,
resolvedEvent.Event.EventId.ToString()
);
return;
}

throw MissingEvent(stream);
}

CompatibilityEvent? ReadPayload(ResolvedEvent resolvedEvent) {
if (resolvedEvent.Event.EventType != CompatibilityContract.EventType)
return null;

var payload = JsonSerializer.Deserialize<CompatibilityEvent>(resolvedEvent.Event.Data.Span, JsonOptions)
?? throw new InvalidDataException("Compatibility event payload is required.");

if (payload.RunId != options.RunId)
return null;

if (string.IsNullOrWhiteSpace(payload.Producer))
throw new InvalidDataException("Compatibility event producer is required.");

return payload;
}

KurrentDBClientSettings CreateClientSettings() =>
KurrentDBClientSettings.Create(options.ServerUri.OriginalString);

KurrentDBClient CreateStreamsClient() => new(CreateClientSettings());

KurrentDBPersistentSubscriptionsClient CreatePersistentSubscriptionsClient() =>
new(CreateClientSettings());

async Task SignalReady(CancellationToken cancellationToken) {
if (options.ReadyFile is { } readyFile)
await readyFile.Signal(cancellationToken);
}

void WriteResult(
string command,
StreamName stream,
GroupName? group,
CompatibilityEvent payload,
string? eventId
) => Console.WriteLine(JsonSerializer.Serialize(
new CompatibilityResult(command, stream.Value, group?.Value, payload.Producer, payload.RunId, eventId),
JsonOptions
));

InvalidDataException MissingEvent(StreamName stream) =>
new($"Stream {stream.Value} does not contain a {CompatibilityContract.EventType} event for run {options.RunId}.");
}

internal sealed record CompatibilityResult(
string Command,
string Stream,
string? Group,
string Producer,
string RunId,
string? EventId
);
Original file line number Diff line number Diff line change
@@ -0,0 +1,49 @@
namespace TrogonEventStore.Client.Compatibility;

internal abstract record CompatibilityCommand(StreamName Stream) {
public sealed record Write(StreamName Stream) : CompatibilityCommand(Stream);
public sealed record BatchWrite(StreamName Stream) : CompatibilityCommand(Stream);
public sealed record Read(StreamName Stream) : CompatibilityCommand(Stream);
public sealed record Subscribe(StreamName Stream) : CompatibilityCommand(Stream);
public sealed record CreatePersistentSubscription(StreamName Stream, GroupName Group) : CompatibilityCommand(Stream);
public sealed record ConsumePersistentSubscription(StreamName Stream, GroupName Group) : CompatibilityCommand(Stream);

public static CompatibilityCommand Parse(IReadOnlyList<string> arguments) => arguments switch {
["write", var stream] => new Write(StreamName.Parse(stream)),
["batch-write", var stream] => new BatchWrite(StreamName.Parse(stream)),
["read", var stream] => new Read(StreamName.Parse(stream)),
["subscribe", var stream] => new Subscribe(StreamName.Parse(stream)),
["create-persistent-subscription", var stream, var group] =>
new CreatePersistentSubscription(StreamName.Parse(stream), GroupName.Parse(group)),
["consume-persistent-subscription", var stream, var group] =>
new ConsumePersistentSubscription(StreamName.Parse(stream), GroupName.Parse(group)),
_ => throw new ArgumentException(
"Expected write <stream>, batch-write <stream>, read <stream>, subscribe <stream>, " +
"create-persistent-subscription <stream> <group>, or " +
"consume-persistent-subscription <stream> <group>."
)
};
}

internal readonly record struct StreamName {
StreamName(string value) => Value = value;

public string Value { get; }

public static StreamName Parse(string value) =>
new(Required(value, "Stream name"));

internal static string Required(string value, string label) =>
!string.IsNullOrWhiteSpace(value)
? value.Trim()
: throw new ArgumentException($"{label} is required.");
}

internal readonly record struct GroupName {
GroupName(string value) => Value = value;

public string Value { get; }

public static GroupName Parse(string value) =>
new(StreamName.Required(value, "Group name"));
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
namespace TrogonEventStore.Client.Compatibility;

internal static class CompatibilityContract {
public const string EventType = "trogon-compatibility";
public const string Producer = "dotnet";
public const string ServiceName = "trogon-eventstore-client-dotnet";
public const string ServerUriName = "TROGON_EVENTSTORE_URI";
public const string RunIdName = "TROGON_EVENTSTORE_RUN_ID";
public const string OtlpEndpointName = "OTEL_EXPORTER_OTLP_ENDPOINT";
public const string ReadyFileName = "TROGON_EVENTSTORE_READY_FILE";
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
using System.Text.Json.Serialization;

namespace TrogonEventStore.Client.Compatibility;

internal sealed record CompatibilityEvent(
[property: JsonPropertyName("producer")] string Producer,
[property: JsonPropertyName("runId")] string RunId
);
Original file line number Diff line number Diff line change
@@ -0,0 +1,58 @@
namespace TrogonEventStore.Client.Compatibility;

internal sealed record CompatibilityOptions(Uri ServerUri, string RunId, Uri OtlpEndpoint, ReadyFilePath? ReadyFile) {
public static CompatibilityOptions Load(Func<string, string?> readEnvironment) {
ArgumentNullException.ThrowIfNull(readEnvironment);

var serverUri = ReadAbsoluteUri(CompatibilityContract.ServerUriName);
var runId = ReadRequired(CompatibilityContract.RunIdName);
var otlpEndpoint = ReadAbsoluteUri(CompatibilityContract.OtlpEndpointName);
var readyFileValue = readEnvironment(CompatibilityContract.ReadyFileName);
ReadyFilePath? readyFile = string.IsNullOrWhiteSpace(readyFileValue)
? null
: ReadyFilePath.Parse(readyFileValue);

return new(serverUri, runId, otlpEndpoint, readyFile);

string ReadRequired(string name) {
var value = readEnvironment(name)?.Trim();
return !string.IsNullOrEmpty(value)
? value
: throw new ArgumentException($"Environment variable {name} is required.");
}

Uri ReadAbsoluteUri(string name) {
var value = ReadRequired(name);
return Uri.TryCreate(value, UriKind.Absolute, out var uri)
? uri
: throw new ArgumentException($"Environment variable {name} must be an absolute URI.");
}
}
}

internal readonly record struct ReadyFilePath {
static readonly byte[] ReadyContent = "ready\n"u8.ToArray();

ReadyFilePath(string fullPath) => FullPath = fullPath;

public string FullPath { get; }

public static ReadyFilePath Parse(string value) {
if (string.IsNullOrWhiteSpace(value) || !Path.IsPathFullyQualified(value))
throw new ArgumentException($"Environment variable {CompatibilityContract.ReadyFileName} must be an absolute path.");

return new(Path.GetFullPath(value));
}

public async Task Signal(CancellationToken cancellationToken) {
await using var file = new FileStream(
FullPath,
FileMode.CreateNew,
FileAccess.Write,
FileShare.Read,
ReadyContent.Length,
FileOptions.Asynchronous
);
await file.WriteAsync(ReadyContent, cancellationToken);
}
}
Loading