All files / packages/core/src/endpoints endpoints-state.ts

96.74% Statements 208/215
92% Branches 115/125
100% Functions 15/15
98.01% Lines 198/202

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                                                                                                                                                                                                                                                    57x   57x         5656x                                                                                                                                                                                                                                                                                     57x         2103x                                     57x 4482x 4482x 4482x 19837x 7450x   4482x 4482x     57x 4482x     57x 6055x     57x                                     57x       8621x 8621x 8621x     20387x 20386x 20386x 20386x 15082x                 8621x 8621x 12618x 3359x                 8621x               57x       42x 42x 42x           42x                 57x             3292x 3292x     3292x 9783x 9783x 2103x 2103x           7680x 7680x 2x       7680x 7680x 7680x 7680x 7680x 7680x 7680x 7680x 7680x 7680x   7680x     2395x 5285x   522x 522x                     3292x 3352x 3352x       3292x 4498x               3292x 3352x                   3292x 3292x 3292x 3292x 3292x   3292x 3292x     3352x                   3292x     57x   153x 3x               150x 150x                       57x           5403x 6x     5403x 5403x 376x               376x   5403x 56x         5403x 5403x   5403x 5403x 5403x   5403x                               57x         767x 767x 767x                             767x     767x 148x               57x         286x 286x 286x 286x     286x       57x         21334x     21334x 5311x     16023x   5656x   5656x 5656x                             464x   445x 440x 440x 440x             440x       440x 440x 440x     9x 9x 9x 9x           9x 6x   3x             2x 2x   2x   2x   3x   1x         8640x         2889x 2852x             2852x       2852x         155x 155x 155x 155x           155x               656x 556x 556x             1770x 1770x 1080x 1080x             1080x   1080x 1080x 102x 102x                     978x     1445x 1445x 110x 110x             110x 110x     5x   520x 520x 520x 520x 520x               520x         765x   284x   150x   1x       1181x   198x 198x 198x 198x 198x               198x         198x 82x 82x   116x     1x     376x 376x   606x       82x                
// The pure half of the endpoints engine: states, context, and a synchronous
// transition with no I/O. All timers, gRPC channels, discovery calls, and the
// clock live in the effectful half (`endpoints-runtime.ts`).
//
// One machine owns the DISCOVERY lifecycle (idle → discovering → ready →
// degraded → closing → closed). Per-endpoint HEALTH is not a machine state —
// it is a sub-state on a registry entry (like the topic reader's PartitionEntry),
// so a 10k-node cluster stays a single machine with two maps, not 10k machines.
//
// The transition rebuilds an immutable RoutingSnapshot (RCU) and emits it as an
// output whenever the routable set changes. The synchronous facade reads that
// snapshot per-RPC without dispatching — see `endpoints-runtime.ts`.
//
// The full transition table lives in packages/core/ARCHITECTURE.md and must be
// updated in the same commit as this dispatch.
 
import type { TransitionResult, TransitionRuntime } from '@ydbjs/fsm'
 
import { EMPTY_SNAPSHOT, buildSnapshot } from './snapshot.js'
import type { RoutingSnapshot } from './snapshot.js'
 
// ── Bridge / multi-pile (2-DC) ──────────────────────────────────────────────
// The generated @ydbjs/api discovery proto carries `pile_states` /
// `bridge_pile_name` (see @ydbjs/api/bridge). The runtime maps proto → DTO; a
// non-bridge cluster returns empty `pile_states`, which `buildSnapshot` treats
// as identity (no pile filter). Bridge clusters filter to usable pile statuses.
export type PileStatus =
	| 'PRIMARY'
	| 'PROMOTED'
	| 'SYNCHRONIZED'
	| 'NOT_SYNCHRONIZED'
	| 'SUSPENDED'
	| 'DISCONNECTED'
	| 'UNSPECIFIED'
 
export type PileState = {
	pileName: string
	status: PileStatus
}
 
// ── Endpoint model (full server-field coverage) ─────────────────────────────
// The domain view the runtime hands the FSM — a superset of the current proto
// EndpointInfo plus the forward-looking bridge field. `loadFactor` is carried
// for observability only; it is never used for routing (the server hardcodes
// it to 0.0 today).
export type DiscoveredEndpoint = {
	nodeId: bigint
	// Raw host (proto EndpointInfo.address), used for dialing.
	host: string
	port: number
	location: string
	loadFactor: number
	sslTargetNameOverride: string
	ipV4: string[]
	ipV6: string[]
	bridgePileName: string
	services: string[]
}
 
