@aipost/mcp-server
Advanced tools
+38
-0
@@ -54,2 +54,40 @@ import { Ed25519Signer } from "./auth/signer.js"; | ||
| } | ||
| export interface MailEvent { | ||
| /** Unique event ID (assigned locally) */ | ||
| eventId: string; | ||
| /** SSE event type (e.g., "message.new", "message.updated", "ping") */ | ||
| eventType: string; | ||
| /** Parsed event data */ | ||
| data: unknown; | ||
| /** Unix timestamp when the event was received */ | ||
| receivedAt: number; | ||
| } | ||
| export declare class EventStream { | ||
| private events; | ||
| private running; | ||
| private reconnectDelay; | ||
| private maxReconnectDelay; | ||
| private maxBufferSize; | ||
| private apiKey; | ||
| private baseUrl; | ||
| private abortController; | ||
| private eventCounter; | ||
| constructor(config: { | ||
| apiKey: string; | ||
| baseUrl?: string; | ||
| }); | ||
| /** Start the background SSE connection. Idempotent — safe to call multiple times. */ | ||
| start(): void; | ||
| /** Stop the background SSE connection. */ | ||
| stop(): void; | ||
| /** Get buffered events and optionally clear the buffer. */ | ||
| getEvents(options?: { | ||
| clear?: boolean; | ||
| }): MailEvent[]; | ||
| /** Return the number of buffered events. */ | ||
| get bufferSize(): number; | ||
| private connect; | ||
| private backoff; | ||
| private handleEvent; | ||
| } | ||
| export declare class ApiError extends Error { | ||
@@ -56,0 +94,0 @@ status: number; |
+144
-0
@@ -142,2 +142,146 @@ const BASE_URL = process.env.AIPOST_BASE_URL || "https://aipost.email"; | ||
| } | ||
| export class EventStream { | ||
| events = []; | ||
| running = false; | ||
| reconnectDelay = 1000; | ||
| maxReconnectDelay = 30_000; | ||
| maxBufferSize = 1000; | ||
| apiKey; | ||
| baseUrl; | ||
| abortController = null; | ||
| eventCounter = 0; | ||
| constructor(config) { | ||
| this.apiKey = config.apiKey; | ||
| this.baseUrl = config.baseUrl || "https://aipost.email"; | ||
| } | ||
| /** Start the background SSE connection. Idempotent — safe to call multiple times. */ | ||
| start() { | ||
| if (this.running) | ||
| return; | ||
| this.running = true; | ||
| this.abortController = new AbortController(); | ||
| console.error("[aipost-mcp] SSE event stream starting..."); | ||
| this.connect(); | ||
| } | ||
| /** Stop the background SSE connection. */ | ||
| stop() { | ||
| this.running = false; | ||
| if (this.abortController) { | ||
| this.abortController.abort(); | ||
| this.abortController = null; | ||
| } | ||
| console.error("[aipost-mcp] SSE event stream stopped"); | ||
| } | ||
| /** Get buffered events and optionally clear the buffer. */ | ||
| getEvents(options) { | ||
| const snapshot = [...this.events]; | ||
| if (options?.clear) { | ||
| this.events = []; | ||
| } | ||
| return snapshot; | ||
| } | ||
| /** Return the number of buffered events. */ | ||
| get bufferSize() { | ||
| return this.events.length; | ||
| } | ||
| async connect() { | ||
| while (this.running) { | ||
| try { | ||
| console.error("[aipost-mcp] SSE connecting..."); | ||
| const response = await fetch(`${this.baseUrl}/v1/mail/events`, { | ||
| headers: { | ||
| Authorization: `Bearer ${this.apiKey}`, | ||
| Accept: "text/event-stream", | ||
| }, | ||
| signal: this.abortController?.signal, | ||
| }); | ||
| if (!response.ok) { | ||
| const text = await response.text().catch(() => ""); | ||
| console.error(`[aipost-mcp] SSE connection failed: ${response.status} ${text}`); | ||
| await this.backoff(); | ||
| continue; | ||
| } | ||
| if (!response.body) { | ||
| console.error("[aipost-mcp] SSE response has no body"); | ||
| await this.backoff(); | ||
| continue; | ||
| } | ||
| console.error("[aipost-mcp] SSE connected, reading stream..."); | ||
| this.reconnectDelay = 1000; // reset on successful connection | ||
| const reader = response.body.getReader(); | ||
| const decoder = new TextDecoder(); | ||
| let buffer = ""; | ||
| let currentEvent = ""; | ||
| while (this.running) { | ||
| const { done, value } = await reader.read(); | ||
| if (done) { | ||
| console.error("[aipost-mcp] SSE stream ended, reconnecting..."); | ||
| break; | ||
| } | ||
| buffer += decoder.decode(value, { stream: true }); | ||
| const lines = buffer.split("\n"); | ||
| // Keep the last partial line in buffer | ||
| buffer = lines.pop() || ""; | ||
| for (const line of lines) { | ||
| if (line.startsWith("event:")) { | ||
| currentEvent = line.slice(6).trim(); | ||
| } | ||
| else if (line.startsWith("data:")) { | ||
| const raw = line.slice(5).trim(); | ||
| this.handleEvent(currentEvent || "message", raw); | ||
| currentEvent = ""; | ||
| } | ||
| else if (line.trim() === "" || line.startsWith(":")) { | ||
| // comment or empty line (heartbeat : ping), reset event type | ||
| if (line.startsWith(": ping")) { | ||
| // server heartbeat — connection is alive | ||
| } | ||
| currentEvent = ""; | ||
| } | ||
| } | ||
| } | ||
| } | ||
| catch (err) { | ||
| if (err instanceof DOMException && err.name === "AbortError") { | ||
| console.error("[aipost-mcp] SSE aborted"); | ||
| return; | ||
| } | ||
| const msg = err instanceof Error ? err.message : String(err); | ||
| console.error(`[aipost-mcp] SSE error: ${msg}`); | ||
| } | ||
| await this.backoff(); | ||
| } | ||
| } | ||
| async backoff() { | ||
| if (!this.running) | ||
| return; | ||
| console.error(`[aipost-mcp] SSE reconnecting in ${this.reconnectDelay}ms...`); | ||
| await sleep(this.reconnectDelay); | ||
| this.reconnectDelay = Math.min(this.reconnectDelay * 2, this.maxReconnectDelay); | ||
| } | ||
| handleEvent(eventType, rawData) { | ||
| let data = rawData; | ||
| try { | ||
| data = JSON.parse(rawData); | ||
| } | ||
| catch { | ||
| // keep as raw string if not JSON | ||
| } | ||
| const event = { | ||
| eventId: `evt_${++this.eventCounter}_${Date.now()}`, | ||
| eventType, | ||
| data, | ||
| receivedAt: Date.now(), | ||
| }; | ||
| console.error(`[aipost-mcp] SSE event: ${eventType}`); | ||
| // Enforce buffer size limit (FIFO) | ||
| while (this.events.length >= this.maxBufferSize) { | ||
| this.events.shift(); | ||
| } | ||
| this.events.push(event); | ||
| } | ||
| } | ||
| function sleep(ms) { | ||
| return new Promise((resolve) => setTimeout(resolve, ms)); | ||
| } | ||
| export class ApiError extends Error { | ||
@@ -144,0 +288,0 @@ status; |
+9
-2
@@ -6,3 +6,3 @@ #!/usr/bin/env node | ||
| import { isInitializeRequest } from "@modelcontextprotocol/sdk/types.js"; | ||
| import { AipostClient } from "./api.js"; | ||
| import { AipostClient, EventStream } from "./api.js"; | ||
| import { createAipostServer } from "./server.js"; | ||
@@ -16,2 +16,9 @@ const PORT = parseInt(process.env.PORT || "3000", 10); | ||
| const BOOT_ID = randomUUID().slice(0, 8); | ||
| // Start shared background SSE event stream for real-time inbox monitoring. | ||
| // Uses the server-level API key. Per-session keys cannot override this | ||
| // because the SSE stream is long-lived and shared across all sessions. | ||
| const sharedEventStream = SERVER_API_KEY | ||
| ? new EventStream({ apiKey: SERVER_API_KEY }) | ||
| : undefined; | ||
| sharedEventStream?.start(); | ||
| console.error(`[aipost-mcp] ========================================`); | ||
@@ -89,3 +96,3 @@ console.error(`[aipost-mcp] HTTP server starting (boot ${BOOT_ID})`); | ||
| const client = new AipostClient({ apiKey: effectiveKey }); | ||
| const server = createAipostServer(client); | ||
| const server = createAipostServer(client, sharedEventStream); | ||
| const transport = new StreamableHTTPServerTransport({ | ||
@@ -92,0 +99,0 @@ sessionIdGenerator: () => randomUUID(), |
+7
-2
| #!/usr/bin/env node | ||
| import { StdioServerTransport } from "@modelcontextprotocol/sdk/server/stdio.js"; | ||
| import { AipostClient } from "./api.js"; | ||
| import { AipostClient, EventStream } from "./api.js"; | ||
| import { Ed25519Signer } from "./auth/signer.js"; | ||
@@ -18,3 +18,8 @@ import { createAipostServer } from "./server.js"; | ||
| const client = new AipostClient({ apiKey: API_KEY, signer }); | ||
| const server = createAipostServer(client); | ||
| // Start background SSE event stream for real-time inbox monitoring | ||
| const eventStream = API_KEY | ||
| ? new EventStream({ apiKey: API_KEY }) | ||
| : undefined; | ||
| eventStream?.start(); | ||
| const server = createAipostServer(client, eventStream); | ||
| async function main() { | ||
@@ -21,0 +26,0 @@ const transport = new StdioServerTransport(); |
+2
-2
| import { Server } from "@modelcontextprotocol/sdk/server/index.js"; | ||
| import { AipostClient } from "./api.js"; | ||
| import { AipostClient, EventStream } from "./api.js"; | ||
| /** | ||
@@ -8,2 +8,2 @@ * Create a configured AIPost MCP Server with all tools registered. | ||
| */ | ||
| export declare function createAipostServer(client: AipostClient): Server; | ||
| export declare function createAipostServer(client: AipostClient, eventStream?: EventStream): Server; |
+24
-1
@@ -138,2 +138,12 @@ import { createRequire } from "module"; | ||
| { | ||
| name: "check_inbox_events", | ||
| description: "Poll real-time inbox events captured via background SSE connection (new mail, status changes). Returns buffered events and optionally clears them.", | ||
| inputSchema: { | ||
| type: "object", | ||
| properties: { | ||
| clear: { type: "boolean", description: "If true, clears the event buffer after returning (default: false)" }, | ||
| }, | ||
| }, | ||
| }, | ||
| { | ||
| name: "get_plans", | ||
@@ -149,3 +159,3 @@ description: "List available subscription plans and pricing.", | ||
| */ | ||
| export function createAipostServer(client) { | ||
| export function createAipostServer(client, eventStream) { | ||
| const server = new Server({ name: "aipost-mcp", version: pkg.version }, { capabilities: { tools: {} } }); | ||
@@ -235,2 +245,15 @@ server.setRequestHandler(ListToolsRequestSchema, async () => ({ tools: TOOLS })); | ||
| break; | ||
| case "check_inbox_events": { | ||
| if (!eventStream) { | ||
| throw new Error("SSE event stream is not running. Set AIPOST_API_KEY to enable background event monitoring."); | ||
| } | ||
| const clear = args.clear; | ||
| const events = eventStream.getEvents({ clear: !!clear }); | ||
| result = { | ||
| count: events.length, | ||
| bufferSize: eventStream.bufferSize, | ||
| events, | ||
| }; | ||
| break; | ||
| } | ||
| case "get_plans": | ||
@@ -237,0 +260,0 @@ result = await client.getPlans(); |
+1
-1
| { | ||
| "name": "@aipost/mcp-server", | ||
| "version": "1.0.21", | ||
| "version": "1.1.0", | ||
| "mcpName": "io.github.AIPOST-EMAIL/mcp-server", | ||
@@ -5,0 +5,0 @@ "description": "MCP Server for AIPost.email — typed, structured messaging for AI agents with Ed25519 identities", |
URL strings
Supply chain riskPackage contains fragments of external URLs or IP addresses, which the package may be accessing at runtime.
URL strings
Supply chain riskPackage contains fragments of external URLs or IP addresses, which the package may be accessing at runtime.
66414
14.34%1361
18.97%2
100%