Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
17 changes: 12 additions & 5 deletions .github/workflows/auto-format.yml
Original file line number Diff line number Diff line change
Expand Up @@ -32,12 +32,19 @@ jobs:
token: ${{ secrets.PARENT_REPO_PAT || github.token }}

- name: Clone dependencies
env:
BRANCH: ${{ github.head_ref || github.ref_name }}
run: |
git clone --depth 1 https://github.com/harmony-labs/meta_plugin_protocol.git
git clone --depth 1 https://github.com/harmony-labs/meta_git_lib.git
git clone --depth 1 https://github.com/harmony-labs/meta_cli.git
git clone --depth 1 https://github.com/harmony-labs/meta_core.git
git clone --depth 1 https://github.com/harmony-labs/loop_lib.git
for repo in meta_plugin_protocol meta_git_lib meta_cli meta_core loop_lib; do
if git clone --depth 1 -b "$BRANCH" "https://github.com/harmony-labs/${repo}.git" 2>/dev/null; then
echo "Cloned ${repo}@${BRANCH}"
elif git clone --depth 1 "https://github.com/harmony-labs/${repo}.git"; then
echo "Branch ${BRANCH} not found for ${repo}; falling back to default branch"
else
echo "::error::Failed to clone ${repo}"
exit 1
fi
done

- name: Create workspace root
run: |
Expand Down
34 changes: 24 additions & 10 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -25,12 +25,19 @@ jobs:

- name: Clone dependencies
shell: bash
env:
BRANCH: ${{ github.head_ref || github.ref_name }}
run: |
git clone --depth 1 https://github.com/harmony-labs/meta_plugin_protocol.git
git clone --depth 1 https://github.com/harmony-labs/meta_git_lib.git
git clone --depth 1 https://github.com/harmony-labs/meta_cli.git
git clone --depth 1 https://github.com/harmony-labs/meta_core.git
git clone --depth 1 https://github.com/harmony-labs/loop_lib.git
for repo in meta_plugin_protocol meta_git_lib meta_cli meta_core loop_lib; do
if git clone --depth 1 -b "$BRANCH" "https://github.com/harmony-labs/${repo}.git" 2>/dev/null; then
echo "Cloned ${repo}@${BRANCH}"
elif git clone --depth 1 "https://github.com/harmony-labs/${repo}.git"; then
echo "Branch ${BRANCH} not found for ${repo}; falling back to default branch"
else
echo "::error::Failed to clone ${repo}"
exit 1
fi
done

- name: Create workspace root
shell: bash
Expand Down Expand Up @@ -59,12 +66,19 @@ jobs:
path: meta_git_cli

- name: Clone dependencies
env:
BRANCH: ${{ github.head_ref || github.ref_name }}
run: |
git clone --depth 1 https://github.com/harmony-labs/meta_plugin_protocol.git
git clone --depth 1 https://github.com/harmony-labs/meta_git_lib.git
git clone --depth 1 https://github.com/harmony-labs/meta_cli.git
git clone --depth 1 https://github.com/harmony-labs/meta_core.git
git clone --depth 1 https://github.com/harmony-labs/loop_lib.git
for repo in meta_plugin_protocol meta_git_lib meta_cli meta_core loop_lib; do
if git clone --depth 1 -b "$BRANCH" "https://github.com/harmony-labs/${repo}.git" 2>/dev/null; then
echo "Cloned ${repo}@${BRANCH}"
elif git clone --depth 1 "https://github.com/harmony-labs/${repo}.git"; then
echo "Branch ${BRANCH} not found for ${repo}; falling back to default branch"
else
echo "::error::Failed to clone ${repo}"
exit 1
fi
done

- name: Create workspace root
run: |
Expand Down
1 change: 1 addition & 0 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ path = "src/lib.rs"

[dependencies]
meta_plugin_protocol = { path = "../meta_plugin_protocol" }
meta_core = { path = "../meta_core" }
meta_git_lib = { path = "../meta_git_lib" }
meta_cli = { path = "../meta_cli", package = "meta" }
loop_lib = { path = "../loop_lib" }
Expand Down
6 changes: 3 additions & 3 deletions src/clone.rs
Original file line number Diff line number Diff line change
@@ -1,8 +1,8 @@
use crate::clone_queue::clone_with_queue;
use crate::clone_queue::CloneQueue;
use crate::clone_worker::clone_with_queue;
use console::style;
use indicatif::MultiProgress;
use meta_cli::config;
use meta_core::config;
use meta_git_lib::clone_queue::CloneQueue;
use meta_plugin_protocol::{CommandResult, PluginRequestOptions};
use std::process::Command;
use std::sync::Arc;
Expand Down
186 changes: 11 additions & 175 deletions src/clone_queue.rs → src/clone_worker.rs
Original file line number Diff line number Diff line change
@@ -1,181 +1,12 @@
use console::style;
use indicatif::{MultiProgress, ProgressBar, ProgressStyle};
use log::debug;
use meta_cli::config;
use std::collections::HashSet;
use std::path::{Path, PathBuf};
use meta_git_lib::clone_queue::{CloneQueue, CloneTask};
use std::process::Command;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::{Arc, Mutex};
use std::time::Duration;

