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'
}
|