feat: Complete native Actions delivery #53

Manually merged
day01 merged 9 commits from feat/0.6-verify-dynamic-git into develop 2026-08-31 07:27:18 +00:00
3 changed files with 31 additions and 12 deletions
Showing only changes of commit 597a9bdc4c - Show all commits
+17 -10
View File
@@ -1,225 +1,232 @@
use axum::extract::{Path, State};
use axum::http::{HeaderMap, Method, StatusCode, header};
use axum::routing::get;
use axum::{Json, Router};
use serde::Serialize;
use syncode_control_node::ActionsReadService;
use syncode_control_node::actions_wire::{
ActionJob, ActionRun, GetJobLogsRequest, GetRunRequest, ListRunsRequest,
};
use syncode_control_runs::RunLog;
use tonic::{Code, Request};
use tower_http::cors::{AllowOrigin, CorsLayer};

use crate::actions_read::{ActionsAuthorization, ActionsRead};

const SESSION_COOKIE: &str = "syncode_identity_session";

pub fn router<L, A>(read: ActionsRead<L, A>, cors_origins: Vec<String>) -> Router
where
L: RunLog + Send + Sync + 'static,
A: ActionsAuthorization + Clone,
{
let mut router = Router::new()
.route(
"/control/actions/repositories/{owner}/{repository}/runs",
get(list_runs::<L, A>),
)
.route("/control/actions/runs/{run}", get(get_run::<L, A>))
.route(
"/control/actions/runs/{run}/jobs/{job}/logs",
get(get_job_logs::<L, A>),
)
.with_state(read);
if !cors_origins.is_empty() {
router = router.layer(
CorsLayer::new()
.allow_credentials(true)
.allow_methods([Method::GET, Method::OPTIONS])
.allow_origin(AllowOrigin::predicate(move |origin, _| {
origin
.to_str()
.is_ok_and(|origin| cors_origins.iter().any(|value| value == origin))
})),
);
}
router
}

async fn list_runs<L, A>(
State(read): State<ActionsRead<L, A>>,
Path((owner, repository)): Path<(String, String)>,
headers: HeaderMap,
) -> Result<Json<Vec<WebRun>>, WebError>
where
L: RunLog + Send + Sync + 'static,
A: ActionsAuthorization,
{
let response = read
.list_runs(Request::new(ListRunsRequest {
session_token: session(&headers)?.to_owned(),
owner,
repository,
}))
.await
.map_err(WebError::from)?
.into_inner();
Ok(Json(response.runs.into_iter().map(WebRun::from).collect()))
}

async fn get_run<L, A>(
State(read): State<ActionsRead<L, A>>,
Path(run): Path<String>,
headers: HeaderMap,
) -> Result<Json<WebRun>, WebError>
where
L: RunLog + Send + Sync + 'static,
A: ActionsAuthorization,
{
read.get_run(Request::new(GetRunRequest {
session_token: session(&headers)?.to_owned(),
run_id: run,
}))
.await
.map(|response| Json(response.into_inner().into()))
.map_err(WebError::from)
}

async fn get_job_logs<L, A>(
State(read): State<ActionsRead<L, A>>,
Path((run, job)): Path<(String, String)>,
headers: HeaderMap,
) -> Result<Json<WebLogs>, WebError>
where
L: RunLog + Send + Sync + 'static,
A: ActionsAuthorization,
{
read.get_job_logs(Request::new(GetJobLogsRequest {
session_token: session(&headers)?.to_owned(),
run_id: run,
job_id: job,
}))
.await
.map(|response| {
Json(WebLogs {
lines: response.into_inner().lines,
})
})
.map_err(WebError::from)
}

fn session(headers: &HeaderMap) -> Result<&str, WebError> {
let cookie = headers
.get(header::COOKIE)
.ok_or_else(WebError::unauthenticated)?
.to_str()
.map_err(|_| WebError::unauthenticated())?;
cookie
.split(';')
.map(str::trim)
.find_map(|pair| pair.strip_prefix(&format!("{SESSION_COOKIE}=")))
.filter(|token| !token.is_empty())
.ok_or_else(WebError::unauthenticated)
}

#[derive(Serialize)]
pub struct WebRun {
id: String,
number: u64,
commit: String,
reference: String,
event: String,
workflow: String,
state: String,
conclusion: String,
jobs: Vec<WebJob>,
}

impl From<ActionRun> for WebRun {
fn from(run: ActionRun) -> Self {
Self {
id: run.id,
number: run.number,
commit: run.commit,
reference: run.reference,
event: run.event,
workflow: run.workflow,
state: run.state,
conclusion: run.conclusion,
jobs: run.jobs.into_iter().map(WebJob::from).collect(),
}
}
}

#[derive(Serialize)]
struct WebJob {
id: String,
key: String,
state: String,
conclusion: String,
}

impl From<ActionJob> for WebJob {
fn from(job: ActionJob) -> Self {
Self {
id: job.id,
key: job.key,
state: job.state,
conclusion: job.conclusion,
}
}
}

#[derive(Serialize)]
struct WebLogs {
lines: Vec<String>,
}

pub struct WebError {
status: StatusCode,
message: String,
}

impl WebError {
fn unauthenticated() -> Self {
Self {
status: StatusCode::UNAUTHORIZED,
message: "a SynCode identity session is required".to_owned(),
}
}
}

impl From<tonic::Status> for WebError {
fn from(error: tonic::Status) -> Self {
let status = match error.code() {
Code::InvalidArgument => StatusCode::BAD_REQUEST,
Code::Unauthenticated => StatusCode::UNAUTHORIZED,
Code::PermissionDenied => StatusCode::FORBIDDEN,
Code::NotFound => StatusCode::NOT_FOUND,
Code::Aborted => StatusCode::CONFLICT,
Code::Unavailable => StatusCode::SERVICE_UNAVAILABLE,
_ => StatusCode::INTERNAL_SERVER_ERROR,
};
Self {
status,
message: error.message().to_owned(),
}
}
}