/// A clone task representing a single repository to clone
#[derive(Debug, Clone)]
pub(crate) struct CloneTask {
/// Display name for progress output
pub name: String,
/// Git URL to clone from
pub url: String,
/// Target path to clone into
pub target_path: PathBuf,
/// Depth level (for display purposes)
pub depth_level: usize,
}

/// Thread-safe queue for managing clone tasks with dynamic discovery
pub(crate) struct CloneQueue {
/// Pending tasks to process
pending: Mutex<Vec<CloneTask>>,
/// Completed task paths (to avoid duplicates)
completed: Mutex<HashSet<PathBuf>>,
/// Failed task paths
failed: Mutex<HashSet<PathBuf>>,
/// Total tasks discovered (for progress display)
total_discovered: AtomicUsize,
/// Total tasks completed
total_completed: AtomicUsize,
/// Git depth argument (if any)
git_depth: Option<String>,
/// Max meta depth for recursion (None = unlimited)
meta_depth: Option<usize>,
}

impl CloneQueue {
pub fn new(git_depth: Option<String>, meta_depth: Option<usize>) -> Self {
Self {
pending: Mutex::new(Vec::new()),
completed: Mutex::new(HashSet::new()),
failed: Mutex::new(HashSet::new()),
total_discovered: AtomicUsize::new(0),
total_completed: AtomicUsize::new(0),
git_depth,
meta_depth,
}
}

/// Add a task to the queue if not already completed or pending
pub fn push(&self, task: CloneTask) -> bool {
let path = task.target_path.clone();

// Check if already completed
{
let completed = self.completed.lock().unwrap_or_else(|e| e.into_inner());
if completed.contains(&path) {
return false;
}
}

// Add to pending
{
let mut pending = self.pending.lock().unwrap_or_else(|e| e.into_inner());
// Check if already in pending
if pending.iter().any(|t| t.target_path == path) {
return false;
}
pending.push(task);
self.total_discovered.fetch_add(1, Ordering::SeqCst);
}

true
}

/// Add multiple tasks from a .meta file
pub fn push_from_meta(&self, base_dir: &Path, depth_level: usize) -> anyhow::Result<usize> {
// Check meta depth limit
if let Some(max_depth) = self.meta_depth {
if depth_level > max_depth {
return Ok(0);
}
}

let Some((meta_path, _format)) = config::find_meta_config_in(base_dir) else {
return Ok(0);
};

let (projects, _) = config::parse_meta_config(&meta_path)?;

let mut added = 0;
for project in projects {
let target_path = base_dir.join(&project.path);

// Skip if already exists
if target_path.exists() {
// But still check if it has a config file for nested discovery
if config::find_meta_config_in(&target_path).is_some() {
// Queue it for discovery even though it's already cloned
added += self.push_from_meta(&target_path, depth_level + 1)?;
}
continue;
}

// Skip projects without a repo URL (cannot clone)
let Some(url) = project.repo else {
continue;
};

let task = CloneTask {
name: project.name,
url,
target_path,
depth_level,
};

if self.push(task) {
added += 1;
}
}

Ok(added)
}

/// Take a single task from the queue (for worker threads)
fn take_one(&self) -> Option<CloneTask> {
let mut pending = self.pending.lock().unwrap_or_else(|e| e.into_inner());
pending.pop()
}

/// Check if queue is finished (no pending and no active workers)
fn is_finished(&self, active_workers: &AtomicUsize) -> bool {
let pending = self.pending.lock().unwrap_or_else(|e| e.into_inner());
pending.is_empty() && active_workers.load(Ordering::SeqCst) == 0
}

/// Drain all pending tasks (for dry-run display)
pub fn drain_all(&self) -> Vec<CloneTask> {
let mut pending = self.pending.lock().unwrap_or_else(|e| e.into_inner());
pending.drain(..).collect()
}

/// Get current counts for display
pub fn get_counts(&self) -> (usize, usize) {
(
self.total_completed.load(Ordering::SeqCst),
self.total_discovered.load(Ordering::SeqCst),
)
}

/// Mark a task as completed and check for nested .meta files
fn mark_completed(&self, task: &CloneTask) -> anyhow::Result<usize> {
self.total_completed.fetch_add(1, Ordering::SeqCst);

{
let mut completed = self.completed.lock().unwrap_or_else(|e| e.into_inner());
completed.insert(task.target_path.clone());
}

// Check for nested .meta file and add children to queue
self.push_from_meta(&task.target_path, task.depth_level + 1)
}

/// Mark a task as failed
fn mark_failed(&self, task: &CloneTask) {
self.total_completed.fetch_add(1, Ordering::SeqCst);

let mut failed = self.failed.lock().unwrap_or_else(|e| e.into_inner());
failed.insert(task.target_path.clone());
}
}

