diff --git a/examples/mastra-weather-agent/weather-agent.ts b/examples/mastra-weather-agent/weather-agent.ts index 7a70576..70bb269 100644 --- a/examples/mastra-weather-agent/weather-agent.ts +++ b/examples/mastra-weather-agent/weather-agent.ts @@ -146,12 +146,13 @@ async function main(): Promise { const agentName = process.env.AGENT_NAME ?? 'mastra-weather-agent'; const apiKey = process.env.LLP_API_KEY ?? ''; const model = process.env.MODEL_NAME ?? 'ollama-cloud/gpt-oss:120b'; + const platformUrl = process.env.LLP_URL; if (!apiKey) { throw new Error('LLP_API_KEY env var is not defined'); } - const client = new LLPClient(agentName, apiKey) + const client = new LLPClient(agentName, apiKey, { url: platformUrl }) .onStart(() => createWeatherAgent(model)) .onMessage(async (agent, msg, annotater) => { const response = await handleMessage(agent, msg, annotater); diff --git a/examples/simple-agent/simple-agent.ts b/examples/simple-agent/simple-agent.ts index fda55ec..ef9a361 100644 --- a/examples/simple-agent/simple-agent.ts +++ b/examples/simple-agent/simple-agent.ts @@ -1,5 +1,5 @@ import { config } from 'dotenv'; -import { LLPClient, type LLPSession, type TextMessage } from '../../src/index.js'; +import { type Annotater, LLPClient, type TextMessage } from '../../src/index.js'; async function main() { config(); @@ -15,13 +15,13 @@ async function main() { throw new Error('LLP_API_KEY env var is not defined'); } - const client = new LLPClient('simple-agent', apiKey, { url: platformUrl }); - - client.onMessage(async (session: LLPSession, msg: TextMessage) => { - const toolCall = msg.toolCall('get_weather', '{"city":"Seattle"}', 'rainy', 1_000); - await session.annotateToolCall(toolCall); - return msg.reply('this is my response'); - }); + const client = new LLPClient('simple-agent', apiKey, { url: platformUrl }).onMessage( + async (_data: unknown, msg: TextMessage, annotater: Annotater) => { + const toolCall = msg.toolCall('get_weather', '{"city":"Seattle"}', 'rainy', 1_000); + await annotater.annotateToolCall(toolCall); + return msg.reply('this is my response'); + }, + ); try { console.log('Connecting to server...'); diff --git a/src/client.ts b/src/client.ts index cb35727..2fe9326 100644 --- a/src/client.ts +++ b/src/client.ts @@ -5,14 +5,13 @@ import { AlreadyClosedError, AlreadyConnectedError, type ErrorCode, + MessageHandlerNotSetError, NotAuthenticatedError, - NotConnectedError, PlatformError, TimeoutError, } from './errors.js'; import { PresenceMessage, TextMessage } from './message.js'; import { ConnectionStatus, PresenceStatus } from './presence.js'; -import { LLPSession } from './session.js'; import type { ToolCall } from './tool_call.js'; export interface LLPClientConfig { @@ -22,27 +21,9 @@ export interface LLPClientConfig { readonly maxQueueSize?: number; // Default: 32 } -// --------------------------------------------------------------------------- -// Legacy handler types (backward compat) -// --------------------------------------------------------------------------- - -export type MessageHandler = ( - session: LLPSession, - msg: TextMessage, -) => Promise; -export type StartHandler = ( - session: LLPSession, - msg: PresenceMessage, -) => TSessionData | undefined | Promise; -export type PresenceHandler = StartHandler; - -// --------------------------------------------------------------------------- -// New chainable API types -// --------------------------------------------------------------------------- - -export type SessionCreator = () => T | Promise; -export type TypedMessageHandler = ( - data: T, +export type SessionCreator = () => T; +export type MessageHandler = ( + agent: T, msg: TextMessage, annotater: Annotater, ) => Promise; @@ -53,15 +34,13 @@ export type StopHandler = (data: T) => void | Promise; * The generic `T` is inferred from the factory return type. */ export interface TypedLLPClient { - onMessage(handler: TypedMessageHandler): TypedLLPClient; + onMessage(handler: MessageHandler): TypedLLPClient; onStop(handler: StopHandler): TypedLLPClient; connect(timeout?: number): Promise; close(): Promise; getStatus(): ConnectionStatus; getSessionId(): string | null; getPresence(): PresenceStatus; - sendMessage(msg: TextMessage, timeout?: number): Promise; - sendAsyncMessage(msg: TextMessage): Promise; annotateToolCall(toolCall: ToolCall): Promise; } @@ -72,22 +51,10 @@ export class LLPClient implements Annotater { private presence: PresenceStatus = PresenceStatus.Unavailable; private readonly outboundQueue: Array = []; - private readonly pending = new Map< - string, - { - resolve: (msg: TextMessage) => void; - reject: (err: Error) => void; - } - >(); - private readonly sessions = new Map>(); + private readonly sessions = new Map(); - // Legacy handlers - private messageHandler: MessageHandler | null = null; - private startHandler: StartHandler | null = null; - - // New API handlers private sessionCreator: SessionCreator | null = null; - private typedMessageHandler: TypedMessageHandler | null = null; + private messageHandler: MessageHandler | null = null; private stopHandler: StopHandler | null = null; private authResolve: (() => void) | null = null; @@ -101,10 +68,6 @@ export class LLPClient implements Annotater { private readonly config: LLPClientConfig = {}, ) {} - // ------------------------------------------------------------------- - // New chainable API - // ------------------------------------------------------------------- - /** * Register a callback that creates session data for each new user session. * The SDK calls this automatically when a user becomes available. @@ -119,48 +82,26 @@ export class LLPClient implements Annotater { * .connect(); * ``` */ - onStart(creator: SessionCreator): TypedLLPClient; - /** @deprecated Use the new overload: `onStart(() => createAgent(...))` */ - onStart(handler: StartHandler): void; - onStart( - creatorOrHandler: SessionCreator | StartHandler, - ): TypedLLPClient | undefined { + onStart(creator: SessionCreator): TypedLLPClient { if (this.status !== ConnectionStatus.Disconnected && this.status !== ConnectionStatus.Closed) { throw new AlreadyConnectedError('Must set onStart callback before connecting'); } - // Distinguish by arity: creator is 0-arg, legacy handler is 2-arg - if (creatorOrHandler.length === 0) { - this.sessionCreator = creatorOrHandler as SessionCreator; - return this as unknown as TypedLLPClient; - } - - this.startHandler = creatorOrHandler as StartHandler; - return undefined; + this.sessionCreator = creator as SessionCreator; + return this as unknown as TypedLLPClient; } /** * Register a message handler (new API — called via chaining after onStart). * The first argument is the session data created by onStart. */ - onMessage(handler: TypedMessageHandler): TypedLLPClient; - /** @deprecated Use the new API: `onStart(...).onMessage((data, msg) => ...)` */ - onMessage(handler: MessageHandler): void; - onMessage( - handler: TypedMessageHandler | MessageHandler, - ): TypedLLPClient | undefined { + onMessage(handler: MessageHandler): TypedLLPClient { if (this.status !== ConnectionStatus.Disconnected && this.status !== ConnectionStatus.Closed) { throw new AlreadyConnectedError('Must set onMessage callback before connecting'); } - // If sessionCreator is set, we're in the new API path - if (this.sessionCreator) { - this.typedMessageHandler = handler as TypedMessageHandler; - return this as unknown as TypedLLPClient; - } - - this.messageHandler = handler as MessageHandler; - return undefined; + this.messageHandler = handler as MessageHandler; + return this as unknown as TypedLLPClient; } /** @@ -174,10 +115,6 @@ export class LLPClient implements Annotater { return this as unknown as TypedLLPClient; } - onPresence(handler: PresenceHandler): void { - this.onStart(handler); - } - // ------------------------------------------------------------------- // Connection lifecycle // ------------------------------------------------------------------- @@ -247,35 +184,6 @@ export class LLPClient implements Annotater { // ------------------------------------------------------------------- // Messaging // ------------------------------------------------------------------- - - async sendMessage(msg: TextMessage, timeout?: number): Promise { - if (this.status !== ConnectionStatus.Authenticated) { - throw new NotAuthenticatedError('Must be authenticated to send messages'); - } - - const timeoutMs = timeout ?? this.config.responseTimeout ?? 10000; - - return new Promise((resolve, reject) => { - const timeoutId = setTimeout(() => { - this.pending.delete(msg.id); - reject(new TimeoutError(`Message ${msg.id} timed out`)); - }, timeoutMs); - - this.pending.set(msg.id, { - resolve: (response) => { - clearTimeout(timeoutId); - resolve(response); - }, - reject: (err) => { - clearTimeout(timeoutId); - reject(err); - }, - }); - - this.enqueue(msg.encode()); - }); - } - async annotateToolCall(toolCall: ToolCall): Promise { if (this.status !== ConnectionStatus.Authenticated) { throw new NotAuthenticatedError('Must be authenticated to annotate tool calls'); @@ -283,7 +191,7 @@ export class LLPClient implements Annotater { this.enqueue(toolCall.encode()); } - async sendAsyncMessage(msg: TextMessage): Promise { + private async sendAsyncMessage(msg: TextMessage): Promise { if (this.status !== ConnectionStatus.Authenticated) { throw new NotAuthenticatedError('Must be authenticated to send messages'); } @@ -375,6 +283,8 @@ export class LLPClient implements Annotater { case 'error': this.handleErrorMessage(json); break; + case 'ack': + break; default: console.warn(`Unknown message type: ${messageType}`); } @@ -396,115 +306,45 @@ export class LLPClient implements Annotater { private async handleTextMessage(json: Record): Promise { const msg = TextMessage.decode(json); - // Check if this is a response to a pending request - const pending = this.pending.get(msg.id); - if (pending) { - this.pending.delete(msg.id); - pending.resolve(msg); - return; - } - - const session = this.getOrCreateSession(msg.sender); - - // New API path: typed message handler - if (this.typedMessageHandler) { - try { - console.log(`[message] received from sender=${msg.sender} id=${msg.id}`); - await session.waitForData(); - const data = session.data; - if (!data) { - console.error( - `[message] no session data for sender=${msg.sender} — was presence received?`, - ); - return; - } - const reply = await this.typedMessageHandler(data, msg, session); - await this.sendAsyncMessage(reply); - } catch (err) { - console.error('Error in message handler:', err); - } + console.log(`[message] received from sender=${msg.sender} id=${msg.id}`); + const data = this.sessions.get(msg.sender); + if (!data) { + console.error(`[message] no session data for sender=${msg.sender} — was presence received?`); return; } - // Legacy API path - if (this.messageHandler) { - try { - const reply = await this.messageHandler(session, msg); - await this.sendAsyncMessage(reply); - } catch (err) { - console.error('Error in message handler:', err); - } + if (!this.messageHandler) { + throw new MessageHandlerNotSetError( + 'Message handler not set, have you called onMessage before connecting?', + ); } + const reply = await this.messageHandler(data, msg, this); + await this.sendAsyncMessage(reply); } private handlePresenceMessage(json: Record): void { const presence = PresenceMessage.decode(json); console.log(`[presence] sender=${presence.sender} status=${presence.status}`); - const session = this.getOrCreateSession(presence.sender); - - // New API path: session creator + stop handler - if (this.sessionCreator) { - if (presence.status === PresenceStatus.Available) { - session.clearData(); - Promise.resolve(this.sessionCreator()) - .then((data) => { - if (data !== undefined) { - session.setData(data as TSessionData); - console.log(`[presence] session ready for sender=${presence.sender}`); - } else { - console.warn(`[presence] onStart returned undefined for sender=${presence.sender}`); - session.failInit(new Error('onStart returned undefined')); - } - }) + + if (presence.status === PresenceStatus.Available) { + if (this.sessionCreator) { + const data = this.sessionCreator() as TSessionData; + this.sessions.set(presence.sender, data); + console.log(`[presence] session ready for sender=${presence.sender}`); + } + } else if (presence.status === PresenceStatus.Unavailable) { + const data = this.sessions.get(presence.sender); + if (data && this.stopHandler) { + Promise.resolve(this.stopHandler(data)) .catch((err) => { - console.error(`[presence] onStart failed for sender=${presence.sender}:`, err); - session.failInit(err instanceof Error ? err : new Error(String(err))); + console.error('Error in onStop handler:', err); + }) + .finally(() => { + this.sessions.delete(presence.sender); }); - } else if (presence.status === PresenceStatus.Unavailable) { - const data = session.data; - if (data && this.stopHandler) { - Promise.resolve(this.stopHandler(data)) - .catch((err) => { - console.error('Error in onStop handler:', err); - }) - .finally(() => { - session.clear(); - this.sessions.delete(presence.sender); - }); - } else { - session.clear(); - this.sessions.delete(presence.sender); - } + } else { + this.sessions.delete(presence.sender); } - return; - } - - // Legacy API path - if (this.startHandler) { - Promise.resolve(this.startHandler(session, presence)) - .then((data) => { - if (presence.status === PresenceStatus.Available) { - session.clearData(); - if (data !== undefined) { - session.setData(data); - } - } - }) - .catch((err) => { - console.error('Error in start handler:', err); - }) - .finally(() => { - if (presence.status === PresenceStatus.Unavailable) { - session.clear(); - this.sessions.delete(presence.sender); - } - }); - return; - } - - if (presence.status === PresenceStatus.Unavailable) { - session.clear(); - this.sessions.delete(presence.sender); } } @@ -515,16 +355,6 @@ export class LLPClient implements Annotater { const error = new PlatformError(code, message, messageId); - // If this is a response to a pending message, reject it - if (messageId) { - const pending = this.pending.get(messageId); - if (pending) { - this.pending.delete(messageId); - pending.reject(error); - return; - } - } - // If we're authenticating, reject the auth promise if (this.authReject) { this.authReject(error); @@ -543,15 +373,6 @@ export class LLPClient implements Annotater { this.status = ConnectionStatus.Disconnected; this.presence = PresenceStatus.Unavailable; this.sessionId = null; - - // Reject all pending messages - for (const pending of this.pending.values()) { - pending.reject(new NotConnectedError('Disconnected from server')); - } - this.pending.clear(); - for (const session of this.sessions.values()) { - session.clear(); - } this.sessions.clear(); } } @@ -565,15 +386,4 @@ export class LLPClient implements Annotater { this.authResolve = null; } } - - private getOrCreateSession(id: string): LLPSession { - const existing = this.sessions.get(id); - if (existing) { - return existing; - } - - const session = new LLPSession(id, this); - this.sessions.set(id, session); - return session; - } } diff --git a/src/errors.ts b/src/errors.ts index 79dde9d..61e4e56 100644 --- a/src/errors.ts +++ b/src/errors.ts @@ -53,3 +53,7 @@ export class TextMessageReplyError extends Error { export class TextMessageEmptyError extends Error { name = 'TextMessageEmptyError'; } + +export class MessageHandlerNotSetError extends Error { + name = 'MessageHandlerNotSetError'; +} diff --git a/src/index.ts b/src/index.ts index 05bd053..55d0c97 100644 --- a/src/index.ts +++ b/src/index.ts @@ -4,14 +4,10 @@ export type { Annotater } from './annotate.js'; export { LLPClient, type LLPClientConfig, - // Legacy handler types (backward compat) type MessageHandler, - type PresenceHandler, type SessionCreator, - type StartHandler, type StopHandler, type TypedLLPClient, - type TypedMessageHandler, } from './client.js'; // Error types and codes export { @@ -34,5 +30,4 @@ export { PresenceStatus, type PresenceStatus as PresenceStatusType, } from './presence.js'; -export { LLPSession, LLPSession as Session } from './session.js'; export { ToolCall } from './tool_call.js'; diff --git a/src/session.ts b/src/session.ts deleted file mode 100644 index 481e4b0..0000000 --- a/src/session.ts +++ /dev/null @@ -1,58 +0,0 @@ -import type { Annotater } from './annotate.js'; -import type { ToolCall } from './tool_call.js'; - -export class LLPSession implements Annotater { - private sessionData: TData | undefined; - private dataReady: { - promise: Promise; - resolve: () => void; - reject: (err: Error) => void; - } | null = null; - - constructor( - public readonly id: string, - private readonly annotater: Annotater, - ) {} - - async annotateToolCall(toolCall: ToolCall): Promise { - await this.annotater.annotateToolCall(toolCall); - } - - get data(): TData | undefined { - return this.sessionData; - } - - setData(value: TData): void { - this.sessionData = value; - this.dataReady?.resolve(); - this.dataReady = null; - } - - failInit(err: Error): void { - this.dataReady?.reject(err); - this.dataReady = null; - } - - clearData(): void { - this.sessionData = undefined; - } - - async waitForData(): Promise { - if (this.sessionData !== undefined) return; - if (!this.dataReady) { - let resolve!: () => void; - let reject!: (err: Error) => void; - const promise = new Promise((res, rej) => { - resolve = res; - reject = rej; - }); - this.dataReady = { promise, resolve, reject }; - } - await this.dataReady.promise; - } - - clear(): void { - this.sessionData = undefined; - this.dataReady = null; - } -} diff --git a/tests/handler.test.ts b/tests/handler.test.ts index 33f3a88..3efcd42 100644 --- a/tests/handler.test.ts +++ b/tests/handler.test.ts @@ -1,76 +1,76 @@ import { describe, expect, it } from 'vitest'; +import type { Annotater } from '../src/annotate.js'; import { PresenceMessage, TextMessage } from '../src/message.js'; import { PresenceStatus } from '../src/presence.js'; -import { LLPSession } from '../src/session.js'; describe('Message Handlers', () => { describe('MessageHandler type', () => { - const session = new LLPSession('alice', { - annotateToolCall: async () => {}, - }); - it('should accept async function that returns TextMessage', async () => { const handler = async ( - _session: LLPSession, + _data: string, msg: TextMessage, + _annotater: Annotater, ): Promise => { return msg.reply('Response'); }; const input = new TextMessage('bob', 'Hello', null, 'msg-1', 'alice'); + const mockAnnotater: Annotater = { annotateToolCall: async () => {} }; - const result = await handler(session, input); + const result = await handler('session-data', input, mockAnnotater); expect(result.prompt).toBe('Response'); expect(result.recipient).toBe('alice'); }); it('should allow handler to process message content', async () => { const handler = async ( - _session: LLPSession, + _data: string, msg: TextMessage, + _annotater: Annotater, ): Promise => { const upperPrompt = msg.prompt.toUpperCase(); return msg.reply(`Echo: ${upperPrompt}`); }; const input = new TextMessage('bob', 'hello world'); + const mockAnnotater: Annotater = { annotateToolCall: async () => {} }; - const result = await handler(session, input); + const result = await handler('session-data', input, mockAnnotater); expect(result.prompt).toBe('Echo: HELLO WORLD'); }); }); describe('PresenceHandler type', () => { it('should accept void function', () => { - const handler = (_session: LLPSession, msg: PresenceMessage): void => { + const handler = (_data: string, msg: PresenceMessage): void => { expect(msg.sender).toBeDefined(); }; const presence = new PresenceMessage(PresenceStatus.Available, 'alice'); - handler(new LLPSession('alice', { annotateToolCall: async () => {} }), presence); + handler('session-data', presence); }); it('should accept async void function', async () => { - const handler = async (_session: LLPSession, msg: PresenceMessage): Promise => { + const handler = async (_data: string, msg: PresenceMessage): Promise => { await new Promise((resolve) => setTimeout(resolve, 1)); expect(msg.status).toBe(PresenceStatus.Unavailable); }; const presence = new PresenceMessage(PresenceStatus.Unavailable, 'bob'); - await handler(new LLPSession('bob', { annotateToolCall: async () => {} }), presence); + await handler('session-data', presence); }); it('should allow side effects in handler', async () => { const log: string[] = []; - const handler = (session: LLPSession, msg: PresenceMessage): void => { - log.push(`${session.id}: ${msg.status}`); + const handler = (data: string, msg: PresenceMessage): void => { + log.push(`${data}: ${msg.status}`); }; const p1 = new PresenceMessage(PresenceStatus.Available, 'alice'); const p2 = new PresenceMessage(PresenceStatus.Unavailable, 'bob'); - handler(new LLPSession('alice', { annotateToolCall: async () => {} }), p1); - handler(new LLPSession('bob', { annotateToolCall: async () => {} }), p2); + handler('alice', p1); + handler('bob', p2); expect(log).toEqual(['alice: available', 'bob: unavailable']); }); diff --git a/tests/integration.test.ts b/tests/integration.test.ts index ab55cae..42d2680 100644 --- a/tests/integration.test.ts +++ b/tests/integration.test.ts @@ -1,8 +1,8 @@ import { beforeEach, describe, expect, it, vi } from 'vitest'; import type WebSocket from 'ws'; import { LLPClient } from '../src/client.js'; -import { ErrorCode, PlatformError, TimeoutError } from '../src/errors.js'; -import { type PresenceMessage, TextMessage } from '../src/message.js'; +import { ErrorCode, PlatformError } from '../src/errors.js'; +import { TextMessage } from '../src/message.js'; import { ConnectionStatus, PresenceStatus } from '../src/presence.js'; // Mock WebSocket @@ -86,43 +86,33 @@ describe('LLPClient Integration Tests', () => { }); describe('Message sending and receiving', () => { - it('should send message and receive response', async () => { - await connectClient(); - - const messageHandler = mockWs.on.mock.calls.find((call) => call[0] === 'message')?.[1]; - - // Send a message - const tm = new TextMessage('bob', 'Hello Bob'); - const sendPromise = client.sendMessage(tm); - - // Verify message was sent - expect(mockWs.send).toHaveBeenCalledTimes(3); // 1 auth + 1 presence + 1 message - - const sentMessage = JSON.parse(mockWs.send.mock.calls[2][0]); - expect(sentMessage.type).toBe('message'); - expect(sentMessage.data.to).toBe('bob'); - - // Simulate response - const response = tm.reply('Hi back!'); - - messageHandler?.(Buffer.from(response.encode())); - - const result = await sendPromise; - expect(result.prompt).toBe('Hi back!'); - }); - it('should handle incoming messages and send replies', async () => { const replies: TextMessage[] = []; - client.onMessage(async (session, msg: TextMessage) => { - expect(session.id).toBe('alice'); - const reply = msg.reply(`Echo: ${msg.prompt}`); - replies.push(reply); - return reply; + const messageCalled = new Promise((resolve, reject) => { + client + .onStart(() => 'alice-session') + .onMessage(async (_agent, msg: TextMessage, _annotater) => { + expect(msg.sender).toBe('alice'); + const reply = msg.reply(`Echo: ${msg.prompt}`); + replies.push(reply); + resolve(true); + return reply; + }); + setTimeout(reject, 100); }); await connectClient(); const messageHandler = mockWs.on.mock.calls.find((call) => call[0] === 'message')?.[1]; + // Simulate presence update from alice to initialize session + const presenceMsg = JSON.stringify({ + type: 'presence', + id: 'pres-alice', + from: 'alice', + data: { status: 'available' }, + }); + messageHandler?.(Buffer.from(presenceMsg)); + // Simulate incoming message const incomingMsg = new TextMessage( 'test-agent', @@ -134,29 +124,24 @@ describe('LLPClient Integration Tests', () => { messageHandler?.(Buffer.from(incomingMsg.encode())); - // Wait for async handler - await new Promise((resolve) => setTimeout(resolve, 10)); + // Wait for async handler, then flush remaining microtasks + const called = await messageCalled; + await new Promise((r) => setTimeout(r, 0)); // Verify reply was sent (1 auth + 1 presence + 1 reply from handler) + expect(called).toBe(true); expect(mockWs.send).toHaveBeenCalledTimes(3); expect(replies[0]?.prompt).toBe('Echo: Hello!'); }); - - it('should timeout if no response received', async () => { - await connectClient(); - - const sendPromise = client.sendMessage(new TextMessage('bob', 'Hello'), 100); // 100ms timeout - - // Don't send a response, let it timeout - await expect(sendPromise).rejects.toThrow(TimeoutError); - }); }); describe('Presence handling', () => { it('should receive and process presence updates', async () => { - const presenceUpdates: Array<{ sessionId: string; msg: PresenceMessage }> = []; - client.onStart((session, msg: PresenceMessage) => { - presenceUpdates.push({ sessionId: session.id, msg }); + const startCalled = new Promise((resolve, reject) => { + client.onStart(() => { + resolve(true); + }); + setTimeout(reject, 10); }); await connectClient(); @@ -175,140 +160,13 @@ describe('LLPClient Integration Tests', () => { messageHandler?.(Buffer.from(presenceMsg)); - // Wait for async handler - await new Promise((resolve) => setTimeout(resolve, 10)); - - expect(presenceUpdates).toHaveLength(1); - expect(presenceUpdates[0]?.sessionId).toBe('alice'); - expect(presenceUpdates[0]?.msg.sender).toBe('alice'); - expect(presenceUpdates[0]?.msg.status).toBe(PresenceStatus.Available); - }); - - it('should handle async presence handlers', async () => { - let handlerCalled = false; - client.onStart(async (session, msg: PresenceMessage) => { - await new Promise((resolve) => setTimeout(resolve, 5)); - handlerCalled = true; - expect(session.id).toBe('bob'); - expect(msg.sender).toBe('bob'); - }); - - await connectClient(); - - const messageHandler = mockWs.on.mock.calls.find((call) => call[0] === 'message')?.[1]; - - const presenceMsg = JSON.stringify({ - type: 'presence', - id: 'pres-2', - from: 'bob', - data: { - status: 'unavailable', - }, - }); - - messageHandler?.(Buffer.from(presenceMsg)); - - // Wait for async handler - await new Promise((resolve) => setTimeout(resolve, 20)); - - expect(handlerCalled).toBe(true); - }); - - it('should isolate session state by sender and clear it on unavailable', async () => { - const seenValues: Array = []; - - client.onStart((_session, msg: PresenceMessage) => { - if (msg.status === PresenceStatus.Available) { - return `agent-for-${msg.sender}`; - } - }); - - client.onMessage(async (session, msg: TextMessage) => { - seenValues.push(session.data); - return msg.reply(`Ack: ${msg.prompt}`); - }); - - await connectClient(); - - const messageHandler = mockWs.on.mock.calls.find((call) => call[0] === 'message')?.[1]; - - const alicePresence = JSON.stringify({ - type: 'presence', - id: 'pres-alice', - from: 'alice', - data: { status: 'available' }, - }); - const bobPresence = JSON.stringify({ - type: 'presence', - id: 'pres-bob', - from: 'bob', - data: { status: 'available' }, - }); - messageHandler?.(Buffer.from(alicePresence)); - messageHandler?.(Buffer.from(bobPresence)); - - // Wait for presence handler microtasks (setData) to resolve - await new Promise((resolve) => setTimeout(resolve, 10)); - - const aliceMsg = new TextMessage('test-agent', 'hello', undefined, 'msg-alice', 'alice'); - const bobMsg = new TextMessage('test-agent', 'hi', undefined, 'msg-bob', 'bob'); - messageHandler?.(Buffer.from(aliceMsg.encode())); - messageHandler?.(Buffer.from(bobMsg.encode())); - - await new Promise((resolve) => setTimeout(resolve, 20)); - - expect(seenValues).toEqual(['agent-for-alice', 'agent-for-bob']); - - const aliceUnavailable = JSON.stringify({ - type: 'presence', - id: 'pres-alice-off', - from: 'alice', - data: { status: 'unavailable' }, - }); - messageHandler?.(Buffer.from(aliceUnavailable)); - await new Promise((resolve) => setTimeout(resolve, 20)); - - const aliceMsgAfter = new TextMessage( - 'test-agent', - 'hello again', - undefined, - 'msg-alice-2', - 'alice', - ); - messageHandler?.(Buffer.from(aliceMsgAfter.encode())); - await new Promise((resolve) => setTimeout(resolve, 20)); - - expect(seenValues).toEqual(['agent-for-alice', 'agent-for-bob', undefined]); + const called = await startCalled; + expect(called).toBe(true); }); }); describe('Error handling', () => { - it('should handle server error responses', async () => { - await connectClient(); - - const messageHandler = mockWs.on.mock.calls.find((call) => call[0] === 'message')?.[1]; - - const msg = new TextMessage('nonexistent', 'Hello'); - - const sendPromise = client.sendMessage(msg); - - // Simulate error response - const errorMsg = JSON.stringify({ - type: 'error', - id: msg.id, - code: ErrorCode.AgentNotFound, - message: 'Agent not found', - }); - - messageHandler?.(Buffer.from(errorMsg)); - - await expect(sendPromise).rejects.toThrow(PlatformError); - await expect(sendPromise).rejects.toMatchObject({ - code: ErrorCode.AgentNotFound, - messageId: msg.id, - }); - }); - + // TODO: client should disconnect when receiving an error it('should reject auth promise on authentication error', async () => { const connectPromise = client.connect(); @@ -333,95 +191,4 @@ describe('LLPClient Integration Tests', () => { }); }); }); - - describe('Queue management', () => { - it('should queue messages when WebSocket is not ready', async () => { - // Create client with small queue - client = new LLPClient('test-agent', 'ws://localhost:4000', 'key', { - maxQueueSize: 3, - }); - - await connectClient(); - - // Mock WebSocket as not open - mockWs.readyState = 0; // CONNECTING - - // These should queue - await expect( - client.sendAsyncMessage(new TextMessage('bob', 'msg1')), - ).resolves.toBeUndefined(); - - await expect( - client.sendAsyncMessage(new TextMessage('bob', 'msg2')), - ).resolves.toBeUndefined(); - }); - - it('should throw when queue is full', async () => { - client = new LLPClient('test-agent', 'key', { - maxQueueSize: 2, - }); - - await connectClient(); - - // Mock WebSocket as not ready - mockWs.readyState = 0; // CONNECTING - - // Fill the queue - await client.sendAsyncMessage(new TextMessage('bob', 'msg1')); - await client.sendAsyncMessage(new TextMessage('bob', 'msg2')); - - // This should overflow - await expect(client.sendAsyncMessage(new TextMessage('bob', 'msg3'))).rejects.toThrow( - 'Outbound queue is full', - ); - }); - }); - - describe('Multiple concurrent requests', () => { - it('should handle multiple pending requests simultaneously', async () => { - await connectClient(); - - const messageHandler = mockWs.on.mock.calls.find((call) => call[0] === 'message')?.[1]; - - // Send three messages concurrently - const msg1 = new TextMessage('alice', 'Hello Alice'); - const msg2 = new TextMessage('bob', 'Hello Bob'); - const msg3 = new TextMessage('charlie', 'Hello Charlie'); - - const promise1 = client.sendMessage(msg1); - const promise2 = client.sendMessage(msg2); - const promise3 = client.sendMessage(msg3); - - // Respond to them in reverse order - const response3 = msg3.reply('Hi from Charlie'); - const response1 = msg1.reply('Hi from Alice'); - const response2 = msg2.reply('Hi from Bob'); - - messageHandler?.(Buffer.from(response3.encode())); - messageHandler?.(Buffer.from(response1.encode())); - messageHandler?.(Buffer.from(response2.encode())); - - const [result1, result2, result3] = await Promise.all([promise1, promise2, promise3]); - - expect(result1.prompt).toBe('Hi from Alice'); - expect(result2.prompt).toBe('Hi from Bob'); - expect(result3.prompt).toBe('Hi from Charlie'); - }); - }); - - describe('Disconnection scenarios', () => { - it('should clean up pending requests on disconnect', async () => { - await connectClient(); - - const sendPromise = client.sendMessage(new TextMessage('bob', 'Hello')); - - // Trigger disconnect - const closeHandler = mockWs.on.mock.calls.find((call) => call[0] === 'close')?.[1]; - closeHandler?.(); - - await expect(sendPromise).rejects.toThrow('Disconnected from server'); - expect(client.getStatus()).toBe(ConnectionStatus.Disconnected); - expect(client.getPresence()).toBe(PresenceStatus.Unavailable); - }); - }); });