use crate::config::Paths;
use crate::desk::{self, Desk, Opened, Origin, Placed};
use crate::render::Kind;
use anyhow::{Context, Result};
use rusqlite::{params, Connection, OptionalExtension};
use serde::Serialize;
use std::fs;
use std::path::PathBuf;
use std::sync::Mutex;
#[derive(Clone, Debug, Serialize)]
pub struct Doc {
pub id: String,
pub project_id: i64,
pub project: String,
pub workflow_id: i64,
pub workflow: String,
pub workflow_title: String,
pub title: String,
pub kind: Kind,
pub lang: Option<String>,
pub size: i64,
pub received_at: i64,
pub source_path: Option<String>,
pub branch: Option<String>,
pub pinned: bool,
pub origin: String,
pub content_hash: String,
pub desk: Option<Origin>,
}
#[derive(Clone, Debug, Serialize)]
pub struct TreeDoc {
pub id: String,
pub title: String,
pub kind: Kind,
pub received_at: i64,
pub pinned: bool,
pub unread: bool,
}
#[derive(Clone, Debug, Serialize)]
pub struct DeskDoc {
pub id: String,
pub title: String,
pub kind: Kind,
pub received_at: i64,
pub unread: bool,
pub pinned: bool,
pub slot: i64,
pub project: String,
pub source_path: Option<String>,
}
#[derive(Clone, Debug, Serialize)]
pub struct TreeWorkflow {
pub id: i64,
pub key: String,
pub title: String,
pub docs: Vec<TreeDoc>,
pub total: i64,
}
#[derive(Clone, Debug, Serialize)]
pub struct TreeProject {
pub id: i64,
pub name: String,
pub root: String,
pub docs: i64,
pub workflows: i64,
}
#[derive(Clone, Debug, Serialize)]
pub struct Hit {
pub id: String,
pub title: String,
pub project: String,
pub workflow_title: String,
pub kind: Kind,
pub received_at: i64,
pub snippet: String,
}
pub struct NewDoc<'a> {
pub project_root: &'a str,
pub project_name: &'a str,
pub workflow_key: &'a str,
pub workflow_title: &'a str,
pub title: &'a str,
pub kind: Kind,
pub lang: Option<&'a str>,
pub source_path: Option<&'a str>,
pub branch: Option<&'a str>,
pub origin: &'a str,
pub sender: &'a str,
pub desk: Option<&'a Origin>,
pub source: &'a [u8],
pub staged: Option<&'a Staged>,
pub search_body: &'a str,
pub html: &'a str,
}
pub struct Store {
conn: Mutex<Connection>,
docs_dir: PathBuf,
}
const SCHEMA: &str = r#"
CREATE TABLE IF NOT EXISTS projects (
id INTEGER PRIMARY KEY,
root TEXT NOT NULL UNIQUE,
name TEXT NOT NULL,
created_at INTEGER NOT NULL,
renamed INTEGER NOT NULL DEFAULT 0
);
CREATE TABLE IF NOT EXISTS workflows (
id INTEGER PRIMARY KEY,
project_id INTEGER NOT NULL REFERENCES projects(id),
key TEXT NOT NULL,
title TEXT NOT NULL,
created_at INTEGER NOT NULL,
UNIQUE(project_id, key)
);
CREATE TABLE IF NOT EXISTS docs (
id TEXT PRIMARY KEY,
project_id INTEGER NOT NULL REFERENCES projects(id),
workflow_id INTEGER NOT NULL REFERENCES workflows(id),
title TEXT NOT NULL,
kind TEXT NOT NULL,
lang TEXT,
size INTEGER NOT NULL,
received_at INTEGER NOT NULL,
source_path TEXT,
branch TEXT,
content_hash TEXT NOT NULL,
pinned INTEGER NOT NULL DEFAULT 0,
origin TEXT NOT NULL DEFAULT 'cli',
unread INTEGER NOT NULL DEFAULT 0,
deleted_at INTEGER NOT NULL DEFAULT 0
);
CREATE INDEX IF NOT EXISTS docs_recv ON docs(received_at DESC);
CREATE INDEX IF NOT EXISTS docs_wf ON docs(workflow_id, received_at);
CREATE VIRTUAL TABLE IF NOT EXISTS docs_fts USING fts5(id UNINDEXED, title, body, tokenize='unicode61');
"#;
const DOC_COLS: &str = "d.id, d.project_id, p.name, d.workflow_id, w.key, w.title, d.title, d.kind, d.lang, d.size, d.received_at, d.source_path, d.branch, d.pinned, d.origin, d.content_hash, d.desk_id, d.desk_name, d.desk_slot";
const DOC_FROM: &str =
"FROM live_docs d JOIN projects p ON p.id = d.project_id JOIN workflows w ON w.id = d.workflow_id";
const HEAD_FROM: &str =
"FROM head_docs d JOIN projects p ON p.id = d.project_id JOIN workflows w ON w.id = d.workflow_id";
impl Store {
pub fn open(paths: &Paths) -> Result<Store> {
fs::create_dir_all(&paths.data_dir).context("creating data dir")?;
fs::create_dir_all(&paths.docs_dir).context("creating docs dir")?;
if let Ok(dir) = fs::read_dir(&paths.docs_dir) {
for e in dir.flatten() {
if e.file_name().to_string_lossy().starts_with(".stage-") {
let _ = fs::remove_file(e.path());
}
}
}
let conn = Connection::open(&paths.db_path).context("opening database")?;
conn.execute_batch(
"PRAGMA journal_mode=WAL; PRAGMA synchronous=NORMAL; PRAGMA foreign_keys=ON;",
)?;
conn.execute_batch(SCHEMA)?;
conn.execute_batch(desk::SCHEMA)?;
conn.execute_batch(
"UPDATE docs SET workflow_id = (
SELECT MIN(w2.id) FROM workflows w2
JOIN workflows w1 ON w1.id = docs.workflow_id
WHERE w2.project_id = w1.project_id AND LOWER(w2.key) = LOWER(w1.key)
);
DELETE FROM workflows WHERE id NOT IN (SELECT DISTINCT workflow_id FROM docs);
UPDATE workflows SET key = LOWER(key) WHERE key <> LOWER(key);",
)
.ok();
for stmt in [
"ALTER TABLE docs ADD COLUMN pinned INTEGER NOT NULL DEFAULT 0",
"ALTER TABLE docs ADD COLUMN origin TEXT NOT NULL DEFAULT 'cli'",
"ALTER TABLE projects ADD COLUMN renamed INTEGER NOT NULL DEFAULT 0",
"ALTER TABLE docs ADD COLUMN unread INTEGER NOT NULL DEFAULT 0",
"ALTER TABLE docs ADD COLUMN deleted_at INTEGER NOT NULL DEFAULT 0",
"ALTER TABLE docs ADD COLUMN sender TEXT NOT NULL DEFAULT ''",
"ALTER TABLE docs ADD COLUMN desk_id INTEGER NOT NULL DEFAULT 0",
"ALTER TABLE docs ADD COLUMN desk_name TEXT NOT NULL DEFAULT ''",
"ALTER TABLE docs ADD COLUMN desk_slot INTEGER NOT NULL DEFAULT 0",
"ALTER TABLE panes ADD COLUMN agent_session TEXT NOT NULL DEFAULT ''",
] {
let _ = conn.execute_batch(stmt);
}
conn.execute_batch(
"CREATE INDEX IF NOT EXISTS docs_unread ON docs(unread, received_at);
CREATE INDEX IF NOT EXISTS docs_path ON docs(project_id, source_path, received_at);
DROP VIEW IF EXISTS live_docs;
CREATE VIEW live_docs AS SELECT rowid AS rowid, * FROM docs WHERE deleted_at = 0;
DROP VIEW IF EXISTS head_docs;
CREATE VIEW head_docs AS SELECT * FROM live_docs d
WHERE d.source_path IS NULL
OR d.rowid = (SELECT d2.rowid FROM live_docs d2
WHERE d2.project_id = d.project_id AND d2.source_path = d.source_path
ORDER BY d2.received_at DESC, d2.rowid DESC LIMIT 1);",
)?;
Ok(Store {
conn: Mutex::new(conn),
docs_dir: paths.docs_dir.clone(),
})
}
pub fn src_path(&self, id: &str) -> PathBuf {
self.docs_dir.join(format!("{id}.src"))
}
fn html_path(&self, id: &str) -> PathBuf {
self.docs_dir.join(format!("{id}.html"))
}
fn outline_path(&self, id: &str) -> PathBuf {
self.docs_dir.join(format!("{id}.outline"))
}
pub fn outline(&self, id: &str) -> Option<String> {
fs::read_to_string(self.outline_path(id)).ok()
}
pub fn set_outline(&self, id: &str, json: &str) -> Result<()> {
Ok(fs::write(self.outline_path(id), json)?)
}
pub fn insert(&self, id: &str, d: NewDoc) -> Result<Doc> {
let now = now();
let id = id.to_string();
let (hash, size) = self.put_source(&id, &d)?;
fs::write(self.html_path(&id), d.html)?;
let mut conn = self.conn.lock().unwrap();
let tx = conn.transaction()?;
tx.execute(
"INSERT INTO projects(root, name, created_at) VALUES(?1, ?2, ?3)
ON CONFLICT(root) DO UPDATE SET name = excluded.name WHERE projects.renamed = 0",
params![d.project_root, d.project_name, now],
)?;
let (project_id, project_name): (i64, String) = tx.query_row(
"SELECT id, name FROM projects WHERE root = ?1",
params![d.project_root],
|r| Ok((r.get(0)?, r.get(1)?)),
)?;
tx.execute(
"INSERT INTO workflows(project_id, key, title, created_at) VALUES(?1, ?2, ?3, ?4)
ON CONFLICT(project_id, key) DO NOTHING",
params![project_id, d.workflow_key, d.workflow_title, now],
)?;
let (workflow_id, workflow_title): (i64, String) = tx.query_row(
"SELECT id, title FROM workflows WHERE project_id = ?1 AND key = ?2",
params![project_id, d.workflow_key],
|r| Ok((r.get(0)?, r.get(1)?)),
)?;
tx.execute(
"INSERT INTO docs(id, project_id, workflow_id, title, kind, lang, size, received_at, source_path, branch, content_hash, pinned, origin, unread, sender, desk_id, desk_name, desk_slot)
VALUES(?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, 0, ?12, 1, ?13, ?14, ?15, ?16)",
params![id, project_id, workflow_id, d.title, d.kind.as_str(), d.lang, size, now, d.source_path, d.branch, hash, d.origin, d.sender,
d.desk.map_or(0, |o| o.id), d.desk.map_or("", |o| o.name.as_str()), d.desk.map_or(0, |o| o.slot)],
)?;
tx.execute(
"INSERT INTO docs_fts(id, title, body) VALUES(?1, ?2, ?3)",
params![id, d.title, d.search_body],
)?;
if let Some(sp) = d.source_path {
tx.execute(
"UPDATE docs SET unread = 0 WHERE project_id = ?1 AND source_path = ?2 AND id != ?3 AND unread = 1",
params![project_id, sp, id],
)?;
}
tx.commit()?;
Ok(Doc {
id,
project_id,
project: project_name,
workflow_id,
workflow: d.workflow_key.to_string(),
workflow_title,
title: d.title.to_string(),
kind: d.kind,
lang: d.lang.map(str::to_string),
size,
received_at: now,
source_path: d.source_path.map(str::to_string),
branch: d.branch.map(str::to_string),
pinned: false,
origin: d.origin.to_string(),
content_hash: hash,
desk: d.desk.cloned(),
})
}
pub fn replace(&self, id: &str, d: NewDoc) -> Result<Doc> {
let now = now();
let (hash, size) = self.put_source(id, &d)?;
fs::write(self.html_path(id), d.html)?;
let _ = fs::remove_file(self.outline_path(id));
let conn = self.conn.lock().unwrap();
conn.execute(
"UPDATE docs SET title = ?2, kind = ?3, lang = ?4, size = ?5, received_at = ?6, branch = ?7, content_hash = ?8 WHERE id = ?1",
params![id, d.title, d.kind.as_str(), d.lang, size, now, d.branch, hash],
)?;
conn.execute("DELETE FROM docs_fts WHERE id = ?1", params![id])?;
conn.execute(
"INSERT INTO docs_fts(id, title, body) VALUES(?1, ?2, ?3)",
params![id, d.title, d.search_body],
)?;
drop(conn);
self.get(id)?.context("replaced document vanished")
}
pub fn replace_html(&self, id: &str, html: &str) -> Result<()> {
Ok(fs::write(self.html_path(id), html)?)
}
pub fn get(&self, id: &str) -> Result<Option<Doc>> {
let conn = self.conn.lock().unwrap();
conn.query_row(
&format!("SELECT {DOC_COLS} {DOC_FROM} WHERE d.id = ?1"),
params![id],
row_to_doc,
)
.optional()
.map_err(Into::into)
}
pub fn html(&self, id: &str) -> Result<String> {
Ok(fs::read_to_string(self.html_path(id))?)
}
pub fn source(&self, id: &str) -> Result<String> {
self.small(id)?;
Ok(fs::read_to_string(self.src_path(id))?)
}
fn small(&self, id: &str) -> Result<()> {
let len = fs::metadata(self.src_path(id))?.len();
if len > crate::receive::MAX_BYTES as u64 {
anyhow::bail!("{id} is {} MB, too large to read into memory", len >> 20);
}
Ok(())
}
fn put_source(&self, id: &str, d: &NewDoc) -> Result<(String, i64)> {
match d.staged {
Some(st) => {
fs::rename(&st.path, self.src_path(id))?;
Ok((st.hash.clone(), st.len as i64))
}
None => {
fs::write(self.src_path(id), d.source)?;
Ok((
blake3::hash(d.source).to_hex().to_string(),
d.source.len() as i64,
))
}
}
}
pub fn stage(&self, from: &std::path::Path, cap: u64) -> Result<Staged> {
use std::io::{Read, Write};
let mut src =
fs::File::open(from).with_context(|| format!("reading {}", from.display()))?;
let gb = |n: u64| format!("{:.1} GB", n as f64 / (1u64 << 30) as f64);
let too_big = |n: u64| {
anyhow::anyhow!(
"{} is {}; snyvi keeps media up to {}. `snyvi browse` on its folder plays it where it is.",
from.display(),
gb(n),
gb(cap).replace(".0 ", " ")
)
};
let len = src.metadata()?.len();
if len > cap {
return Err(too_big(len));
}
let mut st = Staged {
path: self
.docs_dir
.join(format!(".stage-{}", new_id(&from.to_string_lossy()))),
hash: String::new(),
len: 0,
};
let mut out = fs::File::create(&st.path)?;
let mut hasher = blake3::Hasher::new();
let mut buf = vec![0u8; 256 * 1024];
loop {
let n = match src.read(&mut buf) {
Ok(0) => break,
Ok(n) => n,
Err(e) if e.kind() == std::io::ErrorKind::Interrupted => continue,
Err(e) => return Err(e.into()),
};
st.len += n as u64;
if st.len > cap {
return Err(too_big(st.len));
}
hasher.update(&buf[..n]);
out.write_all(&buf[..n])?;
}
st.hash = hasher.finalize().to_hex().to_string();
Ok(st)
}
pub fn previous(&self, doc: &Doc) -> Result<Option<Doc>> {
let conn = self.conn.lock().unwrap();
let earlier = "AND d.id != ?3 \
AND (d.received_at < ?2 OR (d.received_at = ?2 AND d.rowid < (SELECT rowid FROM docs WHERE id = ?3))) \
ORDER BY d.received_at DESC, d.rowid DESC LIMIT 1";
if let Some(path) = &doc.source_path {
return conn
.query_row(
&format!(
"SELECT {DOC_COLS} {DOC_FROM} WHERE d.project_id = ?1 AND d.source_path = ?4 {earlier}"
),
params![doc.project_id, doc.received_at, doc.id, path],
row_to_doc,
)
.optional()
.map_err(Into::into);
}
conn.query_row(
&format!("SELECT {DOC_COLS} {DOC_FROM} WHERE d.workflow_id = ?1 {earlier}"),
params![doc.workflow_id, doc.received_at, doc.id],
row_to_doc,
)
.optional()
.map_err(Into::into)
}
pub fn latest_for_path(&self, project_root: &str, source_path: &str) -> Result<Option<Doc>> {
let conn = self.conn.lock().unwrap();
conn.query_row(
&format!("SELECT {DOC_COLS} {DOC_FROM} WHERE p.root = ?1 AND d.source_path = ?2 ORDER BY d.received_at DESC, d.rowid DESC LIMIT 1"),
params![project_root, source_path],
row_to_doc,
)
.optional()
.map_err(Into::into)
}
pub fn delete(&self, id: &str) -> Result<bool> {
let mut conn = self.conn.lock().unwrap();
let tx = conn.transaction()?;
let at = now();
let found: Option<(i64, Option<String>)> = tx
.query_row(
"SELECT project_id, source_path FROM docs WHERE id = ?1 AND deleted_at = 0",
params![id],
|r| Ok((r.get(0)?, r.get(1)?)),
)
.optional()?;
let Some((project_id, source_path)) = found else {
return Ok(false);
};
let n = match source_path {
Some(sp) => tx.execute(
"UPDATE docs SET deleted_at = ?3 WHERE project_id = ?1 AND source_path = ?2 AND deleted_at = 0",
params![project_id, sp, at],
)?,
None => tx.execute(
"UPDATE docs SET deleted_at = ?2 WHERE id = ?1 AND deleted_at = 0",
params![id, at],
)?,
};
tx.commit()?;
Ok(n > 0)
}
pub fn undelete(&self, id: &str) -> Result<bool> {
let mut conn = self.conn.lock().unwrap();
let tx = conn.transaction()?;
let found: Option<(i64, Option<String>, i64)> = tx
.query_row(
"SELECT project_id, source_path, deleted_at FROM docs WHERE id = ?1 AND deleted_at != 0",
params![id],
|r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)),
)
.optional()?;
let Some((project_id, source_path, at)) = found else {
return Ok(false);
};
let n = match source_path {
Some(sp) => tx.execute(
"UPDATE docs SET deleted_at = 0 WHERE project_id = ?1 AND source_path = ?2 AND deleted_at = ?3",
params![project_id, sp, at],
)?,
None => tx.execute(
"UPDATE docs SET deleted_at = 0 WHERE id = ?1 AND deleted_at != 0",
params![id],
)?,
};
tx.commit()?;
Ok(n > 0)
}
pub fn history(&self, project_id: i64, source_path: &str) -> Result<Vec<Doc>> {
let conn = self.conn.lock().unwrap();
let rows = conn
.prepare(&format!("SELECT {DOC_COLS} {DOC_FROM} WHERE d.project_id = ?1 AND d.source_path = ?2 ORDER BY d.received_at DESC, d.rowid DESC"))?
.query_map(params![project_id, source_path], row_to_doc)?
.collect::<std::result::Result<_, _>>()?;
Ok(rows)
}
pub fn set_pinned(&self, id: &str, pinned: bool) -> Result<bool> {
let conn = self.conn.lock().unwrap();
Ok(conn.execute(
"UPDATE docs SET pinned = ?2 WHERE id = ?1",
params![id, pinned as i64],
)? > 0)
}
pub fn queue(&self, limit: usize) -> Result<Vec<Doc>> {
let conn = self.conn.lock().unwrap();
let rows = conn
.prepare(&format!(
"SELECT {DOC_COLS} {HEAD_FROM} WHERE d.unread = 1 ORDER BY d.received_at, d.rowid LIMIT ?1"
))?
.query_map(params![limit as i64], row_to_doc)?
.collect::<std::result::Result<_, _>>()?;
Ok(rows)
}
pub fn waiting(&self) -> Result<i64> {
let conn = self.conn.lock().unwrap();
Ok(
conn.query_row("SELECT COUNT(*) FROM head_docs WHERE unread = 1", [], |r| {
r.get(0)
})?,
)
}
pub fn mark_read(&self, id: &str) -> Result<bool> {
let conn = self.conn.lock().unwrap();
Ok(conn.execute(
"UPDATE docs SET unread = 0 WHERE id = ?1 AND unread = 1",
params![id],
)? > 0)
}
pub fn mark_all_read(&self) -> Result<Vec<String>> {
let mut conn = self.conn.lock().unwrap();
let tx = conn.transaction()?;
let ids: Vec<String> = tx
.prepare("SELECT id FROM head_docs WHERE unread = 1")?
.query_map([], |r| r.get(0))?
.collect::<std::result::Result<_, _>>()?;
tx.execute("UPDATE docs SET unread = 0 WHERE unread = 1", [])?;
tx.commit()?;
Ok(ids)
}
pub fn project_root(&self, project_id: i64) -> Option<String> {
let conn = self.conn.lock().unwrap();
conn.query_row(
"SELECT root FROM projects WHERE id = ?1",
params![project_id],
|r| r.get::<_, String>(0),
)
.ok()
}
pub fn rename_project(&self, id: i64, name: &str) -> Result<bool> {
let conn = self.conn.lock().unwrap();
Ok(conn.execute(
"UPDATE projects SET name = ?2, renamed = 1 WHERE id = ?1",
params![id, name],
)? > 0)
}
pub fn rename_workflow(&self, id: i64, title: &str) -> Result<bool> {
let conn = self.conn.lock().unwrap();
Ok(conn.execute(
"UPDATE workflows SET title = ?2 WHERE id = ?1",
params![id, title],
)? > 0)
}
pub fn prune(&self, before: i64, dry_run: bool) -> Result<Vec<(String, String)>> {
let mut conn = self.conn.lock().unwrap();
let victims: Vec<(String, String)> = conn
.prepare(
"SELECT id, title FROM docs WHERE deleted_at != 0 OR (pinned = 0 AND received_at < ?1) \
ORDER BY deleted_at, received_at",
)?
.query_map(params![before], |r| Ok((r.get(0)?, r.get(1)?)))?
.collect::<std::result::Result<_, _>>()?;
if dry_run || victims.is_empty() {
return Ok(victims);
}
let tx = conn.transaction()?;
for (id, _) in &victims {
tx.execute("DELETE FROM docs_fts WHERE id = ?1", params![id])?;
tx.execute("DELETE FROM docs WHERE id = ?1", params![id])?;
}
tx.execute(
"DELETE FROM workflows WHERE id NOT IN (SELECT DISTINCT workflow_id FROM docs)",
[],
)?;
tx.execute(
"DELETE FROM projects WHERE id NOT IN (SELECT DISTINCT project_id FROM docs)",
[],
)?;
tx.commit()?;
drop(conn);
for (id, _) in &victims {
let _ = fs::remove_file(self.src_path(id));
let _ = fs::remove_file(self.html_path(id));
let _ = fs::remove_file(self.outline_path(id));
}
Ok(victims)
}
pub fn projects(&self) -> Result<Vec<TreeProject>> {
let conn = self.conn.lock().unwrap();
let rows = conn
.prepare(
"SELECT p.id, p.name, p.root, COUNT(d.id), COUNT(DISTINCT d.workflow_id)
FROM projects p JOIN head_docs d ON d.project_id = p.id
GROUP BY p.id ORDER BY MAX(d.received_at) DESC, p.id DESC",
)?
.query_map([], |r| {
Ok(TreeProject {
id: r.get(0)?,
name: r.get(1)?,
root: r.get(2)?,
docs: r.get(3)?,
workflows: r.get(4)?,
})
})?
.collect::<std::result::Result<_, _>>()?;
Ok(rows)
}
pub fn project_tree(
&self,
project_id: i64,
workflows: usize,
docs: usize,
) -> Result<Vec<TreeWorkflow>> {
let wf_limit = if workflows == 0 { -1 } else { workflows as i64 };
let doc_limit = if docs == 0 { -1 } else { docs as i64 };
let conn = self.conn.lock().unwrap();
let mut wfs: Vec<TreeWorkflow> = conn
.prepare(
"SELECT w.id, w.key, w.title, COUNT(d.id)
FROM workflows w JOIN head_docs d ON d.workflow_id = w.id
WHERE w.project_id = ?1
GROUP BY w.id ORDER BY MAX(d.received_at) DESC, w.id DESC LIMIT ?2",
)?
.query_map(params![project_id, wf_limit], |r| {
Ok(TreeWorkflow {
id: r.get(0)?,
key: r.get(1)?,
title: r.get(2)?,
docs: vec![],
total: r.get(3)?,
})
})?
.collect::<std::result::Result<_, _>>()?;
let mut doc_stmt = conn.prepare(
"SELECT id, title, kind, received_at, pinned, unread FROM head_docs
WHERE workflow_id = ?1 ORDER BY received_at DESC, rowid DESC LIMIT ?2",
)?;
for w in &mut wfs {
w.docs = doc_stmt
.query_map(params![w.id, doc_limit], row_to_tree_doc)?
.collect::<std::result::Result<_, _>>()?;
}
Ok(wfs)
}
pub fn workflow_tree(&self, workflow_id: i64) -> Result<Option<TreeWorkflow>> {
let conn = self.conn.lock().unwrap();
let found = conn
.query_row(
"SELECT id, key, title FROM workflows WHERE id = ?1",
params![workflow_id],
|r| {
Ok(TreeWorkflow {
id: r.get(0)?,
key: r.get(1)?,
title: r.get(2)?,
docs: vec![],
total: 0,
})
},
)
.optional()?;
let Some(mut w) = found else {
return Ok(None);
};
w.docs = conn
.prepare(
"SELECT id, title, kind, received_at, pinned, unread FROM head_docs
WHERE workflow_id = ?1 ORDER BY received_at DESC, rowid DESC",
)?
.query_map(params![workflow_id], row_to_tree_doc)?
.collect::<std::result::Result<Vec<_>, _>>()?;
w.total = w.docs.len() as i64;
Ok((!w.docs.is_empty()).then_some(w))
}
pub fn inbox(&self, limit: usize) -> Result<Vec<Doc>> {
let conn = self.conn.lock().unwrap();
let rows = conn
.prepare(&format!(
"SELECT {DOC_COLS} {HEAD_FROM} ORDER BY d.received_at DESC, d.rowid DESC LIMIT ?1"
))?
.query_map(params![limit as i64], row_to_doc)?
.collect::<std::result::Result<_, _>>()?;
Ok(rows)
}
pub fn search(&self, q: &str, limit: usize) -> Result<Vec<Hit>> {
let mut project: Option<String> = None;
let mut kind: Option<String> = None;
let mut terms: Vec<String> = vec![];
for t in q.split_whitespace() {
if let Some(v) = t.strip_prefix("p:").or_else(|| t.strip_prefix("project:")) {
project = Some(v.to_string());
} else if let Some(v) = t.strip_prefix("kind:").or_else(|| t.strip_prefix("k:")) {
kind = Some(match v {
"md" => "markdown".to_string(),
"txt" => "text".to_string(),
other => other.to_string(),
});
} else {
terms.push(format!("\"{}\"*", t.replace('"', "\"\"")));
}
}
if terms.is_empty() && project.is_none() && kind.is_none() {
return Ok(vec![]);
}
let mut sql =
String::from("SELECT d.id, d.title, p.name, w.title, d.kind, d.received_at, ");
let filtered_only = terms.is_empty();
if filtered_only {
sql.push_str("'' FROM live_docs d JOIN projects p ON p.id = d.project_id JOIN workflows w ON w.id = d.workflow_id WHERE 1=1");
} else {
sql.push_str(
"snippet(docs_fts, 2, '<mark>', '</mark>', '…', 14) FROM docs_fts f JOIN live_docs d ON d.id = f.id \
JOIN projects p ON p.id = d.project_id JOIN workflows w ON w.id = d.workflow_id WHERE docs_fts MATCH ?1",
);
}
let query = terms.join(" ");
let mut args: Vec<Box<dyn rusqlite::ToSql>> = vec![];
if !filtered_only {
args.push(Box::new(query));
}
if let Some(p) = &project {
sql.push_str(&format!(" AND p.name LIKE ?{}", args.len() + 1));
args.push(Box::new(format!("{p}%")));
}
if let Some(k) = &kind {
sql.push_str(&format!(" AND d.kind = ?{}", args.len() + 1));
args.push(Box::new(k.clone()));
}
sql.push_str(if filtered_only {
" ORDER BY d.received_at DESC, d.rowid DESC"
} else {
" ORDER BY bm25(docs_fts, 4.0, 1.0)"
});
sql.push_str(&format!(" LIMIT ?{}", args.len() + 1));
args.push(Box::new(limit as i64));
let conn = self.conn.lock().unwrap();
let rows = conn
.prepare(&sql)?
.query_map(
rusqlite::params_from_iter(args.iter().map(|a| a.as_ref())),
|r| {
Ok(Hit {
id: r.get(0)?,
title: r.get(1)?,
project: r.get(2)?,
workflow_title: r.get(3)?,
kind: Kind::parse(&r.get::<_, String>(4)?).unwrap_or(Kind::Text),
received_at: r.get(5)?,
snippet: r.get(6)?,
})
},
)?
.collect::<std::result::Result<_, _>>()?;
Ok(rows)
}
pub fn count(&self) -> Result<i64> {
let conn = self.conn.lock().unwrap();
Ok(conn.query_row("SELECT COUNT(*) FROM live_docs", [], |r| r.get(0))?)
}
pub fn senders(&self) -> Result<Vec<(String, i64)>> {
let conn = self.conn.lock().unwrap();
let mut st = conn.prepare(
"SELECT sender, MAX(received_at) FROM docs WHERE sender <> '' GROUP BY sender",
)?;
let rows = st.query_map([], |r| Ok((r.get(0)?, r.get(1)?)))?;
Ok(rows.collect::<std::result::Result<_, _>>()?)
}
pub fn desks(&self) -> Result<Vec<Desk>> {
desk::list(&self.conn.lock().unwrap())
}
pub fn desk(&self, id: i64) -> Result<Option<Desk>> {
desk::get(&self.conn.lock().unwrap(), id)
}
pub fn desk_docs(&self, desk_id: i64, limit: usize) -> Result<Vec<DeskDoc>> {
let conn = self.conn.lock().unwrap();
let rows = conn
.prepare(
"SELECT d.id, d.title, d.kind, d.received_at, d.unread, d.pinned, d.desk_slot, p.name, d.source_path
FROM live_docs d JOIN projects p ON p.id = d.project_id
WHERE d.desk_id = ?1
AND (d.source_path IS NULL
OR d.rowid = (SELECT d2.rowid FROM live_docs d2
WHERE d2.desk_id = ?1 AND d2.project_id = d.project_id
AND d2.source_path = d.source_path
ORDER BY d2.received_at DESC, d2.rowid DESC LIMIT 1))
ORDER BY d.received_at DESC, d.rowid DESC LIMIT ?2",
)?
.query_map(params![desk_id, limit as i64], |r| {
Ok(DeskDoc {
id: r.get(0)?,
title: r.get(1)?,
kind: Kind::parse(&r.get::<_, String>(2)?).unwrap_or(Kind::Text),
received_at: r.get(3)?,
unread: r.get::<_, i64>(4)? != 0,
pinned: r.get::<_, i64>(5)? != 0,
slot: r.get(6)?,
project: r.get(7)?,
source_path: r.get(8)?,
})
})?
.collect::<std::result::Result<_, _>>()?;
Ok(rows)
}
pub fn create_desk(&self, root: &str, name: Option<&str>) -> Result<Desk> {
desk::create(&self.conn.lock().unwrap(), root, name, now())
}
pub fn rename_desk(&self, id: i64, name: &str) -> Result<bool> {
desk::rename(&self.conn.lock().unwrap(), id, name)
}
pub fn set_desk_layout(&self, id: i64, col: f64, row: f64) -> Result<bool> {
desk::layout(&self.conn.lock().unwrap(), id, col, row)
}
pub fn delete_desk(&self, id: i64) -> Result<bool> {
desk::delete(&self.conn.lock().unwrap(), id)
}
pub fn pane(&self, id: &str) -> Result<Option<Placed>> {
desk::pane(&self.conn.lock().unwrap(), id)
}
pub fn set_pane_cmd(&self, id: &str, cmd: &str) -> Result<bool> {
desk::set_cmd(&self.conn.lock().unwrap(), id, cmd)
}
pub fn set_pane_session(&self, id: &str, session: &str) -> Result<bool> {
desk::set_agent_session(&self.conn.lock().unwrap(), id, session)
}
pub fn open_pane(&self, desk_id: i64, cwd: &str, cmd: &str) -> Result<Opened> {
desk::open_pane(&mut self.conn.lock().unwrap(), desk_id, cwd, cmd, now())
}
pub fn close_pane(&self, id: &str) -> Result<bool> {
desk::close_pane(&self.conn.lock().unwrap(), id)
}
pub fn panes_open(&self) -> Result<i64> {
desk::panes_open(&self.conn.lock().unwrap())
}
pub fn desk_notes(&self, desk_id: i64) -> Result<Vec<desk::DeskNote>> {
desk::notes(&self.conn.lock().unwrap(), desk_id)
}
pub fn add_desk_note(&self, desk_id: i64, text: &str) -> Result<Option<desk::DeskNote>> {
desk::add_note(&mut self.conn.lock().unwrap(), desk_id, text, now())
}
pub fn set_desk_note(
&self,
desk_id: i64,
id: i64,
text: Option<&str>,
done: Option<bool>,
) -> Result<bool> {
desk::set_note(&self.conn.lock().unwrap(), desk_id, id, text, done, now())
}
pub fn remove_desk_note(&self, desk_id: i64, id: i64) -> Result<bool> {
desk::remove_note(&self.conn.lock().unwrap(), desk_id, id, now())
}
pub fn restore_desk_note(&self, desk_id: i64, id: i64) -> Result<bool> {
desk::restore_note(&self.conn.lock().unwrap(), desk_id, id)
}
pub fn census(&self) -> Result<Census> {
let conn = self.conn.lock().unwrap();
Ok(conn.query_row(
"SELECT COUNT(*), COUNT(DISTINCT project_id), COALESCE(SUM(pinned), 0),
(SELECT COUNT(*) FROM desks) FROM live_docs",
[],
|r| {
Ok(Census {
documents: r.get(0)?,
projects: r.get(1)?,
pinned: r.get(2)?,
desks: r.get(3)?,
})
},
)?)
}
pub fn reset(&self) -> Result<()> {
let conn = self.conn.lock().unwrap();
conn.execute_batch(
"DELETE FROM docs_fts; DELETE FROM docs; DELETE FROM workflows; DELETE FROM projects;",
)?;
desk::clear(&conn)?;
conn.execute_batch("VACUUM;")?;
drop(conn);
if let Ok(entries) = fs::read_dir(&self.docs_dir) {
for e in entries.flatten() {
let _ = fs::remove_file(e.path());
}
}
Ok(())
}
}
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Serialize, serde::Deserialize)]
pub struct Census {
pub documents: i64,
pub projects: i64,
pub pinned: i64,
#[serde(default)]
pub desks: i64,
}
fn row_to_tree_doc(r: &rusqlite::Row) -> rusqlite::Result<TreeDoc> {
Ok(TreeDoc {
id: r.get(0)?,
title: r.get(1)?,
kind: Kind::parse(&r.get::<_, String>(2)?).unwrap_or(Kind::Text),
received_at: r.get(3)?,
pinned: r.get::<_, i64>(4)? != 0,
unread: r.get::<_, i64>(5)? != 0,
})
}
fn row_to_doc(r: &rusqlite::Row) -> rusqlite::Result<Doc> {
Ok(Doc {
id: r.get(0)?,
project_id: r.get(1)?,
project: r.get(2)?,
workflow_id: r.get(3)?,
workflow: r.get(4)?,
workflow_title: r.get(5)?,
title: r.get(6)?,
kind: Kind::parse(&r.get::<_, String>(7)?).unwrap_or(Kind::Text),
lang: r.get(8)?,
size: r.get(9)?,
received_at: r.get(10)?,
source_path: r.get(11)?,
branch: r.get(12)?,
pinned: r.get::<_, i64>(13)? != 0,
origin: r.get(14)?,
content_hash: r.get(15)?,
desk: match r.get::<_, i64>(16)? {
0 => None,
id => Some(Origin {
id,
name: r.get(17)?,
slot: r.get(18)?,
}),
},
})
}
pub fn now() -> i64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_secs() as i64)
.unwrap_or(0)
}
#[derive(Debug)]
pub struct Staged {
path: PathBuf,
pub hash: String,
pub len: u64,
}
impl Drop for Staged {
fn drop(&mut self) {
let _ = fs::remove_file(&self.path);
}
}
pub fn new_id(hash: &str) -> String {
use std::sync::atomic::{AtomicU64, Ordering};
static N: AtomicU64 = AtomicU64::new(0);
let n = N.fetch_add(1, Ordering::Relaxed);
let mixed = blake3::hash(format!("{hash}:{}:{}:{n}", now(), std::process::id()).as_bytes());
mixed.to_hex()[..10].to_string()
}
#[cfg(test)]
mod tests {
use super::*;
fn temp_store() -> (Store, tempdir::Dir) {
let dir = tempdir::Dir::new("snyvi-store");
let paths = Paths {
data_dir: dir.path.clone(),
config_dir: dir.path.clone(),
docs_dir: dir.path.join("docs"),
db_path: dir.path.join("t.db"),
token_path: dir.path.join("token"),
};
(Store::open(&paths).unwrap(), dir)
}
fn new_doc<'a>(title: &'a str, src: &'a str, wf: &'a str) -> NewDoc<'a> {
version_of(src, title, src, wf)
}
fn version_of<'a>(path: &'a str, title: &'a str, src: &'a str, wf: &'a str) -> NewDoc<'a> {
NewDoc {
project_root: "/p",
project_name: "p",
workflow_key: wf,
workflow_title: wf,
title,
kind: Kind::Markdown,
lang: None,
source_path: Some(path),
branch: None,
origin: "cli",
sender: "",
desk: None,
source: src.as_bytes(),
staged: None,
search_body: src,
html: "<p>x</p>",
}
}
#[test]
fn insert_get_previous_search() {
let (s, _d) = temp_store();
let a = s
.insert(
&new_id("a"),
version_of("/p/PLAN.md", "Plan", "# Plan\n\nalpha bravo", "w"),
)
.unwrap();
std::thread::sleep(std::time::Duration::from_millis(1100));
let b = s
.insert(
&new_id("b"),
version_of("/p/PLAN.md", "Plan", "# Plan\n\nalpha charlie", "w"),
)
.unwrap();
assert_ne!(a.id, b.id);
assert_eq!(s.get(&b.id).unwrap().unwrap().title, "Plan");
assert_eq!(s.previous(&b).unwrap().unwrap().id, a.id);
assert!(s.previous(&a).unwrap().is_none());
let hits = s.search("charlie", 10).unwrap();
assert_eq!(hits.len(), 1);
assert_eq!(hits[0].id, b.id);
assert!(
s.search("\"unbalanced", 10).is_ok(),
"punctuation must not break FTS"
);
assert_eq!(
s.latest_for_path("/p", "/p/PLAN.md").unwrap().unwrap().id,
b.id
);
let p = &s.projects().unwrap()[0];
assert_eq!((p.docs, p.workflows), (1, 1));
assert_eq!(s.project_tree(p.id, 10, 10).unwrap()[0].docs.len(), 1);
assert_eq!(s.history(p.id, "/p/PLAN.md").unwrap().len(), 2);
}
#[test]
fn a_chosen_project_name_outranks_the_derived_one() {
let (s, _d) = temp_store();
let a = s.insert(&new_id("a"), new_doc("Plan", "one", "w")).unwrap();
assert_eq!(a.project, "p");
let mut d = new_doc("Plan", "two", "w");
d.project_name = "p-moved";
let b = s.insert(&new_id("b"), d).unwrap();
assert_eq!(b.project, "p-moved");
assert!(s.rename_project(b.project_id, "Auth work").unwrap());
let mut d = new_doc("Plan", "three", "w");
d.project_name = "p-moved-again";
let c = s.insert(&new_id("c"), d).unwrap();
assert_eq!(c.project_id, a.project_id, "the root is still the identity");
assert_eq!(c.project, "Auth work");
assert_eq!(s.projects().unwrap()[0].name, "Auth work");
assert!(!s.rename_project(9999, "nobody").unwrap());
}
#[test]
fn a_renamed_workflow_keeps_the_key_that_sends_find_it_by() {
let (s, _d) = temp_store();
let a = s
.insert(&new_id("a"), new_doc("Plan", "one", "sess-1"))
.unwrap();
assert_eq!(a.workflow_title, "sess-1");
assert!(s.rename_workflow(a.workflow_id, "Auth refactor").unwrap());
let b = s
.insert(&new_id("b"), new_doc("Plan 2", "two", "sess-1"))
.unwrap();
assert_eq!(b.workflow_id, a.workflow_id, "same session, same workflow");
assert_eq!(b.workflow_title, "Auth refactor");
let t = s.project_tree(b.project_id, 10, 10).unwrap();
assert_eq!(t[0].title, "Auth refactor");
assert_eq!(t[0].key, "sess-1");
assert!(!s.rename_workflow(9999, "nobody").unwrap());
}
#[test]
fn same_second_documents_keep_their_order() {
let (s, _d) = temp_store();
let a = s.insert(&new_id("a"), new_doc("A", "first", "w")).unwrap();
let b = s.insert(&new_id("b"), new_doc("B", "second", "w")).unwrap();
let c = s.insert(&new_id("c"), new_doc("C", "third", "w")).unwrap();
let t = a.received_at;
s.conn
.lock()
.unwrap()
.execute("UPDATE docs SET received_at = ?1", params![t])
.unwrap();
let (a, b, c) = (
Doc {
received_at: t,
..a
},
Doc {
received_at: t,
..b
},
Doc {
received_at: t,
..c
},
);
let order: Vec<String> = s.inbox(9).unwrap().into_iter().map(|d| d.title).collect();
assert_eq!(
order,
vec!["C", "B", "A"],
"newest first even within a second"
);
assert_eq!(
s.project_tree(c.project_id, 10, 10).unwrap()[0].docs[0].title,
"C"
);
assert_eq!(s.history(a.project_id, "first").unwrap().len(), 1);
assert!(
s.previous(&b).unwrap().is_none() && s.previous(&c).unwrap().is_none(),
"a different file is not a version of the one before it"
);
let a2 = s
.insert(&new_id("a2"), version_of("first", "A2", "first again", "w"))
.unwrap();
s.conn
.lock()
.unwrap()
.execute("UPDATE docs SET received_at = ?1", params![t])
.unwrap();
let a2 = Doc {
received_at: t,
..a2
};
assert_eq!(s.previous(&a2).unwrap().unwrap().id, a.id);
assert!(
s.previous(&a).unwrap().is_none(),
"the first send of a file"
);
let order: Vec<String> = s.inbox(9).unwrap().into_iter().map(|d| d.title).collect();
assert_eq!(order, vec!["A2", "C", "B"], "and A is behind A2 now");
}
#[test]
fn workflow_keys_ignore_case() {
let (s, _d) = temp_store();
let a = s
.insert(&new_id("a"), new_doc("A", "one", "ksi pivot"))
.unwrap();
let b = s
.insert(&new_id("b"), new_doc("B", "two", "ksi pivot"))
.unwrap();
assert_eq!(a.workflow_id, b.workflow_id, "same key is one workflow");
assert_eq!(s.projects().unwrap()[0].workflows, 1);
}
#[test]
fn pin_and_prune() {
let (s, _d) = temp_store();
let a = s.insert(&new_id("a"), new_doc("A", "aaa", "w")).unwrap();
let b = s.insert(&new_id("b"), new_doc("B", "bbb", "w")).unwrap();
assert!(s.set_pinned(&a.id, true).unwrap());
let dry = s.prune(now() + 10, true).unwrap();
assert_eq!(dry.len(), 1);
assert_eq!(s.count().unwrap(), 2, "dry run deletes nothing");
let gone = s.prune(now() + 10, false).unwrap();
assert_eq!(gone[0].0, b.id);
assert_eq!(s.count().unwrap(), 1);
assert!(s.get(&b.id).unwrap().is_none());
assert!(s.html(&b.id).is_err(), "files removed");
assert!(s.get(&a.id).unwrap().unwrap().pinned);
}
#[test]
fn census_and_reset() {
let (s, d) = temp_store();
let a = s.insert(&new_id("a"), new_doc("A", "aaa", "w")).unwrap();
let b = s.insert(&new_id("b"), new_doc("B", "bbb", "w")).unwrap();
let c = s
.insert(&new_id("c"), new_doc("C", "ccc", "other"))
.unwrap();
assert!(s.set_pinned(&a.id, true).unwrap());
assert!(s.delete(&c.id).unwrap());
assert_eq!(
s.census().unwrap(),
Census {
documents: 2,
projects: 1,
pinned: 1,
desks: 0,
}
);
assert_eq!(std::fs::read_dir(d.path.join("docs")).unwrap().count(), 6);
s.reset().unwrap();
assert_eq!(s.census().unwrap(), Census::default());
assert_eq!(s.count().unwrap(), 0);
assert!(s.get(&a.id).unwrap().is_none());
assert!(!s.undelete(&c.id).unwrap(), "the deleted one went too");
assert!(
s.search("aaa", 10).unwrap().is_empty(),
"and the index with it"
);
assert_eq!(std::fs::read_dir(d.path.join("docs")).unwrap().count(), 0);
let again = s.insert(&new_id("b"), new_doc("B", "bbb", "w")).unwrap();
assert_eq!(s.count().unwrap(), 1);
assert!(s.previous(&again).unwrap().is_none());
let _ = b;
}
#[test]
fn a_desk_is_not_a_document_and_a_reset_still_takes_it() {
let (s, _d) = temp_store();
s.insert(&new_id("a"), new_doc("A", "aaa", "w")).unwrap();
let desk = s.create_desk("/home/p/snyvi", None).unwrap();
assert!(matches!(
s.open_pane(desk.id, "/home/p/snyvi", "").unwrap(),
Opened::Pane(_)
));
assert_eq!(s.count().unwrap(), 1);
assert_eq!(s.census().unwrap().documents, 1);
assert_eq!(s.census().unwrap().desks, 1);
assert!(s.search("snyvi", 10).unwrap().is_empty());
assert_eq!(s.projects().unwrap().len(), 1, "the desk made no project");
assert_eq!(s.inbox(10).unwrap().len(), 1);
assert!(s.get(&desk.id.to_string()).unwrap().is_none());
assert_eq!(s.desks().unwrap().len(), 1);
assert_eq!(s.panes_open().unwrap(), 1);
s.reset().unwrap();
assert_eq!(s.census().unwrap(), Census::default());
assert!(s.desks().unwrap().is_empty());
assert_eq!(s.panes_open().unwrap(), 0);
let again = s.create_desk("/home/p/snyvi", None).unwrap();
assert_eq!(again.name, "snyvi");
}
#[test]
fn a_desk_lists_what_its_panes_sent_newest_first() {
let (s, _d) = temp_store();
let here = Origin {
id: 7,
name: "snyvi".into(),
slot: 2,
};
let elsewhere = Origin {
id: 8,
name: "other".into(),
slot: 1,
};
fn from<'a>(o: &'a Origin, mut d: NewDoc<'a>) -> NewDoc<'a> {
d.desk = Some(o);
d
}
s.insert(&new_id("a"), from(&here, new_doc("First", "aaa", "w")))
.unwrap();
s.insert(
&new_id("b"),
from(&elsewhere, new_doc("Theirs", "bbb", "w")),
)
.unwrap();
s.insert(&new_id("c"), new_doc("From the CLI", "ccc", "w"))
.unwrap();
let last = s
.insert(&new_id("d"), from(&here, new_doc("Second", "ddd", "w")))
.unwrap();
let docs = s.desk_docs(7, 40).unwrap();
assert_eq!(
docs.iter().map(|d| d.title.as_str()).collect::<Vec<_>>(),
["Second", "First"]
);
assert_eq!(docs[0].slot, 2);
assert!(docs[0].unread);
assert_eq!(docs[0].project, "p");
assert_eq!(s.desk_docs(8, 40).unwrap().len(), 1);
assert!(s.desk_docs(9, 40).unwrap().is_empty());
s.mark_read(&last.id).unwrap();
assert!(!s.desk_docs(7, 40).unwrap()[0].unread);
s.delete(&last.id).unwrap();
assert_eq!(s.desk_docs(7, 40).unwrap().len(), 1);
let again = s
.insert(
&new_id("e"),
from(&here, version_of("aaa", "First, again", "eee", "w")),
)
.unwrap();
assert_eq!(
s.desk_docs(7, 40)
.unwrap()
.iter()
.map(|d| d.title.as_str())
.collect::<Vec<_>>(),
["First, again"]
);
s.insert(
&new_id("f"),
version_of("aaa", "First, elsewhere", "fff", "w"),
)
.unwrap();
let listed = s.desk_docs(7, 40).unwrap();
assert_eq!(listed.len(), 1);
assert_eq!(listed[0].id, again.id);
}
#[test]
fn a_file_sent_again_is_one_row_with_its_versions_behind_it() {
let (s, _d) = temp_store();
let first = s
.insert(&new_id("1"), version_of("/p/S.md", "Script", "one", "w"))
.unwrap();
let second = s
.insert(&new_id("2"), version_of("/p/S.md", "Script v2", "two", "w"))
.unwrap();
let third = s
.insert(
&new_id("3"),
version_of("/p/S.md", "Script v3", "three", "w"),
)
.unwrap();
let other = s
.insert(&new_id("o"), new_doc("Notes", "elsewhere", "w"))
.unwrap();
let titles = |v: Vec<Doc>| v.into_iter().map(|d| d.title).collect::<Vec<_>>();
assert_eq!(titles(s.inbox(10).unwrap()), vec!["Notes", "Script v3"]);
let wfs = s.project_tree(first.project_id, 0, 0).unwrap();
assert_eq!(titles(s.queue(10).unwrap()), vec!["Script v3", "Notes"]);
assert_eq!(
wfs[0]
.docs
.iter()
.map(|d| d.title.as_str())
.collect::<Vec<_>>(),
["Notes", "Script v3"]
);
assert_eq!(wfs[0].total, 2, "and it says two, not four");
assert_eq!(s.projects().unwrap()[0].docs, 2);
let hist = s.history(first.project_id, "/p/S.md").unwrap();
assert_eq!(titles(hist), vec!["Script v3", "Script v2", "Script"]);
assert_eq!(s.get(&first.id).unwrap().unwrap().title, "Script");
assert_eq!(
s.latest_for_path("/p", "/p/S.md").unwrap().unwrap().id,
third.id
);
assert_eq!(s.waiting().unwrap(), 2);
assert!(!s.mark_read(&second.id).unwrap(), "already off the queue");
assert!(s.mark_read(&third.id).unwrap());
assert_eq!(s.waiting().unwrap(), 1);
let _ = other;
}
#[test]
fn removing_a_document_takes_its_versions_and_undo_brings_them_back() {
let (s, _d) = temp_store();
let first = s
.insert(&new_id("1"), version_of("/p/S.md", "Script", "one", "w"))
.unwrap();
let newest = s
.insert(&new_id("2"), version_of("/p/S.md", "Script v2", "two", "w"))
.unwrap();
let keep = s
.insert(&new_id("k"), new_doc("Notes", "elsewhere", "w"))
.unwrap();
assert!(s.delete(&newest.id).unwrap());
assert!(s.get(&first.id).unwrap().is_none());
assert_eq!(s.inbox(10).unwrap().len(), 1);
assert!(s.history(first.project_id, "/p/S.md").unwrap().is_empty());
assert_eq!(s.get(&keep.id).unwrap().unwrap().title, "Notes");
assert!(s.undelete(&newest.id).unwrap());
assert_eq!(
s.history(first.project_id, "/p/S.md").unwrap().len(),
2,
"both versions came back, not just the one that was clicked"
);
assert_eq!(s.inbox(10).unwrap().len(), 2);
assert!(s.delete(&first.id).unwrap());
assert!(s.history(first.project_id, "/p/S.md").unwrap().is_empty());
assert!(s.undelete(&first.id).unwrap());
assert_eq!(s.history(first.project_id, "/p/S.md").unwrap().len(), 2);
}
#[test]
fn a_delete_can_be_taken_back_until_prune() {
let (s, _d) = temp_store();
let a = s.insert(&new_id("a"), new_doc("A", "alpha", "w")).unwrap();
let b = s.insert(&new_id("b"), new_doc("B", "bravo", "w")).unwrap();
assert!(s.delete(&b.id).unwrap());
assert!(
!s.delete(&b.id).unwrap(),
"deleting it again changes nothing"
);
assert!(s.get(&b.id).unwrap().is_none());
assert_eq!(s.count().unwrap(), 1);
assert_eq!(s.inbox(10).unwrap().len(), 1);
assert_eq!(s.waiting().unwrap(), 1, "and off the queue");
assert!(s.search("bravo", 10).unwrap().is_empty());
let wfs = s.project_tree(a.project_id, 0, 0).unwrap();
assert_eq!(wfs[0].total, 1);
assert_eq!(wfs[0].docs.len(), 1);
assert_eq!(s.projects().unwrap()[0].docs, 1, "the sidebar's count too");
assert!(s.html(&b.id).is_ok(), "still on disk");
assert!(s.undelete(&b.id).unwrap());
assert!(!s.undelete(&b.id).unwrap(), "undoing twice changes nothing");
assert_eq!(s.get(&b.id).unwrap().unwrap().title, "B");
assert_eq!(s.waiting().unwrap(), 2);
assert_eq!(s.search("bravo", 10).unwrap().len(), 1);
assert!(s.set_pinned(&b.id, true).unwrap());
assert!(s.delete(&b.id).unwrap());
let gone = s.prune(0, false).unwrap();
assert_eq!(gone.len(), 1, "nothing here is old enough but this one");
assert_eq!(gone[0].0, b.id);
assert!(!s.undelete(&b.id).unwrap(), "nothing left to put back");
assert!(s.html(&b.id).is_err(), "files removed");
assert_eq!(s.get(&a.id).unwrap().unwrap().title, "A");
}
#[test]
fn the_queue_is_what_arrived_and_was_not_opened() {
let (s, _d) = temp_store();
assert!(s.queue(10).unwrap().is_empty());
let a = s.insert(&new_id("a"), new_doc("A", "aaa", "w")).unwrap();
let b = s.insert(&new_id("b"), new_doc("B", "bbb", "w")).unwrap();
let ids = |q: Vec<Doc>| q.into_iter().map(|d| d.id).collect::<Vec<_>>();
assert_eq!(
ids(s.queue(10).unwrap()),
vec![a.id.clone(), b.id.clone()],
"oldest first"
);
assert_eq!(s.waiting().unwrap(), 2);
assert!(s.mark_read(&a.id).unwrap());
assert!(!s.mark_read(&a.id).unwrap(), "already read");
assert_eq!(s.waiting().unwrap(), 1);
assert_eq!(ids(s.queue(10).unwrap()), vec![b.id.clone()]);
s.replace(&b.id, new_doc("B2", "bbb2", "w")).unwrap();
assert_eq!(
ids(s.queue(10).unwrap()),
vec![b.id.clone()],
"an overwrite is not an arrival"
);
let c = s.insert(&new_id("c"), new_doc("C", "ccc", "w")).unwrap();
let cleared = s.mark_all_read().unwrap();
assert_eq!(cleared, vec![b.id, c.id]);
assert!(s.queue(10).unwrap().is_empty());
assert!(s.mark_all_read().unwrap().is_empty());
}
#[test]
fn replace_keeps_id_and_updates_index() {
let (s, _d) = temp_store();
let a = s
.insert(&new_id("a"), new_doc("A", "first draft", "w"))
.unwrap();
let r = s
.replace(&a.id, new_doc("A2", "second draft", "w"))
.unwrap();
assert_eq!(r.id, a.id);
assert_eq!(r.title, "A2");
assert_eq!(s.source(&a.id).unwrap(), "second draft");
assert!(s.search("first", 5).unwrap().is_empty());
assert_eq!(s.search("second", 5).unwrap().len(), 1);
}
#[test]
fn an_expanded_project_is_a_screenful_not_a_year() {
let (s, _d) = temp_store();
for w in 0..12 {
for i in 0..12 {
let body = format!("session {w}, document {i}");
s.insert(
&new_id(&format!("d{w}-{i}")),
new_doc("Plan", &body, &format!("sess-{w}")),
)
.unwrap();
}
}
let p = s.projects().unwrap();
assert_eq!(p.len(), 1);
assert_eq!(
(p[0].docs, p[0].workflows),
(144, 12),
"a row carries what is behind it as two numbers, not as rows"
);
let capped = s.project_tree(p[0].id, 10, 10).unwrap();
assert_eq!(capped.len(), 10, "ten of the twelve sessions");
assert!(
capped.iter().all(|w| w.docs.len() == 10 && w.total == 12),
"ten documents each, and each says it holds twelve"
);
let all = s.project_tree(p[0].id, 0, 0).unwrap();
assert_eq!(all.len(), 12, "a zero cap is a reader asking for the rest");
assert!(all.iter().all(|w| w.docs.len() == 12));
let w = s.workflow_tree(capped[0].id).unwrap().unwrap();
assert_eq!((w.docs.len(), w.total), (12, 12));
assert!(s.workflow_tree(9999).unwrap().is_none());
}
}
#[cfg(test)]
pub mod tempdir {
pub struct Dir {
pub path: std::path::PathBuf,
}
impl Dir {
pub fn new(prefix: &str) -> Dir {
static SEQ: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
let seq = SEQ.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
let n = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos();
let path =
std::env::temp_dir().join(format!("{prefix}-{}-{n}-{seq}", std::process::id()));
std::fs::create_dir_all(&path).unwrap();
Dir { path }
}
}
impl Drop for Dir {
fn drop(&mut self) {
let _ = std::fs::remove_dir_all(&self.path);
}
}
}