Skip to content

fix(cache): stop orphaning watches - #18

Merged
krhitesh7 merged 1 commit into
nexusfrom
fix/status-slot-race
Aug 3, 2026
Merged

krhitesh7 merged 1 commit into
nexusfrom
fix/status-slot-race

Conversation

@krhitesh7

@krhitesh7 krhitesh7 commented Jul 31, 2026 •

Copy link
Copy Markdown
Member

key snapshots by node ID, not hash slot

snapshotCache stored snapshots and status in two flat 16384-entry slices indexed by hash.CacheIndexFromKey(nodeID), with the cache-wide mutex commented out. That produced three defects.

Orphaned watches. CreateDeltaWatch and CreateWatch read, then created, then wrote a status slot with no synchronisation. When several streams sharing one node ID opened their first watch simultaneously — what a control-plane restart causes, since every proxy reconnects at once — each goroutine saw a nil statusInfo and installed its own. Only the last write stayed reachable; watches registered into the others were invisible to respondDeltaWatches for the life of the process. The loss was permanent: an orphaned stream blocks on a response that can no longer be sent, so it never ACKs, so the server never re-registers it. Measured in production: 45 of 52 proxies stuck without two clusters that were present in the control plane's own snapshot.

Silent cross-node config. A slot recorded no owner, so two node IDs whose hashes collided shared one snapshot — node A served node B's config. The shipped IDHash default returns len(key), which collides every node ID of equal length.

Lost wakeups. CreateDeltaWatch decided "client is up to date" and registered the watch as two separate steps. An upsert landing between them mutated the snapshot, walked the watches, did not find this one, and skipped the stream.

Replace both slices with 256 shards of node-ID-keyed maps under an RWMutex. Creation is atomic, distinct nodes never share state, and unrelated nodes still proceed concurrently. Shard selection hashes internally rather than through NodeHash.CacheIndexFromKey, so a poor caller hash now costs only distribution. CreateDeltaWatch holds info.mu across the up-to-date check and registration, making it atomic against respondDeltaWatches. Lock order throughout is shard.mu -> info.mu -> snapshot.Mu.

respondDelta's channel send becomes non-blocking. A blocking send under info.mu deadlocks against the stream goroutine, which needs the same lock to re-register, and the ctx.Done() escape never fires because callers pass context.Background(). A full channel now reports "not responded", leaving the watch registered so the next upsert re-diffs and delivers it.

Also fixed, both crashes:

  • NewSnapshotCacheWithHeartbeating panicked on its first tick, dereferencing the ~16382 nil entries in the status slice.
  • CreateWatch's snapshot != nil && version match fell through to the respond path with a nil snapshot and panicked. Upstream leaves a watch open when no snapshot exists; restored that.

BatchUpsertResources' two silent early returns now log. They return a nil error, so callers could not tell a push had been discarded.

Test package restored to compiling — it had rotted against this fork's own API changes (types.Resource gained MarshalVTStrict, NodeHash gained CacheIndex, responses carry VTMarshaledResource), so none of these tests had been running. Adds IndexVTResourcesByName for the pre-marshaled response type.

New regression tests fail against the previous implementation: 18/50 streams served with 32 permanently starved, plus data races at the status slot.

Known remaining, all pre-existing and unrelated: 12 test failures in snapshot consistency and delta handling; TestSnapshotCache is flaky because it byte-compares a proto containing a map; TestLinearConcurrentSetWatch hangs under -race inside linearCache, which this change does not touch.

Summary by CodeRabbit

  • Bug Fixes

    • Improved cache reliability during concurrent watch creation and updates.
    • Prevented distinct nodes from sharing cached state, including when shard collisions occur.
    • Ensured updates are delivered reliably and watches remain registered when responses are delayed or unavailable.
    • Safely handles resource updates when required cache state has not yet been initialized.
  • Tests

    • Added coverage for concurrent watches, cache isolation, resource indexing, and marshaled resource handling.

…h slot

snapshotCache stored snapshots and status in two flat 16384-entry slices
indexed by hash.CacheIndexFromKey(nodeID), with the cache-wide mutex commented
out. That produced three defects.

