Repository navigation
OAuth authentication for the gRPC reader and sink - #257
Conversation
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
…rrors Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…an exception Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…ack hook Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…age version and masking Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Deploying replicator with
|
| Latest commit: |
912dfb5
|
| Status: | ✅ Deploy successful! |
| Preview URL: | https://840b80ba.replicator.pages.dev |
| Branch Preview URL: | https://claude-tyoung-replicator-oau.replicator.pages.dev |
PR Summary by QodoAdd OAuth authentication to the gRPC reader and sink
AI Description
Diagram
High-Level Assessment
Files changed (64)
|
Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
Code Review by Qodo
1.
|
ScavengedEventsFilter passes the auth context's Shutdown token as the caller token to auth.Run, so an in-flight metadata or stream-size gRPC lookup is cancelled at shutdown instead of only the credential wait. GrpcAuthContext.Run is unchanged so writers keep graceful shutdown. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…ts (Qodo #2) The onEvent callback is now bound to its Attempt, and HandleEvent updates the StreamMetaCache only while that attempt is current (pending or published and not dropped). The check and the update happen under the Realtime lock that also guards publish + MarkLive, so a superseded attempt's in-flight event cannot refill the cache with pre-gap metadata after the replacement went live. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Provider free text could carry sensitive data beyond the secrets we know to redact, so the Debug log of error_description is removed along with the now dead redaction code and the retained access-token list. The test now asserts error_description text never appears in logs at any level. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
ConfigRedaction.Display now also masks, under any key, a value containing a URI with user info (scheme://user[:pass]@, matched by pattern so mongodb+srv and multi-host Mongo strings are covered) or DefaultUserCredentials=. A Mongo checkpoint path with credentials is no longer printed at startup. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
GrpcAuthContext.Run and Realtime.HandleDrop logged e.Message / the raw drop exception on a token failure. A gRPC server can put arbitrary text (even an echoed Authorization header) into an Unauthenticated status detail, so this could write a token into the logs. Add AuthFailure.Describe to produce a fixed, safe description instead: OAuthTokenException's own (already-sanitized) message, a fixed string for Unauthenticated/NotAuthenticatedException, or the exception type and gRPC status code otherwise. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
HandleDrop's non-token branch still logged the raw drop exception, and SubscribeLoop's retry warning logged the exception object directly — both paths could surface a server-supplied status detail (e.g. an echoed Authorization header) in the logs. Route every drop and subscribe-failure log line through AuthFailure.Describe instead, for all exceptions, not only token failures. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
…npm 10 pnpm 10 skips dependency build scripts unless allowed, so sharp was missing and the Astro image step failed (Cloudflare Pages build). Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
A brief sink outage (e.g. a KurrentDB container restart) used up the sink pipe's ~0.5 s GreenPipes retry, faulted the writer and stalled replication until the process restarted. GrpcEventWriter now retries append, delete and set-metadata on transient failures (gRPC Unavailable, DeadlineExceeded, ResourceExhausted, Aborted, NotLeaderException, HttpRequestException, IOException, SocketException anywhere in the exception chain) with TokenGate backoff (1 s up to 30 s) on the auth context's TimeProvider. The wait ends on shutdown or the caller's token; the gRPC call itself still only gets the caller's token. Warnings are rate-limited and log only the exception type and status code; one Info line is logged on recovery. Writes use StreamState.Any and fixed event ids, so a retry after an ambiguous failure is safe. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…g (DEV-76) If the writer shovel ended while the replicator was not stopping, the reader kept re-reading the same events every cycle and nothing was ever written. Replicate now stops the reader as soon as the writer ends unexpectedly and, unless the stopping token fires (a writer cancelled by ApplicationStopping gets a short grace period), throws ReplicatorFailedException after flushing the checkpoint. With RestartOnFailure, ReplicatorService rethrows it so the host stops, and sets a non-zero exit code that Program now returns, so a container orchestrator restarts the process. Without RestartOnFailure the service logs and keeps the process and HTTP API up with replication stopped. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
compose/oauth runs idsrv4 and three KurrentDB 26 nodes (basic, OAuth x2) with a generated test CA, plus scenario configs and helper scripts to run Replicator from source against them. License key goes in a git-ignored .env. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…dedupe (Codex review, DEV-76) SetStreamMetadataAsync mints a new event id on every call, so retrying a metadata write whose first attempt landed but whose response was lost appended a second $metadata event. Metadata is now appended as a $metadata event to $$<stream> with the proposed event's stable id (StreamState.Any, per-call credentials, LogPosition returned as before), and KurrentDB deduplicates the repeat. The body is produced by the same converter the readers use to parse $metadata events. Also stops writing "$tb": 0 when the source metadata has no truncate-before (ValueOrNull returned default(StreamPosition) rather than null). Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…Codex review, DEV-76) Any HttpRequestException or IOException in the chain counted as transient, including a TLS certificate or hostname validation failure (an HttpRequestException wrapping AuthenticationException). With a wrong tlsCaFile or certificate the writer retried forever and the host-failure path never ran. AuthenticationException, NotSupportedException, UriFormatException and HTTP/2 version-negotiation failures anywhere in the chain now make the failure non-transient; connection refused/reset, timeouts and DNS failures stay transient. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…CI flake) When the writer died, Replicate threw while its metrics reporter kept running on the host's stopping token, logging "Reporting stopped" only once the host (or test) cancelled it. In tests that log landed in TUnit's per-test output after the test ended, racing TUnit's unlocked StringBuilder.ToString() and failing Writer_failure_ends_replication_with_an_error_instead_of_looping with ArgumentOutOfRangeException (chunkLength). Reproduced at ~1-6 per 1200 runs; 0 per 1200 with this fix. The reporter now runs on its own token linked to the stopping token and is cancelled and awaited before Replicate returns or throws (also if writer start or checkpoint seeding fails). The test sink additionally serialises writes (TUnit's lazy OutputWriter can hand two threads different locks over one builder) and never lets a sink failure reach the logging code. New test checks the reporter is stopped by the time a failed Replicate returns. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…hutdown (Codex review) - The metrics reporter is now cancelled and awaited in a finally around everything after it starts, so a failing final checkpoint flush (file/Mongo) no longer leaves it running after Replicate throws. - Stopping the reporter waits at most 5 s after cancellation, then logs and moves on, so a position query that ignores its token cannot hang shutdown or hide a writer failure (ReplicatorFailedException). - TcpEventReader.GetLastPosition honours its token (WaitAsync on the TCP call, which takes none). - Tests: the background-work test uses an explicit start/completion handshake (reading starts only once the position query is in flight; the query must have finished when Replicate's task completes) instead of racing Task.Run; new tests for a throwing checkpoint flush and a token-ignoring position query. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…V-1903) With scavenge on and a KurrentDB 26 source, the prepare pipe failed until it gave up on the $metadata event of $connectors-mngt/state-projection: the stream has a $maxCount but no events, so the scavenge filter's stream-size read threw StreamNotFoundException (or would index an empty result). - GetStreamSize returns StreamSize.Empty (last event number -1, as the TCP reader already reports) when the read is not found or returns no events, so nothing in such a stream is over its max count. Not-found is detected on enumeration (StreamNotFoundException); ReadState is not awaited because in EventStore.Client 23.3.8 it never completes when the server sends no messages. Every other failure, auth failures included, still propagates (fail-closed). - GetStreamMeta no longer throws for a stream without metadata (no metastream revision), which logged a warning for every such event with scavenge on. - System-stream metadata is not special-cased: EmptyDataFilter already turns every metadata/deletion event into an ignored event before the scavenge filter, so none reaches the sink (no $$$ writes); the filter only had to stop failing. - OAuth test bed scenarios now run with scavenge on; the known-issue notes are gone. Verified live (scenario 1, basic -> OAuth, scavenge on): full replay from the start with no StreamNotFound or warnings, and 50 new events reached the sink. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…iew) The token-ignoring position-query test released the fake query and returned at once, so the abandoned reporter could log "Reporting stopped" after TUnit began reading the test's output. Replicator now raises an internal ReporterExited event (test seam, after the reporter's last log line); the test waits for it, bounded, for its own reader before finishing. Assertions unchanged. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Added: OAuth 2.0 authentication for the gRPC reader and sink (client credentials or token file), configured independently per side, with automatic token refresh and recovery from token failures
Added: Helm values for extra env, env-from, volumes, volume mounts, service account and pod labels (secrets and Azure Workload Identity)
Changed: Connection string and auth setting values are masked as *** when environment variables are printed at startup
Fixed: The realtime scavenge subscription resubscribes reliably after drops, and its metadata cache is rebuilt after a gap
Fixed: Replication no longer hangs on shutdown when the writer has stopped, and the metrics reporter survives a failed position read
Fixed: Transient sink write failures (e.g. the target node restarting) are retried with backoff until they succeed instead of stopping the writer for good, and a writer that fails permanently now exits the process non-zero so it can be restarted (DEV-76)
Fixed: Scavenge no longer stalls replication on streams that have metadata but no events, such as KurrentDB 26 system streams (DEV-1903)
Why
A customer runs KurrentDB with OAuth (Microsoft Entra ID), which turns off basic username/password authentication. Replicator could only authenticate with
user:pass@in the connection string, so it couldn't replicate to or from that cluster.What this does
replicator.reader,replicator.sink) gets an optionalauthsection.typeisconnectionString(the default, behaves as before),oauthClientCredentialsoroauthTokenFile. Any combination works, e.g. a basic-auth source and an OAuth target.additionalParameterscovers providers that needaudienceorresource.user:pass@ortls=false, and is gRPC only.These existing bugs are fixed too, because token failures would have triggered them. They apply to all modes:
Realtimesubscription could stay down after a failed resubscribe, or open twice.Design spec:
docs/superpowers/specs/2026-10-03-oauth-authentication-design.md(reviewed with the Codex spec-review flow). Implementation plan:docs/superpowers/plans/2026-10-03-oauth-authentication.md.Also in this PR: DEV-76 (writer stops for good after a brief sink outage)
Found during local end-to-end testing. When the sink node was briefly unreachable, SinkPipe's ~0.5 s retry ran out, the writer task faulted and was never restarted, and replication stalled until a process restart.
GrpcEventWriternow retries transient errors (Unavailable, DeadlineExceeded, ResourceExhausted, Aborted, network errors) with backoff (1 s, 2 s, 4 s … capped at 30 s) until they succeed or shutdown. Warnings are rate-limited and the error detail is sanitized.Replicatestops and throws instead of looping. WithrestartOnFailure: true(the default) the process exits with code 1, so Kubernetes or Docker restarts it.Local end-to-end test bed
compose/oauth/runs IdentityServer 4 (the IdP used by KurrentDB's own OAuth tests) and three KurrentDB 26.1.1 nodes (basic, OAuth, OAuth), with scenario configs for running Replicator from source. See its README. The license key goes in a git-ignored.env.Results against a licensed KurrentDB 26.1.1:
Scavenge was off for the earlier runs because of a pre-existing bug that stalls replication on KurrentDB 26 system streams (DEV-1903). DEV-1903 is now fixed in this PR, and scenario 1 was re-run with scavenge on: a full replay plus new events, with no errors.
Testing
Kurrent.Replicator.Tests.Auth, covering token sources, quarantine and probe rules, the per-call header through the realEventStoreClientagainst a fake HTTP/2 handler, writer and reader recovery,Realtimeraces, pipeline shutdown, config binding and redaction. They passed 3 times in a row locally.helm template: default output unchanged byte for byte; the populated values render correctly.masteras well (emulated image boot). CI will confirm them.🤖 Generated with Claude Code