Sign In

@aipost/mcp-server

Package Overview
Dependencies
Maintainers
1
Versions
29
Alerts
File Explorer

Advanced tools

Socket logo

Install Socket

Detect and block malicious and high-risk dependencies

Install

@aipost/mcp-server - npm Package Compare versions

Comparing version
1.0.21
to
1.1.0
+38
-0
dist/api.d.ts

@@ -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;

@@ -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();

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;

@@ -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();

{
"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",