diff --git a/apps/stage-web/src/pages/settings/memory/index.vue b/apps/stage-web/src/pages/settings/memory/index.vue
index f49297e8c..7fc7bba10 100644
--- a/apps/stage-web/src/pages/settings/memory/index.vue
+++ b/apps/stage-web/src/pages/settings/memory/index.vue
@@ -1,9 +1,18 @@
diff --git a/packages/memory-pgvector/src/index.ts b/packages/memory-pgvector/src/index.ts
index 4c66daac2..09f1d617c 100644
--- a/packages/memory-pgvector/src/index.ts
+++ b/packages/memory-pgvector/src/index.ts
@@ -1,3 +1,4 @@
+import process from 'node:process'
import { Format, LogLevel, setGlobalFormat, setGlobalLogLevel } from '@guiiai/logg'
import { Client } from '@proj-airi/server-sdk'
import { runUntilSignal } from '@proj-airi/server-sdk/utils/node'
@@ -6,9 +7,17 @@ setGlobalFormat(Format.Pretty)
setGlobalLogLevel(LogLevel.Log)
async function main() {
- const _client = new Client<{ connectionString: string }>({ name: 'memory-pgvector' })
+ const client = new Client<{ connectionString: string }>({
+ name: 'memory-pgvector',
+ })
+
+ client.onEvent('module:configure', (_event) => {
+ })
runUntilSignal()
+
+ process.on('SIGINT', () => client.close())
+ process.on('SIGTERM', () => client.close())
}
main()
diff --git a/packages/server-runtime/src/index.ts b/packages/server-runtime/src/index.ts
index d5d70d465..cb8e13672 100644
--- a/packages/server-runtime/src/index.ts
+++ b/packages/server-runtime/src/index.ts
@@ -9,7 +9,7 @@ import { createApp, createRouter, defineWebSocketHandler } from 'h3'
setGlobalFormat(Format.Pretty)
setGlobalLogLevel(LogLevel.Log)
-function send(peer: Peer, event: WebSocketEvent) {
+function send(peer: Peer, event: WebSocketEvent>) {
peer.send(JSON.stringify(event))
}
@@ -58,20 +58,37 @@ function main() {
peers.set(peer.id, { peer, authenticated: true, name: event.data.name })
return
case 'ui:configure':
- peers.forEach((p) => {
+ if (event.data.moduleName === '') {
+ send(peer, { type: 'error', data: { message: 'the field \'moduleName\' can\'t be empty for event \'ui:configure\'' } })
+ return
+ }
+ if (typeof event.data.moduleIndex !== 'undefined' && typeof event.data.moduleIndex !== 'number') {
+ send(peer, { type: 'error', data: { message: 'the field \'moduleIndex\' must be a number for event \'ui:configure\'' } })
+ return
+ }
+ if (typeof event.data.moduleIndex !== 'undefined' && event.data.moduleIndex < 0) {
+ send(peer, { type: 'error', data: { message: 'the field \'moduleIndex\' must be a positive number for event \'ui:configure\'' } })
+ return
+ }
+
+ for (const [_id, p] of peers.entries()) {
if (p.name === '') {
- return
+ continue
}
- if ((typeof p.index !== 'undefined' && typeof event.data.moduleIndex !== 'undefined' && p.name === event.data.moduleName && p.index === event.data.moduleIndex)) {
- return
- }
- if (p.name !== event.data.moduleName) {
+ if (p.name === event.data.moduleName) {
+ if ((typeof p.index !== 'undefined' && typeof event.data.moduleIndex !== 'undefined' && p.index === event.data.moduleIndex)) {
+ send(p.peer, { type: 'module:configure', data: { config: event.data.config } })
+ return
+ }
+
+ send(p.peer, { type: 'module:configure', data: { config: event.data.config } })
return
}
- p.peer.send(JSON.stringify({ type: 'module:configure', data: { config: event.data.config } } as WebSocketEvent))
- })
+ continue
+ }
+ send(peer, { type: 'error', data: { message: 'module not found, it haven\'t announced it or the name was wrong' } })
return
}
if (!peers.get(peer.id)?.authenticated) {
diff --git a/packages/server-sdk/src/client.ts b/packages/server-sdk/src/client.ts
index 19d3957e0..465c4dcca 100644
--- a/packages/server-sdk/src/client.ts
+++ b/packages/server-sdk/src/client.ts
@@ -22,7 +22,8 @@ export class Client {
private websocket: WebSocket
private eventListeners: Map, WebSocketEvents[keyof WebSocketEvents]>) => void | Promise>> = new Map()
- private authenticateAttempts = 0
+ private reconnectAttempts = 0
+ private shouldClose = false
constructor(options: ClientOptions) {
this.opts = defu>, Required, 'name' | 'token'>>[]>(
@@ -38,53 +39,113 @@ export class Client {
)
if (this.opts.autoConnect) {
- this.connect()
+ try {
+ this.connect()
+ }
+ catch (err) {
+ console.error(err)
+ }
}
}
- connect() {
- if (this.connected)
+ async retryWithExponentialBackoff(fn: () => void | Promise, attempts = 0, maxAttempts = -1) {
+ if (maxAttempts !== -1 && attempts >= maxAttempts) {
+ console.error(`Maximum retry attempts (${maxAttempts}) reached`)
return
+ }
- this.websocket = new WebSocket(this.opts.url)
+ try {
+ await fn()
+ }
+ catch (err) {
+ console.error('Encountered an error when retrying', err)
+ await sleep(2 ** attempts * 1000)
+ await this.retryWithExponentialBackoff(fn, attempts++, maxAttempts)
+ }
+ }
- this.onEvent('module:authenticated', async (event) => {
- const auth = event.data.authenticated
- if (!auth) {
- this.authenticateAttempts++
- await sleep(2 ** this.authenticateAttempts * 1000)
- this.tryAuthenticate()
+ async tryReconnectWithExponentialBackoff() {
+ await this.retryWithExponentialBackoff(() => this._connect(), this.reconnectAttempts)
+ }
+
+ private _connect() {
+ return new Promise((resolve, reject) => {
+ if (this.shouldClose) {
+ resolve()
+ return
}
- else {
- this.tryAnnounce()
+
+ if (this.connected) {
+ resolve()
+ return
+ }
+
+ this.websocket = new WebSocket(this.opts.url)
+
+ this.onEvent('module:authenticated', async (event) => {
+ const auth = event.data.authenticated
+ if (!auth) {
+ this.retryWithExponentialBackoff(() => this.tryAuthenticate())
+ }
+ else {
+ this.tryAnnounce()
+ }
+ })
+
+ this.websocket.onerror = (event) => {
+ this.opts.onError?.(event)
+
+ if ('error' in event && event.error instanceof Error) {
+ if (event.error.message === 'Received network error or non-101 status code.') {
+ this.connected = false
+
+ if (!this.opts.autoReconnect) {
+ this.opts.onError?.(event)
+ this.opts.onClose?.()
+ reject(event.error)
+ return
+ }
+
+ reject(event.error)
+ }
+ }
+ }
+
+ this.websocket.onclose = () => {
+ this.opts.onClose?.()
+ this.connected = false
+
+ if (!this.opts.autoReconnect) {
+ this.opts.onClose?.()
+ }
+ else {
+ this.tryReconnectWithExponentialBackoff()
+ }
+ }
+
+ this.websocket.onmessage = (event) => {
+ this.handleMessage(event)
+ }
+
+ this.websocket.onopen = () => {
+ this.reconnectAttempts = 0
+
+ if (this.opts.token) {
+ this.tryAuthenticate()
+ }
+ else {
+ this.tryAnnounce()
+ }
+
+ this.connected = true
+
+ resolve()
}
})
+ }
- this.websocket.onerror = (event) => {
- this.opts.onError?.(event)
- }
-
- this.websocket.onmessage = this.handleMessage.bind(this)
-
- this.websocket.onopen = () => {
- if (this.opts.token) {
- this.tryAuthenticate()
- }
- else {
- this.tryAnnounce()
- }
-
- this.connected = true
- }
-
- this.websocket.onclose = () => {
- this.connected = false
- this.authenticateAttempts = 0
- this.opts.onClose?.()
- if (this.opts.autoReconnect) {
- this.connect()
- }
- }
+ async connect() {
+ await this.tryReconnectWithExponentialBackoff()
}
private tryAnnounce() {
@@ -104,13 +165,19 @@ export class Client {
}
private async handleMessage(event: any) {
- const data = JSON.parse(event.data) as WebSocketEvent
- const listeners = this.eventListeners.get(data.type)
- if (!listeners)
- return
+ try {
+ const data = JSON.parse(event.data) as WebSocketEvent
+ const listeners = this.eventListeners.get(data.type)
+ if (!listeners)
+ return
- for (const listener of listeners)
- await listener(data)
+ for (const listener of listeners)
+ await listener(data)
+ }
+ catch (err) {
+ console.error('Failed to parse message:', err)
+ this.opts.onError?.(err)
+ }
}
onEvent>(
@@ -133,4 +200,13 @@ export class Client {
sendRaw(data: string | ArrayBufferLike | ArrayBufferView): void {
this.websocket.send(data)
}
+
+ close(): void {
+ this.shouldClose = true
+
+ if (this.connected && this.websocket) {
+ this.websocket.close()
+ this.connected = false
+ }
+ }
}