Skip to content

Commit 85bb569

Browse files
authored
Persist effect result types as ids through a dedicated type store (#254)
* Rename ITypeStore to IFlowTypeStore and its implementations * Persist effect result types as ids through a dedicated type store StoredEffect.ResultType is now a long - the first 8 bytes of the SHA-256 hash of the type's encoded form (UTF-8 of its simple qualified name) - instead of the inline encoded type. The new ITypeStore persists the id -> encoded-type mapping in a {prefix}_dotnet_types table in every store; ids are content-derived, so inserts are idempotent. The registry-wide TypeMapper computes ids without touching the store and tracks which mappings are known-persisted. Every effect write path awaits TypeMapper.EnsurePersisted before handing effects to the store (EffectResults.Flush, control-panel effect writes, staged-message children and CreateFunction's initial effects), so a mapping row is always durable before the first effect referencing it - a crash between the two writes can never leave effects whose results cannot be deserialized. Resolution of an unknown id refreshes the whole (small) mapping table once per process. The flow-type store moves to IFunctionStore.FlowTypeStore, freeing the TypeStore property for the new store. ISerializer no longer serializes types at all: the SerializeType/ResolveType default methods are deleted and the encoding lives in TypeHelper extensions, shared by the TypeMapper and the message paths (where message types remain inline bytes). * Track persistedness by insertion order instead of a separate marker set _serializedTypes now only ever contains mappings that are durable in the type store: an entry is added after its insert has completed (or from a store refresh), never before - making the dictionary itself the persisted marker and removing the separate _persisted set. Concurrent EnsurePersisted calls may insert the same mapping more than once; inserts are idempotent, so duplicates are accepted rather than queued behind a lock. Bytes for a minted-but-not-yet-persisted id are recovered by a reverse lookup over the ids GetTypeId has handed out - both when persisting the mapping and when an effect created in this process is read back before its first flush. * Persist message types as TypeIds and thread the id through the message model TypeId - a readonly struct wrapping the content-derived long - replaces raw longs and inline encoded types everywhere a type travels: an effect's ResultType, SerializedMessage.Type, StoredMessage/StoredDlqMessage's MessageType and the PendingMessages effect-carrier encoding all carry a TypeId now (8 bytes in the packed forms). A null message type marks an empty restart-poke - the poke's empty-content-and-type encoding is gone. Message producers mint ids through the registry's TypeMapper and consumers resolve through it; the ISerializer-independent inline type bytes are gone from the message pipeline. EnsurePersisted becomes parameterless - it persists every id minted by this process that is not yet known durable - so ids buried inside already-encoded payloads (staged-message children, delivered-message captures) are covered by the effect flush without threading them around, and MessageSender persists before appending rows for the same reason. Test-side, DefaultDeserialize takes the mapper it resolves through, and directly-constructed messages mint-and-persist their type id via the new TypeIdTestHelper. * Name the type-mapping table _types and the flow-type table _flowtypes The .NET-type mapping is the primary type table, so it takes the _types name; the flow-type table moves to _flowtypes (matching its IFlowTypeStore rename on main). * Track unpersisted type mappings directly instead of sweeping all minted ids GetTypeId records a newly minted mapping in an unpersisted dictionary, making EnsurePersisted's fast path an emptiness check instead of an iteration over every minted type, and making unpersisted bytes addressable by id - which replaces the reverse lookup over minted ids in both the persist and the read-back-before-first-flush resolution paths. Entries move to the persisted dictionary before they are drained, so a concurrent reader always finds a mapping in at least one of the two. * Cache resolved TypeId -> Type lookups GetTypeId seeds the cache with the Type it already has in hand, and the first resolution of a foreign id fills it lazily - so repeated ResolveType calls are a single dictionary lookup instead of a Type.GetType round-trip over the encoded name. * Resolve type ids asynchronously instead of blocking on the store TypeMapper.ResolveType returns a Task<Type>: the refresh a foreign id falls back on was blocked on with GetAwaiter().GetResult(), and three of its four call chains reached it from inside EffectResults' _sync lock - so one first-resolution parked a thread pool thread on a store round-trip and held every other effect operation on that flow behind it. Nothing required the synchronous signature; every call site bottoms out in an already-async method. EffectResults' three deserializing paths (TryGet, CreateOrGet and the generic InnerCapture) now snapshot the StoredEffect under the lock and resolve outside it - safe because the pending change and the stored effect are immutable records - so no store round-trip happens while the effect state is locked. The new await sits only on the already-completed return path, leaving the write paths' interleaving unchanged. Out params cannot survive async, so EffectResults/Effect.TryGet return (bool Success, T? Value) and Effect.Get returns a Task; both are internal. IdempotencyKeys.Initialize follows its caller into async. The public surface changes are ResolveType, the ResolveResultType extension and the two test-only DefaultDeserialize methods, which now return tasks - the test helper blocks on them rather than rippling awaits through assertion sites, some of which sit inside LINQ predicates. * Admit one type-store refresh at a time Resolving an unknown id refreshes the whole mapping table, so concurrent misses each fetched it in full. A restart reading back a batch of foreign payloads misses on many distinct ids at once, which made the cost scale with the number of unknown types rather than with the one fetch that already covers all of them. A semaphore admits one refresh, and waiters re-check the persisted mappings before reaching the store: the refresh they queued behind fetched every type, so their id is present and they return without further I/O. Only an id that genuinely is not in the store - the path to the TypeLoadException below - refreshes again, so a mapping written between the two attempts is still picked up rather than cached away as a permanent miss.
1 parent dfc5c06 commit 85bb569