export type EndpointSubState = 'active' | 'pessimized' | 'retired' | 'pinned'
 
// Registry record. Mutated synchronously inside the transition only. No clock
// fields (timestamps are runtime concerns) so the transition stays pure and
// model-testable.
export type EndpointEntry = {
	nodeId: bigint
	// Raw host for dialing (proto EndpointInfo.address).
	host: string
	port: number
	// 'host:port' — used for routing keys, diagnostics, and logs (parity with conn.ts).
	address: string
	location: string
	loadFactor: number
	sslTargetNameOverride: string
	ipV4: string[]
	ipV6: string[]
	bridgePileName: string
	services: string[]
	subState: EndpointSubState
	// Bumped every time a pinned entry is (re)pinned; 0 for discovered entries.
	generation: number
}
 
// Lightweight endpoint identity used in outputs (parity with hooks.EndpointInfo).
export type EndpointInfoLite = {
	nodeId: bigint
	address: string
	location: string
	// Bridge pile name; '' on a non-bridge cluster.
	pile: string
}
 
export type EndpointsConfig = {
	// Locality is OPT-IN and soft-only: it only reorders prefer/fallback tiers,
	// never hard-pins to the local DC.
	localityEnabled: boolean
	// Bridge/2DC: prefer PRIMARY/PROMOTED-pile endpoints, fall back to
	// SYNCHRONIZED. Opt-in, soft, and a no-op on a non-bridge cluster. Dominates
	// locality in bridge mode (see snapshot.ts / ARCHITECTURE.md).
	preferPrimaryPile: boolean
	// Force rediscovery when the pessimized fraction exceeds this (0..1).
	degradedThreshold: number
}
 
// Pure logical context — flags, counters, ids, and the two registries. No I/O
// handles, no timers, no clock.
export type EndpointsCtx = {
	byNodeId: Map<bigint, EndpointEntry>
	pinned: Map<bigint, EndpointEntry>
 
	attempts: number
	lastError: unknown
	roundInFlight: boolean
	hasEverDiscovered: boolean
 
	selfLocation: string
	// Empty ⇒ the cluster is not in bridge mode ⇒ the pile filter is identity.
	pileStates: PileState[]
 
	config: EndpointsConfig
}
 
export const DEFAULT_DEGRADED_THRESHOLD = 0.5
 
export let createEndpointsCtx = function createEndpointsCtx(config?: {
	localityEnabled?: boolean | undefined
	preferPrimaryPile?: boolean | undefined
	degradedThreshold?: number | undefined
}): EndpointsCtx {
	return {
		byNodeId: new Map(),
		pinned: new Map(),
		attempts: 0,
		lastError: undefined,
		roundInFlight: false,
		hasEverDiscovered: false,
		selfLocation: '',
		pileStates: [],
		config: {
			localityEnabled: config?.localityEnabled ?? false,
			preferPrimaryPile: config?.preferPrimaryPile ?? false,
			degradedThreshold: config?.degradedThreshold ?? DEFAULT_DEGRADED_THRESHOLD,
		},
	}
}
 
// ── States ──────────────────────────────────────────────────────────────────
export type EndpointsState = 'idle' | 'discovering' | 'ready' | 'degraded' | 'closing' | 'closed'
 
// ── Timers ──────────────────────────────────────────────────────────────────
// No per-endpoint pessimization timer: recovery is optimistic-un-ban-on-rpc-ok
// plus blanket-un-ban-on-discovery. `idle_sweep` reaps genuinely-departed
// retired channels; brief flaps are absorbed (see endpoints-runtime.ts).
export type TimerName = 'discovery_interval' | 'discovery_backoff' | 'idle_sweep' | 'close_deadline'
export type TimerRef = { which: TimerName }
 
