Protobus Host Bridge (#3747)

* WIP host bridge

* Run formatter

* remove tmp impl & rename host grpc client

* gitignore more files

* better layout

* host handler to make other hosts easier to add

* remove adapter pattern

* get host responses correctly

* fix streaming mode for host bridge

* first wip subscription host bridge demo for watching mcp server config

* format, comment

* add cancellation for host grpc stream

* remove unneeded functions from host-grpc-handler

* add a method for canceling request rather than using the registry

* another todo

* use StringRequest for uri.proto

* debounce new file watcher

* remove test setup

* remove some todos and logs

* remove registry use todo

* Revert "remove registry use todo"

This reverts commit 84078d3469.

* fix capitalization of uri.proto

* a better pattern for a callback based bridge without using the requestRegistry directly

---------

Co-authored-by: Andrei Edell <andrei@nugbase.com>
Co-authored-by: Sarah Fortune <sarah.fortune@gmail.com>
Co-authored-by: Andrei Eternal <eternal@cline.bot>
This commit is contained in:
Andrei Eternal
2025-06-04 15:41:30 -07:00
committed by GitHub
co-authored by Andrei Edell Sarah Fortune Andrei Eternal
parent 47899ac52d
commit cbcd17764b
14 changed files with 1113 additions and 19 deletions
+197
View File
@@ -0,0 +1,197 @@
import { v4 as uuidv4 } from "uuid"
import { hostServiceHandlers } from "./host-grpc-service-config"
import { GrpcRequestRegistry } from "../../src/core/controller/grpc-request-registry"
/**
* Type definition for a streaming response handler
*/
export type StreamingResponseHandler = (response: any, isLast?: boolean, sequenceNumber?: number) => Promise<void>
// Registry to track active gRPC requests and their cleanup functions
const requestRegistry = new GrpcRequestRegistry()
/**
* Callback interface for streaming requests
*/
export interface StreamingCallbacks<T = any> {
onResponse: (response: T) => void
onError?: (error: Error) => void
onComplete?: () => void
}
/**
* Handles gRPC requests from the webview
*/
export class GrpcHandler {
constructor() {}
/**
* Handle a gRPC request from the webview
* @param service The service name
* @param method The method name
* @param message The request message
* @param requestId The request ID for response correlation
* @param streamingCallbacks Optional callbacks for streaming responses
* @returns For unary requests: the response message or error. For streaming requests: a cancel function.
*/
async handleRequest<T = any>(
service: string,
method: string,
message: any,
requestId: string,
streamingCallbacks?: StreamingCallbacks<T>,
): Promise<
| {
message?: any
error?: string
request_id: string
}
| (() => void)
> {
// If streaming callbacks are provided, handle as a streaming request
if (streamingCallbacks) {
let completionCalled = false
// Create a response handler that will call the client's callbacks
const responseHandler: StreamingResponseHandler = async (response, isLast = false, sequenceNumber) => {
try {
// Call the client's onResponse callback with the response
streamingCallbacks.onResponse(response)
// If this is the last response, call the onComplete callback
if (isLast && streamingCallbacks.onComplete && !completionCalled) {
completionCalled = true
streamingCallbacks.onComplete()
}
} catch (error) {
// If there's an error in the callback, call the onError callback
if (streamingCallbacks.onError) {
streamingCallbacks.onError(error instanceof Error ? error : new Error(String(error)))
}
}
}
// Register the response handler with the registry
requestRegistry.registerRequest(
requestId,
() => {
console.log(`[DEBUG] Cleaning up streaming request: ${requestId}`)
if (streamingCallbacks.onComplete && !completionCalled) {
completionCalled = true
streamingCallbacks.onComplete()
}
},
{ type: "streaming_request", service, method },
responseHandler,
)
// Call the streaming handler directly
console.log(`[DEBUG] Streaming gRPC host call to ${service}.${method} req:${requestId}`)
try {
await this.handleStreamingRequest(service, method, message, requestId)
} catch (error) {
if (streamingCallbacks.onError) {
streamingCallbacks.onError(error instanceof Error ? error : new Error(String(error)))
}
}
// Return a function to cancel the stream
return () => {
console.log(`[DEBUG] Cancelling streaming request: ${requestId}`)
this.cancelRequest(requestId)
}
}
// Handle as a unary request
try {
// Get the service handler from the config
const serviceConfig = hostServiceHandlers[service]
if (!serviceConfig) {
throw new Error(`Unknown service: ${service}`)
}
// Handle unary request
return {
message: await serviceConfig.requestHandler(method, message),
request_id: requestId,
}
} catch (error) {
return {
error: error instanceof Error ? error.message : String(error),
request_id: requestId,
}
}
}
/**
* Cancel a gRPC request
* @param requestId The request ID to cancel
* @returns True if the request was found and cancelled, false otherwise
*/
public async cancelRequest(requestId: string): Promise<boolean> {
const cancelled = requestRegistry.cancelRequest(requestId)
if (cancelled) {
// Get the registered response handler from the registry
const requestInfo = requestRegistry.getRequestInfo(requestId)
if (requestInfo && requestInfo.responseStream) {
try {
// Send cancellation confirmation using the registered response handler
await requestInfo.responseStream(
{ cancelled: true },
true, // Mark as last message
)
} catch (e) {
console.error(`Error sending cancellation response for ${requestId}:`, e)
}
}
} else {
console.log(`[DEBUG] Request not found for cancellation: ${requestId}`)
}
return cancelled
}
/**
* Handle a streaming gRPC request
* @param service The service name
* @param method The method name
* @param message The request message
* @param requestId The request ID for response correlation
*/
private async handleStreamingRequest(service: string, method: string, message: any, requestId: string): Promise<void> {
// Get the service handler from the config
const serviceConfig = hostServiceHandlers[service]
if (!serviceConfig) {
throw new Error(`Unknown service: ${service}`)
}
// Check if the service supports streaming
if (!serviceConfig.streamingHandler) {
throw new Error(`Service ${service} does not support streaming`)
}
// Get the registered response handler from the registry
const requestInfo = requestRegistry.getRequestInfo(requestId)
if (!requestInfo || !requestInfo.responseStream) {
throw new Error(`No response handler registered for request: ${requestId}`)
}
// Use the registered response handler
const responseStream = requestInfo.responseStream
// Handle streaming request and pass the requestId to all streaming handlers
await serviceConfig.streamingHandler(method, message, responseStream, requestId)
// Don't send a final message here - the stream should stay open for future updates
// The stream will be closed when the client disconnects or when the service explicitly ends it
}
}
/**
* Get the request registry instance
* This allows other parts of the code to access the registry
*/
export function getRequestRegistry(): GrpcRequestRegistry {
return requestRegistry
}
+138
View File
@@ -0,0 +1,138 @@
import { StreamingResponseHandler } from "./host-grpc-handler"
/**
* Generic type for service method handlers
*/
export type ServiceMethodHandler = (message: any) => Promise<any>
/**
* Type for streaming method handlers
*/
export type StreamingMethodHandler = (message: any, responseStream: StreamingResponseHandler, requestId?: string) => Promise<void>
/**
* Method metadata including streaming information
*/
export interface MethodMetadata {
isStreaming: boolean
}
/**
* Generic service registry for gRPC services
*/
export class ServiceRegistry {
private serviceName: string
private methodRegistry: Record<string, ServiceMethodHandler> = {}
private streamingMethodRegistry: Record<string, StreamingMethodHandler> = {}
private methodMetadata: Record<string, MethodMetadata> = {}
/**
* Create a new service registry
* @param serviceName The name of the service (used for logging)
*/
constructor(serviceName: string) {
this.serviceName = serviceName
}
/**
* Register a method handler
* @param methodName The name of the method to register
* @param handler The handler function for the method
* @param metadata Optional metadata about the method
*/
registerMethod(methodName: string, handler: ServiceMethodHandler | StreamingMethodHandler, metadata?: MethodMetadata): void {
const isStreaming = metadata?.isStreaming || false
if (isStreaming) {
this.streamingMethodRegistry[methodName] = handler as StreamingMethodHandler
} else {
this.methodRegistry[methodName] = handler as ServiceMethodHandler
}
this.methodMetadata[methodName] = { isStreaming, ...metadata }
console.log(`Registered ${this.serviceName} method: ${methodName}${isStreaming ? " (streaming)" : ""}`)
}
/**
* Check if a method is a streaming method
* @param method The method name
* @returns True if the method is a streaming method
*/
isStreamingMethod(method: string): boolean {
return this.methodMetadata[method]?.isStreaming || false
}
/**
* Get a streaming method handler
* @param method The method name
* @returns The streaming method handler or undefined if not found
*/
getStreamingHandler(method: string): StreamingMethodHandler | undefined {
return this.streamingMethodRegistry[method]
}
/**
* Handle a service request
* @param method The method name
* @param message The request message
* @returns The response message
*/
async handleRequest(method: string, message: any): Promise<any> {
const handler = this.methodRegistry[method]
if (!handler) {
if (this.isStreamingMethod(method)) {
throw new Error(`Method ${method} is a streaming method and should be handled with handleStreamingRequest`)
}
throw new Error(`Unknown ${this.serviceName} method: ${method}`)
}
return handler(message)
}
/**
* Handle a streaming service request
* @param method The method name
* @param message The request message
* @param responseStream The streaming response handler
* @param requestId The request ID for correlation and cleanup
*/
async handleStreamingRequest(
method: string,
message: any,
responseStream: StreamingResponseHandler,
requestId?: string,
): Promise<void> {
const handler = this.streamingMethodRegistry[method]
if (!handler) {
if (this.methodRegistry[method]) {
throw new Error(`Method ${method} is not a streaming method and should be handled with handleRequest`)
}
throw new Error(`Unknown ${this.serviceName} streaming method: ${method}`)
}
await handler(message, responseStream, requestId)
}
}
/**
* Create a service registry factory function
* @param serviceName The name of the service
* @returns An object with register and handle functions
*/
export function createServiceRegistry(serviceName: string) {
const registry = new ServiceRegistry(serviceName)
return {
registerMethod: (methodName: string, handler: ServiceMethodHandler | StreamingMethodHandler, metadata?: MethodMetadata) =>
registry.registerMethod(methodName, handler, metadata),
handleRequest: (method: string, message: any) => registry.handleRequest(method, message),
handleStreamingRequest: (method: string, message: any, responseStream: StreamingResponseHandler, requestId?: string) =>
registry.handleStreamingRequest(method, message, responseStream, requestId),
isStreamingMethod: (method: string) => registry.isStreamingMethod(method),
}
}
+20
View File
@@ -0,0 +1,20 @@
import * as vscode from "vscode"
import { Uri } from "../../../src/shared/proto/host/uri"
import { StringRequest } from "../../../src/shared/proto/common"
/**
* Creates a file URI from a file path
* @param request The request containing the file path
* @returns A URI object representing the file
*/
export async function file(request: StringRequest): Promise<Uri> {
const uri = vscode.Uri.file(request.value)
return Uri.create({
scheme: uri.scheme,
authority: uri.authority,
path: uri.path,
query: uri.query,
fragment: uri.fragment,
fsPath: uri.fsPath,
})
}
+28
View File
@@ -0,0 +1,28 @@
import * as vscode from "vscode"
import { JoinPathRequest, Uri } from "../../../src/shared/proto/host/uri"
/**
* Joins a URI with additional path segments
* @param request The request containing the base URI and path segments
* @returns A new URI with the path segments joined
*/
export async function joinPath(request: JoinPathRequest): Promise<Uri> {
// Convert proto Uri to vscode.Uri
if (!request.base) {
throw new Error("Base URI is required")
}
const baseUri = vscode.Uri.parse(`${request.base.scheme}://${request.base.authority}${request.base.path}`)
// Join paths
const result = vscode.Uri.joinPath(baseUri, ...request.pathSegments)
// Convert back to proto Uri
return Uri.create({
scheme: result.scheme,
authority: result.authority,
path: result.path,
query: result.query,
fragment: result.fragment,
fsPath: result.fsPath,
})
}
+20
View File
@@ -0,0 +1,20 @@
import * as vscode from "vscode"
import { Uri } from "../../../src/shared/proto/host/uri"
import { StringRequest } from "../../../src/shared/proto/common"
/**
* Parses a string URI into a Uri object
* @param request The request containing the URI string
* @returns A URI object representing the parsed URI
*/
export async function parse(request: StringRequest): Promise<Uri> {
const uri = vscode.Uri.parse(request.value)
return Uri.create({
scheme: uri.scheme,
authority: uri.authority,
path: uri.path,
query: uri.query,
fragment: uri.fragment,
fsPath: uri.fsPath,
})
}
+225
View File
@@ -0,0 +1,225 @@
import * as fs from "fs/promises"
import * as fsSync from "fs"
import { SubscribeToFileRequest, FileChangeEvent, FileChangeEvent_ChangeType } from "../../../src/shared/proto/host/watch"
import { StreamingResponseHandler, getRequestRegistry } from "../host-grpc-handler"
// Debounce configuration
const DEBOUNCE_DELAY = 100 // ms
// Keep track of active file watchers
const fileWatchers = new Map<
string,
{
watcher: fsSync.FSWatcher
subscribers: Set<StreamingResponseHandler>
lastEventTime: Map<FileChangeEvent_ChangeType, number> // Track last event time by event type
}
>()
/**
* Subscribe to file changes
* @param request The request containing the file path
* @param responseStream The streaming response handler
* @param requestId The ID of the request (passed by the gRPC handler)
*/
export async function subscribeToFile(
request: SubscribeToFileRequest,
responseStream: StreamingResponseHandler,
requestId?: string,
): Promise<void> {
const filePath = request.path
console.log(`[DEBUG] Setting up file subscription for ${filePath}`)
try {
// We don't send an initial event to avoid triggering handlers immediately
console.log(`[DEBUG] Now watching file: ${filePath}`)
// Set up or reuse file watcher
if (!fileWatchers.has(filePath)) {
// Create a new watcher for this file using Node.js fs.watch API
// This is more reliable than the VSCode FileSystemWatcher for detecting file saves
const watcher = fsSync.watch(filePath, { persistent: true }, async (eventType, filename) => {
if (eventType === "change") {
try {
const content = await fs.readFile(filePath, "utf8")
console.log(`[DEBUG] File changed: ${filePath}`)
// Get the watcher info
const watcherInfo = fileWatchers.get(filePath)
if (watcherInfo) {
// Check if this event should be debounced
const eventType = FileChangeEvent_ChangeType.CHANGED
const now = Date.now()
const lastTime = watcherInfo.lastEventTime.get(eventType) || 0
if (now - lastTime < DEBOUNCE_DELAY) {
console.log(
`[DEBUG] Debouncing change event for ${filePath} (${now - lastTime}ms since last event)`,
)
return // Skip this event due to debounce
}
// Update the last event time
watcherInfo.lastEventTime.set(eventType, now)
// Notify all subscribers
for (const subscriber of watcherInfo.subscribers) {
try {
await subscriber({
path: filePath,
type: eventType,
content,
})
} catch (error) {
console.error(`Error sending file change event: ${error}`)
watcherInfo.subscribers.delete(subscriber)
}
}
}
} catch (error) {
console.error(`Error reading changed file: ${error}`)
}
} else if (eventType === "rename") {
// In Node.js fs.watch, 'rename' can mean either creation or deletion
// We need to check if the file exists to determine which it is
try {
await fs.access(filePath)
// File exists, so it was created or renamed
const content = await fs.readFile(filePath, "utf8")
console.log(`[DEBUG] File created/renamed: ${filePath}`)
// Get the watcher info
const watcherInfo = fileWatchers.get(filePath)
if (watcherInfo) {
// Check if this event should be debounced
const eventType = FileChangeEvent_ChangeType.CREATED
const now = Date.now()
const lastTime = watcherInfo.lastEventTime.get(eventType) || 0
if (now - lastTime < DEBOUNCE_DELAY) {
console.log(
`[DEBUG] Debouncing creation event for ${filePath} (${now - lastTime}ms since last event)`,
)
return // Skip this event due to debounce
}
// Update the last event time
watcherInfo.lastEventTime.set(eventType, now)
// Notify all subscribers
for (const subscriber of watcherInfo.subscribers) {
try {
await subscriber({
path: filePath,
type: eventType,
content,
})
} catch (error) {
console.error(`Error sending file creation event: ${error}`)
watcherInfo.subscribers.delete(subscriber)
}
}
}
} catch (error) {
// File doesn't exist, so it was deleted
console.log(`[DEBUG] File deleted: ${filePath}`)
// Get the watcher info
const watcherInfo = fileWatchers.get(filePath)
if (watcherInfo) {
// Check if this event should be debounced
const eventType = FileChangeEvent_ChangeType.DELETED
const now = Date.now()
const lastTime = watcherInfo.lastEventTime.get(eventType) || 0
if (now - lastTime < DEBOUNCE_DELAY) {
console.log(
`[DEBUG] Debouncing deletion event for ${filePath} (${now - lastTime}ms since last event)`,
)
return // Skip this event due to debounce
}
// Update the last event time
watcherInfo.lastEventTime.set(eventType, now)
// Notify all subscribers
for (const subscriber of watcherInfo.subscribers) {
try {
await subscriber({
path: filePath,
type: eventType,
content: "",
})
} catch (error) {
console.error(`Error sending file deletion event: ${error}`)
watcherInfo.subscribers.delete(subscriber)
}
}
// Clean up the watcher
cleanupWatcher(filePath)
}
}
}
})
// Set up the watcher info
const watcherInfo = {
watcher,
subscribers: new Set<StreamingResponseHandler>(),
lastEventTime: new Map<FileChangeEvent_ChangeType, number>(),
}
fileWatchers.set(filePath, watcherInfo)
}
// Add this subscriber to the watcher
const watcherInfo = fileWatchers.get(filePath)!
watcherInfo.subscribers.add(responseStream)
// Register cleanup when the connection is closed
const cleanup = () => {
console.log(`[DEBUG] Cleaning up file subscription for ${filePath}`)
const watcherInfo = fileWatchers.get(filePath)
if (watcherInfo) {
watcherInfo.subscribers.delete(responseStream)
// If no subscribers left, clean up the watcher
if (watcherInfo.subscribers.size === 0) {
cleanupWatcher(filePath)
}
}
}
// Register the cleanup function with the request registry
if (requestId) {
getRequestRegistry().registerRequest(
requestId,
cleanup,
{ type: "file_subscription", path: filePath },
responseStream,
)
}
} catch (error) {
console.error(`Error setting up file subscription: ${error}`)
// Send an error response
await responseStream({
path: filePath,
type: FileChangeEvent_ChangeType.DELETED,
content: `Error: ${error instanceof Error ? error.message : String(error)}`,
})
}
}
/**
* Clean up a file watcher
* @param filePath The path of the file to clean up
*/
function cleanupWatcher(filePath: string): void {
const watcherInfo = fileWatchers.get(filePath)
if (watcherInfo) {
watcherInfo.watcher.close()
fileWatchers.delete(filePath)
console.log(`[DEBUG] Removed file watcher for ${filePath}`)
}
}