Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
22 commits
Select commit Hold shift + click to select a range
dace0be
Fix topic acknowledgments and resource budgets
polRk Oct 8, 2026
cf15e2a
Explain topic read credit and flush ordering
polRk Oct 8, 2026
7e0e540
Simplify reader credit to response refunds
polRk Oct 8, 2026
2d9d1de
Preserve channels through endpoint rediscovery
polRk Oct 8, 2026
c5d27ff
Simplify topic state and restore timed batching
polRk Oct 8, 2026
0a4e7aa
Bound retained memory across topic lifecycles
polRk Oct 8, 2026
0814e7f
Keep writer batch readiness checks linear
polRk Oct 8, 2026
e0daa37
Preserve TLS identity in memory test tunnels
polRk Oct 8, 2026
8b55d6e
Restore writer payload byte limits
polRk Oct 9, 2026
3c8a675
Fix lifecycle gaps found in topic audit
polRk Oct 9, 2026
1734c36
Build writer messages in one pass
polRk Oct 9, 2026
3b480f5
Merge main into topic reliability branch
polRk Oct 9, 2026
90ce9b2
Use explicit blocks in topic client changes
polRk Oct 9, 2026
4dc0d91
Check RSS growth in topic memory profiles
polRk Oct 9, 2026
1073365
Keep reader token renewal independent of ACKs
polRk Oct 9, 2026
35583b7
Stop retrying oversized topic read responses
polRk Oct 10, 2026
e3e080a
Represent retained read credit in the FSM
polRk Oct 10, 2026
331d706
Verify active RPCs survive rediscovery
polRk Oct 10, 2026
f73eee0
Harden topic recovery and buffer contracts
polRk Oct 10, 2026
7f91f1d
Cover child delivery after a forced parent stop
polRk Oct 10, 2026
fadc2a2
Bound flush to previously accepted writes
polRk Oct 10, 2026
e12b090
Clarify reconnect deduplication invariant
polRk Oct 10, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions .changeset/stable-endpoint-lifecycle.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
'@ydbjs/core': patch
---

Keep active channels alive when an endpoint repeatedly disappears from and returns to discovery. Reject `ready()` after driver shutdown while preserving the original cause of an initial discovery failure.
5 changes: 5 additions & 0 deletions .changeset/topic-flush-boundary.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
'@ydbjs/topic': patch
---

Make `flush()` wait for acknowledgments of writes accepted before the call. Later writes no longer extend a pending flush. Independent flush boundaries survive reconnect and server deduplication; terminal errors still reject unconfirmed flushes. Use `close()` to stop admission and drain the entire writer.
11 changes: 11 additions & 0 deletions .changeset/topic-lifecycle-audit.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
---
'@ydbjs/core': patch
'@ydbjs/fsm': patch
'@ydbjs/topic': patch
---

Expire unreassigned reader partitions after reconnect even when their commits arrive after session initialization. Reject non-positive reader buffer limits instead of leaving reads stalled.

Keep discovered and explicitly pinned connections separate, close only the affected connection during retirement or invalidation, and drain both during driver shutdown. Preserve the configured TLS server name when discovery provides no override, and omit driver options from debug logging.

Preserve queued items when a paused read is cancelled immediately after resuming.
5 changes: 5 additions & 0 deletions .changeset/topic-linear-batching.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
'@ydbjs/topic': patch
---

Avoid repeatedly scanning a partial writer batch on every `write()`. Track the compressed payload size of the unsent queue so readiness checks remain constant-time while preserving timed flushes, byte/count limits and reconnect recovery.
11 changes: 11 additions & 0 deletions .changeset/topic-memory-ownership.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
---
'@ydbjs/core': patch
'@ydbjs/fsm': patch
'@ydbjs/topic': patch
---

Release cancelled readiness waiters and completed queue items instead of retaining their payloads until the driver or queue closes.

Keep topic reader buffers encoded until `read()` selects messages for delivery. Preserve the unread tail after a clean close and discard it on destruction or terminal failure, including when a read iterator is suspended. Release stopped partition state without pending commits and callbacks or custom codecs after client shutdown. Preserve delivered transaction offsets until the transaction finishes.

Read credit still uses each complete server response's `bytesSize`; no credit is inferred from individual message payloads. `maxBufferBytes` does not bound decompressed data or process RSS. Use `read({ limit })` to bound how many messages are decoded in one batch.
5 changes: 5 additions & 0 deletions .changeset/topic-reader-receive-limit.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
'@ydbjs/topic': patch
---

Report oversized gRPC read responses as terminal errors instead of repeatedly reconnecting to the same unread message. Temporary resource exhaustion remains retryable.
5 changes: 5 additions & 0 deletions .changeset/topic-reader-token-renewal.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
'@ydbjs/topic': patch
---