// ── Events ──────────────────────────────────────────────────────────────────
export type EndpointsEvent =
	| { type: 'endpoints.discovery.start' }
	| { type: 'endpoints.discovery.force' }
	| {
			type: 'endpoints.discovery.round_succeeded'
			endpoints: DiscoveredEndpoint[]
			selfLocation: string
			pileStates: PileState[]
	  }
	| { type: 'endpoints.discovery.round_failed'; error: unknown; retryable: boolean }
	// Per-RPC outcomes (the only per-RPC dispatch; cheap enqueue).
	| { type: 'endpoints.rpc_failed'; nodeId: bigint }
	| { type: 'endpoints.rpc_ok'; nodeId: bigint }
	// Direct-IO pins (server-named node_ids possibly outside ListEndpoints).
	| {
			type: 'endpoints.pin'
			nodeId: bigint
			host: string
			port: number
			location: string
			sslTargetNameOverride: string
			ipV4: string[]
			ipV6: string[]
			generation: number
	  }
	| { type: 'endpoints.invalidate'; nodeId: bigint }
	// The runtime observed a channel become closeable (drained / broken / past
	// grace, or fully drained during close) and asks the FSM to drop it. During
	// closing, dropping the last channel finalizes.
	| { type: 'endpoints.channel_closeable'; nodeId: bigint }
	| { type: 'endpoints.timer.discovery_interval' }
	| { type: 'endpoints.timer.discovery_backoff' }
	| { type: 'endpoints.timer.idle_sweep' }
	| { type: 'endpoints.timer.close_deadline' }
	| { type: 'endpoints.close' }
	| { type: 'endpoints.destroy'; reason?: unknown }
 
// ── Effects ─────────────────────────────────────────────────────────────────
export type EndpointsEffect =
	| { type: 'endpoints.effect.run_discovery_round' }
	| ({ type: 'endpoints.effect.timer.schedule' } & TimerRef)
	| ({ type: 'endpoints.effect.timer.clear' } & TimerRef)
	// Retire-to-drain: keep the channel open, move it to the drain-watch set.
	| { type: 'endpoints.effect.retire_channel'; nodeId: bigint }
	// Physically close and drop the channel. `store` scopes which materialized
	// channel is dropped: 'pinned' closes only the pin (an invalidate must not
	// tear down a discovered channel that shares the same nodeId); 'any' (default)
	// closes whichever store holds it.
	| { type: 'endpoints.effect.close_channel'; nodeId: bigint; store?: 'any' | 'pinned' }
	// Begin the graceful close drain: close idle channels now and wait for
	// in-flight streams to finish (bounded by the close deadline). Runtime-only.
	| { type: 'endpoints.effect.begin_close_drain' }
	// Scan retired channels; dispatch endpoints.channel_closeable for any that
	// are genuinely gone (broken or idle past grace). Pure I/O in the runtime.
	| { type: 'endpoints.effect.idle_sweep' }
	| { type: 'endpoints.effect.finalize'; reason: unknown }
 
// ── Outputs ─────────────────────────────────────────────────────────────────
export type EndpointsOutput =
	| { type: 'endpoints.snapshot'; snapshot: RoutingSnapshot }
	| { type: 'endpoints.ready' }
	| {
			type: 'endpoints.discovery_completed'
			added: EndpointInfoLite[]
			removed: EndpointInfoLite[]
			total: number
			selfLocation: string
	  }
	| { type: 'endpoints.discovery_failed'; error: unknown; attempt: number; retryable: boolean }
	| { type: 'endpoints.added'; nodeId: bigint; address: string; location: string; pile: string }
	| {
			type: 'endpoints.pessimized'
			nodeId: bigint
			address: string
			location: string
			pile: string
	  }
	| {
			type: 'endpoints.unpessimized'
			nodeId: bigint
			address: string
			location: string
			pile: string
	  }
	| {
			type: 'endpoints.retired'
			nodeId: bigint
			address: string
			location: string
			pile: string
			reason: 'stale_active' | 'stale_pessimized'
	  }
	| {
			type: 'endpoints.removed'
			nodeId: bigint
			address: string
			location: string
			pile: string
			reason: 'idle' | 'pool_close'
	  }
	| { type: 'endpoints.failed'; error: unknown }
	| { type: 'endpoints.closed'; reason?: unknown }
 
// The transition-layer runtime handle (state + emit/dispatch). Named distinctly
// from the `EndpointsRuntime` facade in endpoints-runtime.ts to avoid a shadow;
// internal to this module.
type EndpointsTransitionRuntime = TransitionRuntime<EndpointsState, EndpointsEvent, EndpointsOutput>
type Result = TransitionResult<EndpointsState, EndpointsEffect>
 
// ── Pure helpers ────────────────────────────────────────────────────────────
 
