fix: enhance universal fs interface with encoding and withFileTypes support

This commit is contained in:
David Snelling 2025-08-09 14:42:20 -07:00
parent c474f3692a
commit 4594b20450
8 changed files with 967 additions and 22 deletions

View file

@ -0,0 +1,363 @@
/**
* BrainyMCPBroadcast
*
* Enhanced MCP service with real-time WebSocket broadcasting capabilities
* for multi-agent coordination (Jarvis Picasso communication)
*
* Features:
* - WebSocket server for real-time push notifications
* - Subscription management for multiple Claude instances
* - Message broadcasting to all connected agents
* - Works both locally and with cloud deployment
*/
import { WebSocketServer, WebSocket } from 'ws'
import { createServer, IncomingMessage } from 'http'
import { BrainyMCPService } from './brainyMCPService.js'
import { BrainyDataInterface } from '../types/brainyDataInterface.js'
import { MCPServiceOptions } from '../types/mcpTypes.js'
import { v4 as uuidv4 } from '../universal/uuid.js'
interface BroadcastMessage {
id: string
from: string
to?: string | string[]
type: 'message' | 'notification' | 'sync' | 'heartbeat' | 'identify'
event?: string
data: any
timestamp: number
}
interface ConnectedAgent {
id: string
name: string
role: string
socket: WebSocket
lastSeen: number
}
export class BrainyMCPBroadcast extends BrainyMCPService {
private wsServer?: WebSocketServer
private httpServer?: any
private agents: Map<string, ConnectedAgent> = new Map()
private messageHistory: BroadcastMessage[] = []
private maxHistorySize = 100
constructor(
brainyData: BrainyDataInterface,
options: MCPServiceOptions & {
broadcastPort?: number
cloudUrl?: string
} = {}
) {
super(brainyData, options)
}
/**
* Start the WebSocket broadcast server
* @param port Port to listen on (default: 8765)
* @param isCloud Whether this is a cloud deployment
*/
async startBroadcastServer(port = 8765, isCloud = false): Promise<void> {
return new Promise((resolve, reject) => {
try {
// Create HTTP server
this.httpServer = createServer((req, res) => {
// Health check endpoint
if (req.url === '/health') {
res.writeHead(200, { 'Content-Type': 'application/json' })
res.end(JSON.stringify({
status: 'healthy',
agents: Array.from(this.agents.values()).map(a => ({
id: a.id,
name: a.name,
role: a.role,
connected: true
})),
uptime: process.uptime()
}))
} else {
res.writeHead(404)
res.end('Not found')
}
})
// Create WebSocket server
this.wsServer = new WebSocketServer({
server: this.httpServer,
perMessageDeflate: false // Better performance
})
this.wsServer.on('connection', (socket, request) => {
this.handleNewConnection(socket, request)
})
// Start listening
this.httpServer.listen(port, () => {
console.log(`🧠 Brain Jar Broadcast Server running on ${isCloud ? 'cloud' : 'local'} port ${port}`)
console.log(`📡 WebSocket: ws://localhost:${port}`)
console.log(`🔍 Health: http://localhost:${port}/health`)
resolve()
})
// Heartbeat to keep connections alive
setInterval(() => {
this.agents.forEach((agent) => {
if (Date.now() - agent.lastSeen > 30000) {
// Remove inactive agents
this.removeAgent(agent.id)
} else {
// Send heartbeat
this.sendToAgent(agent.id, {
id: uuidv4(),
from: 'server',
type: 'heartbeat',
data: { timestamp: Date.now() },
timestamp: Date.now()
})
}
})
}, 15000)
} catch (error) {
reject(error)
}
})
}
/**
* Handle new WebSocket connection
*/
private handleNewConnection(socket: WebSocket, request: IncomingMessage) {
const agentId = uuidv4()
// Send welcome message
socket.send(JSON.stringify({
id: uuidv4(),
from: 'server',
type: 'notification',
event: 'welcome',
data: {
agentId,
message: 'Connected to Brain Jar Broadcast Server',
agents: Array.from(this.agents.values()).map(a => ({
id: a.id,
name: a.name,
role: a.role
}))
},
timestamp: Date.now()
}))
// Handle messages from this agent
socket.on('message', (data) => {
try {
const message = JSON.parse(data.toString())
this.handleAgentMessage(agentId, message)
} catch (error) {
console.error('Invalid message from agent:', error)
}
})
// Handle disconnection
socket.on('close', () => {
this.removeAgent(agentId)
})
// Handle errors
socket.on('error', (error) => {
console.error(`Agent ${agentId} error:`, error)
})
// Store temporary connection until identified
this.agents.set(agentId, {
id: agentId,
name: 'Unknown',
role: 'Unknown',
socket,
lastSeen: Date.now()
})
}
/**
* Handle message from an agent
*/
private handleAgentMessage(agentId: string, message: any) {
const agent = this.agents.get(agentId)
if (!agent) return
// Update last seen
agent.lastSeen = Date.now()
// Handle identification
if (message.type === 'identify') {
agent.name = message.name || agent.name
agent.role = message.role || agent.role
// Notify all agents about new member
this.broadcast({
id: uuidv4(),
from: 'server',
type: 'notification',
event: 'agent_joined',
data: {
agent: {
id: agent.id,
name: agent.name,
role: agent.role
}
},
timestamp: Date.now()
}, agentId) // Exclude the joining agent
// Send recent history to new agent
if (this.messageHistory.length > 0) {
this.sendToAgent(agentId, {
id: uuidv4(),
from: 'server',
type: 'sync',
data: {
history: this.messageHistory.slice(-20) // Last 20 messages
},
timestamp: Date.now()
})
}
return
}
// Create broadcast message
const broadcastMsg: BroadcastMessage = {
id: message.id || uuidv4(),
from: agent.name,
to: message.to,
type: message.type || 'message',
event: message.event,
data: message.data,
timestamp: Date.now()
}
// Store in history
this.addToHistory(broadcastMsg)
// Broadcast based on recipient
if (message.to) {
// Send to specific agent(s)
const recipients = Array.isArray(message.to) ? message.to : [message.to]
recipients.forEach((recipientName: string) => {
const recipient = Array.from(this.agents.values()).find(
a => a.name === recipientName
)
if (recipient) {
this.sendToAgent(recipient.id, broadcastMsg)
}
})
} else {
// Broadcast to all agents except sender
this.broadcast(broadcastMsg, agentId)
}
}
/**
* Broadcast message to all connected agents
*/
broadcast(message: BroadcastMessage, excludeId?: string) {
const messageStr = JSON.stringify(message)
this.agents.forEach((agent) => {
if (agent.id !== excludeId && agent.socket.readyState === WebSocket.OPEN) {
agent.socket.send(messageStr)
}
})
}
/**
* Send message to specific agent
*/
private sendToAgent(agentId: string, message: BroadcastMessage) {
const agent = this.agents.get(agentId)
if (agent && agent.socket.readyState === WebSocket.OPEN) {
agent.socket.send(JSON.stringify(message))
}
}
/**
* Remove agent from connected list
*/
private removeAgent(agentId: string) {
const agent = this.agents.get(agentId)
if (agent) {
// Notify others about disconnection
this.broadcast({
id: uuidv4(),
from: 'server',
type: 'notification',
event: 'agent_left',
data: {
agent: {
id: agent.id,
name: agent.name,
role: agent.role
}
},
timestamp: Date.now()
})
this.agents.delete(agentId)
}
}
/**
* Add message to history
*/
private addToHistory(message: BroadcastMessage) {
this.messageHistory.push(message)
// Trim history if too large
if (this.messageHistory.length > this.maxHistorySize) {
this.messageHistory = this.messageHistory.slice(-this.maxHistorySize)
}
}
/**
* Stop the broadcast server
*/
async stopBroadcastServer(): Promise<void> {
// Close all agent connections
this.agents.forEach(agent => {
agent.socket.close(1000, 'Server shutting down')
})
this.agents.clear()
// Close WebSocket server
if (this.wsServer) {
this.wsServer.close()
}
// Close HTTP server
if (this.httpServer) {
this.httpServer.close()
}
}
/**
* Get connected agents
*/
getConnectedAgents(): Array<{ id: string; name: string; role: string }> {
return Array.from(this.agents.values()).map(a => ({
id: a.id,
name: a.name,
role: a.role
}))
}
/**
* Get message history
*/
getMessageHistory(): BroadcastMessage[] {
return [...this.messageHistory]
}
}
// Export for both environments
export default BrainyMCPBroadcast

