brokk-mj-controller 2.24.0

Daemon-side controller, session manager, and web server for Mjolnir
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
use super::*;
use mj_client::runtime_feed::RuntimeProjection;
use mj_core::snapshot_map::SnapshotMap;

pub struct RemoteDashboardWorkerPoller {
    pub updates: SessionManagerUpdates,
    pub control: SessionManagerControl,
    pub shutdown: SessionManagerShutdown,
    pub state: tokio::sync::watch::Receiver<RuntimeStateUpdate>,
    /// Reviews the daemon is running for this workspace's sessions.
    pub reviews: tokio::sync::watch::Receiver<Vec<crate::review_host::RuntimeReviewView>>,
    /// Background events the daemon wants reported once, oldest first.
    pub notices: tokio::sync::watch::Receiver<Vec<daemon::RuntimeNotice>>,
    /// What the daemon knows about quota. The daemon is the only prober.
    pub quotas: tokio::sync::watch::Receiver<mj_client::quota::QuotaSnapshot>,
    pub config: tokio::sync::watch::Receiver<mj_core::config::Config>,
    pub health: tokio::sync::watch::Receiver<RuntimeFeedHealth>,
}

/// Records and lifecycle ownership must reach the surface in the same frame.
#[derive(Debug, Clone, Default)]
pub struct RuntimeStateUpdate {
    pub last_subagent_policy: mj_core::subagent::SubagentPolicy,
    pub native_agents: SnapshotMap<String, mj_core::native_agent::NativeAgentView>,
    pub workspace_names: std::collections::BTreeMap<String, String>,
    pub revision: u64,
    pub records: SnapshotMap<String, SessionRecord>,
    pub lifecycles: Vec<daemon::RuntimeLifecycleView>,
    pub moves: SnapshotMap<String, mj_core::state::MoveOperation>,
    /// Parent/child relations for the sessions in `records`, so a surface can
    /// keep a daemon-created child out of the real workspace without a full
    /// state reload.
    pub subagents: SnapshotMap<String, SubagentRecord>,
}

/// What a session looked like the last time a view was published for it.
///
/// The poller compares this before reading anything, so a session that has not
/// moved costs one comparison rather than a full transcript load. Nothing here
/// grows with the transcript: the projection is identified by its ordinal and
/// digest, and the operational state is bounded by the relay's own command and
/// configuration surface.
#[derive(Debug, Clone, PartialEq)]
pub(super) struct PublishedView {
    pub(super) projection_ordinal: u64,
    pub(super) projection_digest: String,
    pub(super) operational: Option<mj_core::relay::RelayOperationalState>,
    pub(super) connected: bool,
    pub(super) error: Option<String>,
}

impl PublishedView {
    pub(super) fn of(runtime: &crate::daemon::RuntimeSessionView) -> Self {
        Self {
            projection_ordinal: runtime.projection_ordinal,
            projection_digest: runtime.projection_digest.clone(),
            operational: runtime.operational.clone(),
            connected: runtime.connected,
            error: runtime.error.as_ref().map(|error| format!("{error:?}")),
        }
    }

    pub(super) fn matches(&self, runtime: &crate::daemon::RuntimeSessionView) -> bool {
        *self == Self::of(runtime)
    }
}

pub(super) const PROJECTION_CONVERGENCE_RETRIES: u8 = 20;
pub(super) const PROJECTION_CONVERGENCE_RETRY_DELAY: Duration = Duration::from_millis(50);

#[derive(Debug, Clone, PartialEq, Eq)]
pub(super) struct ProjectionMismatch {
    pub(super) published_ordinal: u64,
    pub(super) published_digest: String,
    pub(super) durable_ordinal: u64,
    pub(super) durable_digest: String,
}

#[derive(Default)]
pub(super) struct ProjectionConvergence {
    pub(super) attempts: std::collections::BTreeMap<String, (ProjectionMismatch, u8)>,
}

impl ProjectionConvergence {
    pub(super) fn converged(&mut self, session_id: &str) {
        self.attempts.remove(session_id);
    }

    /// Give a lifecycle rollback and the daemon's cached relay view a bounded
    /// window to converge. Repeating the same mismatch eventually reports the
    /// integrity failure instead of hiding it indefinitely.
    pub(super) fn should_retry(&mut self, session_id: &str, mismatch: ProjectionMismatch) -> bool {
        let entry = self
            .attempts
            .entry(session_id.to_owned())
            .or_insert_with(|| (mismatch.clone(), 0));
        if entry.0 != mismatch {
            *entry = (mismatch, 0);
        }
        entry.1 = entry.1.saturating_add(1);
        entry.1 <= PROJECTION_CONVERGENCE_RETRIES
    }
}

/// Read-only updates shared by the dashboard and workspace preview. A snapshot
/// precedes its session views, so consumers can establish membership first.
pub enum RuntimeFeedUpdate {
    Snapshot(Box<RuntimeProjection>),
    Session {
        session_id: String,
        view: Box<ManagedSessionView>,
    },
    /// The daemon no longer runs an actor for this session, so it has no view.
    SessionRemoved(String),
    Error(String),
}

