From a9c281c9f21aff2bb3d51d8e618e87941331eafc Mon Sep 17 00:00:00 2001 From: Neko Ayaka Date: Sun, 29 Mar 2026 23:26:29 +0800 Subject: [PATCH] fix(server-*): added simple message queue for helping selecting consumer Close #1387 --- packages/electron-screen-capture/package.json | 2 +- packages/plugin-protocol/src/types/events.ts | 121 ++++++- packages/server-runtime/src/index.test.ts | 197 ++++++++++ packages/server-runtime/src/index.ts | 336 +++++++++++++++++- packages/server-shared/src/errors.ts | 24 ++ .../src/stores/mods/api/channel-server.ts | 2 + .../src/stores/mods/api/context-bridge.ts | 29 ++ pnpm-lock.yaml | 50 +-- pnpm-workspace.yaml | 2 +- 9 files changed, 718 insertions(+), 45 deletions(-) create mode 100644 packages/server-runtime/src/index.test.ts diff --git a/packages/electron-screen-capture/package.json b/packages/electron-screen-capture/package.json index aeffc4e76..ce1064412 100644 --- a/packages/electron-screen-capture/package.json +++ b/packages/electron-screen-capture/package.json @@ -51,7 +51,7 @@ }, "inlinedDependencies": { "@electron-toolkit/preload": "3.0.2", - "@moeru/eventa": "1.0.0-beta.2", + "@moeru/eventa": "1.0.0-beta.3", "async-mutex": "0.5.0", "nanoid": [ "5.1.6", diff --git a/packages/plugin-protocol/src/types/events.ts b/packages/plugin-protocol/src/types/events.ts index 1ac700b74..20e6cf85b 100644 --- a/packages/plugin-protocol/src/types/events.ts +++ b/packages/plugin-protocol/src/types/events.ts @@ -1,3 +1,4 @@ +import type { Eventa } from '@moeru/eventa' import type { AssistantMessage, CommonContentPart, Message, ToolMessage, UserMessage } from '@xsai/shared-chat' import { defineEventa } from '@moeru/eventa' @@ -394,6 +395,40 @@ export interface ModulePermissionError { recoverable?: boolean } +export type DeliveryMode = 'broadcast' | 'consumer' | 'consumer-group' + +export type DeliverySelectionStrategy = 'first' | 'round-robin' | 'priority' | 'sticky' + +export interface DeliveryConfig { + mode?: DeliveryMode + group?: string + required?: boolean + selection?: DeliverySelectionStrategy + stickyKey?: string +} + +export interface ProtocolEventaMetadata { + delivery?: DeliveryConfig +} + +export interface ProtocolEventaInvokeMetadata { + delivery?: Partial +} + +export type ProtocolEventa

+ = Eventa + +function defineProtocolEventa

( + id: string, + options?: { + inheritFrom?: ProtocolEventa

+ metadata?: ProtocolEventaMetadata + invokeMetadata?: ProtocolEventaInvokeMetadata + }, +): ProtocolEventa

