Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
20 commits
Select commit Hold shift + click to select a range
03a4c8b
docs: add agent termination design spec
Siutan May 6, 2026
28583a7
docs: add agent termination implementation plan
Siutan May 6, 2026
77ccb57
feat(core): add PrefactorTerminatedError exception
Siutan May 6, 2026
214f939
feat(http): add AgentInstance.terminated_reason, AgentStatus terminat…
Siutan May 6, 2026
7058d6c
feat(http): add control signal callback to span create/finish
Siutan May 6, 2026
0da386e
feat(core): add TerminationMonitor with fast + fallback poll paths
Siutan May 6, 2026
0fd43b7
feat(core): wire TerminationMonitor into PrefactorCoreClient
Siutan May 6, 2026
6c514a1
feat(core): AgentInstanceHandle.finish() resets TerminationMonitor be…
Siutan May 6, 2026
254ad40
feat(langchain): add _throw_if_terminated() to middleware hooks
Siutan May 6, 2026
918a16f
feat(langchain): add termination demo script
Siutan May 6, 2026
a028d61
fix(core): fix sync() idempotency + rename terminated_reason → termin…
Siutan May 6, 2026
c93a757
fix(demo): use create_agent API and fix terminate body
Siutan May 6, 2026
0b322a9
feat(langchain): add environment_id to from_config and demo
Siutan May 6, 2026
098c8ac
fix: ruff format and lint
Siutan May 6, 2026
b6d4be8
fix: type errors and poll_task null checks
Siutan May 6, 2026
a577ce6
fix: propagate termination signals
Siutan May 7, 2026
0951c39
docs: align termination cleanup
Siutan May 7, 2026
cef643e
chore: remove outdated agent termination documentation
Siutan May 14, 2026
b73fc61
chore: bump package versions
Siutan May 18, 2026
d12613b
fix: harden termination cleanup
Siutan May 18, 2026
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
4 changes: 4 additions & 0 deletions packages/core/examples/agent_e2e.py
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,8 @@
environment from the token.
"""

from __future__ import annotations

import asyncio
import os

Expand Down Expand Up @@ -186,6 +188,8 @@ async def simulate_retrieval(query: str) -> list[str]:

async def main() -> None:
agent_id = os.environ.get("PREFACTOR_AGENT_ID")
if agent_id is not None:
agent_id = agent_id.strip() or None

config = PrefactorCoreConfig(
http_config=HttpClientConfig(
Expand Down
2 changes: 1 addition & 1 deletion packages/core/pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@ authors = [
]
requires-python = ">=3.11.0, <4.0.0"
dependencies = [
"prefactor-http>=0.1.3",
"prefactor-http>=0.1.4",
"pydantic>=2.0.0",
]

Expand Down
2 changes: 2 additions & 0 deletions packages/core/src/prefactor_core/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@
OperationError,
PrefactorCoreError,
PrefactorTelemetryFailureError,
PrefactorTerminatedError,
SpanNotFoundError,
)
from .managers.agent_instance import AgentInstanceHandle
Expand Down Expand Up @@ -43,6 +44,7 @@
"InstanceNotFoundError",
"SpanNotFoundError",
"PrefactorTelemetryFailureError",
"PrefactorTerminatedError",
# Models
"AgentInstance",
"Span",
Expand Down
2 changes: 1 addition & 1 deletion packages/core/src/prefactor_core/_version.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,5 +3,5 @@
from __future__ import annotations

PACKAGE_NAME = "prefactor-core"
__version__ = "0.2.5"
__version__ = "0.2.6"
PACKAGE_VERSION = __version__
89 changes: 82 additions & 7 deletions packages/core/src/prefactor_core/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,13 +7,18 @@

from __future__ import annotations

import asyncio
import logging
import time
from contextlib import asynccontextmanager
from typing import TYPE_CHECKING, Any

from prefactor_http.client import PrefactorHttpClient
from prefactor_http.exceptions import is_permanent_http_error, is_transient_http_error
from prefactor_http.exceptions import (
PrefactorApiError,
is_permanent_http_error,
is_transient_http_error,
)

from ._version import PACKAGE_NAME as CORE_PACKAGE_NAME
from ._version import PACKAGE_VERSION as CORE_PACKAGE_VERSION
Expand All @@ -26,6 +31,7 @@
)
from .managers.agent_instance import AgentInstanceManager
from .managers.span import SpanManager
from .monitoring.termination_monitor import TerminationMonitor
from .operations import Operation, OperationType
from .queue.base import Queue
from .queue.executor import TaskExecutor
Expand Down Expand Up @@ -86,6 +92,9 @@ def __init__(
self._initialized = False
self._telemetry_failure: PrefactorTelemetryFailureError | None = None
self._telemetry_failure_observed = False
self._termination_monitor: TerminationMonitor | None = None
self._sync_task: asyncio.Task | None = None
self._current_instance_id: str | None = None

def _build_http_sdk_header(self) -> str:
"""Build the effective SDK header for HTTP requests."""
Expand Down Expand Up @@ -155,6 +164,13 @@ async def initialize(self) -> None:

self._initialized = True

self._termination_monitor = TerminationMonitor(
fetch_instance=self._fetch_instance_for_poll,
)
self._sync_task = asyncio.create_task(
self._run_sync_loop(), name="prefactor-termination-sync"
)

async def close(self) -> None:
"""Close the client and cleanup resources.

Expand All @@ -164,6 +180,24 @@ async def close(self) -> None:
if not self._initialized:
return

