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
15 changes: 15 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,21 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

Initial open-source release of FrontierAgent.

### Changed

- **Runtime engine moved to [`apodex-agent-core`](https://pypi.org/project/apodex-agent-core/)
(pinned `==0.12.0`).** The agent loop, loop contracts, tool execution,
compaction, observers, AgentBus, DAG and providers now come from `agent_core`;
`frontier_agent.*` keeps its import paths as `sys.modules` aliases or thin
adapters, so workflows, apodex and benchmarks are unchanged. Product policy is
injected through `AgentLoopHooks` / `ToolExecutionHooks` and the `configure_*`
resolvers (`core/runtime/loop/{agent_loop,tool_exec}.py`,
`components/agent_bus/bus.py`, `infra/openai_client.py`). Behaviour now
follows AgentCore where the fork had diverged, notably: compaction pins the
first user message verbatim and replaces legacy prose spill indexes, and
`Any`-typed tool parameters generate `{"type": "string"}` (`create_file`
now annotates its `rows` / `data` shorthand shapes explicitly).

### Added

- **ReAct workflow**: single stateful agent with tool use, sandboxed execution,
Expand Down
3 changes: 2 additions & 1 deletion benchmarks/public/core/kernel_adapter.py
Original file line number Diff line number Diff line change
Expand Up @@ -209,7 +209,8 @@ async def _bootstrap(self) -> None:
resource_manager = ResourceManager(llm=llm, tools=tools_map)
registry.register(ResourceManager, resource_manager)

agent_comm = AgentComm(event_store, event_bus)
# The OSS EventStore is a no-op sink, so AgentComm runs on its hot queues.
agent_comm = AgentComm(event_store, event_bus) # pyright: ignore[reportArgumentType]
registry.register(AgentComm, agent_comm)
spawn_guard = SpawnGuard(
TaskBudget(max_depth=2, max_parallel=200),
Expand Down
153 changes: 6 additions & 147 deletions frontier_agent/components/agent_bus/agent_comm.py
Original file line number Diff line number Diff line change
@@ -1,150 +1,9 @@
"""Persist typed inter-agent messages and optionally publish them live."""
# pyright: reportWildcardImportFromLibrary=false
"""Persist typed inter-agent messages and optionally publish them live (implemented by ``agent_core.components.agent_bus.agent_comm``)."""

from __future__ import annotations
import sys

import asyncio
import logging
from enum import StrEnum
from typing import Protocol, runtime_checkable
import agent_core.components.agent_bus.agent_comm as _implementation
from agent_core.components.agent_bus.agent_comm import * # noqa: F403

from frontier_agent.core.events import EventType
from frontier_agent.core.protocols import EventReader, EventSink
from frontier_agent.core.runtime.events.bus import EventBus
from frontier_agent.core.types import TaskId
from frontier_agent.models.agent_message import AgentMessage
from frontier_agent.models.event import KernelEvent

logger = logging.getLogger(__name__)


@runtime_checkable
class _AgentCommEventStore(EventSink, EventReader, Protocol):
"""AgentComm needs both append (EventSink) and cursored reads
(EventReader) — combined here because Python lacks an intersection
type. Public Protocols stay minimal in ``core.protocols``.
"""


class DeliveryMode(StrEnum):
"""Message delivery mode.

TRIGGER: persist + hot queue + EventBus broadcast.
Use when receiver must act immediately (e.g., assertion → critic).
QUEUE: persist + hot queue only (no broadcast).
Use for status updates the receiver pulls when ready.
"""

TRIGGER = "trigger"
QUEUE = "queue"


class AgentComm:
"""Sends inter-agent messages and logs them as events.

All messages are persisted to EventStore (truth source).
Hot queues are optional in-memory acceleration.
DeliveryMode controls whether EventBus broadcast fires.
"""

def __init__(
self,
event_store: _AgentCommEventStore,
event_bus: EventBus,
) -> None:
self._event_store = event_store
self._event_bus = event_bus
self._cursors: dict[tuple[str, str | None], int] = {}
self._hot_queues: dict[str, asyncio.Queue[KernelEvent]] = {}

async def send(
self,
msg: AgentMessage,
mode: DeliveryMode = DeliveryMode.QUEUE,
) -> KernelEvent:
"""Send an agent message.

Always persists to EventStore (truth source).
TRIGGER mode additionally broadcasts via EventBus for immediate wakeup.
"""
event = KernelEvent(
task_id=TaskId(msg.task_id),
event_type=EventType.AGENT_MESSAGE,
from_agent=msg.from_agent,
to_agent=msg.to_agent,
message_type=msg.message_type,
correlation_id=msg.content.get("correlation_id"),
payload={
"message_id": msg.id,
"from_agent": msg.from_agent,
"to_agent": msg.to_agent,
"message_type": msg.message_type,
"content": msg.content,
"parent_id": msg.parent_id,
"correlation_id": msg.content.get("correlation_id"),
"delivery_mode": mode.value,
},
)

# 1. Persist (always — EventStore is truth source)
persisted = await self._event_store.append(event)

# 2. Hot queue (always — non-durable acceleration)
queue = self._hot_queues.get(msg.to_agent)
if queue is not None:
try:
queue.put_nowait(persisted)
except asyncio.QueueFull:
logger.debug(
"Hot queue for %s is full; consumer uses EventStore",
msg.to_agent,
)

# 3. Broadcast (TRIGGER only — immediate wakeup signal)
if mode == DeliveryMode.TRIGGER:
await self._event_bus.publish(
event.event_type, event.payload,
)

logger.debug(
"Agent message %s: %s → %s [%s] mode=%s",
msg.id, msg.from_agent, msg.to_agent,
msg.message_type, mode.value,
)
return persisted

async def consume(
self,
agent_id: str,
*,
task_id: str | None = None,
limit: int = 50,
) -> list[KernelEvent]:
"""Consume undelivered messages for an agent using an idempotent cursor.

Cursor is per (agent_id, task_id). Restarting from 0 replays all.
"""
cursor_key = (agent_id, task_id)
after_id = self._cursors.get(cursor_key, 0)
events = await self._event_store.get_events_for_agent(
agent_id,
after_id=after_id,
limit=limit,
task_id=task_id,
)
if events:
self._cursors[cursor_key] = int(events[-1].id)
return events

def reset_cursor(
self, agent_id: str, task_id: str | None = None,
) -> None:
"""Reset cursor for an agent (e.g., after recovery)."""
self._cursors.pop((agent_id, task_id), None)

def hot_queue_for(
self, agent_id: str,
) -> asyncio.Queue[KernelEvent]:
"""Return the non-durable hot queue for an agent."""
return self._hot_queues.setdefault(
agent_id, asyncio.Queue(maxsize=256),
)
sys.modules[__name__] = _implementation
Loading
Loading