Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
25 changes: 25 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.

Expand Down
18 changes: 17 additions & 1 deletion agent_runtime/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -84,6 +91,7 @@
DEFAULT_NATS_STREAM_NAME,
DEFAULT_NATS_STREAM_SUBJECT,
DEFAULT_NATS_SUBJECT_TEMPLATE,
DEFAULT_NATS_V2_SUBJECT_TEMPLATE,
DropEvent,
NATSConsumer,
NATSConsumerConfig,
Expand All @@ -93,6 +101,8 @@
nats_app_token,
nats_token,
render_nats_subject,
render_nats_v2_subject,
v2_app_event_subject,
)

__all__ = [
Expand Down Expand Up @@ -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",
Expand All @@ -132,6 +144,7 @@
"InteractionTransport",
"Event",
"EventEnvelope",
"EventListResponse",
"MCPProviderConfig",
"NATSConsumer",
"NATSConsumerConfig",
Expand Down Expand Up @@ -170,6 +183,7 @@
"WorkspaceSkill",
"ProviderCapability",
"StoreInfo",
"StreamStateSnapshot",
"create_fastapi_command_executor_router",
"create_fastapi_event_callback_router",
"create_fastapi_skill_package_router",
Expand All @@ -183,6 +197,8 @@
"nats_token",
"parse_event_envelope",
"render_nats_subject",
"render_nats_v2_subject",
"v2_app_event_subject",
"verify_bearer_token",
]

Expand Down
35 changes: 34 additions & 1 deletion agent_runtime/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -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):
Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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",
Expand Down Expand Up @@ -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]:
Expand All @@ -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)
Expand Down
2 changes: 2 additions & 0 deletions agent_runtime/constants.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down
22 changes: 22 additions & 0 deletions agent_runtime/events.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)

Expand Down Expand Up @@ -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"))
Expand Down
13 changes: 13 additions & 0 deletions agent_runtime/nats.py
Original file line number Diff line number Diff line change
Expand Up @@ -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):
Expand Down Expand Up @@ -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 (
Expand All @@ -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:
Expand Down
Loading
Loading