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
80 changes: 45 additions & 35 deletions crates/cli/src/daemon/broker/server/socket.rs
Original file line number Diff line number Diff line change
Expand Up @@ -786,45 +786,55 @@ async fn disconnected(
log_control_disconnected(role, &id, reason);
state.sockets.changed.notify_waiters();
drop(_transaction);
tokio::spawn(async move {
tokio::time::sleep_until(deadline).await;
let _transaction = state.sockets.transaction(role, &id).await;
let current = lock(&state.sockets.peers)
.get(&key(role, &id))
.is_some_and(|p| p.generation == generation && p.sender.is_none());
if !current {
return;
}
if role == ComponentRole::Mcp {
let session = lock(&state.mcp_sessions).remove(&id);
if let Some(session) = session
&& let Ok(session_id) = McpSessionId::new(id.clone())
&& let Ok(action) = state.registry.release_mcp(
tokio::spawn(expire_disconnected_peer(
state, role, id, generation, deadline,
));
}

async fn expire_disconnected_peer(
state: Arc<DaemonState>,
role: ComponentRole,
id: String,
generation: String,
deadline: tokio::time::Instant,
) {
tokio::time::sleep_until(deadline).await;
let _transaction = state.sockets.transaction(role, &id).await;
let current = lock(&state.sockets.peers)
.get(&key(role, &id))
.is_some_and(|p| p.generation == generation && p.sender.is_none());
if !current {
return;
}
if role == ComponentRole::Mcp {
let session = lock(&state.mcp_sessions).remove(&id);
if let Some(session) = session
&& let Ok(session_id) = McpSessionId::new(id.clone())
&& let Ok(action) = state.registry.release_mcp(
session.fingerprint,
&session_id,
now_unix_ms().saturating_add(DRAIN_LIFETIME_MS),
)
{
if !session.released {
log_mcp_removed(
session.fingerprint,
&session_id,
now_unix_ms().saturating_add(DRAIN_LIFETIME_MS),
)
{
if !session.released {
log_mcp_removed(
session.fingerprint,
&session_id,
"control_disconnect_timeout",
&action,
);
}
handle_release_action(state.clone(), session.fingerprint, action);
}
lock(&state.pending_directives).remove(&id);
} else {
let session = lock(&state.worker_sessions).remove(&id);
if let Some(session) = session {
cleanup_worker_session(Arc::clone(&state), id.clone(), session).await;
"control_disconnect_timeout",
&action,
);
}
handle_release_action(state.clone(), session.fingerprint, action);
}
lock(&state.sockets.peers).remove(&key(role, &id));
state.sockets.changed.notify_waiters();
});
lock(&state.pending_directives).remove(&id);
} else {
let session = lock(&state.worker_sessions).remove(&id);
if let Some(session) = session {
cleanup_worker_session(Arc::clone(&state), id.clone(), session).await;
}
}
lock(&state.sockets.peers).remove(&key(role, &id));
state.sockets.changed.notify_waiters();
}

/// Complete worker cleanup after its session has been removed from admission state.
Expand Down
125 changes: 68 additions & 57 deletions crates/cli/src/daemon/mcp/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -434,64 +434,10 @@ async fn supervise_route(
BrokerDirective::LaunchWorker { .. } => {
let bootstrap = WorkerBootstrap::from_directive(directive.clone())
.expect("launch directive was matched");
let already_launched = launched
.as_ref()
.is_some_and(|(id, _, _)| id == &bootstrap.activation_id);
if !already_launched {
cleanup_pending_launch(lease, &mut launched).await?;
match launch_worker(&lease.daemon_origin, &bootstrap).await {
Ok(child) => {
launched = Some((
bootstrap.activation_id.clone(),
child,
tokio::time::Instant::now(),
));
}
Err(error) => {
report_activation_failed(
lease,
&bootstrap.activation_id,
error.failure_reason,
&error.source,
)
.await?;
directive = refresh_registration(lease).await?.directive;
continue;
}
}
}
if let Some((activation_id, child, _)) = launched.as_mut()
&& activation_id == &bootstrap.activation_id
&& let Some(status) = child.child.try_wait().map_err(CliError::Io)?
{
let error = CliError::Launch(format!(
"activated worker exited before readiness with {status}"
));
report_activation_failed(
lease,
&bootstrap.activation_id,
WorkerActivationFailureReason::WorkerExitedBeforeReady,
&error,
)
.await?;
directive = refresh_registration(lease).await?.directive;
continue;
}
if launched.as_ref().is_some_and(|(id, _, _)| id == &bootstrap.activation_id)
&& super::common::control::now_unix_ms() >= bootstrap.deadline_unix_ms
if let Some(next_directive) =
supervise_worker_launch(lease, &bootstrap, &mut launched).await?
{
let error = CliError::Launch(
"activated worker exceeded its startup deadline".into(),
);
cleanup_pending_launch(lease, &mut launched).await?;
report_activation_failed(
lease,
&bootstrap.activation_id,
WorkerActivationFailureReason::WorkerReadinessTimeout,
&error,
)
.await?;
directive = refresh_registration(lease).await?.directive;
directive = next_directive;
continue;
}
}
Expand All @@ -514,6 +460,71 @@ async fn supervise_route(
result.and(cleanup)
}

async fn supervise_worker_launch(
lease: &mut McpSession,
bootstrap: &WorkerBootstrap,
launched: &mut PendingLaunch,
) -> Result<Option<BrokerDirective>, CliError> {
let already_launched = launched
.as_ref()
.is_some_and(|(id, _, _)| id == &bootstrap.activation_id);
if !already_launched {
cleanup_pending_launch(lease, launched).await?;
match launch_worker(&lease.daemon_origin, bootstrap).await {
Ok(child) => {
*launched = Some((
bootstrap.activation_id.clone(),
child,
tokio::time::Instant::now(),
));
}
Err(error) => {
report_activation_failed(
lease,
&bootstrap.activation_id,
error.failure_reason,
&error.source,
)
.await?;
return Ok(Some(refresh_registration(lease).await?.directive));
}
}
}
if let Some((activation_id, child, _)) = launched.as_mut()
&& activation_id == &bootstrap.activation_id
&& let Some(status) = child.child.try_wait().map_err(CliError::Io)?
{
let error = CliError::Launch(format!(
"activated worker exited before readiness with {status}"
));
report_activation_failed(
lease,
&bootstrap.activation_id,
WorkerActivationFailureReason::WorkerExitedBeforeReady,
&error,
)
.await?;
return Ok(Some(refresh_registration(lease).await?.directive));
}
if launched
.as_ref()
.is_some_and(|(id, _, _)| id == &bootstrap.activation_id)
&& super::common::control::now_unix_ms() >= bootstrap.deadline_unix_ms
{
let error = CliError::Launch("activated worker exceeded its startup deadline".into());
cleanup_pending_launch(lease, launched).await?;
report_activation_failed(
lease,
&bootstrap.activation_id,
WorkerActivationFailureReason::WorkerReadinessTimeout,
&error,
)
.await?;
return Ok(Some(refresh_registration(lease).await?.directive));
}
Ok(None)
}

async fn next_directive(lease: &mut McpSession) -> Result<BrokerDirective, CliError> {
let event = lease.client.next().await;
receive_directive(lease, event).await
Expand Down
Loading
Loading