brokk-mj-controller 2.29.0

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

/// Derived worker facts, rebuilt on attachment. Records and configuration are
/// read afresh when retrying work after gate contention or a policy deadline.
pub(super) struct BackgroundPolicyState {
    quiet: bool,
    checkpoint_wait: Option<mj_core::activity::CheckpointWait>,
    latest_completed_turn_ordinal: Option<u64>,
    worker_build: Option<String>,
    /// Why the worker is not replaceable while the session is otherwise idle;
    /// empty otherwise. Kept to log a change once, not on every view.
    idle_replacement_blockers: Vec<&'static str>,
}

impl RuntimeState {
    fn observe_background_policy(
        &self,
        session: &SessionRecord,
        config: &Config,
        policy: &BackgroundPolicyState,
    ) {
        if policy.quiet {
            self.worker_upgrade_observer
                .observe(WorkerUpgradeObservation {
                    session: session.clone(),
                    config: config.clone(),
                    worker_build: policy.worker_build.clone(),
                    quiet: true,
                });
        }
        self.recovery_observer.observe(RecoveryObservation {
            checkpoint_wait: policy.checkpoint_wait,
            session: session.clone(),
            config: config.clone(),
            latest_completed_turn_ordinal: policy.latest_completed_turn_ordinal,
        });
    }

    /// Quiet workers need no new view to retry a skipped or delayed operation.
    /// This only queues observations; coordinators perform the actual I/O.
    pub(super) fn refresh_background_policies(&self) {
        let mut owner = self.owner();
        let records = owner.controller().state.sessions.clone();
        let config = owner.controller().config.clone();
        owner.background_policies.retain(|session_id, policy| {
            let Some(session) = records
                .get(session_id)
                .filter(|session| session.state.has_live_worker())
            else {
                return false;
            };
            self.observe_background_policy(session, &config, policy);
            true
        });
    }

    pub async fn reload_controller(&self) -> Result<()> {
        // Serialize installs so an earlier phone publication cannot overwrite
        // a later completed lifecycle with the controller snapshot it loaded.
        let _mutation = self.config_mutation.lock().await;
        if self.committed.is_some() {
            let config = tokio::task::spawn_blocking(Config::load)
                .await
                .context("daemon configuration reader panicked")??;
            self.owner().install_config(config);
        } else {
            let controller_loader = self.controller_loader;
            let controller = tokio::task::spawn_blocking(controller_loader)
                .await
                .context("daemon controller loader panicked")??;
            self.owner().install_controller(controller);
        }
        self.publish_revision();
        Ok(())
    }

    /// Definitive missing-target evidence belongs to the daemon, including
    /// when no terminal is attached. Generic connection failures stay transient.
    pub(super) fn missing_target_record(
        &self,
        session_id: &str,
        view: &ManagedSessionView,
    ) -> Option<(String, String)> {
        let Some(ViewError::TargetMissing(detail)) = &view.error else {
            return None;
        };
        if view.connected {
            return None;
        }
        let controller_owner = self.owner();
        if controller_owner
            .lifecycle
            .get(session_id)
            .is_some_and(|active| active.is_running())
        {
            return None;
        }
        let controller = controller_owner.controller();
        let session = controller.state.sessions.get(session_id)?;
        if !matches!(
            session.state,
            SessionState::Running | SessionState::Disconnected
        ) {
            return None;
        }
        Some((detail.clone(), session.updated_at.clone()))
    }

