fix: Zulip attachment handling — event attachments + all file types
hermes-zulip-plugin (Tanko/Mumuni):
- Add _extract_event_attachments() to handle message.attachments from
Zulip UI file uploads (previously only inline markdown links worked)
- Add shared _cache_attachment() method for images/audio/documents
- Refactor _extract_inline_media() to use _cache_attachment()
- Fix connect() signature for hermes-agent 0.18.0 compat (add is_reconnect)
pi-zulip-extension (Abiba):
- Extend attachment handling beyond images only: text files (decoded
inline), PDFs (pdftotext extraction), binary files (metadata)
- 50K char cap on text/PDF extraction to prevent context flooding
- Classify attachments by extension (image/text/pdf/binary)
pi-mcp-extension (Abiba):
- Detect bridge-side text truncation (… ellipsis marker)
- Rebuild relay message display from structuredContent when truncated
- Add rebuildRelayTextFromStructuredContent() helper
Config:
- Add ZULIP_ROLE=router to ecosystem.abiba.config.cjs (was missing)
This commit is contained in:
@@ -0,0 +1,24 @@
|
||||
# pi Zulip extension — Abiba configuration (flat format for custom parser)
|
||||
# Deploy to: ~/.pi/agent/extensions/config.yaml
|
||||
|
||||
zulip.site: https://chat.sysloggh.net
|
||||
zulip.email: abiba-bot@chat.sysloggh.net
|
||||
zulip.api_key: cKTDMZAPW08dk3zl05sStzO7HRztzyn8
|
||||
zulip.stream: agent-hub
|
||||
zulip.all_bots_user_id: 20
|
||||
|
||||
agent.name: abiba
|
||||
agent.display_name: Abiba
|
||||
agent.zulip_bot_name: abiba-bot
|
||||
agent.owner_email: jerome@sysloggh.com
|
||||
agent.private_topic: abiba
|
||||
|
||||
health_port: 9200
|
||||
|
||||
monitoring.health_endpoint_enabled: true
|
||||
monitoring.log_level: info
|
||||
|
||||
error_handling.timeout_seconds: 30
|
||||
error_handling.retry_count: 3
|
||||
error_handling.retry_delay_seconds: 5
|
||||
error_handling.graceful_message: "Abiba encountered an error processing your request. Please try again later."
|
||||
@@ -0,0 +1,26 @@
|
||||
module.exports = {
|
||||
apps: [
|
||||
{
|
||||
name: "abiba-zulip",
|
||||
script: "/usr/bin/pi",
|
||||
args: "--mode rpc --session-id zulip-service",
|
||||
cwd: "/root",
|
||||
env: {
|
||||
ZULIP_ROLE: "router",
|
||||
ZULIP_EXTENSION_ACTIVE: "true",
|
||||
NODE_OPTIONS: "--max-old-space-size=512",
|
||||
},
|
||||
// Prevent rapid crash-looping: restart with backoff, limit retries
|
||||
min_uptime: "30s",
|
||||
max_restarts: 15,
|
||||
restart_delay: 10000,
|
||||
kill_timeout: 5000,
|
||||
autorestart: true,
|
||||
},
|
||||
{
|
||||
name: "abiba-telegram",
|
||||
script: "/usr/bin/pitg",
|
||||
cwd: "/root",
|
||||
},
|
||||
],
|
||||
};
|
||||
Vendored
+16
@@ -0,0 +1,16 @@
|
||||
/**
|
||||
* pi-zulip-extension — Zulip agent communication plugin for pi agents
|
||||
*
|
||||
* Deploy to: ~/.pi/agent/extensions/zulip/
|
||||
* Config: config.yaml (alongside extension, or env vars)
|
||||
*
|
||||
* DM-first architecture per ADR-001/ADR-002.
|
||||
* Background Zulip event queue poller that injects messages directly into
|
||||
* the current pi session via pi.sendUserMessage(), then captures the LLM
|
||||
* response from agent_end and sends it back to Zulip. No subprocess overhead.
|
||||
*
|
||||
* @see ADR-007: Platform-native plugin contracts
|
||||
* @see docs/ARCHITECTURE.md
|
||||
*/
|
||||
import type { ExtensionAPI } from "@earendil-works/pi-coding-agent";
|
||||
export default function (pi: ExtensionAPI): void;
|
||||
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,460 @@
|
||||
/**
|
||||
* MCP Bridge Extension
|
||||
*
|
||||
* Connects to MCP servers and exposes their tools as pi custom tools.
|
||||
* Supports stdio and streaming-HTTP transports.
|
||||
*
|
||||
* Configuration is read from ~/.pi/agent/extensions/mcp/mcp-servers.json.
|
||||
* See the README in this directory for the format.
|
||||
*/
|
||||
|
||||
import type { ExtensionAPI, ExtensionContext } from "@mariozechner/pi-coding-agent";
|
||||
import { StdioClientTransport } from "@modelcontextprotocol/sdk/client/stdio.js";
|
||||
import { StreamableHTTPClientTransport } from "@modelcontextprotocol/sdk/client/streamableHttp.js";
|
||||
import { Client } from "@modelcontextprotocol/sdk/client/index.js";
|
||||
import {
|
||||
CallToolResultSchema,
|
||||
ListToolsResultSchema,
|
||||
} from "@modelcontextprotocol/sdk/types.js";
|
||||
import { Type } from "typebox";
|
||||
import fs from "node:fs";
|
||||
import path from "node:path";
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Configuration
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
interface McpStdioServer {
|
||||
command: string;
|
||||
args?: string[];
|
||||
env?: Record<string, string>;
|
||||
/** Whitelist: only register these tool names. Takes priority over excludeTools. */
|
||||
tools?: string[];
|
||||
/** Blacklist: exclude these tool names. Ignored if `tools` is set. */
|
||||
excludeTools?: string[];
|
||||
}
|
||||
|
||||
interface McpHttpServer {
|
||||
name: string;
|
||||
description?: string;
|
||||
type: "streaming-http";
|
||||
url: string;
|
||||
timeout?: number;
|
||||
/** Whitelist: only register these tool names. Takes priority over excludeTools. */
|
||||
tools?: string[];
|
||||
/** Blacklist: exclude these tool names. Ignored if `tools` is set. */
|
||||
excludeTools?: string[];
|
||||
}
|
||||
|
||||
type McpServerConfig = McpStdioServer | McpHttpServer;
|
||||
|
||||
interface McpConfig {
|
||||
mcpServers: Record<string, McpServerConfig>;
|
||||
}
|
||||
|
||||
function loadConfig(): McpConfig {
|
||||
const configPath = path.join(
|
||||
path.dirname(import.meta.url.replace("file://", "")),
|
||||
"mcp-servers.json",
|
||||
);
|
||||
if (!fs.existsSync(configPath)) {
|
||||
return { mcpServers: {} };
|
||||
}
|
||||
return JSON.parse(fs.readFileSync(configPath, "utf-8"));
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// MCP client registry
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
interface ServerClient {
|
||||
client: Client;
|
||||
transport: StdioClientTransport | StreamableHTTPClientTransport;
|
||||
name: string;
|
||||
tools: Array<{ name: string; description: string; inputSchema: object }>;
|
||||
}
|
||||
|
||||
const clients = new Map<string, ServerClient>();
|
||||
|
||||
async function connectToServer(
|
||||
name: string,
|
||||
config: McpServerConfig,
|
||||
ctx: ExtensionContext,
|
||||
): Promise<ServerClient> {
|
||||
let client: Client;
|
||||
let transport: StdioClientTransport | StreamableHTTPClientTransport;
|
||||
|
||||
if ("command" in config) {
|
||||
// stdio transport
|
||||
transport = new StdioClientTransport({
|
||||
command: config.command,
|
||||
args: config.args ?? [],
|
||||
env: {
|
||||
...process.env,
|
||||
...config.env,
|
||||
},
|
||||
});
|
||||
client = new Client(
|
||||
{ name: "pi-mcp", version: "1.0.0" },
|
||||
{ capabilities: {} },
|
||||
);
|
||||
await client.connect(transport);
|
||||
} else {
|
||||
// streaming-HTTP transport
|
||||
transport = new StreamableHTTPClientTransport(
|
||||
new URL(config.url),
|
||||
undefined,
|
||||
{
|
||||
requestTimeout: config.timeout ?? 120_000,
|
||||
},
|
||||
);
|
||||
client = new Client(
|
||||
{ name: "pi-mcp", version: "1.0.0" },
|
||||
{ capabilities: {} },
|
||||
);
|
||||
await client.connect(transport);
|
||||
}
|
||||
|
||||
// Discover tools
|
||||
const toolsRes = await client.request(
|
||||
{ method: "tools/list" },
|
||||
ListToolsResultSchema,
|
||||
);
|
||||
|
||||
const tools =
|
||||
toolsRes.tools?.map((t) => ({
|
||||
name: t.name,
|
||||
description: t.description ?? "",
|
||||
inputSchema: t.inputSchema ?? {},
|
||||
})) ?? [];
|
||||
|
||||
return { client, transport, name, tools };
|
||||
}
|
||||
|
||||
async function disconnectClient(server: ServerClient) {
|
||||
try {
|
||||
await server.transport.close();
|
||||
} catch {
|
||||
// ignore close errors
|
||||
}
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Schema builder: convert MCP inputSchema to TypeBox
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
function buildParameters(inputSchema: object) {
|
||||
if (!inputSchema || typeof inputSchema !== "object") {
|
||||
return Type.Object({});
|
||||
}
|
||||
|
||||
const schema = inputSchema as {
|
||||
type?: string;
|
||||
properties?: Record<string, any>;
|
||||
required?: string[];
|
||||
};
|
||||
|
||||
const properties: Record<string, any> = {};
|
||||
|
||||
if (schema.properties) {
|
||||
for (const [key, prop] of Object.entries(schema.properties)) {
|
||||
const p = prop as { type?: string; description?: string };
|
||||
let field: any;
|
||||
switch (p.type) {
|
||||
case "string":
|
||||
field = Type.Optional(Type.String({ description: p.description }));
|
||||
break;
|
||||
case "number":
|
||||
field = Type.Optional(Type.Number());
|
||||
break;
|
||||
case "integer":
|
||||
field = Type.Optional(Type.Number());
|
||||
break;
|
||||
case "boolean":
|
||||
field = Type.Optional(Type.Boolean());
|
||||
break;
|
||||
case "array":
|
||||
field = Type.Optional(
|
||||
Type.Array(Type.String(), {
|
||||
description: p.description,
|
||||
}),
|
||||
);
|
||||
break;
|
||||
case "object":
|
||||
field = Type.Optional(Type.Record(Type.String(), Type.Unknown()));
|
||||
break;
|
||||
default:
|
||||
field = Type.Optional(Type.String());
|
||||
}
|
||||
properties[key] = field;
|
||||
}
|
||||
}
|
||||
|
||||
// Check required fields
|
||||
if (schema.required) {
|
||||
for (const req of schema.required) {
|
||||
if (properties[req]) {
|
||||
// Make required by redefining
|
||||
const def = properties[req];
|
||||
if (def && "origin" in def) {
|
||||
// It's a TypeBox optional - we'd need to rebuild as non-optional
|
||||
// For simplicity, keep all optional since MCP tools often have flexible params
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return Type.Object(properties);
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Tool call forwarder
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
const DEFAULT_MAX_RESULT_CHARS = 15_000;
|
||||
|
||||
/**
|
||||
* Rebuild relay query text from structuredContent when the bridge has
|
||||
* truncated the display text field (ends with Unicode ellipsis character).
|
||||
*/
|
||||
function rebuildRelayTextFromStructuredContent(
|
||||
textContent: string,
|
||||
structuredContent: Record<string, unknown> | undefined,
|
||||
): string | null {
|
||||
if (!structuredContent) return null;
|
||||
|
||||
// Detect bridge-side truncation: the display text ends with …
|
||||
if (!textContent.endsWith("…")) return null;
|
||||
|
||||
const nodes = structuredContent.nodes as Array<Record<string, unknown>> | undefined;
|
||||
if (!nodes || nodes.length === 0) return null;
|
||||
|
||||
const rebuilt: string[] = [];
|
||||
rebuilt.push(`📬 **Relay Query**: ${structuredContent.count ?? nodes.length} result(s).\n`);
|
||||
|
||||
for (const node of nodes) {
|
||||
const id = node.id;
|
||||
const title = node.title ?? "(no title)";
|
||||
const from = (node.metadata as Record<string, unknown>)?.from ?? "?";
|
||||
const to = (node.metadata as Record<string, unknown>)?.to ?? "?";
|
||||
const source = node.source ?? "";
|
||||
const description = node.description ?? "";
|
||||
|
||||
rebuilt.push(`**#${id}** — ${title}`);
|
||||
rebuilt.push(` From: ${from} → To: ${to}`);
|
||||
if (source) {
|
||||
rebuilt.push(` ${source.replace(/\n/g, "\n ")}`);
|
||||
}
|
||||
if (description) {
|
||||
rebuilt.push(` ${description}`);
|
||||
}
|
||||
rebuilt.push("");
|
||||
}
|
||||
|
||||
return rebuilt.join("\n");
|
||||
}
|
||||
|
||||
async function forwardToolCall(
|
||||
server: ServerClient,
|
||||
toolName: string,
|
||||
args: Record<string, unknown>,
|
||||
signal?: AbortSignal,
|
||||
maxResultChars: number = DEFAULT_MAX_RESULT_CHARS,
|
||||
): Promise<{
|
||||
content: Array<{ type: string; text: string }>;
|
||||
isError?: boolean;
|
||||
}> {
|
||||
const result = await server.client.request(
|
||||
{
|
||||
method: "tools/call",
|
||||
params: { name: toolName, arguments: args },
|
||||
},
|
||||
CallToolResultSchema,
|
||||
signal,
|
||||
);
|
||||
|
||||
// Check for bridge-side truncation and rebuild from structuredContent
|
||||
const structuredContent = (result as Record<string, unknown>).structuredContent as Record<string, unknown> | undefined;
|
||||
|
||||
const content: Array<{ type: string; text: string }> = [];
|
||||
let totalChars = 0;
|
||||
let truncated = false;
|
||||
|
||||
if (result.content) {
|
||||
for (const item of result.content) {
|
||||
if (truncated) break;
|
||||
if (item.type === "text") {
|
||||
let itemText = item.text;
|
||||
|
||||
// If bridge-side truncation detected, rebuild from structuredContent
|
||||
const rebuilt = rebuildRelayTextFromStructuredContent(itemText, structuredContent);
|
||||
if (rebuilt) {
|
||||
itemText = rebuilt;
|
||||
}
|
||||
|
||||
totalChars += itemText.length;
|
||||
if (totalChars > maxResultChars) {
|
||||
const keep = maxResultChars - (totalChars - itemText.length);
|
||||
const trimmed = itemText.slice(0, Math.max(0, keep));
|
||||
content.push({
|
||||
type: "text",
|
||||
text: trimmed + `\n\n[TRUNCATED: ${(itemText.length - keep).toLocaleString()} chars removed. Limit: ${maxResultChars.toLocaleString()} chars. Use specific queries or pagination.]`,
|
||||
});
|
||||
truncated = true;
|
||||
} else {
|
||||
content.push({ type: "text", text: itemText });
|
||||
}
|
||||
} else if (item.type === "image") {
|
||||
content.push({
|
||||
type: "text",
|
||||
text: `[MCP image: ${item.mimeType}]`,
|
||||
});
|
||||
} else if (item.type === "resource") {
|
||||
const textParts: string[] = [];
|
||||
if ("resource" in item) {
|
||||
const res = item.resource;
|
||||
if ("text" in res) {
|
||||
textParts.push(res.text);
|
||||
} else if ("blob" in res) {
|
||||
textParts.push(`[binary resource: ${res.mimeType}]`);
|
||||
}
|
||||
}
|
||||
content.push({ type: "text", text: textParts.join("\n") });
|
||||
} else {
|
||||
content.push({
|
||||
type: "text",
|
||||
text: `[MCP ${item.type}: ${JSON.stringify(item).slice(0, 500)}]`,
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return {
|
||||
content,
|
||||
isError: result.isError ?? false,
|
||||
};
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Extension
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
export default async function (pi: ExtensionAPI) {
|
||||
const config = loadConfig();
|
||||
const registeredNames = new Set<string>();
|
||||
|
||||
// Connect to all servers
|
||||
for (const [name, serverConfig] of Object.entries(config.mcpServers)) {
|
||||
try {
|
||||
const client = await connectToServer(name, serverConfig, pi);
|
||||
|
||||
if (client.tools.length === 0) {
|
||||
pi.on("session_start", (_event, ctx) => {
|
||||
ctx.ui.notify(
|
||||
`MCP server "${name}" connected but has no tools.`,
|
||||
"info",
|
||||
);
|
||||
});
|
||||
clients.set(name, client);
|
||||
continue;
|
||||
}
|
||||
|
||||
// Filter tools based on server config
|
||||
let registeredTools = client.tools;
|
||||
if (serverConfig.tools) {
|
||||
registeredTools = client.tools.filter((t) =>
|
||||
serverConfig.tools!.includes(t.name),
|
||||
);
|
||||
} else if (serverConfig.excludeTools) {
|
||||
registeredTools = client.tools.filter(
|
||||
(t) => !serverConfig.excludeTools!.includes(t.name),
|
||||
);
|
||||
}
|
||||
|
||||
pi.on("session_start", (_event, ctx) => {
|
||||
ctx.ui.notify(
|
||||
`MCP: "${name}" connected — ${registeredTools.length}/${client.tools.length} tool(s) registered.`,
|
||||
"info",
|
||||
);
|
||||
});
|
||||
|
||||
// Register each MCP tool as a pi tool
|
||||
for (const tool of registeredTools) {
|
||||
const fullName = `${name}-${tool.name}`;
|
||||
if (registeredNames.has(fullName)) continue;
|
||||
registeredNames.add(fullName);
|
||||
|
||||
const parameters = buildParameters(tool.inputSchema);
|
||||
|
||||
pi.registerTool({
|
||||
name: fullName,
|
||||
label: tool.name,
|
||||
description: `[MCP:${name}] ${tool.description}`,
|
||||
promptSnippet: `[MCP:${name}] ${tool.description}`,
|
||||
promptGuidelines: [
|
||||
`Use ${fullName} to: ${tool.description}.`,
|
||||
],
|
||||
parameters,
|
||||
async execute(
|
||||
_toolCallId,
|
||||
params: Record<string, unknown>,
|
||||
signal,
|
||||
_onUpdate,
|
||||
_ctx,
|
||||
) {
|
||||
const server = clients.get(name);
|
||||
if (!server) {
|
||||
return {
|
||||
content: [{ type: "text", text: `Server "${name}" not connected.` }],
|
||||
isError: true,
|
||||
};
|
||||
}
|
||||
|
||||
const maxChars = serverConfig.maxResultChars ?? DEFAULT_MAX_RESULT_CHARS;
|
||||
return await forwardToolCall(server, tool.name, params, signal, maxChars);
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
clients.set(name, client);
|
||||
} catch (err) {
|
||||
pi.on("session_start", (_event, ctx) => {
|
||||
ctx.ui.notify(
|
||||
`MCP: Failed to connect to "${name}": ${err instanceof Error ? err.message : String(err)}`,
|
||||
"error",
|
||||
);
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
// Register a /mcp-status command
|
||||
pi.registerCommand("mcp-status", {
|
||||
description: "Show status of all MCP servers",
|
||||
handler: async (_args, ctx) => {
|
||||
const lines: string[] = ["=== MCP Server Status ===", ""];
|
||||
for (const [name, client] of clients) {
|
||||
lines.push(` ${name}: connected (${client.tools.length} tools)`);
|
||||
for (const tool of client.tools) {
|
||||
lines.push(` - ${tool.name}: ${tool.description.slice(0, 80)}`);
|
||||
}
|
||||
}
|
||||
const unconnected = Object.keys(config.mcpServers).filter(
|
||||
(n) => !clients.has(n),
|
||||
);
|
||||
for (const name of unconnected) {
|
||||
lines.push(` ${name}: NOT CONNECTED`);
|
||||
}
|
||||
ctx.ui.notify(lines.join("\n"), "info");
|
||||
},
|
||||
});
|
||||
|
||||
// Cleanup on shutdown
|
||||
pi.on("session_shutdown", async (event, _ctx) => {
|
||||
if (event.reason === "quit") {
|
||||
for (const [name, client] of clients) {
|
||||
await disconnectClient(client);
|
||||
delete clients.get(name);
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
Reference in New Issue
Block a user