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
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@
namespace Orleans.Journaling;

/// <summary>
/// Helpers shared by the Orleans binary durable command codecs.
/// Reads and writes Orleans binary command values using an isolated serializer session for each value.
/// </summary>
internal static class OrleansBinaryCommandCodecHelpers
{
Expand All @@ -24,7 +24,15 @@ public static void WriteValue<T>(

public static T ReadValue<T, TInput>(IFieldCodec<T> codec, ref Reader<TInput> reader)
{
var field = reader.ReadFieldHeader();
return codec.ReadValue(ref reader, field)!;
reader.Session.Reset();
try
{
var field = reader.ReadFieldHeader();
return codec.ReadValue(ref reader, field)!;
}
finally
{
reader.Session.Reset();
}
}
}
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
using System.Buffers;
using Orleans.Serialization.Buffers;
using Orleans.Serialization.Codecs;
using Orleans.Serialization.Session;
Expand All @@ -22,12 +23,10 @@ public void WriteSet(TKey key, TValue value, JournalStreamWriter writer)
{
using var entry = writer.BeginEntry();
var output = entry.Writer;
using var session = sessionPool.GetSession();
var payloadWriter = Writer.Create(output, session);
var payloadWriter = Writer.Create(output, session: null!);
payloadWriter.WriteVarUInt32(SetCommand);
keyCodec.WriteField(ref payloadWriter, 0, typeof(TKey), key);
valueCodec.WriteField(ref payloadWriter, 1, typeof(TValue), value);
payloadWriter.Commit();
WriteKeyValue(key, value, output);
entry.Commit();
}

Expand Down Expand Up @@ -67,8 +66,7 @@ public void WriteSnapshot(IReadOnlyCollection<KeyValuePair<TKey, TValue>> items,
foreach (var (key, value) in items)
{
CollectionCodecHelpers.ThrowIfSnapshotItemCountExceeded(count, written);
OrleansBinaryCommandCodecHelpers.WriteValue(keyCodec, key, output, sessionPool);
OrleansBinaryCommandCodecHelpers.WriteValue(valueCodec, value, output, sessionPool);
WriteKeyValue(key, value, output);
written++;
}

Expand Down Expand Up @@ -97,8 +95,7 @@ private void Apply<TInput>(ref Reader<TInput> reader, IDurableDictionaryCommandH
{
case SetCommand:
{
var key = OrleansBinaryCommandCodecHelpers.ReadValue(keyCodec, ref reader);
var value = OrleansBinaryCommandCodecHelpers.ReadValue(valueCodec, ref reader);
var (key, value) = ReadKeyValue(ref reader);
consumer.ApplySet(key, value);
break;
}
Expand All @@ -123,10 +120,34 @@ private void ApplySnapshot<TInput>(ref Reader<TInput> reader, IDurableDictionary
consumer.Reset(count);
for (var i = 0; i < count; i++)
{
var key = OrleansBinaryCommandCodecHelpers.ReadValue(keyCodec, ref reader);
var value = OrleansBinaryCommandCodecHelpers.ReadValue(valueCodec, ref reader);
var (key, value) = ReadKeyValue(ref reader);
consumer.ApplySet(key, value);
}
}

private void WriteKeyValue(TKey key, TValue value, IBufferWriter<byte> output)
{
using var session = sessionPool.GetSession();
var writer = Writer.Create(output, session);
keyCodec.WriteField(ref writer, 0, typeof(TKey), key);
valueCodec.WriteField(ref writer, 1, typeof(TValue), value);
writer.Commit();
}

private (TKey Key, TValue Value) ReadKeyValue<TInput>(ref Reader<TInput> reader)
{
reader.Session.Reset();
try
{
var keyField = reader.ReadFieldHeader();
var key = keyCodec.ReadValue(ref reader, keyField)!;
var valueField = reader.ReadFieldHeader();
var value = valueCodec.ReadValue(ref reader, valueField)!;
return (key, value);
}
finally
{
reader.Session.Reset();
}
}
}
189 changes: 189 additions & 0 deletions test/Orleans.Journaling.Tests/OrleansBinaryCommandCodecTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,87 @@ public void DictionaryCodec_AllCommands_RoundTrip()
Assert.Empty(consumer.SnapshotItems);
}

[Fact]
public void DictionaryCodec_SnapshotPayload_ReplaysEntryReferenceScopes()
{
var codec = new OrleansBinaryDurableDictionaryCommandCodec<string, SnapshotReferenceRecord>(
ValueCodec<string>(),
ValueCodec<SnapshotReferenceRecord>(),
SessionPool);
var items = CreateReferenceItems();
var payload = CodecTestHelpers.WriteEntry(
writer => codec.WriteSnapshot([new("first", items[0]), new("second", items[1])], writer));

Assert.Equal(
"0705400B6669727374214007010203C107E0400D7365636F6E64214007040506C107E0",
Convert.ToHexString(payload));

Comment thread
ReubenBond marked this conversation as resolved.
var consumer = new RecordingDictionaryCommandHandler<string, SnapshotReferenceRecord>();
codec.Apply(CodecTestHelpers.ReadBuffer(payload), consumer);
AssertAliases(consumer.SnapshotItems.Select(static item => item.Value));
}

[Fact]
public void DictionaryCodec_Set_PreservesReferencesBetweenKeyAndValue()
{
var codec = new OrleansBinaryDurableDictionaryCommandCodec<SnapshotReferenceRecord, SnapshotReferenceRecord>(
ValueCodec<SnapshotReferenceRecord>(),
ValueCodec<SnapshotReferenceRecord>(),
SessionPool);
var shared = new byte[] { 1, 2, 3 };
var key = new SnapshotReferenceRecord { Payload = shared, Alias = shared };
var value = new SnapshotReferenceRecord { Payload = shared, Alias = shared };
var payload = CodecTestHelpers.WriteEntry(writer => codec.WriteSet(key, value, writer));

var consumer = new RecordingDictionaryCommandHandler<SnapshotReferenceRecord, SnapshotReferenceRecord>();
codec.Apply(CodecTestHelpers.ReadBuffer(payload), consumer);

var item = Assert.Single(consumer.SnapshotItems);
Assert.Same(item.Key.Payload, item.Key.Alias);
Assert.Same(item.Value.Payload, item.Value.Alias);
Assert.Same(item.Key.Payload, item.Value.Payload);
}

[Fact]
public void DictionaryCodec_Snapshot_PreservesEntryReferencesAndIsolatesEntries()
{
var codec = new OrleansBinaryDurableDictionaryCommandCodec<SnapshotReferenceRecord, SnapshotReferenceRecord>(
ValueCodec<SnapshotReferenceRecord>(),
ValueCodec<SnapshotReferenceRecord>(),
SessionPool);
var shared = new byte[] { 1, 2, 3 };
var payload = CodecTestHelpers.WriteEntry(
writer => codec.WriteSnapshot(
[
new(
new() { Payload = shared, Alias = shared },
new() { Payload = shared, Alias = shared }),
new(
new() { Payload = shared, Alias = shared },
new() { Payload = shared, Alias = shared })
],
writer));

var consumer = new RecordingDictionaryCommandHandler<SnapshotReferenceRecord, SnapshotReferenceRecord>();
codec.Apply(CodecTestHelpers.ReadBuffer(payload), consumer);

Assert.Collection(
consumer.SnapshotItems,
first =>
{
Assert.Same(first.Key.Payload, first.Key.Alias);
Assert.Same(first.Value.Payload, first.Value.Alias);
Assert.Same(first.Key.Payload, first.Value.Payload);
},
second =>
{
Assert.Same(second.Key.Payload, second.Key.Alias);
Assert.Same(second.Value.Payload, second.Value.Alias);
Assert.Same(second.Key.Payload, second.Value.Payload);
Assert.NotSame(consumer.SnapshotItems[0].Key.Payload, second.Key.Payload);
});
}

[Fact]
public void ListCodec_AllCommands_RoundTrip()
{
Expand Down Expand Up @@ -67,6 +148,21 @@ public void ListCodec_AllCommands_RoundTrip()
consumer.Commands);
}

[Fact]
public void ListCodec_SnapshotPayload_ReplaysIndependentReferenceScopes()
{
var codec = new OrleansBinaryDurableListCommandCodec<SnapshotReferenceRecord>(
ValueCodec<SnapshotReferenceRecord>(),
SessionPool);
var payload = CodecTestHelpers.WriteEntry(writer => codec.WriteSnapshot(CreateReferenceItems(), writer));

Assert.Equal("0B05204007010203C105E0204007040506C105E0", Convert.ToHexString(payload));

var consumer = new SnapshotCollectionConsumer();
codec.Apply(CodecTestHelpers.ReadBuffer(payload), consumer);
AssertAliases(consumer.Items);
}

[Fact]
public void QueueCodec_AllCommands_RoundTrip()
{
Expand All @@ -92,6 +188,21 @@ public void QueueCodec_AllCommands_RoundTrip()
consumer.Commands);
}

