1use std::sync::atomic::{AtomicUsize, Ordering};
4use std::sync::{Arc, Mutex as StdMutex, OnceLock, Weak};
5
6use bamboo_agent_core::{Message, PendingQuestion, Session};
7use bamboo_domain::session::runtime_state::{AgentRuntimeState, PlanModeState, PlanModeStatus};
8use bamboo_domain::{
9 latest_response_occurrence, ResponseOccurrence, SessionPermissionMode,
10 CONSUMED_CLARIFICATION_IDS_KEY, CONSUMED_RESPONSE_OCCURRENCES_KEY,
11};
12use bamboo_tools::permission::{PermissionDecisionKind, PermissionDecisionReceipt, PermissionType};
13use chrono::Utc;
14use dashmap::mapref::entry::Entry;
15use dashmap::DashMap;
16use tokio::sync::{Mutex, OwnedMutexGuard};
17
18use super::errors::RespondError;
19use super::execute::mark_startup_handoff;
20use super::provider_model::{derive_model_ref, persist_legacy_model_provider, persist_model_ref};
21use super::repository::SessionAccess;
22use super::types::RespondInput;
23
24const CLARIFICATION_RESUME_PENDING_KEY: &str = "clarification_resume_pending";
25const CONCLUSION_WITH_OPTIONS_RESUME_PENDING_KEY: &str = "conclusion_with_options_resume_pending";
26
27type AppliedPendingResponse = (
28 String,
29 Option<PlanModeTransition>,
30 Vec<(PermissionType, String)>,
31);
32
33struct PendingResponseGate {
40 lock: Arc<Mutex<()>>,
41 waiters: AtomicUsize,
42}
43
44impl PendingResponseGate {
45 fn new() -> Self {
46 Self {
47 lock: Arc::new(Mutex::new(())),
48 waiters: AtomicUsize::new(0),
49 }
50 }
51}
52
53struct PendingResponseWaiter(Arc<PendingResponseGate>);
54
55impl Drop for PendingResponseWaiter {
56 fn drop(&mut self) {
57 self.0.waiters.fetch_sub(1, Ordering::SeqCst);
58 }
59}
60
61fn pending_response_locks() -> &'static DashMap<String, Weak<PendingResponseGate>> {
62 static LOCKS: OnceLock<DashMap<String, Weak<PendingResponseGate>>> = OnceLock::new();
63 LOCKS.get_or_init(DashMap::new)
64}
65
66fn pending_response_lock(session_id: &str) -> Arc<PendingResponseGate> {
67 let locks = pending_response_locks();
68 locks.retain(|_, lock| lock.strong_count() > 0);
69 match locks.entry(session_id.to_string()) {
70 Entry::Occupied(mut entry) => {
71 if let Some(lock) = entry.get().upgrade() {
72 lock
73 } else {
74 let lock = Arc::new(PendingResponseGate::new());
75 entry.insert(Arc::downgrade(&lock));
76 lock
77 }
78 }
79 Entry::Vacant(entry) => {
80 let lock = Arc::new(PendingResponseGate::new());
81 entry.insert(Arc::downgrade(&lock));
82 lock
83 }
84 }
85}
86
87pub struct PendingResponseGuard {
92 session_id: String,
93 _gate: Arc<PendingResponseGate>,
94 _guard: OwnedMutexGuard<()>,
95}
96
97impl PendingResponseGuard {
98 fn ensure_session(&self, session_id: &str) -> Result<(), RespondError> {
99 if self.session_id == session_id {
100 Ok(())
101 } else {
102 Err(RespondError::InvalidResponse(
103 "response transaction guard belongs to a different session".to_string(),
104 ))
105 }
106 }
107}
108
109pub async fn acquire_pending_response_guard(session_id: &str) -> PendingResponseGuard {
111 let gate = pending_response_lock(session_id);
112 gate.waiters.fetch_add(1, Ordering::SeqCst);
113 let waiter = PendingResponseWaiter(gate.clone());
114 let guard = gate.lock.clone().lock_owned().await;
115 drop(waiter);
116 PendingResponseGuard {
117 session_id: session_id.to_string(),
118 _gate: gate,
119 _guard: guard,
120 }
121}
122
123#[doc(hidden)]
126pub fn pending_response_waiter_count(session_id: &str) -> usize {
127 pending_response_locks()
128 .get(session_id)
129 .and_then(|entry| entry.upgrade())
130 .map(|gate| gate.waiters.load(Ordering::SeqCst))
131 .unwrap_or(0)
132}
133
134pub async fn inspect_pending_response_guarded(
138 repo: &dyn SessionAccess,
139 session_id: &str,
140 guard: &PendingResponseGuard,
141) -> Result<Option<Session>, RespondError> {
142 guard.ensure_session(session_id)?;
143 Ok(repo.inspect_for_response(session_id).await?)
144}
145
146pub const PERMISSION_REEXECUTE_METADATA_KEY: &str = "permission.reexecute_tool_call_id";
152pub const PERMISSION_REEXECUTE_GENERATION_METADATA_KEY: &str =
153 "permission.reexecute_request_generation";
154
155#[derive(Debug, Clone, Copy, PartialEq, Eq)]
156pub enum ResponseSource {
157 Human,
158 Gold,
159}
160
161#[derive(Debug, Clone, PartialEq, Eq)]
162pub enum PlanModeTransition {
163 Entered {
164 reason: Option<String>,
165 pre_permission_mode: String,
166 entered_at: chrono::DateTime<chrono::Utc>,
167 status: PlanModeStatus,
168 plan_file_path: Option<String>,
169 },
170 Exited {
171 approved: bool,
172 restored_mode: String,
173 plan: Option<String>,
174 },
175}
176
177pub async fn submit_pending_response(
182 repo: &dyn SessionAccess,
183 input: RespondInput,
184) -> Result<
185 (
186 Session,
187 String,
188 Option<PlanModeTransition>,
189 Vec<(PermissionType, String)>,
190 ),
191 RespondError,
192> {
193 submit_pending_response_checked(repo, input, None).await
194}
195
196pub async fn submit_pending_response_checked(
200 repo: &dyn SessionAccess,
201 input: RespondInput,
202 expected_tool_call_id: Option<String>,
203) -> Result<
204 (
205 Session,
206 String,
207 Option<PlanModeTransition>,
208 Vec<(PermissionType, String)>,
209 ),
210 RespondError,
211> {
212 submit_pending_response_with_source_checked(
213 repo,
214 input,
215 expected_tool_call_id,
216 ResponseSource::Human,
217 )
218 .await
219}
220
221pub async fn submit_pending_response_checked_guarded(
224 repo: &dyn SessionAccess,
225 input: RespondInput,
226 expected_tool_call_id: Option<String>,
227 guard: &PendingResponseGuard,
228) -> Result<
229 (
230 Session,
231 String,
232 Option<PlanModeTransition>,
233 Vec<(PermissionType, String)>,
234 ),
235 RespondError,
236> {
237 submit_pending_response_with_source_checked_guarded(
238 repo,
239 input,
240 expected_tool_call_id,
241 ResponseSource::Human,
242 guard,
243 )
244 .await
245}
246
247pub async fn submit_pending_permission_response_checked_guarded(
250 repo: &dyn SessionAccess,
251 input: RespondInput,
252 expected_tool_call_id: Option<String>,
253 permission_receipt: PermissionDecisionReceipt,
254 guard: &PendingResponseGuard,
255) -> Result<
256 (
257 Session,
258 String,
259 Option<PlanModeTransition>,
260 Vec<(PermissionType, String)>,
261 ),
262 RespondError,
263> {
264 submit_pending_response_with_source_checked_guarded_inner(
265 repo,
266 input,
267 expected_tool_call_id,
268 ResponseSource::Human,
269 Some(permission_receipt),
270 guard,
271 )
272 .await
273}
274
275pub async fn submit_pending_response_with_source(
276 repo: &dyn SessionAccess,
277 input: RespondInput,
278 response_source: ResponseSource,
279) -> Result<
280 (
281 Session,
282 String,
283 Option<PlanModeTransition>,
284 Vec<(PermissionType, String)>,
285 ),
286 RespondError,
287> {
288 submit_pending_response_with_source_checked(repo, input, None, response_source).await
289}
290
291pub async fn submit_pending_response_with_source_checked(
292 repo: &dyn SessionAccess,
293 input: RespondInput,
294 expected_tool_call_id: Option<String>,
295 response_source: ResponseSource,
296) -> Result<
297 (
298 Session,
299 String,
300 Option<PlanModeTransition>,
301 Vec<(PermissionType, String)>,
302 ),
303 RespondError,
304> {
305 let guard = acquire_pending_response_guard(&input.session_id).await;
306 submit_pending_response_with_source_checked_guarded(
307 repo,
308 input,
309 expected_tool_call_id,
310 response_source,
311 &guard,
312 )
313 .await
314}
315
316pub async fn submit_pending_response_with_source_checked_guarded(
320 repo: &dyn SessionAccess,
321 input: RespondInput,
322 expected_tool_call_id: Option<String>,
323 response_source: ResponseSource,
324 guard: &PendingResponseGuard,
325) -> Result<
326 (
327 Session,
328 String,
329 Option<PlanModeTransition>,
330 Vec<(PermissionType, String)>,
331 ),
332 RespondError,
333> {
334 submit_pending_response_with_source_checked_guarded_inner(
335 repo,
336 input,
337 expected_tool_call_id,
338 response_source,
339 None,
340 guard,
341 )
342 .await
343}
344
345async fn submit_pending_response_with_source_checked_guarded_inner(
346 repo: &dyn SessionAccess,
347 input: RespondInput,
348 expected_tool_call_id: Option<String>,
349 response_source: ResponseSource,
350 permission_receipt: Option<PermissionDecisionReceipt>,
351 guard: &PendingResponseGuard,
352) -> Result<
353 (
354 Session,
355 String,
356 Option<PlanModeTransition>,
357 Vec<(PermissionType, String)>,
358 ),
359 RespondError,
360> {
361 guard.ensure_session(&input.session_id)?;
362
363 let applied = Arc::new(StdMutex::new(None));
364 let mutation_outcome = applied.clone();
365 let mutation_input = input.clone();
366 let mutation_expected_tool_call_id = expected_tool_call_id.clone();
367 let mutation_permission_receipt = permission_receipt.clone();
368 let session = repo
369 .mutate_for_response(
370 &input.session_id,
371 Box::new(move |session| {
372 let outcome = apply_pending_response(
373 session,
374 &mutation_input,
375 mutation_expected_tool_call_id.as_deref(),
376 response_source,
377 mutation_permission_receipt.as_ref(),
378 )?;
379 *mutation_outcome
380 .lock()
381 .expect("response outcome lock poisoned") = Some(outcome);
382 Ok(())
383 }),
384 )
385 .await?
386 .ok_or_else(|| RespondError::NotFound(input.session_id.clone()))?;
387 let (user_response, plan_mode_transition, permission_grants) = applied
388 .lock()
389 .expect("response outcome lock poisoned")
390 .take()
391 .expect("successful response mutation records its outcome");
392
393 tracing::info!(
394 "[{}] Response processed successfully, agent loop can resume",
395 input.session_id
396 );
397
398 Ok((
399 session,
400 user_response,
401 plan_mode_transition,
402 permission_grants,
403 ))
404}
405
406fn apply_pending_response(
407 session: &mut Session,
408 input: &RespondInput,
409 expected_tool_call_id: Option<&str>,
410 response_source: ResponseSource,
411 permission_receipt: Option<&PermissionDecisionReceipt>,
412) -> Result<AppliedPendingResponse, RespondError> {
413 let pending = session
414 .pending_question
415 .take()
416 .ok_or(RespondError::NoPendingQuestion)?;
417
418 if session
419 .messages
420 .iter()
421 .rev()
422 .find(|message| message.tool_call_id.as_deref() == Some(pending.tool_call_id.as_str()))
423 .is_some_and(result_payload_has_supervisor_authority)
424 {
425 session.pending_question = Some(pending);
426 return Err(RespondError::InvalidResponse(
427 "Supervisor authority cannot originate in a tool result payload".into(),
428 ));
429 }
430
431 if let Some(expected) = expected_tool_call_id {
432 if pending.tool_call_id != expected {
433 let actual = pending.tool_call_id.clone();
434 session.pending_question = Some(pending);
435 return Err(RespondError::PendingQuestionMismatch {
436 expected: expected.to_string(),
437 actual,
438 });
439 }
440 }
441
442 if let Some(receipt) = permission_receipt {
443 if receipt.session_id != input.session_id
444 || receipt.decision.request_id != pending.tool_call_id
445 {
446 session.pending_question = Some(pending);
447 return Err(RespondError::InvalidResponse(
448 "permission receipt identity does not match the pending question".to_string(),
449 ));
450 }
451 let current_generation = session
452 .messages
453 .iter()
454 .rev()
455 .find(|message| message.tool_call_id.as_deref() == Some(pending.tool_call_id.as_str()))
456 .and_then(permission_request_generation);
457 if current_generation.as_deref() != Some(receipt.decision.request_generation.as_str()) {
458 session.pending_question = Some(pending);
459 return Err(RespondError::InvalidResponse(
460 "permission receipt generation does not match the pending operation".to_string(),
461 ));
462 }
463 }
464
465 if permission_receipt.is_none() {
469 if let Err(error_message) = validate_pending_response(&pending, &input.user_response) {
470 session.pending_question = Some(pending);
472 return Err(RespondError::InvalidResponse(error_message));
473 }
474 }
475
476 let tool_call_id = pending.tool_call_id.clone();
477 tracing::debug!(
478 "[{}] Looking for tool result message with tool_call_id: {}",
479 input.session_id,
480 tool_call_id
481 );
482
483 let reviewed_plan = extract_exit_plan_from_tool_result_message(session, &tool_call_id);
484
485 let typed_permission_approved = permission_receipt.is_some_and(|receipt| {
489 matches!(
490 receipt.decision.decision,
491 PermissionDecisionKind::AllowOnce
492 | PermissionDecisionKind::AllowSession
493 | PermissionDecisionKind::AllowWorkspace
494 | PermissionDecisionKind::AllowGlobal
495 )
496 });
497 let permission_approved = permission_receipt
498 .map(|_| typed_permission_approved)
499 .unwrap_or_else(|| is_permission_approval(&input.user_response));
500 let permission_grants = if permission_approved {
501 extract_permission_grants_from_tool_result_message(session, &tool_call_id)
502 } else {
503 Vec::new()
504 };
505 let should_reexecute = permission_receipt
506 .map(|_| typed_permission_approved)
507 .unwrap_or(!permission_grants.is_empty());
508 if should_reexecute {
509 session.metadata.insert(
513 PERMISSION_REEXECUTE_METADATA_KEY.to_string(),
514 tool_call_id.clone(),
515 );
516 if let Some(receipt) = permission_receipt {
517 session.metadata.insert(
518 PERMISSION_REEXECUTE_GENERATION_METADATA_KEY.to_string(),
519 receipt.decision.request_generation.clone(),
520 );
521 } else {
522 session
523 .metadata
524 .remove(PERMISSION_REEXECUTE_GENERATION_METADATA_KEY);
525 }
526 } else if permission_receipt.is_some() {
527 session.metadata.remove(PERMISSION_REEXECUTE_METADATA_KEY);
529 session
530 .metadata
531 .remove(PERMISSION_REEXECUTE_GENERATION_METADATA_KEY);
532 }
533
534 let found = update_or_append_tool_result_message(
536 session,
537 &tool_call_id,
538 &input.user_response,
539 response_source,
540 );
541 if let Some(receipt) = permission_receipt {
542 if !persist_permission_decision_receipt(session, &tool_call_id, receipt) {
543 return Err(RespondError::InvalidResponse(
544 "permission receipt could not be persisted for the pending operation".to_string(),
545 ));
546 }
547 }
548 if found {
549 tracing::info!(
550 "[{}] Updated existing tool result message",
551 input.session_id
552 );
553 } else {
554 tracing::warn!(
555 "[{}] Tool result message not found for tool_call_id: {}, added fallback message",
556 input.session_id,
557 tool_call_id
558 );
559 }
560
561 let plan_mode_transition =
563 apply_plan_mode_transition(session, &pending, &input.user_response, reviewed_plan);
564
565 session.clear_pending_question();
567 record_consumed_clarification(session, &tool_call_id);
568 session.metadata.remove("runtime.suspend_reason");
569 session.metadata.insert(
570 CLARIFICATION_RESUME_PENDING_KEY.to_string(),
571 "true".to_string(),
572 );
573 session.metadata.insert(
574 CONCLUSION_WITH_OPTIONS_RESUME_PENDING_KEY.to_string(),
575 "true".to_string(),
576 );
577 mark_startup_handoff(session);
578
579 let request_model_ref = derive_model_ref(
581 input.model_ref.as_ref(),
582 input.provider.as_deref(),
583 input.model.as_deref(),
584 );
585 if let Some(model_ref) = request_model_ref.as_ref() {
586 persist_model_ref(session, model_ref);
587 } else {
588 persist_legacy_model_provider(session, input.model.as_deref(), input.provider.as_deref());
589 }
590 if let Some(reasoning_effort) = input.reasoning_effort {
591 session.reasoning_effort = Some(reasoning_effort);
592 }
593
594 Ok((
595 input.user_response.clone(),
596 plan_mode_transition,
597 permission_grants,
598 ))
599}
600
601fn record_consumed_clarification(session: &mut Session, tool_call_id: &str) {
602 let mut legacy_consumed = session
603 .metadata
604 .get(CONSUMED_CLARIFICATION_IDS_KEY)
605 .and_then(|value| serde_json::from_str::<Vec<String>>(value).ok())
606 .unwrap_or_default();
607 legacy_consumed.retain(|existing| existing != tool_call_id);
608 legacy_consumed.push(tool_call_id.to_string());
609 if legacy_consumed.len() > 64 {
610 legacy_consumed.drain(..legacy_consumed.len() - 64);
611 }
612 if let Ok(serialized) = serde_json::to_string(&legacy_consumed) {
613 session
614 .metadata
615 .insert(CONSUMED_CLARIFICATION_IDS_KEY.to_string(), serialized);
616 }
617
618 let Some(occurrence) = latest_response_occurrence(session, tool_call_id) else {
619 return;
620 };
621 let mut consumed = session
622 .metadata
623 .get(CONSUMED_RESPONSE_OCCURRENCES_KEY)
624 .and_then(|value| serde_json::from_str::<Vec<ResponseOccurrence>>(value).ok())
625 .unwrap_or_default();
626 consumed.retain(|existing| existing != &occurrence);
627 consumed.push(occurrence);
628 if consumed.len() > 64 {
629 consumed.drain(..consumed.len() - 64);
630 }
631 if let Ok(serialized) = serde_json::to_string(&consumed) {
632 session
633 .metadata
634 .insert(CONSUMED_RESPONSE_OCCURRENCES_KEY.to_string(), serialized);
635 }
636}
637
638fn apply_plan_mode_transition(
640 session: &mut Session,
641 pending: &PendingQuestion,
642 user_response: &str,
643 reviewed_plan: Option<String>,
644) -> Option<PlanModeTransition> {
645 match pending.tool_name.as_str() {
646 "EnterPlanMode" if user_response.to_lowercase().contains("enter plan mode") => {
647 let pre_mode = session
648 .agent_runtime_state
649 .as_ref()
650 .map(|state| state.effective_permission_mode().as_str().to_string())
651 .unwrap_or_else(|| SessionPermissionMode::Default.as_str().to_string());
652
653 let entered_at = Utc::now();
654 let status = PlanModeStatus::Exploring;
655 let runtime_state = session
656 .agent_runtime_state
657 .get_or_insert_with(|| AgentRuntimeState::new(uuid::Uuid::new_v4().to_string()));
658 runtime_state.plan_mode = Some(PlanModeState {
659 entered_at,
660 pre_permission_mode: pre_mode.clone(),
661 plan_file_path: None,
662 status,
663 });
664 tracing::info!(
665 session_id = %session.id,
666 "Entered plan mode"
667 );
668 Some(PlanModeTransition::Entered {
669 reason: Some(pending.question.clone()),
670 pre_permission_mode: pre_mode,
671 entered_at,
672 status,
673 plan_file_path: None,
674 })
675 }
676 "ExitPlanMode" if is_exit_plan_mode_approved(user_response) => {
677 let restored_mode = session
678 .agent_runtime_state
679 .as_ref()
680 .map(|state| state.effective_permission_mode().as_str().to_string())
681 .unwrap_or_else(|| "default".to_string());
682 if let Some(ref mut runtime_state) = session.agent_runtime_state {
683 runtime_state.plan_mode = None;
688 }
689 tracing::info!(
690 session_id = %session.id,
691 "Exited plan mode"
692 );
693 Some(PlanModeTransition::Exited {
694 approved: true,
695 restored_mode,
696 plan: reviewed_plan,
697 })
698 }
699 _ => None,
700 }
701}
702
703fn is_exit_plan_mode_approved(user_response: &str) -> bool {
705 let lower = user_response.to_lowercase();
706 lower.contains("approve") && !lower.contains("stay in plan mode")
707}
708
709pub fn validate_pending_response(
712 pending: &PendingQuestion,
713 user_response: &str,
714) -> Result<(), String> {
715 if pending.allow_custom {
716 return Ok(());
717 }
718
719 let valid = pending.options.iter().any(|option| option == user_response);
720 if valid {
721 Ok(())
722 } else {
723 let options_str = pending.options.join(", ");
724 Err(format!("Response must be one of: {options_str}"))
725 }
726}
727
728pub fn update_or_append_tool_result_message(
729 session: &mut Session,
730 tool_call_id: &str,
731 user_response: &str,
732 response_source: ResponseSource,
733) -> bool {
734 for message in session.messages.iter_mut().rev() {
735 if message.tool_call_id.as_deref() == Some(tool_call_id) {
736 if result_payload_has_supervisor_authority(message) {
739 return false;
740 }
741 if let Ok(payload) = serde_json::from_str::<serde_json::Value>(&message.content) {
747 if payload.get("status").and_then(serde_json::Value::as_str)
748 == Some("awaiting_permission_approval")
749 {
750 if let Some(request) = payload.get("permission_request") {
751 insert_message_metadata(message, "permission_request", request.clone());
752 }
753 }
754 }
755 message.content = selected_message_content(user_response, response_source);
756 message.tool_success = Some(true);
757 return true;
758 }
759 }
760
761 session.add_message(bamboo_agent_core::Message::tool_result_with_status(
762 tool_call_id,
763 selected_message_content(user_response, response_source),
764 true,
765 ));
766 false
767}
768
769fn result_payload_has_supervisor_authority(message: &Message) -> bool {
770 serde_json::from_str::<serde_json::Value>(&message.content).ok().is_some_and(|payload| {
771 payload.get(bamboo_agent_core::tools::ExecutingSupervisorObservation::PERMISSION_REPLAY_METADATA_KEY).is_some()
772 })
773}
774
775fn insert_message_metadata(message: &mut Message, key: &str, value: serde_json::Value) {
776 let metadata = message
777 .metadata
778 .get_or_insert_with(|| serde_json::Value::Object(Default::default()));
779 if !metadata.is_object() {
780 let previous = std::mem::replace(metadata, serde_json::Value::Object(Default::default()));
781 metadata
782 .as_object_mut()
783 .expect("replacement metadata is an object")
784 .insert("previous_metadata".to_string(), previous);
785 }
786 metadata
787 .as_object_mut()
788 .expect("message metadata is an object")
789 .insert(key.to_string(), value);
790}
791
792fn permission_request_generation(message: &Message) -> Option<String> {
793 message
794 .metadata
795 .as_ref()
796 .and_then(|metadata| metadata.get("permission_request"))
797 .and_then(|request| request.get("request_generation"))
798 .and_then(serde_json::Value::as_str)
799 .map(ToOwned::to_owned)
800 .or_else(|| {
801 serde_json::from_str::<serde_json::Value>(&message.content)
802 .ok()?
803 .get("permission_request")?
804 .get("request_generation")?
805 .as_str()
806 .map(ToOwned::to_owned)
807 })
808}
809
810fn persist_permission_decision_receipt(
811 session: &mut Session,
812 tool_call_id: &str,
813 receipt: &PermissionDecisionReceipt,
814) -> bool {
815 let Some(message) = session
816 .messages
817 .iter_mut()
818 .rev()
819 .find(|message| message.tool_call_id.as_deref() == Some(tool_call_id))
820 else {
821 return false;
822 };
823 if permission_request_generation(message).as_deref()
824 != Some(receipt.decision.request_generation.as_str())
825 {
826 return false;
827 }
828 insert_message_metadata(
829 message,
830 "permission_decision_receipt",
831 serde_json::to_value(receipt).expect("permission receipt is serializable"),
832 );
833 true
834}
835
836fn selected_message_content(user_response: &str, response_source: ResponseSource) -> String {
837 match response_source {
838 ResponseSource::Human => format!("Selected response: {}", user_response),
839 ResponseSource::Gold => format!("Auto-selected response (gold): {}", user_response),
840 }
841}
842
843fn extract_exit_plan_from_tool_result_message(
844 session: &Session,
845 tool_call_id: &str,
846) -> Option<String> {
847 let message = session
848 .messages
849 .iter()
850 .rev()
851 .find(|message| message.tool_call_id.as_deref() == Some(tool_call_id))?;
852 let payload = serde_json::from_str::<serde_json::Value>(&message.content).ok()?;
853 payload
854 .get("plan")
855 .and_then(|value| value.as_str())
856 .map(str::trim)
857 .filter(|value| !value.is_empty())
858 .map(ToOwned::to_owned)
859}
860
861fn is_permission_approval(user_response: &str) -> bool {
866 user_response.trim().eq_ignore_ascii_case("approve")
867}
868
869fn extract_permission_grants_from_tool_result_message(
878 session: &Session,
879 tool_call_id: &str,
880) -> Vec<(PermissionType, String)> {
881 let message = match session
882 .messages
883 .iter()
884 .rev()
885 .find(|message| message.tool_call_id.as_deref() == Some(tool_call_id))
886 {
887 Some(message) => message,
888 None => return Vec::new(),
889 };
890 let payload = match serde_json::from_str::<serde_json::Value>(&message.content) {
891 Ok(payload) => payload,
892 Err(_) => return Vec::new(),
893 };
894 if payload.get("status").and_then(|value| value.as_str())
895 != Some("awaiting_permission_approval")
896 {
897 return Vec::new();
898 }
899
900 let parse_one = |value: &serde_json::Value| -> Option<(PermissionType, String)> {
901 let type_value = value
902 .get("permission_type")
903 .or_else(|| value.get("type"))?
904 .clone();
905 let perm_type: PermissionType = serde_json::from_value(type_value).ok()?;
906 let resource = value.get("resource")?.as_str()?.trim().to_string();
907 if resource.is_empty() {
908 return None;
909 }
910 Some((perm_type, resource))
911 };
912
913 if let Some(array) = payload
914 .get("permissions")
915 .and_then(|value| value.as_array())
916 {
917 array.iter().filter_map(parse_one).collect()
918 } else {
919 parse_one(&payload).into_iter().collect()
920 }
921}
922
923#[cfg(test)]
924mod tests {
925 use super::*;
926 use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
927 use std::sync::Mutex;
928 use tokio::sync::Notify;
929
930 struct TestSessionAccess {
931 session: Mutex<Session>,
932 loads: AtomicUsize,
933 saves: AtomicUsize,
934 block_first_save: AtomicBool,
935 save_started: Notify,
936 release_save: Notify,
937 }
938
939 impl TestSessionAccess {
940 fn with_pending() -> Self {
941 Self::with_pending_and_blocked_save(false)
942 }
943
944 fn with_pending_and_blocked_save(block_first_save: bool) -> Self {
945 let mut session = Session::new("sess-1", "test-model");
946 session.pending_question = Some(make_pending("ConclusionWithOptions"));
947 Self {
948 session: Mutex::new(session),
949 loads: AtomicUsize::new(0),
950 saves: AtomicUsize::new(0),
951 block_first_save: AtomicBool::new(block_first_save),
952 save_started: Notify::new(),
953 release_save: Notify::new(),
954 }
955 }
956 }
957
958 #[async_trait::async_trait]
959 impl SessionAccess for TestSessionAccess {
960 async fn load_session(
961 &self,
962 _id: &str,
963 ) -> Result<Option<Session>, super::super::errors::SessionLoadError> {
964 Ok(Some(self.session.lock().unwrap().clone()))
965 }
966
967 async fn load_or_create(
968 &self,
969 _id: &str,
970 _model: &str,
971 ) -> Result<Session, super::super::errors::SessionLoadError> {
972 Ok(self.session.lock().unwrap().clone())
973 }
974
975 async fn load_merged(
976 &self,
977 _id: &str,
978 ) -> Result<Option<Session>, super::super::errors::SessionLoadError> {
979 self.loads.fetch_add(1, Ordering::SeqCst);
980 Ok(Some(self.session.lock().unwrap().clone()))
981 }
982
983 async fn save_session(
984 &self,
985 session: &mut Session,
986 ) -> Result<(), super::super::errors::SessionSaveError> {
987 if self.block_first_save.swap(false, Ordering::SeqCst) {
988 self.save_started.notify_one();
989 self.release_save.notified().await;
990 }
991 *self.session.lock().unwrap() = session.clone();
992 self.saves.fetch_add(1, Ordering::SeqCst);
993 Ok(())
994 }
995
996 async fn save_and_cache(
997 &self,
998 session: &mut Session,
999 ) -> Result<(), super::super::errors::SessionSaveError> {
1000 self.save_session(session).await
1001 }
1002 }
1003
1004 fn respond_input() -> RespondInput {
1005 RespondInput {
1006 session_id: "sess-1".to_string(),
1007 user_response: "A".to_string(),
1008 model: None,
1009 model_ref: None,
1010 provider: None,
1011 reasoning_effort: None,
1012 }
1013 }
1014
1015 fn make_pending(tool_name: &str) -> PendingQuestion {
1016 PendingQuestion {
1017 tool_call_id: "call-1".to_string(),
1018 tool_name: tool_name.to_string(),
1019 question: "Question?".to_string(),
1020 options: vec!["A".to_string(), "B".to_string()],
1021 allow_custom: false,
1022 source: bamboo_agent_core::PendingQuestionSource::PauseTool,
1023 }
1024 }
1025
1026 #[tokio::test]
1027 async fn expected_tool_call_id_accepts_the_current_question() {
1028 let repo = TestSessionAccess::with_pending();
1029
1030 submit_pending_response_checked(&repo, respond_input(), Some("call-1".to_string()))
1031 .await
1032 .expect("matching identity should submit");
1033
1034 assert_eq!(repo.saves.load(Ordering::SeqCst), 1);
1035 assert!(repo.session.lock().unwrap().pending_question.is_none());
1036 }
1037
1038 #[tokio::test]
1039 async fn stale_tool_call_id_is_rejected_without_consuming_the_question() {
1040 let repo = TestSessionAccess::with_pending();
1041
1042 let error =
1043 submit_pending_response_checked(&repo, respond_input(), Some("stale-call".to_string()))
1044 .await
1045 .expect_err("stale identity must fail");
1046
1047 assert!(matches!(
1048 error,
1049 RespondError::PendingQuestionMismatch {
1050 ref expected,
1051 ref actual,
1052 } if expected == "stale-call" && actual == "call-1"
1053 ));
1054 assert_eq!(repo.saves.load(Ordering::SeqCst), 0);
1055 assert_eq!(
1056 repo.session
1057 .lock()
1058 .unwrap()
1059 .pending_question
1060 .as_ref()
1061 .map(|question| question.tool_call_id.as_str()),
1062 Some("call-1")
1063 );
1064 }
1065
1066 #[tokio::test]
1067 async fn omitted_tool_call_guard_remains_backwards_compatible() {
1068 let repo = TestSessionAccess::with_pending();
1069
1070 submit_pending_response(&repo, respond_input())
1071 .await
1072 .expect("legacy client should remain accepted");
1073
1074 assert_eq!(repo.saves.load(Ordering::SeqCst), 1);
1075 }
1076
1077 #[tokio::test]
1078 async fn concurrent_responses_consume_a_pending_question_once() {
1079 let repo = Arc::new(TestSessionAccess::with_pending_and_blocked_save(true));
1080
1081 let first_repo = repo.clone();
1082 let first = tokio::spawn(async move {
1083 submit_pending_response_checked(
1084 first_repo.as_ref(),
1085 respond_input(),
1086 Some("call-1".to_string()),
1087 )
1088 .await
1089 });
1090 repo.save_started.notified().await;
1091
1092 let second_entered = Arc::new(Notify::new());
1093 let second_repo = repo.clone();
1094 let second_entered_task = second_entered.clone();
1095 let second = tokio::spawn(async move {
1096 second_entered_task.notify_one();
1097 submit_pending_response_checked(
1098 second_repo.as_ref(),
1099 respond_input(),
1100 Some("call-1".to_string()),
1101 )
1102 .await
1103 });
1104 second_entered.notified().await;
1105 tokio::task::yield_now().await;
1106
1107 assert_eq!(
1108 repo.loads.load(Ordering::SeqCst),
1109 1,
1110 "the second responder must wait before loading the pending question"
1111 );
1112 repo.release_save.notify_one();
1113
1114 first.await.unwrap().expect("first response should win");
1115 let error = second
1116 .await
1117 .unwrap()
1118 .expect_err("second response must observe the consumed question");
1119 assert!(matches!(error, RespondError::NoPendingQuestion));
1120 assert_eq!(repo.loads.load(Ordering::SeqCst), 2);
1121 assert_eq!(repo.saves.load(Ordering::SeqCst), 1);
1122 }
1123
1124 #[tokio::test]
1125 async fn cancelled_response_waiter_does_not_leak_diagnostics_or_the_gate() {
1126 let session_id = "cancelled-response-waiter";
1127 let owner = acquire_pending_response_guard(session_id).await;
1128 let waiter = tokio::spawn(async move { acquire_pending_response_guard(session_id).await });
1129 tokio::time::timeout(std::time::Duration::from_secs(1), async {
1130 while pending_response_waiter_count(session_id) != 1 {
1131 tokio::task::yield_now().await;
1132 }
1133 })
1134 .await
1135 .expect("waiter should reach the gate");
1136
1137 waiter.abort();
1138 let _ = waiter.await;
1139 assert_eq!(pending_response_waiter_count(session_id), 0);
1140 drop(owner);
1141
1142 tokio::time::timeout(
1143 std::time::Duration::from_secs(1),
1144 acquire_pending_response_guard(session_id),
1145 )
1146 .await
1147 .expect("cancelled waiter must not retain the gate");
1148 }
1149
1150 #[test]
1151 fn enter_plan_mode_activates_plan_mode_state() {
1152 let mut session = Session::new("sess-1", "test-model");
1153 let pending = make_pending("EnterPlanMode");
1154
1155 apply_plan_mode_transition(&mut session, &pending, "Enter plan mode", None);
1156
1157 assert!(session.agent_runtime_state.is_some());
1158 let state = session.agent_runtime_state.unwrap();
1159 assert!(state.plan_mode.is_some());
1160 let plan = state.plan_mode.unwrap();
1161 assert_eq!(plan.status, PlanModeStatus::Exploring);
1162 assert_eq!(plan.pre_permission_mode, "default");
1163 }
1164
1165 #[test]
1166 fn enter_plan_mode_does_nothing_when_not_approved() {
1167 let mut session = Session::new("sess-1", "test-model");
1168 let pending = make_pending("EnterPlanMode");
1169
1170 apply_plan_mode_transition(&mut session, &pending, "Stay in normal mode", None);
1171
1172 assert!(session.agent_runtime_state.is_none());
1173 }
1174
1175 #[test]
1176 fn plan_mode_preserves_and_restores_typed_auto_request() {
1177 let mut session = Session::new("sess-auto-plan", "test-model");
1178 session
1179 .agent_runtime_state
1180 .get_or_insert_with(|| AgentRuntimeState::new("run-1"))
1181 .set_permission_mode(SessionPermissionMode::Auto);
1182 let enter = make_pending("EnterPlanMode");
1183 let transition =
1184 apply_plan_mode_transition(&mut session, &enter, "Enter plan mode", None).unwrap();
1185 assert!(matches!(
1186 transition,
1187 PlanModeTransition::Entered {
1188 ref pre_permission_mode,
1189 ..
1190 } if pre_permission_mode == "auto"
1191 ));
1192 assert_eq!(
1193 session
1194 .agent_runtime_state
1195 .as_ref()
1196 .unwrap()
1197 .plan_mode
1198 .as_ref()
1199 .unwrap()
1200 .pre_permission_mode,
1201 "auto"
1202 );
1203
1204 let exit = make_pending("ExitPlanMode");
1205 let transition = apply_plan_mode_transition(
1206 &mut session,
1207 &exit,
1208 "Approve (Auto mode)",
1209 Some("Reviewed plan".to_string()),
1210 )
1211 .unwrap();
1212 assert!(matches!(
1213 transition,
1214 PlanModeTransition::Exited {
1215 ref restored_mode,
1216 ..
1217 } if restored_mode == "auto"
1218 ));
1219 assert_eq!(
1220 session
1221 .agent_runtime_state
1222 .as_ref()
1223 .unwrap()
1224 .effective_permission_mode(),
1225 SessionPermissionMode::Auto
1226 );
1227 }
1228
1229 #[test]
1230 fn exit_plan_mode_does_not_restore_over_a_newer_typed_mode() {
1231 let mut session = Session::new("sess-plan-patch", "test-model");
1232 let state = session
1233 .agent_runtime_state
1234 .get_or_insert_with(|| AgentRuntimeState::new("run-1"));
1235 state.set_permission_mode(SessionPermissionMode::Auto);
1236 state.plan_mode = Some(PlanModeState {
1237 entered_at: Utc::now(),
1238 pre_permission_mode: "auto".to_string(),
1239 plan_file_path: None,
1240 status: PlanModeStatus::AwaitingApproval,
1241 });
1242 state.set_permission_mode(SessionPermissionMode::Bypass);
1244
1245 let transition = apply_plan_mode_transition(
1246 &mut session,
1247 &make_pending("ExitPlanMode"),
1248 "Approve (Default mode)",
1249 None,
1250 )
1251 .unwrap();
1252
1253 assert!(matches!(
1254 transition,
1255 PlanModeTransition::Exited {
1256 ref restored_mode,
1257 ..
1258 } if restored_mode == "bypass"
1259 ));
1260 let state = session.agent_runtime_state.unwrap();
1261 assert!(state.plan_mode.is_none());
1262 assert_eq!(
1263 state.effective_permission_mode(),
1264 SessionPermissionMode::Bypass
1265 );
1266 }
1267
1268 #[test]
1269 fn exit_plan_mode_clears_plan_mode_state() {
1270 let mut session = Session::new("sess-1", "test-model");
1271 session.agent_runtime_state = Some(AgentRuntimeState::new("run-1"));
1272 session.agent_runtime_state.as_mut().unwrap().plan_mode = Some(PlanModeState {
1273 entered_at: Utc::now(),
1274 pre_permission_mode: "default".to_string(),
1275 plan_file_path: None,
1276 status: PlanModeStatus::AwaitingApproval,
1277 });
1278 let pending = make_pending("ExitPlanMode");
1279
1280 apply_plan_mode_transition(
1281 &mut session,
1282 &pending,
1283 "Approve (Default mode)",
1284 Some("Reviewed plan".to_string()),
1285 );
1286
1287 assert!(session.agent_runtime_state.unwrap().plan_mode.is_none());
1288 }
1289
1290 #[test]
1291 fn exit_plan_mode_keeps_plan_mode_when_not_approved() {
1292 let mut session = Session::new("sess-1", "test-model");
1293 session.agent_runtime_state = Some(AgentRuntimeState::new("run-1"));
1294 session.agent_runtime_state.as_mut().unwrap().plan_mode = Some(PlanModeState {
1295 entered_at: Utc::now(),
1296 pre_permission_mode: "default".to_string(),
1297 plan_file_path: None,
1298 status: PlanModeStatus::AwaitingApproval,
1299 });
1300 let pending = make_pending("ExitPlanMode");
1301
1302 apply_plan_mode_transition(&mut session, &pending, "Stay in plan mode", None);
1303
1304 assert!(session.agent_runtime_state.unwrap().plan_mode.is_some());
1305 }
1306
1307 #[test]
1308 fn exit_plan_mode_ignores_other_tools() {
1309 let mut session = Session::new("sess-1", "test-model");
1310 let pending = make_pending("ConclusionWithOptions");
1311
1312 apply_plan_mode_transition(&mut session, &pending, "Approve", None);
1313
1314 assert!(session.agent_runtime_state.is_none());
1315 }
1316
1317 #[test]
1318 fn is_exit_plan_mode_approved_detects_approval() {
1319 assert!(is_exit_plan_mode_approved("Approve (Default mode)"));
1320 assert!(is_exit_plan_mode_approved("Approve (Accept edits mode)"));
1321 assert!(!is_exit_plan_mode_approved("Stay in plan mode"));
1322 assert!(!is_exit_plan_mode_approved("Edit plan first"));
1323 }
1324
1325 #[test]
1326 fn extract_exit_plan_from_tool_result_message_reads_plan_payload() {
1327 let mut session = Session::new("sess-1", "test-model");
1328 let mut tool_message = bamboo_agent_core::Message::tool_result(
1329 "call-1",
1330 serde_json::json!({
1331 "plan": "# Plan\n\n1. Step"
1332 })
1333 .to_string(),
1334 );
1335 tool_message.tool_success = Some(true);
1336 session.add_message(tool_message);
1337
1338 let plan = extract_exit_plan_from_tool_result_message(&session, "call-1");
1339 assert_eq!(plan.as_deref(), Some("# Plan\n\n1. Step"));
1340 }
1341
1342 #[test]
1343 fn selected_permission_preserves_typed_request_and_receipt_in_non_visible_metadata() {
1344 let mut session = Session::new("sess-1", "test-model");
1345 session.add_message(bamboo_agent_core::Message::tool_result(
1346 "permission-1",
1347 serde_json::json!({
1348 "status": "awaiting_permission_approval",
1349 "permission_request": {
1350 "request_id": "permission-1",
1351 "request_generation": "generation-1",
1352 "session_id": "sess-1",
1353 "allowed_decisions": ["allow_once", "deny_once"]
1354 }
1355 })
1356 .to_string(),
1357 ));
1358 session
1359 .messages
1360 .last_mut()
1361 .expect("permission tool result")
1362 .metadata = Some(serde_json::json!("legacy-metadata"));
1363
1364 assert!(update_or_append_tool_result_message(
1365 &mut session,
1366 "permission-1",
1367 "Approve",
1368 ResponseSource::Human,
1369 ));
1370 let receipt = PermissionDecisionReceipt {
1371 session_id: "sess-1".to_string(),
1372 decision: bamboo_tools::permission::PermissionDecision {
1373 request_id: "permission-1".to_string(),
1374 request_generation: "generation-1".to_string(),
1375 decision: bamboo_tools::permission::PermissionDecisionKind::AllowOnce,
1376 matcher_id: None,
1377 expected_policy_revision: Some(4),
1378 confirm_global: false,
1379 },
1380 decided_at: Utc::now(),
1381 };
1382 assert!(persist_permission_decision_receipt(
1383 &mut session,
1384 "permission-1",
1385 &receipt
1386 ));
1387
1388 let message = session
1389 .messages
1390 .iter()
1391 .find(|message| message.tool_call_id.as_deref() == Some("permission-1"))
1392 .expect("permission tool result");
1393 assert_eq!(message.content, "Selected response: Approve");
1394 assert_eq!(
1395 message
1396 .metadata
1397 .as_ref()
1398 .and_then(|metadata| metadata.get("permission_request"))
1399 .and_then(|request| request.get("request_id"))
1400 .and_then(serde_json::Value::as_str),
1401 Some("permission-1")
1402 );
1403 assert_eq!(
1404 message
1405 .metadata
1406 .as_ref()
1407 .and_then(|metadata| metadata.get("previous_metadata")),
1408 Some(&serde_json::json!("legacy-metadata"))
1409 );
1410 assert_eq!(
1411 message
1412 .metadata
1413 .as_ref()
1414 .and_then(|metadata| metadata.get("permission_decision_receipt"))
1415 .cloned()
1416 .and_then(|receipt| {
1417 serde_json::from_value::<PermissionDecisionReceipt>(receipt).ok()
1418 }),
1419 Some(receipt)
1420 );
1421 }
1422
1423 #[test]
1424 fn reused_tool_call_id_requires_current_permission_generation() {
1425 let mut session = Session::new("sess-1", "test-model");
1426 for generation in ["generation-old", "generation-current"] {
1427 session.add_message(bamboo_agent_core::Message::tool_result(
1428 "permission-reused",
1429 serde_json::json!({
1430 "status": "awaiting_permission_approval",
1431 "permission_type": "execute_command",
1432 "resource": format!("resource-{generation}"),
1433 "permission_request": {
1434 "request_id": "permission-reused",
1435 "request_generation": generation,
1436 "session_id": "sess-1",
1437 "allowed_decisions": ["allow_once", "deny_once"]
1438 }
1439 })
1440 .to_string(),
1441 ));
1442 }
1443 session.set_pending_question_with_source(
1444 "permission-reused".to_string(),
1445 "Bash".to_string(),
1446 "Allow the current operation?".to_string(),
1447 vec!["Approve".to_string(), "Deny".to_string()],
1448 false,
1449 bamboo_agent_core::PendingQuestionSource::PauseTool,
1450 );
1451 let input = RespondInput {
1452 session_id: "sess-1".to_string(),
1453 user_response: "Approve".to_string(),
1454 model: None,
1455 model_ref: None,
1456 provider: None,
1457 reasoning_effort: None,
1458 };
1459 let receipt = |generation: &str| PermissionDecisionReceipt {
1460 session_id: "sess-1".to_string(),
1461 decision: bamboo_tools::permission::PermissionDecision {
1462 request_id: "permission-reused".to_string(),
1463 request_generation: generation.to_string(),
1464 decision: bamboo_tools::permission::PermissionDecisionKind::AllowOnce,
1465 matcher_id: None,
1466 expected_policy_revision: None,
1467 confirm_global: false,
1468 },
1469 decided_at: Utc::now(),
1470 };
1471
1472 let stale = receipt("generation-old");
1473 assert!(matches!(
1474 apply_pending_response(
1475 &mut session,
1476 &input,
1477 Some("permission-reused"),
1478 ResponseSource::Human,
1479 Some(&stale),
1480 ),
1481 Err(RespondError::InvalidResponse(message))
1482 if message.contains("generation")
1483 ));
1484 assert!(session.pending_question.is_some());
1485
1486 let current = receipt("generation-current");
1487 apply_pending_response(
1488 &mut session,
1489 &input,
1490 Some("permission-reused"),
1491 ResponseSource::Human,
1492 Some(¤t),
1493 )
1494 .expect("current generation resolves the parked operation");
1495 assert!(session.pending_question.is_none());
1496 assert_eq!(
1497 session
1498 .metadata
1499 .get(PERMISSION_REEXECUTE_GENERATION_METADATA_KEY)
1500 .map(String::as_str),
1501 Some("generation-current")
1502 );
1503 assert!(session.messages[0].content.contains("generation-old"));
1504 assert_eq!(session.messages[1].content, "Selected response: Approve");
1505 assert_eq!(
1506 session.messages[1]
1507 .metadata
1508 .as_ref()
1509 .and_then(|metadata| metadata.get("permission_decision_receipt"))
1510 .and_then(|receipt| receipt.get("decision"))
1511 .and_then(|decision| decision.get("request_generation"))
1512 .and_then(serde_json::Value::as_str),
1513 Some("generation-current")
1514 );
1515 }
1516
1517 #[test]
1518 fn typed_permission_receipt_controls_replay_independent_of_display_options() {
1519 let session = || {
1520 let mut session = Session::new("sess-typed", "test-model");
1521 session.add_message(bamboo_agent_core::Message::tool_result(
1522 "permission-localized",
1523 serde_json::json!({
1524 "status": "awaiting_permission_approval",
1525 "permission_type": "execute_command",
1526 "resource": "cargo test --workspace",
1527 "permission_request": {
1528 "request_id": "permission-localized",
1529 "request_generation": "generation-localized",
1530 "session_id": "sess-typed",
1531 "allowed_decisions": ["allow_once", "deny_once"]
1532 }
1533 })
1534 .to_string(),
1535 ));
1536 session.set_pending_question_with_source(
1537 "permission-localized".to_string(),
1538 "Bash".to_string(),
1539 "允许执行吗?".to_string(),
1540 vec!["允许".to_string(), "拒绝".to_string()],
1541 false,
1542 bamboo_agent_core::PendingQuestionSource::PauseTool,
1543 );
1544 session
1545 };
1546 let receipt = |decision| PermissionDecisionReceipt {
1547 session_id: "sess-typed".to_string(),
1548 decision: bamboo_tools::permission::PermissionDecision {
1549 request_id: "permission-localized".to_string(),
1550 request_generation: "generation-localized".to_string(),
1551 decision,
1552 matcher_id: None,
1553 expected_policy_revision: None,
1554 confirm_global: false,
1555 },
1556 decided_at: Utc::now(),
1557 };
1558
1559 let mut allowed = session();
1560 let allow_input = RespondInput {
1561 session_id: "sess-typed".to_string(),
1562 user_response: "已由结构化决定允许".to_string(),
1563 model: None,
1564 model_ref: None,
1565 provider: None,
1566 reasoning_effort: None,
1567 };
1568 apply_pending_response(
1569 &mut allowed,
1570 &allow_input,
1571 Some("permission-localized"),
1572 ResponseSource::Human,
1573 Some(&receipt(PermissionDecisionKind::AllowOnce)),
1574 )
1575 .expect("typed allow must not depend on localized display options");
1576 assert_eq!(
1577 allowed
1578 .metadata
1579 .get(PERMISSION_REEXECUTE_METADATA_KEY)
1580 .map(String::as_str),
1581 Some("permission-localized")
1582 );
1583 assert_eq!(
1584 allowed
1585 .metadata
1586 .get(PERMISSION_REEXECUTE_GENERATION_METADATA_KEY)
1587 .map(String::as_str),
1588 Some("generation-localized")
1589 );
1590
1591 let mut denied = session();
1592 denied.metadata.insert(
1593 PERMISSION_REEXECUTE_METADATA_KEY.to_string(),
1594 "stale-call".to_string(),
1595 );
1596 denied.metadata.insert(
1597 PERMISSION_REEXECUTE_GENERATION_METADATA_KEY.to_string(),
1598 "stale-generation".to_string(),
1599 );
1600 let deny_input = RespondInput {
1601 session_id: "sess-typed".to_string(),
1602 user_response: "Approve".to_string(),
1604 model: None,
1605 model_ref: None,
1606 provider: None,
1607 reasoning_effort: None,
1608 };
1609 apply_pending_response(
1610 &mut denied,
1611 &deny_input,
1612 Some("permission-localized"),
1613 ResponseSource::Human,
1614 Some(&receipt(PermissionDecisionKind::DenyOnce)),
1615 )
1616 .expect("typed deny must not depend on display text");
1617 assert!(!denied
1618 .metadata
1619 .contains_key(PERMISSION_REEXECUTE_METADATA_KEY));
1620 assert!(!denied
1621 .metadata
1622 .contains_key(PERMISSION_REEXECUTE_GENERATION_METADATA_KEY));
1623
1624 let mut legacy = session();
1625 assert!(matches!(
1626 apply_pending_response(
1627 &mut legacy,
1628 &allow_input,
1629 Some("permission-localized"),
1630 ResponseSource::Human,
1631 None,
1632 ),
1633 Err(RespondError::InvalidResponse(_))
1634 ));
1635 assert!(legacy.pending_question.is_some());
1636 }
1637}
1638
1639#[cfg(test)]
1640mod receipt_persistence_tests {
1641 use super::*;
1642 use crate::session_app::approval_replay::{
1643 find_permission_replay_target, restore_permission_replay_authorization,
1644 };
1645 use crate::{read_cached_session, SessionCache, SessionRepository};
1646 use bamboo_agent_core::storage::Storage;
1647 use bamboo_agent_core::tools::{FunctionCall, ToolCall};
1648 use bamboo_storage::{LockedSessionStore, SessionStoreV2};
1649 use bamboo_tools::permission::{
1650 PermissionConfig, PermissionDecision, PermissionMode, PermissionReasonCode,
1651 PermissionRequest, RiskLevel,
1652 };
1653
1654 fn pending_permission(
1655 session_id: &str,
1656 ) -> (Session, PermissionRequest, PermissionDecisionReceipt) {
1657 let mut session = Session::new(session_id, "test-model");
1658 session.agent_runtime_state = Some(AgentRuntimeState::new("receipt-test-run"));
1659 let request = |generation: &str, resource: &str| PermissionRequest {
1660 request_id: "permission-reused".into(),
1661 request_generation: generation.into(),
1662 session_id: session_id.into(),
1663 workspace_path: None,
1664 tool_name: "Bash".into(),
1665 permission_type: PermissionType::ExecuteCommand,
1666 resource: resource.into(),
1667 operation_summary: format!("execute {resource}"),
1668 risk_level: RiskLevel::High,
1669 reason_code: PermissionReasonCode::RiskThreshold,
1670 effective_mode: PermissionMode::Default,
1671 bypass_requested: false,
1672 auto_approve_requested: false,
1673 policy_revision: 4,
1674 matched_rule: None,
1675 allowed_decisions: vec![
1676 PermissionDecisionKind::AllowOnce,
1677 PermissionDecisionKind::DenyOnce,
1678 ],
1679 suggested_matchers: vec![],
1680 };
1681 let current = request("generation-current", "current-command");
1682 for (name, request) in [
1683 ("old", request("generation-old", "old-command")),
1684 ("current", current.clone()),
1685 ] {
1686 session.add_message(Message::assistant(
1687 "",
1688 Some(vec![ToolCall {
1689 id: request.request_id.clone(),
1690 tool_type: "function".into(),
1691 function: FunctionCall {
1692 name: request.tool_name.clone(),
1693 arguments: serde_json::json!({"command": request.resource}).to_string(),
1694 },
1695 }]),
1696 ));
1697 let mut result = Message::tool_result_with_status(
1698 &request.request_id,
1699 serde_json::json!({
1700 "status": "awaiting_permission_approval",
1701 "permission_type": request.permission_type,
1702 "resource": request.resource,
1703 "permission_request": request,
1704 })
1705 .to_string(),
1706 false,
1707 );
1708 result.id = format!("result-{name}");
1709 session.add_message(result);
1710 }
1711 session.set_pending_question_with_source(
1712 current.request_id.clone(),
1713 current.tool_name.clone(),
1714 "允许本次操作?".into(),
1715 vec!["允许".into(), "拒绝".into()],
1716 false,
1717 bamboo_agent_core::PendingQuestionSource::PauseTool,
1718 );
1719 session.metadata.insert(
1720 "runtime.suspend_reason".into(),
1721 "awaiting_clarification".into(),
1722 );
1723 let receipt = PermissionDecisionReceipt {
1724 session_id: session_id.into(),
1725 decision: PermissionDecision {
1726 request_id: current.request_id.clone(),
1727 request_generation: current.request_generation.clone(),
1728 decision: PermissionDecisionKind::AllowOnce,
1729 matcher_id: None,
1730 expected_policy_revision: Some(current.policy_revision),
1731 confirm_global: false,
1732 },
1733 decided_at: Utc::now(),
1734 };
1735 (session, current, receipt)
1736 }
1737
1738 async fn repository(
1739 session: &mut Session,
1740 ) -> (tempfile::TempDir, Arc<SessionStoreV2>, SessionRepository) {
1741 let directory = tempfile::tempdir().unwrap();
1742 let store = Arc::new(
1743 SessionStoreV2::new(directory.path().to_path_buf())
1744 .await
1745 .unwrap(),
1746 );
1747 let storage: Arc<dyn Storage> = store.clone();
1748 let repo = SessionRepository::new(
1749 SessionCache::default(),
1750 storage.clone(),
1751 Arc::new(LockedSessionStore::new(storage)),
1752 );
1753 repo.save(session).await.unwrap();
1754 (directory, store, repo)
1755 }
1756
1757 fn input(session_id: &str) -> RespondInput {
1758 RespondInput {
1759 session_id: session_id.into(),
1760 user_response: "允许".into(),
1761 model: None,
1762 model_ref: None,
1763 provider: None,
1764 reasoning_effort: None,
1765 }
1766 }
1767
1768 #[tokio::test]
1769 async fn typed_response_round_trip_restores_only_the_exact_allow_once() {
1770 eprintln!("debug_assertions={}", cfg!(debug_assertions));
1771 let (mut session, request, receipt) = pending_permission("receipt-round-trip");
1772 let (_directory, store, repo) = repository(&mut session).await;
1773 let old_occurrence = serde_json::to_value(&session.messages[..2]).unwrap();
1774 let guard = acquire_pending_response_guard(&session.id).await;
1775 let (accepted, _, _, _) = submit_pending_permission_response_checked_guarded(
1776 &repo,
1777 input(&session.id),
1778 Some(request.request_id.clone()),
1779 receipt.clone(),
1780 &guard,
1781 )
1782 .await
1783 .expect("public typed response accepts the current occurrence");
1784 let durable = store.load_session(&session.id).await.unwrap().unwrap();
1785 let restarted: Session =
1786 serde_json::from_slice(&serde_json::to_vec(&durable).unwrap()).unwrap();
1787 assert!(accepted.pending_question.is_none() && restarted.pending_question.is_none());
1788 assert_eq!(
1789 serde_json::to_value(&restarted.messages[..2]).unwrap(),
1790 old_occurrence
1791 );
1792 let message = restarted
1793 .messages
1794 .iter()
1795 .find(|message| message.id == "result-current")
1796 .unwrap();
1797 assert_eq!(
1798 message.tool_call_id.as_deref(),
1799 Some(request.request_id.as_str())
1800 );
1801 assert_eq!(message.content, "Selected response: 允许");
1802 assert_eq!(message.tool_success, Some(true));
1803 let metadata = message
1804 .metadata
1805 .as_ref()
1806 .expect("preserved typed request metadata");
1807 assert_eq!(
1808 metadata.get("permission_request"),
1809 Some(&serde_json::to_value(&request).unwrap())
1810 );
1811 assert_eq!(
1812 metadata.get("permission_decision_receipt"),
1813 Some(&serde_json::to_value(&receipt).unwrap()),
1814 "the public response must persist the complete receipt, including decided_at"
1815 );
1816 assert_eq!(
1817 restarted.metadata.get(PERMISSION_REEXECUTE_METADATA_KEY),
1818 Some(&request.request_id)
1819 );
1820 assert_eq!(
1821 restarted
1822 .metadata
1823 .get(PERMISSION_REEXECUTE_GENERATION_METADATA_KEY),
1824 Some(&request.request_generation)
1825 );
1826
1827 let target = find_permission_replay_target(
1828 &restarted,
1829 &request.request_id,
1830 Some(&request.request_generation),
1831 )
1832 .expect("replay resolves the exact current occurrence after reload");
1833 assert_eq!(
1834 target.request_generation(),
1835 Some(request.request_generation.as_str())
1836 );
1837 assert_eq!(
1838 target.tool_call().function.arguments,
1839 serde_json::json!({"command":"current-command"}).to_string()
1840 );
1841 let config = PermissionConfig::new();
1842 restore_permission_replay_authorization(&config, &restarted, &target, "Bash").unwrap();
1843 for (session_id, call_id, generation, resource) in [
1846 (
1847 restarted.id.as_str(),
1848 request.request_id.as_str(),
1849 "generation-old",
1850 request.resource.as_str(),
1851 ),
1852 (
1853 restarted.id.as_str(),
1854 request.request_id.as_str(),
1855 request.request_generation.as_str(),
1856 "other-command",
1857 ),
1858 (
1859 "other-session",
1860 request.request_id.as_str(),
1861 request.request_generation.as_str(),
1862 request.resource.as_str(),
1863 ),
1864 ] {
1865 assert!(!config.consume_once_for_generation(
1866 session_id,
1867 call_id,
1868 generation,
1869 request.permission_type,
1870 resource
1871 ));
1872 }
1873 assert!(config.consume_once_for_generation(
1874 &restarted.id,
1875 &request.request_id,
1876 &request.request_generation,
1877 request.permission_type,
1878 &request.resource
1879 ));
1880 assert!(!config.consume_once_for_generation(
1881 &restarted.id,
1882 &request.request_id,
1883 &request.request_generation,
1884 request.permission_type,
1885 &request.resource
1886 ));
1887 }
1888
1889 #[tokio::test]
1890 async fn typed_deny_round_trip_retains_receipt_without_grants_or_replay() {
1891 eprintln!("debug_assertions={}", cfg!(debug_assertions));
1892 let (mut session, request, mut receipt) = pending_permission("receipt-deny");
1893 receipt.decision.decision = PermissionDecisionKind::DenyOnce;
1894 session.metadata.insert(
1895 PERMISSION_REEXECUTE_METADATA_KEY.into(),
1896 "stale-call".into(),
1897 );
1898 session.metadata.insert(
1899 PERMISSION_REEXECUTE_GENERATION_METADATA_KEY.into(),
1900 "stale-generation".into(),
1901 );
1902 let (_directory, store, repo) = repository(&mut session).await;
1903 let old_occurrence = serde_json::to_value(&session.messages[..2]).unwrap();
1904 let mut response = input(&session.id);
1905 response.user_response = "Approve".into();
1906 let guard = acquire_pending_response_guard(&session.id).await;
1907 let (_, display_response, _, grants) = submit_pending_permission_response_checked_guarded(
1908 &repo,
1909 response,
1910 Some(request.request_id.clone()),
1911 receipt.clone(),
1912 &guard,
1913 )
1914 .await
1915 .expect("typed deny accepts display text without granting permission");
1916 assert_eq!(display_response, "Approve");
1917 assert!(grants.is_empty());
1918 let durable = store.load_session(&session.id).await.unwrap().unwrap();
1919 let restarted: Session =
1920 serde_json::from_slice(&serde_json::to_vec(&durable).unwrap()).unwrap();
1921 assert!(restarted.pending_question.is_none());
1922 assert_eq!(
1923 serde_json::to_value(&restarted.messages[..2]).unwrap(),
1924 old_occurrence
1925 );
1926 let message = restarted.messages.last().unwrap();
1927 assert_eq!(message.id, "result-current");
1928 assert_eq!(message.content, "Selected response: Approve");
1929 let metadata = message.metadata.as_ref().unwrap();
1930 assert_eq!(
1931 metadata.get("permission_request"),
1932 Some(&serde_json::to_value(&request).unwrap())
1933 );
1934 assert_eq!(
1935 metadata.get("permission_decision_receipt"),
1936 Some(&serde_json::to_value(&receipt).unwrap())
1937 );
1938 assert!(!restarted
1939 .metadata
1940 .contains_key(PERMISSION_REEXECUTE_METADATA_KEY));
1941 assert!(!restarted
1942 .metadata
1943 .contains_key(PERMISSION_REEXECUTE_GENERATION_METADATA_KEY));
1944 }
1945
1946 #[tokio::test]
1947 async fn failed_typed_receipt_does_not_publish_a_partial_response() {
1948 eprintln!("debug_assertions={}", cfg!(debug_assertions));
1949 let (mut session, request, receipt) = pending_permission("receipt-conflict");
1950 let current = session.messages.last_mut().unwrap();
1951 current.metadata = Some(serde_json::json!({"permission_request": request}));
1952 let mut payload: serde_json::Value = serde_json::from_str(¤t.content).unwrap();
1953 payload["permission_request"]["request_generation"] =
1954 "generation-conflicting-payload".into();
1955 current.content = payload.to_string();
1956 let (_directory, store, repo) = repository(&mut session).await;
1957 let durable_before = store.load_session(&session.id).await.unwrap().unwrap();
1958 assert_eq!(
1962 permission_request_generation(durable_before.messages.last().unwrap()).as_deref(),
1963 Some(receipt.decision.request_generation.as_str())
1964 );
1965 let cache_before =
1966 serde_json::to_value(read_cached_session(repo.cache(), &session.id).unwrap()).unwrap();
1967 let durable_before = serde_json::to_value(durable_before).unwrap();
1968 let guard = acquire_pending_response_guard(&session.id).await;
1969 let result = submit_pending_permission_response_checked_guarded(
1970 &repo,
1971 input(&session.id),
1972 Some(request.request_id),
1973 receipt,
1974 &guard,
1975 )
1976 .await;
1977 assert!(
1978 matches!(result, Err(RespondError::InvalidResponse(ref message)) if message.contains("receipt")),
1979 "failed receipt persistence must be an explicit response error: {result:?}"
1980 );
1981 assert_eq!(
1982 serde_json::to_value(store.load_session(&session.id).await.unwrap().unwrap()).unwrap(),
1983 durable_before,
1984 "pending question, result and all durable markers stay unchanged"
1985 );
1986 assert_eq!(
1987 serde_json::to_value(read_cached_session(repo.cache(), &session.id).unwrap()).unwrap(),
1988 cache_before,
1989 "a rejected response must not publish a cache snapshot"
1990 );
1991 }
1992}