Orphaned watches. CreateDeltaWatch and CreateWatch read, then created, then
wrote a status slot with no synchronisation. When several streams sharing one
node ID opened their first watch simultaneously — what a control-plane restart
causes, since every proxy reconnects at once — each goroutine saw a nil
statusInfo and installed its own. Only the last write stayed reachable; watches
registered into the others were invisible to respondDeltaWatches for the life
of the process. The loss was permanent: an orphaned stream blocks on a response
that can no longer be sent, so it never ACKs, so the server never re-registers
it. Measured in production: 45 of 52 proxies stuck without two clusters that
were present in the control plane's own snapshot.

Silent cross-node config. A slot recorded no owner, so two node IDs whose
hashes collided shared one snapshot — node A served node B's config. The
shipped IDHash default returns len(key), which collides every node ID of equal
length.

Lost wakeups. CreateDeltaWatch decided "client is up to date" and registered
the watch as two separate steps. An upsert landing between them mutated the
snapshot, walked the watches, did not find this one, and skipped the stream.

Replace both slices with 256 shards of node-ID-keyed maps under an RWMutex.
Creation is atomic, distinct nodes never share state, and unrelated nodes still
proceed concurrently. Shard selection hashes internally rather than through
NodeHash.CacheIndexFromKey, so a poor caller hash now costs only distribution.
CreateDeltaWatch holds info.mu across the up-to-date check and registration,
making it atomic against respondDeltaWatches. Lock order throughout is
shard.mu -> info.mu -> snapshot.Mu.

respondDelta's channel send becomes non-blocking. A blocking send under info.mu
deadlocks against the stream goroutine, which needs the same lock to
re-register, and the ctx.Done() escape never fires because callers pass
context.Background(). A full channel now reports "not responded", leaving the
watch registered so the next upsert re-diffs and delivers it.

Also fixed, both crashes:
- NewSnapshotCacheWithHeartbeating panicked on its first tick, dereferencing
  the ~16382 nil entries in the status slice.
- CreateWatch's `snapshot != nil && version match` fell through to the respond
  path with a nil snapshot and panicked. Upstream leaves a watch open when no
  snapshot exists; restored that.

BatchUpsertResources' two silent early returns now log. They return a nil error,
so callers could not tell a push had been discarded.

Test package restored to compiling — it had rotted against this fork's own API
changes (types.Resource gained MarshalVTStrict, NodeHash gained CacheIndex,
responses carry VTMarshaledResource), so none of these tests had been running.
Adds IndexVTResourcesByName for the pre-marshaled response type.

New regression tests fail against the previous implementation: 18/50 streams
served with 32 permanently starved, plus data races at the status slot.

Known remaining, all pre-existing and unrelated: 12 test failures in snapshot
consistency and delta handling; TestSnapshotCache is flaky because it
byte-compares a proto containing a map; TestLinearConcurrentSetWatch hangs
under -race inside linearCache, which this change does not touch.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
@coderabbitai

coderabbitai Bot commented Jul 31, 2026 •

Copy link
Copy Markdown

Review Change Stack

📝 Walkthrough

Walkthrough

The cache replaces hash-indexed slices with sharded node-keyed maps, synchronizes snapshot and status operations, makes watch registration and response delivery concurrency-safe, and updates tests to use VT-marshaled resources.

Changes

Cache state and watch concurrency