let entryFrom = function entryFrom(
	ep: DiscoveredEndpoint,
	subState: EndpointSubState,
	generation: number
): EndpointEntry {
	return {
		nodeId: ep.nodeId,
		host: ep.host,
		port: ep.port,
		address: `${ep.host}:${ep.port}`,
		location: ep.location,
		loadFactor: ep.loadFactor,
		sslTargetNameOverride: ep.sslTargetNameOverride,
		ipV4: ep.ipV4,
		ipV6: ep.ipV6,
		bridgePileName: ep.bridgePileName,
		services: ep.services,
		subState,
		generation,
	}
}
 
// Fraction of routable (active + pessimized) nodes that are pessimized. Retired
// and pinned entries are excluded from the denominator.
let pessimizedRatio = function pessimizedRatio(ctx: EndpointsCtx): number {
	let active = 0
	let pessimized = 0
	for (let entry of ctx.byNodeId.values()) {
		if (entry.subState === 'active') active++
		else if (entry.subState === 'pessimized') pessimized++
	}
	let total = active + pessimized
	return total === 0 ? 0 : pessimized / total
}
 
let healthState = function healthState(ctx: EndpointsCtx): EndpointsState {
	return pessimizedRatio(ctx) > ctx.config.degradedThreshold ? 'degraded' : 'ready'
}
 
let rebuild = function rebuild(ctx: EndpointsCtx, runtime: EndpointsTransitionRuntime): void {
	runtime.emit({ type: 'endpoints.snapshot', snapshot: buildSnapshot(ctx) })
}
 
let ignored = function ignored(): void {
	// Unhandled (state, event) pair — no-op. Kept explicit for the table.
}
 
export type RetiredInfo = EndpointInfoLite & { reason: 'stale_active' | 'stale_pessimized' }
 
export type RoundDiff = {
	// New OR revived (a retired node reappearing) — both reported as `added` so
	// the add/remove event streams stay balanced.
	added: EndpointInfoLite[]
	// Vanished from discovery while still active/pessimized — reported as `retired`.
	retired: RetiredInfo[]
}
 
// The single source of the round-diff rule. Shared by `applyRound` (which emits
// the FSM's per-node outputs) and the runtime's `publishRoundDiagnostics` (which
// publishes the same events inside the discovery span; it runs before applyRound,
// so it cannot reuse applyRound's result). Pure — computed against the PRE-round
// registry, mutates nothing.
export let computeRoundDiff = function computeRoundDiff(
	byNodeId: ReadonlyMap<bigint, EndpointEntry>,
	endpoints: DiscoveredEndpoint[]
): RoundDiff {
	let discovered = new Set<bigint>()
	let added: EndpointInfoLite[] = []
	for (let ep of endpoints) {
		// A degenerate response may repeat a nodeId — classify the first
		// occurrence only, so `added` never double-counts a node.
		if (discovered.has(ep.nodeId)) continue
		discovered.add(ep.nodeId)
		let existing = byNodeId.get(ep.nodeId)
		if (existing === undefined || existing.subState === 'retired') {
			added.push({
				nodeId: ep.nodeId,
				address: `${ep.host}:${ep.port}`,
				location: ep.location,
				pile: ep.bridgePileName,
			})
		}
	}
 
	let retired: RetiredInfo[] = []
	for (let entry of byNodeId.values()) {
		if (discovered.has(entry.nodeId) || entry.subState === 'retired') continue
		retired.push({
			nodeId: entry.nodeId,
			address: entry.address,
			location: entry.location,
			pile: entry.bridgePileName,
			reason: entry.subState === 'pessimized' ? 'stale_pessimized' : 'stale_active',
		})
	}
 
	return { added, retired }
}
 
// An empty endpoint list is never a usable cluster view — applying it would wipe
// routing (initially: ready() resolves with nothing routable; in steady state:
// every node retires and balanced acquire() throws until the next interval).
// Reject it as a retryable round failure in EVERY state: keep the last snapshot
// and registry untouched, arm the backoff.
let rejectEmptyRound = function rejectEmptyRound(
	ctx: EndpointsCtx,
	runtime: EndpointsTransitionRuntime
): Result {
	ctx.attempts += 1
	ctx.roundInFlight = false
	runtime.emit({
		type: 'endpoints.discovery_failed',
		error: new Error('discovery returned no endpoints'),
		attempt: ctx.attempts,
		retryable: true,
	})
	return {
		effects: [{ type: 'endpoints.effect.timer.schedule', which: 'discovery_backoff' }],
	}
}
 