impl axum::response::IntoResponse for WebError {
fn into_response(self) -> axum::response::Response {
(
self.status,
Json(ErrorBody {
error: self.message,
}),
)
.into_response()
}
}

#[derive(Serialize)]
struct ErrorBody {
error: String,
}
use axum::extract::{Path, State};
use axum::http::{HeaderMap, Method, StatusCode, header};
use axum::routing::get;
use axum::{Extension, Json, Router};
use serde::Serialize;
use std::sync::Arc;
use syncode_control_node::ActionsReadService;
use syncode_control_node::actions_wire::{
ActionJob, ActionRun, GetJobLogsRequest, GetRunRequest, ListRunsRequest,
};
use syncode_control_runs::RunLog;
use tonic::{Code, Request};
use tower_http::cors::{AllowOrigin, CorsLayer};

use crate::actions_read::{ActionsAuthorization, ActionsRead};

pub fn router<L, A>(
read: ActionsRead<L, A>,
cors_origins: Vec<String>,
session_cookie_name: String,
) -> Router
where
L: RunLog + Send + Sync + 'static,
A: ActionsAuthorization + Clone,
{
let mut router = Router::new()
.route(
"/control/actions/repositories/{owner}/{repository}/runs",
get(list_runs::<L, A>),
)
.route("/control/actions/runs/{run}", get(get_run::<L, A>))
.route(
"/control/actions/runs/{run}/jobs/{job}/logs",
get(get_job_logs::<L, A>),
)
.with_state(read)
.layer(Extension(Arc::<str>::from(session_cookie_name)));
if !cors_origins.is_empty() {
router = router.layer(
CorsLayer::new()
.allow_credentials(true)
.allow_methods([Method::GET, Method::OPTIONS])
.allow_origin(AllowOrigin::predicate(move |origin, _| {
origin
.to_str()
.is_ok_and(|origin| cors_origins.iter().any(|value| value == origin))
})),
);
}
router
}

async fn list_runs<L, A>(
State(read): State<ActionsRead<L, A>>,
Extension(session_cookie_name): Extension<Arc<str>>,
Path((owner, repository)): Path<(String, String)>,
headers: HeaderMap,
) -> Result<Json<Vec<WebRun>>, WebError>
where
L: RunLog + Send + Sync + 'static,
A: ActionsAuthorization,
{
let response = read
.list_runs(Request::new(ListRunsRequest {
session_token: session(&headers, &session_cookie_name)?.to_owned(),
owner,
repository,
}))
.await
.map_err(WebError::from)?
.into_inner();
Ok(Json(response.runs.into_iter().map(WebRun::from).collect()))
}

async fn get_run<L, A>(
State(read): State<ActionsRead<L, A>>,
Extension(session_cookie_name): Extension<Arc<str>>,
Path(run): Path<String>,
headers: HeaderMap,
) -> Result<Json<WebRun>, WebError>
where
L: RunLog + Send + Sync + 'static,
A: ActionsAuthorization,
{
read.get_run(Request::new(GetRunRequest {
session_token: session(&headers, &session_cookie_name)?.to_owned(),
run_id: run,
}))
.await
.map(|response| Json(response.into_inner().into()))
.map_err(WebError::from)
}

async fn get_job_logs<L, A>(
State(read): State<ActionsRead<L, A>>,
Extension(session_cookie_name): Extension<Arc<str>>,
Path((run, job)): Path<(String, String)>,
headers: HeaderMap,
) -> Result<Json<WebLogs>, WebError>
where
L: RunLog + Send + Sync + 'static,
A: ActionsAuthorization,
{
read.get_job_logs(Request::new(GetJobLogsRequest {
session_token: session(&headers, &session_cookie_name)?.to_owned(),
run_id: run,
job_id: job,
}))
.await
.map(|response| {
Json(WebLogs {
lines: response.into_inner().lines,
})
})
.map_err(WebError::from)
}

fn session<'a>(headers: &'a HeaderMap, session_cookie_name: &str) -> Result<&'a str, WebError> {
let cookie = headers
.get(header::COOKIE)
.ok_or_else(WebError::unauthenticated)?
.to_str()
.map_err(|_| WebError::unauthenticated())?;
cookie
.split(';')
.map(str::trim)
.find_map(|pair| pair.strip_prefix(&format!("{session_cookie_name}=")))
.filter(|token| !token.is_empty())
.ok_or_else(WebError::unauthenticated)
}

#[derive(Serialize)]
pub struct WebRun {
id: String,
number: u64,
commit: String,
reference: String,
event: String,
workflow: String,
state: String,
conclusion: String,
jobs: Vec<WebJob>,
}

impl From<ActionRun> for WebRun {
fn from(run: ActionRun) -> Self {
Self {
id: run.id,
number: run.number,
commit: run.commit,
reference: run.reference,
event: run.event,
workflow: run.workflow,
state: run.state,
conclusion: run.conclusion,
jobs: run.jobs.into_iter().map(WebJob::from).collect(),
}
}
}

#[derive(Serialize)]
struct WebJob {
id: String,
key: String,
state: String,
conclusion: String,
}

impl From<ActionJob> for WebJob {
fn from(job: ActionJob) -> Self {
Self {
id: job.id,
key: job.key,
state: job.state,
conclusion: job.conclusion,
}
}
}

#[derive(Serialize)]
struct WebLogs {
lines: Vec<String>,
}

pub struct WebError {
status: StatusCode,
message: String,
}

impl WebError {
fn unauthenticated() -> Self {
Self {
status: StatusCode::UNAUTHORIZED,
message: "a SynCode identity session is required".to_owned(),
}
}
}

