use std::path::PathBuf;
use std::sync::atomic::{AtomicBool, AtomicI64, Ordering};
use std::sync::Arc;
use std::time::{Duration, SystemTime, UNIX_EPOCH};
use secrecy::ExposeSecret;
use crate::cloud::worker::license_refused_from;
use crate::cloud::{
body_code_is_site_license, parse_retry_after, CloudError, CloudState, CredentialProvider,
MAX_RETRY_AFTER_SECS,
};
use crate::config::PolicyConfig;
use crate::core::error::{
ERR_BUNDLE_FETCH_FAILED, ERR_BUNDLE_INVALID, ERR_BUNDLE_REJECTED, ERR_BUNDLE_STALE,
ERR_CLOUD_LICENSE_BLOCKED,
};
use crate::core::policy::{project_document, store, PolicyHandle};
use crate::core::protocol::contracts::{
self, CompatibilityDiagnostic, VersionRange, POLICY_BUNDLE_FAMILY,
};
use crate::daemon::{PolicyReadiness, PolicyReadinessPhase};
pub const SUPPORTED_SCHEMA_VERSIONS: &[i64] = &[1, 2];
const CLIENT_CAPABILITIES: &str = "t1_predicate_tree,t2_register_program,exception";
const CLIENT_CAPABILITIES_WITH_HOLDS: &str =
"t1_predicate_tree,t2_register_program,t3_hold,exception";
const BUNDLE_ACCEPT: &str = "application/vnd.openlatch.bundle+json; schema=2";
const BUNDLE_DELIVERY_HEADER: &str = "OpenLatch-Bundle-Delivery";
const BUNDLE_DELIVERY_PENDING_V1: &str = "pending-v1";
const FLOOR_PROBLEM_TYPE: &str = "https://openlatch.ai/problems/client-below-floor";
const BUNDLE_ENDPOINT: &str = "/api/v1/policy/bundle";
use crate::core::cloud::envelope::AGENT_ID_HEADER;
const MACHINE_ID_HEADER: &str = "X-OpenLatch-Machine-Id";
const MIN_POLL_INTERVAL_SECS: u64 = 1;
#[derive(Debug, Clone, PartialEq, Eq)]
enum PollOutcome {
Activated { revision: i64 },
NotModified,
Pending(Duration),
Rejected(&'static str),
NoBundle,
AuthFailed,
FloorBlocked,
CompatibilityUnavailable { family: String },
RateLimited(Option<Duration>),
LicenseBlocked(Option<Duration>),
Failed,
Skipped(&'static str),
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum PolicyRefreshOutcome {
Changed,
Unchanged,
Pending,
Refused,
Offline,
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub struct PolicyRefreshResult {
pub outcome: PolicyRefreshOutcome,
#[serde(skip_serializing_if = "Option::is_none")]
pub response_status: Option<u16>,
#[serde(skip_serializing_if = "Option::is_none")]
pub reason: Option<String>,
pub has_bundle: bool,
#[serde(skip_serializing_if = "Option::is_none")]
pub selected_version: Option<u32>,
#[serde(skip_serializing_if = "Option::is_none")]
pub schema_version: Option<i64>,
#[serde(skip_serializing_if = "Option::is_none")]
pub bundle_revision: Option<i64>,
#[serde(skip_serializing_if = "Option::is_none")]
pub digest: Option<String>,
pub verified_body_received: bool,
#[serde(skip_serializing_if = "Option::is_none")]
pub verified_body_received_at: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub last_platform_contact_at: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub resident_aip_count: Option<usize>,
}
pub struct PolicyRefreshRequest {
pub(crate) reply: tokio::sync::oneshot::Sender<PolicyRefreshResult>,
}
pub type PolicyRefreshRequestReceiver =
Arc<tokio::sync::Mutex<tokio::sync::mpsc::Receiver<PolicyRefreshRequest>>>;
#[allow(clippy::too_many_arguments)]
pub async fn run_policy_poller(
handle: PolicyHandle,
last_fetch_ok: Arc<AtomicBool>,
last_poll_ok_at: Arc<AtomicI64>,
last_download_verified_at: Arc<AtomicI64>,
readiness: Arc<std::sync::RwLock<PolicyReadiness>>,
floor: Arc<std::sync::RwLock<Option<crate::daemon::ClientFloorState>>>,
cloud_state: CloudState,
credentials: Arc<dyn CredentialProvider>,
api_url: String,
policy_config: PolicyConfig,
base_dir: PathBuf,
http_client: crate::egress::ClientHandle,
agent_id: Option<String>,
host_key: Option<String>,
egress: crate::egress::EgressReporter,
refresh_requests: Option<PolicyRefreshRequestReceiver>,
) {
let mut poller = PolicyPoller::new(
handle,
last_fetch_ok,
last_poll_ok_at,
last_download_verified_at,
readiness,
cloud_state,
credentials,
api_url,
policy_config,
base_dir,
http_client,
agent_id,
host_key,
)
.with_floor(floor)
.with_hold_capability()
.with_egress(egress);
poller.seed_poll_clock();
let mut refresh_requests = match refresh_requests {
Some(receiver) => Some(receiver.lock_owned().await),
None => None,
};
let mut outcome = poller.poll_once().await;
poller.check_staleness();
loop {
let mut delay = jittered(poller.config.poll_interval_secs.max(MIN_POLL_INTERVAL_SECS));
if let PollOutcome::RateLimited(Some(retry_after))
| PollOutcome::LicenseBlocked(Some(retry_after)) = outcome
{
if retry_after > delay {
delay = retry_after;
}
}
if let PollOutcome::Pending(retry_after) = outcome {
delay = jittered(retry_after.as_secs().max(1))
.clamp(Duration::from_secs(2), Duration::from_secs(10));
}
let mut manual_reply = None;
tokio::select! {
_ = tokio::time::sleep(delay) => {}
_ = poller.cloud_state.policy_refresh_notify.notified() => {
tracing::debug!(target: "policy", "compatibility refresh requested by event egress");
}
_ = poller.cloud_state.auth_refresh_notify().notified() => {
tracing::debug!(target: "policy", "auth command refreshed the credential source");
}
request = async {
match refresh_requests.as_mut() {
Some(receiver) => receiver.recv().await,
None => std::future::pending().await,
}
} => {
if let Some(request) = request {
manual_reply = Some(request.reply);
tracing::debug!(target: "policy", "manual AIP bundle refresh requested");
}
}
}
let previous_digest = manual_reply.as_ref().and_then(|_| poller.resident_digest());
outcome = poller.poll_once().await;
poller.check_staleness();
if let Some(reply) = manual_reply {
let _ = reply.send(poller.refresh_result(&outcome, previous_digest.as_deref()));
}
}
}
fn problem_range(value: &serde_json::Value) -> Option<VersionRange> {
if let Some(raw) = value.as_str() {
let parsed = contracts::parse_ranges(&format!("peer={raw}")).ok()?;
return parsed.get("peer").copied();
}
if let Some(values) = value.as_array() {
let oldest = u32::try_from(values.first()?.as_u64()?).ok()?;
let newest = u32::try_from(values.get(1)?.as_u64()?).ok()?;
return VersionRange::new(oldest, newest);
}
let oldest = u32::try_from(value.get("oldest")?.as_u64()?).ok()?;
let newest = u32::try_from(value.get("newest")?.as_u64()?).ok()?;
VersionRange::new(oldest, newest)
}
fn sanitize_problem_detail(raw: &str) -> String {
raw.chars()
.filter(|ch| !ch.is_control())
.take(240)
.collect::<String>()
.trim()
.to_string()
}
fn jittered(base: u64) -> Duration {
let b = uuid::Uuid::new_v4().as_bytes()[0] as u64; let span = base.saturating_mul(20) / 100; let offset = (b * span) / 255;
Duration::from_secs(base.saturating_sub(span / 2).saturating_add(offset))
}
fn parse_etag(raw: &str) -> Option<&str> {
let t = raw.trim();
let t = t.strip_prefix("W/").unwrap_or(t).trim();
t.strip_prefix('"')?.strip_suffix('"')
}
struct PolicyPoller {
handle: PolicyHandle,
last_fetch_ok: Arc<AtomicBool>,
last_poll_ok_at: Arc<AtomicI64>,
last_download_verified_at: Arc<AtomicI64>,
readiness: Arc<std::sync::RwLock<PolicyReadiness>>,
cloud_state: CloudState,
credentials: Arc<dyn CredentialProvider>,
url: String,
config: PolicyConfig,
base_dir: PathBuf,
http: crate::egress::ClientHandle,
agent_id: Option<String>,
capabilities: &'static str,
failed_credential: Option<Vec<u8>>,
auth_refresh_generation: u64,
host_key: Option<String>,
floor: Arc<std::sync::RwLock<Option<crate::daemon::ClientFloorState>>>,
egress: crate::egress::EgressReporter,
}
impl PolicyPoller {
#[allow(clippy::too_many_arguments)]
fn new(
handle: PolicyHandle,
last_fetch_ok: Arc<AtomicBool>,
last_poll_ok_at: Arc<AtomicI64>,
last_download_verified_at: Arc<AtomicI64>,
readiness: Arc<std::sync::RwLock<PolicyReadiness>>,
cloud_state: CloudState,
credentials: Arc<dyn CredentialProvider>,
api_url: String,
config: PolicyConfig,
base_dir: PathBuf,
http: crate::egress::ClientHandle,
agent_id: Option<String>,
host_key: Option<String>,
) -> Self {
let agent_id = agent_id.filter(|id| {
reqwest::header::HeaderValue::from_str(id)
.inspect_err(|_| {
tracing::warn!(
target: "policy",
agent_id = ?id,
"configured agent_id is not a usable HTTP header value; polling without it. The platform serves the org bundle with no agent_context, so every scoped rule matches nothing"
);
})
.is_ok()
});
let auth_refresh_generation = cloud_state.auth_refresh_generation();
Self {
handle,
last_fetch_ok,
last_poll_ok_at,
last_download_verified_at,
readiness,
cloud_state,
credentials,
url: format!("{}{}", api_url.trim_end_matches('/'), BUNDLE_ENDPOINT),
config,
base_dir,
http,
agent_id,
host_key,
capabilities: CLIENT_CAPABILITIES,
failed_credential: None,
auth_refresh_generation,
floor: Arc::new(std::sync::RwLock::new(None)),
egress: crate::egress::EgressReporter::direct(),
}
}
fn with_floor(
mut self,
floor: Arc<std::sync::RwLock<Option<crate::daemon::ClientFloorState>>>,
) -> Self {
self.floor = floor;
self
}
fn with_hold_capability(mut self) -> Self {
self.capabilities = CLIENT_CAPABILITIES_WITH_HOLDS;
self
}
fn with_egress(mut self, egress: crate::egress::EgressReporter) -> Self {
self.egress = egress;
self
}
fn resident_digest(&self) -> Option<String> {
self.handle
.load()
.as_ref()
.as_ref()
.and_then(|bundle| bundle.digest.clone())
}
fn refresh_result(
&self,
outcome: &PollOutcome,
previous_digest: Option<&str>,
) -> PolicyRefreshResult {
let current_digest = self.resident_digest();
let (refresh_outcome, response_status, reason) = match outcome {
PollOutcome::Activated { .. } => (
if current_digest.as_deref() == previous_digest {
PolicyRefreshOutcome::Unchanged
} else {
PolicyRefreshOutcome::Changed
},
Some(200),
None,
),
PollOutcome::NotModified => (PolicyRefreshOutcome::Unchanged, Some(304), None),
PollOutcome::Pending(_) => (
PolicyRefreshOutcome::Pending,
Some(202),
Some("install_materializing".to_string()),
),
PollOutcome::Rejected(code) => (
PolicyRefreshOutcome::Refused,
None,
Some((*code).to_string()),
),
PollOutcome::NoBundle => (
PolicyRefreshOutcome::Refused,
Some(404),
Some("bundle_not_found".to_string()),
),
PollOutcome::AuthFailed => (
PolicyRefreshOutcome::Refused,
Some(401),
Some("auth_rejected".to_string()),
),
PollOutcome::FloorBlocked => (
PolicyRefreshOutcome::Refused,
Some(403),
Some("client_below_floor".to_string()),
),
PollOutcome::CompatibilityUnavailable { family } => (
PolicyRefreshOutcome::Refused,
Some(409),
Some(format!("compatibility_unavailable:{family}")),
),
PollOutcome::RateLimited(_) => (
PolicyRefreshOutcome::Pending,
Some(429),
Some("rate_limited".to_string()),
),
PollOutcome::LicenseBlocked(_) => (
PolicyRefreshOutcome::Refused,
None,
Some("license_refused".to_string()),
),
PollOutcome::Failed => (
PolicyRefreshOutcome::Offline,
None,
Some("transport_error".to_string()),
),
PollOutcome::Skipped(
reason @ ("auth_error_latched" | "credential_unchanged_after_auth_failure"),
) => (
PolicyRefreshOutcome::Refused,
None,
Some((*reason).to_string()),
),
PollOutcome::Skipped(reason @ ("no_credential" | "no_egress_route")) => (
PolicyRefreshOutcome::Offline,
None,
Some((*reason).to_string()),
),
PollOutcome::Skipped(reason) => (
PolicyRefreshOutcome::Offline,
None,
Some((*reason).to_string()),
),
};
let resident = self.handle.load_full();
let (has_bundle, schema_version, bundle_revision, digest, resident_aip_count) =
match resident.as_ref() {
Some(bundle) => (
true,
Some(bundle.schema_version),
Some(bundle.revision),
current_digest,
bundle.aip_count(),
),
None => (false, None, None, None, None),
};
let selection_acknowledged = self
.readiness
.read()
.is_ok_and(|state| state.selection_acknowledged);
let selected_version = selection_acknowledged
.then(|| contracts::read_compatibility(&self.base_dir).ok())
.flatten()
.and_then(|state| {
state
.selections
.get(POLICY_BUNDLE_FAMILY)
.map(|selection| selection.version)
});
let verified_body_received_at =
timestamp(self.last_download_verified_at.load(Ordering::Relaxed));
PolicyRefreshResult {
outcome: refresh_outcome,
response_status,
reason,
has_bundle,
selected_version,
schema_version,
bundle_revision,
digest,
verified_body_received: verified_body_received_at.is_some(),
verified_body_received_at,
last_platform_contact_at: timestamp(self.last_poll_ok_at.load(Ordering::Relaxed)),
resident_aip_count,
}
}
fn set_readiness(
&self,
phase: PolicyReadinessPhase,
outcome: &'static str,
reason_code: Option<&'static str>,
) {
if let Ok(mut state) = self.readiness.write() {
state.phase = phase;
state.last_outcome = Some(outcome);
state.reason_code = reason_code;
}
}
fn note_attempt(&self) {
if let Ok(mut state) = self.readiness.write() {
state.last_attempted_at = Some(now_unix_secs());
}
}
fn set_selection_acknowledged(&self, acknowledged: bool) {
if let Ok(mut state) = self.readiness.write() {
state.selection_acknowledged = acknowledged;
}
}
fn set_compatibility_diagnostic_valid(&self, valid: bool) {
if let Ok(mut state) = self.readiness.write() {
state.compatibility_diagnostic_valid = valid;
}
}
fn record_successful_contact(&self, mut meta: Option<store::BundleMeta>) {
let now = crate::install_state::now_rfc3339();
self.last_poll_ok_at.store(
parse_unix_secs(&now).unwrap_or_else(now_unix_secs),
Ordering::Relaxed,
);
if let Some(meta) = meta.as_mut() {
meta.last_poll_ok_at = Some(now);
if let Err(error) = store::write_meta(&self.base_dir, meta) {
tracing::debug!(target: "policy", %error, "could not persist the policy contact clock");
}
}
}
fn seed_poll_clock(&self) {
let Ok(Some(meta)) = store::read_meta(&self.base_dir) else {
return;
};
if let Some(secs) = meta.last_poll_ok_at.as_deref().and_then(parse_unix_secs) {
self.last_poll_ok_at.store(secs, Ordering::Relaxed);
}
self.last_fetch_ok
.store(meta.last_fetch_ok, Ordering::Relaxed);
}
async fn poll_once(&mut self) -> PollOutcome {
let auth_refresh_generation = self.cloud_state.auth_refresh_generation();
if auth_refresh_generation != self.auth_refresh_generation {
self.auth_refresh_generation = auth_refresh_generation;
if self.failed_credential.take().is_some() {
tracing::debug!(
target: "policy",
"explicit auth refresh released the policy credential refusal latch"
);
}
}
if self.cloud_state.is_auth_error() {
tracing::debug!(
target: "policy",
"cloud auth error is latched; skipping the policy poll (the resident bundle keeps enforcing)"
);
self.set_readiness(
PolicyReadinessPhase::TerminalRefusal,
"skipped",
Some("auth_error_latched"),
);
return PollOutcome::Skipped("auth_error_latched");
}
let Some(token) = self.credentials.retrieve() else {
tracing::debug!(
target: "policy",
"no credential available; skipping the policy poll"
);
self.set_readiness(
PolicyReadinessPhase::Retrying,
"skipped",
Some("no_credential"),
);
return PollOutcome::Skipped("no_credential");
};
let credential = token.expose_secret().as_bytes().to_vec();
if self.failed_credential.as_deref() == Some(credential.as_slice()) {
tracing::debug!(
target: "policy",
"credential unchanged since the last 401/403; skipping the policy poll"
);
self.set_readiness(
PolicyReadinessPhase::TerminalRefusal,
"skipped",
Some("auth_rejected"),
);
return PollOutcome::Skipped("credential_unchanged_after_auth_failure");
}
let meta = match store::read_meta(&self.base_dir) {
Ok(meta) => meta,
Err(e) => {
tracing::warn!(
target: "policy",
code = ERR_BUNDLE_FETCH_FAILED,
error = %e,
"could not read bundle.meta.json; polling without an If-None-Match validator"
);
None
}
};
let Some(http) = self.http.current() else {
tracing::debug!(
target: "policy",
code = crate::error::ERR_DIRECT_FORBIDDEN,
"no egress route is permitted; skipping the policy poll (the resident bundle keeps enforcing)"
);
self.set_readiness(
PolicyReadinessPhase::Retrying,
"transport_unavailable",
Some("no_egress_route"),
);
return PollOutcome::Skipped("no_egress_route");
};
let user_agent = format!(
"openlatch-client/{} ({}; {}; mode1)",
env!("CARGO_PKG_VERSION"),
std::env::consts::OS,
std::env::consts::ARCH
);
let mut req = http
.get(&self.url)
.bearer_auth(token.expose_secret())
.header(reqwest::header::USER_AGENT, user_agent)
.header("OpenLatch-Client-Version", env!("CARGO_PKG_VERSION"))
.header("OpenLatch-Client-Capabilities", self.capabilities)
.header(contracts::SUPPORT_HEADER, contracts::support_header_value())
.header(BUNDLE_DELIVERY_HEADER, BUNDLE_DELIVERY_PENDING_V1)
.header(reqwest::header::ACCEPT, BUNDLE_ACCEPT);
let validator = meta
.as_ref()
.filter(|_| self.floor.read().is_ok_and(|floor| floor.is_none()))
.filter(|m| m.state.as_deref() != Some("license_blocked"))
.filter(|m| m.agent_id.as_deref() == self.agent_id.as_deref())
.and_then(|m| m.etag.as_deref());
if let Some(etag) = validator {
req = req.header(reqwest::header::IF_NONE_MATCH, etag);
}
if let Some(agent_id) = self.agent_id.as_deref() {
req = req.header(AGENT_ID_HEADER, agent_id);
}
if let Some(host_key) = self.host_key.as_deref() {
req = req.header(MACHINE_ID_HEADER, host_key);
}
self.note_attempt();
let resp = match req.send().await {
Ok(resp) => {
self.egress.record_ok(&self.url);
resp
}
Err(e) => {
self.egress.record_failure(&self.url, &e);
return self.transport_failure(&format!("policy bundle request failed: {e}"));
}
};
let status = resp.status();
if status.as_u16() != 401 {
self.failed_credential = None;
}
match status.as_u16() {
200 => self.handle_ok(resp, meta).await,
202 => self.handle_pending(resp, meta).await,
304 => self.handle_not_modified(resp, meta),
400 => self.handle_bad_request(resp).await,
401 => {
self.failed_credential = Some(credential);
tracing::error!(
target: "policy",
code = ERR_BUNDLE_FETCH_FAILED,
status = status.as_u16(),
"policy bundle fetch rejected the credential; pausing bundle refresh until the credential changes. The resident bundle keeps enforcing"
);
self.record_poll_failure(&format!(
"{ERR_BUNDLE_FETCH_FAILED} auth rejected ({})",
status.as_u16()
));
self.set_readiness(
PolicyReadinessPhase::TerminalRefusal,
"http_401",
Some("auth_rejected"),
);
PollOutcome::AuthFailed
}
403 => self.handle_forbidden(resp).await,
409 => self.handle_compatibility_unavailable(resp).await,
404 => {
let resident = self.handle.load().is_some();
tracing::warn!(
target: "policy",
code = ERR_BUNDLE_FETCH_FAILED,
resident_bundle = resident,
"policy bundle endpoint returned 404; keeping the last-known-good bundle enforcing"
);
self.record_poll_failure(&format!("{ERR_BUNDLE_FETCH_FAILED} 404 no bundle"));
self.set_readiness(
PolicyReadinessPhase::TerminalRefusal,
"http_404",
Some("bundle_not_found"),
);
PollOutcome::NoBundle
}
429 => {
let retry_after = resp
.headers()
.get(reqwest::header::RETRY_AFTER)
.and_then(|v| v.to_str().ok())
.and_then(parse_retry_after)
.map(|s| Duration::from_secs(s.min(MAX_RETRY_AFTER_SECS)));
tracing::warn!(
target: "policy",
code = ERR_BUNDLE_FETCH_FAILED,
retry_after_secs = retry_after.map(|d| d.as_secs()),
"policy bundle fetch rate limited (429); backing off"
);
self.record_poll_failure(&format!("{ERR_BUNDLE_FETCH_FAILED} rate limited"));
self.set_readiness(
PolicyReadinessPhase::Retrying,
"http_429",
Some("rate_limited"),
);
PollOutcome::RateLimited(retry_after)
}
402 => self.handle_license_refused(resp).await,
503 => {
let headers = resp.headers().clone();
let body = resp.text().await.unwrap_or_default();
if body_code_is_site_license(&body) {
self.record_license_block(&headers, &body)
} else {
self.transport_failure("policy bundle fetch returned 503")
}
}
code => self.transport_failure(&format!(
"policy bundle fetch returned an unusable status {code}"
)),
}
}
async fn handle_ok(
&mut self,
resp: reqwest::Response,
meta: Option<store::BundleMeta>,
) -> PollOutcome {
let acknowledged = match self.selected_policy_version(resp.headers()) {
Ok(version) => version,
Err(outcome) => return outcome,
};
let raw_etag = resp
.headers()
.get(reqwest::header::ETAG)
.and_then(|v| v.to_str().ok())
.map(|s| s.to_string());
let body = match resp.bytes().await {
Ok(b) => b,
Err(e) => {
return self.transport_failure(&format!("policy bundle body was not readable: {e}"))
}
};
let Some(raw_etag) = raw_etag else {
return self.reject(
ERR_BUNDLE_REJECTED,
"200 response carried no ETag header, so the body's digest cannot be verified",
);
};
let Some(tag) = parse_etag(&raw_etag) else {
return self.reject(
ERR_BUNDLE_REJECTED,
"ETag header is not a well-formed entity-tag",
);
};
if let Err(e) = store::verify_digest(&body, tag) {
return self.reject(ERR_BUNDLE_REJECTED, &e.to_string());
}
let value: serde_json::Value = match serde_json::from_slice(&body) {
Ok(v) => v,
Err(e) => {
return self.reject(ERR_BUNDLE_INVALID, &format!("body is not valid JSON: {e}"))
}
};
let declared = value.get("schema_version").and_then(|v| v.as_i64());
let Some(declared) = declared.filter(|version| SUPPORTED_SCHEMA_VERSIONS.contains(version))
else {
return self.reject(
ERR_BUNDLE_INVALID,
&format!("unsupported schema_version {declared:?}"),
);
};
if acknowledged.is_none()
&& !contracts::policy_bundle_can_read(
&value,
u32::try_from(declared).unwrap_or_default(),
)
{
return self.compatibility_reject(
None,
&format!(
"legacy response schema_version {declared} cannot faithfully represent its policy facts"
),
);
}
if acknowledged.is_some_and(|selected| u32::try_from(declared).ok() != Some(selected)) {
let selected = acknowledged.expect("checked Some above");
return self.compatibility_reject(
None,
&format!(
"selected {POLICY_BUNDLE_FAMILY} v{selected} but body declared schema_version {declared}"
),
);
}
if acknowledged.is_some_and(|selected| !contracts::policy_bundle_can_read(&value, selected))
{
let selected = acknowledged.expect("checked Some above");
return self.compatibility_reject(
None,
&format!(
"selected {POLICY_BUNDLE_FAMILY} v{selected} cannot faithfully represent compiled policy facts in the response"
),
);
}
let expected_org: Option<String> = meta
.as_ref()
.map(|m| m.organization_id.clone())
.or_else(|| {
self.handle
.load()
.as_ref()
.as_ref()
.map(|b| b.organization_id.clone())
});
if let Some(expected) = expected_org {
let candidate = value
.get("organization_id")
.and_then(serde_json::Value::as_str);
if candidate != Some(expected.as_str()) {
return self.reject(
ERR_BUNDLE_REJECTED,
&format!(
"organization_id {:?} does not match this host's {expected}",
candidate
),
);
}
}
if value
.get("signature")
.is_some_and(|signature| !signature.is_null())
{
return self.reject(
ERR_BUNDLE_INVALID,
"bundle carries a signature this client cannot verify (D32 lands in v1.1)",
);
}
let (bundle, mut resident) = match project_document(value) {
Ok(projected) => projected,
Err(e) => {
return self.reject(ERR_BUNDLE_INVALID, &format!("body projection failed: {e}"))
}
};
let digest = tag.to_string();
resident.digest = Some(digest.clone());
let mut new_meta = store::BundleMeta::activated(&bundle, digest.clone(), Some(raw_etag));
let verified_at = new_meta
.last_download_verified_at
.as_deref()
.and_then(parse_unix_secs)
.unwrap_or_else(now_unix_secs);
new_meta.agent_id = self.agent_id.clone();
if let Err(e) = store::store(&self.base_dir, &body, &new_meta) {
tracing::warn!(
target: "policy",
code = ERR_BUNDLE_FETCH_FAILED,
error = %e,
"could not persist the policy bundle; activating it in memory anyway"
);
}
let revision = bundle.revision;
let rule_count = resident.command_rules.len();
let request_rule_count = resident.request_rules.len();
let enforcement_enabled = resident.enforcement_enabled;
self.handle.store(Arc::new(Some(resident)));
if let Ok(mut floor) = self.floor.write() {
*floor = None;
}
self.cloud_state.set_policy_license_refusal(None);
self.last_fetch_ok.store(true, Ordering::Relaxed);
self.last_poll_ok_at
.store(now_unix_secs(), Ordering::Relaxed);
self.last_download_verified_at
.store(verified_at, Ordering::Relaxed);
if let Some(selected) = acknowledged {
let persisted = self.record_selection(POLICY_BUNDLE_FAMILY, selected);
self.set_selection_acknowledged(persisted);
self.set_compatibility_diagnostic_valid(false);
self.set_readiness(PolicyReadinessPhase::Resident, "http_200", None);
} else {
self.set_selection_acknowledged(false);
self.set_compatibility_diagnostic_valid(false);
self.set_readiness(
PolicyReadinessPhase::Resident,
"http_200",
Some("legacy_no_ack"),
);
self.record_legacy_response(POLICY_BUNDLE_FAMILY);
}
tracing::info!(
target: "policy",
revision,
digest = %digest,
rules = rule_count,
request_rules = request_rule_count,
enforcement_enabled,
"policy bundle activated"
);
PollOutcome::Activated { revision }
}
async fn handle_pending(
&self,
resp: reqwest::Response,
meta: Option<store::BundleMeta>,
) -> PollOutcome {
#[derive(serde::Deserialize)]
#[serde(deny_unknown_fields)]
struct PendingBody {
status: String,
reason: String,
}
let retry_after = resp
.headers()
.get(reqwest::header::RETRY_AFTER)
.and_then(|value| value.to_str().ok())
.and_then(parse_retry_after)
.unwrap_or(2)
.clamp(2, 10);
let body = match resp.json::<PendingBody>().await {
Ok(body)
if body.status == "bundle_pending" && body.reason == "install_materializing" =>
{
body
}
Ok(_) | Err(_) => {
return self.reject(
ERR_BUNDLE_INVALID,
"202 response did not carry the documented bundle_pending shape",
)
}
};
debug_assert_eq!(body.status, "bundle_pending");
self.last_fetch_ok.store(false, Ordering::Relaxed);
self.record_successful_contact(meta);
self.set_compatibility_diagnostic_valid(false);
self.record_pending_response(POLICY_BUNDLE_FAMILY);
let phase = if self
.readiness
.read()
.is_ok_and(|state| state.last_outcome.is_none())
{
PolicyReadinessPhase::Waiting
} else {
PolicyReadinessPhase::Retrying
};
self.set_readiness(phase, "http_202", Some("install_materializing"));
tracing::debug!(target: "policy", retry_after_secs = retry_after, "policy bundle delivery is pending; retrying promptly");
PollOutcome::Pending(Duration::from_secs(retry_after))
}
async fn handle_bad_request(&self, resp: reqwest::Response) -> PollOutcome {
let value = resp.json::<serde_json::Value>().await.ok();
let code = value
.as_ref()
.and_then(|value| value.pointer("/error/code").or_else(|| value.get("code")))
.and_then(serde_json::Value::as_str);
if code == Some("install_id_required") {
self.record_poll_failure(&format!("{ERR_BUNDLE_FETCH_FAILED} install_id_required"));
self.set_readiness(
PolicyReadinessPhase::TerminalRefusal,
"http_400",
Some("install_id_required"),
);
return PollOutcome::Rejected(ERR_BUNDLE_FETCH_FAILED);
}
self.reject(
ERR_BUNDLE_FETCH_FAILED,
"policy bundle returned an unusable 400",
)
}
async fn handle_forbidden(&self, resp: reqwest::Response) -> PollOutcome {
#[derive(serde::Deserialize)]
struct FloorProblem {
#[serde(rename = "type")]
problem_type: String,
min_client_version: String,
#[serde(default)]
client_version: Option<String>,
}
let problem = match resp.json::<FloorProblem>().await {
Ok(problem) if problem.problem_type == FLOOR_PROBLEM_TYPE => problem,
Ok(_) | Err(_) => {
self.record_poll_failure(&format!(
"{ERR_BUNDLE_FETCH_FAILED} policy bundle fetch returned a non-floor 403"
));
self.set_readiness(
PolicyReadinessPhase::TerminalRefusal,
"http_403",
Some("auth_rejected"),
);
return PollOutcome::Failed;
}
};
let installed_version = problem
.client_version
.unwrap_or_else(|| env!("CARGO_PKG_VERSION").to_string());
if let Ok(mut floor) = self.floor.write() {
*floor = Some(crate::daemon::ClientFloorState {
installed_version: installed_version.clone(),
minimum_version: problem.min_client_version.clone(),
});
}
self.last_fetch_ok.store(false, Ordering::Relaxed);
self.last_poll_ok_at
.store(now_unix_secs(), Ordering::Relaxed);
if let Ok(Some(mut meta)) = store::read_meta(&self.base_dir) {
meta.last_poll_ok_at = Some(crate::install_state::now_rfc3339());
meta.last_fetch_ok = false;
meta.state = Some("floor_blocked".to_string());
meta.min_client_version = Some(problem.min_client_version.clone());
meta.last_error = None;
if let Err(error) = store::write_meta(&self.base_dir, &meta) {
tracing::warn!(target: "policy", code = ERR_BUNDLE_FETCH_FAILED, %error, "could not persist client floor state");
}
}
tracing::warn!(
target: "policy",
code = ERR_BUNDLE_FETCH_FAILED,
installed_version,
minimum_version = %problem.min_client_version,
"client upgrade required; keeping the last-known-good bundle enforcing"
);
self.set_readiness(
PolicyReadinessPhase::TerminalRefusal,
"http_403",
Some("client_below_floor"),
);
PollOutcome::FloorBlocked
}
async fn handle_license_refused(&self, resp: reqwest::Response) -> PollOutcome {
let headers = resp.headers().clone();
let body = resp.text().await.unwrap_or_default();
self.record_license_block(&headers, &body)
}
fn record_license_block(
&self,
headers: &reqwest::header::HeaderMap,
body: &str,
) -> PollOutcome {
let CloudError::LicenseRefused {
code,
licensing_url,
retry_after,
} = license_refused_from(headers, body)
else {
unreachable!("license_refused_from only builds LicenseRefused");
};
self.cloud_state
.set_policy_license_refusal(Some((code.clone(), licensing_url)));
self.last_fetch_ok.store(false, Ordering::Relaxed);
self.last_poll_ok_at
.store(now_unix_secs(), Ordering::Relaxed);
if let Ok(Some(mut meta)) = store::read_meta(&self.base_dir) {
meta.last_poll_ok_at = Some(crate::install_state::now_rfc3339());
meta.last_fetch_ok = false;
meta.state = Some("license_blocked".to_string());
meta.last_error = Some(format!("{ERR_CLOUD_LICENSE_BLOCKED} {code}"));
if let Err(error) = store::write_meta(&self.base_dir, &meta) {
tracing::warn!(target: "policy", code = ERR_CLOUD_LICENSE_BLOCKED, %error, "could not persist the license block");
}
}
tracing::warn!(
target: "policy",
code = ERR_CLOUD_LICENSE_BLOCKED,
license_code = %code,
retry_after_secs = retry_after.map(|d| d.as_secs()),
"policy bundle refused on licensing grounds; the resident bundle keeps enforcing"
);
self.set_readiness(
PolicyReadinessPhase::TerminalRefusal,
"license_refused",
Some("license_refused"),
);
PollOutcome::LicenseBlocked(retry_after)
}
fn handle_not_modified(
&self,
resp: reqwest::Response,
meta: Option<store::BundleMeta>,
) -> PollOutcome {
let acknowledged = match self.selected_policy_version(resp.headers()) {
Ok(version) => version,
Err(outcome) => return outcome,
};
if let Some(selected) = acknowledged {
let cached_schema = meta
.as_ref()
.and_then(|meta| u32::try_from(meta.schema_version).ok());
let activated_selection =
contracts::read_compatibility(&self.base_dir)
.ok()
.and_then(|state| {
state
.selections
.get(POLICY_BUNDLE_FAMILY)
.map(|selection| selection.version)
});
if cached_schema != Some(selected)
|| activated_selection.is_some_and(|version| version != selected)
{
return self.compatibility_reject(
None,
&format!(
"304 acknowledged {POLICY_BUNDLE_FAMILY} v{selected}, but the resident cache is schema {:?} with activated selection {activated_selection:?}",
cached_schema
),
);
}
} else {
let cached_schema = meta
.as_ref()
.and_then(|meta| u32::try_from(meta.schema_version).ok());
if !cached_schema
.is_some_and(|version| contracts::has_reader(POLICY_BUNDLE_FAMILY, version))
{
return self.compatibility_reject(
None,
&format!(
"legacy 304 cannot validate cached schema {cached_schema:?} with a retained reader"
),
);
}
}
self.last_fetch_ok.store(true, Ordering::Relaxed);
self.last_poll_ok_at
.store(now_unix_secs(), Ordering::Relaxed);
if let Some(mut meta) = meta {
meta.last_poll_ok_at = Some(crate::install_state::now_rfc3339());
meta.last_fetch_ok = true;
meta.last_error = None;
if let Err(e) = store::write_meta(&self.base_dir, &meta) {
tracing::warn!(
target: "policy",
code = ERR_BUNDLE_FETCH_FAILED,
error = %e,
"could not update the policy poll clock on disk"
);
}
}
tracing::debug!(
target: "policy",
"policy bundle unchanged (304); the resident bundle stays active"
);
if let Some(selected) = acknowledged {
let persisted = self.record_selection(POLICY_BUNDLE_FAMILY, selected);
self.set_selection_acknowledged(persisted);
self.set_compatibility_diagnostic_valid(false);
self.set_readiness(PolicyReadinessPhase::Resident, "http_304", None);
} else {
self.set_selection_acknowledged(false);
self.set_compatibility_diagnostic_valid(false);
self.set_readiness(
PolicyReadinessPhase::Resident,
"http_304",
Some("legacy_no_ack"),
);
self.record_legacy_response(POLICY_BUNDLE_FAMILY);
}
PollOutcome::NotModified
}
fn selected_policy_version(
&self,
headers: &reqwest::header::HeaderMap,
) -> Result<Option<u32>, PollOutcome> {
let Some(raw) = headers.get(contracts::SELECTED_HEADER) else {
return Ok(None);
};
let raw = match raw.to_str() {
Ok(raw) => raw,
Err(_) => {
return Err(
self.compatibility_reject(None, "selected-map header is not visible ASCII")
)
}
};
let selected = match contracts::parse_selected(raw) {
Ok(selected) => selected,
Err(error) => {
return Err(self
.compatibility_reject(None, &format!("invalid selected-map header: {error}")))
}
};
let Some(version) = selected.get(POLICY_BUNDLE_FAMILY).copied() else {
return Err(self.compatibility_reject(None, "selected-map omitted policy_bundle"));
};
if !contracts::has_reader(POLICY_BUNDLE_FAMILY, version) {
return Err(self.compatibility_reject(
None,
&format!("selected policy_bundle v{version} has no local reader"),
));
}
Ok(Some(version))
}
async fn handle_compatibility_unavailable(&self, resp: reqwest::Response) -> PollOutcome {
let value = match resp.json::<serde_json::Value>().await {
Ok(value) => value,
Err(_) => {
return self
.transport_failure("policy bundle returned malformed compatibility problem")
}
};
if value.get("type").and_then(serde_json::Value::as_str)
!= Some(contracts::COMPATIBILITY_PROBLEM_TYPE)
{
return self.transport_failure("policy bundle returned an unregistered 409 problem");
}
let family = value
.get("family")
.and_then(serde_json::Value::as_str)
.unwrap_or(POLICY_BUNDLE_FAMILY);
let platform_range = value.get("platform_range").and_then(problem_range);
let detail = value
.get("detail")
.and_then(serde_json::Value::as_str)
.map(sanitize_problem_detail)
.filter(|detail| !detail.is_empty())
.unwrap_or_else(|| "no common retained representation".to_string());
let persisted = self.record_compatibility_diagnostic(family, platform_range, &detail);
self.set_compatibility_diagnostic_valid(persisted);
self.last_fetch_ok.store(false, Ordering::Relaxed);
self.set_readiness(
PolicyReadinessPhase::TerminalRefusal,
"http_409",
Some("compatibility_unavailable"),
);
tracing::warn!(
target: "policy",
family,
"policy bundle compatibility unavailable; keeping the last-known-good bundle enforcing"
);
PollOutcome::CompatibilityUnavailable {
family: family.to_string(),
}
}
fn compatibility_reject(
&self,
platform_range: Option<VersionRange>,
detail: &str,
) -> PollOutcome {
self.last_fetch_ok.store(false, Ordering::Relaxed);
self.set_readiness(
PolicyReadinessPhase::TerminalRefusal,
"protocol_rejected",
Some("invalid_response"),
);
let persisted =
self.record_compatibility_diagnostic(POLICY_BUNDLE_FAMILY, platform_range, detail);
self.set_compatibility_diagnostic_valid(persisted);
tracing::warn!(target: "policy", detail, "policy bundle negotiation rejected; resident bytes and activation metadata are unchanged");
PollOutcome::Rejected(ERR_BUNDLE_REJECTED)
}
fn record_compatibility_diagnostic(
&self,
family: &str,
platform_range: Option<VersionRange>,
detail: &str,
) -> bool {
let Some(client_range) = contracts::supported_ranges().get(family).copied() else {
return false;
};
let diagnostic = CompatibilityDiagnostic {
family: family.to_string(),
client_range,
platform_range,
last_selection: None,
observed_at: crate::install_state::now_rfc3339(),
detail: detail.to_string(),
};
match contracts::record_diagnostic(&self.base_dir, diagnostic) {
Ok(()) => true,
Err(error) => {
tracing::debug!(target: "policy", %error, "could not persist compatibility diagnostic");
false
}
}
}
fn record_selection(&self, family: &str, version: u32) -> bool {
match contracts::record_selection(&self.base_dir, family, version) {
Ok(()) => true,
Err(error) => {
tracing::debug!(target: "policy", %error, "could not persist contract selection");
false
}
}
}
fn record_legacy_response(&self, family: &str) {
if let Err(error) = contracts::record_legacy_response(&self.base_dir, family) {
tracing::debug!(target: "policy", %error, "could not clear stale contract selection for a legacy response");
}
}
fn record_pending_response(&self, family: &str) {
if let Err(error) = contracts::record_pending_response(&self.base_dir, family) {
tracing::debug!(target: "policy", %error, "could not clear stale compatibility diagnostic for a pending response");
}
}
fn reject(&self, code: &'static str, detail: &str) -> PollOutcome {
tracing::warn!(
target: "policy",
code,
detail,
"policy bundle rejected; keeping the last-known-good bundle enforcing"
);
self.record_poll_failure(&format!("{code} {detail}"));
self.set_compatibility_diagnostic_valid(false);
self.set_readiness(
PolicyReadinessPhase::TerminalRefusal,
"protocol_rejected",
Some("invalid_response"),
);
PollOutcome::Rejected(code)
}
fn transport_failure(&self, detail: &str) -> PollOutcome {
tracing::warn!(
target: "policy",
code = ERR_BUNDLE_FETCH_FAILED,
detail,
"policy bundle poll failed; keeping the last-known-good bundle enforcing"
);
self.record_poll_failure(&format!("{ERR_BUNDLE_FETCH_FAILED} {detail}"));
self.set_readiness(
PolicyReadinessPhase::Retrying,
"transport_error",
Some("transport_error"),
);
PollOutcome::Failed
}
fn record_poll_failure(&self, message: &str) {
self.last_fetch_ok.store(false, Ordering::Relaxed);
let Ok(Some(mut meta)) = store::read_meta(&self.base_dir) else {
return;
};
meta.last_fetch_ok = false;
meta.last_error = Some(message.to_string());
if let Err(e) = store::write_meta(&self.base_dir, &meta) {
tracing::debug!(
target: "policy",
error = %e,
"could not record the policy poll failure on disk"
);
}
}
fn check_staleness(&self) -> bool {
let last = self.last_poll_ok_at.load(Ordering::Relaxed);
if last <= 0 {
return false;
}
let age = now_unix_secs().saturating_sub(last).max(0) as u64;
if age <= self.config.stale_warn_after_secs {
return false;
}
tracing::warn!(
target: "policy",
code = ERR_BUNDLE_STALE,
stale_seconds = age,
threshold_seconds = self.config.stale_warn_after_secs,
"no successful policy bundle poll within the staleness threshold; the resident bundle keeps enforcing"
);
true
}
}
fn now_unix_secs() -> i64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|d| d.as_secs() as i64)
.unwrap_or(0)
}
fn parse_unix_secs(raw: &str) -> Option<i64> {
chrono::DateTime::parse_from_rfc3339(raw)
.ok()
.map(|dt| dt.timestamp())
}
fn timestamp(seconds: i64) -> Option<String> {
(seconds > 0)
.then(|| chrono::DateTime::from_timestamp(seconds, 0))
.flatten()
.map(|at| at.to_rfc3339_opts(chrono::SecondsFormat::Secs, true))
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::atomic::AtomicUsize;
use std::sync::Mutex;
use secrecy::SecretString;
use crate::core::policy::test_support::wire_rule;
use crate::core::policy::{evaluate_command, new_handle};
use crate::generated::types::{PolicyRule, PolicyRuleMode, PolicyRuleSeverity};
const ORG: &str = "0192f8a1-4c3b-7e2a-9f10-5d8c3b1a7e42";
const OTHER_ORG: &str = "0192f8a1-0000-0000-0000-000000000000";
struct TestCredentialProvider {
key: Mutex<Option<String>>,
calls: AtomicUsize,
}
impl TestCredentialProvider {
fn with_key(key: &str) -> Arc<Self> {
Arc::new(Self {
key: Mutex::new(Some(key.to_string())),
calls: AtomicUsize::new(0),
})
}
fn set_key(&self, key: &str) {
*self.key.lock().expect("lock") = Some(key.to_string());
}
}
impl CredentialProvider for TestCredentialProvider {
fn retrieve(&self) -> Option<SecretString> {
self.calls.fetch_add(1, Ordering::Relaxed);
self.key
.lock()
.ok()
.and_then(|g| g.as_ref().map(|k| SecretString::from(k.clone())))
}
}
fn rule(rule_id: &str, pattern: &str, mode: PolicyRuleMode) -> PolicyRule {
wire_rule(rule_id, pattern, mode, PolicyRuleSeverity::High)
}
fn bundle_json(revision: i64, org: &str, rules: Vec<PolicyRule>) -> serde_json::Value {
serde_json::json!({
"schema_version": 1,
"organization_id": org,
"revision": revision,
"built_at": "2026-07-21T09:00:00Z",
"enforcement_enabled": true,
"rules": rules,
"signature": serde_json::Value::Null,
})
}
fn body(revision: i64) -> Vec<u8> {
serde_json::to_vec(&bundle_json(
revision,
ORG,
vec![rule("OL-CMD-001", "*rm -rf*", PolicyRuleMode::Enforce)],
))
.expect("serialise fixture")
}
fn schema_two_body(revision: i64) -> Vec<u8> {
let mut value: serde_json::Value = serde_json::from_str(include_str!(
"../../tools/policy-seed/fixtures/schema-2.json"
))
.expect("schema-2 fixture");
value["revision"] = serde_json::json!(revision);
value["rules"] = serde_json::json!([
{
"rule_id":"OL-FUTURE-001", "kind":"command", "action":"deny",
"match_pattern":"*curl*", "mode":"enforce", "severity":"high",
"reason":"newer schema", "unknown_key_from_the_future":true
},
{
"rule_id":"OL-CMD-001", "kind":"command", "action":"deny",
"match_pattern":"*rm -rf*", "mode":"enforce", "severity":"high",
"reason":"known rule survives"
}
]);
serde_json::to_vec(&value).expect("serialise fixture")
}
fn etag_for(body: &[u8]) -> String {
format!("\"{}\"", store::digest_of(body))
}
fn refuse_compatibility_writes(base: &std::path::Path) -> PathBuf {
let compatibility_tmp = contracts::compatibility_path(base).with_extension("json.tmp");
std::fs::create_dir(&compatibility_tmp).unwrap();
compatibility_tmp
}
struct Harness {
poller: PolicyPoller,
dir: tempfile::TempDir,
credentials: Arc<TestCredentialProvider>,
cloud_state: CloudState,
handle: PolicyHandle,
last_fetch_ok: Arc<AtomicBool>,
last_poll_ok_at: Arc<AtomicI64>,
last_download_verified_at: Arc<AtomicI64>,
readiness: Arc<std::sync::RwLock<PolicyReadiness>>,
}
const AGENT_ID: &str = "0192f8a1-4c3b-7e2a-9f10-5d8c3b1a7e42";
impl Harness {
fn new(api_url: String) -> Self {
Self::with_agent_id(api_url, Some(AGENT_ID.to_string()))
}
fn with_agent_id(api_url: String, agent_id: Option<String>) -> Self {
let dir = tempfile::tempdir().expect("tempdir");
let handle = new_handle(None);
let last_fetch_ok = Arc::new(AtomicBool::new(false));
let last_poll_ok_at = Arc::new(AtomicI64::new(0));
let last_download_verified_at = Arc::new(AtomicI64::new(0));
let readiness = Arc::new(std::sync::RwLock::new(PolicyReadiness {
phase: PolicyReadinessPhase::Waiting,
last_attempted_at: None,
last_outcome: None,
reason_code: None,
selection_acknowledged: false,
compatibility_diagnostic_valid: false,
}));
let cloud_state = CloudState::new();
let credentials = TestCredentialProvider::with_key("ol_org_test");
let poller = PolicyPoller::new(
handle.clone(),
last_fetch_ok.clone(),
last_poll_ok_at.clone(),
last_download_verified_at.clone(),
readiness.clone(),
cloud_state.clone(),
credentials.clone(),
api_url,
PolicyConfig {
enabled: true,
poll_interval_secs: 300,
stale_warn_after_secs: 86_400,
},
dir.path().to_path_buf(),
crate::egress::ClientHandle::of(crate::egress::client()),
agent_id,
None,
);
Self {
poller,
dir,
credentials,
cloud_state,
handle,
last_fetch_ok,
last_poll_ok_at,
last_download_verified_at,
readiness,
}
}
fn still_enforcing(&self) -> bool {
let loaded = self.handle.load();
let Some(bundle) = loaded.as_ref().as_ref() else {
return false;
};
evaluate_command(bundle, "rm -rf /tmp").is_some_and(|m| !m.shadow)
}
fn resident_revision(&self) -> Option<i64> {
self.handle.load().as_ref().as_ref().map(|b| b.revision)
}
fn meta(&self) -> Option<store::BundleMeta> {
store::read_meta(self.dir.path()).expect("meta readable")
}
}
#[test]
fn parse_etag_strips_quotes_and_the_weak_prefix() {
assert_eq!(parse_etag("\"sha256:abc\""), Some("sha256:abc"));
assert_eq!(parse_etag("W/\"sha256:abc\""), Some("sha256:abc"));
assert_eq!(parse_etag(" W/ \"sha256:abc\" "), Some("sha256:abc"));
assert_ne!(parse_etag("W/\"sha256:abc\""), Some("\"sha256:abc\""));
assert_eq!(parse_etag("sha256:abc"), None);
assert_eq!(parse_etag(""), None);
}
#[test]
fn jitter_stays_within_ten_percent_and_varies() {
let base = 300u64;
let intervals: Vec<u64> = (0..20).map(|_| jittered(base).as_secs()).collect();
for secs in &intervals {
assert!(
(270..=330).contains(secs),
"interval {secs}s outside ±10% of {base}s"
);
}
assert!(
intervals.windows(2).any(|w| w[0] != w[1]),
"at least one consecutive pair must differ: {intervals:?}"
);
}
#[test]
fn jitter_does_not_underflow_on_a_tiny_interval() {
for _ in 0..20 {
let _ = jittered(1);
let _ = jittered(0);
}
}
#[test]
fn a_zero_poll_interval_cannot_produce_a_tight_loop() {
let configured: u64 = 0;
assert_eq!(
jittered(configured),
Duration::ZERO,
"precondition: 0 really is degenerate"
);
for _ in 0..50 {
let delay = jittered(configured.max(MIN_POLL_INTERVAL_SECS));
assert!(
delay >= Duration::from_secs(MIN_POLL_INTERVAL_SECS * 9 / 10),
"clamped interval fell below the floor: {delay:?}"
);
}
}
#[tokio::test]
async fn manual_refresh_distinguishes_changed_and_unchanged_without_changing_resident_facts() {
let mut server = mockito::Server::new_async().await;
let mut document: serde_json::Value =
serde_json::from_slice(&schema_two_body(44)).expect("schema-2 document");
document["artifacts"][0]["policy_id"] = serde_json::json!("policy-44");
let body = serde_json::to_vec(&document).expect("serialise");
let etag = etag_for(&body);
let changed = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(200)
.with_header("ETag", &etag)
.with_header(contracts::SELECTED_HEADER, "policy_bundle=2")
.with_body(&body)
.expect(1)
.create_async()
.await;
let mut harness = Harness::new(server.url());
let outcome = harness.poller.poll_once().await;
let changed_result = harness.poller.refresh_result(&outcome, None);
assert_eq!(changed_result.outcome, PolicyRefreshOutcome::Changed);
assert_eq!(changed_result.response_status, Some(200));
assert_eq!(changed_result.selected_version, Some(2));
assert_eq!(changed_result.schema_version, Some(2));
assert_eq!(changed_result.bundle_revision, Some(44));
assert!(changed_result.digest.is_some());
assert!(changed_result.verified_body_received);
assert!(changed_result.verified_body_received_at.is_some());
assert!(changed_result.last_platform_contact_at.is_some());
assert_eq!(changed_result.resident_aip_count, Some(1));
changed.assert_async().await;
let unchanged = server
.mock("GET", BUNDLE_ENDPOINT)
.match_header(reqwest::header::IF_NONE_MATCH.as_str(), etag.as_str())
.with_status(304)
.with_header(contracts::SELECTED_HEADER, "policy_bundle=2")
.expect(1)
.create_async()
.await;
let outcome = harness.poller.poll_once().await;
let unchanged_result = harness
.poller
.refresh_result(&outcome, changed_result.digest.as_deref());
assert_eq!(unchanged_result.outcome, PolicyRefreshOutcome::Unchanged);
assert_eq!(unchanged_result.response_status, Some(304));
assert_eq!(
unchanged_result.schema_version,
changed_result.schema_version
);
assert_eq!(
unchanged_result.bundle_revision,
changed_result.bundle_revision
);
assert_eq!(unchanged_result.digest, changed_result.digest);
assert_eq!(
unchanged_result.verified_body_received_at, changed_result.verified_body_received_at,
"a 304 must not claim a new verified body"
);
assert_eq!(
unchanged_result.resident_aip_count,
changed_result.resident_aip_count
);
unchanged.assert_async().await;
let downloaded_identical = server
.mock("GET", BUNDLE_ENDPOINT)
.match_header(reqwest::header::IF_NONE_MATCH.as_str(), etag.as_str())
.with_status(200)
.with_header("ETag", &etag)
.with_header(contracts::SELECTED_HEADER, "policy_bundle=2")
.with_body(&body)
.expect(1)
.create_async()
.await;
let before = harness.poller.resident_digest();
let outcome = harness.poller.poll_once().await;
let downloaded_result = harness.poller.refresh_result(&outcome, before.as_deref());
assert_eq!(downloaded_result.outcome, PolicyRefreshOutcome::Unchanged);
assert_eq!(downloaded_result.response_status, Some(200));
assert_eq!(downloaded_result.digest, changed_result.digest);
downloaded_identical.assert_async().await;
}
#[test]
fn manual_refresh_maps_non_activation_outcomes_and_keeps_last_known_good_visible() {
let harness = Harness::new("http://127.0.0.1:9".to_string());
let prior = Arc::new(Some(crate::core::policy::ResidentBundle::from_bundle(
&serde_json::from_value(bundle_json(7, ORG, vec![])).expect("bundle"),
)));
harness.handle.store(prior);
for (outcome, expected, reason) in [
(
PollOutcome::Pending(Duration::from_secs(2)),
PolicyRefreshOutcome::Pending,
"install_materializing",
),
(
PollOutcome::AuthFailed,
PolicyRefreshOutcome::Refused,
"auth_rejected",
),
(
PollOutcome::RateLimited(Some(Duration::from_secs(30))),
PolicyRefreshOutcome::Pending,
"rate_limited",
),
(
PollOutcome::Failed,
PolicyRefreshOutcome::Offline,
"transport_error",
),
] {
let result = harness.poller.refresh_result(&outcome, None);
assert_eq!(result.outcome, expected);
assert_eq!(result.reason.as_deref(), Some(reason));
assert!(result.has_bundle, "last-known-good must remain visible");
assert_eq!(result.bundle_revision, Some(7));
}
}
#[test]
fn a_one_second_identity_bootstrap_cadence_is_not_raised() {
let configured = 1;
for _ in 0..20 {
assert_eq!(
jittered(configured.max(MIN_POLL_INTERVAL_SECS)),
Duration::from_secs(1)
);
}
}
#[tokio::test]
async fn ok_activates_and_stores_the_activation_etag() {
let mut server = mockito::Server::new_async().await;
let body = body(42);
let etag = etag_for(&body);
let mock = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(200)
.with_header("ETag", &etag)
.with_body(&body)
.create_async()
.await;
let mut h = Harness::new(server.url());
assert_eq!(
h.poller.poll_once().await,
PollOutcome::Activated { revision: 42 }
);
assert!(h.still_enforcing(), "the new bundle must be enforcing");
assert_eq!(h.resident_revision(), Some(42));
let meta = h.meta().expect("meta written");
assert_eq!(meta.etag.as_deref(), Some(etag.as_str()));
assert_eq!(meta.digest, store::digest_of(&body));
assert_eq!(meta.organization_id, ORG);
assert!(meta.last_activated_at.is_some());
assert!(h.last_fetch_ok.load(Ordering::Relaxed));
assert!(h.last_poll_ok_at.load(Ordering::Relaxed) > 0);
assert_eq!(
std::fs::read(store::bundle_path(h.dir.path())).expect("body on disk"),
body
);
mock.assert_async().await;
}
#[tokio::test]
async fn weak_etag_is_accepted() {
let mut server = mockito::Server::new_async().await;
let body = body(7);
let mock = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(200)
.with_header("ETag", &format!("W/\"{}\"", store::digest_of(&body)))
.with_body(&body)
.create_async()
.await;
let mut h = Harness::new(server.url());
assert_eq!(
h.poller.poll_once().await,
PollOutcome::Activated { revision: 7 }
);
assert!(h.still_enforcing());
mock.assert_async().await;
}
#[tokio::test]
async fn empty_rule_set_activates_and_keeps_a_revision() {
let mut server = mockito::Server::new_async().await;
let body = serde_json::to_vec(&bundle_json(9, ORG, vec![])).expect("serialise");
let mock = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(200)
.with_header("ETag", &etag_for(&body))
.with_body(&body)
.create_async()
.await;
let mut h = Harness::new(server.url());
assert_eq!(
h.poller.poll_once().await,
PollOutcome::Activated { revision: 9 }
);
assert!(!h.still_enforcing());
assert_eq!(h.resident_revision(), Some(9));
mock.assert_async().await;
}
#[tokio::test]
async fn a_provisioned_install_sends_its_agent_id_on_every_poll() {
let mut server = mockito::Server::new_async().await;
let body = body(11);
let mock = server
.mock("GET", BUNDLE_ENDPOINT)
.match_header("x-openlatch-agent-id", AGENT_ID)
.with_status(200)
.with_header("ETag", &etag_for(&body))
.with_body(&body)
.expect(1)
.create_async()
.await;
let mut h = Harness::new(server.url());
assert_eq!(
h.poller.poll_once().await,
PollOutcome::Activated { revision: 11 }
);
mock.assert_async().await;
}
#[tokio::test]
async fn bundle_request_advertises_only_live_non_hold_capabilities() {
let mut server = mockito::Server::new_async().await;
let body = body(17);
let mock = server
.mock("GET", BUNDLE_ENDPOINT)
.match_header("accept", BUNDLE_ACCEPT)
.match_header(
contracts::SUPPORT_HEADER,
"decision_event=1-2,policy_bundle=1-2",
)
.match_header(BUNDLE_DELIVERY_HEADER, BUNDLE_DELIVERY_PENDING_V1)
.match_header("openlatch-client-version", env!("CARGO_PKG_VERSION"))
.match_header(
"openlatch-client-capabilities",
"t1_predicate_tree,t2_register_program,exception",
)
.match_header(
"user-agent",
format!(
"openlatch-client/{} ({}; {}; mode1)",
env!("CARGO_PKG_VERSION"),
std::env::consts::OS,
std::env::consts::ARCH
)
.as_str(),
)
.with_status(200)
.with_header("ETag", &etag_for(&body))
.with_body(&body)
.create_async()
.await;
let mut h = Harness::new(server.url());
assert!(matches!(
h.poller.poll_once().await,
PollOutcome::Activated { .. }
));
mock.assert_async().await;
}
#[tokio::test]
async fn typed_pending_advances_contact_not_receipt_and_retries_with_opt_in() {
let mut server = mockito::Server::new_async().await;
let pending = server
.mock("GET", BUNDLE_ENDPOINT)
.match_header(BUNDLE_DELIVERY_HEADER, BUNDLE_DELIVERY_PENDING_V1)
.match_header(
contracts::SUPPORT_HEADER,
"decision_event=1-2,policy_bundle=1-2",
)
.with_status(202)
.with_header("Retry-After", "99")
.with_body(r#"{"status":"bundle_pending","reason":"install_materializing"}"#)
.expect(2)
.create_async()
.await;
let mut h = Harness::new(server.url());
assert_eq!(
h.poller.poll_once().await,
PollOutcome::Pending(Duration::from_secs(10))
);
assert!(h.last_poll_ok_at.load(Ordering::Relaxed) > 0);
assert_eq!(h.last_download_verified_at.load(Ordering::Relaxed), 0);
assert_eq!(
h.readiness.read().unwrap().phase,
PolicyReadinessPhase::Waiting
);
assert_eq!(
h.poller.poll_once().await,
PollOutcome::Pending(Duration::from_secs(10))
);
assert_eq!(
h.readiness.read().unwrap().phase,
PolicyReadinessPhase::Retrying
);
pending.assert_async().await;
let body = schema_two_body(48);
let ready = server
.mock("GET", BUNDLE_ENDPOINT)
.match_header(BUNDLE_DELIVERY_HEADER, BUNDLE_DELIVERY_PENDING_V1)
.with_status(200)
.with_header("ETag", &etag_for(&body))
.with_header(contracts::SELECTED_HEADER, "policy_bundle=2")
.with_body(&body)
.create_async()
.await;
assert_eq!(
h.poller.poll_once().await,
PollOutcome::Activated { revision: 48 }
);
assert!(h.last_download_verified_at.load(Ordering::Relaxed) > 0);
assert_eq!(
h.readiness.read().unwrap().phase,
PolicyReadinessPhase::Resident
);
ready.assert_async().await;
}
#[tokio::test]
async fn malformed_pending_and_missing_install_id_are_terminal_and_preserve_lkg() {
let mut server = mockito::Server::new_async().await;
let body = body(49);
let first = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(200)
.with_header("ETag", &etag_for(&body))
.with_body(&body)
.create_async()
.await;
let mut h = Harness::new(server.url());
assert!(matches!(
h.poller.poll_once().await,
PollOutcome::Activated { .. }
));
first.assert_async().await;
let revision = h.resident_revision();
let malformed = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(202)
.with_body(r#"{"status":"something_else","reason":"raw detail"}"#)
.create_async()
.await;
assert_eq!(
h.poller.poll_once().await,
PollOutcome::Rejected(ERR_BUNDLE_INVALID)
);
assert_eq!(h.resident_revision(), revision);
assert_eq!(
h.readiness.read().unwrap().phase,
PolicyReadinessPhase::TerminalRefusal
);
malformed.assert_async().await;
let missing_id = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(400)
.with_body(r#"{"error":{"code":"install_id_required"}}"#)
.create_async()
.await;
assert_eq!(
h.poller.poll_once().await,
PollOutcome::Rejected(ERR_BUNDLE_FETCH_FAILED)
);
assert_eq!(h.resident_revision(), revision);
assert_eq!(
h.readiness.read().unwrap().reason_code,
Some("install_id_required")
);
missing_id.assert_async().await;
}
#[tokio::test]
async fn pending_preserves_historical_selection_without_renewing_it() {
let mut server = mockito::Server::new_async().await;
let body = schema_two_body(52);
let selected = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(200)
.with_header("ETag", &etag_for(&body))
.with_header(contracts::SELECTED_HEADER, "policy_bundle=2")
.with_body(&body)
.create_async()
.await;
let mut h = Harness::new(server.url());
assert!(matches!(
h.poller.poll_once().await,
PollOutcome::Activated { .. }
));
selected.assert_async().await;
h.poller.record_compatibility_diagnostic(
POLICY_BUNDLE_FAMILY,
Some(VersionRange::new(1, 2).unwrap()),
"old refusal",
);
let before = contracts::read_compatibility(h.dir.path())
.unwrap()
.selections[POLICY_BUNDLE_FAMILY]
.clone();
let pending = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(202)
.with_header("Retry-After", "2")
.with_body(r#"{"status":"bundle_pending","reason":"install_materializing"}"#)
.create_async()
.await;
assert_eq!(
h.poller.poll_once().await,
PollOutcome::Pending(Duration::from_secs(2))
);
let after = contracts::read_compatibility(h.dir.path()).unwrap();
assert_eq!(after.selections[POLICY_BUNDLE_FAMILY], before);
assert!(!after.diagnostics.contains_key(POLICY_BUNDLE_FAMILY));
{
let readiness = h.readiness.read().unwrap();
assert!(readiness.selection_acknowledged);
assert!(!readiness.compatibility_diagnostic_valid);
assert_eq!(readiness.last_outcome, Some("http_202"));
}
pending.assert_async().await;
}
#[tokio::test]
async fn failed_version_switch_persistence_invalidates_in_memory_selection() {
let mut server = mockito::Server::new_async().await;
let original = body(53);
let initial = server
.mock("GET", BUNDLE_ENDPOINT)
.match_header("if-none-match", mockito::Matcher::Missing)
.with_status(200)
.with_header("ETag", &etag_for(&original))
.with_header(contracts::SELECTED_HEADER, "policy_bundle=1")
.with_body(&original)
.create_async()
.await;
let mut h = Harness::new(server.url());
assert!(matches!(
h.poller.poll_once().await,
PollOutcome::Activated { .. }
));
initial.assert_async().await;
contracts::record_selection(h.dir.path(), contracts::DECISION_EVENT_FAMILY, 2).unwrap();
let _compatibility_tmp = refuse_compatibility_writes(h.dir.path());
let replacement = schema_two_body(54);
let switched = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(200)
.with_header("ETag", &etag_for(&replacement))
.with_header(contracts::SELECTED_HEADER, "policy_bundle=2")
.with_body(&replacement)
.create_async()
.await;
assert_eq!(
h.poller.poll_once().await,
PollOutcome::Activated { revision: 54 }
);
let disk = contracts::read_compatibility(h.dir.path()).unwrap();
assert_eq!(disk.selections[POLICY_BUNDLE_FAMILY].version, 1);
assert_eq!(disk.selections[contracts::DECISION_EVENT_FAMILY].version, 2);
{
let readiness = h.readiness.read().unwrap();
assert!(!readiness.selection_acknowledged);
assert_eq!(readiness.last_outcome, Some("http_200"));
}
assert_eq!(h.resident_revision(), Some(54));
switched.assert_async().await;
}
#[tokio::test]
async fn in_memory_activation_records_receipt_when_metadata_cannot_persist() {
let mut server = mockito::Server::new_async().await;
let body = body(51);
let response = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(200)
.with_header("ETag", &etag_for(&body))
.with_body(&body)
.create_async()
.await;
let root = tempfile::tempdir().unwrap();
let unusable_base = root.path().join("base-is-a-file");
std::fs::write(&unusable_base, b"not a directory").unwrap();
let receipt = Arc::new(AtomicI64::new(0));
let readiness = Arc::new(std::sync::RwLock::new(PolicyReadiness {
phase: PolicyReadinessPhase::Waiting,
last_attempted_at: None,
last_outcome: None,
reason_code: None,
selection_acknowledged: false,
compatibility_diagnostic_valid: false,
}));
let handle = new_handle(None);
let mut poller = PolicyPoller::new(
handle.clone(),
Arc::new(AtomicBool::new(false)),
Arc::new(AtomicI64::new(0)),
receipt.clone(),
readiness,
CloudState::new(),
TestCredentialProvider::with_key("ol_org_test"),
server.url(),
PolicyConfig::default(),
unusable_base,
crate::egress::ClientHandle::of(crate::egress::client()),
Some(AGENT_ID.to_string()),
None,
);
assert_eq!(
poller.poll_once().await,
PollOutcome::Activated { revision: 51 }
);
assert!(handle.load().is_some());
assert!(receipt.load(Ordering::Relaxed) > 0);
response.assert_async().await;
}
#[tokio::test]
async fn legacy_no_ack_stays_invalid_after_clear_failure_and_transport() {
let mut server = mockito::Server::new_async().await;
let body = body(50);
let etag = etag_for(&body);
let selected = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(200)
.with_header("ETag", &etag)
.with_header(contracts::SELECTED_HEADER, "policy_bundle=1")
.with_body(&body)
.create_async()
.await;
let mut h = Harness::new(server.url());
assert!(matches!(
h.poller.poll_once().await,
PollOutcome::Activated { .. }
));
selected.assert_async().await;
contracts::record_selection(h.dir.path(), contracts::DECISION_EVENT_FAMILY, 2).unwrap();
let compatibility_tmp = refuse_compatibility_writes(h.dir.path());
let legacy = server
.mock("GET", BUNDLE_ENDPOINT)
.match_header("if-none-match", etag.as_str())
.with_status(304)
.expect(1)
.create_async()
.await;
assert_eq!(h.poller.poll_once().await, PollOutcome::NotModified);
let state = contracts::read_compatibility(h.dir.path()).unwrap();
assert!(
state.selections.contains_key(POLICY_BUNDLE_FAMILY),
"the forced write failure leaves stale disk evidence for metrics to mask"
);
assert_eq!(
state.selections[contracts::DECISION_EVENT_FAMILY].version,
2
);
assert!(!h.readiness.read().unwrap().selection_acknowledged);
legacy.assert_async().await;
let serving_url = h.poller.url.clone();
h.poller.url = format!("http://127.0.0.1:1{BUNDLE_ENDPOINT}");
assert_eq!(h.poller.poll_once().await, PollOutcome::Failed);
{
let readiness = h.readiness.read().unwrap();
assert_eq!(readiness.last_outcome, Some("transport_error"));
assert!(
!readiness.selection_acknowledged,
"a later transport result must not resurrect the stale acknowledgement"
);
}
h.poller.url = serving_url;
std::fs::remove_dir(compatibility_tmp).unwrap();
let acknowledged = server
.mock("GET", BUNDLE_ENDPOINT)
.match_header("if-none-match", etag.as_str())
.with_status(304)
.with_header(contracts::SELECTED_HEADER, "policy_bundle=1")
.create_async()
.await;
assert_eq!(h.poller.poll_once().await, PollOutcome::NotModified);
assert!(h.readiness.read().unwrap().selection_acknowledged);
acknowledged.assert_async().await;
}
#[tokio::test]
async fn acknowledged_v1_and_v2_bundles_activate_only_the_named_adapter() {
for (selected, body) in [(1, body(31)), (2, schema_two_body(32))] {
let mut server = mockito::Server::new_async().await;
let mock = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(200)
.with_header("ETag", &etag_for(&body))
.with_header(
contracts::SELECTED_HEADER,
&format!("policy_bundle={selected}"),
)
.with_body(&body)
.create_async()
.await;
let mut h = Harness::new(server.url());
assert!(matches!(
h.poller.poll_once().await,
PollOutcome::Activated { .. }
));
assert_eq!(
contracts::read_compatibility(h.dir.path())
.unwrap()
.selections[POLICY_BUNDLE_FAMILY]
.version,
selected
);
mock.assert_async().await;
}
}
#[tokio::test]
async fn acknowledged_304_records_selection_without_changing_bundle_identity() {
let mut server = mockito::Server::new_async().await;
let body = body(33);
let etag = etag_for(&body);
let first = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(200)
.with_header("ETag", &etag)
.with_body(&body)
.create_async()
.await;
let mut h = Harness::new(server.url());
assert!(matches!(
h.poller.poll_once().await,
PollOutcome::Activated { .. }
));
first.assert_async().await;
let body_before = std::fs::read(store::bundle_path(h.dir.path())).unwrap();
let revalidate = server
.mock("GET", BUNDLE_ENDPOINT)
.match_header("if-none-match", etag.as_str())
.with_status(304)
.with_header(contracts::SELECTED_HEADER, "policy_bundle=1")
.create_async()
.await;
assert_eq!(h.poller.poll_once().await, PollOutcome::NotModified);
assert_eq!(
std::fs::read(store::bundle_path(h.dir.path())).unwrap(),
body_before
);
assert_eq!(
contracts::read_compatibility(h.dir.path())
.unwrap()
.selections[POLICY_BUNDLE_FAMILY]
.version,
1
);
revalidate.assert_async().await;
}
#[tokio::test]
async fn a_304_acknowledgement_must_match_the_resident_schema_and_selection() {
let mut server = mockito::Server::new_async().await;
let body = body(37);
let etag = etag_for(&body);
let first = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(200)
.with_header("ETag", &etag)
.with_header(contracts::SELECTED_HEADER, "policy_bundle=1")
.with_body(&body)
.create_async()
.await;
let mut h = Harness::new(server.url());
assert_eq!(
h.poller.poll_once().await,
PollOutcome::Activated { revision: 37 }
);
first.assert_async().await;
h.poller.record_compatibility_diagnostic(
POLICY_BUNDLE_FAMILY,
Some(VersionRange::new(3, 4).unwrap()),
"existing diagnostic",
);
let body_before = std::fs::read(store::bundle_path(h.dir.path())).unwrap();
let meta_before = std::fs::read(store::meta_path(h.dir.path())).unwrap();
let poll_clock_before = h.last_poll_ok_at.load(Ordering::Relaxed);
let mismatch = server
.mock("GET", BUNDLE_ENDPOINT)
.match_header("if-none-match", etag.as_str())
.with_status(304)
.with_header(contracts::SELECTED_HEADER, "policy_bundle=2")
.create_async()
.await;
assert_eq!(
h.poller.poll_once().await,
PollOutcome::Rejected(ERR_BUNDLE_REJECTED)
);
mismatch.assert_async().await;
assert_eq!(
std::fs::read(store::bundle_path(h.dir.path())).unwrap(),
body_before
);
assert_eq!(
std::fs::read(store::meta_path(h.dir.path())).unwrap(),
meta_before
);
assert_eq!(h.last_poll_ok_at.load(Ordering::Relaxed), poll_clock_before);
let state = contracts::read_compatibility(h.dir.path()).unwrap();
assert_eq!(state.selections[POLICY_BUNDLE_FAMILY].version, 1);
assert!(state.diagnostics.contains_key(POLICY_BUNDLE_FAMILY));
}
#[tokio::test]
async fn no_intersection_is_first_boot_safe_and_a_valid_poll_clears_only_diagnostic() {
let mut server = mockito::Server::new_async().await;
let mismatch = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(409)
.with_body(
serde_json::json!({
"type": contracts::COMPATIBILITY_PROBLEM_TYPE,
"family": POLICY_BUNDLE_FAMILY,
"client_range": {"oldest": 1, "newest": 2},
"platform_range": {"oldest": 3, "newest": 4}
})
.to_string(),
)
.create_async()
.await;
let mut h = Harness::new(server.url());
assert_eq!(
h.poller.poll_once().await,
PollOutcome::CompatibilityUnavailable {
family: POLICY_BUNDLE_FAMILY.to_string()
}
);
assert!(!store::bundle_path(h.dir.path()).exists());
assert!(!store::meta_path(h.dir.path()).exists());
assert!(contracts::read_compatibility(h.dir.path())
.unwrap()
.diagnostics
.contains_key(POLICY_BUNDLE_FAMILY));
mismatch.assert_async().await;
let body = body(34);
let recovery = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(200)
.with_header("ETag", &etag_for(&body))
.with_header(contracts::SELECTED_HEADER, "policy_bundle=1")
.with_body(&body)
.create_async()
.await;
assert_eq!(
h.poller.poll_once().await,
PollOutcome::Activated { revision: 34 }
);
let state = contracts::read_compatibility(h.dir.path()).unwrap();
assert!(!state.diagnostics.contains_key(POLICY_BUNDLE_FAMILY));
assert_eq!(state.selections[POLICY_BUNDLE_FAMILY].version, 1);
recovery.assert_async().await;
let body_before = std::fs::read(store::bundle_path(h.dir.path())).unwrap();
let meta_before = std::fs::read(store::meta_path(h.dir.path())).unwrap();
let mismatch_with_resident = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(409)
.with_body(
serde_json::json!({
"type": contracts::COMPATIBILITY_PROBLEM_TYPE,
"family": POLICY_BUNDLE_FAMILY,
"platform_range": [3, 4]
})
.to_string(),
)
.create_async()
.await;
assert!(matches!(
h.poller.poll_once().await,
PollOutcome::CompatibilityUnavailable { .. }
));
assert_eq!(
std::fs::read(store::bundle_path(h.dir.path())).unwrap(),
body_before
);
assert_eq!(
std::fs::read(store::meta_path(h.dir.path())).unwrap(),
meta_before
);
assert_eq!(h.resident_revision(), Some(34));
mismatch_with_resident.assert_async().await;
}
#[tokio::test]
async fn compatibility_problem_detail_is_bounded_sanitized_and_timestamped() {
let mut server = mockito::Server::new_async().await;
let raw_detail = format!("temporarily unservable\n{}", "x".repeat(300));
let response = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(409)
.with_body(
serde_json::json!({
"type": contracts::COMPATIBILITY_PROBLEM_TYPE,
"family": POLICY_BUNDLE_FAMILY,
"platform_range": {"oldest": 1, "newest": 2},
"detail": raw_detail,
})
.to_string(),
)
.create_async()
.await;
let mut h = Harness::new(server.url());
assert!(matches!(
h.poller.poll_once().await,
PollOutcome::CompatibilityUnavailable { .. }
));
let state = contracts::read_compatibility(h.dir.path()).unwrap();
let diagnostic = &state.diagnostics[POLICY_BUNDLE_FAMILY];
assert!(diagnostic.detail.starts_with("temporarily unservable"));
assert!(diagnostic
.detail
.chars()
.all(|character| !character.is_control()));
assert!(diagnostic.detail.chars().count() <= 240);
assert!(chrono::DateTime::parse_from_rfc3339(&diagnostic.observed_at).is_ok());
assert!(h.readiness.read().unwrap().compatibility_diagnostic_valid);
response.assert_async().await;
}
#[tokio::test]
async fn failed_409_persistence_invalidates_older_disk_diagnostic() {
let mut server = mockito::Server::new_async().await;
let mut h = Harness::new(server.url());
assert!(h.poller.record_compatibility_diagnostic(
POLICY_BUNDLE_FAMILY,
Some(VersionRange::new(1, 2).unwrap()),
"old diagnostic",
));
let _compatibility_tmp = refuse_compatibility_writes(h.dir.path());
let response = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(409)
.with_body(
serde_json::json!({
"type": contracts::COMPATIBILITY_PROBLEM_TYPE,
"family": POLICY_BUNDLE_FAMILY,
"platform_range": {"oldest": 1, "newest": 2},
"detail": "new diagnostic",
})
.to_string(),
)
.create_async()
.await;
assert!(matches!(
h.poller.poll_once().await,
PollOutcome::CompatibilityUnavailable { .. }
));
let disk = contracts::read_compatibility(h.dir.path()).unwrap();
assert_eq!(
disk.diagnostics[POLICY_BUNDLE_FAMILY].detail,
"old diagnostic"
);
{
let readiness = h.readiness.read().unwrap();
assert!(!readiness.compatibility_diagnostic_valid);
assert_eq!(readiness.last_outcome, Some("http_409"));
}
response.assert_async().await;
}
#[tokio::test]
async fn incompatible_ack_preserves_resident_body_and_meta_byte_for_byte() {
let mut server = mockito::Server::new_async().await;
let original = body(35);
let first = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(200)
.with_header("ETag", &etag_for(&original))
.with_body(&original)
.create_async()
.await;
let mut h = Harness::new(server.url());
assert!(matches!(
h.poller.poll_once().await,
PollOutcome::Activated { .. }
));
first.assert_async().await;
let body_before = std::fs::read(store::bundle_path(h.dir.path())).unwrap();
let meta_before = std::fs::read(store::meta_path(h.dir.path())).unwrap();
let candidate = schema_two_body(36);
let rejected = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(200)
.with_header("ETag", &etag_for(&candidate))
.with_header(contracts::SELECTED_HEADER, "policy_bundle=1")
.with_body(&candidate)
.create_async()
.await;
assert_eq!(
h.poller.poll_once().await,
PollOutcome::Rejected(ERR_BUNDLE_REJECTED)
);
assert_eq!(
std::fs::read(store::bundle_path(h.dir.path())).unwrap(),
body_before
);
assert_eq!(
std::fs::read(store::meta_path(h.dir.path())).unwrap(),
meta_before
);
assert_eq!(h.resident_revision(), Some(35));
assert!(contracts::read_compatibility(h.dir.path())
.unwrap()
.diagnostics
.contains_key(POLICY_BUNDLE_FAMILY));
rejected.assert_async().await;
}
#[tokio::test]
async fn schema_one_ack_refuses_compiled_v2_policy_without_silent_lowering() {
let mut server = mockito::Server::new_async().await;
let original = body(38);
let first = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(200)
.with_header("ETag", &etag_for(&original))
.with_body(&original)
.create_async()
.await;
let mut h = Harness::new(server.url());
assert_eq!(
h.poller.poll_once().await,
PollOutcome::Activated { revision: 38 }
);
first.assert_async().await;
let body_before = std::fs::read(store::bundle_path(h.dir.path())).unwrap();
let meta_before = std::fs::read(store::meta_path(h.dir.path())).unwrap();
let mut lossy: serde_json::Value = serde_json::from_slice(&schema_two_body(39)).unwrap();
lossy["schema_version"] = serde_json::json!(1);
let lossy = serde_json::to_vec(&lossy).unwrap();
let rejected = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(200)
.with_header("ETag", &etag_for(&lossy))
.with_header(contracts::SELECTED_HEADER, "policy_bundle=1")
.with_body(&lossy)
.create_async()
.await;
assert_eq!(
h.poller.poll_once().await,
PollOutcome::Rejected(ERR_BUNDLE_REJECTED)
);
rejected.assert_async().await;
assert_eq!(h.resident_revision(), Some(38));
assert_eq!(
std::fs::read(store::bundle_path(h.dir.path())).unwrap(),
body_before
);
assert_eq!(
std::fs::read(store::meta_path(h.dir.path())).unwrap(),
meta_before
);
}
#[tokio::test]
async fn production_hold_source_unlocks_the_tier_three_capability() {
let mut server = mockito::Server::new_async().await;
let body = body(17);
let mock = server
.mock("GET", BUNDLE_ENDPOINT)
.match_header(
"openlatch-client-capabilities",
CLIENT_CAPABILITIES_WITH_HOLDS,
)
.with_status(200)
.with_header("ETag", &etag_for(&body))
.with_body(&body)
.create_async()
.await;
let mut h = Harness::new(server.url());
h.poller = h.poller.with_hold_capability();
assert!(matches!(
h.poller.poll_once().await,
PollOutcome::Activated { .. }
));
mock.assert_async().await;
}
#[tokio::test]
async fn schema_two_network_and_disk_projection_are_identical() {
let mut server = mockito::Server::new_async().await;
let body = schema_two_body(18);
let mock = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(200)
.with_header("ETag", &etag_for(&body))
.with_body(&body)
.create_async()
.await;
let mut h = Harness::new(server.url());
assert_eq!(
h.poller.poll_once().await,
PollOutcome::Activated { revision: 18 }
);
mock.assert_async().await;
let network = h.handle.load_full();
let disk = crate::daemon::PolicyRuntime::load_from_disk(h.dir.path())
.handle
.load_full();
assert_eq!(network, disk);
let resident = network.as_ref().as_ref().expect("schema-2 resident");
assert_eq!(resident.schema_version, 2);
assert_eq!(
resident.digest.as_deref(),
Some(store::digest_of(&body).as_str())
);
let meta = h.meta().expect("activation metadata");
assert_eq!(meta.schema_version, 2);
assert_eq!(meta.digest, resident.digest.as_deref().unwrap());
assert!(network
.as_ref()
.as_ref()
.is_some_and(|bundle| bundle.zone.is_some()));
assert_eq!(
network
.as_ref()
.as_ref()
.map(|bundle| bundle.command_rules.as_slice()),
disk.as_ref()
.as_ref()
.map(|bundle| bundle.command_rules.as_slice())
);
assert_eq!(
network
.as_ref()
.as_ref()
.map(|bundle| bundle.command_rules.len()),
Some(1),
"the forward rule is skipped identically on network and cache paths"
);
}
#[tokio::test]
async fn the_agent_id_header_rides_beside_the_validator() {
let mut server = mockito::Server::new_async().await;
let body = body(12);
let etag = etag_for(&body);
let first = server
.mock("GET", BUNDLE_ENDPOINT)
.match_header("x-openlatch-agent-id", AGENT_ID)
.with_status(200)
.with_header("ETag", &etag)
.with_body(&body)
.expect(1)
.create_async()
.await;
let mut h = Harness::new(server.url());
assert!(matches!(
h.poller.poll_once().await,
PollOutcome::Activated { .. }
));
first.assert_async().await;
let revalidate = server
.mock("GET", BUNDLE_ENDPOINT)
.match_header("if-none-match", etag.as_str())
.match_header("x-openlatch-agent-id", AGENT_ID)
.with_status(304)
.expect(1)
.create_async()
.await;
assert_eq!(h.poller.poll_once().await, PollOutcome::NotModified);
revalidate.assert_async().await;
}
#[tokio::test]
async fn an_unprovisioned_install_omits_the_agent_id_header() {
let mut server = mockito::Server::new_async().await;
let body = body(13);
let mock = server
.mock("GET", BUNDLE_ENDPOINT)
.match_header("x-openlatch-agent-id", mockito::Matcher::Missing)
.with_status(200)
.with_header("ETag", &etag_for(&body))
.with_body(&body)
.expect(1)
.create_async()
.await;
let mut h = Harness::with_agent_id(server.url(), None);
assert_eq!(
h.poller.poll_once().await,
PollOutcome::Activated { revision: 13 }
);
mock.assert_async().await;
}
#[tokio::test]
async fn a_validator_fetched_under_another_identity_is_not_revalidated() {
let mut server = mockito::Server::new_async().await;
let body = body(14);
let etag = etag_for(&body);
let mut h = Harness::new(server.url());
let mut meta = store::BundleMeta::activated(
&serde_json::from_slice(&body).expect("fixture parses"),
store::digest_of(&body),
Some(etag.clone()),
);
meta.agent_id = None;
let mut raw = serde_json::to_value(&meta).expect("meta serialises");
raw.as_object_mut().expect("object").remove("agent_id");
store::write_body(h.dir.path(), &body).expect("write body");
std::fs::write(
store::meta_path(h.dir.path()),
serde_json::to_vec(&raw).expect("serialise"),
)
.expect("write legacy meta");
assert!(
h.meta().expect("meta readable").agent_id.is_none(),
"precondition: the legacy file carries no identity"
);
let full = server
.mock("GET", BUNDLE_ENDPOINT)
.match_header("if-none-match", mockito::Matcher::Missing)
.match_header("x-openlatch-agent-id", AGENT_ID)
.with_status(200)
.with_header("ETag", &etag)
.with_body(&body)
.expect(1)
.create_async()
.await;
assert_eq!(
h.poller.poll_once().await,
PollOutcome::Activated { revision: 14 }
);
full.assert_async().await;
assert_eq!(h.meta().expect("meta").agent_id.as_deref(), Some(AGENT_ID));
let revalidate = server
.mock("GET", BUNDLE_ENDPOINT)
.match_header("if-none-match", etag.as_str())
.match_header("x-openlatch-agent-id", AGENT_ID)
.with_status(304)
.expect(1)
.create_async()
.await;
assert_eq!(h.poller.poll_once().await, PollOutcome::NotModified);
revalidate.assert_async().await;
}
#[tokio::test]
async fn losing_the_agent_id_also_drops_the_validator() {
let mut server = mockito::Server::new_async().await;
let body = body(15);
let etag = etag_for(&body);
let seed = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(200)
.with_header("ETag", &etag)
.with_body(&body)
.expect(1)
.create_async()
.await;
let mut h = Harness::new(server.url());
assert!(matches!(
h.poller.poll_once().await,
PollOutcome::Activated { .. }
));
seed.assert_async().await;
h.poller.agent_id = None;
let full = server
.mock("GET", BUNDLE_ENDPOINT)
.match_header("if-none-match", mockito::Matcher::Missing)
.match_header("x-openlatch-agent-id", mockito::Matcher::Missing)
.with_status(200)
.with_header("ETag", &etag)
.with_body(&body)
.expect(1)
.create_async()
.await;
assert!(matches!(
h.poller.poll_once().await,
PollOutcome::Activated { .. }
));
full.assert_async().await;
}
#[tokio::test]
async fn an_agent_id_that_is_not_a_header_value_is_dropped_not_sent() {
let mut server = mockito::Server::new_async().await;
let body = body(16);
let mock = server
.mock("GET", BUNDLE_ENDPOINT)
.match_header("x-openlatch-agent-id", mockito::Matcher::Missing)
.with_status(200)
.with_header("ETag", &etag_for(&body))
.with_body(&body)
.expect(1)
.create_async()
.await;
let mut h = Harness::with_agent_id(
server.url(),
Some(
"bad
id"
.to_string(),
),
);
assert!(h.poller.agent_id.is_none(), "dropped at construction");
assert_eq!(
h.poller.poll_once().await,
PollOutcome::Activated { revision: 16 }
);
mock.assert_async().await;
}
async fn assert_rejected(
builder: impl FnOnce(mockito::Mock) -> mockito::Mock,
code: &'static str,
) {
let mut server = mockito::Server::new_async().await;
let mock = builder(server.mock("GET", BUNDLE_ENDPOINT))
.create_async()
.await;
let mut h = Harness::new(server.url());
assert_eq!(h.poller.poll_once().await, PollOutcome::Rejected(code));
assert!(
h.handle.load().is_none(),
"a rejected bundle must never activate"
);
assert!(!h.last_fetch_ok.load(Ordering::Relaxed));
assert_eq!(
h.last_poll_ok_at.load(Ordering::Relaxed),
0,
"a rejection is not a successful poll"
);
mock.assert_async().await;
}
#[tokio::test]
async fn missing_etag_on_200_is_rejected() {
let body = body(1);
assert_rejected(
move |m| m.with_status(200).with_body(body),
ERR_BUNDLE_REJECTED,
)
.await;
}
#[tokio::test]
async fn malformed_etag_on_200_is_rejected() {
let body = body(1);
assert_rejected(
move |m| {
m.with_status(200)
.with_header("ETag", "sha256:unquoted")
.with_body(body)
},
ERR_BUNDLE_REJECTED,
)
.await;
}
#[tokio::test]
async fn digest_mismatch_is_rejected() {
let served = body(1);
let other = etag_for(&body(2));
assert_rejected(
move |m| {
m.with_status(200)
.with_header("ETag", &other)
.with_body(served)
},
ERR_BUNDLE_REJECTED,
)
.await;
}
#[tokio::test]
async fn org_mismatch_is_rejected() {
let mut server = mockito::Server::new_async().await;
let first = body(1);
let good = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(200)
.with_header("ETag", &etag_for(&first))
.with_body(&first)
.expect(1)
.create_async()
.await;
let mut h = Harness::new(server.url());
assert!(matches!(
h.poller.poll_once().await,
PollOutcome::Activated { .. }
));
good.assert_async().await;
let intruder = serde_json::to_vec(&bundle_json(
2,
OTHER_ORG,
vec![rule("OL-CMD-999", "*", PolicyRuleMode::Enforce)],
))
.expect("serialise");
let bad = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(200)
.with_header("ETag", &etag_for(&intruder))
.with_body(&intruder)
.create_async()
.await;
assert_eq!(
h.poller.poll_once().await,
PollOutcome::Rejected(ERR_BUNDLE_REJECTED)
);
assert_eq!(h.resident_revision(), Some(1), "never evaluate another org");
bad.assert_async().await;
}
#[tokio::test]
async fn unknown_schema_version_is_rejected() {
let mut doc = bundle_json(1, ORG, vec![]);
doc["schema_version"] = serde_json::json!(3);
let body = serde_json::to_vec(&doc).expect("serialise");
let etag = etag_for(&body);
assert_rejected(
move |m| {
m.with_status(200)
.with_header("ETag", &etag)
.with_body(body)
},
ERR_BUNDLE_INVALID,
)
.await;
}
#[tokio::test]
async fn malformed_json_is_rejected() {
let body = b"{ not json".to_vec();
let etag = etag_for(&body);
assert_rejected(
move |m| {
m.with_status(200)
.with_header("ETag", &etag)
.with_body(body)
},
ERR_BUNDLE_INVALID,
)
.await;
}
#[tokio::test]
async fn non_null_signature_is_rejected() {
let mut doc = bundle_json(1, ORG, vec![]);
doc["signature"] = serde_json::json!("ed25519:deadbeef");
let body = serde_json::to_vec(&doc).expect("serialise");
let etag = etag_for(&body);
assert_rejected(
move |m| {
m.with_status(200)
.with_header("ETag", &etag)
.with_body(body)
},
ERR_BUNDLE_INVALID,
)
.await;
}
#[tokio::test]
async fn activation_failure_rewinds_the_etag() {
let mut server = mockito::Server::new_async().await;
let good = body(1);
let good_etag = etag_for(&good);
let first = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(200)
.with_header("ETag", &good_etag)
.with_body(&good)
.expect(1)
.create_async()
.await;
let mut h = Harness::new(server.url());
assert!(matches!(
h.poller.poll_once().await,
PollOutcome::Activated { .. }
));
first.assert_async().await;
let corrupt = body(2);
let second = server
.mock("GET", BUNDLE_ENDPOINT)
.match_header("if-none-match", good_etag.as_str())
.with_status(200)
.with_header("ETag", &etag_for(b"a different document entirely"))
.with_body(&corrupt)
.expect(1)
.create_async()
.await;
assert_eq!(
h.poller.poll_once().await,
PollOutcome::Rejected(ERR_BUNDLE_REJECTED)
);
second.assert_async().await;
assert_eq!(h.resident_revision(), Some(1), "revision 1 keeps enforcing");
assert!(h.still_enforcing());
assert_eq!(
h.meta().expect("meta").etag.as_deref(),
Some(good_etag.as_str()),
"the stored validator must still be revision 1's"
);
let fixed = body(3);
let third = server
.mock("GET", BUNDLE_ENDPOINT)
.match_header("if-none-match", good_etag.as_str())
.with_status(200)
.with_header("ETag", &etag_for(&fixed))
.with_body(&fixed)
.expect(1)
.create_async()
.await;
assert_eq!(
h.poller.poll_once().await,
PollOutcome::Activated { revision: 3 }
);
third.assert_async().await;
}
#[tokio::test]
async fn bare_304s_do_not_clear_the_validator() {
let mut server = mockito::Server::new_async().await;
let good = body(11);
let etag = etag_for(&good);
let body_on_disk = good.clone();
let download = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(200)
.with_header("ETag", &etag)
.with_body(&good)
.expect(1) .create_async()
.await;
let mut h = Harness::new(server.url());
assert!(matches!(
h.poller.poll_once().await,
PollOutcome::Activated { .. }
));
download.assert_async().await;
let revalidate = server
.mock("GET", BUNDLE_ENDPOINT)
.match_header("if-none-match", etag.as_str())
.with_status(304)
.expect(3)
.create_async()
.await;
for _ in 0..3 {
assert_eq!(h.poller.poll_once().await, PollOutcome::NotModified);
}
revalidate.assert_async().await;
download.assert_async().await;
assert_eq!(
h.meta().expect("meta").etag.as_deref(),
Some(etag.as_str()),
"a 304 must never clear the stored validator"
);
assert_eq!(
std::fs::read(store::bundle_path(h.dir.path())).expect("body"),
body_on_disk,
"a 304 must not touch bundle.json"
);
assert_eq!(h.resident_revision(), Some(11));
}
#[tokio::test]
async fn not_modified_advances_the_poll_clock_only() {
let mut server = mockito::Server::new_async().await;
let good = body(5);
let etag = etag_for(&good);
let download = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(200)
.with_header("ETag", &etag)
.with_body(&good)
.expect(1)
.create_async()
.await;
let mut h = Harness::new(server.url());
assert!(matches!(
h.poller.poll_once().await,
PollOutcome::Activated { .. }
));
download.assert_async().await;
let activated_at = h.meta().expect("meta").last_activated_at;
let before = Arc::as_ptr(&h.handle.load_full());
h.last_poll_ok_at.store(1, Ordering::Relaxed);
let revalidate = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(304)
.with_header("ETag", &etag)
.create_async()
.await;
assert_eq!(h.poller.poll_once().await, PollOutcome::NotModified);
revalidate.assert_async().await;
assert!(h.last_poll_ok_at.load(Ordering::Relaxed) > 1);
assert!(h.last_fetch_ok.load(Ordering::Relaxed));
assert_eq!(
before,
Arc::as_ptr(&h.handle.load_full()),
"the rule set must not be reloaded on a 304"
);
let meta = h.meta().expect("meta");
assert_eq!(
meta.last_activated_at, activated_at,
"the activation tier must not move on a 304"
);
assert!(meta.last_poll_ok_at.is_some());
}
#[tokio::test]
async fn not_found_keeps_the_resident_bundle_enforcing() {
let mut server = mockito::Server::new_async().await;
let good = body(3);
let download = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(200)
.with_header("ETag", &etag_for(&good))
.with_body(&good)
.expect(1)
.create_async()
.await;
let mut h = Harness::new(server.url());
assert!(matches!(
h.poller.poll_once().await,
PollOutcome::Activated { .. }
));
download.assert_async().await;
let polled_at = h.last_poll_ok_at.load(Ordering::Relaxed);
let gone = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(404)
.with_body(r#"{"error":{"code":"not_found","message":"no bundle"}}"#)
.create_async()
.await;
assert_eq!(h.poller.poll_once().await, PollOutcome::NoBundle);
gone.assert_async().await;
assert!(h.still_enforcing(), "a 404 must not disarm the host");
assert_eq!(h.resident_revision(), Some(3));
assert_eq!(
h.last_poll_ok_at.load(Ordering::Relaxed),
polled_at,
"a 404 is neither a 2xx nor a 304"
);
assert!(!h.last_fetch_ok.load(Ordering::Relaxed));
}
#[tokio::test]
async fn server_error_keeps_the_resident_bundle_enforcing() {
let mut server = mockito::Server::new_async().await;
let good = body(4);
let download = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(200)
.with_header("ETag", &etag_for(&good))
.with_body(&good)
.expect(1)
.create_async()
.await;
let mut h = Harness::new(server.url());
assert!(matches!(
h.poller.poll_once().await,
PollOutcome::Activated { .. }
));
download.assert_async().await;
let boom = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(503)
.create_async()
.await;
assert_eq!(h.poller.poll_once().await, PollOutcome::Failed);
boom.assert_async().await;
assert!(h.still_enforcing());
assert!(!h.last_fetch_ok.load(Ordering::Relaxed));
assert_eq!(
h.meta().expect("meta").etag.as_deref(),
Some(etag_for(&good).as_str()),
"a transport failure must not disturb the stored validator"
);
}
#[tokio::test]
async fn rate_limit_honours_retry_after() {
let mut server = mockito::Server::new_async().await;
let mock = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(429)
.with_header("Retry-After", "45")
.create_async()
.await;
let mut h = Harness::new(server.url());
assert_eq!(
h.poller.poll_once().await,
PollOutcome::RateLimited(Some(Duration::from_secs(45)))
);
mock.assert_async().await;
}
#[tokio::test]
async fn rate_limit_without_a_usable_header_falls_back_to_the_interval() {
let mut server = mockito::Server::new_async().await;
let mock = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(429)
.with_header("Retry-After", "when we feel like it")
.create_async()
.await;
let mut h = Harness::new(server.url());
assert_eq!(h.poller.poll_once().await, PollOutcome::RateLimited(None));
mock.assert_async().await;
}
#[tokio::test]
async fn rate_limit_honours_an_http_date_header() {
let future = chrono::Utc::now() + chrono::Duration::seconds(900);
let mut server = mockito::Server::new_async().await;
let mock = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(429)
.with_header(
"Retry-After",
&future.format("%a, %d %b %Y %H:%M:%S GMT").to_string(),
)
.create_async()
.await;
let mut h = Harness::new(server.url());
match h.poller.poll_once().await {
PollOutcome::RateLimited(Some(d)) => assert!(
d >= Duration::from_secs(880) && d <= Duration::from_secs(900),
"expected ~900s from the http-date, got {d:?}"
),
other => panic!("expected a honoured Retry-After, got {other:?}"),
}
mock.assert_async().await;
}
#[tokio::test]
async fn retry_after_is_capped() {
let mut server = mockito::Server::new_async().await;
let mock = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(429)
.with_header("Retry-After", "999999")
.create_async()
.await;
let mut h = Harness::new(server.url());
assert_eq!(
h.poller.poll_once().await,
PollOutcome::RateLimited(Some(Duration::from_secs(MAX_RETRY_AFTER_SECS)))
);
mock.assert_async().await;
}
#[tokio::test]
async fn auth_failure_pauses_until_the_credential_changes() {
let mut server = mockito::Server::new_async().await;
let good = body(6);
let download = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(200)
.with_header("ETag", &etag_for(&good))
.with_body(&good)
.expect(1)
.create_async()
.await;
let mut h = Harness::new(server.url());
assert!(matches!(
h.poller.poll_once().await,
PollOutcome::Activated { .. }
));
download.assert_async().await;
let revoked = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(401)
.expect(1)
.create_async()
.await;
assert_eq!(h.poller.poll_once().await, PollOutcome::AuthFailed);
for _ in 0..3 {
assert_eq!(
h.poller.poll_once().await,
PollOutcome::Skipped("credential_unchanged_after_auth_failure")
);
}
revoked.assert_async().await;
assert!(
h.still_enforcing(),
"a revoked key must not disarm the host"
);
assert!(
!h.cloud_state.is_auth_error(),
"D46: the poller must never write the cloud auth latch"
);
let rotated = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(200)
.with_header("ETag", &etag_for(&good))
.with_body(&good)
.expect(1)
.create_async()
.await;
h.credentials.set_key("ol_org_rotated");
assert!(matches!(
h.poller.poll_once().await,
PollOutcome::Activated { .. }
));
rotated.assert_async().await;
}
#[tokio::test]
async fn auth_refresh_retries_the_same_credential_and_accepts_a_304() {
let mut server = mockito::Server::new_async().await;
let good = body(61);
let download = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(200)
.with_header("ETag", &etag_for(&good))
.with_body(&good)
.expect(1)
.create_async()
.await;
let mut h = Harness::new(server.url());
assert!(matches!(
h.poller.poll_once().await,
PollOutcome::Activated { revision: 61 }
));
download.assert_async().await;
download.remove_async().await;
let rejected = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(401)
.expect(1)
.create_async()
.await;
assert_eq!(h.poller.poll_once().await, PollOutcome::AuthFailed);
rejected.assert_async().await;
rejected.remove_async().await;
h.cloud_state.policy_refresh_notify.notify_one();
assert_eq!(
h.poller.poll_once().await,
PollOutcome::Skipped("credential_unchanged_after_auth_failure"),
"a compatibility refresh must not release an auth refusal"
);
let etag = etag_for(&good);
let revalidated = server
.mock("GET", BUNDLE_ENDPOINT)
.match_header("if-none-match", etag.as_str())
.with_status(304)
.expect(1)
.create_async()
.await;
h.cloud_state.signal_auth_refresh();
assert_eq!(h.poller.poll_once().await, PollOutcome::NotModified);
revalidated.assert_async().await;
assert!(h.poller.failed_credential.is_none());
assert!(h.still_enforcing(), "304 recovery must preserve the LKG");
}
#[tokio::test]
async fn auth_refresh_wakes_the_live_poller_without_waiting_for_the_timer() {
let mut server = mockito::Server::new_async().await;
let rejected = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(401)
.expect(1)
.create_async()
.await;
let dir = tempfile::tempdir().expect("tempdir");
let handle = new_handle(None);
let cloud_state = CloudState::new();
let credentials = TestCredentialProvider::with_key("ol_org_same_key");
let readiness = Arc::new(std::sync::RwLock::new(PolicyReadiness {
phase: PolicyReadinessPhase::Waiting,
last_attempted_at: None,
last_outcome: None,
reason_code: None,
selection_acknowledged: false,
compatibility_diagnostic_valid: false,
}));
let task = tokio::spawn(run_policy_poller(
handle.clone(),
Arc::new(AtomicBool::new(false)),
Arc::new(AtomicI64::new(0)),
Arc::new(AtomicI64::new(0)),
readiness.clone(),
Arc::new(std::sync::RwLock::new(None)),
cloud_state.clone(),
credentials.clone(),
server.url(),
PolicyConfig {
enabled: true,
poll_interval_secs: 300,
stale_warn_after_secs: 86_400,
},
dir.path().to_path_buf(),
crate::egress::ClientHandle::of(crate::egress::client()),
Some(AGENT_ID.to_string()),
None,
crate::egress::EgressReporter::direct(),
None,
));
tokio::time::timeout(Duration::from_secs(2), async {
loop {
if rejected.matched_async().await
&& readiness
.read()
.is_ok_and(|state| state.reason_code == Some("auth_rejected"))
{
break;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
})
.await
.expect("the boot poll must reach the initial 401");
rejected.remove_async().await;
let good = body(62);
let recovery = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(200)
.with_header("ETag", &etag_for(&good))
.with_body(&good)
.expect(1)
.create_async()
.await;
cloud_state.policy_refresh_notify.notify_one();
tokio::time::timeout(Duration::from_secs(2), async {
while credentials.calls.load(Ordering::Relaxed) < 2 {
tokio::time::sleep(Duration::from_millis(10)).await;
}
})
.await
.expect("the compatibility signal must wake one skipped poll tick");
assert!(
!recovery.matched_async().await,
"an unrelated compatibility wake must not retry the rejected credential"
);
cloud_state.signal_auth_refresh();
tokio::time::timeout(Duration::from_secs(2), async {
loop {
if handle
.load()
.as_ref()
.as_ref()
.is_some_and(|bundle| bundle.revision == 62)
{
break;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
})
.await
.expect("the auth refresh must wake the 300 s poll immediately");
task.abort();
recovery.assert_async().await;
}
#[tokio::test]
async fn typed_floor_preserves_lkg_and_does_not_latch_credentials() {
let mut server = mockito::Server::new_async().await;
let good = body(31);
let download = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(200)
.with_header("ETag", &etag_for(&good))
.with_body(&good)
.create_async()
.await;
let mut h = Harness::new(server.url());
assert!(matches!(
h.poller.poll_once().await,
PollOutcome::Activated { .. }
));
download.assert_async().await;
let before_body = std::fs::read(store::bundle_path(h.dir.path())).expect("body");
let before_meta = h.meta().expect("meta");
let floor = server.mock("GET", BUNDLE_ENDPOINT).with_status(403)
.with_body(format!(r#"{{"type":"{FLOOR_PROBLEM_TYPE}","client_version":"{}","min_client_version":"9.0.0"}}"#, env!("CARGO_PKG_VERSION")))
.create_async().await;
assert_eq!(h.poller.poll_once().await, PollOutcome::FloorBlocked);
floor.assert_async().await;
assert!(h.still_enforcing());
assert!(h.poller.failed_credential.is_none());
assert_eq!(
h.poller
.floor
.read()
.expect("floor lock")
.as_ref()
.map(|state| state.minimum_version.as_str()),
Some("9.0.0")
);
assert_eq!(
std::fs::read(store::bundle_path(h.dir.path())).expect("body"),
before_body
);
let after_meta = h.meta().expect("meta");
assert_eq!(after_meta.digest, before_meta.digest);
assert_eq!(after_meta.etag, before_meta.etag);
assert_eq!(after_meta.state.as_deref(), Some("floor_blocked"));
let retry = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(304)
.create_async()
.await;
assert_eq!(h.poller.poll_once().await, PollOutcome::NotModified);
retry.assert_async().await;
assert!(h.poller.floor.read().expect("floor lock").is_some());
let recovery = server
.mock("GET", BUNDLE_ENDPOINT)
.match_header("if-none-match", mockito::Matcher::Missing)
.with_status(200)
.with_header("ETag", &etag_for(&good))
.with_body(&good)
.create_async()
.await;
assert!(matches!(
h.poller.poll_once().await,
PollOutcome::Activated { .. }
));
recovery.assert_async().await;
assert!(h.poller.floor.read().expect("floor lock").is_none());
}
#[tokio::test]
async fn non_floor_forbidden_is_generic_and_retries_normally() {
let mut server = mockito::Server::new_async().await;
let forbidden = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(403)
.with_header("Content-Type", "application/problem+json")
.with_body(r#"{"type":"https://example.test/other"}"#)
.create_async()
.await;
let mut h = Harness::new(server.url());
h.poller.failed_credential = Some(b"older-credential".to_vec());
assert_eq!(h.poller.poll_once().await, PollOutcome::Failed);
forbidden.assert_async().await;
assert!(h.poller.failed_credential.is_none());
let good = body(32);
let retry = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(200)
.with_header("ETag", &etag_for(&good))
.with_body(&good)
.create_async()
.await;
assert!(matches!(
h.poller.poll_once().await,
PollOutcome::Activated { .. }
));
retry.assert_async().await;
}
#[tokio::test]
async fn typed_floor_is_recognized_with_an_alternate_content_type() {
let mut server = mockito::Server::new_async().await;
let response = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(403)
.with_header("Content-Type", "text/plain")
.with_body(format!(
r#"{{"type":"{FLOOR_PROBLEM_TYPE}","min_client_version":"9.0.0"}}"#
))
.create_async()
.await;
let mut h = Harness::new(server.url());
assert_eq!(h.poller.poll_once().await, PollOutcome::FloorBlocked);
response.assert_async().await;
assert!(h.poller.failed_credential.is_none());
}
#[tokio::test]
async fn malformed_forbidden_body_is_a_generic_failure() {
let mut server = mockito::Server::new_async().await;
let response = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(403)
.with_body("not json")
.create_async()
.await;
let mut h = Harness::new(server.url());
assert_eq!(h.poller.poll_once().await, PollOutcome::Failed);
response.assert_async().await;
assert!(h.poller.failed_credential.is_none());
}
const I2_LICENSE_BODY: &str = concat!(
r#"{"error":{"code":"license_expired","message":"This organization's licence has expired.","#,
r#""details":[{"licensing_url":"https://app.openlatch.ai/settings/licensing"}]}}"#
);
#[tokio::test]
async fn a_402_blocks_the_bundle_without_latching_the_credential() {
let mut server = mockito::Server::new_async().await;
let good = body(51);
let first = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(200)
.with_header("ETag", &etag_for(&good))
.with_body(&good)
.expect(1)
.create_async()
.await;
let mut h = Harness::new(server.url());
assert!(matches!(
h.poller.poll_once().await,
PollOutcome::Activated { .. }
));
first.assert_async().await;
let etag_before = h.meta().expect("meta").etag;
h.last_poll_ok_at.store(0, Ordering::Relaxed);
let refused = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(402)
.with_header("Retry-After", "900")
.with_body(I2_LICENSE_BODY)
.expect(1)
.create_async()
.await;
assert_eq!(
h.poller.poll_once().await,
PollOutcome::LicenseBlocked(Some(Duration::from_secs(900)))
);
refused.assert_async().await;
assert!(h.still_enforcing(), "the resident bundle keeps deciding");
assert!(
h.poller.failed_credential.is_none(),
"a licence refusal is never a credential failure"
);
assert!(
h.last_poll_ok_at.load(Ordering::Relaxed) > 0,
"a typed refusal is a completed round trip"
);
let meta = h.meta().expect("meta");
assert_eq!(meta.state.as_deref(), Some("license_blocked"));
assert_eq!(meta.etag, etag_before, "the validator survives");
assert!(meta
.last_error
.as_deref()
.is_some_and(|e| e.contains("OL-1250")));
let refusal = h
.cloud_state
.policy_license_refusal()
.expect("the refusal is reportable");
assert_eq!(refusal.code, "license_expired");
assert_eq!(
refusal.licensing_url.as_deref(),
Some("https://app.openlatch.ai/settings/licensing")
);
assert!(
h.cloud_state.license_gate().is_none(),
"a refused bundle never holds ingest POSTs"
);
}
#[tokio::test]
async fn the_poll_after_a_block_asks_for_a_full_response() {
let mut server = mockito::Server::new_async().await;
let good = body(52);
let first = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(200)
.with_header("ETag", &etag_for(&good))
.with_body(&good)
.expect(1)
.create_async()
.await;
let mut h = Harness::new(server.url());
assert!(matches!(
h.poller.poll_once().await,
PollOutcome::Activated { .. }
));
first.assert_async().await;
let refused = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(402)
.with_body(I2_LICENSE_BODY)
.expect(1)
.create_async()
.await;
assert!(matches!(
h.poller.poll_once().await,
PollOutcome::LicenseBlocked(_)
));
refused.assert_async().await;
let renewed = body(53);
let recovery = server
.mock("GET", BUNDLE_ENDPOINT)
.match_header("if-none-match", mockito::Matcher::Missing)
.with_status(200)
.with_header("ETag", &etag_for(&renewed))
.with_body(&renewed)
.expect(1)
.create_async()
.await;
assert!(matches!(
h.poller.poll_once().await,
PollOutcome::Activated { revision: 53 }
));
recovery.assert_async().await;
assert!(
h.cloud_state.policy_license_refusal().is_none(),
"an activated bundle proves the licence is back"
);
}
#[tokio::test]
async fn only_a_site_coded_503_is_a_licence_block() {
let mut server = mockito::Server::new_async().await;
let plain = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(503)
.with_body("upstream down")
.expect(1)
.create_async()
.await;
let mut h = Harness::new(server.url());
assert_eq!(h.poller.poll_once().await, PollOutcome::Failed);
plain.assert_async().await;
let sited = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(503)
.with_body(r#"{"error":{"code":"clock_regression"}}"#)
.expect(1)
.create_async()
.await;
assert_eq!(
h.poller.poll_once().await,
PollOutcome::LicenseBlocked(None)
);
sited.assert_async().await;
assert!(
h.poller.failed_credential.is_none(),
"a site refusal is not a credential failure either"
);
assert!(h.meta().is_none());
}
#[tokio::test]
async fn latched_cloud_auth_error_skips_the_fetch() {
let mut server = mockito::Server::new_async().await;
let never = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(200)
.expect(0)
.create_async()
.await;
let mut h = Harness::new(server.url());
h.cloud_state.auth_error.store(true, Ordering::Relaxed);
assert_eq!(
h.poller.poll_once().await,
PollOutcome::Skipped("auth_error_latched")
);
assert!(
h.cloud_state.is_auth_error(),
"the poller must not clear a latch it does not own"
);
never.assert_async().await;
}
#[tokio::test]
async fn missing_credential_skips_the_fetch() {
let mut server = mockito::Server::new_async().await;
let never = server
.mock("GET", BUNDLE_ENDPOINT)
.expect(0)
.create_async()
.await;
let mut h = Harness::new(server.url());
*h.credentials.key.lock().expect("lock") = None;
assert_eq!(
h.poller.poll_once().await,
PollOutcome::Skipped("no_credential")
);
never.assert_async().await;
}
#[tokio::test]
async fn staleness_warns_and_keeps_enforcing() {
let mut server = mockito::Server::new_async().await;
let good = body(8);
let etag = etag_for(&good);
let download = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(200)
.with_header("ETag", &etag)
.with_body(&good)
.expect(1)
.create_async()
.await;
let mut h = Harness::new(server.url());
assert!(matches!(
h.poller.poll_once().await,
PollOutcome::Activated { .. }
));
download.assert_async().await;
assert!(!h.poller.check_staleness(), "a fresh poll is not stale");
h.last_poll_ok_at.store(
now_unix_secs() - (h.poller.config.stale_warn_after_secs as i64) - 60,
Ordering::Relaxed,
);
assert!(h.poller.check_staleness(), "OL-1213 must fire");
assert!(
h.still_enforcing(),
"staleness must never stop the host enforcing"
);
let revalidate = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(304)
.create_async()
.await;
assert_eq!(h.poller.poll_once().await, PollOutcome::NotModified);
revalidate.assert_async().await;
assert!(!h.poller.check_staleness());
}
#[test]
fn a_host_that_never_polled_successfully_does_not_warn_as_stale() {
let h = Harness::new("http://127.0.0.1:1".to_string());
assert_eq!(h.last_poll_ok_at.load(Ordering::Relaxed), 0);
assert!(!h.poller.check_staleness());
}
#[tokio::test]
async fn poll_clock_is_seeded_from_disk() {
let mut server = mockito::Server::new_async().await;
let good = body(12);
let download = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(200)
.with_header("ETag", &etag_for(&good))
.with_body(&good)
.expect(1)
.create_async()
.await;
let mut h = Harness::new(server.url());
assert!(matches!(
h.poller.poll_once().await,
PollOutcome::Activated { .. }
));
download.assert_async().await;
let restarted = PolicyPoller::new(
new_handle(None),
Arc::new(AtomicBool::new(false)),
Arc::new(AtomicI64::new(0)),
Arc::new(AtomicI64::new(0)),
Arc::new(std::sync::RwLock::new(PolicyReadiness {
phase: PolicyReadinessPhase::Waiting,
last_attempted_at: None,
last_outcome: None,
reason_code: None,
selection_acknowledged: false,
compatibility_diagnostic_valid: false,
})),
CloudState::new(),
TestCredentialProvider::with_key("ol_org_test"),
server.url(),
PolicyConfig {
enabled: true,
poll_interval_secs: 300,
stale_warn_after_secs: 86_400,
},
h.dir.path().to_path_buf(),
crate::egress::ClientHandle::of(crate::egress::client()),
Some(AGENT_ID.to_string()),
None,
);
restarted.seed_poll_clock();
assert!(restarted.last_poll_ok_at.load(Ordering::Relaxed) > 0);
}
#[tokio::test]
async fn a_client_swapped_into_the_handle_is_used_by_the_next_poll() {
let mut server = mockito::Server::new_async().await;
let good = body(31);
let mock = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(200)
.with_header("ETag", &etag_for(&good))
.with_body(&good)
.expect(1)
.create_async()
.await;
let clients = crate::egress::EgressClients::new(
crate::egress::Timeouts::default(),
crate::egress::Timeouts::default(),
);
let dir = tempfile::tempdir().expect("tempdir");
let mut poller = PolicyPoller::new(
new_handle(None),
Arc::new(AtomicBool::new(false)),
Arc::new(AtomicI64::new(0)),
Arc::new(AtomicI64::new(0)),
Arc::new(std::sync::RwLock::new(PolicyReadiness {
phase: PolicyReadinessPhase::Waiting,
last_attempted_at: None,
last_outcome: None,
reason_code: None,
selection_acknowledged: false,
compatibility_diagnostic_valid: false,
})),
CloudState::new(),
TestCredentialProvider::with_key("ol_org_test"),
server.url(),
PolicyConfig {
enabled: true,
poll_interval_secs: 300,
stale_warn_after_secs: 86_400,
},
dir.path().to_path_buf(),
clients.poller.clone(),
Some(AGENT_ID.to_string()),
None,
);
assert!(
matches!(
poller.poll_once().await,
PollOutcome::Skipped("no_egress_route")
),
"an empty handle must skip the poll rather than reach for a route it does not have"
);
clients
.apply(&crate::egress::EgressConfig::direct())
.expect("install a route");
assert!(
matches!(poller.poll_once().await, PollOutcome::Activated { .. }),
"the next poll must use the newly installed client"
);
mock.assert_async().await;
}
#[tokio::test]
async fn the_bundle_request_carries_the_host_key() {
const HOST: &str = "9262b296baa37a205734509345581251";
let mut server = mockito::Server::new_async().await;
let good = body(37);
let mock = server
.mock("GET", BUNDLE_ENDPOINT)
.match_header("x-openlatch-agent-id", AGENT_ID)
.match_header("x-openlatch-machine-id", HOST)
.with_status(200)
.with_header("ETag", &etag_for(&good))
.with_body(&good)
.expect(1)
.create_async()
.await;
let dir = tempfile::tempdir().expect("tempdir");
let mut poller = PolicyPoller::new(
new_handle(None),
Arc::new(AtomicBool::new(false)),
Arc::new(AtomicI64::new(0)),
Arc::new(AtomicI64::new(0)),
Arc::new(std::sync::RwLock::new(PolicyReadiness {
phase: PolicyReadinessPhase::Waiting,
last_attempted_at: None,
last_outcome: None,
reason_code: None,
selection_acknowledged: false,
compatibility_diagnostic_valid: false,
})),
CloudState::new(),
TestCredentialProvider::with_key("ol_org_test"),
server.url(),
PolicyConfig {
enabled: true,
poll_interval_secs: 300,
stale_warn_after_secs: 86_400,
},
dir.path().to_path_buf(),
crate::egress::ClientHandle::of(crate::egress::client()),
Some(AGENT_ID.to_string()),
Some(HOST.to_string()),
);
assert!(matches!(
poller.poll_once().await,
PollOutcome::Activated { .. }
));
mock.assert_async().await;
}
#[tokio::test]
async fn boot_fetch_happens_before_the_first_tick() {
let mut server = mockito::Server::new_async().await;
let good = body(21);
let mock = server
.mock("GET", BUNDLE_ENDPOINT)
.with_status(200)
.with_header("ETag", &etag_for(&good))
.with_body(&good)
.create_async()
.await;
let dir = tempfile::tempdir().expect("tempdir");
let handle = new_handle(None);
let task = tokio::spawn(run_policy_poller(
handle.clone(),
Arc::new(AtomicBool::new(false)),
Arc::new(AtomicI64::new(0)),
Arc::new(AtomicI64::new(0)),
Arc::new(std::sync::RwLock::new(PolicyReadiness {
phase: PolicyReadinessPhase::Waiting,
last_attempted_at: None,
last_outcome: None,
reason_code: None,
selection_acknowledged: false,
compatibility_diagnostic_valid: false,
})),
Arc::new(std::sync::RwLock::new(None)),
CloudState::new(),
TestCredentialProvider::with_key("ol_org_test"),
server.url(),
PolicyConfig {
enabled: true,
poll_interval_secs: 3_600,
stale_warn_after_secs: 86_400,
},
dir.path().to_path_buf(),
crate::egress::ClientHandle::of(crate::egress::client()),
Some(AGENT_ID.to_string()),
None,
crate::egress::EgressReporter::direct(),
None,
));
let mut activated = false;
for _ in 0..50 {
if handle.load().is_some() {
activated = true;
break;
}
tokio::time::sleep(Duration::from_millis(20)).await;
}
task.abort();
assert!(
activated,
"the poller must fetch on boot, not on the first timer tick"
);
assert_eq!(
handle.load().as_ref().as_ref().map(|b| b.revision),
Some(21)
);
mock.assert_async().await;
}
}