// Apply a fresh discovery result to the registry. Returns the effects to run
// (retire_channel for newly-stale nodes) and emits per-node outputs. Subsumes
// pool.sync() with the retire-reappear fix. The add/retire classification comes
// from `computeRoundDiff`; the loop below only mutates the registry.
let applyRound = function applyRound(
	ctx: EndpointsCtx,
	endpoints: DiscoveredEndpoint[],
	selfLocation: string,
	pileStates: PileState[],
	runtime: EndpointsTransitionRuntime
): EndpointsEffect[] {
	let { added, retired } = computeRoundDiff(ctx.byNodeId, endpoints)
	let effects: EndpointsEffect[] = []
 
	// Mutate: add new entries, revive retired, un-ban pessimized, refresh dial info.
	for (let ep of endpoints) {
		let existing = ctx.byNodeId.get(ep.nodeId)
		if (existing === undefined) {
			ctx.byNodeId.set(ep.nodeId, entryFrom(ep, 'active', 0))
			continue
		}
 
		// A same-nodeId re-registration at a different host:port means the old
		// channel dials a dead address — drop it so the next acquire re-dials.
		// (A brief flap keeps the same address and is absorbed by retire-drain.)
		let newAddress = `${ep.host}:${ep.port}`
		if (existing.address !== newAddress) {
			effects.push({ type: 'endpoints.effect.close_channel', nodeId: ep.nodeId })
		}
 
		// Refresh surface fields (location/pile/load/dial info can change).
		existing.host = ep.host
		existing.port = ep.port
		existing.address = newAddress
		existing.location = ep.location
		existing.loadFactor = ep.loadFactor
		existing.sslTargetNameOverride = ep.sslTargetNameOverride
		existing.ipV4 = ep.ipV4
		existing.ipV6 = ep.ipV6
		existing.bridgePileName = ep.bridgePileName
		existing.services = ep.services
 
		if (existing.subState === 'retired') {
			// Revive in place — keep the draining channel (never close+recreate,
			// that would kill live streams). Reported as `added` by the diff.
			existing.subState = 'active'
		} else if (existing.subState === 'pessimized') {
			// Blanket un-ban on discovery (authoritative recovery). Not an `added`.
			existing.subState = 'active'
			runtime.emit({
				type: 'endpoints.unpessimized',
				nodeId: existing.nodeId,
				address: existing.address,
				location: existing.location,
				pile: existing.bridgePileName,
			})
		}
	}
 
	// Retire the vanished endpoints from the diff. Channel stays open to drain.
	for (let r of retired) {
		ctx.byNodeId.get(r.nodeId)!.subState = 'retired'
		effects.push({ type: 'endpoints.effect.retire_channel', nodeId: r.nodeId })
	}
 
	// Emit the add/retire outputs from the single-sourced diff.
	for (let a of added) {
		runtime.emit({
			type: 'endpoints.added',
			nodeId: a.nodeId,
			address: a.address,
			location: a.location,
			pile: a.pile,
		})
	}
	for (let r of retired) {
		runtime.emit({
			type: 'endpoints.retired',
			nodeId: r.nodeId,
			address: r.address,
			location: r.location,
			pile: r.pile,
			reason: r.reason,
		})
	}
 
	ctx.selfLocation = selfLocation
	ctx.pileStates = pileStates
	ctx.attempts = 0
	ctx.lastError = undefined
	ctx.roundInFlight = false
 
	rebuild(ctx, runtime)
	runtime.emit({
		type: 'endpoints.discovery_completed',
		added,
		removed: retired.map((r) => ({
			nodeId: r.nodeId,
			address: r.address,
			location: r.location,
			pile: r.pile,
		})),
		total: endpoints.length,
		selfLocation,
	})
 
	return effects
}
 