impl From<tonic::Status> for WebError {
fn from(error: tonic::Status) -> Self {
let status = match error.code() {
Code::InvalidArgument => StatusCode::BAD_REQUEST,
Code::Unauthenticated => StatusCode::UNAUTHORIZED,
Code::PermissionDenied => StatusCode::FORBIDDEN,
Code::NotFound => StatusCode::NOT_FOUND,
Code::Aborted => StatusCode::CONFLICT,
Code::Unavailable => StatusCode::SERVICE_UNAVAILABLE,
_ => StatusCode::INTERNAL_SERVER_ERROR,
};
Self {
status,
message: error.message().to_owned(),
}
}
}

impl axum::response::IntoResponse for WebError {
fn into_response(self) -> axum::response::Response {
(
self.status,
Json(ErrorBody {
error: self.message,
}),
)
.into_response()
}
}

#[derive(Serialize)]
struct ErrorBody {
error: String,
}
+4
View File
@@ -1,385 +1,389 @@
use std::error::Error;
use std::net::SocketAddr;
use std::path::PathBuf;
use std::sync::Arc;

use clap::{Parser, Subcommand, ValueEnum};
use syncode_control::action_delivery::router as action_delivery_router;
use syncode_control::action_oci::PinnedOciResolver;
use syncode_control::action_repository::NativeActionRepository;
use syncode_control::action_store::FileActionStore;
use syncode_control::actions::ActionResolver;
use syncode_control::actions_read::ActionsRead;
use syncode_control::actions_read_http::router as actions_read_http_router;
use syncode_control::actions_read_identity::IdentityActionsAuthorization;
use syncode_control::admin::{Admin, router as admin_router};
use syncode_control::check_events::{CheckEventPublisher, CheckEventPublisherConfig};
use syncode_control::checks::Checks;
use syncode_control::maintenance;
use syncode_control::repository::{RepositoryContents, RepositorySecrets};
use syncode_control::repository_credentials::IdentityRepositoryCredentials;
use syncode_control::repository_grpc::NativeRepositoryContents;
use syncode_control::repository_sources::RepositorySources;
use syncode_control::secrets::RuntimeSecrets;
use syncode_control::token::EnrolmentScope;
use syncode_control::webhook::{Intake, router};
use syncode_control_node::{
ArtifactTokenAuthority, CapabilityAuthority, GeneratedActionsReadServer, GeneratedChecksServer,
GeneratedSecretsServer, GeneratedServer, NodeSessionServer, RepositoryCoordinates,
};
use syncode_control_nodes::{Nodes, Scope};
use syncode_control_runs::{
AuditConfiguration, AuditControlMode, AuditEvent, AuditPayload, NodeId, Runs, SchedulerPolicy,
};
use syncode_control_store::Postgres;
use tokio::net::TcpListener;
use tonic::transport::Server;
use url::Url;

#[derive(Debug, Parser)]
#[command(name = "syncode-control", version, about)]
struct Arguments {
#[command(subcommand)]
command: Option<Command>,

/// Where nodes open their session.
#[arg(long, default_value = "127.0.0.1:8090")]
listen: SocketAddr,

/// Where the native repository event feed delivers events.
#[arg(long, default_value = "127.0.0.1:8091")]
listen_events: SocketAddr,

/// The repository service this control plane reads workflows from.
#[arg(long, env = "SYNCODE_REPOSITORY_SOURCE_URL")]
repository_source: Url,

/// A token allowed to read repository contents.
#[arg(long, env = "SYNCODE_REPOSITORY_SOURCE_TOKEN", hide_env_values = true)]
repository_source_token: String,

/// Where native repositories are read over gRPC.
#[arg(long, env = "SYNCODE_REPOSITORY_GRPC_URL")]
repository_grpc: String,

#[arg(long, env = "SYNCODE_IDENTITY_GRPC_URL")]
identity_grpc: String,

#[arg(long, env = "SYNCODE_IDENTITY_SHARED_SECRET", hide_env_values = true)]
identity_shared_secret: String,

#[arg(long, env = "SYNCODE_CONTROL_CORS_ORIGINS", value_delimiter = ',')]
control_cors_origins: Vec<String>,

#[arg(long, env = "SYNCODE_ACTION_MIRROR_URL")]
action_mirror_url: Url,

#[arg(
long,
env = "SYNCODE_ACTION_ALLOWLIST",
value_delimiter = ',',
num_args = 1..
)]
action_allowlist: Vec<Url>,

#[arg(
long,
env = "SYNCODE_ACTION_OCI_REGISTRIES",
value_delimiter = ',',
num_args = 1..
)]
action_oci_registries: Vec<String>,

#[arg(long, env = "SYNCODE_ACTION_STORE")]
action_store: PathBuf,

#[arg(long, env = "SYNCODE_ACTION_ARTIFACT_PUBLIC_URL")]
action_artifact_public_url: Url,

#[arg(long, env = "SYNCODE_SECRET_SOURCE_URL")]
secret_source: Url,

#[arg(long, env = "SYNCODE_SECRET_SOURCE_TOKEN", hide_env_values = true)]
secret_source_token: String,

#[arg(long, env = "SYNCODE_CAPABILITY_SIGNING_KEY", hide_env_values = true)]
capability_signing_key: String,

#[arg(long, env = "SYNCODE_ARTIFACT_SIGNING_KEY", hide_env_values = true)]
artifact_signing_key: String,

/// The secret used to sign the native event feed.
#[arg(long, env = "SYNCODE_EVENT_FEED_SECRET", hide_env_values = true)]
event_feed_secret: String,

#[arg(long, env = "SYNCODE_COLLAB_EVENT_URL")]
collaboration_event_url: Url,

#[arg(long, env = "SYNCODE_COLLAB_EVENT_SECRET", hide_env_values = true)]
collaboration_event_secret: String,

