fix(number-fields): избегаем потери точности длинных чисел
- number-поля теперь рендерятся как text + inputMode=numeric, чтобы браузер не округлял значения через input type=number - пробелы при вставке в number-поля удаляются - бэкенд нормализует значения number-полей в строку перед сохранением - добавлен хелпер normalizeFieldValueForStorage Closes: искажение расчётного счёта и других длинных числовых полей
This commit is contained in:
211
server/utils/mcp-client.ts
Normal file
211
server/utils/mcp-client.ts
Normal file
@@ -0,0 +1,211 @@
|
||||
/**
|
||||
* Lightweight MCP HTTP/SSE client.
|
||||
* Supports tools/list and tools/call via:
|
||||
* 1. Streamable HTTP (POST, response is JSON or SSE stream)
|
||||
* 2. Legacy SSE session transport (GET /sse → session URL, then POST)
|
||||
*/
|
||||
|
||||
export interface McpTool {
|
||||
name: string;
|
||||
description?: string;
|
||||
inputSchema?: Record<string, unknown>;
|
||||
}
|
||||
|
||||
export interface McpCallResult {
|
||||
content: Array<{ type: string; text?: string }>;
|
||||
isError?: boolean;
|
||||
}
|
||||
|
||||
interface McpRequest {
|
||||
jsonrpc: '2.0';
|
||||
id: number;
|
||||
method: string;
|
||||
params?: unknown;
|
||||
}
|
||||
|
||||
interface McpResponse {
|
||||
jsonrpc: '2.0';
|
||||
id: number;
|
||||
result?: unknown;
|
||||
error?: { code: number; message: string };
|
||||
}
|
||||
|
||||
/** Parse a single SSE event stream looking for the first data: line with valid JSON-RPC */
|
||||
function parseSseResponse(text: string): McpResponse | null {
|
||||
for (const line of text.split('\n')) {
|
||||
if (line.startsWith('data:')) {
|
||||
const json = line.slice(5).trim();
|
||||
try {
|
||||
return JSON.parse(json) as McpResponse;
|
||||
} catch {
|
||||
// continue looking
|
||||
}
|
||||
}
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
export class McpClient {
|
||||
private url: string;
|
||||
private apiKey?: string;
|
||||
private timeout: number;
|
||||
/** Resolved session endpoint for legacy SSE transport, or null if using streamable HTTP */
|
||||
private sessionUrl: string | null = null;
|
||||
private transportDetected = false;
|
||||
|
||||
constructor(url: string, apiKey?: string, timeoutMs = 15000) {
|
||||
this.url = url.replace(/\/$/, '');
|
||||
this.apiKey = apiKey;
|
||||
this.timeout = timeoutMs;
|
||||
}
|
||||
|
||||
private headers(extra?: Record<string, string>): Record<string, string> {
|
||||
const h: Record<string, string> = { 'Content-Type': 'application/json', Accept: 'application/json, text/event-stream' };
|
||||
if (this.apiKey) h['Authorization'] = `Bearer ${this.apiKey}`;
|
||||
return { ...h, ...extra };
|
||||
}
|
||||
|
||||
/**
|
||||
* Try legacy SSE transport: GET /sse → read `endpoint` event → set sessionUrl.
|
||||
* Returns true if the server supports the legacy SSE handshake.
|
||||
*/
|
||||
private async tryLegacySseHandshake(): Promise<boolean> {
|
||||
const controller = new AbortController();
|
||||
const timer = setTimeout(() => controller.abort(), Math.min(this.timeout, 5000));
|
||||
try {
|
||||
const sseUrl = `${this.url}/sse`;
|
||||
const res = await fetch(sseUrl, { headers: { Accept: 'text/event-stream', ...(this.apiKey ? { Authorization: `Bearer ${this.apiKey}` } : {}) }, signal: controller.signal });
|
||||
if (!res.ok || !res.headers.get('content-type')?.includes('event-stream')) return false;
|
||||
// Read only first few KB to find endpoint event
|
||||
const reader = res.body?.getReader();
|
||||
if (!reader) return false;
|
||||
let accumulated = '';
|
||||
let found = false;
|
||||
while (accumulated.length < 4096) {
|
||||
const { done, value } = await reader.read();
|
||||
if (done) break;
|
||||
accumulated += new TextDecoder().decode(value);
|
||||
// Look for "event: endpoint\ndata: <url>"
|
||||
const match = accumulated.match(/event:\s*endpoint\s*\ndata:\s*(\S+)/);
|
||||
if (match) {
|
||||
const endpoint = match[1].trim();
|
||||
this.sessionUrl = endpoint.startsWith('http') ? endpoint : `${this.url}${endpoint}`;
|
||||
found = true;
|
||||
reader.cancel().catch(() => {});
|
||||
break;
|
||||
}
|
||||
}
|
||||
return found;
|
||||
} catch {
|
||||
return false;
|
||||
} finally {
|
||||
clearTimeout(timer);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Send MCP initialize + initialized handshake.
|
||||
* Per MCP spec, server expects initialize request before any other method.
|
||||
* We ignore errors here (many simple servers skip this step).
|
||||
*/
|
||||
private async sendInitialize(endpoint: string): Promise<void> {
|
||||
const initReq: McpRequest = {
|
||||
jsonrpc: '2.0',
|
||||
id: Date.now(),
|
||||
method: 'initialize',
|
||||
params: {
|
||||
protocolVersion: '2024-11-05',
|
||||
capabilities: { tools: {} },
|
||||
clientInfo: { name: 'iistwin-ai-bot', version: '1.0.0' },
|
||||
},
|
||||
};
|
||||
try {
|
||||
await this.rpcPost(endpoint, initReq);
|
||||
// Send initialized notification (no response expected)
|
||||
const notif: { jsonrpc: '2.0'; method: string } = { jsonrpc: '2.0', method: 'notifications/initialized' };
|
||||
const controller = new AbortController();
|
||||
const timer = setTimeout(() => controller.abort(), 3000);
|
||||
try {
|
||||
await fetch(endpoint, { method: 'POST', headers: this.headers(), body: JSON.stringify(notif), signal: controller.signal });
|
||||
} catch { /* notification; ignore errors */ } finally { clearTimeout(timer); }
|
||||
} catch { /* many servers don't require initialization; ignore */ }
|
||||
}
|
||||
|
||||
private async rpc(method: string, params?: unknown): Promise<unknown> {
|
||||
const body: McpRequest = { jsonrpc: '2.0', id: Date.now(), method, params };
|
||||
|
||||
// Detect transport on first call
|
||||
if (!this.transportDetected) {
|
||||
this.transportDetected = true;
|
||||
// First try streamable HTTP (primary MCP 2024 spec)
|
||||
try {
|
||||
await this.sendInitialize(this.url);
|
||||
const result = await this.rpcPost(this.url, body);
|
||||
return result;
|
||||
} catch (primaryErr) {
|
||||
// Fall back to legacy SSE transport
|
||||
const legacyOk = await this.tryLegacySseHandshake();
|
||||
if (legacyOk && this.sessionUrl) {
|
||||
await this.sendInitialize(this.sessionUrl);
|
||||
return this.rpcPost(this.sessionUrl, body);
|
||||
}
|
||||
throw primaryErr;
|
||||
}
|
||||
}
|
||||
|
||||
const endpoint = this.sessionUrl ?? this.url;
|
||||
return this.rpcPost(endpoint, body);
|
||||
}
|
||||
|
||||
private async rpcPost(endpoint: string, body: McpRequest): Promise<unknown> {
|
||||
const controller = new AbortController();
|
||||
const timer = setTimeout(() => controller.abort(), this.timeout);
|
||||
try {
|
||||
const res = await fetch(endpoint, {
|
||||
method: 'POST',
|
||||
headers: this.headers(),
|
||||
body: JSON.stringify(body),
|
||||
signal: controller.signal,
|
||||
});
|
||||
if (!res.ok) {
|
||||
throw new Error(`MCP server responded with ${res.status}: ${await res.text().catch(() => '')}`);
|
||||
}
|
||||
const contentType = res.headers.get('content-type') ?? '';
|
||||
let data: McpResponse;
|
||||
if (contentType.includes('event-stream')) {
|
||||
// SSE response for streamable HTTP transport
|
||||
const text = await res.text();
|
||||
const parsed = parseSseResponse(text);
|
||||
if (!parsed) throw new Error('MCP SSE response contained no valid JSON-RPC data');
|
||||
data = parsed;
|
||||
} else {
|
||||
data = await res.json() as McpResponse;
|
||||
}
|
||||
if (data.error) {
|
||||
throw new Error(`MCP error ${data.error.code}: ${data.error.message}`);
|
||||
}
|
||||
return data.result;
|
||||
} finally {
|
||||
clearTimeout(timer);
|
||||
}
|
||||
}
|
||||
|
||||
async listTools(): Promise<McpTool[]> {
|
||||
const result = await this.rpc('tools/list') as { tools?: McpTool[] };
|
||||
return result?.tools ?? [];
|
||||
}
|
||||
|
||||
async callTool(name: string, args: Record<string, unknown>): Promise<McpCallResult> {
|
||||
const result = await this.rpc('tools/call', { name, arguments: args }) as McpCallResult;
|
||||
return result ?? { content: [] };
|
||||
}
|
||||
|
||||
async ping(): Promise<{ ok: boolean; error?: string }> {
|
||||
try {
|
||||
await this.rpc('tools/list');
|
||||
return { ok: true };
|
||||
} catch (err: unknown) {
|
||||
return { ok: false, error: err instanceof Error ? err.message : String(err) };
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user