Persist effect result types as ids through a dedicated type store - #254
Merged
Merged
Conversation
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).
_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.
…e 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.
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).
…ed 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.
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.
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.
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.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Summary
Types persisted inside effect results and messages are no longer stored as inline encoded type strings - they are stored as content-derived ids mapped through a dedicated type store.
TypeId: a readonly record struct wrapping along- the first 8 bytes of the SHA-256 hash of the type's encoded form (UTF-8 of its simple qualified name). Since the id is content-derived it is computed purely, with no store round-trip and no cross-replica coordination.ITypeStorepersists the id → encoded-type mapping ({prefix}_typestable in the in-memory, PostgreSQL, SqlServer and MariaDB stores). Inserts are idempotent: Postgres usesunnestarray parameters withON CONFLICT DO NOTHING, MariaDBINSERT IGNORE, SqlServer a multi-rowVALUESfiltered byNOT EXISTS.TypeMapper(one per registry) mints ids (GetTypeId, cached perType), resolves them back (ResolveType), and enforces the crash invariant via the parameterlessEnsurePersisted(): every id minted by this process that is not yet known-durable is inserted before any durable write referencing it - effect flush, message append, flow creation and control-panel writes all await it first. Persisting all minted ids is what also covers ids buried inside already-encoded payloads (staged-message children written flushlessly), without threading ids through the pipeline._serializedTypesonly ever contains persisted mappings - entries are added after the insert completes - so concurrent writers may insert the same idempotent mapping twice rather than queue behind a lock.ResolveTypereturnsTask<Type>, completing synchronously on a cache hit). Only the first resolution of an id minted by another process reaches the store, and those refreshes are admitted one at a time: a refresh fetches the whole mapping table, so waiters re-check before doing I/O of their own and the first refresh resolves the whole batch a restart misses on. This also keeps the store round-trip out ofEffectResults'_synclock - the deserializing paths snapshot the (immutable) stored effect under the lock and resolve outside it.StoredEffect.ResultTypeis aTypeId?instead of the serializer-encoded type bytes.SerializedMessage.Type,StoredMessage.MessageType(bothTypeId?-nullmarks an empty restart-poke, replacing the empty-content-and-type convention),StoredDlqMessage.MessageTypeand thePendingMessageseffect-carrier encoding all carry the id. Producers mint through the registry'sTypeMapper(MessageSender,QueueManagerstaging, control-panel appends, initial messages); consumers resolve through it (MessageDeserializer,QueueClient,ExistingMessages). Store row/blob formats keep their shape - the type piece is just 8 bytes (or null) now - so no message-table schema changes.ISerializerno longer serializes types: theSerializeType/ResolveTypedefault methods are deleted; the framework-owned encoding lives inTypeHelperextensions used by theTypeMapper.IFunctionStore.FlowTypeStore, freeing theTypeStoreproperty for the new store.Note: pre-existing databases need re-initialization - store
Initialize()is skipped when the main tables already exist, so the new type tables (_typesfor .NET types,_flowtypesfor flow types) are only created on a fresh schema.Public members that became task-returning with the async resolution:
TypeMapper.ResolveType, theStoredEffect.ResolveResultTypeextension, and the twoDefaultDeserializemethods (both already marked//todo remove).Effect'sTryGet/Getalso changed shape -TryGetreturns(bool Success, T? Value)sinceoutparams cannot surviveasync- but they areinternal.Test plan
TypeStoreSunshineScenarioTestcovering insert/get-all/idempotent re-insert in all four stores; the flow-type variant renamed toFlowTypeStoreSunshineScenarioTest.EffectResultTypeTestsresolve persisted types through a freshTypeMapperover the store, exercising the refresh-on-miss path.TypeIdTestHelpermints-and-persists ids for directly-constructed messages;DefaultDeserializeresolves through the mapper.🤖 Generated with Claude Code