brokk-mj-controller 2.31.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
use super::*;

pub fn spawn_dashboard_capacity_poller() -> (
    tokio::sync::watch::Sender<Vec<DeploymentCapacityTarget>>,
    tokio::sync::mpsc::Sender<()>,
    tokio::sync::mpsc::Receiver<CapacityPollUpdate>,
) {
    spawn_capacity_poller_with(|target| async move {
        if let Some(error) = &target.probe_error {
            bail!("capacity probe is unavailable: {error}");
        }
        if target.local {
            let mut usage = collect_local_capacity_with(collect_local_capacity).await?;
            usage.storage = collect_local_storage(&target).await;
            return Ok(Some(usage));
        }
        tokio::time::timeout(RESOURCE_POLL_TIMEOUT, collect_capacity(&target))
            .await
            .context("capacity probe timed out")?
    })
}

pub(super) fn spawn_capacity_poller_with<F, Fut>(
    collect: F,
) -> (
    tokio::sync::watch::Sender<Vec<DeploymentCapacityTarget>>,
    tokio::sync::mpsc::Sender<()>,
    tokio::sync::mpsc::Receiver<CapacityPollUpdate>,
)
where
    F: Fn(DeploymentCapacityTarget) -> Fut + Send + Sync + 'static,
    Fut: Future<Output = Result<Option<DeploymentCapacityUsage>>> + Send + 'static,
{
    let (targets_tx, mut targets_rx) =
        tokio::sync::watch::channel(Vec::<DeploymentCapacityTarget>::new());
    let (updates_tx, updates_rx) = tokio::sync::mpsc::channel(64);
    let (triggers_tx, mut triggers_rx) = tokio::sync::mpsc::channel(1);
    tokio::spawn(async move {
        let mut targets = Vec::new();
        let collect = Arc::new(collect);
        let mut samples = CapacitySamples::default();
        let mut interval = tokio::time::interval_at(
            tokio::time::Instant::now() + CAPACITY_POLL_INTERVAL,
            CAPACITY_POLL_INTERVAL,
        );
        interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
        loop {
            tokio::select! {
                _ = updates_tx.closed() => break,
                _ = interval.tick() => {
                    samples.schedule(targets.iter().cloned(), &collect);
                }
                changed = targets_rx.changed() => {
                    if changed.is_err() {
                        tracing::debug!("capacity poll target feed closed; stopping capacity poller");
                        break;
                    }
                    let updated = targets_rx.borrow_and_update().clone();
                    samples.schedule(
                        updated.iter().filter(|target| !targets.contains(target)).cloned(),
                        &collect,
                    );
                    targets = updated;
                }
                trigger = triggers_rx.recv() => {
                    if trigger.is_none() {
                        break;
                    }
                    samples.schedule(targets.iter().cloned(), &collect);
                }
                completed = samples.tasks.join_next_with_id(), if !samples.tasks.is_empty() => {
                    let (id, result) = match completed.expect("capacity task exists") {
                        Ok((id, result)) => (id, result.map_err(|error| format!("{error:#}"))),
                        Err(error) => (error.id(), Err(format!("capacity probe task failed: {error}"))),
                    };
                    let sampled = samples.targets.remove(&id).expect("capacity task retains its target");
                    if let Err(error) = &result {
                        tracing::warn!(target_id = %sampled.id, %error, "capacity probe failed");
                    }
                    let Ok(permit) = updates_tx.reserve().await else {
                        break;
                    };
                    // A watch update and completion can become ready together.
                    // Revalidate after backpressure, with no await between
                    // reading the latest target and publishing the result.
                    let current = targets_rx.borrow().iter().find(|target| target.id == sampled.id).cloned();
                    let Some(current) = current else {
                        continue;
                    };
                    if current != sampled {
                        // A changed target gets one follow-up; its old result
                        // must not overwrite a reading for the new configuration.
                        // If changed() is still pending, that arm will start it.
                        if targets.contains(&current) {
                            samples.schedule(std::iter::once(current), &collect);
                        }
                        continue;
                    }
                    permit.send(CapacityPollUpdate {
                        target_id: sampled.id,
                        result,
                        sampled_at_epoch_seconds: epoch_seconds(),
                    });
                }
            }
        }
        samples.tasks.abort_all();
        while let Some(completed) = samples.tasks.join_next().await {
            match completed {
                Ok(Err(error)) => tracing::warn!(%error, "capacity probe failed during shutdown"),
                Err(error) if !error.is_cancelled() => {
                    tracing::error!(%error, "capacity probe task failed during shutdown");
                }
                _ => {}
            }
        }
    });
    (targets_tx, triggers_tx, updates_rx)
}

#[derive(Default)]
pub(super) struct CapacitySamples {
    pub(super) tasks: tokio::task::JoinSet<Result<Option<DeploymentCapacityUsage>>>,
    pub(super) targets: std::collections::HashMap<tokio::task::Id, DeploymentCapacityTarget>,
}

