/** * March Agent SDK - Agent * Port of Python march_agent/agent.py */ import { Message } from './message.js' import { Streamer } from './streamer.js' import { HeartbeatManager, type HeartbeatFailureCallback } from './heartbeat.js' import { ConfigurationError } from './exceptions.js' import type { GatewayClient } from './gateway-client.js' import type { ConversationClient } from './conversation-client.js' import type { MemoryClient } from './memory-client.js' import type { AttachmentClient } from './attachment-client.js' import type { AgentRegistrationData, MessageHandler, SenderFilterOptions, StreamerOptions, KafkaMessage, } from './types.js' /** * Filter for matching message senders. */ export class SenderFilter { private readonly _include: Set = new Set() private readonly _exclude: Set = new Set() readonly matchAll: boolean constructor(senders?: string[]) { if (!senders || senders.length === 0) { this.matchAll = true } else { this.matchAll = false for (const sender of senders) { if (sender.startsWith('~')) { this._exclude.add(sender.slice(1)) } else { this._include.add(sender) } } } } /** * Get included senders as an array (for compatibility with Python tests). */ get include(): string[] { return Array.from(this._include) } /** * Get excluded senders as an array (for compatibility with Python tests). */ get exclude(): string[] { return Array.from(this._exclude) } /** * Check if sender matches this filter. */ matches(sender: string): boolean { // If excluded, reject if (this._exclude.has(sender)) { return false } // If match all or explicitly included if (this.matchAll || this._include.size === 0) { return true } return this._include.has(sender) } } // Message handlers stored as tuples: [SenderFilter, MessageHandler] type RegisteredHandler = [SenderFilter, MessageHandler] /** * Core agent class that handles messaging via the Agent Gateway. */ export class Agent { readonly name: string readonly agentData: AgentRegistrationData sendErrorResponses: boolean = true errorMessageTemplate: string private readonly gatewayClient: GatewayClient private readonly conversationClient?: ConversationClient private readonly memoryClient?: MemoryClient private readonly attachmentClient?: AttachmentClient private readonly heartbeatInterval: number private readonly heartbeatFailureThreshold: number private readonly onHeartbeatFailure?: HeartbeatFailureCallback private messageHandlers: RegisteredHandler[] = [] private heartbeatManager?: HeartbeatManager private initialized: boolean = false private running: boolean = false constructor(options: { name: string gatewayClient: GatewayClient agentData: AgentRegistrationData heartbeatInterval?: number heartbeatFailureThreshold?: number onHeartbeatFailure?: HeartbeatFailureCallback conversationClient?: ConversationClient memoryClient?: MemoryClient attachmentClient?: AttachmentClient errorMessageTemplate?: string }) { this.name = options.name this.gatewayClient = options.gatewayClient this.agentData = options.agentData this.heartbeatInterval = options.heartbeatInterval ?? 60 this.heartbeatFailureThreshold = options.heartbeatFailureThreshold ?? 3 this.onHeartbeatFailure = options.onHeartbeatFailure this.conversationClient = options.conversationClient this.memoryClient = options.memoryClient this.attachmentClient = options.attachmentClient this.errorMessageTemplate = options.errorMessageTemplate ?? 'I encountered an error while processing your message. Please try again or contact support if the issue persists.' } /** * Register a message handler. * * Usage: * agent.onMessage(async (message, sender) => { ... }) * agent.onMessage(handler, { senders: ['user'] }) */ onMessage(handler: MessageHandler): void onMessage(handler: MessageHandler, options: SenderFilterOptions): void onMessage( handler: MessageHandler, options?: SenderFilterOptions ): void { const filter = new SenderFilter(options?.senders) this.messageHandlers.push([filter, handler]) } /** * Initialize agent after gateway connection is established. * Can be called multiple times (e.g., after reconnection) to re-register handlers. */ initializeWithGateway(): void { // Register message handler with gateway // This is idempotent - calling it again just updates the handler const topic = `${this.name}.inbox` this.gatewayClient.registerHandler(topic, (msg) => { this.handleKafkaMessage(msg) }) // Start heartbeat if not already created if (!this.heartbeatManager) { this.heartbeatManager = new HeartbeatManager( this.gatewayClient, this.name, this.heartbeatInterval, { failureThreshold: this.heartbeatFailureThreshold, onFailureThresholdExceeded: this.handleHeartbeatFailure.bind(this), } ) } // Start or resume heartbeat (in case it was stopped during disconnect) if (!this.heartbeatManager.isRunning()) { this.heartbeatManager.start() } else if (this.heartbeatManager.isPaused()) { this.heartbeatManager.resume() } this.initialized = true } /** * Handle heartbeat failure threshold exceeded. * This can indicate a stale connection that needs reconnection. */ private handleHeartbeatFailure(consecutiveFailures: number, lastError?: Error): void { console.warn( `Agent '${this.name}' heartbeat failure threshold exceeded (${consecutiveFailures} failures)` ) // Call custom callback if provided if (this.onHeartbeatFailure) { this.onHeartbeatFailure(consecutiveFailures, lastError) } // If gateway is not connected, try to trigger reconnection if (!this.gatewayClient.isConnected() && !this.gatewayClient.isReconnectingNow()) { console.warn('Gateway not connected, attempting force reconnect...') this.gatewayClient.forceReconnect().catch((err) => { console.error('Force reconnect from heartbeat failure failed:', err) }) } } /** * Pause heartbeat sending (e.g., during reconnection). */ pauseHeartbeat(): void { this.heartbeatManager?.pause() } /** * Resume heartbeat sending. */ resumeHeartbeat(): void { this.heartbeatManager?.resume() } /** * Get sender from message headers. */ private getSender(headers: Record): string { return headers.from_ ?? headers.from ?? 'user' } /** * Find first handler that matches the sender. */ private findMatchingHandler(sender: string): MessageHandler | undefined { for (const [filter, handler] of this.messageHandlers) { if (filter.matches(sender)) { return handler } } return undefined } /** * Handle incoming Kafka message. */ private handleKafkaMessage(kafkaMsg: KafkaMessage): void { // Run handler asynchronously this.handleMessageAsync(kafkaMsg).catch((error) => { console.error('Error in message handler:', error) }) } /** * Async message handling with error recovery. */ private async handleMessageAsync(kafkaMsg: KafkaMessage): Promise { let message: Message | undefined try { // Create message from Kafka data message = Message.fromKafkaMessage( kafkaMsg.body, kafkaMsg.headers, { conversationClient: this.conversationClient, memoryClient: this.memoryClient, attachmentClient: this.attachmentClient, agentName: this.name, } ) // Find matching handler const sender = this.getSender(kafkaMsg.headers) const handler = this.findMatchingHandler(sender) if (!handler) { console.warn(`No handler matched for sender: ${sender}`) return } // Call handler await handler(message, sender) } catch (error) { console.error('Error handling message:', error) // Send error response if we have a message if (message && this.sendErrorResponses) { await this.sendErrorResponse(message, error as Error) } } } /** * Send error response to user when handler fails. */ private async sendErrorResponse(message: Message, _error: Error): Promise { try { const streamer = this.streamer(message, { sendTo: 'user' }) streamer.stream(this.errorMessageTemplate) await streamer.finish() } catch (err) { console.error('Failed to send error response:', err) } } /** * Create a new Streamer for streaming responses. */ streamer(message: Message, options: StreamerOptions = {}): Streamer { return new Streamer({ agentName: this.name, originalMessage: message, gatewayClient: this.gatewayClient, conversationClient: this.conversationClient, awaiting: options.awaiting ?? false, sendTo: options.sendTo ?? 'user', }) } /** * Mark agent as ready to consume messages. */ startConsuming(): void { if (this.messageHandlers.length === 0) { throw new ConfigurationError('No message handlers registered') } if (!this.initialized) { throw new ConfigurationError('Agent not initialized with gateway') } this.running = true console.log(`Agent ${this.name} is now consuming messages`) } /** * Shutdown agent gracefully. */ shutdown(): void { this.running = false if (this.heartbeatManager) { this.heartbeatManager.stop() } console.log(`Agent ${this.name} shutdown`) } /** * Check if agent is running. */ isRunning(): boolean { return this.running } }