1use std::collections::{BTreeMap, BTreeSet};
20use std::path::{Path, PathBuf};
21use std::process::{Command, Stdio};
22
23use anyhow::{bail, Context, Result};
24use serde_json::Value;
25
26#[derive(Debug, Clone, PartialEq, Eq, serde::Deserialize)]
28pub struct Scope {
29 pub name: String,
31 pub recipients: Vec<String>,
33 #[serde(default)]
37 pub projects: BTreeMap<String, String>,
38}
39
40impl Scope {
41 #[must_use]
43 pub fn for_project(&self, project: &str) -> String {
44 self.projects
45 .get(project)
46 .cloned()
47 .unwrap_or_else(|| self.name.clone())
48 }
49}
50
51#[derive(Debug, Clone, Default, serde::Deserialize)]
53struct Local {
54 #[serde(default)]
55 default_scope: Option<String>,
56 #[serde(default)]
60 repos: Vec<String>,
61}
62
63fn local() -> Local {
64 std::fs::read_to_string(config_dir().join("sync.toml"))
65 .ok()
66 .and_then(|t| toml::from_str::<Local>(&t).ok())
67 .unwrap_or_default()
68}
69
70fn expand(path: &str) -> PathBuf {
71 match path.strip_prefix("~/") {
72 Some(rest) => std::env::var_os("HOME")
73 .map_or_else(|| PathBuf::from(path), |h| PathBuf::from(h).join(rest)),
74 None => PathBuf::from(path),
75 }
76}
77
78#[must_use]
80pub fn sender_entity(host: &str) -> String {
81 format!("sync:{host}")
82}
83
84#[must_use]
86pub fn host() -> String {
87 std::fs::read_to_string("/proc/sys/kernel/hostname")
88 .map(|h| h.trim().to_string())
89 .ok()
90 .filter(|h| !h.is_empty())
91 .unwrap_or_else(|| "host".into())
92}
93
94fn config_dir() -> PathBuf {
95 std::env::var_os("XDG_CONFIG_HOME")
96 .map(PathBuf::from)
97 .or_else(|| std::env::var_os("HOME").map(|h| PathBuf::from(h).join(".config")))
98 .unwrap_or_else(|| PathBuf::from(".config"))
99 .join("ljos")
100}
101
102fn state_dir() -> PathBuf {
103 std::env::var_os("XDG_STATE_HOME")
104 .map(PathBuf::from)
105 .or_else(|| std::env::var_os("HOME").map(|h| PathBuf::from(h).join(".local/state")))
106 .unwrap_or_else(|| PathBuf::from(".local/state"))
107 .join("ljos")
108 .join("sync")
109}
110
111#[must_use]
113pub fn identity_path() -> PathBuf {
114 config_dir().join("age.key")
115}
116
117pub fn public_key() -> Result<String> {
124 let key = identity_path();
125 if !key.exists() {
126 if let Some(dir) = key.parent() {
127 std::fs::create_dir_all(dir)?;
128 }
129 let made = Command::new("age-keygen")
130 .arg("-o")
131 .arg(&key)
132 .stdin(Stdio::null())
133 .output()
134 .context("age-keygen not on PATH; install age")?;
135 if !made.status.success() {
136 bail!(
137 "age-keygen: {}",
138 String::from_utf8_lossy(&made.stderr).trim()
139 );
140 }
141 }
142 let out = Command::new("age-keygen")
143 .arg("-y")
144 .arg(&key)
145 .output()
146 .context("age-keygen not on PATH; install age")?;
147 if !out.status.success() {
148 bail!(
149 "age-keygen -y: {}",
150 String::from_utf8_lossy(&out.stderr).trim()
151 );
152 }
153 Ok(String::from_utf8_lossy(&out.stdout).trim().to_string())
154}
155
156pub fn scope_of_repo(root: &Path) -> Result<Option<Scope>> {
162 let path = root.join(".ljos").join("sync.toml");
163 match std::fs::read_to_string(&path) {
164 Ok(text) => Ok(Some(
165 toml::from_str(&text).with_context(|| path.display().to_string())?,
166 )),
167 Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(None),
168 Err(e) => Err(e).with_context(|| path.display().to_string()),
169 }
170}
171
172fn default_scope(repo_scope: &str) -> String {
173 local()
174 .default_scope
175 .unwrap_or_else(|| repo_scope.to_string())
176}
177
178#[must_use]
180pub fn atom_scope(atom: &Value, default: &str) -> String {
181 atom["entities"]
182 .as_array()
183 .into_iter()
184 .flatten()
185 .filter_map(Value::as_str)
186 .find_map(|e| e.strip_prefix("scope:"))
187 .map_or_else(|| default.to_string(), str::to_string)
188}
189
190#[must_use]
193pub fn travelling(atom: &Value) -> Value {
194 let mut a = atom.clone();
195 if let Some(map) = a.as_object_mut() {
196 map.remove("embedding");
197 }
198 a
199}
200
201#[must_use]
203pub fn is_own(atom: &Value) -> bool {
204 !atom["entities"]
205 .as_array()
206 .into_iter()
207 .flatten()
208 .filter_map(Value::as_str)
209 .any(|e| e.starts_with("sync:"))
210}
211
212#[must_use]
215pub fn log_plaintext(atoms: &[Value], scope: &str, default: &str) -> String {
216 let mut mine: Vec<Value> = atoms
217 .iter()
218 .filter(|a| is_own(a) && atom_scope(a, default) == scope)
219 .map(travelling)
220 .collect();
221 mine.sort_by(|a, b| a["id"].as_str().cmp(&b["id"].as_str()));
222 mine.iter().map(|a| format!("{a}\n")).collect()
223}
224
225fn log_dir(root: &Path) -> PathBuf {
226 root.join(".ljos").join("atoms")
227}
228
229fn seal(plain: &str, recipients: &[String], out: &Path) -> Result<()> {
231 if recipients.is_empty() {
232 bail!("sync: the scope names no recipients; add this machine's `ljos sync --key` to .ljos/sync.toml");
233 }
234 let mut cmd = Command::new("age");
235 for r in recipients {
236 cmd.arg("-r").arg(r);
237 }
238 let tmp = out.with_extension("age.tmp");
239 let mut child = cmd
240 .arg("-o")
241 .arg(&tmp)
242 .stdin(Stdio::piped())
243 .stderr(Stdio::piped())
244 .spawn()
245 .context("age not on PATH; install age")?;
246 use std::io::Write;
247 child
248 .stdin
249 .take()
250 .context("age: no stdin")?
251 .write_all(plain.as_bytes())?;
252 let done = child.wait_with_output()?;
253 if !done.status.success() {
254 let _ = std::fs::remove_file(&tmp);
255 bail!("age: {}", String::from_utf8_lossy(&done.stderr).trim());
256 }
257 std::fs::rename(&tmp, out)?;
258 Ok(())
259}
260
261fn open(path: &Path, identity: &Path) -> Result<Option<String>> {
264 let out = Command::new("age")
265 .arg("-d")
266 .arg("-i")
267 .arg(identity)
268 .arg(path)
269 .stdin(Stdio::null())
270 .output()
271 .context("age not on PATH; install age")?;
272 if out.status.success() {
273 Ok(Some(String::from_utf8_lossy(&out.stdout).into_owned()))
274 } else if String::from_utf8_lossy(&out.stderr).contains("no identity matched") {
275 Ok(None)
276 } else {
277 bail!(
278 "age -d {}: {}",
279 path.display(),
280 String::from_utf8_lossy(&out.stderr).trim()
281 )
282 }
283}
284
285pub fn export(root: &Path, scope: &Scope, atoms: &[Value]) -> Result<String> {
292 let host = host();
293 let plain = log_plaintext(atoms, &scope.name, &default_scope(&scope.name));
294 let count = plain.lines().count();
295 let dir = log_dir(root);
296 std::fs::create_dir_all(&dir)?;
297 let stamp = dir.join(format!("{host}.digest"));
298 let digest = format!(
299 "{} {}\n",
300 crate::work_id(&plain),
301 crate::work_id(&scope.recipients.join(","))
302 );
303 if std::fs::read_to_string(&stamp).is_ok_and(|d| d == digest) {
304 return Ok(format!(
305 "sync: {count} {} atoms unchanged since the last export\n",
306 scope.name
307 ));
308 }
309 seal(
310 &plain,
311 &scope.recipients,
312 &dir.join(format!("{host}.jsonl.age")),
313 )?;
314 std::fs::write(&stamp, digest)?;
315 Ok(format!(
316 "sync: exported {count} {} atoms to .ljos/atoms/{host}.jsonl.age, sealed to {} recipient{}\n",
317 scope.name,
318 scope.recipients.len(),
319 if scope.recipients.len() == 1 { "" } else { "s" }
320 ))
321}
322
323#[derive(Debug, Default, PartialEq)]
326pub struct Incoming {
327 pub post: Vec<Value>,
328 pub retired: Vec<String>,
329}
330
331#[must_use]
335pub fn incoming(
336 plain: &str,
337 host: &str,
338 before: &BTreeSet<String>,
339 held: &BTreeSet<String>,
340) -> Incoming {
341 let mut post = Vec::new();
342 let mut now = BTreeSet::new();
343 for line in plain.lines().filter(|l| !l.trim().is_empty()) {
344 let Ok(mut atom) = serde_json::from_str::<Value>(line) else {
345 continue;
346 };
347 let text = atom["text"].as_str().unwrap_or("").to_string();
348 if text.is_empty() {
349 continue;
350 }
351 let id = crate::work_id(&text);
352 now.insert(id.clone());
353 if held.contains(&id) {
354 continue;
355 }
356 if let Some(map) = atom.as_object_mut() {
357 let mut entities: Vec<Value> = map
358 .get("entities")
359 .and_then(Value::as_array)
360 .cloned()
361 .unwrap_or_default();
362 let tag = sender_entity(host);
363 if !entities.iter().any(|e| e.as_str() == Some(tag.as_str())) {
364 entities.push(Value::String(tag));
365 }
366 map.insert("entities".into(), Value::Array(entities));
367 map.remove("id");
368 }
369 post.push(atom);
370 }
371 let retired = before.difference(&now).cloned().collect();
372 Incoming { post, retired }
373}
374
375fn taken_path(scope: &str, host: &str) -> PathBuf {
376 state_dir().join(scope).join(format!("{host}.taken"))
377}
378
379fn read_taken(scope: &str, host: &str) -> BTreeSet<String> {
380 std::fs::read_to_string(taken_path(scope, host))
381 .map(|t| t.lines().map(str::to_string).collect())
382 .unwrap_or_default()
383}
384
385fn write_taken(scope: &str, host: &str, plain: &str) -> Result<()> {
386 let path = taken_path(scope, host);
387 if let Some(dir) = path.parent() {
388 std::fs::create_dir_all(dir)?;
389 }
390 let ids: BTreeSet<String> = plain
391 .lines()
392 .filter_map(|l| serde_json::from_str::<Value>(l).ok())
393 .filter_map(|a| a["text"].as_str().map(crate::work_id))
394 .collect();
395 std::fs::write(path, ids.into_iter().map(|i| i + "\n").collect::<String>())?;
396 Ok(())
397}
398
399pub fn import(root: &Path, scope: &Scope) -> Result<String> {
406 let me = host();
407 let identity = identity_path();
408 if !identity.exists() {
409 return Ok("sync: no age identity on this machine; `ljos sync --key` makes one\n".into());
410 }
411 let dir = log_dir(root);
412 let Ok(entries) = std::fs::read_dir(&dir) else {
413 return Ok(format!("sync: no {} logs yet\n", scope.name));
414 };
415 let client = crate::pack()?;
416 let mut out = String::new();
417 let mut paths: Vec<PathBuf> = entries
418 .flatten()
419 .map(|e| e.path())
420 .filter(|p| p.to_string_lossy().ends_with(".jsonl.age"))
421 .collect();
422 paths.sort();
423 for path in paths {
424 let host = path
425 .file_name()
426 .and_then(|n| n.to_str())
427 .and_then(|n| n.strip_suffix(".jsonl.age"))
428 .unwrap_or("")
429 .to_string();
430 if host.is_empty() || host == me {
431 continue;
432 }
433 let Some(plain) = open(&path, &identity)? else {
434 out.push_str(&format!(
435 "sync: {host}'s {} log is not sealed to this machine\n",
436 scope.name
437 ));
438 continue;
439 };
440 let workspace = client.workspace();
441 let tag = sender_entity(&host);
444 let held: BTreeSet<String> = crate::atoms_lean(&client, &workspace)?
445 .iter()
446 .filter(|a| {
447 a["entities"]
448 .as_array()
449 .into_iter()
450 .flatten()
451 .any(|e| e.as_str() == Some(tag.as_str()))
452 })
453 .filter_map(|a| a["text"].as_str().map(crate::work_id))
454 .collect();
455 let got = incoming(&plain, &host, &read_taken(&scope.name, &host), &held);
456 let (mut kept, mut refused) = (0usize, 0usize);
457 for mut atom in got.post {
458 if let Some(map) = atom.as_object_mut() {
459 map.insert("workspace".into(), Value::String(workspace.clone()));
460 }
461 match client.post_atom(&atom) {
462 Ok(_) => kept += 1,
463 Err(_) => refused += 1,
464 }
465 }
466 let mut retired = 0usize;
467 if !got.retired.is_empty() {
468 for atom in crate::atoms_lean(&client, &workspace)? {
469 let text = atom["text"].as_str().unwrap_or("");
470 let from_host = atom["entities"]
471 .as_array()
472 .into_iter()
473 .flatten()
474 .any(|e| e.as_str() == Some(tag.as_str()));
475 if from_host && got.retired.contains(&crate::work_id(text)) {
476 if let Some(id) = atom["id"].as_str() {
477 if crate::packset_forget(id, None).is_ok() {
478 retired += 1;
479 }
480 }
481 }
482 }
483 }
484 write_taken(&scope.name, &host, &plain)?;
485 out.push_str(&format!(
486 "sync: {host}: {kept} {} atoms taken, {refused} refused, {retired} retired\n",
487 scope.name
488 ));
489 }
490 Ok(out)
491}
492
493#[must_use]
495pub fn scope_for_issue(issue: &str) -> Option<String> {
496 let hit = vissue_core::Layout::resolve(None, None)
497 .and_then(vissue_core::Router::load)
498 .and_then(|router| router.find_by_id(issue))
499 .ok()?;
500 let root = hit
501 .path
502 .ancestors()
503 .find(|d| d.join(".git").exists())?
504 .to_path_buf();
505 scope_of_repo(&root)
506 .ok()
507 .flatten()
508 .map(|s| s.for_project(&hit.project))
509}
510
511fn tracker_root() -> Result<PathBuf> {
513 let layout = vissue_core::Layout::resolve(None, None).map_err(anyhow::Error::from)?;
514 Ok(layout.root().to_path_buf())
515}
516
517fn git(root: &Path, args: &[&str]) -> std::io::Result<std::process::Output> {
518 Command::new("git")
519 .arg("-C")
520 .arg(root)
521 .args(args)
522 .stdin(Stdio::null())
523 .output()
524}
525
526fn commit_log(root: &Path, scope: &str) -> String {
529 let host = host();
530 let files = [
531 format!(".ljos/atoms/{host}.jsonl.age"),
532 format!(".ljos/atoms/{host}.digest"),
533 ];
534 let common = git(root, &["rev-parse", "--git-common-dir"])
536 .ok()
537 .filter(|o| o.status.success())
538 .map(|o| root.join(String::from_utf8_lossy(&o.stdout).trim()))
539 .unwrap_or_else(|| root.join(".git"));
540 let held = crate::CommitLock::acquire(&common.join("ljos-commit.lock"));
541 let mut add = vec!["add", "--"];
542 add.extend(files.iter().map(String::as_str));
543 if git(root, &add).map_or(true, |o| !o.status.success()) {
544 return "sync: could not stage the log\n".into();
545 }
546 let staged = git(
547 root,
548 &["diff", "--cached", "--quiet", "--", &files[0], &files[1]],
549 );
550 if staged.is_ok_and(|o| o.status.success()) {
551 return String::new();
552 }
553 let message = format!("chore(sync): {host} {scope} atoms");
554 let mut commit = vec!["commit", "-q", "--only", "-m", message.as_str(), "--"];
555 commit.extend(files.iter().map(String::as_str));
556 match git(root, &commit) {
557 Ok(o) if o.status.success() => {}
558 Ok(o) => {
559 return format!(
560 "sync: commit refused: {}\n",
561 String::from_utf8_lossy(&o.stderr)
562 .lines()
563 .next()
564 .unwrap_or("")
565 )
566 }
567 Err(e) => return format!("sync: git: {e}\n"),
568 }
569 drop(held);
570 let mut pushed = vec![];
571 let pushed_first = git(root, &["push", "-q"]).is_ok_and(|o| o.status.success());
574 let pushed_after_merge = !pushed_first
575 && git(root, &["pull", "-q", "--no-rebase", "--no-edit"]).is_ok_and(|o| o.status.success())
576 && git(root, &["push", "-q"]).is_ok_and(|o| o.status.success());
577 if pushed_first || pushed_after_merge {
578 pushed.push("upstream".to_string());
579 }
580 if let Some(up) = crate::tracker_upstream(root) {
581 for (remote, branch) in crate::tracker_mirrors(root, &up).unwrap_or_default() {
582 let refspec = format!("HEAD:refs/heads/{branch}");
583 if git(root, &["push", "-q", &remote, &refspec]).is_ok_and(|o| o.status.success()) {
584 pushed.push(remote);
585 }
586 }
587 }
588 format!(
589 "sync: committed {message}; pushed to {}\n",
590 if pushed.is_empty() {
591 "nothing".to_string()
592 } else {
593 pushed.join(", ")
594 }
595 )
596}
597
598fn catch_up(root: &Path) -> String {
603 if git(root, &["fetch", "-q", "--all"]).map_or(true, |o| !o.status.success()) {
604 return "sync: fetch failed; reading the logs as they are\n".into();
605 }
606 let Some(up) = crate::tracker_upstream(root) else {
607 return String::new();
608 };
609 let mut refs = vec![up.clone()];
610 refs.extend(
611 crate::tracker_mirrors(root, &up)
612 .unwrap_or_default()
613 .into_iter()
614 .map(|(r, b)| format!("{r}/{b}")),
615 );
616 let mut stuck = Vec::new();
617 for r in refs {
618 let forward = git(root, &["merge", "-q", "--ff-only", &r]);
619 if forward.is_ok_and(|o| o.status.success()) {
620 continue;
621 }
622 let merged = git(root, &["merge", "-q", "--no-edit", &r]);
625 if !merged.is_ok_and(|o| o.status.success()) {
626 let _ = git(root, &["merge", "--abort"]);
627 stuck.push(r);
628 }
629 }
630 if stuck.is_empty() {
631 String::new()
632 } else {
633 format!(
634 "sync: could not merge {}; `ljos doctor` names the split\n",
635 stuck.join(", ")
636 )
637 }
638}
639
640pub fn sync_repo(pull_import: bool, export_log: bool) -> Result<String> {
649 let listed = local().repos;
650 let roots: Vec<PathBuf> = if listed.is_empty() {
651 vec![tracker_root()?]
652 } else {
653 listed.iter().map(|r| expand(r)).collect()
654 };
655 let mut out = String::new();
656 for root in roots {
657 out.push_str(&sync_one(&root, pull_import, export_log)?);
658 }
659 Ok(out)
660}
661
662fn sync_one(root: &Path, pull_import: bool, export_log: bool) -> Result<String> {
663 let root = root.to_path_buf();
664 let Some(scope) = scope_of_repo(&root)? else {
665 return Ok(format!(
666 "sync: {} names no scope; write .ljos/sync.toml with name and recipients (`ljos sync --key` prints this machine's)\n",
667 root.display()
668 ));
669 };
670 let mut out = String::new();
671 if pull_import {
672 out.push_str(&catch_up(&root));
673 out.push_str(&import(&root, &scope)?);
674 }
675 if export_log {
676 let client = crate::pack()?;
677 let atoms = crate::atoms_lean(&client, &client.workspace())?;
678 out.push_str(&export(&root, &scope, &atoms)?);
679 out.push_str(&commit_log(&root, &scope.name));
680 }
681 Ok(out)
682}
683
684#[cfg(test)]
685mod tests {
686 use super::*;
687
688 fn atom(id: &str, text: &str, entities: &[&str]) -> Value {
689 serde_json::json!({"id": id, "text": text, "kind": "lesson",
690 "entities": entities, "embedding": [0.1, 0.2]})
691 }
692
693 #[test]
694 fn a_log_carries_own_atoms_of_its_scope_without_vectors() {
695 let atoms = vec![
696 atom("b", "second", &["seat:x"]),
697 atom("a", "first", &[]),
698 atom("c", "personal", &["scope:personal"]),
699 atom("d", "imported", &["sync:otherhost"]),
700 ];
701 let plain = log_plaintext(&atoms, "surf", "surf");
702 let lines: Vec<Value> = plain
703 .lines()
704 .map(|l| serde_json::from_str(l).unwrap())
705 .collect();
706 assert_eq!(lines.len(), 2, "{plain}");
707 assert_eq!(lines[0]["id"], "a");
708 assert_eq!(lines[1]["id"], "b");
709 assert!(lines.iter().all(|a| a.get("embedding").is_none()));
710 assert_eq!(log_plaintext(&atoms, "personal", "surf").lines().count(), 1);
711 }
712
713 #[test]
714 fn a_project_named_in_the_scope_file_takes_its_own_scope() {
715 let scope: Scope = toml::from_str(
716 "name = \"personal\"\nrecipients = [\"age1x\"]\n[projects]\ntools = \"shared\"\n",
717 )
718 .unwrap();
719 assert_eq!(scope.for_project("tools"), "shared");
720 assert_eq!(scope.for_project("garden"), "personal");
721 let bare: Scope = toml::from_str("name = \"shared\"\nrecipients = []\n").unwrap();
722 assert!(bare.projects.is_empty());
723 }
724
725 #[test]
726 fn an_import_tags_the_sender_and_names_what_it_retired() {
727 let plain = format!(
728 "{}\n{}\n",
729 atom("a", "kept", &["seat:x"]),
730 atom("b", "new", &[])
731 );
732 let before: BTreeSet<String> = [crate::work_id("kept"), crate::work_id("dropped")].into();
733 let got = incoming(&plain, "rglat", &before, &BTreeSet::new());
734 assert_eq!(got.post.len(), 2);
735 let again = incoming(&plain, "rglat", &before, &[crate::work_id("kept")].into());
736 assert_eq!(
737 again.post.len(),
738 1,
739 "a text already held is not posted again"
740 );
741 assert_eq!(again.post[0]["text"], "new");
742 assert!(got.post.iter().all(|a| a.get("id").is_none()));
743 assert!(got.post.iter().all(|a| a["entities"]
744 .as_array()
745 .unwrap()
746 .iter()
747 .any(|e| e == "sync:rglat")));
748 assert_eq!(got.retired, vec![crate::work_id("dropped")]);
749 }
750
751 #[test]
752 fn a_sealed_log_opens_for_a_recipient_and_not_for_another() {
753 let dir = tempfile::tempdir().unwrap();
754 let key = |name: &str| {
755 let path = dir.path().join(name);
756 let made = Command::new("age-keygen")
757 .arg("-o")
758 .arg(&path)
759 .output()
760 .unwrap();
761 assert!(made.status.success(), "age-keygen must be installed");
762 let public = Command::new("age-keygen")
763 .arg("-y")
764 .arg(&path)
765 .output()
766 .unwrap();
767 (
768 path,
769 String::from_utf8_lossy(&public.stdout).trim().to_string(),
770 )
771 };
772 let (mine, my_public) = key("mine.key");
773 let (other, _) = key("other.key");
774 let out = dir.path().join("host.jsonl.age");
775 seal("{\"text\":\"x\"}\n", &[my_public], &out).unwrap();
776 assert_eq!(
777 open(&out, &mine).unwrap().as_deref(),
778 Some("{\"text\":\"x\"}\n")
779 );
780 assert_eq!(open(&out, &other).unwrap(), None);
781 assert!(seal("x", &[], &out).is_err(), "no recipients is refused");
782 }
783
784 #[test]
785 fn an_atoms_scope_is_its_entity_or_the_default() {
786 assert_eq!(
787 atom_scope(&atom("a", "t", &["scope:personal"]), "surf"),
788 "personal"
789 );
790 assert_eq!(atom_scope(&atom("a", "t", &["seat:x"]), "surf"), "surf");
791 }
792}