73 files changed

Lines changed: 1062 additions & 406 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

‎Core/Cleipnir.ResilientFunctions.Tests/InMemoryTests/EffectFlushConcurrencyTests.cs‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,7 @@ private static Effect CreateEffect(StoredId storedId, IFunctionStore functionSto
2626
existingEffects: new List<StoredEffect>(),
2727
functionStore,
2828
DefaultSerializer.Instance,
29+
new TypeMapper(functionStore.TypeStore),
2930
owner: null,
3031
storageSession: null,
3132
clearChildren: true
@@ -133,7 +134,8 @@ public async Task SetEffectResults(StoredId storedId, IReadOnlyList<StoredEffect
133134
await inner.SetEffectResults(storedId, changes, owner, session);
134135
}
135136

136-
public IFlowTypeStore TypeStore => inner.TypeStore;
137+
public IFlowTypeStore FlowTypeStore => inner.FlowTypeStore;
138+
public ITypeStore TypeStore => inner.TypeStore;
137139
public IMessageStore MessageStore => inner.MessageStore;
138140
public IDlqStore DlqStore => inner.DlqStore;
139141
public IReplicaStore ReplicaStore => inner.ReplicaStore;

‎Core/Cleipnir.ResilientFunctions.Tests/InMemoryTests/EffectResultTypeTests.cs‎

Lines changed: 18 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,7 @@ private static Effect CreateEffect(StoredId storedId, IFunctionStore functionSto
2525
existingEffects ?? new List<StoredEffect>(),
2626
functionStore,
2727
DefaultSerializer.Instance,
28+
CreateTypeMapper(functionStore),
2829
owner: null,
2930
storageSession: null,
3031
clearChildren: true
@@ -38,6 +39,14 @@ private static Effect CreateEffect(StoredId storedId, IFunctionStore functionSto
3839
);
3940
}
4041

42+
// A fresh mapper per call: resolving a stored effect's type through it goes via the type store, verifying
43+
// the id -> type mapping was actually persisted alongside the effect.
44+
private static TypeMapper CreateTypeMapper(IFunctionStore functionStore)
45+
=> new(functionStore.TypeStore);
46+
47+
private static Task<Type> ResolveResultType(IFunctionStore store, StoredEffect storedEffect)
48+
=> storedEffect.ResolveResultType(CreateTypeMapper(store));
49+
4150
private static async Task<StoredEffect> GetStoredEffect(IFunctionStore store, StoredId storedId, EffectId effectId)
4251
=> (await store.GetFunction(storedId))!.Effects!.Single(e => e.EffectId == effectId);
4352

@@ -55,7 +64,7 @@ public async Task CapturedResultIsPersistedWithItsType()
5564
await effect.Capture(() => "SomeResult".ToTask());
5665

5766
var storedEffect = await GetSingleStoredEffect(store, storedId);
58-
var resultType = DefaultSerializer.Instance.ResolveType(storedEffect.ResultType!);
67+
var resultType = await ResolveResultType(store, storedEffect);
5968
resultType.ShouldBe(typeof(string));
6069
DefaultSerializer.Instance.Deserialize(storedEffect.Result!, resultType!).ShouldBe("SomeResult");
6170
}
@@ -72,7 +81,7 @@ public async Task UpsertedValueIsPersistedWithItsType()
7281
await effect.Upsert(effectId, value: 42, alias: null, flush: true);
7382

