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 | 10x 10x 20x 20x 20x 5x 1x 1x 1x 24x 4x 20x 20x 20x 13x 30x 30x 64x 24x 24x 20x 20x 6x 6x 6x 6x 6x 10x 10x 10x 2x 8x 8x 3x 8x 8x 2x 2x 1x 6x 6x 10x 10x 6x 6x 7x 7x 7x 10x 10x 10x | import { loggers } from '@ydbjs/debug'
import { LeaderChangedError, ObservationEndedError } from './errors.js'
import { Lease, Semaphore } from './semaphore.js'
let dbg = loggers.coordination.extend('election')
let emptyBytes = new Uint8Array()
export interface LeaderInfo {
data: Uint8Array
}
export interface LeaderState {
data: Uint8Array
isMe: boolean
signal: AbortSignal
}
export class Leadership implements AsyncDisposable {
#semaphore: Semaphore
#lease: Lease
#resigned = false
constructor(lease: Lease, semaphore: Semaphore) {
this.#lease = lease
this.#semaphore = semaphore
}
get signal(): AbortSignal {
return this.#lease.signal
}
async proclaim(data: Uint8Array, signal?: AbortSignal): Promise<void> {
signal?.throwIfAborted()
dbg.log('proclaiming leadership on %s (%d bytes)', this.#semaphore.name, data.byteLength)
return this.#semaphore.update(data, signal)
}
async resign(signal?: AbortSignal): Promise<void> {
if (this.#resigned) {
return
}
this.#resigned = true
dbg.log('resigning from leadership on %s', this.#semaphore.name)
await this.#lease.release(signal)
}
async [Symbol.asyncDispose](): Promise<void> {
await this.resign()
}
}
export class Election {
#semaphore: Semaphore
#sessionId: () => bigint | null
constructor(semaphore: Semaphore, sessionId: () => bigint | null) {
this.#semaphore = semaphore
this.#sessionId = sessionId
}
get name(): string {
return this.#semaphore.name
}
async campaign(data: Uint8Array, signal?: AbortSignal): Promise<Leadership> {
dbg.log('campaigning for leadership on %s', this.name)
let lease = await this.#semaphore.acquire({ count: 1, data }, signal)
dbg.log('won leadership on %s', this.name)
return new Leadership(lease, this.#semaphore)
}
async *observe(signal?: AbortSignal): AsyncIterable<LeaderState> {
dbg.log('observing leadership changes on %s', this.name)
let previousLeader: { sessionId: bigint; orderId: bigint } | null = null
let currentController: AbortController | null = null
try {
for await (let description of this.#semaphore.watch({ owners: true }, signal)) {
let owner = description.owners?.[0]
let currentLeader = owner
? { sessionId: owner.sessionId, orderId: owner.orderId }
: null
if (isSameLeader(previousLeader, currentLeader)) {
continue
}
previousLeader = currentLeader
if (currentController) {
currentController.abort(new LeaderChangedError())
}
currentController = new AbortController()
if (!owner) {
dbg.log('no leader on %s', this.name)
yield { data: emptyBytes, isMe: false, signal: currentController.signal }
continue
}
let sessionId = this.#sessionId()
let isMe = sessionId !== null && owner.sessionId === sessionId
dbg.log(
'leader changed on %s (sessionId=%s, isMe=%s)',
this.name,
owner.sessionId,
isMe
)
yield { data: owner.data, isMe, signal: currentController.signal }
}
} finally {
dbg.log('stopped observing %s', this.name)
currentController?.abort(new ObservationEndedError())
}
}
async leader(signal?: AbortSignal): Promise<LeaderInfo | null> {
let description = await this.#semaphore.describe({ owners: true }, signal)
let owner = description.owners?.[0]
return owner ? { data: owner.data } : null
}
}
let isSameLeader = function isSameLeader(
left: { sessionId: bigint; orderId: bigint } | null,
right: { sessionId: bigint; orderId: bigint } | null
): boolean {
Eif (!left || !right) {
return left === right
}
return left.sessionId === right.sessionId && left.orderId === right.orderId
}
|