bobbin-ai 0.25.2

Local-first context injection engine for AI coding agents
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))]
// Replacing a large repository snapshot is a real remote write transaction, not
// a health probe. Production-scale snapshots can legitimately exceed the old
// 15-second deadline while the remote service remains responsive to small reads.
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;

/// Push a replaceable chunk snapshot to the configured remote Quipu ontology.
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
}

/// quipu `snapshot_upload_id`, which `/knot/stage` recomputes and refuses on mismatch. The
/// graph is NOT part of it: quipu keeps the graph in the immutable upload manifest, so the
/// same content staged for another graph fails closed (aegis-86f2v7).
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}"),
        });
        // Only when set, so a ROOT stage body is byte-identical to before (aegis-86f2v7).
        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);
        // Parts are JSON strings: never split a UTF-8 character or replace its
        // bytes, since promotion verifies the hash of the original snapshot.
        while !turtle.is_char_boundary(end) {
            end -= 1;
        }
        let (part, remaining) = turtle.split_at(end);
        parts.push(part);
        turtle = remaining;
    }
    parts
}

/// Prove the remote store routes `/knot` by graph before staging into `graph`: a write
/// aimed at an unregistered sentinel must be REFUSED as an unknown graph (the key is
/// honoured), and an empty write to `graph` must SUCCEED (it is registered). Anything
/// else refuses the push rather than risk the chunks landing in ROOT.
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}"));
                }
                // A deterministic 4xx cannot improve on retry. 5xx is
                // indeterminate and safe to retry because stage/promote are
                // content-addressed and idempotent.
                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));
    }
}