7483
var storedEffect = await GetStoredEffect(store, storedId, effectId);
75-
var resultType = DefaultSerializer.Instance.ResolveType(storedEffect.ResultType!);
84+
var resultType = await ResolveResultType(store, storedEffect);
7685
resultType.ShouldBe(typeof(int));
7786
DefaultSerializer.Instance.Deserialize(storedEffect.Result!, resultType!).ShouldBe(42);
7887
}
@@ -89,7 +98,7 @@ public async Task CreateOrGetValueIsPersistedWithItsType()
8998
await effect.CreateOrGet(effectId, value: new Person("Peter", 32), alias: null, flush: true);
9099

91100
var storedEffect = await GetStoredEffect(store, storedId, effectId);
92-
var resultType = DefaultSerializer.Instance.ResolveType(storedEffect.ResultType!);
101+
var resultType = await ResolveResultType(store, storedEffect);
93102
resultType.ShouldBe(typeof(Person));
94103
DefaultSerializer.Instance.Deserialize(storedEffect.Result!, resultType!).ShouldBe(new Person("Peter", 32));
95104
}
@@ -106,7 +115,7 @@ public async Task ResultCapturedThroughBaseTypeIsPersistedAndReadBackAsItsActual
106115
await effect.Capture<Animal>(() => Task.FromResult<Animal>(new Dog("Fido", Breed: "Beagle")));
107116

108117
var storedEffect = await GetSingleStoredEffect(store, storedId);
109-
DefaultSerializer.Instance.ResolveType(storedEffect.ResultType!).ShouldBe(typeof(Dog));
118+
(await ResolveResultType(store, storedEffect)).ShouldBe(typeof(Dog));
110119

111120
// Replaying the same capture against the persisted effect returns the instance that was captured -
112121
// not an Animal-shaped shell of it.
@@ -133,7 +142,7 @@ public async Task LazilyTypedSequenceIsMaterializedBeforeItIsPersisted()
133142
await effect.Capture<IEnumerable<string>>(() => Task.FromResult(names.Where(n => n.Length == 5)));
134143

135144
var storedEffect = await GetSingleStoredEffect(store, storedId);
136-
DefaultSerializer.Instance.ResolveType(storedEffect.ResultType!).ShouldBe(typeof(List<string>));
145+
(await ResolveResultType(store, storedEffect)).ShouldBe(typeof(List<string>));
137146

138147
EffectContext.Reset();
139148
var restarted = CreateEffect(storedId, store, existingEffects: [storedEffect]);
@@ -157,7 +166,7 @@ public async Task LazilyTypedSequenceCapturedAsObjectIsMaterializedBeforeItIsPer
157166
await effect.Capture<object>(() => Task.FromResult<object>(numbers.Select(n => n * 2)));
158167

159168
var storedEffect = await GetSingleStoredEffect(store, storedId);
160-
DefaultSerializer.Instance.ResolveType(storedEffect.ResultType!).ShouldBe(typeof(List<int>));
169+
(await ResolveResultType(store, storedEffect)).ShouldBe(typeof(List<int>));
161170

162171
// Without the materialized type the declared type is all there is to go on, and object yields a
163172
// JsonElement rather than the captured sequence.
@@ -181,7 +190,7 @@ public async Task PubliclyNamedCollectionIsPersistedAsIs()
181190
await effect.Capture<IEnumerable<string>>(() => Task.FromResult<IEnumerable<string>>(new[] { "Peter", "Ole" }));
182191

