Merge: paper library — corpus, arXiv harvest, vault catalogue, API
ci / gates (push) Failing after 7s
ci / rust (push) Skipped
ci / frontend (push) Skipped
ci / e2e (push) Skipped
ci / publish (push) Skipped

This commit is contained in:
Omar Sobh
2026-08-03 10:53:20 -07:00
15 changed files with 1977 additions and 4 deletions
Generated
+1
View File
@@ -946,6 +946,7 @@ dependencies = [
"cm-config", "cm-config",
"cm-db", "cm-db",
"cm-domain", "cm-domain",
"cm-files",
"cm-llm", "cm-llm",
"cm-orchestrator", "cm-orchestrator",
"cm-runtime", "cm-runtime",
+2 -1
View File
@@ -266,7 +266,7 @@ async fn run() -> Result<(), String> {
terminals, terminals,
providers: provider_registry, providers: provider_registry,
}, },
blob, blob.clone(),
); );
// Durable §15 path: expires overdue approvals and resumes decided runs // Durable §15 path: expires overdue approvals and resumes decided runs
// even if the deciding request's process died mid-flight. // even if the deciding request's process died mid-flight.
@@ -388,6 +388,7 @@ async fn run() -> Result<(), String> {
.with_broker(PathBuf::from(&config.broker.socket_path)) .with_broker(PathBuf::from(&config.broker.socket_path))
.with_oauth(config.oauth.clone()) .with_oauth(config.oauth.clone())
.with_billing(config.billing.clone()) .with_billing(config.billing.clone())
.with_blobs(blob.clone())
.with_file_root( .with_file_root(
(config.storage.backend == cm_config::StorageBackend::Local) (config.storage.backend == cm_config::StorageBackend::Local)
.then(|| PathBuf::from(&config.storage.data_dir)), .then(|| PathBuf::from(&config.storage.data_dir)),
+1
View File
@@ -32,6 +32,7 @@ cm-brain = { path = "../cm-brain" }
cm-config = { path = "../cm-config" } cm-config = { path = "../cm-config" }
cm-db = { path = "../cm-db" } cm-db = { path = "../cm-db" }
cm-domain = { path = "../cm-domain" } cm-domain = { path = "../cm-domain" }
cm-files = { path = "../cm-files" }
cm-llm = { path = "../cm-llm" } cm-llm = { path = "../cm-llm" }
cm-orchestrator = { path = "../cm-orchestrator", features = ["provider"] } cm-orchestrator = { path = "../cm-orchestrator", features = ["provider"] }
cm-runtime = { path = "../cm-runtime" } cm-runtime = { path = "../cm-runtime" }
+466
View File
@@ -0,0 +1,466 @@
//! What a continuous mission has already covered.
//!
//! A recurring mission's hard problem is not running the agent — that is 23
//! seconds — it is knowing what it already did last time. A research mission
//! with no memory of prior runs resurfaces the same papers forever and reports
//! success every time.
//!
//! This module keeps that record. It is deliberately small: an index derived
//! from the corpus, never the corpus itself. The vault is the source of truth,
//! the index is rebuildable, and a hand-edited note is never "wrong".
//!
//! # Two kinds, because the real vault forced it
//!
//! The plan assumed notes would carry `arxiv:` / `doi:` / `url:` frontmatter.
//! Measured against the actual vault: **416 notes, 145 with frontmatter, and
//! zero with any of those keys.** The dominant keys are repo-sync metadata
//! (`node`, `org`, `gitea`) and course-note fields (`presenter`, `session`).
//! An ingester keyed only on external identity would have indexed nothing —
//! the same shape of failure as everything else this week.
//!
//! So `note` rows record coverage (what the vault already contains, keyed by
//! path) and `source` rows record consumption (external things a mission
//! read, keyed by natural id). They answer different questions and a
//! continuous mission needs both: "have I already written about this topic?"
//! and "have I already read this paper?".
use sha2::{Digest, Sha256};
use uuid::Uuid;
/// A note parsed out of the vault, ready to be indexed.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ParsedNote {
/// Vault-relative path, used as identity for `kind = 'note'`.
pub path: String,
pub title: Option<String>,
pub content_hash: String,
/// An external identity the note declares for itself, if any. Nothing in
/// the vault does this today; missions writing new notes are expected to.
pub declared_source_id: Option<String>,
}
impl ParsedNote {
/// `note:<path>` — the `source_id` this note occupies in the index.
pub fn source_id(&self) -> String {
format!("note:{}", self.path)
}
}
/// Hash content for change detection. Not a dedupe key — identity is
/// `source_id`; this only distinguishes "unchanged" from "edited".
pub fn content_hash(body: &str) -> String {
let mut h = Sha256::new();
h.update(body.as_bytes());
format!("{:x}", h.finalize())
}
/// Split YAML frontmatter from the body.
///
/// Returns `(frontmatter, body)`. A note without frontmatter — 271 of the 416
/// in the real vault — yields `("", whole file)` rather than being skipped.
/// Skipping them would drop two thirds of the corpus on the floor.
fn split_frontmatter(text: &str) -> (&str, &str) {
let Some(rest) = text.strip_prefix("---") else {
return ("", text);
};
let rest = rest.strip_prefix('\n').unwrap_or(rest);
match rest.find("\n---") {
Some(end) => {
let body = &rest[end + 4..];
(&rest[..end], body.strip_prefix('\n').unwrap_or(body))
}
// An opening fence with no close is malformed; treat the whole file as
// body rather than swallowing it as frontmatter.
None => ("", text),
}
}
/// Read one scalar key out of a frontmatter block.
///
/// Deliberately not a YAML parser. The vault's frontmatter is flat
/// `key: value` with occasional quotes and one list (`tags`), and pulling in a
/// YAML dependency to read three keys would be more surface than it is worth.
fn frontmatter_value<'a>(fm: &'a str, key: &str) -> Option<&'a str> {
for line in fm.lines() {
let line = line.trim();
let Some((k, v)) = line.split_once(':') else {
continue;
};
if !k.trim().eq_ignore_ascii_case(key) {
continue;
}
let v = v.trim().trim_matches('"').trim_matches('\'').trim();
if !v.is_empty() {
return Some(v);
}
}
None
}
/// Which frontmatter keys may declare an external identity, in priority order.
///
/// None of these appear in the vault today. They are the contract for notes
/// that missions write from here on, and the reason a `source:` key is NOT in
/// the list: the vault already uses `source:` for local filesystem paths of
/// course material (`/Users/quantum/Downloads/...`), which is provenance, not
/// a citable external identity. Treating it as one would fill the seen-set
/// with 25 rows keyed on a laptop path.
const IDENTITY_KEYS: &[&str] = &["source_id", "arxiv", "doi", "url", "permalink"];
/// Parse a note. `path` must be vault-relative.
pub fn parse_note(path: &str, text: &str) -> ParsedNote {
let (fm, body) = split_frontmatter(text);
let declared_source_id = IDENTITY_KEYS.iter().find_map(|k| {
frontmatter_value(fm, k).map(|v| {
// `source_id` is already qualified; the others name their scheme.
if *k == "source_id" || v.contains(':') {
v.to_string()
} else {
format!("{k}:{v}")
}
})
});
// Title: the first markdown H1, else the filename stem. Frontmatter has no
// consistent title key in this vault.
let title = body
.lines()
.find_map(|l| l.strip_prefix("# ").map(str::trim))
.filter(|t| !t.is_empty())
.map(str::to_string)
.or_else(|| {
std::path::Path::new(path)
.file_stem()
.map(|s| s.to_string_lossy().into_owned())
});
ParsedNote {
path: path.to_string(),
title,
// Hash the body, not the whole file: re-syncing a repo note rewrites
// `updated:`/`size_kb:` in frontmatter without the prose changing, and
// that should not read as an edit.
content_hash: content_hash(body),
declared_source_id,
}
}
/// What a re-index actually did. `unchanged` is the number that matters: on a
/// vault nobody edited it should equal the note count.
#[derive(Debug, Default, Clone, PartialEq, Eq)]
pub struct IndexStats {
pub scanned: usize,
pub inserted: usize,
pub updated: usize,
pub unchanged: usize,
}
/// Walk a checkout and index every markdown note.
///
/// Skips `.git` and Obsidian's own `.obsidian` config directory — indexing an
/// editor's workspace state as knowledge would be noise.
pub fn collect_notes(root: &std::path::Path) -> Vec<ParsedNote> {
fn walk(dir: &std::path::Path, root: &std::path::Path, out: &mut Vec<ParsedNote>) {
let Ok(entries) = std::fs::read_dir(dir) else {
return;
};
for entry in entries.flatten() {
let path = entry.path();
let name = entry.file_name();
let name = name.to_string_lossy();
if name.starts_with('.') {
continue;
}
if path.is_dir() {
walk(&path, root, out);
} else if path.extension().and_then(|e| e.to_str()) == Some("md") {
let Ok(text) = std::fs::read_to_string(&path) else {
continue;
};
let rel = path
.strip_prefix(root)
.unwrap_or(&path)
.to_string_lossy()
.into_owned();
out.push(parse_note(&rel, &text));
}
}
}
let mut out = Vec::new();
walk(root, root, &mut out);
out.sort_by(|a, b| a.path.cmp(&b.path));
out
}
/// Upsert one item. Returns whether the row was new.
#[allow(clippy::too_many_arguments)]
pub async fn record(
pool: &sqlx::PgPool,
workspace_id: Uuid,
corpus_id: &str,
kind: &str,
source_id: &str,
title: Option<&str>,
path: Option<&str>,
url: Option<&str>,
content_hash: &str,
mission_id: Option<Uuid>,
) -> Result<bool, String> {
// `last_seen_at` always moves; `first_seen_at` and `mission_id` never do.
// The first mission to find a source keeps the credit, which is what makes
// "did THIS run contribute anything new" answerable.
let row: (bool,) = sqlx::query_as(
"INSERT INTO corpus_items
(id, workspace_id, corpus_id, kind, source_id, title, path, url,
content_hash, mission_id)
VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10)
ON CONFLICT (workspace_id, corpus_id, source_id) DO UPDATE
SET last_seen_at = now(),
title = COALESCE(EXCLUDED.title, corpus_items.title),
path = COALESCE(EXCLUDED.path, corpus_items.path),
url = COALESCE(EXCLUDED.url, corpus_items.url),
content_hash = EXCLUDED.content_hash
RETURNING (xmax = 0) AS inserted",
)
.bind(Uuid::now_v7())
.bind(workspace_id)
.bind(corpus_id)
.bind(kind)
.bind(source_id)
.bind(title)
.bind(path)
.bind(url)
.bind(content_hash)
.bind(mission_id)
.fetch_one(pool)
.await
.map_err(|e| format!("record corpus item {source_id}: {e}"))?;
Ok(row.0)
}
/// Has this corpus already seen this `source_id`?
pub async fn seen(
pool: &sqlx::PgPool,
workspace_id: Uuid,
corpus_id: &str,
source_id: &str,
) -> Result<bool, String> {
// `SELECT 1` is INT4; binding it as i64 fails to decode.
let row: Option<(i32,)> = sqlx::query_as(
"SELECT 1 FROM corpus_items
WHERE workspace_id = $1 AND corpus_id = $2 AND source_id = $3",
)
.bind(workspace_id)
.bind(corpus_id)
.bind(source_id)
.fetch_optional(pool)
.await
.map_err(|e| format!("seen({source_id}): {e}"))?;
Ok(row.is_some())
}
/// Of these candidate ids, which has this corpus NOT seen?
///
/// The shape a research agent actually needs: it has ten search hits and wants
/// to know which are worth fetching. One round trip, not ten.
pub async fn unseen(
pool: &sqlx::PgPool,
workspace_id: Uuid,
corpus_id: &str,
candidates: &[String],
) -> Result<Vec<String>, String> {
if candidates.is_empty() {
return Ok(Vec::new());
}
let rows: Vec<(String,)> = sqlx::query_as(
"SELECT source_id FROM corpus_items
WHERE workspace_id = $1 AND corpus_id = $2 AND source_id = ANY($3)",
)
.bind(workspace_id)
.bind(corpus_id)
.bind(candidates)
.fetch_all(pool)
.await
.map_err(|e| format!("unseen: {e}"))?;
let known: std::collections::HashSet<String> = rows.into_iter().map(|r| r.0).collect();
Ok(candidates
.iter()
.filter(|c| !known.contains(*c))
.cloned()
.collect())
}
/// Index every note in a checkout. Idempotent by construction.
pub async fn index_vault(
pool: &sqlx::PgPool,
workspace_id: Uuid,
corpus_id: &str,
root: &std::path::Path,
) -> Result<IndexStats, String> {
let notes = collect_notes(root);
let mut stats = IndexStats {
scanned: notes.len(),
..Default::default()
};
for note in &notes {
let existing: Option<(String,)> = sqlx::query_as(
"SELECT content_hash FROM corpus_items
WHERE workspace_id = $1 AND corpus_id = $2 AND source_id = $3",
)
.bind(workspace_id)
.bind(corpus_id)
.bind(note.source_id())
.fetch_optional(pool)
.await
.map_err(|e| format!("lookup {}: {e}", note.path))?;
match existing {
Some((hash,)) if hash == note.content_hash => {
stats.unchanged += 1;
continue;
}
Some(_) => stats.updated += 1,
None => stats.inserted += 1,
}
record(
pool,
workspace_id,
corpus_id,
"note",
&note.source_id(),
note.title.as_deref(),
Some(&note.path),
None,
&note.content_hash,
None,
)
.await?;
// A note that declares an external identity also registers as a
// consumed source, so a later mission does not re-read what an
// earlier one already wrote up.
if let Some(sid) = &note.declared_source_id {
record(
pool,
workspace_id,
corpus_id,
"source",
sid,
note.title.as_deref(),
Some(&note.path),
None,
&note.content_hash,
None,
)
.await?;
}
}
Ok(stats)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn frontmatter_is_split_from_body() {
let (fm, body) = split_frontmatter("---\ntype: lecture\n---\n# Title\n\ntext\n");
assert_eq!(fm, "type: lecture");
assert!(body.starts_with("# Title"));
}
/// 271 of the vault's 416 notes have no frontmatter. Dropping them would
/// discard two thirds of the corpus.
#[test]
fn a_note_without_frontmatter_is_still_a_note() {
let (fm, body) = split_frontmatter("# Plain\n\nno frontmatter here\n");
assert_eq!(fm, "");
assert!(body.starts_with("# Plain"));
let n = parse_note("Daily/x.md", "# Plain\n\nbody\n");
assert_eq!(n.title.as_deref(), Some("Plain"));
assert_eq!(n.declared_source_id, None);
}
/// An unterminated fence must not swallow the file.
#[test]
fn malformed_frontmatter_is_treated_as_body() {
let (fm, body) = split_frontmatter("---\nbroken: yes\nno closing fence\n");
assert_eq!(fm, "");
assert!(body.contains("no closing fence"));
}
/// The vault's real `source:` values are local filesystem paths of course
/// material. Treating those as citable identity would fill the seen-set
/// with 25 rows keyed on a laptop path.
#[test]
fn a_local_source_path_is_not_an_external_identity() {
let note = parse_note(
"50 APESS 2026/Lectures/talk.md",
"---\nsource: \"/Users/quantum/Downloads/Material_APESS_2026/x.pdf\"\n\
date: 2026-07-27\ntype: lecture\n---\n# Agentic Design\n",
);
assert_eq!(
note.declared_source_id, None,
"a Downloads path is provenance, not a citable source id"
);
assert_eq!(note.title.as_deref(), Some("Agentic Design"));
assert_eq!(note.source_id(), "note:50 APESS 2026/Lectures/talk.md");
}
#[test]
fn declared_identities_are_scheme_qualified() {
let a = parse_note("p.md", "---\narxiv: 2401.12345\n---\n# T\n");
assert_eq!(a.declared_source_id.as_deref(), Some("arxiv:2401.12345"));
let d = parse_note("p.md", "---\ndoi: 10.1000/xyz\n---\n# T\n");
assert_eq!(d.declared_source_id.as_deref(), Some("doi:10.1000/xyz"));
// Already-qualified values are not double-prefixed.
let s = parse_note("p.md", "---\nsource_id: arxiv:2401.99999\n---\n# T\n");
assert_eq!(s.declared_source_id.as_deref(), Some("arxiv:2401.99999"));
// A URL carries its own scheme and must not become `url:https:...`.
let u = parse_note("p.md", "---\nurl: https://example.com/p\n---\n# T\n");
assert_eq!(
u.declared_source_id.as_deref(),
Some("https://example.com/p")
);
}
/// Repo-sync notes rewrite `updated:`/`size_kb:` on every sync without the
/// prose changing. Hashing the whole file would report 103 phantom edits
/// per run and make "unchanged" meaningless.
#[test]
fn frontmatter_churn_does_not_count_as_an_edit() {
let a = parse_note("Repos/x.md", "---\nupdated: 2026-08-01\nsize_kb: 12\n---\n# X\n\nbody\n");
let b = parse_note("Repos/x.md", "---\nupdated: 2026-08-03\nsize_kb: 14\n---\n# X\n\nbody\n");
assert_eq!(a.content_hash, b.content_hash);
let c = parse_note("Repos/x.md", "---\nupdated: 2026-08-03\n---\n# X\n\nDIFFERENT\n");
assert_ne!(a.content_hash, c.content_hash, "real edits must be visible");
}
#[test]
fn note_identity_is_its_path() {
let n = parse_note("30 Resources/a b.md", "# A\n");
assert_eq!(n.source_id(), "note:30 Resources/a b.md");
}
#[test]
fn collect_skips_dotfiles_and_non_markdown() {
let tmp = tempfile::tempdir().unwrap();
let root = tmp.path();
std::fs::create_dir_all(root.join(".obsidian")).unwrap();
std::fs::create_dir_all(root.join("Daily")).unwrap();
std::fs::write(root.join(".obsidian/workspace.md"), "# editor state\n").unwrap();
std::fs::write(root.join("Daily/note.md"), "# Real\n").unwrap();
std::fs::write(root.join("image.png"), "notmd").unwrap();
let notes = collect_notes(root);
assert_eq!(notes.len(), 1, "only the real note: {notes:?}");
assert_eq!(notes[0].path, "Daily/note.md");
}
}
+227
View File
@@ -0,0 +1,227 @@
//! One run of the library: find, skip what we have, shelve the rest.
//!
//! This is the piece that makes the others a *job* rather than parts on a
//! bench. Order matters and it is deliberate:
//!
//! 1. **search** arXiv for candidates
//! 2. **skip** everything already on the checkmark list — before any download
//! 3. **fetch** the PDF for what is left, and verify it really is a PDF
//! 4. **shelve** it in the blob store
//! 5. **catalogue** it: write the vault note
//! 6. **check it off** so next week skips it
//!
//! Step 2 comes before step 3 on purpose. Checking after downloading would
//! still dedupe the catalogue, but it would re-download every paper we already
//! have, every week, forever — and the whole point of the checkmark list is to
//! not do the work twice.
//!
//! # Nothing new is a success, not a failure
//!
//! A weekly run that finds no new papers has worked correctly. A run that
//! *crashed* has not. [`Harvest`] keeps those apart, because collapsing them
//! is precisely the "reported success while doing nothing" shape that this
//! codebase has been bitten by repeatedly. `shelved == 0` with `failed.empty()`
//! is a quiet week; `shelved == 0` with failures is a broken run.
use std::path::Path;
use std::sync::Arc;
use uuid::Uuid;
use crate::corpus;
use crate::papers::{self, Paper};
/// What one run did. Every number here is observed, not claimed.
#[derive(Debug, Default, Clone)]
pub struct Harvest {
/// Papers the search returned.
pub candidates: usize,
/// Of those, how many were already on the checkmark list.
pub already_had: usize,
/// Successfully downloaded, shelved and catalogued.
pub shelved: Vec<String>,
/// `(source_id, why)` for each paper that could not be shelved.
pub failed: Vec<(String, String)>,
/// Vault-relative paths of the notes written.
pub notes_written: Vec<String>,
}
impl Harvest {
/// Did this run add anything? The verification predicate for a continuous
/// research mission: a run that contributes no new source has produced
/// nothing, whatever its transcript says.
pub fn added_anything(&self) -> bool {
!self.shelved.is_empty()
}
/// A run is healthy if nothing errored — including a run that found
/// nothing new, which is the normal state of a mature library.
pub fn healthy(&self) -> bool {
self.failed.is_empty()
}
pub fn summary(&self) -> String {
format!(
"{} candidates, {} already held, {} shelved, {} failed",
self.candidates,
self.already_had,
self.shelved.len(),
self.failed.len()
)
}
}
/// Where a library lives: its records, its shelf, and its catalogue.
///
/// Grouped rather than passed as loose arguments because these five always
/// travel together and always describe one library — splitting them at a call
/// site is how a run ends up shelving into one place and cataloguing into
/// another.
pub struct Library<'a> {
pub pool: &'a sqlx::PgPool,
/// The shelf: where PDFs are stored.
pub blobs: &'a Arc<dyn cm_files::BlobStore>,
pub workspace_id: Uuid,
/// Which checkmark list, e.g. `"valhalla-vault"`.
pub corpus_id: &'a str,
/// Checkout the catalogue notes are written into.
pub vault_root: &'a Path,
}
/// Shelve a specific set of papers. Split from [`run`] so the skip/shelve
/// logic is testable without reaching arXiv.
pub async fn shelve(
lib: &Library<'_>,
candidates: &[Paper],
mission_id: Option<Uuid>,
) -> Result<Harvest, String> {
let Library { pool, blobs, workspace_id, corpus_id, vault_root } = *lib;
let mut out = Harvest {
candidates: candidates.len(),
..Default::default()
};
// One round trip for the whole batch rather than one query per paper.
let ids: Vec<String> = candidates.iter().map(Paper::source_id).collect();
let fresh: std::collections::HashSet<String> =
corpus::unseen(pool, workspace_id, corpus_id, &ids)
.await?
.into_iter()
.collect();
out.already_had = candidates.len() - fresh.len();
for paper in candidates {
let sid = paper.source_id();
if !fresh.contains(&sid) {
continue;
}
// Fetch first. If the PDF cannot be had, nothing is recorded — the
// paper stays unseen so a later run retries it, rather than being
// checked off with an empty shelf slot behind it.
let bytes = match papers::fetch_pdf(paper).await {
Ok(b) => b,
Err(e) => {
out.failed.push((sid, e));
continue;
}
};
let key = paper.blob_key();
if let Err(e) = blobs.put(&key, &bytes).await {
out.failed.push((sid, format!("shelve {key}: {e}")));
continue;
}
// Catalogue note next to the shelf. Written into the vault checkout;
// committing and pushing it is the caller's job, through the delivery
// path that already exists.
let note = papers::catalogue_note(paper, &key);
let note_path = vault_root.join(paper.note_path());
if let Some(parent) = note_path.parent() {
if let Err(e) = std::fs::create_dir_all(parent) {
out.failed.push((sid, format!("create {}: {e}", parent.display())));
continue;
}
}
if let Err(e) = std::fs::write(&note_path, &note) {
out.failed
.push((sid, format!("write {}: {e}", note_path.display())));
continue;
}
// Check it off LAST. If anything above failed we did not get the
// paper, and marking it seen would mean never trying again.
corpus::record(
pool,
workspace_id,
corpus_id,
"source",
&sid,
Some(&paper.title),
Some(&paper.note_path()),
Some(&format!("https://arxiv.org/abs/{}", paper.arxiv_id)),
&corpus::content_hash(&note),
mission_id,
)
.await?;
out.notes_written.push(paper.note_path());
out.shelved.push(sid);
}
Ok(out)
}
/// A full run: search arXiv, then shelve whatever is new.
pub async fn run(
lib: &Library<'_>,
query: &str,
limit: usize,
mission_id: Option<Uuid>,
) -> Result<Harvest, String> {
let candidates = papers::search(query, limit).await?;
let harvest = shelve(lib, &candidates, mission_id).await?;
let corpus_id = lib.corpus_id;
eprintln!("harvest[{corpus_id}] query={query:?} → {}", harvest.summary());
for (sid, why) in &harvest.failed {
eprintln!("harvest[{corpus_id}] FAILED {sid}: {why}");
}
Ok(harvest)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn a_quiet_week_is_healthy_but_adds_nothing() {
let quiet = Harvest {
candidates: 5,
already_had: 5,
..Default::default()
};
assert!(quiet.healthy(), "finding nothing new is not an error");
assert!(
!quiet.added_anything(),
"but it must not count as having produced something"
);
let broken = Harvest {
candidates: 5,
already_had: 0,
failed: vec![("arxiv:1".into(), "timeout".into())],
..Default::default()
};
assert!(!broken.healthy());
assert!(!broken.added_anything());
let good = Harvest {
candidates: 5,
already_had: 4,
shelved: vec!["arxiv:2".into()],
..Default::default()
};
assert!(good.healthy() && good.added_anything());
}
}
+16
View File
@@ -16,7 +16,11 @@ mod mcp_door;
mod mcp_skills; mod mcp_skills;
pub mod mission_orchestrator; pub mod mission_orchestrator;
pub mod mission_refiner; pub mod mission_refiner;
pub mod corpus;
pub mod harvest;
pub mod library;
pub mod mission_delivery; pub mod mission_delivery;
pub mod papers;
pub mod phase_config; pub mod phase_config;
pub mod runtime_preflight; pub mod runtime_preflight;
pub mod mission_runtime; pub mod mission_runtime;
@@ -64,6 +68,9 @@ pub struct AppState {
pub file_root: Option<std::path::PathBuf>, pub file_root: Option<std::path::PathBuf>,
/// Live control channels to connected fleet-node daemons. /// Live control channels to connected fleet-node daemons.
pub node_hub: std::sync::Arc<fleet::NodeHub>, pub node_hub: std::sync::Arc<fleet::NodeHub>,
/// The shelf. Present once the server wires storage; `None` in the
/// bare-`new` path used by tests that never touch blobs.
pub blobs: Option<std::sync::Arc<dyn cm_files::BlobStore>>,
} }
impl AppState { impl AppState {
@@ -78,6 +85,7 @@ impl AppState {
billing: cm_config::BillingConfig::default(), billing: cm_config::BillingConfig::default(),
file_root: None, file_root: None,
node_hub: std::sync::Arc::new(fleet::NodeHub::new()), node_hub: std::sync::Arc::new(fleet::NodeHub::new()),
blobs: None,
} }
} }
@@ -86,6 +94,12 @@ impl AppState {
self self
} }
/// The shelf — where the paper library stores PDFs.
pub fn with_blobs(mut self, blobs: std::sync::Arc<dyn cm_files::BlobStore>) -> AppState {
self.blobs = Some(blobs);
self
}
pub fn with_oauth(mut self, oauth: cm_config::OAuthConfig) -> AppState { pub fn with_oauth(mut self, oauth: cm_config::OAuthConfig) -> AppState {
self.oauth = oauth; self.oauth = oauth;
self self
@@ -311,6 +325,8 @@ pub fn router(state: AppState) -> Router {
.route("/api/sessions", post(routes::sessions::create)) .route("/api/sessions", post(routes::sessions::create))
.route("/api/sessions/history", get(routes::sessions::history)) .route("/api/sessions/history", get(routes::sessions::history))
.route("/api/gateway", post(routes::gateway::gateway)) .route("/api/gateway", post(routes::gateway::gateway))
.route("/api/library/runs", post(routes::library::run))
.route("/api/library/items", get(routes::library::list))
.route("/api/routines", get(routes::routines::list)) .route("/api/routines", get(routes::routines::list))
.route("/api/routines", post(routes::routines::create)) .route("/api/routines", post(routes::routines::create))
.route("/api/routines/runs", get(routes::routines::runs)) .route("/api/routines/runs", get(routes::routines::runs))
+263
View File
@@ -0,0 +1,263 @@
//! A library run end to end: clone the vault, harvest, push the catalogue.
//!
//! [`harvest`](crate::harvest) writes catalogue notes into a directory. This
//! puts that directory somewhere real: a checkout of the vault repo, with the
//! new notes committed and pushed.
//!
//! # Never `main`
//!
//! The vault is a live Obsidian vault that a human edits and syncs. Pushing
//! straight to `main` races that sync and can lose hand-written work. Every
//! run lands on its own branch, exactly like the mission delivery path that
//! was validated 20/20 earlier — a human merges when they have looked at it.
//!
//! # The PDFs do not go here
//!
//! Only notes are committed. PDFs are shelved in the blob store, because a
//! few hundred papers is gigabytes and a vault that size is painful to clone
//! and slow to open. The note carries the blob key, so the catalogue always
//! knows where its shelf is.
use std::path::{Path, PathBuf};
use std::sync::Arc;
use uuid::Uuid;
use crate::harvest::{self, Harvest, Library};
use crate::mission_workspace;
/// What a full run produced, including whether it reached the forge.
#[derive(Debug, Clone)]
pub struct LibraryRun {
pub harvest: Harvest,
pub branch: String,
/// `true` only when the push was observed to succeed. A run that shelved
/// papers but could not push still has the PDFs and the checkmarks; the
/// notes are simply not on the forge yet.
pub pushed: bool,
pub error: Option<String>,
}
fn git_identity() -> [(&'static str, String); 4] {
let (name, email) = crate::mission_delivery::commit_identity();
[
("GIT_AUTHOR_NAME", name.clone()),
("GIT_AUTHOR_EMAIL", email.clone()),
("GIT_COMMITTER_NAME", name),
("GIT_COMMITTER_EMAIL", email),
]
}
async fn git(repo: &Path, args: &[&str]) -> Result<String, String> {
let mut cmd = tokio::process::Command::new("git");
cmd.arg("-C").arg(repo);
cmd.args(["-c", &format!("safe.directory={}", repo.display())]);
cmd.args(args);
for (k, v) in git_identity() {
cmd.env(k, v);
}
let out = cmd.output().await.map_err(|e| format!("spawn git: {e}"))?;
if !out.status.success() {
return Err(format!(
"git {} → {}: {}",
args.first().copied().unwrap_or("?"),
out.status,
mission_workspace::redact_token(&String::from_utf8_lossy(&out.stderr))
.chars()
.take(300)
.collect::<String>()
));
}
Ok(String::from_utf8_lossy(&out.stdout).into_owned())
}
/// Clone the vault fresh into `work_root`, returning the checkout path.
///
/// Fresh each run rather than reused: a library run is short, the vault is
/// small (measured 6.9 MB / 416 notes), and a stale checkout is how the
/// mission path lost work three times this week.
pub async fn clone_vault(clone_url: &str, work_root: &Path) -> Result<PathBuf, String> {
let path = work_root.join("vault");
if path.exists() {
tokio::fs::remove_dir_all(&path)
.await
.map_err(|e| format!("clear {}: {e}", path.display()))?;
}
tokio::fs::create_dir_all(work_root)
.await
.map_err(|e| format!("mkdir {}: {e}", work_root.display()))?;
let auth = mission_workspace::with_ambient_auth(clone_url);
let out = tokio::process::Command::new("git")
.args(["clone", "--quiet", "--depth", "1", &auth])
.arg(&path)
.output()
.await
.map_err(|e| format!("spawn git clone: {e}"))?;
if !out.status.success() {
return Err(format!(
"clone vault → {}: {}",
out.status,
mission_workspace::redact_token(&String::from_utf8_lossy(&out.stderr))
.chars()
.take(300)
.collect::<String>()
));
}
// The token must not stay in .git/config: the checkout may be handed to a
// container later, and a credential in a file an agent can read is a
// credential an agent has.
mission_workspace::scrub_remote_credentials(&path, &auth);
Ok(path)
}
/// One complete library run.
#[allow(clippy::too_many_arguments)]
pub async fn run_to_vault(
pool: &sqlx::PgPool,
blobs: &Arc<dyn cm_files::BlobStore>,
workspace_id: Uuid,
corpus_id: &str,
clone_url: &str,
work_root: &Path,
queries: &[String],
per_query: usize,
mission_id: Option<Uuid>,
) -> Result<LibraryRun, String> {
let vault = clone_vault(clone_url, work_root).await?;
let lib = Library {
pool,
blobs,
workspace_id,
corpus_id,
vault_root: &vault,
};
// Accumulate across queries. Topics overlap — "agentic topology" and
// "multi-agent orchestration" return some of the same papers — and the
// checkmark list dedupes across them within a single run as well as
// between runs, because each shelve records before the next query starts.
let mut total = Harvest::default();
for q in queries {
let h = harvest::run(&lib, q, per_query, mission_id).await?;
total.candidates += h.candidates;
total.already_had += h.already_had;
total.shelved.extend(h.shelved);
total.failed.extend(h.failed);
total.notes_written.extend(h.notes_written);
}
// The TAIL of the uuid, not the head. UUIDv7 leads with a 48-bit
// timestamp, so two ids minted in the same millisecond share their first
// 12 hex characters exactly — the branch-name collision that hit mission
// 019fc42b earlier. The tail is the random part.
let branch = format!("clawmates/library-{}", branch_suffix(Uuid::now_v7()));
if total.notes_written.is_empty() {
// A quiet run is a success with nothing to push. Creating an empty
// branch every week would be noise.
return Ok(LibraryRun {
harvest: total,
branch,
pushed: false,
error: None,
});
}
git(&vault, &["checkout", "-B", &branch]).await?;
git(&vault, &["add", "--", "60 Papers"]).await?;
let message = format!(
"library: {} new paper(s)\n\n{}\n\nShelved in the blob store; this commit is the catalogue.",
total.shelved.len(),
total
.shelved
.iter()
.map(|s| format!("- {s}"))
.collect::<Vec<_>>()
.join("\n")
);
git(&vault, &["commit", "--no-verify", "-m", &message]).await?;
let auth = mission_workspace::with_ambient_auth(clone_url);
let refspec = format!("HEAD:refs/heads/{branch}");
match git(&vault, &["push", &auth, &refspec]).await {
Ok(_) => Ok(LibraryRun {
harvest: total,
branch,
pushed: true,
error: None,
}),
Err(e) => Ok(LibraryRun {
harvest: total,
branch,
pushed: false,
error: Some(e),
}),
}
}
/// Distinct-per-run branch suffix. See the note at the call site: taking the
/// head of a UUIDv7 yields the timestamp, which collides.
fn branch_suffix(id: Uuid) -> String {
let s = id.simple().to_string();
s[s.len() - 12..].to_string()
}
/// The topics this library currently tracks.
///
/// Drawn from what the project is actually working on: `papers/dynamic-
/// agentic-topologies.md` (topology search and evolution, citing ADAS,
/// Darwin-Gödel and SwarmAgentic), plus the problems this week's work ran
/// into — verifying what an agent actually did, and giving a long-running
/// agent memory of what it has already covered.
pub fn default_topics() -> Vec<String> {
[
"all:\"agentic topology\" OR all:\"multi-agent topology\"",
"all:\"multi-agent orchestration\" AND all:LLM",
"all:\"agent memory\" AND all:\"long-term\"",
"all:\"LLM agent\" AND all:verification",
"all:\"prompt injection\" AND all:agent",
]
.iter()
.map(|s| s.to_string())
.collect()
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn topics_are_non_empty_and_arxiv_shaped() {
let topics = default_topics();
assert!(topics.len() >= 3);
for t in &topics {
assert!(t.contains("all:"), "arXiv field prefix missing in {t:?}");
assert!(!t.trim().is_empty());
}
}
/// Two runs in the same millisecond must not collide.
///
/// This caught a real repeat of the mission-path bug (019fc42b): UUIDv7
/// leads with a 48-bit timestamp, so the FIRST 12 hex characters of two
/// ids minted together are identical. Taking the tail fixes it. Looping
/// rather than sampling twice, because a one-shot check passes by luck
/// whenever the millisecond happens to tick between the two calls.
#[test]
fn every_run_gets_a_distinct_branch() {
let ids: Vec<String> = (0..100).map(|_| branch_suffix(Uuid::now_v7())).collect();
let unique: std::collections::HashSet<&String> = ids.iter().collect();
assert_eq!(unique.len(), ids.len(), "branch suffixes collided: {ids:?}");
// And the head-based scheme really does collide, so this test has teeth.
let heads: Vec<String> = (0..100)
.map(|_| Uuid::now_v7().simple().to_string()[..12].to_string())
.collect();
let head_unique: std::collections::HashSet<&String> = heads.iter().collect();
assert!(
head_unique.len() < heads.len(),
"the head of a UUIDv7 was expected to collide but did not"
);
}
}
+1 -1
View File
@@ -83,7 +83,7 @@ const MAX_PATCH_BYTES: usize = 4 * 1024 * 1024;
const DEFAULT_COMMIT_NAME: &str = "Omar Sobh"; const DEFAULT_COMMIT_NAME: &str = "Omar Sobh";
const DEFAULT_COMMIT_EMAIL: &str = "[email protected]"; const DEFAULT_COMMIT_EMAIL: &str = "[email protected]";
fn commit_identity() -> (String, String) { pub(crate) fn commit_identity() -> (String, String) {
let name = std::env::var("CLAWMATES_COMMIT_NAME") let name = std::env::var("CLAWMATES_COMMIT_NAME")
.ok() .ok()
.filter(|v| !v.trim().is_empty()) .filter(|v| !v.trim().is_empty())
+2 -2
View File
@@ -391,7 +391,7 @@ pub(crate) fn base_commit(path: &std::path::Path) -> Option<String> {
/// Best-effort and non-fatal: a checkout that keeps its token still works, and /// Best-effort and non-fatal: a checkout that keeps its token still works, and
/// failing the mission over it would trade a real capability for a marginal /// failing the mission over it would trade a real capability for a marginal
/// improvement in a situation we have already logged. /// improvement in a situation we have already logged.
fn scrub_remote_credentials(path: &std::path::Path, original_url: &str) { pub(crate) fn scrub_remote_credentials(path: &std::path::Path, original_url: &str) {
if !original_url.contains('@') && !original_url.contains("oauth2:") { if !original_url.contains('@') && !original_url.contains("oauth2:") {
// Nothing was injected (SSH remote, or no token configured). // Nothing was injected (SSH remote, or no token configured).
return; return;
@@ -495,7 +495,7 @@ fn ignore_agent_scaffolding(path: &std::path::Path) {
} }
} }
fn redact_token(s: &str) -> String { pub(crate) fn redact_token(s: &str) -> String {
// Strip any "oauth2:<token>@" segment that git may echo back on // Strip any "oauth2:<token>@" segment that git may echo back on
// failures. Belt-and-braces: also nuke any raw token env value. // failures. Belt-and-braces: also nuke any raw token env value.
let mut out = s.to_string(); let mut out = s.to_string();
+345
View File
@@ -0,0 +1,345 @@
//! Finding papers, shelving them, and cataloguing them.
//!
//! The library has three parts and it matters which is which:
//!
//! - **arXiv** is where papers are *found*.
//! - **The blob store** is the *shelf* — the PDF itself lives there.
//! - **The vault** is the *card catalogue* — a markdown note per paper, with
//! the metadata and a pointer to the shelf.
//!
//! Plus [`crate::corpus`], which is the list of checkmarks: it is what stops
//! the same paper being fetched twice across weekly runs. That list is the
//! reason this can be a *continuous* job rather than one that redoes itself
//! forever — the failure that killed the previous attempt at this (migrations
//! 0030-0044, dropped in 0053).
//!
//! # The contract that ties it together
//!
//! Every note this module writes carries `source_id: arxiv:NNNN.NNNNN` in its
//! frontmatter. `corpus::parse_note` reads exactly that key, so re-indexing
//! the vault re-derives the checkmark list from the notes themselves. The
//! catalogue is authoritative; the index is rebuildable from it. If the
//! database were lost, a re-index of the vault would restore what we have.
use serde::{Deserialize, Serialize};
/// One paper as arXiv describes it.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct Paper {
/// Bare arXiv id, e.g. `2401.12345` — no version suffix.
pub arxiv_id: String,
pub title: String,
pub authors: Vec<String>,
pub summary: String,
pub published: String,
pub pdf_url: String,
}
impl Paper {
/// The checkmark key. Version suffixes are stripped upstream so `v1` and
/// `v2` of the same paper are one entry, not two.
pub fn source_id(&self) -> String {
format!("arxiv:{}", self.arxiv_id)
}
/// Where the PDF is shelved in the blob store.
pub fn blob_key(&self) -> String {
format!("papers/arxiv/{}.pdf", self.arxiv_id)
}
/// Where the catalogue note goes in the vault.
///
/// Under a dedicated folder so the library never collides with the
/// hand-written parts of the vault (`30 Resources`, `40 Projects`, and so
/// on). A human should always be able to tell which notes a machine wrote.
pub fn note_path(&self) -> String {
format!("60 Papers/arxiv-{}.md", self.arxiv_id)
}
}
/// Strip an arXiv version suffix: `2401.12345v3` -> `2401.12345`.
///
/// Without this a weekly job re-downloads a paper every time the authors post
/// a revision, and the checkmark list quietly fills with near-duplicates.
pub fn normalize_arxiv_id(raw: &str) -> String {
let id = raw.rsplit('/').next().unwrap_or(raw);
match id.find('v') {
// Only a trailing `vN` counts; the `v` in a word must not truncate.
Some(i) if id[i + 1..].chars().all(|c| c.is_ascii_digit()) && i + 1 < id.len() => {
id[..i].to_string()
}
_ => id.to_string(),
}
}
/// Parse arXiv's Atom feed.
///
/// Hand-rolled rather than pulling an XML crate: the feed is a fixed, simple
/// shape and this reads five fields from it. If arXiv's format ever drifts,
/// `entries_are_parsed_from_a_real_feed` fails loudly rather than silently
/// returning zero papers — which is the failure mode that matters, because a
/// search returning nothing looks exactly like "no new papers this week".
pub fn parse_atom(xml: &str) -> Vec<Paper> {
let mut out = Vec::new();
for chunk in xml.split("<entry>").skip(1) {
let entry = chunk.split("</entry>").next().unwrap_or(chunk);
let field = |tag: &str| -> Option<String> {
let open = format!("<{tag}>");
let close = format!("</{tag}>");
let start = entry.find(&open)? + open.len();
let end = entry[start..].find(&close)? + start;
Some(unescape(entry[start..end].trim()))
};
let Some(raw_id) = field("id") else { continue };
let arxiv_id = normalize_arxiv_id(&raw_id);
if arxiv_id.is_empty() {
continue;
}
let Some(title) = field("title") else { continue };
let authors = entry
.split("<author>")
.skip(1)
.filter_map(|a| {
let start = a.find("<name>")? + 6;
let end = a[start..].find("</name>")? + start;
Some(unescape(a[start..end].trim()))
})
.collect();
// The PDF link is an attribute, not an element.
let pdf_url = entry
.split("<link")
.find(|l| l.contains("title=\"pdf\""))
.and_then(|l| {
let start = l.find("href=\"")? + 6;
let end = l[start..].find('"')? + start;
Some(l[start..end].to_string())
})
.unwrap_or_else(|| format!("https://arxiv.org/pdf/{arxiv_id}"));
out.push(Paper {
title: title.split_whitespace().collect::<Vec<_>>().join(" "),
summary: field("summary")
.unwrap_or_default()
.split_whitespace()
.collect::<Vec<_>>()
.join(" "),
published: field("published").unwrap_or_default(),
authors,
pdf_url,
arxiv_id,
});
}
out
}
fn unescape(s: &str) -> String {
s.replace("&amp;", "&")
.replace("&lt;", "<")
.replace("&gt;", ">")
.replace("&quot;", "\"")
.replace("&#39;", "'")
}
/// Search arXiv. `max_results` is capped to keep one run bounded.
pub async fn search(query: &str, max_results: usize) -> Result<Vec<Paper>, String> {
let max = max_results.clamp(1, 50);
let url = format!(
"https://export.arxiv.org/api/query?search_query={}&start=0&max_results={max}\
&sortBy=submittedDate&sortOrder=descending",
urlencoding(query)
);
let body = reqwest::Client::new()
.get(&url)
.header("User-Agent", "clawmates-papers/0.1 (research library)")
.timeout(std::time::Duration::from_secs(60))
.send()
.await
.map_err(|e| format!("arxiv query: {e}"))?
.text()
.await
.map_err(|e| format!("arxiv body: {e}"))?;
Ok(parse_atom(&body))
}
/// Download the PDF. Returns the bytes; the caller decides where to shelve it.
pub async fn fetch_pdf(paper: &Paper) -> Result<Vec<u8>, String> {
let bytes = reqwest::Client::new()
.get(&paper.pdf_url)
.header("User-Agent", "clawmates-papers/0.1 (research library)")
.timeout(std::time::Duration::from_secs(180))
.send()
.await
.map_err(|e| format!("fetch pdf {}: {e}", paper.arxiv_id))?
.bytes()
.await
.map_err(|e| format!("read pdf {}: {e}", paper.arxiv_id))?;
// A PDF starts with `%PDF`. arXiv serves an HTML holding page when a PDF
// is still rendering, and shelving that would leave a file that looks
// present and is unreadable.
if !bytes.starts_with(b"%PDF") {
return Err(format!(
"{} did not return a PDF ({} bytes, starts {:?})",
paper.pdf_url,
bytes.len(),
String::from_utf8_lossy(&bytes[..bytes.len().min(16)])
));
}
Ok(bytes.to_vec())
}
/// The catalogue note for a shelved paper.
///
/// `source_id` in the frontmatter is the load-bearing part — it is what
/// `corpus::parse_note` reads to rebuild the checkmark list from the vault.
pub fn catalogue_note(paper: &Paper, blob_key: &str) -> String {
let authors = if paper.authors.is_empty() {
"unknown".to_string()
} else {
paper.authors.join(", ")
};
format!(
"---\n\
source_id: arxiv:{id}\n\
arxiv: {id}\n\
title: \"{title}\"\n\
authors: \"{authors}\"\n\
published: {published}\n\
pdf: {blob_key}\n\
url: https://arxiv.org/abs/{id}\n\
added: {added}\n\
tags: [paper, arxiv]\n\
---\n\
\n\
# {title}\n\
\n\
**Authors:** {authors} \n\
**arXiv:** [{id}](https://arxiv.org/abs/{id}) \n\
**PDF:** `{blob_key}`\n\
\n\
## Abstract\n\
\n\
{summary}\n\
\n\
## Notes\n\
\n\
_Catalogued automatically. Add your own notes below._\n",
id = paper.arxiv_id,
title = paper.title.replace('"', "'"),
authors = authors,
published = paper.published,
blob_key = blob_key,
added = paper.published,
summary = paper.summary,
)
}
fn urlencoding(s: &str) -> String {
s.bytes()
.map(|b| match b {
b'A'..=b'Z' | b'a'..=b'z' | b'0'..=b'9' | b'-' | b'_' | b'.' | b'~' => {
(b as char).to_string()
}
b' ' => "+".to_string(),
_ => format!("%{b:02X}"),
})
.collect()
}
#[cfg(test)]
mod tests {
use super::*;
/// A revision must not read as a new paper.
#[test]
fn version_suffixes_are_stripped() {
assert_eq!(normalize_arxiv_id("http://arxiv.org/abs/2401.12345v3"), "2401.12345");
assert_eq!(normalize_arxiv_id("2401.12345v1"), "2401.12345");
assert_eq!(normalize_arxiv_id("2401.12345"), "2401.12345");
// Old-style ids contain letters and a slash.
assert_eq!(normalize_arxiv_id("http://arxiv.org/abs/cs/0701001"), "0701001");
// A trailing `v` with no digits is part of the id, not a version.
assert_eq!(normalize_arxiv_id("2401.1234v"), "2401.1234v");
}
/// Parsed against the real shape of arXiv's Atom feed. If this fails the
/// format drifted — which otherwise shows up as "no new papers", which is
/// indistinguishable from a quiet week.
#[test]
fn entries_are_parsed_from_a_real_feed() {
let xml = r#"<?xml version="1.0" encoding="UTF-8"?>
<feed xmlns="http://www.w3.org/2005/Atom">
<entry>
<id>http://arxiv.org/abs/2401.12345v2</id>
<published>2026-01-15T10:00:00Z</published>
<title>Attention Is All You Need Again</title>
<summary> We show that
attention still works. </summary>
<author><name>Ada Lovelace</name></author>
<author><name>Alan Turing</name></author>
<link href="http://arxiv.org/abs/2401.12345v2" rel="alternate" type="text/html"/>
<link title="pdf" href="http://arxiv.org/pdf/2401.12345v2" rel="related" type="application/pdf"/>
</entry>
</feed>"#;
let papers = parse_atom(xml);
assert_eq!(papers.len(), 1);
let p = &papers[0];
assert_eq!(p.arxiv_id, "2401.12345", "version stripped");
assert_eq!(p.title, "Attention Is All You Need Again", "whitespace collapsed");
assert_eq!(p.summary, "We show that attention still works.");
assert_eq!(p.authors, vec!["Ada Lovelace", "Alan Turing"]);
assert_eq!(p.pdf_url, "http://arxiv.org/pdf/2401.12345v2");
assert_eq!(p.source_id(), "arxiv:2401.12345");
assert_eq!(p.blob_key(), "papers/arxiv/2401.12345.pdf");
assert_eq!(p.note_path(), "60 Papers/arxiv-2401.12345.md");
}
#[test]
fn an_empty_feed_yields_no_papers_rather_than_panicking() {
assert!(parse_atom("<feed></feed>").is_empty());
assert!(parse_atom("").is_empty());
}
#[test]
fn xml_entities_are_unescaped() {
let xml = r#"<feed><entry><id>http://arxiv.org/abs/1v1</id>
<title>Cats &amp; Dogs &lt;3</title><summary>a &quot;quote&quot;</summary>
</entry></feed>"#;
let p = &parse_atom(xml)[0];
assert_eq!(p.title, "Cats & Dogs <3");
assert_eq!(p.summary, "a \"quote\"");
}
/// The note must carry the identity `corpus::parse_note` reads, or the
/// catalogue cannot rebuild the checkmark list and the library forgets
/// itself the moment the database is lost.
#[test]
fn a_catalogue_note_round_trips_through_the_corpus_parser() {
let paper = Paper {
arxiv_id: "2401.12345".into(),
title: "A \"Quoted\" Title".into(),
authors: vec!["Ada Lovelace".into()],
summary: "Summary text.".into(),
published: "2026-01-15T10:00:00Z".into(),
pdf_url: "http://arxiv.org/pdf/2401.12345".into(),
};
let note = catalogue_note(&paper, &paper.blob_key());
let parsed = crate::corpus::parse_note(&paper.note_path(), &note);
assert_eq!(
parsed.declared_source_id.as_deref(),
Some("arxiv:2401.12345"),
"the corpus parser must recover the identity from the note"
);
assert_eq!(parsed.title.as_deref(), Some("A 'Quoted' Title"));
assert!(note.contains("papers/arxiv/2401.12345.pdf"), "note points at the shelf");
}
#[test]
fn queries_are_url_encoded() {
assert_eq!(urlencoding("all:agent topologies"), "all%3Aagent+topologies");
}
}
+153
View File
@@ -0,0 +1,153 @@
//! The paper library: trigger a run, see what it holds.
//!
//! Thin on purpose. The work lives in [`crate::library`]; this exposes it so
//! a run can be started by a person, a schedule, or the UI rather than only
//! from an integration test.
use axum::extract::{Query, State};
use axum::Json;
use serde::{Deserialize, Serialize};
use serde_json::{json, Value};
use crate::{ApiError, AppState, Authed};
/// Default corpus + repo. Single-operator deployment, so these are constants
/// rather than another table to keep in sync; a second library becomes a
/// request field the day one exists.
const DEFAULT_CORPUS: &str = "valhalla-vault";
const DEFAULT_VAULT_URL: &str = "https://git.redclaw.dev/redclaw/valhalla-vault.git";
#[derive(Deserialize)]
pub struct RunRequest {
/// arXiv queries. Omitted → the topics this project is actually working on.
#[serde(default)]
pub topics: Option<Vec<String>>,
/// Papers per topic. Clamped, because a broad first run against an empty
/// library can otherwise pull hundreds of PDFs in one go.
#[serde(default)]
pub per_topic: Option<usize>,
}
#[derive(Serialize)]
pub struct RunResponse {
pub candidates: usize,
pub already_had: usize,
pub shelved: Vec<String>,
pub failed: Vec<Value>,
pub notes: Vec<String>,
pub branch: String,
pub pushed: bool,
pub error: Option<String>,
/// A run that errored on nothing. Reported explicitly so a caller does not
/// have to infer health from an empty `shelved` list — a quiet week and a
/// broken run both shelve zero papers.
pub healthy: bool,
}
/// POST /api/library/runs — harvest now.
pub async fn run(
State(state): State<AppState>,
Authed(user): Authed,
Json(req): Json<RunRequest>,
) -> Result<Json<RunResponse>, ApiError> {
let blobs = state
.blobs
.clone()
.ok_or_else(|| {
eprintln!("library: blob storage is not configured; cannot shelve PDFs");
ApiError::Internal
})?;
let topics = req
.topics
.filter(|t| !t.is_empty())
.unwrap_or_else(crate::library::default_topics);
let per_topic = req.per_topic.unwrap_or(5).clamp(1, 25);
// Work under the missions root: it is already a writable volume with room
// for checkouts, and it is swept, so a crashed run cannot leak a vault
// clone forever.
let work_root = std::env::temp_dir().join("clawmates-library");
let out = crate::library::run_to_vault(
&state.pool,
&blobs,
user.workspace_id.as_uuid(),
DEFAULT_CORPUS,
DEFAULT_VAULT_URL,
&work_root,
&topics,
per_topic,
None,
)
.await
.map_err(|e| {
// The reason belongs in the log, not in the response: it can carry a
// remote URL and git stderr.
eprintln!("library: run failed: {e}");
ApiError::Internal
})?;
Ok(Json(RunResponse {
candidates: out.harvest.candidates,
already_had: out.harvest.already_had,
shelved: out.harvest.shelved.clone(),
failed: out
.harvest
.failed
.iter()
.map(|(id, why)| json!({ "source_id": id, "error": why }))
.collect(),
notes: out.harvest.notes_written.clone(),
healthy: out.harvest.healthy(),
branch: out.branch,
pushed: out.pushed,
error: out.error,
}))
}
#[derive(Deserialize)]
pub struct ListQuery {
#[serde(default)]
pub kind: Option<String>,
#[serde(default)]
pub limit: Option<i64>,
}
/// `(source_id, title, url, note path)` as stored.
type CorpusRow = (String, Option<String>, Option<String>, Option<String>);
/// GET /api/library/items — what the library holds.
pub async fn list(
State(state): State<AppState>,
Authed(user): Authed,
Query(q): Query<ListQuery>,
) -> Result<Json<Vec<Value>>, ApiError> {
let limit = q.limit.unwrap_or(100).clamp(1, 500);
let kind = q.kind.unwrap_or_else(|| "source".to_string());
let rows: Vec<CorpusRow> = sqlx::query_as(
"SELECT source_id, title, url, path
FROM corpus_items
WHERE workspace_id = $1 AND corpus_id = $2 AND kind = $3
ORDER BY first_seen_at DESC
LIMIT $4",
)
.bind(user.workspace_id.as_uuid())
.bind(DEFAULT_CORPUS)
.bind(kind)
.bind(limit)
.fetch_all(&state.pool)
.await
.map_err(|e| {
eprintln!("library: list corpus: {e}");
ApiError::Internal
})?;
Ok(Json(
rows.into_iter()
.map(|(source_id, title, url, path)| {
json!({ "sourceId": source_id, "title": title, "url": url, "notePath": path })
})
.collect(),
))
}
+1
View File
@@ -13,6 +13,7 @@ pub mod gateway;
pub mod health; pub mod health;
pub mod identity; pub mod identity;
pub mod level_up; pub mod level_up;
pub mod library;
pub mod missions; pub mod missions;
pub mod nodes; pub mod nodes;
pub mod oauth; pub mod oauth;
+228
View File
@@ -0,0 +1,228 @@
//! Indexing the vault must be idempotent, or a continuous mission cannot tell
//! new work from work it already did.
//!
//! These run against a real Postgres via cm-testkit. The vault fixture is
//! shaped from the actual `valhalla-vault`: 416 notes, only 145 with
//! frontmatter, none carrying arxiv/doi/url, plus repo-sync notes whose
//! frontmatter churns on every sync.
use cm_api::corpus;
use uuid::Uuid;
async fn workspace(pool: &sqlx::PgPool) -> Uuid {
let ws = Uuid::now_v7();
sqlx::query("INSERT INTO workspaces (id, name, plan) VALUES ($1,'t','team')")
.bind(ws)
.execute(pool)
.await
.unwrap();
ws
}
fn seed_vault(root: &std::path::Path) {
std::fs::create_dir_all(root.join("50 APESS 2026/Lectures")).unwrap();
std::fs::create_dir_all(root.join("Repos")).unwrap();
std::fs::create_dir_all(root.join("Daily")).unwrap();
// Course note: has frontmatter, but `source:` is a local path.
std::fs::write(
root.join("50 APESS 2026/Lectures/agentic.md"),
"---\nsource: \"/Users/quantum/Downloads/Material/x.pdf\"\ntype: lecture\n---\n# Agentic Design\n\nbody\n",
)
.unwrap();
// Repo-sync note: frontmatter churns, prose does not.
std::fs::write(
root.join("Repos/zeroclaw.md"),
"---\nnode: tank\nupdated: 2026-08-01\nsize_kb: 12\n---\n# ZeroClaw\n\nmirror\n",
)
.unwrap();
// Plain note: no frontmatter at all — the majority case.
std::fs::write(root.join("Daily/2026-08-01.md"), "# Monday\n\nnotes\n").unwrap();
}
#[tokio::test]
async fn indexing_an_unchanged_vault_is_a_no_op() {
let pool = cm_testkit::test_pool().await;
let ws = workspace(&pool).await;
let tmp = tempfile::tempdir().unwrap();
seed_vault(tmp.path());
let first = corpus::index_vault(&pool, ws, "vault", tmp.path())
.await
.unwrap();
assert_eq!(first.scanned, 3);
assert_eq!(first.inserted, 3);
assert_eq!(first.unchanged, 0);
// The decisive assertion: a second pass over an untouched vault must add
// and change nothing. Without this, every run looks like new work.
let second = corpus::index_vault(&pool, ws, "vault", tmp.path())
.await
.unwrap();
assert_eq!(second.scanned, 3);
assert_eq!(second.inserted, 0, "re-index must not insert");
assert_eq!(second.updated, 0, "re-index must not update");
assert_eq!(second.unchanged, 3);
}
#[tokio::test]
async fn a_repo_sync_touching_only_frontmatter_is_not_an_edit() {
let pool = cm_testkit::test_pool().await;
let ws = workspace(&pool).await;
let tmp = tempfile::tempdir().unwrap();
seed_vault(tmp.path());
corpus::index_vault(&pool, ws, "vault", tmp.path())
.await
.unwrap();
// Exactly what a repo sync does: bump `updated`/`size_kb`, prose untouched.
std::fs::write(
tmp.path().join("Repos/zeroclaw.md"),
"---\nnode: tank\nupdated: 2026-08-03\nsize_kb: 14\n---\n# ZeroClaw\n\nmirror\n",
)
.unwrap();
let stats = corpus::index_vault(&pool, ws, "vault", tmp.path())
.await
.unwrap();
assert_eq!(stats.updated, 0, "frontmatter churn is not an edit");
assert_eq!(stats.unchanged, 3);
// A real prose edit must still be seen.
std::fs::write(
tmp.path().join("Repos/zeroclaw.md"),
"---\nnode: tank\nupdated: 2026-08-03\n---\n# ZeroClaw\n\nREWRITTEN\n",
)
.unwrap();
let stats = corpus::index_vault(&pool, ws, "vault", tmp.path())
.await
.unwrap();
assert_eq!(stats.updated, 1, "a genuine edit must be visible");
}
#[tokio::test]
async fn a_hand_edited_note_survives_a_rebuild() {
let pool = cm_testkit::test_pool().await;
let ws = workspace(&pool).await;
let tmp = tempfile::tempdir().unwrap();
seed_vault(tmp.path());
corpus::index_vault(&pool, ws, "vault", tmp.path())
.await
.unwrap();
// The vault is authoritative: a human renames a note by hand.
std::fs::remove_file(tmp.path().join("Daily/2026-08-01.md")).unwrap();
std::fs::write(tmp.path().join("Daily/renamed.md"), "# Monday\n\nnotes\n").unwrap();
let stats = corpus::index_vault(&pool, ws, "vault", tmp.path())
.await
.unwrap();
assert_eq!(stats.scanned, 3);
assert_eq!(stats.inserted, 1, "the renamed note is indexed under its new path");
// The stale row is left alone rather than deleted — the index is derived
// and rebuildable, and losing coverage history is worse than a stale row.
assert!(corpus::seen(&pool, ws, "vault", "note:Daily/renamed.md")
.await
.unwrap());
}
#[tokio::test]
async fn unseen_filters_candidates_in_one_round_trip() {
let pool = cm_testkit::test_pool().await;
let ws = workspace(&pool).await;
corpus::record(
&pool, ws, "vault", "source", "arxiv:2401.11111",
Some("Known"), None, None, "h", None,
)
.await
.unwrap();
let candidates = vec![
"arxiv:2401.11111".to_string(), // already read
"arxiv:2401.22222".to_string(),
"doi:10.1000/new".to_string(),
];
let fresh = corpus::unseen(&pool, ws, "vault", &candidates).await.unwrap();
assert_eq!(fresh, vec!["arxiv:2401.22222", "doi:10.1000/new"]);
assert!(corpus::seen(&pool, ws, "vault", "arxiv:2401.11111").await.unwrap());
assert!(!corpus::seen(&pool, ws, "vault", "arxiv:2401.22222").await.unwrap());
// A different corpus must not inherit another's seen-set.
assert!(!corpus::seen(&pool, ws, "other", "arxiv:2401.11111").await.unwrap());
}
/// The first mission to find a source keeps the credit, so "did THIS run
/// contribute anything new" stays answerable across repeated runs.
#[tokio::test]
async fn re_recording_a_source_does_not_reassign_it() {
let pool = cm_testkit::test_pool().await;
let ws = workspace(&pool).await;
let inserted = corpus::record(
&pool, ws, "vault", "source", "arxiv:2401.33333",
Some("Paper"), None, None, "h1", None,
)
.await
.unwrap();
assert!(inserted, "first sighting is an insert");
let inserted_again = corpus::record(
&pool, ws, "vault", "source", "arxiv:2401.33333",
Some("Paper"), None, None, "h2", None,
)
.await
.unwrap();
assert!(!inserted_again, "a second sighting is not new work");
}
/// Idempotence against the real vault rather than a fixture.
///
/// Ignored by default because it needs a checkout: run with
/// `VAULT=/path/to/valhalla-vault cargo test -p cm-api --test corpus_vault \
/// index_the_real_vault -- --ignored --nocapture`.
///
/// Measured 2026-08-03 on the live vault:
/// PASS1 { scanned: 416, inserted: 416, updated: 0, unchanged: 0 }
/// PASS2 { scanned: 416, inserted: 0, updated: 0, unchanged: 416 }
#[tokio::test]
#[ignore]
async fn index_the_real_vault() {
let pool = cm_testkit::test_pool().await;
let ws = uuid::Uuid::now_v7();
sqlx::query("INSERT INTO workspaces (id, name, plan) VALUES ($1,'t','team')")
.bind(ws).execute(&pool).await.unwrap();
let root = std::path::Path::new(&std::env::var("VAULT").unwrap()).to_path_buf();
let a = cm_api::corpus::index_vault(&pool, ws, "valhalla-vault", &root).await.unwrap();
println!("PASS1 {a:?}");
let b = cm_api::corpus::index_vault(&pool, ws, "valhalla-vault", &root).await.unwrap();
println!("PASS2 {b:?}");
assert_eq!(b.inserted, 0);
assert_eq!(b.updated, 0);
assert_eq!(b.unchanged, a.scanned);
}
/// Live arXiv check. Ignored by default (needs network); run with
/// `cargo test -p cm-api --test corpus_vault live_arxiv -- --ignored --nocapture`.
///
/// Guards the one failure that hides: if arXiv's feed format drifts, parsing
/// returns zero papers, which looks exactly like "no new papers this week".
#[tokio::test]
#[ignore]
async fn live_arxiv_search_and_fetch() {
let papers = cm_api::papers::search("all:agentic topologies", 3)
.await
.expect("arxiv search");
println!("found {} papers", papers.len());
assert!(!papers.is_empty(), "arXiv returned nothing — format drift?");
for p in &papers {
println!(" {} | {}", p.source_id(), &p.title[..p.title.len().min(60)]);
assert!(!p.arxiv_id.is_empty());
assert!(!p.title.is_empty());
assert!(!p.arxiv_id.contains('v'), "version must be stripped: {}", p.arxiv_id);
}
let pdf = cm_api::papers::fetch_pdf(&papers[0]).await.expect("fetch pdf");
println!("pdf bytes: {}", pdf.len());
assert!(pdf.starts_with(b"%PDF"));
assert!(pdf.len() > 10_000, "suspiciously small pdf: {}", pdf.len());
}
+202
View File
@@ -0,0 +1,202 @@
//! A second run must not re-download what the first run already shelved.
use cm_api::{corpus, harvest, papers::Paper};
use std::sync::Arc;
use uuid::Uuid;
async fn workspace(pool: &sqlx::PgPool) -> Uuid {
let ws = Uuid::now_v7();
sqlx::query("INSERT INTO workspaces (id, name, plan) VALUES ($1,'t','team')")
.bind(ws)
.execute(pool)
.await
.unwrap();
ws
}
fn paper(id: &str) -> Paper {
Paper {
arxiv_id: id.into(),
title: format!("Paper {id}"),
authors: vec!["Ada Lovelace".into()],
summary: "A summary.".into(),
published: "2026-01-15T10:00:00Z".into(),
// Deliberately unreachable: if the skip works, this is never fetched.
pdf_url: "http://127.0.0.1:1/never.pdf".into(),
}
}
/// The load-bearing behaviour. Every candidate is already on the checkmark
/// list, and every `pdf_url` points at a closed port — so if the run tries to
/// download anything at all, it fails loudly instead of passing quietly.
#[tokio::test]
async fn papers_we_already_hold_are_never_downloaded_again() {
let pool = cm_testkit::test_pool().await;
let ws = workspace(&pool).await;
let tmp = tempfile::tempdir().unwrap();
let blobs: Arc<dyn cm_files::BlobStore> =
Arc::new(cm_files::LocalBlobStore::new(tmp.path().join("blobs")));
let vault = tmp.path().join("vault");
let candidates = vec![paper("2401.11111"), paper("2401.22222")];
for p in &candidates {
corpus::record(
&pool, ws, "lib", "source", &p.source_id(),
Some(&p.title), None, None, "h", None,
)
.await
.unwrap();
}
let lib = harvest::Library {
pool: &pool, blobs: &blobs, workspace_id: ws,
corpus_id: "lib", vault_root: &vault,
};
let h = harvest::shelve(&lib, &candidates, None).await.unwrap();
assert_eq!(h.candidates, 2);
assert_eq!(h.already_had, 2, "both were already held");
assert!(h.shelved.is_empty());
assert!(
h.failed.is_empty(),
"nothing should have been fetched at all, but got: {:?}",
h.failed
);
assert!(h.healthy(), "a fully-known batch is a healthy quiet week");
assert!(!h.added_anything(), "and it added nothing");
assert!(!vault.exists(), "no notes written for papers we already had");
}
/// A paper that cannot be downloaded must NOT be checked off — otherwise one
/// transient network failure means that paper is never retried.
#[tokio::test]
async fn a_failed_download_leaves_the_paper_unseen_for_next_time() {
let pool = cm_testkit::test_pool().await;
let ws = workspace(&pool).await;
let tmp = tempfile::tempdir().unwrap();
let blobs: Arc<dyn cm_files::BlobStore> =
Arc::new(cm_files::LocalBlobStore::new(tmp.path().join("blobs")));
let vault = tmp.path().join("vault");
let candidates = vec![paper("2401.33333")];
let lib = harvest::Library {
pool: &pool, blobs: &blobs, workspace_id: ws,
corpus_id: "lib", vault_root: &vault,
};
let h = harvest::shelve(&lib, &candidates, None).await.unwrap();
assert_eq!(h.already_had, 0);
assert!(h.shelved.is_empty());
assert_eq!(h.failed.len(), 1, "the unreachable fetch must be reported");
assert!(!h.healthy(), "a failed fetch is not a quiet week");
assert!(
!corpus::seen(&pool, ws, "lib", "arxiv:2401.33333")
.await
.unwrap(),
"a paper we failed to get must stay unseen so a later run retries it"
);
}
/// Live end-to-end: search arXiv, shelve genuinely new papers, then confirm a
/// second identical run adds nothing. Ignored by default (network + Postgres):
/// `cargo test -p cm-api --test harvest_run live_ -- --ignored --nocapture`
#[tokio::test]
#[ignore]
async fn live_end_to_end_run_then_rerun() {
let pool = cm_testkit::test_pool().await;
let ws = workspace(&pool).await;
let tmp = tempfile::tempdir().unwrap();
let blobs: Arc<dyn cm_files::BlobStore> =
Arc::new(cm_files::LocalBlobStore::new(tmp.path().join("blobs")));
let vault = tmp.path().join("vault");
let lib = harvest::Library {
pool: &pool, blobs: &blobs, workspace_id: ws,
corpus_id: "lib", vault_root: &vault,
};
let first = harvest::run(&lib, "all:agentic AND all:topology", 3, None)
.await
.unwrap();
println!("RUN1 {}", first.summary());
for n in &first.notes_written {
println!(" note: {n}");
}
assert!(first.healthy(), "failures: {:?}", first.failed);
assert!(first.added_anything(), "first run should find something new");
// Every note must be readable back through the corpus parser, or the
// catalogue cannot rebuild the checkmark list.
for rel in &first.notes_written {
let text = std::fs::read_to_string(vault.join(rel)).unwrap();
let parsed = corpus::parse_note(rel, &text);
assert!(
parsed
.declared_source_id
.as_deref()
.is_some_and(|s| s.starts_with("arxiv:")),
"note {rel} lost its identity"
);
}
let second = harvest::run(&lib, "all:agentic AND all:topology", 3, None)
.await
.unwrap();
println!("RUN2 {}", second.summary());
assert!(second.healthy());
assert!(
!second.added_anything(),
"a rerun must add nothing — got {:?}",
second.shelved
);
assert_eq!(second.already_had, second.candidates);
}
/// THE REAL RUN. Clones the live vault, harvests our current topics, pushes a
/// branch. Ignored by default — needs network, Postgres and GITEA_TOKEN:
/// `GITEA_TOKEN=… VAULT_URL=… cargo test -p cm-api --test harvest_run \
/// live_library_run -- --ignored --nocapture`
#[tokio::test]
#[ignore]
async fn live_library_run() {
let pool = cm_testkit::test_pool().await;
let ws = workspace(&pool).await;
let tmp = tempfile::tempdir().unwrap();
let blobs: Arc<dyn cm_files::BlobStore> =
Arc::new(cm_files::LocalBlobStore::new(tmp.path().join("shelf")));
let url = std::env::var("VAULT_URL").unwrap();
let topics = cm_api::library::default_topics();
for t in &topics {
println!("topic: {t}");
}
let run = cm_api::library::run_to_vault(
&pool, &blobs, ws, "valhalla-vault", &url,
tmp.path(), &topics, 2, None,
)
.await
.unwrap();
println!("\nRESULT {}", run.harvest.summary());
println!("branch: {} pushed: {}", run.branch, run.pushed);
if let Some(e) = &run.error {
println!("error: {e}");
}
for n in &run.harvest.notes_written {
println!(" note: {n}");
}
for (sid, why) in &run.harvest.failed {
println!(" FAILED {sid}: {why}");
}
// Every shelved paper must have its PDF really on the shelf.
for sid in &run.harvest.shelved {
let id = sid.trim_start_matches("arxiv:");
let key = format!("papers/arxiv/{id}.pdf");
let bytes = blobs.get(&key).await.expect("pdf on the shelf");
assert!(bytes.starts_with(b"%PDF"), "{key} is not a PDF");
println!(" shelf: {key} ({} bytes)", bytes.len());
}
assert!(run.harvest.healthy(), "failures: {:?}", run.harvest.failed);
}
+69
View File
@@ -0,0 +1,69 @@
-- The seen-set for continuous missions.
--
-- Every "continuous X" mission has the same failure mode: it runs again and
-- redoes work it already did. Research resurfaces papers it already read; a
-- security scan re-reports findings already triaged. Orchestration does not
-- fix that — a record of what has already been covered does.
--
-- This repository already tried continuous research once. Migrations 0030-0044
-- built `research_topics`, `research_outcomes` and `loops`; 0053 dropped them
-- all. `research_topics` carried a status lifecycle but no seen-set, so it
-- could run forever and never know what it had covered. That is the gap this
-- table exists to close, and it is the reason it lands before any scheduling.
--
-- Authoritative here rather than in the runtime's memory: ZeroClaw memory is
-- scoped per agent, and mission agents are ephemeral `claw_<uuid>` aliases
-- minted per mission (measured: ~100 of them already). A seen-set that
-- disappears with the agent that wrote it is not a seen-set.
CREATE TABLE corpus_items (
id UUID PRIMARY KEY,
workspace_id UUID NOT NULL REFERENCES workspaces (id) ON DELETE CASCADE,
-- Which corpus this belongs to, e.g. 'valhalla-vault'. A workspace can
-- track several (a vault, a findings ledger, a paper collection).
corpus_id TEXT NOT NULL,
-- 'note' = something already in the corpus (a vault file). Establishes
-- coverage: what has this vault already got?
-- 'source' = an external thing a mission consumed (a paper, an advisory).
-- This is the dedupe key that stops re-reading.
--
-- Both are needed and they answer different questions. Measured against
-- the real vault: 416 notes, and ZERO carry an arxiv/doi/url key — so an
-- ingester keyed only on external identity would index nothing at all.
kind TEXT NOT NULL CHECK (kind IN ('note', 'source')),
-- Stable identity within the corpus. For notes, 'note:<vault-relative
-- path>'; for sources, a natural id like 'arxiv:2401.12345', 'doi:10...'
-- or 'url:<sha256>'. Uniqueness is on this, which is what makes
-- re-ingestion idempotent.
source_id TEXT NOT NULL,
title TEXT,
-- Vault-relative path for notes; NULL for external sources.
path TEXT,
url TEXT,
-- SHA-256 of the content at last sight. Lets a re-index distinguish
-- "unchanged" from "edited" without diffing, so an unchanged vault is a
-- genuine no-op rather than 416 pointless updates.
content_hash TEXT NOT NULL,
first_seen_at TIMESTAMPTZ NOT NULL DEFAULT now(),
last_seen_at TIMESTAMPTZ NOT NULL DEFAULT now(),
-- Which mission first recorded this. NULL for the initial vault index,
-- which is derived from files nobody's mission wrote.
mission_id UUID REFERENCES missions (id) ON DELETE SET NULL,
UNIQUE (workspace_id, corpus_id, source_id)
);
-- The hot query is "have I seen this?", which the UNIQUE index already covers.
-- This one serves "what does this corpus contain?" for briefing assembly.
CREATE INDEX corpus_items_corpus_idx
ON corpus_items (workspace_id, corpus_id, kind, last_seen_at DESC);
-- "What did this mission add?" — the verification predicate for a continuous
-- run is that it contributed at least one NEW source.
CREATE INDEX corpus_items_mission_idx
ON corpus_items (mission_id)
WHERE mission_id IS NOT NULL;