if self._sync_task is not None:
if not self._sync_task.done():
self._sync_task.cancel()
try:
await self._sync_task
except asyncio.CancelledError:
pass
except Exception:
logger.exception(
"Termination sync loop exited with error during close()"
)
finally:
self._sync_task = None

if self._termination_monitor is not None:
self._termination_monitor.destroy()
self._termination_monitor = None

# Stop executor
if self._executor:
await self._executor.stop()
Expand Down Expand Up @@ -272,12 +306,23 @@ async def _process_operation(self, operation: Operation) -> None:
)

elif operation.type == OperationType.FINISH_AGENT_INSTANCE:
await self._http.agent_instances.finish(
agent_instance_id=operation.payload["instance_id"],
status=operation.payload.get("status", "complete"),
timestamp=operation.timestamp,
idempotency_key=operation.payload.get("idempotency_key"),
)
try:
await self._http.agent_instances.finish(
agent_instance_id=operation.payload["instance_id"],
status=operation.payload.get("status", "complete"),
timestamp=operation.timestamp,
idempotency_key=operation.payload.get("idempotency_key"),
)
except PrefactorApiError as finish_err:
if finish_err.status_code == 409:
logger.debug(
"[prefactor:http] Agent instance %s already in"
" terminal state; skipping finish.",
operation.payload["instance_id"],
)
return
raise

elif operation.type == OperationType.CREATE_SPAN:
await self._http.agent_spans.create(
agent_instance_id=operation.payload["instance_id"],
Expand All @@ -286,6 +331,7 @@ async def _process_operation(self, operation: Operation) -> None:
id=operation.payload.get("span_id"),
parent_span_id=operation.payload.get("parent_span_id"),
payload=operation.payload.get("payload"),
control_signal_callback=self._on_control_signal,
)

elif operation.type == OperationType.FINISH_SPAN:
Expand All @@ -295,6 +341,7 @@ async def _process_operation(self, operation: Operation) -> None:
result_payload=operation.payload.get("result_payload"),
timestamp=operation.timestamp,
idempotency_key=operation.payload.get("idempotency_key"),
control_signal_callback=self._on_control_signal,
)

except Exception as e:
Expand All @@ -307,6 +354,32 @@ async def _process_operation(self, operation: Operation) -> None:
)
raise

async def _fetch_instance_for_poll(self, instance_id: str):
if self._http is None:
return None
return await self._http.agent_instances.get(instance_id)

async def _run_sync_loop(self) -> None:
while True:
await asyncio.sleep(1)
if self._termination_monitor is None:
continue
try:
self._termination_monitor.sync(self._current_instance_id)
except Exception:
logger.exception("Termination sync iteration failed")

def _on_control_signal(self, reason: str | None) -> None:
if self._termination_monitor is not None:
self._termination_monitor.detect_termination(reason)

def _set_current_instance(self, instance_id: str | None) -> None:
self._current_instance_id = instance_id

def _clear_current_instance(self, instance_id: str) -> None:
if self._current_instance_id == instance_id:
self._current_instance_id = None

@property
def instance_manager(self) -> AgentInstanceManager | None:
"""Public accessor for the agent instance manager."""
Expand Down Expand Up @@ -381,6 +454,8 @@ async def create_agent_instance(
environment_id=environment_id,
)

self._set_current_instance(instance_id)

return AgentInstanceHandle(
instance_id=instance_id,
client=self,
Expand Down
18 changes: 18 additions & 0 deletions packages/core/src/prefactor_core/exceptions.py
Original file line number Diff line number Diff line change
Expand Up @@ -58,6 +58,23 @@ def __init__(
self.dropped_operations = dropped_operations


class PrefactorTerminatedError(PrefactorCoreError):
"""Raised when the agent instance has been terminated by p2.

Args:
reason: Optional reason reported by p2 for the termination.
"""

def __init__(self, reason: str | None = None) -> None:
msg = (
f"Agent instance terminated by p2: {reason}"
if reason
else "Agent instance terminated by p2"
)
super().__init__(msg)
self.reason = reason


__all__ = [
"PrefactorCoreError",
"ClientNotInitializedError",
Expand All @@ -66,4 +83,5 @@ def __init__(
"InstanceNotFoundError",
"SpanNotFoundError",
"PrefactorTelemetryFailureError",
"PrefactorTerminatedError",
]
8 changes: 7 additions & 1 deletion packages/core/src/prefactor_core/managers/agent_instance.py
Original file line number Diff line number Diff line change
Expand Up @@ -218,12 +218,18 @@ async def start(self) -> None:
async def finish(self, status: "FinishStatus" = "complete") -> None:
"""Mark the instance as finished.

This queues a finish operation for the instance.
Resets the termination monitor (fence + new event) before enqueueing
the HTTP finish so stale span responses from the dying run cannot
trigger termination on the next run.

Args:
status: Terminal status for the instance — one of ``"complete"``,
``"failed"``, or ``"cancelled"``. Defaults to ``"complete"``.
"""
monitor = getattr(self._client, "_termination_monitor", None)
if monitor is not None:
monitor.reset()
self._client._clear_current_instance(self._instance_id)
manager = self._client.instance_manager
assert manager is not None
await manager.finish_with_idempotency_key(
Expand Down
5 changes: 5 additions & 0 deletions packages/core/src/prefactor_core/monitoring/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
from __future__ import annotations

from prefactor_core.monitoring.termination_monitor import TerminationMonitor

__all__ = ["TerminationMonitor"]
Loading
Loading