返回 last30days-skill
parallel_mcp.py
根目录 / skills / last30days / scripts / lib / parallel_mcp.py
1 """Opt-in Parallel Search MCP adapter using stdlib Streamable HTTP."""
2
3 from __future__ import annotations
4
5 import json
6 import urllib.error
7 import urllib.request
8 from typing import Any
9 from urllib.parse import urlparse
10
11 from . import dates, http, usage
12
13 PARALLEL_MCP_URL = "https://search.parallel.ai/mcp"
14 _PROTOCOL_VERSION = "2025-03-26"
15 _MAX_RESPONSE_BYTES = 4 * 1024 * 1024
16
17
18 class _NoRedirect(urllib.request.HTTPRedirectHandler):
19 def redirect_request(self, req, fp, code, msg, headers, newurl):
20 # The endpoint is fixed. Never forward credentials, session IDs, or
21 # search data to a redirect target, including an HTTPS downgrade.
22 return None
23
24
25 def _matching_response(payload: bytes, request_id: int | None) -> dict[str, Any] | None:
26 value = json.loads(payload)
27 # The negotiated 2025-03-26 protocol permits batched SSE messages.
28 for message in value if isinstance(value, list) else [value]:
29 if not isinstance(message, dict) or message.get("jsonrpc") != "2.0":
30 raise RuntimeError("Parallel MCP returned an invalid JSON-RPC response")
31 if message.get("id") == request_id and ("result" in message or "error" in message):
32 return message
33 return None
34
35
36 def _read_response(response: Any, request_id: int | None) -> dict[str, Any]:
37 if "text/event-stream" not in response.headers.get("Content-Type", "").lower():
38 payload = response.read(_MAX_RESPONSE_BYTES + 1)
39 if len(payload) > _MAX_RESPONSE_BYTES:
40 raise RuntimeError("Parallel MCP response exceeded 4 MiB")
41 if request_id is None and not payload:
42 return {}
43 result = _matching_response(payload, request_id) if payload else None
44 if result is not None:
45 return result
46 else:
47 data = []
48 total = 0
49 while True:
50 line = response.readline(_MAX_RESPONSE_BYTES - total + 1)
51 total += len(line)
52 if total > _MAX_RESPONSE_BYTES:
53 raise RuntimeError("Parallel MCP response exceeded 4 MiB")
54 if not line:
55 break
56 line = line.rstrip(b"\r\n")
57 if not line and data:
58 result = _matching_response(b"\n".join(data), request_id)
59 if result is not None:
60 # A valid response completes the request, even if the
61 # server keeps the stream open or sends more notifications.
62 return result
63 data = []
64 elif line.startswith(b"data:"):
65 value = line[5:]
66 data.append(value[1:] if value.startswith(b" ") else value)
67 raise RuntimeError("Parallel MCP response is missing the requested JSON-RPC result")
68
69
70 def _request(
71 message: dict[str, Any] | None, api_key: str | None, session_id: str | None = None,
72 ) -> tuple[dict[str, Any], str | None]:
73 headers = {
74 "Accept": "application/json, text/event-stream",
75 "Content-Type": "application/json",
76 "User-Agent": http.USER_AGENT,
77 }
78 if api_key:
79 headers["Authorization"] = f"Bearer {api_key}"
80 if session_id:
81 headers["Mcp-Session-Id"] = session_id
82 if message is None or message.get("method") != "initialize":
83 headers["MCP-Protocol-Version"] = _PROTOCOL_VERSION
84 request = urllib.request.Request(
85 PARALLEL_MCP_URL,
86 data=json.dumps(message, separators=(",", ":")).encode("utf-8") if message is not None else None,
87 headers=headers,
88 method="POST" if message is not None else "DELETE",
89 )
90 if api_key and message is not None and message.get("method") == "tools/call":
91 usage.begin("parallel")
92 with urllib.request.build_opener(_NoRedirect()).open(request, timeout=http.DEFAULT_TIMEOUT) as response:
93 result = _read_response(response, message.get("id")) if message is not None else {}
94 return result, response.headers.get("Mcp-Session-Id") or session_id
95
96
97 def _result(response: dict[str, Any]) -> dict[str, Any]:
98 if response.get("error"):
99 error = response["error"]
100 detail = error.get("message") if isinstance(error, dict) else str(error)
101 raise RuntimeError(f"Parallel MCP error: {detail}")
102 result = response.get("result")
103 if not isinstance(result, dict):
104 raise RuntimeError("Parallel MCP response is missing its result")
105 return result
106
107
108 def _search_rows(tool_result: dict[str, Any]) -> list[dict[str, Any]]:
109 payloads = [tool_result.get("structuredContent")]
110 for content in tool_result.get("content") or []:
111 if not isinstance(content, dict) or content.get("type") != "text":
112 continue
113 try:
114 payloads.append(json.loads(content.get("text") or ""))
115 except (TypeError, ValueError):
116 continue
117 for candidates in payloads:
118 if isinstance(candidates, dict) and "data" in candidates and "results" not in candidates:
119 candidates = candidates["data"]
120 if isinstance(candidates, dict):
121 candidates = candidates.get("results")
122 if isinstance(candidates, list):
123 return [row for row in candidates if isinstance(row, dict)]
124 raise RuntimeError("Parallel MCP web_search returned no results array")
125
126
127 def search(
128 query: str, date_range: tuple[str, str], api_key: str | None = None, count: int = 5,
129 ) -> tuple[list[dict[str, Any]], dict[str, Any]]:
130 """Discover and invoke hosted ``web_search`` after explicit dispatcher opt-in."""
131 initialized, session_id = _request(
132 {
133 "jsonrpc": "2.0",
134 "id": 1,
135 "method": "initialize",
136 "params": {
137 "protocolVersion": _PROTOCOL_VERSION,
138 "capabilities": {},
139 "clientInfo": {"name": "last30days-skill", "version": "3"},
140 },
141 },
142 api_key,
143 )
144 try:
145 if _result(initialized).get("protocolVersion") != _PROTOCOL_VERSION:
146 raise RuntimeError("Parallel MCP negotiated an unsupported protocol version")
147 _request(
148 {"jsonrpc": "2.0", "method": "notifications/initialized", "params": {}},
149 api_key,
150 session_id,
151 )
152 request_id = 2
153 params = {}
154 seen_cursors = set()
155 while True:
156 tools_response, _ = _request(
157 {"jsonrpc": "2.0", "id": request_id, "method": "tools/list", "params": params},
158 api_key,
159 session_id,
160 )
161 request_id += 1
162 page = _result(tools_response)
163 tools = page.get("tools") or []
164 if any(isinstance(tool, dict) and tool.get("name") == "web_search" for tool in tools):
165 break
166 cursor = page.get("nextCursor")
167 if not isinstance(cursor, str) or not cursor or cursor in seen_cursors or len(seen_cursors) >= 20:
168 raise RuntimeError("Parallel MCP did not advertise web_search")
169 seen_cursors.add(cursor)
170 params = {"cursor": cursor}
171 objective = f"Find useful public web evidence about {query} from {date_range[0]} through {date_range[1]}."
172 called, _ = _request(
173 {
174 "jsonrpc": "2.0",
175 "id": request_id,
176 "method": "tools/call",
177 "params": {
178 "name": "web_search",
179 "arguments": {"objective": objective, "search_queries": [query]},
180 },
181 },
182 api_key,
183 session_id,
184 )
185 finally:
186 if session_id:
187 try:
188 _request(None, api_key, session_id)
189 except urllib.error.HTTPError as error:
190 error.close()
191 except (OSError, RuntimeError, ValueError):
192 # Cleanup is best-effort: unsupported DELETEs or a network
193 # failure must not hide search results or the original error.
194 pass
195 tool_result = _result(called)
196 if tool_result.get("isError"):
197 raise RuntimeError("Parallel MCP web_search reported an error")
198 items = []
199 for row in _search_rows(tool_result):
200 if len(items) >= count:
201 break
202 url = str(row.get("url") or "")
203 try:
204 parsed_url = urlparse(url)
205 except ValueError:
206 continue
207 if parsed_url.scheme not in ("http", "https") or not parsed_url.netloc:
208 continue
209 raw_date = row.get("publish_date")
210 parsed_date = dates.parse_date(raw_date[:10]) if isinstance(raw_date, str) else None
211 pub_date = parsed_date.date().isoformat() if parsed_date else None
212 # Match the paid grounding backends: only dated evidence inside the
213 # requested window qualifies, including when --as-of is historical.
214 if not pub_date or not date_range[0] <= pub_date <= date_range[1]:
215 continue
216 excerpts = row.get("excerpts") or []
217 if isinstance(excerpts, str):
218 excerpts = [excerpts]
219 snippet = "\n".join(str(value) for value in excerpts if value) if isinstance(excerpts, list) else ""
220 items.append({
221 "id": f"WPM{len(items) + 1}",
222 "title": row.get("title") or parsed_url.netloc,
223 "url": url,
224 "source_domain": parsed_url.netloc.strip().lower(),
225 "snippet": snippet[:500],
226 "date": pub_date,
227 "relevance": 0.8,
228 "why_relevant": "Parallel Search MCP result",
229 })
230 return items, {
231 "label": "parallel-mcp",
232 "webSearchQueries": [query],
233 "resultCount": len(items),
234 }
235
235 lines PYTHON