use std::path::{Path, PathBuf};
use serde_json::{json, Map, Value};
use onevcs::{
ChangeSpec, Error, EventKind, FailureKind, Hosting, Identity, Lifecycle, MergeOutcome,
MergePolicy, PreservedBranch, Provenance, Publication, PublishOutcome, PublishRequest,
Recoverable, Result, Scope, Session, SessionRecord, SessionRequest, SessionToken, Sha, Vcs,
};
use crate::events::{self, Emission};
use crate::remote::DEFAULT_HOST;
use crate::state::{self, VcsState};
use crate::store::{FileStore, MemoryStore, Store};
pub const DEFAULT_BASE: &str = "main";
pub const DEFAULT_PUBLICATION: MergePolicy = MergePolicy::ChangeOpen;
#[derive(Debug)]
pub struct Repository<T> {
store: T,
root: PathBuf,
trees: Trees,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Trees {
Named,
Created,
}
pub type MemoryVcs = Repository<MemoryStore<VcsState>>;
pub type FileVcs = Repository<FileStore<VcsState>>;
impl MemoryVcs {
pub fn new() -> Self {
Self::seeded(VcsState::default())
}
pub fn seeded(state: VcsState) -> Self {
Self {
store: MemoryStore::new(state),
root: std::env::temp_dir().join("onevcs-testing-memory"),
trees: Trees::Named,
}
}
pub fn state(&self) -> VcsState {
self.store
.snapshot()
.expect("an in-memory store always answers")
}
}
impl Default for MemoryVcs {
fn default() -> Self {
Self::new()
}
}
impl FileVcs {
pub fn create(path: impl Into<PathBuf>) -> Result<Self> {
Self::over(FileStore::attach(path, &VcsState::default())?)
}
pub fn seeded(path: impl Into<PathBuf>, state: VcsState) -> Result<Self> {
Self::over(FileStore::replace(path, &state)?)
}
fn over(store: FileStore<VcsState>) -> Result<Self> {
let root = store
.path()
.parent()
.filter(|parent| !parent.as_os_str().is_empty())
.unwrap_or_else(|| Path::new("."))
.join("worktrees");
Ok(Self {
store,
root,
trees: Trees::Created,
})
}
pub fn state(&self) -> Result<VcsState> {
self.store.snapshot()
}
}
impl<T: Store<VcsState>> Vcs for Repository<T> {
fn resolve_identity(&self, origin_or_path: &str) -> Result<Identity> {
self.store.with(|state| {
state::identity_of(state, origin_or_path)
.cloned()
.ok_or_else(|| Error::Invalid {
reason: format!(
"{origin_or_path:?} does not name a repository this provider knows; {}",
state::known(state)
),
})
})
}
fn open_session(&self, req: SessionRequest) -> Result<Session> {
let root = self.root.clone();
let (session, emission) = self.store.with(|state| {
let identity = state::identity_of(state, &req.repo)
.cloned()
.ok_or_else(|| Error::Invalid {
reason: format!(
"{:?} does not name a repository this provider knows; {}",
req.repo,
state::known(state)
),
})?;
let token = SessionToken(format!("s-testing-{}", state.sessions.len() + 1));
let run_root = root.join(&token.0);
let base = req.base.clone().unwrap_or_else(|| DEFAULT_BASE.to_owned());
state::named_branch(&base, "the base")?;
let session = Session {
worktree: run_root.join("worktree"),
branch: state::requested_branch(&req, &token)?,
base,
token: token.clone(),
};
state.sessions.push(session.clone());
state
.session_identities
.insert(token.clone(), identity.origin.clone());
let emission = Emission {
stream: token.0.clone(),
identity: Some(identity.origin.clone()),
kind: EventKind::SessionOpened,
payload: object(json!({
"token": token.0,
"identity": identity.origin,
"branch": session.branch,
"base": session.base,
"worktree": session.worktree.display().to_string(),
"clone": run_root.join("clone").display().to_string(),
"execution_checkout": run_root.join("checkout").display().to_string(),
"publication_checkout": run_root.join("checkout").display().to_string(),
})),
};
Ok((session, emission))
})?;
if self.trees == Trees::Created {
std::fs::create_dir_all(&session.worktree).map_err(|e| Error::Invalid {
reason: format!("cannot create {}: {e}", session.worktree.display()),
})?;
}
events::emit(&emission);
Ok(session)
}
fn adopt_session(&self, token: SessionToken) -> Result<Session> {
self.store.with(|state| {
state::session_of(state, &token)
.cloned()
.ok_or_else(|| Error::Invalid {
reason: format!(
"no session {:?} is open; `onevcs session open` prints a token",
token.0
),
})
})
}
fn preserve(&self, s: &Session, provenance: Provenance) -> Result<PreservedBranch> {
let (branch, emission) = self.store.with(|state| {
let identity = state::identity_for(state, &s.token)?;
let branch = PreservedBranch {
branch: s.branch.clone(),
base: s.base.clone(),
provenance,
change_url: None,
change_base: None,
};
let row = Recoverable {
identity: identity.clone(),
branch: branch.clone(),
checkout: s.worktree.clone(),
stopped_because: format!("session {} was left open", s.token.0),
recover_command: recover_command(&s.branch, &s.worktree, provenance),
};
state.preserved.retain(|kept| {
kept.identity != row.identity || kept.branch.branch != row.branch.branch
});
state.preserved.push(row);
let emission = Emission {
stream: s.token.0.clone(),
identity: None,
kind: EventKind::CommitPreserved,
payload: object(json!({
"branch": s.branch,
"sha": events::stable_sha(&[&s.token.0, &s.branch, spell(provenance)]),
"provenance": spell(provenance),
})),
};
Ok((branch, emission))
})?;
events::emit(&emission);
Ok(branch)
}
fn session(&self, token: &SessionToken) -> Result<SessionRecord> {
self.store.with(|state| {
let session = state::session_of(state, token)
.cloned()
.ok_or_else(|| unknown_session(token))?;
let identity = state::identity_for(state, token)?;
let provenance = state
.preserved
.iter()
.find(|row| row.identity == identity && row.branch.branch == session.branch)
.map_or(Provenance::Complete, |row| row.branch.provenance);
Ok(SessionRecord {
lifecycle: if state.closed_sessions.contains(token) {
Lifecycle::Closed
} else {
Lifecycle::Open
},
session,
identity,
provenance,
})
})
}
fn close_session(&self, token: &SessionToken) -> Result<Session> {
let (session, emission) = self.store.with(|state| {
let session = state::session_of(state, token)
.cloned()
.ok_or_else(|| unknown_session(token))?;
state.closed_sessions.insert(token.clone());
let emission = Emission {
stream: token.0.clone(),
identity: None,
kind: EventKind::SessionClosed,
payload: object(json!({"token": token.0, "branch": session.branch})),
};
Ok((session, emission))
})?;
events::emit(&emission);
Ok(session)
}
fn publish(
&self,
token: &SessionToken,
request: &PublishRequest,
hosting: &dyn Hosting,
) -> Result<Publication> {
let (publication, emissions) = self.store.with(|state| {
let session = state::session_of(state, token)
.cloned()
.ok_or_else(|| unknown_session(token))?;
let identity = state::identity_for(state, token)?;
let resolved = state.policy.unwrap_or(DEFAULT_PUBLICATION);
let policy = match request.policy {
Some(requested) => resolved.narrow(requested)?,
None => resolved,
};
let published = |outcome, emissions| {
(
Publication {
session: token.clone(),
branch: session.branch.clone(),
policy,
outcome,
},
emissions,
)
};
if state.publications.iter().any(|earlier| {
earlier.session == *token && matches!(earlier.outcome, PublishOutcome::Merged(_))
}) {
let (publication, emissions) =
published(PublishOutcome::NothingToPublish, Vec::new());
state.publications.push(publication.clone());
return Ok((publication, emissions));
}
let (outcome, emissions) = if policy == MergePolicy::LocalDirect {
record_local_landing(&identity, &session, token)
} else {
match slug(&identity) {
Some(slug) => match publish_as_change(
hosting, &slug, &identity, &session, policy, request, token,
) {
Ok(published) => published,
Err(error) => (failed(&error), Vec::new()),
},
None => (refusal(&identity), Vec::new()),
}
};
let (publication, emissions) = published(outcome, emissions);
state.publications.push(publication.clone());
Ok((publication, emissions))
})?;
for emission in &emissions {
events::emit(emission);
}
Ok(publication)
}
fn recoverable(&self, scope: Scope) -> Result<Vec<Recoverable>> {
self.store.with(|state| {
let wanted = match &scope {
Scope::All => None,
Scope::Repo(repo) => Some(
state::identity_of(state, repo)
.map(|identity| identity.origin.clone())
.ok_or_else(|| Error::Invalid {
reason: format!(
"{repo:?} does not name a repository this provider knows; {}",
state::known(state)
),
})?,
),
};
Ok(state
.preserved
.iter()
.rev()
.filter(|row| wanted.as_ref().is_none_or(|key| *key == row.identity))
.cloned()
.collect())
})
}
}
fn unknown_session(token: &SessionToken) -> Error {
Error::Invalid {
reason: format!(
"no session {:?} is open; `onevcs session open` prints a token",
token.0
),
}
}
fn slug(identity: &str) -> Option<String> {
let mut parts = identity.split('/');
let (host, owner, name) = (parts.next()?, parts.next()?, parts.next()?);
if parts.next().is_some() || host != DEFAULT_HOST || owner.is_empty() || name.is_empty() {
return None;
}
Some(format!("{owner}/{name}"))
}
fn record_local_landing(
identity: &str,
session: &Session,
token: &SessionToken,
) -> (PublishOutcome, Vec<Emission>) {
let sha = events::stable_sha(&["publish", &token.0, &session.branch]);
let emission = Emission {
stream: token.0.clone(),
identity: Some(identity.to_owned()),
kind: EventKind::MergeCompleted,
payload: object(json!({"identity": identity, "sha": sha, "base": session.base})),
};
(PublishOutcome::Merged(Sha(sha)), vec![emission])
}
fn publish_as_change(
hosting: &dyn Hosting,
slug: &str,
identity: &str,
session: &Session,
policy: MergePolicy,
request: &PublishRequest,
token: &SessionToken,
) -> Result<(PublishOutcome, Vec<Emission>)> {
let host = hosting.for_repo(slug)?;
let author = host.authenticated_user()?;
let existing = host.find_changes(&session.branch, &session.base)?;
let change = match existing.into_iter().next() {
Some(change) => change,
None => host.open_change(ChangeSpec {
head: session.branch.clone(),
base: session.base.clone(),
title: request
.title
.as_deref()
.map_or_else(|| format!("Publish {}", session.branch), str::to_owned),
body: None,
})?,
};
let mut emissions = vec![Emission {
stream: token.0.clone(),
identity: Some(identity.to_owned()),
kind: EventKind::ChangeOpened,
payload: object(json!({
"url": change.url.to_string(),
"host": "github",
"id": change.id.0,
"base": change.base,
"author": author,
})),
}];
if policy == MergePolicy::ChangeOpen {
return Ok((PublishOutcome::ChangeOpen(change.url.clone()), emissions));
}
Ok(match host.merge(&change, policy)? {
MergeOutcome::Merged(sha) => {
emissions.push(Emission {
stream: token.0.clone(),
identity: Some(identity.to_owned()),
kind: EventKind::ChangeMerged,
payload: object(json!({"url": change.url.to_string(), "sha": sha.0})),
});
emissions.push(Emission {
stream: token.0.clone(),
identity: Some(identity.to_owned()),
kind: EventKind::MergeCompleted,
payload: object(json!({"identity": identity, "sha": sha.0})),
});
(PublishOutcome::Merged(sha), emissions)
}
MergeOutcome::Queued => (PublishOutcome::Queued(change.url.clone()), emissions),
MergeOutcome::Open => (PublishOutcome::ChangeOpen(change.url.clone()), emissions),
})
}
fn refusal(identity: &str) -> PublishOutcome {
failed(&if identity.split('/').count() == 3 {
Error::NotImplemented {
operation: "RemoteHost for a host other than github.com",
}
} else {
Error::Invalid {
reason: format!(
"identity {identity:?} is not a hosted repository, so it cannot publish a \
change request; a local identity publishes with local-direct"
),
}
})
}
fn failed(error: &Error) -> PublishOutcome {
PublishOutcome::Failed {
kind: FailureKind::of(error),
reason: error.to_string(),
retained: None,
}
}
fn recover_command(branch: &str, checkout: &Path, provenance: Provenance) -> Vec<String> {
match provenance {
Provenance::IncompleteStep => vec![
"onevcs".to_owned(),
"recover".to_owned(),
branch.to_owned(),
"--repo".to_owned(),
checkout.display().to_string(),
],
Provenance::Complete => vec![
"onevcs".to_owned(),
"integrate".to_owned(),
branch.to_owned(),
],
}
}
fn spell(provenance: Provenance) -> &'static str {
match provenance {
Provenance::Complete => "complete",
Provenance::IncompleteStep => "incomplete-step",
}
}
fn object(value: Value) -> Map<String, Value> {
value.as_object().cloned().unwrap_or_default()
}