use std::path::PathBuf;
use std::time::Duration;
use anyhow::{Context, Result};
use sha2::{Digest, Sha256};
use crate::types::{Chunk, ChunkEdge};
const QUIPU_CLIENT: &str = "ingest-cron";
#[cfg(not(test))]
const TIMEOUT: Duration = Duration::from_secs(30);
#[cfg(not(test))]
const PROMOTE_DEADLINE: Duration = Duration::from_secs(900);
#[cfg(test)]
const TIMEOUT: Duration = Duration::from_millis(100);
#[cfg(test)]
const PROMOTE_DEADLINE: Duration = Duration::from_secs(1);
const PART_BYTES: usize = 256 * 1024;
const MAX_ATTEMPTS: usize = 3;
pub async fn push_chunks_to_remote_quipu(
chunks: &[Chunk],
edges: &[ChunkEdge],
repo_name: &str,
endpoint: &str,
target_graph: Option<&str>,
) -> Result<(i64, usize)> {
let token = quipu_auth_token().context(
"quipu_push_chunks targets a remote ontology but no QUIPU_AUTH_TOKEN or readable token file is available",
)?;
push_with_token(chunks, edges, repo_name, endpoint, &token, target_graph).await
}
fn upload_id(snapshot: &str, content_hash: &str) -> String {
sha256(format!("{snapshot}\n{content_hash}").as_bytes())
}
async fn push_with_token(
chunks: &[Chunk],
edges: &[ChunkEdge],
repo_name: &str,
endpoint: &str,
token: &str,
target_graph: Option<&str>,
) -> Result<(i64, usize)> {
let turtle = super::generate_chunk_turtle(chunks, edges, repo_name);
let snapshot = format!("bobbin-chunks:{repo_name}");
let content_hash = sha256(turtle.as_bytes());
let upload_id = upload_id(&snapshot, &content_hash);
let parts = snapshot_parts(&turtle);
anyhow::ensure!(!parts.is_empty(), "refusing an empty chunk snapshot upload");
let client = reqwest::Client::builder()
.timeout(TIMEOUT)
.build()
.context("building remote Quipu client")?;
let base = endpoint.trim_end_matches('/');
if let Some(graph) = target_graph {
require_remote_graph_routing(&client, base, token, graph).await?;
}
for (part_number, payload) in parts.iter().enumerate() {
let mut body = serde_json::json!({
"upload_id": upload_id, "snapshot": snapshot, "content_hash": content_hash,
"total_parts": parts.len(), "total_bytes": turtle.len(),
"part_number": part_number, "part_hash": sha256(payload.as_bytes()), "payload": payload,
"actor": "bobbin", "source": format!("bobbin chunk index: {repo_name}"),
});
if let Some(graph) = target_graph {
body["graph"] = serde_json::Value::from(graph);
}
post_with_retries(&client, &format!("{base}/knot/stage"), token, &body)
.await
.with_context(|| {
format!(
"staging chunk snapshot for repository {repo_name}, part {part_number}/{} ({} bytes)",
parts.len(), payload.len()
)
})?;
}
let promote_url = format!("{base}/knot/promote");
let promote_body = serde_json::json!({"upload_id": upload_id});
let started = tokio::time::Instant::now();
let result = loop {
let result = post_with_retries(&client, &promote_url, token, &promote_body).await?;
if result.get("pending").and_then(|v| v.as_bool()) != Some(true) {
break result;
}
anyhow::ensure!(
started.elapsed() < PROMOTE_DEADLINE,
"remote Quipu snapshot promotion remained pending for {} seconds",
PROMOTE_DEADLINE.as_secs()
);
tokio::time::sleep(Duration::from_secs(1)).await;
};
if result.get("conforms").and_then(|v| v.as_bool()) == Some(false) {
anyhow::bail!("remote Quipu refused chunk snapshot by SHACL: {result}");
}
if result.get("replaced").and_then(|v| v.as_bool()) != Some(true)
|| result.get("promoted").and_then(|v| v.as_bool()) != Some(true)
|| result.get("content_hash").and_then(|v| v.as_str()) != Some(content_hash.as_str())
{
anyhow::bail!(
"remote Quipu did not confirm the exact promoted snapshot; refusing success: {result}"
);
}
Ok((
result["tx_id"].as_i64().unwrap_or(-1),
result["count"].as_u64().unwrap_or(0) as usize,
))
}
fn snapshot_parts(mut turtle: &str) -> Vec<&str> {
let mut parts = Vec::new();
while !turtle.is_empty() {
let mut end = turtle.len().min(PART_BYTES);
while !turtle.is_char_boundary(end) {
end -= 1;
}
let (part, remaining) = turtle.split_at(end);
parts.push(part);
turtle = remaining;
}
parts
}
async fn require_remote_graph_routing(
client: &reqwest::Client,
base: &str,
token: &str,
graph: &str,
) -> Result<()> {
let knot = format!("{base}/knot");
let body = |g: &str| serde_json::json!({"turtle": "", "actor": "bobbin", "source": "chunk-graph-probe", "graph": g});
let sentinel = format!("{graph}/bobbin-unregistered-sentinel");
match post_with_retries(client, &knot, token, &body(&sentinel)).await {
Ok(_) => anyhow::bail!(
"remote quipu ACCEPTED a write aimed at an unregistered sentinel graph, so it is \
dropping the /knot 'graph' key: chunks meant for <{graph}> would land in ROOT. \
Refusing to push."
),
Err(e) if e.to_string().contains("unknown graph") => {}
Err(e) => anyhow::bail!("remote quipu graph-routing probe failed: {e}"),
}
post_with_retries(client, &knot, token, &body(graph))
.await
.with_context(|| format!("target graph <{graph}> is not usable (register it first)"))?;
Ok(())
}
fn sha256(bytes: &[u8]) -> String {
format!("sha256:{}", hex::encode(Sha256::digest(bytes)))
}
async fn post_with_retries(
client: &reqwest::Client,
url: &str,
token: &str,
body: &serde_json::Value,
) -> Result<serde_json::Value> {
post_with_retries_timeout(client, url, token, body, TIMEOUT).await
}
async fn post_with_retries_timeout(
client: &reqwest::Client,
url: &str,
token: &str,
body: &serde_json::Value,
timeout: Duration,
) -> Result<serde_json::Value> {
let mut last_error = None;
for attempt in 1..=MAX_ATTEMPTS {
match client
.post(url)
.timeout(timeout)
.header("X-Quipu-Client", QUIPU_CLIENT)
.bearer_auth(token)
.json(body)
.send()
.await
{
Ok(response) => {
let status = response.status();
let text = response.text().await.unwrap_or_default();
if status.is_success() {
return serde_json::from_str(&text)
.with_context(|| format!("parsing remote Quipu response from {url}"));
}
if status.is_client_error() {
anyhow::bail!(
"remote Quipu returned HTTP {status}: {}",
text.chars().take(300).collect::<String>()
);
}
last_error = Some(format!(
"HTTP {status}: {}",
text.chars().take(300).collect::<String>()
));
}
Err(error) => last_error = Some(error.to_string()),
}
if attempt < MAX_ATTEMPTS {
tokio::time::sleep(Duration::from_millis(100 * attempt as u64)).await;
}
}
anyhow::bail!(
"POST {url} failed after {MAX_ATTEMPTS} idempotent attempts: {}",
last_error.unwrap_or_else(|| "unknown error".into())
)
}
pub(crate) fn quipu_auth_token() -> Option<String> {
std::env::var("QUIPU_AUTH_TOKEN")
.ok()
.filter(|token| !token.trim().is_empty())
.or_else(|| {
let path = std::env::var_os("QUIPU_AUTH_TOKEN_FILE")
.map(PathBuf::from)
.or_else(|| {
directories::BaseDirs::new()
.map(|dirs| dirs.home_dir().join(".config/aegis/quipu_token"))
})?;
std::fs::read_to_string(path)
.ok()
.map(|token| token.trim().to_string())
.filter(|token| !token.is_empty())
})
}
#[cfg(test)]
#[path = "remote_graph_tests.rs"]
mod graph_tests;
#[cfg(test)]
mod tests {
use super::*;
use axum::{extract::State, http::HeaderMap, routing::post, Json, Router};
use std::sync::{Arc, Mutex};
fn chunk() -> Chunk {
Chunk {
id: "h".into(),
file_path: "docs/a.md".into(),
chunk_type: crate::types::ChunkType::Section,
name: Some("A".into()),
start_line: 1,
end_line: 2,
content: "body".into(),
language: "markdown".into(),
tags: String::new(),
}
}
#[test]
fn snapshot_parts_preserve_bytes_at_multibyte_boundaries() {
for ch in ['é', '界', '🦀'] {
for split in 1..ch.len_utf8() {
let snapshot = format!(
"{}{}{}{}tail",
"a".repeat(PART_BYTES - split),
ch,
"b".repeat(PART_BYTES - ch.len_utf8() - 1),
ch
);
let parts = snapshot_parts(&snapshot);
assert!(parts.len() >= 3);
assert!(parts.iter().all(|p| !p.is_empty() && p.len() <= PART_BYTES));
assert_eq!(parts.concat(), snapshot);
assert_eq!(
sha256(parts.concat().as_bytes()),
sha256(snapshot.as_bytes())
);
}
}
assert!(snapshot_parts("").is_empty());
for size in [
1,
PART_BYTES - 1,
PART_BYTES,
PART_BYTES + 1,
PART_BYTES * 2,
] {
let snapshot = "a".repeat(size);
let parts = snapshot_parts(&snapshot);
assert_eq!(parts.len(), size.div_ceil(PART_BYTES));
assert_eq!(parts.concat(), snapshot);
}
}
#[tokio::test]
async fn snapshot_is_authenticated_and_bounded() {
#[derive(Clone, Default)]
struct Seen(Arc<Mutex<Vec<(HeaderMap, serde_json::Value)>>>);
async fn stage(
State(seen): State<Seen>,
headers: HeaderMap,
Json(body): Json<serde_json::Value>,
) -> Json<serde_json::Value> {
seen.0.lock().unwrap().push((headers, body));
Json(serde_json::json!({"idempotent": false}))
}
async fn promote(
State(seen): State<Seen>,
headers: HeaderMap,
Json(body): Json<serde_json::Value>,
) -> Json<serde_json::Value> {
let content_hash = seen.0.lock().unwrap()[0].1["content_hash"].clone();
seen.0.lock().unwrap().push((headers, body));
Json(serde_json::json!({
"conforms": true, "replaced": true, "promoted": true,
"content_hash": content_hash, "tx_id": 42, "count": 7
}))
}
let seen = Seen::default();
let app = Router::new()
.route("/knot/stage", post(stage))
.route("/knot/promote", post(promote))
.with_state(seen.clone());
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
tokio::spawn(async move { axum::serve(listener, app).await.unwrap() });
let mut large_chunk = chunk();
large_chunk.name = Some("界".repeat(PART_BYTES));
let chunks = [large_chunk];
let expected = super::super::generate_chunk_turtle(&chunks, &[], "repo");
assert_eq!(
push_with_token(
&chunks,
&[],
"repo",
&format!("http://{addr}"),
"secret",
None
)
.await
.unwrap(),
(42, 7)
);
let guard = seen.0.lock().unwrap();
assert!(guard.len() > 2);
let stages = &guard[..guard.len() - 1];
let mut reconstructed = String::new();
for (number, (_, body)) in stages.iter().enumerate() {
let payload = body["payload"].as_str().unwrap();
assert!(payload.len() <= PART_BYTES);
assert_eq!(body["part_number"], number);
assert_eq!(body["total_parts"], stages.len());
assert_eq!(body["total_bytes"], expected.len());
assert_eq!(body["part_hash"], sha256(payload.as_bytes()));
assert_eq!(body["content_hash"], sha256(expected.as_bytes()));
reconstructed.push_str(payload);
}
assert_eq!(reconstructed, expected);
let headers = &guard[0].0;
assert_eq!(headers["authorization"], "Bearer secret");
assert_eq!(headers["x-quipu-client"], QUIPU_CLIENT);
assert_eq!(
guard.last().unwrap().1["upload_id"],
guard[0].1["upload_id"]
);
drop(guard);
async fn stalled() -> Json<serde_json::Value> {
tokio::time::sleep(Duration::from_secs(1)).await;
Json(serde_json::json!({"conforms": true, "replaced": true}))
}
let app = Router::new().route("/knot/stage", post(stalled));
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
tokio::spawn(async move { axum::serve(listener, app).await.unwrap() });
let started = std::time::Instant::now();
assert!(push_with_token(
&[chunk()],
&[],
"repo",
&format!("http://{addr}"),
"secret",
None
)
.await
.is_err());
assert!(started.elapsed() < Duration::from_secs(1));
}
}