/// Dropping a subscription cancels even a pending daemon long poll. The task
/// owns no writer or relay connection; blocking projection reads are bounded.
pub struct RuntimeFeed {
    pub updates: tokio::sync::mpsc::Receiver<RuntimeFeedUpdate>,
    pub(super) task: tokio::task::JoinHandle<()>,
}

impl Drop for RuntimeFeed {
    fn drop(&mut self) {
        self.task.abort();
    }
}

pub(super) type StoredProjection = Option<(MaterializedSession, mj_core::state::ProjectionWindow)>;

pub(super) static PROJECTION_READERS: std::sync::LazyLock<Arc<tokio::sync::Semaphore>> =
    std::sync::LazyLock::new(|| Arc::new(tokio::sync::Semaphore::new(4)));

pub(super) async fn load_runtime_projection(session_id: String) -> Result<StoredProjection> {
    let permit = Arc::clone(&PROJECTION_READERS)
        .acquire_owned()
        .await
        .context("projection readers stopped")?;
    tokio::task::spawn_blocking(move || {
        let _permit = permit;
        let started = Instant::now();
        let result = crate::database::load_materialized_projection_tail(
            &session_id,
            crate::database::PROJECTION_TAIL_ITEMS,
        );
        tracing::debug!(target: "mj_controller::latency", %session_id, elapsed_ms = started.elapsed().as_secs_f64() * 1000.0, "terminal projection loaded");
        // A blocking SQLite read can outlive cancellation of its subscriber.
        if let Err(error) = &result {
            tracing::warn!(%session_id, %error, "could not load runtime projection");
        }
        result
    })
    .await
    .context("projection load task failed")?
}

pub(super) fn spawn_runtime_feed_with<P, PF, L, LF, S>(
    workspace_id: String,
    poll: P,
    load: L,
) -> RuntimeFeed
where
    P: Fn(String, u64) -> PF + Send + 'static,
    PF: Future<Output = Result<S>> + Send,
    S: Into<RuntimeProjection> + Send + 'static,
    L: Fn(String) -> LF + Clone + Send + 'static,
    LF: Future<Output = Result<StoredProjection>> + Send + 'static,
{
    let (tx, updates) = tokio::sync::mpsc::channel(32);
    let task = tokio::spawn(async move {
        let result = run_runtime_feed(workspace_id, poll, load, &tx).await;
        if let Err(error) = result {
            let message = format!("Runtime feed stopped: {error:#}");
            tracing::error!(%message);
            let _ = tx.send(RuntimeFeedUpdate::Error(message)).await;
        }
    });
    RuntimeFeed { updates, task }
}

/// A stopped daemon is a state the person can act on, not an I/O error.
fn refresh_failure_notice(error: &anyhow::Error) -> String {
    if daemon::daemon_not_running(error).is_some() {
        "Mjolnir daemon is stopped; sessions refresh when it starts again.".to_owned()
    } else {
        format!("Could not refresh sessions: {error:#}")
    }
}