Layer / File(s) Summary
VT-marshaled resource test contracts
pkg/cache/v3/resources.go, pkg/cache/v3/*_test.go
Adds VT resource indexing and adapters, then updates cache, delta, linear, and snapshot tests to construct and compare serialized resources.
Sharded node state storage
pkg/cache/v3/simple.go, pkg/cache/v3/simple_test.go
Stores snapshots and statuses in 256 synchronized shards keyed by node ID, updating heartbeat, snapshot, fetch, cleanup, and status enumeration paths.
Atomic watch registration and delivery
pkg/cache/v3/simple.go, pkg/cache/v3/status.go, pkg/cache/v3/simple_race_test.go, pkg/cache/v3/delta_test.go
Makes watch registration atomic with response evaluation, retains watches when snapshots or response channels are unavailable, and adds concurrent registration and collision tests.

Estimated code review effort: 4 (Complex) | ~45 minutes

Sequence Diagram(s)

sequenceDiagram
  participant Client
  participant SimpleCache
  participant NodeStatus
  participant SnapshotStore
  participant DeltaStream
  Client->>SimpleCache: CreateDeltaWatch(nodeID)
  SimpleCache->>NodeStatus: Get or create status
  SimpleCache->>SnapshotStore: Lookup snapshot(nodeID)
  SimpleCache->>NodeStatus: Check version and register watch under lock
  SimpleCache-->>Client: Return watch response or retained watch
  SimpleCache->>SnapshotStore: Store updated snapshot
  SimpleCache->>NodeStatus: Evaluate registered watches
  NodeStatus->>DeltaStream: Non-blocking delta response
Loading

Suggested labels: bugfix, refactor, tests, performance

🚥 Pre-merge checks | ✅ 4 | ❌ 1

❌ Failed checks (1 warning)

Check name Status Explanation Resolution
Docstring Coverage ⚠️ Warning Docstring coverage is 36.36% which is insufficient. The required threshold is 80.00%. Write docstrings for the functions missing them to satisfy the coverage threshold.
✅ Passed checks (4 passed)
Check name Status Explanation
Description Check ✅ Passed Check skipped - CodeRabbit’s high-level summary is enabled.
Title check ✅ Passed The title clearly and concisely describes the primary cache change: preventing watches from becoming orphaned.
Linked Issues check ✅ Passed Check skipped because no linked issues were found for this pull request.
Out of Scope Changes check ✅ Passed Check skipped because no linked issues were found for this pull request.
✨ Finishing Touches 💡 1
📝 Generate docstrings 💡
  • Create stacked PR
  • Commit on current branch
🧪 Generate unit tests (beta)
  • Create PR with unit tests
  • Commit unit tests in branch fix/status-slot-race

Comment @coderabbitai help to get the list of available commands.

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Actionable comments posted: 1

Caution

Some comments are outside the diff and can’t be posted inline due to platform limitations.

⚠️ Outside diff range comments (1)
pkg/cache/v3/delta_test.go (1)

207-215: 🩺 Stability & Availability | 🟡 Minor | ⚡ Quick win

cancel can be nil here, so defer cancel() may panic.

CreateDeltaWatch returns (nil, false) when it responds inline (see pkg/cache/v3/simple.go Line 1005). This test sets snapshots from parallel subtests, so the watch branch can hit the exists path and get a nil cancel func. Guard it.

🛡️ Proposed guard
-					defer cancel()
+					if cancel != nil {
+						defer cancel()
+					}
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@pkg/cache/v3/delta_test.go` around lines 207 - 215, Guard the deferred
cleanup in the test around CreateDeltaWatch so cancel is invoked only when
non-nil. Preserve the existing watch setup and response behavior, while
preventing the inline exists path from panicking when CreateDeltaWatch returns a
nil cancel function.
🧹 Nitpick comments (6)
pkg/cache/v3/simple.go (2)

1144-1150: 📐 Maintainability & Code Quality | 🔵 Trivial | 💤 Low value

Leftover empty block.

♻️ Proposed cleanup
 	all := cache.allStatus()
 	out := make([]string, 0, len(all))
 	for _, statusInfo := range all {
-		{
-			out = append(out, statusInfo.GetNode().GetId())
-		}
+		out = append(out, statusInfo.GetNode().GetId())
 	}
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@pkg/cache/v3/simple.go` around lines 1144 - 1150, Remove the unnecessary
nested empty block inside the loop iterating over allStatus in the surrounding
function, while preserving the existing append of each statusInfo node ID to
out.

293-300: 📐 Maintainability & Code Quality | 🔵 Trivial | 💤 Low value

Stale comment.

allStatus() already filters nils and returns a slice, so the "status is a sparse slice … nil check on the element itself is load-bearing" note no longer describes this code. The remaining info.GetNode() != nil check is what matters.

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@pkg/cache/v3/simple.go` around lines 293 - 300, Remove the stale comment
above the allStatus loop in the heartbeat flow, since allStatus() already
filters nil entries. Keep the existing info.GetNode() != nil guard and
cache.sendHeartbeats(ctx, node) behavior unchanged.
pkg/cache/v3/simple_test.go (1)

179-185: 📐 Maintainability & Code Quality | 🔵 Trivial | 💤 Low value

Duplicated assertion block.

Lines 179-181 and 183-185 assert exactly the same thing. Also, the failure message prints GetResources(typ) while the comparison uses GetResourcesAndTTL(typ), which will mislead on failure.

♻️ Proposed cleanup
-					if !reflect.DeepEqual(cache.IndexVTResourcesByName(out.(*cache.RawResponse).Resources), snapshotWithTTL.GetResourcesAndTTL(typ)) {
-						t.Errorf("get resources %v, want %v", out.(*cache.RawResponse).Resources, snapshotWithTTL.GetResources(typ))
-					}
-
 					if !reflect.DeepEqual(cache.IndexVTResourcesByName(out.(*cache.RawResponse).Resources), snapshotWithTTL.GetResourcesAndTTL(typ)) {
-						t.Errorf("get resources %v, want %v", out.(*cache.RawResponse).Resources, snapshotWithTTL.GetResources(typ))
+						t.Errorf("get resources %v, want %v", out.(*cache.RawResponse).Resources, snapshotWithTTL.GetResourcesAndTTL(typ))
 					}
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@pkg/cache/v3/simple_test.go` around lines 179 - 185, In the test assertion
block, remove the duplicated reflect.DeepEqual check and retain a single
assertion. Update its failure message to print the same
snapshotWithTTL.GetResourcesAndTTL(typ) value used in the comparison.
pkg/cache/v3/simple_race_test.go (1)

102-114: 📐 Maintainability & Code Quality | 🔵 Trivial | 💤 Low value

Replace the fixed sleep with a bounded wait.

UpsertResources responds to delta watches synchronously before returning, so the 500ms sleep is dead time; if delivery ever becomes async this is also the wrong knob. Collect with a per-channel timeout instead.

♻️ Proposed change
-	time.Sleep(500 * time.Millisecond)
 	delivered := 0
+	deadline := time.After(2 * time.Second)
 	for i := 0; i < streams; i++ {
 		select {
 		case <-chans[i]:
 			delivered++
-		default:
+		case <-deadline:
 		}
 	}
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@pkg/cache/v3/simple_race_test.go` around lines 102 - 114, In the delivery
verification around UpsertResources, remove the fixed 500ms sleep and wait for
each channel in chans with its own bounded timeout while collecting deliveries.
Preserve the delivered count and starvation error reporting, and ensure the test
does not block indefinitely if a channel receives no update.
pkg/cache/v3/status.go (1)

173-176: 📐 Maintainability & Code Quality | 🔵 Trivial | ⚖️ Poor tradeoff

NodeHash.CacheIndex/CacheIndexFromKey are now dead surface.

simple.go hashes internally via shardFor, so these two interface methods (Lines 27-32) no longer influence behavior yet every implementor must still provide them. Consider deprecating them (or dropping them in a breaking release) so callers aren't misled into thinking their index function matters.

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@pkg/cache/v3/status.go` around lines 173 - 176, Deprecate the
NodeHash.CacheIndex and CacheIndexFromKey interface methods because simple.go
now determines shards internally through shardFor and no longer uses these
methods. Update the interface documentation and affected implementors or call
sites consistently, while preserving the existing sharding behavior; remove the
methods only if this is an intentional breaking-release change.
pkg/cache/v3/delta_test.go (1)

24-30: 📐 Maintainability & Code Quality | 🔵 Trivial | 💤 Low value

want/got are swapped at every call site.

Callers pass the actual response first and the expected snapshot second (e.g. Line 58), so failures print the expectation as "got" and vice versa. Either swap the parameter order or rename them.

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@pkg/cache/v3/delta_test.go` around lines 24 - 30, Correct
assertResourceMapEqual so its parameter order or names match the existing call
sites, where the actual response is passed first and the expected snapshot
second. Update the cmp.Equal comparison and failure message consistently so
diagnostics label actual data as got and the snapshot as want.
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

Inline comments:
In `@pkg/cache/v3/simple.go`:
- Around line 1060-1079: The delta watch timeout test must match the
non-blocking send contract in the responder flow around the delta response
channel: update TestSnapshotDeltaCacheWatchTimeout to expect SetSnapshot to
succeed when the unbuffered channel is full, with the watch remaining registered
for a later retry. Do not expect context.Canceled to propagate through ADS,
since responder errors are discarded there.

---

Outside diff comments:
In `@pkg/cache/v3/delta_test.go`:
- Around line 207-215: Guard the deferred cleanup in the test around
CreateDeltaWatch so cancel is invoked only when non-nil. Preserve the existing
watch setup and response behavior, while preventing the inline exists path from
panicking when CreateDeltaWatch returns a nil cancel function.

---

Nitpick comments:
In `@pkg/cache/v3/delta_test.go`:
- Around line 24-30: Correct assertResourceMapEqual so its parameter order or
names match the existing call sites, where the actual response is passed first
and the expected snapshot second. Update the cmp.Equal comparison and failure
message consistently so diagnostics label actual data as got and the snapshot as
want.

In `@pkg/cache/v3/simple_race_test.go`:
- Around line 102-114: In the delivery verification around UpsertResources,
remove the fixed 500ms sleep and wait for each channel in chans with its own
bounded timeout while collecting deliveries. Preserve the delivered count and
starvation error reporting, and ensure the test does not block indefinitely if a
channel receives no update.

In `@pkg/cache/v3/simple_test.go`:
- Around line 179-185: In the test assertion block, remove the duplicated
reflect.DeepEqual check and retain a single assertion. Update its failure
message to print the same snapshotWithTTL.GetResourcesAndTTL(typ) value used in
the comparison.

In `@pkg/cache/v3/simple.go`:
- Around line 1144-1150: Remove the unnecessary nested empty block inside the
loop iterating over allStatus in the surrounding function, while preserving the
existing append of each statusInfo node ID to out.
- Around line 293-300: Remove the stale comment above the allStatus loop in the
heartbeat flow, since allStatus() already filters nil entries. Keep the existing
info.GetNode() != nil guard and cache.sendHeartbeats(ctx, node) behavior
unchanged.

In `@pkg/cache/v3/status.go`:
- Around line 173-176: Deprecate the NodeHash.CacheIndex and CacheIndexFromKey
interface methods because simple.go now determines shards internally through
shardFor and no longer uses these methods. Update the interface documentation
and affected implementors or call sites consistently, while preserving the
existing sharding behavior; remove the methods only if this is an intentional
breaking-release change.
🪄 Autofix (Beta)

Fix all unresolved CodeRabbit comments on this PR:

  • Push a commit to this branch (recommended)
  • Create a new PR with the fixes

ℹ️ Review info
⚙️ Run configuration

Configuration used: Repository: ShareChat/coderabbit/.coderabbit.yaml

Review profile: CHILL

Plan: Pro

Run ID: e73bd771-0b36-4763-a2f0-241d5bd26308

📥 Commits

Reviewing files that changed from the base of the PR and between 18a7341 and 4ea6e6a.

📒 Files selected for processing (8)
  • pkg/cache/v3/cache_test.go
  • pkg/cache/v3/delta_test.go
  • pkg/cache/v3/linear_test.go
  • pkg/cache/v3/resources.go
  • pkg/cache/v3/simple.go
  • pkg/cache/v3/simple_race_test.go
  • pkg/cache/v3/simple_test.go
  • pkg/cache/v3/status.go

Comment thread pkg/cache/v3/simple.go

@jensoncs jensoncs left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Verified the metrics and The change looks fine.

@krhitesh7
krhitesh7 merged commit ffab6ee into nexus Aug 3, 2026
1 of 5 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants