use crate::error::{BallError, Result};
use crate::git;
use crate::negotiation::{AttemptClass, FailurePolicy, Protocol};
use crate::participant::{self, Event, EventCtx, Participant, Projection};
use crate::store::Store;
use std::path::{Path, PathBuf};
use std::process::Output;
const REMOTE: &str = "origin";
const STATE_BRANCH: &str = "balls/tasks";
const MAX_RETRIES: usize = 5;
pub fn push_state_classified(dir: &Path) -> Result<AttemptClass> {
let out = git::clean_git_command(dir)
.args(["push", REMOTE, STATE_BRANCH])
.output()?;
Ok(classify_push_output(&out))
}
fn classify_push_output(out: &Output) -> AttemptClass {
if out.status.success() {
return AttemptClass::Ok;
}
let stderr = String::from_utf8_lossy(&out.stderr).to_string();
let l = stderr.to_lowercase();
if l.contains("rejected")
&& (l.contains("non-fast-forward")
|| l.contains("fetch first")
|| l.contains("[rejected]"))
{
return AttemptClass::Conflict;
}
if is_unreachable(&l) {
return AttemptClass::Unreachable(stderr);
}
AttemptClass::Other(stderr)
}
fn is_unreachable(stderr_lower: &str) -> bool {
const UNREACHABLE_MARKERS: &[&str] = &[
"could not resolve",
"could not read from remote",
"connection refused",
"connection timed out",
"connection reset",
"repository not found",
"permission denied",
"does not appear to be a git repository",
"unable to access",
"host key verification failed",
"no such host",
"network is unreachable",
];
UNREACHABLE_MARKERS.iter().any(|m| stderr_lower.contains(m))
}
#[derive(Debug, PartialEq, Eq)]
pub enum SyncedClaimResult {
Pushed,
Lost { winner: String },
}
pub struct GitRemoteClaimProtocol<'a> {
event: Event,
store: &'a Store,
task_id: &'a str,
identity: &'a str,
state_dir: PathBuf,
}
impl<'a> GitRemoteClaimProtocol<'a> {
pub fn new(event: Event, store: &'a Store, task_id: &'a str, identity: &'a str) -> Self {
let state_dir = store.state_worktree_dir();
Self { event, store, task_id, identity, state_dir }
}
}
impl Protocol for GitRemoteClaimProtocol<'_> {
type Outcome = SyncedClaimResult;
fn propose(&mut self) -> Result<AttemptClass> {
push_state_classified(&self.state_dir)
}
fn fetch_remote_view(&mut self) -> Result<()> {
let _ = git::git_fetch(&self.state_dir, REMOTE);
let merge = git::git_merge(&self.state_dir, &format!("{REMOTE}/{STATE_BRANCH}"))?;
if matches!(merge, git::MergeResult::Conflict) {
crate::sync_resolve::auto_resolve_task_conflicts(&self.state_dir)?;
git::git_commit(&self.state_dir, "state: auto-resolve lifecycle conflicts")?;
}
Ok(())
}
fn post_merge(&mut self) -> Result<Option<SyncedClaimResult>> {
if self.event != Event::Claim {
return Ok(None);
}
let claimer = self.store.load_task(self.task_id)?.claimed_by;
if claimer.as_deref() == Some(self.identity) {
return Ok(None);
}
let winner = claimer.unwrap_or_else(|| "(unknown)".into());
let _ = push_state_classified(&self.state_dir);
Ok(Some(SyncedClaimResult::Lost { winner }))
}
fn pushed(&mut self) -> SyncedClaimResult {
SyncedClaimResult::Pushed
}
fn retry_budget(&self) -> usize {
MAX_RETRIES
}
}
pub struct GitRemoteParticipant {
projection: Projection,
subscriptions: Vec<Event>,
}
impl GitRemoteParticipant {
pub fn for_claim() -> Self {
Self::for_lifecycle(&[Event::Claim])
}
pub fn for_lifecycle(events: &[Event]) -> Self {
Self {
projection: Projection::full(),
subscriptions: events.to_vec(),
}
}
}
impl Default for GitRemoteParticipant {
fn default() -> Self {
Self::for_claim()
}
}
impl Participant for GitRemoteParticipant {
type Outcome = SyncedClaimResult;
type Protocol<'a>
= GitRemoteClaimProtocol<'a>
where
Self: 'a;
fn name(&self) -> &'static str {
"git-remote"
}
fn subscriptions(&self) -> &[Event] {
&self.subscriptions
}
fn projection(&self) -> &Projection {
&self.projection
}
fn failure_policy(&self, event: Event) -> FailurePolicy {
if self.subscriptions.contains(&event) {
FailurePolicy::Required
} else {
FailurePolicy::BestEffort
}
}
fn protocol<'a>(
&'a self,
event: Event,
ctx: EventCtx<'a>,
) -> Option<Self::Protocol<'a>> {
match event {
Event::Claim | Event::Review | Event::Close => Some(
GitRemoteClaimProtocol::new(event, ctx.store, ctx.task_id, ctx.identity),
),
_ => None,
}
}
}
pub fn push_claim(
store: &Store,
task_id: &str,
identity: &str,
) -> Result<SyncedClaimResult> {
push_state_for(store, task_id, identity, Event::Claim, "claim --sync")
}
pub fn push_state_for(
store: &Store,
task_id: &str,
identity: &str,
event: Event,
error_prefix: &str,
) -> Result<SyncedClaimResult> {
let state_dir = store.state_worktree_dir();
if !git::git_fetch(&state_dir, REMOTE)? {
return Err(BallError::Other(format!(
"{error_prefix}: cannot reach remote `{REMOTE}` (fetch failed)"
)));
}
let participant = GitRemoteParticipant::for_lifecycle(&[event]);
let ctx = EventCtx { event, store, task_id, identity };
participant::run_strict(&participant, event, ctx)
.map_err(|e| BallError::Other(format!("{error_prefix}: {e}")))
}
#[cfg(test)]
#[path = "claim_sync_tests.rs"]
mod tests;