refactor(server-runtime): split modules
This commit is contained in:
@@ -2,7 +2,7 @@ import type { WebSocketBaseEvent, WebSocketEvents } from '@proj-airi/server-shar
|
||||
|
||||
import { describe, expect, it } from 'vitest'
|
||||
|
||||
import { detectHeartbeatControlFrame, resolveDeliveryConfig, selectConsumerPeerId } from './index'
|
||||
import { heartbeatFrameFrom, resolveEventDelivery, selectConsumerPeerId } from './index'
|
||||
|
||||
function createInputTextEvent(
|
||||
overrides: Partial<WebSocketBaseEvent<'input:text', WebSocketEvents['input:text']>> = {},
|
||||
@@ -27,9 +27,9 @@ function createInputTextEvent(
|
||||
}
|
||||
}
|
||||
|
||||
describe('resolveDeliveryConfig', () => {
|
||||
describe('resolveEventDelivery', () => {
|
||||
it('uses protocol event metadata defaults for input:text', () => {
|
||||
const delivery = resolveDeliveryConfig(createInputTextEvent())
|
||||
const delivery = resolveEventDelivery(createInputTextEvent())
|
||||
|
||||
expect(delivery).toEqual({
|
||||
mode: 'consumer-group',
|
||||
@@ -39,7 +39,7 @@ describe('resolveDeliveryConfig', () => {
|
||||
})
|
||||
|
||||
it('allows route delivery to override protocol defaults', () => {
|
||||
const delivery = resolveDeliveryConfig(createInputTextEvent({
|
||||
const delivery = resolveEventDelivery(createInputTextEvent({
|
||||
route: {
|
||||
delivery: {
|
||||
required: true,
|
||||
@@ -59,7 +59,7 @@ describe('resolveDeliveryConfig', () => {
|
||||
})
|
||||
|
||||
it('returns explicit route delivery for events without protocol defaults', () => {
|
||||
const delivery = resolveDeliveryConfig({
|
||||
const delivery = resolveEventDelivery({
|
||||
type: 'spark:notify',
|
||||
data: {
|
||||
id: 'spark-1',
|
||||
@@ -196,15 +196,15 @@ describe('selectConsumerPeerId', () => {
|
||||
})
|
||||
})
|
||||
|
||||
describe('detectHeartbeatControlFrame', () => {
|
||||
describe('heartbeatFrameFrom', () => {
|
||||
it('recognizes raw websocket control frame text without treating it as protocol JSON', () => {
|
||||
expect(detectHeartbeatControlFrame('ping')).toBe('ping')
|
||||
expect(detectHeartbeatControlFrame('pong')).toBe('pong')
|
||||
expect(heartbeatFrameFrom('ping')).toBe('ping')
|
||||
expect(heartbeatFrameFrom('pong')).toBe('pong')
|
||||
})
|
||||
|
||||
it('ignores non-control payloads', () => {
|
||||
expect(detectHeartbeatControlFrame('')).toBeUndefined()
|
||||
expect(detectHeartbeatControlFrame('🩵')).toBeUndefined()
|
||||
expect(detectHeartbeatControlFrame('{"type":"transport:connection:heartbeat"}')).toBeUndefined()
|
||||
expect(heartbeatFrameFrom('')).toBeUndefined()
|
||||
expect(heartbeatFrameFrom('🩵')).toBeUndefined()
|
||||
expect(heartbeatFrameFrom('{"type":"transport:connection:heartbeat"}')).toBeUndefined()
|
||||
})
|
||||
})
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,78 @@
|
||||
import type { WebSocketEvent } from '@proj-airi/server-shared/types'
|
||||
|
||||
import { stringify } from 'superjson'
|
||||
import { describe, expect, it } from 'vitest'
|
||||
|
||||
import {
|
||||
AiriWebSocketEventFormatError,
|
||||
heartbeatFrameFrom,
|
||||
parseEvent,
|
||||
} from '.'
|
||||
|
||||
describe('airi websocket protocol codec', () => {
|
||||
it('parses superjson encoded events', () => {
|
||||
const event: WebSocketEvent = {
|
||||
type: 'module:authenticate',
|
||||
data: { token: 'secret' },
|
||||
metadata: {
|
||||
source: {
|
||||
kind: 'plugin',
|
||||
id: 'test-plugin-1',
|
||||
plugin: { id: 'test-plugin' },
|
||||
},
|
||||
event: { id: 'event-1' },
|
||||
},
|
||||
}
|
||||
|
||||
expect(parseEvent(stringify(event))).toEqual(event)
|
||||
})
|
||||
|
||||
it('falls back to plain JSON events', () => {
|
||||
const event: WebSocketEvent = {
|
||||
type: 'module:authenticate',
|
||||
data: { token: 'secret' },
|
||||
metadata: {
|
||||
source: {
|
||||
kind: 'plugin',
|
||||
id: 'test-plugin-1',
|
||||
plugin: { id: 'test-plugin' },
|
||||
},
|
||||
event: { id: 'event-1' },
|
||||
},
|
||||
}
|
||||
|
||||
expect(parseEvent(JSON.stringify(event))).toEqual(event)
|
||||
})
|
||||
|
||||
it('rejects payloads without event type', () => {
|
||||
expect(() => parseEvent('null'))
|
||||
.toThrow(AiriWebSocketEventFormatError)
|
||||
expect(() => parseEvent(JSON.stringify({ data: {} })))
|
||||
.toThrow(AiriWebSocketEventFormatError)
|
||||
})
|
||||
|
||||
it('rejects payloads with non-string event type', () => {
|
||||
expect(() => parseEvent(JSON.stringify({ type: 0, data: {} })))
|
||||
.toThrow(AiriWebSocketEventFormatError)
|
||||
})
|
||||
|
||||
it('rejects payloads without object event data', () => {
|
||||
expect(() => parseEvent(JSON.stringify({ type: 'module:authenticate' })))
|
||||
.toThrow(AiriWebSocketEventFormatError)
|
||||
expect(() => parseEvent(JSON.stringify({ type: 'module:authenticate', data: null })))
|
||||
.toThrow(AiriWebSocketEventFormatError)
|
||||
expect(() => parseEvent(JSON.stringify({ type: 'module:authenticate', data: 'secret' })))
|
||||
.toThrow(AiriWebSocketEventFormatError)
|
||||
})
|
||||
|
||||
it('rejects payloads with array event data', () => {
|
||||
expect(() => parseEvent(JSON.stringify({ type: 'module:authenticate', data: [] })))
|
||||
.toThrow(AiriWebSocketEventFormatError)
|
||||
})
|
||||
|
||||
it('classifies raw ping and pong control frames', () => {
|
||||
expect(heartbeatFrameFrom('ping')).toBe('ping')
|
||||
expect(heartbeatFrameFrom('pong')).toBe('pong')
|
||||
expect(heartbeatFrameFrom('{"type":"ping"}')).toBeUndefined()
|
||||
})
|
||||
})
|
||||
@@ -0,0 +1,340 @@
|
||||
import type { DeliveryConfig, MessageHeartbeat, MetadataEventSource, WebSocketBaseEvent, WebSocketEvent } from '@proj-airi/server-shared/types'
|
||||
|
||||
import type {
|
||||
RouteContext,
|
||||
RouteDecision,
|
||||
RouteMiddleware,
|
||||
} from '../../middlewares'
|
||||
import type { Peer } from '../../types'
|
||||
|
||||
import { ServerErrorMessages } from '@proj-airi/server-shared'
|
||||
import {
|
||||
getProtocolEventMetadata,
|
||||
MessageHeartbeatKind,
|
||||
WebSocketEventSource,
|
||||
} from '@proj-airi/server-shared/types'
|
||||
import { nanoid } from 'nanoid'
|
||||
import { parse, stringify } from 'superjson'
|
||||
|
||||
import packageJSON from '../../../package.json'
|
||||
|
||||
import { createEventCodec, createGatewayLifecycle } from '../core'
|
||||
|
||||
const invalidAiriWebSocketEventFormatMessage = 'Invalid WebSocket event format.'
|
||||
|
||||
/**
|
||||
* Close details surfaced by the websocket runtime for AIRI peer shutdown logging.
|
||||
*/
|
||||
export interface AiriServerWsCloseDetails {
|
||||
/** WebSocket close code when the runtime reports one. */
|
||||
code?: number
|
||||
/** WebSocket close reason when the runtime reports one. */
|
||||
reason?: string
|
||||
/** Whether the runtime considers the close clean. */
|
||||
wasClean?: unknown
|
||||
}
|
||||
|
||||
/**
|
||||
* Error thrown when a websocket message parses as JSON but is not an AIRI event envelope.
|
||||
*
|
||||
* Use when:
|
||||
* - The runtime must distinguish malformed event envelopes from invalid JSON text
|
||||
*
|
||||
* Expects:
|
||||
* - Callers convert this to the protocol `invalidEventFormat` response
|
||||
*
|
||||
* Returns:
|
||||
* - A typed error for invalid AIRI websocket event envelopes
|
||||
*/
|
||||
export class AiriWebSocketEventFormatError extends Error {
|
||||
constructor() {
|
||||
super(invalidAiriWebSocketEventFormatMessage)
|
||||
this.name = 'AiriWebSocketEventFormatError'
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Creates the AIRI websocket gateway wrapper.
|
||||
*
|
||||
* Use when:
|
||||
* - `setupApp(...)` needs a gateway object to mount on `/ws`
|
||||
*
|
||||
* Expects:
|
||||
* - `handler` preserves the existing AIRI websocket lifecycle behavior
|
||||
*
|
||||
* Returns:
|
||||
* - A gateway object compatible with H3 `defineWebSocketHandler(...)`
|
||||
*/
|
||||
export function createGateway(input: {
|
||||
handler: {
|
||||
open: (peer: Peer) => void
|
||||
message: (peer: Peer, message: { text: () => string }) => void
|
||||
error: (peer: Peer, error: unknown) => void
|
||||
close: (peer: Peer, details?: AiriServerWsCloseDetails) => void
|
||||
}
|
||||
dispose?: () => void
|
||||
}) {
|
||||
return createGatewayLifecycle({
|
||||
handler: input.handler,
|
||||
dispose: input.dispose,
|
||||
})
|
||||
}
|
||||
|
||||
/**
|
||||
* Creates metadata for events emitted by the AIRI websocket runtime.
|
||||
*
|
||||
* Use when:
|
||||
* - The server sends protocol events to connected peers
|
||||
* - Response events should preserve parent event correlation
|
||||
*
|
||||
* Expects:
|
||||
* - `serverInstanceId` identifies the active server runtime instance
|
||||
*
|
||||
* Returns:
|
||||
* - AIRI protocol metadata with server source and event id
|
||||
*/
|
||||
export function createEventMetadata(
|
||||
serverInstanceId: string,
|
||||
parentId?: string,
|
||||
): { source: MetadataEventSource, event: { id: string, parentId?: string } } {
|
||||
return {
|
||||
event: {
|
||||
id: nanoid(),
|
||||
parentId,
|
||||
},
|
||||
source: {
|
||||
kind: 'plugin',
|
||||
plugin: {
|
||||
id: WebSocketEventSource.Server,
|
||||
version: packageJSON.version,
|
||||
},
|
||||
id: serverInstanceId,
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Creates AIRI server response event factories.
|
||||
*
|
||||
* Use when:
|
||||
* - WebSocket handlers need stable response event shapes
|
||||
*
|
||||
* Expects:
|
||||
* - `serverInstanceId` identifies the current server runtime
|
||||
*
|
||||
* Returns:
|
||||
* - Factory methods for protocol responses emitted by the server
|
||||
*/
|
||||
export function createResponses(serverInstanceId: string) {
|
||||
return {
|
||||
authenticated(parentId?: string) {
|
||||
return {
|
||||
type: 'module:authenticated',
|
||||
data: { authenticated: true },
|
||||
metadata: createEventMetadata(serverInstanceId, parentId),
|
||||
} satisfies WebSocketEvent<Record<string, unknown>>
|
||||
},
|
||||
notAuthenticated(parentId?: string) {
|
||||
return {
|
||||
type: 'error',
|
||||
data: { message: ServerErrorMessages.notAuthenticated },
|
||||
metadata: createEventMetadata(serverInstanceId, parentId),
|
||||
} satisfies WebSocketEvent<Record<string, unknown>>
|
||||
},
|
||||
error(message: string, parentId?: string) {
|
||||
return {
|
||||
type: 'error',
|
||||
data: { message },
|
||||
metadata: createEventMetadata(serverInstanceId, parentId),
|
||||
} satisfies WebSocketEvent<Record<string, unknown>>
|
||||
},
|
||||
heartbeat(kind: MessageHeartbeatKind, message: MessageHeartbeat | string, parentId?: string) {
|
||||
return {
|
||||
type: 'transport:connection:heartbeat',
|
||||
data: { kind, message, at: Date.now() },
|
||||
metadata: createEventMetadata(serverInstanceId, parentId),
|
||||
} satisfies WebSocketEvent<Record<string, unknown>>
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Checks whether an error came from AIRI websocket event envelope validation.
|
||||
*
|
||||
* Use when:
|
||||
* - Message handlers need to map invalid envelopes to protocol errors
|
||||
*
|
||||
* Expects:
|
||||
* - Parser code throws {@link AiriWebSocketEventFormatError} for envelope failures
|
||||
*
|
||||
* Returns:
|
||||
* - `true` when the error should become `ServerErrorMessages.invalidEventFormat`
|
||||
*/
|
||||
export function isAiriWebSocketEventFormatError(error: unknown): error is AiriWebSocketEventFormatError {
|
||||
return error instanceof AiriWebSocketEventFormatError
|
||||
}
|
||||
|
||||
/**
|
||||
* Detects raw websocket heartbeat control frames surfaced as text payloads.
|
||||
*
|
||||
* Use when:
|
||||
* - A websocket runtime forwards ping/pong frames through the normal message callback
|
||||
* - The runtime should ignore transport heartbeats instead of treating them as protocol JSON
|
||||
*
|
||||
* Expects:
|
||||
* - Raw text payloads such as `ping` and `pong`
|
||||
*
|
||||
* Returns:
|
||||
* - The heartbeat kind when the text is a control frame, otherwise `undefined`
|
||||
*/
|
||||
export function heartbeatFrameFrom(text: string): MessageHeartbeatKind | undefined {
|
||||
if (text === MessageHeartbeatKind.Ping || text === MessageHeartbeatKind.Pong) {
|
||||
return text
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Parses one AIRI websocket protocol event.
|
||||
*
|
||||
* Use when:
|
||||
* - Reading text messages from WebSocket peers
|
||||
*
|
||||
* Expects:
|
||||
* - SDK clients may send `superjson.stringify(...)`
|
||||
* - External clients may send plain JSON
|
||||
*
|
||||
* Returns:
|
||||
* - A WebSocket event with a string `type`
|
||||
*/
|
||||
export function parseEvent(text: string): WebSocketEvent {
|
||||
// NOTICE:
|
||||
// SDK clients send events using superjson.stringify, so websocket runtime code must
|
||||
// use superjson.parse instead of message.json() or plain JSON.parse first.
|
||||
// JSON.parse on a superjson-encoded string returns the wrapper object
|
||||
// `{ json: {...}, meta: {...} }` with no protocol `type`, which breaks routing.
|
||||
// Keep this until all AIRI websocket clients share one non-wrapper wire format.
|
||||
let parsed: WebSocketEvent | undefined
|
||||
try {
|
||||
parsed = parse<WebSocketEvent>(text)
|
||||
}
|
||||
catch {
|
||||
parsed = undefined
|
||||
}
|
||||
|
||||
const potentialEvent = (parsed && typeof parsed === 'object' && 'type' in parsed)
|
||||
? parsed
|
||||
: JSON.parse(text)
|
||||
|
||||
if (
|
||||
!potentialEvent
|
||||
|| typeof potentialEvent !== 'object'
|
||||
|| !('type' in potentialEvent)
|
||||
|| typeof potentialEvent.type !== 'string'
|
||||
|| !('data' in potentialEvent)
|
||||
|| !potentialEvent.data
|
||||
|| typeof potentialEvent.data !== 'object'
|
||||
|| Array.isArray(potentialEvent.data)
|
||||
) {
|
||||
throw new AiriWebSocketEventFormatError()
|
||||
}
|
||||
|
||||
return potentialEvent as WebSocketEvent
|
||||
}
|
||||
|
||||
/**
|
||||
* Serializes one AIRI websocket protocol event.
|
||||
*
|
||||
* Use when:
|
||||
* - Sending AIRI events through WebSocket peers
|
||||
*
|
||||
* Expects:
|
||||
* - `event` is already protocol-shaped
|
||||
*
|
||||
* Returns:
|
||||
* - SuperJSON text payload matching existing runtime behavior
|
||||
*/
|
||||
export function stringifyEvent(event: WebSocketBaseEvent<string, unknown> | string) {
|
||||
return typeof event === 'string' ? event : stringify(event)
|
||||
}
|
||||
|
||||
/**
|
||||
* Resolves the effective event delivery policy.
|
||||
*
|
||||
* Use when:
|
||||
* - Protocol defaults should be merged with route-level delivery overrides
|
||||
* - Routing needs to know whether the event should broadcast or target one consumer
|
||||
*
|
||||
* Expects:
|
||||
* - Route delivery to override protocol metadata field-by-field
|
||||
*
|
||||
* Returns:
|
||||
* - The merged broadcast/consumer delivery policy, or `undefined` when unrestricted
|
||||
*/
|
||||
export function resolveEventDelivery(event: WebSocketEvent): DeliveryConfig | undefined {
|
||||
const eventMetadata = getProtocolEventMetadata(event.type)
|
||||
const defaultDelivery = eventMetadata?.delivery
|
||||
const routeDelivery = event.route?.delivery
|
||||
|
||||
if (!defaultDelivery && !routeDelivery) {
|
||||
return undefined
|
||||
}
|
||||
|
||||
return {
|
||||
...defaultDelivery,
|
||||
...routeDelivery,
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Creates event serializer hooks used by server websocket adapters.
|
||||
*
|
||||
* Use when:
|
||||
* - A gateway wants protocol-specific parsing, stringifying, and control-frame detection
|
||||
*
|
||||
* Expects:
|
||||
* - Callers route raw control frames before protocol events
|
||||
*
|
||||
* Returns:
|
||||
* - A reusable `server-ws/core` codec configured for AIRI events
|
||||
*/
|
||||
export function createEventSerializer() {
|
||||
return createEventCodec<WebSocketEvent>({
|
||||
parse: parseEvent,
|
||||
stringify: stringifyEvent,
|
||||
detectControlFrame: heartbeatFrameFrom,
|
||||
})
|
||||
}
|
||||
|
||||
/**
|
||||
* Iterates event middlewares in declaration order until one returns a decision.
|
||||
*
|
||||
* Use when:
|
||||
* - The websocket runtime needs the first route decision from configured middleware
|
||||
*
|
||||
* Expects:
|
||||
* - Middleware functions are ordered by caller policy
|
||||
*
|
||||
* Returns:
|
||||
* - The first route decision, or `undefined` when no middleware decided
|
||||
*/
|
||||
export function forEachEventMiddlewares(input: {
|
||||
event: WebSocketEvent
|
||||
fromPeer: RouteContext['fromPeer']
|
||||
peers: Map<string, RouteContext['fromPeer']>
|
||||
destinations?: RouteContext['destinations']
|
||||
middleware: RouteMiddleware[]
|
||||
}): RouteDecision | undefined {
|
||||
const context: RouteContext = {
|
||||
event: input.event,
|
||||
fromPeer: input.fromPeer,
|
||||
peers: input.peers,
|
||||
destinations: input.destinations,
|
||||
}
|
||||
|
||||
for (const middleware of input.middleware) {
|
||||
const result = middleware(context)
|
||||
if (result) {
|
||||
return result
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,229 @@
|
||||
import type { ServerWsStickyAssignment } from '.'
|
||||
|
||||
import { describe, expect, it } from 'vitest'
|
||||
|
||||
import {
|
||||
createConsumerOrchestrator,
|
||||
selectConsumerPeerId,
|
||||
} from '.'
|
||||
|
||||
describe('server-ws consumer selection', () => {
|
||||
it('selects highest priority then earliest registration', () => {
|
||||
expect(selectConsumerPeerId({
|
||||
eventType: 'event:test',
|
||||
fromPeerId: 'sender',
|
||||
delivery: { mode: 'consumer', selection: 'first' },
|
||||
candidates: [
|
||||
{ peerId: 'late', priority: 1, registeredAt: 2, authenticated: true },
|
||||
{ peerId: 'early', priority: 1, registeredAt: 1, authenticated: true },
|
||||
{ peerId: 'low', priority: 0, registeredAt: 0, authenticated: true },
|
||||
],
|
||||
})).toBe('early')
|
||||
})
|
||||
|
||||
it('skips sender, unauthenticated, and unhealthy candidates', () => {
|
||||
expect(selectConsumerPeerId({
|
||||
eventType: 'event:test',
|
||||
fromPeerId: 'sender',
|
||||
delivery: { mode: 'consumer', selection: 'first' },
|
||||
candidates: [
|
||||
{ peerId: 'sender', priority: 3, registeredAt: 1, authenticated: true },
|
||||
{ peerId: 'unauthenticated', priority: 2, registeredAt: 1, authenticated: false },
|
||||
{ peerId: 'unhealthy', priority: 1, registeredAt: 1, authenticated: true, healthy: false },
|
||||
{ peerId: 'target', priority: 0, registeredAt: 1, authenticated: true },
|
||||
],
|
||||
})).toBe('target')
|
||||
})
|
||||
|
||||
it('preserves sticky assignment for the same sticky key', () => {
|
||||
const stickyAssignments = new Map<string, ServerWsStickyAssignment>()
|
||||
const delivery = { mode: 'consumer-group' as const, group: 'workers', selection: 'sticky' as const, stickyKey: 'job-1' }
|
||||
const candidates = [
|
||||
{ peerId: 'a', priority: 0, registeredAt: 1, authenticated: true },
|
||||
{ peerId: 'b', priority: 0, registeredAt: 2, authenticated: true },
|
||||
]
|
||||
|
||||
expect(selectConsumerPeerId({
|
||||
eventType: 'event:test',
|
||||
fromPeerId: 'sender',
|
||||
delivery,
|
||||
candidates,
|
||||
stickyAssignments,
|
||||
})).toBe('a')
|
||||
expect(selectConsumerPeerId({
|
||||
eventType: 'event:test',
|
||||
fromPeerId: 'sender',
|
||||
delivery,
|
||||
candidates: [...candidates].reverse(),
|
||||
stickyAssignments,
|
||||
})).toBe('a')
|
||||
})
|
||||
|
||||
it('preserves round-robin cursor per event and group', () => {
|
||||
const roundRobinCursor = new Map<string, number>()
|
||||
const delivery = { mode: 'consumer-group' as const, group: 'workers', selection: 'round-robin' as const }
|
||||
const candidates = [
|
||||
{ peerId: 'a', priority: 0, registeredAt: 1, authenticated: true },
|
||||
{ peerId: 'b', priority: 0, registeredAt: 2, authenticated: true },
|
||||
]
|
||||
|
||||
expect(selectConsumerPeerId({ eventType: 'event:test', fromPeerId: 'sender', delivery, candidates, roundRobinCursor })).toBe('a')
|
||||
expect(selectConsumerPeerId({ eventType: 'event:test', fromPeerId: 'sender', delivery, candidates, roundRobinCursor })).toBe('b')
|
||||
expect(selectConsumerPeerId({ eventType: 'event:test', fromPeerId: 'sender', delivery, candidates, roundRobinCursor })).toBe('a')
|
||||
})
|
||||
})
|
||||
|
||||
describe('server-ws consumer registry', () => {
|
||||
it('registers and unregisters consumers', () => {
|
||||
const registry = createConsumerOrchestrator()
|
||||
|
||||
registry.register({ peerId: 'peer-1', event: 'event:test', mode: 'consumer-group', group: 'workers', priority: 2 })
|
||||
expect(registry.listFor({ event: 'event:test', mode: 'consumer-group', group: 'workers' })).toEqual([
|
||||
expect.objectContaining({ peerId: 'peer-1', event: 'event:test', group: 'workers', priority: 2 }),
|
||||
])
|
||||
|
||||
registry.unregister({ peerId: 'peer-1', event: 'event:test', mode: 'consumer-group', group: 'workers' })
|
||||
expect(registry.listFor({ event: 'event:test', mode: 'consumer-group', group: 'workers' })).toEqual([])
|
||||
})
|
||||
|
||||
it('unregisters peer consumers when event and group names contain registry delimiters', () => {
|
||||
const registry = createConsumerOrchestrator()
|
||||
|
||||
// ROOT CAUSE:
|
||||
//
|
||||
// If event or group names contain the previous string key delimiter, peer cleanup
|
||||
// can fail because unregisterPeer reconstructs registry coordinates from a split key.
|
||||
//
|
||||
// Before:
|
||||
// `${event}::${group}` was split back into event/group and missed the original entry.
|
||||
//
|
||||
// After:
|
||||
// Peer cleanup stores structured event/group refs and never decodes registry keys.
|
||||
registry.register({ peerId: 'peer-1', event: 'event::test', mode: 'consumer-group', group: 'group::workers' })
|
||||
registry.unregisterPeer('peer-1')
|
||||
|
||||
expect(registry.listFor({ event: 'event::test', mode: 'consumer-group', group: 'group::workers' })).toEqual([])
|
||||
})
|
||||
|
||||
it('keeps sticky assignments isolated for delimiter-like event and group names', () => {
|
||||
const stickyAssignments = new Map<string, ServerWsStickyAssignment>()
|
||||
const candidates = [
|
||||
{ peerId: 'event::group-target', priority: 0, registeredAt: 1, authenticated: true },
|
||||
{ peerId: 'other-target', priority: 0, registeredAt: 2, authenticated: true },
|
||||
]
|
||||
|
||||
expect(selectConsumerPeerId({
|
||||
eventType: 'event::group',
|
||||
fromPeerId: 'sender',
|
||||
delivery: { mode: 'consumer-group', group: 'target', selection: 'sticky', stickyKey: 'job' },
|
||||
candidates,
|
||||
stickyAssignments,
|
||||
})).toBe('event::group-target')
|
||||
expect(selectConsumerPeerId({
|
||||
eventType: 'event',
|
||||
fromPeerId: 'sender',
|
||||
delivery: { mode: 'consumer-group', group: 'group::target', selection: 'sticky', stickyKey: 'job' },
|
||||
candidates: [
|
||||
{ peerId: 'other-target', priority: 1, registeredAt: 1, authenticated: true },
|
||||
{ peerId: 'event::group-target', priority: 0, registeredAt: 2, authenticated: true },
|
||||
],
|
||||
stickyAssignments,
|
||||
})).toBe('other-target')
|
||||
})
|
||||
|
||||
it('resets round-robin cursor when group membership changes', () => {
|
||||
const registry = createConsumerOrchestrator()
|
||||
registry.register({ peerId: 'a', event: 'event:test', mode: 'consumer-group', group: 'workers' })
|
||||
registry.register({ peerId: 'b', event: 'event:test', mode: 'consumer-group', group: 'workers' })
|
||||
|
||||
expect(registry.select({
|
||||
eventType: 'event:test',
|
||||
fromPeerId: 'sender',
|
||||
delivery: { mode: 'consumer-group', group: 'workers', selection: 'round-robin' },
|
||||
candidates: registry.listFor({ event: 'event:test', mode: 'consumer-group', group: 'workers' }).map(entry => ({
|
||||
peerId: entry.peerId,
|
||||
priority: entry.priority,
|
||||
registeredAt: entry.registeredAt,
|
||||
authenticated: true,
|
||||
})),
|
||||
})).toBe('a')
|
||||
|
||||
registry.unregister({ peerId: 'a', event: 'event:test', mode: 'consumer-group', group: 'workers' })
|
||||
|
||||
expect(registry.select({
|
||||
eventType: 'event:test',
|
||||
fromPeerId: 'sender',
|
||||
delivery: { mode: 'consumer-group', group: 'workers', selection: 'round-robin' },
|
||||
candidates: registry.listFor({ event: 'event:test', mode: 'consumer-group', group: 'workers' }).map(entry => ({
|
||||
peerId: entry.peerId,
|
||||
priority: entry.priority,
|
||||
registeredAt: entry.registeredAt,
|
||||
authenticated: true,
|
||||
})),
|
||||
})).toBe('b')
|
||||
})
|
||||
|
||||
it('keeps round-robin cursor when unregister does not change group membership', () => {
|
||||
const registry = createConsumerOrchestrator()
|
||||
registry.register({ peerId: 'a', event: 'event:test', mode: 'consumer-group', group: 'workers' })
|
||||
registry.register({ peerId: 'b', event: 'event:test', mode: 'consumer-group', group: 'workers' })
|
||||
|
||||
expect(registry.select({
|
||||
eventType: 'event:test',
|
||||
fromPeerId: 'sender',
|
||||
delivery: { mode: 'consumer-group', group: 'workers', selection: 'round-robin' },
|
||||
candidates: registry.listFor({ event: 'event:test', mode: 'consumer-group', group: 'workers' }).map(entry => ({
|
||||
peerId: entry.peerId,
|
||||
priority: entry.priority,
|
||||
registeredAt: entry.registeredAt,
|
||||
authenticated: true,
|
||||
})),
|
||||
})).toBe('a')
|
||||
|
||||
registry.unregister({ peerId: 'missing', event: 'event:test', mode: 'consumer-group', group: 'workers' })
|
||||
|
||||
expect(registry.select({
|
||||
eventType: 'event:test',
|
||||
fromPeerId: 'sender',
|
||||
delivery: { mode: 'consumer-group', group: 'workers', selection: 'round-robin' },
|
||||
candidates: registry.listFor({ event: 'event:test', mode: 'consumer-group', group: 'workers' }).map(entry => ({
|
||||
peerId: entry.peerId,
|
||||
priority: entry.priority,
|
||||
registeredAt: entry.registeredAt,
|
||||
authenticated: true,
|
||||
})),
|
||||
})).toBe('b')
|
||||
})
|
||||
|
||||
it('resets round-robin cursor when group membership grows', () => {
|
||||
const registry = createConsumerOrchestrator()
|
||||
registry.register({ peerId: 'a', event: 'event:test', mode: 'consumer-group', group: 'workers' })
|
||||
registry.register({ peerId: 'b', event: 'event:test', mode: 'consumer-group', group: 'workers' })
|
||||
|
||||
expect(registry.select({
|
||||
eventType: 'event:test',
|
||||
fromPeerId: 'sender',
|
||||
delivery: { mode: 'consumer-group', group: 'workers', selection: 'round-robin' },
|
||||
candidates: registry.listFor({ event: 'event:test', mode: 'consumer-group', group: 'workers' }).map(entry => ({
|
||||
peerId: entry.peerId,
|
||||
priority: entry.priority,
|
||||
registeredAt: entry.registeredAt,
|
||||
authenticated: true,
|
||||
})),
|
||||
})).toBe('a')
|
||||
|
||||
registry.register({ peerId: 'c', event: 'event:test', mode: 'consumer-group', group: 'workers' })
|
||||
|
||||
expect(registry.select({
|
||||
eventType: 'event:test',
|
||||
fromPeerId: 'sender',
|
||||
delivery: { mode: 'consumer-group', group: 'workers', selection: 'round-robin' },
|
||||
candidates: registry.listFor({ event: 'event:test', mode: 'consumer-group', group: 'workers' }).map(entry => ({
|
||||
peerId: entry.peerId,
|
||||
priority: entry.priority,
|
||||
registeredAt: entry.registeredAt,
|
||||
authenticated: true,
|
||||
})),
|
||||
})).toBe('a')
|
||||
})
|
||||
})
|
||||
@@ -0,0 +1,543 @@
|
||||
/**
|
||||
* Delivery settings used by the reusable websocket gateway.
|
||||
*
|
||||
* @param TMode - Delivery mode literals accepted by the adapter.
|
||||
*/
|
||||
export interface ServerWsDeliveryConfig<TMode extends string = 'broadcast' | 'consumer' | 'consumer-group'> {
|
||||
/**
|
||||
* Delivery mode selected by the protocol adapter.
|
||||
*
|
||||
* @default undefined
|
||||
*/
|
||||
mode?: TMode
|
||||
/**
|
||||
* Optional consumer group.
|
||||
*
|
||||
* @default "default" for consumer delivery modes.
|
||||
*/
|
||||
group?: string
|
||||
/**
|
||||
* Selection strategy within the target consumer set.
|
||||
*
|
||||
* @default "first"
|
||||
*/
|
||||
selection?: 'first' | 'priority' | 'sticky' | 'round-robin'
|
||||
/**
|
||||
* Sticky routing key used when `selection` is `sticky`.
|
||||
*
|
||||
* @default undefined
|
||||
*/
|
||||
stickyKey?: string
|
||||
/**
|
||||
* Whether missing consumers should be surfaced as an error by the adapter.
|
||||
*
|
||||
* @default false
|
||||
*/
|
||||
required?: boolean
|
||||
}
|
||||
|
||||
/**
|
||||
* Delivery settings accepted by the reusable consumer registry.
|
||||
*
|
||||
* @param TMode - Consumer delivery mode literals accepted by the adapter.
|
||||
*/
|
||||
export type ServerWsConsumerDeliveryConfig<TMode extends string = 'consumer' | 'consumer-group'> = ServerWsDeliveryConfig<TMode>
|
||||
|
||||
/**
|
||||
* Candidate peer metadata used for consumer selection.
|
||||
*/
|
||||
export interface ServerWsConsumerSelectionCandidate {
|
||||
/** Peer id available to receive the event. */
|
||||
peerId: string
|
||||
/** Higher values are selected before lower values. */
|
||||
priority: number
|
||||
/** Timestamp captured when the peer registered as a consumer. */
|
||||
registeredAt: number
|
||||
/** Whether the peer has completed protocol-level authentication. */
|
||||
authenticated: boolean
|
||||
/** Explicit `false` excludes the peer from selection. */
|
||||
healthy?: boolean
|
||||
}
|
||||
|
||||
/**
|
||||
* Stored consumer registration.
|
||||
*/
|
||||
export interface ServerWsConsumerRegistration {
|
||||
/** Protocol event type consumed by the peer. */
|
||||
event: string
|
||||
/** Normalized consumer group name. */
|
||||
group: string
|
||||
/** Peer id that registered for the event/group pair. */
|
||||
peerId: string
|
||||
/** Higher values are selected before lower values. */
|
||||
priority: number
|
||||
/** Timestamp captured when the peer registered as a consumer. */
|
||||
registeredAt: number
|
||||
}
|
||||
|
||||
/**
|
||||
* Describes protocol-agnostic text encoding and decoding for websocket events.
|
||||
*
|
||||
* @param TEvent - Event envelope shape owned by the protocol adapter.
|
||||
*/
|
||||
export interface ServerWsEventCodec<TEvent> {
|
||||
/** Parses one text payload into a protocol event. */
|
||||
parse: (text: string) => TEvent
|
||||
/** Serializes one protocol event or pre-serialized payload for peer sending. */
|
||||
stringify: (event: TEvent | string) => string
|
||||
/** Detects raw transport control payloads that should not enter protocol routing. */
|
||||
detectControlFrame?: (text: string) => string | undefined
|
||||
}
|
||||
|
||||
/**
|
||||
* Describes a websocket handler object accepted by H3 `defineWebSocketHandler`.
|
||||
*
|
||||
* @param TPeer - Runtime peer object accepted by lifecycle callbacks.
|
||||
* @param TMessage - Runtime message object accepted by the message callback.
|
||||
* @param TCloseDetails - Runtime close details object accepted by the close callback.
|
||||
*/
|
||||
export interface ServerWsGatewayHandler<TPeer = unknown, TMessage = unknown, TCloseDetails = unknown> {
|
||||
/** Called when a peer opens a websocket connection. */
|
||||
open?: (peer: TPeer) => void
|
||||
/** Called when a peer sends one websocket message. */
|
||||
message?: (peer: TPeer, message: TMessage) => void
|
||||
/** Called when the websocket runtime reports an error. */
|
||||
error?: (peer: TPeer, error: unknown) => void
|
||||
/** Called when a peer closes a websocket connection. */
|
||||
close?: (peer: TPeer, details?: TCloseDetails) => void
|
||||
}
|
||||
|
||||
/**
|
||||
* Minimal websocket peer shape used by the reusable gateway.
|
||||
*/
|
||||
export interface ServerWsPeer {
|
||||
/** Stable peer id assigned by the websocket runtime. */
|
||||
get id(): string
|
||||
/** Sends one payload to the peer. */
|
||||
send: (data: unknown, options?: { compress?: boolean }) => number | void | undefined
|
||||
/** Closes the peer connection when the runtime exposes an explicit close hook. */
|
||||
close?: () => void
|
||||
/** WebSocket ready state when exposed by the runtime. */
|
||||
readyState?: number
|
||||
/** Request metadata associated with the websocket upgrade. */
|
||||
request?: {
|
||||
/** Request URL associated with the websocket upgrade. */
|
||||
url?: string
|
||||
/** Request headers associated with the websocket upgrade. */
|
||||
headers?: Headers
|
||||
}
|
||||
/** Remote peer address when exposed by the runtime. */
|
||||
remoteAddress?: string
|
||||
}
|
||||
|
||||
/** Default heartbeat read timeout used by the websocket gateway. */
|
||||
export const serverWsDefaultHeartbeatTtlMs = 60_000
|
||||
|
||||
/** Miss count where a peer becomes unhealthy but remains connected. */
|
||||
export const serverWsHealthCheckMissesUnhealthy = 5
|
||||
|
||||
/** Miss count where a peer is considered dead and should be closed. */
|
||||
export const serverWsHealthCheckMissesDead = serverWsHealthCheckMissesUnhealthy * 2
|
||||
|
||||
const DEFAULT_CONSUMER_GROUP = 'default'
|
||||
|
||||
interface ServerWsConsumerRegistryRef {
|
||||
event: string
|
||||
group: string
|
||||
}
|
||||
|
||||
/**
|
||||
* Sticky consumer assignment stored by the reusable consumer selector.
|
||||
*/
|
||||
export interface ServerWsStickyAssignment {
|
||||
/** Protocol event type the sticky assignment belongs to. */
|
||||
event: string
|
||||
/** Normalized consumer group the sticky assignment belongs to. */
|
||||
group: string
|
||||
/** Peer selected for the sticky key. */
|
||||
peerId: string
|
||||
}
|
||||
|
||||
/**
|
||||
* Creates a websocket event codec from explicit parser and serializer callbacks.
|
||||
*
|
||||
* Use when:
|
||||
* - A protocol adapter wants to plug its own event envelope into `server-ws/core`
|
||||
*
|
||||
* Expects:
|
||||
* - Parser and serializer preserve the adapter's current wire format
|
||||
*
|
||||
* Returns:
|
||||
* - A protocol-agnostic codec object consumed by gateway code
|
||||
*/
|
||||
export function createEventCodec<TEvent>(codec: ServerWsEventCodec<TEvent>) {
|
||||
return codec
|
||||
}
|
||||
|
||||
/**
|
||||
* Wraps websocket lifecycle callbacks and disposal as a reusable mount object.
|
||||
*
|
||||
* Use when:
|
||||
* - Adapters need one stable lifecycle shape for server mounting
|
||||
*
|
||||
* Expects:
|
||||
* - `handler` contains already-bound protocol behavior
|
||||
*
|
||||
* Returns:
|
||||
* - A handler plus idempotent disposal hook
|
||||
*/
|
||||
export function createGatewayLifecycle<TPeer, TMessage, TCloseDetails = unknown>(input: {
|
||||
handler: ServerWsGatewayHandler<TPeer, TMessage, TCloseDetails>
|
||||
dispose?: () => void
|
||||
}) {
|
||||
let disposed = false
|
||||
|
||||
return {
|
||||
handler: input.handler,
|
||||
dispose: () => {
|
||||
if (disposed) {
|
||||
return
|
||||
}
|
||||
|
||||
disposed = true
|
||||
input.dispose?.()
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Resolves the interval used for heartbeat health checks.
|
||||
*
|
||||
* Use when:
|
||||
* - Gateway code needs to convert heartbeat TTL into periodic miss checks
|
||||
*
|
||||
* Expects:
|
||||
* - Very small TTL values should still avoid busy intervals
|
||||
*
|
||||
* Returns:
|
||||
* - Interval in milliseconds
|
||||
*/
|
||||
export function resolveServerWsHealthCheckIntervalMs(heartbeatTtlMs: number) {
|
||||
return Math.max(5_000, Math.floor(heartbeatTtlMs / serverWsHealthCheckMissesUnhealthy))
|
||||
}
|
||||
|
||||
/**
|
||||
* Creates a typed peer store around websocket peer state.
|
||||
*
|
||||
* Use when:
|
||||
* - A gateway needs stable peer lookup, iteration, and cleanup
|
||||
*
|
||||
* Expects:
|
||||
* - `TState` contains protocol-specific peer state
|
||||
*
|
||||
* Returns:
|
||||
* - A small registry over peers keyed by peer id
|
||||
*/
|
||||
export function createServerWsPeerStore<TState extends { peer: ServerWsPeer }>() {
|
||||
const peers = new Map<string, TState>()
|
||||
|
||||
return {
|
||||
peers,
|
||||
get(peerId: string) {
|
||||
return peers.get(peerId)
|
||||
},
|
||||
set(peerId: string, state: TState) {
|
||||
peers.set(peerId, state)
|
||||
return state
|
||||
},
|
||||
delete(peerId: string) {
|
||||
return peers.delete(peerId)
|
||||
},
|
||||
clear() {
|
||||
peers.clear()
|
||||
},
|
||||
values() {
|
||||
return peers.values()
|
||||
},
|
||||
entries() {
|
||||
return peers.entries()
|
||||
},
|
||||
size() {
|
||||
return peers.size
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Checks whether a delivery mode targets the consumer registry.
|
||||
*
|
||||
* Use when:
|
||||
* - A protocol adapter receives broad delivery modes but must call consumer-only APIs
|
||||
*
|
||||
* Expects:
|
||||
* - Non-consumer modes such as `broadcast` should remain outside the consumer registry
|
||||
*
|
||||
* Returns:
|
||||
* - `true` for `consumer` and `consumer-group`
|
||||
*/
|
||||
export function isConsumerDeliveryMode(mode: unknown): mode is ServerWsConsumerDeliveryConfig['mode'] {
|
||||
return mode === 'consumer' || mode === 'consumer-group'
|
||||
}
|
||||
|
||||
/**
|
||||
* Normalizes delivery mode for consumer registration.
|
||||
*
|
||||
* Before:
|
||||
* - undefined with group "workers"
|
||||
*
|
||||
* After:
|
||||
* - "consumer-group"
|
||||
*/
|
||||
export function normalizeConsumerMode(mode: unknown, group?: string): 'consumer' | 'consumer-group' {
|
||||
if (isConsumerDeliveryMode(mode)) {
|
||||
return mode!
|
||||
}
|
||||
|
||||
return group ? 'consumer-group' : 'consumer'
|
||||
}
|
||||
|
||||
/**
|
||||
* Normalizes consumer priority.
|
||||
*
|
||||
* Before:
|
||||
* - NaN
|
||||
*
|
||||
* After:
|
||||
* - 0
|
||||
*/
|
||||
export function normalizeConsumerPriority(priority: unknown) {
|
||||
return typeof priority === 'number' && Number.isFinite(priority)
|
||||
? priority
|
||||
: 0
|
||||
}
|
||||
|
||||
function normalizeConsumerGroup(mode: ServerWsConsumerDeliveryConfig['mode'], group?: string) {
|
||||
if (mode === 'consumer') {
|
||||
return DEFAULT_CONSUMER_GROUP
|
||||
}
|
||||
|
||||
return group || DEFAULT_CONSUMER_GROUP
|
||||
}
|
||||
|
||||
function getConsumerRegistryKey(event: string, group: string) {
|
||||
return JSON.stringify([event, group])
|
||||
}
|
||||
|
||||
function getStickyRegistryKey(event: string, group: string, stickyKey: string) {
|
||||
return JSON.stringify([event, group, stickyKey])
|
||||
}
|
||||
|
||||
function sortConsumers(entries: Array<Pick<ServerWsConsumerSelectionCandidate, 'peerId' | 'priority' | 'registeredAt'>>) {
|
||||
return [...entries].sort((left, right) => {
|
||||
if (right.priority !== left.priority) {
|
||||
return right.priority - left.priority
|
||||
}
|
||||
|
||||
return left.registeredAt - right.registeredAt
|
||||
})
|
||||
}
|
||||
|
||||
/**
|
||||
* Selects a concrete consumer peer for consumer-style delivery modes.
|
||||
*
|
||||
* Use when:
|
||||
* - An event should be sent to exactly one registered consumer
|
||||
* - Sticky or round-robin routing needs to be resolved against live peer metadata
|
||||
*
|
||||
* Expects:
|
||||
* - Candidates already describe authenticated and health state
|
||||
*
|
||||
* Returns:
|
||||
* - The selected peer id, or `undefined` when no eligible consumer is available
|
||||
*/
|
||||
export function selectConsumerPeerId(options: {
|
||||
eventType: string
|
||||
fromPeerId: string
|
||||
delivery?: ServerWsDeliveryConfig
|
||||
candidates: ServerWsConsumerSelectionCandidate[]
|
||||
roundRobinCursor?: Map<string, number>
|
||||
stickyAssignments?: Map<string, ServerWsStickyAssignment>
|
||||
}) {
|
||||
const { candidates, delivery, eventType, fromPeerId } = options
|
||||
if (!delivery || (delivery.mode !== 'consumer' && delivery.mode !== 'consumer-group')) {
|
||||
return
|
||||
}
|
||||
|
||||
const normalizedGroup = normalizeConsumerGroup(delivery.mode, delivery.group)
|
||||
const registryKey = getConsumerRegistryKey(eventType, normalizedGroup)
|
||||
const availableEntries = sortConsumers(
|
||||
candidates
|
||||
.filter(entry => entry.peerId !== fromPeerId)
|
||||
.filter(entry => entry.authenticated && entry.healthy !== false),
|
||||
)
|
||||
|
||||
if (availableEntries.length === 0) {
|
||||
return
|
||||
}
|
||||
|
||||
const selection = delivery.selection ?? 'first'
|
||||
if (selection === 'sticky' && delivery.stickyKey) {
|
||||
const stickyRegistryKey = getStickyRegistryKey(eventType, normalizedGroup, delivery.stickyKey)
|
||||
const stickyAssignment = options.stickyAssignments?.get(stickyRegistryKey)
|
||||
if (stickyAssignment && stickyAssignment.peerId !== fromPeerId) {
|
||||
const stickyCandidate = availableEntries.find(entry => entry.peerId === stickyAssignment.peerId)
|
||||
if (stickyCandidate) {
|
||||
return stickyAssignment.peerId
|
||||
}
|
||||
}
|
||||
|
||||
const selected = availableEntries[0]
|
||||
options.stickyAssignments?.set(stickyRegistryKey, { event: eventType, group: normalizedGroup, peerId: selected.peerId })
|
||||
return selected.peerId
|
||||
}
|
||||
|
||||
if (selection === 'round-robin') {
|
||||
const cursor = options.roundRobinCursor?.get(registryKey) ?? 0
|
||||
const selected = availableEntries[cursor % availableEntries.length]
|
||||
options.roundRobinCursor?.set(registryKey, (cursor + 1) % availableEntries.length)
|
||||
return selected.peerId
|
||||
}
|
||||
|
||||
return availableEntries[0].peerId
|
||||
}
|
||||
|
||||
/**
|
||||
* Creates a reusable consumer delivery orchestrator for websocket peers.
|
||||
*
|
||||
* Use when:
|
||||
* - A protocol adapter supports one-consumer delivery or consumer groups
|
||||
*
|
||||
* Expects:
|
||||
* - Peer liveness is checked by the caller before delivery
|
||||
*
|
||||
* Returns:
|
||||
* - Registration, unregister, listing, selection, and cleanup helpers
|
||||
*/
|
||||
export function createConsumerOrchestrator() {
|
||||
const consumerRegistry = new Map<string, Map<string, Map<string, ServerWsConsumerRegistration>>>()
|
||||
const consumerKeysByPeer = new Map<string, Map<string, ServerWsConsumerRegistryRef>>()
|
||||
const deliveryRoundRobinCursor = new Map<string, number>()
|
||||
const stickyAssignments = new Map<string, ServerWsStickyAssignment>()
|
||||
|
||||
function removeStickyAssignmentsFor(event: string, group: string, peerId?: string) {
|
||||
for (const [stickyKey, assignment] of stickyAssignments.entries()) {
|
||||
if (peerId && assignment.peerId !== peerId) {
|
||||
continue
|
||||
}
|
||||
|
||||
if (assignment.event === event && assignment.group === group) {
|
||||
stickyAssignments.delete(stickyKey)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return {
|
||||
register(input: { peerId: string, event: string, mode: ServerWsConsumerDeliveryConfig['mode'], group?: string, priority?: number }) {
|
||||
const normalizedGroup = normalizeConsumerGroup(input.mode, input.group)
|
||||
const registryKey = getConsumerRegistryKey(input.event, normalizedGroup)
|
||||
let groups = consumerRegistry.get(input.event)
|
||||
if (!groups) {
|
||||
groups = new Map()
|
||||
consumerRegistry.set(input.event, groups)
|
||||
}
|
||||
|
||||
let peersForGroup = groups.get(normalizedGroup)
|
||||
if (!peersForGroup) {
|
||||
peersForGroup = new Map()
|
||||
groups.set(normalizedGroup, peersForGroup)
|
||||
}
|
||||
|
||||
const didGrowMembership = !peersForGroup.has(input.peerId)
|
||||
peersForGroup.set(input.peerId, {
|
||||
event: input.event,
|
||||
group: normalizedGroup,
|
||||
peerId: input.peerId,
|
||||
priority: normalizeConsumerPriority(input.priority),
|
||||
registeredAt: Date.now(),
|
||||
})
|
||||
if (didGrowMembership) {
|
||||
deliveryRoundRobinCursor.delete(registryKey)
|
||||
}
|
||||
|
||||
let registrations = consumerKeysByPeer.get(input.peerId)
|
||||
if (!registrations) {
|
||||
registrations = new Map()
|
||||
consumerKeysByPeer.set(input.peerId, registrations)
|
||||
}
|
||||
registrations.set(registryKey, { event: input.event, group: normalizedGroup })
|
||||
},
|
||||
unregister(input: { peerId: string, event: string, mode: ServerWsConsumerDeliveryConfig['mode'], group?: string }) {
|
||||
const normalizedGroup = normalizeConsumerGroup(input.mode, input.group)
|
||||
const registryKey = getConsumerRegistryKey(input.event, normalizedGroup)
|
||||
const groups = consumerRegistry.get(input.event)
|
||||
const peersForGroup = groups?.get(normalizedGroup)
|
||||
const didDelete = peersForGroup?.delete(input.peerId) ?? false
|
||||
|
||||
if (!didDelete) {
|
||||
return
|
||||
}
|
||||
|
||||
deliveryRoundRobinCursor.delete(registryKey)
|
||||
if (peersForGroup?.size === 0) {
|
||||
groups?.delete(normalizedGroup)
|
||||
}
|
||||
if (groups?.size === 0) {
|
||||
consumerRegistry.delete(input.event)
|
||||
}
|
||||
|
||||
const registrations = consumerKeysByPeer.get(input.peerId)
|
||||
registrations?.delete(registryKey)
|
||||
if (registrations?.size === 0) {
|
||||
consumerKeysByPeer.delete(input.peerId)
|
||||
}
|
||||
|
||||
removeStickyAssignmentsFor(input.event, normalizedGroup, input.peerId)
|
||||
},
|
||||
unregisterPeer(peerId: string) {
|
||||
const registrations = consumerKeysByPeer.get(peerId)
|
||||
if (!registrations?.size) {
|
||||
return
|
||||
}
|
||||
|
||||
for (const registration of registrations.values()) {
|
||||
const { event, group } = registration
|
||||
const groups = consumerRegistry.get(event)
|
||||
const peersForGroup = groups?.get(group)
|
||||
peersForGroup?.delete(peerId)
|
||||
deliveryRoundRobinCursor.delete(getConsumerRegistryKey(event, group))
|
||||
if (peersForGroup?.size === 0) {
|
||||
groups?.delete(group)
|
||||
}
|
||||
if (groups?.size === 0) {
|
||||
consumerRegistry.delete(event)
|
||||
}
|
||||
|
||||
removeStickyAssignmentsFor(event, group, peerId)
|
||||
}
|
||||
|
||||
consumerKeysByPeer.delete(peerId)
|
||||
},
|
||||
listFor(input: { event: string, mode: ServerWsConsumerDeliveryConfig['mode'], group?: string }) {
|
||||
const normalizedGroup = normalizeConsumerGroup(input.mode, input.group)
|
||||
return [...consumerRegistry.get(input.event)?.get(normalizedGroup)?.values() ?? []]
|
||||
},
|
||||
select(input: {
|
||||
eventType: string
|
||||
fromPeerId: string
|
||||
delivery?: ServerWsDeliveryConfig
|
||||
candidates: ServerWsConsumerSelectionCandidate[]
|
||||
}) {
|
||||
return selectConsumerPeerId({
|
||||
...input,
|
||||
roundRobinCursor: deliveryRoundRobinCursor,
|
||||
stickyAssignments,
|
||||
})
|
||||
},
|
||||
clear() {
|
||||
consumerRegistry.clear()
|
||||
consumerKeysByPeer.clear()
|
||||
deliveryRoundRobinCursor.clear()
|
||||
stickyAssignments.clear()
|
||||
},
|
||||
}
|
||||
}
|
||||
@@ -45,7 +45,7 @@ vi.mock('crossws/server', () => ({
|
||||
plugin: vi.fn(() => ({})),
|
||||
}))
|
||||
|
||||
vi.mock('..', () => ({
|
||||
vi.mock('./index', () => ({
|
||||
normalizeLoggerConfig: () => ({
|
||||
appLogFormat: 'pretty',
|
||||
appLogLevel: 'log',
|
||||
|
||||
Reference in New Issue
Block a user