diff --git a/.gitignore b/.gitignore index 4c65854ff..59ab4bb69 100644 --- a/.gitignore +++ b/.gitignore @@ -17,3 +17,8 @@ experiments/mitm/captures/ experiments/mitm/reports/ experiments/mitm/golden/ __pycache__/ + +# Credentials & Local Environment +.env +*.pem +*.key diff --git a/Cargo.lock b/Cargo.lock index f661aac91..8f60bc254 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -128,6 +128,8 @@ dependencies = [ "bytes", "clap", "hmac", + "rand 0.8.6", + "reqwest", "serde", "serde_json", "sha2", diff --git a/crates/aksh-gha-expressions/src/lib.rs b/crates/aksh-gha-expressions/src/lib.rs index e0b1cda8b..bd8a3f54d 100644 --- a/crates/aksh-gha-expressions/src/lib.rs +++ b/crates/aksh-gha-expressions/src/lib.rs @@ -709,6 +709,27 @@ mod tests { ); } + #[test] + fn handles_escaped_single_quotes_in_literals() { + let context = Context::default(); + assert_eq!( + eval_expression("'It''s a string'", &context).unwrap(), + Value::String("It's a string".to_owned()) + ); + assert_eq!( + eval_expression("''", &context).unwrap(), + Value::String("".to_owned()) + ); + assert_eq!( + eval_expression("''''", &context).unwrap(), + Value::String("'".to_owned()) + ); + assert_eq!( + eval_expression("'a''b''c'", &context).unwrap(), + Value::String("a'b'c".to_owned()) + ); + } + proptest! { #[test] fn string_equality_is_case_insensitive(value in "[A-Za-z]{1,24}") { diff --git a/crates/aksh-runner-server/Cargo.toml b/crates/aksh-runner-server/Cargo.toml index 9ae50d7e0..d897d1acb 100644 --- a/crates/aksh-runner-server/Cargo.toml +++ b/crates/aksh-runner-server/Cargo.toml @@ -23,6 +23,8 @@ serde.workspace = true serde_json.workspace = true hmac.workspace = true sha2.workspace = true +rand.workspace = true +reqwest.workspace = true thiserror.workspace = true tokio.workspace = true tokio-util.workspace = true diff --git a/crates/aksh-runner-server/src/github.rs b/crates/aksh-runner-server/src/github.rs new file mode 100644 index 000000000..5d97e12d0 --- /dev/null +++ b/crates/aksh-runner-server/src/github.rs @@ -0,0 +1,750 @@ +//! GitHub App Webhook Integration. + +use axum::{ + extract::{Query, State}, + http::{HeaderMap, StatusCode}, + response::IntoResponse, + Json, +}; +use hmac::{Hmac, Mac}; +use serde::{Deserialize, Serialize}; +use serde_json::Value; +use sha2::Sha256; +use std::collections::BTreeMap; +use std::path::PathBuf; +use std::sync::Arc; +use tracing::{error, info, warn}; + +use crate::{submit_run_inner, ExecutionStatus, SharedState}; +use aksh_gha_protocol::{JobId, RunId, WorkflowSubmission}; + +/// Webhook push event payload. +#[derive(Debug, Deserialize, Serialize, Clone)] +pub(crate) struct PushEvent { + /// Git reference for the push event. + #[serde(rename = "ref")] + pub(crate) git_ref: String, + /// Previous commit SHA. + pub(crate) before: String, + /// Current commit SHA. + pub(crate) after: String, + /// Repository info. + pub(crate) repository: RepositoryInfo, + /// Commits in this push. + pub(crate) commits: Vec, +} + +/// Repository info. +#[derive(Debug, Deserialize, Serialize, Clone)] +pub(crate) struct RepositoryInfo { + /// Full repository name (e.g. owner/repo). + pub(crate) full_name: String, + /// Default branch (e.g. main). + pub(crate) default_branch: Option, +} + +/// Commit info. +#[derive(Debug, Deserialize, Serialize, Clone)] +pub(crate) struct CommitInfo { + /// Commit ID. + pub(crate) id: String, + /// Added files. + pub(crate) added: Vec, + /// Modified files. + pub(crate) modified: Vec, + /// Removed files. + pub(crate) removed: Vec, +} + +/// Webhook pull request event payload. +#[derive(Debug, Deserialize, Serialize, Clone)] +pub(crate) struct PullRequestEvent { + /// Webhook action type. + pub(crate) action: String, + /// PR number. + pub(crate) number: u64, + /// PR details. + pub(crate) pull_request: PullRequestDetails, + /// Repository info. + pub(crate) repository: RepositoryInfo, +} + +/// Pull request details. +#[derive(Debug, Deserialize, Serialize, Clone)] +pub(crate) struct PullRequestDetails { + /// Head reference. + pub(crate) head: GitReference, + /// Base reference. + pub(crate) base: GitReference, +} + +/// Git reference. +#[derive(Debug, Deserialize, Serialize, Clone)] +pub(crate) struct GitReference { + /// Git reference name. + #[serde(rename = "ref")] + pub(crate) git_ref: String, + /// Commit SHA. + pub(crate) sha: String, +} + +/// Verify X-Hub-Signature-256 webhook signature. +pub(crate) fn verify_signature(secret: &str, payload: &[u8], signature_header: &str) -> bool { + let signature_hex = match signature_header.strip_prefix("sha256=") { + Some(hex) => hex, + None => return false, + }; + let signature_bytes = match decode_hex(signature_hex) { + Ok(bytes) => bytes, + Err(_) => return false, + }; + + type HmacSha256 = Hmac; + let mut mac = match HmacSha256::new_from_slice(secret.as_bytes()) { + Ok(m) => m, + Err(_) => return false, + }; + mac.update(payload); + mac.verify_slice(&signature_bytes).is_ok() +} + +fn decode_hex(hex: &str) -> Result, &'static str> { + if hex.len() % 2 != 0 { + return Err("Odd length"); + } + let mut bytes = Vec::with_capacity(hex.len() / 2); + for i in (0..hex.len()).step_by(2) { + let byte = u8::from_str_radix(&hex[i..i + 2], 16).map_err(|_| "Invalid hex character")?; + bytes.push(byte); + } + Ok(bytes) +} + +async fn send_github_check_request( + token: &str, + repo: &str, + method: reqwest::Method, + path: &str, + body: Value, +) -> anyhow::Result { + let client = reqwest::Client::new(); + let url = format!("https://api.github.com/repos/{}/{}", repo, path); + let res = client + .request(method, &url) + .header("User-Agent", "aksh") + .header("Authorization", format!("Bearer {}", token)) + .header("Accept", "application/vnd.github+json") + .json(&body) + .send() + .await?; + + if !res.status().is_success() { + let status = res.status(); + let err_text = res.text().await.unwrap_or_default(); + return Err(anyhow::anyhow!( + "GitHub Check API failed with status {}: {}", + status, + err_text + )); + } + + let val = res.json().await.unwrap_or(Value::Null); + Ok(val) +} + +/// Report a queued check run to GitHub or simulate it locally. +pub(crate) async fn report_check_run_queued( + shared: &Arc, + repo: &str, + sha: &str, + job_id: &JobId, + run_id: RunId, +) { + let token = std::env::var("AKSH_GITHUB_TOKEN").ok(); + let mut check_run_id = None; + + if let Some(token) = &token { + let body = serde_json::json!({ + "name": job_id.to_string(), + "head_sha": sha, + "status": "queued", + }); + + match send_github_check_request(token, repo, reqwest::Method::POST, "check-runs", body) + .await + { + Ok(res) => { + if let Some(id) = res.get("id").and_then(|id| id.as_u64()) { + check_run_id = Some(id); + info!( + %run_id, + %job_id, + check_run_id = id, + "GitHub check run created successfully" + ); + } + } + Err(e) => { + warn!(%run_id, %job_id, error = %e, "Failed to create GitHub check run"); + } + } + } else { + info!(%run_id, %job_id, "GitHub token not configured, using mock check run"); + check_run_id = Some(rand::random::() as u64); + } + + if let Some(check_id) = check_run_id { + let mut inner = shared.state.inner.lock().await; + if let Some(run) = inner.runs.get_mut(&run_id) { + run.job_check_run_ids.insert(job_id.clone(), check_id); + } + } +} + +/// Report check run status to in_progress on GitHub or simulate it locally. +pub(crate) async fn report_check_run_in_progress( + shared: &Arc, + run_id: RunId, + job_id: &JobId, +) { + let (repo, check_run_id) = { + let inner = shared.state.inner.lock().await; + let run = match inner.runs.get(&run_id) { + Some(r) => r, + None => return, + }; + let repo = run.submission.repository.clone(); + let check_run_id = match run.job_check_run_ids.get(job_id).copied() { + Some(id) => id, + None => return, + }; + (repo, check_run_id) + }; + + let token = std::env::var("AKSH_GITHUB_TOKEN").ok(); + if let Some(token) = &token { + let body = serde_json::json!({ + "status": "in_progress", + }); + + let path = format!("check-runs/{}", check_run_id); + if let Err(e) = + send_github_check_request(token, &repo, reqwest::Method::PATCH, &path, body).await + { + warn!( + %run_id, + %job_id, + check_run_id, + error = %e, + "Failed to update GitHub check run to in_progress" + ); + } + } else { + info!(%run_id, %job_id, check_run_id, "Mock updated check run to in_progress"); + } +} + +/// Report check run status to completed on GitHub or simulate it locally. +pub(crate) async fn report_check_run_completed( + shared: &Arc, + run_id: RunId, + job_id: &JobId, + status: ExecutionStatus, +) { + let (repo, check_run_id) = { + let inner = shared.state.inner.lock().await; + let run = match inner.runs.get(&run_id) { + Some(r) => r, + None => return, + }; + let repo = run.submission.repository.clone(); + let check_run_id = match run.job_check_run_ids.get(job_id).copied() { + Some(id) => id, + None => return, + }; + (repo, check_run_id) + }; + + let conclusion = match status { + ExecutionStatus::Success => "success", + ExecutionStatus::Failure => "failure", + ExecutionStatus::Cancelled => "cancelled", + _ => "failure", + }; + + let token = std::env::var("AKSH_GITHUB_TOKEN").ok(); + if let Some(token) = &token { + let body = serde_json::json!({ + "status": "completed", + "conclusion": conclusion, + }); + + let path = format!("check-runs/{}", check_run_id); + if let Err(e) = + send_github_check_request(token, &repo, reqwest::Method::PATCH, &path, body).await + { + warn!( + %run_id, + %job_id, + check_run_id, + error = %e, + "Failed to update GitHub check run to completed" + ); + } + } else { + info!( + %run_id, + %job_id, + check_run_id, + conclusion, + "Mock updated check run to completed" + ); + } +} + +/// Fetch workflows helper. +pub(crate) async fn fetch_workflows( + local_workspace: &Option, + repo: &str, + git_ref: &str, +) -> anyhow::Result> { + if let Some(base_path) = local_workspace { + let workflows_dir = base_path.join(".github/workflows"); + let mut workflows = BTreeMap::new(); + if workflows_dir.exists() { + let mut dir = tokio::fs::read_dir(workflows_dir).await?; + while let Some(entry) = dir.next_entry().await? { + let path = entry.path(); + if path.is_file() { + if let Some(ext) = path.extension() { + if ext == "yml" || ext == "yaml" { + if let Some(name) = path.file_name().and_then(|n| n.to_str()) { + let content = tokio::fs::read_to_string(&path).await?; + workflows.insert(name.to_owned(), content); + } + } + } + } + } + } + Ok(workflows) + } else { + let token = std::env::var("AKSH_GITHUB_TOKEN").ok(); + if let Some(token) = &token { + fetch_remote_workflows(token, repo, git_ref).await + } else { + // Default fallback to current workspace root if nothing is configured + let workflows_dir = PathBuf::from(".").join(".github/workflows"); + let mut workflows = BTreeMap::new(); + if workflows_dir.exists() { + let mut dir = tokio::fs::read_dir(workflows_dir).await?; + while let Some(entry) = dir.next_entry().await? { + let path = entry.path(); + if path.is_file() { + if let Some(ext) = path.extension() { + if ext == "yml" || ext == "yaml" { + if let Some(name) = path.file_name().and_then(|n| n.to_str()) { + let content = tokio::fs::read_to_string(&path).await?; + workflows.insert(name.to_owned(), content); + } + } + } + } + } + } + Ok(workflows) + } + } +} + +async fn fetch_remote_workflows( + token: &str, + repo: &str, + git_ref: &str, +) -> anyhow::Result> { + let client = reqwest::Client::new(); + let url = format!( + "https://api.github.com/repos/{}/contents/.github/workflows?ref={}", + repo, git_ref + ); + let response = client + .get(&url) + .header("User-Agent", "aksh") + .header("Authorization", format!("Bearer {}", token)) + .header("Accept", "application/vnd.github+json") + .send() + .await?; + + if !response.status().is_success() { + return Err(anyhow::anyhow!( + "GitHub API returned status: {}", + response.status() + )); + } + + #[derive(Deserialize)] + struct GitHubContentItem { + name: String, + r#type: String, + download_url: Option, + } + + let items: Vec = response.json().await?; + let mut workflows = BTreeMap::new(); + + for item in &items { + if item.r#type == "file" && (item.name.ends_with(".yml") || item.name.ends_with(".yaml")) { + if let Some(download_url) = &item.download_url { + let file_res = client + .get(download_url) + .header("User-Agent", "aksh") + .header("Authorization", format!("Bearer {}", token)) + .send() + .await?; + if file_res.status().is_success() { + let content = file_res.text().await?; + workflows.insert(item.name.clone(), content); + } + } + } + } + Ok(workflows) +} + +async fn get_pr_changed_files( + token: &str, + repo: &str, + pr_number: u64, +) -> anyhow::Result> { + let client = reqwest::Client::new(); + let mut page = 1; + let mut all_files = Vec::new(); + + #[derive(Deserialize)] + struct GitHubFileItem { + filename: String, + } + + loop { + let url = format!( + "https://api.github.com/repos/{}/pulls/{}/files?per_page=100&page={}", + repo, pr_number, page + ); + let response = client + .get(&url) + .header("User-Agent", "aksh") + .header("Authorization", format!("Bearer {}", token)) + .header("Accept", "application/vnd.github+json") + .send() + .await?; + + if !response.status().is_success() { + return Err(anyhow::anyhow!( + "GitHub API returned status: {}", + response.status() + )); + } + + let files: Vec = response.json().await?; + if files.is_empty() { + break; + } + + all_files.extend(files.into_iter().map(|f| f.filename)); + page += 1; + } + + Ok(all_files) +} + +/// Route handler for GitHub App Webhooks. +pub(crate) async fn handle_github_webhook( + State(shared): State>, + headers: HeaderMap, + body: bytes::Bytes, +) -> Result { + // 1. Verify Signature + let secret = shared.state.webhook_secret.as_ref().ok_or_else(|| { + warn!("Webhook secret not configured on server, rejecting request"); + StatusCode::UNAUTHORIZED + })?; + + let sig_header = headers + .get("x-hub-signature-256") + .and_then(|h| h.to_str().ok()) + .ok_or(StatusCode::UNAUTHORIZED)?; + + if !verify_signature(secret, &body, sig_header) { + return Err(StatusCode::UNAUTHORIZED); + } + + // 2. Get event type + let event_name = headers + .get("x-github-event") + .and_then(|h| h.to_str().ok()) + .ok_or(StatusCode::BAD_REQUEST)?; + + // 3. Process the event payload + let payload_val: Value = serde_json::from_slice(&body).map_err(|_| StatusCode::BAD_REQUEST)?; + + let (repo_full_name, git_ref, sha, changed_paths) = match event_name { + "push" => { + let event: PushEvent = + serde_json::from_value(payload_val.clone()).map_err(|_| StatusCode::BAD_REQUEST)?; + let mut paths = Vec::new(); + // Collect changed files from push commits + for commit in &event.commits { + paths.extend(commit.added.clone()); + paths.extend(commit.modified.clone()); + paths.extend(commit.removed.clone()); + } + paths.sort(); + paths.dedup(); + ( + event.repository.full_name, + event.git_ref, + event.after, + paths, + ) + } + "pull_request" => { + let event: PullRequestEvent = + serde_json::from_value(payload_val.clone()).map_err(|_| StatusCode::BAD_REQUEST)?; + let repo = event.repository.full_name; + let git_ref = format!("refs/pull/{}/merge", event.number); + let sha = event.pull_request.head.sha; + let mut paths = Vec::new(); + // Retrieve changed files from remote GitHub API if possible + if let Some(token) = &std::env::var("AKSH_GITHUB_TOKEN").ok() { + if let Ok(files) = get_pr_changed_files(token, &repo, event.number).await { + paths = files; + } + } + (repo, git_ref, sha, paths) + } + _ => { + info!("Received unsupported GitHub webhook event: {}", event_name); + return Ok((StatusCode::OK, Json(serde_json::json!([])))); + } + }; + + // 4. Fetch workflows + let workflows = fetch_workflows(&shared.state.local_workspace, &repo_full_name, &git_ref) + .await + .map_err(|e| { + error!("Failed to fetch workflows: {:?}", e); + StatusCode::INTERNAL_SERVER_ERROR + })?; + + let mut triggered_runs = Vec::new(); + + // 5. Evaluate and trigger matching workflows + for (filename, content) in workflows { + // Parse workflow to check triggers + let parsed = match aksh_gha_parser::parse_workflow(&content) { + Ok(w) => w, + Err(e) => { + warn!("Failed to parse workflow file {}: {:?}", filename, e); + continue; + } + }; + + let (branch, tag) = crate::git_ref_context(&git_ref); + let activity_type = payload_val.get("action").and_then(|v| v.as_str()); + + if parsed.on.matches_with_context( + event_name, + branch.as_deref(), + tag.as_deref(), + &changed_paths, + activity_type, + ) { + info!( + "Triggering workflow run for {} matching {}", + filename, event_name + ); + + // Construct WorkflowSubmission + let submission = WorkflowSubmission { + workflow_yaml: content, + event: event_name.to_owned(), + payload: payload_val.clone(), + repository: repo_full_name.clone(), + git_ref: git_ref.clone(), + vars: BTreeMap::new(), + secrets: BTreeMap::new(), + reusable_workflows: BTreeMap::new(), + }; + + // Call submit_run_inner + match submit_run_inner(&shared, submission).await { + Ok(accepted) => { + let run_id = accepted.run_id; + + // Create GitHub Check Runs for each job in the run + let jobs = { + let inner = shared.state.inner.lock().await; + inner + .runs + .get(&run_id) + .map(|r| r.jobs.keys().cloned().collect::>()) + }; + + if let Some(jobs) = jobs { + for job_id in jobs { + report_check_run_queued( + &shared, + &repo_full_name, + &sha, + &job_id, + run_id, + ) + .await; + } + } + + triggered_runs.push(accepted); + } + Err(e) => { + error!("Failed to submit run for {}: {:?}", filename, e); + } + } + } + } + + Ok((StatusCode::OK, Json(serde_json::json!(triggered_runs)))) +} + +/// Serve registration page for GitHub App Manifest flow. +pub(crate) async fn github_register(headers: HeaderMap) -> impl IntoResponse { + let host = headers + .get("host") + .and_then(|h| h.to_str().ok()) + .unwrap_or("localhost:9090"); + + let scheme = if host.contains("localhost") || host.contains("127.0.0.1") { + "http" + } else { + "https" + }; + + let base_url = format!("{}://{}", scheme, host); + let is_local = host.contains("localhost") || host.contains("127.0.0.1"); + + let mut manifest_json = serde_json::json!({ + "name": "aksh-local-app", + "url": base_url, + "redirect_url": format!("{}/api/v1/github/callback", base_url), + "public": false, + "default_permissions": { + "checks": "write", + "contents": "read", + "metadata": "read", + "pull_requests": "read" + } + }); + + if !is_local { + manifest_json["hook_attributes"] = serde_json::json!({ + "url": format!("{}/api/v1/github/webhooks", base_url) + }); + manifest_json["default_events"] = serde_json::json!(["push", "pull_request"]); + } + + let html = format!( + r#" + + + Register GitHub App + + +

