1use 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
39pub struct FetchRemote {
51 repo_path: PathBuf,
52 remote_name: String,
53 refspecs: Vec<String>,
54 auth: GitAuth,
55}
56
57impl FetchRemote {
58 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 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
122pub struct PushRemote {
134 repo_path: PathBuf,
135 remote_name: String,
136 refspecs: Vec<String>,
137 auth: GitAuth,
138}
139
140impl PushRemote {
141 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 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
205pub struct RemotePrune {
220 repo_path: PathBuf,
221 remote_name: String,
222 auth: GitAuth,
223}
224
225impl RemotePrune {
226 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 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
279pub struct RemoteDefaultBranch {
293 repo_path: PathBuf,
294 remote_name: String,
295 auth: GitAuth,
296}
297
298impl RemoteDefaultBranch {
299 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 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 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 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 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}