use anyhow::{Result, anyhow, bail};
use serde::Deserialize;
use std::path::{Path, PathBuf};
use std::process::Command;
#[derive(Debug, Clone, Deserialize)]
pub struct ProjectAdvanced {
pub project_id: String,
#[serde(default)]
pub head: Option<String>,
}
#[derive(Debug, Clone)]
pub struct EpicReplica {
pub project_id: String,
pub remote_url: String,
pub local: PathBuf,
}
impl EpicReplica {
pub fn should_sync(&self, event: &ProjectAdvanced, local_head: Option<&str>) -> bool {
if event.project_id != self.project_id {
return false;
}
match (event.head.as_deref(), local_head) {
(Some(remote), Some(local)) => remote != local, (Some(_), None) => true, (None, _) => true, }
}
pub fn local_head(&self) -> Option<String> {
git(&self.local, &["rev-parse", "HEAD"])
.ok()
.map(|s| s.trim().to_string())
.filter(|s| !s.is_empty())
}
pub fn is_cloned(&self) -> bool {
self.local.join(".git").exists()
}
pub fn ensure_cloned(&self) -> Result<()> {
if self.is_cloned() {
return Ok(());
}
if let Some(parent) = self.local.parent() {
std::fs::create_dir_all(parent)?;
}
run(Command::new("git")
.args(["clone", "--recursive", &self.remote_url])
.arg(&self.local))?;
Ok(())
}
pub fn sync(&self) -> Result<String> {
git(&self.local, &["pull", "--ff-only"])?;
let _ = git(
&self.local,
&["submodule", "update", "--init", "--recursive"],
);
self.local_head()
.ok_or_else(|| anyhow!("no HEAD after sync of {:?}", self.local))
}
pub fn on_event(&self, event: &ProjectAdvanced) -> Result<Option<String>> {
self.ensure_cloned()?;
if self.should_sync(event, self.local_head().as_deref()) {
Ok(Some(self.sync()?))
} else {
Ok(None)
}
}
}
pub fn advanced_subject(subject_prefix: &str, project_id: &str) -> String {
format!("{subject_prefix}.project.{project_id}.advanced")
}
pub async fn run_sync_loop(
nats: &async_nats::Client,
subject_prefix: &str,
replica: &EpicReplica,
) -> Result<()> {
use futures::StreamExt;
replica.ensure_cloned()?;
let subject = advanced_subject(subject_prefix, &replica.project_id);
let mut sub = nats
.subscribe(subject.clone())
.await
.map_err(|e| anyhow!("subscribe {subject}: {e}"))?;
tracing::info!(subject = %subject, local = ?replica.local, "epic sync loop started");
while let Some(msg) = sub.next().await {
match serde_json::from_slice::<ProjectAdvanced>(&msg.payload) {
Ok(event) => match replica.on_event(&event) {
Ok(Some(head)) => {
tracing::info!(project = %event.project_id, head = %head, "epic synced")
}
Ok(None) => {}
Err(e) => tracing::warn!(error = %e, "epic sync failed; will retry on next event"),
},
Err(e) => tracing::warn!(error = %e, "unparseable project_advanced event"),
}
}
Ok(())
}
fn git(dir: &Path, args: &[&str]) -> Result<String> {
run(Command::new("git").arg("-C").arg(dir).args(args))
}
fn run(cmd: &mut Command) -> Result<String> {
cmd.env_remove("GIT_DIR")
.env_remove("GIT_WORK_TREE")
.env_remove("GIT_INDEX_FILE");
let out = cmd.output()?;
if !out.status.success() {
bail!(
"git failed: {}",
String::from_utf8_lossy(&out.stderr).trim()
);
}
Ok(String::from_utf8_lossy(&out.stdout).trim().to_string())
}
#[cfg(test)]
mod tests {
use super::*;
fn adv(project: &str, head: Option<&str>) -> ProjectAdvanced {
ProjectAdvanced {
project_id: project.to_string(),
head: head.map(str::to_string),
}
}
fn replica(project: &str, remote: &Path, local: &Path) -> EpicReplica {
EpicReplica {
project_id: project.to_string(),
remote_url: remote.to_string_lossy().into_owned(),
local: local.to_path_buf(),
}
}
#[test]
fn advanced_subject_is_project_scoped() {
assert_eq!(
advanced_subject("nsed", "root-sha"),
"nsed.project.root-sha.advanced"
);
}
#[test]
fn should_sync_only_for_this_project_and_a_new_head() {
let r = replica("mine", Path::new("/x"), Path::new("/y"));
assert!(r.should_sync(&adv("mine", Some("h2")), Some("h1")));
assert!(!r.should_sync(&adv("mine", Some("h1")), Some("h1")));
assert!(r.should_sync(&adv("mine", Some("h1")), None));
assert!(r.should_sync(&adv("mine", None), Some("h1")));
assert!(!r.should_sync(&adv("theirs", Some("h9")), Some("h1")));
}
fn g(dir: &Path, args: &[&str]) {
let o = Command::new("git")
.env_remove("GIT_DIR")
.env_remove("GIT_WORK_TREE")
.env_remove("GIT_INDEX_FILE")
.arg("-C")
.arg(dir)
.args(args)
.output()
.unwrap();
assert!(
o.status.success(),
"git {args:?}: {}",
String::from_utf8_lossy(&o.stderr)
);
}
fn rev(dir: &Path) -> String {
let o = Command::new("git")
.env_remove("GIT_DIR")
.env_remove("GIT_WORK_TREE")
.env_remove("GIT_INDEX_FILE")
.arg("-C")
.arg(dir)
.args(["rev-parse", "HEAD"])
.output()
.unwrap();
String::from_utf8_lossy(&o.stdout).trim().to_string()
}
#[test]
fn clone_then_pull_advances_the_replica() {
let sbx = std::env::temp_dir().join(format!("qr-sync-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&sbx);
let origin = sbx.join("epic.git");
std::fs::create_dir_all(&origin).unwrap();
g(&origin, &["init", "-q", "-b", "main", "--bare"]);
let work = sbx.join("work");
g(
&sbx,
&[
"clone",
"-q",
origin.to_str().unwrap(),
work.to_str().unwrap(),
],
);
g(&work, &["config", "user.email", "a@b"]);
g(&work, &["config", "user.name", "a"]);
std::fs::write(work.join("kpi.md"), "A\n").unwrap();
g(&work, &["add", "-A"]);
g(&work, &["commit", "-qm", "A"]);
g(&work, &["push", "-q", "origin", "main"]);
let a = rev(&work);
let r = replica("proj", &origin, &sbx.join("replica"));
assert!(!r.is_cloned());
r.ensure_cloned().unwrap();
assert!(r.is_cloned());
assert_eq!(r.local_head().as_deref(), Some(a.as_str()), "replica at A");
std::fs::write(work.join("kpi.md"), "B\n").unwrap();
g(&work, &["add", "-A"]);
g(&work, &["commit", "-qm", "B"]);
g(&work, &["push", "-q", "origin", "main"]);
let b = rev(&work);
let synced = r.on_event(&adv("proj", Some(&b))).unwrap();
assert_eq!(synced.as_deref(), Some(b.as_str()), "on_event pulled to B");
assert_eq!(
r.local_head().as_deref(),
Some(b.as_str()),
"replica now at B"
);
assert_eq!(
r.on_event(&adv("proj", Some(&b))).unwrap(),
None,
"same head → skip"
);
let _ = std::fs::remove_dir_all(&sbx);
}
}