pub(super) async fn run_runtime_feed<P, PF, L, LF, S>(
    workspace_id: String,
    poll: P,
    load: L,
    tx: &tokio::sync::mpsc::Sender<RuntimeFeedUpdate>,
) -> Result<()>
where
    P: Fn(String, u64) -> PF,
    PF: Future<Output = Result<S>>,
    S: Into<RuntimeProjection>,
    L: Fn(String) -> LF + Clone + Send + 'static,
    LF: Future<Output = Result<StoredProjection>> + Send + 'static,
{
    let mut revision = 0;
    let mut convergence = ProjectionConvergence::default();
    let mut published = std::collections::BTreeMap::<String, PublishedView>::new();
    let mut previous_sessions = SnapshotMap::new();
    let mut pending_ids = std::collections::BTreeSet::new();
    loop {
        let mut snapshot = match poll(workspace_id.clone(), revision).await {
            Ok(snapshot) => snapshot.into(),
            Err(error) => {
                if tx
                    .send(RuntimeFeedUpdate::Error(refresh_failure_notice(&error)))
                    .await
                    .is_err()
                {
                    return Ok(());
                }
                tokio::time::sleep(Duration::from_millis(250)).await;
                continue;
            }
        };
        let snapshot_revision = snapshot.revision;
        let sessions = std::mem::take(&mut snapshot.sessions);
        let mut removed = Vec::new();
        for (id, runtime) in previous_sessions.changes(&sessions) {
            match runtime {
                Some(runtime) if !published.get(id).is_some_and(|last| last.matches(runtime)) => {
                    pending_ids.insert(id.clone());
                }
                None => {
                    published.remove(id);
                    pending_ids.remove(id);
                    convergence.attempts.remove(id);
                    removed.push(id.clone());
                }
                Some(_) => {}
            }
        }
        previous_sessions = sessions.clone();
        if tx
            .send(RuntimeFeedUpdate::Snapshot(Box::new(snapshot)))
            .await
            .is_err()
        {
            return Ok(());
        }
        // Forgetting a view here is announced, so a consumer never keeps a
        // copy this loop will not send again.
        for session_id in removed {
            if tx
                .send(RuntimeFeedUpdate::SessionRemoved(session_id))
                .await
                .is_err()
            {
                return Ok(());
            }
        }
        let mut pending = pending_ids
            .iter()
            .filter_map(|id| sessions.get(id).cloned())
            .collect::<std::collections::VecDeque<_>>();
        let mut tasks = tokio::task::JoinSet::new();
        let mut retry = false;
        while !pending.is_empty() || !tasks.is_empty() {
            // Independent session reads overlap, without flooding SQLite or
            // leaving an unbounded number of blocking reads after cancellation.
            while tasks.len() < 4 {
                let Some(runtime) = pending.pop_front() else {
                    break;
                };
                let load = load.clone();
                tasks.spawn(async move {
                    let stored = if runtime.operational.is_some() {
                        load(runtime.session_id.clone()).await
                    } else {
                        Ok(None)
                    };
                    (runtime, stored)
                });
            }
            let Some(result) = tasks.join_next().await else {
                break;
            };
            let (runtime, stored) = result.context("join runtime projection reader")?;
            let session_id = runtime.session_id.clone();
            let fingerprint = PublishedView::of(&runtime);
            let expects_projection = runtime.operational.is_some();
            let Some(view) = runtime_projection_view(runtime, stored, &mut convergence) else {
                retry = true;
                continue;
            };
            if view.snapshot.is_some() || !expects_projection {
                published.insert(session_id.clone(), fingerprint);
                pending_ids.remove(&session_id);
            } else {
                published.remove(&session_id);
            }
            if tx
                .send(RuntimeFeedUpdate::Session {
                    session_id,
                    view: Box::new(view),
                })
                .await
                .is_err()
            {
                return Ok(());
            }
        }
        if retry {
            tokio::time::sleep(PROJECTION_CONVERGENCE_RETRY_DELAY).await;
        } else {
            revision = revision.max(snapshot_revision);
        }
    }
}

pub(super) fn runtime_projection_view(
    runtime: daemon::RuntimeSessionView,
    stored: Result<StoredProjection>,
    convergence: &mut ProjectionConvergence,
) -> Option<ManagedSessionView> {
    let Some(operational) = runtime.operational else {
        return Some(ManagedSessionView {
            snapshot: None,
            connected: runtime.connected,
            error: runtime.error,
        });
    };
    let detail = match stored {
        Ok(Some((materialized, window)))
            if materialized.applied_event_ordinal > runtime.projection_ordinal
                || (materialized.applied_event_ordinal == runtime.projection_ordinal
                    && materialized.applied_event_digest == runtime.projection_digest) =>
        {
            convergence.converged(&runtime.session_id);
            return Some(ManagedSessionView {
                snapshot: Some(ManagedSessionSnapshot {
                    materialized,
                    window,
                    operational,
                    latest_credential_sync_signal: runtime.latest_credential_sync_signal,
                    worker_build: None,
                    subagent_requests: Vec::new(),
                    subagent_results: Vec::new(),
                }),
                connected: runtime.connected,
                error: runtime.error,
            });
        }
        Ok(Some((materialized, _))) => {
            let mismatch = ProjectionMismatch {
                published_ordinal: runtime.projection_ordinal,
                published_digest: runtime.projection_digest,
                durable_ordinal: materialized.applied_event_ordinal,
                durable_digest: materialized.applied_event_digest.clone(),
            };
            if convergence.should_retry(&runtime.session_id, mismatch) {
                return None;
            }
            if materialized.applied_event_ordinal < runtime.projection_ordinal {
                format!(
                    "daemon published projection {} but SQLite contains only {} after a bounded convergence retry",
                    runtime.projection_ordinal, materialized.applied_event_ordinal
                )
            } else {
                format!(
                    "daemon and SQLite projection digests differ at ordinal {} after a bounded convergence retry",
                    runtime.projection_ordinal
                )
            }
        }
        Ok(None) => "daemon published a session with no durable projection".into(),
        Err(error) => format!("load daemon-owned projection: {error:#}"),
    };
    Some(ManagedSessionView {
        snapshot: None,
        connected: false,
        error: Some(ViewError::ProjectionIntegrity(detail)),
    })
}

#[cfg(test)]
mod refresh_notice_tests {
    use super::*;

    #[test]
    fn a_stopped_daemon_is_reported_as_stopped_not_as_a_file_error() {
        let stopped = anyhow::Error::new(daemon::DaemonNotRunning {
            metadata_path: "/data/daemon.json".into(),
        });
        assert_eq!(
            refresh_failure_notice(&stopped),
            "Mjolnir daemon is stopped; sessions refresh when it starts again."
        );
        assert_eq!(
            refresh_failure_notice(&anyhow::anyhow!("unexpected end of file")),
            "Could not refresh sessions: unexpected end of file"
        );
    }
}