Skip to main content

bamboo_engine/session_app/
respond.rs

1//! Respond use case: submit a user response to a pending question.
2
3use 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
33/// Process-local serialization for consuming one pending question. The
34/// persistence layer serializes individual loads and saves, but releasing its
35/// lock between those operations would let two responders both observe and
36/// consume the same question. Every response entrypoint reaches this use case,
37/// so holding this gate across load -> validate -> save makes one consumer win
38/// and forces later callers to reload the already-consumed state.
39struct 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
87/// Process-local ownership for one session's complete response transaction.
88/// High-level HTTP, Connect, and Gold paths hold this from authoritative
89/// preflight through successor dispatch, so a stale duplicate cannot reserve
90/// and later cancel a phantom runner after the real successor has completed.
91pub 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
109/// Acquire the shared response single-flight used by every response source.
110pub 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/// Number of callers currently waiting to enter a session's response gate.
124/// Exposed for deterministic concurrency tests and lightweight diagnostics.
125#[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
134/// Reload the exact snapshot a response CAS would mutate while holding the
135/// shared response single-flight. This must happen before successor
136/// reservation so an already-consumed/stale response allocates no runner.
137pub 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
146/// Session-metadata key marking a tool call that was approved through a permission
147/// prompt and must be RE-EXECUTED on resume. The gated tool never actually ran
148/// (the permission gate intercepted it before execution), so on approval the
149/// server resume adapter re-runs it and writes the real output back — instead of
150/// leaving the model to infer/fabricate it. Value = the tool_call_id.
151pub 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
177/// Submit a pending response: load session, validate, update messages,
178/// apply plan mode transitions, persist, and return the updated session.
179///
180/// The caller (handler) is responsible for auto-resume triggering.
181pub 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
196/// Submit a response guarded by the exact pending tool-call identity displayed
197/// to a typed client. The separate parameter keeps the established public
198/// [`RespondInput`] struct source-compatible for SDK/in-process callers.
199pub 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
221/// Human-response variant for callers that already hold the shared response
222/// transaction guard across preflight, reservation, CAS, and dispatch.
223pub 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
247/// Typed-permission variant that persists the exact decision receipt in the
248/// same durable session mutation that consumes the pending question.
249pub 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
316/// Guarded form for high-level response + resume transactions. The caller
317/// must retain `guard` until the exact successor reservation has been handed
318/// to its detached execution owner.
319pub 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 let Some(expected) = expected_tool_call_id {
419        if pending.tool_call_id != expected {
420            let actual = pending.tool_call_id.clone();
421            session.pending_question = Some(pending);
422            return Err(RespondError::PendingQuestionMismatch {
423                expected: expected.to_string(),
424                actual,
425            });
426        }
427    }
428
429    if let Some(receipt) = permission_receipt {
430        if receipt.session_id != input.session_id
431            || receipt.decision.request_id != pending.tool_call_id
432        {
433            session.pending_question = Some(pending);
434            return Err(RespondError::InvalidResponse(
435                "permission receipt identity does not match the pending question".to_string(),
436            ));
437        }
438        let current_generation = session
439            .messages
440            .iter()
441            .rev()
442            .find(|message| message.tool_call_id.as_deref() == Some(pending.tool_call_id.as_str()))
443            .and_then(permission_request_generation);
444        if current_generation.as_deref() != Some(receipt.decision.request_generation.as_str()) {
445            session.pending_question = Some(pending);
446            return Err(RespondError::InvalidResponse(
447                "permission receipt generation does not match the pending operation".to_string(),
448            ));
449        }
450    }
451
452    // Typed permission control flow comes exclusively from the structured
453    // receipt. Display strings and localized options remain transcript-only.
454    // Legacy clarifications still validate their selected display option.
455    if permission_receipt.is_none() {
456        if let Err(error_message) = validate_pending_response(&pending, &input.user_response) {
457            // Put the pending question back when validation fails.
458            session.pending_question = Some(pending);
459            return Err(RespondError::InvalidResponse(error_message));
460        }
461    }
462
463    let tool_call_id = pending.tool_call_id.clone();
464    tracing::debug!(
465        "[{}] Looking for tool result message with tool_call_id: {}",
466        input.session_id,
467        tool_call_id
468    );
469
470    let reviewed_plan = extract_exit_plan_from_tool_result_message(session, &tool_call_id);
471
472    // Permission grants implied by approving a permission prompt. Read from the
473    // (still-unmodified) synthesized tool-result payload, BEFORE it is overwritten
474    // by the user's selection below.
475    let typed_permission_approved = permission_receipt.is_some_and(|receipt| {
476        matches!(
477            receipt.decision.decision,
478            PermissionDecisionKind::AllowOnce
479                | PermissionDecisionKind::AllowSession
480                | PermissionDecisionKind::AllowWorkspace
481                | PermissionDecisionKind::AllowGlobal
482        )
483    });
484    let permission_approved = permission_receipt
485        .map(|_| typed_permission_approved)
486        .unwrap_or_else(|| is_permission_approval(&input.user_response));
487    let permission_grants = if permission_approved {
488        extract_permission_grants_from_tool_result_message(session, &tool_call_id)
489    } else {
490        Vec::new()
491    };
492    let should_reexecute = permission_receipt
493        .map(|_| typed_permission_approved)
494        .unwrap_or(!permission_grants.is_empty());
495    if should_reexecute {
496        // Approved a permission prompt: mark the gated tool call for re-execution
497        // on resume so the operation actually runs (real output) rather than the
498        // model inferring it. Consumed by the server resume adapter.
499        session.metadata.insert(
500            PERMISSION_REEXECUTE_METADATA_KEY.to_string(),
501            tool_call_id.clone(),
502        );
503        if let Some(receipt) = permission_receipt {
504            session.metadata.insert(
505                PERMISSION_REEXECUTE_GENERATION_METADATA_KEY.to_string(),
506                receipt.decision.request_generation.clone(),
507            );
508        } else {
509            session
510                .metadata
511                .remove(PERMISSION_REEXECUTE_GENERATION_METADATA_KEY);
512        }
513    } else if permission_receipt.is_some() {
514        // A typed deny cannot inherit replay markers from an older occurrence.
515        session.metadata.remove(PERMISSION_REEXECUTE_METADATA_KEY);
516        session
517            .metadata
518            .remove(PERMISSION_REEXECUTE_GENERATION_METADATA_KEY);
519    }
520
521    // ---- Update or append tool result message ----
522    let found = update_or_append_tool_result_message(
523        session,
524        &tool_call_id,
525        &input.user_response,
526        response_source,
527    );
528    if let Some(receipt) = permission_receipt {
529        debug_assert!(persist_permission_decision_receipt(
530            session,
531            &tool_call_id,
532            receipt
533        ));
534    }
535    if found {
536        tracing::info!(
537            "[{}] Updated existing tool result message",
538            input.session_id
539        );
540    } else {
541        tracing::warn!(
542            "[{}] Tool result message not found for tool_call_id: {}, added fallback message",
543            input.session_id,
544            tool_call_id
545        );
546    }
547
548    // ---- Plan mode state transitions ----
549    let plan_mode_transition =
550        apply_plan_mode_transition(session, &pending, &input.user_response, reviewed_plan);
551
552    // ---- Clear pending question and set resume marker ----
553    session.clear_pending_question();
554    record_consumed_clarification(session, &tool_call_id);
555    session.metadata.remove("runtime.suspend_reason");
556    session.metadata.insert(
557        CLARIFICATION_RESUME_PENDING_KEY.to_string(),
558        "true".to_string(),
559    );
560    session.metadata.insert(
561        CONCLUSION_WITH_OPTIONS_RESUME_PENDING_KEY.to_string(),
562        "true".to_string(),
563    );
564    mark_startup_handoff(session);
565
566    // ---- Merge model/reasoning from request ----
567    let request_model_ref = derive_model_ref(
568        input.model_ref.as_ref(),
569        input.provider.as_deref(),
570        input.model.as_deref(),
571    );
572    if let Some(model_ref) = request_model_ref.as_ref() {
573        persist_model_ref(session, model_ref);
574    } else {
575        persist_legacy_model_provider(session, input.model.as_deref(), input.provider.as_deref());
576    }
577    if let Some(reasoning_effort) = input.reasoning_effort {
578        session.reasoning_effort = Some(reasoning_effort);
579    }
580
581    Ok((
582        input.user_response.clone(),
583        plan_mode_transition,
584        permission_grants,
585    ))
586}
587
588fn record_consumed_clarification(session: &mut Session, tool_call_id: &str) {
589    let mut legacy_consumed = session
590        .metadata
591        .get(CONSUMED_CLARIFICATION_IDS_KEY)
592        .and_then(|value| serde_json::from_str::<Vec<String>>(value).ok())
593        .unwrap_or_default();
594    legacy_consumed.retain(|existing| existing != tool_call_id);
595    legacy_consumed.push(tool_call_id.to_string());
596    if legacy_consumed.len() > 64 {
597        legacy_consumed.drain(..legacy_consumed.len() - 64);
598    }
599    if let Ok(serialized) = serde_json::to_string(&legacy_consumed) {
600        session
601            .metadata
602            .insert(CONSUMED_CLARIFICATION_IDS_KEY.to_string(), serialized);
603    }
604
605    let Some(occurrence) = latest_response_occurrence(session, tool_call_id) else {
606        return;
607    };
608    let mut consumed = session
609        .metadata
610        .get(CONSUMED_RESPONSE_OCCURRENCES_KEY)
611        .and_then(|value| serde_json::from_str::<Vec<ResponseOccurrence>>(value).ok())
612        .unwrap_or_default();
613    consumed.retain(|existing| existing != &occurrence);
614    consumed.push(occurrence);
615    if consumed.len() > 64 {
616        consumed.drain(..consumed.len() - 64);
617    }
618    if let Ok(serialized) = serde_json::to_string(&consumed) {
619        session
620            .metadata
621            .insert(CONSUMED_RESPONSE_OCCURRENCES_KEY.to_string(), serialized);
622    }
623}
624
625/// Apply plan mode state transitions based on the pending question tool and user response.
626fn apply_plan_mode_transition(
627    session: &mut Session,
628    pending: &PendingQuestion,
629    user_response: &str,
630    reviewed_plan: Option<String>,
631) -> Option<PlanModeTransition> {
632    match pending.tool_name.as_str() {
633        "EnterPlanMode" if user_response.to_lowercase().contains("enter plan mode") => {
634            let pre_mode = session
635                .agent_runtime_state
636                .as_ref()
637                .map(|state| state.effective_permission_mode().as_str().to_string())
638                .unwrap_or_else(|| SessionPermissionMode::Default.as_str().to_string());
639
640            let entered_at = Utc::now();
641            let status = PlanModeStatus::Exploring;
642            let runtime_state = session
643                .agent_runtime_state
644                .get_or_insert_with(|| AgentRuntimeState::new(uuid::Uuid::new_v4().to_string()));
645            runtime_state.plan_mode = Some(PlanModeState {
646                entered_at,
647                pre_permission_mode: pre_mode.clone(),
648                plan_file_path: None,
649                status,
650            });
651            tracing::info!(
652                session_id = %session.id,
653                "Entered plan mode"
654            );
655            Some(PlanModeTransition::Entered {
656                reason: Some(pending.question.clone()),
657                pre_permission_mode: pre_mode,
658                entered_at,
659                status,
660                plan_file_path: None,
661            })
662        }
663        "ExitPlanMode" if is_exit_plan_mode_approved(user_response) => {
664            let restored_mode = session
665                .agent_runtime_state
666                .as_ref()
667                .map(|state| state.effective_permission_mode().as_str().to_string())
668                .unwrap_or_else(|| "default".to_string());
669            if let Some(ref mut runtime_state) = session.agent_runtime_state {
670                // The typed requested mode remains live while Plan is active and
671                // may have been changed by a newer PATCH. Exiting Plan clears
672                // only the overlay; the old pre-mode is event history, never a
673                // write authority that may roll back the newer request.
674                runtime_state.plan_mode = None;
675            }
676            tracing::info!(
677                session_id = %session.id,
678                "Exited plan mode"
679            );
680            Some(PlanModeTransition::Exited {
681                approved: true,
682                restored_mode,
683                plan: reviewed_plan,
684            })
685        }
686        _ => None,
687    }
688}
689
690/// Check if the user response approves exiting plan mode.
691fn is_exit_plan_mode_approved(user_response: &str) -> bool {
692    let lower = user_response.to_lowercase();
693    lower.contains("approve") && !lower.contains("stay in plan mode")
694}
695
696// ---- Internal helpers ----
697
698pub fn validate_pending_response(
699    pending: &PendingQuestion,
700    user_response: &str,
701) -> Result<(), String> {
702    if pending.allow_custom {
703        return Ok(());
704    }
705
706    let valid = pending.options.iter().any(|option| option == user_response);
707    if valid {
708        Ok(())
709    } else {
710        let options_str = pending.options.join(", ");
711        Err(format!("Response must be one of: {options_str}"))
712    }
713}
714
715pub fn update_or_append_tool_result_message(
716    session: &mut Session,
717    tool_call_id: &str,
718    user_response: &str,
719    response_source: ResponseSource,
720) -> bool {
721    for message in session.messages.iter_mut().rev() {
722        if message.tool_call_id.as_deref() == Some(tool_call_id) {
723            // Preserve the server-issued typed permission contract outside the
724            // model-visible content before replacing the synthetic waiting
725            // payload with the selected answer. This lets an exact durable
726            // decision receipt be reconstructed after a daemon restart without
727            // trusting display strings or replaying an already-consumed run.
728            if let Ok(payload) = serde_json::from_str::<serde_json::Value>(&message.content) {
729                if payload.get("status").and_then(serde_json::Value::as_str)
730                    == Some("awaiting_permission_approval")
731                {
732                    if let Some(request) = payload.get("permission_request") {
733                        insert_message_metadata(message, "permission_request", request.clone());
734                    }
735                }
736            }
737            message.content = selected_message_content(user_response, response_source);
738            message.tool_success = Some(true);
739            return true;
740        }
741    }
742
743    session.add_message(bamboo_agent_core::Message::tool_result_with_status(
744        tool_call_id,
745        selected_message_content(user_response, response_source),
746        true,
747    ));
748    false
749}
750
751fn insert_message_metadata(message: &mut Message, key: &str, value: serde_json::Value) {
752    let metadata = message
753        .metadata
754        .get_or_insert_with(|| serde_json::Value::Object(Default::default()));
755    if !metadata.is_object() {
756        let previous = std::mem::replace(metadata, serde_json::Value::Object(Default::default()));
757        metadata
758            .as_object_mut()
759            .expect("replacement metadata is an object")
760            .insert("previous_metadata".to_string(), previous);
761    }
762    metadata
763        .as_object_mut()
764        .expect("message metadata is an object")
765        .insert(key.to_string(), value);
766}
767
768fn permission_request_generation(message: &Message) -> Option<String> {
769    message
770        .metadata
771        .as_ref()
772        .and_then(|metadata| metadata.get("permission_request"))
773        .and_then(|request| request.get("request_generation"))
774        .and_then(serde_json::Value::as_str)
775        .map(ToOwned::to_owned)
776        .or_else(|| {
777            serde_json::from_str::<serde_json::Value>(&message.content)
778                .ok()?
779                .get("permission_request")?
780                .get("request_generation")?
781                .as_str()
782                .map(ToOwned::to_owned)
783        })
784}
785
786fn persist_permission_decision_receipt(
787    session: &mut Session,
788    tool_call_id: &str,
789    receipt: &PermissionDecisionReceipt,
790) -> bool {
791    let Some(message) = session
792        .messages
793        .iter_mut()
794        .rev()
795        .find(|message| message.tool_call_id.as_deref() == Some(tool_call_id))
796    else {
797        return false;
798    };
799    if permission_request_generation(message).as_deref()
800        != Some(receipt.decision.request_generation.as_str())
801    {
802        return false;
803    }
804    insert_message_metadata(
805        message,
806        "permission_decision_receipt",
807        serde_json::to_value(receipt).expect("permission receipt is serializable"),
808    );
809    true
810}
811
812fn selected_message_content(user_response: &str, response_source: ResponseSource) -> String {
813    match response_source {
814        ResponseSource::Human => format!("Selected response: {}", user_response),
815        ResponseSource::Gold => format!("Auto-selected response (gold): {}", user_response),
816    }
817}
818
819fn extract_exit_plan_from_tool_result_message(
820    session: &Session,
821    tool_call_id: &str,
822) -> Option<String> {
823    let message = session
824        .messages
825        .iter()
826        .rev()
827        .find(|message| message.tool_call_id.as_deref() == Some(tool_call_id))?;
828    let payload = serde_json::from_str::<serde_json::Value>(&message.content).ok()?;
829    payload
830        .get("plan")
831        .and_then(|value| value.as_str())
832        .map(str::trim)
833        .filter(|value| !value.is_empty())
834        .map(ToOwned::to_owned)
835}
836
837/// Detect whether the user response approves a pending permission request.
838///
839/// Permission prompts (synthesized by the permission gate, and the
840/// `request_permissions` tool) offer exactly `["Approve", "Deny"]`.
841fn is_permission_approval(user_response: &str) -> bool {
842    user_response.trim().eq_ignore_ascii_case("approve")
843}
844
845/// Extract the permission grants implied by an approved permission prompt.
846///
847/// Reads the pending tool-result message (still the synthesized
848/// `awaiting_permission_approval` payload, before it is overwritten by the
849/// user's selection) and returns the `(PermissionType, resource)` pairs the
850/// caller should grant for the session. Handles both the single-gated-tool shape
851/// (top-level `permission_type` + `resource`) and the `request_permissions` shape
852/// (a `permissions` array).
853fn extract_permission_grants_from_tool_result_message(
854    session: &Session,
855    tool_call_id: &str,
856) -> Vec<(PermissionType, String)> {
857    let message = match session
858        .messages
859        .iter()
860        .rev()
861        .find(|message| message.tool_call_id.as_deref() == Some(tool_call_id))
862    {
863        Some(message) => message,
864        None => return Vec::new(),
865    };
866    let payload = match serde_json::from_str::<serde_json::Value>(&message.content) {
867        Ok(payload) => payload,
868        Err(_) => return Vec::new(),
869    };
870    if payload.get("status").and_then(|value| value.as_str())
871        != Some("awaiting_permission_approval")
872    {
873        return Vec::new();
874    }
875
876    let parse_one = |value: &serde_json::Value| -> Option<(PermissionType, String)> {
877        let type_value = value
878            .get("permission_type")
879            .or_else(|| value.get("type"))?
880            .clone();
881        let perm_type: PermissionType = serde_json::from_value(type_value).ok()?;
882        let resource = value.get("resource")?.as_str()?.trim().to_string();
883        if resource.is_empty() {
884            return None;
885        }
886        Some((perm_type, resource))
887    };
888
889    if let Some(array) = payload
890        .get("permissions")
891        .and_then(|value| value.as_array())
892    {
893        array.iter().filter_map(parse_one).collect()
894    } else {
895        parse_one(&payload).into_iter().collect()
896    }
897}
898
899#[cfg(test)]
900mod tests {
901    use super::*;
902    use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
903    use std::sync::Mutex;
904    use tokio::sync::Notify;
905
906    struct TestSessionAccess {
907        session: Mutex<Session>,
908        loads: AtomicUsize,
909        saves: AtomicUsize,
910        block_first_save: AtomicBool,
911        save_started: Notify,
912        release_save: Notify,
913    }
914
915    impl TestSessionAccess {
916        fn with_pending() -> Self {
917            Self::with_pending_and_blocked_save(false)
918        }
919
920        fn with_pending_and_blocked_save(block_first_save: bool) -> Self {
921            let mut session = Session::new("sess-1", "test-model");
922            session.pending_question = Some(make_pending("ConclusionWithOptions"));
923            Self {
924                session: Mutex::new(session),
925                loads: AtomicUsize::new(0),
926                saves: AtomicUsize::new(0),
927                block_first_save: AtomicBool::new(block_first_save),
928                save_started: Notify::new(),
929                release_save: Notify::new(),
930            }
931        }
932    }
933
934    #[async_trait::async_trait]
935    impl SessionAccess for TestSessionAccess {
936        async fn load_session(
937            &self,
938            _id: &str,
939        ) -> Result<Option<Session>, super::super::errors::SessionLoadError> {
940            Ok(Some(self.session.lock().unwrap().clone()))
941        }
942
943        async fn load_or_create(
944            &self,
945            _id: &str,
946            _model: &str,
947        ) -> Result<Session, super::super::errors::SessionLoadError> {
948            Ok(self.session.lock().unwrap().clone())
949        }
950
951        async fn load_merged(
952            &self,
953            _id: &str,
954        ) -> Result<Option<Session>, super::super::errors::SessionLoadError> {
955            self.loads.fetch_add(1, Ordering::SeqCst);
956            Ok(Some(self.session.lock().unwrap().clone()))
957        }
958
959        async fn save_session(
960            &self,
961            session: &mut Session,
962        ) -> Result<(), super::super::errors::SessionSaveError> {
963            if self.block_first_save.swap(false, Ordering::SeqCst) {
964                self.save_started.notify_one();
965                self.release_save.notified().await;
966            }
967            *self.session.lock().unwrap() = session.clone();
968            self.saves.fetch_add(1, Ordering::SeqCst);
969            Ok(())
970        }
971
972        async fn save_and_cache(
973            &self,
974            session: &mut Session,
975        ) -> Result<(), super::super::errors::SessionSaveError> {
976            self.save_session(session).await
977        }
978    }
979
980    fn respond_input() -> RespondInput {
981        RespondInput {
982            session_id: "sess-1".to_string(),
983            user_response: "A".to_string(),
984            model: None,
985            model_ref: None,
986            provider: None,
987            reasoning_effort: None,
988        }
989    }
990
991    fn make_pending(tool_name: &str) -> PendingQuestion {
992        PendingQuestion {
993            tool_call_id: "call-1".to_string(),
994            tool_name: tool_name.to_string(),
995            question: "Question?".to_string(),
996            options: vec!["A".to_string(), "B".to_string()],
997            allow_custom: false,
998            source: bamboo_agent_core::PendingQuestionSource::PauseTool,
999        }
1000    }
1001
1002    #[tokio::test]
1003    async fn expected_tool_call_id_accepts_the_current_question() {
1004        let repo = TestSessionAccess::with_pending();
1005
1006        submit_pending_response_checked(&repo, respond_input(), Some("call-1".to_string()))
1007            .await
1008            .expect("matching identity should submit");
1009
1010        assert_eq!(repo.saves.load(Ordering::SeqCst), 1);
1011        assert!(repo.session.lock().unwrap().pending_question.is_none());
1012    }
1013
1014    #[tokio::test]
1015    async fn stale_tool_call_id_is_rejected_without_consuming_the_question() {
1016        let repo = TestSessionAccess::with_pending();
1017
1018        let error =
1019            submit_pending_response_checked(&repo, respond_input(), Some("stale-call".to_string()))
1020                .await
1021                .expect_err("stale identity must fail");
1022
1023        assert!(matches!(
1024            error,
1025            RespondError::PendingQuestionMismatch {
1026                ref expected,
1027                ref actual,
1028            } if expected == "stale-call" && actual == "call-1"
1029        ));
1030        assert_eq!(repo.saves.load(Ordering::SeqCst), 0);
1031        assert_eq!(
1032            repo.session
1033                .lock()
1034                .unwrap()
1035                .pending_question
1036                .as_ref()
1037                .map(|question| question.tool_call_id.as_str()),
1038            Some("call-1")
1039        );
1040    }
1041
1042    #[tokio::test]
1043    async fn omitted_tool_call_guard_remains_backwards_compatible() {
1044        let repo = TestSessionAccess::with_pending();
1045
1046        submit_pending_response(&repo, respond_input())
1047            .await
1048            .expect("legacy client should remain accepted");
1049
1050        assert_eq!(repo.saves.load(Ordering::SeqCst), 1);
1051    }
1052
1053    #[tokio::test]
1054    async fn concurrent_responses_consume_a_pending_question_once() {
1055        let repo = Arc::new(TestSessionAccess::with_pending_and_blocked_save(true));
1056
1057        let first_repo = repo.clone();
1058        let first = tokio::spawn(async move {
1059            submit_pending_response_checked(
1060                first_repo.as_ref(),
1061                respond_input(),
1062                Some("call-1".to_string()),
1063            )
1064            .await
1065        });
1066        repo.save_started.notified().await;
1067
1068        let second_entered = Arc::new(Notify::new());
1069        let second_repo = repo.clone();
1070        let second_entered_task = second_entered.clone();
1071        let second = tokio::spawn(async move {
1072            second_entered_task.notify_one();
1073            submit_pending_response_checked(
1074                second_repo.as_ref(),
1075                respond_input(),
1076                Some("call-1".to_string()),
1077            )
1078            .await
1079        });
1080        second_entered.notified().await;
1081        tokio::task::yield_now().await;
1082
1083        assert_eq!(
1084            repo.loads.load(Ordering::SeqCst),
1085            1,
1086            "the second responder must wait before loading the pending question"
1087        );
1088        repo.release_save.notify_one();
1089
1090        first.await.unwrap().expect("first response should win");
1091        let error = second
1092            .await
1093            .unwrap()
1094            .expect_err("second response must observe the consumed question");
1095        assert!(matches!(error, RespondError::NoPendingQuestion));
1096        assert_eq!(repo.loads.load(Ordering::SeqCst), 2);
1097        assert_eq!(repo.saves.load(Ordering::SeqCst), 1);
1098    }
1099
1100    #[tokio::test]
1101    async fn cancelled_response_waiter_does_not_leak_diagnostics_or_the_gate() {
1102        let session_id = "cancelled-response-waiter";
1103        let owner = acquire_pending_response_guard(session_id).await;
1104        let waiter = tokio::spawn(async move { acquire_pending_response_guard(session_id).await });
1105        tokio::time::timeout(std::time::Duration::from_secs(1), async {
1106            while pending_response_waiter_count(session_id) != 1 {
1107                tokio::task::yield_now().await;
1108            }
1109        })
1110        .await
1111        .expect("waiter should reach the gate");
1112
1113        waiter.abort();
1114        let _ = waiter.await;
1115        assert_eq!(pending_response_waiter_count(session_id), 0);
1116        drop(owner);
1117
1118        tokio::time::timeout(
1119            std::time::Duration::from_secs(1),
1120            acquire_pending_response_guard(session_id),
1121        )
1122        .await
1123        .expect("cancelled waiter must not retain the gate");
1124    }
1125
1126    #[test]
1127    fn enter_plan_mode_activates_plan_mode_state() {
1128        let mut session = Session::new("sess-1", "test-model");
1129        let pending = make_pending("EnterPlanMode");
1130
1131        apply_plan_mode_transition(&mut session, &pending, "Enter plan mode", None);
1132
1133        assert!(session.agent_runtime_state.is_some());
1134        let state = session.agent_runtime_state.unwrap();
1135        assert!(state.plan_mode.is_some());
1136        let plan = state.plan_mode.unwrap();
1137        assert_eq!(plan.status, PlanModeStatus::Exploring);
1138        assert_eq!(plan.pre_permission_mode, "default");
1139    }
1140
1141    #[test]
1142    fn enter_plan_mode_does_nothing_when_not_approved() {
1143        let mut session = Session::new("sess-1", "test-model");
1144        let pending = make_pending("EnterPlanMode");
1145
1146        apply_plan_mode_transition(&mut session, &pending, "Stay in normal mode", None);
1147
1148        assert!(session.agent_runtime_state.is_none());
1149    }
1150
1151    #[test]
1152    fn plan_mode_preserves_and_restores_typed_auto_request() {
1153        let mut session = Session::new("sess-auto-plan", "test-model");
1154        session
1155            .agent_runtime_state
1156            .get_or_insert_with(|| AgentRuntimeState::new("run-1"))
1157            .set_permission_mode(SessionPermissionMode::Auto);
1158        let enter = make_pending("EnterPlanMode");
1159        let transition =
1160            apply_plan_mode_transition(&mut session, &enter, "Enter plan mode", None).unwrap();
1161        assert!(matches!(
1162            transition,
1163            PlanModeTransition::Entered {
1164                ref pre_permission_mode,
1165                ..
1166            } if pre_permission_mode == "auto"
1167        ));
1168        assert_eq!(
1169            session
1170                .agent_runtime_state
1171                .as_ref()
1172                .unwrap()
1173                .plan_mode
1174                .as_ref()
1175                .unwrap()
1176                .pre_permission_mode,
1177            "auto"
1178        );
1179
1180        let exit = make_pending("ExitPlanMode");
1181        let transition = apply_plan_mode_transition(
1182            &mut session,
1183            &exit,
1184            "Approve (Auto mode)",
1185            Some("Reviewed plan".to_string()),
1186        )
1187        .unwrap();
1188        assert!(matches!(
1189            transition,
1190            PlanModeTransition::Exited {
1191                ref restored_mode,
1192                ..
1193            } if restored_mode == "auto"
1194        ));
1195        assert_eq!(
1196            session
1197                .agent_runtime_state
1198                .as_ref()
1199                .unwrap()
1200                .effective_permission_mode(),
1201            SessionPermissionMode::Auto
1202        );
1203    }
1204
1205    #[test]
1206    fn exit_plan_mode_does_not_restore_over_a_newer_typed_mode() {
1207        let mut session = Session::new("sess-plan-patch", "test-model");
1208        let state = session
1209            .agent_runtime_state
1210            .get_or_insert_with(|| AgentRuntimeState::new("run-1"));
1211        state.set_permission_mode(SessionPermissionMode::Auto);
1212        state.plan_mode = Some(PlanModeState {
1213            entered_at: Utc::now(),
1214            pre_permission_mode: "auto".to_string(),
1215            plan_file_path: None,
1216            status: PlanModeStatus::AwaitingApproval,
1217        });
1218        // Simulate a newer PATCH while the Plan overlay is still active.
1219        state.set_permission_mode(SessionPermissionMode::Bypass);
1220
1221        let transition = apply_plan_mode_transition(
1222            &mut session,
1223            &make_pending("ExitPlanMode"),
1224            "Approve (Default mode)",
1225            None,
1226        )
1227        .unwrap();
1228
1229        assert!(matches!(
1230            transition,
1231            PlanModeTransition::Exited {
1232                ref restored_mode,
1233                ..
1234            } if restored_mode == "bypass"
1235        ));
1236        let state = session.agent_runtime_state.unwrap();
1237        assert!(state.plan_mode.is_none());
1238        assert_eq!(
1239            state.effective_permission_mode(),
1240            SessionPermissionMode::Bypass
1241        );
1242    }
1243
1244    #[test]
1245    fn exit_plan_mode_clears_plan_mode_state() {
1246        let mut session = Session::new("sess-1", "test-model");
1247        session.agent_runtime_state = Some(AgentRuntimeState::new("run-1"));
1248        session.agent_runtime_state.as_mut().unwrap().plan_mode = Some(PlanModeState {
1249            entered_at: Utc::now(),
1250            pre_permission_mode: "default".to_string(),
1251            plan_file_path: None,
1252            status: PlanModeStatus::AwaitingApproval,
1253        });
1254        let pending = make_pending("ExitPlanMode");
1255
1256        apply_plan_mode_transition(
1257            &mut session,
1258            &pending,
1259            "Approve (Default mode)",
1260            Some("Reviewed plan".to_string()),
1261        );
1262
1263        assert!(session.agent_runtime_state.unwrap().plan_mode.is_none());
1264    }
1265
1266    #[test]
1267    fn exit_plan_mode_keeps_plan_mode_when_not_approved() {
1268        let mut session = Session::new("sess-1", "test-model");
1269        session.agent_runtime_state = Some(AgentRuntimeState::new("run-1"));
1270        session.agent_runtime_state.as_mut().unwrap().plan_mode = Some(PlanModeState {
1271            entered_at: Utc::now(),
1272            pre_permission_mode: "default".to_string(),
1273            plan_file_path: None,
1274            status: PlanModeStatus::AwaitingApproval,
1275        });
1276        let pending = make_pending("ExitPlanMode");
1277
1278        apply_plan_mode_transition(&mut session, &pending, "Stay in plan mode", None);
1279
1280        assert!(session.agent_runtime_state.unwrap().plan_mode.is_some());
1281    }
1282
1283    #[test]
1284    fn exit_plan_mode_ignores_other_tools() {
1285        let mut session = Session::new("sess-1", "test-model");
1286        let pending = make_pending("ConclusionWithOptions");
1287
1288        apply_plan_mode_transition(&mut session, &pending, "Approve", None);
1289
1290        assert!(session.agent_runtime_state.is_none());
1291    }
1292
1293    #[test]
1294    fn is_exit_plan_mode_approved_detects_approval() {
1295        assert!(is_exit_plan_mode_approved("Approve (Default mode)"));
1296        assert!(is_exit_plan_mode_approved("Approve (Accept edits mode)"));
1297        assert!(!is_exit_plan_mode_approved("Stay in plan mode"));
1298        assert!(!is_exit_plan_mode_approved("Edit plan first"));
1299    }
1300
1301    #[test]
1302    fn extract_exit_plan_from_tool_result_message_reads_plan_payload() {
1303        let mut session = Session::new("sess-1", "test-model");
1304        let mut tool_message = bamboo_agent_core::Message::tool_result(
1305            "call-1",
1306            serde_json::json!({
1307                "plan": "# Plan\n\n1. Step"
1308            })
1309            .to_string(),
1310        );
1311        tool_message.tool_success = Some(true);
1312        session.add_message(tool_message);
1313
1314        let plan = extract_exit_plan_from_tool_result_message(&session, "call-1");
1315        assert_eq!(plan.as_deref(), Some("# Plan\n\n1. Step"));
1316    }
1317
1318    #[test]
1319    fn selected_permission_preserves_typed_request_and_receipt_in_non_visible_metadata() {
1320        let mut session = Session::new("sess-1", "test-model");
1321        session.add_message(bamboo_agent_core::Message::tool_result(
1322            "permission-1",
1323            serde_json::json!({
1324                "status": "awaiting_permission_approval",
1325                "permission_request": {
1326                    "request_id": "permission-1",
1327                    "request_generation": "generation-1",
1328                    "session_id": "sess-1",
1329                    "allowed_decisions": ["allow_once", "deny_once"]
1330                }
1331            })
1332            .to_string(),
1333        ));
1334        session
1335            .messages
1336            .last_mut()
1337            .expect("permission tool result")
1338            .metadata = Some(serde_json::json!("legacy-metadata"));
1339
1340        assert!(update_or_append_tool_result_message(
1341            &mut session,
1342            "permission-1",
1343            "Approve",
1344            ResponseSource::Human,
1345        ));
1346        let receipt = PermissionDecisionReceipt {
1347            session_id: "sess-1".to_string(),
1348            decision: bamboo_tools::permission::PermissionDecision {
1349                request_id: "permission-1".to_string(),
1350                request_generation: "generation-1".to_string(),
1351                decision: bamboo_tools::permission::PermissionDecisionKind::AllowOnce,
1352                matcher_id: None,
1353                expected_policy_revision: Some(4),
1354                confirm_global: false,
1355            },
1356            decided_at: Utc::now(),
1357        };
1358        assert!(persist_permission_decision_receipt(
1359            &mut session,
1360            "permission-1",
1361            &receipt
1362        ));
1363
1364        let message = session
1365            .messages
1366            .iter()
1367            .find(|message| message.tool_call_id.as_deref() == Some("permission-1"))
1368            .expect("permission tool result");
1369        assert_eq!(message.content, "Selected response: Approve");
1370        assert_eq!(
1371            message
1372                .metadata
1373                .as_ref()
1374                .and_then(|metadata| metadata.get("permission_request"))
1375                .and_then(|request| request.get("request_id"))
1376                .and_then(serde_json::Value::as_str),
1377            Some("permission-1")
1378        );
1379        assert_eq!(
1380            message
1381                .metadata
1382                .as_ref()
1383                .and_then(|metadata| metadata.get("previous_metadata")),
1384            Some(&serde_json::json!("legacy-metadata"))
1385        );
1386        assert_eq!(
1387            message
1388                .metadata
1389                .as_ref()
1390                .and_then(|metadata| metadata.get("permission_decision_receipt"))
1391                .cloned()
1392                .and_then(|receipt| {
1393                    serde_json::from_value::<PermissionDecisionReceipt>(receipt).ok()
1394                }),
1395            Some(receipt)
1396        );
1397    }
1398
1399    #[test]
1400    fn reused_tool_call_id_requires_current_permission_generation() {
1401        let mut session = Session::new("sess-1", "test-model");
1402        for generation in ["generation-old", "generation-current"] {
1403            session.add_message(bamboo_agent_core::Message::tool_result(
1404                "permission-reused",
1405                serde_json::json!({
1406                    "status": "awaiting_permission_approval",
1407                    "permission_type": "execute_command",
1408                    "resource": format!("resource-{generation}"),
1409                    "permission_request": {
1410                        "request_id": "permission-reused",
1411                        "request_generation": generation,
1412                        "session_id": "sess-1",
1413                        "allowed_decisions": ["allow_once", "deny_once"]
1414                    }
1415                })
1416                .to_string(),
1417            ));
1418        }
1419        session.set_pending_question_with_source(
1420            "permission-reused".to_string(),
1421            "Bash".to_string(),
1422            "Allow the current operation?".to_string(),
1423            vec!["Approve".to_string(), "Deny".to_string()],
1424            false,
1425            bamboo_agent_core::PendingQuestionSource::PauseTool,
1426        );
1427        let input = RespondInput {
1428            session_id: "sess-1".to_string(),
1429            user_response: "Approve".to_string(),
1430            model: None,
1431            model_ref: None,
1432            provider: None,
1433            reasoning_effort: None,
1434        };
1435        let receipt = |generation: &str| PermissionDecisionReceipt {
1436            session_id: "sess-1".to_string(),
1437            decision: bamboo_tools::permission::PermissionDecision {
1438                request_id: "permission-reused".to_string(),
1439                request_generation: generation.to_string(),
1440                decision: bamboo_tools::permission::PermissionDecisionKind::AllowOnce,
1441                matcher_id: None,
1442                expected_policy_revision: None,
1443                confirm_global: false,
1444            },
1445            decided_at: Utc::now(),
1446        };
1447
1448        let stale = receipt("generation-old");
1449        assert!(matches!(
1450            apply_pending_response(
1451                &mut session,
1452                &input,
1453                Some("permission-reused"),
1454                ResponseSource::Human,
1455                Some(&stale),
1456            ),
1457            Err(RespondError::InvalidResponse(message))
1458                if message.contains("generation")
1459        ));
1460        assert!(session.pending_question.is_some());
1461
1462        let current = receipt("generation-current");
1463        apply_pending_response(
1464            &mut session,
1465            &input,
1466            Some("permission-reused"),
1467            ResponseSource::Human,
1468            Some(&current),
1469        )
1470        .expect("current generation resolves the parked operation");
1471        assert!(session.pending_question.is_none());
1472        assert_eq!(
1473            session
1474                .metadata
1475                .get(PERMISSION_REEXECUTE_GENERATION_METADATA_KEY)
1476                .map(String::as_str),
1477            Some("generation-current")
1478        );
1479        assert!(session.messages[0].content.contains("generation-old"));
1480        assert_eq!(session.messages[1].content, "Selected response: Approve");
1481        assert_eq!(
1482            session.messages[1]
1483                .metadata
1484                .as_ref()
1485                .and_then(|metadata| metadata.get("permission_decision_receipt"))
1486                .and_then(|receipt| receipt.get("decision"))
1487                .and_then(|decision| decision.get("request_generation"))
1488                .and_then(serde_json::Value::as_str),
1489            Some("generation-current")
1490        );
1491    }
1492
1493    #[test]
1494    fn typed_permission_receipt_controls_replay_independent_of_display_options() {
1495        let session = || {
1496            let mut session = Session::new("sess-typed", "test-model");
1497            session.add_message(bamboo_agent_core::Message::tool_result(
1498                "permission-localized",
1499                serde_json::json!({
1500                    "status": "awaiting_permission_approval",
1501                    "permission_type": "execute_command",
1502                    "resource": "cargo test --workspace",
1503                    "permission_request": {
1504                        "request_id": "permission-localized",
1505                        "request_generation": "generation-localized",
1506                        "session_id": "sess-typed",
1507                        "allowed_decisions": ["allow_once", "deny_once"]
1508                    }
1509                })
1510                .to_string(),
1511            ));
1512            session.set_pending_question_with_source(
1513                "permission-localized".to_string(),
1514                "Bash".to_string(),
1515                "允许执行吗?".to_string(),
1516                vec!["允许".to_string(), "拒绝".to_string()],
1517                false,
1518                bamboo_agent_core::PendingQuestionSource::PauseTool,
1519            );
1520            session
1521        };
1522        let receipt = |decision| PermissionDecisionReceipt {
1523            session_id: "sess-typed".to_string(),
1524            decision: bamboo_tools::permission::PermissionDecision {
1525                request_id: "permission-localized".to_string(),
1526                request_generation: "generation-localized".to_string(),
1527                decision,
1528                matcher_id: None,
1529                expected_policy_revision: None,
1530                confirm_global: false,
1531            },
1532            decided_at: Utc::now(),
1533        };
1534
1535        let mut allowed = session();
1536        let allow_input = RespondInput {
1537            session_id: "sess-typed".to_string(),
1538            user_response: "已由结构化决定允许".to_string(),
1539            model: None,
1540            model_ref: None,
1541            provider: None,
1542            reasoning_effort: None,
1543        };
1544        apply_pending_response(
1545            &mut allowed,
1546            &allow_input,
1547            Some("permission-localized"),
1548            ResponseSource::Human,
1549            Some(&receipt(PermissionDecisionKind::AllowOnce)),
1550        )
1551        .expect("typed allow must not depend on localized display options");
1552        assert_eq!(
1553            allowed
1554                .metadata
1555                .get(PERMISSION_REEXECUTE_METADATA_KEY)
1556                .map(String::as_str),
1557            Some("permission-localized")
1558        );
1559        assert_eq!(
1560            allowed
1561                .metadata
1562                .get(PERMISSION_REEXECUTE_GENERATION_METADATA_KEY)
1563                .map(String::as_str),
1564            Some("generation-localized")
1565        );
1566
1567        let mut denied = session();
1568        denied.metadata.insert(
1569            PERMISSION_REEXECUTE_METADATA_KEY.to_string(),
1570            "stale-call".to_string(),
1571        );
1572        denied.metadata.insert(
1573            PERMISSION_REEXECUTE_GENERATION_METADATA_KEY.to_string(),
1574            "stale-generation".to_string(),
1575        );
1576        let deny_input = RespondInput {
1577            session_id: "sess-typed".to_string(),
1578            // Deliberately approval-looking: the receipt enum must win.
1579            user_response: "Approve".to_string(),
1580            model: None,
1581            model_ref: None,
1582            provider: None,
1583            reasoning_effort: None,
1584        };
1585        apply_pending_response(
1586            &mut denied,
1587            &deny_input,
1588            Some("permission-localized"),
1589            ResponseSource::Human,
1590            Some(&receipt(PermissionDecisionKind::DenyOnce)),
1591        )
1592        .expect("typed deny must not depend on display text");
1593        assert!(!denied
1594            .metadata
1595            .contains_key(PERMISSION_REEXECUTE_METADATA_KEY));
1596        assert!(!denied
1597            .metadata
1598            .contains_key(PERMISSION_REEXECUTE_GENERATION_METADATA_KEY));
1599
1600        let mut legacy = session();
1601        assert!(matches!(
1602            apply_pending_response(
1603                &mut legacy,
1604                &allow_input,
1605                Some("permission-localized"),
1606                ResponseSource::Human,
1607                None,
1608            ),
1609            Err(RespondError::InvalidResponse(_))
1610        ));
1611        assert!(legacy.pending_question.is_some());
1612    }
1613}