返回 last30days-skill
hosted.py
根目录 / skills / last30days / scripts / lib / hosted.py
1 """Remote API client for last30days (optional hosted-backend mode).
2
3 When both LAST30DAYS_API_KEY and LAST30DAYS_API_BASE are set, the engine
4 submits the topic to the configured remote API, polls until the run reaches a
5 terminal status, streams narration progress to stderr, and renders the
6 server's report. No local provider keys are required in this mode. The
7 endpoint comes only from LAST30DAYS_API_BASE - there is no built-in default.
8
9 Contract (API v1):
10 POST {base}/search Authorization: Bearer <key>
11 {"query": ..., "depth": "quick"|"default"|"deep",
12 "register"?: "exec"|"dev"|"creator"|"eli5"}
13 -> 200 {"search_id": "<uuid>", "status": "running"}
14 -> 200 clarify payload {"needs_clarification": true, ...}
15 -> 401 {"error"} / 402 {"error","requires_credits",
16 "balance","needed"} / 429 {"error"}
17 GET {base}/search?id=<uuid> same auth header; poll until status is
18 terminal ("complete" | "error"). Running rows carry
19 "stderr" (narration + engine lines) and "eta_ms";
20 terminal complete rows carry "synthesis_text" and
21 "raw_markdown" (stderr stripped).
22
23 This module carries ZERO pricing, rate-card, cost, or billing logic.
24 Balance/credit numbers are only ever displayed verbatim from API responses.
25 The API key is never printed, logged, or persisted by this module.
26 """
27
28 from __future__ import annotations
29
30 import json
31 import os
32 import re
33 import sys
34 import time
35
36 from . import env, http, usage
37 from .log import source_log
38
39 # Distinct exit code for the clarify gate so the invoking model can tell
40 # "re-run with a chosen angle" apart from a plain failure (1).
41 EXIT_CLARIFY = 3
42
43 POLL_INITIAL_DELAY = 3.0
44 POLL_MAX_DELAY = 10.0
45 POLL_TIMEOUT_SECONDS = 15 * 60
46 # GET is idempotent: retry a few times across network blips before giving up.
47 POLL_NETWORK_RETRIES = 3
48 # Cadence for the compact elapsed/eta progress line (seconds).
49 PROGRESS_LINE_INTERVAL = 15.0
50
51 NARRATE_PREFIX = "[narrate] step="
52 TERMINAL_STATUSES = {"complete", "error"}
53
54
55 def _err(msg: str) -> None:
56 source_log("hosted", msg, tty_only=False)
57
58
59 def _api_base() -> str:
60 # Endpoint comes only from the environment - no built-in default. Hosted
61 # mode is gated on this being set (see last30days.py), so by the time this
62 # is called it is populated; an empty value means "not configured".
63 return (os.environ.get("LAST30DAYS_API_BASE") or "").rstrip("/")
64
65
66 def _billing_url() -> str:
67 """Derive a billing link from the configured base, so no URL is hardcoded.
68 Convention: the base is the API-version root (e.g. ends in /api/v1); drop
69 that segment and point at the account's billing page."""
70 base = _api_base()
71 root = re.sub(r"/api/v\d+$", "", base)
72 return f"{root}/dashboard/billing"
73
74
75 def _auth_headers() -> dict[str, str]:
76 # Key is read at call time and placed only in the request header;
77 # it must never be interpolated into any log or output line.
78 key = env.read_secret_env("LAST30DAYS_API_KEY") or ""
79 return {"Authorization": f"Bearer {key}"}
80
81
82 def submit(query: str, depth: str, register: str = "default") -> dict:
83 """POST the search. retries=1: a blind POST retry could double-submit."""
84 payload = {"query": query, "depth": depth}
85 if register != "default":
86 payload["register"] = register
87 usage.begin("hosted")
88 return http.post(
89 f"{_api_base()}/search",
90 json_data=payload,
91 headers=_auth_headers(),
92 retries=1,
93 )
94
95
96 def poll(search_id: str) -> dict:
97 """GET the search row once. Callers own the retry loop (GET is idempotent)."""
98 return http.get(
99 f"{_api_base()}/search",
100 headers=_auth_headers(),
101 params={"id": search_id},
102 retries=1,
103 )
104
105
106 def _parse_error_body(exc: http.HTTPError) -> dict:
107 if not exc.body:
108 return {}
109 try:
110 parsed = json.loads(exc.body)
111 except (json.JSONDecodeError, TypeError):
112 return {}
113 return parsed if isinstance(parsed, dict) else {}
114
115
116 def _handle_http_error(exc: http.HTTPError) -> int:
117 body = _parse_error_body(exc)
118 if exc.status_code == 401:
119 _err(
120 "API key rejected: invalid or revoked. Check "
121 "LAST30DAYS_API_KEY (and LAST30DAYS_API_BASE), or unset them "
122 "to fall back to local sources."
123 )
124 return 1
125 if exc.status_code == 402:
126 _err(f"API: {body.get('error') or 'insufficient credits.'}")
127 if body.get("balance") is not None or body.get("needed") is not None:
128 _err(
129 f"Balance: {body.get('balance')} credits. "
130 f"Needed for this search: {body.get('needed')} credits."
131 )
132 _err(f"Add credits at {_billing_url()}")
133 return 1
134 if exc.status_code == 429:
135 _err(
136 f"API rate limit hit: "
137 f"{body.get('error') or 'too many requests.'} "
138 "Wait a minute and re-run."
139 )
140 return 1
141 _err(f"API request failed: {exc}")
142 return 1
143
144
145 def _handle_clarify(resp: dict) -> int:
146 question = resp.get("question") or "The API needs a clarification before searching."
147 options = resp.get("options") or []
148 _err(f"Clarification needed before this search runs: {question}")
149 for index, option in enumerate(options, 1):
150 label = option if isinstance(option, str) else json.dumps(option)
151 sys.stderr.write(f" {index}. {label}\n")
152 sys.stderr.flush()
153 _err(
154 "No search was started. Re-run last30days with the chosen angle "
155 "folded into the topic text."
156 )
157 return EXIT_CLARIFY
158
159
160 def _print_new_narration(stderr_blob: str, seen: set[str]) -> bool:
161 """Print each '[narrate] step=' line once, verbatim. Returns True if any new."""
162 printed = False
163 for line in stderr_blob.splitlines():
164 if line.startswith(NARRATE_PREFIX) and line not in seen:
165 seen.add(line)
166 sys.stderr.write(f"{line}\n")
167 printed = True
168 if printed:
169 sys.stderr.flush()
170 return printed
171
172
173 def _print_progress_line(elapsed: float, eta_ms) -> None:
174 line = f"elapsed {int(elapsed)}s"
175 if isinstance(eta_ms, (int, float)) and eta_ms > 0:
176 line += f", eta ~{int(eta_ms / 1000)}s"
177 _err(line)
178
179
180 def _poll_with_retry(search_id: str) -> dict | None:
181 """Poll once, retrying transient network failures. None means give up
182 (a user-facing message has already been printed)."""
183 last_error: http.HTTPError | None = None
184 for attempt in range(POLL_NETWORK_RETRIES):
185 try:
186 return poll(search_id)
187 except http.HTTPError as exc:
188 if exc.status_code is not None and 400 <= exc.status_code < 500 and exc.status_code != 429:
189 _handle_http_error(exc)
190 return None
191 # Network blip / timeout / 5xx / 429: GET is idempotent, retry.
192 last_error = exc
193 if attempt < POLL_NETWORK_RETRIES - 1:
194 time.sleep(POLL_INITIAL_DELAY)
195 _err(
196 f"API unreachable while polling search {search_id} "
197 f"after {POLL_NETWORK_RETRIES} attempts: {last_error}"
198 )
199 return None
200
201
202 def _slugify(value: str) -> str:
203 slug = re.sub(r"[^a-z0-9]+", "-", value.lower()).strip("-")
204 return slug or "last30days"
205
206
207 def _save_output(topic: str, content: str, emit: str, save_dir: str, suffix: str):
208 """Mirror local save_output() naming: <slug>-raw[-suffix].<ext>."""
209 from datetime import datetime
210 from pathlib import Path
211
212 path = Path(save_dir).expanduser().resolve()
213 path.mkdir(parents=True, exist_ok=True)
214 slug = _slugify(topic)
215 extension = "json" if emit == "json" else "md"
216 suffix_part = f"-{suffix}" if suffix else ""
217 base = path / f"{slug}-raw{suffix_part}.{extension}"
218 date_str = datetime.now().strftime('%Y-%m-%d')
219 candidates = [base]
220 candidates.append(path / f"{slug}-raw{suffix_part}-{date_str}.{extension}")
221 for i in range(1, 100):
222 candidates.append(path / f"{slug}-raw{suffix_part}-{date_str}-{i}.{extension}")
223 encoded = content.encode("utf-8")
224 for candidate in candidates:
225 try:
226 fd = os.open(candidate, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o644)
227 except FileExistsError:
228 continue
229 with os.fdopen(fd, "wb") as f:
230 f.write(encoded)
231 return candidate
232 # Fallback: all 101 candidates existed (extremely unlikely).
233 raise RuntimeError(
234 f"_save_output: could not find a unique filename after 101 attempts in {path}"
235 )
236
237
238 def _render_complete(row: dict, topic: str, emit: str, save_dir, save_suffix: str) -> int:
239 synthesis = row.get("synthesis_text") or ""
240 raw_markdown = row.get("raw_markdown") or ""
241 if emit == "json":
242 payload = {
243 key: row.get(key)
244 for key in ("id", "status", "synthesis_text", "raw_markdown")
245 if key in row
246 }
247 rendered = json.dumps(payload, indent=2, sort_keys=True)
248 save_content = rendered
249 else:
250 # The server report is the content source; it already synthesized.
251 # All markdown-ish emit modes print the synthesis text as-is.
252 rendered = synthesis or raw_markdown
253 save_content = raw_markdown or synthesis
254 if save_dir:
255 out_path = _save_output(topic, save_content, emit, save_dir, save_suffix)
256 sys.stderr.write(f"[last30days] Saved output to {out_path}\n")
257 sys.stderr.flush()
258 print(rendered)
259 return 0
260
261
262 def run_hosted(
263 topic: str,
264 depth: str,
265 *,
266 emit: str = "compact",
267 save_dir=None,
268 save_suffix: str = "",
269 register: str = "default",
270 ) -> int:
271 """Submit topic to the remote API, poll to terminal status, render report."""
272 _err(f"Running via last30days API ({_api_base()}), depth={depth}")
273 try:
274 resp = submit(topic, depth, register=register)
275 except http.HTTPError as exc:
276 return _handle_http_error(exc)
277
278 if resp.get("needs_clarification"):
279 return _handle_clarify(resp)
280
281 search_id = resp.get("search_id")
282 if not search_id:
283 _err(f"Unexpected API response (no search_id): {json.dumps(resp)[:200]}")
284 return 1
285 _err(f"Search submitted (id: {search_id}). Polling for results...")
286
287 started = time.monotonic()
288 delay = POLL_INITIAL_DELAY
289 seen_narration: set[str] = set()
290 last_progress_line = 0.0
291 while True:
292 elapsed = time.monotonic() - started
293 if elapsed > POLL_TIMEOUT_SECONDS:
294 _err(
295 f"Search did not finish within "
296 f"{POLL_TIMEOUT_SECONDS // 60} minutes (id: {search_id}). "
297 "It may still complete server-side; check the dashboard."
298 )
299 return 1
300 time.sleep(delay)
301 delay = min(delay * 2, POLL_MAX_DELAY)
302
303 row = _poll_with_retry(search_id)
304 if row is None:
305 return 1
306
307 status = row.get("status")
308 narrated = _print_new_narration(row.get("stderr") or "", seen_narration)
309 elapsed = time.monotonic() - started
310 if status not in TERMINAL_STATUSES and (
311 narrated or elapsed - last_progress_line >= PROGRESS_LINE_INTERVAL or last_progress_line == 0.0
312 ):
313 _print_progress_line(elapsed, row.get("eta_ms"))
314 last_progress_line = elapsed
315
316 if status == "error":
317 _err(f"Search failed: {row.get('error') or 'unknown server error'}")
318 return 1
319 if status == "complete":
320 _err(f"Search complete in {int(elapsed)}s.")
321 return _render_complete(row, topic, emit, save_dir, save_suffix)
322 # pending | running -> keep polling
323
323 lines PYTHON