magi-code 0.96.1

Repository-aware CLI coding agent for terminal work
Documentation
use super::*;

impl AgentSession {
    pub(super) fn build_turn_request_and_emit_context_usage(
        &self,
        state: &mut PrintTurnState<'_>,
        run: &mut AgentRunRequest<'_, '_>,
    ) -> anyhow::Result<(ProviderRequest, u64, usize, usize)> {
        let mut request = self.request_from_shared_conversation(
            Arc::clone(&state.base_conversation),
            state.turn_state.request_items_slice(),
            run.semantic_progress_timeout,
            state.prompt_cache_key.as_deref(),
            &state.tool_request_configuration,
        );
        if let Some(conversation_id) = state.conversation_id.as_deref() {
            request = request.with_conversation_id(conversation_id);
        }
        if self.provider_id == crate::providers::OPENAI_CODEX_PROVIDER {
            request.codex_turn_context = Some(state.codex_turn_context.clone());
        }
        let projection = if request.has_skill_read_provenance() {
            project_provider_request_input_tokens(&self.provider_id, &request)
        } else {
            state
                .request_projection_cache
                .project_deterministic_turn_items(state.turn_state.request_items_slice(), &[])
        };
        if self.provider_id == crate::providers::CLAUDE_SUBSCRIPTION_PROVIDER {
            // The CLI transport also checks a local estimate. Surface that
            // overflow here so the runner can compact before starting it.
            self.ensure_request_context_fits(projection, ContextBudgetPhase::Continuation)?;
        }
        let calibrated = state
            .request_projection_cache
            .apply_usage_calibration(projection);
        let context_projection =
            self.ensure_request_context_fits(calibrated, ContextBudgetPhase::Continuation)?;
        let current_request_sequence = next_request_sequence(&state.request_sequence);
        if let Some(sink) = run.output_sink.as_deref_mut() {
            sink.output_event(OutputEvent::ContextUsage {
                current_tokens: context_projection.tokens,
                max_tokens: self.context_budget.max_tokens,
                reasoning_tokens: None,
                source: context_projection.source,
                request_sequence: current_request_sequence,
            })?;
        }
        Ok((
            request,
            current_request_sequence,
            context_projection.tokens,
            projection.tokens,
        ))
    }

    #[expect(clippy::too_many_arguments)]
    pub(super) fn collect_provider_turn<P: Provider + ?Sized>(
        &self,
        provider: &P,
        request: ProviderRequest,
        state: &mut PrintTurnState<'_>,
        run: &mut AgentRunRequest<'_, '_>,
        cancellation: &AgentCancellation,
        request_sequence: u64,
        request_input_tokens: usize,
    ) -> Result<ProviderStreamResult, ProviderStreamFailure> {
        collect_provider_events(
            provider,
            request,
            ProviderEventCollector {
                cancellation,
                output_sink: &mut run.output_sink,
                output: &mut state.output,
                turn_state: &mut state.turn_state,
                assistant_chunk_batch: &mut state.assistant_chunk_batch,
                session_persistence: &mut state.session_persistence,
                context_budget_max_tokens: self.context_budget.max_tokens,
                provider_id: &self.provider_id,
                model: &self.model,
                request_sequence,
                request_input_tokens,
            },
        )
    }

    pub(super) fn project_prompt_for_preflight(
        &self,
        prompt: &str,
        session: &Session,
        tools: Option<&ToolRuntime>,
    ) -> anyhow::Result<PreflightRequestProjection> {
        let scoped_tools = self.tools_for_task_scope(Some(session), tools)?;
        let tools = scoped_tools.as_ref().or(tools);
        let initial = self.build_initial_conversation(prompt, Some(session), tools)?;
        let request = self.request_from_shared_conversation(
            Arc::from(initial.conversation),
            &[],
            None,
            None,
            &tool_request_configuration(tools),
        );
        let projection = project_provider_request_input_tokens(&self.provider_id, &request);
        Ok(PreflightRequestProjection {
            provider_id: self.provider_id.clone(),
            session_id: session.id().to_string(),
            session_path: session.path().to_path_buf(),
            replay_generation: session.replay_generation(),
            request,
            projection,
        })
    }

    pub(super) fn project_prompt_input_token_count(
        &self,
        prompt: &str,
        session: Option<&Session>,
        tools: Option<&ToolRuntime>,
    ) -> anyhow::Result<ContextTokenCount> {
        let scoped_tools = self.tools_for_task_scope(session, tools)?;
        let tools = scoped_tools.as_ref().or(tools);
        let conversation = self.build_initial_conversation(prompt, session, tools)?;
        let request = self.request_from_conversation(conversation.conversation, None, None, tools);
        Ok(project_provider_request_input_tokens(
            &self.provider_id,
            &request,
        ))
    }

