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 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    // Typed permission control flow comes exclusively from the structured
466    // receipt. Display strings and localized options remain transcript-only.
467    // Legacy clarifications still validate their selected display option.
468    if permission_receipt.is_none() {
469        if let Err(error_message) = validate_pending_response(&pending, &input.user_response) {
470            // Put the pending question back when validation fails.
471            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    // Permission grants implied by approving a permission prompt. Read from the
486    // (still-unmodified) synthesized tool-result payload, BEFORE it is overwritten
487    // by the user's selection below.
488    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        // Approved a permission prompt: mark the gated tool call for re-execution
510        // on resume so the operation actually runs (real output) rather than the
511        // model inferring it. Consumed by the server resume adapter.
512        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        // A typed deny cannot inherit replay markers from an older occurrence.
528        session.metadata.remove(PERMISSION_REEXECUTE_METADATA_KEY);
529        session
530            .metadata
531            .remove(PERMISSION_REEXECUTE_GENERATION_METADATA_KEY);
532    }
533
534    // ---- Update or append tool result message ----
535    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    // ---- Plan mode state transitions ----
562    let plan_mode_transition =
563        apply_plan_mode_transition(session, &pending, &input.user_response, reviewed_plan);
564
565    // ---- Clear pending question and set resume marker ----
566    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    // ---- Merge model/reasoning from request ----
580    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
638/// Apply plan mode state transitions based on the pending question tool and user response.
639fn 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                // The typed requested mode remains live while Plan is active and
684                // may have been changed by a newer PATCH. Exiting Plan clears
685                // only the overlay; the old pre-mode is event history, never a
686                // write authority that may roll back the newer request.
687                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
703/// Check if the user response approves exiting plan mode.
704fn 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
709// ---- Internal helpers ----
710
711pub 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            // Retain invalid input for fail-closed inspection; replacing the
737            // payload must not disguise forged authority as absent legacy data.
738            if result_payload_has_supervisor_authority(message) {
739                return false;
740            }
741            // Preserve the server-issued typed permission contract outside the
742            // model-visible content before replacing the synthetic waiting
743            // payload with the selected answer. This lets an exact durable
744            // decision receipt be reconstructed after a daemon restart without
745            // trusting display strings or replaying an already-consumed run.
746            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
861/// Detect whether the user response approves a pending permission request.
862///
863/// Permission prompts (synthesized by the permission gate, and the
864/// `request_permissions` tool) offer exactly `["Approve", "Deny"]`.
865fn is_permission_approval(user_response: &str) -> bool {
866    user_response.trim().eq_ignore_ascii_case("approve")
867}
868
869/// Extract the permission grants implied by an approved permission prompt.
870///
871/// Reads the pending tool-result message (still the synthesized
872/// `awaiting_permission_approval` payload, before it is overwritten by the
873/// user's selection) and returns the `(PermissionType, resource)` pairs the
874/// caller should grant for the session. Handles both the single-gated-tool shape
875/// (top-level `permission_type` + `resource`) and the `request_permissions` shape
876/// (a `permissions` array).
877fn 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        // Simulate a newer PATCH while the Plan overlay is still active.
1243        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(&current),
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            // Deliberately approval-looking: the receipt enum must win.
1603            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        // Try mismatches before consuming the valid grant, so these assertions
1844        // cannot pass merely because the one-shot grant was already exhausted.
1845        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(&current.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        // The real preflight reads metadata A, while display migration will
1959        // preserve payload B. This reaches a failed receipt write without a
1960        // test-only writer hook or a new storage failure protocol.
1961        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}