brokk-mj-controller 2.23.2

Daemon-side controller, session manager, and web server for Mjolnir
use super::*;

/// How long creating a parent's report root on its target may take.
const REPORT_ROOT_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(60);

impl RuntimeState {
    pub(super) async fn start_create_session(
        self: &Arc<Self>,
        request: CreateSessionRequest,
    ) -> Result<RegisteredSession> {
        self.start_create_session_inner(request, CreateSessionControl::default(), None)
            .await
    }

    /// Register and start a child worker on its parent's existing target.
    pub async fn start_subagent_session(
        self: &Arc<Self>,
        request: crate::controller::RegisterSubagentRequest,
    ) -> Result<mj_core::subagent::SubagentRecord> {
        let _upgrade_work = crate::upgrade::activity("subagent admission")?;
        let (relation, needs_provisioning) = blocking(move || {
            let mut controller = Controller::load()?;
            if let Some(existing) = crate::database::lookup_subagent_request(
                &request.parent_session_id,
                &request.request_key,
            )? {
                let needs_provisioning = controller
                    .state
                    .sessions
                    .get(&existing.child_session_id)
                    .is_some_and(|record| {
                        record.state == mj_core::state::SessionState::Provisioning
                    });
                return Ok((existing, needs_provisioning));
            }
            let mut request = request;
            // The parent's report root is made before the child exists, so the
            // child's first prompt can name its own directory under it. A
            // request run again after a restart finds the same root.
            let executor = CancellableProcessExecutor::with_timeout(REPORT_ROOT_TIMEOUT);
            request.report_root = Some(
                controller
                    .prepare_subagent_report_root(&request.parent_session_id, &executor)
                    .context("create the sub-agent report directory")?,
            );
            controller
                .register_subagent(request)
                .map(|relation| (relation, true))
        })
        .await?;
        if !needs_provisioning {
            return Ok(relation);
        }
        let session_id = relation.child_session_id.clone();
        self.start_or_join_lifecycle_controlled(
            session_id.clone(),
            LifecycleKind::Create,
            None,
            Some(relation.request_key.clone()),
            None,
            move |state, session_id, cancelled| async move {
                let mut controller = tokio::task::spawn_blocking(Controller::load)
                    .await
                    .context("load controller for sub-agent startup")??;
                // The lifecycle owner now excludes competing resume/close work.
                // Recovery may already have settled this registration before
                // admission; replay never reinstalls an existing worker.
                if controller
                    .state
                    .sessions
                    .get(&session_id)
                    .is_none_or(|record| record.state != mj_core::state::SessionState::Provisioning)
                {
                    return Ok(DaemonLifecycleResult::Done);
                }
                let executor = DaemonStageReportingExecutor::new(
                    CancellableProcessExecutor::new(cancelled),
                    state,
                    session_id.clone(),
                );
                controller
                    .provision_subagent_session_controlled(&session_id, &executor)
                    .await?;
                Ok(DaemonLifecycleResult::Done)
            },
        )?;
        self.reload_controller().await?;
        Ok(relation)
    }

    pub async fn start_create_session_controlled(
        self: &Arc<Self>,
        request: CreateSessionRequest,
        control: CreateSessionControl,
        publication: tokio::sync::oneshot::Receiver<std::result::Result<(), String>>,
    ) -> Result<RegisteredSession> {
        self.start_create_session_inner(request, control, Some(publication))
            .await
    }

