All files / packages/coordination/src election.ts

98.07% Statements 51/52
85% Branches 17/20
100% Functions 11/11
98.07% Lines 51/52

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
}