use std::collections::{BTreeMap, BTreeSet};
use atelier_sdk_diff::{Diff, diff_listings};
use futures::{AsyncReadExt, StreamExt};
use jj_lib::backend::{CommitId, Signature, Timestamp, TreeValue};
use jj_lib::config::{ConfigLayer, ConfigSource, StackedConfig};
use jj_lib::default_backend_factories::{
default_backend_factories, default_working_copy_factories, default_working_copy_factory,
};
use jj_lib::git::{self, GitImportOptions};
use jj_lib::gitignore::GitIgnoreFile;
use jj_lib::matchers::{EverythingMatcher, NothingMatcher};
use jj_lib::merged_tree::MergedTree;
use jj_lib::object_id::ObjectId;
use jj_lib::op_store::RefTarget;
use jj_lib::ref_name::{RefName, WorkspaceNameBuf};
use jj_lib::repo::{ReadonlyRepo, Repo};
use jj_lib::repo_path::RepoPath;
use jj_lib::rewrite::rebase_commit;
use jj_lib::settings::UserSettings;
use jj_lib::working_copy::SnapshotOptions;
use jj_lib::workspace::{LockedWorkspace, Workspace as JjWorkspace};
use pollster::block_on;
use std::fs;
use std::os::unix::fs::PermissionsExt;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use crate::config::Actor;
use crate::error::{Error, config_err, engine_err};
use crate::workspace::SKIP_NAMES;
const NEW_FILE_SIZE_MAX: u64 = 50 * 1024 * 1024;
pub(crate) const LADDER_FILE_SIZE_MAX: u64 = 8 * 1024 * 1024;
const _: () = assert!(LADDER_FILE_SIZE_MAX <= NEW_FILE_SIZE_MAX);
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Snapshot {
pub id: String,
pub actor: String,
pub at_ms: i64,
pub parents: Vec<String>,
}
pub(crate) struct DiffSides {
before: MergedTree,
after: MergedTree,
}
pub(crate) struct FileBlob {
pub id: String,
pub bytes: Vec<u8>,
}
pub(crate) enum Side {
Absent,
TooLarge,
Blob(FileBlob),
}
pub(crate) enum StepBack {
Stepped { restored: String },
AlreadyStepped,
LineMoved { head: String },
}
pub(crate) enum LandOutcome {
Landed {
snapshot: String,
},
Conflicted,
}
enum SnapshotStyle {
Stack,
Amend,
}
pub(crate) struct Engine {
jj: JjWorkspace,
repo: Arc<ReadonlyRepo>,
_settings: UserSettings,
boundary: Vec<String>,
}
impl Engine {
pub fn init(root: &Path, actor: &Actor, boundary: &[String]) -> Result<Self, Error> {
let settings = build_settings(actor)?;
let (jj, repo) = block_on(JjWorkspace::init_colocated_git(
&settings,
root,
gix_hash::Kind::Sha1,
))
.map_err(engine_err)?;
Ok(Self {
jj,
repo,
_settings: settings,
boundary: boundary.to_vec(),
})
}
pub fn open(root: &Path, actor: &Actor, boundary: &[String]) -> Result<Self, Error> {
let settings = build_settings(actor)?;
let jj = JjWorkspace::load(
&settings,
root,
&default_backend_factories(),
&default_working_copy_factories(),
)
.map_err(engine_err)?;
let repo = block_on(jj.repo_loader().load_at_head()).map_err(engine_err)?;
Ok(Self {
jj,
repo,
_settings: settings,
boundary: boundary.to_vec(),
})
}
pub fn refresh(&mut self) -> Result<(), Error> {
self.repo = block_on(self.jj.repo_loader().load_at_head()).map_err(engine_err)?;
Ok(())
}
pub fn adopt_git(root: &Path, actor: &Actor, boundary: &[String]) -> Result<Self, Error> {
block_on(Self::adopt_git_async(root, actor, boundary))
}
async fn adopt_git_async(
root: &Path,
actor: &Actor,
boundary: &[String],
) -> Result<Self, Error> {
let settings = build_settings(actor)?;
let (mut jj, repo) = JjWorkspace::init_external_git(&settings, root, &root.join(".git"))
.await
.map_err(engine_err)?;
let mut tx = repo.start_transaction();
git::import_head(tx.repo_mut()).await.map_err(engine_err)?;
let options = GitImportOptions {
abandon_unreachable_commits: false,
record_synthetic_predecessors: false,
remote_auto_track_bookmarks: std::collections::HashMap::new(),
};
git::import_refs(tx.repo_mut(), &options)
.await
.map_err(engine_err)?;
let head = tx.repo_mut().view().git_head().as_normal().cloned();
let name = jj.workspace_name().to_owned();
let wc_commit = match head {
Some(head_id) => {
let head = tx
.repo_mut()
.store()
.get_commit(&head_id)
.map_err(engine_err)?;
let wc_commit = tx
.repo_mut()
.new_commit(vec![head_id], head.tree())
.set_author(signature(actor))
.write()
.await
.map_err(engine_err)?;
tx.repo_mut()
.set_wc_commit(name, wc_commit.id().clone())
.map_err(engine_err)?;
tx.repo_mut()
.rebase_descendants()
.await
.map_err(engine_err)?;
Some(wc_commit)
}
None => None,
};
if let Some(wc_commit) = &wc_commit {
git::reset_head(tx.repo_mut(), wc_commit)
.await
.map_err(engine_err)?;
}
let repo = tx.commit("adopt git repo").await.map_err(engine_err)?;
if let Some(wc_commit) = &wc_commit {
jj.check_out(repo.op_id().clone(), None, wc_commit)
.await
.map_err(engine_err)?;
}
Ok(Self {
jj,
repo,
_settings: settings,
boundary: boundary.to_vec(),
})
}
pub fn snapshot(&mut self) -> Result<Option<String>, Error> {
block_on(self.snapshot_with(&SnapshotStyle::Stack))
}
pub fn snapshot_amend(&mut self) -> Result<Option<String>, Error> {
block_on(self.snapshot_with(&SnapshotStyle::Amend))
}
async fn snapshot_with(&mut self, style: &SnapshotStyle) -> Result<Option<String>, Error> {
let name = self.jj.workspace_name().to_owned();
let wc_id = match self.repo.view().get_wc_commit_id(&name) {
Some(id) => id.clone(),
None => return Err(Error::Engine("no working-copy commit".to_owned())),
};
let options = snapshot_options(base_ignores(&self.boundary)?);
let mut locked = self
.jj
.start_working_copy_mutation()
.await
.map_err(engine_err)?;
let (new_tree, stats) = match locked.locked_wc().snapshot(&options).await {
Ok(result) => result,
Err(err) => {
release_at_old_operation(locked).await?;
return Err(engine_err(err));
}
};
if !stats.invalid_utf8_paths.is_empty() {
release_at_old_operation(locked).await?;
return Err(Error::Engine(
"working copy has paths with invalid utf-8 names".to_owned(),
));
}
let wc_commit = self.repo.store().get_commit(&wc_id).map_err(engine_err)?;
if new_tree.tree_ids() == wc_commit.tree_ids() {
release_at_old_operation(locked).await?;
return Ok(None);
}
let mut tx = self.repo.start_transaction();
tx.set_is_snapshot(true);
let new_commit = match style {
SnapshotStyle::Stack => tx
.repo_mut()
.new_commit(vec![wc_id], new_tree)
.write()
.await
.map_err(engine_err)?,
SnapshotStyle::Amend => tx
.repo_mut()
.rewrite_commit(&wc_commit)
.set_tree(new_tree)
.write()
.await
.map_err(engine_err)?,
};
let new_id = new_commit.id().clone();
tx.repo_mut()
.set_wc_commit(name, new_id.clone())
.map_err(engine_err)?;
tx.repo_mut()
.rebase_descendants()
.await
.map_err(engine_err)?;
if let SnapshotStyle::Stack = style {
git::reset_head(tx.repo_mut(), &new_commit)
.await
.map_err(engine_err)?;
}
let repo = tx.commit("snapshot").await.map_err(engine_err)?;
locked
.finish(repo.op_id().clone())
.await
.map_err(engine_err)?;
self.repo = repo;
Ok(Some(new_id.hex()))
}
pub fn log(&self, limit: usize) -> Result<Vec<Snapshot>, Error> {
let name = self.jj.workspace_name().to_owned();
let wc_id = match self.repo.view().get_wc_commit_id(&name) {
Some(id) => id.clone(),
None => return Err(Error::Engine("no working-copy commit".to_owned())),
};
let root = self.repo.store().root_commit_id().clone();
let mut out = Vec::new();
let mut current = Some(wc_id);
while let Some(id) = current {
if out.len() >= limit {
break;
}
let commit = self.repo.store().get_commit(&id).map_err(engine_err)?;
let parents: Vec<String> = commit
.parent_ids()
.iter()
.filter(|parent| **parent != root)
.map(ObjectId::hex)
.collect();
let author = commit.author();
out.push(Snapshot {
id: id.hex(),
actor: author.name.clone(),
at_ms: author.timestamp.timestamp.0,
parents,
});
current = commit
.parent_ids()
.iter()
.find(|parent| **parent != root)
.cloned();
}
Ok(out)
}
pub fn head(&self) -> Result<String, Error> {
Ok(self.wc_commit_id()?.hex())
}
pub fn parent_of(&self, id: &str) -> Result<String, Error> {
let commit = self.commit_at(id)?;
match commit.parent_ids().first() {
Some(parent) => Ok(parent.hex()),
None => Err(Error::Engine(format!("snapshot {id} has no parent"))),
}
}
pub fn create_session_workspace(
&mut self,
root: &Path,
name: &str,
actor: &Actor,
) -> Result<String, Error> {
block_on(self.create_session_workspace_async(root, name, actor))
}
async fn create_session_workspace_async(
&mut self,
root: &Path,
name: &str,
actor: &Actor,
) -> Result<String, Error> {
let head_id = self.wc_commit_id()?;
let head = self.repo.store().get_commit(&head_id).map_err(engine_err)?;
fs::create_dir_all(root)?;
let (mut session_ws, repo) = JjWorkspace::init_workspace_with_existing_repo(
root,
self.jj.repo_path(),
&self.repo,
&*default_working_copy_factory(),
WorkspaceNameBuf::from(name),
)
.await
.map_err(engine_err)?;
let mut tx = repo.start_transaction();
let wc_commit = tx
.repo_mut()
.new_commit(vec![head_id], head.tree())
.set_author(signature(actor))
.write()
.await
.map_err(engine_err)?;
tx.repo_mut()
.edit(WorkspaceNameBuf::from(name), &wc_commit)
.await
.map_err(engine_err)?;
tx.repo_mut()
.rebase_descendants()
.await
.map_err(engine_err)?;
let repo = tx.commit("open session").await.map_err(engine_err)?;
session_ws
.check_out(repo.op_id().clone(), None, &wc_commit)
.await
.map_err(engine_err)?;
self.repo = repo;
Ok(wc_commit.change_id().hex())
}
pub fn land(&mut self, tip: &str, bookmark: &str) -> Result<LandOutcome, Error> {
block_on(self.land_async(tip, bookmark))
}
async fn land_async(&mut self, tip: &str, bookmark: &str) -> Result<LandOutcome, Error> {
let name = self.jj.workspace_name().to_owned();
let head_id = self.wc_commit_id()?;
let tip_commit = self.commit_at(tip)?;
let mut tx = self.repo.start_transaction();
let rebased = rebase_commit(tx.repo_mut(), tip_commit, vec![head_id])
.await
.map_err(engine_err)?;
if rebased.has_conflict() {
return Ok(LandOutcome::Conflicted);
}
tx.repo_mut()
.set_wc_commit(name, rebased.id().clone())
.map_err(engine_err)?;
tx.repo_mut()
.rebase_descendants()
.await
.map_err(engine_err)?;
git::reset_head(tx.repo_mut(), &rebased)
.await
.map_err(engine_err)?;
tx.repo_mut().set_local_bookmark_target(
RefName::new(bookmark),
RefTarget::normal(rebased.id().clone()),
);
let exported = git::export_refs(tx.repo_mut()).map_err(engine_err)?;
if !exported.failed_bookmarks.is_empty() {
return Err(Error::Engine(format!(
"bookmark {bookmark:?} failed to export: {:?}",
exported.failed_bookmarks
)));
}
let repo = tx.commit("land").await.map_err(engine_err)?;
self.repo = repo;
self.jj
.check_out(self.repo.op_id().clone(), None, &rebased)
.await
.map_err(engine_err)?;
Ok(LandOutcome::Landed {
snapshot: rebased.id().hex(),
})
}
pub fn step_back(&mut self, landed: &str, bookmark: &str) -> Result<StepBack, Error> {
block_on(self.step_back_async(landed, bookmark))
}
async fn step_back_async(&mut self, landed: &str, bookmark: &str) -> Result<StepBack, Error> {
let name = self.jj.workspace_name().to_owned();
let head = self.wc_commit_id()?;
let landed_commit = self.commit_at(landed)?;
let Some(parent_id) = landed_commit.parent_ids().first() else {
return Err(Error::Engine(format!(
"the landed snapshot {landed} has no parent to step back to"
)));
};
if head == *parent_id {
return Ok(StepBack::AlreadyStepped);
}
if head.hex() != landed {
return Ok(StepBack::LineMoved { head: head.hex() });
}
let parent = self
.repo
.store()
.get_commit(parent_id)
.map_err(engine_err)?;
let mut tx = self.repo.start_transaction();
tx.repo_mut()
.set_wc_commit(name, parent_id.clone())
.map_err(engine_err)?;
git::reset_head(tx.repo_mut(), &parent)
.await
.map_err(engine_err)?;
tx.repo_mut().set_local_bookmark_target(
RefName::new(bookmark),
RefTarget::normal(parent_id.clone()),
);
let exported = git::export_refs(tx.repo_mut()).map_err(engine_err)?;
if !exported.failed_bookmarks.is_empty() {
return Err(Error::Engine(format!(
"bookmark {bookmark:?} failed to export: {:?}",
exported.failed_bookmarks
)));
}
let repo = tx.commit("undo").await.map_err(engine_err)?;
self.repo = repo;
self.jj
.check_out(self.repo.op_id().clone(), None, &parent)
.await
.map_err(engine_err)?;
Ok(StepBack::Stepped {
restored: parent_id.hex(),
})
}
pub fn tree_changed(&self, id: &str) -> Result<bool, Error> {
let commit = self.commit_at(id)?;
let Some(parent_id) = commit.parent_ids().first() else {
return Ok(true);
};
let parent = self
.repo
.store()
.get_commit(parent_id)
.map_err(engine_err)?;
Ok(commit.tree_ids() != parent.tree_ids())
}
fn wc_commit_id(&self) -> Result<CommitId, Error> {
let name = self.jj.workspace_name().to_owned();
match self.repo.view().get_wc_commit_id(&name) {
Some(id) => Ok(id.clone()),
None => Err(Error::Engine("no working-copy commit".to_owned())),
}
}
fn commit_at(&self, id: &str) -> Result<jj_lib::commit::Commit, Error> {
let Some(commit_id) = CommitId::try_from_hex(id) else {
return Err(Error::Engine(format!("not a snapshot id: {id}")));
};
self.repo.store().get_commit(&commit_id).map_err(engine_err)
}
pub fn diff_latest(&self) -> Result<(Diff, DiffSides), Error> {
block_on(self.diff_latest_async())
}
async fn diff_latest_async(&self) -> Result<(Diff, DiffSides), Error> {
let name = self.jj.workspace_name().to_owned();
let wc_id = match self.repo.view().get_wc_commit_id(&name) {
Some(id) => id.clone(),
None => return Err(Error::Engine("no working-copy commit".to_owned())),
};
let root = self.repo.store().root_commit_id().clone();
let commit = self.repo.store().get_commit(&wc_id).map_err(engine_err)?;
let new_tree = commit.tree();
let parent = commit
.parent_ids()
.iter()
.find(|parent| **parent != root)
.cloned();
let old_tree = match parent {
Some(parent_id) => self
.repo
.store()
.get_commit(&parent_id)
.map_err(engine_err)?
.tree(),
None => self.empty_tree()?,
};
self.tree_diff(old_tree, new_tree).await
}
pub fn diff_between(&self, before: &str, after: &str) -> Result<(Diff, DiffSides), Error> {
block_on(self.diff_between_async(before, after))
}
async fn diff_between_async(
&self,
before: &str,
after: &str,
) -> Result<(Diff, DiffSides), Error> {
let old_tree = self.tree_at(before)?;
let new_tree = self.tree_at(after)?;
self.tree_diff(old_tree, new_tree).await
}
pub fn read_file_sides(&self, sides: &DiffSides, path: &str) -> Result<(Side, Side), Error> {
block_on(async {
let path = RepoPath::from_internal_string(path).map_err(engine_err)?;
let before = self.file_blob(&sides.before, path).await?;
let after = self.file_blob(&sides.after, path).await?;
Ok((before, after))
})
}
async fn file_blob(&self, tree: &MergedTree, path: &RepoPath) -> Result<Side, Error> {
let value = tree.path_value(path).await.map_err(engine_err)?;
let Some(Some(TreeValue::File { id, .. })) = value.as_resolved() else {
return Ok(Side::Absent);
};
let reader = self
.repo
.store()
.read_file(path, id)
.await
.map_err(engine_err)?;
let mut bytes = Vec::new();
reader
.take(LADDER_FILE_SIZE_MAX + 1)
.read_to_end(&mut bytes)
.await
.map_err(engine_err)?;
if bytes.len() as u64 > LADDER_FILE_SIZE_MAX {
return Ok(Side::TooLarge);
}
Ok(Side::Blob(FileBlob {
id: id.hex(),
bytes,
}))
}
async fn tree_diff(
&self,
old_tree: MergedTree,
new_tree: MergedTree,
) -> Result<(Diff, DiffSides), Error> {
let mut before = BTreeMap::new();
let mut after = BTreeMap::new();
let mut stream = old_tree.diff_stream(&new_tree, &EverythingMatcher);
while let Some(entry) = stream.next().await {
let path = entry.path.as_internal_file_string().to_owned();
let values = entry.values.map_err(engine_err)?;
if values.before.is_present() {
before.insert(path.clone(), format!("{:?}", values.before));
}
if values.after.is_present() {
after.insert(path, format!("{:?}", values.after));
}
}
drop(stream);
let diff = diff_listings(&before, &after);
Ok((
diff,
DiffSides {
before: old_tree,
after: new_tree,
},
))
}
fn tree_at(&self, id: &str) -> Result<MergedTree, Error> {
Ok(self.commit_at(id)?.tree())
}
pub fn export_tree(&self, id: &str, dest: &Path) -> Result<(), Error> {
block_on(self.export_tree_async(id, dest))
}
async fn export_tree_async(&self, id: &str, dest: &Path) -> Result<(), Error> {
let empty = self.empty_tree()?;
let tree = self.tree_at(id)?;
let mut kept: BTreeSet<String> = BTreeSet::new();
let mut stream = empty.diff_stream(&tree, &EverythingMatcher);
while let Some(entry) = stream.next().await {
let values = entry.values.map_err(engine_err)?;
let rel = entry.path.as_internal_file_string().to_owned();
let Some(value) = values.after.as_resolved() else {
return Err(Error::Engine(format!("conflicted tree entry at {rel}")));
};
let Some(value) = value else {
continue;
};
let target = dest.join(&rel);
if let Some(parent) = target.parent() {
fs::create_dir_all(parent)?;
}
match value {
TreeValue::File { id, executable, .. } => {
let mut reader = self
.repo
.store()
.read_file(&entry.path, id)
.await
.map_err(engine_err)?;
let mut bytes = Vec::new();
reader.read_to_end(&mut bytes).await.map_err(engine_err)?;
fs::write(&target, &bytes)?;
if *executable {
let mut permissions = fs::metadata(&target)?.permissions();
permissions.set_mode(0o755);
fs::set_permissions(&target, permissions)?;
}
}
TreeValue::Symlink(id) => {
let link = self
.repo
.store()
.read_symlink(&entry.path, id)
.await
.map_err(engine_err)?;
if target.symlink_metadata().is_ok() {
fs::remove_file(&target)?;
}
std::os::unix::fs::symlink(&link, &target)?;
}
other => {
return Err(Error::Engine(format!("cannot export {rel}: {other:?}")));
}
}
kept.insert(rel);
}
drop(stream);
remove_unkept(dest, &kept)
}
fn empty_tree(&self) -> Result<MergedTree, Error> {
let root = self.repo.store().root_commit_id().clone();
let commit = self.repo.store().get_commit(&root).map_err(engine_err)?;
Ok(commit.tree())
}
}
fn build_settings(actor: &Actor) -> Result<UserSettings, Error> {
#[derive(serde::Serialize)]
struct UserConfig<'a> {
user: UserSection<'a>,
}
#[derive(serde::Serialize)]
struct UserSection<'a> {
name: &'a str,
email: String,
}
let mut config = StackedConfig::with_defaults();
let text = toml::to_string(&UserConfig {
user: UserSection {
name: &actor.name,
email: format!("{}@atelier.local", actor.name),
},
})
.map_err(config_err)?;
let layer = ConfigLayer::parse(ConfigSource::User, &text).map_err(config_err)?;
config.add_layer(layer);
UserSettings::from_config(config).map_err(config_err)
}
fn signature(actor: &Actor) -> Signature {
Signature {
name: actor.name.clone(),
email: format!("{}@atelier.local", actor.name),
timestamp: Timestamp::now(),
}
}
async fn release_at_old_operation(mut locked: LockedWorkspace<'_>) -> Result<(), Error> {
let operation = locked.locked_wc().old_operation_id().clone();
locked.finish(operation).await.map_err(engine_err)?;
Ok(())
}
fn snapshot_options(base_ignores: Arc<GitIgnoreFile>) -> SnapshotOptions<'static> {
SnapshotOptions {
base_ignores,
progress: None,
start_tracking_matcher: &EverythingMatcher,
force_tracking_matcher: &NothingMatcher,
max_new_file_size: NEW_FILE_SIZE_MAX,
}
}
fn base_ignores(boundary: &[String]) -> Result<Arc<GitIgnoreFile>, Error> {
let mut rules = String::from(".atelier/\n.git/\n.jj/\n");
for name in boundary {
rules.push('/');
rules.push_str(name);
rules.push_str("/\n");
}
GitIgnoreFile::empty()
.chain(RepoPath::root(), Path::new(".gitignore"), rules.as_bytes())
.map_err(engine_err)
}
fn remove_unkept(root: &Path, kept: &BTreeSet<String>) -> Result<(), Error> {
let mut directories: Vec<PathBuf> = Vec::new();
let mut pending = vec![root.to_path_buf()];
while let Some(dir) = pending.pop() {
for entry in fs::read_dir(&dir)? {
let entry = entry?;
let name = entry.file_name();
let Some(name) = name.to_str() else {
return Err(Error::Engine(format!(
"cannot mirror over a non-utf8 name at {}",
entry.path().display()
)));
};
if SKIP_NAMES.contains(&name) {
continue;
}
let path = entry.path();
if entry.file_type()?.is_dir() {
directories.push(path.clone());
pending.push(path);
} else {
let rel = path
.strip_prefix(root)
.map_err(engine_err)?
.components()
.map(|component| component.as_os_str().to_string_lossy())
.collect::<Vec<_>>()
.join("/");
if !kept.contains(&rel) {
fs::remove_file(&path)?;
}
}
}
}
directories.sort_by_key(|dir| std::cmp::Reverse(dir.components().count()));
for dir in directories {
if fs::read_dir(&dir)?.next().is_none() {
fs::remove_dir(&dir)?;
}
}
Ok(())
}