From 093b4007f99910b32a3d8c21116c2d13740405a4 Mon Sep 17 00:00:00 2001 From: Amad Ali Date: Wed, 29 Jul 2026 19:31:39 +0000 Subject: [PATCH] feat: add Python SDK v2 event protocol support --- README.md | 25 ++++++++ agent_runtime/__init__.py | 18 +++++- agent_runtime/client.py | 35 ++++++++++- agent_runtime/constants.py | 2 + agent_runtime/events.py | 22 +++++++ agent_runtime/nats.py | 13 ++++ tests/test_client.py | 123 +++++++++++++++++++++++++++++++++++++ tests/test_integrations.py | 6 ++ 8 files changed, 242 insertions(+), 2 deletions(-) diff --git a/README.md b/README.md index 67f4f1e..73b949b 100644 --- a/README.md +++ b/README.md @@ -30,6 +30,31 @@ for event in client.iter_run_events(run.id): print(event.sequence_no, event.type) ``` +The durable v2 event contract is opt-in. Existing clients omit the protocol +header and continue to use v1 unchanged. Hosts can opt in when their Agent +Runtime app configuration is ready for v2: + +```python +from agent_runtime import AgentRuntimeClient + +client = AgentRuntimeClient( + base_url="https://agent-runtime.internal", + app_id="host_app", + service_token="service-token", + event_protocol="v2", +) + +replay = client.list_v2_events(run.id, after_sequence=last_sequence) +for event in replay.events: + print(event.sequence_no, event.segment_id, event.data) + +snapshot = client.get_v2_stream_state(run.id) +print(snapshot.through_sequence, snapshot.state) +``` + +For JetStream consumers, use `v2_app_event_subject(app_id)` with a separate +durable consumer. The existing `NATSConsumerConfig` defaults remain on v1. + Codex ChatGPT device-code auth can be driven through the SDK when a Codex run pauses for authentication. diff --git a/agent_runtime/__init__.py b/agent_runtime/__init__.py index ea5fd2b..44e5e8b 100644 --- a/agent_runtime/__init__.py +++ b/agent_runtime/__init__.py @@ -9,16 +9,23 @@ verify_bearer_token, ) from . import constants as constants -from .client import AgentRuntimeClient, AgentRuntimeError, AgentRuntimeHTTPError +from .client import ( + EVENT_PROTOCOL_HEADER, + AgentRuntimeClient, + AgentRuntimeError, + AgentRuntimeHTTPError, +) from .constants import * # noqa: F403 from .events import ( AssistantMessageEventData, CodexAuthStateEventData, Event, EventEnvelope, + EventListResponse, PlanUpdatedEventData, ReasoningMessageEventData, RunPlanStep, + StreamStateSnapshot, ToolCallEventData, UsageCheckpointEventData, parse_event_envelope, @@ -84,6 +91,7 @@ DEFAULT_NATS_STREAM_NAME, DEFAULT_NATS_STREAM_SUBJECT, DEFAULT_NATS_SUBJECT_TEMPLATE, + DEFAULT_NATS_V2_SUBJECT_TEMPLATE, DropEvent, NATSConsumer, NATSConsumerConfig, @@ -93,6 +101,8 @@ nats_app_token, nats_token, render_nats_subject, + render_nats_v2_subject, + v2_app_event_subject, ) __all__ = [ @@ -120,9 +130,11 @@ "CommandExecutionRequest", "CommandExecutionResponse", "DurableInfo", + "EVENT_PROTOCOL_HEADER", "DEFAULT_NATS_STREAM_NAME", "DEFAULT_NATS_STREAM_SUBJECT", "DEFAULT_NATS_SUBJECT_TEMPLATE", + "DEFAULT_NATS_V2_SUBJECT_TEMPLATE", "DropEvent", "EventCallbackConfig", "FinalizeWorkspaceRequest", @@ -132,6 +144,7 @@ "InteractionTransport", "Event", "EventEnvelope", + "EventListResponse", "MCPProviderConfig", "NATSConsumer", "NATSConsumerConfig", @@ -170,6 +183,7 @@ "WorkspaceSkill", "ProviderCapability", "StoreInfo", + "StreamStateSnapshot", "create_fastapi_command_executor_router", "create_fastapi_event_callback_router", "create_fastapi_skill_package_router", @@ -183,6 +197,8 @@ "nats_token", "parse_event_envelope", "render_nats_subject", + "render_nats_v2_subject", + "v2_app_event_subject", "verify_bearer_token", ] diff --git a/agent_runtime/client.py b/agent_runtime/client.py index c9bc963..808b0d4 100644 --- a/agent_runtime/client.py +++ b/agent_runtime/client.py @@ -26,7 +26,10 @@ ToolCall, ) from .constants import RESUME_INTENT_APPROVE, RESUME_INTENT_REQUEST_CHANGES -from .events import EventEnvelope, parse_event_envelope +from .events import EventEnvelope, EventListResponse, StreamStateSnapshot, parse_event_envelope + + +EVENT_PROTOCOL_HEADER = "X-Agent-Runtime-Event-Protocol" class AgentRuntimeError(RuntimeError): @@ -60,12 +63,14 @@ def __init__( app_id: str, service_token: Optional[str] = None, client: Optional[httpx.Client] = None, + event_protocol: Optional[str] = None, ) -> None: self.base_url = base_url.strip().rstrip("/") if not self.base_url: raise ValueError("agent runtime base URL is required") self.app_id = app_id.strip() self.service_token = service_token + self.event_protocol = (event_protocol or "").strip().lower() or None self.client = client or httpx.Client(timeout=30.0) def close(self) -> None: @@ -245,6 +250,29 @@ def list_run_events(self, run_id: str) -> List[EventEnvelope]: ) return [EventEnvelope(**item) for item in data] + def list_v2_events(self, run_id: str, after_sequence: int = 0) -> EventListResponse: + """Return durable ordered v2 events after a per-run sequence cursor.""" + + params: Dict[str, Any] = self._app_params() + if after_sequence > 0: + params["after_sequence"] = after_sequence + data = self._request( + "GET", + self._v2_run_path(run_id, "/events"), + params=params, + ) + return EventListResponse(**data) + + def get_v2_stream_state(self, run_id: str) -> StreamStateSnapshot: + """Return the authoritative materialized v2 stream snapshot.""" + + data = self._request( + "GET", + self._v2_run_path(run_id, "/stream-state"), + params=self._app_params(), + ) + return StreamStateSnapshot(**data) + def get_run_execution(self, run_id: str) -> RunExecutionInfo: data = self._request( "GET", @@ -416,6 +444,8 @@ def _headers(self, headers: Optional[Dict[str, str]] = None) -> Dict[str, str]: result = dict(headers or {}) if self.service_token and "Authorization" not in result: result["Authorization"] = f"Bearer {self.service_token}" + if self.event_protocol: + result[EVENT_PROTOCOL_HEADER] = self.event_protocol return result def _app_params(self) -> Dict[str, str]: @@ -427,6 +457,9 @@ def _path_id(self, value: str) -> str: def _run_path(self, run_id: str, suffix: str = "") -> str: return f"/v1/runs/{self._path_id(run_id)}{suffix}" + def _v2_run_path(self, run_id: str, suffix: str = "") -> str: + return f"/v2/runs/{self._path_id(run_id)}{suffix}" + def _dump(self, value: Any) -> Dict[str, Any]: if isinstance(value, dict): return dict(value) diff --git a/agent_runtime/constants.py b/agent_runtime/constants.py index e595038..32e168b 100644 --- a/agent_runtime/constants.py +++ b/agent_runtime/constants.py @@ -58,6 +58,8 @@ REPOSITORY_FINALIZE_PUSH_BRANCH = "push_branch" REPOSITORY_FINALIZE_OPEN_PR = "open_pr" +EVENT_SCHEMA_VERSION_V2 = "2" + EVENT_RUN_QUEUED = "run.queued" EVENT_RUN_STARTED = "run.started" EVENT_RUN_RESUMED = "run.resumed" diff --git a/agent_runtime/events.py b/agent_runtime/events.py index 0711545..3c720a9 100644 --- a/agent_runtime/events.py +++ b/agent_runtime/events.py @@ -39,6 +39,11 @@ class EventEnvelope(BaseModel): app_id: str = "" run_id: str = "" host_run_id: Optional[str] = None + schema_version: str = "" + turn_id: str = "" + segment_id: str = "" + revision: int = 0 + base_revision: int = 0 type: str = "" data: Dict[str, Any] = Field(default_factory=dict) @@ -142,6 +147,23 @@ class CodexAuthStateEventData(BaseModel): error: Optional[str] = None +class StreamStateSnapshot(BaseModel): + """Authoritative materialized state for a v2 run stream.""" + + schema_version: str + run_id: str + through_sequence: int + state: Dict[str, Any] = Field(default_factory=dict) + + +class EventListResponse(BaseModel): + """Durable ordered v2 events and the replay cursor that follows them.""" + + events: List[EventEnvelope] = Field(default_factory=list) + next_sequence_no: int = 0 + stream_state_snapshot: Optional[StreamStateSnapshot] = None + + def parse_event_envelope(payload: bytes | str | Dict[str, Any]) -> EventEnvelope: if isinstance(payload, bytes): value = json.loads(payload.decode("utf-8")) diff --git a/agent_runtime/nats.py b/agent_runtime/nats.py index f91c976..1acb8bb 100644 --- a/agent_runtime/nats.py +++ b/agent_runtime/nats.py @@ -12,6 +12,7 @@ DEFAULT_NATS_STREAM_NAME = "AGENT_RUNTIME_EVENTS" DEFAULT_NATS_STREAM_SUBJECT = "agent-runtime.events.>" DEFAULT_NATS_SUBJECT_TEMPLATE = "agent-runtime.events.{app_id}.{run_id}.{event_type}" +DEFAULT_NATS_V2_SUBJECT_TEMPLATE = "agent-runtime.events.v2.{app_id}.{run_id}.{event_type}" class RetryEvent(Exception): @@ -210,6 +211,12 @@ def app_event_subject(app_id: str) -> str: return f"agent-runtime.events.{nats_app_token(app_id)}.>" +def v2_app_event_subject(app_id: str) -> str: + """Return the isolated v2 subject family for one host app.""" + + return f"agent-runtime.events.v2.{nats_app_token(app_id)}.>" + + def render_nats_subject(template: str, event: EventEnvelope) -> str: value = template.strip() or DEFAULT_NATS_SUBJECT_TEMPLATE return ( @@ -220,6 +227,12 @@ def render_nats_subject(template: str, event: EventEnvelope) -> str: ) +def render_nats_v2_subject(event: EventEnvelope) -> str: + """Render an event on the versioned v2 subject family.""" + + return render_nats_subject(DEFAULT_NATS_V2_SUBJECT_TEMPLATE, event) + + def nats_token(value: str) -> str: value = str(value or "").strip() if not value: diff --git a/tests/test_client.py b/tests/test_client.py index f191dfd..477975c 100644 --- a/tests/test_client.py +++ b/tests/test_client.py @@ -25,10 +25,17 @@ TURN_POLICY_PAUSE_AFTER_ASSISTANT, USAGE_SEMANTIC_CUMULATIVE, AppendMessageRequest, + DEFAULT_NATS_V2_SUBJECT_TEMPLATE, + EVENT_PROTOCOL_HEADER, + EVENT_SCHEMA_VERSION_V2, EventEnvelope, + EventListResponse, EVENT_CODEX_AUTH_STATE_CHANGED, + StreamStateSnapshot, Usage, + render_nats_v2_subject, parse_event_envelope, + v2_app_event_subject, WorkspaceSkill, verify_bearer_token, ) @@ -87,6 +94,7 @@ def test_client_sends_v1_requests_with_auth(self): def handler(request): calls.append(request) self.assertEqual(request.headers["authorization"], "Bearer secret") + self.assertNotIn(EVENT_PROTOCOL_HEADER, request.headers) self.assertEqual(request.url.path, "/v1/runs") self.assertEqual(request.url.params["app_id"], "app-a") return httpx.Response(200, json=[run_payload()]) @@ -104,6 +112,20 @@ def handler(request): self.assertEqual(runs[0].id, "run-1") self.assertEqual(len(calls), 1) + def test_client_sends_opt_in_v2_protocol_header(self): + def handler(request): + self.assertEqual(request.headers[EVENT_PROTOCOL_HEADER], "v2") + return httpx.Response(202, json=run_payload()) + + client = AgentRuntimeClient( + "https://runtime.internal", + "app-a", + event_protocol=" V2 ", + client=httpx.Client(transport=httpx.MockTransport(handler)), + ) + + self.assertEqual(client.start_run({"agent_id": "agent-1"}).id, "run-1") + def test_start_run_defaults_app_id(self): def handler(request): body = json.loads(request.content) @@ -539,6 +561,107 @@ def test_event_envelope_requires_identity(self): with self.assertRaisesRegex(ValueError, "requires app_id"): parse_event_envelope({"type": "run.completed"}) + def test_v2_event_envelope_preserves_metadata_and_delta_whitespace(self): + envelope = parse_event_envelope({ + "event_id": "event-v2", + "app_id": "app-a", + "run_id": "run-1", + "schema_version": EVENT_SCHEMA_VERSION_V2, + "sequence_no": 7, + "turn_id": "turn-1", + "segment_id": "message-1", + "revision": 3, + "base_revision": 2, + "type": "assistant_message_delta", + "data": { + "message_id": "message-1", + "text": " world", + "content": " world", + }, + }) + + self.assertEqual(envelope.schema_version, "2") + self.assertEqual(envelope.sequence_no, 7) + self.assertEqual(envelope.turn_id, "turn-1") + self.assertEqual(envelope.segment_id, "message-1") + self.assertEqual(envelope.revision, 3) + self.assertEqual(envelope.base_revision, 2) + assistant, ok = envelope.assistant_message() + self.assertTrue(ok) + self.assertEqual(assistant.text, " world") + self.assertEqual(assistant.content, " world") + + def test_v2_replay_and_stream_state_methods(self): + requests = [] + snapshot = { + "schema_version": "2", + "run_id": "run/with spaces", + "through_sequence": 7, + "state": { + "events": [{ + "event_id": "event-v2", + "app_id": "app-a", + "run_id": "run/with spaces", + "schema_version": "2", + "sequence_no": 7, + "turn_id": "turn-1", + "segment_id": "message-1", + "revision": 1, + "type": "assistant_message_delta", + "data": {"message_id": "message-1", "content": " next"}, + }], + }, + } + + def handler(request): + requests.append(request) + self.assertEqual(request.headers[EVENT_PROTOCOL_HEADER], "v2") + self.assertEqual(request.url.params["app_id"], "app-a") + if request.url.path.endswith("/events"): + self.assertEqual(request.url.params["after_sequence"], "4") + return httpx.Response(200, json={ + "events": snapshot["state"]["events"], + "next_sequence_no": 7, + "stream_state_snapshot": snapshot, + }) + return httpx.Response(200, json=snapshot) + + client = AgentRuntimeClient( + "https://runtime.internal", + "app-a", + event_protocol="v2", + client=httpx.Client(transport=httpx.MockTransport(handler)), + ) + replay = client.list_v2_events("run/with spaces", after_sequence=4) + state = client.get_v2_stream_state("run/with spaces") + + self.assertIsInstance(replay, EventListResponse) + self.assertEqual(replay.next_sequence_no, 7) + self.assertEqual(replay.events[0].segment_id, "message-1") + self.assertEqual(replay.events[0].data["content"], " next") + self.assertIsInstance(replay.stream_state_snapshot, StreamStateSnapshot) + self.assertIsInstance(state, StreamStateSnapshot) + self.assertEqual(state.through_sequence, 7) + self.assertEqual(len(requests), 2) + self.assertTrue(all(b"run%2Fwith%20spaces" in request.url.raw_path for request in requests)) + + def test_v2_nats_subject_helpers_match_go_sdk_contract(self): + event = EventEnvelope( + app_id="helpin.stage", + run_id="run/1", + type="assistant_message_delta", + ) + + self.assertEqual( + DEFAULT_NATS_V2_SUBJECT_TEMPLATE, + "agent-runtime.events.v2.{app_id}.{run_id}.{event_type}", + ) + self.assertEqual(v2_app_event_subject("helpin.stage"), "agent-runtime.events.v2.helpin_stage.>") + self.assertEqual( + render_nats_v2_subject(event), + "agent-runtime.events.v2.helpin_stage.run_1.assistant_message_delta", + ) + def test_codex_auth_event_contract(self): envelope = parse_event_envelope({ "event_id": "event-auth", diff --git a/tests/test_integrations.py b/tests/test_integrations.py index fd452f5..50da625 100644 --- a/tests/test_integrations.py +++ b/tests/test_integrations.py @@ -13,6 +13,7 @@ from agent_runtime.events import EventEnvelope from agent_runtime.nats import ( DEFAULT_NATS_STREAM_NAME, + DEFAULT_NATS_V2_SUBJECT_TEMPLATE, DropEvent, NATSConsumer, NATSConsumerConfig, @@ -22,6 +23,8 @@ nats_app_token, nats_token, render_nats_subject, + render_nats_v2_subject, + v2_app_event_subject, ) pytestmark = pytest.mark.optional @@ -96,7 +99,10 @@ def test_nats_subject_helpers_and_defaults(): assert nats_token("run/1") == "run_1" assert nats_app_token("helpin.stage") == "helpin_stage" assert app_event_subject("helpin.stage") == "agent-runtime.events.helpin_stage.>" + assert v2_app_event_subject("helpin.stage") == "agent-runtime.events.v2.helpin_stage.>" assert render_nats_subject("", event) == "agent-runtime.events.helpin_stage.run_1.run.completed" + assert DEFAULT_NATS_V2_SUBJECT_TEMPLATE == "agent-runtime.events.v2.{app_id}.{run_id}.{event_type}" + assert render_nats_v2_subject(event) == "agent-runtime.events.v2.helpin_stage.run_1.run.completed" config = NATSConsumerConfig(app_id="helpin") assert config.stream == DEFAULT_NATS_STREAM_NAME