let toClosing = function toClosing(ctx: EndpointsCtx, runtime: EndpointsTransitionRuntime): Result {
	// Nothing registered → finalize immediately (no channels can exist).
	if (ctx.byNodeId.size === 0 && ctx.pinned.size === 0) {
		return terminate(ctx, new Error('Endpoints closed'), runtime)
	}
 
	// Freeze routing so no new RPC is dialed while draining, then hand the drain
	// to the runtime: `begin_close_drain` closes idle channels immediately and
	// dispatches `drained` once in-flight streams finish; `close_deadline` is the
	// hard cap that force-closes whatever is left. This avoids waiting the full
	// deadline when there is nothing (or nothing busy) to drain.
	runtime.emit({ type: 'endpoints.snapshot', snapshot: EMPTY_SNAPSHOT })
	return {
		state: 'closing',
		effects: [
			{ type: 'endpoints.effect.timer.clear', which: 'discovery_interval' },
			{ type: 'endpoints.effect.timer.clear', which: 'discovery_backoff' },
			{ type: 'endpoints.effect.timer.clear', which: 'idle_sweep' },
			{ type: 'endpoints.effect.begin_close_drain' },
			{ type: 'endpoints.effect.timer.schedule', which: 'close_deadline' },
		],
	}
}
 
let terminate = function terminate(
	ctx: EndpointsCtx,
	reason: unknown,
	runtime: EndpointsTransitionRuntime,
	failure?: unknown
): Result {
	if (failure !== undefined) {
		runtime.emit({ type: 'endpoints.failed', error: failure })
	}
 
	let effects: EndpointsEffect[] = []
	for (let entry of ctx.byNodeId.values()) {
		runtime.emit({
			type: 'endpoints.removed',
			nodeId: entry.nodeId,
			address: entry.address,
			location: entry.location,
			pile: entry.bridgePileName,
			reason: 'pool_close',
		})
		effects.push({ type: 'endpoints.effect.close_channel', nodeId: entry.nodeId })
	}
	for (let entry of ctx.pinned.values()) {
		effects.push({ type: 'endpoints.effect.close_channel', nodeId: entry.nodeId })
	}
 
	// Empty the read plane: a post-close acquire() then selects nothing and throws
	// instead of vending a fresh channel no one will ever close.
	runtime.emit({ type: 'endpoints.snapshot', snapshot: EMPTY_SNAPSHOT })
	runtime.emit({ type: 'endpoints.closed', reason })
 
	ctx.byNodeId.clear()
	ctx.pinned.clear()
	ctx.roundInFlight = false
 
	return {
		state: 'closed',
		final: { reason },
		effects: [
			...effects,
			{ type: 'endpoints.effect.timer.clear', which: 'discovery_interval' },
			{ type: 'endpoints.effect.timer.clear', which: 'discovery_backoff' },
			{ type: 'endpoints.effect.timer.clear', which: 'idle_sweep' },
			{ type: 'endpoints.effect.timer.clear', which: 'close_deadline' },
			{ type: 'endpoints.effect.finalize', reason },
		],
	}
}
 
// Handle a pin/invalidate uniformly across live states. Returns a rebuild result
// when the pinned set changed, else void.
let applyPin = function applyPin(
	ctx: EndpointsCtx,
	event: Extract<EndpointsEvent, { type: 'endpoints.pin' }>,
	runtime: EndpointsTransitionRuntime
): Result | void {
	let address = `${event.host}:${event.port}`
	let prev = ctx.pinned.get(event.nodeId)
	ctx.pinned.set(event.nodeId, {
		nodeId: event.nodeId,
		host: event.host,
		port: event.port,
		address,
		location: event.location,
		loadFactor: 0,
		sslTargetNameOverride: event.sslTargetNameOverride,
		ipV4: event.ipV4,
		ipV6: event.ipV6,
		bridgePileName: '',
		services: [],
		subState: 'pinned',
		generation: event.generation,
	})
	rebuild(ctx, runtime)
	// Re-pinning the same node to a new address/generation must drop the old
	// pinned channel so the next acquire dials the new target.
	if (prev !== undefined && (prev.address !== address || prev.generation !== event.generation)) {
		return {
			effects: [
				{ type: 'endpoints.effect.close_channel', nodeId: event.nodeId, store: 'pinned' },
			],
		}
	}
}
 
let applyInvalidate = function applyInvalidate(
	ctx: EndpointsCtx,
	nodeId: bigint,
	runtime: EndpointsTransitionRuntime
): Result | void {
	let entry = ctx.pinned.get(nodeId)
	Iif (entry === undefined) return
	ctx.pinned.delete(nodeId)
	rebuild(ctx, runtime)
	// Close only the pinned channel — a discovered channel sharing this nodeId
	// stays live (invalidating a pin must not abort healthy discovered streams).
	return { effects: [{ type: 'endpoints.effect.close_channel', nodeId, store: 'pinned' }] }
}
 