    pub(super) async fn start_create_session_inner(
        self: &Arc<Self>,
        request: CreateSessionRequest,
        control: CreateSessionControl,
        publication: Option<tokio::sync::oneshot::Receiver<std::result::Result<(), String>>>,
    ) -> Result<RegisteredSession> {
        let mut request = request;
        let parent = request.profile_id.clone();
        let supplied = request.subagents.clone();
        let policy = blocking(move || {
            let controller = Controller::load()?;
            let kind = controller
                .config
                .enabled_profile(&parent)
                .context("parent profile unavailable")?
                .kind;
            let policy = supplied.unwrap_or_else(|| {
                if kind.supports_delegation_tools() {
                    controller.state.last_subagent_policy.clone()
                } else {
                    Default::default()
                }
            });
            anyhow::ensure!(
                kind.supports_delegation_tools()
                    || policy == mj_core::subagent::SubagentPolicy::Native,
                "subagent policies are supported only by Claude and Codex"
            );
            Ok(policy)
        })
        .await?;
        if let mj_core::subagent::SubagentPolicy::SingleModel { model, .. } = &policy {
            let options = crate::controller::profile_config::subagent_options(
                request.profile_id.clone(),
                Some(model.clone()),
            )
            .await?;
            options.validate(&policy).map_err(anyhow::Error::msg)?;
        }
        // Discovery is restartable preparation, not admitted lifecycle work.
        let _upgrade_work = crate::upgrade::activity("session admission")?;
        request.subagents = Some(policy);
        let path_cancelled = control.cancelled.clone();
        let registered = blocking(move || {
            let mut controller = Controller::load()?;
            let path_executor = crate::targets::CancellableProcessExecutor::new(path_cancelled)
                .with_deadline(Duration::from_secs(30));
            let project_directory = request
                .project_directory
                .as_deref()
                .map(|path| {
                    controller.resolve_project_directory(
                        &request.target_template_id,
                        path,
                        &path_executor,
                    )
                })
                .transpose()?;
            let session_id = controller.register_session_with_resources(
                &request.profile_id,
                &request.bundle_id,
                &request.target_template_id,
                request.title,
                SessionLaunchOptions {
                    create_managed_worktree: request.create_managed_worktree,
                    launch_base: request.launch_base,
                    launch_branch: request.launch_branch,
                    checkout: request.checkout,
                    expected_runtime_identity: request.expected_runtime_identity,
                    subagents: request.subagents,
                    initial_prompt: request.initial_prompt,
                    workspace_id: request.workspace_id,
                    additional_mounts: request.additional_mounts,
                    resource_allocation: request.resource_allocation,
                    project_directory,
                    session_title_override: request.session_title_override,
                },
            )?;
            // Every surface creates sessions through here, so the dashboard,
            // the phone, `mj new`, and `mj acp` all leave a default pair
            // behind for the next caller that names none. A preference that
            // cannot be written does not undo a session that was created.
            if let Err(error) = mj_core::go::GoPreferences::remember_first_pair(
                &mj_core::go::GoPreferences::path(),
                &request.profile_id,
                &request.target_template_id,
            ) {
                tracing::warn!(%error, "could not save the default profile and target");
            }
            let session = controller
                .state
                .sessions
                .get(&session_id)
                .expect("newly registered session exists")
                .clone();
            let remembered_container_size = controller
                .config
                .targets
                .get(&request.target_template_id)
                .and_then(mj_core::config::container_size_host)
                .and_then(|host| {
                    controller
                        .state
                        .container_sizes
                        .get(host)
                        .copied()
                        .map(|size| (host.to_owned(), size))
                });
            Ok(RegisteredSession {
                session,
                remembered_container_size,
            })
        })
        .await?;
        let session_id = registered.session.id.clone();
        self.start_or_join_lifecycle_controlled(
            session_id,
            LifecycleKind::Create,
            None,
            None,
            Some(control.clone()),
            move |state, session_id, cancelled| async move {
                let mut controller = tokio::task::spawn_blocking(Controller::load)
                    .await
                    .context("load controller for daemon create task")??;
                let publication_error = if let Some(publication) = publication {
                    let published = tokio::select! {
                        result = publication => result.context("session publication owner stopped")
                            .and_then(|result| result.map_err(anyhow::Error::msg)),
                        () = async {
                            while !cancelled.load(Ordering::Acquire) {
                                tokio::time::sleep(Duration::from_millis(25)).await;
                            }
                        } => Err(anyhow!("session creation cancelled before publication")),
                    };
                    published.err()
                } else {
                    None
                };
                if publication_error.is_some() {
                    control.request_cancel();
                }
                let executor = DaemonStageReportingExecutor::new(
                    CancellableProcessExecutor::new(cancelled),
                    state,
                    session_id.clone(),
                );
                let provision = controller
                    .provision_session_controlled_with_commit(&session_id, &executor, || {
                        ensure!(
                            control.grant_commit(),
                            "session creation cancelled before commit"
                        );
                        Ok(())
                    })
                    .await;
                if let Some(error) = publication_error {
                    return match provision {
                        Ok(()) => Err(error),
                        Err(rollback) => {
                            Err(error.context(format!("discard unpublished session: {rollback:#}")))
                        }
                    };
                }
                provision?;
                Ok(DaemonLifecycleResult::Done)
            },
        )?;
        self.reload_controller().await?;
        Ok(registered)
    }