Keep refreshing reader credentials when the server does not acknowledge an unchanged token.
5 changes: 5 additions & 0 deletions .changeset/topic-stream-recovery-contracts.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
'@ydbjs/topic': patch
---

Resume pending reads, commits, and writes after transient stream deadlines and send buffered partial batches after reconnect without restarting their batching delay. Preserve reader source settings after caller mutation and contain rejected async acknowledgment observers. Clarify RAW buffer ownership.
17 changes: 17 additions & 0 deletions .changeset/topic-stream-reliability.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,17 @@
---
'@ydbjs/topic': patch
---

Preserve topic delivery guarantees and resource budgets across concurrent calls and reconnects.

- Prevent an earlier `flush()` completion from resolving a later flush before its messages are acknowledged.
- Release payload references recovered through reconnect deduplication, and always close writer streams and timers after an internal failure.
- Snapshot metadata and dates, avoid retaining oversized backing buffers, and reject invalid dates or sequence numbers before accepting a message.
- Preserve the reader's retained-byte budget across reconnects and account for releases during backoff without granting excess credit to a new stream.
- Honor read cancellation and partition revocation between slices of a response while preserving unread messages and releasing each response's credit once.
- Restore the `TopicTxWriter` type export for consumers that import the transactional writer contract.
- Restore timed batching: send full batches immediately and partial batches on `flushIntervalMs`; explicit `flush()` and `close()` drain without waiting for the interval. Preserve an expired flush deadline while the in-flight window is occupied.
- Keep graceful writer shutdown on the normal connection lifecycle, waiting for every reconnect handshake and cancelling obsolete retry timers after initialization.
- Ignore late stream completions and token refreshes from replaced connections, and prevent opening a stream after cancellation during driver readiness.
- Preserve transaction cancellation when updating read offsets, keep registered reader callbacks stable, and retain falsy destruction reasons in terminal errors.
- Generate default producer identities with `randomUUID()` and make manual flush sequence numbers independent of initialization timing.
15 changes: 7 additions & 8 deletions e2e/topic/read-write.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,6 @@ import {
TopicServiceDefinition,
} from '@ydbjs/api/topic'
import { create } from '@bufbuild/protobuf'
import { once } from 'node:events'