310
src/mcp/brainyMCPClient.ts Normal file
View file

@ -0,0 +1,310 @@
/**
* BrainyMCPClient
*
* Client for connecting Claude instances to the Brain Jar Broadcast Server
* Utilizes Brainy for persistent memory and vector search capabilities
*/
import WebSocket from 'ws'
import { BrainyData } from '../brainyData.js'
import { v4 as uuidv4 } from '../universal/uuid.js'
interface ClientOptions {
name: string // e.g., 'Jarvis' or 'Picasso'
role: string // e.g., 'Backend Systems' or 'Frontend Design'
serverUrl?: string // Default: ws://localhost:8765
autoReconnect?: boolean
useBrainyMemory?: boolean // Store messages in Brainy for persistence
}
interface Message {
id: string
from: string
to?: string | string[]
type: 'message' | 'notification' | 'sync' | 'heartbeat' | 'identify'
event?: string
data: any
timestamp: number
}
export class BrainyMCPClient {
private socket?: WebSocket
private options: Required<ClientOptions>
private brainy?: BrainyData
private messageHandlers: Map<string, (message: Message) => void> = new Map()
private reconnectTimeout?: NodeJS.Timeout
private isConnected = false
constructor(options: ClientOptions) {
this.options = {
serverUrl: 'ws://localhost:8765',
autoReconnect: true,
useBrainyMemory: true,
...options
}
}
/**
* Initialize Brainy for persistent memory
*/
private async initBrainy() {
if (this.options.useBrainyMemory && !this.brainy) {
this.brainy = new BrainyData({
storage: {
requestPersistentStorage: true
}
})
await this.brainy.init()
console.log(`🧠 Brainy memory initialized for ${this.options.name}`)
}
}
/**
* Connect to the broadcast server
*/
async connect(): Promise<void> {
// Initialize Brainy first
await this.initBrainy()
return new Promise((resolve, reject) => {
try {
this.socket = new WebSocket(this.options.serverUrl)
this.socket.on('open', () => {
console.log(`${this.options.name} connected to Brain Jar Broadcast`)
this.isConnected = true
// Identify ourselves
this.send({
type: 'identify',
data: {
name: this.options.name,
role: this.options.role
}
})
resolve()
})
this.socket.on('message', async (data) => {
try {
const message = JSON.parse(data.toString()) as Message
await this.handleMessage(message)
} catch (error) {
console.error('Error parsing message:', error)
}
})
this.socket.on('close', () => {
console.log(`${this.options.name} disconnected from Brain Jar`)
this.isConnected = false
if (this.options.autoReconnect) {
this.scheduleReconnect()
}
})
this.socket.on('error', (error) => {
console.error(`Connection error for ${this.options.name}:`, error)
reject(error)
})
} catch (error) {
reject(error)
}
})
}
/**
* Handle incoming message
*/
private async handleMessage(message: Message) {
// Store in Brainy for persistent memory
if (this.brainy && message.type === 'message') {
try {
await this.brainy.add({
text: `${message.from}: ${JSON.stringify(message.data)}`,
metadata: {
messageId: message.id,
from: message.from,
to: message.to,
timestamp: message.timestamp,
type: message.type,
event: message.event
}
})
} catch (error) {
console.error('Error storing message in Brainy:', error)
}
}
// Handle sync messages (receive history)
if (message.type === 'sync' && message.data.history) {
console.log(`📜 ${this.options.name} received ${message.data.history.length} historical messages`)
// Store history in Brainy
if (this.brainy) {
for (const histMsg of message.data.history) {
await this.brainy.add({
text: `${histMsg.from}: ${JSON.stringify(histMsg.data)}`,
metadata: histMsg
})
}
}
}
// Call registered handlers
const handler = this.messageHandlers.get(message.type)
if (handler) {
handler(message)
}
// Call universal handler
const universalHandler = this.messageHandlers.get('*')
if (universalHandler) {
universalHandler(message)
}
}
/**
* Send a message
*/
send(message: Partial<Message>) {
if (!this.socket || this.socket.readyState !== WebSocket.OPEN) {
console.error(`${this.options.name} is not connected`)
return
}
const fullMessage: Message = {
id: message.id || uuidv4(),
from: this.options.name,
type: message.type || 'message',
data: message.data || {},
timestamp: Date.now(),
...message
}
this.socket.send(JSON.stringify(fullMessage))
}
/**
* Send a message to specific agent(s)
*/
sendTo(recipient: string | string[], data: any) {
this.send({
to: recipient,
type: 'message',
data
})
}
/**
* Broadcast to all agents
*/
broadcast(data: any) {
this.send({
type: 'message',
data
})
}
/**
* Register a message handler
*/
on(type: string, handler: (message: Message) => void) {
this.messageHandlers.set(type, handler)
}
/**
* Remove a message handler
*/
off(type: string) {
this.messageHandlers.delete(type)
}
/**
* Search historical messages using Brainy's vector search
*/
async searchMemory(query: string, limit = 10): Promise<any[]> {
if (!this.brainy) {
console.warn('Brainy memory not initialized')
return []
}
const results = await this.brainy.search(query, limit)
return results.map(r => ({
...r.metadata,
relevance: r.score
}))
}
/**
* Get recent messages from Brainy memory
*/
async getRecentMessages(limit = 20): Promise<any[]> {
if (!this.brainy) {
console.warn('Brainy memory not initialized')
return []
}
// Search for recent activity
const results = await this.brainy.search('recent messages communication', limit)
return results
.map(r => r.metadata)
.sort((a: any, b: any) => b.timestamp - a.timestamp)
}
/**
* Schedule reconnection attempt
*/
private scheduleReconnect() {
if (this.reconnectTimeout) {
clearTimeout(this.reconnectTimeout)
}
this.reconnectTimeout = setTimeout(() => {
console.log(`🔄 ${this.options.name} attempting to reconnect...`)
this.connect().catch(error => {
console.error('Reconnection failed:', error)
this.scheduleReconnect()
})
}, 5000)
}
/**
* Disconnect from server
*/
disconnect() {
if (this.reconnectTimeout) {
clearTimeout(this.reconnectTimeout)
}
if (this.socket) {
this.socket.close(1000, 'Client disconnecting')
this.socket = undefined
}
this.isConnected = false
}
/**
* Check if connected
*/
getIsConnected(): boolean {
return this.isConnected
}
/**
* Get agent info
*/
getAgentInfo() {
return {
name: this.options.name,
role: this.options.role,
connected: this.isConnected
}
}
}
// Export for both environments
export default BrainyMCPClient