#[arg(long, env = "SYNCODE_CONTROL_SOURCE_NODE_ID")]
source_node_id: uuid::Uuid,

#[arg(
long,
env = "SYNCODE_CHECK_EVENT_INTERVAL_SECONDS",
default_value_t = 1
)]
check_event_interval_seconds: u64,

#[arg(long, env = "SYNCODE_CHECK_EVENT_BATCH", default_value_t = 100)]
check_event_batch: i64,

/// Bearer token protecting the operational API.
#[arg(long, env = "SYNCODE_CONTROL_ADMIN_TOKEN", hide_env_values = true)]
admin_token: Option<String>,

/// Where the run log is kept.
#[arg(long, env = "SYNCODE_DATABASE_URL", hide_env_values = true)]
database: String,

#[arg(long, default_value_t = 8)]
database_connections: u32,

#[arg(long, env = "SYNCODE_ORGANIZATION_CONCURRENCY", default_value_t = 100)]
organization_concurrency: u32,

#[arg(long, env = "SYNCODE_REPOSITORY_CONCURRENCY", default_value_t = 20)]
repository_concurrency: u32,

#[arg(long, env = "SYNCODE_PRINCIPAL_CONCURRENCY", default_value_t = 20)]
principal_concurrency: u32,

#[arg(
long,
env = "SYNCODE_ORGANIZATION_QUEUE_QUOTA",
default_value_t = 10_000
)]
organization_queue_quota: u32,

#[arg(long, env = "SYNCODE_PRINCIPAL_QUEUE_QUOTA", default_value_t = 1_000)]
principal_queue_quota: u32,

#[arg(long, env = "SYNCODE_AUDIT_RETENTION_DAYS", default_value_t = 90)]
audit_retention_days: u32,

/// Whether compiled shadow runs may be assigned to nodes.
#[arg(long, value_enum, default_value_t = Mode::Shadow)]
mode: Mode,
}

#[derive(Clone, Copy, Debug, Eq, PartialEq, ValueEnum)]
enum Mode {
Shadow,
Active,
}

#[derive(Debug, Subcommand)]
enum Command {
/// Issue an enrolment token a node can spend for an identity.
///
/// The token is stored before it is printed, so one that reaches an
/// operator is one the control plane will honour.
IssueToken {
/// How far the token reaches: `instance`, `organisation:<name>` or
/// `repository:<owner>/<name>`.
#[arg(long, default_value = "instance")]
scope: String,
},
/// Revoke a node identity immediately.
RevokeNode {
/// The node UUID printed at enrolment.
node: String,
},
}

#[tokio::main]
async fn main() -> Result<(), Box<dyn Error>> {
let arguments = Arguments::parse();
// The log is read before anything is served: a node that reconnects has to
// meet the run it was holding, not a control plane that has forgotten it.
let store = Postgres::connect(&arguments.database, arguments.database_connections).await?;
let policy = SchedulerPolicy::new(
arguments.organization_concurrency,
arguments.repository_concurrency,
arguments.principal_concurrency,
arguments.organization_queue_quota,
arguments.principal_queue_quota,
);
let runs = Runs::restored_with_policy(store.clone(), policy).await?;
let nodes = Nodes::restored(store.clone()).await?;

if let Some(command) = arguments.command {
match command {
Command::IssueToken { scope } => {
let scope: Scope = scope.parse::<EnrolmentScope>()?.into();
let token = nodes.issue_token(scope).await?;
println!("{}", token.secret().expose());
}
Command::RevokeNode { node } => {
nodes.revoke(node.parse::<NodeId>()?).await?;
}
}
return Ok(());
}

store
.record_audit(AuditEvent::new("control.configuration", "applied").after(
AuditPayload::Configuration(AuditConfiguration {
mode: match arguments.mode {
Mode::Shadow => AuditControlMode::Shadow,
Mode::Active => AuditControlMode::Active,
},
organization_concurrency: arguments.organization_concurrency,
repository_concurrency: arguments.repository_concurrency,
principal_concurrency: arguments.principal_concurrency,
organization_queue_quota: arguments.organization_queue_quota,
principal_queue_quota: arguments.principal_queue_quota,
audit_retention_days: arguments.audit_retention_days,
}),
))
.await?;
tokio::spawn(maintenance::sweep(runs.clone(), nodes.clone()));
tokio::spawn(maintenance::retain_audit(
store.clone(),
arguments.audit_retention_days,
));

let action_mirror = arguments
.action_mirror_url
.to_string()
.parse::<syncode_workflow::RepositoryUrl>()?;
let action_allowlist = arguments
.action_allowlist
.into_iter()
.map(|repository| {
repository
.to_string()
.parse::<syncode_workflow::RepositoryUrl>()
})
.collect::<Result<Vec<_>, _>>()?;
let action_store = FileActionStore::new(arguments.action_store);
let action_repository = NativeActionRepository::connect(
arguments.repository_grpc.clone(),
arguments.identity_grpc.clone(),
arguments.identity_shared_secret.clone(),
)
.await?;
let action_resolver = ActionResolver::new(
action_repository,
action_store.clone(),
PinnedOciResolver::new(arguments.action_oci_registries)?,
action_mirror,
action_allowlist,
);
let native_repository =
NativeRepositoryContents::connect(arguments.repository_grpc.clone()).await?;
let repository_sources = RepositorySources::new(
native_repository,
RepositoryContents::new(
arguments.repository_source,
arguments.repository_source_token,
),
);
let intake = Arc::new(Intake::new(
repository_sources,
runs.clone(),
action_resolver,
arguments.event_feed_secret,
));
let admin_token = arguments
.admin_token
.filter(|token| !token.is_empty())
.ok_or("SYNCODE_CONTROL_ADMIN_TOKEN is required while serving")?;
let artifact_authority = ArtifactTokenAuthority::new(arguments.artifact_signing_key)?;
let actions_read = ActionsRead::new(
runs.clone(),
IdentityActionsAuthorization::connect(
arguments.identity_grpc.clone(),
arguments.identity_shared_secret.clone(),
)
.await?,
);
let http = router(intake)
.merge(admin_router(Arc::new(
Admin::new(runs.clone(), nodes.clone(), admin_token).with_operations(store.clone()),
)))
.merge(action_delivery_router(
action_store,
artifact_authority.clone(),
runs.clone(),
))
.merge(actions_read_http_router(
actions_read.clone(),
arguments.control_cors_origins,
));
let events = TcpListener::bind(arguments.listen_events).await?;
let check_events = CheckEventPublisher::new(
store.clone(),
CheckEventPublisherConfig {
endpoint: arguments.collaboration_event_url.as_str(),
secret: arguments.collaboration_event_secret,
source_node_id: arguments.source_node_id,
interval: std::time::Duration::from_secs(arguments.check_event_interval_seconds),
batch: arguments.check_event_batch,
},
)?;

let authority = CapabilityAuthority::new(arguments.capability_signing_key)?;
let repositories = RepositoryCoordinates::new(
arguments.identity_grpc.clone(),
arguments.identity_shared_secret.clone(),
)?;
let secret_service = RuntimeSecrets::new(
runs.clone(),
nodes.clone(),
RepositorySecrets::new(arguments.secret_source, arguments.secret_source_token),
IdentityRepositoryCredentials::connect(
arguments.identity_grpc.clone(),
arguments.identity_shared_secret.clone(),
)
.await?,
authority.clone(),
);
// Both ends of the control plane run for as long as the other does: without
// events there is nothing to assign, and without sessions there is nobody to
// assign it to. Whichever stops first takes the process down with it.
let node_service = match arguments.mode {
Mode::Shadow => NodeSessionServer::shadow(
runs.clone(),
nodes,
authority,
artifact_authority,
arguments.action_artifact_public_url,
repositories.clone(),
),
Mode::Active => NodeSessionServer::new(
runs.clone(),
nodes,
authority,
artifact_authority,
arguments.action_artifact_public_url,
repositories,
),
};
tokio::select! {
served = axum::serve(events, http).into_future() => served?,
published = check_events.run() => published?,
served = Server::builder()
.add_service(GeneratedServer::new(node_service))
.add_service(GeneratedSecretsServer::new(secret_service))
.add_service(GeneratedChecksServer::new(Checks::new(runs.clone())))
.add_service(GeneratedActionsReadServer::new(actions_read))
.serve_with_shutdown(arguments.listen, shutdown()) => served?,
}

Ok(())
}