    pub async fn wait_create_session(&self, session_id: &str) -> Result<()> {
        let result = {
            let lifecycle_owner = self.owner();
            let lifecycle = &lifecycle_owner.lifecycle;
            let active = lifecycle
                .get(session_id)
                .with_context(|| format!("no create operation exists for session {session_id}"))?;
            ensure!(
                active.kind == LifecycleKind::Create,
                "session {session_id} is no longer being created"
            );
            active.result.clone()
        };
        let channel = result.clone();
        let outcome = Self::wait_lifecycle_result(result).await;
        self.remove_completed_lifecycle(&channel);
        match outcome? {
            DaemonLifecycleResult::Done => Ok(()),
            DaemonLifecycleResult::Move(_)
            | DaemonLifecycleResult::Park(_)
            | DaemonLifecycleResult::Superseded => {
                unreachable!("cleanup cannot return a move outcome")
            }
            DaemonLifecycleResult::DeferredCleanup => {
                unreachable!("session creation cannot schedule target cleanup")
            }
        }
    }
}

#[cfg(test)]
mod delegation_replay_tests {
    use crate::controller::test_support::{IsolatedTest, test_name};

    #[tokio::test]
    async fn replayed_spawn_does_not_reprovision_an_existing_child() {
        const CHILD: &str = "MJ_TEST_SPAWN_REPLAY";
        if std::env::var_os(CHILD).is_none() {
            let root = tempfile::tempdir().unwrap();
            IsolatedTest::new(test_name(
                module_path!(),
                "replayed_spawn_does_not_reprovision_an_existing_child",
            ))
            .env(CHILD, "1")
            .env("MJ_INSTANCE", "concurrency-sweep-spawn")
            .isolated_store(root.path())
            .run();
            return;
        }
        let _writer = crate::database::install_isolated_test_writer();
        let workspace = crate::database::create_workspace("spawn replay").unwrap();
        let parent = crate::daemon::tests::runtime_test_session(
            "parent",
            &workspace.id,
            mj_core::state::SessionState::Running,
        );
        crate::database::save_session(&parent).unwrap();
        for (index, status) in [
            mj_core::state::SessionState::Running,
            mj_core::state::SessionState::Parked,
            mj_core::state::SessionState::Stopped,
        ]
        .into_iter()
        .enumerate()
        {
            let child_id = format!("child-{index}");
            let child =
                crate::daemon::tests::runtime_test_session(&child_id, &workspace.id, status);
            let relation = crate::daemon::tests::runtime_test_subagent(&child_id, "parent");
            crate::database::save_subagent_session(&child, &relation).unwrap();
            let runtime = crate::daemon::tests::test_runtime_state();
            let replayed = runtime
                .start_subagent_session(crate::controller::RegisterSubagentRequest {
                    parent_session_id: "parent".into(),
                    task_name: "must reuse original".into(),
                    profile_id: "missing-profile".into(),
                    model: None,
                    effort: None,
                    working_directory: Default::default(),
                    initial_prompt: "must not resend".into(),
                    request_key: relation.request_key.clone(),
                    report_root: None,
                })
                .await
                .unwrap();
            assert_eq!(replayed, relation);
            assert!(runtime.active_lifecycles().is_empty());
            assert_eq!(
                crate::database::load_session_record(&child_id)
                    .unwrap()
                    .unwrap()
                    .state,
                status
            );
        }
    }
}