[Fact]
public void QueueCodec_SnapshotPayload_ReplaysIndependentReferenceScopes()
{
var codec = new OrleansBinaryDurableQueueCommandCodec<SnapshotReferenceRecord>(
ValueCodec<SnapshotReferenceRecord>(),
SessionPool);
var payload = CodecTestHelpers.WriteEntry(writer => codec.WriteSnapshot(CreateReferenceItems(), writer));

Assert.Equal("0705204007010203C105E0204007040506C105E0", Convert.ToHexString(payload));

var consumer = new SnapshotCollectionConsumer();
codec.Apply(CodecTestHelpers.ReadBuffer(payload), consumer);
AssertAliases(consumer.Items);
}

[Fact]
public void SetCodec_AllCommands_RoundTrip()
{
Expand All @@ -117,6 +228,21 @@ public void SetCodec_AllCommands_RoundTrip()
consumer.Commands);
}

[Fact]
public void SetCodec_SnapshotPayload_ReplaysIndependentReferenceScopes()
{
var codec = new OrleansBinaryDurableSetCommandCodec<SnapshotReferenceRecord>(
ValueCodec<SnapshotReferenceRecord>(),
SessionPool);
var payload = CodecTestHelpers.WriteEntry(writer => codec.WriteSnapshot(CreateReferenceItems(), writer));

Assert.Equal("0705204007010203C105E0204007040506C105E0", Convert.ToHexString(payload));

var consumer = new SnapshotCollectionConsumer();
codec.Apply(CodecTestHelpers.ReadBuffer(payload), consumer);
AssertAliases(consumer.Items);
}