183192
var storedEffect = await GetSingleStoredEffect(store, storedId);
184-
DefaultSerializer.Instance.ResolveType(storedEffect.ResultType!).ShouldBe(typeof(string[]));
193+
(await ResolveResultType(store, storedEffect)).ShouldBe(typeof(string[]));
185194
}
186195

187196
[TestMethod]
@@ -197,7 +206,7 @@ public async Task DictionaryIsPersistedAsIs()
197206
await effect.Capture<IDictionary<string, int>>(() => Task.FromResult<IDictionary<string, int>>(dictionary));
198207

199208
var storedEffect = await GetSingleStoredEffect(store, storedId);
200-
DefaultSerializer.Instance.ResolveType(storedEffect.ResultType!).ShouldBe(typeof(Dictionary<string, int>));
209+
(await ResolveResultType(store, storedEffect)).ShouldBe(typeof(Dictionary<string, int>));
201210

202211
EffectContext.Reset();
203212
var restarted = CreateEffect(storedId, store, existingEffects: [storedEffect]);
@@ -221,7 +230,7 @@ public async Task NonVisibleReadOnlyDictionaryIsMaterializedIntoADictionaryBefor
221230

222231
// Materialized as a dictionary - not a list of pairs - so the payload keeps its JSON-object shape.
223232
var storedEffect = await GetSingleStoredEffect(store, storedId);
224-
DefaultSerializer.Instance.ResolveType(storedEffect.ResultType!).ShouldBe(typeof(Dictionary<string, int>));
233+
(await ResolveResultType(store, storedEffect)).ShouldBe(typeof(Dictionary<string, int>));
225234

226235
EffectContext.Reset();
227236
var restarted = CreateEffect(storedId, store, existingEffects: [storedEffect]);

‎Core/Cleipnir.ResilientFunctions.Tests/InMemoryTests/SerializationTests.cs‎

Lines changed: 2 additions & 28 deletions
Original file line numberDiff line numberDiff line change
@@ -18,8 +18,8 @@ public void ConcreteTypeOfEventIsSerializedAndDeserializedByDefaultSerializer()
1818
var serializer = DefaultSerializer.Instance;
1919
Parent @event = new Child("Hello World");
2020
var content = serializer.Serialize(@event, @event.GetType());
21-
var type = serializer.SerializeType(@event.GetType());
22-
var deserialized = serializer.Deserialize(content, serializer.ResolveType(type)!);
21+
var type = @event.GetType().SerializeType();
22+
var deserialized = serializer.Deserialize(content, type.ResolveType()!);
2323
if (deserialized is not Child child)
2424
throw new Exception("Expected event to be of child-type");
2525

@@ -65,30 +65,4 @@ public void ExceptionCanBeConvertedToAndFromFatalWorkflowException()
6565

6666
public record Parent;
6767
public record Child(string Value) : Parent;
68-
69-
[TestMethod]
70-
public void ImplementingClassCanOverrideResolveTypeDefaultMethod()
71-
{
72-
ISerializer defaultSerializer = DefaultSerializer.Instance;
73-
ISerializer customSerializer = new CustomResolveTypeSerializer();
74-
75-
// Default implementation uses Type.GetType
76-
defaultSerializer.ResolveType(typeof(string).SimpleQualifiedName().ToUtf8Bytes()).ShouldBe(typeof(string));
77-
78-
// Custom implementation always returns typeof(int) regardless of input
79-
customSerializer.ResolveType(typeof(string).SimpleQualifiedName().ToUtf8Bytes()).ShouldBe(typeof(int));
80-
customSerializer.ResolveType("anything".ToUtf8Bytes()).ShouldBe(typeof(int));
81-
}
82-
83-
private class CustomResolveTypeSerializer : ISerializer
84-
{
85-
public byte[] Serialize(object value, Type type)
86-
=> throw new NotImplementedException();
87-
88-
public object Deserialize(byte[] bytes, Type type)
89-
=> throw new NotImplementedException();
90-
91-
// Override the default interface method
92-
public Type? ResolveType(byte[] type) => typeof(int);
93-
}
9468
}

