Skip to main content

mj_controller/pollers/
quota.rs

1use super::*;
2
3pub fn quota_refresh_profiles(controller: &Controller) -> Vec<QuotaRefreshRequest> {
4    let cwd = std::env::current_dir().unwrap_or_else(|_| PathBuf::from("."));
5    controller
6        .config
7        .enabled_profiles()
8        .map(|(id, profile)| QuotaRefreshRequest::for_profile(id, profile, cwd.clone()))
9        .collect()
10}
11
12/// Reads a profile's stored report, the rebuildable copy the daemon keeps in
13/// the store. Called on a blocking thread.
14pub type QuotaCacheLoader =
15    Arc<dyn Fn(&QuotaRefreshRequest) -> Option<crate::quota::ProfileQuota> + Send + Sync>;
16
17/// When a profile is next due for a probe, in epoch seconds: one interval
18/// after its last report, or the end of a rate-limit hold when that is later.
19/// A profile with no report is due at once (zero).
20pub(super) fn next_probe_at(report: Option<&crate::quota::ProfileQuota>) -> u64 {
21    report.map_or(0, |report| {
22        (report.refreshed_at_epoch_seconds + QUOTA_REFRESH_INTERVAL.as_secs())
23            .max(report.rate_limited_until_epoch_seconds.unwrap_or(0))
24    })
25}
26
27/// Whether a provider told this profile to wait, and the wait is not over.
28fn on_hold(report: Option<&crate::quota::ProfileQuota>, now: u64) -> bool {
29    report
30        .and_then(|report| report.rate_limited_until_epoch_seconds)
31        .is_some_and(|until| until > now)
32}
33
34/// The daemon's quota poller, the only process that asks a provider.
35///
36/// Each profile has its own schedule. A profile whose stored report is younger
37/// than [`QUOTA_REFRESH_INTERVAL`] is published as it is and probed when that
38/// report ages out; anything else is probed now. A batch with `refresh` set
39/// probes every profile.
40pub fn spawn_quota_refresher(
41    cache: QuotaCacheLoader,
42) -> (
43    tokio::sync::watch::Sender<QuotaRefreshBatch>,
44    tokio::sync::mpsc::Receiver<QuotaUpdate>,
45) {
46    let (profiles_tx, mut profiles_rx) = tokio::sync::watch::channel(QuotaRefreshBatch::default());
47    let (updates_tx, updates_rx) = tokio::sync::mpsc::channel(32);
48    tokio::spawn(async move {
49        let mut quotas = QuotaManager::default();
50        let mut batch = QuotaRefreshBatch::default();
51        // The cache identity each profile was last adopted under. A changed
52        // identity is a changed profile, whose report no longer applies.
53        let mut identities: std::collections::BTreeMap<String, String> = Default::default();
54        loop {
55            let now = epoch_seconds();
56            let wake = batch
57                .profiles
58                .iter()
59                .map(|request| next_probe_at(quotas.report(&request.profile_id)))
60                .min()
61                .map(|due| Duration::from_secs(due.saturating_sub(now).max(5)));
62            tokio::select! {
63                _ = tokio::time::sleep(wake.unwrap_or_default()), if wake.is_some() => {}
64                changed = profiles_rx.changed() => {
65                    if changed.is_err() {
66                        tracing::debug!("quota profile target feed closed; stopping quota refresher");
67                        break;
68                    }
69                    batch = profiles_rx.borrow_and_update().clone();
70                    if !adopt_profiles(&mut quotas, &mut identities, &batch, &cache, &updates_tx).await {
71                        break;
72                    }
73                }
74            }
75            let now = epoch_seconds();
76            let due = batch
77                .profiles
78                .iter()
79                .filter(|request| {
80                    let report = quotas.report(&request.profile_id);
81                    // A hold outranks even an explicit refresh: probing an
82                    // endpoint that just said 429 only extends the limit.
83                    !on_hold(report, now) && (batch.refresh || next_probe_at(report) <= now)
84                })
85                .cloned()
86                .collect::<Vec<_>>();
87            // A batch that only changed the profile set and finds nothing due
88            // is not a cycle. A requested refresh always is, so the person's
89            // notice can end.
90            if due.is_empty() && !batch.refresh {
91                continue;
92            }
93            let generation = batch.generation;
94            batch.refresh = false;
95            if !refresh_profile_quotas(&mut quotas, generation, &due, &updates_tx).await {
96                break;
97            }
98        }
99        quotas.shutdown().await;
100    });
101    (profiles_tx, updates_rx)
102}
103
104/// Bring the manager in line with a new batch: drop profiles that left or
105/// changed, and adopt the stored report of each new one when it is still
106/// current. Reports whether the consumer is still listening.
107async fn adopt_profiles(
108    quotas: &mut QuotaManager,
109    identities: &mut std::collections::BTreeMap<String, String>,
110    batch: &QuotaRefreshBatch,
111    cache: &QuotaCacheLoader,
112    updates: &tokio::sync::mpsc::Sender<QuotaUpdate>,
113) -> bool {
114    let keep = batch
115        .profiles
116        .iter()
117        .map(|request| request.profile_id.clone())
118        .collect::<std::collections::BTreeSet<_>>();
119    identities.retain(|id, _| keep.contains(id));
120    quotas.retain_profiles(&keep).await;
121    for request in &batch.profiles {
122        let identity = request.cache_identity();
123        if identities.get(&request.profile_id) == Some(&identity) {
124            continue;
125        }
126        identities.insert(request.profile_id.clone(), identity);
127        quotas.forget(&request.profile_id);
128        let load = cache.clone();
129        let for_request = request.clone();
130        let stored = match tokio::task::spawn_blocking(move || load(&for_request)).await {
131            Ok(stored) => stored,
132            Err(error) => {
133                tracing::warn!(%error, "stored quota read task failed");
134                None
135            }
136        };
137        let Some(stored) = stored.filter(|report| {
138            report.error.is_none() && next_probe_at(Some(report)) > epoch_seconds()
139        }) else {
140            continue;
141        };
142        tracing::debug!(
143            profile_id = %request.profile_id,
144            refreshed_at = stored.refreshed_at_epoch_seconds,
145            "using the stored quota report; the next probe waits for it to age"
146        );
147        quotas.seed(stored.clone());
148        let outcome = QuotaRefreshOutcome {
149            report: stored,
150            credentials_changed: false,
151        };
152        if updates.send(QuotaUpdate::Report(outcome)).await.is_err() {
153            return false;
154        }
155    }
156    true
157}
158
159pub(super) async fn refresh_profile_quotas(
160    quotas: &mut QuotaManager,
161    generation: u64,
162    profiles: &[QuotaRefreshRequest],
163    updates: &tokio::sync::mpsc::Sender<QuotaUpdate>,
164) -> bool {
165    let ids = profiles
166        .iter()
167        .map(|profile| profile.profile_id.clone())
168        .collect::<Vec<_>>();
169    if updates
170        .send(QuotaUpdate::Refreshing { profile_ids: ids })
171        .await
172        .is_err()
173    {
174        tracing::debug!("quota update consumer closed before refresh started");
175        return false;
176    }
177    // Keep draining even if the UI is gone so codex clients return to the
178    // manager for a clean shutdown; just stop sending.
179    let delivered = AtomicBool::new(true);
180    quotas
181        .probe(profiles.to_vec(), |quota| {
182            let delivered = &delivered;
183            async move {
184                if delivered.load(Ordering::Acquire)
185                    && updates.send(QuotaUpdate::Report(quota)).await.is_err()
186                {
187                    tracing::debug!("quota update consumer closed while reporting a profile");
188                    delivered.store(false, Ordering::Release);
189                }
190            }
191        })
192        .await;
193    if !delivered.into_inner() {
194        return false;
195    }
196    if updates
197        .send(QuotaUpdate::Finished { generation })
198        .await
199        .is_err()
200    {
201        tracing::debug!(
202            generation,
203            "quota update consumer closed before refresh completed"
204        );
205        false
206    } else {
207        true
208    }
209}
210
211/// What the daemon wants to tell the user about a background image download.
212///
213/// The refresher logs every detail; these are the few moments worth a notice
214/// in the dashboard, because the person's first session waits on them.
215#[derive(Debug, Clone, PartialEq, Eq)]
216pub enum ImageRefreshReport {
217    /// A download has just begun.
218    Started { host: String, image: String },
219    /// A download finished and left the host with the image.
220    Pulled { host: String, image: String },
221    /// A download failed. Reported once per distinct error, not once an hour.
222    Failed {
223        host: String,
224        image: String,
225        error: String,
226    },
227}
228
229/// Download every configured container image the host lacks, and keep the
230/// ones that track a moving tag current, away from any session launch.
231///
232/// The first pass runs shortly after the daemon starts, which is what spares
233/// the person's first session a multi-gigabyte download. `plan` is called on
234/// every tick rather than once, so a config reload changes what gets
235/// downloaded without a daemon restart. Hosts refresh concurrently; each host
236/// runs its own commands in order.
237///
238/// `report` is how the daemon speaks: it is called from the blocking download
239/// threads as well as from this task, so it must be cheap and must not block.
240pub fn spawn_image_refresher(
241    plan: impl Fn() -> Vec<ImageRefresh> + Send + 'static,
242    report: impl Fn(ImageRefreshReport) + Send + Sync + 'static,
243    cancellation: tokio_util::sync::CancellationToken,
244) -> tokio::task::JoinHandle<()> {
245    let report: Arc<dyn Fn(ImageRefreshReport) + Send + Sync> = Arc::new(report);
246    tokio::spawn(async move {
247        let mut interval = tokio::time::interval_at(
248            tokio::time::Instant::now() + IMAGE_REFRESH_DELAY,
249            IMAGE_REFRESH_INTERVAL,
250        );
251        // A refresh slower than the interval collapses the ticks it missed
252        // instead of stacking a second pull behind the first.
253        interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
254        // The last error reported for each host and image. An unreachable
255        // host must not produce the same notice every hour.
256        let mut last_failures: BTreeMap<String, String> = BTreeMap::new();
257        loop {
258            tokio::select! {
259                // Quitting wins over a tick that came due during a long
260                // refresh, so shutdown never starts one more pull.
261                biased;
262                _ = cancellation.cancelled() => return,
263                _ = interval.tick() => {
264                    refresh_images(plan(), &report, &mut last_failures, &cancellation).await;
265                }
266            }
267        }
268    })
269}
270
271/// The key that identifies one host's copy of one image, for failure
272/// suppression.
273fn refresh_key(host: &str, image: &str) -> String {
274    format!("{host}|{image}")
275}
276
277/// Report a failed download once, and stay quiet while it keeps failing the
278/// same way.
279///
280/// A host that is simply offline fails identically every hour, and a notice
281/// an hour would be noise. A different error is new information, and so is a
282/// failure after a success, which is why success clears the record.
283pub(super) fn record_refresh_result(
284    last_failures: &mut BTreeMap<String, String>,
285    host: &str,
286    image: &str,
287    error: Option<String>,
288    report: &dyn Fn(ImageRefreshReport),
289) {
290    let key = refresh_key(host, image);
291    let Some(error) = error else {
292        last_failures.remove(&key);
293        return;
294    };
295    if last_failures.get(&key) == Some(&error) {
296        return;
297    }
298    last_failures.insert(key, error.clone());
299    report(ImageRefreshReport::Failed {
300        host: host.to_owned(),
301        image: image.to_owned(),
302        error,
303    });
304}
305
306/// Whether a local image host's engine can be run at all.
307///
308/// Only local hosts are checked: a remote host's engine lives on the other
309/// side of ssh, and a failure there is real news about that host.
310pub(super) fn local_engine_installed(host: &ImageHost, path: Option<&std::ffi::OsStr>) -> bool {
311    match host {
312        ImageHost::LocalPodman | ImageHost::LocalDocker | ImageHost::AppleContainer => {
313            crate::targets::program_on_path(host.engine(), path)
314        }
315        ImageHost::SshPodman(_) | ImageHost::SshDocker(_) => true,
316    }
317}
318
319pub(super) async fn refresh_images(
320    plan: Vec<ImageRefresh>,
321    report: &Arc<dyn Fn(ImageRefreshReport) + Send + Sync>,
322    last_failures: &mut BTreeMap<String, String>,
323    cancellation: &tokio_util::sync::CancellationToken,
324) {
325    if plan.is_empty() {
326        return;
327    }
328    // Pulls are restartable preparation. Shutdown cancels their subprocesses;
329    // the next daemon discovers any image still missing or out of date.
330    // One flag for every host, so quitting kills the pulls in flight instead of
331    // waiting out a multi-gigabyte download.
332    let cancelled = Arc::new(AtomicBool::new(false));
333    let mut hosts = tokio::task::JoinSet::new();
334    for refresh in plan {
335        // The default configuration names a podman, a docker, and on macOS an
336        // Apple container target whether or not the engine is installed. An
337        // engine that is not on this machine is not a failed download, and it
338        // must not become a notice on every start.
339        if !local_engine_installed(&refresh.host, std::env::var_os("PATH").as_deref()) {
340            tracing::debug!(
341                host = refresh.host.label(),
342                image = refresh.image,
343                "container engine is not installed; skipping the image refresh"
344            );
345            continue;
346        }
347        // ProcessExecutor is synchronous, and a pull is long: it belongs on a
348        // blocking thread, never on the runtime.
349        let executor = CancellableProcessExecutor::new(cancelled.clone());
350        let report = report.clone();
351        hosts.spawn_blocking(move || {
352            let host = refresh.host.label();
353            // A launch that needs this image waits on the same lock, so it
354            // never starts a second download of what this tick is fetching.
355            let lock = crate::image_pull_gate::image_pull_mutex(&refresh.host, &refresh.image);
356            let held =
357                crate::image_pull_gate::hold_image_pull(&lock, || executor.is_cancelled(), || {});
358            let outcome = held.and_then(|guard| {
359                let outcome = refresh_host_image(&refresh, &executor, &*report);
360                drop(guard);
361                outcome
362            });
363            let error = match outcome {
364                Ok(_) => None,
365                Err(error) if executor.is_cancelled() => {
366                    // The daemon is leaving. That is not a fault of the host,
367                    // and it is not news for the user either.
368                    tracing::debug!(
369                        host,
370                        image = refresh.image,
371                        error = format!("{error:#}"),
372                        "container image refresh cancelled"
373                    );
374                    return None;
375                }
376                Err(error) => {
377                    tracing::warn!(
378                        host,
379                        image = refresh.image,
380                        error = format!("{error:#}"),
381                        "could not refresh a container image"
382                    );
383                    Some(format!("{error:#}"))
384                }
385            };
386            Some((host, refresh.image, error))
387        });
388    }
389    let mut cancelling = false;
390    loop {
391        tokio::select! {
392            biased;
393            _ = cancellation.cancelled(), if !cancelling => {
394                cancelling = true;
395                cancelled.store(true, Ordering::Release);
396            }
397            joined = hosts.join_next() => match joined {
398                None => return,
399                Some(Ok(None)) => {}
400                Some(Ok(Some((host, image, error)))) => {
401                    record_refresh_result(last_failures, &host, &image, error, &**report);
402                }
403                Some(Err(error)) => {
404                    tracing::warn!(%error, "container image refresh task failed");
405                }
406            },
407        }
408    }
409}
410
411/// What one host's refresh of one image did.
412#[derive(Debug, Clone, PartialEq, Eq)]
413pub(super) enum ImageRefreshOutcome {
414    /// The host already had the image and this target only wants it present,
415    /// so nothing was downloaded.
416    Present,
417    /// A pull ran and the host's copy did not change.
418    Unchanged,
419    /// A pull ran and left the host with a different image.
420    Pulled { id: String },
421}
422
423/// Pull one image on one host, then drop whatever that unlinked.
424///
425/// The image id before and after says whether the pull actually changed
426/// anything, which is the only part worth an `info` line. An image that is
427/// only downloaded when absent skips the pull, and the prune with it, as soon
428/// as the host reports a copy.
429///
430/// `report` is told when a download actually starts and when one leaves the
431/// host with a new image, so the user hears about the wait they are in rather
432/// than about every hourly check.
433pub(super) fn refresh_host_image(
434    refresh: &ImageRefresh,
435    executor: &impl CommandExecutor,
436    report: &dyn Fn(ImageRefreshReport),
437) -> Result<ImageRefreshOutcome> {
438    let host = refresh.host.label();
439    let cached = image_id(&refresh.image_id, executor);
440    if refresh.when == RefreshWhen::WhenAbsent && cached.is_some() {
441        tracing::debug!(
442            host,
443            image = refresh.image,
444            "the host already has this container image"
445        );
446        return Ok(ImageRefreshOutcome::Present);
447    }
448    // Only a host with no copy is about to download for real. An hourly
449    // refresh of a moving tag usually finds nothing newer, and announcing it
450    // every hour would be noise.
451    if cached.is_none() {
452        report(ImageRefreshReport::Started {
453            host: host.clone(),
454            image: refresh.image.clone(),
455        });
456    }
457    run_refresh_command(&refresh.pull, executor)?;
458    let pulled = image_id(&refresh.image_id, executor);
459    let outcome = if pulled.is_some() && (cached.is_none() || pulled != cached) {
460        let id = pulled.unwrap_or_default();
461        tracing::info!(
462            host,
463            image = refresh.image,
464            id,
465            "pulled a newer container image"
466        );
467        report(ImageRefreshReport::Pulled {
468            host,
469            image: refresh.image.clone(),
470        });
471        ImageRefreshOutcome::Pulled { id }
472    } else {
473        tracing::debug!(
474            host,
475            image = refresh.image,
476            "container image is already current"
477        );
478        ImageRefreshOutcome::Unchanged
479    };
480    if let Some(prune) = &refresh.prune {
481        run_refresh_command(prune, executor)?;
482    }
483    Ok(outcome)
484}
485
486/// The host's id for an image, or `None` when it has no copy of it yet. A
487/// missing image is the ordinary first-pull case, not a fault.
488pub(super) fn image_id(command: &CommandSpec, executor: &impl CommandExecutor) -> Option<String> {
489    let output = executor.execute(command).ok()?;
490    if output.status != 0 {
491        return None;
492    }
493    let id = String::from_utf8_lossy(&output.stdout).trim().to_owned();
494    (!id.is_empty()).then_some(id)
495}
496
497pub(super) fn run_refresh_command(
498    command: &CommandSpec,
499    executor: &impl CommandExecutor,
500) -> Result<()> {
501    let output = executor.execute(command)?;
502    if output.status != 0 {
503        bail!(
504            "{} failed with status {}: {}",
505            command.purpose,
506            output.status,
507            String::from_utf8_lossy(&output.stderr).trim()
508        );
509    }
510    Ok(())
511}
512
513/// Whether a refresh the person asked for is over. `pending_cycles` is the
514/// daemon's finished-cycle count at the moment of the request, and
515/// `finished_cycles` is its count now. A cycle already running when the
516/// request arrived can end first and complete the notice a few seconds early;
517/// the report itself is what the row shows, so that is harmless.
518pub fn complete_manual_quota_refresh(
519    pending_cycles: &mut Option<u64>,
520    finished_cycles: u64,
521) -> bool {
522    if !pending_cycles.is_some_and(|pending| finished_cycles > pending) {
523        return false;
524    }
525    *pending_cycles = None;
526    true
527}