@@ -1,7 +1,7 @@
|
||||
import { fromEnv } from './env'
|
||||
|
||||
export function optionOrEnv<T extends string>(option: T | undefined, envKey: string, envDefault: T, options?: { validator: (value: string) => value is T }): T
|
||||
export function optionOrEnv<T extends string>(option: T | undefined, envKey: string, envDefault?: undefined, options?: { validator: (value: string) => value is T }): T | undefined
|
||||
export function optionOrEnv<T extends string>(option: T | undefined, envKey: string, envDefault?: undefined, options?: { validator: (value: string) => value is T }): undefined | T
|
||||
export function optionOrEnv<T extends string>(option: T | undefined, envKey: string, envDefault?: T, options?: { validator: (value: string) => value is T }): T | undefined {
|
||||
if (option !== undefined) {
|
||||
return option
|
||||
|
||||
@@ -12,22 +12,22 @@ function createInputTextEvent(
|
||||
overrides: Partial<WebSocketBaseEvent<'input:text', WebSocketEvents['input:text']>> = {},
|
||||
): WebSocketBaseEvent<'input:text', WebSocketEvents['input:text']> {
|
||||
return {
|
||||
type: 'input:text',
|
||||
data: {
|
||||
text: 'hello',
|
||||
...overrides.data,
|
||||
},
|
||||
metadata: overrides.metadata ?? {
|
||||
source: {
|
||||
kind: 'plugin',
|
||||
plugin: { id: 'discord' },
|
||||
id: 'discord-instance',
|
||||
},
|
||||
event: {
|
||||
id: 'event-1',
|
||||
},
|
||||
source: {
|
||||
id: 'discord-instance',
|
||||
kind: 'plugin',
|
||||
plugin: { id: 'discord' },
|
||||
},
|
||||
},
|
||||
route: overrides.route,
|
||||
type: 'input:text',
|
||||
}
|
||||
}
|
||||
|
||||
@@ -36,8 +36,8 @@ describe('resolveEventDelivery', () => {
|
||||
const delivery = resolveEventDelivery(createInputTextEvent())
|
||||
|
||||
expect(delivery).toEqual({
|
||||
group: 'chat-ingestion',
|
||||
mode: 'consumer-group',
|
||||
group: 'chat-ingestion',
|
||||
selection: 'first',
|
||||
})
|
||||
})
|
||||
@@ -54,8 +54,8 @@ describe('resolveEventDelivery', () => {
|
||||
}))
|
||||
|
||||
expect(delivery).toEqual({
|
||||
group: 'chat-ingestion',
|
||||
mode: 'consumer-group',
|
||||
group: 'chat-ingestion',
|
||||
required: true,
|
||||
selection: 'sticky',
|
||||
stickyKey: 'discord-dm-user-1',
|
||||
@@ -64,22 +64,23 @@ describe('resolveEventDelivery', () => {
|
||||
|
||||
it('returns explicit route delivery for events without protocol defaults', () => {
|
||||
const delivery = resolveEventDelivery({
|
||||
type: 'spark:notify',
|
||||
data: {
|
||||
destinations: ['module:character'],
|
||||
eventId: 'spark-notify-1',
|
||||
headline: 'hello',
|
||||
id: 'spark-1',
|
||||
eventId: 'spark-notify-1',
|
||||
kind: 'ping',
|
||||
urgency: 'soon',
|
||||
headline: 'hello',
|
||||
destinations: ['module:character'],
|
||||
},
|
||||
metadata: {
|
||||
event: {
|
||||
id: 'event-2',
|
||||
},
|
||||
source: {
|
||||
id: 'stage-web-instance',
|
||||
kind: 'plugin',
|
||||
plugin: { id: 'stage-web' },
|
||||
id: 'stage-web-instance',
|
||||
},
|
||||
event: {
|
||||
id: 'event-2',
|
||||
},
|
||||
},
|
||||
route: {
|
||||
@@ -88,7 +89,6 @@ describe('resolveEventDelivery', () => {
|
||||
required: true,
|
||||
},
|
||||
},
|
||||
type: 'spark:notify',
|
||||
})
|
||||
|
||||
expect(delivery).toEqual({
|
||||
@@ -101,36 +101,36 @@ describe('resolveEventDelivery', () => {
|
||||
describe('selectConsumerPeerId', () => {
|
||||
it('selects the highest-priority healthy consumer in the delivery group', () => {
|
||||
const selectedPeerId = selectConsumerPeerId({
|
||||
eventType: 'input:text',
|
||||
fromPeerId: 'discord-instance',
|
||||
delivery: {
|
||||
mode: 'consumer-group',
|
||||
group: 'chat-ingestion',
|
||||
selection: 'priority',
|
||||
},
|
||||
candidates: [
|
||||
{
|
||||
authenticated: true,
|
||||
healthy: true,
|
||||
peerId: 'stage-window-a',
|
||||
priority: 10,
|
||||
registeredAt: 2,
|
||||
},
|
||||
{
|
||||
authenticated: true,
|
||||
healthy: true,
|
||||
},
|
||||
{
|
||||
peerId: 'stage-window-b',
|
||||
priority: 20,
|
||||
registeredAt: 3,
|
||||
authenticated: true,
|
||||
healthy: true,
|
||||
},
|
||||
{
|
||||
authenticated: true,
|
||||
healthy: false,
|
||||
peerId: 'stage-window-c',
|
||||
priority: 30,
|
||||
registeredAt: 1,
|
||||
authenticated: true,
|
||||
healthy: false,
|
||||
},
|
||||
],
|
||||
delivery: {
|
||||
group: 'chat-ingestion',
|
||||
mode: 'consumer-group',
|
||||
selection: 'priority',
|
||||
},
|
||||
eventType: 'input:text',
|
||||
fromPeerId: 'discord-instance',
|
||||
})
|
||||
|
||||
expect(selectedPeerId).toBe('stage-window-b')
|
||||
@@ -140,58 +140,58 @@ describe('selectConsumerPeerId', () => {
|
||||
const stickyAssignments = new Map<string, ConsumerStickyAssignment>()
|
||||
|
||||
const firstSelectedPeerId = selectConsumerPeerId({
|
||||
candidates: [
|
||||
{
|
||||
authenticated: true,
|
||||
healthy: true,
|
||||
peerId: 'stage-window-a',
|
||||
priority: 10,
|
||||
registeredAt: 1,
|
||||
},
|
||||
{
|
||||
authenticated: true,
|
||||
healthy: true,
|
||||
peerId: 'stage-window-b',
|
||||
priority: 10,
|
||||
registeredAt: 2,
|
||||
},
|
||||
],
|
||||
eventType: 'input:text',
|
||||
fromPeerId: 'discord-instance',
|
||||
delivery: {
|
||||
group: 'chat-ingestion',
|
||||
mode: 'consumer-group',
|
||||
group: 'chat-ingestion',
|
||||
selection: 'sticky',
|
||||
stickyKey: 'discord-dm-user-1',
|
||||
},
|
||||
eventType: 'input:text',
|
||||
fromPeerId: 'discord-instance',
|
||||
candidates: [
|
||||
{
|
||||
peerId: 'stage-window-a',
|
||||
priority: 10,
|
||||
registeredAt: 1,
|
||||
authenticated: true,
|
||||
healthy: true,
|
||||
},
|
||||
{
|
||||
peerId: 'stage-window-b',
|
||||
priority: 10,
|
||||
registeredAt: 2,
|
||||
authenticated: true,
|
||||
healthy: true,
|
||||
},
|
||||
],
|
||||
stickyAssignments,
|
||||
})
|
||||
|
||||
const secondSelectedPeerId = selectConsumerPeerId({
|
||||
candidates: [
|
||||
{
|
||||
authenticated: true,
|
||||
healthy: true,
|
||||
peerId: 'stage-window-a',
|
||||
priority: 10,
|
||||
registeredAt: 1,
|
||||
},
|
||||
{
|
||||
authenticated: true,
|
||||
healthy: true,
|
||||
peerId: 'stage-window-b',
|
||||
priority: 10,
|
||||
registeredAt: 2,
|
||||
},
|
||||
],
|
||||
eventType: 'input:text',
|
||||
fromPeerId: 'discord-instance',
|
||||
delivery: {
|
||||
group: 'chat-ingestion',
|
||||
mode: 'consumer-group',
|
||||
group: 'chat-ingestion',
|
||||
selection: 'sticky',
|
||||
stickyKey: 'discord-dm-user-1',
|
||||
},
|
||||
eventType: 'input:text',
|
||||
fromPeerId: 'discord-instance',
|
||||
candidates: [
|
||||
{
|
||||
peerId: 'stage-window-a',
|
||||
priority: 10,
|
||||
registeredAt: 1,
|
||||
authenticated: true,
|
||||
healthy: true,
|
||||
},
|
||||
{
|
||||
peerId: 'stage-window-b',
|
||||
priority: 10,
|
||||
registeredAt: 2,
|
||||
authenticated: true,
|
||||
healthy: true,
|
||||
},
|
||||
],
|
||||
stickyAssignments,
|
||||
})
|
||||
|
||||
|
||||
@@ -66,26 +66,6 @@ import {
|
||||
resolveEventDelivery,
|
||||
} from './server-ws/airi/routing'
|
||||
|
||||
export interface AppOptions {
|
||||
auth?: {
|
||||
token: string
|
||||
}
|
||||
heartbeat?: {
|
||||
message?: MessageHeartbeat | string
|
||||
readTimeout?: number
|
||||
}
|
||||
instanceId?: string
|
||||
logger?: {
|
||||
app?: { format?: Format, level?: LogLevelString }
|
||||
websocket?: { format?: Format, level?: LogLevelString }
|
||||
}
|
||||
routing?: {
|
||||
allowBypass?: boolean
|
||||
middleware?: RouteMiddleware[]
|
||||
policy?: RoutingPolicy
|
||||
}
|
||||
}
|
||||
|
||||
interface AiriWsMessage {
|
||||
text: () => string
|
||||
}
|
||||
@@ -94,6 +74,89 @@ interface AiriWsPeerState {
|
||||
rawPeer: CrossWsPeer
|
||||
}
|
||||
|
||||
function airiPeerFromRaw(rawPeer: CrossWsPeer): Peer {
|
||||
// CrossWS peers expose the connection fields AIRI historically used directly
|
||||
// (`id`, `send`, `close`, `remoteAddress`, and `request`). Keep the cast in
|
||||
// this adapter boundary so protocol code below still depends on the AIRI peer
|
||||
// contract instead of the concrete transport type.
|
||||
return rawPeer as Peer
|
||||
}
|
||||
|
||||
function rawPeerFrom(wsPeer: WsPeer<AiriWsMessage, AiriWsPeerState>): Peer | undefined {
|
||||
const rawPeer = wsPeer.state?.rawPeer
|
||||
if (!rawPeer) {
|
||||
return undefined
|
||||
}
|
||||
|
||||
return airiPeerFromRaw(rawPeer)
|
||||
}
|
||||
|
||||
/**
|
||||
* Constant-time string comparison that prevents timing attacks (CWE-208).
|
||||
*
|
||||
* Compares two strings in constant time to prevent attackers from learning
|
||||
* information about the target string through timing side-channels.
|
||||
*
|
||||
* Use when:
|
||||
* - Comparing authentication tokens or secrets
|
||||
* - Any security-sensitive string comparison
|
||||
*
|
||||
* Expects:
|
||||
* - Both strings are available (no lazy evaluation)
|
||||
*
|
||||
* Returns:
|
||||
* - `true` if the strings are equal, `false` otherwise
|
||||
*/
|
||||
function timingSafeCompare(a: string, b: string): boolean {
|
||||
const bufA = Buffer.from(a)
|
||||
const bufB = Buffer.from(b)
|
||||
|
||||
// Normalize attacker-controlled input to the expected length
|
||||
// so timingSafeEqual always performs a real comparison.
|
||||
const paddedA = Buffer.alloc(bufB.length)
|
||||
|
||||
bufA.copy(
|
||||
paddedA,
|
||||
0,
|
||||
0,
|
||||
Math.min(bufA.length, bufB.length),
|
||||
)
|
||||
|
||||
return (
|
||||
timingSafeEqual(paddedA, bufB)
|
||||
&& bufA.length === bufB.length
|
||||
)
|
||||
}
|
||||
|
||||
/**
|
||||
* Sends an event to a specific peer.
|
||||
* Converts the event to JSON format before transmission.
|
||||
* @internal
|
||||
*/
|
||||
function send(peer: Peer, event: WebSocketBaseEvent<string, unknown> | string) {
|
||||
peer.send(stringifyEvent(event))
|
||||
}
|
||||
|
||||
export interface AppOptions {
|
||||
instanceId?: string
|
||||
auth?: {
|
||||
token: string
|
||||
}
|
||||
logger?: {
|
||||
app?: { level?: LogLevelString, format?: Format }
|
||||
websocket?: { level?: LogLevelString, format?: Format }
|
||||
}
|
||||
routing?: {
|
||||
middleware?: RouteMiddleware[]
|
||||
allowBypass?: boolean
|
||||
policy?: RoutingPolicy
|
||||
}
|
||||
heartbeat?: {
|
||||
readTimeout?: number
|
||||
message?: MessageHeartbeat | string
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Normalizes logger settings from explicit options and environment variables.
|
||||
*
|
||||
@@ -114,10 +177,10 @@ export function normalizeLoggerConfig(options?: AppOptions) {
|
||||
const websocketLogFormat = options?.logger?.websocket?.format || appLogFormat || Format.Pretty
|
||||
|
||||
return {
|
||||
appLogFormat,
|
||||
appLogLevel,
|
||||
websocketLogFormat,
|
||||
appLogFormat,
|
||||
websocketLogLevel,
|
||||
websocketLogFormat,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -154,7 +217,7 @@ export function setupApp(options?: AppOptions): { app: H3, closeAllPeers: () =>
|
||||
const instanceId = options?.instanceId || optionOrEnv(undefined, 'SERVER_INSTANCE_ID', nanoid())
|
||||
const authToken = optionOrEnv(options?.auth?.token, 'AUTHENTICATION_TOKEN', '')
|
||||
|
||||
const { appLogFormat, appLogLevel, websocketLogFormat, websocketLogLevel } = normalizeLoggerConfig(options)
|
||||
const { appLogLevel, appLogFormat, websocketLogLevel, websocketLogFormat } = normalizeLoggerConfig(options)
|
||||
|
||||
const appLogger = useLogg('@proj-airi/server-runtime').withLogLevel(logLevelStringToLogLevelMap[appLogLevel]).withFormat(appLogFormat)
|
||||
const logger = useLogg('@proj-airi/server-runtime:websocket').withLogLevel(logLevelStringToLogLevelMap[websocketLogLevel]).withFormat(websocketLogFormat)
|
||||
@@ -188,31 +251,31 @@ export function setupApp(options?: AppOptions): { app: H3, closeAllPeers: () =>
|
||||
}
|
||||
|
||||
broadcastToAuthenticated({
|
||||
data: { identity: peerInfo.identity, index: peerInfo.index, name: peerInfo.name },
|
||||
metadata: createEventMetadata(instanceId, parentId),
|
||||
type: 'registry:modules:health:healthy',
|
||||
data: { name: peerInfo.name, index: peerInfo.index, identity: peerInfo.identity },
|
||||
metadata: createEventMetadata(instanceId, parentId),
|
||||
})
|
||||
}
|
||||
|
||||
function broadcastPeerUnhealthy(peerInfo: AuthenticatedPeer, reason: string) {
|
||||
if (peerInfo.name && peerInfo.identity) {
|
||||
broadcastToAuthenticated({
|
||||
data: { identity: peerInfo.identity, index: peerInfo.index, name: peerInfo.name, reason },
|
||||
metadata: createEventMetadata(instanceId),
|
||||
type: 'registry:modules:health:unhealthy',
|
||||
data: { name: peerInfo.name, index: peerInfo.index, identity: peerInfo.identity, reason },
|
||||
metadata: createEventMetadata(instanceId),
|
||||
})
|
||||
}
|
||||
|
||||
for (const module of peerInfo.extensionModules?.values() ?? []) {
|
||||
broadcastToAuthenticated({
|
||||
data: { identity: module.identity, name: module.name, reason },
|
||||
metadata: createEventMetadata(instanceId),
|
||||
type: 'registry:modules:health:unhealthy',
|
||||
data: { name: module.name, identity: module.identity, reason },
|
||||
metadata: createEventMetadata(instanceId),
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
function markPeerAlive(peerInfo: AuthenticatedPeer, options?: { logMessage?: string, parentId?: string }) {
|
||||
function markPeerAlive(peerInfo: AuthenticatedPeer, options?: { parentId?: string, logMessage?: string }) {
|
||||
peerInfo.lastHeartbeatAt = Date.now()
|
||||
peerInfo.missedHeartbeats = 0
|
||||
|
||||
@@ -278,11 +341,11 @@ export function setupApp(options?: AppOptions): { app: H3, closeAllPeers: () =>
|
||||
}
|
||||
|
||||
function registerConsumer(peerId: string, event: string, mode: ReturnType<typeof normalizeConsumerMode>, group?: string, priority?: number) {
|
||||
consumers.register({ event, group, mode, peerId, priority })
|
||||
consumers.register({ peerId, event, mode, group, priority })
|
||||
}
|
||||
|
||||
function unregisterConsumer(peerId: string, event: string, mode: ReturnType<typeof normalizeConsumerMode>, group?: string) {
|
||||
consumers.unregister({ event, group, mode, peerId })
|
||||
consumers.unregister({ peerId, event, mode, group })
|
||||
}
|
||||
|
||||
function unregisterPeerConsumers(peerId: string) {
|
||||
@@ -295,20 +358,20 @@ export function setupApp(options?: AppOptions): { app: H3, closeAllPeers: () =>
|
||||
}
|
||||
|
||||
const selectedPeerId = consumers.select({
|
||||
eventType: event.type,
|
||||
fromPeerId,
|
||||
delivery,
|
||||
candidates: consumers.listFor({
|
||||
event: event.type,
|
||||
group: delivery?.group,
|
||||
mode: delivery?.mode,
|
||||
group: delivery?.group,
|
||||
}).map(entry => ({
|
||||
authenticated: Boolean(peers.get(entry.peerId)?.authenticated),
|
||||
healthy: peers.get(entry.peerId)?.healthy,
|
||||
peerId: entry.peerId,
|
||||
priority: entry.priority,
|
||||
registeredAt: entry.registeredAt,
|
||||
authenticated: Boolean(peers.get(entry.peerId)?.authenticated),
|
||||
healthy: peers.get(entry.peerId)?.healthy,
|
||||
})),
|
||||
delivery,
|
||||
eventType: event.type,
|
||||
fromPeerId,
|
||||
})
|
||||
|
||||
if (!selectedPeerId) {
|
||||
@@ -341,9 +404,9 @@ export function setupApp(options?: AppOptions): { app: H3, closeAllPeers: () =>
|
||||
// broadcast extension:module:de-announced to all authenticated peers
|
||||
if (peerInfo.identity) {
|
||||
broadcastToAuthenticated({
|
||||
data: { identity: peerInfo.identity, name: peerInfo.name, possibleEvents: [], reason: options?.reason },
|
||||
metadata: createEventMetadata(instanceId),
|
||||
type: 'extension:module:de-announced',
|
||||
data: { name: peerInfo.name, identity: peerInfo.identity, possibleEvents: [], reason: options?.reason },
|
||||
metadata: createEventMetadata(instanceId),
|
||||
})
|
||||
}
|
||||
|
||||
@@ -369,9 +432,9 @@ export function setupApp(options?: AppOptions): { app: H3, closeAllPeers: () =>
|
||||
|
||||
peerInfo.extensionModules?.delete(module.identity.id)
|
||||
broadcastToAuthenticated({
|
||||
data: { identity: module.identity, name: module.name, possibleEvents: [], reason },
|
||||
metadata: createEventMetadata(instanceId),
|
||||
type: 'extension:module:de-announced',
|
||||
data: { name: module.name, identity: module.identity, possibleEvents: [], reason },
|
||||
metadata: createEventMetadata(instanceId),
|
||||
})
|
||||
}
|
||||
|
||||
@@ -397,15 +460,15 @@ export function setupApp(options?: AppOptions): { app: H3, closeAllPeers: () =>
|
||||
const legacyModules = Array.from(peers.values())
|
||||
.filter(peerInfo => peerInfo.name && peerInfo.identity)
|
||||
.map(peerInfo => ({
|
||||
identity: peerInfo.identity!,
|
||||
index: peerInfo.index,
|
||||
name: peerInfo.name,
|
||||
index: peerInfo.index,
|
||||
identity: peerInfo.identity!,
|
||||
}))
|
||||
|
||||
const extensionModules = Array.from(peers.values()).flatMap(peerInfo =>
|
||||
Array.from(peerInfo.extensionModules?.values() ?? []).map(module => ({
|
||||
identity: module.identity,
|
||||
name: module.name,
|
||||
identity: module.identity,
|
||||
})),
|
||||
)
|
||||
|
||||
@@ -432,9 +495,9 @@ export function setupApp(options?: AppOptions): { app: H3, closeAllPeers: () =>
|
||||
// === Broadcasting & Registry Synchronization ===
|
||||
function sendRegistrySync(peer: Peer, parentId?: string) {
|
||||
send(peer, {
|
||||
type: 'registry:modules:sync',
|
||||
data: { modules: listKnownModules() },
|
||||
metadata: createEventMetadata(instanceId, parentId),
|
||||
type: 'registry:modules:sync',
|
||||
})
|
||||
}
|
||||
|
||||
@@ -457,14 +520,14 @@ export function setupApp(options?: AppOptions): { app: H3, closeAllPeers: () =>
|
||||
// === WebSocket Server Handlers ===
|
||||
// Handles AIRI peer lifecycle: open, message, error, close.
|
||||
const wsServer = createWsServer<AiriWsMessage, AiriWsPeerState>({
|
||||
peers: {
|
||||
unhealthyTimeout: heartbeatTtlMs,
|
||||
closeTimeout: heartbeatTtlMs * 2,
|
||||
},
|
||||
heartbeat: {
|
||||
interval: healthCheckIntervalMs,
|
||||
timeout: heartbeatTtlMs,
|
||||
},
|
||||
peers: {
|
||||
closeTimeout: heartbeatTtlMs * 2,
|
||||
unhealthyTimeout: heartbeatTtlMs,
|
||||
},
|
||||
})
|
||||
wsServer.onPeerOpen(({ peer: wsPeer }) => {
|
||||
const peer = rawPeerFrom(wsPeer)
|
||||
@@ -472,18 +535,18 @@ export function setupApp(options?: AppOptions): { app: H3, closeAllPeers: () =>
|
||||
return
|
||||
|
||||
if (authToken) {
|
||||
peers.set(peer.id, { authenticated: false, lastHeartbeatAt: Date.now(), name: '', peer })
|
||||
peers.set(peer.id, { peer, authenticated: false, name: '', lastHeartbeatAt: Date.now() })
|
||||
}
|
||||
else {
|
||||
send(peer, RESPONSES.authenticated())
|
||||
peers.set(peer.id, { authenticated: true, lastHeartbeatAt: Date.now(), name: '', peer })
|
||||
peers.set(peer.id, { peer, authenticated: true, name: '', lastHeartbeatAt: Date.now() })
|
||||
sendRegistrySync(peer)
|
||||
}
|
||||
|
||||
logger.withFields({ activePeers: peers.size, peer: peer.id }).log('connected')
|
||||
logger.withFields({ peer: peer.id, activePeers: peers.size }).log('connected')
|
||||
})
|
||||
|
||||
wsServer.onMessage(({ message, peer: wsPeer }) => {
|
||||
wsServer.onMessage(({ peer: wsPeer, message }) => {
|
||||
const peer = rawPeerFrom(wsPeer)
|
||||
if (!peer)
|
||||
return
|
||||
@@ -536,114 +599,19 @@ export function setupApp(options?: AppOptions): { app: H3, closeAllPeers: () =>
|
||||
}
|
||||
|
||||
switch (event.type) {
|
||||
case 'extension:announce': {
|
||||
const p = peers.get(peer.id)
|
||||
if (!p) {
|
||||
return
|
||||
}
|
||||
|
||||
if (authToken && !p.authenticated) {
|
||||
send(peer, RESPONSES.error(ServerErrorMessages.mustAuthenticateBeforeAnnouncing))
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
if (!isExtensionIdentity(event.data.identity)) {
|
||||
send(peer, RESPONSES.error(ServerErrorMessages.moduleAnnounceIdentityInvalid))
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
p.extensionIdentity = event.data.identity
|
||||
|
||||
send(peer, {
|
||||
data: event.data,
|
||||
metadata: createEventMetadata(instanceId, event.metadata?.event.id),
|
||||
type: 'extension:announced',
|
||||
})
|
||||
|
||||
for (const other of peers.values()) {
|
||||
if (other.authenticated && !(other.peer.id === peer.id)) {
|
||||
send(other.peer, {
|
||||
data: event.data,
|
||||
metadata: createEventMetadata(instanceId, event.metadata?.event.id),
|
||||
type: 'extension:announced',
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
case 'extension:authenticate': {
|
||||
const clientToken = typeof event.data.token === 'string' ? event.data.token : ''
|
||||
if (authToken && !timingSafeCompare(clientToken, authToken)) {
|
||||
logger.withFields({ peer: peer.id, peerRemote: peer.remoteAddress, peerRequest: peer.request?.url }).log('extension authentication failed')
|
||||
send(peer, RESPONSES.error(ServerErrorMessages.invalidToken, event.metadata?.event.id))
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
case 'transport:connection:heartbeat': {
|
||||
const p = peers.get(peer.id)
|
||||
if (p) {
|
||||
p.authenticated = true
|
||||
p.extensionIdentity = event.data.identity
|
||||
markPeerAlive(p, {
|
||||
parentId: event.metadata?.event.id,
|
||||
logMessage: 'heartbeat recovered, marking healthy',
|
||||
})
|
||||
|
||||
// recover from unhealthy → healthy
|
||||
}
|
||||
|
||||
send(peer, RESPONSES.extensionAuthenticated(event.data.identity, event.metadata?.event.id))
|
||||
sendRegistrySync(peer, event.metadata?.event.id)
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
case 'extension:module:announce': {
|
||||
const p = peers.get(peer.id)
|
||||
if (!p) {
|
||||
return
|
||||
}
|
||||
|
||||
if (authToken && !p.authenticated) {
|
||||
send(peer, RESPONSES.error(ServerErrorMessages.mustAuthenticateBeforeAnnouncing))
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
const { identity, name } = event.data
|
||||
if (!name || typeof name !== 'string') {
|
||||
send(peer, RESPONSES.error(ServerErrorMessages.moduleAnnounceNameInvalid))
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
if (!isExtensionModuleIdentity(identity)) {
|
||||
send(peer, RESPONSES.error(ServerErrorMessages.moduleAnnounceIdentityInvalid))
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
if (p.extensionIdentity && identity.extension.id !== p.extensionIdentity.id) {
|
||||
send(peer, RESPONSES.error(ServerErrorMessages.moduleAnnounceIdentityInvalid))
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
p.extensionIdentity = identity.extension
|
||||
registerExtensionModulePeer(p, { identity, name })
|
||||
|
||||
send(peer, {
|
||||
data: event.data,
|
||||
metadata: createEventMetadata(instanceId, event.metadata?.event.id),
|
||||
type: 'extension:module:announced',
|
||||
})
|
||||
|
||||
for (const other of peers.values()) {
|
||||
if (other.authenticated && !(other.peer.id === peer.id)) {
|
||||
send(other.peer, {
|
||||
data: event.data,
|
||||
metadata: createEventMetadata(instanceId, event.metadata?.event.id),
|
||||
type: 'extension:module:announced',
|
||||
})
|
||||
}
|
||||
if (event.data.kind === MessageHeartbeatKind.Ping) {
|
||||
send(peer, RESPONSES.heartbeat(MessageHeartbeatKind.Pong, heartbeatMessage, event.metadata?.event.id))
|
||||
}
|
||||
|
||||
return
|
||||
@@ -669,57 +637,6 @@ export function setupApp(options?: AppOptions): { app: H3, closeAllPeers: () =>
|
||||
return
|
||||
}
|
||||
|
||||
case 'module:consumer:register': {
|
||||
const p = peers.get(peer.id)
|
||||
if (!p?.authenticated) {
|
||||
send(peer, RESPONSES.notAuthenticated(event.metadata?.event.id))
|
||||
return
|
||||
}
|
||||
|
||||
const data = event.data as {
|
||||
event?: string
|
||||
group?: string
|
||||
mode?: 'consumer' | 'consumer-group'
|
||||
priority?: number
|
||||
}
|
||||
|
||||
if (!data.event || typeof data.event !== 'string') {
|
||||
send(peer, RESPONSES.error(ServerErrorMessages.moduleConsumerEventInvalid, event.metadata?.event.id))
|
||||
return
|
||||
}
|
||||
|
||||
registerConsumer(
|
||||
peer.id,
|
||||
data.event,
|
||||
normalizeConsumerMode(data.mode, data.group),
|
||||
data.group,
|
||||
normalizeConsumerPriority(data.priority),
|
||||
)
|
||||
return
|
||||
}
|
||||
|
||||
case 'module:consumer:unregister': {
|
||||
const p = peers.get(peer.id)
|
||||
if (!p?.authenticated) {
|
||||
send(peer, RESPONSES.notAuthenticated(event.metadata?.event.id))
|
||||
return
|
||||
}
|
||||
|
||||
const data = event.data as {
|
||||
event?: string
|
||||
group?: string
|
||||
mode?: 'consumer' | 'consumer-group'
|
||||
}
|
||||
|
||||
if (!data.event || typeof data.event !== 'string') {
|
||||
send(peer, RESPONSES.error(ServerErrorMessages.moduleConsumerEventInvalid, event.metadata?.event.id))
|
||||
return
|
||||
}
|
||||
|
||||
unregisterConsumer(peer.id, data.event, normalizeConsumerMode(data.mode, data.group), data.group)
|
||||
return
|
||||
}
|
||||
|
||||
case 'peer:authenticate': {
|
||||
const clientToken = typeof event.data.token === 'string' ? event.data.token : ''
|
||||
if (authToken && !timingSafeCompare(clientToken, authToken)) {
|
||||
@@ -744,19 +661,114 @@ export function setupApp(options?: AppOptions): { app: H3, closeAllPeers: () =>
|
||||
return
|
||||
}
|
||||
|
||||
case 'transport:connection:heartbeat': {
|
||||
const p = peers.get(peer.id)
|
||||
if (p) {
|
||||
markPeerAlive(p, {
|
||||
logMessage: 'heartbeat recovered, marking healthy',
|
||||
parentId: event.metadata?.event.id,
|
||||
})
|
||||
case 'extension:authenticate': {
|
||||
const clientToken = typeof event.data.token === 'string' ? event.data.token : ''
|
||||
if (authToken && !timingSafeCompare(clientToken, authToken)) {
|
||||
logger.withFields({ peer: peer.id, peerRemote: peer.remoteAddress, peerRequest: peer.request?.url }).log('extension authentication failed')
|
||||
send(peer, RESPONSES.error(ServerErrorMessages.invalidToken, event.metadata?.event.id))
|
||||
|
||||
// recover from unhealthy → healthy
|
||||
return
|
||||
}
|
||||
|
||||
if (event.data.kind === MessageHeartbeatKind.Ping) {
|
||||
send(peer, RESPONSES.heartbeat(MessageHeartbeatKind.Pong, heartbeatMessage, event.metadata?.event.id))
|
||||
const p = peers.get(peer.id)
|
||||
if (p) {
|
||||
p.authenticated = true
|
||||
p.extensionIdentity = event.data.identity
|
||||
}
|
||||
|
||||
send(peer, RESPONSES.extensionAuthenticated(event.data.identity, event.metadata?.event.id))
|
||||
sendRegistrySync(peer, event.metadata?.event.id)
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
case 'extension:announce': {
|
||||
const p = peers.get(peer.id)
|
||||
if (!p) {
|
||||
return
|
||||
}
|
||||
|
||||
if (authToken && !p.authenticated) {
|
||||
send(peer, RESPONSES.error(ServerErrorMessages.mustAuthenticateBeforeAnnouncing))
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
if (!isExtensionIdentity(event.data.identity)) {
|
||||
send(peer, RESPONSES.error(ServerErrorMessages.moduleAnnounceIdentityInvalid))
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
p.extensionIdentity = event.data.identity
|
||||
|
||||
send(peer, {
|
||||
type: 'extension:announced',
|
||||
data: event.data,
|
||||
metadata: createEventMetadata(instanceId, event.metadata?.event.id),
|
||||
})
|
||||
|
||||
for (const other of peers.values()) {
|
||||
if (other.authenticated && !(other.peer.id === peer.id)) {
|
||||
send(other.peer, {
|
||||
type: 'extension:announced',
|
||||
data: event.data,
|
||||
metadata: createEventMetadata(instanceId, event.metadata?.event.id),
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
case 'extension:module:announce': {
|
||||
const p = peers.get(peer.id)
|
||||
if (!p) {
|
||||
return
|
||||
}
|
||||
|
||||
if (authToken && !p.authenticated) {
|
||||
send(peer, RESPONSES.error(ServerErrorMessages.mustAuthenticateBeforeAnnouncing))
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
const { name, identity } = event.data
|
||||
if (!name || typeof name !== 'string') {
|
||||
send(peer, RESPONSES.error(ServerErrorMessages.moduleAnnounceNameInvalid))
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
if (!isExtensionModuleIdentity(identity)) {
|
||||
send(peer, RESPONSES.error(ServerErrorMessages.moduleAnnounceIdentityInvalid))
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
if (p.extensionIdentity && identity.extension.id !== p.extensionIdentity.id) {
|
||||
send(peer, RESPONSES.error(ServerErrorMessages.moduleAnnounceIdentityInvalid))
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
p.extensionIdentity = identity.extension
|
||||
registerExtensionModulePeer(p, { name, identity })
|
||||
|
||||
send(peer, {
|
||||
type: 'extension:module:announced',
|
||||
data: event.data,
|
||||
metadata: createEventMetadata(instanceId, event.metadata?.event.id),
|
||||
})
|
||||
|
||||
for (const other of peers.values()) {
|
||||
if (other.authenticated && !(other.peer.id === peer.id)) {
|
||||
send(other.peer, {
|
||||
type: 'extension:module:announced',
|
||||
data: event.data,
|
||||
metadata: createEventMetadata(instanceId, event.metadata?.event.id),
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
return
|
||||
@@ -764,10 +776,10 @@ export function setupApp(options?: AppOptions): { app: H3, closeAllPeers: () =>
|
||||
|
||||
case 'ui:configure': {
|
||||
const data = event.data as {
|
||||
config?: Record<string, unknown>
|
||||
identity?: MetadataEventSource
|
||||
moduleIndex?: number
|
||||
moduleName?: string
|
||||
moduleIndex?: number
|
||||
identity?: MetadataEventSource
|
||||
config?: Record<string, unknown>
|
||||
}
|
||||
const moduleName = data.moduleName ?? (isExtensionModuleIdentity(data.identity) ? data.identity.id : '') ?? ''
|
||||
const moduleIndex = data.moduleIndex
|
||||
@@ -789,10 +801,10 @@ export function setupApp(options?: AppOptions): { app: H3, closeAllPeers: () =>
|
||||
const target = findModulePeer(moduleName, moduleIndex, data.identity)
|
||||
if (target) {
|
||||
send(target.peer, {
|
||||
type: 'module:configure',
|
||||
data: { config: config || {} },
|
||||
// NOTICE: this will forward the original event metadata as-is
|
||||
metadata: event.metadata,
|
||||
type: 'module:configure',
|
||||
})
|
||||
}
|
||||
else {
|
||||
@@ -801,6 +813,57 @@ export function setupApp(options?: AppOptions): { app: H3, closeAllPeers: () =>
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
case 'module:consumer:register': {
|
||||
const p = peers.get(peer.id)
|
||||
if (!p?.authenticated) {
|
||||
send(peer, RESPONSES.notAuthenticated(event.metadata?.event.id))
|
||||
return
|
||||
}
|
||||
|
||||
const data = event.data as {
|
||||
event?: string
|
||||
mode?: 'consumer' | 'consumer-group'
|
||||
group?: string
|
||||
priority?: number
|
||||
}
|
||||
|
||||
if (!data.event || typeof data.event !== 'string') {
|
||||
send(peer, RESPONSES.error(ServerErrorMessages.moduleConsumerEventInvalid, event.metadata?.event.id))
|
||||
return
|
||||
}
|
||||
|
||||
registerConsumer(
|
||||
peer.id,
|
||||
data.event,
|
||||
normalizeConsumerMode(data.mode, data.group),
|
||||
data.group,
|
||||
normalizeConsumerPriority(data.priority),
|
||||
)
|
||||
return
|
||||
}
|
||||
|
||||
case 'module:consumer:unregister': {
|
||||
const p = peers.get(peer.id)
|
||||
if (!p?.authenticated) {
|
||||
send(peer, RESPONSES.notAuthenticated(event.metadata?.event.id))
|
||||
return
|
||||
}
|
||||
|
||||
const data = event.data as {
|
||||
event?: string
|
||||
mode?: 'consumer' | 'consumer-group'
|
||||
group?: string
|
||||
}
|
||||
|
||||
if (!data.event || typeof data.event !== 'string') {
|
||||
send(peer, RESPONSES.error(ServerErrorMessages.moduleConsumerEventInvalid, event.metadata?.event.id))
|
||||
return
|
||||
}
|
||||
|
||||
unregisterConsumer(peer.id, data.event, normalizeConsumerMode(data.mode, data.group), data.group)
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
// default case
|
||||
@@ -819,22 +882,22 @@ export function setupApp(options?: AppOptions): { app: H3, closeAllPeers: () =>
|
||||
const delivery = shouldBypass ? undefined : resolveEventDelivery(event)
|
||||
const effectiveRoutingMiddleware = shouldBypass ? [] : routingMiddleware
|
||||
const decision = forEachEventMiddlewares({
|
||||
destinations,
|
||||
event,
|
||||
fromPeer: p,
|
||||
middleware: effectiveRoutingMiddleware,
|
||||
peers,
|
||||
destinations,
|
||||
middleware: effectiveRoutingMiddleware,
|
||||
})
|
||||
|
||||
if (decision?.type === 'drop') {
|
||||
logger.withFields({ event, peer: peer.id, peerName: p.name }).debug('routing dropped event')
|
||||
logger.withFields({ peer: peer.id, peerName: p.name, event }).debug('routing dropped event')
|
||||
return
|
||||
}
|
||||
|
||||
const selectedConsumer = selectConsumer(event, peer.id, delivery)
|
||||
if (delivery && (delivery.mode === 'consumer' || delivery.mode === 'consumer-group')) {
|
||||
if (!selectedConsumer) {
|
||||
logger.withFields({ delivery, event, peer: peer.id, peerName: p.name }).warn('no consumer registered for event delivery')
|
||||
logger.withFields({ peer: peer.id, peerName: p.name, event, delivery }).warn('no consumer registered for event delivery')
|
||||
if (delivery.required) {
|
||||
send(peer, RESPONSES.error(ServerErrorMessages.noConsumerRegistered, event.metadata?.event.id))
|
||||
}
|
||||
@@ -843,24 +906,24 @@ export function setupApp(options?: AppOptions): { app: H3, closeAllPeers: () =>
|
||||
|
||||
try {
|
||||
logger.withFields({
|
||||
delivery,
|
||||
event,
|
||||
fromPeer: peer.id,
|
||||
fromPeerName: p.name,
|
||||
toPeer: selectedConsumer.peer.id,
|
||||
toPeerName: selectedConsumer.name,
|
||||
event,
|
||||
delivery,
|
||||
}).debug('sending event to selected consumer')
|
||||
|
||||
selectedConsumer.peer.send(payload)
|
||||
}
|
||||
catch (err) {
|
||||
logger.withFields({
|
||||
delivery,
|
||||
event,
|
||||
fromPeer: peer.id,
|
||||
fromPeerName: p.name,
|
||||
toPeer: selectedConsumer.peer.id,
|
||||
toPeerName: selectedConsumer.name,
|
||||
event,
|
||||
delivery,
|
||||
}).withError(err).error('failed to send event to selected consumer, removing peer')
|
||||
|
||||
removeFailedPeer(selectedConsumer, 'consumer send failed')
|
||||
@@ -871,16 +934,16 @@ export function setupApp(options?: AppOptions): { app: H3, closeAllPeers: () =>
|
||||
const targetIds = decision?.type === 'targets' ? decision.targetIds : undefined
|
||||
const shouldBroadcast = decision?.type === 'broadcast' || !targetIds
|
||||
|
||||
logger.withFields({ event, peer: peer.id, peerName: p.name }).debug('broadcasting event to peers')
|
||||
logger.withFields({ peer: peer.id, peerName: p.name, event }).debug('broadcasting event to peers')
|
||||
|
||||
for (const [id, other] of peers.entries()) {
|
||||
if (id === peer.id) {
|
||||
logger.withFields({ event, peer: peer.id, peerName: p.name }).debug('not sending event to self')
|
||||
logger.withFields({ peer: peer.id, peerName: p.name, event }).debug('not sending event to self')
|
||||
continue
|
||||
}
|
||||
|
||||
if (!other.authenticated) {
|
||||
logger.withFields({ event, fromPeer: peer.id, toPeer: other.peer.id, toPeerName: other.name }).debug('not sending event to unauthenticated peer')
|
||||
logger.withFields({ fromPeer: peer.id, toPeer: other.peer.id, toPeerName: other.name, event }).debug('not sending event to unauthenticated peer')
|
||||
continue
|
||||
}
|
||||
|
||||
@@ -893,11 +956,11 @@ export function setupApp(options?: AppOptions): { app: H3, closeAllPeers: () =>
|
||||
}
|
||||
|
||||
try {
|
||||
logger.withFields({ event, fromPeer: peer.id, fromPeerName: p.name, toPeer: other.peer.id, toPeerName: other.name }).debug('sending event to peer')
|
||||
logger.withFields({ fromPeer: peer.id, fromPeerName: p.name, toPeer: other.peer.id, toPeerName: other.name, event }).debug('sending event to peer')
|
||||
other.peer.send(payload)
|
||||
}
|
||||
catch (err) {
|
||||
logger.withFields({ event, fromPeer: peer.id, fromPeerName: p.name, toPeer: other.peer.id, toPeerName: other.name }).withError(err).error('failed to send event to peer, removing peer')
|
||||
logger.withFields({ fromPeer: peer.id, fromPeerName: p.name, toPeer: other.peer.id, toPeerName: other.name, event }).withError(err).error('failed to send event to peer, removing peer')
|
||||
logger.withFields({ peer: peer.id, peerName: other.name }).debug('removing closed peer')
|
||||
removeFailedPeer(other, 'send failed')
|
||||
}
|
||||
@@ -938,25 +1001,25 @@ export function setupApp(options?: AppOptions): { app: H3, closeAllPeers: () =>
|
||||
}
|
||||
|
||||
logger.withFields({
|
||||
activePeers: peers.size,
|
||||
peer: peer.id,
|
||||
peerRemote: peer.remoteAddress,
|
||||
details,
|
||||
closeCode,
|
||||
closeReason,
|
||||
closeWasClean,
|
||||
details,
|
||||
healthCheckIntervalMs,
|
||||
activePeers: peers.size,
|
||||
peerAuthenticated: p?.authenticated,
|
||||
peerName,
|
||||
peerIndex,
|
||||
peerHealthy,
|
||||
peerMissedHeartbeats: p?.missedHeartbeats,
|
||||
peerSilentFor,
|
||||
heartbeatLastSeenAt,
|
||||
heartbeatSilentForMs,
|
||||
heartbeatTtlMs,
|
||||
healthCheckIntervalMs,
|
||||
likelyHeartbeatExpiry,
|
||||
likelySilentNetworkClose,
|
||||
peer: peer.id,
|
||||
peerAuthenticated: p?.authenticated,
|
||||
peerHealthy,
|
||||
peerIndex,
|
||||
peerMissedHeartbeats: p?.missedHeartbeats,
|
||||
peerName,
|
||||
peerRemote: peer.remoteAddress,
|
||||
peerSilentFor,
|
||||
}).log('closed')
|
||||
}
|
||||
|
||||
@@ -970,7 +1033,7 @@ export function setupApp(options?: AppOptions): { app: H3, closeAllPeers: () =>
|
||||
handlePeerClose(peerInfo.peer, { reason })
|
||||
}
|
||||
|
||||
wsServer.onPeerClose(({ details, peerId }) => {
|
||||
wsServer.onPeerClose(({ peerId, details }) => {
|
||||
const peerInfo = peers.get(peerId)
|
||||
if (!peerInfo) {
|
||||
return
|
||||
@@ -979,7 +1042,7 @@ export function setupApp(options?: AppOptions): { app: H3, closeAllPeers: () =>
|
||||
handlePeerClose(peerInfo.peer, details)
|
||||
})
|
||||
|
||||
wsServer.onPeerHealthChange(({ healthy, peer, silentFor }) => {
|
||||
wsServer.onPeerHealthChange(({ peer, healthy, silentFor }) => {
|
||||
const peerInfo = peers.get(peer.id)
|
||||
if (!peerInfo) {
|
||||
return
|
||||
@@ -1038,15 +1101,15 @@ export function setupApp(options?: AppOptions): { app: H3, closeAllPeers: () =>
|
||||
}
|
||||
|
||||
app.get('/ws', toH3Handler(wsServer, {
|
||||
error({ error, peer, rawPeer }) {
|
||||
handlePeerError(peer ? rawPeerFrom(peer) ?? airiPeerFromRaw(rawPeer) : airiPeerFromRaw(rawPeer), error)
|
||||
},
|
||||
readMessage(message: CrossWsMessage) {
|
||||
return { text: () => message.text() }
|
||||
},
|
||||
state(rawPeer: CrossWsPeer) {
|
||||
return { rawPeer }
|
||||
},
|
||||
error({ peer, rawPeer, error }) {
|
||||
handlePeerError(peer ? rawPeerFrom(peer) ?? airiPeerFromRaw(rawPeer) : airiPeerFromRaw(rawPeer), error)
|
||||
},
|
||||
}))
|
||||
|
||||
function closeAllPeers() {
|
||||
@@ -1108,66 +1171,3 @@ export function setupApp(options?: AppOptions): { app: H3, closeAllPeers: () =>
|
||||
dispose,
|
||||
}
|
||||
}
|
||||
|
||||
function airiPeerFromRaw(rawPeer: CrossWsPeer): Peer {
|
||||
// CrossWS peers expose the connection fields AIRI historically used directly
|
||||
// (`id`, `send`, `close`, `remoteAddress`, and `request`). Keep the cast in
|
||||
// this adapter boundary so protocol code below still depends on the AIRI peer
|
||||
// contract instead of the concrete transport type.
|
||||
return rawPeer as Peer
|
||||
}
|
||||
|
||||
function rawPeerFrom(wsPeer: WsPeer<AiriWsMessage, AiriWsPeerState>): Peer | undefined {
|
||||
const rawPeer = wsPeer.state?.rawPeer
|
||||
if (!rawPeer) {
|
||||
return undefined
|
||||
}
|
||||
|
||||
return airiPeerFromRaw(rawPeer)
|
||||
}
|
||||
|
||||
/**
|
||||
* Sends an event to a specific peer.
|
||||
* Converts the event to JSON format before transmission.
|
||||
* @internal
|
||||
*/
|
||||
function send(peer: Peer, event: string | WebSocketBaseEvent<string, unknown>) {
|
||||
peer.send(stringifyEvent(event))
|
||||
}
|
||||
|
||||
/**
|
||||
* Constant-time string comparison that prevents timing attacks (CWE-208).
|
||||
*
|
||||
* Compares two strings in constant time to prevent attackers from learning
|
||||
* information about the target string through timing side-channels.
|
||||
*
|
||||
* Use when:
|
||||
* - Comparing authentication tokens or secrets
|
||||
* - Any security-sensitive string comparison
|
||||
*
|
||||
* Expects:
|
||||
* - Both strings are available (no lazy evaluation)
|
||||
*
|
||||
* Returns:
|
||||
* - `true` if the strings are equal, `false` otherwise
|
||||
*/
|
||||
function timingSafeCompare(a: string, b: string): boolean {
|
||||
const bufA = Buffer.from(a)
|
||||
const bufB = Buffer.from(b)
|
||||
|
||||
// Normalize attacker-controlled input to the expected length
|
||||
// so timingSafeEqual always performs a real comparison.
|
||||
const paddedA = Buffer.alloc(bufB.length)
|
||||
|
||||
bufA.copy(
|
||||
paddedA,
|
||||
0,
|
||||
0,
|
||||
Math.min(bufA.length, bufB.length),
|
||||
)
|
||||
|
||||
return (
|
||||
timingSafeEqual(paddedA, bufB)
|
||||
&& bufA.length === bufB.length
|
||||
)
|
||||
}
|
||||
|
||||
@@ -7,78 +7,78 @@ import { describe, expect, it } from 'vitest'
|
||||
import { collectDestinations, createPolicyMiddleware, isDevtoolsPeer, matchesDestinations } from './route'
|
||||
import { matchesLabelSelector, matchesLabelSelectors, matchesRouteExpression } from './route/match-expression'
|
||||
|
||||
function createPeer(options: {
|
||||
id: string
|
||||
name: string
|
||||
peerIds?: string[]
|
||||
extensionLabels?: Record<string, string>
|
||||
extension?: string
|
||||
instanceId?: string
|
||||
labels?: Record<string, string>
|
||||
authenticated?: boolean
|
||||
}): AuthenticatedPeer {
|
||||
return {
|
||||
peer: {
|
||||
id: options.id,
|
||||
send: () => 0,
|
||||
request: { url: 'http://localhost', headers: new Headers() },
|
||||
remoteAddress: '127.0.0.1',
|
||||
},
|
||||
authenticated: options.authenticated ?? true,
|
||||
peerIds: options.peerIds ? new Set(options.peerIds) : undefined,
|
||||
name: options.name,
|
||||
identity: options.extension && options.instanceId
|
||||
? { id: options.instanceId, extension: { id: options.extension }, labels: options.labels }
|
||||
: undefined,
|
||||
extensionIdentity: options.extensionLabels
|
||||
? { id: options.name, sessionId: `${options.id}-session`, labels: options.extensionLabels }
|
||||
: undefined,
|
||||
}
|
||||
}
|
||||
|
||||
function createExtensionModulePeer(): AuthenticatedPeer {
|
||||
const peer = createPeer({
|
||||
extension: 'airi-extension-chess',
|
||||
id: 'peer-extension',
|
||||
instanceId: 'extension-session-1',
|
||||
name: 'airi-extension-chess',
|
||||
extension: 'airi-extension-chess',
|
||||
instanceId: 'extension-session-1',
|
||||
})
|
||||
|
||||
peer.extensionModules = new Map([
|
||||
['chess-gamelet', {
|
||||
name: 'character',
|
||||
identity: {
|
||||
id: 'chess-gamelet',
|
||||
extension: {
|
||||
id: 'airi-extension-chess',
|
||||
sessionId: 'extension-session-1',
|
||||
},
|
||||
id: 'chess-gamelet',
|
||||
},
|
||||
name: 'character',
|
||||
}],
|
||||
])
|
||||
|
||||
return peer
|
||||
}
|
||||
|
||||
function createPeer(options: {
|
||||
authenticated?: boolean
|
||||
extension?: string
|
||||
extensionLabels?: Record<string, string>
|
||||
id: string
|
||||
instanceId?: string
|
||||
labels?: Record<string, string>
|
||||
name: string
|
||||
peerIds?: string[]
|
||||
}): AuthenticatedPeer {
|
||||
return {
|
||||
authenticated: options.authenticated ?? true,
|
||||
extensionIdentity: options.extensionLabels
|
||||
? { id: options.name, labels: options.extensionLabels, sessionId: `${options.id}-session` }
|
||||
: undefined,
|
||||
identity: options.extension && options.instanceId
|
||||
? { extension: { id: options.extension }, id: options.instanceId, labels: options.labels }
|
||||
: undefined,
|
||||
name: options.name,
|
||||
peer: {
|
||||
id: options.id,
|
||||
remoteAddress: '127.0.0.1',
|
||||
request: { headers: new Headers(), url: 'http://localhost' },
|
||||
send: () => 0,
|
||||
},
|
||||
peerIds: options.peerIds ? new Set(options.peerIds) : undefined,
|
||||
}
|
||||
}
|
||||
|
||||
function createSparkNotifyEvent(overrides: Partial<WebSocketEventOf<'spark:notify'>> = {}): WebSocketBaseEvent<'spark:notify', WebSocketEvents['spark:notify'], any> {
|
||||
const data: WebSocketEvents['spark:notify'] = {
|
||||
destinations: ['module:character'],
|
||||
eventId: 'spark-1',
|
||||
headline: 'hello',
|
||||
id: 'evt-1',
|
||||
eventId: 'spark-1',
|
||||
kind: 'ping',
|
||||
urgency: 'soon',
|
||||
headline: 'hello',
|
||||
destinations: ['module:character'],
|
||||
...overrides.data,
|
||||
}
|
||||
|
||||
return {
|
||||
type: 'spark:notify',
|
||||
data,
|
||||
metadata: overrides.metadata ?? {
|
||||
source: { id: 'test', extension: { id: 'server-runtime' } },
|
||||
event: { id: data.id },
|
||||
source: { extension: { id: 'server-runtime' }, id: 'test' },
|
||||
},
|
||||
route: overrides.route,
|
||||
type: 'spark:notify',
|
||||
} as WebSocketBaseEvent<'spark:notify', WebSocketEvents['spark:notify'], any>
|
||||
}
|
||||
|
||||
@@ -98,17 +98,17 @@ describe('match-expression', () => {
|
||||
|
||||
it('matches route expressions', () => {
|
||||
const peer = createPeer({
|
||||
extension: 'stage-ui',
|
||||
id: 'peer-1',
|
||||
name: 'stage-ui',
|
||||
extension: 'stage-ui',
|
||||
instanceId: 'stage-ui-1',
|
||||
labels: { env: 'prod' },
|
||||
name: 'stage-ui',
|
||||
})
|
||||
|
||||
const expression: RouteTargetExpression = { selectors: ['env=prod'], type: 'label' }
|
||||
const expression: RouteTargetExpression = { type: 'label', selectors: ['env=prod'] }
|
||||
expect(matchesRouteExpression(expression, peer)).toBe(true)
|
||||
|
||||
const globExpression: RouteTargetExpression = { glob: 'stage-*', type: 'glob' }
|
||||
const globExpression: RouteTargetExpression = { type: 'glob', glob: 'stage-*' }
|
||||
expect(matchesRouteExpression(globExpression, peer)).toBe(true)
|
||||
})
|
||||
})
|
||||
@@ -117,12 +117,12 @@ describe('route middleware', () => {
|
||||
it('collects destinations from route before data', () => {
|
||||
const event = createSparkNotifyEvent({
|
||||
data: {
|
||||
destinations: ['module:character'],
|
||||
eventId: 'spark-2',
|
||||
headline: 'hello',
|
||||
id: 'evt-2',
|
||||
eventId: 'spark-2',
|
||||
kind: 'ping',
|
||||
urgency: 'soon',
|
||||
headline: 'hello',
|
||||
destinations: ['module:character'],
|
||||
},
|
||||
route: { destinations: ['label:env=prod'] },
|
||||
})
|
||||
@@ -132,12 +132,12 @@ describe('route middleware', () => {
|
||||
it('treats an explicit empty route destination list as the override', () => {
|
||||
const event = createSparkNotifyEvent({
|
||||
data: {
|
||||
destinations: ['module:character'],
|
||||
eventId: 'spark-override',
|
||||
headline: 'hello',
|
||||
id: 'evt-override',
|
||||
eventId: 'spark-override',
|
||||
kind: 'ping',
|
||||
urgency: 'soon',
|
||||
headline: 'hello',
|
||||
destinations: ['module:character'],
|
||||
},
|
||||
route: { destinations: [] },
|
||||
})
|
||||
@@ -148,12 +148,12 @@ describe('route middleware', () => {
|
||||
it('treats an explicit empty data destination list as the override', () => {
|
||||
const event = createSparkNotifyEvent({
|
||||
data: {
|
||||
destinations: [],
|
||||
eventId: 'spark-data-empty',
|
||||
headline: 'hello',
|
||||
id: 'evt-data-empty',
|
||||
eventId: 'spark-data-empty',
|
||||
kind: 'ping',
|
||||
urgency: 'soon',
|
||||
headline: 'hello',
|
||||
destinations: [],
|
||||
},
|
||||
route: undefined,
|
||||
})
|
||||
@@ -163,13 +163,13 @@ describe('route middleware', () => {
|
||||
|
||||
it('ignores primitive data payloads when checking destinations', () => {
|
||||
const event = {
|
||||
type: 'spark:notify',
|
||||
data: 'not-an-object',
|
||||
metadata: {
|
||||
source: { id: 'test', extension: { id: 'server-runtime' } },
|
||||
event: { id: 'evt-primitive' },
|
||||
source: { extension: { id: 'server-runtime' }, id: 'test' },
|
||||
},
|
||||
route: undefined,
|
||||
type: 'spark:notify',
|
||||
} as unknown as WebSocketBaseEvent<'spark:notify', WebSocketEvents['spark:notify'], any>
|
||||
|
||||
expect(collectDestinations(event)).toBeUndefined()
|
||||
@@ -177,11 +177,11 @@ describe('route middleware', () => {
|
||||
|
||||
it('matches destinations by label selector', () => {
|
||||
const peer = createPeer({
|
||||
extension: 'telegram-bot',
|
||||
id: 'peer-2',
|
||||
name: 'telegram-bot',
|
||||
extension: 'telegram-bot',
|
||||
instanceId: 'telegram-1',
|
||||
labels: { app: 'telegram', env: 'prod' },
|
||||
name: 'telegram-bot',
|
||||
})
|
||||
|
||||
expect(matchesDestinations(['label:app=telegram'], peer)).toBe(true)
|
||||
@@ -194,13 +194,13 @@ describe('route middleware', () => {
|
||||
*/
|
||||
it('matches destinations by extension identity labels', () => {
|
||||
const peer = createPeer({
|
||||
extensionLabels: { surface: 'websocket-extension' },
|
||||
id: 'peer-extension-labels',
|
||||
name: 'airi-extension',
|
||||
extensionLabels: { surface: 'websocket-extension' },
|
||||
})
|
||||
|
||||
expect(matchesDestinations(['label:surface=websocket-extension'], peer)).toBe(true)
|
||||
expect(matchesRouteExpression({ selectors: ['surface=websocket-extension'], type: 'label' }, peer)).toBe(true)
|
||||
expect(matchesRouteExpression({ type: 'label', selectors: ['surface=websocket-extension'] }, peer)).toBe(true)
|
||||
expect(matchesDestinations(['label:surface=legacy-plugin'], peer)).toBe(false)
|
||||
})
|
||||
|
||||
@@ -216,7 +216,7 @@ describe('route middleware', () => {
|
||||
})
|
||||
|
||||
expect(matchesDestinations(['peer:stage-window'], peer)).toBe(true)
|
||||
expect(matchesDestinations([{ ids: ['stage-window'], type: 'ids' }], peer)).toBe(true)
|
||||
expect(matchesDestinations([{ type: 'ids', ids: ['stage-window'] }], peer)).toBe(true)
|
||||
expect(matchesDestinations(['peer:missing'], peer)).toBe(false)
|
||||
})
|
||||
|
||||
@@ -236,16 +236,16 @@ describe('route middleware', () => {
|
||||
|
||||
it('policy middleware filters targets', () => {
|
||||
const peers = new Map<string, AuthenticatedPeer>([
|
||||
['peer-1', createPeer({ extension: 'telegram-bot', id: 'peer-1', instanceId: 'telegram-1', labels: { env: 'prod' }, name: 'telegram' })],
|
||||
['peer-2', createPeer({ extension: 'stage-ui', id: 'peer-2', instanceId: 'stage-ui-1', labels: { env: 'dev' }, name: 'stage-ui' })],
|
||||
['peer-1', createPeer({ id: 'peer-1', name: 'telegram', extension: 'telegram-bot', instanceId: 'telegram-1', labels: { env: 'prod' } })],
|
||||
['peer-2', createPeer({ id: 'peer-2', name: 'stage-ui', extension: 'stage-ui', instanceId: 'stage-ui-1', labels: { env: 'dev' } })],
|
||||
])
|
||||
|
||||
const policy = createPolicyMiddleware({ allowLabels: ['env=prod'] })
|
||||
const decision = policy({
|
||||
destinations: undefined,
|
||||
event: createSparkNotifyEvent(),
|
||||
fromPeer: peers.get('peer-1')!,
|
||||
peers,
|
||||
destinations: undefined,
|
||||
})
|
||||
|
||||
expect(decision).toBeDefined()
|
||||
@@ -261,16 +261,16 @@ describe('route middleware', () => {
|
||||
|
||||
it('policy middleware excludes unauthenticated peers', () => {
|
||||
const peers = new Map<string, AuthenticatedPeer>([
|
||||
['peer-1', createPeer({ extension: 'telegram-bot', id: 'peer-1', instanceId: 'telegram-1', labels: { env: 'prod' }, name: 'telegram' })],
|
||||
['peer-2', createPeer({ authenticated: false, extension: 'stage-ui', id: 'peer-2', instanceId: 'stage-ui-1', labels: { env: 'prod' }, name: 'stage-ui' })],
|
||||
['peer-1', createPeer({ id: 'peer-1', name: 'telegram', extension: 'telegram-bot', instanceId: 'telegram-1', labels: { env: 'prod' } })],
|
||||
['peer-2', createPeer({ id: 'peer-2', name: 'stage-ui', extension: 'stage-ui', instanceId: 'stage-ui-1', labels: { env: 'prod' }, authenticated: false })],
|
||||
])
|
||||
|
||||
const policy = createPolicyMiddleware({ allowLabels: ['env=prod'] })
|
||||
const decision = policy({
|
||||
destinations: undefined,
|
||||
event: createSparkNotifyEvent(),
|
||||
fromPeer: peers.get('peer-1')!,
|
||||
peers,
|
||||
destinations: undefined,
|
||||
})
|
||||
|
||||
expect(decision).toBeDefined()
|
||||
@@ -282,16 +282,16 @@ describe('route middleware', () => {
|
||||
|
||||
it('policy middleware does not authorize bypass by itself', () => {
|
||||
const peers = new Map<string, AuthenticatedPeer>([
|
||||
['peer-1', createPeer({ extension: 'telegram-bot', id: 'peer-1', instanceId: 'telegram-1', labels: { env: 'prod' }, name: 'telegram' })],
|
||||
['peer-2', createPeer({ extension: 'stage-ui', id: 'peer-2', instanceId: 'stage-ui-1', labels: { env: 'dev' }, name: 'stage-ui' })],
|
||||
['peer-1', createPeer({ id: 'peer-1', name: 'telegram', extension: 'telegram-bot', instanceId: 'telegram-1', labels: { env: 'prod' } })],
|
||||
['peer-2', createPeer({ id: 'peer-2', name: 'stage-ui', extension: 'stage-ui', instanceId: 'stage-ui-1', labels: { env: 'dev' } })],
|
||||
])
|
||||
|
||||
const policy = createPolicyMiddleware({ allowLabels: ['env=prod'] })
|
||||
const decision = policy({
|
||||
destinations: undefined,
|
||||
event: createSparkNotifyEvent({ route: { bypass: true } }),
|
||||
fromPeer: peers.get('peer-1')!,
|
||||
peers,
|
||||
destinations: undefined,
|
||||
})
|
||||
|
||||
expect(decision).toBeDefined()
|
||||
@@ -303,11 +303,11 @@ describe('route middleware', () => {
|
||||
|
||||
it('devtools peer detection uses label', () => {
|
||||
const peer = createPeer({
|
||||
extension: 'debug-ui',
|
||||
id: 'peer-3',
|
||||
name: 'debug-ui',
|
||||
extension: 'debug-ui',
|
||||
instanceId: 'debug-ui-1',
|
||||
labels: { devtools: 'true' },
|
||||
name: 'debug-ui',
|
||||
})
|
||||
|
||||
expect(isDevtoolsPeer(peer)).toBe(true)
|
||||
|
||||
@@ -4,82 +4,33 @@ import type { AuthenticatedPeer } from '../types'
|
||||
|
||||
import { matchesDestinations, matchesLabelSelectors } from './route/match-expression'
|
||||
|
||||
export interface RouteContext {
|
||||
destinations?: Array<RouteTargetExpression | string>
|
||||
event: WebSocketEvent
|
||||
fromPeer: AuthenticatedPeer
|
||||
peers: Map<string, AuthenticatedPeer>
|
||||
}
|
||||
|
||||
export type RouteDecision
|
||||
= | { targetIds: Set<string>, type: 'targets' }
|
||||
= | { type: 'drop' }
|
||||
| { type: 'broadcast' }
|
||||
| { type: 'drop' }
|
||||
|
||||
export type RouteMiddleware = (context: RouteContext) => RouteDecision | void
|
||||
| { type: 'targets', targetIds: Set<string> }
|
||||
|
||||
export interface RoutingPolicy {
|
||||
allowExtensions?: string[]
|
||||
allowLabels?: string[]
|
||||
denyExtensions?: string[]
|
||||
allowLabels?: string[]
|
||||
denyLabels?: string[]
|
||||
}
|
||||
|
||||
type DestinationList = Array<RouteTargetExpression | string>
|
||||
|
||||
/**
|
||||
* Collects explicit route destinations from the route envelope or event payload.
|
||||
*
|
||||
* Use when:
|
||||
* - Routing middleware needs the effective destination override for a websocket event
|
||||
* - Callers must preserve explicit empty destination lists instead of falling back to broadcast
|
||||
*
|
||||
* Expects:
|
||||
* - A websocket event whose `route.destinations` or `data.destinations` may be present
|
||||
* - `data.destinations` is only treated as valid when it is an array-shaped override
|
||||
*
|
||||
* Returns:
|
||||
* - The explicit destination list when present
|
||||
* - `undefined` when no destination override was provided
|
||||
*/
|
||||
export function collectDestinations(
|
||||
event: (Omit<WebSocketEvent, 'metadata'> & Partial<Pick<WebSocketEvent, 'metadata'>>) | WebSocketEvent,
|
||||
): DestinationList | undefined {
|
||||
if (event.route && 'destinations' in event.route) {
|
||||
return event.route.destinations
|
||||
}
|
||||
|
||||
const data = event.data as unknown
|
||||
if (typeof data === 'object' && data !== null && 'destinations' in data && Array.isArray(data.destinations)) {
|
||||
return data.destinations
|
||||
}
|
||||
|
||||
return undefined
|
||||
export interface RouteContext {
|
||||
event: WebSocketEvent
|
||||
fromPeer: AuthenticatedPeer
|
||||
peers: Map<string, AuthenticatedPeer>
|
||||
destinations?: Array<string | RouteTargetExpression>
|
||||
}
|
||||
|
||||
/**
|
||||
* Creates a routing middleware from a static allow/deny policy.
|
||||
*
|
||||
* Use when:
|
||||
* - Server-wide routing rules should be applied consistently
|
||||
* - Destination filtering should be derived from peer metadata instead of event payloads
|
||||
*
|
||||
* Expects:
|
||||
* - Route bypass authorization to be handled by the caller, not by the policy itself
|
||||
*
|
||||
* Returns:
|
||||
* - A middleware that narrows delivery to the peers allowed by the policy
|
||||
*/
|
||||
export function createPolicyMiddleware(policy: RoutingPolicy): RouteMiddleware {
|
||||
return ({ peers }) => {
|
||||
const targetIds = new Set<string>()
|
||||
for (const [id, peer] of peers.entries()) {
|
||||
if (peerMatchesPolicy(peer, policy)) {
|
||||
targetIds.add(id)
|
||||
}
|
||||
}
|
||||
export type RouteMiddleware = (context: RouteContext) => RouteDecision | void
|
||||
|
||||
return { targetIds, type: 'targets' }
|
||||
type DestinationList = Array<string | RouteTargetExpression>
|
||||
|
||||
function getPeerLabels(peer: AuthenticatedPeer) {
|
||||
return {
|
||||
...peer.extensionIdentity?.labels,
|
||||
...peer.identity?.labels,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -142,11 +93,60 @@ export function peerMatchesPolicy(peer: AuthenticatedPeer, policy: RoutingPolicy
|
||||
return true
|
||||
}
|
||||
|
||||
function getPeerLabels(peer: AuthenticatedPeer) {
|
||||
return {
|
||||
...peer.extensionIdentity?.labels,
|
||||
...peer.identity?.labels,
|
||||
/**
|
||||
* Creates a routing middleware from a static allow/deny policy.
|
||||
*
|
||||
* Use when:
|
||||
* - Server-wide routing rules should be applied consistently
|
||||
* - Destination filtering should be derived from peer metadata instead of event payloads
|
||||
*
|
||||
* Expects:
|
||||
* - Route bypass authorization to be handled by the caller, not by the policy itself
|
||||
*
|
||||
* Returns:
|
||||
* - A middleware that narrows delivery to the peers allowed by the policy
|
||||
*/
|
||||
export function createPolicyMiddleware(policy: RoutingPolicy): RouteMiddleware {
|
||||
return ({ peers }) => {
|
||||
const targetIds = new Set<string>()
|
||||
for (const [id, peer] of peers.entries()) {
|
||||
if (peerMatchesPolicy(peer, policy)) {
|
||||
targetIds.add(id)
|
||||
}
|
||||
}
|
||||
|
||||
return { type: 'targets', targetIds }
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Collects explicit route destinations from the route envelope or event payload.
|
||||
*
|
||||
* Use when:
|
||||
* - Routing middleware needs the effective destination override for a websocket event
|
||||
* - Callers must preserve explicit empty destination lists instead of falling back to broadcast
|
||||
*
|
||||
* Expects:
|
||||
* - A websocket event whose `route.destinations` or `data.destinations` may be present
|
||||
* - `data.destinations` is only treated as valid when it is an array-shaped override
|
||||
*
|
||||
* Returns:
|
||||
* - The explicit destination list when present
|
||||
* - `undefined` when no destination override was provided
|
||||
*/
|
||||
export function collectDestinations(
|
||||
event: WebSocketEvent | (Omit<WebSocketEvent, 'metadata'> & Partial<Pick<WebSocketEvent, 'metadata'>>),
|
||||
): DestinationList | undefined {
|
||||
if (event.route && 'destinations' in event.route) {
|
||||
return event.route.destinations
|
||||
}
|
||||
|
||||
const data = event.data as unknown
|
||||
if (typeof data === 'object' && data !== null && 'destinations' in data && Array.isArray(data.destinations)) {
|
||||
return data.destinations
|
||||
}
|
||||
|
||||
return undefined
|
||||
}
|
||||
|
||||
export { matchesDestinations }
|
||||
|
||||
@@ -2,44 +2,19 @@ import type { RouteTargetExpression } from '@proj-airi/server-shared/types'
|
||||
|
||||
import type { AuthenticatedPeer } from '../../types'
|
||||
|
||||
export function matchesDestination(destination: RouteTargetExpression | string, peer: AuthenticatedPeer) {
|
||||
if (typeof destination !== 'string') {
|
||||
return matchesRouteExpression(destination, peer)
|
||||
}
|
||||
function globToRegExp(glob: string) {
|
||||
const escaped = glob.replace(/[.+^${}()|[\]\\]/g, '\\$&')
|
||||
|
||||
if (destination === '*') {
|
||||
return true
|
||||
}
|
||||
|
||||
const [prefix, rawValue] = destination.split(':', 2)
|
||||
const value = rawValue ?? ''
|
||||
|
||||
switch (prefix) {
|
||||
case 'instance':
|
||||
return peer.identity?.id === value
|
||||
case 'label':
|
||||
return matchesLabelSelectors([value], getPeerLabels(peer))
|
||||
case 'module':
|
||||
return peer.name === value || matchesExtensionModule(peer, value)
|
||||
case 'peer':
|
||||
return matchesPeerId(peer, value)
|
||||
case 'plugin':
|
||||
return getPeerExtensionId(peer) === value
|
||||
case 'source':
|
||||
return peer.name === value
|
||||
default: {
|
||||
const extensionId = getPeerExtensionId(peer)
|
||||
// REVIEW: Bare/glob destination matching is kept for existing event payloads that do not use module:<name>.
|
||||
return matchesGlob(destination, peer.name)
|
||||
|| matchesGlob(destination, extensionId)
|
||||
|| matchesGlob(destination, peer.identity?.id)
|
||||
|| matchesExtensionModuleGlob(peer, destination)
|
||||
}
|
||||
}
|
||||
const pattern = `^${escaped.replace(/\*/g, '.*')}$`
|
||||
return new RegExp(pattern)
|
||||
}
|
||||
|
||||
export function matchesDestinations(destinations: Array<RouteTargetExpression | string>, peer: AuthenticatedPeer) {
|
||||
return destinations.some(destination => matchesDestination(destination, peer))
|
||||
function matchesGlob(glob: string, value?: string) {
|
||||
if (!value) {
|
||||
return false
|
||||
}
|
||||
|
||||
return globToRegExp(glob).test(value)
|
||||
}
|
||||
|
||||
export function matchesLabelSelector(selector: string, labels: Record<string, string>) {
|
||||
@@ -62,10 +37,38 @@ export function matchesLabelSelectors(selectors: string[], labels: Record<string
|
||||
return selectors.every(selector => matchesLabelSelector(selector, labels))
|
||||
}
|
||||
|
||||
function getPeerLabels(peer: AuthenticatedPeer) {
|
||||
return {
|
||||
...peer.extensionIdentity?.labels,
|
||||
...peer.identity?.extension.labels,
|
||||
...peer.identity?.labels,
|
||||
}
|
||||
}
|
||||
|
||||
function getPeerExtensionId(peer: AuthenticatedPeer) {
|
||||
return peer.identity?.extension.id ?? peer.extensionIdentity?.id
|
||||
}
|
||||
|
||||
function matchesExtensionModule(peer: AuthenticatedPeer, moduleName: string) {
|
||||
return [...peer.extensionModules?.values() ?? []]
|
||||
.some(module => module.name === moduleName || module.identity.id === moduleName)
|
||||
}
|
||||
|
||||
function matchesExtensionModuleGlob(peer: AuthenticatedPeer, glob: string) {
|
||||
return [...peer.extensionModules?.values() ?? []]
|
||||
.some(module => matchesGlob(glob, module.name) || matchesGlob(glob, module.identity.id))
|
||||
}
|
||||
|
||||
function matchesPeerId(peer: AuthenticatedPeer, peerId: string) {
|
||||
return peer.peer.id === peerId || Boolean(peer.peerIds?.has(peerId))
|
||||
}
|
||||
|
||||
export function matchesRouteExpression(expression: RouteTargetExpression, peer: AuthenticatedPeer): boolean {
|
||||
switch (expression.type) {
|
||||
case 'and':
|
||||
return expression.all.every(expr => matchesRouteExpression(expr, peer))
|
||||
case 'or':
|
||||
return expression.any.some(expr => matchesRouteExpression(expr, peer))
|
||||
case 'glob': {
|
||||
const extensionId = getPeerExtensionId(peer)
|
||||
const matched = matchesGlob(expression.glob, peer.name)
|
||||
@@ -78,6 +81,10 @@ export function matchesRouteExpression(expression: RouteTargetExpression, peer:
|
||||
const matched = expression.ids.some(peerId => matchesPeerId(peer, peerId))
|
||||
return expression.inverted ? !matched : matched
|
||||
}
|
||||
case 'plugin': {
|
||||
const matched = expression.plugins.includes(getPeerExtensionId(peer) ?? '')
|
||||
return expression.inverted ? !matched : matched
|
||||
}
|
||||
case 'instance': {
|
||||
const matched = expression.instances.includes(peer.identity?.id ?? '')
|
||||
return expression.inverted ? !matched : matched
|
||||
@@ -90,12 +97,6 @@ export function matchesRouteExpression(expression: RouteTargetExpression, peer:
|
||||
const matched = expression.modules.some(module => peer.name === module || matchesExtensionModule(peer, module))
|
||||
return expression.inverted ? !matched : matched
|
||||
}
|
||||
case 'or':
|
||||
return expression.any.some(expr => matchesRouteExpression(expr, peer))
|
||||
case 'plugin': {
|
||||
const matched = expression.plugins.includes(getPeerExtensionId(peer) ?? '')
|
||||
return expression.inverted ? !matched : matched
|
||||
}
|
||||
case 'source': {
|
||||
const matched = expression.sources.includes(peer.name)
|
||||
return expression.inverted ? !matched : matched
|
||||
@@ -105,43 +106,42 @@ export function matchesRouteExpression(expression: RouteTargetExpression, peer:
|
||||
}
|
||||
}
|
||||
|
||||
function getPeerExtensionId(peer: AuthenticatedPeer) {
|
||||
return peer.identity?.extension.id ?? peer.extensionIdentity?.id
|
||||
}
|
||||
export function matchesDestination(destination: string | RouteTargetExpression, peer: AuthenticatedPeer) {
|
||||
if (typeof destination !== 'string') {
|
||||
return matchesRouteExpression(destination, peer)
|
||||
}
|
||||
|
||||
function getPeerLabels(peer: AuthenticatedPeer) {
|
||||
return {
|
||||
...peer.extensionIdentity?.labels,
|
||||
...peer.identity?.extension.labels,
|
||||
...peer.identity?.labels,
|
||||
if (destination === '*') {
|
||||
return true
|
||||
}
|
||||
|
||||
const [prefix, rawValue] = destination.split(':', 2)
|
||||
const value = rawValue ?? ''
|
||||
|
||||
switch (prefix) {
|
||||
case 'plugin':
|
||||
return getPeerExtensionId(peer) === value
|
||||
case 'instance':
|
||||
return peer.identity?.id === value
|
||||
case 'label':
|
||||
return matchesLabelSelectors([value], getPeerLabels(peer))
|
||||
case 'peer':
|
||||
return matchesPeerId(peer, value)
|
||||
case 'module':
|
||||
return peer.name === value || matchesExtensionModule(peer, value)
|
||||
case 'source':
|
||||
return peer.name === value
|
||||
default: {
|
||||
const extensionId = getPeerExtensionId(peer)
|
||||
// REVIEW: Bare/glob destination matching is kept for existing event payloads that do not use module:<name>.
|
||||
return matchesGlob(destination, peer.name)
|
||||
|| matchesGlob(destination, extensionId)
|
||||
|| matchesGlob(destination, peer.identity?.id)
|
||||
|| matchesExtensionModuleGlob(peer, destination)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
function globToRegExp(glob: string) {
|
||||
const escaped = glob.replace(/[.+^${}()|[\]\\]/g, '\\$&')
|
||||
|
||||
const pattern = `^${escaped.replace(/\*/g, '.*')}$`
|
||||
return new RegExp(pattern)
|
||||
}
|
||||
|
||||
function matchesExtensionModule(peer: AuthenticatedPeer, moduleName: string) {
|
||||
return [...peer.extensionModules?.values() ?? []]
|
||||
.some(module => module.name === moduleName || module.identity.id === moduleName)
|
||||
}
|
||||
|
||||
function matchesExtensionModuleGlob(peer: AuthenticatedPeer, glob: string) {
|
||||
return [...peer.extensionModules?.values() ?? []]
|
||||
.some(module => matchesGlob(glob, module.name) || matchesGlob(glob, module.identity.id))
|
||||
}
|
||||
|
||||
function matchesGlob(glob: string, value?: string) {
|
||||
if (!value) {
|
||||
return false
|
||||
}
|
||||
|
||||
return globToRegExp(glob).test(value)
|
||||
}
|
||||
|
||||
function matchesPeerId(peer: AuthenticatedPeer, peerId: string) {
|
||||
return peer.peer.id === peerId || Boolean(peer.peerIds?.has(peerId))
|
||||
export function matchesDestinations(destinations: Array<string | RouteTargetExpression>, peer: AuthenticatedPeer) {
|
||||
return destinations.some(destination => matchesDestination(destination, peer))
|
||||
}
|
||||
|
||||
@@ -13,16 +13,16 @@ import {
|
||||
describe('airi websocket codec', () => {
|
||||
it('parses superjson encoded events', () => {
|
||||
const event: WebSocketEvent = {
|
||||
type: 'module:authenticate',
|
||||
data: { token: 'secret' },
|
||||
metadata: {
|
||||
event: { id: 'event-1' },
|
||||
source: {
|
||||
id: 'test-plugin-1',
|
||||
kind: 'plugin',
|
||||
id: 'test-plugin-1',
|
||||
plugin: { id: 'test-plugin' },
|
||||
},
|
||||
event: { id: 'event-1' },
|
||||
},
|
||||
type: 'module:authenticate',
|
||||
}
|
||||
|
||||
expect(parseEvent(stringifySuperJson(event))).toEqual(event)
|
||||
@@ -30,16 +30,16 @@ describe('airi websocket codec', () => {
|
||||
|
||||
it('falls back to plain JSON events', () => {
|
||||
const event: WebSocketEvent = {
|
||||
type: 'module:authenticate',
|
||||
data: { token: 'secret' },
|
||||
metadata: {
|
||||
event: { id: 'event-1' },
|
||||
source: {
|
||||
id: 'test-plugin-1',
|
||||
kind: 'plugin',
|
||||
id: 'test-plugin-1',
|
||||
plugin: { id: 'test-plugin' },
|
||||
},
|
||||
event: { id: 'event-1' },
|
||||
},
|
||||
type: 'module:authenticate',
|
||||
}
|
||||
|
||||
expect(parseEvent(JSON.stringify(event))).toEqual(event)
|
||||
@@ -50,20 +50,20 @@ describe('airi websocket codec', () => {
|
||||
.toThrow(InvalidEventError)
|
||||
expect(() => parseEvent(JSON.stringify({ data: {} })))
|
||||
.toThrow(InvalidEventError)
|
||||
expect(() => parseEvent(JSON.stringify({ data: {}, type: 0 })))
|
||||
expect(() => parseEvent(JSON.stringify({ type: 0, data: {} })))
|
||||
.toThrow(InvalidEventError)
|
||||
expect(() => parseEvent(JSON.stringify({ type: 'module:authenticate' })))
|
||||
.toThrow(InvalidEventError)
|
||||
expect(() => parseEvent(JSON.stringify({ data: null, type: 'module:authenticate' })))
|
||||
expect(() => parseEvent(JSON.stringify({ type: 'module:authenticate', data: null })))
|
||||
.toThrow(InvalidEventError)
|
||||
expect(() => parseEvent(JSON.stringify({ data: 'secret', type: 'module:authenticate' })))
|
||||
expect(() => parseEvent(JSON.stringify({ type: 'module:authenticate', data: 'secret' })))
|
||||
.toThrow(InvalidEventError)
|
||||
expect(() => parseEvent(JSON.stringify({ data: [], type: 'module:authenticate' })))
|
||||
expect(() => parseEvent(JSON.stringify({ type: 'module:authenticate', data: [] })))
|
||||
.toThrow(InvalidEventError)
|
||||
})
|
||||
|
||||
it('keeps validation cause and source on invalid event errors', () => {
|
||||
const source = { data: 'secret', type: 'module:authenticate' }
|
||||
const source = { type: 'module:authenticate', data: 'secret' }
|
||||
|
||||
try {
|
||||
parseEvent(JSON.stringify(source))
|
||||
|
||||
@@ -15,8 +15,8 @@ const eventDataSchema = pipe(
|
||||
)
|
||||
|
||||
const eventEnvelopeSchema = objectWithRest({
|
||||
data: eventDataSchema,
|
||||
type: string(),
|
||||
data: eventDataSchema,
|
||||
}, unknown())
|
||||
|
||||
interface InvalidEventErrorOptions {
|
||||
@@ -35,6 +35,11 @@ export class InvalidEventError extends Error {
|
||||
}
|
||||
}
|
||||
|
||||
/** Checks whether an error came from AIRI websocket event envelope validation. */
|
||||
export function isInvalidEventError(error: unknown): error is InvalidEventError {
|
||||
return error instanceof InvalidEventError
|
||||
}
|
||||
|
||||
/** Detects raw ping/pong text frames that should not enter the event protocol. */
|
||||
export function heartbeatFrameFrom(text: string): MessageHeartbeatKind | undefined {
|
||||
if (text === MessageHeartbeatKind.Ping || text === MessageHeartbeatKind.Pong) {
|
||||
@@ -42,11 +47,6 @@ export function heartbeatFrameFrom(text: string): MessageHeartbeatKind | undefin
|
||||
}
|
||||
}
|
||||
|
||||
/** Checks whether an error came from AIRI websocket event envelope validation. */
|
||||
export function isInvalidEventError(error: unknown): error is InvalidEventError {
|
||||
return error instanceof InvalidEventError
|
||||
}
|
||||
|
||||
/** Parses one AIRI websocket protocol event from SuperJSON or plain JSON text. */
|
||||
export function parseEvent(text: string): WebSocketEvent {
|
||||
// NOTICE:
|
||||
@@ -55,7 +55,7 @@ export function parseEvent(text: string): WebSocketEvent {
|
||||
// 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: undefined | WebSocketEvent
|
||||
let parsed: WebSocketEvent | undefined
|
||||
try {
|
||||
parsed = parse<WebSocketEvent>(text)
|
||||
}
|
||||
@@ -76,6 +76,6 @@ export function parseEvent(text: string): WebSocketEvent {
|
||||
}
|
||||
|
||||
/** Serializes one AIRI websocket protocol event with the existing SuperJSON wire format. */
|
||||
export function stringifyEvent(event: string | WebSocketBaseEvent<string, unknown>) {
|
||||
export function stringifyEvent(event: WebSocketBaseEvent<string, unknown> | string) {
|
||||
return typeof event === 'string' ? event : stringify(event)
|
||||
}
|
||||
|
||||
@@ -10,66 +10,66 @@ import {
|
||||
describe('airi websocket consumer selection', () => {
|
||||
it('selects highest priority then earliest registration', () => {
|
||||
expect(selectConsumerPeerId({
|
||||
candidates: [
|
||||
{ authenticated: true, peerId: 'late', priority: 1, registeredAt: 2 },
|
||||
{ authenticated: true, peerId: 'early', priority: 1, registeredAt: 1 },
|
||||
{ authenticated: true, peerId: 'low', priority: 0, registeredAt: 0 },
|
||||
],
|
||||
delivery: { mode: 'consumer', selection: 'first' },
|
||||
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({
|
||||
candidates: [
|
||||
{ authenticated: true, peerId: 'sender', priority: 3, registeredAt: 1 },
|
||||
{ authenticated: false, peerId: 'unauthenticated', priority: 2, registeredAt: 1 },
|
||||
{ authenticated: true, healthy: false, peerId: 'unhealthy', priority: 1, registeredAt: 1 },
|
||||
{ authenticated: true, peerId: 'target', priority: 0, registeredAt: 1 },
|
||||
],
|
||||
delivery: { mode: 'consumer', selection: 'first' },
|
||||
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, ConsumerStickyAssignment>()
|
||||
const delivery = { group: 'workers', mode: 'consumer-group' as const, selection: 'sticky' as const, stickyKey: 'job-1' }
|
||||
const delivery = { mode: 'consumer-group' as const, group: 'workers', selection: 'sticky' as const, stickyKey: 'job-1' }
|
||||
const candidates = [
|
||||
{ authenticated: true, peerId: 'a', priority: 0, registeredAt: 1 },
|
||||
{ authenticated: true, peerId: 'b', priority: 0, registeredAt: 2 },
|
||||
{ peerId: 'a', priority: 0, registeredAt: 1, authenticated: true },
|
||||
{ peerId: 'b', priority: 0, registeredAt: 2, authenticated: true },
|
||||
]
|
||||
|
||||
expect(selectConsumerPeerId({
|
||||
candidates,
|
||||
delivery,
|
||||
eventType: 'event:test',
|
||||
fromPeerId: 'sender',
|
||||
delivery,
|
||||
candidates,
|
||||
stickyAssignments,
|
||||
})).toBe('a')
|
||||
expect(selectConsumerPeerId({
|
||||
candidates: [...candidates].reverse(),
|
||||
delivery,
|
||||
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 = { group: 'workers', mode: 'consumer-group' as const, selection: 'round-robin' as const }
|
||||
const delivery = { mode: 'consumer-group' as const, group: 'workers', selection: 'round-robin' as const }
|
||||
const candidates = [
|
||||
{ authenticated: true, peerId: 'a', priority: 0, registeredAt: 1 },
|
||||
{ authenticated: true, peerId: 'b', priority: 0, registeredAt: 2 },
|
||||
{ peerId: 'a', priority: 0, registeredAt: 1, authenticated: true },
|
||||
{ peerId: 'b', priority: 0, registeredAt: 2, authenticated: true },
|
||||
]
|
||||
|
||||
expect(selectConsumerPeerId({ candidates, delivery, eventType: 'event:test', fromPeerId: 'sender', roundRobinCursor })).toBe('a')
|
||||
expect(selectConsumerPeerId({ candidates, delivery, eventType: 'event:test', fromPeerId: 'sender', roundRobinCursor })).toBe('b')
|
||||
expect(selectConsumerPeerId({ candidates, delivery, eventType: 'event:test', fromPeerId: 'sender', roundRobinCursor })).toBe('a')
|
||||
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')
|
||||
})
|
||||
})
|
||||
|
||||
@@ -77,13 +77,13 @@ describe('airi websocket consumer registry', () => {
|
||||
it('registers and unregisters consumers', () => {
|
||||
const registry = createConsumerOrchestrator()
|
||||
|
||||
registry.register({ event: 'event:test', group: 'workers', mode: 'consumer-group', peerId: 'peer-1', priority: 2 })
|
||||
expect(registry.listFor({ event: 'event:test', group: 'workers', mode: 'consumer-group' })).toEqual([
|
||||
expect.objectContaining({ event: 'event:test', group: 'workers', peerId: 'peer-1', priority: 2 }),
|
||||
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({ event: 'event:test', group: 'workers', mode: 'consumer-group', peerId: 'peer-1' })
|
||||
expect(registry.listFor({ event: 'event:test', group: 'workers', mode: 'consumer-group' })).toEqual([])
|
||||
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', () => {
|
||||
@@ -99,131 +99,131 @@ describe('airi websocket consumer registry', () => {
|
||||
//
|
||||
// After:
|
||||
// Peer cleanup stores structured event/group refs and never decodes registry keys.
|
||||
registry.register({ event: 'event::test', group: 'group::workers', mode: 'consumer-group', peerId: 'peer-1' })
|
||||
registry.register({ peerId: 'peer-1', event: 'event::test', mode: 'consumer-group', group: 'group::workers' })
|
||||
registry.unregisterPeer('peer-1')
|
||||
|
||||
expect(registry.listFor({ event: 'event::test', group: 'group::workers', mode: 'consumer-group' })).toEqual([])
|
||||
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, ConsumerStickyAssignment>()
|
||||
const candidates = [
|
||||
{ authenticated: true, peerId: 'event::group-target', priority: 0, registeredAt: 1 },
|
||||
{ authenticated: true, peerId: 'other-target', priority: 0, registeredAt: 2 },
|
||||
{ peerId: 'event::group-target', priority: 0, registeredAt: 1, authenticated: true },
|
||||
{ peerId: 'other-target', priority: 0, registeredAt: 2, authenticated: true },
|
||||
]
|
||||
|
||||
expect(selectConsumerPeerId({
|
||||
candidates,
|
||||
delivery: { group: 'target', mode: 'consumer-group', selection: 'sticky', stickyKey: 'job' },
|
||||
eventType: 'event::group',
|
||||
fromPeerId: 'sender',
|
||||
delivery: { mode: 'consumer-group', group: 'target', selection: 'sticky', stickyKey: 'job' },
|
||||
candidates,
|
||||
stickyAssignments,
|
||||
})).toBe('event::group-target')
|
||||
expect(selectConsumerPeerId({
|
||||
candidates: [
|
||||
{ authenticated: true, peerId: 'other-target', priority: 1, registeredAt: 1 },
|
||||
{ authenticated: true, peerId: 'event::group-target', priority: 0, registeredAt: 2 },
|
||||
],
|
||||
delivery: { group: 'group::target', mode: 'consumer-group', selection: 'sticky', stickyKey: 'job' },
|
||||
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({ event: 'event:test', group: 'workers', mode: 'consumer-group', peerId: 'a' })
|
||||
registry.register({ event: 'event:test', group: 'workers', mode: 'consumer-group', peerId: 'b' })
|
||||
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({
|
||||
candidates: registry.listFor({ event: 'event:test', group: 'workers', mode: 'consumer-group' }).map(entry => ({
|
||||
authenticated: true,
|
||||
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,
|
||||
})),
|
||||
delivery: { group: 'workers', mode: 'consumer-group', selection: 'round-robin' },
|
||||
eventType: 'event:test',
|
||||
fromPeerId: 'sender',
|
||||
})).toBe('a')
|
||||
|
||||
registry.unregister({ event: 'event:test', group: 'workers', mode: 'consumer-group', peerId: 'a' })
|
||||
registry.unregister({ peerId: 'a', event: 'event:test', mode: 'consumer-group', group: 'workers' })
|
||||
|
||||
expect(registry.select({
|
||||
candidates: registry.listFor({ event: 'event:test', group: 'workers', mode: 'consumer-group' }).map(entry => ({
|
||||
authenticated: true,
|
||||
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,
|
||||
})),
|
||||
delivery: { group: 'workers', mode: 'consumer-group', selection: 'round-robin' },
|
||||
eventType: 'event:test',
|
||||
fromPeerId: 'sender',
|
||||
})).toBe('b')
|
||||
})
|
||||
|
||||
it('keeps round-robin cursor when unregister does not change group membership', () => {
|
||||
const registry = createConsumerOrchestrator()
|
||||
registry.register({ event: 'event:test', group: 'workers', mode: 'consumer-group', peerId: 'a' })
|
||||
registry.register({ event: 'event:test', group: 'workers', mode: 'consumer-group', peerId: 'b' })
|
||||
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({
|
||||
candidates: registry.listFor({ event: 'event:test', group: 'workers', mode: 'consumer-group' }).map(entry => ({
|
||||
authenticated: true,
|
||||
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,
|
||||
})),
|
||||
delivery: { group: 'workers', mode: 'consumer-group', selection: 'round-robin' },
|
||||
eventType: 'event:test',
|
||||
fromPeerId: 'sender',
|
||||
})).toBe('a')
|
||||
|
||||
registry.unregister({ event: 'event:test', group: 'workers', mode: 'consumer-group', peerId: 'missing' })
|
||||
registry.unregister({ peerId: 'missing', event: 'event:test', mode: 'consumer-group', group: 'workers' })
|
||||
|
||||
expect(registry.select({
|
||||
candidates: registry.listFor({ event: 'event:test', group: 'workers', mode: 'consumer-group' }).map(entry => ({
|
||||
authenticated: true,
|
||||
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,
|
||||
})),
|
||||
delivery: { group: 'workers', mode: 'consumer-group', selection: 'round-robin' },
|
||||
eventType: 'event:test',
|
||||
fromPeerId: 'sender',
|
||||
})).toBe('b')
|
||||
})
|
||||
|
||||
it('resets round-robin cursor when group membership grows', () => {
|
||||
const registry = createConsumerOrchestrator()
|
||||
registry.register({ event: 'event:test', group: 'workers', mode: 'consumer-group', peerId: 'a' })
|
||||
registry.register({ event: 'event:test', group: 'workers', mode: 'consumer-group', peerId: 'b' })
|
||||
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({
|
||||
candidates: registry.listFor({ event: 'event:test', group: 'workers', mode: 'consumer-group' }).map(entry => ({
|
||||
authenticated: true,
|
||||
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,
|
||||
})),
|
||||
delivery: { group: 'workers', mode: 'consumer-group', selection: 'round-robin' },
|
||||
eventType: 'event:test',
|
||||
fromPeerId: 'sender',
|
||||
})).toBe('a')
|
||||
|
||||
registry.register({ event: 'event:test', group: 'workers', mode: 'consumer-group', peerId: 'c' })
|
||||
registry.register({ peerId: 'c', event: 'event:test', mode: 'consumer-group', group: 'workers' })
|
||||
|
||||
expect(registry.select({
|
||||
candidates: registry.listFor({ event: 'event:test', group: 'workers', mode: 'consumer-group' }).map(entry => ({
|
||||
authenticated: true,
|
||||
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,
|
||||
})),
|
||||
delivery: { group: 'workers', mode: 'consumer-group', selection: 'round-robin' },
|
||||
eventType: 'event:test',
|
||||
fromPeerId: 'sender',
|
||||
})).toBe('a')
|
||||
})
|
||||
})
|
||||
|
||||
@@ -2,6 +2,27 @@ import type { DeliveryConfig } from '@proj-airi/server-shared/types'
|
||||
|
||||
const DEFAULT_CONSUMER_GROUP = 'default'
|
||||
|
||||
interface ConsumerRegistryRef {
|
||||
event: string
|
||||
group: string
|
||||
}
|
||||
|
||||
/**
|
||||
* Candidate peer metadata used for AIRI consumer selection.
|
||||
*/
|
||||
export interface ConsumerSelectionCandidate {
|
||||
/** 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 AIRI consumer registration.
|
||||
*/
|
||||
@@ -18,22 +39,6 @@ export interface ConsumerRegistration {
|
||||
registeredAt: number
|
||||
}
|
||||
|
||||
/**
|
||||
* Candidate peer metadata used for AIRI consumer selection.
|
||||
*/
|
||||
export interface ConsumerSelectionCandidate {
|
||||
/** Whether the peer has completed protocol-level authentication. */
|
||||
authenticated: boolean
|
||||
/** Explicit `false` excludes the peer from selection. */
|
||||
healthy?: boolean
|
||||
/** 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
|
||||
}
|
||||
|
||||
/**
|
||||
* Sticky AIRI consumer assignment stored by the consumer selector.
|
||||
*/
|
||||
@@ -46,9 +51,118 @@ export interface ConsumerStickyAssignment {
|
||||
peerId: string
|
||||
}
|
||||
|
||||
interface ConsumerRegistryRef {
|
||||
event: string
|
||||
group: string
|
||||
/**
|
||||
* Checks whether a delivery mode targets the AIRI consumer registry.
|
||||
*/
|
||||
export function isConsumerDeliveryMode(mode: unknown): mode is 'consumer' | 'consumer-group' {
|
||||
return mode === 'consumer' || mode === 'consumer-group'
|
||||
}
|
||||
|
||||
/**
|
||||
* Normalizes delivery mode for AIRI 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 AIRI consumer priority.
|
||||
*
|
||||
* Before:
|
||||
* - NaN
|
||||
*
|
||||
* After:
|
||||
* - 0
|
||||
*/
|
||||
export function normalizeConsumerPriority(priority: unknown) {
|
||||
return typeof priority === 'number' && Number.isFinite(priority)
|
||||
? priority
|
||||
: 0
|
||||
}
|
||||
|
||||
function normalizeConsumerGroup(mode: 'consumer' | 'consumer-group', group?: string) {
|
||||
if (mode === 'consumer') {
|
||||
return DEFAULT_CONSUMER_GROUP
|
||||
}
|
||||
|
||||
return group || DEFAULT_CONSUMER_GROUP
|
||||
}
|
||||
|
||||
function sortConsumers(entries: Array<Pick<ConsumerSelectionCandidate, '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 peer for AIRI consumer-style delivery modes.
|
||||
*
|
||||
* Sticky and round-robin state are keyed with structured JSON tuples so event,
|
||||
* group, and sticky key values may contain delimiter-like text safely.
|
||||
*/
|
||||
export function selectConsumerPeerId(options: {
|
||||
eventType: string
|
||||
fromPeerId: string
|
||||
delivery?: DeliveryConfig
|
||||
candidates: ConsumerSelectionCandidate[]
|
||||
roundRobinCursor?: Map<string, number>
|
||||
stickyAssignments?: Map<string, ConsumerStickyAssignment>
|
||||
}) {
|
||||
const { candidates, delivery, eventType, fromPeerId } = options
|
||||
if (!delivery || !isConsumerDeliveryMode(delivery.mode)) {
|
||||
return
|
||||
}
|
||||
|
||||
const normalizedGroup = normalizeConsumerGroup(delivery.mode, delivery.group)
|
||||
const registryKey = JSON.stringify([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 = JSON.stringify([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
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -76,17 +190,7 @@ export function createConsumerOrchestrator() {
|
||||
}
|
||||
|
||||
return {
|
||||
clear() {
|
||||
consumerRegistry.clear()
|
||||
consumerKeysByPeer.clear()
|
||||
deliveryRoundRobinCursor.clear()
|
||||
stickyAssignments.clear()
|
||||
},
|
||||
listFor(input: { event: string, group?: string, mode: 'consumer' | 'consumer-group' }) {
|
||||
const normalizedGroup = normalizeConsumerGroup(input.mode, input.group)
|
||||
return [...consumerRegistry.get(input.event)?.get(normalizedGroup)?.values() ?? []]
|
||||
},
|
||||
register(input: { event: string, group?: string, mode: 'consumer' | 'consumer-group', peerId: string, priority?: number }) {
|
||||
register(input: { peerId: string, event: string, mode: 'consumer' | 'consumer-group', group?: string, priority?: number }) {
|
||||
const normalizedGroup = normalizeConsumerGroup(input.mode, input.group)
|
||||
const registryKey = JSON.stringify([input.event, normalizedGroup])
|
||||
let groups = consumerRegistry.get(input.event)
|
||||
@@ -120,19 +224,7 @@ export function createConsumerOrchestrator() {
|
||||
}
|
||||
registrations.set(registryKey, { event: input.event, group: normalizedGroup })
|
||||
},
|
||||
select(input: {
|
||||
candidates: ConsumerSelectionCandidate[]
|
||||
delivery?: DeliveryConfig
|
||||
eventType: string
|
||||
fromPeerId: string
|
||||
}) {
|
||||
return selectConsumerPeerId({
|
||||
...input,
|
||||
roundRobinCursor: deliveryRoundRobinCursor,
|
||||
stickyAssignments,
|
||||
})
|
||||
},
|
||||
unregister(input: { event: string, group?: string, mode: 'consumer' | 'consumer-group', peerId: string }) {
|
||||
unregister(input: { peerId: string, event: string, mode: 'consumer' | 'consumer-group', group?: string }) {
|
||||
const normalizedGroup = normalizeConsumerGroup(input.mode, input.group)
|
||||
const registryKey = JSON.stringify([input.event, normalizedGroup])
|
||||
const groups = consumerRegistry.get(input.event)
|
||||
@@ -183,119 +275,27 @@ export function createConsumerOrchestrator() {
|
||||
|
||||
consumerKeysByPeer.delete(peerId)
|
||||
},
|
||||
listFor(input: { event: string, mode: 'consumer' | 'consumer-group', group?: string }) {
|
||||
const normalizedGroup = normalizeConsumerGroup(input.mode, input.group)
|
||||
return [...consumerRegistry.get(input.event)?.get(normalizedGroup)?.values() ?? []]
|
||||
},
|
||||
select(input: {
|
||||
eventType: string
|
||||
fromPeerId: string
|
||||
delivery?: DeliveryConfig
|
||||
candidates: ConsumerSelectionCandidate[]
|
||||
}) {
|
||||
return selectConsumerPeerId({
|
||||
...input,
|
||||
roundRobinCursor: deliveryRoundRobinCursor,
|
||||
stickyAssignments,
|
||||
})
|
||||
},
|
||||
clear() {
|
||||
consumerRegistry.clear()
|
||||
consumerKeysByPeer.clear()
|
||||
deliveryRoundRobinCursor.clear()
|
||||
stickyAssignments.clear()
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Checks whether a delivery mode targets the AIRI consumer registry.
|
||||
*/
|
||||
export function isConsumerDeliveryMode(mode: unknown): mode is 'consumer' | 'consumer-group' {
|
||||
return mode === 'consumer' || mode === 'consumer-group'
|
||||
}
|
||||
|
||||
/**
|
||||
* Normalizes delivery mode for AIRI 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 AIRI consumer priority.
|
||||
*
|
||||
* Before:
|
||||
* - NaN
|
||||
*
|
||||
* After:
|
||||
* - 0
|
||||
*/
|
||||
export function normalizeConsumerPriority(priority: unknown) {
|
||||
return typeof priority === 'number' && Number.isFinite(priority)
|
||||
? priority
|
||||
: 0
|
||||
}
|
||||
|
||||
/**
|
||||
* Selects a concrete peer for AIRI consumer-style delivery modes.
|
||||
*
|
||||
* Sticky and round-robin state are keyed with structured JSON tuples so event,
|
||||
* group, and sticky key values may contain delimiter-like text safely.
|
||||
*/
|
||||
export function selectConsumerPeerId(options: {
|
||||
candidates: ConsumerSelectionCandidate[]
|
||||
delivery?: DeliveryConfig
|
||||
eventType: string
|
||||
fromPeerId: string
|
||||
roundRobinCursor?: Map<string, number>
|
||||
stickyAssignments?: Map<string, ConsumerStickyAssignment>
|
||||
}) {
|
||||
const { candidates, delivery, eventType, fromPeerId } = options
|
||||
if (!delivery || !isConsumerDeliveryMode(delivery.mode)) {
|
||||
return
|
||||
}
|
||||
|
||||
const normalizedGroup = normalizeConsumerGroup(delivery.mode, delivery.group)
|
||||
const registryKey = JSON.stringify([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 = JSON.stringify([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
|
||||
}
|
||||
|
||||
function normalizeConsumerGroup(mode: 'consumer' | 'consumer-group', group?: string) {
|
||||
if (mode === 'consumer') {
|
||||
return DEFAULT_CONSUMER_GROUP
|
||||
}
|
||||
|
||||
return group || DEFAULT_CONSUMER_GROUP
|
||||
}
|
||||
|
||||
function sortConsumers(entries: Array<Pick<ConsumerSelectionCandidate, 'peerId' | 'priority' | 'registeredAt'>>) {
|
||||
return [...entries].sort((left, right) => {
|
||||
if (right.priority !== left.priority) {
|
||||
return right.priority - left.priority
|
||||
}
|
||||
|
||||
return left.registeredAt - right.registeredAt
|
||||
})
|
||||
}
|
||||
|
||||
@@ -10,12 +10,12 @@ describe('airi websocket responses', () => {
|
||||
const metadata = createEventMetadata('server-1', 'parent-event-1')
|
||||
|
||||
expect(metadata.source).toEqual({
|
||||
id: 'server-1',
|
||||
kind: 'plugin',
|
||||
plugin: {
|
||||
id: WebSocketEventSource.Server,
|
||||
version: packageJSON.version,
|
||||
},
|
||||
id: 'server-1',
|
||||
})
|
||||
expect(metadata.event).toEqual({
|
||||
id: expect.any(String),
|
||||
@@ -27,14 +27,12 @@ describe('airi websocket responses', () => {
|
||||
const responses = createResponses('server-1')
|
||||
|
||||
expect(responses.peerAuthenticated('peer-1', 'event-1')).toMatchObject({
|
||||
type: 'peer:authenticated',
|
||||
data: {
|
||||
authenticated: true,
|
||||
peerId: 'peer-1',
|
||||
},
|
||||
metadata: {
|
||||
event: {
|
||||
parentId: 'event-1',
|
||||
},
|
||||
source: {
|
||||
id: 'server-1',
|
||||
plugin: {
|
||||
@@ -42,10 +40,13 @@ describe('airi websocket responses', () => {
|
||||
version: packageJSON.version,
|
||||
},
|
||||
},
|
||||
event: {
|
||||
parentId: 'event-1',
|
||||
},
|
||||
},
|
||||
type: 'peer:authenticated',
|
||||
})
|
||||
expect(responses.extensionAuthenticated({ id: 'airi-extension-chess' }, 'event-2')).toMatchObject({
|
||||
type: 'extension:authenticated',
|
||||
data: {
|
||||
authenticated: true,
|
||||
identity: {
|
||||
@@ -53,9 +54,6 @@ describe('airi websocket responses', () => {
|
||||
},
|
||||
},
|
||||
metadata: {
|
||||
event: {
|
||||
parentId: 'event-2',
|
||||
},
|
||||
source: {
|
||||
id: 'server-1',
|
||||
plugin: {
|
||||
@@ -63,8 +61,10 @@ describe('airi websocket responses', () => {
|
||||
version: packageJSON.version,
|
||||
},
|
||||
},
|
||||
event: {
|
||||
parentId: 'event-2',
|
||||
},
|
||||
},
|
||||
type: 'extension:authenticated',
|
||||
})
|
||||
})
|
||||
})
|
||||
|
||||
@@ -10,19 +10,19 @@ import packageJSON from '../../../package.json'
|
||||
export function createEventMetadata(
|
||||
serverInstanceId: string,
|
||||
parentId?: string,
|
||||
): { event: { id: string, parentId?: string }, source: MetadataEventSource } {
|
||||
): { source: MetadataEventSource, event: { id: string, parentId?: string } } {
|
||||
return {
|
||||
event: {
|
||||
id: nanoid(),
|
||||
parentId,
|
||||
},
|
||||
source: {
|
||||
id: serverInstanceId,
|
||||
kind: 'plugin',
|
||||
plugin: {
|
||||
id: WebSocketEventSource.Server,
|
||||
version: packageJSON.version,
|
||||
},
|
||||
id: serverInstanceId,
|
||||
},
|
||||
}
|
||||
}
|
||||
@@ -32,44 +32,44 @@ export function createResponses(serverInstanceId: string) {
|
||||
return {
|
||||
authenticated(parentId?: string) {
|
||||
return {
|
||||
type: 'module:authenticated',
|
||||
data: { authenticated: true },
|
||||
metadata: createEventMetadata(serverInstanceId, parentId),
|
||||
type: 'module:authenticated',
|
||||
} satisfies WebSocketEvent<Record<string, unknown>>
|
||||
},
|
||||
error(message: string, parentId?: string) {
|
||||
return {
|
||||
data: { message },
|
||||
metadata: createEventMetadata(serverInstanceId, parentId),
|
||||
type: 'error',
|
||||
} satisfies WebSocketEvent<Record<string, unknown>>
|
||||
},
|
||||
extensionAuthenticated(identity: ExtensionIdentity, parentId?: string) {
|
||||
return {
|
||||
data: { authenticated: true, identity },
|
||||
metadata: createEventMetadata(serverInstanceId, parentId),
|
||||
type: 'extension:authenticated',
|
||||
} satisfies WebSocketEvent<Record<string, unknown>>
|
||||
},
|
||||
heartbeat(kind: MessageHeartbeatKind, message: MessageHeartbeat | string, parentId?: string) {
|
||||
return {
|
||||
data: { at: Date.now(), kind, message },
|
||||
metadata: createEventMetadata(serverInstanceId, parentId),
|
||||
type: 'transport:connection:heartbeat',
|
||||
} satisfies WebSocketEvent<Record<string, unknown>>
|
||||
},
|
||||
notAuthenticated(parentId?: string) {
|
||||
return {
|
||||
data: { message: ServerErrorMessages.notAuthenticated },
|
||||
metadata: createEventMetadata(serverInstanceId, parentId),
|
||||
type: 'error',
|
||||
} satisfies WebSocketEvent<Record<string, unknown>>
|
||||
},
|
||||
peerAuthenticated(peerId: string, parentId?: string) {
|
||||
return {
|
||||
type: 'peer:authenticated',
|
||||
data: { authenticated: true, peerId },
|
||||
metadata: createEventMetadata(serverInstanceId, parentId),
|
||||
type: 'peer:authenticated',
|
||||
} satisfies WebSocketEvent<Record<string, unknown>>
|
||||
},
|
||||
extensionAuthenticated(identity: ExtensionIdentity, parentId?: string) {
|
||||
return {
|
||||
type: 'extension:authenticated',
|
||||
data: { identity, 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>>
|
||||
},
|
||||
}
|
||||
|
||||
@@ -4,40 +4,6 @@ import type { RouteContext, RouteDecision, RouteMiddleware } from '../../middlew
|
||||
|
||||
import { getProtocolEventMetadata } from '@proj-airi/server-shared/types'
|
||||
|
||||
/**
|
||||
* 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: {
|
||||
destinations?: RouteContext['destinations']
|
||||
event: WebSocketEvent
|
||||
fromPeer: RouteContext['fromPeer']
|
||||
middleware: RouteMiddleware[]
|
||||
peers: Map<string, RouteContext['fromPeer']>
|
||||
}): RouteDecision | undefined {
|
||||
const context: RouteContext = {
|
||||
destinations: input.destinations,
|
||||
event: input.event,
|
||||
fromPeer: input.fromPeer,
|
||||
peers: input.peers,
|
||||
}
|
||||
|
||||
for (const middleware of input.middleware) {
|
||||
const result = middleware(context)
|
||||
if (result) {
|
||||
return result
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Resolves the effective event delivery policy.
|
||||
*
|
||||
@@ -65,3 +31,37 @@ export function resolveEventDelivery(event: WebSocketEvent): DeliveryConfig | un
|
||||
...routeDelivery,
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 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
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -37,8 +37,8 @@ vi.mock('h3', () => ({
|
||||
get = vi.fn()
|
||||
},
|
||||
serve: vi.fn(() => ({
|
||||
close: serveMocks.closeCall,
|
||||
serve: serveMocks.serveCall,
|
||||
close: serveMocks.closeCall,
|
||||
})),
|
||||
}))
|
||||
|
||||
@@ -110,11 +110,11 @@ describe('createServer', async () => {
|
||||
it('merges nested config updates instead of replacing sibling logger settings', async () => {
|
||||
const server = createServer({
|
||||
hostname: '127.0.0.1',
|
||||
port: 6121,
|
||||
logger: {
|
||||
app: { level: LogLevelString.Log },
|
||||
websocket: { format: Format.Pretty },
|
||||
},
|
||||
port: 6121,
|
||||
})
|
||||
|
||||
server.updateConfig({
|
||||
@@ -130,8 +130,8 @@ describe('createServer', async () => {
|
||||
expect(serveMocks.setupAppCall).toHaveBeenCalledWith(expect.objectContaining({
|
||||
logger: {
|
||||
app: {
|
||||
format: Format.Pretty,
|
||||
level: LogLevelString.Log,
|
||||
format: Format.Pretty,
|
||||
},
|
||||
websocket: {
|
||||
format: Format.Pretty,
|
||||
|
||||
@@ -12,28 +12,86 @@ import { serve } from 'h3'
|
||||
|
||||
import { normalizeLoggerConfig, setupApp } from '..'
|
||||
|
||||
export interface Server {
|
||||
getConnectionHost: () => string[]
|
||||
restart: () => Promise<void>
|
||||
start: () => Promise<void>
|
||||
stop: () => Promise<void>
|
||||
updateConfig: (newOptions: ServerOptions) => void
|
||||
}
|
||||
|
||||
export interface ServerOptions extends AppOptions {
|
||||
hostname?: string
|
||||
port?: number
|
||||
tlsConfig?: null | {
|
||||
hostname?: string
|
||||
tlsConfig?: {
|
||||
cert?: string
|
||||
key?: string
|
||||
passphrase?: string
|
||||
}
|
||||
} | null
|
||||
}
|
||||
|
||||
interface ServerInstance {
|
||||
close: (closeActiveConnections?: boolean) => Promise<void>
|
||||
}
|
||||
|
||||
export interface Server {
|
||||
getConnectionHost: () => string[]
|
||||
start: () => Promise<void>
|
||||
stop: () => Promise<void>
|
||||
restart: () => Promise<void>
|
||||
updateConfig: (newOptions: ServerOptions) => void
|
||||
}
|
||||
|
||||
function isAddressInUseError(error: unknown) {
|
||||
return typeof error === 'object'
|
||||
&& error !== null
|
||||
&& 'code' in error
|
||||
&& (error as NodeJS.ErrnoException).code === 'EADDRINUSE'
|
||||
}
|
||||
|
||||
/**
|
||||
* Collects local IP addresses that can be used to reach the server from the LAN.
|
||||
*
|
||||
* Use when:
|
||||
* - Building connection hints for `0.0.0.0` listeners
|
||||
* - Showing reachable addresses in logs or UI
|
||||
*
|
||||
* Expects:
|
||||
* - Virtual interfaces should be ignored to reduce noisy or misleading addresses
|
||||
*
|
||||
* Returns:
|
||||
* - A de-duplicated list of valid IP addresses discovered from the host network interfaces
|
||||
*/
|
||||
export function getLocalIPs(): string[] {
|
||||
const interfaces = networkInterfaces()
|
||||
const addresses = new Set<string>()
|
||||
|
||||
const VIRTUAL_INTERFACE_PREFIXES = [
|
||||
'vboxnet',
|
||||
'vmnet',
|
||||
'docker',
|
||||
'br-',
|
||||
'veth',
|
||||
'utun',
|
||||
'wg',
|
||||
'tap',
|
||||
'tun',
|
||||
]
|
||||
const isVirtualInterface = (name: string) =>
|
||||
VIRTUAL_INTERFACE_PREFIXES.some(prefix => name.startsWith(prefix))
|
||||
|
||||
for (const [name, entries] of Object.entries(interfaces)) {
|
||||
if (!entries)
|
||||
continue
|
||||
if (isVirtualInterface(name))
|
||||
continue
|
||||
|
||||
for (const entry of entries) {
|
||||
const rawAddress = entry.address
|
||||
if (!rawAddress)
|
||||
continue
|
||||
|
||||
const address = rawAddress.includes('%') ? rawAddress.split('%')[0] : rawAddress
|
||||
if (isIP(address))
|
||||
addresses.add(address)
|
||||
}
|
||||
}
|
||||
|
||||
return [...addresses]
|
||||
}
|
||||
|
||||
/**
|
||||
* Creates the websocket server controller for the AIRI runtime.
|
||||
*
|
||||
@@ -48,12 +106,12 @@ interface ServerInstance {
|
||||
* - Lifecycle helpers for starting, stopping, restarting, and updating server options
|
||||
*/
|
||||
export function createServer(opts?: ServerOptions): Server {
|
||||
let options = merge<ServerOptions>({ hostname: '127.0.0.1', port: 6121 }, opts)
|
||||
let options = merge<ServerOptions>({ port: 6121, hostname: '127.0.0.1' }, opts)
|
||||
|
||||
const { appLogFormat, appLogLevel } = normalizeLoggerConfig(options)
|
||||
const log = useLogg('@proj-airi/server-runtime/server').withLogLevelString(appLogLevel).withFormat(appLogFormat)
|
||||
let serverInstance: null | ServerInstance = null
|
||||
let startTask: null | Promise<void> = null
|
||||
let serverInstance: ServerInstance | null = null
|
||||
let startTask: Promise<void> | null = null
|
||||
|
||||
log.withFields({ hasTlsConfig: !!options?.tlsConfig }).log('creating server channel')
|
||||
|
||||
@@ -103,17 +161,17 @@ export function createServer(opts?: ServerOptions): Server {
|
||||
const hostname = options.hostname
|
||||
|
||||
const instance = serve(h3App.app, {
|
||||
plugins: [createH3CrossWsPlugin(crossWsApp)],
|
||||
port,
|
||||
hostname,
|
||||
tls: options?.tlsConfig || undefined,
|
||||
reusePort: true,
|
||||
silent: true,
|
||||
manual: true,
|
||||
gracefulShutdown: {
|
||||
forceTimeout: 0.5,
|
||||
gracefulTimeout: 0.5,
|
||||
},
|
||||
hostname,
|
||||
manual: true,
|
||||
plugins: [createH3CrossWsPlugin(crossWsApp)],
|
||||
port,
|
||||
reusePort: true,
|
||||
silent: true,
|
||||
tls: options?.tlsConfig || undefined,
|
||||
})
|
||||
|
||||
try {
|
||||
@@ -177,67 +235,9 @@ export function createServer(opts?: ServerOptions): Server {
|
||||
|
||||
return getLocalIPs()
|
||||
},
|
||||
restart,
|
||||
start,
|
||||
stop,
|
||||
restart,
|
||||
updateConfig,
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Collects local IP addresses that can be used to reach the server from the LAN.
|
||||
*
|
||||
* Use when:
|
||||
* - Building connection hints for `0.0.0.0` listeners
|
||||
* - Showing reachable addresses in logs or UI
|
||||
*
|
||||
* Expects:
|
||||
* - Virtual interfaces should be ignored to reduce noisy or misleading addresses
|
||||
*
|
||||
* Returns:
|
||||
* - A de-duplicated list of valid IP addresses discovered from the host network interfaces
|
||||
*/
|
||||
export function getLocalIPs(): string[] {
|
||||
const interfaces = networkInterfaces()
|
||||
const addresses = new Set<string>()
|
||||
|
||||
const VIRTUAL_INTERFACE_PREFIXES = [
|
||||
'vboxnet',
|
||||
'vmnet',
|
||||
'docker',
|
||||
'br-',
|
||||
'veth',
|
||||
'utun',
|
||||
'wg',
|
||||
'tap',
|
||||
'tun',
|
||||
]
|
||||
const isVirtualInterface = (name: string) =>
|
||||
VIRTUAL_INTERFACE_PREFIXES.some(prefix => name.startsWith(prefix))
|
||||
|
||||
for (const [name, entries] of Object.entries(interfaces)) {
|
||||
if (!entries)
|
||||
continue
|
||||
if (isVirtualInterface(name))
|
||||
continue
|
||||
|
||||
for (const entry of entries) {
|
||||
const rawAddress = entry.address
|
||||
if (!rawAddress)
|
||||
continue
|
||||
|
||||
const address = rawAddress.includes('%') ? rawAddress.split('%')[0] : rawAddress
|
||||
if (isIP(address))
|
||||
addresses.add(address)
|
||||
}
|
||||
}
|
||||
|
||||
return [...addresses]
|
||||
}
|
||||
|
||||
function isAddressInUseError(error: unknown) {
|
||||
return typeof error === 'object'
|
||||
&& error !== null
|
||||
&& 'code' in error
|
||||
&& (error as NodeJS.ErrnoException).code === 'EADDRINUSE'
|
||||
}
|
||||
|
||||
@@ -8,18 +8,18 @@ import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
|
||||
import { setupApp } from './index'
|
||||
|
||||
interface TestWebSocketHandler {
|
||||
close?: (peer: Peer, details?: { code?: number, reason?: string, wasClean?: unknown }) => void
|
||||
message?: (peer: Peer, message: { text: () => string }) => void
|
||||
open?: (peer: Peer) => void
|
||||
message?: (peer: Peer, message: { text: () => string }) => void
|
||||
close?: (peer: Peer, details?: { code?: number, reason?: string, wasClean?: unknown }) => void
|
||||
}
|
||||
|
||||
interface TestWsServer {
|
||||
accept: (
|
||||
adapter: { close?: () => void, id: string, send: (message: { text: () => string }) => number | void },
|
||||
adapter: { id: string, send: (message: { text: () => string }) => void | number, close?: () => void },
|
||||
options: { state: { rawPeer: Peer } },
|
||||
) => void
|
||||
peers: {
|
||||
get: (peerId: string) => undefined | { receive: (message: { text: () => string }) => void }
|
||||
get: (peerId: string) => { receive: (message: { text: () => string }) => void } | undefined
|
||||
}
|
||||
remove: (peerId: string, details?: { code?: number, reason?: string, wasClean?: unknown }) => void
|
||||
}
|
||||
@@ -38,52 +38,24 @@ vi.mock('h3', () => ({
|
||||
|
||||
vi.mock('@proj-airi/better-ws/server/h3', () => ({
|
||||
toH3Handler: vi.fn((server: TestWsServer, options: { state: (peer: Peer) => { rawPeer: Peer } }) => ({
|
||||
close(peer: Peer, details?: { code?: number, reason?: string, wasClean?: unknown }) {
|
||||
server.remove(peer.id, details)
|
||||
},
|
||||
message(peer: Peer, message: { text: () => string }) {
|
||||
server.peers.get(peer.id)?.receive(message)
|
||||
},
|
||||
open(peer: Peer) {
|
||||
server.accept({
|
||||
close: () => peer.close?.(),
|
||||
id: peer.id,
|
||||
send: message => peer.send(message.text()),
|
||||
close: () => peer.close?.(),
|
||||
}, {
|
||||
state: options.state(peer),
|
||||
})
|
||||
},
|
||||
message(peer: Peer, message: { text: () => string }) {
|
||||
server.peers.get(peer.id)?.receive(message)
|
||||
},
|
||||
close(peer: Peer, details?: { code?: number, reason?: string, wasClean?: unknown }) {
|
||||
server.remove(peer.id, details)
|
||||
},
|
||||
})),
|
||||
}))
|
||||
|
||||
function createExtensionModuleAnnounceEvent(): WebSocketEvent {
|
||||
return {
|
||||
data: {
|
||||
identity: {
|
||||
extension: {
|
||||
id: 'extension-1',
|
||||
},
|
||||
id: 'memory-module-1',
|
||||
},
|
||||
name: 'memory',
|
||||
possibleEvents: [],
|
||||
},
|
||||
metadata: {
|
||||
event: {
|
||||
id: 'announce-1',
|
||||
},
|
||||
source: {
|
||||
id: 'extension-1',
|
||||
kind: 'plugin',
|
||||
plugin: {
|
||||
id: 'extension-1',
|
||||
},
|
||||
},
|
||||
},
|
||||
type: 'extension:module:announce',
|
||||
}
|
||||
}
|
||||
|
||||
function createPeer(id: string) {
|
||||
const sent: string[] = []
|
||||
const send: Peer['send'] = (data) => {
|
||||
@@ -92,18 +64,23 @@ function createPeer(id: string) {
|
||||
|
||||
return {
|
||||
peer: {
|
||||
close: vi.fn(),
|
||||
id,
|
||||
remoteAddress: '127.0.0.1',
|
||||
request: { url: `/ws?id=${id}` },
|
||||
send: vi.fn(send),
|
||||
close: vi.fn(),
|
||||
request: { url: `/ws?id=${id}` },
|
||||
remoteAddress: '127.0.0.1',
|
||||
} satisfies Peer,
|
||||
sent,
|
||||
}
|
||||
}
|
||||
|
||||
function decodeEvents(sent: string[]) {
|
||||
return sent.map(message => parse<WebSocketEvent>(message))
|
||||
function wsHandler() {
|
||||
const handler = h3Mocks.handlers.get('/ws') as TestWebSocketHandler | undefined
|
||||
if (!handler) {
|
||||
throw new Error('Expected setupApp to register a /ws websocket handler.')
|
||||
}
|
||||
|
||||
return handler
|
||||
}
|
||||
|
||||
function sendEvent(
|
||||
@@ -114,13 +91,36 @@ function sendEvent(
|
||||
handler.message?.(peer, { text: () => stringify(event) })
|
||||
}
|
||||
|
||||
function wsHandler() {
|
||||
const handler = h3Mocks.handlers.get('/ws') as TestWebSocketHandler | undefined
|
||||
if (!handler) {
|
||||
throw new Error('Expected setupApp to register a /ws websocket handler.')
|
||||
}
|
||||
function decodeEvents(sent: string[]) {
|
||||
return sent.map(message => parse<WebSocketEvent>(message))
|
||||
}
|
||||
|
||||
return handler
|
||||
function createExtensionModuleAnnounceEvent(): WebSocketEvent {
|
||||
return {
|
||||
type: 'extension:module:announce',
|
||||
data: {
|
||||
name: 'memory',
|
||||
possibleEvents: [],
|
||||
identity: {
|
||||
id: 'memory-module-1',
|
||||
extension: {
|
||||
id: 'extension-1',
|
||||
},
|
||||
},
|
||||
},
|
||||
metadata: {
|
||||
source: {
|
||||
kind: 'plugin',
|
||||
id: 'extension-1',
|
||||
plugin: {
|
||||
id: 'extension-1',
|
||||
},
|
||||
},
|
||||
event: {
|
||||
id: 'announce-1',
|
||||
},
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
describe('setupApp websocket liveness', () => {
|
||||
@@ -148,17 +148,17 @@ describe('setupApp websocket liveness', () => {
|
||||
|
||||
expect(decodeEvents(observer.sent)).toEqual(expect.arrayContaining([
|
||||
expect.objectContaining({
|
||||
type: 'registry:modules:health:unhealthy',
|
||||
data: {
|
||||
name: 'memory',
|
||||
identity: {
|
||||
id: 'memory-module-1',
|
||||
extension: {
|
||||
id: 'extension-1',
|
||||
},
|
||||
id: 'memory-module-1',
|
||||
},
|
||||
name: 'memory',
|
||||
reason: 'heartbeat late',
|
||||
},
|
||||
type: 'registry:modules:health:unhealthy',
|
||||
}),
|
||||
]))
|
||||
|
||||
@@ -184,11 +184,11 @@ describe('setupApp websocket liveness', () => {
|
||||
expect(modulePeer.peer.close).toHaveBeenCalledOnce()
|
||||
expect(decodeEvents(observer.sent)).toEqual(expect.arrayContaining([
|
||||
expect.objectContaining({
|
||||
type: 'extension:module:de-announced',
|
||||
data: expect.objectContaining({
|
||||
name: 'memory',
|
||||
reason: 'heartbeat expired',
|
||||
}),
|
||||
type: 'extension:module:de-announced',
|
||||
}),
|
||||
]))
|
||||
|
||||
@@ -211,11 +211,11 @@ describe('setupApp websocket liveness', () => {
|
||||
|
||||
expect(decodeEvents(observer.sent)).toEqual(expect.arrayContaining([
|
||||
expect.objectContaining({
|
||||
type: 'extension:module:de-announced',
|
||||
data: expect.objectContaining({
|
||||
name: 'memory',
|
||||
reason: 'connection closed',
|
||||
}),
|
||||
type: 'extension:module:de-announced',
|
||||
}),
|
||||
]))
|
||||
|
||||
|
||||
@@ -1,5 +1,41 @@
|
||||
import type { ExtensionIdentity, ExtensionModuleIdentity } from '@proj-airi/server-shared/types'
|
||||
|
||||
export interface Peer {
|
||||
/**
|
||||
* Unique random [uuid v4](https://developer.mozilla.org/en-US/docs/Glossary/UUID) identifier for the peer.
|
||||
*/
|
||||
get id(): string
|
||||
send: (data: unknown, options?: {
|
||||
compress?: boolean
|
||||
}) => number | void | undefined
|
||||
close?: () => void
|
||||
/**
|
||||
* WebSocket lifecycle state (mirrors WebSocket.readyState)
|
||||
*/
|
||||
readyState?: number
|
||||
request?: {
|
||||
url?: string
|
||||
headers?: Headers
|
||||
}
|
||||
remoteAddress?: string
|
||||
}
|
||||
|
||||
export interface NamedPeer {
|
||||
name: string
|
||||
index?: number
|
||||
peer: Peer
|
||||
}
|
||||
|
||||
/**
|
||||
* Tracks one module announced by an extension over a websocket peer.
|
||||
*/
|
||||
export interface RegisteredExtensionModule {
|
||||
/** Human-readable module name used by registry sync and legacy routing lookup. */
|
||||
name: string
|
||||
/** Module identity scoped to the owning extension session. */
|
||||
identity: ExtensionModuleIdentity
|
||||
}
|
||||
|
||||
export enum WebSocketReadyState {
|
||||
CONNECTING = 0,
|
||||
OPEN = 1,
|
||||
@@ -9,53 +45,17 @@ export enum WebSocketReadyState {
|
||||
|
||||
export interface AuthenticatedPeer extends NamedPeer {
|
||||
authenticated: boolean
|
||||
/** Caller-supplied peer ids acknowledged during manual peer authentication. */
|
||||
peerIds?: Set<string>
|
||||
identity?: ExtensionModuleIdentity
|
||||
extensionIdentity?: ExtensionIdentity
|
||||
extensionModules?: Map<string, RegisteredExtensionModule>
|
||||
healthy?: boolean
|
||||
identity?: ExtensionModuleIdentity
|
||||
lastHeartbeatAt?: number
|
||||
healthy?: boolean
|
||||
/**
|
||||
* REVIEW: Legacy field name kept during the better-ws migration.
|
||||
* The value now stores peer silence duration in milliseconds, not a miss count.
|
||||
* Rename this with the server-runtime peer state cleanup.
|
||||
*/
|
||||
missedHeartbeats?: number
|
||||
/** Caller-supplied peer ids acknowledged during manual peer authentication. */
|
||||
peerIds?: Set<string>
|
||||
}
|
||||
|
||||
export interface NamedPeer {
|
||||
index?: number
|
||||
name: string
|
||||
peer: Peer
|
||||
}
|
||||
|
||||
export interface Peer {
|
||||
close?: () => void
|
||||
/**
|
||||
* Unique random [uuid v4](https://developer.mozilla.org/en-US/docs/Glossary/UUID) identifier for the peer.
|
||||
*/
|
||||
get id(): string
|
||||
/**
|
||||
* WebSocket lifecycle state (mirrors WebSocket.readyState)
|
||||
*/
|
||||
readyState?: number
|
||||
remoteAddress?: string
|
||||
request?: {
|
||||
headers?: Headers
|
||||
url?: string
|
||||
}
|
||||
send: (data: unknown, options?: {
|
||||
compress?: boolean
|
||||
}) => number | undefined | void
|
||||
}
|
||||
|
||||
/**
|
||||
* Tracks one module announced by an extension over a websocket peer.
|
||||
*/
|
||||
export interface RegisteredExtensionModule {
|
||||
/** Module identity scoped to the owning extension session. */
|
||||
identity: ExtensionModuleIdentity
|
||||
/** Human-readable module name used by registry sync and legacy routing lookup. */
|
||||
name: string
|
||||
}
|
||||
|
||||
@@ -1,13 +1,13 @@
|
||||
import { defineConfig } from 'tsdown'
|
||||
|
||||
export default defineConfig({
|
||||
clean: true,
|
||||
dts: true,
|
||||
entry: {
|
||||
'bin/run': 'src/bin/run.ts',
|
||||
'index': 'src/index.ts',
|
||||
'server': 'src/server/index.ts',
|
||||
'bin/run': 'src/bin/run.ts',
|
||||
},
|
||||
outDir: 'dist',
|
||||
target: 'node18',
|
||||
outDir: 'dist',
|
||||
clean: true,
|
||||
dts: true,
|
||||
})
|
||||
|
||||
Reference in New Issue
Block a user