brokk-mj-controller 2.29.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
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
use super::*;

/// Stop the detached worker daemon at `worker_root` without deleting its files.
///
/// The script signals the worker's process group so a wedged ACP child dies
/// with it. Checkpoint then restarts the daemon against the same relay root.
pub(in crate::controller) fn stop_worker(
    owner: &crate::worker_lifecycle::WorkerPermit,
    executor: &impl CommandExecutor,
    locator: &targets::TargetLocator,
    worker_root: &str,
) -> Result<()> {
    ensure!(
        targets::worker_root(locator, owner.session_id())? == worker_root,
        "worker owner does not cover the selected root"
    );
    execute_checked(executor, stop_worker_command(locator, worker_root))?;
    Ok(())
}

/// Restore a stopped Podman target before signaling its worker. Checkpoint
/// recovery uses this instead of assuming every persisted target is running.
pub(in crate::controller) fn stop_worker_after_target_recovery(
    executor: &impl CommandExecutor,
    locator: &targets::TargetLocator,
    session_id: &str,
    worker_root: &str,
) -> Result<()> {
    let owner = crate::worker_lifecycle::require(session_id)?;
    let target = targets::target_recovery_plan(locator, session_id)?;
    targets::ensure_recovery_target_running(executor, target.as_ref())
        .context("restore Mjolnir worker target")?;
    stop_worker(&owner, executor, locator, worker_root)
}

pub(super) fn stop_worker_command(
    locator: &targets::TargetLocator,
    worker_root: &str,
) -> CommandSpec {
    let script = targets::stop_worker_daemon_script(worker_root);
    targets::locator_command(locator, vec!["sh".into(), "-c".into(), script])
        .purpose("stop Mjolnir worker daemon")
}

pub(in crate::controller) fn worker_liveness_command(
    locator: &targets::TargetLocator,
    worker_root: &str,
) -> CommandSpec {
    let script = targets::worker_daemon_liveness_script(worker_root);
    targets::locator_command(locator, vec!["sh".into(), "-c".into(), script])
        .purpose("probe Mjolnir worker daemon liveness")
}

/// Print the worker's exit record, or nothing when it left none. Reads only,
/// so it works on a target whose disk is full.
pub(in crate::controller) fn worker_exit_record_command(
    locator: &targets::TargetLocator,
    worker_root: &str,
) -> CommandSpec {
    let path = targets::join_remote_command(&[format!(
        "{worker_root}/{}",
        mj_core::relay::WORKER_EXIT_FILE
    )]);
    targets::locator_command(
        locator,
        vec![
            "sh".into(),
            "-c".into(),
            format!("cat -- {path} 2>/dev/null; exit 0"),
        ],
    )
    .purpose("read the Mjolnir worker's exit record")
}

/// The reason a worker recorded for its exit, and when (epoch seconds), from
/// the output of [`worker_exit_record_command`]. A reserved record that says
/// `null` means the worker left no reason.
pub(crate) fn recorded_exit_reason(output: &[u8]) -> Option<(String, Option<u64>)> {
    serde_json::from_slice::<Option<WorkerExitRecord>>(output)
        .ok()
        .flatten()
        .map(|record| {
            let at = record
                .at
                .as_deref()
                .and_then(|at| chrono::DateTime::parse_from_rfc3339(at).ok())
                .and_then(|at| u64::try_from(at.timestamp()).ok());
            (record.reason, at)
        })
}

pub(in crate::controller) fn start_worker(
    owner: &crate::worker_lifecycle::WorkerPermit,
    executor: &impl CommandExecutor,
    locator: &targets::TargetLocator,
    worker_root: &str,
) -> Result<()> {
    ensure!(
        targets::worker_root(locator, owner.session_id())? == worker_root,
        "worker owner does not cover the selected root"
    );
    execute_checked(executor, start_worker_command(locator, worker_root))?;
    Ok(())
}

