orcs 0.0.8

Microservices monorepo orchestration tool
Documentation
use crate::project::Project;
use crate::service::{get_all_services, Service};
use crate::Error;
use clap::ArgMatches;
use dep_graph::{DepGraph, Node};
use rayon::prelude::*;
use std::sync::Arc;
use tracing::{debug, info, instrument};

#[instrument(skip(subcommand))]
pub fn run(project_path: &str, subcommand: &ArgMatches) -> Result<(), Error> {
    info!("Command run for project {}", project_path);
    // Retrieve the project
    let project = Project::from_path(&project_path)?;

    // Retrieve stage name
    let stage_name = subcommand
        .value_of("stage")
        .expect("missing 'stage' argument");

    // Run all stages
    if stage_name == "all" {
        match subcommand.value_of("service") {
            // Run for all stages in a specific service
            Some(name) => {
                let service = Service::from_name(project.clone(), name)?;
                run_service(project, service)?;
            }
            // Run for all stages and services
            None => {
                run_all(project)?;
            }
        }
    } else if !project.stages.contains_key(stage_name) {
        return Err(Error::RunStageError(
            "stage does not exist for the project",
            stage_name.to_string(),
        ));
    } else {
        match subcommand.value_of("service") {
            // Run for a specific stage and service
            Some(name) => {
                let service = Service::from_name(project.clone(), name)?;
                run_service_stage(project, service, stage_name)?;
            }
            // Run for all services in a specific stage
            None => {
                run_stage(project, stage_name)?;
            }
        }
    }

    Ok(())
}

#[instrument(skip(project))]
fn run_all(project: Arc<Project>) -> Result<(), Error> {
    debug!("Run all stages/services in project {}", (*project).path);
    let services = Arc::new(get_all_services(project.clone())?);

    // Create a dependency graph between stages
    let stage_nodes: Vec<Node<String>> = (&project.stages)
        .iter()
        .filter_map(|(name, stage)| {
            // Stage that should be skipped when running all stages
            if stage.skip {
                return None;
            }
            let mut n = Node::new(name.clone());
            for dep_name in &stage.depends_on {
                n.add_dep(dep_name.clone());
            }
            Some(n)
        })
        .collect();

    // Iterate over all stages
    // TODO: Used nested parallel iterators
    DepGraph::new(&stage_nodes)
        .into_iter()
        .map(|stage_name| {
            // Create a dependency graph between services
            let service_nodes: Vec<Node<Arc<Service>>> = services
                .clone()
                .iter()
                .map(|(_, service)| {
                    let mut n = Node::new(service.clone());
                    for dep_name in service.depends_on(&stage_name) {
                        n.add_dep(services.get(&dep_name).unwrap().clone());
                    }
                    n
                })
                .collect();

            // Run for all services
            DepGraph::new(&service_nodes)
                .into_par_iter()
                .map(|service| (*service).run(&stage_name))
                .collect::<Result<Vec<()>, Error>>()?;

            Ok(())
        })
        .collect::<Result<Vec<()>, Error>>()?;

    Ok(())
}

#[instrument(skip(project, service))]
fn run_service(project: Arc<Project>, service: Arc<Service>) -> Result<(), Error> {
    debug!(
        "Run all stages in project {} for service {}",
        (*project).path,
        service.name
    );
    // Create a dependency graph between stages
    let stage_nodes: Vec<Node<String>> = (&project.stages)
        .iter()
        .filter_map(|(name, stage)| {
            // Stages that should be skipped when running all stages
            if stage.skip {
                return None;
            }
            let mut n = Node::new(name.clone());
            for dep_name in &stage.depends_on {
                n.add_dep(dep_name.clone());
            }
            Some(n)
        })
        .collect();

    DepGraph::new(&stage_nodes)
        .into_par_iter()
        .map(|stage_name| service.run(&stage_name))
        .collect::<Result<Vec<()>, Error>>()?;

    Ok(())
}

#[instrument(skip(project, stage_name))]
fn run_stage(project: Arc<Project>, stage_name: &str) -> Result<(), Error> {
    debug!(
        "Run all services in project {} for stage {}",
        (*project).path,
        stage_name
    );
    let services = get_all_services(project)?;

    // Create a dependency graph between services
    let service_nodes: Vec<Node<Arc<Service>>> = services
        .iter()
        .map(|(_, service)| {
            let mut n = Node::new(service.clone());
            for dep_name in service.depends_on(stage_name) {
                n.add_dep(services.get(&dep_name).unwrap().clone());
            }
            n
        })
        .collect();

    // Run for all services
    DepGraph::new(&service_nodes)
        .into_par_iter()
        .map(|service| {
            if !service.check_dep_stages(stage_name)? {
                println!("ono");
                return Err(Error::RunStageError(
                    "needs to run another stage before this one",
                    stage_name.to_string(),
                ));
            }

            (*service).run(stage_name)
        })
        .collect::<Result<Vec<()>, Error>>()?;

    Ok(())
}

#[instrument(skip(project, service, stage_name))]
fn run_service_stage(
    project: Arc<Project>,
    service: Arc<Service>,
    stage_name: &str,
) -> Result<(), Error> {
    debug!(
        "Run project {} service {} stage {}",
        (*project).path,
        service.name,
        stage_name
    );

    if !service.check_dep_stages(stage_name)? {
        return Err(Error::RunStageError(
            "needs to run another stage before this one",
            stage_name.to_string(),
        ));
    }

    match service.check_deps(stage_name)? {
        false => {
            return Err(Error::RunStageError(
                "needs to run another service before this one",
                stage_name.to_string(),
            ))
        }
        true => (),
    };
    service.run(stage_name)?;

    Ok(())
}