impl CapacitySamples {
    pub(super) fn schedule<F, Fut>(
        &mut self,
        targets: impl IntoIterator<Item = DeploymentCapacityTarget>,
        collect: &Arc<F>,
    ) where
        F: Fn(DeploymentCapacityTarget) -> Fut + Send + Sync + 'static,
        Fut: Future<Output = Result<Option<DeploymentCapacityUsage>>> + Send + 'static,
    {
        for target in targets {
            if self.targets.values().any(|running| running.id == target.id) {
                continue;
            }
            let collect = collect.clone();
            let sampled = target.clone();
            let task = self.tasks.spawn(async move {
                let started = Instant::now();
                let target_id = sampled.id.clone();
                let result = collect(sampled).await;
                tracing::debug!(
                    %target_id,
                    elapsed_ms = started.elapsed().as_millis() as u64,
                    success = result.is_ok(),
                    "capacity probe completed",
                );
                result
            });
            self.targets.insert(task.id(), target);
        }
    }
}

pub(super) async fn collect_capacity(
    target: &DeploymentCapacityTarget,
) -> Result<Option<DeploymentCapacityUsage>> {
    if let Some(error) = &target.probe_error {
        anyhow::bail!("capacity probe is unavailable: {error}");
    }
    match target.kind {
        DeploymentCapacityKind::Host => {
            let mut last_error = None;
            for command in &target.probes {
                match execute_resource_command(command).await {
                    Ok(output) => {
                        return crate::targets::parse_host_capacity(
                            &output.stdout,
                            &probe_storage_host(command),
                        )
                        .map(Some);
                    }
                    Err(error) => last_error = Some(error),
                }
            }
            Err(last_error.unwrap_or_else(|| anyhow::anyhow!("no host probe is configured")))
        }
        DeploymentCapacityKind::AwsFleet => {
            if target.probes.is_empty() {
                return Ok(None);
            }
            let mut tasks = tokio::task::JoinSet::new();
            for command in target.probes.clone() {
                tasks.spawn(async move {
                    let output = execute_resource_command(&command).await?;
                    crate::targets::parse_aws_allocated_capacity(
                        &output.stdout,
                        &probe_storage_host(&command),
                    )
                });
            }
            let mut usages = Vec::new();
            while let Some(result) = tasks.join_next().await {
                usages.push(result.context("join EC2 capacity probe")??);
            }
            aggregate_aws_capacity(&usages).map(Some)
        }
    }
}

pub fn aggregate_aws_capacity(
    usages: &[DeploymentCapacityUsage],
) -> Result<DeploymentCapacityUsage> {
    let mut total = DeploymentCapacityUsage {
        cpu_percent: None,
        memory_used_bytes: 0,
        memory_total_bytes: 0,
        logical_cores: 0,
        disk_total_bytes: Some(0),
        storage: Vec::new(),
    };
    for usage in usages {
        total.storage.extend(usage.storage.iter().cloned());
        total.memory_total_bytes = total
            .memory_total_bytes
            .checked_add(usage.memory_total_bytes)
            .context("aggregate EC2 RAM overflow")?;
        total.logical_cores = total
            .logical_cores
            .checked_add(usage.logical_cores)
            .context("aggregate EC2 core count overflow")?;
        total.disk_total_bytes = Some(
            total
                .disk_total_bytes
                .unwrap_or(0)
                .checked_add(usage.disk_total_bytes.unwrap_or(0))
                .context("aggregate EC2 disk overflow")?,
        );
    }
    Ok(total)
}

pub(super) fn collect_local_capacity() -> Result<DeploymentCapacityUsage> {
    let mut system = sysinfo::System::new();
    system.refresh_memory();
    // Frequency is unused and scans every core in parallel on each refresh.
    system.refresh_cpu_usage();
    std::thread::sleep(sysinfo::MINIMUM_CPU_UPDATE_INTERVAL);
    system.refresh_cpu_usage();
    Ok(DeploymentCapacityUsage {
        cpu_percent: Some(system.global_cpu_usage().round().clamp(0.0, 100.0) as u8),
        memory_used_bytes: system
            .total_memory()
            .saturating_sub(system.available_memory()),
        memory_total_bytes: system.total_memory(),
        logical_cores: system
            .cpus()
            .len()
            .try_into()
            .context("logical CPU count overflow")?,
        disk_total_bytes: None,
        storage: Vec::new(),
    })
}

/// The storage owner's name for the machine a probe ran on.
fn probe_storage_host(command: &CommandSpec) -> String {
    mj_core::targets::storage::storage_host_of_destination(command.ssh_destination.as_deref())
}