‎Core/Cleipnir.ResilientFunctions.Tests/InMemoryTests/StoreTests.cs‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -142,6 +142,10 @@ public override Task BulkScheduleWithEmptyCollectionReturnsZero()
142142
public override Task DifferentTypesAreFetchedByGetExpiredFunctionsCall()
143143
=> DifferentTypesAreFetchedByGetExpiredFunctionsCall(FunctionStoreFactory.Create());
144144

145+
[TestMethod]
146+
public override Task FlowTypeStoreSunshineScenarioTest()
147+
=> FlowTypeStoreSunshineScenarioTest(FunctionStoreFactory.Create());
148+
145149
[TestMethod]
146150
public override Task TypeStoreSunshineScenarioTest()
147151
=> TypeStoreSunshineScenarioTest(FunctionStoreFactory.Create());

‎Core/Cleipnir.ResilientFunctions.Tests/InMemoryTests/StoredEffectSerializationTests.cs‎

Lines changed: 5 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,5 @@
11
using System;
2+
using System.Threading.Tasks;
23
using Cleipnir.ResilientFunctions.CoreRuntime.Serialization;
34
using Cleipnir.ResilientFunctions.Domain;
45
using Cleipnir.ResilientFunctions.Storage;
@@ -199,11 +200,12 @@ public void StoredEffectWithEmptyIdCanBeSerializedAndDeserialized()
199200
}
200201

201202
[TestMethod]
202-
public void ResultTypeIsSerializedAndDeserializedAlongsideResult()
203+
public async Task ResultTypeIsSerializedAndDeserializedAlongsideResult()
203204
{
204205
var effectId = new EffectId([1]);
205206
var result = DefaultSerializer.Instance.Serialize("SomeResult", typeof(string));
206-
var resultType = DefaultSerializer.Instance.SerializeType(typeof(string));
207+
var typeMapper = new TypeMapper(new InMemoryTypeStore());
208+
var resultType = typeMapper.GetTypeId(typeof(string));
207209
var storedEffect = StoredEffect.CreateCompleted(effectId, result, resultType, alias: null);
208210

209211
var serialized = storedEffect.Serialize();
@@ -212,7 +214,7 @@ public void ResultTypeIsSerializedAndDeserializedAlongsideResult()
212214
deserialized.Result.ShouldBe(result);
213215
deserialized.ResultType.ShouldBe(resultType);
214216
DefaultSerializer.Instance
215-
.Deserialize(deserialized.Result!, DefaultSerializer.Instance.ResolveType(deserialized.ResultType!)!)
217+
.Deserialize(deserialized.Result!, await typeMapper.ResolveType(deserialized.ResultType!.Value))
216218
.ShouldBe("SomeResult");
217219
}
218220

