All files / packages/coordination/src client.ts

90.47% Statements 38/42
87.5% Branches 7/8
100% Functions 10/10
90.24% Lines 37/41

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                          10x                           11x 11x       8x       80x       2x       81x               108x   108x 108x   108x 108x 106x 106x                         4x 4x 5x         5x 5x     3x 3x 2x     1x                   3x   3x 3x   3x         10x       3x   3x         3x 1x   2x   3x   3x 2x     1x    
import { linkSignals } from '@ydbjs/abortable'
import type { Driver } from '@ydbjs/core'
import { loggers } from '@ydbjs/debug'
 
import { createDeferred } from './runtime/session-registry.js'
 
import {
	type CoordinationNodeConfig,
	type CoordinationNodeDescription,
	CoordinationNodeRuntime,
} from './node.js'
import { CoordinationSession } from './session.js'
 
let dbg = loggers.coordination.extend('client')
 
export interface SessionOptions {
	description?: string
	recoveryWindow?: number
	startTimeout?: number
	retryBackoff?: number
}
 
export class CoordinationClient {
	#driver: Driver
	#nodeRuntime: CoordinationNodeRuntime
 
	constructor(driver: Driver) {
		this.#driver = driver
		this.#nodeRuntime = new CoordinationNodeRuntime(driver)
	}
 
	describeNode(path: string, signal?: AbortSignal): Promise<CoordinationNodeDescription> {
		return this.#nodeRuntime.describe(path, signal)
	}
 
	createNode(path: string, config: CoordinationNodeConfig, signal?: AbortSignal): Promise<void> {
		return this.#nodeRuntime.create(path, config, signal)
	}
 
	alterNode(path: string, config: CoordinationNodeConfig, signal?: AbortSignal): Promise<void> {
		return this.#nodeRuntime.alter(path, config, signal)
	}
 
	dropNode(path: string, signal?: AbortSignal): Promise<void> {
		return this.#nodeRuntime.drop(path, signal)
	}
 
	async createSession(
		path: string,
		options?: SessionOptions,
		signal?: AbortSignal
	): Promise<CoordinationSession> {
		signal?.throwIfAborted()
 
		dbg.log('creating session on %s', path)
		let session = new CoordinationSession(this.#driver, { path, ...options }, signal)
 
		try {
			await session.waitReady(signal)
			dbg.log('session ready on %s (id=%s)', path, session.sessionId)
			return session
		} catch (error) {
			dbg.log('failed to open session on %s: %O', path, error)
			session.destroy(error)
			throw error
		}
	}
 
	async *openSession(
		path: string,
		options?: SessionOptions,
		signal?: AbortSignal
	): AsyncIterable<CoordinationSession> {
		dbg.log('opening persistent session on %s', path)
		for (;;) {
			Iif (signal?.aborted) {
				return
			}
 
			// oxlint-disable-next-line no-await-in-loop
			let session = await this.createSession(path, options, signal)
			yield session
 
			// oxlint-disable-next-line no-await-in-loop
			let shouldOpenNext = await shouldOpenNextSession(session, signal)
			if (!shouldOpenNext) {
				return
			}
 
			dbg.log('session expired on %s, reopening', path)
		}
	}
 
	async withSession<T>(
		path: string,
		callback: (session: CoordinationSession) => Promise<T>,
		options?: SessionOptions,
		signal?: AbortSignal
	): Promise<T> {
		let session = await this.createSession(path, options, signal)
 
		try {
			return await callback(session)
		} finally {
			await session.close(signal)
		}
	}
}
 
let shouldOpenNextSession = async function shouldOpenNextSession(
	session: CoordinationSession,
	externalSignal?: AbortSignal
): Promise<boolean> {
	let deferred = createDeferred<void>()
 
	using combined = linkSignals(session.signal, externalSignal)
	// linkSignals aborts `combined` synchronously when an input is already
	// aborted, without ever firing a future 'abort' event — attaching only a
	// listener here would wait forever if the session already expired (or the
	// caller already cancelled) before this function runs.
	if (combined.signal.aborted) {
		deferred.resolve()
	} else {
		combined.signal.addEventListener('abort', () => deferred.resolve(), { once: true })
	}
	await deferred.promise
 
	if (externalSignal?.aborted) {
		return false
	}
 
	return session.status === 'expired'
}