use std::sync::Arc;
use tokio::sync::broadcast;
use bamboo_agent_core::tools::ToolExecutionContext;
use bamboo_agent_core::{AgentEvent, Session};
use bamboo_engine::execution::{
create_event_forwarder, get_or_create_event_sender, reserve_session_execution,
SessionExecutionReserveOutcome,
};
use bamboo_engine::runtime::execution::agent_spawn::{
spawn_session_execution, SessionExecutionArgs,
};
use bamboo_engine::session_app::approval_replay::{
apply_permission_replay_result, find_permission_replay_target, refresh_approval_replay_posture,
repark_permission_replay, restore_permission_replay_authorization,
validate_permission_replay_authority, ApprovalReplayDecision, PermissionReplayTarget,
};
use bamboo_engine::session_app::execute::consume_pending_clarification_resume;
use bamboo_engine::session_app::resolution::resolve_resume_config_snapshot;
use bamboo_engine::session_app::respond::{
acquire_pending_response_guard, inspect_pending_response_guarded,
submit_pending_response_checked_guarded, validate_pending_response,
PERMISSION_REEXECUTE_GENERATION_METADATA_KEY, PERMISSION_REEXECUTE_METADATA_KEY,
};
use bamboo_engine::session_app::resume::{ResumeExecutionPort, ResumeSpawnRequest};
use bamboo_engine::session_app::types::RespondInput;
use bamboo_engine::{ModelRoster, RoleModel};
use super::bridge::ConnectContext;
use super::platform::{Button, OutboundMessage, Platform, PlatformResult, ReplyCtx};
use super::render::PendingAsk;
const BUTTON_LABEL_MAX_CHARS: usize = 48;
#[derive(Debug, Clone)]
pub struct ParkedAsk {
pub nonce: String,
pub session_id: String,
pub tool_call_id: String,
pub tool_name: String,
pub question: String,
pub options: Vec<String>,
pub allow_custom: bool,
}
impl ParkedAsk {
pub fn new(nonce: String, session_id: String, ask: &PendingAsk) -> Option<Self> {
Some(Self {
nonce,
session_id,
tool_call_id: ask.tool_call_id.clone()?,
tool_name: ask.tool_name.clone(),
question: ask.question.clone(),
options: ask.options.clone(),
allow_custom: ask.allow_custom,
})
}
}
pub fn new_nonce() -> String {
let raw = uuid::Uuid::new_v4().to_string();
raw.split('-').next().unwrap_or(&raw).to_string()
}
fn truncate_label(text: &str) -> String {
if text.chars().count() <= BUTTON_LABEL_MAX_CHARS {
return text.to_string();
}
let mut out: String = text.chars().take(BUTTON_LABEL_MAX_CHARS - 1).collect();
out.push('…');
out
}
fn format_ask_text(ask: &ParkedAsk) -> String {
let mut text = ask.question.clone();
if !ask.options.is_empty() {
text.push_str("\n\n");
for (index, option) in ask.options.iter().enumerate() {
text.push_str(&format!("{}. {}\n", index + 1, option));
}
}
if ask.allow_custom {
text.push_str("\n(or reply with your own answer)");
}
text
}
pub async fn render_ask(
platform: &Arc<dyn Platform>,
reply_ctx: &ReplyCtx,
ask: &ParkedAsk,
buttons_capable: bool,
) -> PlatformResult<()> {
let text = format_ask_text(ask);
let outbound = if buttons_capable && !ask.options.is_empty() {
let rows: Vec<Vec<Button>> = ask
.options
.iter()
.enumerate()
.map(|(index, option)| {
vec![Button::new(
truncate_label(option),
format!("{}:{index}", ask.nonce),
)]
})
.collect();
OutboundMessage::text(text).with_buttons(rows)
} else {
OutboundMessage::text(text)
};
platform.reply(reply_ctx, outbound).await.map(|_| ())
}
pub async fn render_read_only_ask(
platform: &Arc<dyn Platform>,
reply_ctx: &ReplyCtx,
ask: &PendingAsk,
reason: &str,
) -> PlatformResult<()> {
let mut text = ask.question.clone();
if !ask.options.is_empty() {
text.push_str("\n\n");
for (index, option) in ask.options.iter().enumerate() {
text.push_str(&format!("{}. {}\n", index + 1, option));
}
}
text.push_str("\n(Response unavailable: ");
text.push_str(reason);
text.push(')');
platform
.reply(reply_ctx, OutboundMessage::text(text))
.await
.map(|_| ())
}
const AFFIRMATIVE_KEYWORDS: &[&str] = &[
"允许", "同意", "确定", "是", "yes", "allow", "approve", "ok",
];
const NEGATIVE_KEYWORDS: &[&str] = &["拒绝", "不", "否", "no", "deny", "reject", "stay"];
fn classify_intent(text: &str) -> Option<bool> {
let lower = text.trim().to_lowercase();
if AFFIRMATIVE_KEYWORDS.iter().any(|keyword| lower == *keyword) {
return Some(true);
}
if NEGATIVE_KEYWORDS.iter().any(|keyword| lower == *keyword) {
return Some(false);
}
None
}
fn pick_option_by_intent(options: &[String], affirmative: bool) -> Option<String> {
let keywords: &[&str] = if affirmative {
AFFIRMATIVE_KEYWORDS
} else {
NEGATIVE_KEYWORDS
};
if let Some(option) = options.iter().find(|option| {
let lower = option.to_lowercase();
keywords.iter().any(|keyword| lower.contains(keyword))
}) {
return Some(option.clone());
}
if options.len() == 2 {
return Some(if affirmative {
options[0].clone()
} else {
options[1].clone()
});
}
None
}
pub fn match_text_answer(ask: &ParkedAsk, text: &str) -> Option<String> {
let trimmed = text.trim();
if trimmed.is_empty() {
return None;
}
if let Ok(index) = trimmed.parse::<usize>() {
if index >= 1 && index <= ask.options.len() {
return Some(ask.options[index - 1].clone());
}
}
if let Some(option) = ask
.options
.iter()
.find(|option| option.eq_ignore_ascii_case(trimmed))
{
return Some(option.clone());
}
if !ask.allow_custom {
if let Some(intent) = classify_intent(trimmed) {
if let Some(option) = pick_option_by_intent(&ask.options, intent) {
return Some(option);
}
}
}
if ask.allow_custom {
return Some(trimmed.to_string());
}
None
}
pub fn match_callback_data(ask: &ParkedAsk, data: &str) -> Option<String> {
let (nonce, index_str) = data.split_once(':')?;
if nonce != ask.nonce {
return None;
}
let index: usize = index_str.parse().ok()?;
ask.options.get(index).cloned()
}
pub enum RespondAndResumeOutcome {
Resumed(broadcast::Receiver<AgentEvent>),
NotResumed(String),
}
#[derive(Debug, thiserror::Error)]
pub enum ResponderError {
#[error("session not found")]
NotFound,
#[error("no pending question waiting for a response")]
NoPendingQuestion,
#[error("the pending question changed; this action has expired")]
PendingQuestionChanged,
#[error("invalid response: {0}")]
InvalidResponse(String),
#[error("{0}")]
Other(String),
}
#[async_trait::async_trait]
pub trait Responder: Send + Sync {
async fn respond_and_resume(
&self,
session_id: &str,
expected_tool_call_id: Option<&str>,
answer: String,
) -> Result<RespondAndResumeOutcome, ResponderError>;
}
fn map_respond_error(error: bamboo_engine::session_app::errors::RespondError) -> ResponderError {
use bamboo_engine::session_app::errors::RespondError;
match error {
RespondError::NotFound(_) => ResponderError::NotFound,
RespondError::NoPendingQuestion => ResponderError::NoPendingQuestion,
RespondError::PendingQuestionMismatch { .. } => ResponderError::PendingQuestionChanged,
RespondError::InvalidResponse(message) => ResponderError::InvalidResponse(message),
other => ResponderError::Other(other.to_string()),
}
}
pub struct EngineResponder {
ctx: ConnectContext,
}
impl EngineResponder {
pub fn new(ctx: ConnectContext) -> Self {
Self { ctx }
}
}
#[async_trait::async_trait]
impl Responder for EngineResponder {
async fn respond_and_resume(
&self,
session_id: &str,
expected_tool_call_id: Option<&str>,
answer: String,
) -> Result<RespondAndResumeOutcome, ResponderError> {
let response_guard = acquire_pending_response_guard(session_id).await;
let current =
inspect_pending_response_guarded(&self.ctx.session_repo, session_id, &response_guard)
.await
.map_err(map_respond_error)?
.ok_or(ResponderError::NotFound)?;
let pending = current
.pending_question
.as_ref()
.ok_or(ResponderError::NoPendingQuestion)?;
if expected_tool_call_id.is_some_and(|expected| expected != pending.tool_call_id) {
return Err(ResponderError::PendingQuestionChanged);
}
validate_pending_response(pending, &answer).map_err(ResponderError::InvalidResponse)?;
let port = ConnectResumePort {
ctx: self.ctx.clone(),
};
let handoff = bamboo_engine::session_app::resume::reserve_response_resume_handoff(
&port,
session_id,
std::time::Duration::from_secs(15),
)
.await
.map_err(|_| {
ResponderError::Other(
"the suspending run is still finalizing; answer not consumed".to_string(),
)
})?;
let receiver = handoff.subscribe();
let config_snapshot = self.ctx.config.read().await.clone();
let input = RespondInput {
session_id: session_id.to_string(),
user_response: answer,
model: None,
model_ref: None,
provider: None,
reasoning_effort: None,
};
let submission = submit_pending_response_checked_guarded(
&self.ctx.session_repo,
input,
expected_tool_call_id.map(str::to_string),
&response_guard,
)
.await;
let (session, _submitted_answer, plan_mode_transition, permission_grants) = match submission
{
Ok(submission) => submission,
Err(error) => {
handoff.abandon().await;
return Err(map_respond_error(error));
}
};
for (perm_type, resource) in &permission_grants {
if let Some(request_id) = session.metadata.get(PERMISSION_REEXECUTE_METADATA_KEY) {
self.ctx.permission_checker.grant_once(
session_id,
request_id,
*perm_type,
resource.clone(),
);
}
}
if let Some(event) = plan_mode_transition_event(session_id, plan_mode_transition.as_ref()) {
handoff.publish_event(event);
}
let resume_config = resolve_resume_config_snapshot(
&config_snapshot,
&self.ctx.provider_registry,
&session,
None,
);
let outcome = bamboo_engine::session_app::resume::resume_session_execution_with_handoff(
&port,
session_id,
session,
resume_config,
handoff,
)
.await;
drop(response_guard);
match outcome {
bamboo_engine::session_app::types::ResumeOutcome::Started { .. } => {
Ok(RespondAndResumeOutcome::Resumed(receiver))
}
bamboo_engine::session_app::types::ResumeOutcome::AlreadyRunning { .. } => Ok(
RespondAndResumeOutcome::NotResumed("this session is already running".to_string()),
),
bamboo_engine::session_app::types::ResumeOutcome::Completed => Ok(
RespondAndResumeOutcome::NotResumed("nothing left to resume".to_string()),
),
bamboo_engine::session_app::types::ResumeOutcome::NotFound => Ok(
RespondAndResumeOutcome::NotResumed("session no longer exists".to_string()),
),
}
}
}
fn plan_mode_transition_event(
session_id: &str,
transition: Option<&bamboo_engine::session_app::respond::PlanModeTransition>,
) -> Option<AgentEvent> {
use bamboo_engine::session_app::respond::PlanModeTransition;
transition.map(|transition| match transition {
PlanModeTransition::Entered {
reason,
pre_permission_mode,
entered_at,
status,
plan_file_path,
} => AgentEvent::PlanModeEntered {
session_id: session_id.to_string(),
reason: reason.clone(),
pre_permission_mode: pre_permission_mode.clone(),
entered_at: *entered_at,
status: *status,
plan_file_path: plan_file_path.clone(),
},
PlanModeTransition::Exited {
approved,
restored_mode,
plan,
} => AgentEvent::PlanModeExited {
session_id: session_id.to_string(),
approved: *approved,
restored_mode: restored_mode.clone(),
plan: plan.clone(),
},
})
}
struct ConnectResumePort {
ctx: ConnectContext,
}
#[async_trait::async_trait]
impl ResumeExecutionPort for ConnectResumePort {
async fn load_session(&self, session_id: &str) -> Option<Session> {
self.ctx.session_repo.load_merged(session_id).await
}
async fn save_and_cache_session(&self, session: &mut Session) {
self.ctx.session_repo.save_and_cache(session).await;
}
async fn reserve_session_execution(
&self,
session_id: &str,
event_sender: &broadcast::Sender<AgentEvent>,
) -> SessionExecutionReserveOutcome {
reserve_session_execution(
&self.ctx.agent,
&self.ctx.agent_runners,
&self.ctx.session_event_senders,
session_id,
event_sender,
)
.await
}
async fn get_or_create_event_sender(&self, session_id: &str) -> broadcast::Sender<AgentEvent> {
get_or_create_event_sender(&self.ctx.session_event_senders, session_id).await
}
fn dispatch_resume_execution(
&self,
request: ResumeSpawnRequest,
) -> Result<(), ResumeSpawnRequest> {
let owner = ConnectResumePort {
ctx: self.ctx.clone(),
};
tokio::spawn(async move {
ResumeExecutionPort::spawn_resume_execution(&owner, request).await;
});
Ok(())
}
async fn spawn_resume_execution(&self, request: ResumeSpawnRequest) {
let ResumeSpawnRequest {
session_id,
mut session,
mut execution_reservation,
event_sender,
config,
} = request;
if let Err(error) = execution_reservation.ensure_registered().await {
tracing::warn!(
%session_id,
run_id = %execution_reservation.run_id(),
%error,
"cannot resume connect session without exact router ownership"
);
return;
}
let model = session.model.clone();
let reasoning_effort = session.reasoning_effort;
let model_roster = ModelRoster {
model: Some(model),
provider_name: Some(config.provider_name.clone()),
provider_type: config.provider_type.clone(),
fast: RoleModel::from_parts(config.fast_model.clone(), None),
background: RoleModel::from_parts(
config.background_model.clone(),
config.background_model_provider.clone(),
),
summarization: RoleModel::from_parts(
config.summarization_model.clone(),
config.summarization_model_provider.clone(),
),
};
let (mpsc_tx, _forwarder_handle) = create_event_forwarder(
session_id.clone(),
execution_reservation.run_id().to_string(),
event_sender,
self.ctx.agent_runners.clone(),
self.ctx.account_feed_inbox.clone(),
);
let reexecute_tool_call_id = session
.metadata
.get(PERMISSION_REEXECUTE_METADATA_KEY)
.cloned();
let reexecute_request_generation = session
.metadata
.get(PERMISSION_REEXECUTE_GENERATION_METADATA_KEY)
.cloned();
let Some(reexecute_tool_call_id) = reexecute_tool_call_id else {
if reexecute_request_generation.is_some() {
tracing::error!(
%session_id,
"connect found orphaned permission replay generation marker; refusing to resume"
);
return;
}
consume_pending_clarification_resume(&mut session);
spawn_session_execution(SessionExecutionArgs {
agent: self.ctx.agent.clone(),
session_id,
session,
execution_reservation,
tools_override: Some(self.ctx.tools.clone()),
provider_override: None,
model_roster,
reasoning_effort,
reasoning_effort_source: "connect_resume".to_string(),
auxiliary_model_resolver: None,
disabled_filter_resolver: None,
disabled_tools: Some(config.disabled_tools.clone()),
disabled_skill_ids: Some(config.disabled_skill_ids.clone()),
selected_skill_ids: None,
selected_skill_mode: None,
mpsc_tx,
image_fallback: config.image_fallback.clone(),
gold_config: config.gold_config.clone(),
guardian_config: None,
guardian_spawner: None,
bash_resume_hook: None,
bash_completion_sink: None,
app_data_dir: self.ctx.app_data_dir.clone(),
run_budget: None,
runners: self.ctx.agent_runners.clone(),
sessions_cache: self.ctx.session_repo.cache().clone(),
on_complete: None,
child_completion_handler: None,
});
return;
};
let ctx = self.ctx.clone();
tokio::spawn(async move {
let mut session = session;
if let Some(replay_target) = find_pending_tool_call(
&session,
&reexecute_tool_call_id,
reexecute_request_generation.as_deref(),
) {
if reexecute_request_generation.is_none()
&& replay_target.request_generation().is_some()
{
tracing::error!(
%session_id,
tool_call_id = %reexecute_tool_call_id,
"connect typed permission replay is missing its generation marker; refusing to resume"
);
return;
}
let tool_call = replay_target.tool_call().clone();
let tool_name = tool_call.function.name.clone();
let executor = ctx.tools.clone();
let replay_owner = bamboo_domain::resolve_tool_reference_name(&tool_name, |name| {
executor.owns_exact_tool(name)
})
.unwrap_or_else(|| tool_name.clone());
let executing_supervisor = match validate_permission_replay_authority(
&session,
&replay_target,
&replay_owner,
) {
Ok(observation) => observation,
Err(error) => {
tracing::error!(%session_id, %error, "Supervisor approval replay binding failed closed");
return;
}
};
let configured_mode = ctx
.permission_checker
.permission_config()
.map(|config| config.mode())
.unwrap_or_default();
let decision = match refresh_approval_replay_posture(
ctx.session_repo.storage().as_ref(),
&mut session,
configured_mode,
&tool_name,
)
.await
{
Ok(decision) => decision,
Err(error) => {
tracing::error!(
%session_id,
tool_call_id = %reexecute_tool_call_id,
%error,
"connect approval replay posture refresh failed closed"
);
return;
}
};
session.metadata.remove(PERMISSION_REEXECUTE_METADATA_KEY);
session
.metadata
.remove(PERMISSION_REEXECUTE_GENERATION_METADATA_KEY);
let (content, success) = match decision {
ApprovalReplayDecision::BlockedByPlan(_) => (
format!(
"Plan mode blocked approved mutating tool '{tool_name}'; the stale approval was not executed"
),
false,
),
ApprovalReplayDecision::Execute(flags) => {
let Some(permission_config) =
ctx.permission_checker.permission_config()
else {
tracing::error!(
%session_id,
tool_call_id = %reexecute_tool_call_id,
"connect typed approval replay has no permission configuration; refusing to resume"
);
return;
};
if let Err(error) = restore_permission_replay_authorization(
permission_config.as_ref(),
&session,
&replay_target,
&replay_owner,
) {
tracing::error!(
%session_id,
tool_call_id = %reexecute_tool_call_id,
%error,
"connect typed approval replay authorization recovery failed closed"
);
return;
}
let is_mutating = bamboo_tools::orchestrator::classify_tool(&tool_name)
== bamboo_tools::orchestrator::ToolMutability::Mutating;
let mut emitter = bamboo_tools::ToolEmitter::new(
&tool_call.id,
&tool_name,
is_mutating,
);
emitter.set_auto_approved(true);
let _ = mpsc_tx
.send(emitter.begin().clone().into_agent_event())
.await;
let exec_result = bamboo_tools::permission::with_permission_replay_generation(
session.id.as_str(),
reexecute_tool_call_id.as_str(),
reexecute_request_generation.as_deref(),
executor.execute_exact_with_context_outcome(
&tool_call,
&replay_owner,
ToolExecutionContext {
executing_supervisor,
session_id: Some(session.id.as_str()),
root_session_id: Some(
if session.root_session_id.trim().is_empty() {
session.id.as_str()
} else {
session.root_session_id.as_str()
},
),
tool_call_id: reexecute_tool_call_id.as_str(),
event_tx: Some(&mpsc_tx),
available_tool_schemas: None,
bypass_permissions: flags.bypass_permissions,
auto_approve_permissions: flags.auto_approve_permissions,
plan_read_only: flags.plan_read_only,
can_async_resume: false,
bash_completion_sink: None,
pre_parsed_args: None,
},
),
)
.await.map(bamboo_agent_core::tools::ToolOutcome::into_tool_result);
match exec_result {
Ok(tool_result) => {
match repark_permission_replay(
&mut session,
&replay_target,
&tool_result,
&replay_owner,
) {
Ok(Some(reparked)) => {
let _ = mpsc_tx
.send(
emitter
.finish(Some(
"Awaiting additional permission approval"
.to_string(),
))
.clone()
.into_agent_event(),
)
.await;
let _ = mpsc_tx
.send(AgentEvent::ToolComplete {
tool_call_id: tool_call.id.clone(),
result: tool_result,
})
.await;
let _ = mpsc_tx
.send(AgentEvent::NeedClarification {
question: reparked.question,
options: (!reparked.options.is_empty())
.then_some(reparked.options),
tool_call_id: Some(tool_call.id.clone()),
tool_name: Some(tool_name.clone()),
allow_custom: reparked.allow_custom,
source: Some(
bamboo_agent_core::PendingQuestionSource::PauseTool,
),
})
.await;
ctx.session_repo.save_and_cache(&mut session).await;
return;
}
Ok(None) => {}
Err(error) => {
tracing::error!(
%session_id,
tool_call_id = %reexecute_tool_call_id,
%error,
"connect additional permission replay could not be re-parked; refusing to resume"
);
return;
}
}
let _ = mpsc_tx
.send(
emitter
.finish(Some(
"Re-executed after approval".to_string(),
))
.clone()
.into_agent_event(),
)
.await;
let _ = mpsc_tx
.send(AgentEvent::ToolComplete {
tool_call_id: tool_call.id.clone(),
result: tool_result.clone(),
})
.await;
(tool_result.result, tool_result.success)
}
Err(error) => {
let message =
format!("Tool re-execution after approval failed: {error}");
let _ = mpsc_tx
.send(
emitter.error(message.clone()).clone().into_agent_event(),
)
.await;
(message, false)
}
}
}
};
tracing::info!(
"[{}] connect: resolved approved tool replay '{}' ({}) -> success={}",
session_id,
tool_name,
reexecute_tool_call_id,
success
);
if !apply_tool_result(&mut session, &replay_target, content, success) {
tracing::error!(
%session_id,
tool_call_id = %reexecute_tool_call_id,
"connect approved tool replay result target changed unexpectedly; refusing to resume"
);
return;
}
ctx.session_repo.save_and_cache(&mut session).await;
} else {
tracing::error!(
%session_id,
tool_call_id = %reexecute_tool_call_id,
request_generation = ?reexecute_request_generation,
"connect permission replay target missing or generation-mismatched; markers retained and resume refused"
);
return;
}
consume_pending_clarification_resume(&mut session);
spawn_session_execution(SessionExecutionArgs {
agent: ctx.agent.clone(),
session_id,
session,
execution_reservation,
tools_override: Some(ctx.tools.clone()),
provider_override: None,
model_roster,
reasoning_effort,
reasoning_effort_source: "connect_resume".to_string(),
auxiliary_model_resolver: None,
disabled_filter_resolver: None,
disabled_tools: Some(config.disabled_tools.clone()),
disabled_skill_ids: Some(config.disabled_skill_ids.clone()),
selected_skill_ids: None,
selected_skill_mode: None,
mpsc_tx,
image_fallback: config.image_fallback.clone(),
gold_config: config.gold_config.clone(),
guardian_config: None,
guardian_spawner: None,
bash_resume_hook: None,
bash_completion_sink: None,
app_data_dir: ctx.app_data_dir.clone(),
run_budget: None,
runners: ctx.agent_runners.clone(),
sessions_cache: ctx.session_repo.cache().clone(),
on_complete: None,
child_completion_handler: None,
});
});
}
}
fn find_pending_tool_call(
session: &Session,
tool_call_id: &str,
request_generation: Option<&str>,
) -> Option<PermissionReplayTarget> {
find_permission_replay_target(session, tool_call_id, request_generation)
}
fn apply_tool_result(
session: &mut Session,
target: &PermissionReplayTarget,
content: String,
success: bool,
) -> bool {
apply_permission_replay_result(session, target, content, success)
}
#[cfg(test)]
mod tests {
use bamboo_agent_core::tools::{FunctionCall, ToolCall};
use bamboo_agent_core::Message;
use super::*;
fn supervisor_context(state: &crate::app_state::AppState) -> ConnectContext {
ConnectContext {
agent: state.agent.clone(),
tools: state.tools_for(crate::tools::ToolSurface::Root),
session_repo: state.session_repo.clone(),
agent_runners: state.agent_runners.clone(),
session_event_senders: state.session_event_senders.clone(),
account_feed_inbox: None,
app_data_dir: Some(state.app_data_dir.clone()),
config: state.config.clone(),
provider_registry: state.provider_registry.clone(),
project_store: state.project_store.clone(),
workspace_resolver: state.workspace_resolver.clone(),
project_ids_by_platform: Arc::new(Default::default()),
permission_checker: state.permission_checker.clone(),
}
}
#[tokio::test]
async fn connect_supervisor_typed_replay_restores_identity_and_rejects_corruption() {
use crate::app_state::resume_adapter::supervisor_tests::Fixture;
for corrupt in [false, true] {
let fixture = Box::pin(Fixture::pending()).await;
fixture.prepare_workspace_catalog().await;
fixture.formal_approve().await;
if corrupt {
fixture.corrupt().await;
}
let ctx = supervisor_context(&fixture.state);
let session = fixture.reload().await;
let config = resolve_resume_config_snapshot(
&*ctx.config.read().await,
&ctx.provider_registry,
&session,
None,
);
let outcome = bamboo_engine::session_app::resume::resume_session_execution(
&ConnectResumePort { ctx },
&session.id,
config,
)
.await;
assert!(matches!(
outcome,
bamboo_engine::session_app::types::ResumeOutcome::Started { .. }
));
fixture.settled(usize::from(!corrupt), corrupt).await;
}
}
#[tokio::test]
async fn connect_supervisor_text_answer_cannot_replace_a_typed_receipt() {
use crate::app_state::resume_adapter::supervisor_tests::{Fixture, CALL};
let fixture = Box::pin(Fixture::pending()).await;
let responder = EngineResponder::new(supervisor_context(&fixture.state));
let outcome = responder
.respond_and_resume(&fixture.original.session_id, Some(CALL), "Approve".into())
.await
.unwrap();
assert!(matches!(outcome, RespondAndResumeOutcome::Resumed(_)));
fixture.settled(0, true).await;
let session = fixture.reload().await;
let result = session
.messages
.iter()
.rev()
.find(|message| message.tool_call_id.as_deref() == Some(CALL))
.unwrap();
assert!(result
.metadata
.as_ref()
.unwrap()
.get("permission_decision_receipt")
.is_none());
}
fn append_permission_round(
session: &mut Session,
call_id: &str,
generation: &str,
arguments: &str,
result_id: &str,
) {
session.add_message(Message::assistant(
"",
Some(vec![ToolCall {
id: call_id.to_string(),
tool_type: "function".to_string(),
function: FunctionCall {
name: "Write".to_string(),
arguments: arguments.to_string(),
},
}]),
));
let mut result = Message::tool_result(
call_id,
serde_json::json!({
"status": "awaiting_permission_approval",
"permission_request": { "request_generation": generation }
})
.to_string(),
);
result.id = result_id.to_string();
session.add_message(result);
}
fn ask(options: Vec<&str>, allow_custom: bool) -> ParkedAsk {
ParkedAsk {
nonce: "abc12345".to_string(),
session_id: "sess-1".to_string(),
tool_call_id: "call-1".to_string(),
tool_name: "conclusion_with_options".to_string(),
question: "Approve?".to_string(),
options: options.into_iter().map(str::to_string).collect(),
allow_custom,
}
}
#[test]
fn connect_replay_targets_current_generation_when_provider_reuses_id() {
let mut session = Session::new("session", "model");
append_permission_round(
&mut session,
"reused",
"generation-old",
r#"{"content":"old"}"#,
"result-old",
);
append_permission_round(
&mut session,
"reused",
"generation-current",
r#"{"content":"current"}"#,
"result-current",
);
let target = find_pending_tool_call(&session, "reused", Some("generation-current"))
.expect("current generation target");
assert_eq!(
target.tool_call().function.arguments,
r#"{"content":"current"}"#
);
assert!(apply_tool_result(
&mut session,
&target,
"executed current".to_string(),
true,
));
assert!(session.messages[1].content.contains("generation-old"));
assert_eq!(session.messages[3].content, "executed current");
}
#[test]
fn new_nonce_is_short_and_hex_like() {
let nonce = new_nonce();
assert!(!nonce.is_empty());
assert!(nonce.len() <= 16);
assert!(nonce.chars().all(|c| c.is_ascii_hexdigit()));
}
#[test]
fn match_text_answer_numeric_index_selects_option() {
let pending = ask(vec!["Approve", "Deny"], false);
assert_eq!(
match_text_answer(&pending, "1"),
Some("Approve".to_string())
);
assert_eq!(match_text_answer(&pending, "2"), Some("Deny".to_string()));
assert_eq!(match_text_answer(&pending, "3"), None);
}
#[test]
fn match_text_answer_exact_text_is_case_insensitive() {
let pending = ask(vec!["Approve", "Deny"], false);
assert_eq!(
match_text_answer(&pending, "approve"),
Some("Approve".to_string())
);
}
#[test]
fn match_text_answer_binary_keyword_mapping() {
let pending = ask(vec!["Approve", "Deny"], false);
assert_eq!(
match_text_answer(&pending, "允许"),
Some("Approve".to_string())
);
assert_eq!(
match_text_answer(&pending, "yes"),
Some("Approve".to_string())
);
assert_eq!(
match_text_answer(&pending, "deny"),
Some("Deny".to_string())
);
assert_eq!(match_text_answer(&pending, "no"), Some("Deny".to_string()));
}
#[test]
fn match_text_answer_exact_option_named_stay_beats_negative_keyword_fallback() {
let pending = ask(vec!["Stay", "Leave"], false);
assert_eq!(
match_text_answer(&pending, "stay"),
Some("Stay".to_string())
);
let plan_pending = ask(vec!["Approve", "Stay in plan mode"], false);
assert_eq!(
match_text_answer(&plan_pending, "stay"),
Some("Stay in plan mode".to_string())
);
}
#[test]
fn match_text_answer_closed_ask_non_matching_text_falls_through() {
let pending = ask(vec!["Approve", "Deny"], false);
assert_eq!(match_text_answer(&pending, "banana"), None);
}
#[test]
fn match_text_answer_open_question_accepts_any_free_text() {
let pending = ask(vec!["OK", "Need changes"], true);
assert_eq!(
match_text_answer(&pending, "please add tests too"),
Some("please add tests too".to_string())
);
}
#[test]
fn match_text_answer_empty_text_never_matches() {
let pending = ask(vec!["OK", "Need changes"], true);
assert_eq!(match_text_answer(&pending, " "), None);
}
#[test]
fn match_callback_data_requires_the_exact_nonce() {
let pending = ask(vec!["Approve", "Deny"], false);
assert_eq!(
match_callback_data(&pending, "abc12345:0"),
Some("Approve".to_string())
);
assert_eq!(match_callback_data(&pending, "stale-nonce:0"), None);
}
#[test]
fn match_callback_data_rejects_out_of_range_index() {
let pending = ask(vec!["Approve", "Deny"], false);
assert_eq!(match_callback_data(&pending, "abc12345:9"), None);
}
#[test]
fn match_callback_data_rejects_malformed_data() {
let pending = ask(vec!["Approve", "Deny"], false);
assert_eq!(match_callback_data(&pending, "not-a-valid-shape"), None);
assert_eq!(match_callback_data(&pending, "abc12345:not-a-number"), None);
}
#[test]
fn format_ask_text_numbers_every_option() {
let pending = ask(vec!["Approve", "Deny"], false);
let text = format_ask_text(&pending);
assert!(text.contains("1. Approve"));
assert!(text.contains("2. Deny"));
assert!(!text.contains("reply with your own answer"));
}
#[test]
fn format_ask_text_open_question_mentions_free_text() {
let pending = ask(vec!["OK", "Need changes"], true);
assert!(format_ask_text(&pending).contains("reply with your own answer"));
}
}