‎Core/Cleipnir.ResilientFunctions.Tests/Messaging/InMemoryTests/MessageClearerTests.cs‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -137,7 +137,7 @@ await messageStore.AppendMessages([
137137
}
138138

139139
private static StoredMessage Message(StoredId storedId)
140-
=> new(storedId, MessageContent: new byte[] { 1 }, MessageType: new byte[] { 2 }, Position: 0, Replica: ReplicaId.Empty);
140+
=> new(storedId, MessageContent: new byte[] { 1 }, MessageType: new TypeId(2), Position: 0, Replica: ReplicaId.Empty);
141141

142142
private static MessageClearer CreateClearer(
143143
IMessageStore messageStore,

‎Core/Cleipnir.ResilientFunctions.Tests/Messaging/TestTemplates/DlqManagerTests.cs‎

Lines changed: 22 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -47,10 +47,10 @@ protected async Task DeadLetteredMessagesCanBeRedriven(Task<IFunctionStore> func
4747
var untouchedId = rFunc.MapToStoredId("untouched".ToFlowInstance());
4848

4949
await functionStore.DlqStore.Append([
50-
CreateMessage(byPositionId, "byPositionMsg"),
51-
CreateMessage(byStoredIdId, "byStoredIdMsg"),
52-
CreateMessage(byFlowIdId, "byFlowIdMsg"),
53-
CreateMessage(untouchedId, "untouchedMsg")
50+
CreateMessage(functionStore, byPositionId, "byPositionMsg"),
51+
CreateMessage(functionStore, byStoredIdId, "byStoredIdMsg"),
52+
CreateMessage(functionStore, byFlowIdId, "byFlowIdMsg"),
53+
CreateMessage(functionStore, untouchedId, "untouchedMsg")
5454
]);
5555

5656
var dlq = functionsRegistry.DeadLetterQueue;
@@ -67,7 +67,7 @@ await functionStore.DlqStore.Append([
6767

6868
// Only the redriven messages leave the dead letter queue.
6969
(await dlq.GetMessages([byPositionId, byStoredIdId, byFlowIdId])).ShouldBeEmpty();
70-
(await dlq.GetMessages([untouchedId])).Single().DefaultDeserialize().ShouldBe("untouchedMsg");
70+
(await dlq.GetMessages([untouchedId])).Single().DefaultDeserialize(functionStore).ShouldBe("untouchedMsg");
7171

7272
unhandledExceptionCatcher.ShouldNotHaveExceptions();
7373
}
@@ -99,8 +99,8 @@ protected async Task AllDeadLetteredMessagesForFlowAreRedriven(Task<IFunctionSto
9999
var storedId = rFunc.MapToStoredId("instanceId".ToFlowInstance());
100100

101101
await functionStore.DlqStore.Append([
102-
CreateMessage(storedId, "first"),
103-
CreateMessage(storedId, "second")
102+
CreateMessage(functionStore, storedId, "first"),
103+
CreateMessage(functionStore, storedId, "second")
104104
]);
105105

106106
var dlq = functionsRegistry.DeadLetterQueue;
@@ -130,7 +130,7 @@ protected async Task RedrivenMessageRetainsIdempotencyKeySenderAndReceiver(Task<
130130
// leaves the appended row in place for inspection.
131131
var storedId = TestStoredId.Create();
132132
await functionStore.DlqStore.Append([
133-
CreateMessage(storedId, "hello world", idempotencyKey: "idempotencyKey1", sender: "sender1", receiver: "receiver1")
133+
CreateMessage(functionStore, storedId, "hello world", idempotencyKey: "idempotencyKey1", sender: "sender1", receiver: "receiver1")
134134
]);
135135

136136
var dlq = functionsRegistry.DeadLetterQueue;
@@ -139,7 +139,7 @@ await functionStore.DlqStore.Append([
139139
var messages = await functionStore.MessageStore.GetMessages(storedId);
140140
messages.Count.ShouldBe(1);
141141
var message = messages.Single();
142-
message.DefaultDeserialize().ShouldBe("hello world");
142+
message.DefaultDeserialize(functionStore).ShouldBe("hello world");
143143
message.IdempotencyKey.ShouldBe("idempotencyKey1");
144144
message.Sender.ShouldBe("sender1");
145145
message.Receiver.ShouldBe("receiver1");
@@ -173,8 +173,8 @@ protected async Task DeletedDeadLetteredMessagesAreNotRedelivered(Task<IFunction
173173
var otherId = rFunc.MapToStoredId("otherInstance".ToFlowInstance());
174174

175175
await functionStore.DlqStore.Append([
176-
CreateMessage(storedId, "deleted"),
177-
CreateMessage(otherId, "kept")
176+
CreateMessage(functionStore, storedId, "deleted"),
177+
CreateMessage(functionStore, otherId, "kept")
178178
]);
179179

180180
var dlq = functionsRegistry.DeadLetterQueue;
@@ -183,7 +183,7 @@ await functionStore.DlqStore.Append([
183183

184184
// Unlike redrive, delete must not put the message back into the message store - the flow stays waiting.
185185
(await dlq.GetMessages([storedId])).ShouldBeEmpty();
186-
(await dlq.GetMessages([otherId])).Single().DefaultDeserialize().ShouldBe("kept");
186+
(await dlq.GetMessages([otherId])).Single().DefaultDeserialize(functionStore).ShouldBe("kept");
187187
(await functionStore.MessageStore.GetMessages(storedId)).ShouldBeEmpty();
188188
await Should.ThrowAsync<TimeoutException>(() => scheduled.Completion(timeout: TimeSpan.FromSeconds(1)));
189189

@@ -205,7 +205,7 @@ protected async Task DeadLetteredMessagesCanBePagedThrough(Task<IFunctionStore>
205205
await functionStore.DlqStore.Append(
206206
Enumerable
207207
.Range(0, 5)
208-
.Select(i => CreateMessage(storedId, $"msg{i}"))
208+
.Select(i => CreateMessage(functionStore, storedId, $"msg{i}"))
209209
.ToList()
210210
);
211211

@@ -214,18 +214,18 @@ await functionStore.DlqStore.Append(
214214

215215
var firstPage = await dlq.GetMessages(limit: 2);
216216
firstPage.Count.ShouldBe(2);
217-
firstPage[0].DefaultDeserialize().ShouldBe("msg0");
218-
firstPage[1].DefaultDeserialize().ShouldBe("msg1");
217+
firstPage[0].DefaultDeserialize(functionStore).ShouldBe("msg0");
218+
firstPage[1].DefaultDeserialize(functionStore).ShouldBe("msg1");
219219

220220
//the offset is exclusive - paging is done by passing the last returned position as the next offset
221221
var secondPage = await dlq.GetMessages(offset: firstPage[1].Position, limit: 2);
222222
secondPage.Count.ShouldBe(2);
223-
secondPage[0].DefaultDeserialize().ShouldBe("msg2");
224-
secondPage[1].DefaultDeserialize().ShouldBe("msg3");
223+
secondPage[0].DefaultDeserialize(functionStore).ShouldBe("msg2");
224+
secondPage[1].DefaultDeserialize(functionStore).ShouldBe("msg3");
225225

226226
var thirdPage = await dlq.GetMessages(offset: secondPage[1].Position, limit: 2);
227227
thirdPage.Count.ShouldBe(1);
228-
thirdPage[0].DefaultDeserialize().ShouldBe("msg4");
228+
thirdPage[0].DefaultDeserialize(functionStore).ShouldBe("msg4");
229229

230230
(await dlq.GetMessages(offset: thirdPage[0].Position)).ShouldBeEmpty();
231231

@@ -244,7 +244,7 @@ protected async Task RedrivingAndDeletingUnknownOrEmptyInputIsANoOp(Task<IFuncti
244244
);
245245

246246
var storedId = TestStoredId.Create();
247-
await functionStore.DlqStore.Append([CreateMessage(storedId, "untouched")]);
247+
await functionStore.DlqStore.Append([CreateMessage(functionStore, storedId, "untouched")]);
248248

249249
var dlq = functionsRegistry.DeadLetterQueue;
250250

@@ -258,16 +258,16 @@ protected async Task RedrivingAndDeletingUnknownOrEmptyInputIsANoOp(Task<IFuncti
258258
await dlq.Delete(new List<long>());
259259
await dlq.Delete([999_999L]);
260260

261-
(await dlq.GetMessages()).Single().DefaultDeserialize().ShouldBe("untouched");
261+
(await dlq.GetMessages()).Single().DefaultDeserialize(functionStore).ShouldBe("untouched");
262262

263263
unhandledExceptionCatcher.ShouldNotHaveExceptions();
264264
}
265265

266-
private static StoredMessage CreateMessage(StoredId storedId, string content, string? idempotencyKey = null, string? sender = null, string? receiver = null)
266+
private static StoredMessage CreateMessage(IFunctionStore functionStore, StoredId storedId, string content, string? idempotencyKey = null, string? sender = null, string? receiver = null)
267267
=> new(
268268
storedId,
269269
content.ToJson().ToUtf8Bytes(),
270-
typeof(string).SimpleQualifiedName().ToUtf8Bytes(),
270+
functionStore.GetTypeId(typeof(string)),
271271
Position: 0,
272272
Replica: ReplicaId.Empty,
273273
idempotencyKey,

0 commit comments

Comments
 (0)