[Fact]
public void ValueStateAndTcsCodecs_AllCommands_RoundTrip()
{
Expand Down Expand Up @@ -492,4 +618,67 @@ private static void AssertTrailingDataRejected(Action action)
var exception = Assert.Throws<InvalidOperationException>(action);
Assert.Contains("trailing data", exception.Message);
}

private static SnapshotReferenceRecord[] CreateReferenceItems()
{
var first = new byte[] { 1, 2, 3 };
var second = new byte[] { 4, 5, 6 };
return
[
new() { Payload = first, Alias = first },
new() { Payload = second, Alias = second }
];
}

private static void AssertAliases(IEnumerable<SnapshotReferenceRecord> items)
{
Assert.Collection(
items,
item =>
{
Assert.Equal([1, 2, 3], item.Payload);
Assert.Same(item.Payload, item.Alias);
},
item =>
{
Assert.Equal([4, 5, 6], item.Payload);
Assert.Same(item.Payload, item.Alias);
});
}
}

[GenerateSerializer]
internal sealed class SnapshotReferenceRecord
{
[Id(0)]
public required byte[] Payload { get; init; }

[Id(1)]
public required byte[] Alias { get; init; }
}

internal sealed class SnapshotCollectionConsumer :
IDurableListCommandHandler<SnapshotReferenceRecord>,
IDurableQueueCommandHandler<SnapshotReferenceRecord>,
IDurableSetCommandHandler<SnapshotReferenceRecord>
{
public List<SnapshotReferenceRecord> Items { get; } = [];

public void ApplyAdd(SnapshotReferenceRecord item) => Items.Add(item);

public void ApplySet(int index, SnapshotReferenceRecord item) => throw new NotSupportedException();

public void ApplyInsert(int index, SnapshotReferenceRecord item) => throw new NotSupportedException();

public void ApplyRemoveAt(int index) => throw new NotSupportedException();

public void ApplyClear() => Items.Clear();

public void Reset(int capacityHint) => Items.Clear();

public void ApplyEnqueue(SnapshotReferenceRecord item) => Items.Add(item);

public void ApplyDequeue() => throw new NotSupportedException();

public void ApplyRemove(SnapshotReferenceRecord item) => throw new NotSupportedException();
}
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
HEX: 0117000000080000000705400B616C70686100054009626574610009
HEX: 0117000000080000000705400B616C70686101054009626574610109
DISASSEMBLY:
[entry 0] frame-version=1 length=23 streamId=8 payload-bytes=19
0000 07 05 40 0B 61 6C 70 68 61 00 05 40 09 62 65 74 |..@.alpha..@.bet|
0010 61 00 09 |a..|
0000 07 05 40 0B 61 6C 70 68 61 01 05 40 09 62 65 74 |..@.alpha..@.bet|
0010 61 01 09 |a..|
Original file line number Diff line number Diff line change
@@ -1,8 +1,8 @@
HEX: 0153000000080000000705400B616C70686120E8401F736E617073686F742D7265636F7264011D21E84007746167011A03E0E040096265746120E84015616C742D7265636F7264013521E8400F616C742D7461670145E0E0
HEX: 0153000000080000000705400B616C70686121E8401F736E617073686F742D7265636F7264011D21E84007746167011A03E0E040096265746121E84015616C742D7265636F7264013521E8400F616C742D7461670145E0E0
DISASSEMBLY:
[entry 0] frame-version=1 length=83 streamId=8 payload-bytes=79
0000 07 05 40 0B 61 6C 70 68 61 20 E8 40 1F 73 6E 61 |..@.alpha .@.sna|
0000 07 05 40 0B 61 6C 70 68 61 21 E8 40 1F 73 6E 61 |..@.alpha!.@.sna|
0010 70 73 68 6F 74 2D 72 65 63 6F 72 64 01 1D 21 E8 |pshot-record..!.|
0020 40 07 74 61 67 01 1A 03 E0 E0 40 09 62 65 74 61 |@.tag.....@.beta|
0030 20 E8 40 15 61 6C 74 2D 72 65 63 6F 72 64 01 35 | .@.alt-record.5|
0030 21 E8 40 15 61 6C 74 2D 72 65 63 6F 72 64 01 35 |!.@.alt-record.5|
0040 21 E8 40 0F 61 6C 74 2D 74 61 67 01 45 E0 E0 |!.@.alt-tag.E..|