/// Every accepted boot is protected by the durable replacement state machine.
/// This includes provisioning and resume, not just automatic upgrades.
pub(in crate::controller) fn start_worker_durably(
    owner: &crate::worker_lifecycle::WorkerPermit,
    target: &mj_core::state::TargetLocator,
    executor: &impl CommandExecutor,
    locator: &targets::TargetLocator,
    worker_root: &str,
) -> Result<()> {
    owner.verify_target(target)?;
    let id = owner.session_id();
    if let Some(intent) = crate::database::load_worker_restart(id)?
        && intent.operation_id != owner.operation_id()
    {
        let output = executor.execute(&worker_liveness_command(locator, worker_root))?;
        ensure!(
            output.status == 0 && String::from_utf8_lossy(&output.stdout).trim() == "dead",
            "another worker replacement is still booting"
        );
        owner.settle_dead_restart(target)?;
    }
    if crate::database::load_worker_restart(id)?.is_none() {
        owner.begin_restart(target, String::new())?;
    }
    let intent =
        crate::database::load_worker_restart(id)?.context("worker boot intent disappeared")?;
    ensure!(
        intent.operation_id == owner.operation_id(),
        "another worker replacement is still pending"
    );
    if crate::database::worker_restart_phase(id)?
        == Some(crate::database::WorkerRestartPhase::Prepared)
    {
        crate::database::advance_worker_restart(
            id,
            owner.operation_id(),
            crate::database::WorkerRestartPhase::Swapping,
        )?;
    }
    ensure!(
        crate::database::worker_restart_phase(id)?
            == Some(crate::database::WorkerRestartPhase::Swapping),
        "worker replacement has already launched"
    );
    start_worker(owner, executor, locator, worker_root)?;
    crate::database::advance_worker_restart(
        id,
        owner.operation_id(),
        crate::database::WorkerRestartPhase::AwaitingReadiness,
    )
}

pub(super) fn start_worker_command(
    locator: &targets::TargetLocator,
    worker_root: &str,
) -> CommandSpec {
    let binary = format!("{worker_root}/hel");
    let config = format!("{worker_root}/launch.json");
    // A launch attempt cannot touch the incumbent's diagnostics or socket.
    // The worker installs worker.log only after claiming worker.lock.
    let attempt_log = format!(
        "hel_launch_log=$(mktemp {}/worker-launch.XXXXXXXX) || exit $?; ",
        targets::posix_quote(worker_root),
    );
    let detached_script = format!(
        "{attempt_log}nohup {} >\"$hel_launch_log\" 2>&1 </dev/null &",
        targets::join_remote_command(&[
            binary.clone(),
            "worker".into(),
            "run".into(),
            "--root".into(),
            worker_root.into(),
            "--config".into(),
            config.clone(),
        ]),
    );
    // Even a loader failure is retained in the attempt-specific log.
    let exec_script = format!(
        "{attempt_log}exec {} >\"$hel_launch_log\" 2>&1",
        targets::join_remote_command(&[
            binary.clone(),
            "worker".into(),
            "run".into(),
            "--root".into(),
            worker_root.into(),
            "--config".into(),
            config.clone(),
        ]),
    );
    match locator {
        targets::TargetLocator::LocalBare { .. } => {
            // The worker this launches outlives the launch, so completion must
            // not signal the process group it inherited from this shell.
            let mut spec = CommandSpec::new("sh", ["-c", &detached_script]);
            spec.detaches = true;
            spec
        }
        targets::TargetLocator::LocalPodman { container_id, .. } => CommandSpec::new(
            "podman",
            ["exec", "--detach", container_id, "sh", "-c", &exec_script],
        ),
        targets::TargetLocator::LocalDocker { container_id, .. } => CommandSpec::new(
            "docker",
            ["exec", "--detach", container_id, "sh", "-c", &exec_script],
        ),
        targets::TargetLocator::AppleContainer { container_id, .. } => CommandSpec::new(
            "container",
            ["exec", "--detach", container_id, "sh", "-c", &exec_script],
        ),
        targets::TargetLocator::AwsEc2 { ssh, .. }
        | targets::TargetLocator::SshBare { ssh, .. } => {
            crate::targets::ssh_command(ssh, ["sh", "-c", &detached_script])
        }
        targets::TargetLocator::SshPodman {
            ssh, container_id, ..
        } => crate::targets::ssh_command(
            ssh,
            [
                "podman",
                "exec",
                "--detach",
                container_id,
                "sh",
                "-c",
                &exec_script,
            ],
        ),
        targets::TargetLocator::SshDocker {
            ssh, container_id, ..
        } => crate::targets::ssh_command(
            ssh,
            [
                "docker",
                "exec",
                "--detach",
                container_id,
                "sh",
                "-c",
                &exec_script,
            ],
        ),
    }
    .purpose("start detached Mjolnir worker")
    // Everything before this moves data into the target and reports as Sync.
    // Start begins here, with the daemon launch.
    .stage(ProvisionStage::Starting)
}