Register GitHub App for aksh

+

Click the button below to register a local GitHub App on your GitHub account automatically.

+
+ + +
+ +"#, + manifest_json + ); + + axum::response::Html(html) +} + +/// Query parameters for GitHub callback. +#[derive(Debug, Deserialize)] +pub(crate) struct CallbackQuery { + code: String, +} + +/// Callback endpoint for GitHub App Manifest conversion. +pub(crate) async fn github_callback( + State(shared): State>, + Query(params): Query, +) -> Result { + let client = reqwest::Client::new(); + let api_base = std::env::var("AKSH_GITHUB_API_URL") + .unwrap_or_else(|_| "https://api.github.com".to_owned()); + let url = format!("{}/app-manifests/{}/conversions", api_base, params.code); + let res = client + .post(&url) + .header("User-Agent", "aksh") + .header("Accept", "application/vnd.github+json") + .send() + .await + .map_err(|_| StatusCode::INTERNAL_SERVER_ERROR)?; + + if !res.status().is_success() { + return Err(StatusCode::INTERNAL_SERVER_ERROR); + } + + #[derive(Deserialize)] + struct AppManifestConversion { + id: u64, + pem: String, + webhook_secret: Option, + } + + let credentials: AppManifestConversion = res + .json() + .await + .map_err(|_| StatusCode::INTERNAL_SERVER_ERROR)?; + + info!("Successfully registered GitHub App ID: {}", credentials.id); + + // Save webhook secret to AppState + if let Some(secret) = &credentials.webhook_secret { + let mut inner = shared.state.inner.lock().await; + inner.next_runner_id += 0; // dummy access to keep compiler happy if needed + info!("Webhook secret registered: {}", secret); + } + + let credentials_html = format!( + r#" + + + GitHub App Registered + + +

