Skip to main content

ironflow_ops_git/
fetch.rs

1//! Fetch and push operations.
2//!
3//! Every operation here talks to a remote. Over HTTPS they authenticate with
4//! the token held in the `git_token` secret (username `oauth2`), both
5//! overridable with `token_secret` and `username`. Without that secret they
6//! fall back to the SSH agent and the git credential helper.
7
8use std::path::PathBuf;
9
10use async_trait::async_trait;
11use git2::{Direction, FetchOptions, PushOptions, Repository};
12use ironflow_core::error::OperationError;
13use ironflow_core::operation::{Operation, OperationContext, TypedOperation};
14use serde::{Deserialize, Serialize};
15use serde_json::Value;
16
17use crate::helpers::{
18    GitAuth, auth_builders, blocking_authenticated, credentials_callbacks, to_value,
19};
20
21#[derive(Debug, Clone, Serialize, Deserialize)]
22pub struct FetchPushOutput {
23    pub remote: String,
24    pub refspecs: Vec<String>,
25}
26
27#[derive(Debug, Clone, Serialize, Deserialize)]
28pub struct RemotePruneOutput {
29    pub remote: String,
30    pub pruned: bool,
31}
32
33#[derive(Debug, Clone, Serialize, Deserialize)]
34pub struct RemoteDefaultBranchOutput {
35    pub remote: String,
36    pub default_branch: Option<String>,
37}
38
39/// Fetch from a remote.
40///
41/// # Examples
42///
43/// ```no_run
44/// use ironflow_ops_git::fetch::FetchRemote;
45/// use ironflow_core::operation::Operation;
46///
47/// let op = FetchRemote::new("/path/to/repo", "origin", vec!["main"]);
48/// assert_eq!(op.kind(), "git");
49/// ```
50pub struct FetchRemote {
51    repo_path: PathBuf,
52    remote_name: String,
53    refspecs: Vec<String>,
54    auth: GitAuth,
55}
56
57impl FetchRemote {
58    /// Create a new fetch operation.
59    pub fn new(
60        repo_path: impl Into<PathBuf>,
61        remote_name: impl Into<String>,
62        refspecs: Vec<impl Into<String>>,
63    ) -> Self {
64        Self {
65            repo_path: repo_path.into(),
66            remote_name: remote_name.into(),
67            refspecs: refspecs.into_iter().map(Into::into).collect(),
68            auth: GitAuth::default(),
69        }
70    }
71
72    /// Execute and return a typed result.
73    ///
74    /// # Errors
75    ///
76    /// Returns [`OperationError::Secret`] if the secret store fails, and
77    /// [`OperationError::External`] if the fetch fails. The token never
78    /// appears in the message.
79    pub async fn run(&self, ctx: &OperationContext) -> Result<FetchPushOutput, OperationError> {
80        let repo_path = self.repo_path.clone();
81        let remote_name = self.remote_name.clone();
82        let refspecs = self.refspecs.clone();
83        blocking_authenticated(ctx, &self.auth, move |creds| {
84            let repo = Repository::open(&repo_path)?;
85            let mut remote = repo.find_remote(&remote_name)?;
86            let refs: Vec<&str> = refspecs.iter().map(String::as_str).collect();
87            let mut fetch_opts = FetchOptions::new();
88            fetch_opts.remote_callbacks(credentials_callbacks(creds));
89            remote.fetch(&refs, Some(&mut fetch_opts), None)?;
90            Ok(FetchPushOutput {
91                remote: remote_name,
92                refspecs,
93            })
94        })
95        .await
96    }
97}
98
99#[async_trait]
100impl Operation for FetchRemote {
101    fn kind(&self) -> &str {
102        "git"
103    }
104    async fn execute(&self, ctx: &OperationContext) -> Result<Value, OperationError> {
105        to_value(&self.run(ctx).await?)
106    }
107    fn input(&self) -> Option<Value> {
108        Some(serde_json::json!({ "repo_path": self.repo_path, "remote": self.remote_name }))
109    }
110}
111
112impl TypedOperation for FetchRemote {
113    type Output = FetchPushOutput;
114}
115
116auth_builders!(
117    FetchRemote,
118    "fetch",
119    "\"/path/to/repo\", \"origin\", vec![\"main\"]"
120);
121
122/// Push to a remote.
123///
124/// # Examples
125///
126/// ```no_run
127/// use ironflow_ops_git::fetch::PushRemote;
128/// use ironflow_core::operation::Operation;
129///
130/// let op = PushRemote::new("/path/to/repo", "origin", vec!["refs/heads/main"]);
131/// assert_eq!(op.kind(), "git");
132/// ```
133pub struct PushRemote {
134    repo_path: PathBuf,
135    remote_name: String,
136    refspecs: Vec<String>,
137    auth: GitAuth,
138}
139
140impl PushRemote {
141    /// Create a new push operation.
142    pub fn new(
143        repo_path: impl Into<PathBuf>,
144        remote_name: impl Into<String>,
145        refspecs: Vec<impl Into<String>>,
146    ) -> Self {
147        Self {
148            repo_path: repo_path.into(),
149            remote_name: remote_name.into(),
150            refspecs: refspecs.into_iter().map(Into::into).collect(),
151            auth: GitAuth::default(),
152        }
153    }
154
155    /// Execute and return a typed result.
156    ///
157    /// # Errors
158    ///
159    /// Returns [`OperationError::Secret`] if the secret store fails, and
160    /// [`OperationError::External`] if the push fails. The token never
161    /// appears in the message.
162    pub async fn run(&self, ctx: &OperationContext) -> Result<FetchPushOutput, OperationError> {
163        let repo_path = self.repo_path.clone();
164        let remote_name = self.remote_name.clone();
165        let refspecs = self.refspecs.clone();
166        blocking_authenticated(ctx, &self.auth, move |creds| {
167            let repo = Repository::open(&repo_path)?;
168            let mut remote = repo.find_remote(&remote_name)?;
169            let refs: Vec<&str> = refspecs.iter().map(String::as_str).collect();
170            let mut push_opts = PushOptions::new();
171            push_opts.remote_callbacks(credentials_callbacks(creds));
172            remote.push(&refs, Some(&mut push_opts))?;
173            Ok(FetchPushOutput {
174                remote: remote_name,
175                refspecs,
176            })
177        })
178        .await
179    }
180}
181
182#[async_trait]
183impl Operation for PushRemote {
184    fn kind(&self) -> &str {
185        "git"
186    }
187    async fn execute(&self, ctx: &OperationContext) -> Result<Value, OperationError> {
188        to_value(&self.run(ctx).await?)
189    }
190    fn input(&self) -> Option<Value> {
191        Some(serde_json::json!({ "repo_path": self.repo_path, "remote": self.remote_name }))
192    }
193}
194
195impl TypedOperation for PushRemote {
196    type Output = FetchPushOutput;
197}
198
199auth_builders!(
200    PushRemote,
201    "fetch",
202    "\"/path/to/repo\", \"origin\", vec![\"refs/heads/main\"]"
203);
204
205/// Prune stale remote-tracking branches.
206///
207/// Connects to the remote to list its branches, then deletes the
208/// remote-tracking refs whose branch no longer exists there.
209///
210/// # Examples
211///
212/// ```no_run
213/// use ironflow_ops_git::fetch::RemotePrune;
214/// use ironflow_core::operation::Operation;
215///
216/// let op = RemotePrune::new("/path/to/repo", "origin");
217/// assert_eq!(op.kind(), "git");
218/// ```
219pub struct RemotePrune {
220    repo_path: PathBuf,
221    remote_name: String,
222    auth: GitAuth,
223}
224
225impl RemotePrune {
226    /// Create a new prune operation.
227    pub fn new(repo_path: impl Into<PathBuf>, remote_name: impl Into<String>) -> Self {
228        Self {
229            repo_path: repo_path.into(),
230            remote_name: remote_name.into(),
231            auth: GitAuth::default(),
232        }
233    }
234
235    /// Execute and return a typed result.
236    ///
237    /// # Errors
238    ///
239    /// Returns [`OperationError::Secret`] if the secret store fails, and
240    /// [`OperationError::External`] if the connection or the prune fails. The
241    /// token never appears in the message.
242    pub async fn run(&self, ctx: &OperationContext) -> Result<RemotePruneOutput, OperationError> {
243        let repo_path = self.repo_path.clone();
244        let remote_name = self.remote_name.clone();
245        blocking_authenticated(ctx, &self.auth, move |creds| {
246            let repo = Repository::open(&repo_path)?;
247            let mut remote = repo.find_remote(&remote_name)?;
248            let mut connection =
249                remote.connect_auth(Direction::Fetch, Some(credentials_callbacks(creds)), None)?;
250            connection.remote().prune(None)?;
251            Ok(RemotePruneOutput {
252                remote: remote_name,
253                pruned: true,
254            })
255        })
256        .await
257    }
258}
259
260#[async_trait]
261impl Operation for RemotePrune {
262    fn kind(&self) -> &str {
263        "git"
264    }
265    async fn execute(&self, ctx: &OperationContext) -> Result<Value, OperationError> {
266        to_value(&self.run(ctx).await?)
267    }
268    fn input(&self) -> Option<Value> {
269        Some(serde_json::json!({ "repo_path": self.repo_path, "remote": self.remote_name }))
270    }
271}
272
273impl TypedOperation for RemotePrune {
274    type Output = RemotePruneOutput;
275}
276
277auth_builders!(RemotePrune, "fetch", "\"/path/to/repo\", \"origin\"");
278
279/// Get the default branch of a remote.
280///
281/// Connects to the remote and reads the branch its `HEAD` points to.
282///
283/// # Examples
284///
285/// ```no_run
286/// use ironflow_ops_git::fetch::RemoteDefaultBranch;
287/// use ironflow_core::operation::Operation;
288///
289/// let op = RemoteDefaultBranch::new("/path/to/repo", "origin");
290/// assert_eq!(op.kind(), "git");
291/// ```
292pub struct RemoteDefaultBranch {
293    repo_path: PathBuf,
294    remote_name: String,
295    auth: GitAuth,
296}
297
298impl RemoteDefaultBranch {
299    /// Create a new default-branch query operation.
300    pub fn new(repo_path: impl Into<PathBuf>, remote_name: impl Into<String>) -> Self {
301        Self {
302            repo_path: repo_path.into(),
303            remote_name: remote_name.into(),
304            auth: GitAuth::default(),
305        }
306    }
307
308    /// Execute and return a typed result.
309    ///
310    /// # Errors
311    ///
312    /// Returns [`OperationError::Secret`] if the secret store fails, and
313    /// [`OperationError::External`] if the connection fails or the remote
314    /// has no default branch. The token never appears in the message.
315    pub async fn run(
316        &self,
317        ctx: &OperationContext,
318    ) -> Result<RemoteDefaultBranchOutput, OperationError> {
319        let repo_path = self.repo_path.clone();
320        let remote_name = self.remote_name.clone();
321        blocking_authenticated(ctx, &self.auth, move |creds| {
322            let repo = Repository::open(&repo_path)?;
323            let mut remote = repo.find_remote(&remote_name)?;
324            let connection =
325                remote.connect_auth(Direction::Fetch, Some(credentials_callbacks(creds)), None)?;
326            let default = connection.default_branch()?;
327            let name = default.as_str().map(String::from);
328            Ok(RemoteDefaultBranchOutput {
329                remote: remote_name,
330                default_branch: name,
331            })
332        })
333        .await
334    }
335}
336
337#[async_trait]
338impl Operation for RemoteDefaultBranch {
339    fn kind(&self) -> &str {
340        "git"
341    }
342    async fn execute(&self, ctx: &OperationContext) -> Result<Value, OperationError> {
343        to_value(&self.run(ctx).await?)
344    }
345    fn input(&self) -> Option<Value> {
346        Some(serde_json::json!({ "repo_path": self.repo_path, "remote": self.remote_name }))
347    }
348}
349
350impl TypedOperation for RemoteDefaultBranch {
351    type Output = RemoteDefaultBranchOutput;
352}
353
354auth_builders!(
355    RemoteDefaultBranch,
356    "fetch",
357    "\"/path/to/repo\", \"origin\""
358);
359
360#[cfg(test)]
361mod tests {
362    use git2::Repository;
363
364    use super::*;
365    use crate::test_helpers::{ctx, init_repo};
366
367    // Push and fetch over a local bare remote (file://) never hit the
368    // credentials callback, but exercising them proves that wiring
369    // PushOptions/FetchOptions with RemoteCallbacks does not break the normal
370    // flow. The SSH-agent path (the literal issue symptom) needs a live SSH
371    // server and is covered in Out of Test Scope.
372    fn setup_with_bare_remote() -> (tempfile::TempDir, tempfile::TempDir) {
373        let work = tempfile::tempdir().unwrap();
374        init_repo(work.path());
375        let bare = tempfile::tempdir().unwrap();
376        Repository::init_bare(bare.path()).unwrap();
377        let repo = Repository::open(work.path()).unwrap();
378        repo.remote("origin", bare.path().to_str().unwrap())
379            .unwrap();
380        (work, bare)
381    }
382
383    #[tokio::test]
384    async fn push_to_local_bare_remote() {
385        let (work, bare) = setup_with_bare_remote();
386        let result = PushRemote::new(
387            work.path(),
388            "origin",
389            vec!["refs/heads/master:refs/heads/master"],
390        )
391        .run(&ctx())
392        .await
393        .unwrap();
394        assert_eq!(result.remote, "origin");
395        // The ref must exist on the remote side after the push.
396        let bare_repo = Repository::open(bare.path()).unwrap();
397        assert!(bare_repo.find_reference("refs/heads/master").is_ok());
398    }
399
400    #[tokio::test]
401    async fn fetch_from_local_bare_remote() {
402        let (work, bare) = setup_with_bare_remote();
403        PushRemote::new(
404            work.path(),
405            "origin",
406            vec!["refs/heads/master:refs/heads/master"],
407        )
408        .run(&ctx())
409        .await
410        .unwrap();
411        // A second clone fetching from the same bare remote must succeed.
412        let clone_dir = tempfile::tempdir().unwrap();
413        let clone = Repository::init(clone_dir.path()).unwrap();
414        clone
415            .remote("origin", bare.path().to_str().unwrap())
416            .unwrap();
417        let result = FetchRemote::new(clone_dir.path(), "origin", vec!["master"])
418            .run(&ctx())
419            .await
420            .unwrap();
421        assert_eq!(result.remote, "origin");
422    }
423
424    #[tokio::test]
425    async fn push_missing_remote_fails() {
426        let work = tempfile::tempdir().unwrap();
427        init_repo(work.path());
428        let result = PushRemote::new(work.path(), "nope", vec!["refs/heads/master"])
429            .run(&ctx())
430            .await;
431        assert!(result.is_err());
432    }
433}