refactor(cardinal): move client and inter-shard communication into pkg/transport - #1059
winton-library wants to merge 4 commits into
Conversation
❌ 2 Tests Failed:
View the top 2 failed test(s) by shortest run time
To view more test analytics, go to the Test Analytics Dashboard |
There was a problem hiding this comment.
2 issues found across 19 files
Prompt for AI agents (unresolved issues)
Check if these issues are valid — if so, understand the root cause of each and fix them. When an issue isn't valid or won't be fixed in this PR, reply in its thread with the reason and then resolve the thread. If appropriate, use sub-agents to investigate and fix each issue separately.
<file name="pkg/transport/server.go">
<violation number="1" location="pkg/transport/server.go:377">
P1: `SendTo` targets one user, but this lookup selects every waiter for the event name. Another user's `SendCommandWithReply` can therefore receive the targeted payload; scope waiters by user as well.</violation>
</file>
<file name="pkg/transport/intershard_internal_test.go">
<violation number="1" location="pkg/transport/intershard_internal_test.go:79">
P2: This 500 ms wall-clock assertion can fail under CI scheduling contention, especially in a parallel test. Use a synchronization-based check with a more generous deadline to verify `Flush` returns while the handler remains blocked.</violation>
</file>
Shadow auto-approve: would not auto-approve because issues were found.
Re-trigger cubic
| } | ||
| } | ||
| } | ||
| waiters := append([]chan *iscv1.Event(nil), s.replyWaiters[eventPb.GetName()]...) |
There was a problem hiding this comment.
P1: SendTo targets one user, but this lookup selects every waiter for the event name. Another user's SendCommandWithReply can therefore receive the targeted payload; scope waiters by user as well.
Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. When an issue isn't valid or won't be fixed in this PR, reply in its thread with the reason and then resolve the thread. At pkg/transport/server.go, line 377:
<comment>`SendTo` targets one user, but this lookup selects every waiter for the event name. Another user's `SendCommandWithReply` can therefore receive the targeted payload; scope waiters by user as well.</comment>
<file context>
@@ -0,0 +1,421 @@
+ }
+ }
+ }
+ waiters := append([]chan *iscv1.Event(nil), s.replyWaiters[eventPb.GetName()]...)
+ s.mu.RUnlock()
+ span.SetAttributes(attrEventSubscribers.Int(len(subscribers)), attrEventWaiters.Int(len(waiters)))
</file context>
| fixtureA.tr.Enqueue(context.Background(), fixtureB.address, testutils.SimpleCommand{Value: flush}) | ||
| start := time.Now() | ||
| fixtureA.tr.Flush() | ||
| assert.Less(t, time.Since(start), 500*time.Millisecond, "flush waited on a send") |
There was a problem hiding this comment.
P2: This 500 ms wall-clock assertion can fail under CI scheduling contention, especially in a parallel test. Use a synchronization-based check with a more generous deadline to verify Flush returns while the handler remains blocked.
Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. When an issue isn't valid or won't be fixed in this PR, reply in its thread with the reason and then resolve the thread. At pkg/transport/intershard_internal_test.go, line 79:
<comment>This 500 ms wall-clock assertion can fail under CI scheduling contention, especially in a parallel test. Use a synchronization-based check with a more generous deadline to verify `Flush` returns while the handler remains blocked.</comment>
<file context>
@@ -0,0 +1,136 @@
+ fixtureA.tr.Enqueue(context.Background(), fixtureB.address, testutils.SimpleCommand{Value: flush})
+ start := time.Now()
+ fixtureA.tr.Flush()
+ assert.Less(t, time.Since(start), 500*time.Millisecond, "flush waited on a send")
+ }
+
</file context>
Start now flushes the NATS connection after registering endpoints, so a command sent right after Start returns is not answered with no responders and dropped; this was the cause of the intermittent CI failures in the cardinal inter-shard tests. SendCommandWithReply registers its reply waiter before dispatching, so a handler that publishes its reply before returning, as a non-ECS service may, is delivered. A second Start call is rejected instead of leaking the first connection, and Start is split into startNATS and startHTTP to stay within the linter's length limit. On the cardinal side, NewWorld and the end-to-end harness now build the transport through one newTransport helper instead of two copies, the debug service test goes through startTransport so it fails if startup stops finalizing the catalog, and the end-to-end harness sends the X-Email header that Dev auth has required since passthrough auth was removed.
There was a problem hiding this comment.
0 issues found across 7 files (changes from recent commits).
Shadow auto-approve: would not auto-approve. Auto-approval blocked by 3 unresolved issues from previous reviews.
Re-trigger cubic
…ork on Stop The transport no longer opens its own NATS connection or HTTP server. The app passes its micro.Client to Start and serves the handler returned by Handler on its own server, so one connection can be shared and the app decides when connections open and close. Cardinal's World now opens one NATS client, shares it with the JetStream snapshot store instead of letting the store open a second connection, serves the transport and debug service on its own mux, and closes the client last on shutdown. Enqueue and Flush are now safe to call from any goroutine, so a service that handles commands concurrently can send to other shards from inside its handlers. Stop now rejects new commands as unavailable and unsubscribes the NATS endpoints, waits for running handlers, sends what they flushed, and then ends open event streams and reply waits, so a reply committed by a running handler is no longer lost and an open stream no longer holds the HTTP server's shutdown until the deadline. Commands flushed after Stop are dropped and logged.
There was a problem hiding this comment.
3 issues found across 12 files (changes from recent commits).
Prompt for AI agents (unresolved issues)
Check if these issues are valid — if so, understand the root cause of each and fix them. When an issue isn't valid or won't be fixed in this PR, reply in its thread with the reason and then resolve the thread. If appropriate, use sub-agents to investigate and fix each issue separately.
<file name="pkg/transport/intershard_internal_test.go">
<violation number="1" location="pkg/transport/intershard_internal_test.go:82">
P2: This assertion lets a routing regression send all 400 commands to one target and still pass. Assert the expected payload set for each target, including that each value appears exactly once.</violation>
</file>
<file name="pkg/cardinal/service.go">
<violation number="1" location="pkg/cardinal/service.go:56">
P2: A later `NewWorld` failure can strand this connection: JetStream snapshot initialization returns an error without exposing the `World`, so no shutdown can close `w.client`. Close the client on construction failure.</violation>
</file>
<file name="pkg/cardinal/snapshot/storage_jetstream.go">
<violation number="1" location="pkg/cardinal/snapshot/storage_jetstream.go:108">
P3: `Logger` is now a no-op: `NewJetStreamStorage` no longer uses it to configure a NATS client, but Cardinal still supplies `tel.GetLogger("snapshot")`. Remove this option and its call-site initializer; configure logging on the shared client instead.</violation>
</file>
Shadow auto-approve: would not auto-approve because issues were found.
Tip: Review your code locally with the cubic CLI to iterate faster.
Re-trigger cubic
| total++ | ||
| } | ||
| } | ||
| assert.Equal(t, goroutines*perGoroutine, total) |
There was a problem hiding this comment.
P2: This assertion lets a routing regression send all 400 commands to one target and still pass. Assert the expected payload set for each target, including that each value appears exactly once.
Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. When an issue isn't valid or won't be fixed in this PR, reply in its thread with the reason and then resolve the thread. At pkg/transport/intershard_internal_test.go, line 82:
<comment>This assertion lets a routing regression send all 400 commands to one target and still pass. Assert the expected payload set for each target, including that each value appears exactly once.</comment>
<file context>
@@ -39,6 +40,48 @@ func TestInterShard_StopSendsFlushedInOrder(t *testing.T) {
+ total++
+ }
+ }
+ assert.Equal(t, goroutines*perGoroutine, total)
+}
+
</file context>
| } | ||
| s.microService = microService | ||
| s.interShard = newInterShard(s.world.address, client, &s.world.commands, s.log) | ||
| w.client = client |
There was a problem hiding this comment.
P2: A later NewWorld failure can strand this connection: JetStream snapshot initialization returns an error without exposing the World, so no shutdown can close w.client. Close the client on construction failure.
Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. When an issue isn't valid or won't be fixed in this PR, reply in its thread with the reason and then resolve the thread. At pkg/cardinal/service.go, line 56:
<comment>A later `NewWorld` failure can strand this connection: JetStream snapshot initialization returns an error without exposing the `World`, so no shutdown can close `w.client`. Close the client on construction failure.</comment>
<file context>
@@ -39,12 +39,90 @@ func (w *World) newTransport() (*transport.Transport, error) {
+ if err != nil {
+ return nil, eris.Wrap(err, "failed to connect to NATS")
+ }
+ w.client = client
+ return client, nil
+}
</file context>
| Logger zerolog.Logger | ||
| NATSConfig *micro.NATSConfig // Optional NATS config override (nil = use env/defaults) | ||
| Address *micro.ServiceAddress | ||
| Logger zerolog.Logger |
There was a problem hiding this comment.
P3: Logger is now a no-op: NewJetStreamStorage no longer uses it to configure a NATS client, but Cardinal still supplies tel.GetLogger("snapshot"). Remove this option and its call-site initializer; configure logging on the shared client instead.
Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. When an issue isn't valid or won't be fixed in this PR, reply in its thread with the reason and then resolve the thread. At pkg/cardinal/snapshot/storage_jetstream.go, line 108:
<comment>`Logger` is now a no-op: `NewJetStreamStorage` no longer uses it to configure a NATS client, but Cardinal still supplies `tel.GetLogger("snapshot")`. Remove this option and its call-site initializer; configure logging on the shared client instead.</comment>
<file context>
@@ -113,9 +104,9 @@ func (j *JetStreamStorage) Load(ctx context.Context) ([]byte, error) {
- Logger zerolog.Logger
- NATSConfig *micro.NATSConfig // Optional NATS config override (nil = use env/defaults)
+ Address *micro.ServiceAddress
+ Logger zerolog.Logger
+ Client *micro.Client // NATS connection; the caller owns it
</file context>
…returns Publish picks the subscribed streams under the lock, releases it, and then sends to each one. If a stream's handler returned in between, because the client disconnected or Stop ended the stream, connect-go finished the stream while the send was still writing to it, which the race detector reports. The tick made this rare in Cardinal, but Stop ending every stream at once and services that publish from concurrent handlers make it likely. Each stream subscriber now has a closed flag guarded by the lock that send already holds around every write. The stream handler sets it before returning, which waits for a send in progress, and later sends skip the stream. A race test publishes in a loop while streams open and close.
There was a problem hiding this comment.
1 issue found across 2 files (changes from recent commits).
Prompt for AI agents (unresolved issues)
Check if these issues are valid — if so, understand the root cause of each and fix them. When an issue isn't valid or won't be fixed in this PR, reply in its thread with the reason and then resolve the thread. If appropriate, use sub-agents to investigate and fix each issue separately.
<file name="pkg/transport/stop_internal_test.go">
<violation number="1" location="pkg/transport/stop_internal_test.go:185">
P2: This test can pass without exercising an open stream if the initial receive fails. Require the initial message before testing the publish/close race.</violation>
</file>
Shadow auto-approve: would not auto-approve because issues were found.
Tip: Review your code locally with the cubic CLI to iterate faster.
Re-trigger cubic
| req.Header().Set("X-Email", "user-"+strconv.Itoa(i)) | ||
| stream, err := client.StartEventStream(ctx, req) | ||
| require.NoError(t, err) | ||
| stream.Receive() |
There was a problem hiding this comment.
P2: This test can pass without exercising an open stream if the initial receive fails. Require the initial message before testing the publish/close race.
Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. When an issue isn't valid or won't be fixed in this PR, reply in its thread with the reason and then resolve the thread. At pkg/transport/stop_internal_test.go, line 185:
<comment>This test can pass without exercising an open stream if the initial receive fails. Require the initial message before testing the publish/close race.</comment>
<file context>
@@ -140,3 +142,48 @@ func TestTransport_FlushAfterStop(t *testing.T) {
+ req.Header().Set("X-Email", "user-"+strconv.Itoa(i))
+ stream, err := client.StartEventStream(ctx, req)
+ require.NoError(t, err)
+ stream.Receive()
+ time.Sleep(time.Millisecond) // Let some events reach the stream before it closes
+ cancel()
</file context>
| stream.Receive() | |
| require.True(t, stream.Receive(), "no initial message: %v", stream.Err()) |

Cardinal's ConnectRPC service, auth, event streams and inter-shard send/receive move into a new pkg/transport that knows nothing about ECS. Modules register command handlers with Handle and send through Publish, Enqueue and Flush, so a non-ECS service can use the same transport. Cardinal now wires its command manager, event dispatch and debug service onto it; wire format and span names are unchanged.
Stack created with GitHub Stacks CLI • Give Feedback 💬
Summary by cubic
Moves Cardinal's client-facing ConnectRPC service, auth, event streams, and inter-shard send/receive over NATS into a new
pkg/transportthat is independent of the ECS, so non-ECS services can reuse the same transport.Refactors
micro.ClienttoStartand serves the handler fromHandleron its own server, closing the client last on shutdown; Cardinal shares one NATS client with the JetStream snapshot store instead of letting the store open a second connection.Startflushes the NATS connection after registering endpoints so a command sent right after start is not dropped; a secondStartis rejected.SendCommandWithReplyregisters its reply waiter before dispatching, so a handler that publishes its reply before returning receives it;EnqueueandFlushare safe to call from any goroutine.Stopno longer races with the stream handler's return.Stoprejects new commands as unavailable, unsubscribes the NATS endpoints, waits for running handlers, sends what they flushed, then ends open event streams and reply waits; commands flushed after stop are dropped and logged.newTransporthelper; the harness sends theX-Emailheader Dev auth requires.Handleand send throughPublish,Enqueue, andFlush; wire format and span names are unchanged, so dashboard and alert keys stay stable.Written for commit 2928f4b. Summary will update on new commits.