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>,
}
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;
fn automatic(origin: &str) -> bool {
matches!(origin, "hook" | "watch")
}
pub fn receive(store: &Store, renderer: &Renderer, p: Payload) -> Result<Received> {
let mut staged: Option<Staged> = None;
let (body, text, path) = match (&p.content, &p.path) {
(Some(c), _) => (c.clone().into_bytes(), c.clone(), p.path.clone()),
(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() {
staged = Some(store.stage(&path, MAX_MEDIA_BYTES)?);
(Vec::new(), String::new(), Some(sp))
} 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()
};
(bytes, text, Some(sp))
}
}
(None, None) => bail!("send_document needs either `path` or `content`"),
};
if body.len() > MAX_BYTES {
bail!("content is larger than {} MB", MAX_BYTES / 1024 / 1024);
}
let origin = p.origin.as_deref().unwrap_or("cli");
let from = p
.pane
.as_deref()
.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,
});
let anchor: PathBuf = p
.cwd
.as_deref()
.map(PathBuf::from)
.or_else(|| path.as_deref().map(PathBuf::from))
.unwrap_or_else(|| std::env::current_dir().unwrap_or_else(|_| PathBuf::from("/")));
let proj = project::resolve(&anchor);
let root = proj.root.to_string_lossy().to_string();
let branch = project::branch(&proj.root);
let hash = match &staged {
Some(st) => st.hash.clone(),
None => blake3::hash(&body).to_hex().to_string(),
};
let latest_same_path = match &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) = match renderer.detect(path.as_deref(), p.lang.as_deref(), &text) {
_ if staged.is_some() => {
let ext = render::ext_of(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 !text.is_empty() && render::looks_binary(&body) => (Kind::Binary, lang),
(Kind::Image, lang) => (Kind::Image, lang),
_ if text.is_empty() && !body.is_empty() => (Kind::Binary, None),
other => other,
};
let title = render::title_for(p.title.as_deref(), kind, path.as_deref(), &text);
let (wf_name, wf_title) = match (&p.workflow, &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.clone()),
_ => ("manual".to_string(), "Sent manually".to_string()),
};
let wf_key = wf_name.to_lowercase();
let _ = &wf_name;
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 body_src = if kind == Kind::Markdown {
render::strip_leading_h1(&text, &title)
} else {
None
};
let file_base = path.as_ref().map(|_| format!("/files/{id}/"));
let html = 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(path.as_deref().unwrap_or_default()),
),
Kind::Binary => render::placeholder(&render::describe_bytes(&title, body.len() as u64)),
_ => renderer.render_with_base(
kind,
lang.as_deref(),
body_src.as_deref().unwrap_or(&text),
file_base.as_deref(),
),
};
let needs_full_highlight = kind == Kind::Code && 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: path.as_deref(),
branch: branch.as_deref(),
origin,
sender: p.sender.as_deref().unwrap_or(""),
desk: from.as_ref(),
source: &body,
staged: staged.as_ref(),
search_body: &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 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 missing_input_is_an_error() {
let (s, r, _d) = setup();
assert!(receive(&s, &r, Payload::default()).is_err());
}
}