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}