Repository navigation
Conversation
…ontext items Persistent subscriptions stashed the ResolvedEvent (a struct, so boxed) and the PersistentSubscription in the consume context items for every event, only for Ack and Nack to read them back. That boxed the event and forced the lazily allocated items dictionary into existence on each message. Both are now passed to Ack and Nack as parameters. The per-event Logger.Configure call is replaced with Logger.Current = Log. The base class Log is built from the same subscription id and logger factory, so it is equivalent, and it avoids building a new LogContext when the gRPC callback does not carry the subscription's AsyncLocal. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
… async context A context without a payload never enters the consume pipe: it is marked ignored and, where positions are checkpointed, acknowledged so the checkpoint can move past it. It still paid for everything around the pipe: the logging scope and its array, the activity check, and on checkpointed subscriptions an AsyncConsumeContext wrapper that only existed to carry the ack. Handler now branches on the payload first and sends payload-less contexts to a small helper that only ignores and acknowledges, with the same failure handling as before. HandleInternal calls the same helper with the run's own ack, so the wrapper is no longer allocated. Measured on a checkpointed subscription: 609 -> 457 bytes per payload-less event. Events with a payload are unchanged. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
The reader loop of ConcurrentChannelWorker, which AsyncHandlingFilter runs handlers on, read with the worker's cancellable token. A bounded channel serves a read on an empty channel from an operation it keeps and reuses only when the token cannot be cancelled; with a cancellable one it allocates a new operation and a cancellation registration each time. An empty channel is the normal state of a caught-up subscription, so that was paid per event. The loop now reads with CancellationToken.None and keeps checking the worker token between elements. Measured in isolation on an idle bounded channel, one item at a time (net10.0): 182 -> 82 B/item with a single reader, and 182 -> 155 B/item with four readers, where only one parked reader at a time can hold the reusable operation. WaitToReadAsync + TryRead does the same for one reader but wakes every parked reader per write and came out at 350 B/item with four, so it is not used. The allocation harness feeds events as fast as it can, the channel is rarely empty there, and its four scenarios are unchanged within noise. Shutdown now relies on an invariant that already held: the worker token is never cancelled before the channel is completed. An idle reader is woken by completion alone, no longer by the token. The token is private to ChannelWorkerBase and is only cancelled by Stop's drain deadline and StopWorker's finally, both after Stop has completed the channel. Elements still queued when the token is cancelled are left unprocessed as before. Tests pin the worker's shutdown directly: an idle worker stops promptly, queued elements are drained, several writers feed several readers, and a drain that outlives its deadline leaves the rest unprocessed. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…thout LINQ FirstBeforeGap walks the set once with the struct enumerator instead of Zip/Skip/FirstOrDefault, and the commit handler trims the committed prefix with a RemoveUpTo helper instead of RemoveWhere with a closure. Behaviour is unchanged, including the gap diagnostic. Tests pin the gap cases and the prefix removal; the gap tests also pass on the old implementation. The harness shows payload-less events at about 409 bytes/event against a baseline of about 460; the other scenarios are within noise. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…per event Logger.Current skips the AsyncLocal write when the value is unchanged, but on the consume path the write sat inside async methods (DelayedConsume, HandleInternal). An async method's ExecutionContext is restored when it returns, so the calling loop never saw the value and the next event wrote it again: a new ExecutionContext and value map, 72 bytes, per write. The write now lands in the loop's own context. AsyncHandlingFilter sets it in a non-async shim in front of the unchanged async body, so the channel reader keeps it and writes again only when the next message carries another LogContext (a pipe shared by two subscriptions). The KurrentDB $all pump sets it before its loop, and the stream subscription's callback sets it in a non-async shim, which puts it in the client's read loop. SQL and Redis already set it from a synchronous method inside their polling loops. Measured with the allocation harness, bytes per event: typed async 3331-3360 -> 3245-3274 from the filter alone, and 3208-3216 once the calling loop holds the context as the $all pump now does; payload-less 415 -> 346 on the same condition. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
When a handler really suspends, every async method above it boxes its own state machine. Two layers on the consume path were async only to forward the task of the next one: EventHandler.HandleEvent awaited the typed handler and returned its status unchanged. It now returns the handler's task, or the cached Ignored one for a message type without a handler. TracingFilter.Send has nothing to do after the next filter when there is no activity: none started here because nobody listens or the span is sampled out, and no current one to reuse. It now returns the next filter's task in that case, and the part that carries an activity moved to an async helper. The activity is created before the split and started inside the helper, so it is current for exactly the same scope as before. With the dummy listener in place an activity always exists, so this only pays off once OpenTelemetry tracing is wired and drops the span. Allocation harness, bytes per event with a handler that suspends: 3280-3350 before, 3125-3150 after. In a copy of the harness with the dummy listener removed and a suspending handler, the tracing pipe goes from 2547 to 2340-2350. The rows with a synchronous handler and the traced row with a listener are unchanged, within noise. An exception thrown before the first await (a context without a message) now leaves HandleEvent directly instead of in the returned task. The callers await inside their try blocks, so the outcome is the same. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Test Results 39 files 39 suites 16m 23s ⏱️ For more details on these failures, see this check. Results for commit 6e97a3d. ♻️ This comment has been updated with latest results. |
The hand-written comparison returned the same -1, 0 or 1 that ulong.CompareTo gives for the two sequences. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
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
Per-event allocation cuts on the subscription consume path. Stacked on Eventuous#601 (
perf/hot-path-quick-wins); this PR shows only the six commits on top of it. No public signature changes. Each commit is one change and builds on its own.perf(kurrentdb)pass the resolved event to ack and nackResolvedEventinto context items, andLogger.Configureper event becameLogger.Current = Log.perf(subscriptions)payload-less eventsAsyncConsumeContextwrapper.perf(subscriptions)channel waitCancellationToken.None, so the bounded channel reuses its pooled operation.perf(subscriptions)checkpoint sequenceRemoveWherewith a closure.perf(subscriptions)logger contextAsyncLocalwrite moved out ofasyncmethods, where it was lost on return and repeated for every event.perf(subscriptions)forwarded tasksEventHandler.HandleEventand the no-activity path ofTracingFilterreturn the next task instead of awaiting it.Allocated bytes per event
Checkpointed in-memory subscription, concurrency 1, diagnostics on, no OpenTelemetry, net10.0, 200,000 events after warm-up. The harness is a console app outside the repo.
TracingFilterandTracedEventHandlerObservable behaviour changes
"resolvedEvent"and"subscription"tocontext.Items. Nothing in the repo reads them.SubscriptionId,StreamandMessageType. The message templates still carry type, stream and position.ChannelWorkerBaseand is only cancelled afterStophas completed the channel; two comments record the invariant.Logger.Current. It stays set between messages in the handling filter's reader loop, in the KurrentDB$allread loop, and in the KurrentDB client's loop that calls the stream subscription callback. It used to be null there.DeserializeMetalogs throughLogger.Current. In the two KurrentDB catch-up loops a metadata deserialization failure would have hit a null context and thrownNullReferenceException; it should now log and followThrowOnError.EventHandler.HandleEvent(a context without a message), or inTracingFilter.Sendwhen there is no activity, leaves the method directly instead of through the returned task. Every caller awaits inside atry, so the outcome is the same. Stack traces of exceptions raised after a suspension lose those frames.Not changed, but found
If acknowledging a payload-less event fails, the failure is only logged, even under
ThrowOnError: the context already holds an "ignored" result under the subscription id, so theNackunder the same name is dropped. The new tests pin this existing behaviour.Testing
net10.0, locally: Subscriptions 192 (28 new), Core 36, Diagnostics 13, SQLite 53, KurrentDB 83. All pass.
Not verified:
TracingFilterhas no unit test, because the process-wide dummy listener can't be removed reversibly in the test project.One new test waits out the hardcoded 10-second drain deadline of the channel worker, so the Subscriptions project now takes about 11 seconds.
🤖 Generated with Claude Code