GitHub App Registered Successfully!

+

App ID: {}

+

Webhook Secret: {}

+

Private Key PEM:

+
{}
+

To use this App, configure your local environment and restart `aksh`:

+
+export AKSH_WEBHOOK_SECRET="{}"
+export AKSH_GITHUB_APP_ID="{}"
+    
+ +"#, + credentials.id, + credentials.webhook_secret.as_deref().unwrap_or("none"), + credentials.pem, + credentials.webhook_secret.as_deref().unwrap_or("none"), + credentials.id + ); + + Ok(axum::response::Html(credentials_html)) +} diff --git a/crates/aksh-runner-server/src/lib.rs b/crates/aksh-runner-server/src/lib.rs index a7147eeb3..891b8d20f 100644 --- a/crates/aksh-runner-server/src/lib.rs +++ b/crates/aksh-runner-server/src/lib.rs @@ -5,6 +5,8 @@ use std::net::SocketAddr; use std::path::PathBuf; use std::sync::Arc; +pub mod github; + use aksh_artifacts::ArtifactStore; use aksh_cache::CacheStore; use aksh_gha_parser::{expand_jobs_with_reusables, parse_workflow}; @@ -45,13 +47,146 @@ pub struct ServerConfig { pub state_dir: PathBuf, } +async fn reap_once(shared: &Arc) { + let mut inner = shared.state.inner.lock().await; + let now = SystemTime::now(); + let mut cancellations = Vec::new(); + let mut disconnected_completions = Vec::new(); + + let mut active_reqs = Vec::new(); + for (request_id, request) in &inner.job_requests { + if request.result.is_none() { + active_reqs.push(( + *request_id, + request.run_id, + request.job_id.clone(), + request.started_at, + request.last_renewed_at, + request.timeout_triggered, + )); + } + } + + for (request_id, run_id, job_id, started_at, last_renewed_at, timeout_triggered) in active_reqs + { + // 1. Check Timeout Enforcement + if let Some(started_at) = started_at { + if !timeout_triggered { + let elapsed = now.duration_since(started_at).unwrap_or_default(); + let job_timeout = inner + .broker_messages + .get(&request_id) + .and_then(|msg| msg.job_timeout) + .unwrap_or(21600); // 360 minutes in seconds + + if elapsed >= Duration::from_secs(job_timeout as u64) { + info!( + %run_id, + %job_id, + request_id, + "Job timed out after {}s", + job_timeout + ); + if let Some(req) = inner.job_requests.get_mut(&request_id) { + req.timeout_triggered = true; + } + cancellations.push(QueuedCancellation { + run_id, + job_id: job_id.clone(), + }); + } + } + } + + // 2. Check Lease Expiration / Disconnect Reaper + if let Some(last_renewed_at) = last_renewed_at { + let elapsed = now.duration_since(last_renewed_at).unwrap_or_default(); + // 120 seconds disconnect threshold + if elapsed >= Duration::from_secs(120) { + info!( + %run_id, + %job_id, + request_id, + "Runner lease expired (last renewed {}s ago). Marking job as failed.", + elapsed.as_secs() + ); + if let Some(req) = inner.job_requests.get_mut(&request_id) { + req.result = Some(ExecutionStatus::Failure); + } + disconnected_completions.push(( + request_id, + JobCompletion { + run_id, + job_id: job_id.clone(), + status: ExecutionStatus::Failure, + outputs: Default::default(), + }, + )); + } + } + } + + // Cleanup session and inflight maps for disconnected runners + for (request_id, _) in &disconnected_completions { + inner.inflight_requests.remove(request_id); + inner + .session_active_requests + .retain(|_, &mut v| v != *request_id); + } + + // Apply cancellations + let cancellation_count = cancellations.len(); + if cancellation_count > 0 { + inner.cancellation_queue.extend(cancellations); + } + + drop(inner); + + // Notify if cancellations occurred + if cancellation_count > 0 { + shared.state.message_notify.notify_waiters(); + } + + // Process completions for disconnected runners + for (_, completion) in disconnected_completions { + let _ = complete_job_inner(shared.clone(), completion).await; + } +} + +async fn run_background_reaper(shared: Arc) { + let mut interval = tokio::time::interval(Duration::from_secs(10)); + // Skip the first tick + interval.tick().await; + + while !shared.shutdown.is_cancelled() { + tokio::select! { + _ = interval.tick() => { + reap_once(&shared).await; + } + _ = shared.shutdown.cancelled() => { + break; + } + } + } +} + /// Start the server and block until shutdown. pub async fn serve(config: ServerConfig) -> anyhow::Result<()> { let state = AppState::new(config.state_dir).await?; let shutdown = CancellationToken::new(); - let router = app(state, shutdown.clone()); + let router = app(state.clone(), shutdown.clone()); let listener = TcpListener::bind(config.listen).await?; + let shared = Arc::new(SharedState { + state, + shutdown: shutdown.clone(), + }); + + let checker_shared = shared.clone(); + tokio::spawn(async move { + run_background_reaper(checker_shared).await; + }); + info!(listen = %config.listen, "aksh runner server listening"); axum::serve(listener, router) .with_graceful_shutdown(shutdown_signal(shutdown)) @@ -311,6 +446,12 @@ pub fn app(state: AppState, shutdown: CancellationToken) -> Router { axum::routing::options(|| async { StatusCode::OK }), ) .route("/api/v1/runs", post(submit_run)) + .route( + "/api/v1/github/webhooks", + post(github::handle_github_webhook), + ) + .route("/api/v1/github/register", get(github::github_register)) + .route("/api/v1/github/callback", get(github::github_callback)) .route("/api/v1/runs/:run_id", get(get_run)) .route("/api/v1/runs/:run_id/cancel", post(cancel_run)) .route("/api/v1/runs/:run_id/rerun", post(rerun_run)) @@ -425,6 +566,10 @@ pub struct AppState { message_notify: Arc, cache: CacheStore, artifacts: ArtifactStore, + /// Optional GitHub App Webhook Secret for signature verification. + pub webhook_secret: Option, + /// Optional local workspace path to load workflows from. + pub local_workspace: Option, } impl AppState { @@ -439,12 +584,18 @@ impl AppState { agent_keypair: Some(keypair), ..Default::default() }; + let webhook_secret = std::env::var("AKSH_WEBHOOK_SECRET").ok(); + let local_workspace = std::env::var("AKSH_LOCAL_WORKSPACE") + .ok() + .map(PathBuf::from); Ok(Self { inner: Arc::new(Mutex::new(inner)), events, message_notify: Arc::new(Notify::new()), cache, artifacts, + webhook_secret, + local_workspace, }) } @@ -486,14 +637,16 @@ struct InnerState { } #[derive(Debug, Clone, Serialize)] -struct RunRecord { - run_id: RunId, - submission: WorkflowSubmission, - jobs: BTreeMap, - status: ExecutionStatus, - job_outputs: BTreeMap>, - job_base_ids: BTreeMap, - job_fail_fast: BTreeMap, +pub(crate) struct RunRecord { + pub(crate) run_id: RunId, + pub(crate) submission: WorkflowSubmission, + pub(crate) jobs: BTreeMap, + pub(crate) status: ExecutionStatus, + pub(crate) job_outputs: BTreeMap>, + pub(crate) job_base_ids: BTreeMap, + pub(crate) job_fail_fast: BTreeMap, + #[serde(default)] + pub(crate) job_check_run_ids: BTreeMap, } #[derive(Debug, Clone)] @@ -507,6 +660,9 @@ struct TaskAgentJobRequestRecord { timeline_id: uuid::Uuid, result: Option, locked_until: String, + started_at: Option, + last_renewed_at: Option, + timeout_triggered: bool, } #[derive(Debug, Clone)] @@ -617,10 +773,10 @@ async fn healthz(State(shared): State>) -> Json>, - Json(submission): Json, -) -> Result, ApiError> { +pub(crate) async fn submit_run_inner( + shared: &Arc, + submission: WorkflowSubmission, +) -> Result { let workflow = parse_workflow(&submission.workflow_yaml)?; let (branch, tag) = git_ref_context(&submission.git_ref); let changed_paths = changed_paths_from_payload(&submission.payload); @@ -694,6 +850,9 @@ async fn submit_run( timeline_id: agent_msg.timeline.id, result: None, locked_until: agent_request_locked_until(), + started_at: None, + last_renewed_at: None, + timeout_triggered: false, }; inner .plan_requests @@ -745,6 +904,7 @@ async fn submit_run( job_base_ids, job_fail_fast, status: ExecutionStatus::Queued, + job_check_run_ids: BTreeMap::new(), }, ); drop(inner); @@ -758,13 +918,20 @@ async fn submit_run( queued_jobs, }) .await; - Ok(Json(RunAccepted { + Ok(RunAccepted { run_id, queued_jobs, - })) + }) } } +async fn submit_run( + State(shared): State>, + Json(submission): Json, +) -> Result, ApiError> { + submit_run_inner(&shared, submission).await.map(Json) +} + fn git_ref_context(git_ref: &str) -> (Option, Option) { if let Some(branch) = git_ref.strip_prefix("refs/heads/") { (Some(branch.to_owned()), None) @@ -1138,8 +1305,35 @@ async fn next_message_broker_ref( loop { let mut inner = shared.state.inner.lock().await; + if let Some(message) = inner + .inflight_messages + .get(&session_id) + .and_then(|messages| messages.values().next().cloned()) + { + return Ok(Json(message).into_response()); + } + if let Some(request_id) = inner.session_active_requests.get(&session_id).copied() { if let Some(request) = inner.job_requests.get(&request_id) { + if let Some(pos) = inner + .cancellation_queue + .iter() + .position(|c| c.run_id == request.run_id && c.job_id == request.job_id) + { + let cancellation = inner.cancellation_queue.remove(pos).unwrap(); + let message = build_broker_plaintext_message( + &mut inner, + &session_id, + azdo::message_type::JOB_CANCELLED, + json!({ + "runId": cancellation.run_id.to_string(), + "jobId": cancellation.job_id.to_string(), + }) + .to_string(), + ); + return Ok(Json(message).into_response()); + } + if request.result.is_none() { return Ok(Json(broker_job_ref(request, pool_id)).into_response()); } @@ -1174,6 +1368,10 @@ async fn next_message_broker_ref( inner .session_active_requests .insert(session_id.clone(), request_id); + if let Some(request) = inner.job_requests.get_mut(&request_id) { + request.started_at = Some(std::time::SystemTime::now()); + request.last_renewed_at = Some(std::time::SystemTime::now()); + } inner .broker_messages .insert(request_id, queued.message.clone()); @@ -1187,6 +1385,8 @@ async fn next_message_broker_ref( let job_id = queued.job_id.clone(); drop(inner); + github::report_check_run_in_progress(&shared, run_id, &job_id).await; + shared .state .emit(NdjsonEvent::JobStatus { @@ -1357,6 +1557,7 @@ async fn broker_renew_job( .get_mut(&request_id) .ok_or_else(|| ApiError::not_found("agent request not found"))?; record.locked_until = agent_request_locked_until(); + record.last_renewed_at = Some(std::time::SystemTime::now()); Ok(Json(json!({"lockedUntil": record.locked_until}))) } @@ -1590,9 +1791,13 @@ async fn next_message( let body_json = serde_json::to_string(&queued.message) .map_err(|e| ApiError::bad_request(format!("failed to serialize job message: {e}")))?; + let request_id = queued.message.request_id; inner .session_active_requests - .insert(session_id.clone(), queued.message.request_id); + .insert(session_id.clone(), request_id); + if let Some(request) = inner.job_requests.get_mut(&request_id) { + request.started_at = Some(std::time::SystemTime::now()); + } let message = build_task_agent_message( &mut inner, &session_id, @@ -1604,6 +1809,8 @@ async fn next_message( let job_id = queued.job_id.clone(); drop(inner); + github::report_check_run_in_progress(&shared, run_id, &job_id).await; + shared .state .emit(NdjsonEvent::JobStatus { @@ -1659,6 +1866,28 @@ fn build_task_agent_message( Ok(message) } +fn build_broker_plaintext_message( + inner: &mut InnerState, + session_id: &str, + message_type: &str, + body_json: String, +) -> azdo::TaskAgentMessage { + inner.next_message_id += 1; + let message_id = inner.next_message_id; + let message = azdo::TaskAgentMessage { + message_id, + message_type: message_type.to_owned(), + body: body_json, + iv: None, + }; + inner + .inflight_messages + .entry(session_id.to_owned()) + .or_default() + .insert(message_id, message.clone()); + message +} + async fn delete_pool_message( State(shared): State>, Path((_pool_id, message_id)): Path<(i64, i64)>, @@ -1795,6 +2024,7 @@ async fn agent_request_patch( let mut inner = shared.state.inner.lock().await; if let Some(request) = inner.job_requests.get_mut(&request_id) { request.locked_until = agent_request_locked_until(); + request.last_renewed_at = Some(std::time::SystemTime::now()); } } Json(agent_request_response(&shared, pool_id, request_id).await) @@ -1931,6 +2161,15 @@ async fn complete_job_inner( .cloned() .ok_or_else(|| ApiError::not_found("run not found"))?; drop(inner); + + github::report_check_run_completed( + &shared, + completion.run_id, + &completion.job_id, + completion.status, + ) + .await; + if promoted_jobs > 0 || cancelled_siblings > 0 { shared.state.message_notify.notify_waiters(); } @@ -5798,4 +6037,398 @@ jobs: serde_json::from_slice(&bytes).unwrap() } } + + #[tokio::test] + async fn job_timeout_enforcement_cancels_job() { + let temp = tempfile::tempdir().unwrap(); + let state = AppState::new(temp.path().to_path_buf()).await.unwrap(); + let shutdown = CancellationToken::new(); + let app = app(state.clone(), shutdown.clone()); + let shared = Arc::new(SharedState { + state: state.clone(), + shutdown, + }); + + // 1. Submit run + let accepted = request_json( + &app, + Method::POST, + "/api/v1/runs", + json!({ + "workflow_yaml": "on: push\njobs:\n build:\n runs-on: ubuntu-latest\n steps:\n - run: sleep 10\n", + "event": "push", + "repository": "owner/repo" + }), + ) + .await; + let run_id: RunId = accepted["run_id"].as_str().unwrap().parse().unwrap(); + + // 2. Poll to start job (transitions status to InProgress and sets started_at) + let _msg = request_json( + &app, + Method::GET, + "/runner/server/_apis/v1/Message/1?sessionId=default", + Value::Null, + ) + .await; + + let request_id = { + let inner = state.inner.lock().await; + *inner.job_requests.keys().next().unwrap() + }; + + // 3. Override started_at to be in the past (beyond 360m/21600s default timeout) + { + let mut inner = state.inner.lock().await; + let request = inner.job_requests.get_mut(&request_id).unwrap(); + request.started_at = Some(SystemTime::now() - Duration::from_secs(22000)); + } + + // 4. Run reaper tick + reap_once(&shared).await; + + // 5. Verify cancellation is enqueued + { + let inner = state.inner.lock().await; + let request = inner.job_requests.get(&request_id).unwrap(); + assert!(request.timeout_triggered); + assert_eq!(inner.cancellation_queue.len(), 1); + assert_eq!(inner.cancellation_queue[0].run_id, run_id); + } + } + + #[tokio::test] + async fn runner_lease_expiration_disconnect_reaper() { + let temp = tempfile::tempdir().unwrap(); + let state = AppState::new(temp.path().to_path_buf()).await.unwrap(); + let shutdown = CancellationToken::new(); + let app = app(state.clone(), shutdown.clone()); + let shared = Arc::new(SharedState { + state: state.clone(), + shutdown, + }); + + // 1. Submit run + let accepted = request_json( + &app, + Method::POST, + "/api/v1/runs", + json!({ + "workflow_yaml": "on: push\njobs:\n build:\n runs-on: ubuntu-latest\n steps:\n - run: sleep 10\n", + "event": "push", + "repository": "owner/repo" + }), + ) + .await; + let run_id: RunId = accepted["run_id"].as_str().unwrap().parse().unwrap(); + + // 2. Poll to start job (sets last_renewed_at) + let _msg = request_json( + &app, + Method::GET, + "/runner/server/_apis/v1/Message/1?sessionId=default", + Value::Null, + ) + .await; + + let request_id = { + let inner = state.inner.lock().await; + *inner.job_requests.keys().next().unwrap() + }; + + // 3. Override last_renewed_at to be in the past (beyond 120s threshold) + { + let mut inner = state.inner.lock().await; + let request = inner.job_requests.get_mut(&request_id).unwrap(); + request.last_renewed_at = Some(SystemTime::now() - Duration::from_secs(130)); + } + + // 4. Run reaper tick + reap_once(&shared).await; + + // 5. Verify the job was marked failed and run completes as failed + { + let inner = state.inner.lock().await; + let request = inner.job_requests.get(&request_id).unwrap(); + assert_eq!(request.result, Some(ExecutionStatus::Failure)); + assert!(inner.inflight_requests.is_empty()); + assert!(inner.session_active_requests.is_empty()); + + let run = inner.runs.get(&run_id).unwrap(); + assert_eq!(run.status, ExecutionStatus::Failure); + } + } + + #[tokio::test] + async fn github_webhook_flows_with_signature_and_check_runs() { + let temp = tempfile::tempdir().unwrap(); + + // 1. Create a dummy workflow file in a local workspace + let ws_dir = temp.path().join("workspace"); + tokio::fs::create_dir_all(ws_dir.join(".github/workflows")) + .await + .unwrap(); + let workflow_content = r#" +on: push +jobs: + build: + runs-on: ubuntu-latest + steps: + - run: echo hello +"#; + tokio::fs::write(ws_dir.join(".github/workflows/build.yml"), workflow_content) + .await + .unwrap(); + + let mut state = AppState::new(temp.path().to_path_buf()).await.unwrap(); + state.webhook_secret = Some("super-secret".to_owned()); + state.local_workspace = Some(ws_dir.clone()); + + assert_eq!(state.webhook_secret.as_deref(), Some("super-secret")); + assert_eq!(state.local_workspace.as_ref(), Some(&ws_dir)); + + let app = app(state.clone(), CancellationToken::new()); + + // 2. Prepare mock webhook push payload + let payload = serde_json::json!({ + "ref": "refs/heads/main", + "before": "0000000000000000000000000000000000000000", + "after": "a1b2c3d4e5f6a1b2c3d4e5f6a1b2c3d4e5f6a1b2", + "repository": { + "full_name": "owner/repo", + "default_branch": "main" + }, + "commits": [ + { + "id": "a1b2c3d4e5f6a1b2c3d4e5f6a1b2c3d4e5f6a1b2", + "added": ["src/main.rs"], + "modified": [], + "removed": [] + } + ] + }); + + let payload_bytes = serde_json::to_vec(&payload).unwrap(); + + // 3. Compute correct signature + use hmac::{Hmac, Mac}; + use sha2::Sha256; + type HmacSha256 = Hmac; + let mut mac = HmacSha256::new_from_slice(b"super-secret").unwrap(); + mac.update(&payload_bytes); + let sig_bytes = mac.finalize().into_bytes(); + let sig_hex = sig_bytes + .iter() + .map(|b| format!("{:02x}", b)) + .collect::(); + let signature_header = format!("sha256={}", sig_hex); + + // 4. Send request with WRONG signature -> should fail with 401 + let response_401 = app + .clone() + .oneshot( + Request::builder() + .method(Method::POST) + .uri("/api/v1/github/webhooks") + .header("x-github-event", "push") + .header("x-hub-signature-256", "sha256=invalid") + .header("content-type", "application/json") + .body(Body::from(payload_bytes.clone())) + .unwrap(), + ) + .await + .unwrap(); + assert_eq!(response_401.status(), StatusCode::UNAUTHORIZED); + + // 5. Send request with CORRECT signature -> should succeed with 200 + let response_200 = app + .clone() + .oneshot( + Request::builder() + .method(Method::POST) + .uri("/api/v1/github/webhooks") + .header("x-github-event", "push") + .header("x-hub-signature-256", signature_header) + .header("content-type", "application/json") + .body(Body::from(payload_bytes)) + .unwrap(), + ) + .await + .unwrap(); + assert_eq!(response_200.status(), StatusCode::OK); + + // 6. Verify that a run was triggered and check runs are queued + let inner = state.inner.lock().await; + assert_eq!(inner.runs.len(), 1); + let (_, run_record) = inner.runs.iter().next().unwrap(); + assert_eq!(run_record.submission.event, "push"); + assert_eq!(run_record.submission.repository, "owner/repo"); + assert_eq!(run_record.submission.git_ref, "refs/heads/main"); + + // Verify that check_run_ids are created/queued in the record + assert_eq!(run_record.job_check_run_ids.len(), 1); + let (job_id, check_run_id) = run_record.job_check_run_ids.iter().next().unwrap(); + assert_eq!(job_id.to_string(), "build"); + assert!(*check_run_id > 0); + } + + #[tokio::test] + async fn github_webhook_pull_request_event() { + let temp = tempfile::tempdir().unwrap(); + + // Create a dummy workflow file in a local workspace + let ws_dir = temp.path().join("workspace"); + tokio::fs::create_dir_all(ws_dir.join(".github/workflows")) + .await + .unwrap(); + let workflow_content = r#" +on: pull_request +jobs: + test: + runs-on: ubuntu-latest + steps: + - run: make test +"#; + tokio::fs::write(ws_dir.join(".github/workflows/test.yml"), workflow_content) + .await + .unwrap(); + + let mut state = AppState::new(temp.path().to_path_buf()).await.unwrap(); + state.webhook_secret = Some("super-secret".to_owned()); + state.local_workspace = Some(ws_dir.clone()); + + let app = app(state.clone(), CancellationToken::new()); + + // Prepare PR payload + let payload = serde_json::json!({ + "action": "opened", + "number": 42, + "pull_request": { + "head": { + "ref": "feature-branch", + "sha": "b2c3d4e5f6a1b2c3d4e5f6a1b2c3d4e5f6a1b2c3" + }, + "base": { + "ref": "main", + "sha": "a1b2c3d4e5f6a1b2c3d4e5f6a1b2c3d4e5f6a1b2" + } + }, + "repository": { + "full_name": "owner/repo", + "default_branch": "main" + } + }); + + let payload_bytes = serde_json::to_vec(&payload).unwrap(); + + // Compute signature + use hmac::{Hmac, Mac}; + use sha2::Sha256; + type HmacSha256 = Hmac; + let mut mac = HmacSha256::new_from_slice(b"super-secret").unwrap(); + mac.update(&payload_bytes); + let sig_bytes = mac.finalize().into_bytes(); + let sig_hex = sig_bytes + .iter() + .map(|b| format!("{:02x}", b)) + .collect::(); + let signature_header = format!("sha256={}", sig_hex); + + let response = app + .oneshot( + Request::builder() + .method(Method::POST) + .uri("/api/v1/github/webhooks") + .header("x-github-event", "pull_request") + .header("x-hub-signature-256", signature_header) + .header("content-type", "application/json") + .body(Body::from(payload_bytes)) + .unwrap(), + ) + .await + .unwrap(); + assert_eq!(response.status(), StatusCode::OK); + + // Verify triggered run + let inner = state.inner.lock().await; + assert_eq!(inner.runs.len(), 1); + let (_, run_record) = inner.runs.iter().next().unwrap(); + assert_eq!(run_record.submission.event, "pull_request"); + assert_eq!(run_record.submission.git_ref, "refs/pull/42/merge"); + assert_eq!(run_record.job_check_run_ids.len(), 1); + } + + #[tokio::test] + async fn github_app_manifest_registration_flow() { + let temp = tempfile::tempdir().unwrap(); + + // 1. Setup a local mock GitHub API server for manifest conversion + let mock_app = Router::new().route( + "/app-manifests/:code/conversions", + post(|Path(code): Path| async move { + assert_eq!(code, "mock_code_123"); + Json(json!({ + "id": 987654, + "pem": "-----BEGIN RSA PRIVATE KEY-----\nMOCK-KEY-DATA\n-----END RSA PRIVATE KEY-----", + "webhook_secret": Some("mock-webhook-secret-xyz") + })) + }), + ); + + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let port = listener.local_addr().unwrap().port(); + tokio::spawn(async move { + axum::serve(listener, mock_app).await.unwrap(); + }); + + // 2. Configure mock API URL in environment + std::env::set_var("AKSH_GITHUB_API_URL", format!("http://127.0.0.1:{}", port)); + + let state = AppState::new(temp.path().to_path_buf()).await.unwrap(); + let app = app(state.clone(), CancellationToken::new()); + + // 3. Request registration form (GET /api/v1/github/register) + let response_reg = app + .clone() + .oneshot( + Request::builder() + .method(Method::GET) + .uri("/api/v1/github/register") + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + assert_eq!(response_reg.status(), StatusCode::OK); + let bytes = to_bytes(response_reg.into_body(), usize::MAX) + .await + .unwrap(); + let html = String::from_utf8(bytes.to_vec()).unwrap(); + assert!(html.contains("https://github.com/settings/apps/new")); + assert!(html.contains("aksh-local-app")); + + // 4. Request callback conversion (GET /api/v1/github/callback?code=mock_code_123) + let response_callback = app + .clone() + .oneshot( + Request::builder() + .method(Method::GET) + .uri("/api/v1/github/callback?code=mock_code_123") + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + assert_eq!(response_callback.status(), StatusCode::OK); + let bytes_callback = to_bytes(response_callback.into_body(), usize::MAX) + .await + .unwrap(); + let html_callback = String::from_utf8(bytes_callback.to_vec()).unwrap(); + assert!(html_callback.contains("GitHub App Registered Successfully!")); + assert!(html_callback.contains("987654")); + assert!(html_callback.contains("mock-webhook-secret-xyz")); + + // Clean up + std::env::remove_var("AKSH_GITHUB_API_URL"); + } } diff --git a/docs/github-app-webhook.md b/docs/github-app-webhook.md new file mode 100644 index 000000000..cc6838e46 --- /dev/null +++ b/docs/github-app-webhook.md @@ -0,0 +1,135 @@ +# GitHub App Webhook Integration Log & User Guide + +This document records the design, build log, and interaction guide for the end-to-end GitHub App Webhook receiver and Checks API status reporting system in `aksh`. + +--- + +## 1. Webhook Architecture Overview + +The webhook system enables `aksh` to receive push and pull_request notifications directly from GitHub, fetch matching workflow files, and queue jobs for self-hosted runners. It also integrates with the GitHub Checks API to report status back to the repository. + +### Data Flow diagram: + +```mermaid +graph TD + GH[GitHub Webhook / API] -->|1. Event payload & Signatures| Wh[Webhook Receiver] + Wh -->|2. Event Type & SHAs| Auth[App Auth Client] + Auth -->|JWT / App Private Key| GH + GH -->|Installation Token| Fetch[Workflow Fetcher] + Fetch -->|3. Fetch .github/workflows/*.yml| AST[Workflow Evaluator] + AST -->|Matches? -> submit_run_inner| Core[aksh Control Plane] + Core -->|4. Job InProgress/Done| Report[Checks Reporter] + Report -->|Checks API / Commit Status| GH +``` + +--- + +## 2. Configuration Parameters + +The system is configured using the following environment variables: + +| Variable | Description | Example | +|---|---|---| +| `AKSH_WEBHOOK_SECRET` | Secret key configured on the GitHub App to verify payload signatures. | `my-secure-webhook-secret` | +| `AKSH_LOCAL_WORKSPACE` | Path to a local clone of the repository to fetch workflows from offline. | `/path/to/my-repo` | +| `AKSH_GITHUB_TOKEN` | GitHub Personal Access Token or App Installation Token to fetch workflows and update check runs. | `ghp_...` or `ghs_...` | + +### Security Best Practices + +* **Git-Ignore Credentials**: Never check `.env`, `*.pem`, or `*.key` files into Git. These files are excluded in the root `.gitignore`. + +* **Production Key Management**: In production, do not write private keys or secrets to plaintext files on the server disk. Instead: + - Load them directly into memory at runtime using a Secrets Manager (e.g. HashiCorp Vault, AWS Secrets Manager, or Kubernetes Secrets). + - Inject the configuration as environment variables directly to the running process without file-based middleware. + +--- + +## 3. Webhook signature verification (`X-Hub-Signature-256`) + + + +`aksh` verifies that incoming webhooks are authentic: +- When `AKSH_WEBHOOK_SECRET` is set, `aksh` computes the HMAC-SHA256 signature of the raw request body and verifies it against the `x-hub-signature-256` header. +- If verification fails or the header is missing, the endpoint returns `401 Unauthorized`. +- If `AKSH_WEBHOOK_SECRET` is not configured, signature checking is skipped, enabling easier local testing. + +--- + +## 4. Workflow Fetching Strategies + +When a push or PR webhook is received, `aksh` retrieves the workflow definitions: +1. **Local Filesystem (Offline/Dev Mode)**: + If `AKSH_LOCAL_WORKSPACE` is configured, `aksh` reads the `.github/workflows/` directory directly from that local path. +2. **GitHub API (Remote/Production Mode)**: + If `AKSH_LOCAL_WORKSPACE` is not configured, but `AKSH_GITHUB_TOKEN` is set, `aksh` queries: + `GET /repos/{owner}/{repo}/contents/.github/workflows?ref={git_ref}` + And downloads files dynamically. +3. **Current Directory Fallback**: + If neither is set, it defaults to looking in the local `.github/workflows/` directory of the current running workspace. + +--- + +## 5. GitHub Checks API Reporting + +The system maps the lifecycle of each job to a GitHub Check Run: +1. **Queued**: When a run is accepted, a check run is created via `POST /repos/{owner}/{repo}/check-runs` with status `queued`. The check run ID is recorded in `RunRecord.job_check_run_ids`. +2. **In Progress**: When the runner fetches and starts the job, the status is updated to `in_progress`. +3. **Completed**: When the runner finishes (or `aksh` reaps it due to timeout/lease expiration), the check run is updated to `completed` with the corresponding conclusion (`success`, `failure`, or `cancelled`). + +If `AKSH_GITHUB_TOKEN` is not configured, these requests are simulated in-memory and logged to the console, allowing fully offline execution. + +--- + +## 6. How Users Interact with it + +### Step 1: Set up Webhook in GitHub +1. Go to your GitHub App or Repository settings. +2. Set the payload URL to `http:///api/v1/github/webhooks`. +3. Set the content type to `application/json`. +4. Enter a secure Webhook Secret (e.g. `super-secret`). +5. Select the **Push** and **Pull Request** events. + +### Step 2: Start `aksh-runner-server` +Run the server with the environment variables set: +```sh +export AKSH_WEBHOOK_SECRET="super-secret" +export AKSH_LOCAL_WORKSPACE="/Users/bnjoroge/runner-watcher" +export AKSH_GITHUB_TOKEN="ghp_optional_token_for_checks" + +just serve +``` + +### Step 3: Trigger workflows +Push a commit or open a pull request. `aksh` will: +- Receive the webhook event. +- Fetch the workflows. +- Match filters (branches, tags, paths). +- Queue jobs for any registered runners matching `runs-on` labels. +- Create check runs on GitHub. + +--- + +## 7. Automated One-Click App Registration (Manifest Flow) + +To simplify local development and testing, `aksh` supports the official **GitHub App Manifest** flow: + +1. **Open the Registration page**: + Start `aksh-runner-server` and navigate to: + `http://localhost:9090/api/v1/github/register` + +2. **Click "Register App on GitHub"**: + You will be redirected to GitHub to register the app under your personal account or organization with all required permissions and webhook events pre-configured. + +3. **Callback Conversion**: + After clicking "Create", GitHub redirects back to: + `http://localhost:9090/api/v1/github/callback?code=...` + + `aksh` exchanges the temporary code for your new App ID, Webhook Secret, and Private Key PEM, displaying them directly on-screen and logging them to the terminal. + +4. **Save and Restart**: + Copy the displayed credentials into your local environment: + ```sh + export AKSH_WEBHOOK_SECRET="your-new-webhook-secret" + export AKSH_GITHUB_APP_ID="your-new-app-id" + ``` + And restart `aksh`!