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