Skip to main content

doiget_core/store/
mod.rs

1//! Filesystem-backed metadata store.
2//!
3//! Binding spec: [`docs/STORE.md`](../../../../docs/STORE.md) (NORMATIVE shared
4//! spec for layout, schema, lock protocol, atomic write, normalization).
5//! Public API surface: `docs/PUBLIC_API.md` §2 (Store trait), §3 (Metadata).
6//!
7//! ## Entry points
8//!
9//! - [`Store`] — the trait surface implementations expose.
10//! - [`FsStore`] — filesystem-backed implementation rooted at a configurable
11//!   directory (default `./papers`, under the cwd; ADR-0036).
12//! - [`Metadata`] / [`DoigetExtension`] — the on-disk schema, mirrored from
13//!   `docs/STORE.md` §2.
14//!
15//! ## Other writers
16//!
17//! A store may also be written by another tool -- BiblioFetch.jl stores stay
18//! readable, though the shared contract is retired (ADR-0060). Writers follow
19//! the lock protocol in `docs/STORE.md` §4 and the atomic-write sequence in §5. Per §6, doiget MUST NOT overwrite reserved
20//! top-level fields previously written by another tool — see [`FsStore::write`].
21
22pub mod citekey;
23mod fs_store;
24pub mod metadata;
25pub mod render;
26
27pub use fs_store::FsStore;
28
29/// Crash-consistent write (tmp + fsync + rename), shared with the resolver
30/// cache so both write the same way. See `docs/STORE.md` §5.
31pub(crate) use fs_store::atomic_write;
32pub use metadata::{DoigetExtension, Metadata, ORIGIN_USER_SUPPLIED};
33pub use render::{to_bibtex, to_csl_array};
34
35/// Run a synchronous [`Store`] call from async code without stalling the
36/// runtime (#590).
37///
38/// Every `Store` method does blocking filesystem I/O: `write` can poll the
39/// advisory lock for up to 5 s (`LOCK_TIMEOUT`, `std::thread::sleep` in
40/// 50 ms steps) and then `fsync`s, which on a Dropbox / OneDrive / SMB store
41/// root costs hundreds of milliseconds uncontended. Called directly from an
42/// `async fn`, that holds a tokio worker for the whole duration, delaying
43/// every other task on it -- including the rate limiter's timers and other
44/// in-flight MCP tool calls.
45///
46/// On a multi-thread runtime (`#[tokio::main]`, `doiget serve`) this runs
47/// `f` under [`tokio::task::block_in_place`], which hands the worker's other
48/// tasks to another thread first. Elsewhere -- no runtime, or a
49/// current-thread runtime, where `block_in_place` would panic -- `f` runs
50/// inline, as before.
51///
52/// `block_in_place` rather than `spawn_blocking` because the orchestrator
53/// holds the store as `&dyn Store` (`docs/PUBLIC_API.md` §2): moving it into
54/// a `'static` closure would change the public signature of `fetch_paper`,
55/// which is the cost the issue ruled out for making the trait async.
56pub fn blocking_section<T>(f: impl FnOnce() -> T) -> T {
57    match tokio::runtime::Handle::try_current().map(|h| h.runtime_flavor()) {
58        Ok(tokio::runtime::RuntimeFlavor::MultiThread) => tokio::task::block_in_place(f),
59        _ => f(),
60    }
61}
62
63use camino::Utf8Path;
64use serde::Serialize;
65use thiserror::Error;
66
67use crate::Safekey;
68
69/// Brief summary of a stored entry; returned by
70/// [`Store::list_recent`] / [`Store::search`].
71///
72/// `non_exhaustive` so adding new summary fields (e.g. `doi`, `authors`) in a
73/// later revision is non-breaking. Pattern-match with a wildcard arm.
74///
75/// `Serialize` enables `list-recent --mode json` / `search --mode json`
76/// (#204) — the wire form is the obvious field-name JSON: `{"safekey":
77/// "...", "title": "...", "year": 2024, "fetched_at": "2026-05-20T…Z"}`,
78/// with `null` for absent optionals.
79///
80/// # Wire-format stability (post-#208 self-review §1)
81///
82/// Once a release ships with the \[`Serialize`\] derive, the field
83/// **names** below become part of the public API: a downstream consumer
84/// (CLI agent, MCP tool, BiblioFetch.jl, third-party script) MAY bind
85/// to them. Renaming a field is then a semver minor bump and warrants
86/// a CHANGELOG \[BREAKING\] note. Adding new fields is still safe
87/// (per `#[non_exhaustive]`).
88#[derive(Debug, Clone, Serialize)]
89#[non_exhaustive]
90pub struct EntryInfo {
91    /// The safekey of the entry. See `docs/SAFEKEY.md`.
92    pub safekey: Safekey,
93    /// Title from the entry's reserved `title` field.
94    pub title: String,
95    /// Year, if any, from the entry's reserved `year` field.
96    pub year: Option<i32>,
97    /// `fetched_at` from the `[doiget]` table, if any.
98    pub fetched_at: Option<chrono::DateTime<chrono::Utc>>,
99    /// `size_bytes` from the `[doiget]` table: the size of the stored PDF,
100    /// `0` for a metadata-only entry, `None` when the entry has no
101    /// `[doiget]` table at all.
102    ///
103    /// #481: without it the inventory commands could not tell a fetched
104    /// paper from a metadata-only stub. Every other surface could -- the
105    /// fetch itself failed loudly, the TOML omits `pdf_path`, the
106    /// provenance log carries an `err` row, `doiget info` shows
107    /// `size_bytes = 0` -- and the one command that answers "what do I
108    /// have?" without knowing the ref in advance was the one that dropped
109    /// it. Fifty refs with ten blocked listed as fifty identical rows.
110    pub size_bytes: Option<u64>,
111}
112
113impl EntryInfo {
114    /// Whether a PDF was actually stored for this entry.
115    ///
116    /// `false` for a metadata-only entry (`size_bytes == 0`) and for one
117    /// with no `[doiget]` table. Deliberately not "is this entry useful" --
118    /// a metadata-only entry is a legitimate result, it is just a different
119    /// one, and #118 is the standing rule that the two must not be
120    /// presented alike.
121    #[must_use]
122    pub fn has_pdf(&self) -> bool {
123        self.size_bytes.is_some_and(|n| n > 0)
124    }
125}
126
127/// Errors emitted by [`Store`] implementations.
128#[derive(Debug, Error)]
129#[non_exhaustive]
130pub enum StoreError {
131    /// Underlying I/O failure.
132    #[error("io error: {0}")]
133    Io(#[from] std::io::Error),
134    /// Malformed TOML or schema mismatch on read.
135    #[error("toml deserialize error: {0}")]
136    Deserialize(#[from] toml::de::Error),
137    /// Failed to serialize a [`Metadata`] to TOML.
138    #[error("toml serialize error: {0}")]
139    Serialize(#[from] toml::ser::Error),
140    /// Could not acquire the advisory `flock` within the 5 s budget named in
141    /// `docs/STORE.md` §4.
142    #[error("flock timeout (5s) on {path}")]
143    LockTimeout {
144        /// The lock-file path that was contended.
145        path: camino::Utf8PathBuf,
146    },
147    /// The on-disk `schema_version` is a future major; per `docs/STORE.md` §3
148    /// the entry is read-only for this build.
149    #[error("schema_version too new: {theirs} > {ours}; entry is read-only")]
150    SchemaTooNew {
151        /// Schema version observed on disk.
152        theirs: String,
153        /// Schema version this build supports.
154        ours: String,
155    },
156    /// A reserved field that the spec marks as required is missing.
157    #[error("required field missing: {field}")]
158    MissingField {
159        /// The name of the missing reserved field.
160        field: &'static str,
161    },
162    /// The supplied [`Safekey`] resolves to a path outside the store root.
163    /// Defense-in-depth check; `Safekey` construction already enforces the
164    /// `[A-Za-z0-9._-]`-only charset per `docs/SAFEKEY.md`.
165    #[error("path is outside the store root: {path}")]
166    PathTraversal {
167        /// The offending resolved path.
168        path: camino::Utf8PathBuf,
169    },
170}
171
172/// Who is authoritative for the user-authored `[doiget]` fields
173/// (`tags`, `collections`, `annotation`) on a write.
174///
175/// A fetch never authors them: every `DoigetExtension` the orchestrator
176/// builds hard-codes `Vec::new()` / `None`. Letting that win silently
177/// discarded a user's tags on any re-fetch, which is the loss ADR-0056
178/// closed for `oa_status` / `license` and left open here.
179///
180/// The distinction has to be explicit rather than "is the incoming value
181/// empty", because `doiget tag --remove` and `doiget annotate --clear`
182/// legitimately mean the empty value.
183#[derive(Debug, Clone, Copy, PartialEq, Eq)]
184#[non_exhaustive]
185pub enum UserFields {
186    /// The caller did not author them; keep whatever is on disk. Fetches.
187    Preserve,
188    /// The caller means exactly what it wrote, empty included. `doiget tag`,
189    /// `doiget annotate`, and their MCP equivalents.
190    Authored,
191}
192
193/// Filesystem-shaped metadata store, semver-locked per `docs/PUBLIC_API.md`
194/// §2.
195///
196/// Implementations are responsible for honoring:
197///
198/// - `docs/STORE.md` §4 lock protocol (advisory `flock` on
199///   `<safekey>.toml.lock` with a 5 s timeout).
200/// - `docs/STORE.md` §5 atomic-write sequence (`tmp` → fsync → rename →
201///   fsync parent).
202/// - `docs/STORE.md` §6 doiget write discipline: never overwrite reserved
203///   top-level fields previously written by another tool.
204/// - `docs/STORE.md` §7 TOML normalization (alphabetical key order, `\n`
205///   line endings, trailing newline).
206pub trait Store: Send + Sync {
207    /// Read the entry keyed by `key`.
208    ///
209    /// Returns `Ok(None)` if no entry exists. Returns `Err` on I/O failure,
210    /// malformed TOML, or unrecoverable schema mismatch (e.g. future major).
211    fn read(&self, key: &Safekey) -> Result<Option<Metadata>, StoreError>;
212
213    /// Write or update the entry keyed by `key`.
214    ///
215    /// If `pdf` is `Some`, the file at that path is copied to
216    /// `<root>/<safekey>.pdf` via the same atomic-rename dance as the
217    /// metadata file. The caller is responsible for emitting the
218    /// `event=store_write` provenance row (see `docs/PROVENANCE_LOG.md` §3).
219    /// Write a fetch result. User-authored `[doiget]` fields already on disk
220    /// are preserved ([`UserFields::Preserve`]) -- a fetch does not author
221    /// them, and silently dropping them is data loss.
222    fn write(&self, key: &Safekey, m: &Metadata, pdf: Option<&Utf8Path>) -> Result<(), StoreError>;
223
224    /// Write on behalf of a caller that DID author the user fields, so an
225    /// empty `tags` / `collections` or a `None` annotation means exactly that
226    /// ([`UserFields::Authored`]). `doiget tag` / `doiget annotate` only.
227    fn write_user_authored(
228        &self,
229        key: &Safekey,
230        m: &Metadata,
231        pdf: Option<&Utf8Path>,
232    ) -> Result<(), StoreError>;
233
234    /// Return up to `limit` entries, most-recent first by `[doiget].fetched_at`.
235    fn list_recent(&self, limit: usize) -> Result<Vec<EntryInfo>, StoreError>;
236
237    /// Return up to `limit` entries whose title / authors / venue / publisher
238    /// case-insensitively contain `query`.
239    fn search(&self, query: &str, limit: usize) -> Result<Vec<EntryInfo>, StoreError>;
240}
241
242#[cfg(test)]
243#[allow(clippy::expect_used, clippy::unwrap_used)]
244mod tests {
245    use std::sync::atomic::{AtomicUsize, Ordering};
246    use std::sync::Arc;
247    use std::time::Duration;
248
249    /// Ticks a spawned task managed during `hold` of blocking work, run on
250    /// the ONLY worker of a one-worker multi-thread runtime.
251    async fn ticks_during(hold: Duration, through_helper: bool) -> usize {
252        let ticks = Arc::new(AtomicUsize::new(0));
253        let running = Arc::new(std::sync::atomic::AtomicBool::new(true));
254        let (t, r) = (Arc::clone(&ticks), Arc::clone(&running));
255        let ticker = tokio::spawn(async move {
256            while r.load(Ordering::SeqCst) {
257                t.fetch_add(1, Ordering::SeqCst);
258                tokio::time::sleep(Duration::from_millis(5)).await;
259            }
260        });
261        // Let the ticker start, then measure only the blocked window.
262        tokio::time::sleep(Duration::from_millis(20)).await;
263        let before = ticks.load(Ordering::SeqCst);
264        let blocker = tokio::spawn(async move {
265            if through_helper {
266                super::blocking_section(|| std::thread::sleep(hold));
267            } else {
268                std::thread::sleep(hold);
269            }
270        });
271        blocker.await.expect("blocker");
272        let during = ticks.load(Ordering::SeqCst) - before;
273        running.store(false, Ordering::SeqCst);
274        ticker.await.expect("ticker");
275        during
276    }
277
278    /// #590: a store call that blocks (lock poll, fsync on a synced
279    /// folder) must not stall the other tasks on its worker.
280    #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
281    async fn a_blocking_store_call_leaves_the_runtime_responsive() {
282        let hold = Duration::from_millis(500);
283        // Control: the same wait called directly starves the ticker, so the
284        // assertion below measures the helper and not the scheduler's luck.
285        let direct = ticks_during(hold, false).await;
286        let wrapped = ticks_during(hold, true).await;
287        assert!(
288            direct <= 3,
289            "control: direct blocking let {direct} ticks through"
290        );
291        // Loose on purpose: Windows' default timer granularity (~15.6 ms)
292        // caps a 5 ms ticker near 32 ticks in 500 ms, before CI load. The
293        // control above is what makes the bound meaningful.
294        assert!(
295            wrapped >= 8,
296            "blocking_section let only {wrapped} ticks through"
297        );
298    }
299
300    #[test]
301    fn blocking_section_runs_inline_without_a_multi_thread_runtime() {
302        assert_eq!(super::blocking_section(|| 7), 7);
303        let rt = tokio::runtime::Builder::new_current_thread()
304            .build()
305            .expect("runtime");
306        assert_eq!(rt.block_on(async { super::blocking_section(|| 8) }), 8);
307    }
308
309    /// The call sites cannot express the convention in the type system
310    /// (the trait is sync on purpose), so this pins it in the source: no
311    /// `Store` method is called from the orchestrator except through
312    /// `blocking_section`.
313    #[test]
314    fn every_orchestrator_store_call_goes_through_blocking_section() {
315        let src = include_str!("../orchestrator.rs");
316        assert_eq!(unwrapped_store_calls(src), Vec::<String>::new());
317        assert!(src.contains("blocking_section(|| store.write("));
318    }
319
320    /// Every `store.<method>(` in the non-test part of `src` that is not the
321    /// body of a `blocking_section` closure. Whitespace is stripped first, so
322    /// neither rustfmt re-wrapping a call nor a method chain split across
323    /// lines (`store\n    .read(`) hides one. Calls on a receiver not named
324    /// `store` are out of its reach; the call sites keep that name.
325    pub(crate) fn unwrapped_store_calls(src: &str) -> Vec<String> {
326        let body = src.split("\nmod tests {").next().unwrap_or(src);
327        let flat: String = body.chars().filter(|c| !c.is_whitespace()).collect();
328        let mut out = Vec::new();
329        for m in STORE_METHODS {
330            let calls = flat.match_indices(&*format!("store.{m}(")).count();
331            let wrapped = flat
332                .match_indices(&*format!("blocking_section(||store.{m}("))
333                .count()
334                + flat
335                    .match_indices(&*format!("blocking_section(||{{store.{m}("))
336                    .count();
337            if calls != wrapped {
338                out.push(format!("store.{m}: {calls} calls, {wrapped} wrapped"));
339            }
340        }
341        out
342    }
343
344    /// The `Store` trait's methods and `FsStore`'s inherent search.
345    pub(crate) const STORE_METHODS: &[&str] = &[
346        "read",
347        "write",
348        "write_user_authored",
349        "list_recent",
350        "search",
351        "search_by_tag",
352    ];
353}