Skip to main content

kanade_shared/
bootstrap.rs

1//! Idempotent JetStream bootstrap (Sprint 6.x follow-up).
2//!
3//! Lists every NATS JetStream resource the kanade fleet expects —
4//! streams, KV buckets, Object Stores — and asks the broker to
5//! create-or-update them. v0.25.0 switched from `create_*` to
6//! `create_or_update_*`: the old form returned error 10058 ("name
7//! already in use with a different configuration") when a release
8//! widened a stream's subjects or changed its retention policy on
9//! a broker that still held the older config. With the new form the
10//! broker reconciles its definition to the one in this file, so
11//! version bumps no longer require operator-side data wipes.
12//!
13//! Centralising the list here means a future "we added a new
14//! bucket" change touches one place and both the operator CLI +
15//! the auto-bootstrap path pick it up.
16
17use std::time::Duration;
18
19use anyhow::{Context, Result};
20use async_nats::jetstream::{
21    self,
22    kv::Config as KvConfig,
23    object_store::Config as ObjectStoreConfig,
24    stream::{Config as StreamConfig, DiscardPolicy},
25};
26use tracing::{info, warn};
27
28/// How long the broker retains a command on `STREAM_EXEC`. Agents size their
29/// admission tombstones from this: a record must outlive the last moment the
30/// broker could redeliver the command it guards.
31pub const COMMAND_STREAM_MAX_AGE: Duration = Duration::from_secs(7 * 24 * 60 * 60);
32
33use crate::kv::{
34    BUCKET_AGENT_CONFIG, BUCKET_AGENT_GROUPS, BUCKET_AGENT_GROUPS_DERIVED, BUCKET_AGENT_META,
35    BUCKET_AGENTS_STATE, BUCKET_FLEET_CONFIG, BUCKET_GROUP_CONTACTS, BUCKET_JOBS, BUCKET_JOBS_YAML,
36    BUCKET_NOTIFICATIONS_READ, BUCKET_SCHEDULES, BUCKET_SCHEDULES_YAML, BUCKET_SCRIPT_CURRENT,
37    BUCKET_SCRIPT_STATUS, BUCKET_SERVER_SETTINGS, OBJECT_AGENT_RELEASES, OBJECT_APP_PACKAGES,
38    OBJECT_COLLECTIONS, OBJECT_RESULT_OUTPUT, OBJECT_SCRIPTS, STREAM_AUDIT, STREAM_EVENTS,
39    STREAM_EXEC, STREAM_INVENTORY, STREAM_NOTIFICATIONS, STREAM_OBS_EVENTS, STREAM_RESULTS,
40};
41use crate::wire::{
42    DEFAULT_AGENT_RELEASES_CAP_MIB, DEFAULT_APP_PACKAGES_CAP_MIB, DEFAULT_COLLECT_RETENTION_DAYS,
43    DEFAULT_COLLECTIONS_CAP_MIB, DEFAULT_RESULT_OUTPUT_CAP_MIB,
44    DEFAULT_RESULT_OUTPUT_RETENTION_DAYS, DEFAULT_SCRIPTS_CAP_MIB,
45};
46
47/// Create-or-update an Object Store, but never let it wedge backend
48/// startup. `create_object_store` neither reconciles an existing
49/// store's config nor has a `create_or_update` form in async-nats
50/// 0.49, so a store whose desired config drifted — e.g. the #518
51/// `max_bytes` cap added after the bucket was first created uncapped,
52/// which the broker then rejects with error 10058 ("stream name
53/// already in use with a different configuration") — would otherwise
54/// fail `ensure_jetstream_resources` and crash the backend on boot
55/// (production outage 2026-06-11). Fall back to the existing store
56/// (uncapped, as it already was) and warn. #506 tracks real
57/// reconciliation of object-store config.
58async fn ensure_object_store(js: &jetstream::Context, cfg: ObjectStoreConfig) -> Result<()> {
59    let name = cfg.bucket.clone();
60    if let Err(e) = js.create_object_store(cfg).await {
61        // The fallback is deliberately broad — any create error is
62        // tolerated AS LONG AS the store already exists, because the
63        // alternative is a wedged backend and "never crash on boot"
64        // wins over "surface this specific error". The expected error
65        // is 10058 (config drift, the incident), but auth/network
66        // blips on an already-bootstrapped broker take this path too;
67        // they remain visible via the `warn!`. Only a genuine
68        // "can't create AND doesn't exist" is fatal.
69        if js.get_object_store(&name).await.is_err() {
70            return Err(e).with_context(|| {
71                format!("create_object_store {name} (and no existing store to fall back to)")
72            });
73        }
74        warn!(
75            store = %name, error = %e,
76            "object store exists with a different config; using it as-is (cap not reconciled)",
77        );
78    }
79    info!(store = %name, "ready");
80    Ok(())
81}
82
83/// Idempotently create every NATS JetStream resource the kanade
84/// fleet relies on. Calling repeatedly is safe — `create_*` returns
85/// the existing resource if it's already configured.
86///
87/// Returns once every resource is in place. The function is async
88/// so backends can `await` it as part of their startup sequence
89/// (one round-trip per resource — ~10 RTTs total).
90pub async fn ensure_jetstream_resources(js: &jetstream::Context) -> Result<()> {
91    // ── Streams ──────────────────────────────────────────────────
92    // #518: every stream carries a `max_bytes` cap with
93    // `Discard::Old` on top of its `max_age` window. Within their
94    // age windows the streams used to be unbounded by size, and
95    // JetStream's file store shares a disk with SQLite on the
96    // backend host — one job printing 200 KB per run fleet-wide
97    // could exhaust the store, at which point EVERY publish fails
98    // (results, obs, audit, KV puts). With the caps, worst-case
99    // degradation is "shorter history on the offending stream"
100    // instead of "broker down".
101    //
102    // Sizing: JetStream RESERVES each `max_bytes` against its
103    // available storage (min of max_file_store and free disk) at
104    // create/update time and fails with error 10047 when the sum
105    // doesn't fit, so these must stay small enough for modest
106    // hosts. That's fine: every stream here is a transport +
107    // replay buffer — the durable record is the backend's SQLite
108    // (results/inventory/obs/audit are all projected within
109    // seconds) — so the caps are runaway-output backstops, not
110    // history budgets. Total reservation ≈ 5.3 GiB including the
111    // result_output object store below.
112    const MIB: i64 = 1024 * 1024;
113    const GIB: i64 = 1024 * MIB;
114
115    // INVENTORY — 90-day rolling history (spec §2.3.1).
116    js.create_or_update_stream(StreamConfig {
117        name: STREAM_INVENTORY.into(),
118        subjects: vec!["inventory.>".into()],
119        max_age: Duration::from_secs(90 * 24 * 60 * 60),
120        max_bytes: GIB,
121        discard: DiscardPolicy::Old,
122        ..Default::default()
123    })
124    .await
125    .with_context(|| format!("create_or_update_stream {STREAM_INVENTORY}"))?;
126    info!(stream = STREAM_INVENTORY, "ready");
127
128    // RESULTS — 30-day rolling history. The biggest producer by
129    // far (every job run on every PC, with up to 256 KB of inline
130    // stdout/stderr per message), so it gets the largest slice of
131    // the disk budget.
132    js.create_or_update_stream(StreamConfig {
133        name: STREAM_RESULTS.into(),
134        subjects: vec!["results.>".into()],
135        max_age: Duration::from_secs(30 * 24 * 60 * 60),
136        max_bytes: 2 * GIB,
137        discard: DiscardPolicy::Old,
138        ..Default::default()
139    })
140    .await
141    .with_context(|| format!("create_or_update_stream {STREAM_RESULTS}"))?;
142    info!(stream = STREAM_RESULTS, "ready");
143
144    // EXEC — latest-per-subject only (spec §2.6 Layer 1). v0.22.1:
145    // catch the existing `commands.{all,group.X,pc.Y}` subjects so a
146    // single backend publish lands in BOTH the agent's live core
147    // subscription AND the stream's retention store. Reconnecting
148    // agents catch up via a durable consumer with
149    // `DeliverPolicy::LastPerSubject` — they receive the most
150    // recent Command per subject they care about, no matter how
151    // long they were offline (within `max_age`).
152    js.create_or_update_stream(StreamConfig {
153        name: STREAM_EXEC.into(),
154        subjects: vec!["commands.>".into()],
155        max_messages_per_subject: 1,
156        max_age: COMMAND_STREAM_MAX_AGE,
157        // Latest-per-subject keeps this tiny (one Command per
158        // group/pc subject); the cap is a backstop against subject
159        // cardinality bugs, not a working budget.
160        max_bytes: 64 * MIB,
161        discard: DiscardPolicy::Old,
162        ..Default::default()
163    })
164    .await
165    .with_context(|| format!("create_or_update_stream {STREAM_EXEC}"))?;
166    info!(stream = STREAM_EXEC, "ready");
167
168    // EVENTS — short-lived broadcast bus for kill / revoke / etc.
169    // 7-day window matches the EXEC spec window.
170    js.create_or_update_stream(StreamConfig {
171        name: STREAM_EVENTS.into(),
172        subjects: vec!["events.>".into()],
173        max_age: Duration::from_secs(7 * 24 * 60 * 60),
174        max_bytes: 256 * MIB,
175        discard: DiscardPolicy::Old,
176        ..Default::default()
177    })
178    .await
179    .with_context(|| format!("create_or_update_stream {STREAM_EVENTS}"))?;
180    info!(stream = STREAM_EVENTS, "ready");
181
182    // AUDIT — operator-action record (spec §2.3.1). The DURABLE
183    // copy is the backend's SQLite `audit_log` table (the projector
184    // INSERTs each message, idempotently since #501; 365-day
185    // retention since #486) — the stream is transport + replay
186    // buffer, not the archive, so it can be bounded like the rest.
187    // 90 days / 512 MiB is far more than the projector ever lags;
188    // previously this stream had NO limits at all, making it an
189    // unbounded disk leak on the broker host.
190    js.create_or_update_stream(StreamConfig {
191        name: STREAM_AUDIT.into(),
192        subjects: vec!["audit.>".into()],
193        max_age: Duration::from_secs(90 * 24 * 60 * 60),
194        max_bytes: 512 * MIB,
195        discard: DiscardPolicy::Old,
196        ..Default::default()
197    })
198    .await
199    .with_context(|| format!("create_or_update_stream {STREAM_AUDIT}"))?;
200    info!(stream = STREAM_AUDIT, "ready");
201
202    // OBS_EVENTS — per-PC observability timeline (Issue #246). The
203    // 90-day window matches `obs_events` table retention so a
204    // backend bootstrapping after long downtime can catch up but
205    // doesn't carry data the table will discard anyway. Subject
206    // filter `obs.>` catches every PC without a per-PC subscription.
207    //
208    // Days-to-seconds is spelt out once instead of `90 * 24 * 60 *
209    // 60` open-coded across bootstrap + cleanup; the matching prune
210    // window in `kanade-backend::cleanup` quotes the same number
211    // separately (SQLite-relative string syntax there, not a
212    // duration), so it can't share a constant — but a single
213    // arithmetic spell-out here makes the relationship grep-able.
214    const SECS_PER_DAY: u64 = 24 * 60 * 60;
215    const OBS_EVENTS_RETENTION_DAYS: u64 = 90;
216    js.create_or_update_stream(StreamConfig {
217        name: STREAM_OBS_EVENTS.into(),
218        subjects: vec!["obs.>".into()],
219        max_age: Duration::from_secs(OBS_EVENTS_RETENTION_DAYS * SECS_PER_DAY),
220        max_bytes: 512 * MIB,
221        discard: DiscardPolicy::Old,
222        ..Default::default()
223    })
224    .await
225    .with_context(|| format!("create_or_update_stream {STREAM_OBS_EVENTS}"))?;
226    info!(stream = STREAM_OBS_EVENTS, "ready");
227
228    // NOTIFICATIONS — end-user notification history (SPEC §2.3.1 /
229    // Phase E). 90-day window matches INVENTORY: a Client App that
230    // connects after a notification was sent fetches the missed ones
231    // via KLP `notifications.list`. Subject filter `notifications.>`
232    // catches every fan-out target (`all` / `group.X` / `pc.Y`) with
233    // one stream. Retains all messages per subject — each notification
234    // is its own history entry, not a latest-only state like EXEC.
235    // #518: 512 MiB cap + DiscardPolicy::Old, matching the other
236    // 90-day streams (AUDIT / OBS_EVENTS) — notification payloads are
237    // small, so this is generous headroom while still bounding the
238    // broker's disk lease.
239    js.create_or_update_stream(StreamConfig {
240        name: STREAM_NOTIFICATIONS.into(),
241        subjects: vec!["notifications.>".into()],
242        max_age: Duration::from_secs(90 * 24 * 60 * 60),
243        max_bytes: 512 * MIB,
244        discard: DiscardPolicy::Old,
245        ..Default::default()
246    })
247    .await
248    .with_context(|| format!("create_or_update_stream {STREAM_NOTIFICATIONS}"))?;
249    info!(stream = STREAM_NOTIFICATIONS, "ready");
250
251    // ── KV buckets ───────────────────────────────────────────────
252    // script_current — cmd_id → version (spec §2.6 Layer 2).
253    js.create_or_update_key_value(KvConfig {
254        bucket: BUCKET_SCRIPT_CURRENT.into(),
255        history: 5,
256        ..Default::default()
257    })
258    .await
259    .with_context(|| format!("create_or_update_key_value {BUCKET_SCRIPT_CURRENT}"))?;
260    info!(bucket = BUCKET_SCRIPT_CURRENT, "ready");
261
262    // script_status — cmd_id → ACTIVE / REVOKED.
263    js.create_or_update_key_value(KvConfig {
264        bucket: BUCKET_SCRIPT_STATUS.into(),
265        history: 5,
266        ..Default::default()
267    })
268    .await
269    .with_context(|| format!("create_or_update_key_value {BUCKET_SCRIPT_STATUS}"))?;
270    info!(bucket = BUCKET_SCRIPT_STATUS, "ready");
271
272    // agents_state — pc_id → latest hw snapshot (history=1).
273    js.create_or_update_key_value(KvConfig {
274        bucket: BUCKET_AGENTS_STATE.into(),
275        history: 1,
276        ..Default::default()
277    })
278    .await
279    .with_context(|| format!("create_or_update_key_value {BUCKET_AGENTS_STATE}"))?;
280    info!(bucket = BUCKET_AGENTS_STATE, "ready");
281
282    // agent_config — Sprint 6 layered scopes (global / groups.* /
283    // pcs.*) plus the legacy target_version key.
284    // history: 1 — agents only ever read the current value (the watch is
285    // DeliverPolicy::New + an initial_sync get(), never kv.history()).
286    // Retained old revisions only fed reconnect history-replay, which
287    // flapped self-update backward (#828). Operator change-history lives
288    // in the audit log, so keeping one revision loses nothing. (#830)
289    js.create_or_update_key_value(KvConfig {
290        bucket: BUCKET_AGENT_CONFIG.into(),
291        history: 1,
292        ..Default::default()
293    })
294    .await
295    .with_context(|| format!("create_or_update_key_value {BUCKET_AGENT_CONFIG}"))?;
296    info!(bucket = BUCKET_AGENT_CONFIG, "ready");
297
298    // agent_groups — Sprint 5 per-pc group membership.
299    // history: 1 — same reasoning as agent_config above: agents only need
300    // the current membership; replayed history just churned subscriptions
301    // through stale sets on every reconnect (a transient wrong membership,
302    // #830). One revision makes that replay material non-existent. (#830)
303    js.create_or_update_key_value(KvConfig {
304        bucket: BUCKET_AGENT_GROUPS.into(),
305        history: 1,
306        ..Default::default()
307    })
308    .await
309    .with_context(|| format!("create_or_update_key_value {BUCKET_AGENT_GROUPS}"))?;
310    info!(bucket = BUCKET_AGENT_GROUPS, "ready");
311
312    // agent_groups_derived — #1032①: per-pc DERIVED group membership written by
313    // the backend group-materializer (resolved from GroupDefs). Agents watch it
314    // alongside agent_groups and union. history: 1 for the same reason as
315    // agent_groups (#830 — agents read only the current value; replayed history
316    // churns subscriptions). Separate bucket so the materializer (sole writer
317    // here) never clobbers operator-set membership in agent_groups.
318    js.create_or_update_key_value(KvConfig {
319        bucket: BUCKET_AGENT_GROUPS_DERIVED.into(),
320        history: 1,
321        ..Default::default()
322    })
323    .await
324    .with_context(|| format!("create_or_update_key_value {BUCKET_AGENT_GROUPS_DERIVED}"))?;
325    info!(bucket = BUCKET_AGENT_GROUPS_DERIVED, "ready");
326
327    // agent_meta — per-PC operator-managed key/value annotations
328    // (edited via the SPA agent detail page / the `kanade meta` CLI, both
329    // through the backend API, and typically bulk-populated by an operator AD-sync job). history: 5 —
330    // a few revisions of operator edit history, like group_contacts;
331    // nothing replays it, so the exact depth is not load-bearing.
332    js.create_or_update_key_value(KvConfig {
333        bucket: BUCKET_AGENT_META.into(),
334        history: 5,
335        ..Default::default()
336    })
337    .await
338    .with_context(|| format!("create_or_update_key_value {BUCKET_AGENT_META}"))?;
339    info!(bucket = BUCKET_AGENT_META, "ready");
340
341    // group_contacts — per-group notification email addresses
342    // (operator-managed via the SPA Groups page).
343    js.create_or_update_key_value(KvConfig {
344        bucket: BUCKET_GROUP_CONTACTS.into(),
345        history: 5,
346        ..Default::default()
347    })
348    .await
349    .with_context(|| format!("create_or_update_key_value {BUCKET_GROUP_CONTACTS}"))?;
350    info!(bucket = BUCKET_GROUP_CONTACTS, "ready");
351
352    // schedules — admin-API CRUD'd cron table (spec §2.5.3).
353    // Backend's scheduler.rs also creates this on startup; calling
354    // twice is harmless.
355    js.create_or_update_key_value(KvConfig {
356        bucket: BUCKET_SCHEDULES.into(),
357        history: 5,
358        ..Default::default()
359    })
360    .await
361    .with_context(|| format!("create_or_update_key_value {BUCKET_SCHEDULES}"))?;
362    info!(bucket = BUCKET_SCHEDULES, "ready");
363
364    // jobs — v0.15 operator-registered Manifest catalog. Schedules
365    // reference rows here by id; editing a job rewrites what future
366    // schedule fires exec.
367    js.create_or_update_key_value(KvConfig {
368        bucket: BUCKET_JOBS.into(),
369        history: 5,
370        ..Default::default()
371    })
372    .await
373    .with_context(|| format!("create_or_update_key_value {BUCKET_JOBS}"))?;
374    info!(bucket = BUCKET_JOBS, "ready");
375
376    // fleet_config — #418 Phase 5 fleet-wide singletons (the global
377    // change-freeze under KEY_FREEZE). history: 1 — only the current
378    // state matters; both schedulers watch it.
379    js.create_or_update_key_value(KvConfig {
380        bucket: BUCKET_FLEET_CONFIG.into(),
381        history: 1,
382        ..Default::default()
383    })
384    .await
385    .with_context(|| format!("create_or_update_key_value {BUCKET_FLEET_CONFIG}"))?;
386    info!(bucket = BUCKET_FLEET_CONFIG, "ready");
387
388    // server_settings — backend-side operator-editable settings (SPA
389    // Settings page "server settings" tab). A single JSON document under
390    // KEY_SERVER_SETTINGS; history: 1 since only the current state
391    // matters. First consumer is the cleanup task's dead-agent prune
392    // window.
393    js.create_or_update_key_value(KvConfig {
394        bucket: BUCKET_SERVER_SETTINGS.into(),
395        history: 1,
396        ..Default::default()
397    })
398    .await
399    .with_context(|| format!("create_or_update_key_value {BUCKET_SERVER_SETTINGS}"))?;
400    info!(bucket = BUCKET_SERVER_SETTINGS, "ready");
401
402    // notifications_read — per-(pc, user, notification) read/ack state
403    // (SPEC §2.3.2 / Phase E). The agent writes here on KLP
404    // `notifications.ack`; `notifications.list` reads it back to filter
405    // the unread bucket. history: 1 — only the latest ack per key
406    // matters.
407    js.create_or_update_key_value(KvConfig {
408        bucket: BUCKET_NOTIFICATIONS_READ.into(),
409        history: 1,
410        ..Default::default()
411    })
412    .await
413    .with_context(|| format!("create_or_update_key_value {BUCKET_NOTIFICATIONS_READ}"))?;
414    info!(bucket = BUCKET_NOTIFICATIONS_READ, "ready");
415
416    // jobs_yaml / schedules_yaml — operator source-of-truth YAML
417    // alongside the JSON catalogs above. Same key shape (manifest id
418    // / schedule id), but the value is the raw YAML bytes so the
419    // SPA's YAML editor preserves comments + script block-scalar
420    // indentation across edits. Agents/scheduler don't read these.
421    js.create_or_update_key_value(KvConfig {
422        bucket: BUCKET_JOBS_YAML.into(),
423        history: 5,
424        ..Default::default()
425    })
426    .await
427    .with_context(|| format!("create_or_update_key_value {BUCKET_JOBS_YAML}"))?;
428    info!(bucket = BUCKET_JOBS_YAML, "ready");
429
430    js.create_or_update_key_value(KvConfig {
431        bucket: BUCKET_SCHEDULES_YAML.into(),
432        history: 5,
433        ..Default::default()
434    })
435    .await
436    .with_context(|| format!("create_or_update_key_value {BUCKET_SCHEDULES_YAML}"))?;
437    info!(bucket = BUCKET_SCHEDULES_YAML, "ready");
438
439    // ── Object Store ─────────────────────────────────────────────
440    // #1247: every bucket is born with a `max_bytes` cap. The first
441    // three had none at all (unbounded disk leak); result_output /
442    // collections had caps in code since #518, but pre-existing
443    // buckets never received them (create-drift tolerated by
444    // `ensure_object_store`) — the backend's boot reconcile
445    // (`reconcile_object_store_max_bytes` driven by
446    // `ServerSettings::object_store_caps`) is what delivers these to
447    // a live broker. The shared DEFAULT_*_CAP_MIB constants keep the
448    // fresh-bucket value and the SPA-visible default in one place.
449    // agent_releases — one object per version, raw exe bytes.
450    ensure_object_store(
451        js,
452        ObjectStoreConfig {
453            bucket: OBJECT_AGENT_RELEASES.into(),
454            max_bytes: DEFAULT_AGENT_RELEASES_CAP_MIB as i64 * MIB,
455            ..Default::default()
456        },
457    )
458    .await?;
459
460    // app_packages — generic operator-uploaded binary distribution
461    // (kanade-client today; third-party installers like Webex /
462    // Teams once those flows land). Object keys are
463    // `<name>/<version>`; see `kanade-shared::kv::OBJECT_APP_PACKAGES`
464    // for the full rationale.
465    ensure_object_store(
466        js,
467        ObjectStoreConfig {
468            bucket: OBJECT_APP_PACKAGES.into(),
469            max_bytes: DEFAULT_APP_PACKAGES_CAP_MIB as i64 * MIB,
470            ..Default::default()
471        },
472    )
473    .await?;
474
475    // scripts — manifest script bodies referenced by
476    // `Execute::script_object` (SPEC §2.4.1). Sibling of
477    // `app_packages`; see `kanade-shared::kv::OBJECT_SCRIPTS` for
478    // the bucket-split rationale (smaller payloads + manifest-
479    // coupled lifecycle vs operator-curated installers).
480    ensure_object_store(
481        js,
482        ObjectStoreConfig {
483            bucket: OBJECT_SCRIPTS.into(),
484            max_bytes: DEFAULT_SCRIPTS_CAP_MIB as i64 * MIB,
485            ..Default::default()
486        },
487    )
488    .await?;
489
490    // result_output — overflow stdout / stderr blobs for the
491    // `ExecResult` wire kind (#227). Anything larger than the agent's
492    // 256 KB inline threshold gets uploaded here under
493    // `<request_id>/{stdout,stderr}`; the backend's results
494    // projector derefs the pointer fields before INSERT so SQLite
495    // + the SPA see the full text inline.
496    //
497    // The window is DELIBERATELY shorter than STREAM_RESULTS's 30 days, and
498    // used to match it on the reasoning that "a row still resolvable in
499    // execution_results never points at a missing blob". That had the
500    // dependency backwards: the row does not point at the blob at all once
501    // projected — the projector derefs it into SQLite first. The only reader
502    // left is a `-WipeDb` replay, so this is a RECOVERY window, not a data
503    // lifetime, and thirty days of copies nobody reads is what walked the
504    // bucket into its cap (#1321: a hard wall, not eviction).
505    // Operator-tunable via `ServerSettings::result_output_retention_days`;
506    // applied to EXISTING buckets by `reconcile_object_store_max_age`, since
507    // `create_object_store` does not reconcile (#506).
508    // #518: capped like the streams — a job whose output overflows
509    // the inline threshold writes blobs HERE instead of
510    // STREAM_RESULTS, so without its own cap this store bypasses
511    // the stream budget entirely and can still fill the file store.
512    // The projector derefs blobs within seconds of publish, so
513    // eviction only ever hits already-projected (or expired)
514    // output.
515    ensure_object_store(
516        js,
517        ObjectStoreConfig {
518            bucket: OBJECT_RESULT_OUTPUT.into(),
519            max_age: Duration::from_secs(
520                SECS_PER_DAY * DEFAULT_RESULT_OUTPUT_RETENTION_DAYS as u64,
521            ),
522            max_bytes: DEFAULT_RESULT_OUTPUT_CAP_MIB as i64 * MIB,
523            ..Default::default()
524        },
525    )
526    .await?;
527
528    // #219: collected file bundles. A `collect:` job's agent zips the
529    // script's listed files and uploads the archive here under
530    // `<pc_id>/<job_id>/<rfc3339>.zip`; the SPA Collect page lists /
531    // downloads them. Default max_age = DEFAULT_COLLECT_RETENTION_DAYS —
532    // bundles are debugging / audit artifacts (not curated config like
533    // app_packages / scripts), so they auto-expire and the bucket doesn't
534    // grow unbounded. Capped at 5 GiB (DiscardPolicy::Old evicts oldest
535    // first) so a fleet's worth of bundles can't fill the file store.
536    //
537    // This is only the value a FRESH bucket is born with; the window is
538    // operator-tunable from the SPA (`ServerSettings::collect_retention_days`)
539    // and the backend reconciles the live bucket's max_age to the configured
540    // value at boot and on save — see [`reconcile_collect_retention`].
541    ensure_object_store(
542        js,
543        ObjectStoreConfig {
544            bucket: OBJECT_COLLECTIONS.into(),
545            max_age: Duration::from_secs(SECS_PER_DAY * DEFAULT_COLLECT_RETENTION_DAYS as u64),
546            max_bytes: DEFAULT_COLLECTIONS_CAP_MIB as i64 * MIB,
547            ..Default::default()
548        },
549    )
550    .await?;
551
552    Ok(())
553}
554
555/// NATS names the stream backing an Object Store `OBJ_<bucket>` (mirroring
556/// `KV_<bucket>` for key-value stores). We reconcile the collect bucket's
557/// retention through this stream because async-nats 0.49 has no
558/// create-or-update / reconcile form for Object Stores themselves (the same
559/// gap [`ensure_object_store`] works around) — but the underlying stream
560/// *does* support `update_stream`.
561fn object_store_stream_name(bucket: &str) -> String {
562    format!("OBJ_{bucket}")
563}
564
565/// Reconcile the `collections` Object Store's retention window to
566/// `retention_days` by updating the `max_age` on its backing stream.
567///
568/// Why this exists: the bucket is created once (at bootstrap) with the
569/// built-in default, and `create_object_store` neither has a
570/// create-or-update form nor reconciles config in async-nats 0.49. So to
571/// honour an operator's `ServerSettings::collect_retention_days` change on an
572/// already-provisioned bucket, we read the backing stream's config, patch
573/// **only** `max_age` (a read-modify-write that leaves every object-store-
574/// specific stream setting untouched), and `update_stream`. `max_bytes`
575/// and the discard policy are deliberately left as-is, so extending the
576/// window never lifts the 5 GiB disk ceiling.
577///
578/// Idempotent: if the stream's `max_age` already matches, it's a no-op
579/// (skips the update round-trip and returns `false`). A missing stream (the
580/// bucket was never provisioned — e.g. a broker that predates this feature
581/// and hasn't run bootstrap) is a soft error the caller can log-and-continue:
582/// bootstrap runs before this on the backend boot path, so in practice the
583/// stream is always present.
584///
585/// Returns `Ok(true)` when it actually changed the stream, `Ok(false)` when
586/// already in sync.
587pub async fn reconcile_collect_retention(
588    js: &jetstream::Context,
589    retention_days: u32,
590) -> Result<bool> {
591    reconcile_object_store_max_age(js, OBJECT_COLLECTIONS, retention_days).await
592}
593
594/// Reconcile ONE Object Store's retention window to `retention_days`, by
595/// updating `max_age` on its backing `OBJ_<bucket>` stream.
596///
597/// Generalised from the collect-only version when `result_output` needed the
598/// same treatment (#1321 fallout). Two buckets reconciling their `max_age`
599/// through two near-identical read-modify-write functions is the shape that
600/// drifts — and a drift here means one bucket silently keeps the old window.
601///
602/// Read-modify-write patching ONLY `max_age`, so object-store-specific stream
603/// settings and the `max_bytes` cap stay untouched.
604///
605/// Returns `Ok(true)` when it actually changed the stream, `Ok(false)` when
606/// already in sync.
607pub async fn reconcile_object_store_max_age(
608    js: &jetstream::Context,
609    bucket: &str,
610    retention_days: u32,
611) -> Result<bool> {
612    const SECS_PER_DAY: u64 = 24 * 60 * 60;
613    let desired = Duration::from_secs(SECS_PER_DAY * retention_days as u64);
614    let stream_name = object_store_stream_name(bucket);
615
616    let mut stream = js
617        .get_stream(&stream_name)
618        .await
619        .with_context(|| format!("get_stream {stream_name} for max_age reconcile"))?;
620    let info = stream
621        .info()
622        .await
623        .with_context(|| format!("stream info {stream_name}"))?;
624    if info.config.max_age == desired {
625        return Ok(false);
626    }
627    let mut cfg = info.config.clone();
628    cfg.max_age = desired;
629    js.update_stream(cfg)
630        .await
631        .with_context(|| format!("update_stream {stream_name} max_age"))?;
632    info!(
633        stream = %stream_name,
634        bucket,
635        retention_days,
636        "reconciled Object Store max_age",
637    );
638    Ok(true)
639}
640
641/// Reconcile one Object Store's disk cap to `cap_mib` by updating
642/// `max_bytes` on its backing `OBJ_<bucket>` stream (#1247).
643///
644/// Same shape as [`reconcile_collect_retention`] (read-modify-write
645/// patching ONLY the one field, so object-store-specific stream
646/// settings and the `max_age` window stay untouched), but for the
647/// field that bounds disk. This is what delivers caps to buckets
648/// created BEFORE the cap existed — `ensure_object_store` deliberately
649/// tolerates the 10058 config-drift error instead of reconciling, so
650/// without this a live bucket keeps whatever it was born with forever
651/// (how `OBJ_result_output` reached 6.76 GB against a nominal 1 GiB).
652///
653/// Idempotent: returns `Ok(false)` when already in sync. A missing
654/// stream is a soft error the caller logs and continues past —
655/// bootstrap runs before this on the boot path. Eviction is the
656/// broker's job once the cap lands (`DiscardPolicy::Old`, the object
657/// store default); shrink-to-fit is not instantaneous but needs no
658/// operator action.
659pub async fn reconcile_object_store_max_bytes(
660    js: &jetstream::Context,
661    bucket: &str,
662    cap_mib: u32,
663) -> Result<bool> {
664    const MIB: i64 = 1024 * 1024;
665    let desired = cap_mib as i64 * MIB;
666    let stream_name = object_store_stream_name(bucket);
667
668    let mut stream = js
669        .get_stream(&stream_name)
670        .await
671        .with_context(|| format!("get_stream {stream_name} for max_bytes reconcile"))?;
672    let info = stream
673        .info()
674        .await
675        .with_context(|| format!("stream info {stream_name}"))?;
676    if info.config.max_bytes == desired {
677        return Ok(false);
678    }
679    let mut cfg = info.config.clone();
680    cfg.max_bytes = desired;
681    js.update_stream(cfg)
682        .await
683        .with_context(|| format!("update_stream {stream_name} max_bytes"))?;
684    info!(
685        stream = %stream_name,
686        cap_mib,
687        "object store: reconciled max_bytes",
688    );
689    Ok(true)
690}
691
692#[cfg(test)]
693mod tests {
694    use super::*;
695    use std::process::Stdio;
696
697    /// Throwaway `nats-server -js` on a random port, like the
698    /// kv_cas_live / offline_boot harnesses. Ignored tests only.
699    struct Broker {
700        js: jetstream::Context,
701        _server: tokio::process::Child,
702        _storage: tempfile::TempDir,
703    }
704
705    async fn spawn_broker() -> Broker {
706        let port = portpicker::pick_unused_port().expect("pick port");
707        let storage = tempfile::TempDir::new().expect("storage tempdir");
708        let server = tokio::process::Command::new("nats-server")
709            .arg("-js")
710            .arg("-p")
711            .arg(port.to_string())
712            .arg("-sd")
713            .arg(storage.path())
714            .stdout(Stdio::null())
715            .stderr(Stdio::null())
716            .kill_on_drop(true)
717            .spawn()
718            .expect("spawn nats-server (is it in PATH?)");
719        let url = format!("nats://127.0.0.1:{port}");
720        let mut client = None;
721        for _ in 0..50 {
722            if let Ok(c) = async_nats::connect(&url).await {
723                client = Some(c);
724                break;
725            }
726            tokio::time::sleep(Duration::from_millis(100)).await;
727        }
728        Broker {
729            js: jetstream::new(client.expect("nats-server did not come up in 5s")),
730            _server: server,
731            _storage: storage,
732        }
733    }
734
735    /// #506 / 2026-06-11 incident: `create_object_store` neither
736    /// reconciles config nor has a create-or-update form, so adding
737    /// the #518 `max_bytes` cap to a store first created uncapped made
738    /// the broker reject the create (error 10058 "name already in use
739    /// with a different configuration") and crashed the backend on
740    /// boot. `ensure_object_store` must instead accept the existing
741    /// store and let startup proceed.
742    #[tokio::test]
743    #[ignore = "requires nats-server in PATH; cargo test -- --ignored"]
744    async fn ensure_object_store_accepts_config_drift() {
745        let b = spawn_broker().await;
746        // First create: uncapped, as the pre-#518 backend did.
747        ensure_object_store(
748            &b.js,
749            ObjectStoreConfig {
750                bucket: "result_output".into(),
751                ..Default::default()
752            },
753        )
754        .await
755        .expect("fresh create");
756
757        // Second create with a conflicting config (now capped) must
758        // NOT error — it accepts the existing store.
759        ensure_object_store(
760            &b.js,
761            ObjectStoreConfig {
762                bucket: "result_output".into(),
763                max_bytes: 1024 * 1024 * 1024,
764                ..Default::default()
765            },
766        )
767        .await
768        .expect("config drift must not wedge startup");
769
770        // The store is still usable.
771        let store = b.js.get_object_store("result_output").await.expect("store");
772        store
773            .put("k", &mut &b"hi"[..])
774            .await
775            .expect("put after drift");
776    }
777
778    /// A fresh create with a cap succeeds on a broker with room (the
779    /// normal first-boot path).
780    #[tokio::test]
781    #[ignore = "requires nats-server in PATH; cargo test -- --ignored"]
782    async fn ensure_object_store_fresh_create_with_cap() {
783        let b = spawn_broker().await;
784        ensure_object_store(
785            &b.js,
786            ObjectStoreConfig {
787                bucket: "fresh".into(),
788                max_bytes: 64 * 1024 * 1024,
789                ..Default::default()
790            },
791        )
792        .await
793        .expect("fresh capped create");
794        b.js.get_object_store("fresh").await.expect("exists");
795    }
796
797    /// The fatal path: when create fails for a store that ALSO does
798    /// not exist, the error must propagate (we only swallow errors we
799    /// can fall back from). An invalid bucket name fails create's
800    /// charset validation and never creates a store to fall back to.
801    #[tokio::test]
802    #[ignore = "requires nats-server in PATH; cargo test -- --ignored"]
803    async fn ensure_object_store_propagates_when_no_fallback() {
804        let b = spawn_broker().await;
805        let err = ensure_object_store(
806            &b.js,
807            ObjectStoreConfig {
808                // Spaces / '!' are rejected by the object-store name
809                // rules, so create fails and get also finds nothing.
810                bucket: "bad name!".into(),
811                ..Default::default()
812            },
813        )
814        .await
815        .expect_err("a create failure with no existing store must be fatal");
816        assert!(
817            err.to_string()
818                .contains("no existing store to fall back to"),
819            "unexpected error: {err:#}",
820        );
821    }
822
823    /// `result_output` reconciles through the SAME generalised function, so
824    /// this pins that the generalisation actually reaches the second bucket
825    /// — the failure it prevents is one bucket silently keeping the old
826    /// window while the other moves.
827    ///
828    /// Provisions the bucket the way production was BORN (30-day max_age),
829    /// then reconciles to the new 7-day default and asserts the live stream
830    /// moved. That is the case that matters: `create_object_store` does not
831    /// reconcile (#506), so every existing deployment is still at thirty days
832    /// and only this path can change it.
833    #[tokio::test]
834    #[ignore = "requires nats-server in PATH; cargo test -- --ignored"]
835    async fn reconcile_reaches_result_output_too() {
836        use crate::kv::OBJECT_RESULT_OUTPUT;
837        const SECS_PER_DAY: u64 = 24 * 60 * 60;
838        let b = spawn_broker().await;
839
840        ensure_object_store(
841            &b.js,
842            ObjectStoreConfig {
843                bucket: OBJECT_RESULT_OUTPUT.into(),
844                max_age: Duration::from_secs(SECS_PER_DAY * 30),
845                max_bytes: 1024 * 1024 * 1024,
846                ..Default::default()
847            },
848        )
849        .await
850        .expect("fresh result_output bucket at the old 30-day window");
851
852        let changed = reconcile_object_store_max_age(
853            &b.js,
854            OBJECT_RESULT_OUTPUT,
855            DEFAULT_RESULT_OUTPUT_RETENTION_DAYS,
856        )
857        .await
858        .expect("reconcile result_output max_age");
859        assert!(changed, "30d -> 7d must report a change");
860
861        let name = object_store_stream_name(OBJECT_RESULT_OUTPUT);
862        let mut stream = b.js.get_stream(&name).await.expect("get stream");
863        let info = stream.info().await.expect("stream info");
864        assert_eq!(
865            info.config.max_age,
866            Duration::from_secs(SECS_PER_DAY * DEFAULT_RESULT_OUTPUT_RETENTION_DAYS as u64)
867        );
868        // The cap is NOT collateral damage: read-modify-write must patch only
869        // max_age, or reconciling retention would silently drop the #1247 cap.
870        assert_eq!(info.config.max_bytes, 1024 * 1024 * 1024);
871
872        // Idempotent.
873        assert!(
874            !reconcile_object_store_max_age(
875                &b.js,
876                OBJECT_RESULT_OUTPUT,
877                DEFAULT_RESULT_OUTPUT_RETENTION_DAYS
878            )
879            .await
880            .expect("second reconcile"),
881            "already in sync must report no change"
882        );
883    }
884
885    /// when already in sync) and that `max_bytes` survives the update.
886    /// `reconcile_collect_retention` must change the live bucket's `max_age`
887    /// (broker-side retention) without disturbing the other stream config —
888    /// the mechanism the SPA relies on to extend collect retention past the
889    /// 30-day default. Also asserts the idempotent no-op path (`Ok(false)`
890    /// when already in sync).
891    #[tokio::test]
892    #[ignore = "requires nats-server in PATH; cargo test -- --ignored"]
893    async fn reconcile_collect_retention_updates_max_age() {
894        use crate::kv::OBJECT_COLLECTIONS;
895        const SECS_PER_DAY: u64 = 24 * 60 * 60;
896        let b = spawn_broker().await;
897
898        // Provision the collections bucket the way bootstrap does: 30-day
899        // default max_age, 5 GiB cap.
900        ensure_object_store(
901            &b.js,
902            ObjectStoreConfig {
903                bucket: OBJECT_COLLECTIONS.into(),
904                max_age: Duration::from_secs(SECS_PER_DAY * 30),
905                max_bytes: 5 * 1024 * 1024 * 1024,
906                ..Default::default()
907            },
908        )
909        .await
910        .expect("fresh collections bucket");
911
912        let stream_name = object_store_stream_name(OBJECT_COLLECTIONS);
913
914        // Extend to 90 days — first call changes the stream.
915        assert!(
916            reconcile_collect_retention(&b.js, 90)
917                .await
918                .expect("reconcile to 90d"),
919            "first reconcile should report a change",
920        );
921        let mut stream = b.js.get_stream(&stream_name).await.expect("stream");
922        let info = stream.info().await.expect("info");
923        assert_eq!(
924            info.config.max_age,
925            Duration::from_secs(SECS_PER_DAY * 90),
926            "max_age must be extended to 90 days",
927        );
928        assert_eq!(
929            info.config.max_bytes,
930            5 * 1024 * 1024 * 1024,
931            "the size cap must survive the max_age-only update",
932        );
933
934        // Re-applying the same value is a no-op (no revision-bumping update).
935        assert!(
936            !reconcile_collect_retention(&b.js, 90)
937                .await
938                .expect("idempotent reconcile"),
939            "second reconcile with the same value should be a no-op",
940        );
941    }
942
943    /// `reconcile_object_store_max_bytes` must deliver a cap to a bucket
944    /// that was provisioned WITHOUT one — the exact drift case #1247
945    /// fixes (`OBJ_result_output` at 6.76 GB against a nominal 1 GiB,
946    /// because a cap added in code after the bucket existed never reached
947    /// the broker). Asserts the backing stream's `max_bytes` changes, the
948    /// `max_age` window survives the max_bytes-only update, and the
949    /// idempotent no-op path. Mirrors
950    /// [`reconcile_collect_retention_updates_max_age`].
951    #[tokio::test]
952    #[ignore = "requires nats-server in PATH; cargo test -- --ignored"]
953    async fn reconcile_object_store_max_bytes_updates_cap() {
954        use crate::kv::OBJECT_RESULT_OUTPUT;
955        const SECS_PER_DAY: u64 = 24 * 60 * 60;
956        let b = spawn_broker().await;
957
958        // Provision the way the pre-#518 bucket was born: a 30-day window
959        // but NO size cap (the drifted production state).
960        ensure_object_store(
961            &b.js,
962            ObjectStoreConfig {
963                bucket: OBJECT_RESULT_OUTPUT.into(),
964                max_age: Duration::from_secs(SECS_PER_DAY * 30),
965                ..Default::default()
966            },
967        )
968        .await
969        .expect("fresh uncapped result_output bucket");
970
971        let stream_name = object_store_stream_name(OBJECT_RESULT_OUTPUT);
972
973        // Apply the 1 GiB cap — first call changes the stream.
974        assert!(
975            reconcile_object_store_max_bytes(&b.js, OBJECT_RESULT_OUTPUT, 1024)
976                .await
977                .expect("reconcile to 1024 MiB"),
978            "first reconcile should report a change",
979        );
980        let mut stream = b.js.get_stream(&stream_name).await.expect("stream");
981        let info = stream.info().await.expect("info");
982        assert_eq!(
983            info.config.max_bytes,
984            1024 * 1024 * 1024,
985            "max_bytes must land on the previously-uncapped stream",
986        );
987        assert_eq!(
988            info.config.max_age,
989            Duration::from_secs(SECS_PER_DAY * 30),
990            "the retention window must survive the max_bytes-only update",
991        );
992
993        // Re-applying the same value is a no-op (no revision-bumping update).
994        assert!(
995            !reconcile_object_store_max_bytes(&b.js, OBJECT_RESULT_OUTPUT, 1024)
996                .await
997                .expect("idempotent reconcile"),
998            "second reconcile with the same value should be a no-op",
999        );
1000
1001        // A different value reconciles again (operator tune-down/tune-up).
1002        assert!(
1003            reconcile_object_store_max_bytes(&b.js, OBJECT_RESULT_OUTPUT, 2048)
1004                .await
1005                .expect("reconcile to 2048 MiB"),
1006            "changing the cap should report a change",
1007        );
1008    }
1009}