async fn shutdown() {
if let Err(error) = tokio::signal::ctrl_c().await {
eprintln!("cannot listen for shutdown, keeping the service running: {error}");
std::future::pending::<()>().await;
}
}
use std::error::Error;
use std::net::SocketAddr;
use std::path::PathBuf;
use std::sync::Arc;

use clap::{Parser, Subcommand, ValueEnum};
use syncode_control::action_delivery::router as action_delivery_router;
use syncode_control::action_oci::PinnedOciResolver;
use syncode_control::action_repository::NativeActionRepository;
use syncode_control::action_store::FileActionStore;
use syncode_control::actions::ActionResolver;
use syncode_control::actions_read::ActionsRead;
use syncode_control::actions_read_http::router as actions_read_http_router;
use syncode_control::actions_read_identity::IdentityActionsAuthorization;
use syncode_control::admin::{Admin, router as admin_router};
use syncode_control::check_events::{CheckEventPublisher, CheckEventPublisherConfig};
use syncode_control::checks::Checks;
use syncode_control::maintenance;
use syncode_control::repository::{RepositoryContents, RepositorySecrets};
use syncode_control::repository_credentials::IdentityRepositoryCredentials;
use syncode_control::repository_grpc::NativeRepositoryContents;
use syncode_control::repository_sources::RepositorySources;
use syncode_control::secrets::RuntimeSecrets;
use syncode_control::token::EnrolmentScope;
use syncode_control::webhook::{Intake, router};
use syncode_control_node::{
ArtifactTokenAuthority, CapabilityAuthority, GeneratedActionsReadServer, GeneratedChecksServer,
GeneratedSecretsServer, GeneratedServer, NodeSessionServer, RepositoryCoordinates,
};
use syncode_control_nodes::{Nodes, Scope};
use syncode_control_runs::{
AuditConfiguration, AuditControlMode, AuditEvent, AuditPayload, NodeId, Runs, SchedulerPolicy,
};
use syncode_control_store::Postgres;
use tokio::net::TcpListener;
use tonic::transport::Server;
use url::Url;

#[derive(Debug, Parser)]
#[command(name = "syncode-control", version, about)]
struct Arguments {
#[command(subcommand)]
command: Option<Command>,

/// Where nodes open their session.
#[arg(long, default_value = "127.0.0.1:8090")]
listen: SocketAddr,

/// Where the native repository event feed delivers events.
#[arg(long, default_value = "127.0.0.1:8091")]
listen_events: SocketAddr,

/// The repository service this control plane reads workflows from.
#[arg(long, env = "SYNCODE_REPOSITORY_SOURCE_URL")]
repository_source: Url,

/// A token allowed to read repository contents.
#[arg(long, env = "SYNCODE_REPOSITORY_SOURCE_TOKEN", hide_env_values = true)]
repository_source_token: String,

/// Where native repositories are read over gRPC.
#[arg(long, env = "SYNCODE_REPOSITORY_GRPC_URL")]
repository_grpc: String,

#[arg(long, env = "SYNCODE_IDENTITY_GRPC_URL")]
identity_grpc: String,

#[arg(long, env = "SYNCODE_IDENTITY_SHARED_SECRET", hide_env_values = true)]
identity_shared_secret: String,

#[arg(long, env = "SYNCODE_CONTROL_CORS_ORIGINS", value_delimiter = ',')]
control_cors_origins: Vec<String>,

#[arg(long, env = "SYNCODE_IDENTITY_SESSION_COOKIE_NAME")]
identity_session_cookie_name: String,

#[arg(long, env = "SYNCODE_ACTION_MIRROR_URL")]
action_mirror_url: Url,

#[arg(
long,
env = "SYNCODE_ACTION_ALLOWLIST",
value_delimiter = ',',
num_args = 1..
)]
action_allowlist: Vec<Url>,