    pub(super) async fn persist_missing_target(
        self: &Arc<Self>,
        session_id: &str,
        detail: String,
        observed_updated_at: String,
    ) -> Result<()> {
        let changed = blocking({
            let session_id = session_id.to_owned();
            let detail = detail.clone();
            move || {
                crate::database::mark_session_target_missing_if_current(
                    &session_id,
                    &detail,
                    &chrono::Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Secs, true),
                    &observed_updated_at,
                )
            }
        })
        .await?;
        if changed.is_some() {
            self.reload_controller().await?;
            // A lost session has no checkpoint and no target, so its record
            // offers nothing but Destroy. Report it once and discard it
            // instead of leaving a tombstone behind.
            if changed == Some(SessionState::Lost) {
                self.discard_lost_session(session_id.to_owned()).await;
            } else {
                self.push_notice(session_id, detail);
            }
        }
        Ok(())
    }

    /// One pass of the harness readiness wait: which live sessions have run
    /// out of time to become usable.
    ///
    /// Locks are taken one at a time and nothing here touches the store, so
    /// this belongs on the daemon's background tick. The write it leads to
    /// does not; see [`Self::fail_unready_session`].
    pub(super) fn sessions_without_a_usable_harness(&self) -> Vec<UnreadySession> {
        let now = chrono::Utc::now();
        let observations = {
            let owner = self.owner();
            owner
                .indexes
                .running
                .keys()
                .filter(|id| {
                    !owner
                        .lifecycle
                        .get(*id)
                        .is_some_and(|active| active.is_running())
                })
                .filter(|id| !owner.close_requested.contains(*id))
                .filter_map(|id| owner.controller().state.sessions.get(id))
                .map(|record| {
                    let view = owner.sessions.get(&record.id);
                    ReadinessObservation {
                        harness_ready: view.is_some_and(|view| {
                            view.operational.as_ref().is_some_and(
                                mj_core::relay::RelayOperationalState::native_session_is_ready,
                            )
                        }),
                        record_age: record_age(&record.updated_at, now),
                        updated_at: record.updated_at.clone(),
                        detail: view
                            .and_then(|view| view.error.as_ref())
                            .map(|error| error.detail().to_owned()),
                        session_id: record.id.clone(),
                    }
                })
                .collect::<Vec<_>>()
        };
        self.harness_readiness
            .lock()
            .unwrap_or_else(PoisonError::into_inner)
            .observe(std::time::Instant::now(), observations)
    }

    /// Fail one session whose harness never became usable, so it stops looking
    /// like a session a driver can wait for.
    pub(super) async fn fail_unready_session(&self, unready: UnreadySession) {
        let cause = unready.cause();
        let session_id = unready.session_id.clone();
        tracing::warn!(%session_id, waited_seconds = unready.waited.as_secs(), "{cause}");
        let applied = blocking({
            let session_id = session_id.clone();
            let cause = cause.clone();
            move || {
                let mut controller = Controller::load()?;
                controller.fail_unready_session(&session_id, &cause, &unready.observed_updated_at)
            }
        })
        .await;
        match applied {
            Ok(true) => {
                if let Err(error) = self.reload_controller().await {
                    tracing::warn!(%session_id, error = format!("{error:#}"), "could not reload state after failing an unready session");
                }
                self.push_notice(&session_id, cause);
                self.publish_revision();
            }
            Ok(false) => {}
            Err(error) => tracing::warn!(
                %session_id,
                error = format!("{error:#}"),
                "could not record that a session's harness never became usable"
            ),
        }
    }

    /// Persist a worker-reported preparation failure against the session
    /// revision observed by the upgrade coordinator. Lifecycle work that has
    /// since changed that revision remains the owner of the outcome.
    pub(super) async fn fail_harness_preparation(
        &self,
        session_id: String,
        failure: crate::controller::HarnessPreparationFailure,
        observed_updated_at: String,
    ) {
        let cause = failure.to_string();
        let applied = blocking({
            let session_id = session_id.clone();
            let cause = cause.clone();
            move || {
                let mut controller = Controller::load()?;
                controller.fail_unready_session(&session_id, &cause, &observed_updated_at)
            }
        })
        .await;
        match applied {
            Ok(true) => {
                if let Err(error) = self.reload_controller().await {
                    tracing::warn!(%session_id, error = format!("{error:#}"), "could not reload state after harness preparation failed");
                }
                self.push_notice(&session_id, cause);
                self.publish_revision();
            }
            Ok(false) => {}
            Err(error) => tracing::warn!(
                %session_id,
                error = format!("{error:#}"),
                "could not record worker-reported harness preparation failure"
            ),
        }
    }

    pub(super) async fn publish_session(
        &self,
        session_id: String,
        view: ManagedSessionView,
    ) -> Result<()> {
        let connected = view.connected;
        let has_snapshot = view.snapshot.is_some();
        tracing::debug!(
            %session_id,
            connected,
            has_snapshot,
            "daemon received a session view"
        );
        {
            let mut owner = self.owner();
            let controller = owner.controller();
            if !controller.state.sessions.contains_key(&session_id)
                || !owner.runs_relay_actor(&session_id)
            {
                return Ok(());
            }
            if view.connected
                && let Some(snapshot) = view.snapshot.as_ref()
                && let Some(session) = controller.state.sessions.get(&session_id)
                && session.state.has_live_worker()
            {
                let quiet = snapshot.operational.quiet();
                if !quiet.is_yes() {
                    tracing::debug!(
                        session_id = %session_id,
                        reason = quiet.reason(),
                        "session is not quiet; upgrade, checkpoint and move wait"
                    );
                }
                let facts = snapshot.operational.facts();
                // A busy session is not replaced, and that needs no comment.
                // An idle one that is still not replaceable would otherwise
                // be skipped without a trace (B-5).
                let idle_replacement_blockers = if mj_core::activity::can_submit(&facts)
                    && !mj_core::activity::safe_to_replace(&facts, session.harness_kind)
                {
                    mj_core::activity::replacement_blockers(&facts, session.harness_kind)
                } else {
                    Vec::new()
                };
                if !idle_replacement_blockers.is_empty()
                    && owner
                        .background_policies
                        .get(&session_id)
                        .is_none_or(|previous| {
                            previous.idle_replacement_blockers != idle_replacement_blockers
                        })
                {
                    tracing::info!(
                        %session_id,
                        harness = ?session.harness_kind,
                        worker_build = ?snapshot.worker_build,
                        reasons = ?idle_replacement_blockers,
                        "an idle worker is not replaced"
                    );
                }
                let policy = BackgroundPolicyState {
                    idle_replacement_blockers,
                    quiet: snapshot.operational.safe_to_replace(session.harness_kind),
                    checkpoint_wait: snapshot
                        .operational
                        .routine_checkpoint_wait(session.harness_kind),
                    latest_completed_turn_ordinal: snapshot.latest_completed_turn_ordinal(),
                    worker_build: snapshot.worker_build.clone(),
                };
                self.observe_background_policy(session, &controller.config, &policy);
                owner.background_policies.insert(session_id.clone(), policy);
            } else {
                owner.background_policies.remove(&session_id);
            }
            owner.publish_view(session_id, view);
        }
        reach_test_hook("relay_projection_before_revision_publication").await?;
        self.publish_revision();
        Ok(())
    }
}

/// How long ago a durable record was written, when its timestamp can be read.
///
/// A record written in the future — a clock that moved backwards, a store
/// written by another host — reports no age rather than a wrapped one, so the
/// readiness wait leaves those sessions alone.
fn record_age(updated_at: &str, now: chrono::DateTime<chrono::Utc>) -> Option<Duration> {
    let written = chrono::DateTime::parse_from_rfc3339(updated_at).ok()?;
    now.signed_duration_since(written.with_timezone(&chrono::Utc))
        .to_std()
        .ok()
}