Press n or j to go to the next uncovered block, b, p or k for the previous block.
| 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299 300 301 302 303 304 305 306 307 308 309 310 311 312 313 314 315 316 317 318 319 320 321 322 323 324 325 326 327 328 329 330 331 332 333 334 335 336 337 338 339 340 341 342 343 344 345 346 347 348 349 350 351 352 353 354 355 356 357 358 359 360 361 362 363 364 365 366 367 368 369 370 371 372 373 374 375 376 377 378 379 380 381 382 383 384 385 386 387 388 389 390 391 392 393 394 395 396 397 398 399 400 401 402 403 404 405 406 407 408 409 410 411 412 413 414 415 416 417 418 419 420 421 422 423 424 425 426 427 428 429 430 431 432 433 434 435 436 437 438 439 440 441 442 443 444 445 446 447 448 449 450 451 452 453 454 455 456 457 458 459 460 461 462 463 464 465 466 467 468 469 470 471 472 473 474 475 476 477 478 479 480 481 482 483 484 485 486 487 488 489 490 491 492 493 494 495 496 497 498 499 500 501 502 503 504 505 506 507 508 509 510 511 512 513 514 515 516 517 518 519 520 521 522 523 524 525 526 527 528 529 530 531 532 533 534 535 536 537 538 539 540 541 542 543 544 545 546 547 548 549 550 551 552 553 554 555 556 557 558 559 560 561 562 563 564 565 566 567 568 569 570 571 572 573 574 575 576 577 578 579 580 581 582 583 584 585 586 587 588 589 590 591 592 593 594 595 596 597 598 599 600 601 602 603 604 605 606 607 608 609 610 611 612 613 614 615 616 617 618 619 620 621 622 623 624 625 626 627 628 629 630 631 632 633 634 635 636 637 638 639 640 641 642 643 644 645 646 647 648 649 650 651 652 653 654 655 656 657 658 659 660 661 662 663 664 665 666 667 668 669 670 671 672 673 674 675 676 677 678 679 680 681 682 683 684 685 686 687 688 689 690 691 692 693 694 695 696 697 698 699 700 701 702 703 704 705 706 707 708 709 710 711 712 713 714 715 716 717 718 719 720 721 722 723 724 725 726 727 728 729 730 731 732 733 734 735 736 737 738 739 740 741 742 743 744 745 746 747 748 749 750 751 752 753 754 755 756 757 758 759 760 761 762 763 764 765 766 767 768 769 770 771 772 773 774 775 776 777 778 779 780 781 782 783 784 785 786 787 788 789 790 791 792 793 794 795 796 797 798 799 800 801 802 803 804 805 806 807 808 809 810 811 812 813 814 815 816 817 818 819 820 821 822 823 824 825 826 827 828 829 830 831 832 833 834 835 836 837 838 839 840 841 842 843 844 845 846 847 848 849 850 851 852 853 854 855 856 857 858 859 860 861 862 863 864 865 866 867 868 869 870 871 872 873 874 875 876 877 878 879 880 881 882 883 884 885 886 887 888 889 890 891 892 893 894 895 896 897 898 899 900 901 902 903 904 905 906 907 908 909 910 911 912 913 914 915 916 917 918 919 | 12x 12x 12x 12x 52331x 12x 1729x 704x 1025x 4x 1021x 6x 1015x 12x 1025x 12x 4505x 12x 1048553x 12x 9599x 919x 919x 12x 2172x 2172x 2172x 1820x 12x 53408x 12x 1054856x 1054856x 12x 7860x 7860x 7860x 7860x 1055611x 1055611x 1055611x 452x 1055159x 303x 1054856x 1053816x 1053816x 1054856x 1054856x 7860x 7860x 7860x 7860x 7860x 12x 1207x 1207x 879x 269x 879x 1207x 1207x 131x 1207x 12x 1207x 1207x 1207x 1207x 1207x 416x 416x 131x 285x 285x 285x 1207x 1207x 1207x 1207x 1207x 12x 6220x 6220x 1052948x 6220x 6220x 6220x 6220x 1053171x 1053171x 1053171x 225x 1052946x 1052946x 1052946x 1052946x 6220x 6220x 12x 6220x 6220x 1x 6219x 6219x 6219x 12x 1108207x 1108207x 52031x 1108207x 1554x 1108207x 1108207x 12x 52302x 532x 532x 52302x 52302x 52302x 12x 52302x 52302x 52302x 52302x 52302x 12x 1040693x 1032833x 7860x 7860x 7860x 268x 7860x 12x 1059x 1059x 1059x 1059x 1059x 12x 1582x 689x 1582x 1582x 1582x 1x 1582x 12x 847x 232x 615x 12x 560x 560x 148x 412x 12x 967x 12x 2267963x 2267963x 51538x 2216425x 52307x 52296x 11x 11x 71194x 68589x 68589x 319x 319x 938x 284x 730x 101x 629x 1x 136x 197x 2085238x 1036664x 1036664x 1036664x 1040158x 6082x 6082x 6081x 6082x 6082x 6082x 976x 976x 220x 209x 769x 100x 669x 160x 5956x 2943x 2943x 873x 873x 121x 895x 895x 2x 89x 40x 89x 264x 769x 1729x 138x 138x 138x 138x 138x 668x 133x 535x 4x 4x 148x 148x 148x 217x 303x 302x 1x 32x 219x 26x 193x 1x | import { create } from '@bufbuild/protobuf'
import { timestampFromDate } from '@bufbuild/protobuf/wkt'
import { StatusIds_StatusCode } from '@ydbjs/api/operation'
import {
type StreamWriteMessage_WriteRequest_MessageData,
StreamWriteMessage_WriteRequest_MessageDataSchema,
} from '@ydbjs/api/topic'
import { loggers } from '@ydbjs/debug'
import { YDBError } from '@ydbjs/error'
import type { TransitionResult, TransitionRuntime } from '@ydbjs/fsm'
import { isRetryableError, isRetryableStreamError } from '@ydbjs/retry'
import { ClientError, Status } from 'nice-grpc'
import type { AckStatus, WriteAck } from './types.js'
// The pure half of the writer: states, context, and a synchronous transition
// with no I/O. Everything here mutates `ctx` in place and returns the next
// state + a list of effects for the runtime to execute — see writer-runtime.ts
// for the I/O side. The buffer is a sliding window (see WriterCtx) and byte
// budgeting lives in the facade, so the transition only counts messages.
//
// The full transition map (table + diagram) lives in packages/topic/ARCHITECTURE.md —
// update it in the same commit when you change this dispatch.
let dbg = loggers.topic.extend('writer')
// Hard service limits (bytes).
export const MAX_BATCH_BYTES = 48n * 1024n * 1024n // one WriteRequest frame stays under 48MiB
export const MAX_PAYLOAD_BYTES = 48n * 1024n * 1024n // single message payload cap
// ── State / context ─────────────────────────────────────────────────────────────
// `closed` = graceful/destroyed terminal; `errored` = fatal terminal. Both are final.
export type WriterState =
| 'idle'
| 'connecting'
| 'ready'
| 'reconnecting'
| 'closing'
| 'closed'
| 'errored'
// A message living in the sliding window before/while it is on the wire.
// In auto mode `seqNo` stays 0n until the message is actually sent (assigned in `pump`);
// in manual mode it is set at enqueue. This is why buffered auto messages never
// need renumbering on reconnect — they simply have no number yet.
export type BufferedMessage = {
// Already compressed with the writer's codec (identity for RAW).
data: Uint8Array
// Original (pre-compression) payload size reported to the server.
uncompressedSize: bigint
seqNo: bigint
createdAt: Date
metadataItems?: Record<string, Uint8Array>
}
export type WriterLimits = {
maxInflightCount: number
maxBatchBytes: bigint
}
// Pure logical context — mutated synchronously inside the transition only.
// The single message array is a sliding window: [garbage | inflight | buffer].
// garbage = [0, inflightStart) (acked, awaiting compaction)
// inflight = [inflightStart, bufferStart) (sent, awaiting ack)
// buffer = [bufferStart, messages.length) (not yet sent)
export type WriterCtx = {
// connection identity
sessionId: string
hasEverConnected: boolean
// seqNo bookkeeping
seqNoMode: 'auto' | 'manual' | null
lastSeqNo: bigint
// reconnect bookkeeping
attempts: number
lastError: unknown
// When set, a SCHEME_ERROR (e.g. the topic does not exist yet) is retried instead
// of being fatal — the writer waits until the topic is created.
retryOnSchemeError: boolean
// Terminal reconnect deadline (ms). Infinity = unbounded (reconnect forever); the
// transition owns whether to arm the `recovery_window` timer based on this.
recoveryWindowMs: number
// flush barrier
flushRequested: boolean
// sliding-window buffer (see the diagram above)
messages: BufferedMessage[]
bufferStart: number
bufferLength: number
inflightStart: number
inflightLength: number
limits: WriterLimits
}
// Timer names shared with the reader, plus the writer-only flush cadence. The
// writer has no partition-scoped timers, so a TimerRef is just the global name —
// the reader's TimerRef adds a partition-keyed variant.
export type GlobalTimerName =
| 'start_timeout'
| 'retry_backoff'
| 'recovery_window'
| 'update_token'
| 'graceful_timeout'
| 'flush_tick'
export type TimerRef = { which: GlobalTimerName }
export type WriterEvent =
// user (dispatched by the facade)
| { type: 'writer.start' }
| { type: 'writer.write'; message: BufferedMessage }
| { type: 'writer.flush' }
| { type: 'writer.close' }
| { type: 'writer.destroy'; reason?: unknown }
// internal self-dispatch — fsm has no `always`/`after`, so the send loop is an explicit event
| { type: 'writer.pump' }
// transport → writer (ingested from the transport FSM output)
| {
type: 'writer.stream.init_response'
sessionId: string
lastSeqNo: bigint
partitionId?: bigint
}
| { type: 'writer.stream.write_response'; acks: WriteAck[] }
| { type: 'writer.stream.token_response' }
| { type: 'writer.stream.disconnected'; error?: unknown }
// timers
| { type: 'writer.timer.start_timeout' }
| { type: 'writer.timer.retry_backoff' }
| { type: 'writer.timer.recovery_window' }
| { type: 'writer.timer.flush_tick' }
| { type: 'writer.timer.update_token' }
| { type: 'writer.timer.graceful_timeout' }
// transport.* = socket lifecycle only; send.* = anything written to the stream.
export type WriterEffect =
| { type: 'writer.effect.transport.connect'; getLastSeqNo: boolean }
| { type: 'writer.effect.transport.close' }
| {
type: 'writer.effect.send.write_request'
messages: StreamWriteMessage_WriteRequest_MessageData[]
}
| { type: 'writer.effect.send.update_token' }
| ({ type: 'writer.effect.timer.schedule' } & TimerRef)
| ({ type: 'writer.effect.timer.clear' } & TimerRef)
| { type: 'writer.effect.finalize'; reason: unknown }
export type WriterOutput =
| { type: 'writer.session'; sessionId: string; lastSeqNo: bigint; nextSeqNo: bigint }
// `freedBytes` = compressed bytes that left the window with this ack batch, so
// the facade can decrement its byte budget without re-tracking message sizes.
| {
type: 'writer.acknowledgments'
acknowledgments: Map<bigint, AckStatus>
freedBytes: bigint
}
| { type: 'writer.flushed'; lastSeqNo: bigint }
| { type: 'writer.reconnecting'; attempt: number; error?: unknown }
| { type: 'writer.error'; error: unknown }
| { type: 'writer.closed'; reason?: unknown }
type WriterRuntime = TransitionRuntime<WriterState, WriterEvent, WriterOutput>
// ── Helpers ─────────────────────────────────────────────────────────────────────
export let createWriterCtx = function createWriterCtx(
limits: WriterLimits,
options?: { retryOnSchemeError?: boolean; recoveryWindowMs?: number }
): WriterCtx {
return {
sessionId: '',
hasEverConnected: false,
seqNoMode: null,
lastSeqNo: 0n,
attempts: 0,
lastError: undefined,
retryOnSchemeError: options?.retryOnSchemeError ?? false,
recoveryWindowMs: options?.recoveryWindowMs ?? Infinity,
flushRequested: false,
messages: [],
bufferStart: 0,
bufferLength: 0,
inflightStart: 0,
inflightLength: 0,
limits,
}
}
// A stream error is retryable when the writer should reconnect transparently.
// Topic writes are idempotent (dedup by producerId+seqNo), so we use the
// idempotent classification — unlike the plain stream classifier, this retries
// the "conditionally" YDB statuses (SESSION_EXPIRED, UNDETERMINED, TIMEOUT).
// A clean stream end with no error object is also retryable (server-side reconnect).
// SCHEME_ERROR is fatal unless `retryOnSchemeError` is set (wait for topic creation).
export let isRetryableWriterError = function isRetryableWriterError(
error: unknown,
retryOnSchemeError = false
): boolean {
if (error === undefined || error === null) {
return true
}
if (isPayloadTooLargeError(error)) {
return false
}
if (
retryOnSchemeError &&
error instanceof YDBError &&
error.code === StatusIds_StatusCode.SCHEME_ERROR
) {
return true
}
return isRetryableStreamError(error) || isRetryableError(error, true)
}
// A size-limit rejection is deterministic — resending the same oversized frame can
// only fail again, so it must be fatal (Go demotes this case explicitly). Every
// size rejection observed against a real server (tests/writer-protocol.test.ts) is
// a gRPC ClientError RESOURCE_EXHAUSTED whose details carry a size complaint:
// server frame cap: 'Received message larger than max (66060326 vs. 64000000)'
// client send cap: 'Attempted to send message with a size larger than 67108864'
// (grpc-js receive paths use the same 'larger than' wording). The code alone is not
// enough — RESOURCE_EXHAUSTED also covers genuine throttling, which SHOULD be
// retried — so the details text narrows it. Everything else the server could send
// (e.g. a YDBError BAD_REQUEST issue) is already non-retryable via the generic
// classifier and needs no special case here.
let isPayloadTooLargeError = function isPayloadTooLargeError(error: unknown): boolean {
return (
error instanceof ClientError &&
error.code === Status.RESOURCE_EXHAUSTED &&
/larger than/i.test(error.details)
)
}
let allDrained = function allDrained(ctx: WriterCtx): boolean {
return ctx.bufferLength === 0 && ctx.inflightLength === 0
}
// The window has work and headroom: something is buffered and inflight has room.
let canSend = function canSend(ctx: WriterCtx): boolean {
return ctx.bufferLength > 0 && ctx.inflightLength < ctx.limits.maxInflightCount
}
// Resolve a pending flush the moment the window is empty. Every path that can
// drain the buffer (a write_response ack, or a reconnect whose init dedups all
// in-flight messages) must call this — otherwise a flush that drains via the
// dedup path never emits writer.flushed and the caller hangs forever.
let resolveFlushIfDrained = function resolveFlushIfDrained(
ctx: WriterCtx,
runtime: WriterRuntime
): void {
if (ctx.flushRequested && allDrained(ctx)) {
ctx.flushRequested = false
runtime.emit({ type: 'writer.flushed', lastSeqNo: ctx.lastSeqNo })
}
}
// Record a flush request. Honored in every live state — a flush issued while the
// writer is still connecting must resolve once messages drain after init, not be
// dropped. Resolves immediately when there is nothing pending.
let requestFlush = function requestFlush(ctx: WriterCtx, runtime: WriterRuntime): void {
ctx.flushRequested = true
resolveFlushIfDrained(ctx, runtime)
// Still pending — kick the send loop to drain it.
if (ctx.flushRequested) {
runtime.dispatch({ type: 'writer.pump' })
}
}
// One stream attempt: open the transport and arm its watchdog. Shared by every
// (re)connect site so the pair can never drift apart.
let connectEffects = function connectEffects(ctx: WriterCtx): WriterEffect[] {
return [
// Re-request last_seq_no only until it was recovered once — a retry that
// races the very first connect must not resume at seqNo 0 and silently
// collide with already-persisted messages.
{ type: 'writer.effect.transport.connect', getLastSeqNo: !ctx.hasEverConnected },
{ type: 'writer.effect.timer.schedule', which: 'start_timeout' },
]
}
// ── Window & batch helpers ──────────────────────────────────────────────────────
// Build the on-wire MessageData for one buffered message. Pure — no I/O.
let toMessageData = function toMessageData(
message: BufferedMessage
): StreamWriteMessage_WriteRequest_MessageData {
let metadataItems = message.metadataItems
? Object.entries(message.metadataItems).map(([key, value]) => ({ key, value }))
: []
return create(StreamWriteMessage_WriteRequest_MessageDataSchema, {
data: message.data,
seqNo: message.seqNo,
createdAt: timestampFromDate(message.createdAt),
metadataItems,
uncompressedSize: message.uncompressedSize,
})
}
// Form the next batch: take from the front of the buffer up to the inflight and
// batch-byte limits, assigning auto seqNos as we go. Mutates the window in place
// (buffer → inflight) and returns the on-wire messages. Synchronous by design.
let formBatch = function formBatch(ctx: WriterCtx): StreamWriteMessage_WriteRequest_MessageData[] {
let batch: StreamWriteMessage_WriteRequest_MessageData[] = []
let batchBytes = 0n
let end = ctx.bufferStart + ctx.bufferLength
for (let i = ctx.bufferStart; i < end; i++) {
let message = ctx.messages[i]!
let size = BigInt(message.data.length)
if (batch.length > 0 && batchBytes + size > ctx.limits.maxBatchBytes) {
break
}
if (ctx.inflightLength + batch.length >= ctx.limits.maxInflightCount) {
break
}
// Auto mode: the seqNo is assigned now, at send time, from the high-water mark.
if (message.seqNo === 0n) {
ctx.lastSeqNo += 1n
message.seqNo = ctx.lastSeqNo
}
batch.push(toMessageData(message))
batchBytes += size
}
let count = batch.length
ctx.bufferStart += count
ctx.bufferLength -= count
ctx.inflightLength += count
return batch
}
// Apply a server init: recover the seqNo high-water mark once (auto numbering),
// then drop any server-persisted in-flight messages and rewind the rest for resend.
//
// The dedup runs on EVERY init, including reconnects: YDB reports last_seq_no even
// when get_last_seq_no is false (proven in tests/writer-protocol.test.ts), so we
// skip resending messages the server already has — like the Java SDK. We only
// request get_last_seq_no on the first connect (like Go) to avoid its cost. If a
// reconnect ever reported 0, dropAckedAndRewind drops nothing and we resend
// everything; the server dedups by producerId+seqNo — correct either way, just
// less efficient. So this is an optimization, not a correctness dependency.
let applyInit = function applyInit(
ctx: WriterCtx,
sessionId: string,
serverLastSeqNo: bigint,
runtime: WriterRuntime
): void {
ctx.sessionId = sessionId
if (!ctx.hasEverConnected) {
// Trust the server's high-water mark exactly once. Manual mode keeps the
// user's numbers; auto mode continues above the recovered value.
if (ctx.seqNoMode !== 'manual' && serverLastSeqNo > ctx.lastSeqNo) {
ctx.lastSeqNo = serverLastSeqNo
}
ctx.hasEverConnected = true
}
let { recovered, freedBytes } = dropAckedAndRewind(ctx, serverLastSeqNo)
if (recovered.size > 0) {
runtime.emit({ type: 'writer.acknowledgments', acknowledgments: recovered, freedBytes })
}
runtime.emit({
type: 'writer.session',
sessionId,
lastSeqNo: ctx.lastSeqNo,
nextSeqNo: ctx.lastSeqNo + 1n,
})
}
// Drop in-flight messages the server already persisted (seqNo <= serverLastSeqNo),
// surfacing them as `skipped` (deduplicated), and move the remaining unacked
// in-flight messages back to the front of the buffer to be resent in order.
// Only scans the in-flight range; buffered (unsent, unnumbered) messages are untouched.
let dropAckedAndRewind = function dropAckedAndRewind(
ctx: WriterCtx,
serverLastSeqNo: bigint
): { recovered: Map<bigint, AckStatus>; freedBytes: bigint } {
let recovered = new Map<bigint, AckStatus>()
let freedBytes = 0n
// In-flight seqNos are strictly increasing (assigned in order in formBatch), so
// `seqNo <= serverLastSeqNo` splits the in-flight range at one boundary — walk
// the acked prefix, exactly like acknowledge() walks the acked prefix.
let inflightEnd = ctx.bufferStart
let i = ctx.inflightStart
while (i < inflightEnd) {
let message = ctx.messages[i]!
if (message.seqNo === 0n || message.seqNo > serverLastSeqNo) {
break
}
freedBytes += BigInt(message.data.length)
recovered.set(message.seqNo, 'skipped')
i += 1
}
ctx.bufferStart = i
ctx.bufferLength = ctx.messages.length - i
ctx.inflightStart = i
ctx.inflightLength = 0
return { recovered, freedBytes }
}
// Move server-acknowledged messages out of the in-flight window into garbage.
// The server acks the in-flight prefix in order, so we walk from the head and
// stop at the first unacked message. We report only the messages actually removed
// from the window, so the emitted acks and the freed-byte total can never drift
// from the window even if a stream ever delivered a non-prefix ack set.
let acknowledge = function acknowledge(
ctx: WriterCtx,
acks: WriteAck[]
): { acknowledgments: Map<bigint, AckStatus>; freedBytes: bigint } {
let status = new Map<bigint, AckStatus>()
for (let ack of acks) {
status.set(ack.seqNo, ack.status)
}
let acknowledgments = new Map<bigint, AckStatus>()
let freedBytes = 0n
let inflightEnd = ctx.bufferStart
while (ctx.inflightStart < inflightEnd) {
let message = ctx.messages[ctx.inflightStart]!
let messageStatus = status.get(message.seqNo)
if (messageStatus === undefined) {
break
}
acknowledgments.set(message.seqNo, messageStatus)
freedBytes += BigInt(message.data.length)
ctx.inflightStart += 1
ctx.inflightLength -= 1
}
compactGarbage(ctx)
return { acknowledgments, freedBytes }
}
// Reclaim the garbage prefix by splicing it out and rebasing the window pointers.
let compactGarbage = function compactGarbage(ctx: WriterCtx): void {
let garbageLength = ctx.inflightStart
if (garbageLength === 0) {
return
}
ctx.messages.splice(0, garbageLength)
ctx.inflightStart = 0
ctx.bufferStart -= garbageLength
}
// Append a message to the buffer. Total by design — seqNo-mode validation
// (which must throw synchronously to the caller) lives in the facade, so the
// transition never throws and can never accidentally destroy the machine.
// A non-zero seqNo means the facade already validated a manual message; a zero
// seqNo is an auto message that gets its number at send time (see formBatch).
let enqueue = function enqueue(ctx: WriterCtx, message: BufferedMessage): void {
let providedSeqNo = message.seqNo !== 0n
if (ctx.seqNoMode === null) {
ctx.seqNoMode = providedSeqNo ? 'manual' : 'auto'
}
// Manual mode: lastSeqNo tracks the user's high-water mark for resend/recovery.
if (providedSeqNo && message.seqNo > ctx.lastSeqNo) {
ctx.lastSeqNo = message.seqNo
}
ctx.messages.push(message)
ctx.bufferLength += 1
}
// ── Terminal / transitions ──────────────────────────────────────────────────────
// Terminal transition into `closed` or `errored`: emit the lifecycle output,
// tear the transport down and finalize. Reused from many states.
let terminate = function terminate(
ctx: WriterCtx,
state: 'closed' | 'errored',
reason: unknown,
runtime: WriterRuntime
): TransitionResult<WriterState, WriterEffect> {
if (state === 'errored') {
ctx.lastError = reason
runtime.emit({ type: 'writer.error', error: reason })
}
runtime.emit({ type: 'writer.closed', reason })
// Drop any still-buffered/in-flight messages so their payloads can be GC'd —
// on a terminal stop they will never be sent or acknowledged.
releaseState(ctx)
return {
state,
// Terminal: the runtime seals itself after the finalize effect runs, so the
// buffered lifecycle outputs (writer.closed / writer.error) are delivered first.
final: { reason },
effects: [
{ type: 'writer.effect.transport.close' },
// No per-timer clears — the finalize handler clears the whole timer map.
{ type: 'writer.effect.finalize', reason },
],
}
}
// Free the message window. Called on terminal stop to release payload memory.
let releaseState = function releaseState(ctx: WriterCtx): void {
ctx.messages = []
ctx.bufferStart = 0
ctx.bufferLength = 0
ctx.inflightStart = 0
ctx.inflightLength = 0
}
// ready + (write/pump/flush_tick): drain buffer → inflight, one batch per event.
let pump = function pump(
ctx: WriterCtx,
runtime: WriterRuntime
): TransitionResult<WriterState, WriterEffect> | void {
if (!canSend(ctx)) {
return
}
let messages = formBatch(ctx)
Iif (messages.length === 0) {
return
}
// More to send and room to send it — keep pumping on the next tick.
if (canSend(ctx)) {
runtime.dispatch({ type: 'writer.pump' })
}
return { effects: [{ type: 'writer.effect.send.write_request', messages }] }
}
// Enter `ready` on a successful init — from `connecting`, or from `reconnecting`
// when a slow init lands after start_timeout already moved us there. Recover the
// seqNo state, resolve any pending flush the recovery just drained, and resume.
let toReady = function toReady(
ctx: WriterCtx,
event: Extract<WriterEvent, { type: 'writer.stream.init_response' }>,
runtime: WriterRuntime
): TransitionResult<WriterState, WriterEffect> {
ctx.attempts = 0
applyInit(ctx, event.sessionId, event.lastSeqNo, runtime)
resolveFlushIfDrained(ctx, runtime)
runtime.dispatch({ type: 'writer.pump' })
return {
state: 'ready',
effects: [
{ type: 'writer.effect.timer.clear', which: 'start_timeout' },
{ type: 'writer.effect.timer.clear', which: 'retry_backoff' },
{ type: 'writer.effect.timer.clear', which: 'recovery_window' },
{ type: 'writer.effect.timer.schedule', which: 'flush_tick' },
{ type: 'writer.effect.timer.schedule', which: 'update_token' },
],
}
}
let toReconnecting = function toReconnecting(
ctx: WriterCtx,
error: unknown,
runtime: WriterRuntime
): TransitionResult<WriterState, WriterEffect> {
if (error !== undefined) {
ctx.lastError = error
}
runtime.emit({
type: 'writer.reconnecting',
attempt: ctx.attempts,
...(error !== undefined && { error }),
})
// No transport.close here — the transport already closed its own stream on
// disconnect, and the reconnect happens via transport.connect (which reopens).
let effects: WriterEffect[] = [
{ type: 'writer.effect.timer.clear', which: 'start_timeout' },
{ type: 'writer.effect.timer.clear', which: 'flush_tick' },
{ type: 'writer.effect.timer.clear', which: 'update_token' },
{ type: 'writer.effect.timer.schedule', which: 'retry_backoff' },
]
// Arm the terminal deadline only when recovery is bounded. Unbounded (Infinity)
// means reconnect forever — the transition owns that policy so the emitted effects
// reflect it (model-testable), instead of the runtime silently dropping the timer.
if (Number.isFinite(ctx.recoveryWindowMs)) {
effects.push({ type: 'writer.effect.timer.schedule', which: 'recovery_window' })
}
return { state: 'reconnecting', effects }
}
// Closing drain gate: finalize once the window is empty, otherwise keep pumping.
let closeWhenDrained = function closeWhenDrained(
ctx: WriterCtx,
runtime: WriterRuntime
): TransitionResult<WriterState, WriterEffect> | void {
if (allDrained(ctx)) {
return terminate(ctx, 'closed', new Error('Writer closed'), runtime)
}
runtime.dispatch({ type: 'writer.pump' })
}
// Enter graceful shutdown. If nothing is pending, finalize now; otherwise drain
// the buffer (over the live stream, or over a reconnect if one is already
// scheduled) bounded by the graceful-shutdown timeout. Retry/recovery timers are
// intentionally preserved so a close issued while reconnecting still flushes.
let toClosing = function toClosing(
ctx: WriterCtx,
runtime: WriterRuntime
): TransitionResult<WriterState, WriterEffect> {
let drained = closeWhenDrained(ctx, runtime)
if (drained) {
return drained
}
return {
state: 'closing',
effects: [
// Closing may be entered while reconnecting — cancel a stale recovery_window
// so it cannot cut the graceful drain short (close is bounded by
// graceful_timeout, like the reader). start_timeout / retry_backoff are left
// armed: the closing state uses them to keep reconnecting to finish the drain.
{ type: 'writer.effect.timer.clear', which: 'recovery_window' },
{ type: 'writer.effect.timer.clear', which: 'flush_tick' },
{ type: 'writer.effect.timer.clear', which: 'update_token' },
{ type: 'writer.effect.timer.schedule', which: 'graceful_timeout' },
],
}
}
// ── Transition ──────────────────────────────────────────────────────────────────
// Deliberately-ignored (state, event) pairs route through here so an unhandled
// event shows up in debug logs instead of vanishing.
let ignored = function ignored(state: WriterState, event: WriterEvent): void {
dbg.log('ignoring %s in state %s', event.type, state)
}
export let writerTransition = function writerTransition(
ctx: WriterCtx,
event: WriterEvent,
runtime: WriterRuntime
): TransitionResult<WriterState, WriterEffect> | void {
let state = runtime.state
// Global: hard destroy from any non-terminal state.
if (state !== 'closed' && state !== 'errored' && event.type === 'writer.destroy') {
return terminate(ctx, 'closed', event.reason ?? new Error('Writer destroyed'), runtime)
}
switch (state) {
case 'idle': {
switch (event.type) {
case 'writer.start':
return { state: 'connecting', effects: connectEffects(ctx) }
case 'writer.write':
enqueue(ctx, event.message)
return
case 'writer.close':
return terminate(
ctx,
'closed',
new Error('Writer closed before start'),
runtime
)
default:
return ignored(state, event)
}
}
case 'connecting': {
switch (event.type) {
case 'writer.write':
enqueue(ctx, event.message)
return
case 'writer.flush':
requestFlush(ctx, runtime)
return
case 'writer.stream.init_response':
return toReady(ctx, event, runtime)
case 'writer.timer.start_timeout':
return toReconnecting(ctx, undefined, runtime)
case 'writer.stream.disconnected':
if (!isRetryableWriterError(event.error, ctx.retryOnSchemeError)) {
return terminate(ctx, 'errored', event.error, runtime)
}
return toReconnecting(ctx, event.error, runtime)
// The recovery window is armed while reconnecting and can elapse during a
// connect attempt — without this the terminal bound would never fire.
case 'writer.timer.recovery_window':
return terminate(
ctx,
'errored',
ctx.lastError ?? new Error('Writer recovery window expired'),
runtime
)
case 'writer.close':
return toClosing(ctx, runtime)
default:
return ignored(state, event)
}
}
case 'ready': {
switch (event.type) {
case 'writer.write':
enqueue(ctx, event.message)
runtime.dispatch({ type: 'writer.pump' })
return
case 'writer.pump':
case 'writer.timer.flush_tick':
return pump(ctx, runtime)
case 'writer.stream.write_response': {
let { acknowledgments, freedBytes } = acknowledge(ctx, event.acks)
if (acknowledgments.size > 0) {
runtime.emit({
type: 'writer.acknowledgments',
acknowledgments,
freedBytes,
})
}
resolveFlushIfDrained(ctx, runtime)
runtime.dispatch({ type: 'writer.pump' })
return
}
case 'writer.flush':
requestFlush(ctx, runtime)
return
case 'writer.timer.update_token':
return { effects: [{ type: 'writer.effect.send.update_token' }] }
case 'writer.stream.token_response':
return
case 'writer.stream.disconnected':
if (!isRetryableWriterError(event.error, ctx.retryOnSchemeError)) {
return terminate(ctx, 'errored', event.error, runtime)
}
return toReconnecting(ctx, event.error, runtime)
case 'writer.close':
return toClosing(ctx, runtime)
default:
return ignored(state, event)
}
}
case 'reconnecting': {
switch (event.type) {
case 'writer.write':
enqueue(ctx, event.message)
return
case 'writer.flush':
requestFlush(ctx, runtime)
return
// A connect attempt whose init lands here (start_timeout fired just before
// the init was dequeued, so the stream is still open) is a live session —
// honor it rather than dropping it and forcing a wasted reconnect.
case 'writer.stream.init_response':
return toReady(ctx, event, runtime)
case 'writer.timer.retry_backoff':
ctx.attempts += 1
return { state: 'connecting', effects: connectEffects(ctx) }
case 'writer.timer.recovery_window':
return terminate(
ctx,
'errored',
ctx.lastError ?? new Error('Writer recovery window expired'),
runtime
)
case 'writer.stream.disconnected':
// Already backing off — record the reason but stay put.
if (event.error !== undefined) {
ctx.lastError = event.error
}
return
case 'writer.close':
return toClosing(ctx, runtime)
default:
return ignored(state, event)
}
}
case 'closing': {
switch (event.type) {
case 'writer.stream.write_response': {
let { acknowledgments, freedBytes } = acknowledge(ctx, event.acks)
Eif (acknowledgments.size > 0) {
runtime.emit({
type: 'writer.acknowledgments',
acknowledgments,
freedBytes,
})
}
resolveFlushIfDrained(ctx, runtime)
return closeWhenDrained(ctx, runtime)
}
case 'writer.pump':
case 'writer.timer.flush_tick':
// Never assign auto seqNos from an unrecovered high-water mark: a close()
// racing the first init would number from 0n and the server would dedup
// the whole batch as already-written — silent data loss on a clean close.
// The init_response handler below resumes the drain once the mark is known.
if (!ctx.hasEverConnected) {
return
}
return pump(ctx, runtime)
case 'writer.flush':
requestFlush(ctx, runtime)
return
// A reconnect completed mid-close — recover and keep draining.
case 'writer.stream.init_response':
applyInit(ctx, event.sessionId, event.lastSeqNo, runtime)
resolveFlushIfDrained(ctx, runtime)
return (
closeWhenDrained(ctx, runtime) ?? {
effects: [
{ type: 'writer.effect.timer.clear', which: 'start_timeout' },
],
}
)
case 'writer.timer.retry_backoff':
return { effects: connectEffects(ctx) }
case 'writer.timer.graceful_timeout':
// Forced shutdown with messages still pending is a failure to flush —
// surface it (as `errored`) so close() rejects instead of silently
// dropping undelivered writes (critical for tx commit integrity).
if (!allDrained(ctx)) {
return terminate(
ctx,
'errored',
new Error('Graceful shutdown timed out with undelivered messages'),
runtime
)
}
return closeWhenDrained(ctx, runtime)
case 'writer.timer.start_timeout':
// Retry the drain over a fresh stream (bounded by graceful_timeout).
return {
effects: [
{ type: 'writer.effect.timer.clear', which: 'start_timeout' },
{ type: 'writer.effect.timer.schedule', which: 'retry_backoff' },
],
}
case 'writer.stream.disconnected':
// Retry the drain over a fresh stream (bounded by graceful_timeout);
// give up terminally on a fatal error.
if (!isRetryableWriterError(event.error, ctx.retryOnSchemeError)) {
return terminate(ctx, 'errored', event.error, runtime)
}
return {
effects: [
{ type: 'writer.effect.timer.clear', which: 'start_timeout' },
{ type: 'writer.effect.timer.schedule', which: 'retry_backoff' },
],
}
// New writes are rejected once closing (facade throws before dispatch),
// so ignore anything else.
default:
return ignored(state, event)
}
}
case 'closed':
case 'errored':
return ignored(state, event)
}
}
|