#[arg(
long,
env = "SYNCODE_ACTION_OCI_REGISTRIES",
value_delimiter = ',',
num_args = 1..
)]
action_oci_registries: Vec<String>,

#[arg(long, env = "SYNCODE_ACTION_STORE")]
action_store: PathBuf,

#[arg(long, env = "SYNCODE_ACTION_ARTIFACT_PUBLIC_URL")]
action_artifact_public_url: Url,

#[arg(long, env = "SYNCODE_SECRET_SOURCE_URL")]
secret_source: Url,

#[arg(long, env = "SYNCODE_SECRET_SOURCE_TOKEN", hide_env_values = true)]
secret_source_token: String,

#[arg(long, env = "SYNCODE_CAPABILITY_SIGNING_KEY", hide_env_values = true)]
capability_signing_key: String,

#[arg(long, env = "SYNCODE_ARTIFACT_SIGNING_KEY", hide_env_values = true)]
artifact_signing_key: String,

/// The secret used to sign the native event feed.
#[arg(long, env = "SYNCODE_EVENT_FEED_SECRET", hide_env_values = true)]
event_feed_secret: String,

#[arg(long, env = "SYNCODE_COLLAB_EVENT_URL")]
collaboration_event_url: Url,

#[arg(long, env = "SYNCODE_COLLAB_EVENT_SECRET", hide_env_values = true)]
collaboration_event_secret: String,

#[arg(long, env = "SYNCODE_CONTROL_SOURCE_NODE_ID")]
source_node_id: uuid::Uuid,

#[arg(
long,
env = "SYNCODE_CHECK_EVENT_INTERVAL_SECONDS",
default_value_t = 1
)]
check_event_interval_seconds: u64,

#[arg(long, env = "SYNCODE_CHECK_EVENT_BATCH", default_value_t = 100)]
check_event_batch: i64,

/// Bearer token protecting the operational API.
#[arg(long, env = "SYNCODE_CONTROL_ADMIN_TOKEN", hide_env_values = true)]
admin_token: Option<String>,

/// Where the run log is kept.
#[arg(long, env = "SYNCODE_DATABASE_URL", hide_env_values = true)]
database: String,

#[arg(long, default_value_t = 8)]
database_connections: u32,

#[arg(long, env = "SYNCODE_ORGANIZATION_CONCURRENCY", default_value_t = 100)]
organization_concurrency: u32,

#[arg(long, env = "SYNCODE_REPOSITORY_CONCURRENCY", default_value_t = 20)]
repository_concurrency: u32,

#[arg(long, env = "SYNCODE_PRINCIPAL_CONCURRENCY", default_value_t = 20)]
principal_concurrency: u32,

#[arg(
long,
env = "SYNCODE_ORGANIZATION_QUEUE_QUOTA",
default_value_t = 10_000
)]
organization_queue_quota: u32,

#[arg(long, env = "SYNCODE_PRINCIPAL_QUEUE_QUOTA", default_value_t = 1_000)]
principal_queue_quota: u32,

#[arg(long, env = "SYNCODE_AUDIT_RETENTION_DAYS", default_value_t = 90)]
audit_retention_days: u32,

/// Whether compiled shadow runs may be assigned to nodes.
#[arg(long, value_enum, default_value_t = Mode::Shadow)]
mode: Mode,
}

#[derive(Clone, Copy, Debug, Eq, PartialEq, ValueEnum)]
enum Mode {
Shadow,
Active,
}

#[derive(Debug, Subcommand)]
enum Command {
/// Issue an enrolment token a node can spend for an identity.
///
/// The token is stored before it is printed, so one that reaches an
/// operator is one the control plane will honour.
IssueToken {
/// How far the token reaches: `instance`, `organisation:<name>` or
/// `repository:<owner>/<name>`.
#[arg(long, default_value = "instance")]
scope: String,
},
/// Revoke a node identity immediately.
RevokeNode {
/// The node UUID printed at enrolment.
node: String,
},
}

