use crate::config::Paths;
use crate::desk::{self, Desk, Opened, Origin, Placed};
use crate::peer;
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>,
#[serde(skip_serializing_if = "String::is_empty")]
pub sender: String,
}
#[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,
pub size: i64,
}
#[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,
pub latest: 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 FTS_INSERT: &str =
"INSERT INTO docs_fts(rowid, id, title, body) SELECT rowid, ?1, ?2, ?3 FROM docs WHERE id = ?1";
const FTS_DELETE: &str =
"DELETE FROM docs_fts WHERE rowid = (SELECT rowid FROM docs WHERE id = ?1)";
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, d.sender";
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";
const MIGRATIONS: &[(i64, &str)] = &[
(
1,
"ALTER TABLE docs ADD COLUMN pinned INTEGER NOT NULL DEFAULT 0",
),
(
1,
"ALTER TABLE docs ADD COLUMN origin TEXT NOT NULL DEFAULT 'cli'",
),
(
1,
"ALTER TABLE projects ADD COLUMN renamed INTEGER NOT NULL DEFAULT 0",
),
(
1,
"ALTER TABLE docs ADD COLUMN unread INTEGER NOT NULL DEFAULT 0",
),
(
1,
"ALTER TABLE docs ADD COLUMN deleted_at INTEGER NOT NULL DEFAULT 0",
),
(
1,
"ALTER TABLE docs ADD COLUMN sender TEXT NOT NULL DEFAULT ''",
),
(
1,
"ALTER TABLE docs ADD COLUMN desk_id INTEGER NOT NULL DEFAULT 0",
),
(
1,
"ALTER TABLE docs ADD COLUMN desk_name TEXT NOT NULL DEFAULT ''",
),
(
1,
"ALTER TABLE docs ADD COLUMN desk_slot INTEGER NOT NULL DEFAULT 0",
),
(
1,
"ALTER TABLE panes ADD COLUMN agent_session TEXT NOT NULL DEFAULT ''",
),
(
1,
"ALTER TABLE desk_notes ADD COLUMN done_by TEXT NOT NULL DEFAULT ''",
),
(
1,
"ALTER TABLE panes ADD COLUMN resume_next INTEGER NOT NULL DEFAULT 0",
),
(
1,
"ALTER TABLE panes ADD COLUMN name TEXT NOT NULL DEFAULT ''",
),
(
1,
"ALTER TABLE desks ADD COLUMN full_slot INTEGER NOT NULL DEFAULT 0",
),
(
1,
"ALTER TABLE desk_notes ADD COLUMN done_commit TEXT NOT NULL DEFAULT ''",
),
(
1,
"ALTER TABLE desk_notes ADD COLUMN done_doc TEXT NOT NULL DEFAULT ''",
),
(
1,
"ALTER TABLE desks ADD COLUMN closed_at INTEGER NOT NULL DEFAULT 0",
),
(
1,
"ALTER TABLE desks ADD COLUMN left_off TEXT NOT NULL DEFAULT ''",
),
(
1,
"ALTER TABLE desks ADD COLUMN left_off_at INTEGER NOT NULL DEFAULT 0",
),
(
1,
"ALTER TABLE desks ADD COLUMN left_off_by TEXT NOT NULL DEFAULT ''",
),
(
1,
"ALTER TABLE desks ADD COLUMN left_off_about TEXT NOT NULL DEFAULT ''",
),
(
1,
"ALTER TABLE desk_notes ADD COLUMN done_evidence TEXT NOT NULL DEFAULT ''",
),
(
1,
"ALTER TABLE desk_notes ADD COLUMN suggested_by TEXT NOT NULL DEFAULT ''",
),
(
1,
"ALTER TABLE desks ADD COLUMN visited_at INTEGER NOT NULL DEFAULT 0",
),
(
1,
"ALTER TABLE desks ADD COLUMN parked_at INTEGER NOT NULL DEFAULT 0",
),
(
1,
"ALTER TABLE desks ADD COLUMN parked_next TEXT NOT NULL DEFAULT ''",
),
(
1,
"ALTER TABLE docs ADD COLUMN desk_off INTEGER NOT NULL DEFAULT 0",
),
(
1,
"ALTER TABLE desk_notes ADD COLUMN images TEXT NOT NULL DEFAULT ''",
),
(
1,
"ALTER TABLE desk_notes ADD COLUMN stage TEXT NOT NULL DEFAULT ''",
),
(
1,
"ALTER TABLE desk_notes ADD COLUMN stage_by TEXT NOT NULL DEFAULT ''",
),
(
1,
"ALTER TABLE desk_notes ADD COLUMN stage_doc TEXT NOT NULL DEFAULT ''",
),
(
1,
"ALTER TABLE desk_notes ADD COLUMN stage_at INTEGER NOT NULL DEFAULT 0",
),
(
1,
"ALTER TABLE desk_notes ADD COLUMN stage_pane TEXT NOT NULL DEFAULT ''",
),
(
1,
"ALTER TABLE desk_notes ADD COLUMN stage_session TEXT NOT NULL DEFAULT ''",
),
(
1,
"ALTER TABLE desk_notes ADD COLUMN done_pane TEXT NOT NULL DEFAULT ''",
),
(
1,
"ALTER TABLE desks ADD COLUMN left_off_pane TEXT NOT NULL DEFAULT ''",
),
(2, desk::POS_COLUMN),
(2, "UPDATE desks SET pos = id"),
(3, desk::KIND_COLUMNS[0]),
(3, desk::KIND_COLUMNS[1]),
(4, desk::RETIRE_STUDIO),
(
5,
"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);",
),
(
5,
"CREATE TEMP TABLE fts_at AS
SELECT d.rowid AS r, MAX(f.rowid) AS fr FROM docs_fts f JOIN docs d ON d.id = f.id GROUP BY d.rowid;
CREATE TEMP TABLE fts_rows AS
SELECT t.r AS r, f.id AS id, f.title AS title, substr(f.body, 1, 524288) AS body
FROM fts_at t JOIN docs_fts f ON f.rowid = t.fr;
DELETE FROM docs_fts;
INSERT INTO docs_fts(rowid, id, title, body) SELECT r, id, title, body FROM fts_rows;
DROP TABLE fts_rows;
DROP TABLE fts_at;",
),
(
6,
"ALTER TABLE docs ADD COLUMN is_head INTEGER NOT NULL DEFAULT 1",
),
(
6,
"UPDATE docs SET is_head = CASE
WHEN rowid = (SELECT d2.rowid FROM docs d2
WHERE d2.project_id = docs.project_id AND d2.source_path = docs.source_path
AND d2.deleted_at = 0
ORDER BY d2.received_at DESC, d2.rowid DESC LIMIT 1) THEN 1
ELSE 0 END
WHERE source_path IS NOT NULL",
),
];
fn desk_docs_sql(off: bool) -> String {
format!(
"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.desk_off {}
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",
if off { "> 0" } else { "= 0" },
if off { "d.desk_off DESC," } else { "" },
)
}
fn rehead(conn: &Connection, project_id: i64, source_path: &str) -> rusqlite::Result<usize> {
conn.execute(
"UPDATE docs SET is_head = CASE
WHEN rowid = (SELECT d2.rowid FROM docs d2
WHERE d2.project_id = ?1 AND d2.source_path = ?2 AND d2.deleted_at = 0
ORDER BY d2.received_at DESC, d2.rowid DESC LIMIT 1) THEN 1
ELSE 0 END
WHERE project_id = ?1 AND source_path = ?2",
params![project_id, source_path],
)
}
fn migrate(conn: &Connection) -> Result<()> {
let have: i64 = conn
.query_row("PRAGMA user_version", [], |r| r.get(0))
.context("reading the schema version")?;
let latest = MIGRATIONS.last().map(|(v, _)| *v).unwrap_or(0);
if have >= latest {
return Ok(());
}
let tx = conn.unchecked_transaction()?;
for (version, stmt) in MIGRATIONS.iter().filter(|(v, _)| *v > have) {
match tx.execute_batch(stmt) {
Ok(()) => {}
Err(e) if *version == 1 && e.to_string().contains("duplicate column name") => {}
Err(e) => return Err(e).with_context(|| format!("schema step {version}: {stmt}")),
}
}
tx.execute_batch(&format!("PRAGMA user_version = {latest}"))?;
tx.commit()?;
Ok(())
}
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(peer::SCHEMA)?;
migrate(&conn)?;
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);
CREATE INDEX IF NOT EXISTS docs_desk ON docs(desk_id, desk_off, received_at);
CREATE INDEX IF NOT EXISTS docs_head ON docs(project_id, is_head, deleted_at, 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.is_head = 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, is_head)
VALUES(?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, 0, ?12, 1, ?13, ?14, ?15, ?16, 1)",
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(FTS_INSERT, 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],
)?;
rehead(&tx, project_id, sp)?;
}
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(),
sender: d.sender.to_string(),
})
}
pub fn replace(&self, id: &str, d: NewDoc) -> Result<Doc> {
let now = now();
let (hash, size) = self.put_source(id, &d)?;
let _ = fs::remove_file(self.outline_path(id));
let mut conn = self.conn.lock().unwrap();
fs::write(self.html_path(id), d.html)?;
let tx = conn.transaction()?;
let before: Option<(String, i64, Option<String>)> = tx
.query_row(
"SELECT content_hash, project_id, source_path FROM docs WHERE id = ?1",
params![id],
|r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)),
)
.optional()?;
let Some((old_hash, project_id, source_path)) = before else {
anyhow::bail!("replaced document vanished");
};
tx.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],
)?;
if old_hash != hash {
tx.execute(FTS_DELETE, params![id])?;
tx.execute(FTS_INSERT, params![id, d.title, d.search_body])?;
}
if let Some(sp) = &source_path {
rehead(&tx, project_id, sp)?;
}
tx.commit()?;
drop(conn);
self.get(id)?.context("replaced document vanished")
}
pub fn replace_html_if(&self, id: &str, html: &str, hash: &str) -> Result<bool> {
let conn = self.conn.lock().unwrap();
let now: Option<String> = conn
.query_row(
"SELECT content_hash FROM docs WHERE id = ?1",
params![id],
|r| r.get(0),
)
.optional()?;
if now.as_deref() != Some(hash) {
return Ok(false);
}
fs::write(self.html_path(id), html)?;
Ok(true)
}
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)
}
#[cfg(test)]
pub fn delete(&self, id: &str) -> Result<bool> {
Ok(self.delete_versions(id)? > 0)
}
pub fn delete_versions(&self, id: &str) -> Result<usize> {
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(0);
};
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],
)?,
};
if let Some(sp) = &source_path {
rehead(&tx, project_id, sp)?;
}
tx.commit()?;
Ok(n)
}
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],
)?,
};
if let Some(sp) = &source_path {
rehead(&tx, project_id, sp)?;
}
tx.commit()?;
Ok(n > 0)
}
pub fn removed(&self, limit: usize, desks: bool) -> Result<Vec<Removed>> {
let conn = self.conn.lock().unwrap();
let mut out: Vec<Removed> = conn
.prepare(
"SELECT id, title, name, deleted_at, n FROM (
SELECT d.id, d.title, p.name, d.deleted_at, COUNT(*) OVER w AS n,
ROW_NUMBER() OVER (w ORDER BY d.received_at DESC, d.rowid DESC) AS k
FROM docs d JOIN projects p ON p.id = d.project_id WHERE d.deleted_at != 0
WINDOW w AS (PARTITION BY d.project_id, COALESCE(d.source_path, d.id), d.deleted_at))
WHERE k = 1 ORDER BY deleted_at DESC LIMIT ?1",
)?
.query_map(params![limit as i64], |r| {
let id: String = r.get(0)?;
Ok(Removed {
kind: "doc",
restore: format!("/api/docs/{id}/undelete"),
id,
title: r.get(1)?,
from: r.get(2)?,
desk: None,
at: r.get(3)?,
versions: r.get::<_, i64>(4)? as usize,
})
})?
.collect::<std::result::Result<_, _>>()?;
if desks {
for (desk, id, name, text, at) in desk::removed_notes(&conn, limit)? {
out.push(Removed {
kind: "note",
id: id.to_string(),
title: text,
from: name,
desk: Some(desk),
at,
versions: 1,
restore: format!("/api/desks/{desk}/notes/{id}/restore"),
});
}
for (id, name, root, at) in desk::closed_desks(&conn, limit)? {
out.push(Removed {
kind: "desk",
id: id.to_string(),
title: name,
from: root,
desk: Some(id),
at,
versions: 1,
restore: format!("/api/desks/{id}/reopen"),
});
}
for (id, desk, name, label, at) in desk::closed_panes(&conn, limit)? {
out.push(Removed {
kind: "panel",
restore: format!("/api/panes/{id}/restore"),
id,
title: label,
from: name,
desk: Some(desk),
at,
versions: 1,
});
}
}
out.sort_by_key(|a| std::cmp::Reverse(a.at));
out.truncate(limit);
Ok(out)
}
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 mark_unread(&self, ids: &[String]) -> Result<Vec<String>> {
let mut conn = self.conn.lock().unwrap();
let tx = conn.transaction()?;
let mut back = Vec::new();
for id in ids {
if tx.execute(
"UPDATE docs SET unread = 1 WHERE id = ?1 AND unread = 0 AND deleted_at = 0",
params![id],
)? > 0
{
back.push(id.clone());
}
}
tx.commit()?;
Ok(back)
}
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_project_by_root(&self, root: &str, name: &str) -> Result<bool> {
let conn = self.conn.lock().unwrap();
Ok(conn.execute(
"UPDATE projects SET name = ?2 WHERE root = ?1 AND renamed = 0",
params![root, 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()?;
let mut files: std::collections::BTreeSet<(i64, String)> = Default::default();
for (id, _) in &victims {
let file: Option<(i64, Option<String>)> = tx
.query_row(
"SELECT project_id, source_path FROM docs WHERE id = ?1",
params![id],
|r| Ok((r.get(0)?, r.get(1)?)),
)
.optional()?;
if let Some((pid, Some(sp))) = file {
files.insert((pid, sp));
}
tx.execute(FTS_DELETE, params![id])?;
tx.execute("DELETE FROM docs WHERE id = ?1", params![id])?;
}
for (pid, sp) in &files {
rehead(&tx, *pid, sp)?;
}
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), MAX(d.received_at)
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)?,
latest: r.get(5)?,
})
})?
.collect::<std::result::Result<_, _>>()?;
Ok(rows)
}
pub fn project_row(&self, project_id: i64) -> Result<Option<TreeProject>> {
let conn = self.conn.lock().unwrap();
conn.query_row(
"SELECT p.id, p.name, p.root, COUNT(d.id), COUNT(DISTINCT d.workflow_id), MAX(d.received_at)
FROM projects p JOIN head_docs d ON d.project_id = p.id
WHERE p.id = ?1 GROUP BY p.id",
params![project_id],
|r| {
Ok(TreeProject {
id: r.get(0)?,
name: r.get(1)?,
root: r.get(2)?,
docs: r.get(3)?,
workflows: r.get(4)?,
latest: r.get(5)?,
})
},
)
.optional()
.map_err(Into::into)
}
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, size 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, size 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, off: bool) -> Result<Vec<DeskDoc>> {
let conn = self.conn.lock().unwrap();
let rows = conn
.prepare(&desk_docs_sql(off))?
.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 set_desk_doc_off(&self, desk_id: i64, doc: &str, off: bool) -> Result<bool> {
let n = self.conn.lock().unwrap().execute(
"UPDATE docs SET desk_off = ?3 WHERE id = ?2 AND desk_id = ?1 AND deleted_at = 0",
params![desk_id, doc, if off { now() } else { 0 }],
)?;
Ok(n > 0)
}
pub fn create_desk(&self, root: &str, name: Option<&str>) -> Result<Desk> {
desk::create(&self.conn.lock().unwrap(), root, name, now())
}
pub fn reorder_desks(&self, ids: &[i64]) -> Result<bool> {
desk::reorder(&mut self.conn.lock().unwrap(), ids)
}
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, full: Option<i64>) -> Result<bool> {
desk::layout(&self.conn.lock().unwrap(), id, col, row, full)
}
pub fn close_desk(&self, id: i64) -> Result<Option<Vec<String>>> {
desk::close(&mut self.conn.lock().unwrap(), id, now())
}
pub fn reopen_desk(&self, id: i64) -> Result<bool> {
desk::reopen(&mut self.conn.lock().unwrap(), id)
}
pub fn prune_desks(&self, before: i64, dry_run: bool) -> Result<Vec<(i64, String)>> {
desk::prune_desks(&self.conn.lock().unwrap(), before, dry_run)
}
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 mark_panes_resume(&self, ids: &[String]) -> Result<usize> {
desk::mark_resume(&self.conn.lock().unwrap(), ids)
}
pub fn panes_resume(&self) -> Result<(Vec<String>, Vec<String>)> {
desk::read_resume(&self.conn.lock().unwrap())
}
pub fn clear_panes_resume(&self) -> Result<()> {
desk::clear_resume(&self.conn.lock().unwrap())
}
pub fn offer_panes_resume(&self, ids: &[String]) -> Result<usize> {
desk::mark_offer(&self.conn.lock().unwrap(), ids)
}
pub fn peers(&self) -> Result<Vec<peer::Peer>> {
peer::list(&self.conn.lock().unwrap())
}
pub fn peer(&self, id: i64) -> Result<Option<peer::Peer>> {
peer::get(&self.conn.lock().unwrap(), id)
}
pub fn peer_by_key(&self, key: &str) -> Result<Option<peer::Peer>> {
peer::by_sign_key(&self.conn.lock().unwrap(), key)
}
pub fn peer_by_name(&self, name: &str) -> Result<Option<peer::Peer>> {
peer::by_name(&self.conn.lock().unwrap(), name)
}
pub fn pin_peer(&self, p: &peer::Peer) -> Result<peer::Peer> {
peer::pin(&self.conn.lock().unwrap(), p, now())
}
pub fn rename_peer(&self, id: i64, name: &str) -> Result<bool> {
peer::rename(&self.conn.lock().unwrap(), id, name)
}
pub fn mute_peer(&self, id: i64, muted: bool) -> Result<bool> {
peer::mute(&self.conn.lock().unwrap(), id, muted)
}
pub fn remove_peer(&self, id: i64) -> Result<bool> {
peer::remove(&self.conn.lock().unwrap(), id, now())
}
pub fn restore_peer(&self, id: i64) -> Result<bool> {
peer::restore(&self.conn.lock().unwrap(), id)
}
pub fn touch_peer(&self, id: i64, from: bool) -> Result<()> {
peer::touch(&self.conn.lock().unwrap(), id, from, now())
}
pub fn peer_queue(&self, p: &peer::Peer, doc_id: &str) -> Result<String> {
peer::queue(&self.conn.lock().unwrap(), p, doc_id, now())
}
pub fn peer_unsent(&self) -> Result<Vec<(String, i64, String, i64)>> {
peer::unsent(&self.conn.lock().unwrap())
}
pub fn peer_sent(&self, id: &str) -> Result<()> {
peer::sent(&self.conn.lock().unwrap(), id, now())
}
pub fn peer_failed(&self, id: &str, why: &str) -> Result<()> {
peer::failed(&self.conn.lock().unwrap(), id, why)
}
pub fn peer_notes(&self) -> Result<Vec<peer::PeerNote>> {
peer::notes_waiting(&self.conn.lock().unwrap())
}
pub fn peer_note_arrived(&self, peer_id: i64, text: &str) -> Result<i64> {
peer::note_arrived(&self.conn.lock().unwrap(), peer_id, text, now())
}
pub fn settle_peer_note(&self, id: i64, what: &str) -> Result<bool> {
peer::settle_note(&self.conn.lock().unwrap(), id, what, now())
}
pub fn peer_offers(&self) -> Result<Vec<peer::Offer>> {
peer::offers_open(&self.conn.lock().unwrap())
}
pub fn peer_offer(&self, peer_id: i64, doc_id: &str, pane: &str, by: &str) -> Result<i64> {
peer::offer(&self.conn.lock().unwrap(), peer_id, doc_id, pane, by, now())
}
pub fn peer_offer_get(&self, id: i64) -> Result<Option<peer::Offer>> {
peer::offer_get(&self.conn.lock().unwrap(), id)
}
pub fn answer_peer_offer(&self, id: i64, sent: bool) -> Result<bool> {
peer::answer_offer(&self.conn.lock().unwrap(), id, sent, now())
}
pub fn drop_peer_offers_of(&self, pane: &str) -> Result<usize> {
peer::drop_offers_of(&self.conn.lock().unwrap(), pane, now())
}
pub fn peer_taken(&self, id: &str) -> Result<bool> {
peer::taken(&self.conn.lock().unwrap(), id)
}
pub fn peer_take(&self, id: &str) -> Result<()> {
peer::take(&self.conn.lock().unwrap(), id, now())
}
pub fn prune_peer_taken(&self) -> Result<usize> {
peer::prune_taken(&self.conn.lock().unwrap(), now() - peer::TAKEN_KEPT)
}
pub fn set_pane_cwd(&self, id: &str, cwd: &str) -> Result<bool> {
desk::set_cwd(&self.conn.lock().unwrap(), id, cwd)
}
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 move_pane(&self, desk_id: i64, from: i64, to: i64) -> Result<bool> {
let mut conn = self.conn.lock().unwrap();
let tx = conn.transaction()?;
if !desk::move_pane(&tx, desk_id, from, to)? {
return Ok(false);
}
tx.execute(
"UPDATE docs SET desk_slot = CASE desk_slot WHEN ?2 THEN ?3 ELSE ?2 END
WHERE desk_id = ?1 AND desk_slot IN (?2, ?3)",
params![desk_id, from, to],
)?;
tx.commit()?;
Ok(true)
}
pub fn close_pane(&self, id: &str) -> Result<bool> {
let mut conn = self.conn.lock().unwrap();
let tx = conn.transaction()?;
let Some((desk_id, slot)) = desk::close_pane(&tx, id, now())? else {
return Ok(false);
};
tx.execute(
"UPDATE docs SET desk_slot = CASE WHEN desk_slot = ?2 THEN 0 ELSE desk_slot - 1 END
WHERE desk_id = ?1 AND desk_slot >= ?2",
params![desk_id, slot],
)?;
tx.commit()?;
Ok(true)
}
pub fn restore_pane(&self, id: &str) -> Result<desk::Restored> {
desk::restore_pane(&mut self.conn.lock().unwrap(), id)
}
pub fn rename_pane(&self, id: &str, name: &str) -> Result<bool> {
desk::rename_pane(&self.conn.lock().unwrap(), id, name)
}
pub fn prune_panes(&self, before: i64, dry_run: bool) -> Result<Vec<(String, String)>> {
desk::prune_closed(&self.conn.lock().unwrap(), before, dry_run)
}
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 desk_keys(&self, desk_id: i64) -> Result<Vec<desk::DeskKey>> {
desk::keys(&self.conn.lock().unwrap(), desk_id)
}
pub fn add_desk_key(&self, desk_id: i64, name: &str, provider: &str) -> Result<()> {
desk::add_key(&self.conn.lock().unwrap(), desk_id, name, provider, now())
}
pub fn remove_desk_key(&self, desk_id: i64, name: &str) -> Result<bool> {
desk::remove_key(&self.conn.lock().unwrap(), desk_id, name)
}
pub fn touch_desk_keys(&self, keys: &[desk::DeskKey]) -> Result<()> {
desk::touch_keys(&self.conn.lock().unwrap(), keys, now())
}
pub fn prune_desk_keys(&self, before: i64, dry_run: bool) -> Result<Vec<(i64, String)>> {
desk::prune_keys(&self.conn.lock().unwrap(), before, dry_run)
}
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 tick_desk_note(&self, desk_id: i64, id: i64, tick: &desk::Tick) -> Result<bool> {
desk::tick_note(&self.conn.lock().unwrap(), desk_id, id, tick, now())
}
pub fn removed_desk_notes_since(&self, desk_id: i64, since: i64) -> Result<Vec<(i64, String)>> {
desk::removed_since(&self.conn.lock().unwrap(), desk_id, since)
}
pub fn mark_desk_note(&self, desk_id: i64, id: i64, mark: &desk::Mark) -> Result<bool> {
desk::mark_note(&self.conn.lock().unwrap(), desk_id, id, mark, now())
}
pub fn add_note_image(&self, desk_id: i64, id: i64, name: &str) -> Result<Option<Vec<String>>> {
desk::add_note_image(&self.conn.lock().unwrap(), desk_id, id, name)
}
pub fn set_note_images(&self, desk_id: i64, id: i64, images: &[String]) -> Result<bool> {
desk::set_note_images(&self.conn.lock().unwrap(), desk_id, id, images)
}
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 suggest_desk_note(&self, desk_id: i64, text: &str, by: &str) -> Result<desk::Suggested> {
desk::suggest_note(&mut self.conn.lock().unwrap(), desk_id, text, by, now())
}
pub fn keep_desk_note(&self, desk_id: i64, id: i64) -> Result<bool> {
desk::keep_note(&self.conn.lock().unwrap(), desk_id, id)
}
pub fn visit_desk(&self, id: i64) -> Result<bool> {
desk::visit(&self.conn.lock().unwrap(), id, now())
}
pub fn park_desk(&self, id: i64, next: Option<&str>) -> Result<Option<Option<desk::Parked>>> {
let to = next.map(|n| desk::Parked {
at: now(),
next: n.to_string(),
});
desk::park(&self.conn.lock().unwrap(), id, to.as_ref())
}
pub fn desks_done_since(&self, since: i64) -> Result<Vec<desk::Done>> {
desk::done_since(&self.conn.lock().unwrap(), since)
}
pub fn desks_sent_since(&self, since: i64) -> Result<Vec<(i64, String, String, i64)>> {
let conn = self.conn.lock().unwrap();
let mut st = conn.prepare(
"SELECT desk_id, id, title, received_at FROM live_docs
WHERE desk_id != 0 AND received_at >= ?1 ORDER BY received_at, rowid",
)?;
let rows = st.query_map(params![since], |r| {
Ok((r.get(0)?, r.get(1)?, r.get(2)?, r.get(3)?))
})?;
Ok(rows.collect::<std::result::Result<_, _>>()?)
}
pub fn set_left_off(
&self,
desk_id: i64,
to: &desk::LeftOff,
) -> Result<Option<Option<desk::LeftOff>>> {
let mut to = to.clone();
if to.at == 0 {
to.at = now();
}
desk::set_left_off(&self.conn.lock().unwrap(), desk_id, &to)
}
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 WHERE closed_at = 0) 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)?;
peer::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,
size: r.get(6)?,
})
}
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)?,
}),
},
sender: r.get(19)?,
})
}
#[derive(Debug, Serialize)]
pub struct Removed {
pub kind: &'static str,
pub id: String,
pub title: String,
pub from: String,
pub desk: Option<i64>,
pub at: i64,
pub versions: usize,
pub restore: String,
}
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;
#[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);
}
}
}