// ── Transition ──────────────────────────────────────────────────────────────
export let endpointsTransition = function endpointsTransition(
	ctx: EndpointsCtx,
	event: EndpointsEvent,
	runtime: EndpointsTransitionRuntime
): Result | void {
	let state = runtime.state
 
	// Global hard destroy from any non-terminal state.
	if (state !== 'closed' && event.type === 'endpoints.destroy') {
		return terminate(ctx, event.reason ?? new Error('Endpoints destroyed'), runtime)
	}
 
	switch (state) {
		case 'idle':
			switch (event.type) {
				case 'endpoints.discovery.start':
					ctx.roundInFlight = true
					return {
						state: 'discovering',
						effects: [{ type: 'endpoints.effect.run_discovery_round' }],
					}
				case 'endpoints.pin':
					return applyPin(ctx, event, runtime)
				case 'endpoints.invalidate':
					return applyInvalidate(ctx, event.nodeId, runtime)
				case 'endpoints.close':
					return terminate(ctx, new Error('Endpoints closed'), runtime)
				default:
					return ignored()
			}
 
		case 'discovering':
			switch (event.type) {
				case 'endpoints.discovery.round_succeeded': {
					if (event.endpoints.length === 0) return rejectEmptyRound(ctx, runtime)
					let firstReady = !ctx.hasEverDiscovered
					ctx.hasEverDiscovered = true
					let effects = applyRound(
						ctx,
						event.endpoints,
						event.selfLocation,
						event.pileStates,
						runtime
					)
					effects.push({
						type: 'endpoints.effect.timer.schedule',
						which: 'discovery_interval',
					})
					effects.push({ type: 'endpoints.effect.timer.schedule', which: 'idle_sweep' })
					Eif (firstReady) runtime.emit({ type: 'endpoints.ready' })
					return { state: healthState(ctx), effects }
				}
				case 'endpoints.discovery.round_failed': {
					ctx.attempts += 1
					ctx.lastError = event.error
					ctx.roundInFlight = false
					runtime.emit({
						type: 'endpoints.discovery_failed',
						error: event.error,
						attempt: ctx.attempts,
						retryable: event.retryable,
					})
					if (!event.retryable) {
						return terminate(ctx, event.error, runtime, event.error)
					}
					return {
						effects: [
							{ type: 'endpoints.effect.timer.schedule', which: 'discovery_backoff' },
						],
					}
				}
				case 'endpoints.timer.discovery_backoff':
					ctx.roundInFlight = true
					return { effects: [{ type: 'endpoints.effect.run_discovery_round' }] }
				case 'endpoints.pin':
					return applyPin(ctx, event, runtime)
				case 'endpoints.invalidate':
					return applyInvalidate(ctx, event.nodeId, runtime)
				case 'endpoints.close':
					return toClosing(ctx, runtime)
				default:
					return ignored()
			}
 
		case 'ready':
		case 'degraded':
			switch (event.type) {
				case 'endpoints.discovery.round_succeeded': {
					// Same guard as in `discovering`: applying an empty round here
					// would retire every node and black-hole balanced RPCs until the
					// next interval while the state still reads 'ready'.
					if (event.endpoints.length === 0) return rejectEmptyRound(ctx, runtime)
					let effects = applyRound(
						ctx,
						event.endpoints,
						event.selfLocation,
						event.pileStates,
						runtime
					)
					effects.push({
						type: 'endpoints.effect.timer.schedule',
						which: 'discovery_interval',
					})
					return { state: healthState(ctx), effects }
				}
				case 'endpoints.discovery.round_failed':
					// Background failure is never terminal — keep serving the last
					// snapshot; the interval/backoff retries.
					ctx.attempts += 1
					ctx.lastError = event.error
					ctx.roundInFlight = false
					runtime.emit({
						type: 'endpoints.discovery_failed',
						error: event.error,
						attempt: ctx.attempts,
						retryable: event.retryable,
					})
					return {
						effects: [
							{ type: 'endpoints.effect.timer.schedule', which: 'discovery_backoff' },
						],
					}
				case 'endpoints.discovery.force':
				case 'endpoints.timer.discovery_interval':
				case 'endpoints.timer.discovery_backoff':
					if (ctx.roundInFlight) return ignored()
					ctx.roundInFlight = true
					return {
						effects: [
							{ type: 'endpoints.effect.timer.clear', which: 'discovery_backoff' },
							{ type: 'endpoints.effect.run_discovery_round' },
						],
					}
				case 'endpoints.rpc_failed': {
					let entry = ctx.byNodeId.get(event.nodeId)
					if (entry === undefined || entry.subState !== 'active') return ignored()
					entry.subState = 'pessimized'
					runtime.emit({
						type: 'endpoints.pessimized',
						nodeId: entry.nodeId,
						address: entry.address,
						location: entry.location,
						pile: entry.bridgePileName,
					})
					rebuild(ctx, runtime)
					// Cross into degraded and force a round when too many are down.
					let next = healthState(ctx)
					if (next === 'degraded' && !ctx.roundInFlight) {
						ctx.roundInFlight = true
						return {
							state: 'degraded',
							effects: [
								{
									type: 'endpoints.effect.timer.clear',
									which: 'discovery_backoff',
								},
								{ type: 'endpoints.effect.run_discovery_round' },
							],
						}
					}
					return { state: next }
				}
				case 'endpoints.rpc_ok': {
					let entry = ctx.byNodeId.get(event.nodeId)
					if (entry === undefined || entry.subState !== 'pessimized') return ignored()
					entry.subState = 'active'
					runtime.emit({
						type: 'endpoints.unpessimized',
						nodeId: entry.nodeId,
						address: entry.address,
						location: entry.location,
						pile: entry.bridgePileName,
					})
					rebuild(ctx, runtime)
					return { state: healthState(ctx) }
				}
				case 'endpoints.timer.idle_sweep':
					return { effects: [{ type: 'endpoints.effect.idle_sweep' }] }
				case 'endpoints.channel_closeable': {
					let entry = ctx.byNodeId.get(event.nodeId)
					Iif (entry === undefined || entry.subState !== 'retired') return ignored()
					ctx.byNodeId.delete(event.nodeId)
					rebuild(ctx, runtime)
					runtime.emit({
						type: 'endpoints.removed',
						nodeId: entry.nodeId,
						address: entry.address,
						location: entry.location,
						pile: entry.bridgePileName,
						reason: 'idle',
					})
					return {
						effects: [{ type: 'endpoints.effect.close_channel', nodeId: event.nodeId }],
					}
				}
				case 'endpoints.pin':
					return applyPin(ctx, event, runtime)
				case 'endpoints.invalidate':
					return applyInvalidate(ctx, event.nodeId, runtime)
				case 'endpoints.close':
					return toClosing(ctx, runtime)
				default:
					return ignored()
			}
 
		case 'closing':
			switch (event.type) {
				case 'endpoints.channel_closeable': {
					let entry = ctx.byNodeId.get(event.nodeId) ?? ctx.pinned.get(event.nodeId)
					Iif (entry === undefined) return ignored()
					ctx.byNodeId.delete(event.nodeId)
					ctx.pinned.delete(event.nodeId)
					runtime.emit({
						type: 'endpoints.removed',
						nodeId: entry.nodeId,
						address: entry.address,
						location: entry.location,
						pile: entry.bridgePileName,
						reason: 'pool_close',
					})
					let close: EndpointsEffect = {
						type: 'endpoints.effect.close_channel',
						nodeId: event.nodeId,
					}
					// Last channel drained → finalize, keeping this node's close effect.
					if (ctx.byNodeId.size === 0 && ctx.pinned.size === 0) {
						let term = terminate(ctx, new Error('Endpoints closed'), runtime)
						return { ...term, effects: [close, ...(term.effects ?? [])] }
					}
					return { effects: [close] }
				}
				case 'endpoints.timer.close_deadline':
					return terminate(ctx, new Error('Endpoints closed'), runtime)
				case 'endpoints.discovery.round_succeeded':
				case 'endpoints.discovery.round_failed':
					ctx.roundInFlight = false
					return ignored()
				default:
					return ignored()
			}
 
		case 'closed':
			return ignored()
 
		/* node:coverage ignore start -- EndpointsState is exhaustive; default unreachable */
		default:
			return ignored()
		/* node:coverage ignore stop */
	}
}