1
Fork 0
mirror of https://github.com/thegeneralist01/archivr synced 2026-10-09 21:03:17 +02:00

feat: capture, summaries, search, and yt-dlp reliability (#38)

* ui: show spinner for pending captures

* feat(core): add text capture path with title + Markdown/plain body

- Add downloader/text.rs module with save() function that stages and hashes text content
- Support text/markdown and text/plain MIME types with .md and .txt extensions
- Add perform_text_capture() function for capturing user-supplied text
- Validates title (non-empty, max 500 chars) and body (non-empty, max 2 MiB)
- Creates blob records and entries with source_kind='text', entity_kind='document'
- Includes comprehensive unit tests for markdown, plain text, and validation

Co-Authored-By: Claude Haiku 4.5 <noreply@anthropic.com>

* feat(server): add POST /api/archives/:archive_id/captures/text

- Add CaptureTextBody struct for title, body, and optional MIME type
- Implement capture_text_handler with validation for empty fields and MIME type
- Route text submissions to perform_text_capture() in background
- Reuse existing capture job tracking and polling infrastructure
- Default MIME type to text/markdown when not specified
- Include route tests covering happy path, validation, auth, and error cases

Co-Authored-By: Claude Haiku 4.5 <noreply@anthropic.com>

* feat(frontend): add text-capture form to CaptureDialog

- Add submitTextCapture API client function with same error handling as submitCapture
- Create makeTextItem() factory for text capture state
- Implement CaptureTextRow component with title, body textarea, and MIME selector
- Add 'Add text' button in capture dialog toolbar
- Update handleArchive to filter and route text submissions
- Modify submitBgJob to detect and submit text items via submitTextCapture
- Skip probe and conflict checks for text items
- Reuse job tracking and batch settlement for text captures

Co-Authored-By: Claude Haiku 4.5 <noreply@anthropic.com>

* feat(core): add entry_summaries schema + summarizer trait/providers

Per-entry LLM summaries as a regenerable child record, not a column on
archived_entries and not an on-disk artifact: an entry may carry several
summaries (one per provider/model/prompt version), any of which can be
discarded and recomputed. Generation is manual-only — nothing in capture.rs
calls into this module.

- database.rs: entry_summaries table + index, EntrySummaryRecord, and
  upsert/update/find/latest helpers mirroring the capture_jobs style.
  provider_model is stored as '' rather than NULL because SQLite treats
  NULLs as distinct inside a UNIQUE index, which would stop the CLI
  providers (no model) from ever deduping on the cache key.
- summarizer.rs: SummaryProvider trait with four implementations —
  Anthropic Messages API, OpenAI-compatible chat completions, `claude -p`
  and `codex exec -`. Configuration comes from env vars only (never TOML),
  matching how yt-dlp / single-file / tweet-scraper are resolved, which
  also keeps API keys out of anything the archive persists.
- archive.rs: EntryDetail gains latest_summary, populated by one extra
  LIMIT 1 query in get_entry_detail. EntrySummaryView aliases the DB row
  rather than duplicating it.

Implementation notes:
- No tokio in core. CLI timeouts are enforced structurally: stdout is
  drained on its own thread and handed back over a channel so the calling
  thread can recv_timeout and kill an overrunning child; stdin is written
  on a third thread so a 48 KB prompt cannot deadlock against a child
  waiting for us to read.
- HTML is reduced with regex rather than a parser: html5ever is not in the
  tree, and a model tolerates imperfect whitespace. Paired tags are spelled
  out per tag because Rust's regex engine has no backreferences by design.
- reqwest is declared with only the `blocking` feature here, so bodies are
  serialized via .body(value.to_string()) instead of widening the
  workspace dependency for .json().
- input_sha256 holds a SHA3-256 digest via hash::hash_bytes, the tree's one
  hashing primitive; the content is truncated to 48 KB *before* hashing so
  the cache key describes exactly the bytes the model saw.

Tests: no mockito/wiremock in dev-deps, and adding a mock HTTP server for
one JSON shape is a poor trade, so the two halves that can actually break
are tested directly — request-body builders and response parsers — leaving
only reqwest's own transport uncovered. Plus schema idempotency, cache-key
dedupe, cascade-on-delete, provider_from_env happy/missing-var paths, HTML
and tweet extraction, output normalization, and the CLI runner's stdin
round-trip, timeout kill, and nonzero-exit paths.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>

* feat(server): add GET/POST /api/archives/:id/entries/:uid/summary

GET is read-only and gated exactly like entry detail, so a guest can read a
summary only for an entry whose content they could already read. POST
requires ROLE_USER, matching capture / tags / patch / rearchive; no auth
roles change.

Both the provider config and the content extraction resolve on the request
thread, before spawn_blocking. That is what lets a missing env var come back
as a synchronous 400 naming the exact variable, and an unsummarizable
artifact (video, audio) as a 400 saying so, rather than becoming a
background job the caller must poll only to learn about a config typo.

The pending row is claimed before spawning so the 202 can name a summary_uid
the client can poll immediately. summarize_entry owns the
pending → running → completed/failed transitions for that same row — the
cache key is identical, so both upserts resolve to one row — leaving the
handler to catch only the case where it fails before recording anything.
When !force and an identical cache key already completed, the existing row
comes back as a 200 with no new work.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>

* feat(frontend): render Summary rail section + provider selector

New "Summary" rail section between the URL/Preview controls and .meta-list.
A completed summary renders as bold tl;dr, body paragraph, tag chips, and a
provider · model footer; missing or failed shows Generate; pending/running
shows an inline spinner and polls GET every 1500 ms until terminal.

- api.js: fetchEntrySummary + requestEntrySummary. The POST helper unwraps
  ApiError's { "error": ... } body so the missing-env-var message reaches
  the user verbatim rather than as a bare status code.
- ContextRail.jsx: state seeds from detail.latest_summary so the section
  renders immediately on selection. Polling is anchored on the summary
  status rather than started inside the click handler, so a job still
  running when the user navigates away and back is picked up again. A
  transient poll failure is swallowed — the next tick retries, and a real
  failure arrives as status === 'failed'.
- Regenerate passes force:true only when a completed summary is already
  shown; otherwise the request can take the server's 200 cache-hit path.
- Provider choice persists in sessionStorage under archivr:summary:provider,
  with try/catch around both accessors for private-mode browsers.
- Public sessions never see the selector or the Generate button, and the
  section renders at all only when a completed summary made it through the
  server's visibility gate.
- styles.css: .rail-summary-* only; spacing and the action button reuse
  .rail-section and .rail-rearchive-btn. The spinner honours
  prefers-reduced-motion — the text alone conveys the state.
- AGENTS.md: document the summary env vars alongside the existing
  external-tool convention.

Smoke-tested end to end against a scratch archive with a seeded markdown
entry: claude_cli produced a real summary (pending → running → completed in
~11s); a local mock server exercised the openai_compatible transport and
confirmed the Bearer header, model, and system/user role split on the wire;
unconfigured providers return 400 naming the exact variable; a video entry
returns 400 "v1 unsupported"; a repeat POST returns 200 from cache without
adding a row.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>

* fix(core): codex_cli — auto-discover binary + use --output-last-message

Two related fixes for the codex_cli summary provider:

1. Executable discovery. `ARCHIVR_CODEX_CLI` was already respected, but
   without it the code resolved to bare `codex` and relied on PATH.
   The ChatGPT desktop app installs codex at
   `/Applications/ChatGPT.app/Contents/Resources/codex` and does not
   put it on PATH, so users who only have the desktop app saw
   'No such file or directory' with no hint. `resolve_cli` now walks
   env override → a small set of well-known absolute paths → HOME
   /.local/bin/<bare> → bare fallback. Same treatment applied to
   claude_cli for symmetry (/opt/homebrew/bin/claude, /usr/local/bin/
   claude, HOME/.local/bin/claude).

2. Clean output. `codex exec -` writes a runtime header ("OpenAI
   Codex vX", session id, sandbox, model), the assistant reply, and a
   footer ("tokens used", replay of the reply) to stdout. The JSON
   extractor took the first '{' from the *user prompt echo* and the
   last '}' from the trailing replay, producing invalid text that
   fell through to the "raw text under summary" fallback path. Now
   uses `--output-last-message <tempfile>` and reads only the final
   assistant message. Fallback (positional prompt) uses the same
   flag. Tempfile is cleaned up on all paths, incl. spawn failure.

* fix(frontend): give text-capture row real CSS

The text row shipped with semantic classnames (`capture-text-inputs`,
`capture-text-title`, `capture-text-body`, `capture-text-mime`,
`capture-text-icon`) but no CSS rules. Falling through to the parent
`.capture-row-main` flex-row (`display: flex; align-items: center`)
meant the title, textarea, and mime-select stacked as intrinsic-width
boxes centered on the tall body, producing a layout where the body
floated to the top-right, the title box appeared BELOW it, and the mime
selector rendered as an unstyled OS dropdown.

Fix:
- `.capture-text-row .capture-row-main` uses `align-items: flex-start`
  so the leading icon and trailing × pin to the top of the block.
- `.capture-text-inputs` is now a full-width column-flex container with
  proper gaps.
- `.capture-text-title` reuses the 44px input height and typography of
  `.capture-input`; `.capture-text-body` gets a 140px min-height,
  vertical resize, and matching border/focus treatment.
- `.capture-text-mime` is styled as a small chip with a custom caret
  so it matches `.capture-quality` and stops looking like a raw
  `<select>`. Sits in a right-aligned footer under the body.
- `.capture-text-icon` gets a 44px column so it aligns with the title
  input; remove button gets a small top-margin for the same reason.

Rebuilt static bundle bumped as well (`index-BLxoi9rt.css`,
`index-CQcpPA_I.js`).

* chore(static): rebuild bundle after text-row CSS merge

* fix(core): summarize tweets + walk all tweets in a thread

Tweet and tweet_thread entries store their payload under artifact_role
`raw_tweet_json`, not `primary_media`. `build_summary_input` filtered
strictly for `primary_media LIMIT 1`, so both cases silently failed
with 'entry X has no primary_media artifact to summarize'.

Threads compound the problem: the tweet scraper writes ONE json file
per status, so even a fixed lookup that took the first row would
summarize only the initial tweet and lose the rest of the conversation.

Fixes:
- New `load_summary_artifacts` helper returns every artifact for a
  role in insertion order.
- For entity_kind `tweet` / `tweet_thread`, load all
  `raw_tweet_json` artifacts (falling back to `primary_media` for
  archives predating that role convention).
- Iterate artifacts, extract text per file with the existing
  markdown/html/json branches, then join thread pieces with a
  `---` separator so the model sees a real paragraph break between
  statuses instead of one flowing document.

Single-tweet entries produce one piece and the separator never
renders. Non-tweet entries behave exactly as before.

* feat(frontend): preview text-capture entries (.md / .txt)

Text captures land as `.md` (Markdown) or `.txt` (plain) blobs, but
PreviewPanel only dispatched on video/audio/image/pdf/html extensions,
so opening a text entry hit the 'No preview available' fallback with
the raw artifact path exposed.

- New `TextPreview` component fetches the primary artifact as text,
  renders it in a monospace `<pre>` with word-wrap, and shows the
  entry title on top and the MIME as a small trailing tag. Handles
  loading/error states.
- `PreviewPanel` gains a `TEXT_EXTS` set + a branch that dispatches
  to `TextPreview` for `md` / `markdown` / `txt`.
- CSS is padded and centered to ~780px so a text note reads like a
  document rather than an edge-to-edge terminal dump.

v1 intentionally does NOT parse Markdown: keeping frontend deps at
react+react-dom only. Bump to a real Markdown renderer if we start
capturing Markdown-authored notes.

* fix(nix): pin yt-dlp from its own release + wire into server wrapper

Two independent problems, one commit:

1. Stale binary. nixpkgs-provided `pkgs.yt-dlp` on the pinned
   nixos-unstable rev is 2026.03.17 (Mar 2026). yt-dlp itself
   releases days-to-weeks, and YouTube frequently rotates the
   player-signature / client surfaces the older builds request
   (`android_vr` is the current casualty), which returns HTTP 403
   mid-download for the format specs archivr passes (`-f
   bestvideo+bestaudio/best`). Even bumping the nixpkgs input would
   leave us dependent on that channel's yt-dlp cadence.

   Fetch the upstream zipapp directly instead
   (github.com/yt-dlp/yt-dlp/releases/download/<ver>/yt-dlp), wrap so
   `python3` and `ffmpeg` are on PATH, and pin version+hash in one
   place. Bumping is: change version, replace hash from
   `nix hash file <url>`.

2. Missing pin in server wrapper. `archivr-cli` was already wrapped
   with `--set ARCHIVR_YT_DLP` + a PATH prefix; `archivr-server`
   was NOT — it only pinned single-file, chrome, and the tweet
   scraper, silently falling back to whatever `yt-dlp` the user
   happened to have on PATH. Server captures therefore inherited
   the user's (often stale) system yt-dlp regardless of the flake
   pin. Same wrapper flags now apply to both binaries.

devShell keeps `pkgs.yt-dlp` for now: the dev shell is a
convenience, not a release surface, and matching wouldn't fit in this
commit without duplicating the derivation across let-scopes.

* chore(static): rebuild bundle for round-3 fixes

* chore(nix): pin python 3.12 for yt-dlp zipapp (avoid py3.14 libffi crash on darwin/arm64)

* feat(core): resolve_yt_dlp picks the newer of pinned vs state-dir

The nix flake wrapper pins a yt-dlp via ARCHIVR_YT_DLP, but yt-dlp rots
fast — extractors break within weeks of a pin. Add a resolver that probes
`--version` on both the pinned binary and a user-installed copy under the
mutable state dir, and runs whichever is newer.

Version strings are YYYY.MM.DD, so plain string ordering is chronological.
Ties resolve toward the state dir: a user who installed it there did so
deliberately. ARCHIVR_YT_DLP_FORCE bypasses the comparison entirely, and
with no candidate at all we fall back to bare `yt-dlp` on PATH — exactly
the previous behaviour.

Resolution is cached in a OnceLock so `--version` costs one subprocess per
process, and all four inline env::var lookups now go through it.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>

* ci: bump yt-dlp from upstream releases, not nixpkgs

The flake no longer takes yt-dlp from nixpkgs; a dedicated `ytDlp`
derivation fetches the upstream release binary directly and pins both
`version` and an SRI `hash`. That makes the previous workflow inert: it
ran `nix flake update nixpkgs` and compared `nixpkgs#yt-dlp.version`
before and after, so it could churn the lockfile forever without ever
moving the version we actually ship.

The workflow now reads the pinned version straight out of the `ytDlp`
block in flake.nix, asks the GitHub API for yt-dlp's latest release tag,
short-circuits when they already match, downloads the new release to
recompute its SRI hash (required — the hash is part of the derivation's
identity, so the URL cannot be changed alone), and rewrites the three
pinned fields under a sed range address scoped to that block so sibling
pins like ublockLite and isdcac are untouched. It asserts only flake.nix
changed and that the new version appears exactly twice before opening
the PR.

* feat(cli): add `archivr yt-dlp update|status` subcommand

`update` fetches the latest release tag from the GitHub API (or takes
--version), downloads the cross-platform python zipapp, and installs it
into archivr's state dir. The install is atomic — staged as yt-dlp.new,
chmod +x'd, then renamed over the target — so a concurrently running
capture never sees a half-written binary. A sibling .version file makes a
repeat update a no-op instead of a 3MB re-download.

The download is checked for the python3 shebang before install, which
catches the usual failure mode of getting an HTML error page back. python3
itself is only warned about, not required: the server may run under a nix
wrapper with its own PATH.

`status` prints all three candidates (env / state-dir / PATH fallback) with
their versions and stars whichever the resolver picks, so it is obvious
which yt-dlp a capture will actually use.

reqwest is pulled from the existing workspace dependency; the GitHub JSON is
parsed with serde_json so the "json" feature is not needed.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>

* docs(maintainer): document summarizer, text capture, and yt-dlp lifecycle

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>

* docs(readme): document LLM summaries, text notes, yt-dlp resolver + bump paths

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>

* docs(plan): specify X Article, image summaries, summary search

* fix: summarize X Article text

* feat: search completed summary tags

* fix: preserve X Article block order

* feat: model opt-in summary images

* feat: attach opted-in images to summaries

* feat: accept image summary requests

* feat: add summary image consent control

* docs: explain X Article, image summaries and summary search

* chore(static): rebuild bundle for summary image consent

* fix: show civilized unsupported summary errors

* chore(static): rebuild bundle for civilized summary errors

* fix: keep text previews and summary state scoped to entry

* fix: show forced yt-dlp candidate in status

* fix: infer trusted MIME for tweet images

* fix: preserve text capture bytes and hide synthetic URL

* fix: recover and preserve summary attempts

* fix: bound codex fallback and record resolved model

* fix: protect public summary diagnostics

* fix: guard summary callbacks during entry render

* fix: scope summary callbacks to selected entry

* test: cover terminal newline in text capture

* fix: preserve missing summary entry status

* test: cover tweet image summary selection

* docs: record summary lifecycle and review hardening

* chore(static): rebuild bundle for Sol review fixes

* fix: retain completed summary during regeneration display

* fix: hide superseded summary attempts

* chore(static): rebuild bundle after regeneration display fix

* fix: allow full-size text capture requests

* fix: allow escaped full-size text captures

* fix: preserve text draft whitespace in capture UI

* chore(static): rebuild bundle for text whitespace fix
This commit is contained in:
TheGeneralist 2026-08-24 18:30:01 +02:00 • committed by GitHub
parent b4b4e67157
commit 79ac44834e
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
35 changed files with 7201 additions and 218 deletions

View file

@ -36,6 +36,14 @@ pub struct EntrySummary {
pub cacheable_bytes: i64,
}
/// One stored LLM summary, as exposed over the API.
///
/// Aliased rather than redefined: the DB row is already the exact shape the
/// frontend needs, and a second near-identical struct would only add a mapping
/// step to keep in sync. The `View` name exists because `EntrySummary` in this
/// module is the *entry listing* row, an unrelated thing.
pub use crate::database::EntrySummaryRecord as EntrySummaryView;
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
pub struct EntryDetail {
pub summary: EntrySummary,
@ -43,6 +51,12 @@ pub struct EntryDetail {
pub source_metadata_json: String,
pub display_metadata_json: Option<String>,
pub artifacts: Vec<EntryArtifactSummary>,
/// Most recent completed summary for this entry. Always `None` on a fresh
/// capture — summarization is manual.
pub latest_summary: Option<EntrySummaryView>,
/// Latest non-completed generation attempt, kept separate so a replacement
/// never displaces readable completed content.
pub summary_attempt: Option<EntrySummaryView>,
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
@ -343,12 +357,17 @@ pub fn get_entry_detail(
})?
.collect::<rusqlite::Result<Vec<_>>>()?;
let latest_summary = database::latest_completed_entry_summary(conn, entry_id)?;
let summary_attempt = database::latest_entry_summary_attempt(conn, entry_id)?;
Ok(Some(EntryDetail {
summary,
structured_root_relpath,
source_metadata_json,
display_metadata_json,
artifacts,
latest_summary,
summary_attempt,
}))
}
@ -759,7 +778,14 @@ pub fn search_entries(
sql.push_str(&format!(
" AND (LOWER(e.title) LIKE ?{n} OR LOWER(si.canonical_url) LIKE ?{n} \
OR LOWER(e.entry_uid) LIKE ?{n} OR LOWER(e.source_kind) LIKE ?{n} \
OR LOWER(e.entity_kind) LIKE ?{n} OR LOWER(e.visibility) LIKE ?{n})"
OR LOWER(e.entity_kind) LIKE ?{n} OR LOWER(e.visibility) LIKE ?{n} \
OR LOWER(COALESCE((\
SELECT s.summary_text FROM entry_summaries s \
WHERE s.entry_id = e.id AND s.status = 'completed' \
AND s.summary_text IS NOT NULL \
ORDER BY s.completed_at DESC, s.updated_at DESC, s.id DESC \
LIMIT 1\
), '')) LIKE ?{n})"
));
params.push(term);
}
@ -1483,6 +1509,192 @@ mod tests {
assert_eq!(results.len(), 1);
}
fn entry_id_by_title(conn: &rusqlite::Connection, title: &str) -> i64 {
conn.query_row(
"SELECT id FROM archived_entries WHERE title = ?1",
[title],
|row| row.get(0),
)
.unwrap()
}
fn complete_summary(
conn: &rusqlite::Connection,
entry_id: i64,
cache_key: &str,
summary_text: &str,
) -> String {
let uid = database::upsert_pending_entry_summary(
conn,
entry_id,
"test_provider",
None,
"v1",
cache_key,
)
.unwrap();
database::update_entry_summary_status(conn, &uid, "completed", Some(summary_text), None)
.unwrap();
uid
}
fn set_summary_timestamps(
conn: &rusqlite::Connection,
summary_uid: &str,
timestamp: &str,
) {
conn.execute(
"UPDATE entry_summaries SET completed_at = ?1, updated_at = ?1 WHERE summary_uid = ?2",
rusqlite::params![timestamp, summary_uid],
)
.unwrap();
}
#[test]
fn search_summary_json_tags_match_and_unrelated_text_is_absent() {
let conn = make_test_db_with_entries();
let entry_id = entry_id_by_title(&conn, "Resume Templates");
complete_summary(
&conn,
entry_id,
"tags",
r#"{"tags":["skincare","dermatology"]}"#,
);
let matches = search_entries(
&conn,
&SearchEntriesQuery {
q: Some("skincare".to_string()),
..Default::default()
},
)
.unwrap();
assert_eq!(matches.len(), 1);
assert_eq!(matches[0].title.as_deref(), Some("Resume Templates"));
let unrelated = search_entries(
&conn,
&SearchEntriesQuery {
q: Some("neurology".to_string()),
..Default::default()
},
)
.unwrap();
assert!(unrelated.is_empty());
}
#[test]
fn search_summary_uses_only_the_newest_completed_row() {
let conn = make_test_db_with_entries();
let entry_id = entry_id_by_title(&conn, "Resume Templates");
let older = complete_summary(&conn, entry_id, "older", "legacy-skincare-term");
set_summary_timestamps(&conn, &older, "2026-01-01T00:00:00Z");
let newer = complete_summary(&conn, entry_id, "newer", "current-dermatology-term");
set_summary_timestamps(&conn, &newer, "2026-02-01T00:00:00Z");
let old_matches = search_entries(
&conn,
&SearchEntriesQuery {
q: Some("legacy-skincare-term".to_string()),
..Default::default()
},
)
.unwrap();
assert!(old_matches.is_empty());
let current_matches = search_entries(
&conn,
&SearchEntriesQuery {
q: Some("current-dermatology-term".to_string()),
..Default::default()
},
)
.unwrap();
assert_eq!(current_matches.len(), 1);
}
#[test]
fn search_summary_keeps_latest_completed_when_newer_rows_are_pending_or_failed() {
let conn = make_test_db_with_entries();
let entry_id = entry_id_by_title(&conn, "Resume Templates");
let completed = complete_summary(&conn, entry_id, "completed", "retained-skincare-term");
set_summary_timestamps(&conn, &completed, "2026-01-01T00:00:00Z");
let pending = database::upsert_pending_entry_summary(
&conn,
entry_id,
"test_provider",
None,
"v1",
"pending",
)
.unwrap();
set_summary_timestamps(&conn, &pending, "2026-03-01T00:00:00Z");
let failed = database::upsert_pending_entry_summary(
&conn,
entry_id,
"test_provider",
None,
"v1",
"failed",
)
.unwrap();
database::update_entry_summary_status(&conn, &failed, "failed", None, Some("boom")).unwrap();
set_summary_timestamps(&conn, &failed, "2026-04-01T00:00:00Z");
let matches = search_entries(
&conn,
&SearchEntriesQuery {
q: Some("retained-skincare-term".to_string()),
..Default::default()
},
)
.unwrap();
assert_eq!(matches.len(), 1);
}
#[test]
fn search_summary_preserves_prefix_collection_and_visibility_scope() {
let conn = make_test_db_with_entries();
let entry_id = entry_id_by_title(&conn, "Polymarket tweet");
complete_summary(&conn, entry_id, "scoped", "scoped-skincare-term");
let tag = create_tag(&conn, "/summary-scope").unwrap();
database::assign_entry_to_tag(
&conn,
entry_id,
database::get_tag_by_uid(&conn, &tag.tag_uid).unwrap().unwrap().id,
)
.unwrap();
let collection = database::create_collection(&conn, "Summary scope", "summary-scope", 2, false)
.unwrap();
database::add_entry_to_collection(&conn, collection.id, entry_id, 2).unwrap();
let query = SearchEntriesQuery {
q: Some("scoped-skincare-term".to_string()),
source_kind: Some("x".to_string()),
entity_kind: Some("tweet".to_string()),
url: Some("x.com".to_string()),
title: Some("polymarket".to_string()),
after: Some("2020-01-01T00:00:00Z".to_string()),
before: Some("9999-01-01T00:00:00Z".to_string()),
tag: Some("/summary-scope".to_string()),
caller_bits: 1,
collection_id: Some(collection.id),
};
assert!(search_entries(&conn, &query).unwrap().is_empty());
let matches = search_entries(
&conn,
&SearchEntriesQuery {
caller_bits: 2,
..query
},
)
.unwrap();
assert_eq!(matches.len(), 1);
assert_eq!(matches[0].title.as_deref(), Some("Polymarket tweet"));
}
// ---- tag API tests ----
fn make_tag_test_db() -> (rusqlite::Connection, i64, i64) {

View file

@ -951,7 +951,10 @@ fn register_tweet_artifacts(
})?;
for (role, raw_relpath) in tweet_raw_artifacts(&json_str)? {
let raw_path = PathBuf::from(&raw_relpath);
let blob = blob_record_for_raw_relpath(store_path, &raw_path)?;
let mut blob = blob_record_for_raw_relpath(store_path, &raw_path)?;
if role == "media" {
blob.mime_type = tweet_media_image_mime(blob.extension.as_deref());
}
let blob_id = database::upsert_blob(conn, &blob)?;
database::add_entry_artifact(
conn,
@ -1030,6 +1033,22 @@ fn record_tweet_entry(
Ok(entry)
}
/// Trusted image MIME types emitted by the X downloader's media paths.
///
/// Tweet JSON has no MIME field for ordinary downloaded media. Restricting this
/// inference to the explicit image extensions keeps binary video/audio and
/// unknown extensions out of multimodal summary input.
fn tweet_media_image_mime(extension: Option<&str>) -> Option<String> {
match extension?.to_ascii_lowercase().as_str() {
"jpg" | "jpeg" => Some("image/jpeg".to_string()),
"png" => Some("image/png".to_string()),
"webp" => Some("image/webp".to_string()),
"gif" => Some("image/gif".to_string()),
"avif" => Some("image/avif".to_string()),
_ => None,
}
}
fn tweet_raw_artifacts(tweet_json: &str) -> Result<Vec<(String, String)>> {
let regex = regex::Regex::new(r#""(avatar_local_path|local_path)": "([^"\n]+)""#)?;
let mut seen = HashSet::new();
@ -1562,8 +1581,8 @@ pub fn perform_capture(
let (rewritten, fonts) =
downloader::font_extractor::extract_and_rewrite(&content, store_path, aid)
.unwrap_or_else(|_| (content.clone(), vec![])); // non-fatal
// Extract title after font-stripping so the title tag is not buried
// behind multi-MB embedded font data that would exceed the 256 KiB window.
// Extract title after font-stripping so the title tag is not buried
// behind multi-MB embedded font data that would exceed the 256 KiB window.
let title = downloader::singlefile::extract_html_title_str(&rewritten);
fs::write(&temp_html, rewritten.as_bytes())
.with_context(|| "failed to write rewritten HTML")?;
@ -1918,6 +1937,171 @@ pub fn perform_capture(
})
}
/// Archives user-supplied plain text or Markdown content.
///
/// # Arguments
/// * `archive_paths` - Path configuration for the archive
/// * `title` - User-supplied title (non-empty, trimmed, capped at 500 chars)
/// * `body` - Text content (non-empty, capped at 2 MiB)
/// * `mime` - MIME type: "text/plain" or "text/markdown"
/// * `archive_id` - Optional archive ID (used for job tracking if provided)
///
/// # Returns
/// * `CaptureResult` with the run UID and status
///
/// # Errors
/// * Empty or oversized title/body
/// * Unsupported MIME type
/// * Database or file system errors
pub fn perform_text_capture(
archive_paths: &ArchivePaths,
title: &str,
body: &str,
mime: &str,
_archive_id: Option<&str>,
) -> Result<CaptureResult> {
// Validate title
let title = title.trim();
if title.is_empty() {
anyhow::bail!("title must not be empty");
}
if title.len() > 500 {
anyhow::bail!("title must not exceed 500 characters");
}
// Validate body
if body.trim().is_empty() {
anyhow::bail!("body must not be empty");
}
if body.len() > 2 * 1024 * 1024 {
anyhow::bail!("body must not exceed 2 MiB");
}
// Validate MIME type
if mime != "text/plain" && mime != "text/markdown" {
anyhow::bail!("unsupported MIME type: {mime}. Must be 'text/plain' or 'text/markdown'");
}
// Generate timestamp
let timestamp = format!(
"{}-{}",
Local::now().format("%Y-%m-%dT%H-%M-%S%.3f"),
Uuid::new_v4().simple(),
);
let store_path = &archive_paths.store_path;
// Initialize database
let conn = database::open_or_initialize(&archive_paths.archive_path)?;
let user_id = database::ensure_default_user(&conn)?;
// Create run and item
let run = database::create_archive_run(&conn, user_id, 1)?;
let source_kind = "text";
let entity_kind = "document";
let item = database::create_archive_run_item(
&conn,
run.id,
None,
0,
&format!("text:{}", title),
None,
source_kind,
entity_kind,
)?;
// Stage the text content
let staged_text = downloader::text::save(body.as_bytes(), mime, store_path, &timestamp)?;
// Check if hash already exists
let file_extension = format!(".{}", staged_text.extension);
let hash_exists = hash_exists(&staged_text.hash, &file_extension, store_path)?;
if !hash_exists {
// Move staged file to raw storage
move_temp_to_raw(&staged_text.staged_path, &staged_text.hash, store_path)?;
}
// Clean up temp directory
let _ = fs::remove_dir_all(store_path.join("temp").join(&timestamp));
// Create blob record
let raw_relpath = raw_relative_path_from_hash(&staged_text.hash, &file_extension)?;
let blob = database::BlobRecord {
sha256: staged_text.hash.clone(),
byte_size: staged_text.byte_size as i64,
mime_type: Some(mime.to_string()),
extension: Some(staged_text.extension.clone()),
raw_relpath: path_to_store_string(&raw_relpath),
};
let blob_id = database::upsert_blob(&conn, &blob)?;
// Create source identity
let canonical_locator = format!("text:{}", staged_text.hash);
let source_identity_id = database::upsert_source_identity(
&conn,
source_kind,
entity_kind,
None,
None,
&canonical_locator,
)?;
// Create entry
let entry = database::create_archived_entry(
&conn,
&database::NewEntry {
source_identity_id,
archive_run_id: run.id,
parent_entry_id: None,
root_entry_id: None,
created_by_user_id: user_id,
owned_by_user_id: user_id,
source_kind: source_kind.to_string(),
entity_kind: entity_kind.to_string(),
title: Some(title.to_string()),
visibility: "private".to_string(),
representation_kind: "text".to_string(),
source_metadata_json: json!({
"requested_locator": format!("text:{}", title),
"canonical_locator": canonical_locator,
"mime_type": mime
})
.to_string(),
display_metadata_json: None,
},
)?;
// Create structured root directory
create_structured_root(store_path, &entry)?;
// Create primary_media artifact
database::add_entry_artifact(
&conn,
&database::NewArtifact {
entry_id: entry.id,
artifact_role: "primary_media".to_string(),
storage_area: "raw".to_string(),
relpath: blob.raw_relpath,
blob_id: Some(blob_id),
logical_path: None,
metadata_json: None,
},
)?;
// Complete the run item
database::complete_archive_run_item(&conn, item.id, entry.id)?;
database::refresh_entry_cached_bytes(&conn, entry.id)?;
database::finish_archive_run(&conn, run.id)?;
Ok(CaptureResult {
run_uid: run.run_uid.clone(),
status: "completed".to_string(),
completed_child_count: 0,
ublock_skipped: false,
cookie_ext_skipped: false,
})
}
/// Result of a tweet re-archive operation.
#[derive(Debug, serde::Serialize)]
pub struct RearchiveResult {
@ -2628,6 +2812,325 @@ mod tests {
);
}
#[test]
fn test_text_capture_markdown() {
let base_path = env::temp_dir().join(format!(
"archivr-text-test-{}",
Local::now().format("%Y%m%d%H%M%S%3f")
));
let _ = fs::remove_dir_all(&base_path);
fs::create_dir_all(&base_path).unwrap();
let store_path = base_path.join("store");
let archive_path = base_path.join(".archivr");
// Initialize archive structure
archive::initialize_store_directories(&store_path).unwrap();
fs::create_dir_all(&archive_path).unwrap();
fs::write(archive_path.join("name"), "test-archive").unwrap();
fs::write(
archive_path.join("store_path"),
store_path.to_str().unwrap(),
)
.unwrap();
let archive_paths = ArchivePaths {
archive_path: archive_path.clone(),
store_path: store_path.clone(),
name: "test-archive".to_string(),
};
let title = "My Markdown Note";
let body = "# Heading\n\nSome **bold** text.";
let mime = "text/markdown";
let result = perform_text_capture(&archive_paths, title, body, mime, None).unwrap();
assert_eq!(result.status, "completed");
assert_eq!(result.completed_child_count, 0);
assert!(!result.ublock_skipped);
assert!(!result.cookie_ext_skipped);
// Verify entry was created
let conn = database::open_or_initialize(&archive_path).unwrap();
let default_coll_id = database::ensure_default_collection(&conn).unwrap();
let entries =
archive::list_entries_for_collection(&conn, default_coll_id, 0xFFFFFFFF).unwrap();
assert_eq!(entries.len(), 1);
let entry = &entries[0];
assert_eq!(entry.title, Some(title.to_string()));
assert_eq!(entry.source_kind, "text");
assert_eq!(entry.entity_kind, "document");
// Clean up
let _ = fs::remove_dir_all(&base_path);
}
#[test]
fn test_text_capture_plain() {
let base_path = env::temp_dir().join(format!(
"archivr-text-plain-test-{}",
Local::now().format("%Y%m%d%H%M%S%3f")
));
let _ = fs::remove_dir_all(&base_path);
fs::create_dir_all(&base_path).unwrap();
let store_path = base_path.join("store");
let archive_path = base_path.join(".archivr");
// Initialize archive structure
archive::initialize_store_directories(&store_path).unwrap();
fs::create_dir_all(&archive_path).unwrap();
fs::write(archive_path.join("name"), "test-archive").unwrap();
fs::write(
archive_path.join("store_path"),
store_path.to_str().unwrap(),
)
.unwrap();
let archive_paths = ArchivePaths {
archive_path: archive_path.clone(),
store_path: store_path.clone(),
name: "test-archive".to_string(),
};
let title = "Plain Text Note";
let body = "Just plain text content.";
let mime = "text/plain";
let result = perform_text_capture(&archive_paths, title, body, mime, None).unwrap();
assert_eq!(result.status, "completed");
// Verify entry was created
let conn = database::open_or_initialize(&archive_path).unwrap();
let default_coll_id = database::ensure_default_collection(&conn).unwrap();
let entries =
archive::list_entries_for_collection(&conn, default_coll_id, 0xFFFFFFFF).unwrap();
assert_eq!(entries.len(), 1);
let entry = &entries[0];
assert_eq!(entry.title, Some(title.to_string()));
// Clean up
let _ = fs::remove_dir_all(&base_path);
}
#[test]
fn test_text_capture_preserves_intentional_whitespace_in_stored_artifact() {
let base_path = env::temp_dir().join(format!(
"archivr-text-whitespace-test-{}",
Local::now().format("%Y%m%d%H%M%S%3f")
));
let _ = fs::remove_dir_all(&base_path);
fs::create_dir_all(&base_path).unwrap();
let store_path = base_path.join("store");
let archive_path = base_path.join(".archivr");
archive::initialize_store_directories(&store_path).unwrap();
fs::create_dir_all(&archive_path).unwrap();
fs::write(archive_path.join("name"), "test-archive").unwrap();
fs::write(archive_path.join("store_path"), store_path.to_str().unwrap()).unwrap();
let archive_paths = ArchivePaths {
archive_path: archive_path.clone(),
store_path: store_path.clone(),
name: "test-archive".to_string(),
};
let body = " \n# Heading\n\nContent with a final newline\n\t \n";
perform_text_capture(
&archive_paths,
"Whitespace Note",
body,
"text/markdown",
None,
)
.unwrap();
let conn = database::open_or_initialize(&archive_path).unwrap();
let raw_relpath: String = conn
.query_row(
"SELECT b.raw_relpath
FROM entry_artifacts ea
JOIN blobs b ON b.id = ea.blob_id
WHERE ea.artifact_role = 'primary_media'",
[],
|row| row.get(0),
)
.unwrap();
assert_eq!(fs::read(store_path.join(raw_relpath)).unwrap(), body.as_bytes());
let _ = fs::remove_dir_all(&base_path);
}
#[test]
fn test_text_capture_hides_synthetic_url_but_reuses_source_identity() {
let base_path = env::temp_dir().join(format!(
"archivr-text-identity-test-{}",
Local::now().format("%Y%m%d%H%M%S%3f")
));
let _ = fs::remove_dir_all(&base_path);
fs::create_dir_all(&base_path).unwrap();
let store_path = base_path.join("store");
let archive_path = base_path.join(".archivr");
archive::initialize_store_directories(&store_path).unwrap();
fs::create_dir_all(&archive_path).unwrap();
fs::write(archive_path.join("name"), "test-archive").unwrap();
fs::write(archive_path.join("store_path"), store_path.to_str().unwrap()).unwrap();
let archive_paths = ArchivePaths {
archive_path: archive_path.clone(),
store_path,
name: "test-archive".to_string(),
};
let body = "Same body, same text identity.";
perform_text_capture(&archive_paths, "First title", body, "text/plain", None).unwrap();
perform_text_capture(&archive_paths, "Second title", body, "text/plain", None).unwrap();
let conn = database::open_or_initialize(&archive_path).unwrap();
let default_coll_id = database::ensure_default_collection(&conn).unwrap();
let entries = archive::list_entries_for_collection(&conn, default_coll_id, 0xFFFFFFFF).unwrap();
assert_eq!(entries.len(), 2);
assert!(entries.iter().all(|entry| entry.original_url.is_none()));
let expected_locator = format!("text:{}", crate::hash::hash_bytes(body.as_bytes()));
let (canonical_url, normalized_locator): (Option<String>, String) = conn
.query_row(
"SELECT canonical_url, normalized_locator
FROM source_identities
WHERE source_kind = 'text' AND normalized_locator = ?1",
[&expected_locator],
|row| Ok((row.get(0)?, row.get(1)?)),
)
.unwrap();
assert_eq!(canonical_url, None);
assert_eq!(normalized_locator, expected_locator);
let source_identity_count: i64 = conn
.query_row(
"SELECT COUNT(*) FROM source_identities WHERE source_kind = 'text'",
[],
|row| row.get(0),
)
.unwrap();
assert_eq!(source_identity_count, 1);
let _ = fs::remove_dir_all(&base_path);
}
#[test]
fn test_text_capture_rejects_empty_title() {
let base_path = env::temp_dir().join(format!(
"archivr-text-empty-title-test-{}",
Local::now().format("%Y%m%d%H%M%S%3f")
));
let _ = fs::remove_dir_all(&base_path);
fs::create_dir_all(&base_path).unwrap();
let store_path = base_path.join("store");
let archive_path = base_path.join(".archivr");
archive::initialize_store_directories(&store_path).unwrap();
fs::create_dir_all(&archive_path).unwrap();
fs::write(archive_path.join("name"), "test-archive").unwrap();
fs::write(
archive_path.join("store_path"),
store_path.to_str().unwrap(),
)
.unwrap();
let archive_paths = ArchivePaths {
archive_path,
store_path,
name: "test-archive".to_string(),
};
let result = perform_text_capture(&archive_paths, "", "Some body", "text/plain", None);
assert!(result.is_err());
assert!(result
.unwrap_err()
.to_string()
.contains("title must not be empty"));
// Clean up
let _ = fs::remove_dir_all(&base_path);
}
#[test]
fn test_text_capture_rejects_empty_body() {
let base_path = env::temp_dir().join(format!(
"archivr-text-empty-body-test-{}",
Local::now().format("%Y%m%d%H%M%S%3f")
));
let _ = fs::remove_dir_all(&base_path);
fs::create_dir_all(&base_path).unwrap();
let store_path = base_path.join("store");
let archive_path = base_path.join(".archivr");
archive::initialize_store_directories(&store_path).unwrap();
fs::create_dir_all(&archive_path).unwrap();
fs::write(archive_path.join("name"), "test-archive").unwrap();
fs::write(
archive_path.join("store_path"),
store_path.to_str().unwrap(),
)
.unwrap();
let archive_paths = ArchivePaths {
archive_path,
store_path,
name: "test-archive".to_string(),
};
let result = perform_text_capture(&archive_paths, "Some Title", "", "text/plain", None);
assert!(result.is_err());
assert!(result
.unwrap_err()
.to_string()
.contains("body must not be empty"));
// Clean up
let _ = fs::remove_dir_all(&base_path);
}
#[test]
fn test_text_capture_rejects_bad_mime() {
let base_path = env::temp_dir().join(format!(
"archivr-text-bad-mime-test-{}",
Local::now().format("%Y%m%d%H%M%S%3f")
));
let _ = fs::remove_dir_all(&base_path);
fs::create_dir_all(&base_path).unwrap();
let store_path = base_path.join("store");
let archive_path = base_path.join(".archivr");
archive::initialize_store_directories(&store_path).unwrap();
fs::create_dir_all(&archive_path).unwrap();
fs::write(archive_path.join("name"), "test-archive").unwrap();
fs::write(
archive_path.join("store_path"),
store_path.to_str().unwrap(),
)
.unwrap();
let archive_paths = ArchivePaths {
archive_path,
store_path,
name: "test-archive".to_string(),
};
let result = perform_text_capture(&archive_paths, "Title", "Body", "text/html", None);
assert!(result.is_err());
assert!(result
.unwrap_err()
.to_string()
.contains("unsupported MIME type"));
// Clean up
let _ = fs::remove_dir_all(&base_path);
}
#[test]
fn test_initialize_store_directories() {
let store_path = env::temp_dir().join(format!(
@ -2648,12 +3151,15 @@ mod tests {
#[test]
fn test_record_tweet_entry_links_json_and_raw_artifacts() {
let store_path = env::temp_dir().join(format!(
"archivr-tweet-db-test-{}",
Local::now().format("%Y%m%d%H%M%S%3f")
));
let _ = fs::remove_dir_all(&store_path);
archive::initialize_store_directories(&store_path).unwrap();
let temp = tempfile::tempdir().unwrap();
let archive_paths = archive::initialize_archive(
temp.path(),
&temp.path().join("store"),
"Tweet summary test",
false,
)
.unwrap();
let store_path = &archive_paths.store_path;
fs::create_dir_all(store_path.join("raw").join("a").join("b")).unwrap();
fs::create_dir_all(store_path.join("raw").join("c").join("d")).unwrap();
fs::write(
@ -2670,21 +3176,21 @@ mod tests {
.join("raw")
.join("c")
.join("d")
.join("cdef01.mp4"),
.join("cdef01.jpg"),
b"media",
)
.unwrap();
fs::write(
store_path.join("raw_tweets").join("tweet-123.json"),
r#"{
"full_text": "Tweet body for summary selection.",
"author": { "avatar_local_path": "raw/a/b/abcdef.jpg" },
"entities": { "media": [{ "local_path": "raw/c/d/cdef01.mp4" }] }
"entities": { "media": [{ "local_path": "raw/c/d/cdef01.jpg" }] }
}"#,
)
.unwrap();
let conn = rusqlite::Connection::open_in_memory().unwrap();
database::initialize_schema(&conn).unwrap();
let conn = database::open_or_initialize(&archive_paths.archive_path).unwrap();
let user_id = database::ensure_default_user(&conn).unwrap();
let run = database::create_archive_run(&conn, user_id, 1).unwrap();
let item = database::create_archive_run_item(
@ -2733,10 +3239,26 @@ mod tests {
assert_eq!(artifact_count, 3);
assert_eq!(blob_count, 2);
let media_mime: Option<String> = conn
.query_row(
"SELECT b.mime_type FROM entry_artifacts ea JOIN blobs b ON b.id = ea.blob_id WHERE ea.entry_id = ?1 AND ea.artifact_role = 'media'",
[entry.id],
|row| row.get(0),
)
.unwrap();
assert_eq!(media_mime.as_deref(), Some("image/jpeg"));
let summary_input = crate::summarizer::build_summary_input(
&archive_paths,
&entry.entry_uid,
crate::summarizer::SummaryBuildOptions {
include_images: true,
},
)
.unwrap();
assert_eq!(summary_input.request.images.len(), 1);
assert_eq!(summary_input.request.images[0].mime_type, "image/jpeg");
assert_eq!(run_status, "completed");
assert!(store_path.join(&entry.structured_root_relpath).is_dir());
let _ = fs::remove_dir_all(store_path);
}
mod title_tests {

View file

@ -110,6 +110,30 @@ pub struct CaptureJobRecord {
pub updated_at: String,
}
/// One row of `entry_summaries` — a regenerable LLM summary of an entry.
///
/// `provider_model` is stored as `''` (not NULL) when a provider has no explicit
/// model, because SQLite treats NULLs as distinct inside a UNIQUE index and a
/// NULL model would defeat the `(entry_id, provider_kind, provider_model,
/// prompt_version, input_sha256)` dedupe key. Readers map `''` back to `None`.
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
pub struct EntrySummaryRecord {
pub summary_uid: String,
pub entry_uid: String,
pub provider_kind: String,
/// Model selected by the provider after resolving the requested cache-key alias.
pub resolved_model: Option<String>,
pub provider_model: Option<String>,
pub prompt_version: String,
pub input_sha256: String,
pub status: String,
pub summary_text: Option<String>,
pub error_text: Option<String>,
pub created_at: String,
pub updated_at: String,
pub completed_at: Option<String>,
}
#[derive(Debug, Clone, serde::Serialize)]
pub struct UserSummary {
pub user_uid: String,
@ -322,6 +346,26 @@ pub fn initialize_schema(conn: &Connection) -> Result<()> {
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS entry_summaries (
id INTEGER PRIMARY KEY,
summary_uid TEXT NOT NULL UNIQUE,
entry_id INTEGER NOT NULL REFERENCES archived_entries(id) ON DELETE CASCADE,
provider_kind TEXT NOT NULL,
provider_model TEXT NOT NULL DEFAULT '',
resolved_model TEXT,
prompt_version TEXT NOT NULL,
input_sha256 TEXT NOT NULL,
status TEXT NOT NULL CHECK(status IN ('pending','running','completed','failed')),
summary_text TEXT,
error_text TEXT,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL,
completed_at TEXT
);
CREATE INDEX IF NOT EXISTS idx_entry_summaries_entry_updated
ON entry_summaries(entry_id, updated_at DESC);
CREATE INDEX IF NOT EXISTS idx_entry_summaries_cache_lookup
ON entry_summaries(entry_id, provider_kind, provider_model, prompt_version, input_sha256, status);
CREATE INDEX IF NOT EXISTS idx_capture_jobs_status ON capture_jobs(status);
CREATE INDEX IF NOT EXISTS idx_archive_run_items_run_id ON archive_run_items(run_id);
CREATE INDEX IF NOT EXISTS idx_archived_entries_source_identity_id ON archived_entries(source_identity_id);
@ -436,6 +480,12 @@ pub fn initialize_schema(conn: &Connection) -> Result<()> {
// Migration: add notes_json column to existing capture_jobs tables.
// Silently ignored when the column already exists (idempotent).
let _ = conn.execute("ALTER TABLE capture_jobs ADD COLUMN notes_json TEXT", []);
// Provider responses may resolve a requested alias to a concrete model.
// Keep that display-only value outside the cache key.
let _ = conn.execute(
"ALTER TABLE entry_summaries ADD COLUMN resolved_model TEXT",
[],
);
// Migration: add requires_auth column to existing collections tables.
// Silently ignored when the column already exists (idempotent).
let _ = conn.execute(
@ -443,6 +493,54 @@ pub fn initialize_schema(conn: &Connection) -> Result<()> {
[],
);
// Summary attempts used to be unique by cache key, which meant forced
// regeneration erased the last completed result. Rebuild that small table
// without the cache-key constraint while retaining all existing rows.
let summary_table_sql: Option<String> = conn
.query_row(
"SELECT sql FROM sqlite_master WHERE type = 'table' AND name = 'entry_summaries'",
[],
|row| row.get(0),
)
.optional()?;
if summary_table_sql.as_deref().is_some_and(|sql| {
sql.contains(
"UNIQUE(entry_id, provider_kind, provider_model, prompt_version, input_sha256)",
)
}) {
conn.execute_batch(
"BEGIN;
CREATE TABLE entry_summaries_rebuilt (
id INTEGER PRIMARY KEY,
summary_uid TEXT NOT NULL UNIQUE,
entry_id INTEGER NOT NULL REFERENCES archived_entries(id) ON DELETE CASCADE,
provider_kind TEXT NOT NULL,
provider_model TEXT NOT NULL DEFAULT '',
resolved_model TEXT,
prompt_version TEXT NOT NULL,
input_sha256 TEXT NOT NULL,
status TEXT NOT NULL CHECK(status IN ('pending','running','completed','failed')),
summary_text TEXT,
error_text TEXT,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL,
completed_at TEXT
);
INSERT INTO entry_summaries_rebuilt
SELECT id, summary_uid, entry_id, provider_kind, COALESCE(provider_model, ''),
resolved_model, prompt_version, input_sha256, status, summary_text, error_text,
created_at, updated_at, completed_at
FROM entry_summaries;
DROP TABLE entry_summaries;
ALTER TABLE entry_summaries_rebuilt RENAME TO entry_summaries;
CREATE INDEX idx_entry_summaries_entry_updated
ON entry_summaries(entry_id, updated_at DESC);
CREATE INDEX idx_entry_summaries_cache_lookup
ON entry_summaries(entry_id, provider_kind, provider_model, prompt_version, input_sha256, status);
COMMIT;",
)?;
}
Ok(())
}
@ -1354,6 +1452,265 @@ pub fn fail_stalled_capture_jobs(conn: &Connection) -> Result<usize> {
Ok(n)
}
/// Marks summary attempts left pending or running by a previous process as
/// failed. The message is deliberately stable and does not expose internals.
pub fn fail_stalled_entry_summaries(conn: &Connection) -> Result<usize> {
let now = now_timestamp();
conn.execute(
"UPDATE entry_summaries
SET status = 'failed',
error_text = 'Summary generation was interrupted by a server restart.',
updated_at = ?1
WHERE status IN ('pending', 'running')",
[now],
)
.map_err(Into::into)
}
// ── entry_summaries ────────────────────────────────────────────────────────
//
// Summaries are a regenerable child record of an entry, never a column on
// `archived_entries`: an entry can carry several (one per provider/model/prompt
// version), and any of them can be thrown away and recomputed. Generation is
// manual-only — nothing in `capture.rs` writes here.
/// `SELECT` list shared by every `entry_summaries` read, so all readers build
/// an identical `EntrySummaryRecord` from the same column ordering.
const ENTRY_SUMMARY_COLS: &str =
"SELECT s.summary_uid, e.entry_uid, s.provider_kind, s.provider_model, s.resolved_model,
s.prompt_version, s.input_sha256, s.status, s.summary_text,
s.error_text, s.created_at, s.updated_at, s.completed_at
FROM entry_summaries s
JOIN archived_entries e ON e.id = s.entry_id";
fn map_entry_summary(row: &rusqlite::Row<'_>) -> rusqlite::Result<EntrySummaryRecord> {
Ok(EntrySummaryRecord {
summary_uid: row.get(0)?,
entry_uid: row.get(1)?,
provider_kind: row.get(2)?,
// '' is the stored stand-in for "this provider has no model"; see the
// doc comment on EntrySummaryRecord for why it is not NULL.
provider_model: row.get::<_, Option<String>>(3)?.filter(|m| !m.is_empty()),
resolved_model: row.get(4)?,
prompt_version: row.get(5)?,
input_sha256: row.get(6)?,
status: row.get(7)?,
summary_text: row.get(8)?,
error_text: row.get(9)?,
created_at: row.get(10)?,
updated_at: row.get(11)?,
completed_at: row.get(12)?,
})
}
/// Resolves `entry_uid` to its integer primary key. `Ok(None)` if absent.
pub fn entry_id_for_uid(conn: &Connection, entry_uid: &str) -> Result<Option<i64>> {
conn.query_row(
"SELECT id FROM archived_entries WHERE entry_uid = ?1",
[entry_uid],
|row| row.get(0),
)
.optional()
.map_err(Into::into)
}
/// Creates a fresh pending summary attempt for one cache key.
///
/// Attempts are intentionally not unique by cache key: a forced regeneration
/// must leave an older completed result available while the new attempt runs.
pub fn upsert_pending_entry_summary(
conn: &Connection,
entry_id: i64,
provider_kind: &str,
provider_model: Option<&str>,
prompt_version: &str,
input_sha256: &str,
) -> Result<String> {
let summary_uid = format!("sum_{}", &Uuid::new_v4().simple().to_string()[..10]);
let now = now_timestamp();
let model = provider_model.unwrap_or("");
conn.execute(
"INSERT INTO entry_summaries
(summary_uid, entry_id, provider_kind, provider_model, prompt_version,
input_sha256, status, summary_text, error_text, created_at, updated_at, completed_at)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, 'pending', NULL, NULL, ?7, ?7, NULL)",
rusqlite::params![
summary_uid,
entry_id,
provider_kind,
model,
prompt_version,
input_sha256,
now
],
)?;
Ok(summary_uid)
}
/// Completes a summary while retaining the provider's concrete response model
/// for display. The requested provider model remains the cache-key identity.
pub fn update_entry_summary_completed(
conn: &Connection,
summary_uid: &str,
summary_text: &str,
resolved_model: Option<&str>,
) -> Result<()> {
let now = now_timestamp();
conn.execute(
"UPDATE entry_summaries
SET status = 'completed', summary_text = ?1, error_text = NULL,
resolved_model = ?2, completed_at = ?3, updated_at = ?3
WHERE summary_uid = ?4",
rusqlite::params![summary_text, resolved_model, now, summary_uid],
)?;
Ok(())
}
/// Moves a summary row through `running` → `completed` / `failed`.
/// `completed_at` is stamped only on the terminal `completed` transition.
pub fn update_entry_summary_status(
conn: &Connection,
summary_uid: &str,
status: &str,
summary_text: Option<&str>,
error_text: Option<&str>,
) -> Result<()> {
let now = now_timestamp();
let completed_at = if status == "completed" {
Some(now.clone())
} else {
None
};
conn.execute(
"UPDATE entry_summaries
SET status = ?1, summary_text = ?2, error_text = ?3,
completed_at = ?4, updated_at = ?5
WHERE summary_uid = ?6",
rusqlite::params![
status,
summary_text,
error_text,
completed_at,
now,
summary_uid
],
)?;
Ok(())
}
/// Returns one summary by its public uid.
pub fn get_entry_summary_by_uid(
conn: &Connection,
summary_uid: &str,
) -> Result<Option<EntrySummaryRecord>> {
conn.query_row(
&format!("{ENTRY_SUMMARY_COLS} WHERE s.summary_uid = ?1"),
[summary_uid],
map_entry_summary,
)
.optional()
.map_err(Into::into)
}
/// Looks up the row for one exact cache key — used to short-circuit a POST when
/// a completed summary for identical input already exists and `force` is false.
pub fn find_entry_summary(
conn: &Connection,
entry_id: i64,
provider_kind: &str,
provider_model: Option<&str>,
prompt_version: &str,
input_sha256: &str,
) -> Result<Option<EntrySummaryRecord>> {
conn.query_row(
&format!(
"{ENTRY_SUMMARY_COLS}
WHERE s.entry_id = ?1 AND s.provider_kind = ?2 AND s.provider_model = ?3
AND s.prompt_version = ?4 AND s.input_sha256 = ?5
ORDER BY CASE WHEN s.status = 'completed' THEN 0 ELSE 1 END,
s.completed_at DESC, s.updated_at DESC, s.id DESC"
),
rusqlite::params![
entry_id,
provider_kind,
provider_model.unwrap_or(""),
prompt_version,
input_sha256
],
map_entry_summary,
)
.optional()
.map_err(Into::into)
}
/// Most recently touched summary for an entry, whatever its status.
/// Used where a caller explicitly needs the most recent attempt regardless of
/// whether it has completed.
pub fn latest_entry_summary(
conn: &Connection,
entry_id: i64,
) -> Result<Option<EntrySummaryRecord>> {
conn.query_row(
&format!(
"{ENTRY_SUMMARY_COLS} WHERE s.entry_id = ?1
ORDER BY s.updated_at DESC, s.id DESC LIMIT 1"
),
[entry_id],
map_entry_summary,
)
.optional()
.map_err(Into::into)
}
/// The most recent completed summary, excluding in-flight and failed attempts.
/// This is the stable summary shown in entry detail and used by free-text search.
pub fn latest_completed_entry_summary(
conn: &Connection,
entry_id: i64,
) -> Result<Option<EntrySummaryRecord>> {
conn.query_row(
&format!(
"{ENTRY_SUMMARY_COLS} WHERE s.entry_id = ?1 AND s.status = 'completed'
ORDER BY s.completed_at DESC, s.updated_at DESC, s.id DESC LIMIT 1"
),
[entry_id],
map_entry_summary,
)
.optional()
.map_err(Into::into)
}
/// The newest non-completed replacement attempt after the retained completed
/// result. This includes failed rows so authenticated callers can show a recent
/// failure beside readable content, but suppresses historical failures after a
/// newer successful regeneration.
pub fn latest_entry_summary_attempt(
conn: &Connection,
entry_id: i64,
) -> Result<Option<EntrySummaryRecord>> {
conn.query_row(
&format!(
"{ENTRY_SUMMARY_COLS} WHERE s.entry_id = ?1 AND s.status != 'completed'
AND (
NOT EXISTS (
SELECT 1 FROM entry_summaries c
WHERE c.entry_id = s.entry_id AND c.status = 'completed'
)
OR (s.updated_at, s.id) > (
SELECT c.updated_at, c.id FROM entry_summaries c
WHERE c.entry_id = s.entry_id AND c.status = 'completed'
ORDER BY c.completed_at DESC, c.updated_at DESC, c.id DESC LIMIT 1
)
)
ORDER BY s.updated_at DESC, s.id DESC LIMIT 1"
),
[entry_id],
map_entry_summary,
)
.optional()
.map_err(Into::into)
}
pub fn create_archive_run(
conn: &Connection,
created_by_user_id: i64,
@ -1445,7 +1802,11 @@ pub fn finish_archive_run(conn: &Connection, run_id: i64) -> Result<()> {
[run_id],
|row| row.get(0),
)?;
let status = if failed_count > 0 { "failed" } else { "completed" };
let status = if failed_count > 0 {
"failed"
} else {
"completed"
};
conn.execute(
"UPDATE archive_runs SET status = ?1, finished_at = ?2 WHERE id = ?3",
params![status, now_timestamp(), run_id],
@ -1681,9 +2042,8 @@ pub fn delete_entry(conn: &Connection, entry_uid: &str) -> Result<bool> {
// (no grandchildren), so without `id = ?1` the set would be empty and
// cascade_cached_bytes_after_subtree_delete would not recalculate shared-blob totals.
let subtree_ids: Vec<i64> = {
let mut stmt = conn.prepare(
"SELECT id FROM archived_entries WHERE id = ?1 OR root_entry_id = ?1",
)?;
let mut stmt =
conn.prepare("SELECT id FROM archived_entries WHERE id = ?1 OR root_entry_id = ?1")?;
stmt.query_map([entry_id], |row| row.get(0))?
.collect::<rusqlite::Result<_>>()?
};
@ -2334,7 +2694,6 @@ pub fn visibility_to_bits(visibility: &str) -> u32 {
}
}
/// Returns the id of the '_default_' collection, creating it if absent.
pub fn ensure_default_collection(conn: &Connection) -> Result<i64> {
let now = now_timestamp();
@ -2437,16 +2796,20 @@ pub fn get_collection_by_slug(conn: &Connection, slug: &str) -> Result<Option<Co
"SELECT id, collection_uid, name, slug, default_visibility_bits, created_at, requires_auth \
FROM collections WHERE slug = ?1",
[slug],
|row| Ok(CollectionRecord {
id: row.get(0)?,
collection_uid: row.get(1)?,
name: row.get(2)?,
slug: row.get(3)?,
default_visibility_bits: row.get::<_, i64>(4)? as u32,
created_at: row.get(5)?,
requires_auth: row.get::<_, i64>(6)? != 0,
}),
).optional().map_err(Into::into)
|row| {
Ok(CollectionRecord {
id: row.get(0)?,
collection_uid: row.get(1)?,
name: row.get(2)?,
slug: row.get(3)?,
default_visibility_bits: row.get::<_, i64>(4)? as u32,
created_at: row.get(5)?,
requires_auth: row.get::<_, i64>(6)? != 0,
})
},
)
.optional()
.map_err(Into::into)
}
/// Adds an entry to a collection with given visibility_bits. Idempotent (INSERT OR IGNORE).
@ -2505,7 +2868,12 @@ pub fn get_entry_collection_memberships(
WHERE ce.entry_id = ?1",
)?;
stmt.query_map([entry_id], |row| {
Ok((row.get(0)?, row.get(1)?, row.get(2)?, row.get::<_, i64>(3)? as u32))
Ok((
row.get(0)?,
row.get(1)?,
row.get(2)?,
row.get::<_, i64>(3)? as u32,
))
})?
.collect::<Result<_, _>>()
.map_err(Into::into)
@ -3063,7 +3431,10 @@ mod tests {
|row| row.get::<_, i64>(0).map(|v| v as u32),
)
.unwrap();
assert_eq!(default_bits, 2, "default collection should start USER-visible");
assert_eq!(
default_bits, 2,
"default collection should start USER-visible"
);
// Create an entry with visibility = "private" (what capture always passes).
let entry = create_entry_fixture(&conn, "private", None, None);
@ -3111,7 +3482,10 @@ mod tests {
|row| row.get::<_, i64>(0).map(|v| v as u32),
)
.unwrap();
assert_eq!(public_bits, 3, "collection default=public should produce bits=3");
assert_eq!(
public_bits, 3,
"collection default=public should produce bits=3"
);
// Child entries must NOT use the collection default — they keep
// visibility_to_bits(entry.visibility) so parent-child visibility
@ -3124,7 +3498,10 @@ mod tests {
|row| row.get::<_, i64>(0).map(|v| v as u32),
)
.unwrap();
assert_eq!(child_bits, 0, "child entries must use visibility_to_bits, not collection default");
assert_eq!(
child_bits, 0,
"child entries must use visibility_to_bits, not collection default"
);
}
#[test]
@ -4153,27 +4530,414 @@ mod tests {
// Root item (parent_item_id IS NULL) — mirrors what record_container_entry does.
let root_item = create_archive_run_item(
&c, run.id, None, 0, "https://example.com/pl", None, "youtube", "playlist",
).unwrap();
&c,
run.id,
None,
0,
"https://example.com/pl",
None,
"youtube",
"playlist",
)
.unwrap();
// Child item (parent_item_id IS NOT NULL).
let child_item = create_archive_run_item(
&c, run.id, Some(root_item.id), 1, "https://example.com/v1", None, "youtube", "video",
).unwrap();
&c,
run.id,
Some(root_item.id),
1,
"https://example.com/v1",
None,
"youtube",
"video",
)
.unwrap();
// Complete both — marks archive_runs.completed_count = 2.
c.execute(
"UPDATE archive_run_items SET status = 'completed' WHERE id IN (?1, ?2)",
rusqlite::params![root_item.id, child_item.id],
).unwrap();
)
.unwrap();
refresh_run_counters(&c, run.id).unwrap();
let total: i64 = c.query_row(
"SELECT completed_count FROM archive_runs WHERE id = ?1", [run.id], |r| r.get(0),
).unwrap();
let total: i64 = c
.query_row(
"SELECT completed_count FROM archive_runs WHERE id = ?1",
[run.id],
|r| r.get(0),
)
.unwrap();
assert_eq!(total, 2, "both items completed: DB counter must be 2");
let child_count = get_run_completed_child_count(&c, run.id).unwrap();
assert_eq!(child_count, 1, "only the child item must be counted");
}
// ── entry_summaries ────────────────────────────────────────────────────
#[test]
fn initialize_schema_is_idempotent_for_entry_summaries() {
// initialize_schema runs on *every* open_or_initialize, so re-running it
// against a populated DB must be a no-op, not an error.
let c = conn();
let entry = create_entry_fixture(&c, "private", None, None);
upsert_pending_entry_summary(&c, entry.id, "claude_cli", None, "v1", "abc").unwrap();
initialize_schema(&c).unwrap();
initialize_schema(&c).unwrap();
let exists: i64 = c
.query_row(
"SELECT COUNT(*) FROM sqlite_master WHERE type='table' AND name='entry_summaries'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(exists, 1);
assert!(latest_entry_summary(&c, entry.id).unwrap().is_some());
}
#[test]
fn initialize_schema_migrates_resolved_model_for_existing_summary_table() {
let c = Connection::open_in_memory().unwrap();
c.execute_batch(
"CREATE TABLE entry_summaries (
id INTEGER PRIMARY KEY,
summary_uid TEXT NOT NULL UNIQUE,
entry_id INTEGER NOT NULL,
provider_kind TEXT NOT NULL,
provider_model TEXT,
prompt_version TEXT NOT NULL,
input_sha256 TEXT NOT NULL,
status TEXT NOT NULL,
summary_text TEXT,
error_text TEXT,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL,
completed_at TEXT,
UNIQUE(entry_id, provider_kind, provider_model, prompt_version, input_sha256)
);",
)
.unwrap();
initialize_schema(&c).unwrap();
let has_resolved_model: i64 = c
.query_row(
"SELECT COUNT(*) FROM pragma_table_info('entry_summaries') WHERE name = 'resolved_model'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(has_resolved_model, 1);
}
#[test]
fn entry_summary_lifecycle_pending_running_completed() {
let c = conn();
let entry = create_entry_fixture(&c, "private", None, None);
let uid =
upsert_pending_entry_summary(&c, entry.id, "anthropic_http", Some("m1"), "v1", "sha1")
.unwrap();
assert!(uid.starts_with("sum_"), "got {uid}");
assert_eq!(uid.len(), "sum_".len() + 10);
let rec = get_entry_summary_by_uid(&c, &uid).unwrap().unwrap();
assert_eq!(rec.status, "pending");
assert_eq!(rec.entry_uid, entry.entry_uid);
assert_eq!(rec.provider_model.as_deref(), Some("m1"));
assert!(rec.completed_at.is_none());
update_entry_summary_status(&c, &uid, "running", None, None).unwrap();
assert_eq!(
get_entry_summary_by_uid(&c, &uid).unwrap().unwrap().status,
"running"
);
update_entry_summary_status(&c, &uid, "completed", Some("{\"summary\":\"s\"}"), None)
.unwrap();
let rec = get_entry_summary_by_uid(&c, &uid).unwrap().unwrap();
assert_eq!(rec.status, "completed");
assert_eq!(rec.summary_text.as_deref(), Some("{\"summary\":\"s\"}"));
assert!(rec.completed_at.is_some(), "completed rows must be stamped");
}
#[test]
fn entry_summary_keeps_requested_model_for_cache_and_resolved_model_for_display() {
let c = conn();
let entry = create_entry_fixture(&c, "private", None, None);
let uid = upsert_pending_entry_summary(
&c,
entry.id,
"anthropic_http",
Some("claude-3-5-sonnet-latest"),
"v1",
"sha1",
)
.unwrap();
update_entry_summary_completed(&c, &uid, "summary", Some("claude-3-5-sonnet-20241022"))
.unwrap();
let rec = get_entry_summary_by_uid(&c, &uid).unwrap().unwrap();
assert_eq!(
rec.provider_model.as_deref(),
Some("claude-3-5-sonnet-latest")
);
assert_eq!(
rec.resolved_model.as_deref(),
Some("claude-3-5-sonnet-20241022")
);
assert_eq!(
find_entry_summary(
&c,
entry.id,
"anthropic_http",
Some("claude-3-5-sonnet-latest"),
"v1",
"sha1",
)
.unwrap()
.unwrap()
.summary_uid,
uid,
);
}
#[test]
fn entry_summary_failure_records_error_and_no_completed_at() {
let c = conn();
let entry = create_entry_fixture(&c, "private", None, None);
let uid = upsert_pending_entry_summary(&c, entry.id, "codex_cli", None, "v1", "s").unwrap();
update_entry_summary_status(&c, &uid, "failed", None, Some("boom")).unwrap();
let rec = get_entry_summary_by_uid(&c, &uid).unwrap().unwrap();
assert_eq!(rec.status, "failed");
assert_eq!(rec.error_text.as_deref(), Some("boom"));
assert!(rec.completed_at.is_none());
}
#[test]
fn regenerating_a_completed_summary_creates_a_new_attempt_and_preserves_completion() {
let c = conn();
let entry = create_entry_fixture(&c, "private", None, None);
let first =
upsert_pending_entry_summary(&c, entry.id, "claude_cli", None, "v1", "sha").unwrap();
update_entry_summary_status(&c, &first, "completed", Some("old"), None).unwrap();
// Regeneration must retain the finished result while a new attempt runs.
let second =
upsert_pending_entry_summary(&c, entry.id, "claude_cli", None, "v1", "sha").unwrap();
assert_ne!(first, second, "regeneration needs a distinct attempt uid");
let n: i64 = c
.query_row("SELECT COUNT(*) FROM entry_summaries", [], |r| r.get(0))
.unwrap();
assert_eq!(n, 2);
let completed = get_entry_summary_by_uid(&c, &first).unwrap().unwrap();
assert_eq!(completed.status, "completed");
assert_eq!(completed.summary_text.as_deref(), Some("old"));
let rec = get_entry_summary_by_uid(&c, &second).unwrap().unwrap();
assert_eq!(rec.status, "pending");
update_entry_summary_status(&c, &second, "failed", None, Some("boom")).unwrap();
assert_eq!(
latest_completed_entry_summary(&c, entry.id)
.unwrap()
.unwrap()
.summary_uid,
first,
"a failed regeneration must not replace the previous completed result"
);
}
#[test]
fn latest_summary_attempt_is_separate_from_the_retained_completed_summary() {
let c = conn();
let entry = create_entry_fixture(&c, "private", None, None);
let completed =
upsert_pending_entry_summary(&c, entry.id, "claude_cli", None, "v1", "old").unwrap();
update_entry_summary_status(&c, &completed, "completed", Some("previous"), None).unwrap();
let pending =
upsert_pending_entry_summary(&c, entry.id, "claude_cli", None, "v1", "new").unwrap();
assert_eq!(
latest_completed_entry_summary(&c, entry.id).unwrap().unwrap().summary_uid,
completed
);
assert_eq!(
latest_entry_summary_attempt(&c, entry.id).unwrap().unwrap().summary_uid,
pending
);
update_entry_summary_status(&c, &pending, "failed", None, Some("boom")).unwrap();
assert_eq!(
latest_completed_entry_summary(&c, entry.id).unwrap().unwrap().summary_text.as_deref(),
Some("previous")
);
assert_eq!(
latest_entry_summary_attempt(&c, entry.id).unwrap().unwrap().status,
"failed"
);
let successful_replacement =
upsert_pending_entry_summary(&c, entry.id, "claude_cli", None, "v1", "newer").unwrap();
update_entry_summary_status(
&c,
&successful_replacement,
"completed",
Some("replacement"),
None,
)
.unwrap();
assert_eq!(
latest_completed_entry_summary(&c, entry.id)
.unwrap()
.unwrap()
.summary_uid,
successful_replacement
);
assert!(
latest_entry_summary_attempt(&c, entry.id).unwrap().is_none(),
"a failed attempt predating a successful replacement must not remain visible"
);
}
#[test]
fn fail_stalled_entry_summaries_marks_pending_and_running_with_restart_message() {
let c = conn();
let entry = create_entry_fixture(&c, "private", None, None);
let pending =
upsert_pending_entry_summary(&c, entry.id, "claude_cli", None, "v1", "a").unwrap();
let running =
upsert_pending_entry_summary(&c, entry.id, "codex_cli", None, "v1", "b").unwrap();
let completed =
upsert_pending_entry_summary(&c, entry.id, "anthropic_http", None, "v1", "c").unwrap();
update_entry_summary_status(&c, &running, "running", None, None).unwrap();
update_entry_summary_status(&c, &completed, "completed", Some("done"), None).unwrap();
assert_eq!(fail_stalled_entry_summaries(&c).unwrap(), 2);
for uid in [&pending, &running] {
let row = get_entry_summary_by_uid(&c, uid).unwrap().unwrap();
assert_eq!(row.status, "failed");
assert_eq!(
row.error_text.as_deref(),
Some("Summary generation was interrupted by a server restart.")
);
assert!(!row.updated_at.is_empty());
}
assert_eq!(
get_entry_summary_by_uid(&c, &completed)
.unwrap()
.unwrap()
.status,
"completed"
);
}
#[test]
fn entry_summaries_with_no_model_preserve_none_across_attempts() {
// CLI providers have no explicit model, so the stored empty-string
// sentinel must always map back to None even when attempts accumulate.
let c = conn();
let entry = create_entry_fixture(&c, "private", None, None);
upsert_pending_entry_summary(&c, entry.id, "claude_cli", None, "v1", "sha").unwrap();
upsert_pending_entry_summary(&c, entry.id, "claude_cli", None, "v1", "sha").unwrap();
let n: i64 = c
.query_row("SELECT COUNT(*) FROM entry_summaries", [], |r| r.get(0))
.unwrap();
assert_eq!(n, 2);
// …and it reads back as None, not as an empty-string model.
let rec = latest_entry_summary(&c, entry.id).unwrap().unwrap();
assert_eq!(rec.provider_model, None);
}
#[test]
fn different_cache_keys_produce_separate_summaries() {
let c = conn();
let entry = create_entry_fixture(&c, "private", None, None);
upsert_pending_entry_summary(&c, entry.id, "claude_cli", None, "v1", "sha").unwrap();
upsert_pending_entry_summary(&c, entry.id, "codex_cli", None, "v1", "sha").unwrap();
upsert_pending_entry_summary(&c, entry.id, "claude_cli", None, "v2", "sha").unwrap();
upsert_pending_entry_summary(&c, entry.id, "claude_cli", None, "v1", "other").unwrap();
let n: i64 = c
.query_row("SELECT COUNT(*) FROM entry_summaries", [], |r| r.get(0))
.unwrap();
assert_eq!(n, 4);
}
#[test]
fn find_entry_summary_matches_only_the_exact_cache_key() {
let c = conn();
let entry = create_entry_fixture(&c, "private", None, None);
let uid =
upsert_pending_entry_summary(&c, entry.id, "openai_compatible", Some("gpt"), "v1", "s")
.unwrap();
let hit = find_entry_summary(&c, entry.id, "openai_compatible", Some("gpt"), "v1", "s")
.unwrap()
.unwrap();
assert_eq!(hit.summary_uid, uid);
assert!(
find_entry_summary(
&c,
entry.id,
"openai_compatible",
Some("gpt"),
"v1",
"changed"
)
.unwrap()
.is_none(),
"a changed input digest must miss the cache"
);
}
#[test]
fn latest_entry_summary_returns_the_most_recently_updated_row() {
let c = conn();
let entry = create_entry_fixture(&c, "private", None, None);
let a = upsert_pending_entry_summary(&c, entry.id, "claude_cli", None, "v1", "a").unwrap();
let b = upsert_pending_entry_summary(&c, entry.id, "codex_cli", None, "v1", "b").unwrap();
update_entry_summary_status(&c, &a, "completed", Some("first"), None).unwrap();
update_entry_summary_status(&c, &b, "completed", Some("second"), None).unwrap();
// Same-timestamp ties break on id DESC, so the later insert wins either way.
assert_eq!(
latest_entry_summary(&c, entry.id)
.unwrap()
.unwrap()
.summary_uid,
b
);
}
#[test]
fn latest_entry_summary_is_none_for_an_unsummarized_entry() {
let c = conn();
let entry = create_entry_fixture(&c, "private", None, None);
assert!(latest_entry_summary(&c, entry.id).unwrap().is_none());
}
#[test]
fn deleting_an_entry_cascades_to_its_summaries() {
let c = conn();
c.pragma_update(None, "foreign_keys", "ON").unwrap();
let entry = create_entry_fixture(&c, "private", None, None);
upsert_pending_entry_summary(&c, entry.id, "claude_cli", None, "v1", "a").unwrap();
assert!(delete_entry(&c, &entry.entry_uid).unwrap());
let n: i64 = c
.query_row("SELECT COUNT(*) FROM entry_summaries", [], |r| r.get(0))
.unwrap();
assert_eq!(n, 0);
}
#[test]
fn entry_id_for_uid_resolves_and_misses() {
let c = conn();
let entry = create_entry_fixture(&c, "private", None, None);
assert_eq!(
entry_id_for_uid(&c, &entry.entry_uid).unwrap(),
Some(entry.id)
);
assert_eq!(entry_id_for_uid(&c, "ent_nope").unwrap(), None);
}
}

View file

@ -7,3 +7,4 @@ pub mod metadata;
pub mod http;
pub mod singlefile;
pub mod font_extractor;
pub mod text;

View file

@ -0,0 +1,118 @@
use anyhow::{bail, Result};
use std::path::{Path, PathBuf};
use crate::hash::hash_bytes;
/// Represents a staged text file ready to be moved into the raw store.
#[derive(Debug)]
pub struct StagedText {
pub staged_path: PathBuf,
pub hash: String,
pub extension: String,
pub byte_size: u64,
}
/// Stages a text body (plain or Markdown) in the temp directory and computes its hash.
///
/// # Arguments
/// * `body` - The raw bytes of the text content
/// * `mime` - MIME type, must be "text/plain" or "text/markdown"
/// * `store_path` - Root store path where temp/ subdirectory will be created
/// * `timestamp` - Timestamp string used in the staged file name
///
/// # Returns
/// * `StagedText` with the staged path, hash, extension, and byte size
///
/// # Errors
/// * Rejects MIME types other than "text/plain" or "text/markdown"
/// * IO errors during directory creation or file writing
pub fn save(body: &[u8], mime: &str, store_path: &Path, timestamp: &str) -> Result<StagedText> {
// Validate MIME type
let extension = match mime {
"text/markdown" => ".md",
"text/plain" => ".txt",
_ => bail!("unsupported MIME type: {mime}. Must be 'text/plain' or 'text/markdown'"),
};
// Create temp directory
let temp_dir = store_path.join("temp").join(timestamp);
std::fs::create_dir_all(&temp_dir)?;
// Stage under temp/<timestamp>/<timestamp><ext>
let staged_path = temp_dir.join(format!("{timestamp}{extension}"));
// Write the content
std::fs::write(&staged_path, body)?;
// Compute SHA3 hash
let hash = hash_bytes(body);
let byte_size = body.len() as u64;
let extension_str = extension.trim_start_matches('.').to_string();
Ok(StagedText {
staged_path,
hash,
extension: extension_str,
byte_size,
})
}
#[cfg(test)]
mod tests {
use super::*;
use tempfile::TempDir;
#[test]
fn test_save_markdown() {
let temp_dir = TempDir::new().unwrap();
let store_path = temp_dir.path();
let content = b"# Hello\n\nThis is markdown.";
let mime = "text/markdown";
let result = save(content, mime, store_path, "2024-01-01T12-00-00.000-abc123").unwrap();
assert_eq!(result.extension, "md");
assert_eq!(result.byte_size, content.len() as u64);
assert!(result.staged_path.exists());
assert_eq!(std::fs::read(&result.staged_path).unwrap(), content);
}
#[test]
fn test_save_plain_text() {
let temp_dir = TempDir::new().unwrap();
let store_path = temp_dir.path();
let content = b"Plain text content";
let mime = "text/plain";
let result = save(content, mime, store_path, "2024-01-01T12-00-00.000-abc123").unwrap();
assert_eq!(result.extension, "txt");
assert_eq!(result.byte_size, content.len() as u64);
assert!(result.staged_path.exists());
}
#[test]
fn test_save_rejects_unsupported_mime() {
let temp_dir = TempDir::new().unwrap();
let store_path = temp_dir.path();
let content = b"test";
let result = save(content, "text/html", store_path, "2024-01-01T12-00-00.000-abc123");
assert!(result.is_err());
assert!(result.unwrap_err().to_string().contains("unsupported MIME type"));
}
#[test]
fn test_save_hash_is_consistent() {
let temp_dir = TempDir::new().unwrap();
let store_path = temp_dir.path();
let content = b"archivr text";
let result1 = save(content, "text/plain", store_path, "2024-01-01T12-00-00.000-abc123").unwrap();
let temp_dir2 = TempDir::new().unwrap();
let result2 = save(content, "text/plain", temp_dir2.path(), "2024-01-01T12-00-01.000-def456").unwrap();
assert_eq!(result1.hash, result2.hash);
}
}

View file

@ -4,6 +4,7 @@ use std::{
env,
path::{Path, PathBuf},
process::Command,
sync::OnceLock,
};
use uuid::Uuid;
use serde_json;
@ -11,6 +12,112 @@ use serde_json;
use crate::downloader::cookies::{domain_from_url, write_netscape_cookie_file};
use crate::hash::hash_file;
/// Env var that force-pins a specific yt-dlp binary, bypassing version comparison.
pub const YT_DLP_FORCE_ENV: &str = "ARCHIVR_YT_DLP_FORCE";
/// Env var set by the nix flake wrapper, pointing at the pinned yt-dlp.
pub const YT_DLP_ENV: &str = "ARCHIVR_YT_DLP";
/// Override for the mutable state directory (used by `archivr yt-dlp` and tests).
pub const STATE_DIR_ENV: &str = "ARCHIVR_STATE_DIR";
static RESOLVED_YT_DLP: OnceLock<PathBuf> = OnceLock::new();
/// Mutable per-user state directory for archivr.
///
/// `ARCHIVR_STATE_DIR` wins if set. Otherwise this mirrors what `dirs::state_dir()`
/// would give us without taking on the dependency: `~/Library/Application Support`
/// on macOS, `$XDG_STATE_HOME` (default `~/.local/state`) elsewhere.
pub fn state_dir() -> Option<PathBuf> {
if let Some(dir) = env::var_os(STATE_DIR_ENV) {
if !dir.is_empty() {
return Some(PathBuf::from(dir));
}
}
let home = PathBuf::from(env::var_os("HOME").filter(|h| !h.is_empty())?);
if cfg!(target_os = "macos") {
Some(home.join("Library").join("Application Support").join("archivr"))
} else {
let base = env::var_os("XDG_STATE_HOME")
.filter(|d| !d.is_empty())
.map(PathBuf::from)
.unwrap_or_else(|| home.join(".local").join("state"));
Some(base.join("archivr"))
}
}
/// Path of the user-installed (self-updated) yt-dlp inside the state dir.
pub fn state_dir_yt_dlp() -> Option<PathBuf> {
state_dir().map(|d| d.join("yt-dlp").join("yt-dlp"))
}
/// The explicit yt-dlp override, if it points to a file on disk.
pub fn forced_yt_dlp() -> Option<PathBuf> {
let p = PathBuf::from(env::var_os(YT_DLP_FORCE_ENV).filter(|v| !v.is_empty())?);
p.is_file().then_some(p)
}
/// The nix-pinned yt-dlp advertised via `ARCHIVR_YT_DLP`, if it exists on disk.
pub fn pinned_yt_dlp() -> Option<PathBuf> {
let p = PathBuf::from(env::var_os(YT_DLP_ENV).filter(|v| !v.is_empty())?);
p.is_file().then_some(p)
}
/// Runs `<binary> --version` and returns the trimmed stdout.
///
/// yt-dlp versions are `YYYY.MM.DD`, so plain string ordering is chronological
/// ordering — no semver parsing needed.
pub fn probe_version(binary: &Path) -> Option<String> {
let out = Command::new(binary).arg("--version").output().ok()?;
if !out.status.success() {
return None;
}
let version = String::from_utf8_lossy(&out.stdout).trim().to_string();
(!version.is_empty()).then_some(version)
}
/// The candidate yt-dlp binaries, in priority order for tie-breaking
/// (later entries win ties, so the deliberately-installed state-dir copy is last).
pub fn yt_dlp_candidates() -> Vec<(&'static str, PathBuf)> {
let mut candidates = Vec::new();
if let Some(p) = pinned_yt_dlp() {
candidates.push(("env (ARCHIVR_YT_DLP)", p));
}
if let Some(p) = state_dir_yt_dlp() {
if p.is_file() {
candidates.push(("state-dir", p));
}
}
candidates
}
/// Picks the yt-dlp binary to run, without consulting the process-wide cache.
///
/// Priority: `ARCHIVR_YT_DLP_FORCE` > newest of (pinned, state-dir) by version
/// string > bare `yt-dlp` (PATH lookup, the historical behaviour).
pub fn resolve_yt_dlp_uncached() -> PathBuf {
if let Some(forced) = forced_yt_dlp() {
return forced;
}
yt_dlp_candidates()
.into_iter()
.filter_map(|(_, path)| probe_version(&path).map(|v| (v, path)))
// `max_by` keeps the *last* maximum, and the state-dir candidate is last,
// so an exact version tie resolves in favour of the user's own install.
.max_by(|(a, _), (b, _)| a.cmp(b))
.map(|(_, path)| path)
.unwrap_or_else(|| PathBuf::from("yt-dlp"))
}
/// Cached [`resolve_yt_dlp_uncached`] — `--version` is spawned at most once
/// per process no matter how many yt-dlp calls the run makes.
pub fn resolve_yt_dlp() -> PathBuf {
RESOLVED_YT_DLP
.get_or_init(resolve_yt_dlp_uncached)
.clone()
}
/// A single item in a flat playlist listing from `fetch_playlist_info`.
#[derive(Debug)]
pub struct PlaylistItem {
@ -164,7 +271,7 @@ pub fn download(
) -> Result<(String, String)> {
println!("Downloading with yt-dlp: {path}");
let ytdlp = env::var("ARCHIVR_YT_DLP").unwrap_or_else(|_| "yt-dlp".to_string());
let ytdlp = resolve_yt_dlp();
let is_audio = quality == Some("audio");
let temp_dir = store_path.join("temp").join(timestamp);
@ -207,7 +314,7 @@ pub fn download(
.arg("-o")
.arg(&out_template)
.output()
.with_context(|| format!("failed to spawn {ytdlp} process"));
.with_context(|| format!("failed to spawn {} process", ytdlp.display()));
// Remove cookie file immediately regardless of outcome.
if let Some(cf) = &cookie_file {
@ -253,7 +360,7 @@ fn find_downloaded_file(temp_dir: &Path, timestamp: &str) -> Result<PathBuf> {
/// On failure (non-zero exit or no stdout), prints the captured stderr
/// to stderr (for debugging) then returns `None` so callers can proceed.
pub fn fetch_metadata(path: &str, cookies: &HashMap<String, String>) -> Option<String> {
let ytdlp = std::env::var("ARCHIVR_YT_DLP").unwrap_or_else(|_| "yt-dlp".to_string());
let ytdlp = resolve_yt_dlp();
// Write a temp cookie file if needed; UUID-named to avoid collisions.
let cookie_file: Option<PathBuf> = if !cookies.is_empty() {
@ -342,7 +449,7 @@ fn normalize_item_url(
/// Returns an error if yt-dlp fails, the output is not valid JSON, or
/// the root `_type` is not `"playlist"`.
pub fn fetch_playlist_info(url: &str, cookies: &HashMap<String, String>) -> Result<PlaylistInfo> {
let ytdlp = std::env::var("ARCHIVR_YT_DLP").unwrap_or_else(|_| "yt-dlp".to_string());
let ytdlp = resolve_yt_dlp();
let cookie_file: Option<PathBuf> = if !cookies.is_empty() {
let domain = domain_from_url(url);
@ -366,7 +473,7 @@ pub fn fetch_playlist_info(url: &str, cookies: &HashMap<String, String>) -> Resu
if let Some(cf) = &cookie_file {
let _ = std::fs::remove_file(cf);
}
let out = out.with_context(|| format!("failed to spawn {ytdlp}"))?;
let out = out.with_context(|| format!("failed to spawn {}", ytdlp.display()))?;
if !out.status.success() {
let stderr = String::from_utf8_lossy(&out.stderr);
bail!("yt-dlp -J --flat-playlist failed for {url}: {stderr}");
@ -425,7 +532,7 @@ pub fn probe_playlist_qualities(
url: &str,
cookies: &HashMap<String, String>,
) -> Result<PlaylistProbeResult> {
let ytdlp = std::env::var("ARCHIVR_YT_DLP").unwrap_or_else(|_| "yt-dlp".to_string());
let ytdlp = resolve_yt_dlp();
let cookie_file: Option<PathBuf> = if !cookies.is_empty() {
let domain = domain_from_url(url);
@ -449,7 +556,7 @@ pub fn probe_playlist_qualities(
if let Some(cf) = &cookie_file {
let _ = std::fs::remove_file(cf);
}
let out = out.with_context(|| format!("failed to spawn {ytdlp}"))?;
let out = out.with_context(|| format!("failed to spawn {}", ytdlp.display()))?;
if !out.status.success() {
let stderr = String::from_utf8_lossy(&out.stderr);
bail!("yt-dlp -J failed for {url}: {stderr}");
@ -501,7 +608,131 @@ pub fn probe_playlist_qualities(
#[cfg(test)]
mod tests {
use super::{available_video_heights, has_audio_track, quality_format};
use super::{
available_video_heights, has_audio_track, quality_format, resolve_yt_dlp_uncached,
state_dir, STATE_DIR_ENV, YT_DLP_ENV, YT_DLP_FORCE_ENV,
};
use std::path::{Path, PathBuf};
use std::sync::{Mutex, MutexGuard};
/// Env vars are process-global, so resolver tests take turns.
static ENV_LOCK: Mutex<()> = Mutex::new(());
/// Clears every env var the resolver reads and hands back the serialising guard.
fn env_guard() -> MutexGuard<'static, ()> {
let guard = ENV_LOCK.lock().unwrap_or_else(|e| e.into_inner());
for key in [YT_DLP_FORCE_ENV, YT_DLP_ENV, STATE_DIR_ENV] {
unsafe { std::env::remove_var(key) };
}
guard
}
/// Writes an executable stub that reports `version` when asked for `--version`.
fn fake_yt_dlp(path: &Path, version: &str) {
std::fs::create_dir_all(path.parent().unwrap()).unwrap();
std::fs::write(path, format!("#!/bin/sh\necho {version}\n")).unwrap();
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o755)).unwrap();
}
}
#[test]
fn resolve_yt_dlp_prefers_state_dir_when_newer() {
let _guard = env_guard();
let tmp = tempfile::tempdir().unwrap();
let pinned = tmp.path().join("nix/yt-dlp");
fake_yt_dlp(&pinned, "2026.08.19");
let state = tmp.path().join("state");
fake_yt_dlp(&state.join("yt-dlp/yt-dlp"), "2026.09.01");
unsafe {
std::env::set_var(YT_DLP_ENV, &pinned);
std::env::set_var(STATE_DIR_ENV, &state);
}
assert_eq!(resolve_yt_dlp_uncached(), state.join("yt-dlp/yt-dlp"));
}
#[test]
fn resolve_yt_dlp_prefers_pinned_when_newer() {
let _guard = env_guard();
let tmp = tempfile::tempdir().unwrap();
let pinned = tmp.path().join("nix/yt-dlp");
fake_yt_dlp(&pinned, "2026.09.15");
let state = tmp.path().join("state");
fake_yt_dlp(&state.join("yt-dlp/yt-dlp"), "2026.08.19");
unsafe {
std::env::set_var(YT_DLP_ENV, &pinned);
std::env::set_var(STATE_DIR_ENV, &state);
}
assert_eq!(resolve_yt_dlp_uncached(), pinned);
}
#[test]
fn resolve_yt_dlp_breaks_version_ties_toward_state_dir() {
let _guard = env_guard();
let tmp = tempfile::tempdir().unwrap();
let pinned = tmp.path().join("nix/yt-dlp");
fake_yt_dlp(&pinned, "2026.09.01");
let state = tmp.path().join("state");
fake_yt_dlp(&state.join("yt-dlp/yt-dlp"), "2026.09.01");
unsafe {
std::env::set_var(YT_DLP_ENV, &pinned);
std::env::set_var(STATE_DIR_ENV, &state);
}
assert_eq!(resolve_yt_dlp_uncached(), state.join("yt-dlp/yt-dlp"));
}
#[test]
fn resolve_yt_dlp_honours_force_override_regardless_of_version() {
let _guard = env_guard();
let tmp = tempfile::tempdir().unwrap();
let forced = tmp.path().join("forced/yt-dlp");
fake_yt_dlp(&forced, "2020.01.01");
let state = tmp.path().join("state");
fake_yt_dlp(&state.join("yt-dlp/yt-dlp"), "2026.09.01");
unsafe {
std::env::set_var(YT_DLP_FORCE_ENV, &forced);
std::env::set_var(STATE_DIR_ENV, &state);
}
assert_eq!(resolve_yt_dlp_uncached(), forced);
}
#[test]
fn resolve_yt_dlp_falls_back_to_bare_when_no_candidate_exists() {
let _guard = env_guard();
let tmp = tempfile::tempdir().unwrap();
unsafe {
std::env::set_var(YT_DLP_ENV, tmp.path().join("missing/yt-dlp"));
std::env::set_var(STATE_DIR_ENV, tmp.path().join("empty-state"));
}
assert_eq!(resolve_yt_dlp_uncached(), PathBuf::from("yt-dlp"));
}
#[test]
fn state_dir_override_wins_over_platform_default() {
let _guard = env_guard();
unsafe { std::env::set_var(STATE_DIR_ENV, "/tmp/archivr-state-override") };
assert_eq!(state_dir(), Some(PathBuf::from("/tmp/archivr-state-override")));
}
#[test]
fn quality_format_audio() {

View file

@ -4,3 +4,4 @@ pub mod database;
pub mod downloader;
pub mod hash;
pub mod twitter;
pub mod summarizer;

File diff suppressed because it is too large Load diff