{ + return defineEventa(id, options) +} + export type RouteTargetExpression = | { type: 'and', all: RouteTargetExpression[] } | { type: 'or', any: RouteTargetExpression[] } @@ -408,6 +443,7 @@ export type RouteTargetExpression export interface RouteConfig { destinations?: Array bypass?: boolean + delivery?: DeliveryConfig } export enum MessageHeartbeatKind { @@ -892,6 +928,19 @@ interface ModuleConfigureEvent { config: C | Record } +interface ModuleConsumerRegisterEvent { + event: string + mode?: Exclude + group?: string + priority?: number +} + +interface ModuleConsumerUnregisterEvent { + event: string + mode?: Exclude + group?: string +} + interface UiConfigureEvent { moduleName: string moduleIndex?: number @@ -1047,26 +1096,62 @@ export const moduleContributeCapabilityConfigurationCommitStatus = defineEventa< export const moduleContributeCapabilityConfigurationConfigured = defineEventa('module:contribute:capability:configuration:configured') export const moduleContributeCapabilityActivated = defineEventa('module:contribute:capability:activated') -export const moduleStatusChange = defineEventa('module:status:change') +export const moduleStatusChange = defineProtocolEventa('module:status:change') -export const moduleConfigure = defineEventa('module:configure') +export const moduleConfigure = defineProtocolEventa('module:configure') +export const moduleConsumerRegister = defineProtocolEventa('module:consumer:register') +export const moduleConsumerUnregister = defineProtocolEventa('module:consumer:unregister') -export const uiConfigure = defineEventa('ui:configure') +export const uiConfigure = defineProtocolEventa('ui:configure') -export const inputText = defineEventa('input:text') -export const inputTextVoice = defineEventa('input:text:voice') -export const inputVoice = defineEventa('input:voice') +export const inputText = defineProtocolEventa('input:text', { + metadata: { + delivery: { + mode: 'consumer-group', + group: 'chat-ingestion', + selection: 'first', + }, + }, +}) +export const inputTextVoice = defineProtocolEventa('input:text:voice', { + metadata: { + delivery: { + mode: 'consumer-group', + group: 'chat-ingestion', + selection: 'first', + }, + }, +}) +export const inputVoice = defineProtocolEventa('input:voice', { + metadata: { + delivery: { + mode: 'consumer-group', + group: 'chat-ingestion', + selection: 'first', + }, + }, +}) -export const outputGenAiChatToolCall = defineEventa('output:gen-ai:chat:tool-call') -export const outputGenAiChatMessage = defineEventa('output:gen-ai:chat:message') -export const outputGenAiChatComplete = defineEventa('output:gen-ai:chat:complete') +export const outputGenAiChatToolCall = defineProtocolEventa('output:gen-ai:chat:tool-call') +export const outputGenAiChatMessage = defineProtocolEventa('output:gen-ai:chat:message') +export const outputGenAiChatComplete = defineProtocolEventa('output:gen-ai:chat:complete') -export const sparkNotify = defineEventa('spark:notify') -export const sparkEmit = defineEventa('spark:emit') -export const sparkCommand = defineEventa('spark:command') +export const sparkNotify = defineProtocolEventa('spark:notify') +export const sparkEmit = defineProtocolEventa('spark:emit') +export const sparkCommand = defineProtocolEventa('spark:command') -export const transportConnectionHeartbeat = defineEventa('transport:connection:heartbeat') -export const contextUpdate = defineEventa('context:update') +export const transportConnectionHeartbeat = defineProtocolEventa('transport:connection:heartbeat') +export const contextUpdate = defineProtocolEventa('context:update') + +export const protocolEventMetadataByType = { + [inputText.id]: inputText.metadata, + [inputTextVoice.id]: inputTextVoice.metadata, + [inputVoice.id]: inputVoice.metadata, +} satisfies Partial> + +export function getProtocolEventMetadata(eventType: keyof ProtocolEvents | string) { + return protocolEventMetadataByType[eventType as keyof typeof protocolEventMetadataByType] +} // Thanks to: // @@ -1202,6 +1287,14 @@ export interface ProtocolEvents { * Push configuration down to module (host → module). */ 'module:configure': ModuleConfigureEvent + /** + * Register the current module instance as a consumer for an event or event group. + */ + 'module:consumer:register': ModuleConsumerRegisterEvent + /** + * Unregister the current module instance from an event consumer registration. + */ + 'module:consumer:unregister': ModuleConsumerUnregisterEvent 'ui:configure': UiConfigureEvent diff --git a/packages/server-runtime/src/index.test.ts b/packages/server-runtime/src/index.test.ts new file mode 100644 index 000000000..5834db207 --- /dev/null +++ b/packages/server-runtime/src/index.test.ts @@ -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']> { + 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() + + 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') + }) +}) diff --git a/packages/server-runtime/src/index.ts b/packages/server-runtime/src/index.ts index c3a867c98..89f93a7d3 100644 --- a/packages/server-runtime/src/index.ts +++ b/packages/server-runtime/src/index.ts @@ -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 WebSocketEvent>> 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> | 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>) { + 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 + stickyAssignments?: Map +}) { + 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() const peersByModule = new Map>() + const consumerRegistry = new Map>>() + const consumerKeysByPeer = new Map>() + const deliveryRoundRobinCursor = new Map() + const stickyAssignments = new Map() 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) diff --git a/packages/server-shared/src/errors.ts b/packages/server-shared/src/errors.ts index 3ba1dd4e3..6830da8fa 100644 --- a/packages/server-shared/src/errors.ts +++ b/packages/server-shared/src/errors.ts @@ -5,7 +5,9 @@ export const ServerErrorMessages = { moduleAnnounceIdentityInvalid: 'module identity must include kind=plugin and a plugin id for event \'module:announce\'', moduleAnnounceIndexInvalid: 'the field \'index\' must be a non-negative integer for event \'module:announce\'', moduleAnnounceNameInvalid: 'the field \'name\' must be a non-empty string for event \'module:announce\'', + moduleConsumerEventInvalid: 'the field \'event\' must be a non-empty string for event consumer registration', moduleNotFound: 'module not found, it hasn\'t announced itself or the name is incorrect', + noConsumerRegistered: 'no consumer registered for requested event delivery', notAuthenticated: 'not authenticated', uiConfigureModuleIndexInvalid: 'the field \'moduleIndex\' must be a non-negative integer for event \'ui:configure\'', uiConfigureModuleNameInvalid: 'the field \'moduleName\' can\'t be empty for event \'ui:configure\'', @@ -18,8 +20,10 @@ export type ServerErrorCode | 'module-announce-identity-invalid' | 'module-announce-index-invalid' | 'module-announce-name-invalid' + | 'module-consumer-event-invalid' | 'module-not-found' | 'must-authenticate-before-announcing' + | 'no-consumer-registered' | 'not-authenticated' | 'ui-configure-module-index-invalid' | 'ui-configure-module-name-invalid' @@ -128,6 +132,26 @@ export function parseServerErrorMessage(message: string): ParsedServerErrorMessa } } + if (message === ServerErrorMessages.moduleConsumerEventInvalid) { + return { + authentication: false, + code: 'module-consumer-event-invalid', + message, + recoverable: false, + terminal: false, + } + } + + if (message === ServerErrorMessages.noConsumerRegistered) { + return { + authentication: false, + code: 'no-consumer-registered', + message, + recoverable: true, + terminal: false, + } + } + if (message === ServerErrorMessages.uiConfigureModuleNameInvalid) { return { authentication: false, diff --git a/packages/stage-ui/src/stores/mods/api/channel-server.ts b/packages/stage-ui/src/stores/mods/api/channel-server.ts index 507701e35..ff47531d8 100644 --- a/packages/stage-ui/src/stores/mods/api/channel-server.ts +++ b/packages/stage-ui/src/stores/mods/api/channel-server.ts @@ -26,6 +26,8 @@ export const useModsServerChannelStore = defineStore('mods:channels:proj-airi:se 'error', 'module:announce', 'module:configure', + 'module:consumer:register', + 'module:consumer:unregister', 'module:authenticated', 'spark:notify', 'spark:emit', diff --git a/packages/stage-ui/src/stores/mods/api/context-bridge.ts b/packages/stage-ui/src/stores/mods/api/context-bridge.ts index ef8c2814a..93609c2da 100644 --- a/packages/stage-ui/src/stores/mods/api/context-bridge.ts +++ b/packages/stage-ui/src/stores/mods/api/context-bridge.ts @@ -34,6 +34,11 @@ export function normalizeContextSnapshot { + const consumerRegistrationEvents = [ + 'input:text', + 'input:text:voice', + 'input:voice', + ] as const const mutex = new Mutex() const chatOrchestrator = useChatOrchestratorStore() @@ -55,6 +60,19 @@ export const useContextBridgeStore = defineStore('mods:api:context-bridge', () = await mutex.acquire() try { + await serverChannelStore.ensureConnected() + + for (const consumerEvent of consumerRegistrationEvents) { + serverChannelStore.send({ + type: 'module:consumer:register', + data: { + event: consumerEvent, + mode: 'consumer-group', + group: 'chat-ingestion', + }, + }) + } + let isProcessingRemoteStream = false const { stop } = watch(incomingContext, (event) => { @@ -348,6 +366,17 @@ export const useContextBridgeStore = defineStore('mods:api:context-bridge', () = await mutex.acquire() try { + for (const consumerEvent of consumerRegistrationEvents) { + serverChannelStore.send({ + type: 'module:consumer:unregister', + data: { + event: consumerEvent, + mode: 'consumer-group', + group: 'chat-ingestion', + }, + }) + } + for (const fn of disposeHookFns.value) { fn() } diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index ecc406428..1c265e96f 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -55,8 +55,8 @@ catalogs: specifier: 0.1.0-beta.15 version: 0.1.0-beta.15 '@moeru/eventa': - specifier: 1.0.0-beta.2 - version: 1.0.0-beta.2 + specifier: 1.0.0-beta.3 + version: 1.0.0-beta.3 '@moeru/std': specifier: 0.1.0-beta.17 version: 0.1.0-beta.17 @@ -529,7 +529,7 @@ importers: version: 1.3.0(@hono/node-server@1.19.11(hono@4.11.3))(bufferutil@4.1.0)(hono@4.11.3)(utf-8-validate@5.0.10) '@moeru/eventa': specifier: 'catalog:' - version: 1.0.0-beta.2(electron@41.0.3) + version: 1.0.0-beta.3(electron@41.0.3) '@moeru/std': specifier: 'catalog:' version: 0.1.0-beta.17 @@ -659,7 +659,7 @@ importers: version: 3.8.1 '@moeru/eventa': specifier: 'catalog:' - version: 1.0.0-beta.2(electron@41.0.3) + version: 1.0.0-beta.3(electron@41.0.3) '@moeru/std': specifier: 'catalog:' version: 0.1.0-beta.17 @@ -1056,7 +1056,7 @@ importers: version: 11.3.0 '@moeru/eventa': specifier: 'catalog:' - version: 1.0.0-beta.2(electron@40.8.3)(h3@2.0.1-rc.19(crossws@0.4.4(srvx@0.11.13(patch_hash=c761d25a70e22a0925c88fe91a9dd1c3153dd341d8fd4f260166678156c4df2d)))) + version: 1.0.0-beta.3(electron@40.8.3)(h3@2.0.1-rc.19(crossws@0.4.4(srvx@0.11.13(patch_hash=c761d25a70e22a0925c88fe91a9dd1c3153dd341d8fd4f260166678156c4df2d)))) '@moeru/std': specifier: 'catalog:' version: 0.1.0-beta.17 @@ -1480,7 +1480,7 @@ importers: version: 3.8.1 '@moeru/eventa': specifier: 'catalog:' - version: 1.0.0-beta.2(electron@41.0.3) + version: 1.0.0-beta.3(electron@41.0.3) '@moeru/std': specifier: 'catalog:' version: 0.1.0-beta.17 @@ -2081,7 +2081,7 @@ importers: dependencies: '@moeru/eventa': specifier: 'catalog:' - version: 1.0.0-beta.2(electron@40.8.3)(h3@2.0.1-rc.19(crossws@0.4.4(srvx@0.11.13(patch_hash=c761d25a70e22a0925c88fe91a9dd1c3153dd341d8fd4f260166678156c4df2d)))) + version: 1.0.0-beta.3(electron@40.8.3)(h3@2.0.1-rc.19(crossws@0.4.4(srvx@0.11.13(patch_hash=c761d25a70e22a0925c88fe91a9dd1c3153dd341d8fd4f260166678156c4df2d)))) builder-util-runtime: specifier: 'catalog:' version: 9.5.1 @@ -2109,7 +2109,7 @@ importers: version: 3.0.2(electron@39.7.0) '@moeru/eventa': specifier: 'catalog:' - version: 1.0.0-beta.2(electron@39.7.0) + version: 1.0.0-beta.3(electron@39.7.0) '@moeru/std': specifier: 'catalog:' version: 0.1.0-beta.17 @@ -2127,7 +2127,7 @@ importers: dependencies: '@moeru/eventa': specifier: 'catalog:' - version: 1.0.0-beta.2(electron@40.8.3)(h3@2.0.1-rc.19(crossws@0.4.4(srvx@0.11.13(patch_hash=c761d25a70e22a0925c88fe91a9dd1c3153dd341d8fd4f260166678156c4df2d)))) + version: 1.0.0-beta.3(electron@40.8.3)(h3@2.0.1-rc.19(crossws@0.4.4(srvx@0.11.13(patch_hash=c761d25a70e22a0925c88fe91a9dd1c3153dd341d8fd4f260166678156c4df2d)))) '@moeru/std': specifier: 'catalog:' version: 0.1.0-beta.17 @@ -2220,7 +2220,7 @@ importers: dependencies: '@moeru/eventa': specifier: 'catalog:' - version: 1.0.0-beta.2(electron@41.0.3) + version: 1.0.0-beta.3(electron@41.0.3) '@moeru/std': specifier: 'catalog:' version: 0.1.0-beta.17 @@ -2232,7 +2232,7 @@ importers: dependencies: '@moeru/eventa': specifier: 'catalog:' - version: 1.0.0-beta.2(electron@41.0.3) + version: 1.0.0-beta.3(electron@41.0.3) '@xsai/shared-chat': specifier: 'catalog:' version: 0.4.0-beta.13 @@ -2241,7 +2241,7 @@ importers: dependencies: '@moeru/eventa': specifier: 'catalog:' - version: 1.0.0-beta.2(electron@41.0.3) + version: 1.0.0-beta.3(electron@41.0.3) '@proj-airi/plugin-protocol': specifier: workspace:* version: link:../plugin-protocol @@ -2310,7 +2310,7 @@ importers: dependencies: '@moeru/eventa': specifier: 'catalog:' - version: 1.0.0-beta.2(electron@41.0.3) + version: 1.0.0-beta.3(electron@41.0.3) packages/server-shared: dependencies: @@ -2322,7 +2322,7 @@ importers: dependencies: '@moeru/eventa': specifier: 'catalog:' - version: 1.0.0-beta.2(electron@41.0.3) + version: 1.0.0-beta.3(electron@41.0.3) '@moeru/std': specifier: 'catalog:' version: 0.1.0-beta.17 @@ -2425,7 +2425,7 @@ importers: dependencies: '@moeru/eventa': specifier: 'catalog:' - version: 1.0.0-beta.2(electron@41.0.3) + version: 1.0.0-beta.3(electron@41.0.3) '@moeru/std': specifier: 'catalog:' version: 0.1.0-beta.17 @@ -2559,7 +2559,7 @@ importers: version: 3.0.2(electron@41.0.3) '@moeru/eventa': specifier: 'catalog:' - version: 1.0.0-beta.2(electron@41.0.3) + version: 1.0.0-beta.3(electron@41.0.3) '@nekopaw/tempora': specifier: 'catalog:' version: 0.4.0-alpha.1 @@ -2583,7 +2583,7 @@ importers: version: 3.8.1 '@moeru/eventa': specifier: 'catalog:' - version: 1.0.0-beta.2(electron@41.0.3) + version: 1.0.0-beta.3(electron@41.0.3) '@proj-airi/audio': specifier: workspace:^ version: link:../audio @@ -3041,7 +3041,7 @@ importers: dependencies: '@moeru/eventa': specifier: 'catalog:' - version: 1.0.0-beta.2(electron@41.0.3) + version: 1.0.0-beta.3(electron@41.0.3) '@pixiv/three-vrm': specifier: ^3.5.1 version: 3.5.1(@types/three@0.183.1)(three@0.183.2) @@ -3311,7 +3311,7 @@ importers: version: 1.2.4 '@moeru/eventa': specifier: 'catalog:' - version: 1.0.0-beta.2(electron@41.0.3) + version: 1.0.0-beta.3(electron@41.0.3) '@proj-airi/server-sdk': specifier: workspace:^ version: link:../../packages/server-sdk @@ -6171,8 +6171,8 @@ packages: web-worker: optional: true - '@moeru/eventa@1.0.0-beta.2': - resolution: {integrity: sha512-Sr4GMRaFNURj72coeK3ZoPz8bvOTSJopTmRcHiVjaae2U3bzrz1gNZXcaN3dRSvd+3Vxm7VPrsUGFHkcPZ2c0w==} + '@moeru/eventa@1.0.0-beta.3': + resolution: {integrity: sha512-ug/jrF2MNiJBVkqgxWUMAePRH6LDW8pItRVkrT744Rz4CCE0A7BJBGYg4QJnumSZP0fXkNxCnZ2Q+7BCjtUNdA==} peerDependencies: electron: '>=30' h3: 2.0.0-beta.1 @@ -20520,14 +20520,14 @@ snapshots: optionalDependencies: electron: 41.0.3 - '@moeru/eventa@1.0.0-beta.2(electron@39.7.0)': + '@moeru/eventa@1.0.0-beta.3(electron@39.7.0)': dependencies: nanoid: 5.1.7 picomatch: 4.0.3 optionalDependencies: electron: 39.7.0 - '@moeru/eventa@1.0.0-beta.2(electron@40.8.3)(h3@2.0.1-rc.19(crossws@0.4.4(srvx@0.11.13(patch_hash=c761d25a70e22a0925c88fe91a9dd1c3153dd341d8fd4f260166678156c4df2d))))': + '@moeru/eventa@1.0.0-beta.3(electron@40.8.3)(h3@2.0.1-rc.19(crossws@0.4.4(srvx@0.11.13(patch_hash=c761d25a70e22a0925c88fe91a9dd1c3153dd341d8fd4f260166678156c4df2d))))': dependencies: nanoid: 5.1.7 picomatch: 4.0.3 @@ -20535,7 +20535,7 @@ snapshots: electron: 40.8.3 h3: 2.0.1-rc.19(crossws@0.4.4(srvx@0.11.13(patch_hash=c761d25a70e22a0925c88fe91a9dd1c3153dd341d8fd4f260166678156c4df2d))) - '@moeru/eventa@1.0.0-beta.2(electron@41.0.3)': + '@moeru/eventa@1.0.0-beta.3(electron@41.0.3)': dependencies: nanoid: 5.1.7 picomatch: 4.0.3 @@ -23557,7 +23557,7 @@ snapshots: '@typescript-eslint/project-service@8.56.1(typescript@5.9.3)': dependencies: - '@typescript-eslint/tsconfig-utils': 8.56.1(typescript@5.9.3) + '@typescript-eslint/tsconfig-utils': 8.57.1(typescript@5.9.3) '@typescript-eslint/types': 8.57.1 debug: 4.4.3 typescript: 5.9.3 diff --git a/pnpm-workspace.yaml b/pnpm-workspace.yaml index 16ca0b410..b73747f05 100644 --- a/pnpm-workspace.yaml +++ b/pnpm-workspace.yaml @@ -43,7 +43,7 @@ catalog: '@intlify/core': ^11.3.0 '@modelcontextprotocol/sdk': ^1.27.1 '@moeru/eslint-config': 0.1.0-beta.15 - '@moeru/eventa': 1.0.0-beta.2 + '@moeru/eventa': 1.0.0-beta.3 '@moeru/std': 0.1.0-beta.17 '@nekopaw/tempora': 0.4.0-alpha.1 '@pinia/testing': ^1.0.3