    pub(crate) fn project_prompt_input_tokens(
        &self,
        prompt: &str,
        session: &Session,
        tools: Option<&ToolRuntime>,
    ) -> anyhow::Result<usize> {
        Ok(self
            .project_prompt_input_token_count(prompt, Some(session), tools)?
            .tokens)
    }

    pub(crate) fn project_prompt_input_tokens_without_session(
        &self,
        prompt: &str,
        tools: Option<&ToolRuntime>,
    ) -> anyhow::Result<usize> {
        Ok(self
            .project_prompt_input_token_count(prompt, None, tools)?
            .tokens)
    }

    pub(crate) fn ensure_prompt_context_fits(
        &self,
        prompt: &str,
        session: &Session,
        tools: Option<&ToolRuntime>,
    ) -> anyhow::Result<usize> {
        Ok(self
            .ensure_request_context_fits(
                self.project_prompt_input_token_count(prompt, Some(session), tools)?,
                ContextBudgetPhase::Initial,
            )?
            .tokens)
    }

    pub(crate) fn context_enabled(&self) -> bool {
        self.context_budget.enabled
    }

    pub(crate) fn context_threshold_tokens(&self) -> usize {
        self.context_budget.threshold_tokens()
    }

    pub(crate) fn context_max_tokens(&self) -> usize {
        self.context_budget.max_tokens
    }

    pub(super) fn task_scope_for_run(
        &self,
        session: Option<&Session>,
        tools: Option<&ToolRuntime>,
    ) -> anyhow::Result<Option<crate::sessions::task_scope::TaskScope>> {
        let memory = tools.and_then(|tools| tools.task_scope.as_ref());
        if let Some(session) = session {
            let scope = self
                .replay_cache
                .lock()
                .map_err(|_| anyhow::anyhow!("session replay cache was poisoned"))?
                .task_scope(session)?;
            anyhow::ensure!(
                scope.is_some() || memory.is_none(),
                "missing required delegated task scope; execution blocked"
            );
            Ok(scope)
        } else {
            memory
                .map(|scope| {
                    scope
                        .lock()
                        .map(|scope| scope.clone())
                        .map_err(|_| anyhow::anyhow!("task scope lock poisoned"))
                })
                .transpose()
        }
    }

    pub(super) fn tools_for_task_scope(
        &self,
        session: Option<&Session>,
        tools: Option<&ToolRuntime>,
    ) -> anyhow::Result<Option<ToolRuntime>> {
        let Some(scope) = self.task_scope_for_run(session, tools)? else {
            return Ok(None);
        };
        let tools = tools.ok_or_else(|| {
            anyhow::anyhow!("delegated task requires its capability-enforced tool runtime")
        })?;
        // Durable sessions rebuild the derived state; memory-only children share committed steering.
        if session.is_none() {
            return Ok(Some(tools.clone()));
        }
        Ok(Some(tools.clone().with_task_scope(scope)?))
    }

    pub(super) fn build_initial_conversation(
        &self,
        prompt: &str,
        session: Option<&Session>,
        tools: Option<&ToolRuntime>,
    ) -> anyhow::Result<InitialConversation> {
        let replay = self
            .replay_cache
            .lock()
            .map_err(|_| anyhow::anyhow!("session replay cache was poisoned"))?
            .replay(session)?;
        let _legacy_lossy_events = replay.legacy_lossy_events;
        let mut conversation = Vec::with_capacity(replay.items.len() + 2);
        conversation.push(ProviderConversationItem::Message(ChatMessage::system(
            self.system_prompt.clone(),
        )));
        let scope = if session.is_some() {
            replay.task_scope.map_err(anyhow::Error::msg)?
        } else {
            self.task_scope_for_run(None, tools)?
        };
        if let Some(scope) = scope {
            conversation.push(ProviderConversationItem::Message(ChatMessage::system(
                crate::sessions::task_scope::SCOPE_RULES,
            )));
            conversation.push(ProviderConversationItem::Message(ChatMessage::user(
                scope.provider_note()?,
            )));
        } else if tools.is_some_and(ToolRuntime::is_subagent)
            || session.is_some_and(|session| {
                session
                    .path()
                    .parent()
                    .and_then(|path| path.file_name())
                    .is_some_and(|name| name == "subagents")
            })
        {
            conversation.push(ProviderConversationItem::Message(ChatMessage::system(
                crate::sessions::task_scope::LEGACY_SCOPE_NOTE,
            )));
        }
        conversation.extend(replay.items);
        if self.reasoning_updates_enabled && session.is_some() {
            conversation.push(ProviderConversationItem::ReasoningSelection {
                provider: self.provider_id.clone(),
                model: self.model.clone(),
                effort: self.thinking_level,
            });
        }
        conversation.push(ProviderConversationItem::Message(ChatMessage::user(
            prompt.to_string(),
        )));
        Ok(InitialConversation {
            conversation,
            session_read_diagnostics: replay.session_read_diagnostics,
            replay_warnings: replay.warnings,
        })
    }