#[tokio::main]
async fn main() -> Result<(), Box<dyn Error>> {
let arguments = Arguments::parse();
// The log is read before anything is served: a node that reconnects has to
// meet the run it was holding, not a control plane that has forgotten it.
let store = Postgres::connect(&arguments.database, arguments.database_connections).await?;
let policy = SchedulerPolicy::new(
arguments.organization_concurrency,
arguments.repository_concurrency,
arguments.principal_concurrency,
arguments.organization_queue_quota,
arguments.principal_queue_quota,
);
let runs = Runs::restored_with_policy(store.clone(), policy).await?;
let nodes = Nodes::restored(store.clone()).await?;

if let Some(command) = arguments.command {
match command {
Command::IssueToken { scope } => {
let scope: Scope = scope.parse::<EnrolmentScope>()?.into();
let token = nodes.issue_token(scope).await?;
println!("{}", token.secret().expose());
}
Command::RevokeNode { node } => {
nodes.revoke(node.parse::<NodeId>()?).await?;
}
}
return Ok(());
}

store
.record_audit(AuditEvent::new("control.configuration", "applied").after(
AuditPayload::Configuration(AuditConfiguration {
mode: match arguments.mode {
Mode::Shadow => AuditControlMode::Shadow,
Mode::Active => AuditControlMode::Active,
},
organization_concurrency: arguments.organization_concurrency,
repository_concurrency: arguments.repository_concurrency,
principal_concurrency: arguments.principal_concurrency,
organization_queue_quota: arguments.organization_queue_quota,
principal_queue_quota: arguments.principal_queue_quota,
audit_retention_days: arguments.audit_retention_days,
}),
))
.await?;
tokio::spawn(maintenance::sweep(runs.clone(), nodes.clone()));
tokio::spawn(maintenance::retain_audit(
store.clone(),
arguments.audit_retention_days,
));

let action_mirror = arguments
.action_mirror_url
.to_string()
.parse::<syncode_workflow::RepositoryUrl>()?;
let action_allowlist = arguments
.action_allowlist
.into_iter()
.map(|repository| {
repository
.to_string()
.parse::<syncode_workflow::RepositoryUrl>()
})
.collect::<Result<Vec<_>, _>>()?;
let action_store = FileActionStore::new(arguments.action_store);
let action_repository = NativeActionRepository::connect(
arguments.repository_grpc.clone(),
arguments.identity_grpc.clone(),
arguments.identity_shared_secret.clone(),
)
.await?;
let action_resolver = ActionResolver::new(
action_repository,
action_store.clone(),
PinnedOciResolver::new(arguments.action_oci_registries)?,
action_mirror,
action_allowlist,
);
let native_repository =
NativeRepositoryContents::connect(arguments.repository_grpc.clone()).await?;
let repository_sources = RepositorySources::new(
native_repository,
RepositoryContents::new(
arguments.repository_source,
arguments.repository_source_token,
),
);
let intake = Arc::new(Intake::new(
repository_sources,
runs.clone(),
action_resolver,
arguments.event_feed_secret,
));
let admin_token = arguments
.admin_token
.filter(|token| !token.is_empty())
.ok_or("SYNCODE_CONTROL_ADMIN_TOKEN is required while serving")?;
let artifact_authority = ArtifactTokenAuthority::new(arguments.artifact_signing_key)?;
let actions_read = ActionsRead::new(
runs.clone(),
IdentityActionsAuthorization::connect(
arguments.identity_grpc.clone(),
arguments.identity_shared_secret.clone(),
)
.await?,
);
let http = router(intake)
.merge(admin_router(Arc::new(
Admin::new(runs.clone(), nodes.clone(), admin_token).with_operations(store.clone()),
)))
.merge(action_delivery_router(
action_store,
artifact_authority.clone(),
runs.clone(),
))
.merge(actions_read_http_router(
actions_read.clone(),
arguments.control_cors_origins,
arguments.identity_session_cookie_name,
));
let events = TcpListener::bind(arguments.listen_events).await?;
let check_events = CheckEventPublisher::new(
store.clone(),
CheckEventPublisherConfig {
endpoint: arguments.collaboration_event_url.as_str(),
secret: arguments.collaboration_event_secret,
source_node_id: arguments.source_node_id,
interval: std::time::Duration::from_secs(arguments.check_event_interval_seconds),
batch: arguments.check_event_batch,
},
)?;

let authority = CapabilityAuthority::new(arguments.capability_signing_key)?;
let repositories = RepositoryCoordinates::new(
arguments.identity_grpc.clone(),
arguments.identity_shared_secret.clone(),
)?;
let secret_service = RuntimeSecrets::new(
runs.clone(),
nodes.clone(),
RepositorySecrets::new(arguments.secret_source, arguments.secret_source_token),
IdentityRepositoryCredentials::connect(
arguments.identity_grpc.clone(),
arguments.identity_shared_secret.clone(),
)
.await?,
authority.clone(),
);
// Both ends of the control plane run for as long as the other does: without
// events there is nothing to assign, and without sessions there is nobody to
// assign it to. Whichever stops first takes the process down with it.
let node_service = match arguments.mode {
Mode::Shadow => NodeSessionServer::shadow(
runs.clone(),
nodes,
authority,
artifact_authority,
arguments.action_artifact_public_url,
repositories.clone(),
),
Mode::Active => NodeSessionServer::new(
runs.clone(),
nodes,
authority,
artifact_authority,
arguments.action_artifact_public_url,
repositories,
),
};
tokio::select! {
served = axum::serve(events, http).into_future() => served?,
published = check_events.run() => published?,
served = Server::builder()
.add_service(GeneratedServer::new(node_service))
.add_service(GeneratedSecretsServer::new(secret_service))
.add_service(GeneratedChecksServer::new(Checks::new(runs.clone())))
.add_service(GeneratedActionsReadServer::new(actions_read))
.serve_with_shutdown(arguments.listen, shutdown()) => served?,
}

Ok(())
}

async fn shutdown() {
if let Err(error) = tokio::signal::ctrl_c().await {
eprintln!("cannot listen for shutdown, keeping the service running: {error}");
std::future::pending::<()>().await;
}
}
+10 -2
View File
@@ -1,108 +1,116 @@
#![allow(clippy::expect_used)]

use std::error::Error;

use axum::body::{Body, to_bytes};
use axum::http::{Request, StatusCode, header};
use syncode_control::actions_read::{ActionsAuthorization, ActionsRead, AuthorizationError};
use syncode_control::actions_read_http;
use syncode_control_runs::{Forgotten, JobId, Origin, Runs};
use tower::ServiceExt;

type TestResult<T = ()> = Result<T, Box<dyn Error + Send + Sync>>;

const REPOSITORY_ID: &str = "018f47e2-b2c4-7f19-8a6d-13ef76c89214";

#[derive(Clone)]
struct Authorization;