/// Clone repositories using a worker pool where each worker continuously pulls from the queue
pub(crate) fn clone_with_queue(
queue: Arc<CloneQueue>,
Expand Down Expand Up @@ -204,14 +35,16 @@ pub(crate) fn clone_with_queue(

std::thread::spawn(move || {
loop {
// Mark worker as active BEFORE taking a task to prevent
// a race where is_finished() sees pending=empty, active=0
// while a worker is between take_one() and starting work.
active.fetch_add(1, Ordering::SeqCst);

// Try to get a task
let task = queue.take_one();

match task {
Some(task) => {
// Mark worker as active
active.fetch_add(1, Ordering::SeqCst);

// Create progress bar for this task
let (completed, total) = queue.get_counts();
let pb = mp.add(ProgressBar::new_spinner());
Expand All @@ -233,7 +66,10 @@ pub(crate) fn clone_with_queue(
cvar.notify_all();
}
None => {
// No task available - check if we should terminate
// No task available - mark worker as inactive
active.fetch_sub(1, Ordering::SeqCst);

// Check if we should terminate
if queue.is_finished(&active) {
break;
}
Expand Down Expand Up @@ -282,7 +118,7 @@ fn clone_single_repo(task: &CloneTask, queue: &Arc<CloneQueue>, pb: &ProgressBar
// Build git clone command
let mut cmd = Command::new("git");
cmd.arg("clone").arg(&task.url).arg(&task.target_path);
if let Some(ref d) = queue.git_depth {
if let Some(d) = queue.git_depth() {
cmd.arg("--depth").arg(d);
}

Expand Down
4 changes: 2 additions & 2 deletions src/commands/worktree/create.rs
Original file line number Diff line number Diff line change
Expand Up @@ -315,14 +315,14 @@ pub(crate) fn handle_create(
/// 2. Resolves transitive dependencies via provides/depends_on from .meta.yaml
fn resolve_repos_with_dependencies(
meta_dir: &std::path::Path,
projects: &[meta_cli::config::ProjectInfo],
projects: &[meta_core::config::ProjectInfo],
repo_specs: &[meta_git_lib::worktree::RepoSpec],
worktree_name: &str,
branch_flag: Option<&str>,
verbose: bool,
) -> Result<Vec<(String, std::path::PathBuf, String)>> {
// Build dependency graph from projects
let project_deps: Vec<_> = projects.iter().map(|p| p.to_dependencies()).collect();
let project_deps: Vec<_> = projects.iter().map(|p| p.clone().into()).collect();
let graph = DependencyGraph::build(project_deps)?;

// Collect all repos to include (using HashSet for deduplication)
Expand Down
8 changes: 4 additions & 4 deletions src/commands/worktree/prune.rs
Original file line number Diff line number Diff line change
Expand Up @@ -34,15 +34,15 @@ fn check_repo_orphaned(
entry: &WorktreeStoreEntry,
config_cache: &mut std::collections::HashMap<
String,
Option<Vec<meta_cli::config::ProjectInfo>>,
Option<Vec<meta_core::config::ProjectInfo>>,
>,
) -> Option<String> {
let config = config_cache
.entry(entry.project.clone())
.or_insert_with(|| {
let project_path = Path::new(&entry.project);
meta_cli::config::find_meta_config_in(project_path).and_then(|(meta_path, _)| {
meta_cli::config::parse_meta_config(&meta_path)
meta_core::config::find_meta_config_in(project_path).and_then(|(meta_path, _)| {
meta_core::config::parse_meta_config(&meta_path)
.ok()
.map(|(projects, _)| projects)
})
Expand Down Expand Up @@ -93,7 +93,7 @@ pub(crate) fn handle_prune(
let mut to_remove: Vec<PruneEntry> = Vec::new();
let mut config_cache: std::collections::HashMap<
String,
Option<Vec<meta_cli::config::ProjectInfo>>,
Option<Vec<meta_core::config::ProjectInfo>>,
> = std::collections::HashMap::new();

for (path_key, entry) in &store.worktrees {
Expand Down
2 changes: 1 addition & 1 deletion src/commit.rs
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
use console::style;
use meta_cli::config;
use meta_core::config;
use meta_plugin_protocol::{CommandResult, PlannedCommand, PluginRequestOptions};
use std::process::Command;

Expand Down
2 changes: 1 addition & 1 deletion src/helpers.rs
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
use meta_cli::config;
use meta_core::config;
use std::path::Path;

/// Get project directories - uses passed-in list if non-empty, otherwise reads local .meta
Expand Down
Loading