use std::path::{Path, PathBuf};
use std::time::{SystemTime, UNIX_EPOCH};
use super::{decode_tripwires, GovernedTripwire, Provenance};
pub const TRIPWIRE_QUERY: &str = "\
PREFIX aegis: <http://aegis.gastown.local/ontology/>
PREFIX rdfs: <http://www.w3.org/2000/01/rdf-schema#>
SELECT ?policy ?name ?appliesTo ?effect ?claim ?constraintClass ?verificationPoint
?backoffFormula ?selector ?predicate WHERE {
?policy a aegis:Policy ;
aegis:boundary \"action\" ;
aegis:appliesTo ?appliesTo .
OPTIONAL { ?policy aegis:selector ?selector }
OPTIONAL { ?policy aegis:predicate ?predicate }
OPTIONAL { ?policy rdfs:label ?name }
OPTIONAL { ?policy aegis:effect ?effect }
OPTIONAL { ?policy aegis:claim ?claim }
OPTIONAL { ?policy aegis:constraintClass ?constraintClass }
OPTIONAL { ?policy aegis:verificationPoint ?verificationPoint }
OPTIONAL { ?policy aegis:backoffFormula ?backoffFormula }
}";
const TTL_SECS: u64 = 300;
const TIMEOUT_SECS: u64 = 2;
#[must_use]
pub fn endpoint(config: &crate::config::Config) -> Option<String> {
std::env::var("BOBBIN_QUIPU_REMOTE")
.ok()
.or_else(|| config.quipu_endpoint.clone())
.map(|s| s.trim().trim_end_matches('/').to_string())
.filter(|s| !s.is_empty())
}
#[must_use]
pub fn cache_path(repo_root: &Path) -> PathBuf {
repo_root.join(".bobbin").join("tripwire-cache.json")
}
fn now_secs() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|d| d.as_secs())
.unwrap_or(0)
}
#[derive(serde::Serialize, serde::Deserialize)]
struct CacheFile {
endpoint: String,
fetched_at: u64,
body: String,
}
pub async fn load_tripwires(
config: &crate::config::Config,
repo_root: &Path,
) -> Option<(Vec<GovernedTripwire>, Provenance)> {
let endpoint = endpoint(config)?;
let path = cache_path(repo_root);
let cached = std::fs::read_to_string(&path)
.ok()
.and_then(|s| serde_json::from_str::<CacheFile>(&s).ok())
.filter(|c| c.endpoint == endpoint);
let now = now_secs();
if let Some(ref c) = cached {
let age = now.saturating_sub(c.fetched_at);
if age < TTL_SECS {
if let Ok(wires) = decode_tripwires(&c.body) {
return Some((
wires,
Provenance::Cached {
endpoint,
age_secs: age,
refresh_error: None,
},
));
}
}
}
match fetch(&endpoint).await {
Ok(body) => match decode_tripwires(&body) {
Ok(wires) => {
store(&path, &endpoint, now, &body);
Some((wires, Provenance::Live { endpoint }))
}
Err(e) => degrade(
cached,
endpoint,
now,
format!("undecodable response: {e:#}"),
),
},
Err(e) => degrade(cached, endpoint, now, format!("{e:#}")),
}
}
fn degrade(
cached: Option<CacheFile>,
endpoint: String,
now: u64,
reason: String,
) -> Option<(Vec<GovernedTripwire>, Provenance)> {
let c = cached?;
let wires = decode_tripwires(&c.body).ok()?;
Some((
wires,
Provenance::Cached {
endpoint,
age_secs: now.saturating_sub(c.fetched_at),
refresh_error: Some(reason),
},
))
}
fn store(path: &Path, endpoint: &str, fetched_at: u64, body: &str) {
let file = CacheFile {
endpoint: endpoint.to_string(),
fetched_at,
body: body.to_string(),
};
if let Some(parent) = path.parent() {
let _ = std::fs::create_dir_all(parent);
}
if let Ok(json) = serde_json::to_string(&file) {
let _ = std::fs::write(path, json);
}
}
async fn fetch(endpoint: &str) -> anyhow::Result<String> {
let client = reqwest::Client::builder()
.timeout(std::time::Duration::from_secs(TIMEOUT_SECS))
.build()?;
let resp = client
.post(format!("{endpoint}/query"))
.json(&serde_json::json!({ "query": TRIPWIRE_QUERY }))
.send()
.await?;
let status = resp.status();
let text = resp.text().await.unwrap_or_default();
if !status.is_success() {
anyhow::bail!(
"quipu /query returned HTTP {status}: {}",
text.chars().take(200).collect::<String>()
);
}
Ok(text)
}