Skip to content

Commit 923d4fd

Browse files
committed
feat: SSE streaming for live pipeline output (Phase 4)
mmcp_cloud/streaming.py: - SSE stream_chat() async generator - Proxies OpenRouter streaming API - Emits start/token/usage/done/error events - POST /v1/chat/completions/stream endpoint mmcp_core/stream_executor.py: - stream_execute() async generator for plan execution - Real-time token streaming per step - Context chaining between steps - Error recovery (retry/skip/fallback) - Tool action support Event types: plan_start, step_start, step_token, step_done, step_error, plan_done
1 parent 84d5253 commit 923d4fd

3 files changed

Lines changed: 399 additions & 0 deletions

File tree

‎python/mmcp_cloud/server.py‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -50,6 +50,10 @@ async def lifespan(app: FastAPI):
5050
lifespan=lifespan,
5151
)
5252

53+
# Register SSE streaming routes
54+
from .streaming import register_stream_routes
55+
register_stream_routes(app)
56+
5357

5458
# ── Auth helpers ────────────────────────────────────────────────────────────
5559

‎python/mmcp_cloud/streaming.py‎

Lines changed: 183 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,183 @@
1+
"""
2+
MMCP Cloud — Server-Sent Events (SSE) streaming for live pipeline output.
3+
4+
Usage:
5+
POST /v1/chat/completions/stream — streaming proxy with SSE
6+
7+
Events:
8+
data: {"type": "start", "model": "...", "step": 1}
9+
data: {"type": "token", "content": "...", "index": 0}
10+
data: {"type": "usage", "tokens": 150, "cost": 0.001}
11+
data: {"type": "done"}
12+
data: {"type": "error", "message": "..."}
13+
"""
14+
from __future__ import annotations
15+
import json
16+
import os
17+
import time
18+
from typing import AsyncGenerator
19+
20+
import httpx
21+
from fastapi import HTTPException, Request
22+
from starlette.responses import StreamingResponse
23+
24+
from .database import get_user_by_key, log_usage
25+
from .billing import apply_markup, check_rate_limit, get_usage, get_month_start
26+
27+
28+
OPENROUTER_URL = "https://openrouter.ai/api/v1/chat/completions"
29+
30+
31+
def _sse_event(data: dict) -> str:
32+
"""Format a dict as an SSE event string."""
33+
return f"data: {json.dumps(data)}\n\n"
34+
35+
36+
async def stream_chat(
37+
model: str,
38+
messages: list[dict],
39+
max_tokens: int,
40+
temperature: float,
41+
user: dict,
42+
) -> AsyncGenerator[str, None]:
43+
"""Stream chat completions via SSE.
44+
45+
Yields SSE-formatted events as the upstream model generates tokens.
46+
"""
47+
or_key = os.environ.get("OPENROUTER_API_KEY")
48+
if not or_key:
49+
yield _sse_event({"type": "error", "message": "Server misconfigured: no upstream API key"})
50+
return
51+
52+
# Check rate limit
53+
month_start = get_month_start()
54+
usage = get_usage(user["user_id"], since=month_start)
55+
limit_check = check_rate_limit(usage["runs"], user["plan"])
56+
57+
if not limit_check["allowed"]:
58+
yield _sse_event({
59+
"type": "error",
60+
"message": f"Monthly limit reached ({limit_check['limit']} runs). Upgrade your plan.",
61+
})
62+
return
63+
64+
# Start event
65+
yield _sse_event({"type": "start", "model": model})
66+
67+
start = time.time()
68+
total_content = ""
69+
total_tokens = 0
70+
71+
try:
72+
async with httpx.AsyncClient(timeout=120.0) as client:
73+
async with client.stream(
74+
"POST",
75+
OPENROUTER_URL,
76+
headers={
77+
"Content-Type": "application/json",
78+
"Authorization": f"Bearer {or_key}",
79+
"HTTP-Referer": "https://mmcp.dev",
80+
"X-Title": "MMCP Cloud",
81+
},
82+
json={
83+
"model": model,
84+
"max_tokens": max_tokens,
85+
"temperature": temperature,
86+
"messages": messages,
87+
"stream": True,
88+
},
89+
) as resp:
90+
if resp.status_code != 200:
91+
body = await resp.aread()
92+
yield _sse_event({"type": "error", "message": f"Upstream error: {resp.status_code}"})
93+
return
94+
95+
async for line in resp.aiter_lines():
96+
if not line.startswith("data: "):
97+
continue
98+
99+
payload = line[6:].strip()
100+
if payload == "[DONE]":
101+
break
102+
103+
try:
104+
chunk = json.loads(payload)
105+
delta = chunk.get("choices", [{}])[0].get("delta", {})
106+
content = delta.get("content", "")
107+
108+
if content:
109+
total_content += content
110+
yield _sse_event({"type": "token", "content": content})
111+
112+
# Track token usage from final chunk
113+
if "usage" in chunk:
114+
total_tokens = chunk["usage"].get("total_tokens", 0)
115+
116+
except json.JSONDecodeError:
117+
continue
118+
119+
except httpx.ReadTimeout:
120+
yield _sse_event({"type": "error", "message": "Upstream timeout"})
121+
return
122+
except Exception as e:
123+
yield _sse_event({"type": "error", "message": str(e)})
124+
return
125+
126+
# Calculate cost
127+
duration_ms = int((time.time() - start) * 1000)
128+
# Estimate cost if not provided (rough estimate based on model)
129+
estimated_cost = total_tokens * 0.000001 # $1/1M tokens baseline
130+
billed_cost = apply_markup(estimated_cost, user["plan"])
131+
132+
# Log usage
133+
log_usage(
134+
user_id=user["user_id"],
135+
tokens=total_tokens,
136+
cost_usd=billed_cost,
137+
model=model,
138+
pipeline="stream",
139+
)
140+
141+
# Final usage event
142+
yield _sse_event({
143+
"type": "usage",
144+
"tokens": total_tokens,
145+
"cost_usd": round(billed_cost, 6),
146+
"duration_ms": duration_ms,
147+
"runs_remaining": limit_check["remaining"],
148+
})
149+
150+
# Done event
151+
yield _sse_event({"type": "done"})
152+
153+
154+
def register_stream_routes(app):
155+
"""Register streaming endpoints on the FastAPI app."""
156+
157+
@app.post("/v1/chat/completions/stream")
158+
async def chat_completions_stream(request: Request):
159+
"""Stream chat completions via Server-Sent Events."""
160+
# Auth
161+
auth = request.headers.get("Authorization", "")
162+
if not auth.startswith("Bearer mmcp_"):
163+
raise HTTPException(401, "Invalid API key")
164+
api_key = auth.replace("Bearer ", "")
165+
user = get_user_by_key(api_key)
166+
if not user:
167+
raise HTTPException(401, "Invalid or revoked API key")
168+
169+
body = await request.json()
170+
model = body.get("model", "anthropic/claude-3.5-haiku")
171+
messages = body.get("messages", [])
172+
max_tokens = body.get("max_tokens", 4096)
173+
temperature = body.get("temperature", 0.7)
174+
175+
return StreamingResponse(
176+
stream_chat(model, messages, max_tokens, temperature, user),
177+
media_type="text/event-stream",
178+
headers={
179+
"Cache-Control": "no-cache",
180+
"Connection": "keep-alive",
181+
"X-Accel-Buffering": "no", # Disable nginx buffering
182+
},
183+
)

0 commit comments

Comments
 (0)