#[tonic::async_trait]
impl ActionsAuthorization for Authorization {
async fn repository(
&self,
session_token: &str,
owner: &str,
name: &str,
) -> Result<String, AuthorizationError> {
if session_token != "valid-session" {
return Err(AuthorizationError::Unauthenticated);
}
if owner == "syncode" && name == "control" {
Ok(REPOSITORY_ID.to_owned())
} else {
Err(AuthorizationError::NotFound)
}
}

async fn authorize(
&self,
session_token: &str,
repository_id: &str,
) -> Result<(), AuthorizationError> {
if session_token == "valid-session" && repository_id == REPOSITORY_ID {
Ok(())
} else {
Err(AuthorizationError::Denied)
}
}
}

#[tokio::test]
async fn lists_runs_for_a_browser_identity_session() -> TestResult {
let runs = Runs::restored(Forgotten::default()).await?;
let run = runs
.queue(
JobId::fresh(),
Origin::new(
REPOSITORY_ID.to_owned(),
"abc123".to_owned(),
"refs/heads/main".to_owned(),
"push".to_owned(),
".syncode/workflows/ci.yml".to_owned(),
),
b"plan".to_vec(),
)
.await?;
let router = actions_read_http::router(
ActionsRead::new(runs, Authorization),
vec!["https://new.dev.syncode.sh".to_owned()],
);

let response = router
.oneshot(
Request::builder()
.uri("/control/actions/repositories/syncode/control/runs")
.header(header::COOKIE, "syncode_identity_session=valid-session")
.header(header::ORIGIN, "https://new.dev.syncode.sh")
.body(Body::empty())?,
)
.await?;

assert_eq!(StatusCode::OK, response.status());
assert_eq!(
"https://new.dev.syncode.sh",
response.headers()[header::ACCESS_CONTROL_ALLOW_ORIGIN]
);
let body: serde_json::Value =
serde_json::from_slice(&to_bytes(response.into_body(), 4096).await?)?;
assert_eq!(run.to_string(), body[0]["id"]);
assert_eq!("abc123", body[0]["commit"]);
Ok(())
}

#[tokio::test]
async fn rejects_a_request_without_the_identity_cookie() -> TestResult {
let runs = Runs::restored(Forgotten::default()).await?;
let router = actions_read_http::router(ActionsRead::new(runs, Authorization), Vec::new());

let response = router
.oneshot(
Request::builder()
.uri("/control/actions/repositories/syncode/control/runs")
.body(Body::empty())?,
)
.await?;

assert_eq!(StatusCode::UNAUTHORIZED, response.status());
Ok(())
}
#![allow(clippy::expect_used)]

use std::error::Error;

use axum::body::{Body, to_bytes};
use axum::http::{Request, StatusCode, header};
use syncode_control::actions_read::{ActionsAuthorization, ActionsRead, AuthorizationError};
use syncode_control::actions_read_http;
use syncode_control_runs::{Forgotten, JobId, Origin, Runs};
use tower::ServiceExt;

type TestResult<T = ()> = Result<T, Box<dyn Error + Send + Sync>>;

const REPOSITORY_ID: &str = "018f47e2-b2c4-7f19-8a6d-13ef76c89214";

#[derive(Clone)]
struct Authorization;

#[tonic::async_trait]
impl ActionsAuthorization for Authorization {
async fn repository(
&self,
session_token: &str,
owner: &str,
name: &str,
) -> Result<String, AuthorizationError> {
if session_token != "valid-session" {
return Err(AuthorizationError::Unauthenticated);
}
if owner == "syncode" && name == "control" {
Ok(REPOSITORY_ID.to_owned())
} else {
Err(AuthorizationError::NotFound)
}
}

async fn authorize(
&self,
session_token: &str,
repository_id: &str,
) -> Result<(), AuthorizationError> {
if session_token == "valid-session" && repository_id == REPOSITORY_ID {
Ok(())
} else {
Err(AuthorizationError::Denied)
}
}
}

#[tokio::test]
async fn lists_runs_for_a_browser_identity_session() -> TestResult {
let runs = Runs::restored(Forgotten::default()).await?;
let run = runs
.queue(
JobId::fresh(),
Origin::new(
REPOSITORY_ID.to_owned(),
"abc123".to_owned(),
"refs/heads/main".to_owned(),
"push".to_owned(),
".syncode/workflows/ci.yml".to_owned(),
),
b"plan".to_vec(),
)
.await?;
let router = actions_read_http::router(
ActionsRead::new(runs, Authorization),
vec!["https://new.dev.syncode.sh".to_owned()],
"syncode_dev_identity_session".to_owned(),
);

let response = router
.oneshot(
Request::builder()
.uri("/control/actions/repositories/syncode/control/runs")
.header(
header::COOKIE,
"syncode_identity_session=prod-session; syncode_dev_identity_session=valid-session",
)
.header(header::ORIGIN, "https://new.dev.syncode.sh")
.body(Body::empty())?,
)
.await?;

assert_eq!(StatusCode::OK, response.status());
assert_eq!(
"https://new.dev.syncode.sh",
response.headers()[header::ACCESS_CONTROL_ALLOW_ORIGIN]
);
let body: serde_json::Value =
serde_json::from_slice(&to_bytes(response.into_body(), 4096).await?)?;
assert_eq!(run.to_string(), body[0]["id"]);
assert_eq!("abc123", body[0]["commit"]);
Ok(())
}

#[tokio::test]
async fn rejects_a_request_without_the_identity_cookie() -> TestResult {
let runs = Runs::restored(Forgotten::default()).await?;
let router = actions_read_http::router(
ActionsRead::new(runs, Authorization),
Vec::new(),
"syncode_dev_identity_session".to_owned(),
);

let response = router
.oneshot(
Request::builder()
.uri("/control/actions/repositories/syncode/control/runs")
.body(Body::empty())?,
)
.await?;

assert_eq!(StatusCode::UNAUTHORIZED, response.status());
Ok(())
}