From 4487409da4d3f3866b1bece00b0427182a67d0ab Mon Sep 17 00:00:00 2001 From: Przemyslaw Olszewski Date: Sun, 30 Aug 2026 19:25:39 +0200 Subject: [PATCH] fix: authorize native checkout credentials --- diff --git a/src/trigger.rs b/src/trigger.rs --- a/src/trigger.rs +++ b/src/trigger.rs @@ -1,280 +1,287 @@ -use std::sync::Arc; - -use syncode_control_runs::{ - Architecture, JobId, MatrixPolicy, OperatingSystem, Origin, Priority, QueuedJob, Requirements, - RunId, RunLog, Runs, RunsError, -}; -use syncode_workflow::{ - BooleanValue, Event, PositiveIntegerValue, StepKind, Value, VersionedPlan, WorkflowCompiler, - WorkflowDialect, WorkflowSource, expand, -}; -use syncode_workflow_github_actions::compiler::{GithubActionsCompileError, GithubActionsCompiler}; -use syncode_workflow_github_actions::expression::{ContextName, EvaluationContext}; -use syncode_workflow_github_actions::template::render; -use syncode_workflow_github_actions::workflow::{Workflow, parse}; -use thiserror::Error; - -#[derive(Debug, Error)] -pub enum TriggerError { - #[error(transparent)] - Compile(#[from] GithubActionsCompileError), - - #[error("the workflow does not lower into a plan: {0}")] - Plan(String), - - #[error("the compiled plan cannot be encoded: {0}")] - Encode(#[from] serde_json::Error), - - #[error("the matrix strategy cannot be evaluated by the control plane: {0}")] - Strategy(String), - - #[error("the runner requirement cannot be evaluated by the control plane: {0}")] - Requirement(String), - - #[error(transparent)] - Runs(#[from] RunsError), - - #[error("the workflow file is not text: {0}")] - NotText(#[from] std::str::Utf8Error), - - #[error("the workflow contains an invalid secret reference: {0}")] - SecretReference(String), - - #[error(transparent)] - Action(#[from] crate::actions::ActionResolutionError), -} - -/// What an event did to one workflow. A workflow that does not declare the -/// event is not an error and not a run; it simply has nothing to say about it. -#[derive(Debug, Eq, PartialEq)] -pub enum Triggered { - Runs(Vec), - NotForThisEvent, -} - -pub async fn trigger( - runs: &Runs, - source: &[u8], - event: &Event, - origin: &Origin, - actions: &A, -) -> Result { - trigger_with_event(runs, source, Some(event), origin, actions).await -} - -pub(crate) async fn trigger_resolved( - runs: &Runs, - workflow: Workflow, - origin: &Origin, - actions: &A, -) -> Result { - let hir = syncode_workflow_github_actions::compiler::lower::workflow(workflow)?; - queue_hir(runs, hir, origin, actions).await -} - -async fn trigger_with_event( - runs: &Runs, - source: &[u8], - event: Option<&Event>, - origin: &Origin, - actions: &A, -) -> Result { - let document = parse(std::str::from_utf8(source)?).map_err(GithubActionsCompileError::from)?; - let workflow = Workflow::from_node(&document).map_err(GithubActionsCompileError::from)?; - if event.is_some_and(|event| !workflow.triggers.fire_on(event)) { - return Ok(Triggered::NotForThisEvent); - } - - let hir = GithubActionsCompiler.compile(&WorkflowSource::new( - WorkflowDialect::GitHubActions, - source.to_vec(), - ))?; - queue_hir(runs, hir, origin, actions).await -} - -async fn queue_hir( - runs: &Runs, - hir: syncode_workflow::WorkflowHir< - syncode_workflow_github_actions::expression::ExpressionProgram, - >, - origin: &Origin, - actions: &A, -) -> Result { - let plans = - syncode_workflow::plans(hir).map_err(|error| TriggerError::Plan(error.to_string()))?; - - let mut queued = Vec::new(); - for plan in &plans { - for mut combination in expand(plan) { - actions.resolve(&mut combination, origin).await?; - let key = combination.job().key().as_ref().to_owned(); - let needs = combination - .job() - .needs() - .iter() - .map(|need| need.as_ref().to_owned()) - .collect(); - let matrix = matrix_policy(&combination)?; - let requirements = requirements(&combination)?; - let priority = priority(origin); - let secrets = crate::secret_references::collect(&combination) - .map_err(TriggerError::SecretReference)?; - let encoded = serde_json::to_vec(&VersionedPlan::new(combination))?; - queued.push( - QueuedJob::new(JobId::fresh(), key, needs, matrix, encoded) - .scheduled(priority, requirements) - .referencing(secrets), - ); - } - } - Ok(Triggered::Runs(vec![ - runs.queue_run(origin.clone(), queued).await?, - ])) -} - -fn priority(origin: &Origin) -> Priority { - match origin.event() { - "workflow_dispatch" | "manual" => Priority::High, - "schedule" => Priority::Low, - _ => Priority::Normal, - } -} - -fn requirements( - plan: &syncode_workflow::ExecutionPlan< - syncode_workflow_github_actions::expression::ExpressionProgram, - >, -) -> Result { - let mut context = EvaluationContext::default(); - context.values_mut().insert( - ContextName::Matrix, - Value::Object(Arc::new(plan.job().matrix().clone())), - ); - let mut labels: Vec = plan - .job() - .runner() - .labels() - .iter() - .map(|label| { - render(label, &context) - .map(String::from) - .map_err(|error| TriggerError::Requirement(error.to_string())) - }) - .collect::>()?; - if let Some(group) = plan.job().runner().group() { - labels.push( - render(group, &context) - .map(String::from) - .map_err(|error| TriggerError::Requirement(error.to_string()))?, - ); - } - let mut required = Requirements::new(labels.clone()); - for label in &labels { - let normalized = label.to_ascii_lowercase(); - if matches!(normalized.as_str(), "amd64" | "x64" | "x86_64") { - required = required.architecture(Architecture::Amd64); - } else if matches!(normalized.as_str(), "arm64" | "aarch64") { - required = required.architecture(Architecture::Arm64); - } else if normalized == "linux" || normalized.starts_with("ubuntu-") { - required = required.operating_system(OperatingSystem::Linux); - } else if normalized == "windows" || normalized.starts_with("windows-") { - required = required.operating_system(OperatingSystem::Windows); - } else if matches!(normalized.as_str(), "macos" | "macos-latest") { - required = required.operating_system(OperatingSystem::MacOs); - } - } - let mut images = Vec::new(); - if let Some(container) = plan.job().container() { - images.push( - render(container.image(), &context) - .map(String::from) - .map_err(|error| TriggerError::Requirement(error.to_string()))?, - ); - } - for service in plan.job().services() { - images.push( - render(service.container().image(), &context) - .map(String::from) - .map_err(|error| TriggerError::Requirement(error.to_string()))?, - ); - } - let actions = plan - .job() - .steps() - .iter() - .filter_map(|step| match step.kind() { - StepKind::Action(action) => Some(action), - StepKind::Shell(_) => None, - }) - .map(|action| action_requirement(action, &context)) - .collect::, _>>()?; - Ok(required.prefer(images, actions)) -} - -fn action_requirement( - action: &syncode_workflow::ActionStep< - syncode_workflow_github_actions::expression::ExpressionProgram, - >, - context: &EvaluationContext, -) -> Result { - if let Some(reference) = action.reference() { - return render(reference, context) - .map(String::from) - .map_err(|error| TriggerError::Requirement(error.to_string())); - } - match action.source().resolved_source() { - Some(syncode_workflow::ResolvedAction::Local(local)) => Ok(format!("./{}", local.path())), - Some(syncode_workflow::ResolvedAction::Remote(remote)) => { - Ok(remote.requested_reference().to_string()) - } - Some(syncode_workflow::ResolvedAction::Oci(oci)) => { - Ok(oci.requested_reference().to_string()) - } - None => Err(TriggerError::Requirement( - "action source is neither unresolved nor resolved".to_owned(), - )), - } -} - -fn matrix_policy( - plan: &syncode_workflow::ExecutionPlan< - syncode_workflow_github_actions::expression::ExpressionProgram, - >, -) -> Result { - let mut context = EvaluationContext::default(); - context.values_mut().insert( - ContextName::Matrix, - Value::Object(Arc::new(plan.job().matrix().clone())), - ); - let fail_fast = match plan.job().strategy().fail_fast() { - BooleanValue::Literal(value) => *value, - BooleanValue::Expression(expression) => expression - .evaluate_condition(&context) - .map_err(|error| TriggerError::Strategy(error.to_string()))?, - }; - let max_parallel = match plan.job().strategy().max_parallel() { - None => None, - Some(PositiveIntegerValue::Literal(value)) => Some(*value), - Some(PositiveIntegerValue::Expression(expression)) => { - let value = expression - .evaluate(&context) - .map_err(|error| TriggerError::Strategy(error.to_string()))?; - match value { - Value::Number(value) - if value.is_finite() - && value > 0.0 - && value.fract() == 0.0 - && value <= u64::MAX as f64 => - { - Some(value as u64) - } - value => { - return Err(TriggerError::Strategy(format!( - "max-parallel evaluated to {value:?}, not a positive integer" - ))); - } - } - } - }; - Ok(MatrixPolicy::new(fail_fast, max_parallel)) -} +use std::sync::Arc; + +use syncode_control_runs::{ + Architecture, JobId, MatrixPolicy, OperatingSystem, Origin, Priority, QueuedJob, Requirements, + RunId, RunLog, Runs, RunsError, +}; +use syncode_workflow::{ + BooleanValue, Event, PositiveIntegerValue, StepKind, Value, VersionedPlan, WorkflowCompiler, + WorkflowDialect, WorkflowSource, expand, +}; +use syncode_workflow_github_actions::compiler::{GithubActionsCompileError, GithubActionsCompiler}; +use syncode_workflow_github_actions::expression::{ContextName, EvaluationContext}; +use syncode_workflow_github_actions::template::render; +use syncode_workflow_github_actions::workflow::{Workflow, parse}; +use thiserror::Error; + +#[derive(Debug, Error)] +pub enum TriggerError { + #[error(transparent)] + Compile(#[from] GithubActionsCompileError), + + #[error("the workflow does not lower into a plan: {0}")] + Plan(String), + + #[error("the compiled plan cannot be encoded: {0}")] + Encode(#[from] serde_json::Error), + + #[error("the matrix strategy cannot be evaluated by the control plane: {0}")] + Strategy(String), + + #[error("the runner requirement cannot be evaluated by the control plane: {0}")] + Requirement(String), + + #[error(transparent)] + Runs(#[from] RunsError), + + #[error("the workflow file is not text: {0}")] + NotText(#[from] std::str::Utf8Error), + + #[error("the workflow contains an invalid secret reference: {0}")] + SecretReference(String), + + #[error(transparent)] + Action(#[from] crate::actions::ActionResolutionError), +} + +/// What an event did to one workflow. A workflow that does not declare the +/// event is not an error and not a run; it simply has nothing to say about it. +#[derive(Debug, Eq, PartialEq)] +pub enum Triggered { + Runs(Vec), + NotForThisEvent, +} + +pub async fn trigger( + runs: &Runs, + source: &[u8], + event: &Event, + origin: &Origin, + actions: &A, +) -> Result { + trigger_with_event(runs, source, Some(event), origin, actions).await +} + +pub(crate) async fn trigger_resolved( + runs: &Runs, + workflow: Workflow, + origin: &Origin, + actions: &A, +) -> Result { + let hir = syncode_workflow_github_actions::compiler::lower::workflow(workflow)?; + queue_hir(runs, hir, origin, actions).await +} + +async fn trigger_with_event( + runs: &Runs, + source: &[u8], + event: Option<&Event>, + origin: &Origin, + actions: &A, +) -> Result { + let document = parse(std::str::from_utf8(source)?).map_err(GithubActionsCompileError::from)?; + let workflow = Workflow::from_node(&document).map_err(GithubActionsCompileError::from)?; + if event.is_some_and(|event| !workflow.triggers.fire_on(event)) { + return Ok(Triggered::NotForThisEvent); + } + + let hir = GithubActionsCompiler.compile(&WorkflowSource::new( + WorkflowDialect::GitHubActions, + source.to_vec(), + ))?; + queue_hir(runs, hir, origin, actions).await +} + +async fn queue_hir( + runs: &Runs, + hir: syncode_workflow::WorkflowHir< + syncode_workflow_github_actions::expression::ExpressionProgram, + >, + origin: &Origin, + actions: &A, +) -> Result { + let plans = + syncode_workflow::plans(hir).map_err(|error| TriggerError::Plan(error.to_string()))?; + + let mut queued = Vec::new(); + for plan in &plans { + for mut combination in expand(plan) { + actions.resolve(&mut combination, origin).await?; + let key = combination.job().key().as_ref().to_owned(); + let needs = combination + .job() + .needs() + .iter() + .map(|need| need.as_ref().to_owned()) + .collect(); + let matrix = matrix_policy(&combination)?; + let requirements = requirements(&combination)?; + let priority = priority(origin); + let mut secrets = crate::secret_references::collect(&combination) + .map_err(TriggerError::SecretReference)?; + if uuid::Uuid::parse_str(origin.repository()).is_ok() + && !secrets + .iter() + .any(|name| name == syncode_control_node::REPOSITORY_TOKEN_SECRET) + { + secrets.push(syncode_control_node::REPOSITORY_TOKEN_SECRET.to_owned()); + } + let encoded = serde_json::to_vec(&VersionedPlan::new(combination))?; + queued.push( + QueuedJob::new(JobId::fresh(), key, needs, matrix, encoded) + .scheduled(priority, requirements) + .referencing(secrets), + ); + } + } + Ok(Triggered::Runs(vec![ + runs.queue_run(origin.clone(), queued).await?, + ])) +} + +fn priority(origin: &Origin) -> Priority { + match origin.event() { + "workflow_dispatch" | "manual" => Priority::High, + "schedule" => Priority::Low, + _ => Priority::Normal, + } +} + +fn requirements( + plan: &syncode_workflow::ExecutionPlan< + syncode_workflow_github_actions::expression::ExpressionProgram, + >, +) -> Result { + let mut context = EvaluationContext::default(); + context.values_mut().insert( + ContextName::Matrix, + Value::Object(Arc::new(plan.job().matrix().clone())), + ); + let mut labels: Vec = plan + .job() + .runner() + .labels() + .iter() + .map(|label| { + render(label, &context) + .map(String::from) + .map_err(|error| TriggerError::Requirement(error.to_string())) + }) + .collect::>()?; + if let Some(group) = plan.job().runner().group() { + labels.push( + render(group, &context) + .map(String::from) + .map_err(|error| TriggerError::Requirement(error.to_string()))?, + ); + } + let mut required = Requirements::new(labels.clone()); + for label in &labels { + let normalized = label.to_ascii_lowercase(); + if matches!(normalized.as_str(), "amd64" | "x64" | "x86_64") { + required = required.architecture(Architecture::Amd64); + } else if matches!(normalized.as_str(), "arm64" | "aarch64") { + required = required.architecture(Architecture::Arm64); + } else if normalized == "linux" || normalized.starts_with("ubuntu-") { + required = required.operating_system(OperatingSystem::Linux); + } else if normalized == "windows" || normalized.starts_with("windows-") { + required = required.operating_system(OperatingSystem::Windows); + } else if matches!(normalized.as_str(), "macos" | "macos-latest") { + required = required.operating_system(OperatingSystem::MacOs); + } + } + let mut images = Vec::new(); + if let Some(container) = plan.job().container() { + images.push( + render(container.image(), &context) + .map(String::from) + .map_err(|error| TriggerError::Requirement(error.to_string()))?, + ); + } + for service in plan.job().services() { + images.push( + render(service.container().image(), &context) + .map(String::from) + .map_err(|error| TriggerError::Requirement(error.to_string()))?, + ); + } + let actions = plan + .job() + .steps() + .iter() + .filter_map(|step| match step.kind() { + StepKind::Action(action) => Some(action), + StepKind::Shell(_) => None, + }) + .map(|action| action_requirement(action, &context)) + .collect::, _>>()?; + Ok(required.prefer(images, actions)) +} + +fn action_requirement( + action: &syncode_workflow::ActionStep< + syncode_workflow_github_actions::expression::ExpressionProgram, + >, + context: &EvaluationContext, +) -> Result { + if let Some(reference) = action.reference() { + return render(reference, context) + .map(String::from) + .map_err(|error| TriggerError::Requirement(error.to_string())); + } + match action.source().resolved_source() { + Some(syncode_workflow::ResolvedAction::Local(local)) => Ok(format!("./{}", local.path())), + Some(syncode_workflow::ResolvedAction::Remote(remote)) => { + Ok(remote.requested_reference().to_string()) + } + Some(syncode_workflow::ResolvedAction::Oci(oci)) => { + Ok(oci.requested_reference().to_string()) + } + None => Err(TriggerError::Requirement( + "action source is neither unresolved nor resolved".to_owned(), + )), + } +} + +fn matrix_policy( + plan: &syncode_workflow::ExecutionPlan< + syncode_workflow_github_actions::expression::ExpressionProgram, + >, +) -> Result { + let mut context = EvaluationContext::default(); + context.values_mut().insert( + ContextName::Matrix, + Value::Object(Arc::new(plan.job().matrix().clone())), + ); + let fail_fast = match plan.job().strategy().fail_fast() { + BooleanValue::Literal(value) => *value, + BooleanValue::Expression(expression) => expression + .evaluate_condition(&context) + .map_err(|error| TriggerError::Strategy(error.to_string()))?, + }; + let max_parallel = match plan.job().strategy().max_parallel() { + None => None, + Some(PositiveIntegerValue::Literal(value)) => Some(*value), + Some(PositiveIntegerValue::Expression(expression)) => { + let value = expression + .evaluate(&context) + .map_err(|error| TriggerError::Strategy(error.to_string()))?; + match value { + Value::Number(value) + if value.is_finite() + && value > 0.0 + && value.fract() == 0.0 + && value <= u64::MAX as f64 => + { + Some(value as u64) + } + value => { + return Err(TriggerError::Strategy(format!( + "max-parallel evaluated to {value:?}, not a positive integer" + ))); + } + } + } + }; + Ok(MatrixPolicy::new(fail_fast, max_parallel)) +} diff --git a/tests/trigger.rs b/tests/trigger.rs --- a/tests/trigger.rs +++ b/tests/trigger.rs @@ -1,413 +1,446 @@ -#![allow(clippy::expect_used, clippy::panic, clippy::unwrap_used)] - -#[path = "support/actions.rs"] -mod actions; - -use std::collections::BTreeMap; -use std::error::Error; - -use syncode_control::trigger::{Triggered, trigger}; -use syncode_control_runs::{Conclusion, Forgotten, NodeId, Origin, Runs}; -use syncode_workflow::{Event, EventKind, GitReference, PlanSchemaVersion, Value, VersionedPlan}; -use syncode_workflow_github_actions::expression::ExpressionProgram; - -use actions::FIXTURE_ACTIONS; - -type TestResult = Result>; - -fn push(branch: &str, paths: &[&str]) -> Event { - Event::new( - EventKind::Push, - GitReference::Branch(branch.to_owned()), - paths.iter().map(|path| (*path).to_owned()).collect(), - ) -} - -fn origin() -> Origin { - Origin::new( - "syncode/meta".to_owned(), - "a-commit".to_owned(), - "refs/heads/main".to_owned(), - "push".to_owned(), - ".gitea/workflows/ci.yml".to_owned(), - ) -} - -fn queued(triggered: Triggered) -> Vec { - match triggered { - Triggered::Runs(runs) => runs, - Triggered::NotForThisEvent => Vec::new(), - } -} - -const MATRIX: &str = r#" -name: CI -on: - push: - branches: [main] -jobs: - build: - runs-on: [self-hosted, linux] - strategy: - matrix: - rust: ["1.95.0", "nightly"] - steps: - - run: cargo +${{ matrix.rust }} test -"#; - -async fn assigned(runs: &Runs) -> TestResult>> { - let mut plans = Vec::new(); - while let Some(assignment) = runs.take_next(NodeId::fresh()).await? { - plans.push(serde_json::from_slice(assignment.plan())?); - } - Ok(plans) -} - -#[tokio::test] -async fn a_matrix_workflow_becomes_one_run_with_each_combination() -> TestResult { - let runs = Runs::restored(Forgotten::default()).await?; - - let triggered = queued( - trigger( - &runs, - MATRIX.as_bytes(), - &push("main", &[]), - &origin(), - &FIXTURE_ACTIONS, - ) - .await?, - ); - - assert_eq!(triggered.len(), 1, "one workflow event, one run"); - let plans = assigned(&runs).await?; - assert_eq!(plans.len(), 2); - - let mut versions: Vec = plans - .iter() - .map( - |plan| match plan.plan().job().strategy().matrix().property("rust") { - Some(Value::String(value)) => value.clone(), - other => panic!("expected a single value, got {other:?}"), - }, - ) - .collect(); - versions.sort(); - assert_eq!(versions, vec!["1.95.0".to_owned(), "nightly".to_owned()]); - - for plan in &plans { - assert_eq!(plan.schema(), PlanSchemaVersion::CURRENT); - assert_eq!(plan.plan().job().key().as_ref(), "build"); - } - - Ok(()) -} - -#[tokio::test] -async fn a_workflow_without_a_matrix_becomes_one_run() -> TestResult { - let runs = Runs::restored(Forgotten::default()).await?; - - let triggered = queued( - trigger( - &runs, - br#" -on: [push] -jobs: - build: - runs-on: ubuntu-latest - steps: - - run: echo one -"#, - &push("main", &[]), - &origin(), - &FIXTURE_ACTIONS, - ) - .await?, - ); - - assert_eq!(triggered.len(), 1); - assert_eq!(assigned(&runs).await?.len(), 1); - - Ok(()) -} - -#[tokio::test] -async fn a_broken_workflow_is_refused_when_the_run_is_triggered() -> TestResult { - let runs = Runs::restored(Forgotten::default()).await?; - - let error = trigger( - &runs, - br#" -on: [push] -jobs: - build: - runs-on: ubuntu-latest - steps: - - run: echo "${{ github. }}" -"#, - &push("main", &[]), - &origin(), - &FIXTURE_ACTIONS, - ) - .await - .expect_err("a workflow that does not compile must not produce a run"); - - assert!( - !error.to_string().is_empty(), - "the reason must reach whoever triggered it" - ); - assert!( - runs.take_next(NodeId::fresh()).await?.is_none(), - "nothing may be queued from a workflow that did not compile" - ); - - Ok(()) -} - -#[tokio::test] -async fn a_workflow_that_does_not_want_the_event_queues_nothing() -> TestResult { - let runs = Runs::restored(Forgotten::default()).await?; - - let triggered = trigger( - &runs, - MATRIX.as_bytes(), - &push("wip/x", &[]), - &origin(), - &FIXTURE_ACTIONS, - ) - .await?; - - assert_eq!(triggered, Triggered::NotForThisEvent); - assert!( - runs.take_next(NodeId::fresh()).await?.is_none(), - "an event the workflow does not declare must queue nothing" - ); - - Ok(()) -} - -#[tokio::test] -async fn every_run_is_numbered_and_states_where_it_came_from() -> TestResult { - let runs = Runs::restored(Forgotten::default()).await?; - - queued( - trigger( - &runs, - MATRIX.as_bytes(), - &push("main", &[]), - &origin(), - &FIXTURE_ACTIONS, - ) - .await?, - ); - - let mut numbers = Vec::new(); - while let Some(assignment) = runs.take_next(NodeId::fresh()).await? { - assert_eq!(assignment.origin(), &origin()); - numbers.push(assignment.number().get()); - } - - assert_eq!( - numbers, - vec![1, 1], - "jobs of one run carry the same run number" - ); - - Ok(()) -} - -#[tokio::test] -async fn a_job_waits_until_every_job_it_needs_has_succeeded() -> TestResult { - let runs = Runs::restored(Forgotten::default()).await?; - queued( - trigger( - &runs, - br#" -on: [push] -jobs: - build: - runs-on: ubuntu-latest - steps: - - run: echo build - publish: - needs: build - runs-on: ubuntu-latest - steps: - - run: echo publish -"#, - &push("main", &[]), - &origin(), - &FIXTURE_ACTIONS, - ) - .await?, - ); - let build_node = NodeId::fresh(); - let build = runs - .take_next(build_node) - .await? - .ok_or("build was not queued")?; - let plan: VersionedPlan = serde_json::from_slice(build.plan())?; - assert_eq!(plan.plan().job().key().as_ref(), "build"); - assert!(runs.take_next(NodeId::fresh()).await?.is_none()); - - runs.finished_job( - build.run(), - build.job(), - build_node, - build.fence(), - Conclusion::Success, - ) - .await?; - let publish = runs - .take_next(NodeId::fresh()) - .await? - .ok_or("publish did not become ready")?; - let plan: VersionedPlan = serde_json::from_slice(publish.plan())?; - assert_eq!(plan.plan().job().key().as_ref(), "publish"); - Ok(()) -} - -#[tokio::test] -async fn an_always_job_is_assigned_with_the_failed_need() -> TestResult { - let runs = Runs::restored(Forgotten::default()).await?; - queued( - trigger( - &runs, - br#" -on: [push] -jobs: - build: - runs-on: ubuntu-latest - steps: - - run: exit 1 - cleanup: - if: always() - needs: build - runs-on: ubuntu-latest - steps: - - run: echo cleanup -"#, - &push("main", &[]), - &origin(), - &FIXTURE_ACTIONS, - ) - .await?, - ); - let build_node = NodeId::fresh(); - let build = runs - .take_next(build_node) - .await? - .ok_or("build was not queued")?; - runs.finished_job_with_outputs( - build.run(), - build.job(), - build_node, - build.fence(), - Conclusion::Failure, - BTreeMap::from([("artifact".to_owned(), "bundle.tar".to_owned())]), - ) - .await?; - - let cleanup = runs - .take_next(NodeId::fresh()) - .await? - .ok_or("always job did not become ready")?; - assert_eq!(cleanup.needs().len(), 1); - assert_eq!(cleanup.needs()[0].key(), "build"); - assert_eq!(cleanup.needs()[0].conclusion(), Conclusion::Failure); - assert_eq!( - cleanup.needs()[0] - .outputs() - .get("artifact") - .map(String::as_str), - Some("bundle.tar") - ); - Ok(()) -} - -#[tokio::test] -async fn matrix_max_parallel_is_enforced_by_the_queue() -> TestResult { - let runs = Runs::restored(Forgotten::default()).await?; - queued( - trigger( - &runs, - br#" -on: [push] -jobs: - build: - runs-on: ubuntu-latest - strategy: - max-parallel: 1 - matrix: - shard: [1, 2, 3] - steps: - - run: echo shard -"#, - &push("main", &[]), - &origin(), - &FIXTURE_ACTIONS, - ) - .await?, - ); - let first_node = NodeId::fresh(); - let first = runs - .take_next(first_node) - .await? - .ok_or("first matrix job was not queued")?; - assert!(runs.take_next(NodeId::fresh()).await?.is_none()); - runs.finished_job( - first.run(), - first.job(), - first_node, - first.fence(), - Conclusion::Success, - ) - .await?; - assert!(runs.take_next(NodeId::fresh()).await?.is_some()); - Ok(()) -} - -#[tokio::test] -async fn matrix_fail_fast_stops_unassigned_siblings() -> TestResult { - let runs = Runs::restored(Forgotten::default()).await?; - let run = queued( - trigger( - &runs, - br#" -on: [push] -jobs: - build: - runs-on: ubuntu-latest - strategy: - max-parallel: 1 - matrix: - shard: [1, 2, 3] - steps: - - run: exit 1 -"#, - &push("main", &[]), - &origin(), - &FIXTURE_ACTIONS, - ) - .await?, - )[0]; - let first_node = NodeId::fresh(); - let first = runs - .take_next(first_node) - .await? - .ok_or("first matrix job was not queued")?; - runs.finished_job( - first.run(), - first.job(), - first_node, - first.fence(), - Conclusion::Failure, - ) - .await?; - - assert!(runs.take_next(NodeId::fresh()).await?.is_none()); - assert_eq!( - runs.state_of(run).await?, - syncode_control_runs::RunState::Finished(Conclusion::Failure) - ); - Ok(()) -} +#![allow(clippy::expect_used, clippy::panic, clippy::unwrap_used)] + +#[path = "support/actions.rs"] +mod actions; + +use std::collections::BTreeMap; +use std::error::Error; + +use syncode_control::trigger::{Triggered, trigger}; +use syncode_control_node::REPOSITORY_TOKEN_SECRET; +use syncode_control_runs::{Conclusion, Forgotten, NodeId, Origin, Runs}; +use syncode_workflow::{Event, EventKind, GitReference, PlanSchemaVersion, Value, VersionedPlan}; +use syncode_workflow_github_actions::expression::ExpressionProgram; + +use actions::FIXTURE_ACTIONS; + +type TestResult = Result>; + +fn push(branch: &str, paths: &[&str]) -> Event { + Event::new( + EventKind::Push, + GitReference::Branch(branch.to_owned()), + paths.iter().map(|path| (*path).to_owned()).collect(), + ) +} + +fn origin() -> Origin { + Origin::new( + "syncode/meta".to_owned(), + "a-commit".to_owned(), + "refs/heads/main".to_owned(), + "push".to_owned(), + ".gitea/workflows/ci.yml".to_owned(), + ) +} + +fn queued(triggered: Triggered) -> Vec { + match triggered { + Triggered::Runs(runs) => runs, + Triggered::NotForThisEvent => Vec::new(), + } +} + +const MATRIX: &str = r#" +name: CI +on: + push: + branches: [main] +jobs: + build: + runs-on: [self-hosted, linux] + strategy: + matrix: + rust: ["1.95.0", "nightly"] + steps: + - run: cargo +${{ matrix.rust }} test +"#; + +async fn assigned(runs: &Runs) -> TestResult>> { + let mut plans = Vec::new(); + while let Some(assignment) = runs.take_next(NodeId::fresh()).await? { + plans.push(serde_json::from_slice(assignment.plan())?); + } + Ok(plans) +} + +#[tokio::test] +async fn a_matrix_workflow_becomes_one_run_with_each_combination() -> TestResult { + let runs = Runs::restored(Forgotten::default()).await?; + + let triggered = queued( + trigger( + &runs, + MATRIX.as_bytes(), + &push("main", &[]), + &origin(), + &FIXTURE_ACTIONS, + ) + .await?, + ); + + assert_eq!(triggered.len(), 1, "one workflow event, one run"); + let plans = assigned(&runs).await?; + assert_eq!(plans.len(), 2); + + let mut versions: Vec = plans + .iter() + .map( + |plan| match plan.plan().job().strategy().matrix().property("rust") { + Some(Value::String(value)) => value.clone(), + other => panic!("expected a single value, got {other:?}"), + }, + ) + .collect(); + versions.sort(); + assert_eq!(versions, vec!["1.95.0".to_owned(), "nightly".to_owned()]); + + for plan in &plans { + assert_eq!(plan.schema(), PlanSchemaVersion::CURRENT); + assert_eq!(plan.plan().job().key().as_ref(), "build"); + } + + Ok(()) +} + +#[tokio::test] +async fn a_workflow_without_a_matrix_becomes_one_run() -> TestResult { + let runs = Runs::restored(Forgotten::default()).await?; + + let triggered = queued( + trigger( + &runs, + br#" +on: [push] +jobs: + build: + runs-on: ubuntu-latest + steps: + - run: echo one +"#, + &push("main", &[]), + &origin(), + &FIXTURE_ACTIONS, + ) + .await?, + ); + + assert_eq!(triggered.len(), 1); + assert_eq!(assigned(&runs).await?.len(), 1); + + Ok(()) +} + +#[tokio::test] +async fn a_native_job_authorizes_its_repository_token() -> TestResult { + let runs = Runs::restored(Forgotten::default()).await?; + let native_origin = origin().with_repository(uuid::Uuid::new_v4().to_string()); + + queued( + trigger( + &runs, + br#" +on: [push] +jobs: + build: + runs-on: ubuntu-latest + steps: + - uses: actions/checkout@v4 +"#, + &push("main", &[]), + &native_origin, + &FIXTURE_ACTIONS, + ) + .await?, + ); + + let assignment = runs + .take_next(NodeId::fresh()) + .await? + .ok_or("nothing was queued")?; + + assert_eq!(assignment.secrets(), [REPOSITORY_TOKEN_SECRET]); + Ok(()) +} + +#[tokio::test] +async fn a_broken_workflow_is_refused_when_the_run_is_triggered() -> TestResult { + let runs = Runs::restored(Forgotten::default()).await?; + + let error = trigger( + &runs, + br#" +on: [push] +jobs: + build: + runs-on: ubuntu-latest + steps: + - run: echo "${{ github. }}" +"#, + &push("main", &[]), + &origin(), + &FIXTURE_ACTIONS, + ) + .await + .expect_err("a workflow that does not compile must not produce a run"); + + assert!( + !error.to_string().is_empty(), + "the reason must reach whoever triggered it" + ); + assert!( + runs.take_next(NodeId::fresh()).await?.is_none(), + "nothing may be queued from a workflow that did not compile" + ); + + Ok(()) +} + +#[tokio::test] +async fn a_workflow_that_does_not_want_the_event_queues_nothing() -> TestResult { + let runs = Runs::restored(Forgotten::default()).await?; + + let triggered = trigger( + &runs, + MATRIX.as_bytes(), + &push("wip/x", &[]), + &origin(), + &FIXTURE_ACTIONS, + ) + .await?; + + assert_eq!(triggered, Triggered::NotForThisEvent); + assert!( + runs.take_next(NodeId::fresh()).await?.is_none(), + "an event the workflow does not declare must queue nothing" + ); + + Ok(()) +} + +#[tokio::test] +async fn every_run_is_numbered_and_states_where_it_came_from() -> TestResult { + let runs = Runs::restored(Forgotten::default()).await?; + + queued( + trigger( + &runs, + MATRIX.as_bytes(), + &push("main", &[]), + &origin(), + &FIXTURE_ACTIONS, + ) + .await?, + ); + + let mut numbers = Vec::new(); + while let Some(assignment) = runs.take_next(NodeId::fresh()).await? { + assert_eq!(assignment.origin(), &origin()); + numbers.push(assignment.number().get()); + } + + assert_eq!( + numbers, + vec![1, 1], + "jobs of one run carry the same run number" + ); + + Ok(()) +} + +#[tokio::test] +async fn a_job_waits_until_every_job_it_needs_has_succeeded() -> TestResult { + let runs = Runs::restored(Forgotten::default()).await?; + queued( + trigger( + &runs, + br#" +on: [push] +jobs: + build: + runs-on: ubuntu-latest + steps: + - run: echo build + publish: + needs: build + runs-on: ubuntu-latest + steps: + - run: echo publish +"#, + &push("main", &[]), + &origin(), + &FIXTURE_ACTIONS, + ) + .await?, + ); + let build_node = NodeId::fresh(); + let build = runs + .take_next(build_node) + .await? + .ok_or("build was not queued")?; + let plan: VersionedPlan = serde_json::from_slice(build.plan())?; + assert_eq!(plan.plan().job().key().as_ref(), "build"); + assert!(runs.take_next(NodeId::fresh()).await?.is_none()); + + runs.finished_job( + build.run(), + build.job(), + build_node, + build.fence(), + Conclusion::Success, + ) + .await?; + let publish = runs + .take_next(NodeId::fresh()) + .await? + .ok_or("publish did not become ready")?; + let plan: VersionedPlan = serde_json::from_slice(publish.plan())?; + assert_eq!(plan.plan().job().key().as_ref(), "publish"); + Ok(()) +} + +#[tokio::test] +async fn an_always_job_is_assigned_with_the_failed_need() -> TestResult { + let runs = Runs::restored(Forgotten::default()).await?; + queued( + trigger( + &runs, + br#" +on: [push] +jobs: + build: + runs-on: ubuntu-latest + steps: + - run: exit 1 + cleanup: + if: always() + needs: build + runs-on: ubuntu-latest + steps: + - run: echo cleanup +"#, + &push("main", &[]), + &origin(), + &FIXTURE_ACTIONS, + ) + .await?, + ); + let build_node = NodeId::fresh(); + let build = runs + .take_next(build_node) + .await? + .ok_or("build was not queued")?; + runs.finished_job_with_outputs( + build.run(), + build.job(), + build_node, + build.fence(), + Conclusion::Failure, + BTreeMap::from([("artifact".to_owned(), "bundle.tar".to_owned())]), + ) + .await?; + + let cleanup = runs + .take_next(NodeId::fresh()) + .await? + .ok_or("always job did not become ready")?; + assert_eq!(cleanup.needs().len(), 1); + assert_eq!(cleanup.needs()[0].key(), "build"); + assert_eq!(cleanup.needs()[0].conclusion(), Conclusion::Failure); + assert_eq!( + cleanup.needs()[0] + .outputs() + .get("artifact") + .map(String::as_str), + Some("bundle.tar") + ); + Ok(()) +} + +#[tokio::test] +async fn matrix_max_parallel_is_enforced_by_the_queue() -> TestResult { + let runs = Runs::restored(Forgotten::default()).await?; + queued( + trigger( + &runs, + br#" +on: [push] +jobs: + build: + runs-on: ubuntu-latest + strategy: + max-parallel: 1 + matrix: + shard: [1, 2, 3] + steps: + - run: echo shard +"#, + &push("main", &[]), + &origin(), + &FIXTURE_ACTIONS, + ) + .await?, + ); + let first_node = NodeId::fresh(); + let first = runs + .take_next(first_node) + .await? + .ok_or("first matrix job was not queued")?; + assert!(runs.take_next(NodeId::fresh()).await?.is_none()); + runs.finished_job( + first.run(), + first.job(), + first_node, + first.fence(), + Conclusion::Success, + ) + .await?; + assert!(runs.take_next(NodeId::fresh()).await?.is_some()); + Ok(()) +} + +#[tokio::test] +async fn matrix_fail_fast_stops_unassigned_siblings() -> TestResult { + let runs = Runs::restored(Forgotten::default()).await?; + let run = queued( + trigger( + &runs, + br#" +on: [push] +jobs: + build: + runs-on: ubuntu-latest + strategy: + max-parallel: 1 + matrix: + shard: [1, 2, 3] + steps: + - run: exit 1 +"#, + &push("main", &[]), + &origin(), + &FIXTURE_ACTIONS, + ) + .await?, + )[0]; + let first_node = NodeId::fresh(); + let first = runs + .take_next(first_node) + .await? + .ok_or("first matrix job was not queued")?; + runs.finished_job( + first.run(), + first.job(), + first_node, + first.fence(), + Conclusion::Failure, + ) + .await?; + + assert!(runs.take_next(NodeId::fresh()).await?.is_none()); + assert_eq!( + runs.state_of(run).await?, + syncode_control_runs::RunState::Finished(Conclusion::Failure) + ); + Ok(()) +} diff --git a/tests/webhook.rs b/tests/webhook.rs --- a/tests/webhook.rs +++ b/tests/webhook.rs @@ -1,748 +1,759 @@ -#![allow(clippy::expect_used, clippy::panic, clippy::unwrap_used)] - -#[path = "support/actions.rs"] -mod actions; -#[path = "support/repository.rs"] -mod repository; - -use std::error::Error; -use std::sync::Arc; -use std::time::Duration; - -use base64::Engine; -use base64::engine::general_purpose::STANDARD; -use hmac::{Hmac, KeyInit, Mac}; -use sha2::Sha256; -use syncode_control::repository::RepositoryContents; -use syncode_control::repository_grpc::NativeRepositoryContents; -use syncode_control::repository_sources::RepositorySources; -use syncode_control::webhook::{Intake, router}; -use syncode_control_node::identity_wire::identity_server::{Identity, IdentityServer}; -use syncode_control_node::identity_wire::{ - CheckCapabilityRequest, CheckCapabilityResponse, GetRepositoryCoordinatesRequest, - GetRepositoryCoordinatesResponse, IssueWorkflowRepositoryTokenRequest, - IssueWorkflowRepositoryTokenResponse, ResolveRepositoryRequest, ResolveRepositoryResponse, - ValidateSessionRequest, ValidateSessionResponse, -}; -use syncode_control_node::wire::{ - Capabilities, Capacity, ControlMessage, Enrol, Hello, NodeMessage, control_message, - node_message, -}; -use syncode_control_node::{ - ArtifactTokenAuthority, CapabilityAuthority, GeneratedServer, NodeSessionClient, - NodeSessionServer, ProjectionClient, -}; -use syncode_control_nodes::{Ephemeral, Nodes, Scope}; -use syncode_control_runs::{Forgotten, NodeId, Runs}; -use syncode_workflow::VersionedPlan; -use syncode_workflow_github_actions::expression::ExpressionProgram; -use tokio::io::{AsyncReadExt, AsyncWriteExt}; -use tokio::net::TcpListener; -use tokio::sync::mpsc; -use tokio_stream::StreamExt; -use tokio_stream::wrappers::{ReceiverStream, TcpListenerStream}; -use tonic::transport::Server; -use tonic::{Request, Response, Status, Streaming}; -use url::Url; - -use actions::FIXTURE_ACTIONS; - -type TestResult = Result>; - -struct ProjectionIdentity; - -#[tonic::async_trait] -impl Identity for ProjectionIdentity { - async fn issue_workflow_repository_token( - &self, - _request: Request, - ) -> Result, Status> { - Err(Status::unimplemented("issue_workflow_repository_token")) - } - - async fn validate_session( - &self, - _request: Request, - ) -> Result, Status> { - Err(Status::unimplemented("validate_session")) - } - - async fn check_capability( - &self, - _request: Request, - ) -> Result, Status> { - Err(Status::unimplemented("check_capability")) - } - - async fn resolve_repository( - &self, - _request: Request, - ) -> Result, Status> { - Err(Status::unimplemented("resolve_repository")) - } - - async fn get_repository_coordinates( - &self, - _request: Request, - ) -> Result, Status> { - Ok(Response::new(GetRepositoryCoordinatesResponse { - owner: "syncode".to_owned(), - name: "fixture".to_owned(), - })) - } -} - -const SECRET: &str = "the secret the hook was configured with"; -const COMMIT: &str = "9f2c1e4a7b3d5f6081a2c3d4e5f60718293a4b5c"; - -const WORKFLOW: &str = r#" -name: CI -on: - push: - branches: [main] - tags: ["v*"] - pull_request: - branches: [main] - paths: - - "src/**" - workflow_dispatch: - schedule: - - cron: "0 6 * * *" -jobs: - build: - runs-on: [self-hosted, linux] - steps: - - run: echo "the run came from a push" -"#; - -const OTHER_MANUAL_WORKFLOW: &str = r#" -name: Other manual workflow -on: - push: - workflow_dispatch: -jobs: - other: - runs-on: [self-hosted, linux] - steps: - - run: echo "the other workflow must not run" -"#; - -/// A forge that serves one workflow directory and has nothing in the other. -async fn forge() -> TestResult { - let listener = TcpListener::bind("127.0.0.1:0").await?; - let base = Url::parse(&format!("http://{}/", listener.local_addr()?))?; - let listing = format!( - r#"[{{"path":".gitea/workflows/ci.yml","type":"file","content":"{}"}},{{"path":".gitea/workflows/other.yml","type":"file","content":"{}"}}]"#, - STANDARD.encode(WORKFLOW), - STANDARD.encode(OTHER_MANUAL_WORKFLOW) - ); - - tokio::spawn(async move { - // Answers by what was asked for rather than in a fixed order: a push - // never asks for a file list, a pull request does. - let files = r#"[{"filename":"src/main.rs"}]"#.to_owned(); - loop { - let Ok((mut stream, _)) = listener.accept().await else { - return; - }; - let mut buffer = vec![0_u8; 4096]; - let read = stream.read(&mut buffer).await.unwrap_or(0); - let request = String::from_utf8_lossy(&buffer[..read]).to_string(); - let asked = request.lines().next().unwrap_or_default().to_owned(); - - let (status, body) = if asked.contains("/pulls/") { - (200, files.clone()) - } else if asked.contains(".gitea%2Fworkflows") || asked.contains(".gitea/workflows") { - (200, listing.clone()) - } else { - (404, String::from("{}")) - }; - - let response = format!( - "HTTP/1.1 {status} X\r\ncontent-type: application/json\r\ncontent-length: {}\r\nconnection: close\r\n\r\n{body}", - body.len() - ); - let _ = stream.write_all(response.as_bytes()).await; - let _ = stream.shutdown().await; - } - }); - - Ok(base) -} - -struct Harness { - runs: Runs, - nodes: Nodes, - session: String, - events: String, -} - -/// Both ends of the control plane over one set of registries, exactly as the -/// binary wires them. -async fn start() -> TestResult { - start_mode(false).await -} - -async fn start_mode(shadow: bool) -> TestResult { - let runs = Runs::restored(Forgotten::default()).await?; - let nodes = Nodes::restored(Ephemeral::default()).await?; - let projection = projection_client().await?; - - let sessions = TcpListener::bind("127.0.0.1:0").await?; - let session = format!("http://{}", sessions.local_addr()?); - let served = (runs.clone(), nodes.clone()); - tokio::spawn(async move { - let session = if shadow { - NodeSessionServer::shadow( - served.0, - served.1, - capability_authority(), - artifact_authority(), - "http://127.0.0.1:1/".parse().expect("artifact URL"), - projection, - ) - } else { - NodeSessionServer::new( - served.0, - served.1, - capability_authority(), - artifact_authority(), - "http://127.0.0.1:1/".parse().expect("artifact URL"), - projection, - ) - }; - let _ = Server::builder() - .add_service(GeneratedServer::new(session)) - .serve_with_incoming(TcpListenerStream::new(sessions)) - .await; - }); - - let repository_listener = TcpListener::bind("127.0.0.1:0").await?; - let repository_endpoint = format!("http://{}", repository_listener.local_addr()?); - tokio::spawn(async move { - let _ = Server::builder() - .add_service(repository::service()) - .serve_with_incoming(TcpListenerStream::new(repository_listener)) - .await; - }); - let native = NativeRepositoryContents::connect(repository_endpoint).await?; - - let deliveries = TcpListener::bind("127.0.0.1:0").await?; - let events = format!("http://{}", deliveries.local_addr()?); - let intake = Arc::new(Intake::new( - RepositorySources::new( - native, - RepositoryContents::new(forge().await?, "a repository token".to_owned()), - ), - runs.clone(), - FIXTURE_ACTIONS, - SECRET.to_owned(), - )); - tokio::spawn(async move { - let _ = axum::serve(deliveries, router(intake)).await; - }); - - Ok(Harness { - runs, - nodes, - session, - events, - }) -} - -fn capability_authority() -> CapabilityAuthority { - CapabilityAuthority::new("test-capability-key").expect("capability key") -} - -fn artifact_authority() -> ArtifactTokenAuthority { - ArtifactTokenAuthority::new("test-artifact-key").expect("artifact key") -} - -async fn projection_client() -> TestResult { - let listener = TcpListener::bind("127.0.0.1:0").await?; - let base = Url::parse(&format!("http://{}/", listener.local_addr()?))?; - tokio::spawn(async move { - let app = axum::Router::new().fallback(|| async { axum::http::StatusCode::NO_CONTENT }); - if let Err(error) = axum::serve(listener, app).await { - eprintln!("projection fixture failed: {error}"); - } - }); - let identity_listener = TcpListener::bind("127.0.0.1:0").await?; - let identity_endpoint = format!("http://{}", identity_listener.local_addr()?); - tokio::spawn(async move { - let _ = Server::builder() - .add_service(IdentityServer::new(ProjectionIdentity)) - .serve_with_incoming(TcpListenerStream::new(identity_listener)) - .await; - }); - Ok(ProjectionClient::new( - base, - String::new(), - identity_endpoint, - String::new(), - )?) -} - -async fn open_node( - harness: &Harness, -) -> TestResult<(mpsc::Sender, Streaming)> { - let (sender, receiver) = mpsc::channel(8); - let mut client = NodeSessionClient::connect(harness.session.clone()).await?; - let mut inbound = client - .open(ReceiverStream::new(receiver)) - .await? - .into_inner(); - let token = harness.nodes.issue_token(Scope::Instance).await?; - sender - .send(NodeMessage { - sequence: 1, - message_id: "node-message-1".to_owned(), - idempotency_key: "node-message-1".to_owned(), - body: Some(node_message::Body::Enrol(Enrol { - token: token.secret().expose().to_owned(), - })), - }) - .await?; - let message = next(&mut inbound).await?; - let Some(control_message::Body::Enrolled(enrolled)) = message.body else { - return Err("expected an enrolled node".into()); - }; - sender - .send(NodeMessage { - sequence: 2, - message_id: "node-message-2".to_owned(), - idempotency_key: "node-message-2".to_owned(), - body: Some(node_message::Body::Hello(Hello { - node: enrolled.node, - credential: enrolled.credential, - capabilities: Some(Capabilities { - architecture: "arm64".to_owned(), - operating_system: "linux".to_owned(), - container_runtime: "docker".to_owned(), - container_runtime_version: "28.6.1".to_owned(), - cores: 2, - memory_bytes: 8 * 1024 * 1024 * 1024, - labels: vec!["self-hosted".to_owned(), "linux".to_owned()], - }), - capacity: Some(Capacity { - build_volume_free_bytes: 60 * 1024 * 1024 * 1024, - layer_store_bytes: 12 * 1024 * 1024 * 1024, - cache_volume_present: true, - cache_volume_total_bytes: 100, - cache_volume_used_bytes: 40, - cache_volume_path: "/var/cache/syncode".to_owned(), - cached_images: Vec::new(), - cached_actions: Vec::new(), - }), - max_parallel: 1, - })), - }) - .await?; - let welcome = next(&mut inbound).await?; - if !matches!(welcome.body, Some(control_message::Body::Welcome(_))) { - return Err("expected a welcome".into()); - } - Ok((sender, inbound)) -} - -fn push_body() -> String { - format!( - r#"{{"message_id":"018f47e2-b2c4-7f19-8a6d-13ef76c89211","protocol_version":"1.0","schema_version":1,"source_node_id":"018f47e2-b2c4-7f19-8a6d-13ef76c89212","repository_id":"{}","message_type":"repository.ref.updated","sequence":2,"term":1,"payload":{{"request":{{"principal_id":"018f47e2-b2c4-7f19-8a6d-13ef76c89213"}},"changes":[{{"name_hex":"726566732f68656164732f6d61696e","old":{{"kind":"object","object_id":"1111111111111111111111111111111111111111"}},"new":{{"kind":"object","object_id":"{COMMIT}"}}}}]}}}}"#, - repository::REPOSITORY - ) -} - -fn tag_body() -> String { - format!( - r#"{{"message_id":"018f47e2-b2c4-7f19-8a6d-13ef76c89211","protocol_version":"1.0","schema_version":1,"source_node_id":"018f47e2-b2c4-7f19-8a6d-13ef76c89212","repository_id":"{}","message_type":"repository.ref.updated","sequence":2,"term":1,"payload":{{"request":{{"principal_id":"018f47e2-b2c4-7f19-8a6d-13ef76c89213"}},"changes":[{{"name_hex":"726566732f746167732f76302e352e30","old":null,"new":{{"kind":"object","object_id":"{}"}}}}]}}}}"#, - repository::REPOSITORY, - repository::TAG_OBJECT - ) -} - -fn legacy_push_body() -> String { - format!( - r#"{{"ref":"refs/heads/main","after":"{COMMIT}","repository":{{"full_name":"syncode/demo"}},"commits":[{{"modified":["src/main.rs"]}}]}}"# - ) -} - -fn signature(body: &str) -> String { - let mut mac = Hmac::::new_from_slice(SECRET.as_bytes()).expect("key"); - mac.update(body.as_bytes()); - mac.finalize() - .into_bytes() - .iter() - .map(|byte| format!("{byte:02x}")) - .collect() -} - -async fn deliver(harness: &Harness, event: &str, body: &str, signed: bool) -> TestResult { - let signature = if signed { - signature(body) - } else { - signature("something else entirely") - }; - let response = reqwest::Client::new() - .post(format!("{}/events", harness.events)) - .header("x-syncode-event", event) - .header("x-syncode-signature", signature) - .header("x-syncode-delivery", format!("{event}-delivery")) - .header( - "x-syncode-workflow", - if event == "repository.ref.updated" { - ".gitea/workflows/ci.yml" - } else { - "" - }, - ) - .body(body.to_owned()) - .send() - .await?; - Ok(response.status().as_u16()) -} - -async fn deliver_native_pull_request(harness: &Harness, body: &str) -> TestResult { - let response = reqwest::Client::new() - .post(format!("{}/events", harness.events)) - .header("x-syncode-event", "pull_request") - .header("x-syncode-signature", signature(body)) - .header("x-syncode-delivery", "native-pull-request-delivery") - .header("x-syncode-repository-id", repository::REPOSITORY) - .header( - "x-syncode-before", - "1111111111111111111111111111111111111111", - ) - .header("x-syncode-after", COMMIT) - .body(body.to_owned()) - .send() - .await?; - Ok(response.status().as_u16()) -} - -async fn next(stream: &mut Streaming) -> TestResult { - Ok(stream - .next() - .await - .ok_or("the control plane closed the stream")??) -} - -#[tokio::test] -async fn a_push_becomes_a_plan_the_node_is_handed() -> TestResult { - let harness = start().await?; - - let body = push_body(); - assert_eq!( - deliver(&harness, "repository.ref.updated", &body, true).await?, - 202 - ); - - let (sender, receiver) = mpsc::channel(8); - let mut client = NodeSessionClient::connect(harness.session.clone()).await?; - let mut inbound = client - .open(ReceiverStream::new(receiver)) - .await? - .into_inner(); - - let token = harness.nodes.issue_token(Scope::Instance).await?; - sender - .send(NodeMessage { - sequence: 1, - message_id: "node-message-1".to_owned(), - idempotency_key: "node-message-1".to_owned(), - body: Some(node_message::Body::Enrol(Enrol { - token: token.secret().expose().to_owned(), - })), - }) - .await?; - let message = next(&mut inbound).await?; - let Some(control_message::Body::Enrolled(enrolled)) = message.body else { - panic!("expected an identity, got {:?}", message.body); - }; - - sender - .send(NodeMessage { - sequence: 2, - message_id: "node-message-2".to_owned(), - idempotency_key: "node-message-2".to_owned(), - body: Some(node_message::Body::Hello(Hello { - node: enrolled.node, - credential: enrolled.credential, - capabilities: Some(Capabilities { - architecture: "arm64".to_owned(), - operating_system: "linux".to_owned(), - container_runtime: "docker".to_owned(), - container_runtime_version: "28.6.1".to_owned(), - cores: 2, - memory_bytes: 8 * 1024 * 1024 * 1024, - labels: vec!["self-hosted".to_owned(), "linux".to_owned()], - }), - capacity: Some(Capacity { - build_volume_free_bytes: 60 * 1024 * 1024 * 1024, - layer_store_bytes: 12 * 1024 * 1024 * 1024, - cache_volume_present: true, - cache_volume_total_bytes: 100, - cache_volume_used_bytes: 40, - cache_volume_path: "/var/cache/syncode".to_owned(), - cached_images: Vec::new(), - cached_actions: Vec::new(), - }), - max_parallel: 1, - })), - }) - .await?; - - let welcome = next(&mut inbound).await?; - assert!(matches!( - welcome.body, - Some(control_message::Body::Welcome(_)) - )); - - let assignment = next(&mut inbound).await?; - let Some(control_message::Body::Assignment(assignment)) = assignment.body else { - panic!("expected the plan the push produced, got {assignment:?}"); - }; - // The node is handed a plan, not the workflow file the push carried. - let plan: VersionedPlan = serde_json::from_slice(&assignment.plan)?; - assert_eq!(plan.schema(), syncode_workflow::PlanSchemaVersion::CURRENT); - - // A plan says nothing about what it is being built from, so the assignment - // states it: without this a job cannot check anything out. - let origin = assignment.origin.ok_or("the assignment stated no origin")?; - assert_eq!(origin.number, 1); - assert_eq!(origin.repository, "syncode/fixture"); - assert_eq!(origin.commit, COMMIT); - assert_eq!(origin.reference, "refs/heads/main"); - assert_eq!(origin.event, "push"); - assert_eq!( - assignment.secrets, - vec![syncode_control_node::REPOSITORY_TOKEN_SECRET] - ); - assert!(!assignment.secret_capability.is_empty()); - - Ok(()) -} - -#[tokio::test] -async fn an_annotated_tag_push_uses_the_target_commit_and_tag_reference() -> TestResult { - let harness = start().await?; - let body = tag_body(); - - assert_eq!( - deliver(&harness, "repository.ref.updated", &body, true).await?, - 202 - ); - let assignment = harness - .runs - .take_next(NodeId::fresh()) - .await? - .ok_or("the tag push produced no run")?; - - assert_eq!(assignment.origin().commit(), COMMIT); - assert_eq!(assignment.origin().reference(), "refs/tags/v0.5.0"); - assert_eq!(assignment.origin().event(), "push"); - Ok(()) -} - -#[tokio::test] -async fn a_delivery_nobody_signed_queues_nothing() -> TestResult { - let harness = start().await?; - - let body = push_body(); - assert_eq!( - deliver(&harness, "repository.ref.updated", &body, false).await?, - 401 - ); - assert!( - harness - .runs - .take_next(syncode_control_runs::NodeId::fresh()) - .await? - .is_none() - ); - - Ok(()) -} - -#[tokio::test] -async fn a_retried_delivery_does_not_duplicate_its_run() -> TestResult { - let harness = start().await?; - let body = push_body(); - assert_eq!( - deliver(&harness, "repository.ref.updated", &body, true).await?, - 202 - ); - assert_eq!( - deliver(&harness, "repository.ref.updated", &body, true).await?, - 202 - ); - - let first = harness - .runs - .take_next(NodeId::fresh()) - .await? - .ok_or("the delivery produced no run")?; - assert!(harness.runs.take_next(NodeId::fresh()).await?.is_none()); - assert_eq!(first.origin().event(), "push"); - Ok(()) -} - -#[tokio::test] -async fn an_event_the_workflow_does_not_declare_is_accepted_and_ignored() -> TestResult { - let harness = start().await?; - - let body = push_body(); - assert_eq!(deliver(&harness, "issues", &body, true).await?, 202); - assert!( - harness - .runs - .take_next(syncode_control_runs::NodeId::fresh()) - .await? - .is_none() - ); - - Ok(()) -} - -#[tokio::test] -async fn a_legacy_forge_push_is_accepted_without_a_duplicate_run() -> TestResult { - let harness = start().await?; - let body = legacy_push_body(); - assert_eq!(deliver(&harness, "push", &body, true).await?, 202); - assert!(harness.runs.take_next(NodeId::fresh()).await?.is_none()); - Ok(()) -} - -fn pull_request_body(action: &str) -> String { - format!( - r#"{{"action":"{action}","number":7,"repository":{{"full_name":"syncode/demo"}},"pull_request":{{"head":{{"sha":"{COMMIT}"}},"base":{{"ref":"main"}}}}}}"# - ) -} - -#[tokio::test] -async fn a_pull_request_becomes_a_run_using_the_files_it_touches() -> TestResult { - let harness = start().await?; - - let body = pull_request_body("opened"); - assert_eq!(deliver(&harness, "pull_request", &body, true).await?, 202); - - // The workflow filters on `paths: src/**`. The webhook body carries no file - // list, so unless the control plane asks the forge for one, this run does - // not exist and nothing says why. - let assignment = harness - .runs - .take_next(syncode_control_runs::NodeId::fresh()) - .await? - .ok_or("the pull request produced no run")?; - let plan: VersionedPlan = serde_json::from_slice(assignment.plan())?; - assert_eq!(plan.schema(), syncode_workflow::PlanSchemaVersion::CURRENT); - - // A pull request head sits on no branch this control plane knows, but the - // forge publishes it under the pull request, which is what a checkout can - // fetch. - assert_eq!(assignment.origin().reference(), "refs/pull/7/head"); - assert_eq!(assignment.origin().event(), "pull_request"); - Ok(()) -} - -#[tokio::test] -async fn a_repository_plane_pull_request_reads_content_over_grpc() -> TestResult { - let harness = start().await?; - - let body = pull_request_body("opened"); - assert_eq!(deliver_native_pull_request(&harness, &body).await?, 202); - - let assignment = harness - .runs - .take_next(syncode_control_runs::NodeId::fresh()) - .await? - .ok_or("the native pull request produced no run")?; - assert_eq!(assignment.origin().repository(), "syncode/demo"); - assert_eq!(assignment.origin().reference(), "refs/pull/7/head"); - Ok(()) -} - -#[tokio::test] -async fn closing_a_pull_request_starts_nothing() -> TestResult { - let harness = start().await?; - - let body = pull_request_body("closed"); - assert_eq!(deliver(&harness, "pull_request", &body, true).await?, 202); - assert!( - harness - .runs - .take_next(syncode_control_runs::NodeId::fresh()) - .await? - .is_none() - ); - Ok(()) -} - -#[tokio::test] -async fn manual_and_scheduled_deliveries_become_runs() -> TestResult { - for event in ["workflow_dispatch", "schedule"] { - let harness = start().await?; - let body = format!( - r#"{{"ref":"refs/heads/main","after":"{COMMIT}","workflow":".gitea/workflows/ci.yml","repository":{{"full_name":"syncode/demo"}}}}"# - ); - assert_eq!(deliver(&harness, event, &body, true).await?, 202); - let assignment = harness - .runs - .take_next(syncode_control_runs::NodeId::fresh()) - .await? - .ok_or("event produced no run")?; - assert_eq!(assignment.origin().event(), event); - assert_eq!(assignment.origin().reference(), "refs/heads/main"); - assert!( - harness - .runs - .take_next(syncode_control_runs::NodeId::fresh()) - .await? - .is_none(), - "a requested event started a workflow other than the selected file" - ); - } - Ok(()) -} - -#[tokio::test] -async fn a_manual_delivery_can_select_a_tag() -> TestResult { - let harness = start().await?; - let body = format!( - r#"{{"ref":"refs/tags/v0.4.0","after":"{COMMIT}","workflow":".gitea/workflows/ci.yml","repository":{{"full_name":"syncode/demo"}}}}"# - ); - - assert_eq!( - deliver(&harness, "workflow_dispatch", &body, true).await?, - 202 - ); - let assignment = harness - .runs - .take_next(NodeId::fresh()) - .await? - .ok_or("manual tag event produced no run")?; - assert_eq!(assignment.origin().reference(), "refs/tags/v0.4.0"); - Ok(()) -} - -#[tokio::test] -async fn shadow_mode_persists_the_run_but_never_assigns_it() -> TestResult { - let harness = start_mode(true).await?; - let body = push_body(); - assert_eq!( - deliver(&harness, "repository.ref.updated", &body, true).await?, - 202 - ); - let (_sender, mut inbound) = open_node(&harness).await?; - - assert!( - tokio::time::timeout(Duration::from_millis(100), inbound.next()) - .await - .is_err(), - "shadow mode must not assign the run" - ); - assert!( - harness.runs.take_next(NodeId::fresh()).await?.is_some(), - "the shadow run must still exist in the aggregate" - ); - Ok(()) -} +#![allow(clippy::expect_used, clippy::panic, clippy::unwrap_used)] + +#[path = "support/actions.rs"] +mod actions; +#[path = "support/repository.rs"] +mod repository; + +use std::error::Error; +use std::sync::Arc; +use std::time::Duration; + +use base64::Engine; +use base64::engine::general_purpose::STANDARD; +use hmac::{Hmac, KeyInit, Mac}; +use sha2::Sha256; +use syncode_control::repository::RepositoryContents; +use syncode_control::repository_grpc::NativeRepositoryContents; +use syncode_control::repository_sources::RepositorySources; +use syncode_control::webhook::{Intake, router}; +use syncode_control_node::identity_wire::identity_server::{Identity, IdentityServer}; +use syncode_control_node::identity_wire::{ + CheckCapabilityRequest, CheckCapabilityResponse, GetRepositoryCoordinatesRequest, + GetRepositoryCoordinatesResponse, IssueWorkflowRepositoryTokenRequest, + IssueWorkflowRepositoryTokenResponse, ResolveRepositoryRequest, ResolveRepositoryResponse, + ValidateSessionRequest, ValidateSessionResponse, +}; +use syncode_control_node::wire::{ + Capabilities, Capacity, ControlMessage, Enrol, Hello, NodeMessage, control_message, + node_message, +}; +use syncode_control_node::{ + ArtifactTokenAuthority, CapabilityAuthority, GeneratedServer, NodeSessionClient, + NodeSessionServer, ProjectionClient, +}; +use syncode_control_nodes::{Ephemeral, Nodes, Scope}; +use syncode_control_runs::{Fence, Forgotten, NodeId, Runs}; +use syncode_workflow::VersionedPlan; +use syncode_workflow_github_actions::expression::ExpressionProgram; +use tokio::io::{AsyncReadExt, AsyncWriteExt}; +use tokio::net::TcpListener; +use tokio::sync::mpsc; +use tokio_stream::StreamExt; +use tokio_stream::wrappers::{ReceiverStream, TcpListenerStream}; +use tonic::transport::Server; +use tonic::{Request, Response, Status, Streaming}; +use url::Url; + +use actions::FIXTURE_ACTIONS; + +type TestResult = Result>; + +struct ProjectionIdentity; + +#[tonic::async_trait] +impl Identity for ProjectionIdentity { + async fn issue_workflow_repository_token( + &self, + _request: Request, + ) -> Result, Status> { + Err(Status::unimplemented("issue_workflow_repository_token")) + } + + async fn validate_session( + &self, + _request: Request, + ) -> Result, Status> { + Err(Status::unimplemented("validate_session")) + } + + async fn check_capability( + &self, + _request: Request, + ) -> Result, Status> { + Err(Status::unimplemented("check_capability")) + } + + async fn resolve_repository( + &self, + _request: Request, + ) -> Result, Status> { + Err(Status::unimplemented("resolve_repository")) + } + + async fn get_repository_coordinates( + &self, + _request: Request, + ) -> Result, Status> { + Ok(Response::new(GetRepositoryCoordinatesResponse { + owner: "syncode".to_owned(), + name: "fixture".to_owned(), + })) + } +} + +const SECRET: &str = "the secret the hook was configured with"; +const COMMIT: &str = "9f2c1e4a7b3d5f6081a2c3d4e5f60718293a4b5c"; + +const WORKFLOW: &str = r#" +name: CI +on: + push: + branches: [main] + tags: ["v*"] + pull_request: + branches: [main] + paths: + - "src/**" + workflow_dispatch: + schedule: + - cron: "0 6 * * *" +jobs: + build: + runs-on: [self-hosted, linux] + steps: + - run: echo "the run came from a push" +"#; + +const OTHER_MANUAL_WORKFLOW: &str = r#" +name: Other manual workflow +on: + push: + workflow_dispatch: +jobs: + other: + runs-on: [self-hosted, linux] + steps: + - run: echo "the other workflow must not run" +"#; + +/// A forge that serves one workflow directory and has nothing in the other. +async fn forge() -> TestResult { + let listener = TcpListener::bind("127.0.0.1:0").await?; + let base = Url::parse(&format!("http://{}/", listener.local_addr()?))?; + let listing = format!( + r#"[{{"path":".gitea/workflows/ci.yml","type":"file","content":"{}"}},{{"path":".gitea/workflows/other.yml","type":"file","content":"{}"}}]"#, + STANDARD.encode(WORKFLOW), + STANDARD.encode(OTHER_MANUAL_WORKFLOW) + ); + + tokio::spawn(async move { + // Answers by what was asked for rather than in a fixed order: a push + // never asks for a file list, a pull request does. + let files = r#"[{"filename":"src/main.rs"}]"#.to_owned(); + loop { + let Ok((mut stream, _)) = listener.accept().await else { + return; + }; + let mut buffer = vec![0_u8; 4096]; + let read = stream.read(&mut buffer).await.unwrap_or(0); + let request = String::from_utf8_lossy(&buffer[..read]).to_string(); + let asked = request.lines().next().unwrap_or_default().to_owned(); + + let (status, body) = if asked.contains("/pulls/") { + (200, files.clone()) + } else if asked.contains(".gitea%2Fworkflows") || asked.contains(".gitea/workflows") { + (200, listing.clone()) + } else { + (404, String::from("{}")) + }; + + let response = format!( + "HTTP/1.1 {status} X\r\ncontent-type: application/json\r\ncontent-length: {}\r\nconnection: close\r\n\r\n{body}", + body.len() + ); + let _ = stream.write_all(response.as_bytes()).await; + let _ = stream.shutdown().await; + } + }); + + Ok(base) +} + +struct Harness { + runs: Runs, + nodes: Nodes, + session: String, + events: String, +} + +/// Both ends of the control plane over one set of registries, exactly as the +/// binary wires them. +async fn start() -> TestResult { + start_mode(false).await +} + +async fn start_mode(shadow: bool) -> TestResult { + let runs = Runs::restored(Forgotten::default()).await?; + let nodes = Nodes::restored(Ephemeral::default()).await?; + let projection = projection_client().await?; + + let sessions = TcpListener::bind("127.0.0.1:0").await?; + let session = format!("http://{}", sessions.local_addr()?); + let served = (runs.clone(), nodes.clone()); + tokio::spawn(async move { + let session = if shadow { + NodeSessionServer::shadow( + served.0, + served.1, + capability_authority(), + artifact_authority(), + "http://127.0.0.1:1/".parse().expect("artifact URL"), + projection, + ) + } else { + NodeSessionServer::new( + served.0, + served.1, + capability_authority(), + artifact_authority(), + "http://127.0.0.1:1/".parse().expect("artifact URL"), + projection, + ) + }; + let _ = Server::builder() + .add_service(GeneratedServer::new(session)) + .serve_with_incoming(TcpListenerStream::new(sessions)) + .await; + }); + + let repository_listener = TcpListener::bind("127.0.0.1:0").await?; + let repository_endpoint = format!("http://{}", repository_listener.local_addr()?); + tokio::spawn(async move { + let _ = Server::builder() + .add_service(repository::service()) + .serve_with_incoming(TcpListenerStream::new(repository_listener)) + .await; + }); + let native = NativeRepositoryContents::connect(repository_endpoint).await?; + + let deliveries = TcpListener::bind("127.0.0.1:0").await?; + let events = format!("http://{}", deliveries.local_addr()?); + let intake = Arc::new(Intake::new( + RepositorySources::new( + native, + RepositoryContents::new(forge().await?, "a repository token".to_owned()), + ), + runs.clone(), + FIXTURE_ACTIONS, + SECRET.to_owned(), + )); + tokio::spawn(async move { + let _ = axum::serve(deliveries, router(intake)).await; + }); + + Ok(Harness { + runs, + nodes, + session, + events, + }) +} + +fn capability_authority() -> CapabilityAuthority { + CapabilityAuthority::new("test-capability-key").expect("capability key") +} + +fn artifact_authority() -> ArtifactTokenAuthority { + ArtifactTokenAuthority::new("test-artifact-key").expect("artifact key") +} + +async fn projection_client() -> TestResult { + let listener = TcpListener::bind("127.0.0.1:0").await?; + let base = Url::parse(&format!("http://{}/", listener.local_addr()?))?; + tokio::spawn(async move { + let app = axum::Router::new().fallback(|| async { axum::http::StatusCode::NO_CONTENT }); + if let Err(error) = axum::serve(listener, app).await { + eprintln!("projection fixture failed: {error}"); + } + }); + let identity_listener = TcpListener::bind("127.0.0.1:0").await?; + let identity_endpoint = format!("http://{}", identity_listener.local_addr()?); + tokio::spawn(async move { + let _ = Server::builder() + .add_service(IdentityServer::new(ProjectionIdentity)) + .serve_with_incoming(TcpListenerStream::new(identity_listener)) + .await; + }); + Ok(ProjectionClient::new( + base, + String::new(), + identity_endpoint, + String::new(), + )?) +} + +async fn open_node( + harness: &Harness, +) -> TestResult<(mpsc::Sender, Streaming)> { + let (sender, receiver) = mpsc::channel(8); + let mut client = NodeSessionClient::connect(harness.session.clone()).await?; + let mut inbound = client + .open(ReceiverStream::new(receiver)) + .await? + .into_inner(); + let token = harness.nodes.issue_token(Scope::Instance).await?; + sender + .send(NodeMessage { + sequence: 1, + message_id: "node-message-1".to_owned(), + idempotency_key: "node-message-1".to_owned(), + body: Some(node_message::Body::Enrol(Enrol { + token: token.secret().expose().to_owned(), + })), + }) + .await?; + let message = next(&mut inbound).await?; + let Some(control_message::Body::Enrolled(enrolled)) = message.body else { + return Err("expected an enrolled node".into()); + }; + sender + .send(NodeMessage { + sequence: 2, + message_id: "node-message-2".to_owned(), + idempotency_key: "node-message-2".to_owned(), + body: Some(node_message::Body::Hello(Hello { + node: enrolled.node, + credential: enrolled.credential, + capabilities: Some(Capabilities { + architecture: "arm64".to_owned(), + operating_system: "linux".to_owned(), + container_runtime: "docker".to_owned(), + container_runtime_version: "28.6.1".to_owned(), + cores: 2, + memory_bytes: 8 * 1024 * 1024 * 1024, + labels: vec!["self-hosted".to_owned(), "linux".to_owned()], + }), + capacity: Some(Capacity { + build_volume_free_bytes: 60 * 1024 * 1024 * 1024, + layer_store_bytes: 12 * 1024 * 1024 * 1024, + cache_volume_present: true, + cache_volume_total_bytes: 100, + cache_volume_used_bytes: 40, + cache_volume_path: "/var/cache/syncode".to_owned(), + cached_images: Vec::new(), + cached_actions: Vec::new(), + }), + max_parallel: 1, + })), + }) + .await?; + let welcome = next(&mut inbound).await?; + if !matches!(welcome.body, Some(control_message::Body::Welcome(_))) { + return Err("expected a welcome".into()); + } + Ok((sender, inbound)) +} + +fn push_body() -> String { + format!( + r#"{{"message_id":"018f47e2-b2c4-7f19-8a6d-13ef76c89211","protocol_version":"1.0","schema_version":1,"source_node_id":"018f47e2-b2c4-7f19-8a6d-13ef76c89212","repository_id":"{}","message_type":"repository.ref.updated","sequence":2,"term":1,"payload":{{"request":{{"principal_id":"018f47e2-b2c4-7f19-8a6d-13ef76c89213"}},"changes":[{{"name_hex":"726566732f68656164732f6d61696e","old":{{"kind":"object","object_id":"1111111111111111111111111111111111111111"}},"new":{{"kind":"object","object_id":"{COMMIT}"}}}}]}}}}"#, + repository::REPOSITORY + ) +} + +fn tag_body() -> String { + format!( + r#"{{"message_id":"018f47e2-b2c4-7f19-8a6d-13ef76c89211","protocol_version":"1.0","schema_version":1,"source_node_id":"018f47e2-b2c4-7f19-8a6d-13ef76c89212","repository_id":"{}","message_type":"repository.ref.updated","sequence":2,"term":1,"payload":{{"request":{{"principal_id":"018f47e2-b2c4-7f19-8a6d-13ef76c89213"}},"changes":[{{"name_hex":"726566732f746167732f76302e352e30","old":null,"new":{{"kind":"object","object_id":"{}"}}}}]}}}}"#, + repository::REPOSITORY, + repository::TAG_OBJECT + ) +} + +fn legacy_push_body() -> String { + format!( + r#"{{"ref":"refs/heads/main","after":"{COMMIT}","repository":{{"full_name":"syncode/demo"}},"commits":[{{"modified":["src/main.rs"]}}]}}"# + ) +} + +fn signature(body: &str) -> String { + let mut mac = Hmac::::new_from_slice(SECRET.as_bytes()).expect("key"); + mac.update(body.as_bytes()); + mac.finalize() + .into_bytes() + .iter() + .map(|byte| format!("{byte:02x}")) + .collect() +} + +async fn deliver(harness: &Harness, event: &str, body: &str, signed: bool) -> TestResult { + let signature = if signed { + signature(body) + } else { + signature("something else entirely") + }; + let response = reqwest::Client::new() + .post(format!("{}/events", harness.events)) + .header("x-syncode-event", event) + .header("x-syncode-signature", signature) + .header("x-syncode-delivery", format!("{event}-delivery")) + .header( + "x-syncode-workflow", + if event == "repository.ref.updated" { + ".gitea/workflows/ci.yml" + } else { + "" + }, + ) + .body(body.to_owned()) + .send() + .await?; + Ok(response.status().as_u16()) +} + +async fn deliver_native_pull_request(harness: &Harness, body: &str) -> TestResult { + let response = reqwest::Client::new() + .post(format!("{}/events", harness.events)) + .header("x-syncode-event", "pull_request") + .header("x-syncode-signature", signature(body)) + .header("x-syncode-delivery", "native-pull-request-delivery") + .header("x-syncode-repository-id", repository::REPOSITORY) + .header( + "x-syncode-before", + "1111111111111111111111111111111111111111", + ) + .header("x-syncode-after", COMMIT) + .body(body.to_owned()) + .send() + .await?; + Ok(response.status().as_u16()) +} + +async fn next(stream: &mut Streaming) -> TestResult { + Ok(stream + .next() + .await + .ok_or("the control plane closed the stream")??) +} + +#[tokio::test] +async fn a_push_becomes_a_plan_the_node_is_handed() -> TestResult { + let harness = start().await?; + + let body = push_body(); + assert_eq!( + deliver(&harness, "repository.ref.updated", &body, true).await?, + 202 + ); + + let (sender, receiver) = mpsc::channel(8); + let mut client = NodeSessionClient::connect(harness.session.clone()).await?; + let mut inbound = client + .open(ReceiverStream::new(receiver)) + .await? + .into_inner(); + + let token = harness.nodes.issue_token(Scope::Instance).await?; + sender + .send(NodeMessage { + sequence: 1, + message_id: "node-message-1".to_owned(), + idempotency_key: "node-message-1".to_owned(), + body: Some(node_message::Body::Enrol(Enrol { + token: token.secret().expose().to_owned(), + })), + }) + .await?; + let message = next(&mut inbound).await?; + let Some(control_message::Body::Enrolled(enrolled)) = message.body else { + panic!("expected an identity, got {:?}", message.body); + }; + let node = enrolled.node.parse()?; + + sender + .send(NodeMessage { + sequence: 2, + message_id: "node-message-2".to_owned(), + idempotency_key: "node-message-2".to_owned(), + body: Some(node_message::Body::Hello(Hello { + node: enrolled.node, + credential: enrolled.credential, + capabilities: Some(Capabilities { + architecture: "arm64".to_owned(), + operating_system: "linux".to_owned(), + container_runtime: "docker".to_owned(), + container_runtime_version: "28.6.1".to_owned(), + cores: 2, + memory_bytes: 8 * 1024 * 1024 * 1024, + labels: vec!["self-hosted".to_owned(), "linux".to_owned()], + }), + capacity: Some(Capacity { + build_volume_free_bytes: 60 * 1024 * 1024 * 1024, + layer_store_bytes: 12 * 1024 * 1024 * 1024, + cache_volume_present: true, + cache_volume_total_bytes: 100, + cache_volume_used_bytes: 40, + cache_volume_path: "/var/cache/syncode".to_owned(), + cached_images: Vec::new(), + cached_actions: Vec::new(), + }), + max_parallel: 1, + })), + }) + .await?; + + let welcome = next(&mut inbound).await?; + assert!(matches!( + welcome.body, + Some(control_message::Body::Welcome(_)) + )); + + let assignment = next(&mut inbound).await?; + let Some(control_message::Body::Assignment(assignment)) = assignment.body else { + panic!("expected the plan the push produced, got {assignment:?}"); + }; + // The node is handed a plan, not the workflow file the push carried. + let plan: VersionedPlan = serde_json::from_slice(&assignment.plan)?; + assert_eq!(plan.schema(), syncode_workflow::PlanSchemaVersion::CURRENT); + + // A plan says nothing about what it is being built from, so the assignment + // states it: without this a job cannot check anything out. + let origin = assignment.origin.ok_or("the assignment stated no origin")?; + assert_eq!(origin.number, 1); + assert_eq!(origin.repository, "syncode/fixture"); + assert_eq!(origin.commit, COMMIT); + assert_eq!(origin.reference, "refs/heads/main"); + assert_eq!(origin.event, "push"); + assert_eq!( + assignment.secrets, + vec![syncode_control_node::REPOSITORY_TOKEN_SECRET] + ); + assert!(!assignment.secret_capability.is_empty()); + harness + .runs + .authorize_secret( + assignment.run.parse()?, + assignment.job.parse()?, + node, + Fence::from(assignment.fence), + syncode_control_node::REPOSITORY_TOKEN_SECRET, + ) + .await?; + + Ok(()) +} + +#[tokio::test] +async fn an_annotated_tag_push_uses_the_target_commit_and_tag_reference() -> TestResult { + let harness = start().await?; + let body = tag_body(); + + assert_eq!( + deliver(&harness, "repository.ref.updated", &body, true).await?, + 202 + ); + let assignment = harness + .runs + .take_next(NodeId::fresh()) + .await? + .ok_or("the tag push produced no run")?; + + assert_eq!(assignment.origin().commit(), COMMIT); + assert_eq!(assignment.origin().reference(), "refs/tags/v0.5.0"); + assert_eq!(assignment.origin().event(), "push"); + Ok(()) +} + +#[tokio::test] +async fn a_delivery_nobody_signed_queues_nothing() -> TestResult { + let harness = start().await?; + + let body = push_body(); + assert_eq!( + deliver(&harness, "repository.ref.updated", &body, false).await?, + 401 + ); + assert!( + harness + .runs + .take_next(syncode_control_runs::NodeId::fresh()) + .await? + .is_none() + ); + + Ok(()) +} + +#[tokio::test] +async fn a_retried_delivery_does_not_duplicate_its_run() -> TestResult { + let harness = start().await?; + let body = push_body(); + assert_eq!( + deliver(&harness, "repository.ref.updated", &body, true).await?, + 202 + ); + assert_eq!( + deliver(&harness, "repository.ref.updated", &body, true).await?, + 202 + ); + + let first = harness + .runs + .take_next(NodeId::fresh()) + .await? + .ok_or("the delivery produced no run")?; + assert!(harness.runs.take_next(NodeId::fresh()).await?.is_none()); + assert_eq!(first.origin().event(), "push"); + Ok(()) +} + +#[tokio::test] +async fn an_event_the_workflow_does_not_declare_is_accepted_and_ignored() -> TestResult { + let harness = start().await?; + + let body = push_body(); + assert_eq!(deliver(&harness, "issues", &body, true).await?, 202); + assert!( + harness + .runs + .take_next(syncode_control_runs::NodeId::fresh()) + .await? + .is_none() + ); + + Ok(()) +} + +#[tokio::test] +async fn a_legacy_forge_push_is_accepted_without_a_duplicate_run() -> TestResult { + let harness = start().await?; + let body = legacy_push_body(); + assert_eq!(deliver(&harness, "push", &body, true).await?, 202); + assert!(harness.runs.take_next(NodeId::fresh()).await?.is_none()); + Ok(()) +} + +fn pull_request_body(action: &str) -> String { + format!( + r#"{{"action":"{action}","number":7,"repository":{{"full_name":"syncode/demo"}},"pull_request":{{"head":{{"sha":"{COMMIT}"}},"base":{{"ref":"main"}}}}}}"# + ) +} + +#[tokio::test] +async fn a_pull_request_becomes_a_run_using_the_files_it_touches() -> TestResult { + let harness = start().await?; + + let body = pull_request_body("opened"); + assert_eq!(deliver(&harness, "pull_request", &body, true).await?, 202); + + // The workflow filters on `paths: src/**`. The webhook body carries no file + // list, so unless the control plane asks the forge for one, this run does + // not exist and nothing says why. + let assignment = harness + .runs + .take_next(syncode_control_runs::NodeId::fresh()) + .await? + .ok_or("the pull request produced no run")?; + let plan: VersionedPlan = serde_json::from_slice(assignment.plan())?; + assert_eq!(plan.schema(), syncode_workflow::PlanSchemaVersion::CURRENT); + + // A pull request head sits on no branch this control plane knows, but the + // forge publishes it under the pull request, which is what a checkout can + // fetch. + assert_eq!(assignment.origin().reference(), "refs/pull/7/head"); + assert_eq!(assignment.origin().event(), "pull_request"); + Ok(()) +} + +#[tokio::test] +async fn a_repository_plane_pull_request_reads_content_over_grpc() -> TestResult { + let harness = start().await?; + + let body = pull_request_body("opened"); + assert_eq!(deliver_native_pull_request(&harness, &body).await?, 202); + + let assignment = harness + .runs + .take_next(syncode_control_runs::NodeId::fresh()) + .await? + .ok_or("the native pull request produced no run")?; + assert_eq!(assignment.origin().repository(), "syncode/demo"); + assert_eq!(assignment.origin().reference(), "refs/pull/7/head"); + Ok(()) +} + +#[tokio::test] +async fn closing_a_pull_request_starts_nothing() -> TestResult { + let harness = start().await?; + + let body = pull_request_body("closed"); + assert_eq!(deliver(&harness, "pull_request", &body, true).await?, 202); + assert!( + harness + .runs + .take_next(syncode_control_runs::NodeId::fresh()) + .await? + .is_none() + ); + Ok(()) +} + +#[tokio::test] +async fn manual_and_scheduled_deliveries_become_runs() -> TestResult { + for event in ["workflow_dispatch", "schedule"] { + let harness = start().await?; + let body = format!( + r#"{{"ref":"refs/heads/main","after":"{COMMIT}","workflow":".gitea/workflows/ci.yml","repository":{{"full_name":"syncode/demo"}}}}"# + ); + assert_eq!(deliver(&harness, event, &body, true).await?, 202); + let assignment = harness + .runs + .take_next(syncode_control_runs::NodeId::fresh()) + .await? + .ok_or("event produced no run")?; + assert_eq!(assignment.origin().event(), event); + assert_eq!(assignment.origin().reference(), "refs/heads/main"); + assert!( + harness + .runs + .take_next(syncode_control_runs::NodeId::fresh()) + .await? + .is_none(), + "a requested event started a workflow other than the selected file" + ); + } + Ok(()) +} + +#[tokio::test] +async fn a_manual_delivery_can_select_a_tag() -> TestResult { + let harness = start().await?; + let body = format!( + r#"{{"ref":"refs/tags/v0.4.0","after":"{COMMIT}","workflow":".gitea/workflows/ci.yml","repository":{{"full_name":"syncode/demo"}}}}"# + ); + + assert_eq!( + deliver(&harness, "workflow_dispatch", &body, true).await?, + 202 + ); + let assignment = harness + .runs + .take_next(NodeId::fresh()) + .await? + .ok_or("manual tag event produced no run")?; + assert_eq!(assignment.origin().reference(), "refs/tags/v0.4.0"); + Ok(()) +} + +#[tokio::test] +async fn shadow_mode_persists_the_run_but_never_assigns_it() -> TestResult { + let harness = start_mode(true).await?; + let body = push_body(); + assert_eq!( + deliver(&harness, "repository.ref.updated", &body, true).await?, + 202 + ); + let (_sender, mut inbound) = open_node(&harness).await?; + + assert!( + tokio::time::timeout(Duration::from_millis(100), inbound.next()) + .await + .is_err(), + "shadow mode must not assign the run" + ); + assert!( + harness.runs.take_next(NodeId::fresh()).await?.is_some(), + "the shadow run must still exist in the aggregate" + ); + Ok(()) +} -- SynCode