/// Enrich an opaque handshake failure by running the installed worker binary
/// directly in the target. This surfaces loader errors (for example a
/// glibc-linked worker inside an older-glibc container) that a detached start
/// swallows.
pub(in crate::controller) fn worker_probe_diagnosis(
    executor: &impl CommandExecutor,
    locator: &targets::TargetLocator,
    worker_root: &str,
    error: anyhow::Error,
) -> anyhow::Error {
    let error = match worker_binary_probe_failure(executor, locator, worker_root) {
        Some(failure) => error.context(failure),
        None => error,
    };
    match probe_worker(executor, locator, worker_root) {
        Ok(probe) => error.context(probe.to_string()),
        Err(probe_error) => {
            error.context(format!("the worker could not be probed: {probe_error:#}"))
        }
    }
}

pub(super) fn worker_binary_probe_failure(
    executor: &impl CommandExecutor,
    locator: &targets::TargetLocator,
    worker_root: &str,
) -> Option<String> {
    let binary = format!("{worker_root}/hel");
    let command = targets::locator_command(locator, vec![binary.clone(), "--version".into()])
        .purpose("probe installed worker binary");
    match executor.execute(&command) {
        Ok(output) if output.status == 0 => None,
        Ok(output) => {
            let stderr = String::from_utf8_lossy(&output.stderr);
            let stdout = String::from_utf8_lossy(&output.stdout);
            let detail = if !stderr.trim().is_empty() {
                stderr.trim()
            } else if !stdout.trim().is_empty() {
                stdout.trim()
            } else {
                "the process exited unsuccessfully without output"
            };
            Some(format!(
                "worker binary {binary} fails to run in the target: {detail}; \
                 if this is a loader/glibc error, provide a musl worker \
                 (cargo build --release --target <arch>-unknown-linux-musl \
                  -p brokk-mj-worker --bin mj-worker, \
                 or set MJ_WORKER_BINARY/MJ_WORKER_DIR)"
            ))
        }
        Err(probe_error) => Some(format!("worker probe failed: {probe_error:#}")),
    }
}

/// What one probe of a starting or dead worker found.
///
/// The probe script prints exactly this document: the worker's own startup and
/// exit records, inserted unchanged, and the live processes for its root. All
/// three facts come from one command, because on a container or SSH target
/// every probe costs a round trip.
#[derive(Debug, Clone, PartialEq, Eq, serde::Deserialize)]
pub(in crate::controller) struct WorkerProbe {
    /// `worker-startup.json`, when the worker wrote one.
    pub startup: Option<WorkerStartupRecord>,
    /// `worker-exit.json`: the worker recorded its own death.
    pub exit: Option<WorkerExitRecord>,
    /// Live worker processes for this root, the recorded one first.
    pub pids: Vec<u32>,
}

#[derive(Debug, Clone, PartialEq, Eq, serde::Deserialize)]
pub(in crate::controller) struct WorkerStartupRecord {
    /// The latest step the worker reached.
    pub step: String,
}

#[derive(Debug, Clone, PartialEq, Eq, serde::Deserialize)]
pub(in crate::controller) struct WorkerExitRecord {
    pub reason: String,
    /// A sentence the worker wrote for whoever asked, when it stopped on a
    /// precondition the caller can fix rather than on an internal failure.
    #[serde(default)]
    pub refusal: Option<String>,
    /// When the worker stopped, as RFC 3339.
    #[serde(default)]
    pub at: Option<String>,
}

impl WorkerProbe {
    pub fn alive(&self) -> bool {
        !self.pids.is_empty()
    }

    pub fn step(&self) -> Option<&str> {
        self.startup.as_ref().map(|record| record.step.as_str())
    }
}

/// One sentence naming what the worker did: its recorded exit reason, or
/// whether it is running and the step it last reached.
impl std::fmt::Display for WorkerProbe {
    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
        if let Some(exit) = &self.exit {
            return write!(formatter, "the worker exited: {}", exit.reason);
        }
        match (self.alive(), self.step()) {
            (true, Some(step)) => write!(
                formatter,
                "the worker is running; its last startup step was {step:?}"
            ),
            (true, None) => write!(
                formatter,
                "the worker is running and recorded no startup step"
            ),
            (false, Some(step)) => write!(
                formatter,
                "the worker process is gone; it reached the startup step {step:?} \
                 and left no exit record"
            ),
            (false, None) => write!(
                formatter,
                "the worker process is gone and recorded no startup step"
            ),
        }
    }
}

