Skip to content

Commit 5b7c7ed

Browse files
feat(mcp-server): switch channel from SSE to WebSocket for DO hibernation (#165)
* feat(mcp-server): switch channel connection from SSE to WebSocket SSE接続ではDOがhibernationできず、duration(GB-sec)が無料枠を超過していた。 local-mcpのWebSocket実装パターンに倣い、mcp-serverのチャンネル接続をネイティブWebSocketに切り替え。 25秒keepalive(ping送信)、exponential backoff自動再接続、token refresh on close code 1008/4401を含む。 eventsource依存を除去し、Node.js 18+ネイティブWebSocket APIのみ使用。 Refs #163 Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com> * docs: update installation.md SSE references to WebSocket Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com> --------- Co-authored-by: liplus-lin-lay <liplus-lin-lay@users.noreply.github.com> Co-authored-by: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
1 parent d9ca0d1 commit 5b7c7ed

5 files changed

Lines changed: 77 additions & 63 deletions

File tree

README.md

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -8,16 +8,16 @@ Real-time GitHub webhook notifications for Claude via Cloudflare Worker + Durabl
88
GitHub ──POST──▶ Cloudflare Worker ──▶ Durable Object (SQLite)
99
1010
├── MCP tools (Streamable HTTP)
11-
├── SSE real-time stream
11+
├── WebSocket real-time stream
1212
1313
┌────────────────┘
1414
1515
Desktop / Codex: .mcpb local bridge ──▶ polling via MCP tools
16-
Claude Code CLI: .mcpb local bridge ──▶ SSE → channel notifications
16+
Claude Code CLI: .mcpb local bridge ──▶ WebSocket → channel notifications
1717
```
1818

1919
- **Cloudflare Worker** receives GitHub webhooks, verifies signatures, stores events in a Durable Object with SQLite.
20-
- **Local MCP bridge** (.mcpb) proxies tool calls to the Worker and optionally listens to SSE for real-time channel notifications.
20+
- **Local MCP bridge** (.mcpb) proxies tool calls to the Worker and optionally connects via WebSocket for real-time channel notifications.
2121
- No local webhook receiver or tunnel required.
2222

2323
## Prerequisites

docs/0-requirements.md

Lines changed: 12 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -10,7 +10,7 @@ GitHub Webhook イベントを Cloudflare Worker で受信・永続化し、MCP
1010
- GitHub Webhook は Cloudflare Worker に直接到達する(ローカルサーバー不要)
1111
- イベントは Durable Object 内の SQLite に永続化される
1212
- AI エージェントは MCP stdio トランスポート(ローカルブリッジ)または Streamable HTTP(リモート)で接続する
13-
- ローカルブリッジは Worker にツール呼び出しをプロキシし、SSE でリアルタイム通知を中継する
13+
- ローカルブリッジは Worker にツール呼び出しをプロキシし、WebSocket でリアルタイム通知を中継する
1414

1515
## Architecture
1616

@@ -21,12 +21,12 @@ GitHub ──POST──▶ Cloudflare Worker ──▶ TenantRegistry DO
2121
│ ▼
2222
│ WebhookStore DO (SQLite) [per-tenant]
2323
│ │
24-
├── /mcp (Streamable HTTP) ├── SSE real-time stream
24+
├── /mcp (Streamable HTTP) ├── WebSocket real-time stream
2525
│ WebhookMcpAgent DO └── REST endpoints
2626
│ [per-tenant] /pending-status
2727
│ └── tools → WebhookStore /pending-events
2828
│ /webhook-events
29-
├── /events (SSE) /event
29+
├── /events (WebSocket/SSE) /event
3030
│ └── WebhookStore DO /mark-processed
3131
3232
└── /webhooks/github (POST)
@@ -36,15 +36,15 @@ GitHub ──POST──▶ Cloudflare Worker ──▶ TenantRegistry DO
3636
│ Local MCP Bridge (.mcpb) │
3737
│ stdio ← Claude Desktop/CLI │
3838
│ → proxy tool calls to /mcp │
39-
│ → SSE listener → channel
39+
│ → WebSocket listener → channel │
4040
└─────────────────────────────┘
4141
```
4242

4343
システムは四つのコンポーネントで構成される:
4444

4545
1. **Cloudflare Worker** — webhook 受信、署名検証、テナントルーティング
4646
2. **TenantRegistry Durable Object** — installation_id → account_id マッピング管理、テナント単位クォータ管理(単一インスタンス)
47-
3. **WebhookStore Durable Object** — SQLite によるイベント永続化、REST/SSE エンドポイント(テナント別インスタンス: `store-{accountId}`
47+
3. **WebhookStore Durable Object** — SQLite によるイベント永続化、REST/WebSocket/SSE エンドポイント(テナント別インスタンス: `store-{accountId}`
4848
4. **WebhookMcpAgent Durable Object** — MCP Streamable HTTP サーバー、ツール定義(テナント別インスタンス: `tenant-{accountId}`
4949

5050
ローカルブリッジ(mcp-server/)は Worker に対するプロキシであり、データを保持しない。
@@ -115,21 +115,21 @@ WebhookMcpAgent DO が以下のツールセットを提供する。ローカル
115115
}
116116
```
117117

118-
### F4. SSE リアルタイムイベント配信
118+
### F4. リアルタイムイベント配信(WebSocket / SSE)
119119

120120
| ID | 要件 |
121121
|----|------|
122-
| F4.1 | `GET /events` で SSE ストリームを提供する |
123-
| F4.2 | webhook ingest 時に接続中の全 SSE クライアントにイベントサマリーをブロードキャストする |
124-
| F4.3 | 30 秒間隔でハートビートを送信する |
122+
| F4.1 | `GET /events`WebSocket および SSE ストリームを提供する(Upgrade ヘッダで切り替え) |
123+
| F4.2 | webhook ingest 時に接続中の全クライアントにイベントサマリーをブロードキャストする |
124+
| F4.3 | 30 秒間隔でハートビート(WebSocket: ping、SSE: heartbeat コメント)を送信する |
125125
| F4.4 | クライアント切断時にクリーンアップする |
126126

127127
### F5. チャンネル通知(ローカルブリッジ)
128128

129129
| ID | 要件 |
130130
|----|------|
131131
| F5.1 | ローカルブリッジが Claude Code の `claude/channel` experimental capability を宣言する |
132-
| F5.2 | Worker の SSE エンドポイントに接続し、新規イベント検出時に `notifications/claude/channel` を送信する |
132+
| F5.2 | Worker の WebSocket エンドポイントに接続し、新規イベント検出時に `notifications/claude/channel` を送信する |
133133
| F5.3 | 通知内容はイベントサマリー(type, repo, action, title, sender)を含む |
134134
| F5.4 | `meta` フィールドに `chat_id`, `message_id`, `user`, `ts` を付与する |
135135
| F5.5 | `WEBHOOK_CHANNEL=0` 環境変数でチャンネル通知を無効化できる(デフォルト: 有効) |
@@ -191,7 +191,7 @@ WebhookMcpAgent DO が以下のツールセットを提供する。ローカル
191191
|----|------|
192192
| N3.1 | WebhookStore / McpAgent DO はテナント別インスタンス(`idFromName("store-{accountId}")` / `getAgentByName("tenant-{accountId}")`)で動作する。TenantRegistry DO は単一インスタンスで全テナントの installation-account マッピングを管理する |
193193
| N3.4 | OAuth コールバック時に `GET /user/installations` で取得した accessible_account_ids(ユーザー + org)を GitHubUserProps に保存し、McpAgent が複数 store を並列クエリして結果をマージする。これにより org インストールのイベントもメンバーの MCP セッションから参照できる |
194-
| N3.2 | SSE 接続は DO のメモリ内で管理される(DO eviction 時に切断) |
194+
| N3.2 | WebSocket / SSE 接続は DO のメモリ内で管理される(DO eviction 時に切断) |
195195
| N3.3 | ローカルブリッジはツール呼び出しごとに Worker セッションを再利用する(セッション失効時は自動リトライ) |
196196

197197
## Dependencies
@@ -209,7 +209,6 @@ WebhookMcpAgent DO が以下のツールセットを提供する。ローカル
209209
| パッケージ | 用途 |
210210
|-----------|------|
211211
| @modelcontextprotocol/sdk | MCP SDK(`Server` クラス直接使用) |
212-
| eventsource | SSE クライアント |
213212

214213
Node.js >= 18.0.0 が必要。
215214

@@ -245,7 +244,7 @@ manifest.json のバージョンも一致させる。
245244
| コンポーネント | 用途 |
246245
|---------------|------|
247246
| Cloudflare Worker | webhook 受信 + MCP サーバー |
248-
| Cloudflare Durable Objects | イベント永続化 (SQLite) + SSE |
247+
| Cloudflare Durable Objects | イベント永続化 (SQLite) + WebSocket/SSE |
249248
| GitHub Webhook | イベント送信元 |
250249
| MCPB | Claude Desktop 向けローカルブリッジ配布 |
251250
| npx | CLI/Codex 向けローカルブリッジ配布 |

docs/installation.md

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -223,7 +223,7 @@ Webhook エンドポイントへのアクセスを GitHub の IP 範囲に制限
223223

224224
### 8. チャンネル通知(オプション)
225225

226-
ローカル MCP ブリッジは Claude Code の `claude/channel` 機能をサポートしています。有効にすると、新しい Webhook イベントが SSE 経由でリアルタイムにセッションにプッシュされます。Claude Code CLI でのみ利用可能です。
226+
ローカル MCP ブリッジは Claude Code の `claude/channel` 機能をサポートしています。有効にすると、新しい Webhook イベントが WebSocket 経由でリアルタイムにセッションにプッシュされます。Claude Code CLI でのみ利用可能です。
227227

228228
MCP クライアント設定で `WEBHOOK_CHANNEL=1` を設定し(上記 [Claude Code CLI](#claude-code-cli--npx) 参照)、チャンネルをロード:
229229

@@ -237,7 +237,7 @@ claude --dangerously-load-development-channels server:github-webhook-mcp
237237

238238
1. **Webhook 受信テスト:** GitHub App の設定ページ → **Advanced****Recent Deliveries** で配信状況を確認
239239
2. **MCP 接続テスト:** MCP クライアントから `get_pending_status` ツールを呼び出して応答を確認
240-
3. **SSE テスト:** `curl -N https://<your-worker>/events` でストリーム接続を確認
240+
3. **WebSocket テスト:** `wscat -c wss://<your-worker>/events` でストリーム接続を確認(SSE: `curl -N https://<your-worker>/events`
241241

242242
### トラブルシューティング
243243

mcp-server/package.json

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -18,8 +18,7 @@
1818
"pack:mcpb": "mcpb pack"
1919
},
2020
"dependencies": {
21-
"@modelcontextprotocol/sdk": "^1.0.0",
22-
"eventsource": "^2.0.2"
21+
"@modelcontextprotocol/sdk": "^1.0.0"
2322
},
2423
"devDependencies": {
2524
"@anthropic-ai/mcpb": "^2.1.0"

mcp-server/server/index.js

Lines changed: 59 additions & 43 deletions
Original file line numberDiff line numberDiff line change
@@ -4,7 +4,8 @@
44
*
55
* Thin stdio MCP server that proxies tool calls to a remote
66
* Cloudflare Worker + Durable Object backend via Streamable HTTP.
7-
* Optionally listens to SSE for real-time channel notifications.
7+
* Optionally listens via WebSocket for real-time channel notifications
8+
* (enables DO hibernation on the Worker side).
89
* Authenticates via OAuth 2.1 with PKCE (localhost callback).
910
*
1011
* Discord MCP pattern: data lives in the cloud, local MCP is a thin bridge.
@@ -547,67 +548,75 @@ server.setRequestHandler(CallToolRequestSchema, async (req) => {
547548
}
548549
});
549550

550-
// ── SSE Listener → Channel Notifications ─────────────────────────────────────
551+
// ── WebSocket Listener → Channel Notifications ──────────────────────────────
551552

552553
/** Track whether OAuth has been established (first successful tool call). */
553554
let _oauthEstablished = false;
554555

555556
function markOAuthEstablished() {
556557
if (!_oauthEstablished) {
557558
_oauthEstablished = true;
558-
if (CHANNEL_ENABLED && !_sseConnected) {
559-
process.stderr.write("[github-webhook-mcp] OAuth established, starting SSE connection\n");
560-
connectSSE();
559+
if (CHANNEL_ENABLED && !_wsConnected) {
560+
process.stderr.write("[github-webhook-mcp] OAuth established, starting WebSocket connection\n");
561+
connectWebSocket();
561562
}
562563
}
563564
}
564565

565-
let _sseConnected = false;
566+
let _wsConnected = false;
566567

567-
async function connectSSE() {
568-
let EventSourceImpl;
569-
try {
570-
EventSourceImpl = (await import("eventsource")).default;
571-
} catch {
572-
// eventsource not installed — skip SSE
573-
return;
574-
}
568+
async function connectWebSocket() {
569+
const wsUrl = WORKER_URL.replace(/^http/, "ws") + "/events";
575570

576-
_sseConnected = true;
571+
_wsConnected = true;
577572
let retryCount = 0;
578573
const MAX_RETRY_DELAY = 60_000; // 60 seconds
579574
const BASE_RETRY_DELAY = 1_000; // 1 second
580575

581-
async function attemptConnection() {
576+
async function connect() {
582577
let token;
583578
try {
584579
token = await getAccessToken();
585580
} catch (err) {
586-
process.stderr.write(`[github-webhook-mcp] SSE: failed to get access token: ${err}\n`);
581+
process.stderr.write(`[github-webhook-mcp] WebSocket: failed to get access token: ${err}\n`);
587582
scheduleRetry();
588583
return;
589584
}
590585

591586
if (!token) {
592-
process.stderr.write("[github-webhook-mcp] SSE: no access token available, will retry\n");
587+
process.stderr.write("[github-webhook-mcp] WebSocket: no access token available, will retry\n");
593588
scheduleRetry();
594589
return;
595590
}
596591

597-
const sseUrl = `${WORKER_URL}/events`;
598-
const es = new EventSourceImpl(sseUrl, {
599-
headers: { Authorization: `Bearer ${token}` },
600-
});
592+
let ws;
593+
let pingTimer = null;
594+
595+
try {
596+
ws = new WebSocket(wsUrl, { headers: { Authorization: `Bearer ${token}` } });
597+
} catch (err) {
598+
process.stderr.write(`[github-webhook-mcp] WebSocket: failed to create connection: ${err}\n`);
599+
scheduleRetry();
600+
return;
601+
}
601602

602-
es.onopen = () => {
603+
ws.addEventListener("open", () => {
603604
retryCount = 0; // Reset backoff on successful connection
604-
process.stderr.write("[github-webhook-mcp] SSE: connected\n");
605-
};
605+
process.stderr.write("[github-webhook-mcp] WebSocket: connected\n");
606+
// Send periodic pings to keep connection alive (25s keepalive)
607+
pingTimer = setInterval(() => {
608+
if (ws.readyState === WebSocket.OPEN) {
609+
ws.send("ping");
610+
}
611+
}, 25_000);
612+
});
606613

607-
es.onmessage = (event) => {
614+
ws.addEventListener("message", (event) => {
608615
try {
609-
const data = JSON.parse(event.data);
610-
if ("heartbeat" in data || "status" in data) return;
616+
const data = JSON.parse(typeof event.data === "string" ? event.data : event.data.toString());
617+
618+
// Skip status, pong, heartbeat messages
619+
if ("status" in data || "pong" in data || "heartbeat" in data) return;
611620
if (!data.summary) return;
612621

613622
const s = data.summary;
@@ -634,29 +643,36 @@ async function connectSSE() {
634643
} catch {
635644
// Ignore parse errors
636645
}
637-
};
646+
});
638647

639-
es.onerror = (err) => {
640-
const status = err && err.status;
641-
if (status === 401) {
642-
process.stderr.write("[github-webhook-mcp] SSE: 401 unauthorized, closing and retrying with fresh token\n");
643-
_cachedTokens = null; // Force token refresh on next attempt
648+
ws.addEventListener("close", (event) => {
649+
if (pingTimer) clearInterval(pingTimer);
650+
pingTimer = null;
651+
const code = event.code;
652+
if (code === 1008 || code === 4401) {
653+
// Policy violation or unauthorized — refresh token
654+
process.stderr.write(`[github-webhook-mcp] WebSocket: closed with code ${code}, refreshing token\n`);
655+
_cachedTokens = null;
644656
} else {
645-
process.stderr.write(`[github-webhook-mcp] SSE: connection error${status ? ` (status ${status})` : ""}\n`);
657+
process.stderr.write(`[github-webhook-mcp] WebSocket: closed (code ${code})\n`);
646658
}
647-
es.close();
648659
scheduleRetry();
649-
};
660+
});
661+
662+
ws.addEventListener("error", () => {
663+
process.stderr.write("[github-webhook-mcp] WebSocket: connection error\n");
664+
// Will trigger close event, which handles reconnect
665+
});
650666
}
651667

652668
function scheduleRetry() {
653669
const delay = Math.min(BASE_RETRY_DELAY * Math.pow(2, retryCount), MAX_RETRY_DELAY);
654670
retryCount++;
655-
process.stderr.write(`[github-webhook-mcp] SSE: retrying in ${Math.round(delay / 1000)}s (attempt ${retryCount})\n`);
656-
setTimeout(() => attemptConnection(), delay);
671+
process.stderr.write(`[github-webhook-mcp] WebSocket: retrying in ${Math.round(delay / 1000)}s (attempt ${retryCount})\n`);
672+
setTimeout(() => void connect(), delay);
657673
}
658674

659-
attemptConnection();
675+
await connect();
660676
}
661677

662678
// ── Start ────────────────────────────────────────────────────────────────────
@@ -666,15 +682,15 @@ await server.connect(transport);
666682

667683
if (CHANNEL_ENABLED) {
668684
// Check if tokens already exist from a previous session.
669-
// If so, start SSE immediately. Otherwise, defer until after the first
685+
// If so, start WebSocket immediately. Otherwise, defer until after the first
670686
// successful tool call establishes OAuth.
671687
loadTokens().then((tokens) => {
672688
if (tokens && tokens.access_token) {
673689
_cachedTokens = tokens;
674690
markOAuthEstablished();
675691
}
676-
// If no tokens, SSE will start after the first successful tool call
692+
// If no tokens, WebSocket will start after the first successful tool call
677693
}).catch(() => {
678-
// Token load failed, SSE will start after first tool call
694+
// Token load failed, WebSocket will start after first tool call
679695
});
680696
}

0 commit comments

Comments
 (0)