    pub(super) fn request_from_conversation(
        &self,
        conversation: Vec<ProviderConversationItem>,
        semantic_progress_timeout: Option<Duration>,
        prompt_cache_key: Option<&str>,
        tools: Option<&ToolRuntime>,
    ) -> ProviderRequest {
        let tool_configuration = tool_request_configuration(tools);
        self.request_from_conversation_with_configuration(
            conversation,
            semantic_progress_timeout,
            prompt_cache_key,
            &tool_configuration,
        )
    }

    pub(super) fn request_from_conversation_with_configuration(
        &self,
        conversation: Vec<ProviderConversationItem>,
        semantic_progress_timeout: Option<Duration>,
        prompt_cache_key: Option<&str>,
        tool_configuration: &ToolRequestConfiguration,
    ) -> ProviderRequest {
        self.configure_request(
            ProviderRequest::from_conversation(self.model.clone(), conversation),
            semantic_progress_timeout,
            prompt_cache_key,
            tool_configuration,
        )
    }

    pub(super) fn request_from_shared_conversation(
        &self,
        base_conversation: Arc<[ProviderConversationItem]>,
        turn_items: &[ProviderConversationItem],
        semantic_progress_timeout: Option<Duration>,
        prompt_cache_key: Option<&str>,
        tool_configuration: &ToolRequestConfiguration,
    ) -> ProviderRequest {
        self.configure_request(
            ProviderRequest::from_shared_conversation(
                self.model.clone(),
                base_conversation,
                turn_items,
            ),
            semantic_progress_timeout,
            prompt_cache_key,
            tool_configuration,
        )
    }

    pub(super) fn configure_request(
        &self,
        request: ProviderRequest,
        semantic_progress_timeout: Option<Duration>,
        prompt_cache_key: Option<&str>,
        tool_configuration: &ToolRequestConfiguration,
    ) -> ProviderRequest {
        // Stable ordering also permits exact preflight/execution request comparison.
        let mut disabled_tool_names: Vec<_> = tool_configuration
            .disabled_tool_names
            .iter()
            .cloned()
            .collect();
        disabled_tool_names.sort_unstable();
        let mut request = request
            .with_thinking_level(self.thinking_level)
            .with_text_verbosity(self.text_verbosity)
            .with_default_reasoning_summary(self.send_default_reasoning_summary)
            .with_subagents_tool_enabled(tool_configuration.subagents_tool_enabled)
            .with_disabled_tool_names(disabled_tool_names)
            .with_dynamic_tool_definitions(tool_configuration.dynamic_tool_definitions.clone());
        if let Some(catalog) = &tool_configuration.resolved_catalog {
            request = request.with_resolved_tool_catalog(catalog.clone());
        }
        if self.reasoning_updates_enabled {
            request = request.with_reasoning_updates(&self.provider_id);
        }
        if let Some(prompt_cache_key) = prompt_cache_key {
            request = request.with_prompt_cache_key(prompt_cache_key);
        }
        match semantic_progress_timeout {
            Some(timeout) => request.with_semantic_progress_timeout(timeout),
            None => request,
        }
    }
    pub(super) fn ensure_request_context_fits(
        &self,
        projection: ContextTokenCount,
        phase: ContextBudgetPhase,
    ) -> anyhow::Result<ContextTokenCount> {
        if !self.context_budget.enabled {
            return Ok(projection);
        }
        let estimated_tokens = projection.tokens;
        let threshold = self.context_budget.threshold_tokens();
        if estimated_tokens <= threshold {
            return Ok(projection);
        }
        Err(ContextBudgetError::new(
            phase,
            estimated_tokens,
            threshold,
            self.context_budget.max_tokens,
            self.context_budget.reserve_tokens,
        )
        .into())
    }
}