fix(server-*): added simple message queue for helping selecting consumer

Close #1387
This commit is contained in:
Neko Ayaka
2026-03-29 23:26:29 +08:00
parent 8946abdfac
commit a9c281c9f2
9 changed files with 718 additions and 45 deletions
+197
View File
@@ -0,0 +1,197 @@
import type { WebSocketBaseEvent, WebSocketEvents } from '@proj-airi/server-shared/types'
import { describe, expect, it } from 'vitest'
import { resolveDeliveryConfig, selectConsumerPeerId } from './index'
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',
},
},
route: overrides.route,
}
}
describe('resolveDeliveryConfig', () => {
it('uses protocol event metadata defaults for input:text', () => {
const delivery = resolveDeliveryConfig(createInputTextEvent())
expect(delivery).toEqual({
mode: 'consumer-group',
group: 'chat-ingestion',
selection: 'first',
})
})
it('allows route delivery to override protocol defaults', () => {
const delivery = resolveDeliveryConfig(createInputTextEvent({
route: {
delivery: {
required: true,
selection: 'sticky',
stickyKey: 'discord-dm-user-1',
},
},
}))
expect(delivery).toEqual({
mode: 'consumer-group',
group: 'chat-ingestion',
required: true,
selection: 'sticky',
stickyKey: 'discord-dm-user-1',
})
})
it('returns explicit route delivery for events without protocol defaults', () => {
const delivery = resolveDeliveryConfig({
type: 'spark:notify',
data: {
id: 'spark-1',
eventId: 'spark-notify-1',
kind: 'ping',
urgency: 'soon',
headline: 'hello',
destinations: ['module:character'],
},
metadata: {
source: {
kind: 'plugin',
plugin: { id: 'stage-web' },
id: 'stage-web-instance',
},
event: {
id: 'event-2',
},
},
route: {
delivery: {
mode: 'consumer',
required: true,
},
},
})
expect(delivery).toEqual({
mode: 'consumer',
required: true,
})
})
})
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: [
{
peerId: 'stage-window-a',
priority: 10,
registeredAt: 2,
authenticated: true,
healthy: true,
},
{
peerId: 'stage-window-b',
priority: 20,
registeredAt: 3,
authenticated: true,
healthy: true,
},
{
peerId: 'stage-window-c',
priority: 30,
registeredAt: 1,
authenticated: true,
healthy: false,
},
],
})
expect(selectedPeerId).toBe('stage-window-b')
})
it('keeps sticky delivery on the same consumer when available', () => {
const stickyAssignments = new Map<string, string>()
const firstSelectedPeerId = selectConsumerPeerId({
eventType: 'input:text',
fromPeerId: 'discord-instance',
delivery: {
mode: 'consumer-group',
group: 'chat-ingestion',
selection: 'sticky',
stickyKey: 'discord-dm-user-1',
},
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({
eventType: 'input:text',
fromPeerId: 'discord-instance',
delivery: {
mode: 'consumer-group',
group: 'chat-ingestion',
selection: 'sticky',
stickyKey: 'discord-dm-user-1',
},
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,
})
expect(firstSelectedPeerId).toBe('stage-window-a')
expect(secondSelectedPeerId).toBe('stage-window-a')
})
})
+332 -4
View File
@@ -1,4 +1,8 @@
import type { MetadataEventSource, WebSocketEvent } from '@proj-airi/server-shared/types'
import type {
DeliveryConfig,
MetadataEventSource,
WebSocketEvent,
} from '@proj-airi/server-shared/types'
import type {
RouteContext,
@@ -16,7 +20,12 @@ import {
createInvalidJsonServerErrorMessage,
ServerErrorMessages,
} from '@proj-airi/server-shared'
import { MessageHeartbeat, MessageHeartbeatKind, WebSocketEventSource } from '@proj-airi/server-shared/types'
import {
getProtocolEventMetadata,
MessageHeartbeat,
MessageHeartbeatKind,
WebSocketEventSource,
} from '@proj-airi/server-shared/types'
import { defineWebSocketHandler, H3 } from 'h3'
import { nanoid } from 'nanoid'
import { parse, stringify } from 'superjson'
@@ -95,12 +104,117 @@ const RESPONSES = {
} satisfies Record<string, (...args: unknown[]) => WebSocketEvent<Record<string, unknown>>>
const DEFAULT_HEARTBEAT_TTL_MS = 60_000
const DEFAULT_CONSUMER_GROUP = 'default'
interface ConsumerRegistration {
event: string
group: string
peerId: string
priority: number
registeredAt: number
}
export interface ConsumerSelectionCandidate {
peerId: string
priority: number
registeredAt: number
authenticated: boolean
healthy?: boolean
}
// helper send function
function send(peer: Peer, event: WebSocketEvent<Record<string, unknown>> | string) {
peer.send(typeof event === 'string' ? event : stringify(event))
}
export function resolveDeliveryConfig(event: WebSocketEvent): DeliveryConfig | undefined {
const eventMetadata = getProtocolEventMetadata(event.type)
const defaultDelivery = eventMetadata?.delivery
const routeDelivery = event.route?.delivery
if (!defaultDelivery && !routeDelivery) {
return undefined
}
return {
...defaultDelivery,
...routeDelivery,
}
}
function getConsumerRegistryKey(event: string, group: string) {
return `${event}::${group}`
}
function normalizeConsumerGroup(mode: DeliveryConfig['mode'], 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
})
}
export function selectConsumerPeerId(options: {
eventType: string
fromPeerId: string
delivery?: DeliveryConfig
candidates: ConsumerSelectionCandidate[]
roundRobinCursor?: Map<string, number>
stickyAssignments?: Map<string, string>
}) {
const { candidates, delivery, eventType, fromPeerId } = options
if (!delivery || (delivery.mode !== 'consumer' && delivery.mode !== 'consumer-group')) {
return
}
const normalizedGroup = normalizeConsumerGroup(delivery.mode, delivery.group)
const registryKey = getConsumerRegistryKey(eventType, normalizedGroup)
const availableEntries = sortConsumers(
candidates
.filter(entry => entry.peerId !== fromPeerId)
.filter(entry => entry.authenticated && entry.healthy !== false),
)
if (availableEntries.length === 0) {
return
}
const selection = delivery.selection ?? 'first'
if (selection === 'sticky' && delivery.stickyKey) {
const stickyRegistryKey = `${registryKey}::${delivery.stickyKey}`
const stickyPeerId = options.stickyAssignments?.get(stickyRegistryKey)
if (stickyPeerId && stickyPeerId !== fromPeerId) {
const stickyCandidate = availableEntries.find(entry => entry.peerId === stickyPeerId)
if (stickyCandidate) {
return stickyPeerId
}
}
const selected = availableEntries[0]
options.stickyAssignments?.set(stickyRegistryKey, 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
}
export interface AppOptions {
instanceId?: string
auth?: {
@@ -150,6 +264,10 @@ export function setupApp(options?: AppOptions): { app: H3, closeAllPeers: () =>
const peers = new Map<string, AuthenticatedPeer>()
const peersByModule = new Map<string, Map<number | undefined, AuthenticatedPeer>>()
const consumerRegistry = new Map<string, Map<string, Map<string, ConsumerRegistration>>>()
const consumerKeysByPeer = new Map<string, Set<string>>()
const deliveryRoundRobinCursor = new Map<string, number>()
const stickyAssignments = new Map<string, string>()
const heartbeatTtlMs = options?.heartbeat?.readTimeout ?? DEFAULT_HEARTBEAT_TTL_MS
const heartbeatMessage = options?.heartbeat?.message ?? MessageHeartbeat.Pong
const routingMiddleware = [
@@ -169,7 +287,6 @@ export function setupApp(options?: AppOptions): { app: H3, closeAllPeers: () =>
}
const elapsed = now - peerInfo.lastHeartbeatAt
if (elapsed > healthCheckIntervalMs) {
peerInfo.missedHeartbeats = (peerInfo.missedHeartbeats ?? 0) + 1
}
@@ -186,6 +303,7 @@ export function setupApp(options?: AppOptions): { app: H3, closeAllPeers: () =>
catch (error) {
logger.withFields({ peer: id, peerName: peerInfo.name }).withError(error as Error).debug('failed to close expired peer')
}
peers.delete(id)
unregisterModulePeer(peerInfo, 'heartbeat expired')
}
@@ -218,7 +336,133 @@ export function setupApp(options?: AppOptions): { app: H3, closeAllPeers: () =>
broadcastRegistrySync()
}
function registerConsumer(peerId: string, event: string, mode: DeliveryConfig['mode'], group?: string, priority?: number) {
const normalizedGroup = normalizeConsumerGroup(mode, group)
const registryKey = getConsumerRegistryKey(event, normalizedGroup)
let groups = consumerRegistry.get(event)
if (!groups) {
groups = new Map()
consumerRegistry.set(event, groups)
}
let peersForGroup = groups.get(normalizedGroup)
if (!peersForGroup) {
peersForGroup = new Map()
groups.set(normalizedGroup, peersForGroup)
}
peersForGroup.set(peerId, {
event,
group: normalizedGroup,
peerId,
priority: priority ?? 0,
registeredAt: Date.now(),
})
let registrations = consumerKeysByPeer.get(peerId)
if (!registrations) {
registrations = new Set()
consumerKeysByPeer.set(peerId, registrations)
}
registrations.add(registryKey)
}
function unregisterConsumer(peerId: string, event: string, mode: DeliveryConfig['mode'], group?: string) {
const normalizedGroup = normalizeConsumerGroup(mode, group)
const registryKey = getConsumerRegistryKey(event, normalizedGroup)
const groups = consumerRegistry.get(event)
const peersForGroup = groups?.get(normalizedGroup)
peersForGroup?.delete(peerId)
if (peersForGroup?.size === 0) {
groups?.delete(normalizedGroup)
deliveryRoundRobinCursor.delete(registryKey)
}
if (groups?.size === 0) {
consumerRegistry.delete(event)
}
const registrations = consumerKeysByPeer.get(peerId)
registrations?.delete(registryKey)
if (registrations?.size === 0) {
consumerKeysByPeer.delete(peerId)
}
for (const [stickyKey, stickyPeerId] of stickyAssignments.entries()) {
if (stickyPeerId === peerId && stickyKey.startsWith(`${registryKey}::`)) {
stickyAssignments.delete(stickyKey)
}
}
}
function unregisterPeerConsumers(peerId: string) {
const registrations = consumerKeysByPeer.get(peerId)
if (!registrations?.size) {
return
}
for (const registration of registrations) {
const [event, group] = registration.split('::', 2)
const groups = consumerRegistry.get(event)
const peersForGroup = groups?.get(group)
peersForGroup?.delete(peerId)
if (peersForGroup?.size === 0) {
groups?.delete(group)
deliveryRoundRobinCursor.delete(registration)
}
if (groups?.size === 0) {
consumerRegistry.delete(event)
}
}
for (const [stickyKey, stickyPeerId] of stickyAssignments.entries()) {
if (stickyPeerId === peerId) {
stickyAssignments.delete(stickyKey)
}
}
consumerKeysByPeer.delete(peerId)
}
function isEligibleConsumer(peerId: string) {
const candidate = peers.get(peerId)
return Boolean(
candidate
&& candidate.authenticated
&& candidate.healthy !== false,
)
}
function selectConsumer(event: WebSocketEvent, fromPeerId: string, delivery?: DeliveryConfig) {
const entries = consumerRegistry
.get(event.type)
?.get(normalizeConsumerGroup(delivery?.mode, delivery?.group))
const selectedPeerId = selectConsumerPeerId({
eventType: event.type,
fromPeerId,
delivery,
candidates: Array.from(entries?.values() ?? [], entry => ({
peerId: entry.peerId,
priority: entry.priority,
registeredAt: entry.registeredAt,
authenticated: Boolean(peers.get(entry.peerId)?.authenticated),
healthy: peers.get(entry.peerId)?.healthy,
})),
roundRobinCursor: deliveryRoundRobinCursor,
stickyAssignments,
})
if (!selectedPeerId || !isEligibleConsumer(selectedPeerId)) {
return
}
return peers.get(selectedPeerId)
}
function unregisterModulePeer(p: AuthenticatedPeer, reason?: string) {
unregisterPeerConsumers(p.peer.id)
if (!p.name)
return
@@ -480,6 +724,51 @@ 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(instanceId, 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, instanceId, event.metadata?.event.id))
return
}
registerConsumer(peer.id, data.event, data.mode ?? (data.group ? 'consumer-group' : 'consumer'), data.group, data.priority)
return
}
case 'module:consumer:unregister': {
const p = peers.get(peer.id)
if (!p?.authenticated) {
send(peer, RESPONSES.notAuthenticated(instanceId, 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, instanceId, event.metadata?.event.id))
return
}
unregisterConsumer(peer.id, data.event, data.mode ?? (data.group ? 'consumer-group' : 'consumer'), data.group)
return
}
}
// default case
@@ -495,6 +784,7 @@ export function setupApp(options?: AppOptions): { app: H3, closeAllPeers: () =>
const allowBypass = options?.routing?.allowBypass !== false
const shouldBypass = Boolean(event.route?.bypass && allowBypass && isDevtoolsPeer(p))
const destinations = shouldBypass ? undefined : collectDestinations(event)
const delivery = shouldBypass ? undefined : resolveDeliveryConfig(event)
const routingContext: RouteContext = {
event,
fromPeer: p,
@@ -516,6 +806,44 @@ export function setupApp(options?: AppOptions): { app: H3, closeAllPeers: () =>
return
}
const selectedConsumer = selectConsumer(event, peer.id, delivery)
if (delivery && (delivery.mode === 'consumer' || delivery.mode === 'consumer-group')) {
if (!selectedConsumer) {
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, instanceId, event.metadata?.event.id))
}
return
}
try {
logger.withFields({
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({
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')
peers.delete(selectedConsumer.peer.id)
unregisterModulePeer(selectedConsumer, 'consumer send failed')
}
return
}
const targetIds = decision?.type === 'targets' ? decision.targetIds : undefined
const shouldBroadcast = decision?.type === 'broadcast' || !targetIds
@@ -540,7 +868,7 @@ export function setupApp(options?: AppOptions): { app: H3, closeAllPeers: () =>
other.peer.send(payload)
}
catch (err) {
logger.withFields({ fromPeer: peer.id, fromPeerName: p.name, toPeer: other.peer.id, toPeerName: other.name, event }).withError(err as Error).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')
peers.delete(id)