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
3 changes: 3 additions & 0 deletions codex-rs/Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

5 changes: 4 additions & 1 deletion codex-rs/agent-message-board-client/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -15,14 +15,17 @@ workspace = true
anyhow = { workspace = true }
chrono = { workspace = true, features = ["serde"] }
codex-agent-message-board-extension = { workspace = true }
codex-extension-api = { workspace = true }
codex-http-client = { workspace = true }
codex-protocol = { workspace = true }
eventsource-stream = { workspace = true }
futures = { workspace = true }
http = { workspace = true }
serde = { workspace = true, features = ["derive"] }
serde_json = { workspace = true }
tokio = { workspace = true, features = ["time"] }
tokio = { workspace = true, features = ["rt", "time"] }
tokio-util = { workspace = true }
tracing = { workspace = true }
url = { workspace = true }

[dev-dependencies]
Expand Down
53 changes: 30 additions & 23 deletions codex-rs/agent-message-board-client/src/client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -169,6 +169,11 @@ impl RemoteAgentMessageBoard {
.map_err(transport_error)?
.map_err(transport_error)?;
if !response.status().is_success() {
tracing::warn!(
%caller, %turn_id,
http_status = response.status().as_u16(),
"Remote board notification request rejected"
);
tokio::time::timeout_at(deadline, decode::<()>(response))
.await
.map_err(transport_error)??;
Expand Down Expand Up @@ -245,15 +250,7 @@ impl RemoteAgentMessageBoard {
"notification stream did not acknowledge readiness",
));
}
let stream = events
.map(|event| {
let event = event.map_err(transport_error)?;
if event.event != "notification" {
return Err(transport_error("invalid board notification"));
}
serde_json::from_str(&event.data).map_err(CodexErr::from)
})
.boxed();
let stream = events.map(|event| event.map_err(transport_error)).boxed();
Ok(BoardNotifications {
caller,
turn_id,
Expand Down Expand Up @@ -347,25 +344,35 @@ impl AgentMessageBoard for RemoteAgentMessageBoard {
pub struct BoardNotifications {
caller: ThreadId,
turn_id: String,
stream: BoxStream<'static, Result<BoardNotification>>,
stream: BoxStream<'static, Result<eventsource_stream::Event>>,
}

impl BoardNotifications {
/// Skips invalid individual notices. Transport and framing failures remain errors.
pub async fn next(&mut self) -> Result<Option<BoardNotification>> {
let Some(notice) = self.stream.next().await.transpose()? else {
return Ok(None);
};
if notice.recipient != self.caller || notice.turn_id != self.turn_id {
return Err(transport_error(
"notification belongs to another agent or turn",
));
}
if notice.post.text_preview.chars().count() > 150 {
return Err(transport_error(
"board notification preview exceeds 150 characters",
));
while let Some(event) = self.stream.next().await.transpose()? {
let error_kind = if event.event != "notification" {
"unexpected_event"
} else {
match serde_json::from_str::<BoardNotification>(&event.data) {
Ok(notice) => {
if notice.recipient != self.caller || notice.turn_id != self.turn_id {
"wrong_recipient_or_turn"
} else if notice.post.text_preview.chars().count() > 150 {
"oversized_preview"
} else {
return Ok(Some(notice));
}
}
Err(_) => "invalid_json",
}
};
tracing::warn!(
caller = %self.caller, turn_id = %self.turn_id, error_kind,
"Skipping invalid remote board notification"
);
}
Ok(Some(notice))
Ok(None)
}
}

Expand Down
2 changes: 2 additions & 0 deletions codex-rs/agent-message-board-client/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2,9 +2,11 @@
//! Runtime credentials are scoped to one session; tools never see them.
mod client;
mod notifications;
mod protocol;

pub use client::BoardNotifications;
pub use client::RemoteAgentMessageBoard;
pub use notifications::install_notifications;
pub use protocol::AccessToken;
pub use protocol::BoardNotification;
122 changes: 122 additions & 0 deletions codex-rs/agent-message-board-client/src/notifications.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,122 @@
//! Receives live board previews for an active turn using the extension's turn hooks.
//! Provisioning and durable board ownership belong to the research host.
use crate::RemoteAgentMessageBoard;
use codex_agent_message_board_extension::MessageBoardHost;
use codex_extension_api::ExtensionData;
use codex_extension_api::ExtensionFuture;
use codex_extension_api::ExtensionRegistryBuilder;
use codex_extension_api::TurnAbortInput;
use codex_extension_api::TurnErrorInput;
use codex_extension_api::TurnLifecycleContributor;
use codex_extension_api::TurnStartInput;
use codex_extension_api::TurnStartPhase;
use codex_extension_api::TurnStopInput;
use codex_protocol::ThreadId;
use codex_protocol::error::CodexErrKind;
use futures::channel::oneshot;
use std::sync::Arc;
use std::time::Duration;
use tokio_util::task::AbortOnDropHandle;

type Binding = (Arc<RemoteAgentMessageBoard>, Arc<dyn MessageBoardHost>);
type ResolveBinding = dyn Fn(&TurnStartInput<'_>) -> Option<Binding> + Send + Sync;

struct Notifications(Box<ResolveBinding>);
struct Receiver {
_task: AbortOnDropHandle<()>,
}

impl TurnLifecycleContributor for Notifications {
fn turn_start_phase(&self, _thread_store: &ExtensionData) -> TurnStartPhase {
TurnStartPhase::RegularTaskStart
}

fn on_turn_start<'a>(&'a self, input: TurnStartInput<'a>) -> ExtensionFuture<'a, ()> {
Box::pin(async move {
let Some((board, host)) = (self.0)(&input) else {
return;
};
let Ok(caller) = ThreadId::from_string(input.thread_store.level_id()) else {
return;
};
let turn_id = input.turn_id.to_owned();
let (first_attempt_tx, first_attempt_rx) = oneshot::channel();
input.turn_store.insert(Receiver {
_task: AbortOnDropHandle::new(tokio::spawn(async move {
let mut first_attempt_tx = Some(first_attempt_tx);
loop {
let connection = board.notifications(caller, turn_id.clone()).await;
if let Some(first_attempt_tx) = first_attempt_tx.take() {
let _ = first_attempt_tx.send(());
}
match connection {
Ok(mut stream) => loop {
match stream.next().await {
Ok(Some(notice)) => {
if let Err(error) = host.notify(caller, notice.post).await {
tracing::warn!(
%caller, %turn_id,
error_kind = ?CodexErrKind::from(&error),
"Remote board notification delivery failed"
);
}
}
Ok(None) => break,
Err(error) => {
tracing::warn!(
%caller, %turn_id,
error_kind = ?CodexErrKind::from(&error),
"Remote board notification stream failed"
);
break;
}
}
},
Err(error) => {
tracing::warn!(
%caller, %turn_id,
error_kind = ?CodexErrKind::from(&error),
"Remote board notification connection failed"
);
}
}
// Reopen for future previews; the service does not replay missed ones.
tokio::time::sleep(Duration::from_secs(/*secs*/ 1)).await;
}
})),
});
// A healthy service is listening before inference. A failed first attempt
// releases startup while the same turn-owned receiver keeps retrying.
let _ = first_attempt_rx.await;
})
}

fn on_turn_stop<'a>(&'a self, input: TurnStopInput<'a>) -> ExtensionFuture<'a, ()> {
Box::pin(async move {
input.turn_store.remove::<Receiver>();
})
}

fn on_turn_abort<'a>(&'a self, input: TurnAbortInput<'a>) -> ExtensionFuture<'a, ()> {
Box::pin(async move {
input.turn_store.remove::<Receiver>();
})
}

fn on_turn_error<'a>(&'a self, input: TurnErrorInput<'a>) -> ExtensionFuture<'a, ()> {
Box::pin(async move {
input.turn_store.remove::<Receiver>();
})
}
}

/// Installs live notification delivery. The host resolves the existing remote
/// client and a delivery host bound to this exact turn; idle agents are never woken.
/// The turn store owns the receiver and cancels it when the turn ends.
pub fn install_notifications<C: Sync>(
registry: &mut ExtensionRegistryBuilder<C>,
resolve: impl Fn(&TurnStartInput<'_>) -> Option<Binding> + Send + Sync + 'static,
) {
registry.turn_lifecycle_contributor(Arc::new(Notifications(Box::new(resolve))));
}
3 changes: 1 addition & 2 deletions codex-rs/agent-message-board-client/tests/remote_board.rs
Original file line number Diff line number Diff line change
Expand Up @@ -200,8 +200,7 @@ async fn remote_board_preserves_the_host_contract() -> anyhow::Result<()> {
.await;
let mut receiver = client.notifications(caller, "turn-1".into()).await?;
assert_eq!(receiver.next().await?, Some(notice));
assert!(receiver.next().await.is_err());
assert!(receiver.next().await.is_err());
assert_eq!(receiver.next().await?, None);
for body in [
format!(
"event: ready\nid: {}\ndata: {{}}\n\n",
Expand Down
34 changes: 32 additions & 2 deletions codex-rs/core/src/agent_message_board.rs
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@ use codex_agent_message_board_extension::MessageBoardHost;
use codex_agent_message_board_extension::NotificationDelivery;
use codex_agent_message_board_extension::PostPreview;
use codex_extension_api::ExtensionRegistryBuilder;
use codex_extension_api::ThreadStartInput;
use codex_features::Feature;
use codex_http_client::ClientRouteClass;
use codex_protocol::AgentPath;
Expand All @@ -39,12 +40,24 @@ pub fn install_agent_message_board(
registry: &mut ExtensionRegistryBuilder<Config>,
manager: Weak<ThreadManager>,
) {
codex_agent_message_board_client::install_notifications(registry, |input| {
let board = input.thread_store.get::<RemoteAgentMessageBoard>()?;
let host = input.thread_store.get::<LocalBoardHost>()?;
Some((
board,
Arc::new(LocalBoardHost {
expected_turn_id: Some(input.turn_id.into()),
..host.as_ref().clone()
}),
))
});
let in_memory_boards = Arc::new(InMemoryMessageBoards::default());
codex_agent_message_board_extension::install(
registry,
MULTI_AGENT_V2_NAMESPACE_DESCRIPTION,
|config: &Config| config.multi_agent_v2.tool_namespace.clone(),
move |config: &Config, tree, caller| {
move |input: &ThreadStartInput<'_, Config>, tree, caller| {
let config = input.config;
let in_memory = config.multi_agent_v2.message_board_in_memory;
// MAv2 supplies tree paths; ephemeral runtimes must not open local SQLite.
if !config.features.enabled(Feature::AgentMessageBoard)
Expand All @@ -63,6 +76,7 @@ pub fn install_agent_message_board(
manager: manager.clone(),
tree,
caller,
expected_turn_id: None,
});
Box::pin(async move {
let board: Arc<dyn AgentMessageBoard> = if let Some(remote) = remote {
Expand All @@ -79,13 +93,14 @@ pub fn install_agent_message_board(
let http = http_factory
.build_client(&remote.url, ClientRouteClass::Api)
.map_err(|err| CodexErr::Io(std::io::Error::other(err)))?;
input.thread_store.insert(host.as_ref().clone());
let board = RemoteAgentMessageBoard::new(http, &remote.url, tree, token)
.map_err(|err| CodexErr::InvalidRequest(err.to_string()))?
.with_clock(move |caller| {
let host = host.clone();
Box::pin(async move { host.current_time(caller).await })
});
Arc::new(board)
input.thread_store.get_or_init(|| board)
} else if in_memory {
Arc::new(in_memory_boards.open(tree, host).await)
} else {
Expand All @@ -97,10 +112,12 @@ pub fn install_agent_message_board(
);
}

#[derive(Clone)]
struct LocalBoardHost {
manager: Weak<ThreadManager>,
tree: SessionId,
caller: ThreadId,
expected_turn_id: Option<String>,
}

impl LocalBoardHost {
Expand Down Expand Up @@ -214,13 +231,26 @@ impl MessageBoardHost for LocalBoardHost {
notice.render(),
/*trigger_turn*/ false,
);
// Remote metadata is untrusted. Skip an oversized notice without closing
// the receiver; the post remains available through the board tools.
if self.expected_turn_id.is_some()
&& communication.content.len()
+ communication.author.as_str().len()
+ communication.recipient.as_str().len()
> 1024
{
return Err(CodexErr::InvalidRequest(
"remote board notice exceeds the context budget".into(),
));
}
Ok(
if recipient
.session
.input_queue
.deliver_mailbox_communication_to_current_turn(
&recipient.session.active_turn,
communication,
self.expected_turn_id.as_deref(),
)
.await
{
Expand Down
9 changes: 8 additions & 1 deletion codex-rs/core/src/session/input_queue.rs
Original file line number Diff line number Diff line change
Expand Up @@ -143,9 +143,16 @@ impl InputQueue {
&self,
active_turn: &Mutex<Option<ActiveTurn>>,
communication: InterAgentCommunication,
expected_turn_id: Option<&str>,
) -> bool {
let active = active_turn.lock().await;
let Some(active_turn) = active.as_ref().filter(|turn| turn.task.is_some()) else {
let Some(active_turn) = active.as_ref().filter(|turn| {
turn.task.as_ref().is_some_and(|task| {
expected_turn_id.is_none_or(|id| {
task.turn_context.sub_id == id && !task.cancellation_token.is_cancelled()
})
})
}) else {
return false;
};
let mut turn_state = active_turn.turn_state.lock().await;
Expand Down
Loading
Loading