#!/usr/bin/env python3 """Manual SSE parser (no SDK): accumulates text, tool input JSON and thinking; handles ping/error/unknown events. STATUS: LIVE_VERIFIED 2026-09-18 (stdlib urllib via scripts/live.py). Run: .venv/bin/python examples/anthropic/streaming/manual_sse_parser.py """ import json import sys from pathlib import Path sys.path.insert(0, str(Path(__file__).resolve().parents[3])) from scripts.live import anthropic_request # noqa: E402 def parse_sse(lines): """Yield (event_name, data_dict). SSE frames are separated by blank lines; we only use event:/data: fields.""" name, data = None, [] for line in lines: if line == "": if data: yield name, json.loads("\n".join(data)) name, data = None, [] elif line.startswith("event:"): name = line[6:].strip() elif line.startswith("data:"): data.append(line[5:].strip()) if data: yield name, json.loads("\n".join(data)) body = {"model": "claude-haiku-4-5-20251001", "max_tokens": 64, "stream": True, "tools": [{"name": "get_weather", "description": "Weather for a city", "input_schema": {"type": "object", "properties": {"city": {"type": "string"}}, "required": ["city"]}}], "tool_choice": {"type": "tool", "name": "get_weather"}, "messages": [{"role": "user", "content": "Weather in Paris?"}]} status, lines, headers = anthropic_request("POST", "/v1/messages", body, stream=True, est_cost_usd=0.0009, note="example manual_sse_parser.py") assert status == 200, status blocks: dict[int, dict] = {} message, order = {}, [] for name, data in parse_sse(lines): order.append(name) t = data.get("type") if t == "message_start": message = data["message"] elif t == "content_block_start": b = dict(data["content_block"]); b["_json"] = ""; blocks[data["index"]] = b elif t == "content_block_delta": d, b = data["delta"], blocks[data["index"]] if d["type"] == "text_delta": b["text"] += d["text"] elif d["type"] == "input_json_delta": b["_json"] += d["partial_json"] elif d["type"] == "thinking_delta": b["thinking"] += d["thinking"] elif d["type"] == "signature_delta": b["signature"] = d["signature"] elif d["type"] == "citations_delta": b.setdefault("citations", []).append(d["citation"]) elif t == "content_block_stop": b = blocks[data["index"]] if b["type"] in ("tool_use", "server_tool_use"): b["input"] = json.loads(b.pop("_json") or "{}") # guard the parse when using eager_input_streaming else: b.pop("_json", None) elif t == "message_delta": message.update({k: v for k, v in data["delta"].items()}); message["usage"].update(data["usage"]) # usage is cumulative elif t == "ping": pass elif t == "error": raise RuntimeError(data["error"]) # e.g. overloaded_error mid-stream elif t == "message_stop": break else: pass # unknown event types must be ignored (versioning policy) message["content"] = [blocks[i] for i in sorted(blocks)] print("order :", " ".join(order)) print("stop :", message["stop_reason"], "| usage:", message["usage"]) print("blocks:", json.dumps(message["content"])[:200])