/// Read the worker's startup record, exit record and live processes from the
/// target. The process state distinguishes a worker that died early from one
/// that is still running but never accepted a relay connection, so it must be
/// read before the caller stops the worker.
///
/// A target that cannot be asked, or a record that does not parse, is an
/// error: the caller must not mistake it for a worker with nothing to report.
pub(in crate::controller) fn probe_worker(
    executor: &impl CommandExecutor,
    locator: &targets::TargetLocator,
    worker_root: &str,
) -> Result<WorkerProbe> {
    let command = targets::locator_command(
        locator,
        vec!["sh".into(), "-c".into(), worker_probe_script(worker_root)],
    )
    .purpose("probe worker state");
    let output = executor.execute(&command)?;
    ensure!(
        output.status == 0,
        "the worker probe exited with status {}: {}",
        output.status,
        String::from_utf8_lossy(&output.stderr).trim()
    );
    serde_json::from_slice(&output.stdout).with_context(|| {
        format!(
            "read the worker probe: {}",
            String::from_utf8_lossy(&output.stdout).trim()
        )
    })
}

/// The shell half of [`probe_worker`]. The records are JSON the worker wrote,
/// so they are inserted as they are; everything else it prints is a number.
fn worker_probe_script(worker_root: &str) -> String {
    format!(
        r#"{identity}
hel_record() {{
    if [ -s "$1" ]; then cat "$1"; else printf null; fi
}}
printf '{{"startup":'
hel_record "$hel_root/{startup_file}"
printf ',"exit":'
hel_record "$hel_root/{exit_file}"
printf ',"pids":['
if hel_pid=$(hel_recorded_worker); then
    printf '%s' "$hel_pid"
else
    hel_separator=
    while read -r hel_pid hel_args; do
        case "$hel_pid" in
            '' | *[!0-9]*) continue ;;
        esac
        [ "$hel_pid" -eq $$ ] && continue
        case "$hel_args" in
            *"$hel_match"*|*"$hel_match_home"*)
                printf '%s%s' "$hel_separator" "$hel_pid"
                hel_separator=,
                ;;
        esac
    done <<MJ_PS
$(hel_ps -eo pid=,args=)
MJ_PS
fi
printf ']}}\n'
"#,
        identity = targets::worker_daemon_identity_script(worker_root),
        startup_file = mj_core::relay::WORKER_STARTUP_FILE,
        exit_file = mj_core::relay::WORKER_EXIT_FILE,
    )
}

/// A failure on one line, for a session's error field: the first line of each
/// cause. A worker's exit reason can carry its bridge's stderr after the first
/// line; the caller logs the full chain.
pub(in crate::controller) fn failure_line(error: &anyhow::Error) -> String {
    error
        .chain()
        .map(|cause| {
            cause
                .to_string()
                .lines()
                .next()
                .unwrap_or_default()
                .trim()
                .to_owned()
        })
        .filter(|line| !line.is_empty())
        .collect::<Vec<_>>()
        .join(": ")
}

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

    /// R8-2 (cli/048): a resume whose worker exited stored the whole probe
    /// as the session's error, 96 lines long. The error keeps one line: the
    /// reason the worker recorded, then the rest of the chain.
    // Hard-won: 2fd45e3c62e5: a failed resume stored the 96-line worker probe instead of a one-line cause.
    #[test]
    fn a_failure_carrying_a_worker_exit_is_one_line_naming_the_workers_reason() {
        let probe = WorkerProbe {
            startup: Some(WorkerStartupRecord {
                step: "acp-initialized".into(),
            }),
            exit: Some(WorkerExitRecord {
                reason: "select required ACP execution mode auto: Cannot set permission mode \
                         to auto: auto mode unavailable for this model\n\
                         ACP bridge stderr:\n[session/create] phase=register durationMs=1"
                    .into(),
                refusal: None,
                at: None,
            }),
            pids: vec![],
        };
        let error = anyhow::Error::new(std::io::Error::from(std::io::ErrorKind::BrokenPipe))
            .context("write relay history_requests request")
            .context(probe.to_string());

        assert_eq!(
            failure_line(&error),
            "the worker exited: select required ACP execution mode auto: Cannot set \
             permission mode to auto: auto mode unavailable for this model: write relay \
             history_requests request: broken pipe"
        );
    }
}