hyperhive/swarm-queue-client/src/status.rs
atlas 7b5f383b05 fix(#3297): reject a publish template that names no hive
`--hive-publish-subject` exists to put a second stream inside one hive's
namespace. A template with no `{hive}` in it expands to the same subject
for every hive, so the option whose only purpose is scoping becomes the
way to remove it — silently, and only in the deployment that set it.

`Policy::new` returns a `Result` rather than checking at the call site:
that makes an unscoped policy unconstructible instead of merely
unlikely, the same reason `grant` takes its permissions by value. The
error names the offending template and says what goes wrong with it,
because an operator meets it at boot with no other context.

Also documents what the prefix match does not do. A client id is a hive
here because it starts with the configured prefix, not because it
appears in the roster — the responder runs in a container and cannot see
`swarm.hives`. Passing the roster in would close that and would also be
a second place deciding who may connect as what, which `introspect`'s
docs argue against for the same reason admission lives in one place.

The two intra-doc links to `open_or_create` become plain backticks.
Un-gating the `status` module means its module doc now renders in builds
without the `kv` feature, where the item it linked does not exist.
2026-08-17 17:34:27 +02:00

75 lines
3.4 KiB
Rust

//! The hive-status KV bucket: its name, and the shape it is created with.
//!
//! Two processes touch this bucket from opposite ends — a hive writes its
//! own key, the swarm controller reads every key — and they live in
//! different crates. That is the whole reason this module exists rather
//! than a `const` on each side: **the two ends must agree, and a literal
//! repeated across crates is an agreement nothing checks.**
//!
//! The name is the obvious half. The sharper half is the *config*: both
//! ends open the bucket with `open_or_create`, because either may
//! arrive first on a fresh swarm and neither can assume the other has
//! run. If the two ends passed different `Config`s, whichever created it
//! would win and the other's `get_key_value` would succeed against a
//! bucket it did not ask for — no error, no log, just a retention policy
//! nobody chose. Sharing the constructor makes the race have one outcome
//! instead of two.
//!
//! The bucket *name* is unconditional; only `open_or_create` is behind the
//! `kv` feature. (Named in backticks rather than linked: with `kv` off the item
//! does not exist, and an intra-doc link to it fails the rustdoc gate in
//! exactly the configuration this split exists to support.) A third end names the bucket without ever opening it — the
//! auth-callout responder, which derives the subjects a hive may publish to
//! from it — and it speaks neither `jetstream` nor `kv`. Gating the name too
//! would have made that consumer choose between a JetStream stack it does not
//! use and a copied literal, and a copied literal is precisely the agreement
//! nothing checks.
#[cfg(feature = "kv")]
use crate::Error;
/// The KV bucket hives publish their status snapshots into, one key per
/// hive keyed by `hiveName`.
///
/// A constant and not an option: reader and writer must name the same
/// bucket, and an option is a way for two deployments to disagree about
/// which one that is. Nothing about a bucket name is site-specific.
pub const BUCKET: &str = "hive-status";
/// Open the status bucket, creating it if nothing has yet.
///
/// `history: 1` is the shape: every consumer reads *the last thing each
/// hive said*, and retaining more would be storage bought for a query
/// nobody makes.
///
/// Creating rather than requiring a provisioning step is deliberate — the
/// controller and the hives come up in no particular order, and a bucket
/// that must pre-exist turns "the swarm was deployed in the wrong order"
/// into a permanent, silent absence of data.
#[cfg(feature = "kv")]
pub async fn open_or_create(
client: &async_nats::Client,
) -> Result<async_nats::jetstream::kv::Store, Error> {
let js = async_nats::jetstream::new(client.clone());
match js.get_key_value(BUCKET).await {
Ok(store) => Ok(store),
Err(e) => {
tracing::info!(
bucket = BUCKET,
reason = %e,
"status bucket not available, creating it"
);
js.create_key_value(async_nats::jetstream::kv::Config {
bucket: BUCKET.to_owned(),
description: "Last status snapshot offered by each hive".to_owned(),
history: 1,
..Default::default()
})
.await
.map_err(|source| Error::CreateBucket {
bucket: BUCKET,
source,
})
}
}
}