// #region setup
declare module 'vitest' {
Expand Down Expand Up @@ -92,6 +91,7 @@ test('writes and reads concurrently', { timeout: 60_000 }, async (tc) => {
await using writer = createTopicWriter(driver, {
topic: testTopicName,
producer: testProducerName,
maxBufferBytes: BigInt(TOTAL_TRAFFIC) + 16n * 1024n * 1024n,
maxInflightCount: TOTAL_BATCHES * BATCH_SIZE,
})

Expand All @@ -105,11 +105,11 @@ test('writes and reads concurrently', { timeout: 60_000 }, async (tc) => {
let rb = 0
let ctrl = new AbortController()
// linkSignals, not AbortSignal.any (composite signals accumulate listeners).
using combined = linkSignals(tc.signal, ctrl.signal, AbortSignal.timeout(25_000))
using combined = linkSignals(tc.signal, ctrl.signal)
let signal = combined.signal

// Producer.
void (async () => {
let producerTask = (async () => {
while (wb < TOTAL_TRAFFIC) {
if (signal.aborted) break

Expand All @@ -121,28 +121,27 @@ test('writes and reads concurrently', { timeout: 60_000 }, async (tc) => {
}

let start = performance.now()
await writer.flush()
await writer.flush(tc.signal)
console.log(`write took ${performance.now() - start} ms`)
})()

// Consumer.
void (async () => {
let consumerTask = (async () => {
for await (let batch of reader.read({ signal })) {
let promise = reader.commit(batch)
await reader.commit(batch)
rb += MESSAGE_SIZE * batch.length

// >=, not ==: at-least-once redelivery can overshoot the total, and an
// exact-equality trigger would then never fire and hang the test.
if (rb >= TOTAL_TRAFFIC) {
await promise
ctrl.abort()
break
}
}
})()

let start = Date.now()
await once(ctrl.signal, 'abort')
await Promise.all([producerTask, consumerTask])
await writer.close()
await reader.close()

Expand Down
13 changes: 5 additions & 8 deletions packages/core/ARCHITECTURE.md
Original file line number Diff line number Diff line change
Expand Up @@ -26,9 +26,9 @@ registry, and health. On any change to the routable set it rebuilds an **immutab
`acquireNode()` — reads the latest snapshot reference (swapped by `#consume`),
selects a `RoutingSnapshot` ref (pure), and lazily materializes a channel. The only
per-RPC dispatch is a fire-and-forget `penalize()` / `recover()` (enqueue only;
handled off the hot path). Reads never dispatch; writes never happen inline. `#channels` / `#retired` /
`#pinned` are facade-owned I/O the transition never touches — that is what makes the
sync hot path race-free against the async FSM.
handled off the hot path). Reads never dispatch; writes never happen inline.

The runtime keeps discovered physical channels in one map across retirement and revival. The registry owns their active/retired state; retirement does not move the connection between stores. Explicitly pinned channels have a separate cache, even when a pin and discovery share a nodeId. Discovery retirement and pin invalidation each close only their own channel; pool shutdown drains and closes both.

### States (discovery lifecycle)

Expand All @@ -47,7 +47,7 @@ Global guard: `endpoints.destroy` from any non-terminal state → `closed` (imme
| idle | discovery.start | discovering | run first round |
| idle | pin | idle | rebuild snapshot |
| idle | close | closed | close-before-start |
| discovering | round_succeeded | ready / degraded | apply round, ready-latch, arm interval+idle_sweep |
| discovering | round_succeeded | ready / degraded | apply round, settle readiness waiters, arm interval+idle_sweep |
| discovering | round_succeeded (0 endpoints) | discovering | rejected as retryable failure — arm backoff |
| discovering | round_failed (retryable) | discovering | arm backoff, stay |
| discovering | round_failed (non-retryable) | closed | emit `failed` (only terminal-failure path) |
Expand Down Expand Up @@ -129,10 +129,7 @@ hard-pin, and the pile-relaxed last-resort tiers still route if every pile is un

A connection dropped from discovery is **not** torn down while it works: live streams
drain on it and a brief flap does not close it. New RPCs are simply not routed there.
The `idle_sweep` effect closes a retired channel only on genuine breakage
(`SHUTDOWN` / sustained `TRANSIENT_FAILURE`) or after `retiredGraceMs` idle with no
reappearance; a returning node revives the **same** channel. Still-discovered channels
are never proactively closed — grpc-js manages their idle socket.
The `idle_sweep` effect keeps retired channels in `READY`, removes `SHUTDOWN` channels immediately, and reaps other connectivity states after `retiredGraceMs` from retirement. A returning node reuses the same channel; another retirement starts a new grace interval. Still-discovered channels are never proactively closed — grpc-js manages their idle socket.

### Direct topic IO

Expand Down
4 changes: 3 additions & 1 deletion packages/core/src/conn.ts
Original file line number Diff line number Diff line change
Expand Up @@ -52,7 +52,9 @@ export class GrpcConnection implements Connection {
...channelOptions,
// Required when the TLS certificate CN doesn't match the gRPC endpoint
// address (common in YDB deployments behind a load balancer).
'grpc.ssl_target_name_override': endpoint.sslTargetNameOverride,
...(endpoint.sslTargetNameOverride && {
'grpc.ssl_target_name_override': endpoint.sslTargetNameOverride,
}),
})
}

Expand Down
117 changes: 117 additions & 0 deletions packages/core/src/driver.discovery.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,117 @@
import { channel } from 'node:diagnostics_channel'
import { setTimeout as delay } from 'node:timers/promises'

import { create } from '@bufbuild/protobuf'
import { anyPack } from '@bufbuild/protobuf/wkt'
import { DiscoveryServiceDefinition, ListEndpointsResultSchema } from '@ydbjs/api/discovery'
import { StatusIds_StatusCode } from '@ydbjs/api/operation'
import { createServer } from 'nice-grpc'
import { expect, test } from 'vitest'

import { Driver } from './driver.ts'

// Controlled discovery can remove a node while its real RPC remains live.
test.for([false, true])(
'rediscovery preserves the active call and reroutes new calls (GOAWAY: %s)',
{ timeout: 10_000 },
async (goaway, tc) => {
let entered = Promise.withResolvers<void>()
let release = Promise.withResolvers<void>()
await using servers = {
first: createServer(),
second: createServer(),
async [Symbol.asyncDispose]() {
release.resolve()
await Promise.all([this.first.shutdown(), this.second.shutdown()])
},
}
let callsOnOriginal = 0
let endpoints: { nodeId: number; address: string; port: number }[] = []
servers.first.add(
{
listEndpoints: DiscoveryServiceDefinition.listEndpoints,
whoAmI: DiscoveryServiceDefinition.whoAmI,
},
{
async listEndpoints() {
return {
operation: {
ready: true,
status: StatusIds_StatusCode.SUCCESS,
result: anyPack(
ListEndpointsResultSchema,
create(ListEndpointsResultSchema, { endpoints })
),
},
}
},
async whoAmI() {
if (++callsOnOriginal === 1) {
entered.resolve()
await release.promise
}
return { operation: { id: 'original' } }
},
}
)
servers.second.add(
{ whoAmI: DiscoveryServiceDefinition.whoAmI },
{
async whoAmI() {
return { operation: { id: 'replacement' } }
},
}
)
let firstPort = await servers.first.listen('127.0.0.1:0')
let secondPort = await servers.second.listen('127.0.0.1:0')
endpoints = [{ nodeId: 1, address: '127.0.0.1', port: firstPort }]
using driver = new Driver(`grpc://127.0.0.1:${firstPort}/local`, {
'ydb.sdk.discovery_interval_ms': 200,
'ydb.sdk.discovery_timeout_ms': 100,
'ydb.sdk.connection_idle_interval_ms': 10,
'ydb.sdk.connection_idle_timeout_ms': 0,
})
await driver.ready(tc.signal)
let retired = channel('ydb:driver.connection.retired')
let onRetired = (event: unknown) => {
let info = event as { driver: unknown; nodeId: bigint }
if (info.driver === driver.identity && info.nodeId === 1n) {
observed.retired = true
}
}
using observed = {
retired: false,
[Symbol.dispose]() {
retired.unsubscribe(onRetired)
},
}
retired.subscribe(onRetired)

let client = driver.createClient(DiscoveryServiceDefinition)
let completed = false
let original = client.whoAmI({}, { signal: tc.signal })
void original.then(
() => {
completed = true
return true
},
() => {
completed = true
return true
}
)
await entered.promise
endpoints = [{ nodeId: 2, address: '127.0.0.1', port: secondPort }]
await expect.poll(() => observed.retired, { timeout: 5000 }).toBe(true)

let replacement = await client.whoAmI({}, { signal: tc.signal })
expect(replacement.operation?.id).toBe('replacement')
let shutdown = goaway ? servers.first.shutdown() : undefined
// Let retirement sweeps run after GOAWAY while the existing RPC is still active.
await delay(100, undefined, { signal: tc.signal })
expect(completed).toBe(false)
release.resolve()
await expect(original).resolves.toMatchObject({ operation: { id: 'original' } })
await shutdown
}
)
31 changes: 31 additions & 0 deletions packages/core/src/driver.pinning.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -217,3 +217,34 @@ test('keeps pin confirmations ordered when an earlier client is disposed', async
expect(response.operation?.id).toBe('discovery')
expect(fixture.calls).toBe(0)
})

test.for([false, true])(
'routes pinned and discovered clients independently (pin first: %s)',
async (pinFirst, tc) => {
await using fixture = await pinFixture()
await fixture.driver.ready(tc.signal)

let discovered = fixture.driver.createClient(DiscoveryServiceDefinition, 1n)
using pinned = fixture.driver.createClient(DiscoveryServiceDefinition, {
...fixture.target,
nodeId: 1n,
})

let clients = pinFirst ? [pinned, discovered] : [discovered, pinned]
let responses = []

for (let client of clients) {
// oxlint-disable-next-line no-await-in-loop
let response = await client.whoAmI({}, { signal: tc.signal })
responses.push(response.operation?.id)
}

expect(responses).toEqual(pinFirst ? ['direct', 'discovery'] : ['discovery', 'direct'])

pinned[Symbol.dispose]()
await setImmediate()

let response = await discovered.whoAmI({}, { signal: tc.signal })
expect(response.operation?.id).toBe('discovery')
}
)
34 changes: 34 additions & 0 deletions packages/core/src/driver.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -372,3 +372,37 @@ test('survives background rediscovery failure and keeps serving requests', async
await server.shutdown()
}
})

test('rejects ready after close with discovery disabled', async (tc) => {
using driver = new Driver('grpc://127.0.0.1:1/local', {
'ydb.sdk.enable_discovery': false,
})
await driver.ready(tc.signal)
driver.close()
await expect(driver.ready(tc.signal)).rejects.toThrow(/closed/i)
})

test('rejects ready after close with discovery enabled', async (tc) => {
await using server = await startBadDiscovery({
status: StatusIds_StatusCode.SUCCESS,
ready: true,
result: anyPack(
ListEndpointsResultSchema,
create(ListEndpointsResultSchema, {
endpoints: [{ nodeId: 1, address: '127.0.0.1', port: 2136 }],
})
),
})
using driver = new Driver(`grpc://127.0.0.1:${server.port}/local`)
await driver.ready(tc.signal)
driver.close()
await expect(driver.ready(tc.signal)).rejects.toThrow(/closed/i)
})

test('honors an already aborted ready signal with discovery disabled', async () => {
using driver = new Driver('grpc://127.0.0.1:1/local', {
'ydb.sdk.enable_discovery': false,
})
let reason = new Error('caller cancelled readiness')
await expect(driver.ready(AbortSignal.abort(reason))).rejects.toBe(reason)
})
Loading
Loading