diff --git a/src/Orleans.Journaling/Formats/OrleansBinary/OrleansBinaryCommandCodecHelpers.cs b/src/Orleans.Journaling/Formats/OrleansBinary/OrleansBinaryCommandCodecHelpers.cs
index db4e508e986..5b1e4757fe5 100644
--- a/src/Orleans.Journaling/Formats/OrleansBinary/OrleansBinaryCommandCodecHelpers.cs
+++ b/src/Orleans.Journaling/Formats/OrleansBinary/OrleansBinaryCommandCodecHelpers.cs
@@ -6,7 +6,7 @@
namespace Orleans.Journaling;
///
-/// Helpers shared by the Orleans binary durable command codecs.
+/// Reads and writes Orleans binary command values using an isolated serializer session for each value.
///
internal static class OrleansBinaryCommandCodecHelpers
{
@@ -24,7 +24,15 @@ public static void WriteValue(
public static T ReadValue(IFieldCodec codec, ref Reader 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();
+ }
}
}
diff --git a/src/Orleans.Journaling/Formats/OrleansBinary/OrleansBinaryDurableDictionaryCommandCodec.cs b/src/Orleans.Journaling/Formats/OrleansBinary/OrleansBinaryDurableDictionaryCommandCodec.cs
index 52ca5126bb6..0b78d6ed78d 100644
--- a/src/Orleans.Journaling/Formats/OrleansBinary/OrleansBinaryDurableDictionaryCommandCodec.cs
+++ b/src/Orleans.Journaling/Formats/OrleansBinary/OrleansBinaryDurableDictionaryCommandCodec.cs
@@ -1,3 +1,4 @@
+using System.Buffers;
using Orleans.Serialization.Buffers;
using Orleans.Serialization.Codecs;
using Orleans.Serialization.Session;
@@ -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();
}
@@ -67,8 +66,7 @@ public void WriteSnapshot(IReadOnlyCollection> 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++;
}
@@ -97,8 +95,7 @@ private void Apply(ref Reader 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;
}
@@ -123,10 +120,34 @@ private void ApplySnapshot(ref Reader 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 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(ref Reader 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();
+ }
+ }
}
diff --git a/test/Orleans.Journaling.Tests/OrleansBinaryCommandCodecTests.cs b/test/Orleans.Journaling.Tests/OrleansBinaryCommandCodecTests.cs
index 89b93d80e6c..4623248fcfc 100644
--- a/test/Orleans.Journaling.Tests/OrleansBinaryCommandCodecTests.cs
+++ b/test/Orleans.Journaling.Tests/OrleansBinaryCommandCodecTests.cs
@@ -38,6 +38,87 @@ public void DictionaryCodec_AllCommands_RoundTrip()
Assert.Empty(consumer.SnapshotItems);
}
+ [Fact]
+ public void DictionaryCodec_SnapshotPayload_ReplaysEntryReferenceScopes()
+ {
+ var codec = new OrleansBinaryDurableDictionaryCommandCodec(
+ ValueCodec(),
+ ValueCodec(),
+ 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));
+
+ var consumer = new RecordingDictionaryCommandHandler();
+ codec.Apply(CodecTestHelpers.ReadBuffer(payload), consumer);
+ AssertAliases(consumer.SnapshotItems.Select(static item => item.Value));
+ }
+
+ [Fact]
+ public void DictionaryCodec_Set_PreservesReferencesBetweenKeyAndValue()
+ {
+ var codec = new OrleansBinaryDurableDictionaryCommandCodec(
+ ValueCodec(),
+ ValueCodec(),
+ 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();
+ 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(
+ ValueCodec(),
+ ValueCodec(),
+ 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();
+ 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()
{
@@ -67,6 +148,21 @@ public void ListCodec_AllCommands_RoundTrip()
consumer.Commands);
}
+ [Fact]
+ public void ListCodec_SnapshotPayload_ReplaysIndependentReferenceScopes()
+ {
+ var codec = new OrleansBinaryDurableListCommandCodec(
+ ValueCodec(),
+ 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()
{
@@ -92,6 +188,21 @@ public void QueueCodec_AllCommands_RoundTrip()
consumer.Commands);
}
+ [Fact]
+ public void QueueCodec_SnapshotPayload_ReplaysIndependentReferenceScopes()
+ {
+ var codec = new OrleansBinaryDurableQueueCommandCodec(
+ ValueCodec(),
+ 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()
{
@@ -117,6 +228,21 @@ public void SetCodec_AllCommands_RoundTrip()
consumer.Commands);
}
+ [Fact]
+ public void SetCodec_SnapshotPayload_ReplaysIndependentReferenceScopes()
+ {
+ var codec = new OrleansBinaryDurableSetCommandCodec(
+ ValueCodec(),
+ 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()
{
@@ -492,4 +618,67 @@ private static void AssertTrailingDataRejected(Action action)
var exception = Assert.Throws(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 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,
+ IDurableQueueCommandHandler,
+ IDurableSetCommandHandler
+{
+ public List 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();
}
diff --git a/test/Orleans.Journaling.Tests/snapshots/OrleansBinaryCodecSnapshotTests.Dictionary_Snapshot_Primitives.verified.txt b/test/Orleans.Journaling.Tests/snapshots/OrleansBinaryCodecSnapshotTests.Dictionary_Snapshot_Primitives.verified.txt
index fc4ab392858..85033a3abfa 100644
--- a/test/Orleans.Journaling.Tests/snapshots/OrleansBinaryCodecSnapshotTests.Dictionary_Snapshot_Primitives.verified.txt
+++ b/test/Orleans.Journaling.Tests/snapshots/OrleansBinaryCodecSnapshotTests.Dictionary_Snapshot_Primitives.verified.txt
@@ -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..|
diff --git a/test/Orleans.Journaling.Tests/snapshots/OrleansBinaryCodecSnapshotTests.Dictionary_Snapshot_Record.verified.txt b/test/Orleans.Journaling.Tests/snapshots/OrleansBinaryCodecSnapshotTests.Dictionary_Snapshot_Record.verified.txt
index 2bd1c7b61de..3d45b787db5 100644
--- a/test/Orleans.Journaling.Tests/snapshots/OrleansBinaryCodecSnapshotTests.Dictionary_Snapshot_Record.verified.txt
+++ b/test/Orleans.Journaling.Tests/snapshots/OrleansBinaryCodecSnapshotTests.Dictionary_Snapshot_Record.verified.txt
@@ -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..|