use crate::project;
use crate::render::{self, Kind, Renderer, HIGHLIGHT_CAP};
use crate::store::{Doc, NewDoc, Staged, Store};
use anyhow::{bail, Context, Result};
use serde::{Deserialize, Serialize};
use std::path::{Path, PathBuf};
#[derive(Clone, Debug, Default, Deserialize, Serialize)]
pub struct Payload {
pub path: Option<String>,
pub content: Option<String>,
pub title: Option<String>,
pub workflow: Option<String>,
pub lang: Option<String>,
pub cwd: Option<String>,
pub session: Option<String>,
pub origin: Option<String>,
#[serde(default)]
pub sender: Option<String>,
#[serde(default)]
pub pane: Option<String>,
#[serde(skip)]
pub peer: Option<FromPeer>,
}
#[derive(Clone, Debug, Default)]
pub struct FromPeer {
pub name: String,
pub sign_key: String,
pub bytes: Vec<u8>,
pub file: Option<String>,
pub desk: Option<(crate::desk::Origin, String)>,
pub root: Option<String>,
pub lineage: Option<String>,
}
pub struct Received {
pub doc: Doc,
pub needs_full_highlight: bool,
pub existing: bool,
pub supersedes: Option<String>,
}
pub const MAX_BYTES: usize = 32 * 1024 * 1024;
pub const MAX_MEDIA_BYTES: u64 = 1024 * 1024 * 1024;
const COALESCE_SECS: i64 = 180;
pub const SEARCH_CAP: usize = 512 * 1024;
pub fn search_text(text: &str) -> &str {
if text.len() <= SEARCH_CAP {
return text;
}
let mut end = SEARCH_CAP;
while !text.is_char_boundary(end) {
end -= 1;
}
&text[..end]
}
fn automatic(origin: &str) -> bool {
matches!(origin, "hook" | "watch")
}
struct Body {
bytes: Vec<u8>,
text: String,
path: Option<String>,
staged: Option<Staged>,
}
fn read(store: &Store, p: &Payload) -> Result<Body> {
let body = match (&p.content, &p.path) {
_ if p.peer.is_some() => {
let fp = p.peer.as_ref().unwrap();
let name = fp.file.as_deref().unwrap_or("");
let ext = render::ext_of(name);
let opaque = render::is_image_ext(&ext)
|| render::preview_kind(&ext) == Some("pdf")
|| render::looks_binary(&fp.bytes);
Body {
text: if opaque {
String::new()
} else {
String::from_utf8_lossy(&fp.bytes).into_owned()
},
bytes: fp.bytes.clone(),
path: fp
.lineage
.clone()
.or_else(|| fp.file.clone())
.filter(|f| !f.trim().is_empty()),
staged: None,
}
}
(Some(c), _) => Body {
bytes: c.clone().into_bytes(),
text: c.clone(),
path: p.path.clone(),
staged: None,
},
(None, Some(path)) => {
let path = absolutize(path, p.cwd.as_deref());
let sp = path.to_string_lossy().to_string();
if render::media_kind(&render::ext_of(&sp)).is_some() {
Body {
bytes: Vec::new(),
text: String::new(),
path: Some(sp),
staged: Some(store.stage(&path, MAX_MEDIA_BYTES)?),
}
} else {
let bytes =
std::fs::read(&path).with_context(|| format!("reading {}", path.display()))?;
if bytes.len() > MAX_BYTES {
bail!("file is larger than {} MB", MAX_BYTES / 1024 / 1024);
}
let ext = render::ext_of(&sp);
let opaque = render::is_image_ext(&ext)
|| render::preview_kind(&ext) == Some("pdf")
|| render::looks_binary(&bytes);
let text = if opaque {
String::new()
} else {
String::from_utf8_lossy(&bytes).into_owned()
};
Body {
bytes,
text,
path: Some(sp),
staged: None,
}
}
}
(None, None) => bail!("send_document needs either `path` or `content`"),
};
if body.bytes.len() > MAX_BYTES {
bail!("content is larger than {} MB", MAX_BYTES / 1024 / 1024);
}
Ok(body)
}
fn classify(renderer: &Renderer, b: &Body, lang: Option<&str>) -> (Kind, Option<String>) {
match renderer.detect(b.path.as_deref(), lang, &b.text) {
_ if b.staged.is_some() => {
let ext = render::ext_of(b.path.as_deref().unwrap_or_default());
match render::media_kind(&ext) {
Some("video") => (Kind::Video, Some(ext)),
_ => (Kind::Audio, Some(ext)),
}
}
(Kind::Video | Kind::Audio, _) => (Kind::Text, None),
(_, lang) if !b.text.is_empty() && render::looks_binary(&b.bytes) => (Kind::Binary, lang),
(Kind::Image, lang) => (Kind::Image, lang),
_ if b.text.is_empty() && !b.bytes.is_empty() => (Kind::Binary, None),
other => other,
}
}
fn page(
renderer: &Renderer,
b: &Body,
text: &str,
id: &str,
kind: Kind,
lang: Option<&str>,
title: &str,
) -> String {
let body_src = if kind == Kind::Markdown {
render::strip_leading_h1(text, title)
} else {
None
};
let file_base = b.path.as_ref().map(|_| format!("/files/{id}/"));
match kind {
Kind::Image => render::image_body(&format!("/api/docs/{id}/blob"), title),
Kind::Video | Kind::Audio => render::media_body(
&format!("/api/docs/{id}/blob"),
&render::ext_of(b.path.as_deref().unwrap_or_default()),
),
Kind::Binary => render::placeholder(&render::describe_bytes(title, b.bytes.len() as u64)),
_ => renderer.render_with_files(
kind,
lang,
body_src.as_deref().unwrap_or(text),
file_base.as_deref(),
b.path.as_deref().and_then(|p| Path::new(p).parent()),
),
}
}
fn workflow(
p: &Payload,
origin: &str,
from: Option<&crate::desk::Origin>,
project: &str,
title: &str,
) -> (String, String) {
let plan_home = (origin == "plan" && p.workflow.is_none()).then(|| {
from.map(|o| o.name.clone())
.unwrap_or_else(|| project.to_string())
});
if let Some(fp) = &p.peer {
return ("sent".to_string(), format!("Sent by {}", fp.name));
}
match (plan_home.as_ref().or(p.workflow.as_ref()), &p.session) {
(Some(w), _) if !w.trim().is_empty() => (w.trim().to_string(), w.trim().to_string()),
(_, Some(s)) if !s.trim().is_empty() => (s.trim().to_string(), title.to_string()),
_ => ("manual".to_string(), "Sent manually".to_string()),
}
}
fn place(p: &Payload, b: &Body) -> (String, String, Option<String>) {
if let Some(root) = p.peer.as_ref().and_then(|fp| fp.root.as_ref()) {
return desk_project(root);
}
if let Some((_, root)) = p.peer.as_ref().and_then(|fp| fp.desk.as_ref()) {
return desk_project(root);
}
if let Some(fp) = &p.peer {
return (
format!("peer:{}", fp.sign_key),
format!("From {}", fp.name),
None,
);
}
let anchor: PathBuf = p
.cwd
.as_deref()
.map(PathBuf::from)
.or_else(|| b.path.as_deref().map(PathBuf::from))
.unwrap_or_else(|| std::env::current_dir().unwrap_or_else(|_| PathBuf::from("/")));
let proj = project::resolve(&anchor);
let branch = project::branch(&proj.root);
(proj.root.to_string_lossy().to_string(), proj.name, branch)
}
pub fn desk_project(root: &str) -> (String, String, Option<String>) {
let proj = project::resolve(Path::new(root));
let branch = project::branch(&proj.root);
(proj.root.to_string_lossy().to_string(), proj.name, branch)
}
pub fn pane_origin(store: &Store, pane: Option<&str>) -> Option<crate::desk::Origin> {
pane.filter(|id| crate::pane::valid_id(id))
.and_then(|id| store.pane(id).ok().flatten())
.map(|placed| crate::desk::Origin {
id: placed.desk_id,
name: placed.desk_name,
slot: placed.pane.slot,
})
}
pub fn receive(store: &Store, renderer: &Renderer, p: Payload) -> Result<Received> {
let b = read(store, &p)?;
let origin = p.origin.as_deref().unwrap_or("cli");
let from = pane_origin(store, p.pane.as_deref()).or_else(|| {
p.peer
.as_ref()
.and_then(|fp| fp.desk.as_ref().map(|(o, _)| o.clone()))
});
let (root, proj_name, branch) = place(&p, &b);
let hash = match &b.staged {
Some(st) => st.hash.clone(),
None => blake3::hash(&b.bytes).to_hex().to_string(),
};
let latest_same_path = match &b.path {
Some(sp) => store.latest_for_path(&root, sp)?,
None => None,
};
if let Some(existing) = &latest_same_path {
if existing.content_hash == hash {
return Ok(Received {
doc: existing.clone(),
needs_full_highlight: false,
existing: true,
supersedes: None,
});
}
}
let (kind, lang) = classify(renderer, &b, p.lang.as_deref());
let title = render::title_for(p.title.as_deref(), kind, b.path.as_deref(), &b.text);
let (wf_name, wf_title) = workflow(&p, origin, from.as_ref(), &proj_name, &title);
let wf_key = wf_name.to_lowercase();
let coalesce_into = match (origin, &latest_same_path) {
(o, Some(prev))
if automatic(o)
&& automatic(&prev.origin)
&& prev.workflow == wf_key
&& crate::store::now() - prev.received_at < COALESCE_SECS =>
{
Some(prev.id.clone())
}
_ => None,
};
let id = coalesce_into
.clone()
.unwrap_or_else(|| crate::store::new_id(&hash));
let shown = match &p.peer {
Some(fp) if kind == Kind::Markdown => Some(render::stayed_with(&b.text, &fp.name)),
_ => None,
};
let html = page(
renderer,
&b,
shown.as_deref().unwrap_or(&b.text),
&id,
kind,
lang.as_deref(),
&title,
);
let needs_full_highlight = kind == Kind::Code && b.text.len() > HIGHLIGHT_CAP;
let new_doc = NewDoc {
project_root: &root,
project_name: &proj_name,
workflow_key: &wf_key,
workflow_title: &wf_title,
title: &title,
kind,
lang: lang.as_deref(),
source_path: b.path.as_deref(),
branch: branch.as_deref(),
origin,
sender: p.sender.as_deref().unwrap_or(""),
desk: from.as_ref(),
source: &b.bytes,
staged: b.staged.as_ref(),
search_body: search_text(&b.text),
html: &html,
};
if coalesce_into.is_some() {
let doc = store.replace(&id, new_doc)?;
return Ok(Received {
doc,
needs_full_highlight,
existing: true,
supersedes: None,
});
}
let supersedes = latest_same_path.as_ref().map(|prev| prev.id.clone());
let doc = store.insert(&id, new_doc)?;
Ok(Received {
doc,
needs_full_highlight,
existing: false,
supersedes,
})
}
fn absolutize(path: &str, cwd: Option<&str>) -> PathBuf {
let p = Path::new(path);
if p.is_absolute() {
return p.to_path_buf();
}
match cwd {
Some(c) => Path::new(c).join(p),
None => std::env::current_dir()
.map(|d| d.join(p))
.unwrap_or_else(|_| p.to_path_buf()),
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::config::Paths;
use crate::store::tempdir::Dir;
fn setup() -> (Store, Renderer, Dir) {
let dir = Dir::new("snyvi-recv");
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(), Renderer::new(), dir)
}
#[test]
fn the_indexed_text_is_capped_on_a_character() {
let short = "a plan";
assert_eq!(search_text(short), short);
let long = "é".repeat(SEARCH_CAP);
let cut = search_text(&long);
assert!(cut.len() <= SEARCH_CAP);
assert!(cut.len() >= SEARCH_CAP - 1);
assert!(cut.chars().all(|c| c == 'é'));
}
#[test]
fn inline_markdown_gets_title_from_h1_and_session_workflow() {
let (s, r, _d) = setup();
let p = Payload {
content: Some("# My Plan\n\nbody".into()),
session: Some("sess-1".into()),
cwd: Some("/tmp".into()),
..Default::default()
};
let got = receive(&s, &r, p).unwrap();
assert_eq!(got.doc.title, "My Plan");
assert_eq!(got.doc.workflow, "sess-1");
assert_eq!(got.doc.workflow_title, "My Plan");
assert!(!got.existing);
let html = s.html(&got.doc.id).unwrap();
assert!(
!html.contains("<h1"),
"leading H1 stripped from body: {html}"
);
assert!(html.contains("<p>body</p>"));
}
#[test]
fn path_send_dedups_identical_content_and_coalesces_hook_edits() {
let (s, r, d) = setup();
let file = d.path.join("NOTES.md");
std::fs::write(&file, "# Notes\n\nv1").unwrap();
let path = file.to_string_lossy().to_string();
let cwd = d.path.to_string_lossy().to_string();
let mk = |origin: &str| Payload {
path: Some(path.clone()),
cwd: Some(cwd.clone()),
session: Some("s".into()),
origin: Some(origin.into()),
..Default::default()
};
let first = receive(&s, &r, mk("hook")).unwrap();
let again = receive(&s, &r, mk("mcp")).unwrap();
assert!(again.existing);
assert_eq!(
again.doc.id, first.doc.id,
"identical bytes are not stored twice"
);
std::fs::write(&file, "# Notes\n\nv2").unwrap();
let edited = receive(&s, &r, mk("hook")).unwrap();
assert!(edited.existing, "hook edit within the window overwrites");
assert_eq!(edited.doc.id, first.doc.id);
assert_eq!(s.source(&first.doc.id).unwrap(), "# Notes\n\nv2");
std::fs::write(&file, "# Notes\n\nv3").unwrap();
let explicit = receive(&s, &r, mk("mcp")).unwrap();
assert!(
!explicit.existing,
"an explicit send is always a new version"
);
assert_ne!(explicit.doc.id, first.doc.id);
assert_eq!(s.count().unwrap(), 2);
assert_eq!(explicit.supersedes.as_deref(), Some(first.doc.id.as_str()));
assert!(
again.supersedes.is_none(),
"the same bytes supersede nothing"
);
assert!(
edited.supersedes.is_none(),
"nor does an overwrite in place"
);
std::fs::write(&file, "# Notes\n\nv4").unwrap();
let watched = receive(&s, &r, mk("watch")).unwrap();
assert!(!watched.existing, "a watch send after an mcp send is new");
assert_eq!(
watched.supersedes.as_deref(),
Some(explicit.doc.id.as_str())
);
std::fs::write(&file, "# Notes\n\nv5").unwrap();
let again = receive(&s, &r, mk("watch")).unwrap();
assert!(again.existing);
assert_eq!(again.doc.id, watched.doc.id);
std::fs::write(&file, "# Notes\n\nv6").unwrap();
let hooked = receive(&s, &r, mk("hook")).unwrap();
assert!(hooked.existing, "hook and watch coalesce with each other");
assert_eq!(hooked.doc.id, watched.doc.id);
assert_eq!(s.source(&watched.doc.id).unwrap(), "# Notes\n\nv6");
assert_eq!(s.count().unwrap(), 3);
}
#[test]
fn images_keep_their_bytes_and_binaries_are_described() {
let (s, r, d) = setup();
let cwd = d.path.to_string_lossy().to_string();
let png: &[u8] = b"\x89PNG\r\n\x1a\n\x00\x00\x00\rIHDR\x00\x00\x00\x01";
let img = d.path.join("shot.png");
std::fs::write(&img, png).unwrap();
let got = receive(
&s,
&r,
Payload {
path: Some(img.to_string_lossy().to_string()),
cwd: Some(cwd.clone()),
..Default::default()
},
)
.unwrap();
assert_eq!(got.doc.kind, Kind::Image);
assert_eq!(
std::fs::read(s.src_path(&got.doc.id)).unwrap(),
png,
"the stored bytes are the file, not a lossy decode"
);
let html = s.html(&got.doc.id).unwrap();
assert!(html.contains("<img"), "an image document displays: {html}");
assert!(html.contains(&format!("/api/docs/{}/blob", got.doc.id)));
let blob = d.path.join("sheet.xlsx");
std::fs::write(&blob, b"PK\x03\x04\x00\x00rest of a zip").unwrap();
let got = receive(
&s,
&r,
Payload {
path: Some(blob.to_string_lossy().to_string()),
cwd: Some(cwd),
..Default::default()
},
)
.unwrap();
assert_eq!(got.doc.kind, Kind::Binary);
let html = s.html(&got.doc.id).unwrap();
assert!(
html.contains("binary file"),
"described, not decoded: {html}"
);
assert!(
!html.contains("PK"),
"the bytes never reach the page: {html}"
);
}
#[test]
fn media_is_staged_into_the_store_and_played() {
let (s, r, d) = setup();
let cwd = d.path.to_string_lossy().to_string();
let bytes: Vec<u8> = (0..300_000u32).map(|i| (i % 251) as u8).collect();
let clip = d.path.join("clip.mp4");
std::fs::write(&clip, &bytes).unwrap();
let send = |p: &std::path::Path| Payload {
path: Some(p.to_string_lossy().to_string()),
cwd: Some(cwd.clone()),
lang: Some("rust".into()),
..Default::default()
};
let got = receive(&s, &r, send(&clip)).unwrap();
assert_eq!(
got.doc.kind,
Kind::Video,
"a language hint does not unmake a video"
);
assert_eq!(got.doc.size, bytes.len() as i64);
assert_eq!(
got.doc.content_hash,
blake3::hash(&bytes).to_hex().to_string()
);
assert_eq!(std::fs::read(s.src_path(&got.doc.id)).unwrap(), bytes);
let html = s.html(&got.doc.id).unwrap();
assert!(
html.contains("<video") && html.contains(&format!("/api/docs/{}/blob", got.doc.id)),
"{html}"
);
let again = receive(&s, &r, send(&clip)).unwrap();
assert!(again.existing);
let stray = std::fs::read_dir(d.path.join("docs"))
.unwrap()
.flatten()
.filter(|e| e.file_name().to_string_lossy().starts_with(".stage-"))
.count();
assert_eq!(stray, 0, "an unused copy is removed");
let song = d.path.join("take.ogg");
std::fs::write(&song, b"OggS\x00\x02").unwrap();
let got = receive(&s, &r, send(&song)).unwrap();
assert_eq!(got.doc.kind, Kind::Audio);
assert!(s.html(&got.doc.id).unwrap().contains("<audio"));
let err = s.stage(&clip, 1000).unwrap_err().to_string();
assert!(
err.contains("clip.mp4") && err.contains("snyvi browse"),
"{err}"
);
let got = receive(
&s,
&r,
Payload {
content: Some("notes about the cut".into()),
path: Some(clip.to_string_lossy().to_string()),
cwd: Some(cwd.clone()),
..Default::default()
},
)
.unwrap();
assert_eq!(got.doc.kind, Kind::Text);
}
#[test]
fn a_friends_document_lands_in_their_project() {
let (s, r, _d) = setup();
let from = FromPeer {
name: "Trapti".into(),
sign_key: "KEY".into(),
bytes: b"# Garden\n\nbeans".to_vec(),
file: Some("garden.md".into()),
desk: None,
root: None,
lineage: None,
};
let got = receive(
&s,
&r,
Payload {
title: Some("Garden".into()),
origin: Some("peer".into()),
sender: Some("Trapti".into()),
peer: Some(from.clone()),
..Default::default()
},
)
.unwrap();
assert_eq!(got.doc.project, "From Trapti");
assert_eq!(got.doc.workflow_title, "Sent by Trapti");
assert_eq!(got.doc.origin, "peer");
assert_eq!(got.doc.sender, "Trapti");
assert_eq!(got.doc.kind, Kind::Markdown);
assert_eq!(
s.project_root(got.doc.project_id).as_deref(),
Some("peer:KEY")
);
assert!(s.html(&got.doc.id).unwrap().contains("beans"));
let row = &s.projects().unwrap()[0];
assert!(row.friend);
assert_eq!(row.root, "");
let json = serde_json::to_value(row).unwrap();
assert_eq!(json["friend"], true);
let again = receive(
&s,
&r,
Payload {
origin: Some("peer".into()),
peer: Some(from.clone()),
..Default::default()
},
)
.unwrap();
assert!(again.existing);
let mut v2 = from.clone();
v2.bytes = b"# Garden\n\npeas".to_vec();
let newer = receive(
&s,
&r,
Payload {
origin: Some("peer".into()),
peer: Some(v2),
..Default::default()
},
)
.unwrap();
assert_eq!(newer.supersedes.as_deref(), Some(got.doc.id.as_str()));
let png = FromPeer {
bytes: b"\x89PNG\r\n\x1a\n\x00\x00\x00\rIHDR".to_vec(),
file: Some("shot.png".into()),
..from
};
let pic = receive(
&s,
&r,
Payload {
origin: Some("peer".into()),
peer: Some(png),
..Default::default()
},
)
.unwrap();
assert_eq!(pic.doc.kind, Kind::Image);
}
#[test]
fn a_friends_document_lands_on_their_desk_when_they_have_one() {
let (s, r, _d) = setup();
let garden = Dir::new("snyvi-recv-garden");
std::fs::create_dir_all(garden.path.join(".git")).unwrap();
let at = crate::desk::Origin {
id: 3,
name: "Garden".into(),
slot: 0,
};
let got = receive(
&s,
&r,
Payload {
origin: Some("peer".into()),
sender: Some("Trapti".into()),
peer: Some(FromPeer {
name: "Trapti".into(),
sign_key: "KEY".into(),
bytes: b"# Seeds\n\nbeans".to_vec(),
file: Some("seeds.md".into()),
desk: Some((at.clone(), garden.path.to_string_lossy().to_string())),
root: None,
lineage: None,
}),
..Default::default()
},
)
.unwrap();
let (root, _, _) = desk_project(&garden.path.to_string_lossy());
assert_eq!(
s.project_root(got.doc.project_id).as_deref(),
Some(root.as_str())
);
assert_eq!(got.doc.desk, Some(at));
assert_eq!(got.doc.origin, "peer");
assert_eq!(got.doc.sender, "Trapti");
assert!(!s.projects().unwrap()[0].friend, "the desk's own row");
}
#[test]
fn a_friends_document_lands_in_the_folder_both_have() {
let (s, r, _d) = setup();
let mine = Dir::new("snyvi-recv-mine");
let theirs = Dir::new("snyvi-recv-theirs");
for d in [&mine, &theirs] {
std::fs::create_dir_all(d.path.join(".git")).unwrap();
}
let (root, _, _) = desk_project(&mine.path.to_string_lossy());
let send = |bytes: &[u8], lineage: &str| {
receive(
&s,
&r,
Payload {
title: Some("Plan".into()),
origin: Some("peer".into()),
sender: Some("Trapti".into()),
peer: Some(FromPeer {
name: "Trapti".into(),
sign_key: "KEY".into(),
bytes: bytes.to_vec(),
file: Some("PLAN.md".into()),
desk: Some((
crate::desk::Origin {
id: 9,
name: "Theirs".into(),
slot: 0,
},
theirs.path.to_string_lossy().to_string(),
)),
root: Some(root.clone()),
lineage: Some(lineage.into()),
}),
..Default::default()
},
)
.unwrap()
};
let one = send(b"# Plan\n\none", "docs/PLAN.md");
assert_eq!(
s.project_root(one.doc.project_id).as_deref(),
Some(root.as_str())
);
assert_eq!(one.doc.source_path.as_deref(), Some("docs/PLAN.md"));
assert_eq!(one.doc.kind, Kind::Markdown);
let two = send(b"# Plan\n\ntwo", "docs/PLAN.md");
assert_eq!(
two.supersedes.as_deref(),
Some(one.doc.id.as_str()),
"a version"
);
let other = send(b"# Plan\n\nelse", "web/PLAN.md");
assert_eq!(
other.supersedes, None,
"the same name elsewhere is its own row"
);
let scratch = send(b"# Plan\n\nscratch", "0123456789abcdef/PLAN.md");
assert_eq!(scratch.supersedes, None);
assert_eq!(
send(b"# Plan\n\nscratch 2", "0123456789abcdef/PLAN.md")
.supersedes
.as_deref(),
Some(scratch.doc.id.as_str()),
"a file outside the repository is versioned by its key"
);
}
#[test]
fn missing_input_is_an_error() {
let (s, r, _d) = setup();
assert!(receive(&s, &r, Payload::default()).is_err());
}
}