/// Measure local free space over the local target's storage paths. A failed
/// measurement leaves storage unknown; CPU and memory still publish.
async fn collect_local_storage(
    target: &DeploymentCapacityTarget,
) -> Vec<mj_core::targets::storage::HostStorageSample> {
    let paths = target.local_storage_paths.clone();
    let measured = tokio::task::spawn_blocking(move || {
        crate::targets::measure_local_storage(
            &paths,
            &crate::targets::BoundedProcessExecutor::new(RESOURCE_POLL_TIMEOUT),
        )
    })
    .await
    .context("join the local free-space probe")
    .and_then(|measured| measured);
    match measured {
        Ok((_, filesystems)) if filesystems.is_empty() => Vec::new(),
        Ok((home, filesystems)) => vec![mj_core::targets::storage::HostStorageSample {
            host: mj_core::targets::storage::LOCAL_STORAGE_HOST.to_owned(),
            home,
            filesystems,
        }],
        Err(error) => {
            tracing::warn!(
                error = format!("{error:#}"),
                "local free-space probe failed"
            );
            Vec::new()
        }
    }
}

pub(super) async fn collect_local_capacity_with(
    collect: impl FnOnce() -> Result<DeploymentCapacityUsage> + Send + 'static,
) -> Result<DeploymentCapacityUsage> {
    // A blocking sample cannot be cancelled. Keep its slot occupied
    // until it exits, even when the deadline has elapsed.
    let mut sample = tokio::task::spawn_blocking(move || {
        let result = collect();
        // Shutdown can drop the awaiting future before this thread exits.
        if let Err(error) = &result {
            tracing::warn!(%error, "local capacity sample failed");
        }
        result
    });
    match tokio::time::timeout(RESOURCE_POLL_TIMEOUT, &mut sample).await {
        Ok(result) => result.context("join local capacity probe")?,
        Err(_) => {
            match sample.await {
                Ok(Ok(_)) => {}
                Ok(Err(error)) => tracing::warn!(%error, "timed-out capacity probe failed"),
                Err(error) => tracing::error!(%error, "timed-out capacity probe task failed"),
            }
            bail!("capacity probe timed out")
        }
    }
}

/// Run one resource probe. A probe the SSH server turned away never ran, so
/// it is retried with the same backoff the executor and the relay use (launch
/// finding R3-5); any other failure is reported at once.
pub(super) async fn execute_resource_command(command: &CommandSpec) -> Result<CommandOutput> {
    use crate::targets::{SSH_RETRY_ATTEMPTS, SshRefusal, ssh_refusal, ssh_retry_delay};
    let attempts = match command.ssh_destination {
        Some(_) => SSH_RETRY_ATTEMPTS,
        None => 1,
    };
    for attempt in 1..=attempts {
        let (output, ssh_session) = run_resource_command(command).await?;
        if output.status == 0 {
            return Ok(output);
        }
        let stderr = String::from_utf8_lossy(&output.stderr);
        if let Some(destination) = command.ssh_destination.as_deref()
            && let Some(refusal) = ssh_refusal(output.status, &stderr)
        {
            // A refused session found its master alive; a connection closed
            // before authentication may mean the master is gone.
            if refusal == SshRefusal::BeforeAuthentication
                && let Some(lease) = &ssh_session
            {
                lease.invalidate();
            }
            // Only Unix leases own a shared session slot; free it before backoff.
            #[cfg(unix)]
            drop(ssh_session);
            if attempt < attempts {
                let delay = ssh_retry_delay(attempt);
                refusal.log_retry(destination, &command.purpose, attempt, delay, stderr.trim());
                tokio::time::sleep(delay).await;
                continue;
            }
            refusal.log_exhausted(destination, &command.purpose, stderr.trim());
        }
        bail!(
            "{} failed with status {}: {}",
            command.purpose,
            output.status,
            stderr.trim()
        );
    }
    unreachable!("the last attempt always returns")
}

/// Spawn one probe and collect its output, with the SSH session lease it ran
/// on. The lease is returned so it outlives the child.
async fn run_resource_command(
    command: &CommandSpec,
) -> Result<(CommandOutput, Option<crate::targets::SshSessionLease>)> {
    // A remote probe runs as one session on a shared SSH connection; the
    // lease is held until the probe has exited.
    let (command, ssh_session) = if command.ssh_session.is_some() {
        let requested = command.clone();
        let executor = crate::targets::CancellableProcessExecutor::with_timeout(
            crate::targets::SSH_MASTER_OPEN_TIMEOUT,
        );
        let _cancel_preparation = executor.cancel_on_drop();
        tokio::task::spawn_blocking(move || {
            requested
                .open_ssh_session(&executor)
                .map(crate::targets::SessionCommand::into_parts)
        })
        .await
        .context("join the SSH session lease for a resource probe")??
    } else {
        (command.clone(), None)
    };
    let command = &command;
    let mut process = tokio::process::Command::new(&command.program);
    process.args(&command.args).envs(&command.env);
    let output =
        mj_core::subprocess::run_bounded(&mut process, 8 * 1024 * 1024, RESOURCE_POLL_TIMEOUT)
            .await
            .with_context(|| format!("wait for {}", command.purpose))?;
    let command_output = CommandOutput {
        status: output.status.code().unwrap_or(-1),
        stdout: output.stdout,
        stderr: output.stderr,
    };
    Ok((command_output, ssh_session))
}