From 408359f365838f0393d91db640cf8903071369c4 Mon Sep 17 00:00:00 2001 From: Trezy Date: Wed, 01 Jul 2026 18:46:02 +0000 Subject: [PATCH] feat: add support for async jobs Signed-off-by: Trezy --- src/lib.rs | 1 + src/main.rs | 9 +++++++++ tests/e2e_jobs.rs | 491 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ web/playwright.config.ts | 2 ++ migrations/postgres/20260701000000_create_jobs.sql | 17 +++++++++++++++++ migrations/postgres/20260702000000_add_inherit_auth_to_jobs.sql | 1 + migrations/sqlite/20260701000000_create_jobs.sql | 17 +++++++++++++++++ migrations/sqlite/20260702000000_add_inherit_auth_to_jobs.sql | 1 + src/admin/jobs.rs | 124 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ src/admin/mod.rs | 6 ++++++ src/admin/permissions.rs | 34 ++++++++++++++++++++++++++++++++++ src/jobs/db.rs | 299 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ src/jobs/mod.rs | 20 ++++++++++++++++++++ src/jobs/worker.rs | 340 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ src/lua/execute.rs | 26 ++++++++++++++++++++++++++ src/lua/jobs_api.rs | 297 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ src/lua/mod.rs | 9 +++++---- src/lua/scripts.rs | 2 ++ tests/common/db.rs | 3 ++- web/src/components/app-sidebar.tsx | 7 +++++++ web/src/lib/api.ts | 33 +++++++++++++++++++++++++++++++++ web/src/types/jobs.ts | 18 ++++++++++++++++++ web/src/types/scripts.ts | 32 +++++++++++++++++++++++++++++++- web/tests/e2e/jobs.spec.ts | 216 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ web/tests/e2e/script-job.spec.ts | 90 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ packages/docs/content/blog/happyview-2.10.md | 6 +++--- packages/docs/content/blog/happyview-2.9.md | 8 ++++---- web/src/app/dashboard/jobs/page.tsx | 515 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ web/src/app/dashboard/settings/scripts/script-form.tsx | 196 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--------------------------------------- web/src/app/dashboard/settings/scripts/[id]/script-detail.tsx | 32 +++++++++++++++++++++++++++----- web/src/app/dashboard/settings/scripts/new/page.tsx | 147 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++----------------------- 31 file(s) changed, 2919 insertion(s)(+), 80 deletion(s)(-) diff --git a/src/lib.rs b/src/lib.rs --- a/src/lib.rs +++ b/src/lib.rs @@ -14,6 +14,7 @@ pub mod feature_middleware; pub mod http_retry; pub mod jetstream; +pub mod jobs; pub mod labeler; pub mod lexicon; pub mod lua; diff --git a/src/main.rs b/src/main.rs --- a/src/main.rs +++ b/src/main.rs @@ -681,6 +681,15 @@ happyview::admin::backfill::resume_backfill_jobs(&state).await; + // Resume interrupted jobs and start the job worker + happyview::jobs::worker::resume_interrupted_jobs(&state).await; + { + let job_state = state.clone(); + tokio::spawn(async move { + happyview::jobs::worker::run_worker(job_state).await; + }); + } + { let state = state.clone(); tokio::spawn(async move { diff --git a/tests/e2e_jobs.rs b/tests/e2e_jobs.rs new file mode 100644 --- /dev/null +++ b/tests/e2e_jobs.rs @@ -0,0 +1,491 @@ +mod common; + +use axum::body::Body; +use axum::http::{Request, StatusCode}; +use happyview::db::adapt_sql; +use http_body_util::BodyExt; +use serde_json::{Value, json}; +use serial_test::serial; +use tower::ServiceExt; +use uuid::Uuid; + +use common::app::TestApp; + +async fn json_body(resp: axum::response::Response) -> Value { + let body = resp.into_body().collect().await.unwrap().to_bytes(); + serde_json::from_slice(&body).unwrap() +} + +fn admin_get( + uri: &str, + cookie: (axum::http::HeaderName, axum::http::HeaderValue), +) -> Request { + Request::builder() + .uri(uri) + .header(cookie.0, cookie.1) + .body(Body::empty()) + .unwrap() +} + +fn admin_post( + uri: &str, + cookie: (axum::http::HeaderName, axum::http::HeaderValue), + body: &Value, +) -> Request { + Request::builder() + .method("POST") + .uri(uri) + .header(cookie.0, cookie.1) + .header("content-type", "application/json") + .body(Body::from(serde_json::to_vec(body).unwrap())) + .unwrap() +} + +async fn seed_job(app: &TestApp, job_type: &str, status: &str) -> String { + let id = Uuid::new_v4().to_string(); + let now = happyview::db::now_rfc3339(); + let input = serde_json::to_string(&json!({"test": true})).unwrap(); + + let sql = adapt_sql( + "INSERT INTO happyview_jobs (id, job_type, status, input, progress, created_by, created_at) VALUES (?, ?, ?, ?, '{}', ?, ?)", + app.state.db_backend, + ); + sqlx::query(&sql) + .bind(&id) + .bind(job_type) + .bind(status) + .bind(&input) + .bind(&app.admin_did) + .bind(&now) + .execute(&app.state.db) + .await + .expect("seed_job: insert failed"); + + id +} + +async fn set_job_status(app: &TestApp, id: &str, status: &str) { + let sql = adapt_sql( + "UPDATE happyview_jobs SET status = ? WHERE id = ?", + app.state.db_backend, + ); + sqlx::query(&sql) + .bind(status) + .bind(id) + .execute(&app.state.db) + .await + .expect("set_job_status failed"); +} + +// --------------------------------------------------------------------------- +// List jobs +// --------------------------------------------------------------------------- + +#[tokio::test] +#[serial] +async fn list_jobs_empty() { + common::require_db!(); + let app = TestApp::new().await; + + let resp = app + .router + .clone() + .oneshot(admin_get("/admin/jobs", app.admin_cookie())) + .await + .unwrap(); + + assert_eq!(resp.status(), StatusCode::OK); + let body = json_body(resp).await; + assert_eq!(body["jobs"].as_array().unwrap().len(), 0); + assert_eq!(body["cursor"], Value::Null); +} + +#[tokio::test] +#[serial] +async fn list_jobs_returns_seeded_jobs() { + common::require_db!(); + let app = TestApp::new().await; + + seed_job(&app, "test.export", "pending").await; + seed_job(&app, "test.import", "running").await; + + let resp = app + .router + .clone() + .oneshot(admin_get("/admin/jobs", app.admin_cookie())) + .await + .unwrap(); + + assert_eq!(resp.status(), StatusCode::OK); + let body = json_body(resp).await; + let jobs = body["jobs"].as_array().unwrap(); + assert_eq!(jobs.len(), 2); +} + +#[tokio::test] +#[serial] +async fn list_jobs_filters_by_status() { + common::require_db!(); + let app = TestApp::new().await; + + seed_job(&app, "test.export", "pending").await; + seed_job(&app, "test.import", "running").await; + seed_job(&app, "test.cleanup", "completed").await; + + let resp = app + .router + .clone() + .oneshot(admin_get("/admin/jobs?status=running", app.admin_cookie())) + .await + .unwrap(); + + assert_eq!(resp.status(), StatusCode::OK); + let body = json_body(resp).await; + let jobs = body["jobs"].as_array().unwrap(); + assert_eq!(jobs.len(), 1); + assert_eq!(jobs[0]["job_type"], "test.import"); + assert_eq!(jobs[0]["status"], "running"); +} + +// --------------------------------------------------------------------------- +// Get job +// --------------------------------------------------------------------------- + +#[tokio::test] +#[serial] +async fn get_job_returns_details() { + common::require_db!(); + let app = TestApp::new().await; + + let id = seed_job(&app, "test.export", "pending").await; + + let resp = app + .router + .clone() + .oneshot(admin_get(&format!("/admin/jobs/{id}"), app.admin_cookie())) + .await + .unwrap(); + + assert_eq!(resp.status(), StatusCode::OK); + let body = json_body(resp).await; + assert_eq!(body["id"], id); + assert_eq!(body["job_type"], "test.export"); + assert_eq!(body["status"], "pending"); + assert_eq!(body["input"]["test"], true); +} + +#[tokio::test] +#[serial] +async fn get_job_not_found() { + common::require_db!(); + let app = TestApp::new().await; + + let fake_id = Uuid::new_v4(); + let resp = app + .router + .clone() + .oneshot(admin_get( + &format!("/admin/jobs/{fake_id}"), + app.admin_cookie(), + )) + .await + .unwrap(); + + assert_eq!(resp.status(), StatusCode::NOT_FOUND); +} + +// --------------------------------------------------------------------------- +// Cancel job +// --------------------------------------------------------------------------- + +#[tokio::test] +#[serial] +async fn cancel_pending_job_sets_cancelled() { + common::require_db!(); + let app = TestApp::new().await; + + let id = seed_job(&app, "test.export", "pending").await; + + let resp = app + .router + .clone() + .oneshot(admin_post( + &format!("/admin/jobs/{id}/cancel"), + app.admin_cookie(), + &json!({}), + )) + .await + .unwrap(); + + assert_eq!(resp.status(), StatusCode::OK); + let body = json_body(resp).await; + assert_eq!(body["status"], "cancelled"); +} + +#[tokio::test] +#[serial] +async fn cancel_running_job_sets_cancelling() { + common::require_db!(); + let app = TestApp::new().await; + + let id = seed_job(&app, "test.export", "running").await; + + let resp = app + .router + .clone() + .oneshot(admin_post( + &format!("/admin/jobs/{id}/cancel"), + app.admin_cookie(), + &json!({}), + )) + .await + .unwrap(); + + assert_eq!(resp.status(), StatusCode::OK); + let body = json_body(resp).await; + assert_eq!(body["status"], "cancelling"); +} + +#[tokio::test] +#[serial] +async fn cancel_paused_job_sets_cancelled() { + common::require_db!(); + let app = TestApp::new().await; + + let id = seed_job(&app, "test.export", "paused").await; + + let resp = app + .router + .clone() + .oneshot(admin_post( + &format!("/admin/jobs/{id}/cancel"), + app.admin_cookie(), + &json!({}), + )) + .await + .unwrap(); + + assert_eq!(resp.status(), StatusCode::OK); + let body = json_body(resp).await; + assert_eq!(body["status"], "cancelled"); +} + +#[tokio::test] +#[serial] +async fn cancel_completed_job_returns_400() { + common::require_db!(); + let app = TestApp::new().await; + + let id = seed_job(&app, "test.export", "completed").await; + + let resp = app + .router + .clone() + .oneshot(admin_post( + &format!("/admin/jobs/{id}/cancel"), + app.admin_cookie(), + &json!({}), + )) + .await + .unwrap(); + + assert_eq!(resp.status(), StatusCode::BAD_REQUEST); +} + +// --------------------------------------------------------------------------- +// Pause job +// --------------------------------------------------------------------------- + +#[tokio::test] +#[serial] +async fn pause_running_job_sets_pausing() { + common::require_db!(); + let app = TestApp::new().await; + + let id = seed_job(&app, "test.export", "running").await; + + let resp = app + .router + .clone() + .oneshot(admin_post( + &format!("/admin/jobs/{id}/pause"), + app.admin_cookie(), + &json!({}), + )) + .await + .unwrap(); + + assert_eq!(resp.status(), StatusCode::OK); + let body = json_body(resp).await; + assert_eq!(body["status"], "pausing"); +} + +#[tokio::test] +#[serial] +async fn pause_pending_job_returns_400() { + common::require_db!(); + let app = TestApp::new().await; + + let id = seed_job(&app, "test.export", "pending").await; + + let resp = app + .router + .clone() + .oneshot(admin_post( + &format!("/admin/jobs/{id}/pause"), + app.admin_cookie(), + &json!({}), + )) + .await + .unwrap(); + + assert_eq!(resp.status(), StatusCode::BAD_REQUEST); +} + +// --------------------------------------------------------------------------- +// Resume job +// --------------------------------------------------------------------------- + +#[tokio::test] +#[serial] +async fn resume_paused_job_sets_pending() { + common::require_db!(); + let app = TestApp::new().await; + + let id = seed_job(&app, "test.export", "paused").await; + + let resp = app + .router + .clone() + .oneshot(admin_post( + &format!("/admin/jobs/{id}/resume"), + app.admin_cookie(), + &json!({}), + )) + .await + .unwrap(); + + assert_eq!(resp.status(), StatusCode::OK); + let body = json_body(resp).await; + assert_eq!(body["status"], "pending"); +} + +#[tokio::test] +#[serial] +async fn resume_running_job_returns_400() { + common::require_db!(); + let app = TestApp::new().await; + + let id = seed_job(&app, "test.export", "running").await; + + let resp = app + .router + .clone() + .oneshot(admin_post( + &format!("/admin/jobs/{id}/resume"), + app.admin_cookie(), + &json!({}), + )) + .await + .unwrap(); + + assert_eq!(resp.status(), StatusCode::BAD_REQUEST); +} + +// --------------------------------------------------------------------------- +// Auth: unauthenticated requests +// --------------------------------------------------------------------------- + +#[tokio::test] +#[serial] +async fn list_jobs_without_auth_returns_401() { + common::require_db!(); + let app = TestApp::new().await; + + let resp = app + .router + .clone() + .oneshot( + Request::builder() + .uri("/admin/jobs") + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + + assert_eq!(resp.status(), StatusCode::UNAUTHORIZED); +} + +// --------------------------------------------------------------------------- +// Full lifecycle: pending → running → pausing → paused → pending → cancel +// --------------------------------------------------------------------------- + +#[tokio::test] +#[serial] +async fn full_job_lifecycle() { + common::require_db!(); + let app = TestApp::new().await; + + let id = seed_job(&app, "test.lifecycle", "pending").await; + + // Simulate worker claiming → running + set_job_status(&app, &id, "running").await; + + // Pause the running job + let resp = app + .router + .clone() + .oneshot(admin_post( + &format!("/admin/jobs/{id}/pause"), + app.admin_cookie(), + &json!({}), + )) + .await + .unwrap(); + assert_eq!(resp.status(), StatusCode::OK); + assert_eq!(json_body(resp).await["status"], "pausing"); + + // Simulate worker acknowledging pause + set_job_status(&app, &id, "paused").await; + + // Resume the paused job + let resp = app + .router + .clone() + .oneshot(admin_post( + &format!("/admin/jobs/{id}/resume"), + app.admin_cookie(), + &json!({}), + )) + .await + .unwrap(); + assert_eq!(resp.status(), StatusCode::OK); + assert_eq!(json_body(resp).await["status"], "pending"); + + // Cancel the pending job + let resp = app + .router + .clone() + .oneshot(admin_post( + &format!("/admin/jobs/{id}/cancel"), + app.admin_cookie(), + &json!({}), + )) + .await + .unwrap(); + assert_eq!(resp.status(), StatusCode::OK); + assert_eq!(json_body(resp).await["status"], "cancelled"); + + // Verify final state + let resp = app + .router + .clone() + .oneshot(admin_get(&format!("/admin/jobs/{id}"), app.admin_cookie())) + .await + .unwrap(); + assert_eq!(resp.status(), StatusCode::OK); + let job = json_body(resp).await; + assert_eq!(job["status"], "cancelled"); + assert!(job["completed_at"].is_string()); +} diff --git a/web/playwright.config.ts b/web/playwright.config.ts --- a/web/playwright.config.ts +++ b/web/playwright.config.ts @@ -30,9 +30,11 @@ "lexicon-services.spec.ts", "lexicon-delete.spec.ts", "script-delete.spec.ts", + "script-job.spec.ts", "record-delete.spec.ts", "proxy-config.spec.ts", "spaces.spec.ts", + "jobs.spec.ts", ], dependencies: ["setup"], use: { browserName: "chromium" }, diff --git a/migrations/postgres/20260701000000_create_jobs.sql b/migrations/postgres/20260701000000_create_jobs.sql new file mode 100644 --- /dev/null +++ b/migrations/postgres/20260701000000_create_jobs.sql @@ -0,0 +1,17 @@ +CREATE TABLE happyview_jobs ( + id TEXT PRIMARY KEY, + job_type TEXT NOT NULL, + status TEXT NOT NULL DEFAULT 'pending', + input TEXT NOT NULL DEFAULT '{}', + progress TEXT NOT NULL DEFAULT '{}', + result TEXT, + error TEXT, + created_by TEXT NOT NULL, + started_at TEXT, + completed_at TEXT, + created_at TEXT NOT NULL +); + +CREATE INDEX idx_happyview_jobs_status ON happyview_jobs (status); +CREATE INDEX idx_happyview_jobs_job_type ON happyview_jobs (job_type); +CREATE INDEX idx_happyview_jobs_created_by ON happyview_jobs (created_by); diff --git a/migrations/postgres/20260702000000_add_inherit_auth_to_jobs.sql b/migrations/postgres/20260702000000_add_inherit_auth_to_jobs.sql new file mode 100644 --- /dev/null +++ b/migrations/postgres/20260702000000_add_inherit_auth_to_jobs.sql @@ -0,0 +1,1 @@ +ALTER TABLE happyview_jobs ADD COLUMN inherit_auth BOOLEAN NOT NULL DEFAULT FALSE; diff --git a/migrations/sqlite/20260701000000_create_jobs.sql b/migrations/sqlite/20260701000000_create_jobs.sql new file mode 100644 --- /dev/null +++ b/migrations/sqlite/20260701000000_create_jobs.sql @@ -0,0 +1,17 @@ +CREATE TABLE happyview_jobs ( + id TEXT PRIMARY KEY, + job_type TEXT NOT NULL, + status TEXT NOT NULL DEFAULT 'pending', + input TEXT NOT NULL DEFAULT '{}', + progress TEXT NOT NULL DEFAULT '{}', + result TEXT, + error TEXT, + created_by TEXT NOT NULL, + started_at TEXT, + completed_at TEXT, + created_at TEXT NOT NULL DEFAULT (datetime('now')) +); + +CREATE INDEX idx_happyview_jobs_status ON happyview_jobs (status); +CREATE INDEX idx_happyview_jobs_job_type ON happyview_jobs (job_type); +CREATE INDEX idx_happyview_jobs_created_by ON happyview_jobs (created_by); diff --git a/migrations/sqlite/20260702000000_add_inherit_auth_to_jobs.sql b/migrations/sqlite/20260702000000_add_inherit_auth_to_jobs.sql new file mode 100644 --- /dev/null +++ b/migrations/sqlite/20260702000000_add_inherit_auth_to_jobs.sql @@ -0,0 +1,1 @@ +ALTER TABLE happyview_jobs ADD COLUMN inherit_auth BOOLEAN NOT NULL DEFAULT 0; diff --git a/src/admin/jobs.rs b/src/admin/jobs.rs new file mode 100644 --- /dev/null +++ b/src/admin/jobs.rs @@ -0,0 +1,124 @@ +use axum::Json; +use axum::extract::{Path, Query, State}; +use serde::Deserialize; + +use crate::AppState; +use crate::error::AppError; +use crate::jobs; + +use super::auth::UserAuth; +use super::permissions::Permission; + +#[derive(Deserialize)] +pub struct ListJobsQuery { + pub status: Option, + pub limit: Option, + pub cursor: Option, +} + +pub async fn list_jobs( + State(state): State, + auth: UserAuth, + Query(query): Query, +) -> Result, AppError> { + auth.require(Permission::JobsRead).await?; + + let limit = query.limit.unwrap_or(50).min(100); + let (jobs_list, cursor) = jobs::db::list_jobs( + &state, + query.status.as_deref(), + limit, + query.cursor.as_deref(), + ) + .await?; + + Ok(Json(serde_json::json!({ + "jobs": jobs_list, + "cursor": cursor, + }))) +} + +pub async fn get_job( + State(state): State, + auth: UserAuth, + Path(id): Path, +) -> Result, AppError> { + auth.require(Permission::JobsRead).await?; + + let job = jobs::db::get_job(&state, &id) + .await? + .ok_or_else(|| AppError::NotFound("job not found".into()))?; + + Ok(Json(serde_json::to_value(job).unwrap())) +} + +pub async fn cancel_job( + State(state): State, + auth: UserAuth, + Path(id): Path, +) -> Result, AppError> { + auth.require(Permission::JobsManage).await?; + + let job = jobs::db::get_job(&state, &id) + .await? + .ok_or_else(|| AppError::NotFound("job not found".into()))?; + + match job.status.as_str() { + "running" => { + jobs::db::set_status(&state, &id, "cancelling").await?; + Ok(Json(serde_json::json!({ "status": "cancelling" }))) + } + "pending" | "paused" => { + jobs::db::set_status(&state, &id, "cancelled").await?; + Ok(Json(serde_json::json!({ "status": "cancelled" }))) + } + _ => Err(AppError::BadRequest(format!( + "cannot cancel job with status: {}", + job.status + ))), + } +} + +pub async fn pause_job( + State(state): State, + auth: UserAuth, + Path(id): Path, +) -> Result, AppError> { + auth.require(Permission::JobsManage).await?; + + let job = jobs::db::get_job(&state, &id) + .await? + .ok_or_else(|| AppError::NotFound("job not found".into()))?; + + if job.status != "running" { + return Err(AppError::BadRequest(format!( + "cannot pause job with status: {}", + job.status + ))); + } + + jobs::db::set_status(&state, &id, "pausing").await?; + Ok(Json(serde_json::json!({ "status": "pausing" }))) +} + +pub async fn resume_job( + State(state): State, + auth: UserAuth, + Path(id): Path, +) -> Result, AppError> { + auth.require(Permission::JobsManage).await?; + + let job = jobs::db::get_job(&state, &id) + .await? + .ok_or_else(|| AppError::NotFound("job not found".into()))?; + + if job.status != "paused" { + return Err(AppError::BadRequest(format!( + "cannot resume job with status: {}", + job.status + ))); + } + + jobs::db::set_status(&state, &id, "pending").await?; + Ok(Json(serde_json::json!({ "status": "pending" }))) +} diff --git a/src/admin/mod.rs b/src/admin/mod.rs --- a/src/admin/mod.rs +++ b/src/admin/mod.rs @@ -6,6 +6,7 @@ mod domains; mod events; mod feature_flags; +mod jobs; mod labelers; mod lexicons; mod network_lexicons; @@ -62,6 +63,11 @@ "/backfill/{id}/details", delete(backfill::flush_backfill_details), ) + .route("/jobs", get(jobs::list_jobs)) + .route("/jobs/{id}", get(jobs::get_job)) + .route("/jobs/{id}/cancel", post(jobs::cancel_job)) + .route("/jobs/{id}/pause", post(jobs::pause_job)) + .route("/jobs/{id}/resume", post(jobs::resume_job)) .route("/events", get(events::list_events)) .route("/users", post(users::create_user).get(users::list_users)) .route("/users/transfer-super", post(users::transfer_super)) diff --git a/src/admin/permissions.rs b/src/admin/permissions.rs --- a/src/admin/permissions.rs +++ b/src/admin/permissions.rs @@ -113,6 +113,13 @@ ScriptsRead, #[serde(rename = "scripts:manage")] ScriptsManage, + + #[serde(rename = "jobs:read")] + JobsRead, + #[serde(rename = "jobs:create")] + JobsCreate, + #[serde(rename = "jobs:manage")] + JobsManage, } impl Permission { @@ -162,6 +169,9 @@ Self::SpacesManageCredentials => "spaces:manage-credentials", Self::ScriptsRead => "scripts:read", Self::ScriptsManage => "scripts:manage", + Self::JobsRead => "jobs:read", + Self::JobsCreate => "jobs:create", + Self::JobsManage => "jobs:manage", } } @@ -425,6 +435,24 @@ description: "Create, update, and delete trigger-keyed scripts", category: "Scripts", }, + Self::JobsRead => PermissionInfo { + key: "jobs:read", + name: "View Jobs", + description: "View background job status and progress", + category: "Jobs", + }, + Self::JobsCreate => PermissionInfo { + key: "jobs:create", + name: "Create Jobs", + description: "Queue new background jobs", + category: "Jobs", + }, + Self::JobsManage => PermissionInfo { + key: "jobs:manage", + name: "Manage Jobs", + description: "Cancel, pause, and resume background jobs", + category: "Jobs", + }, } } @@ -474,6 +502,9 @@ Self::SpacesManageCredentials, Self::ScriptsRead, Self::ScriptsManage, + Self::JobsRead, + Self::JobsCreate, + Self::JobsManage, ]) } } @@ -525,6 +556,9 @@ SpacesManageInvites, SpacesManageRecords, SpacesManageCredentials, + JobsRead, + JobsCreate, + JobsManage, ] .iter() .map(|p| p.info()) diff --git a/src/jobs/db.rs b/src/jobs/db.rs new file mode 100644 --- /dev/null +++ b/src/jobs/db.rs @@ -0,0 +1,299 @@ +use serde_json::Value; +use uuid::Uuid; + +use crate::AppState; +use crate::db::{adapt_sql, now_rfc3339}; +use crate::error::AppError; + +use super::Job; + +type JobRow = ( + String, + String, + String, + String, + String, + Option, + Option, + String, + Option, + Option, + String, + bool, +); + +fn row_to_job( + ( + id, + job_type, + status, + input, + progress, + result, + error, + created_by, + started_at, + completed_at, + created_at, + inherit_auth, + ): JobRow, +) -> Job { + Job { + id, + job_type, + status, + input: serde_json::from_str(&input).unwrap_or(Value::Null), + progress: serde_json::from_str(&progress).unwrap_or(Value::Null), + result: result.and_then(|r| serde_json::from_str(&r).ok()), + error, + created_by, + started_at, + completed_at, + created_at, + inherit_auth, + } +} + +pub async fn create_job( + state: &AppState, + job_type: &str, + input: &Value, + created_by: &str, + inherit_auth: bool, +) -> Result { + let id = Uuid::new_v4().to_string(); + let now = now_rfc3339(); + let input_str = serde_json::to_string(input) + .map_err(|e| AppError::Internal(format!("failed to serialize job input: {e}")))?; + + let sql = adapt_sql( + "INSERT INTO happyview_jobs (id, job_type, status, input, created_by, created_at, inherit_auth) VALUES (?, ?, 'pending', ?, ?, ?, ?)", + state.db_backend, + ); + sqlx::query(&sql) + .bind(&id) + .bind(job_type) + .bind(&input_str) + .bind(created_by) + .bind(&now) + .bind(inherit_auth) + .execute(&state.db) + .await + .map_err(|e| AppError::Internal(format!("failed to create job: {e}")))?; + + Ok(id) +} + +pub async fn get_job(state: &AppState, id: &str) -> Result, AppError> { + let sql = adapt_sql( + "SELECT * FROM happyview_jobs WHERE id = ?", + state.db_backend, + ); + let row: Option = sqlx::query_as(&sql) + .bind(id) + .fetch_optional(&state.db) + .await + .map_err(|e| AppError::Internal(format!("failed to fetch job: {e}")))?; + + Ok(row.map(row_to_job)) +} + +pub async fn list_jobs( + state: &AppState, + status_filter: Option<&str>, + limit: i64, + cursor: Option<&str>, +) -> Result<(Vec, Option), AppError> { + let sql = if status_filter.is_some() { + let base = if cursor.is_some() { + "SELECT * FROM happyview_jobs WHERE status = ? AND created_at < ? ORDER BY created_at DESC LIMIT ?" + } else { + "SELECT * FROM happyview_jobs WHERE status = ? ORDER BY created_at DESC LIMIT ?" + }; + adapt_sql(base, state.db_backend) + } else { + let base = if cursor.is_some() { + "SELECT * FROM happyview_jobs WHERE created_at < ? ORDER BY created_at DESC LIMIT ?" + } else { + "SELECT * FROM happyview_jobs ORDER BY created_at DESC LIMIT ?" + }; + adapt_sql(base, state.db_backend) + }; + + let mut query = sqlx::query_as::<_, JobRow>(&sql); + + if let Some(status) = status_filter { + query = query.bind(status); + } + if let Some(cursor) = cursor { + query = query.bind(cursor); + } + query = query.bind(limit + 1); + + let rows = query + .fetch_all(&state.db) + .await + .map_err(|e| AppError::Internal(format!("failed to list jobs: {e}")))?; + + let has_more = rows.len() as i64 > limit; + let jobs: Vec = rows + .into_iter() + .take(limit as usize) + .map(row_to_job) + .collect(); + + let next_cursor = if has_more { + jobs.last().map(|j| j.created_at.clone()) + } else { + None + }; + + Ok((jobs, next_cursor)) +} + +pub async fn set_status(state: &AppState, id: &str, status: &str) -> Result<(), AppError> { + let now = now_rfc3339(); + let sql = match status { + "running" => adapt_sql( + "UPDATE happyview_jobs SET status = ?, started_at = ? WHERE id = ?", + state.db_backend, + ), + "completed" | "failed" | "cancelled" => adapt_sql( + "UPDATE happyview_jobs SET status = ?, completed_at = ? WHERE id = ?", + state.db_backend, + ), + _ => adapt_sql( + "UPDATE happyview_jobs SET status = ? WHERE id = ? AND 1=1", + state.db_backend, + ), + }; + + match status { + "running" | "completed" | "failed" | "cancelled" => { + sqlx::query(&sql) + .bind(status) + .bind(&now) + .bind(id) + .execute(&state.db) + .await + .map_err(|e| AppError::Internal(format!("failed to update job status: {e}")))?; + } + _ => { + let sql = adapt_sql( + "UPDATE happyview_jobs SET status = ? WHERE id = ?", + state.db_backend, + ); + sqlx::query(&sql) + .bind(status) + .bind(id) + .execute(&state.db) + .await + .map_err(|e| AppError::Internal(format!("failed to update job status: {e}")))?; + } + } + + Ok(()) +} + +pub async fn update_progress(state: &AppState, id: &str, progress: &Value) -> Result<(), AppError> { + let progress_str = serde_json::to_string(progress) + .map_err(|e| AppError::Internal(format!("failed to serialize progress: {e}")))?; + let sql = adapt_sql( + "UPDATE happyview_jobs SET progress = ? WHERE id = ?", + state.db_backend, + ); + sqlx::query(&sql) + .bind(&progress_str) + .bind(id) + .execute(&state.db) + .await + .map_err(|e| AppError::Internal(format!("failed to update job progress: {e}")))?; + Ok(()) +} + +pub async fn set_result(state: &AppState, id: &str, result: &Value) -> Result<(), AppError> { + let result_str = serde_json::to_string(result) + .map_err(|e| AppError::Internal(format!("failed to serialize result: {e}")))?; + let now = now_rfc3339(); + let sql = adapt_sql( + "UPDATE happyview_jobs SET status = 'completed', result = ?, completed_at = ? WHERE id = ?", + state.db_backend, + ); + sqlx::query(&sql) + .bind(&result_str) + .bind(&now) + .bind(id) + .execute(&state.db) + .await + .map_err(|e| AppError::Internal(format!("failed to set job result: {e}")))?; + Ok(()) +} + +pub async fn set_error(state: &AppState, id: &str, error: &str) -> Result<(), AppError> { + let now = now_rfc3339(); + let sql = adapt_sql( + "UPDATE happyview_jobs SET status = 'failed', error = ?, completed_at = ? WHERE id = ?", + state.db_backend, + ); + sqlx::query(&sql) + .bind(error) + .bind(&now) + .bind(id) + .execute(&state.db) + .await + .map_err(|e| AppError::Internal(format!("failed to set job error: {e}")))?; + Ok(()) +} + +/// Check if a job should stop (status changed to cancelling or pausing). +/// Same cooperative cancellation pattern as the backfill system. +pub async fn should_stop(state: &AppState, id: &str) -> Option<&'static str> { + let sql = adapt_sql( + "SELECT status FROM happyview_jobs WHERE id = ?", + state.db_backend, + ); + let status = sqlx::query_as::<_, (String,)>(&sql) + .bind(id) + .fetch_optional(&state.db) + .await + .ok() + .flatten() + .map(|(s,)| s); + match status.as_deref() { + Some("cancelling") => Some("cancelling"), + Some("pausing") => Some("pausing"), + _ => None, + } +} + +/// Find jobs that were interrupted by a server restart. +pub async fn find_interrupted_jobs(state: &AppState) -> Vec { + let sql = adapt_sql( + "SELECT * FROM happyview_jobs WHERE status IN ('running', 'cancelling', 'pausing')", + state.db_backend, + ); + let rows: Vec = sqlx::query_as(&sql) + .fetch_all(&state.db) + .await + .unwrap_or_default(); + + rows.into_iter().map(row_to_job).collect() +} + +/// Pick the next pending job and atomically set it to running. +pub async fn claim_next_job(state: &AppState) -> Result, AppError> { + let now = now_rfc3339(); + + let sql = adapt_sql( + "UPDATE happyview_jobs SET status = 'running', started_at = ? WHERE id = (SELECT id FROM happyview_jobs WHERE status = 'pending' ORDER BY created_at ASC LIMIT 1) RETURNING *", + state.db_backend, + ); + + let row: Option = sqlx::query_as(&sql) + .bind(&now) + .fetch_optional(&state.db) + .await + .map_err(|e| AppError::Internal(format!("failed to claim job: {e}")))?; + + Ok(row.map(row_to_job)) +} diff --git a/src/jobs/mod.rs b/src/jobs/mod.rs new file mode 100644 --- /dev/null +++ b/src/jobs/mod.rs @@ -0,0 +1,20 @@ +pub(crate) mod db; +pub mod worker; + +use serde::{Deserialize, Serialize}; + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct Job { + pub id: String, + pub job_type: String, + pub status: String, + pub input: serde_json::Value, + pub progress: serde_json::Value, + pub result: Option, + pub error: Option, + pub created_by: String, + pub started_at: Option, + pub completed_at: Option, + pub created_at: String, + pub inherit_auth: bool, +} diff --git a/src/jobs/worker.rs b/src/jobs/worker.rs new file mode 100644 --- /dev/null +++ b/src/jobs/worker.rs @@ -0,0 +1,340 @@ +use std::sync::Arc; +use std::time::Duration; + +use mlua::LuaSerdeExt; + +use crate::AppState; +use crate::db::adapt_sql; +use crate::event_log::{EventLog, Severity, log_event}; +use crate::lua::{sandbox, scripts}; +use crate::repo; + +use super::db; + +const POLL_INTERVAL: Duration = Duration::from_secs(5); + +/// Start the background job worker. Polls for pending jobs and +/// executes them one at a time. +pub async fn run_worker(state: AppState) { + tracing::info!("job worker started"); + + loop { + match db::claim_next_job(&state).await { + Ok(Some(job)) => { + tracing::info!(job_id = %job.id, job_type = %job.job_type, "executing job"); + execute_job(&state, &job).await; + } + Ok(None) => { + tokio::time::sleep(POLL_INTERVAL).await; + } + Err(e) => { + tracing::error!(error = %e, "job worker: failed to claim job"); + tokio::time::sleep(POLL_INTERVAL).await; + } + } + } +} + +/// Resume jobs that were interrupted by a server restart. +pub async fn resume_interrupted_jobs(state: &AppState) { + let jobs = db::find_interrupted_jobs(state).await; + + for job in jobs { + match job.status.as_str() { + "cancelling" => { + tracing::info!(job_id = %job.id, "finalising cancelled job from previous run"); + let _ = db::set_status(state, &job.id, "cancelled").await; + } + "pausing" => { + tracing::info!(job_id = %job.id, "finalising paused job from previous run"); + let _ = db::set_status(state, &job.id, "paused").await; + } + "running" => { + tracing::info!(job_id = %job.id, "re-queuing interrupted job"); + let _ = db::set_status(state, &job.id, "pending").await; + } + _ => {} + } + } +} + +async fn execute_job(state: &AppState, job: &super::Job) { + let backend = state.db_backend; + + log_event( + &state.db, + EventLog { + event_type: "job.started".to_string(), + severity: Severity::Info, + actor_did: Some(job.created_by.clone()), + subject: Some(job.job_type.clone()), + detail: serde_json::json!({ + "job_id": job.id, + "job_type": job.job_type, + }), + }, + backend, + ) + .await; + + let trigger_id = format!("job.run:{}", job.job_type); + let script = match scripts::resolve(state, &trigger_id).await { + Some(s) => s, + None => { + let error = format!("no script found for trigger: {trigger_id}"); + tracing::error!(job_id = %job.id, %error); + let _ = db::set_error(state, &job.id, &error).await; + log_event( + &state.db, + EventLog { + event_type: "job.failed".to_string(), + severity: Severity::Error, + actor_did: Some(job.created_by.clone()), + subject: Some(job.job_type.clone()), + detail: serde_json::json!({ + "job_id": job.id, + "error": error, + }), + }, + backend, + ) + .await; + return; + } + }; + + let (claims, pds_auth_arc) = if job.inherit_auth { + let pds_auth = match repo::get_oauth_session(state, &job.created_by).await { + Ok(session) => repo::PdsAuth::OAuth(Arc::new(session)), + Err(e) => { + let error = format!("failed to obtain PDS auth for {}: {e}", job.created_by); + tracing::error!(job_id = %job.id, %error); + let _ = db::set_error(state, &job.id, &error).await; + log_event( + &state.db, + EventLog { + event_type: "job.failed".to_string(), + severity: Severity::Error, + actor_did: Some(job.created_by.clone()), + subject: Some(job.job_type.clone()), + detail: serde_json::json!({ + "job_id": job.id, + "error": error, + }), + }, + backend, + ) + .await; + return; + } + }; + ( + Some(Arc::new(crate::auth::Claims::internal( + job.created_by.clone(), + ))), + Some(Arc::new(pds_auth)), + ) + } else { + (None, None) + }; + + let lua = match sandbox::create_sandbox() { + Ok(l) => l, + Err(e) => { + let error = format!("failed to create Lua VM: {e}"); + let _ = db::set_error(state, &job.id, &error).await; + return; + } + }; + + lua.remove_hook(); + + let state_arc = Arc::new(state.clone()); + + if let Err(e) = crate::lua::db_api::register_db_api(&lua, state_arc.clone()) { + let _ = db::set_error(state, &job.id, &format!("db api: {e}")).await; + return; + } + if let Err(e) = crate::lua::http_api::register_http_api(&lua, state_arc.clone()) { + let _ = db::set_error(state, &job.id, &format!("http api: {e}")).await; + return; + } + if let Err(e) = crate::lua::xrpc_api::register_xrpc_api( + &lua, + state_arc.clone(), + Some(job.created_by.clone()), + ) { + let _ = db::set_error(state, &job.id, &format!("xrpc api: {e}")).await; + return; + } + if let Err(e) = crate::lua::atproto_api::register_atproto_api( + &lua, + state_arc.clone(), + Some(&job.created_by), + ) { + let _ = db::set_error(state, &job.id, &format!("atproto api: {e}")).await; + return; + } + if let (Some(c), Some(p)) = (&claims, &pds_auth_arc) + && let Err(e) = crate::lua::atproto_api::register_atproto_blob_api( + &lua, + state_arc.clone(), + c.clone(), + p.clone(), + ) + { + let _ = db::set_error(state, &job.id, &format!("blob api: {e}")).await; + return; + } + if let Err(e) = crate::lua::jobs_api::register_jobs_api( + &lua, + state_arc.clone(), + Some(job.created_by.clone()), + ) { + let _ = db::set_error(state, &job.id, &format!("jobs api: {e}")).await; + return; + } + if let Err(e) = + crate::lua::record::register_record_api(&lua, state_arc.clone(), claims, pds_auth_arc, None) + { + let _ = db::set_error(state, &job.id, &format!("record api: {e}")).await; + return; + } + if let Err(e) = crate::lua::scripts::register_log_event_api( + &lua, + &state_arc, + &trigger_id, + Some(&job.created_by), + ) { + let _ = db::set_error(state, &job.id, &format!("log api: {e}")).await; + return; + } + if let Err(e) = crate::lua::jobs_api::register_job_context( + &lua, + state_arc.clone(), + job.id.clone(), + job.input.clone(), + ) { + let _ = db::set_error(state, &job.id, &format!("job context: {e}")).await; + return; + } + + let env_vars = load_env_vars(&state.db, backend).await; + if let Err(e) = crate::lua::context::set_env_context(&lua, &env_vars) { + let _ = db::set_error(state, &job.id, &format!("env context: {e}")).await; + return; + } + + if let Err(e) = lua.globals().set("caller_did", job.created_by.as_str()) { + let _ = db::set_error(state, &job.id, &format!("caller_did: {e}")).await; + return; + } + + if let Err(e) = lua.load(script.body.as_str()).exec() { + let error = format!("script load failed: {e}"); + let _ = db::set_error(state, &job.id, &error).await; + return; + } + + let handle: mlua::Function = match lua.globals().get("handle") { + Ok(f) => f, + Err(e) => { + let _ = db::set_error(state, &job.id, &format!("missing handle(): {e}")).await; + return; + } + }; + + match handle.call_async::(()).await { + Ok(result) => { + let json_result: serde_json::Value = + lua.from_value(result).unwrap_or(serde_json::json!(null)); + + match db::should_stop(state, &job.id).await { + Some("pausing") => { + let _ = db::set_status(state, &job.id, "paused").await; + tracing::info!(job_id = %job.id, "job paused"); + log_event( + &state.db, + EventLog { + event_type: "job.paused".to_string(), + severity: Severity::Info, + actor_did: Some(job.created_by.clone()), + subject: Some(job.job_type.clone()), + detail: serde_json::json!({ "job_id": job.id }), + }, + backend, + ) + .await; + } + Some("cancelling") => { + let _ = db::set_status(state, &job.id, "cancelled").await; + tracing::info!(job_id = %job.id, "job cancelled"); + log_event( + &state.db, + EventLog { + event_type: "job.cancelled".to_string(), + severity: Severity::Info, + actor_did: Some(job.created_by.clone()), + subject: Some(job.job_type.clone()), + detail: serde_json::json!({ "job_id": job.id }), + }, + backend, + ) + .await; + } + _ => { + let _ = db::set_result(state, &job.id, &json_result).await; + tracing::info!(job_id = %job.id, "job completed"); + log_event( + &state.db, + EventLog { + event_type: "job.completed".to_string(), + severity: Severity::Info, + actor_did: Some(job.created_by.clone()), + subject: Some(job.job_type.clone()), + detail: serde_json::json!({ + "job_id": job.id, + "result": json_result, + }), + }, + backend, + ) + .await; + } + } + } + Err(e) => { + let error = format!("{e}"); + tracing::error!(job_id = %job.id, %error, "job script failed"); + let _ = db::set_error(state, &job.id, &error).await; + log_event( + &state.db, + EventLog { + event_type: "job.failed".to_string(), + severity: Severity::Error, + actor_did: Some(job.created_by.clone()), + subject: Some(job.job_type.clone()), + detail: serde_json::json!({ + "job_id": job.id, + "error": error, + }), + }, + backend, + ) + .await; + } + } +} + +async fn load_env_vars( + db: &sqlx::AnyPool, + backend: crate::db::DatabaseBackend, +) -> std::collections::HashMap { + let sql = adapt_sql("SELECT key, value FROM happyview_script_variables", backend); + sqlx::query_as::<_, (String, String)>(&sql) + .fetch_all(db) + .await + .unwrap_or_default() + .into_iter() + .collect() +} diff --git a/src/lua/execute.rs b/src/lua/execute.rs --- a/src/lua/execute.rs +++ b/src/lua/execute.rs @@ -268,6 +268,32 @@ return Err(AppError::Internal(error_message)); } + if let Err(e) = + super::jobs_api::register_jobs_api(&lua, state_arc.clone(), Some(claims.did().to_string())) + { + let error_message = format!("failed to register jobs API: {e}"); + log_event( + &state.db, + EventLog { + event_type: "script.error".to_string(), + severity: Severity::Error, + actor_did: Some(claims.did().to_string()), + subject: Some(method.to_string()), + detail: serde_json::json!({ + "error": error_message, + "script_source": script_source, + "input": input_json, + "caller_did": claims.did(), + "method": method, + "duration_ms": start.elapsed().as_millis() as u64, + }), + }, + backend, + ) + .await; + return Err(AppError::Internal(error_message)); + } + if let Err(e) = record::register_record_api( &lua, state_arc.clone(), diff --git a/src/lua/jobs_api.rs b/src/lua/jobs_api.rs new file mode 100644 --- /dev/null +++ b/src/lua/jobs_api.rs @@ -0,0 +1,297 @@ +use mlua::{Lua, LuaSerdeExt, Result as LuaResult}; +use std::sync::Arc; + +use crate::AppState; +use crate::jobs; + +/// Register the `jobs` table for queuing jobs from scripts. +/// Available in all script contexts (procedure, query, record-event). +pub fn register_jobs_api( + lua: &Lua, + state: Arc, + caller_did: Option, +) -> LuaResult<()> { + let jobs_table = lua.create_table()?; + + // jobs.create(job_type, input[, opts]) -> job_id string + // opts.auth: boolean (default false) — inherit caller's PDS auth + { + let state = state.clone(); + let caller_did = caller_did.clone(); + let create_fn = lua.create_async_function( + move |lua, (job_type, input, opts): (String, mlua::Value, Option)| { + let state = state.clone(); + let caller_did = caller_did.clone(); + + let input_json: serde_json::Value = + lua.from_value(input).unwrap_or(serde_json::json!({})); + + let inherit_auth = opts + .and_then(|t| t.get::("auth").ok()) + .unwrap_or(false); + + async move { + let caller = caller_did.as_deref().ok_or_else(|| { + mlua::Error::runtime("jobs.create requires an authenticated caller") + })?; + + let job_id = + jobs::db::create_job(&state, &job_type, &input_json, caller, inherit_auth) + .await + .map_err(|e| { + mlua::Error::runtime(format!("jobs.create failed: {e}")) + })?; + + Ok(job_id) + } + }, + )?; + jobs_table.set("create", create_fn)?; + } + + lua.globals().set("jobs", jobs_table)?; + Ok(()) +} + +/// Register the `job` context table for use inside job scripts. +/// Provides access to job input, progress reporting, cooperative +/// cancellation, and sleep/wait. +/// +/// Called by the job worker, not by the normal script execution path. +pub fn register_job_context( + lua: &Lua, + state: Arc, + job_id: String, + input: serde_json::Value, +) -> LuaResult<()> { + let job_table = lua.create_table()?; + + // job.input — the JSONB input passed to jobs.create() + let input_value = lua.to_value(&input)?; + job_table.set("input", input_value)?; + + // job.id — the job's UUID + job_table.set("id", job_id.clone())?; + + // job.progress(data) — persist progress to DB + { + let state = state.clone(); + let job_id = job_id.clone(); + let progress_fn = lua.create_async_function(move |lua, data: mlua::Value| { + let state = state.clone(); + let job_id = job_id.clone(); + let json_data: serde_json::Value = + lua.from_value(data).unwrap_or(serde_json::json!({})); + async move { + jobs::db::update_progress(&state, &job_id, &json_data) + .await + .map_err(|e| mlua::Error::runtime(format!("job.progress failed: {e}")))?; + Ok(()) + } + })?; + job_table.set("progress", progress_fn)?; + } + + // job.should_stop() -> boolean + { + let state = state.clone(); + let job_id = job_id.clone(); + let should_stop_fn = lua.create_async_function(move |_lua, ()| { + let state = state.clone(); + let job_id = job_id.clone(); + async move { + let result = jobs::db::should_stop(&state, &job_id).await; + Ok(result.is_some()) + } + })?; + job_table.set("should_stop", should_stop_fn)?; + } + + // job.wait(seconds) — yield execution for the given duration + { + let wait_fn = lua.create_async_function(move |_lua, seconds: f64| async move { + let duration = std::time::Duration::from_secs_f64(seconds.clamp(0.0, 3600.0)); + tokio::time::sleep(duration).await; + Ok(()) + })?; + job_table.set("wait", wait_fn)?; + } + + lua.globals().set("job", job_table)?; + Ok(()) +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::config::Config; + use crate::db::DatabaseBackend; + use crate::lexicon::LexiconRegistry; + use tokio::sync::watch; + + fn test_state() -> AppState { + let config = Config { + host: "127.0.0.1".into(), + port: 3000, + database_url: String::new(), + database_backend: crate::db::DatabaseBackend::Sqlite, + public_url: String::new(), + session_secret: "test-secret".into(), + jetstream_url: String::new(), + relay_url: String::new(), + plc_url: String::new(), + static_dir: String::new(), + base_path: None, + event_log_retention_days: 30, + app_name: None, + logo_uri: None, + tos_uri: None, + policy_uri: None, + token_encryption_key: None, + default_rate_limit_capacity: 100, + default_rate_limit_refill_rate: 2.0, + }; + let (tx, _) = watch::channel(vec![]); + let (labeler_tx, _) = watch::channel(()); + sqlx::any::install_default_drivers(); + let test_db = sqlx::AnyPool::connect_lazy("sqlite::memory:").unwrap(); + let atrium_http = std::sync::Arc::new(atrium_oauth::DefaultHttpClient::default()); + let did_resolver = atrium_identity::did::CommonDidResolver::new( + atrium_identity::did::CommonDidResolverConfig { + plc_directory_url: "https://plc.directory".into(), + http_client: std::sync::Arc::clone(&atrium_http), + }, + ); + let handle_resolver = atrium_identity::handle::AtprotoHandleResolver::new( + atrium_identity::handle::AtprotoHandleResolverConfig { + dns_txt_resolver: crate::dns::NativeDnsResolver::new(), + http_client: atrium_http, + }, + ); + let oauth = atrium_oauth::OAuthClient::new(atrium_oauth::OAuthClientConfig { + client_metadata: atrium_oauth::AtprotoLocalhostClientMetadata { + redirect_uris: Some(vec!["http://127.0.0.1:0/auth/callback".into()]), + scopes: Some(vec![atrium_oauth::Scope::Known( + atrium_oauth::KnownScope::Atproto, + )]), + }, + keys: None, + state_store: crate::auth::oauth_store::DbStateStore::new( + test_db.clone(), + crate::db::DatabaseBackend::Sqlite, + ), + session_store: crate::auth::oauth_store::DbSessionStore::new( + test_db.clone(), + crate::db::DatabaseBackend::Sqlite, + ), + resolver: atrium_oauth::OAuthResolverConfig { + did_resolver, + handle_resolver, + authorization_server_metadata: Default::default(), + protected_resource_metadata: Default::default(), + }, + }) + .expect("Failed to create test OAuth client"); + AppState { + config, + http: reqwest::Client::new(), + db: test_db.clone(), + backfill_db: test_db.clone(), + db_backend: DatabaseBackend::Sqlite, + domain_cache: crate::domain::DomainCache::new(), + lexicons: LexiconRegistry::new(), + collections_tx: tx, + labeler_subscriptions_tx: labeler_tx, + rate_limiter: crate::rate_limit::RateLimiter::new( + crate::rate_limit::RateLimitDefaults { + query_cost: 1, + procedure_cost: 1, + proxy_cost: 1, + }, + ), + oauth: std::sync::Arc::new(crate::auth::OAuthClientRegistry::new(std::sync::Arc::new( + oauth, + ))), + oauth_state_store: crate::auth::oauth_store::DbStateStore::new( + test_db.clone(), + crate::db::DatabaseBackend::Sqlite, + ), + cookie_key: axum_extra::extract::cookie::Key::derive_from( + b"test-secret-for-tests-only-not-production", + ), + plugin_registry: std::sync::Arc::new(crate::plugin::PluginRegistry::new()), + wasm_runtime: std::sync::Arc::new( + crate::plugin::WasmRuntime::new().expect("wasm runtime"), + ), + attestation_signer: None, + official_registry: std::sync::Arc::new(tokio::sync::RwLock::new( + crate::plugin::official_registry::OfficialRegistryState::default(), + )), + official_registry_config: crate::plugin::official_registry::RegistryConfig::production( + ), + proxy_config: std::sync::Arc::new(arc_swap::ArcSwap::new(std::sync::Arc::new( + crate::proxy_config::ProxyConfig::default(), + ))), + backfill_events_tx: tokio::sync::broadcast::channel(16).0, + verbose_event_logging: std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false)), + } + } + + #[tokio::test] + async fn jobs_api_is_registered() { + let lua = crate::lua::sandbox::create_sandbox().unwrap(); + let state = test_state(); + register_jobs_api(&lua, Arc::new(state), Some("did:plc:test".into())).unwrap(); + + let has_create: bool = lua + .load("return type(jobs.create) == 'function'") + .eval_async() + .await + .unwrap(); + assert!(has_create); + } + + #[tokio::test] + async fn job_context_exposes_input() { + let lua = crate::lua::sandbox::create_sandbox().unwrap(); + let state = test_state(); + let input = serde_json::json!({ "game_uri": "at://did:plc:test/game/123" }); + register_job_context(&lua, Arc::new(state), "test-job-id".into(), input).unwrap(); + + let game_uri: String = lua + .load("return job.input.game_uri") + .eval_async() + .await + .unwrap(); + assert_eq!(game_uri, "at://did:plc:test/game/123"); + + let job_id: String = lua.load("return job.id").eval_async().await.unwrap(); + assert_eq!(job_id, "test-job-id"); + } + + #[tokio::test] + async fn job_context_has_required_functions() { + let lua = crate::lua::sandbox::create_sandbox().unwrap(); + let state = test_state(); + register_job_context( + &lua, + Arc::new(state), + "test-id".into(), + serde_json::json!({}), + ) + .unwrap(); + + let result: bool = lua + .load( + r#" + return type(job.progress) == 'function' + and type(job.should_stop) == 'function' + and type(job.wait) == 'function' + "#, + ) + .eval_async() + .await + .unwrap(); + assert!(result); + } +} diff --git a/src/lua/mod.rs b/src/lua/mod.rs --- a/src/lua/mod.rs +++ b/src/lua/mod.rs @@ -1,13 +1,14 @@ -mod atproto_api; -mod context; +pub(crate) mod atproto_api; +pub(crate) mod context; pub mod db_api; mod execute; -mod http_api; +pub(crate) mod http_api; +pub(crate) mod jobs_api; pub mod record; pub(crate) mod sandbox; pub mod scripts; pub(crate) mod tid; -mod xrpc_api; +pub(crate) mod xrpc_api; #[allow(unused_imports)] pub(crate) use context::SpaceContext; diff --git a/src/lua/scripts.rs b/src/lua/scripts.rs --- a/src/lua/scripts.rs +++ b/src/lua/scripts.rs @@ -725,6 +725,8 @@ .map_err(|e| format!("xrpc api: {e}"))?; atproto_api::register_atproto_api(lua, state.clone(), None) .map_err(|e| format!("atproto api: {e}"))?; + super::jobs_api::register_jobs_api(lua, state.clone(), caller_did.map(String::from)) + .map_err(|e| format!("jobs api: {e}"))?; record::register_record_api_no_auth(lua, state.clone()) .map_err(|e| format!("record api: {e}"))?; register_log_event_api(lua, state, trigger_id, caller_did)?; diff --git a/tests/common/db.rs b/tests/common/db.rs --- a/tests/common/db.rs +++ b/tests/common/db.rs @@ -48,7 +48,7 @@ match backend { DatabaseBackend::Postgres => { sqlx::query( - "TRUNCATE happyview_records, happyview_lexicons, happyview_backfill_jobs, happyview_users, happyview_user_permissions, happyview_api_keys, happyview_event_logs, happyview_script_variables, happyview_scripts, happyview_dead_letter_scripts, happyview_dead_letter_hooks, happyview_record_refs, happyview_labeler_subscriptions, happyview_labels, happyview_instance_settings, happyview_domains, happyview_dpop_sessions, happyview_dpop_keys, happyview_api_clients, happyview_delegated_accounts, happyview_account_delegates, happyview_service_identity, happyview_service_entries, happyview_service_entry_xrpcs RESTART IDENTITY CASCADE", + "TRUNCATE happyview_records, happyview_lexicons, happyview_backfill_jobs, happyview_users, happyview_user_permissions, happyview_api_keys, happyview_event_logs, happyview_script_variables, happyview_scripts, happyview_dead_letter_scripts, happyview_dead_letter_hooks, happyview_record_refs, happyview_labeler_subscriptions, happyview_labels, happyview_instance_settings, happyview_domains, happyview_dpop_sessions, happyview_dpop_keys, happyview_api_clients, happyview_delegated_accounts, happyview_account_delegates, happyview_service_identity, happyview_service_entries, happyview_service_entry_xrpcs, happyview_jobs RESTART IDENTITY CASCADE", ) .execute(pool) .await @@ -80,6 +80,7 @@ "happyview_labels", "happyview_instance_settings", "happyview_domains", + "happyview_jobs", ]; for table in tables { sqlx::query(&format!("DELETE FROM {table}")) diff --git a/web/src/components/app-sidebar.tsx b/web/src/components/app-sidebar.tsx --- a/web/src/components/app-sidebar.tsx +++ b/web/src/components/app-sidebar.tsx @@ -22,6 +22,7 @@ IconSkull, IconFlask, IconFingerprint, + IconPlayerPlay, } from "@tabler/icons-react"; import Image from "next/image"; import Link from "next/link"; @@ -59,6 +60,12 @@ { title: "Lexicons", url: "/dashboard/lexicons", icon: IconFileDescription }, { title: "Records", url: "/dashboard/records", icon: IconTable }, { title: "Backfill", url: "/dashboard/backfill", icon: IconDatabase }, + { + title: "Jobs", + url: "/dashboard/jobs", + icon: IconPlayerPlay, + requiredPermissions: ["jobs:read"], + }, { title: "Dead Letters", url: "/dashboard/dead-letters", diff --git a/web/src/lib/api.ts b/web/src/lib/api.ts --- a/web/src/lib/api.ts +++ b/web/src/lib/api.ts @@ -3,6 +3,7 @@ import type { LexiconSummary, LexiconDetail } from "@/types/lexicons"; import type { NetworkLexiconSummary } from "@/types/network-lexicons"; import type { BackfillJob, BackfillReposResponse, PdsSummaryResponse } from "@/types/backfill"; +import type { Job, JobsListResponse } from "@/types/jobs"; import type { UserSummary } from "@/types/users"; import type { AdminListRecordsResponse } from "@/types/records"; import type { EventsListResponse } from "@/types/events"; @@ -252,6 +253,38 @@ export function flushAllBackfillDetails() { return apiFetch(`/admin/backfill/details`, { method: "DELETE" }); +} + +// Jobs +export function getJobs(params: { status?: string; limit?: number; cursor?: string } = {}) { + const qs = new URLSearchParams(); + if (params.status) qs.set("status", params.status); + if (params.limit) qs.set("limit", String(params.limit)); + if (params.cursor) qs.set("cursor", params.cursor); + const query = qs.toString(); + return apiFetch(`/admin/jobs${query ? `?${query}` : ""}`); +} + +export function getJob(id: string) { + return apiFetch(`/admin/jobs/${id}`); +} + +export function cancelJob(id: string) { + return apiFetch<{ status: string }>(`/admin/jobs/${id}/cancel`, { + method: "POST", + }); +} + +export function pauseJob(id: string) { + return apiFetch<{ status: string }>(`/admin/jobs/${id}/pause`, { + method: "POST", + }); +} + +export function resumeJob(id: string) { + return apiFetch<{ status: string }>(`/admin/jobs/${id}/resume`, { + method: "POST", + }); } // Users diff --git a/web/src/types/jobs.ts b/web/src/types/jobs.ts new file mode 100644 --- /dev/null +++ b/web/src/types/jobs.ts @@ -0,0 +1,18 @@ +export interface Job { + id: string; + job_type: string; + status: string; + input: Record; + progress: Record; + result: Record | null; + error: string | null; + created_by: string; + started_at: string | null; + completed_at: string | null; + created_at: string; +} + +export interface JobsListResponse { + jobs: Job[]; + cursor: string | null; +} diff --git a/web/src/types/scripts.ts b/web/src/types/scripts.ts --- a/web/src/types/scripts.ts +++ b/web/src/types/scripts.ts @@ -65,6 +65,7 @@ | "xrpc.query" | "xrpc.procedure" | "labeler.apply" + | "job.run" /** Display labels for each trigger kind. */ export const TRIGGER_KIND_LABELS: Record = { @@ -75,21 +76,24 @@ "xrpc.query": "XRPC query", "xrpc.procedure": "XRPC procedure", "labeler.apply": "Label arrival", + "job.run": "Job runner", } /** Top-level grouping for the Scripts list page. */ -export type TriggerFamily = "record" | "xrpc" | "labeler" +export type TriggerFamily = "record" | "xrpc" | "labeler" | "job" export const TRIGGER_FAMILY_LABELS: Record = { record: "Record events", xrpc: "XRPC handlers", labeler: "Label arrivals", + job: "Job runners", } /** Map a trigger kind to its top-level family. */ export function familyOf(kind: TriggerKind): TriggerFamily { if (kind.startsWith("record.")) return "record" if (kind.startsWith("xrpc.")) return "xrpc" + if (kind.startsWith("job.")) return "job" return "labeler" } @@ -113,6 +117,7 @@ "xrpc.query", "xrpc.procedure", "labeler.apply", + "job.run", ] as const).find((k) => k === prefix) if (!kind) return null return { kind, suffix } @@ -131,5 +136,30 @@ function handle() log("script fired") return event +end +` + +export const DEFAULT_JOB_SCRIPT_BODY = `-- Job runner: executes as a background job. +-- +-- Available globals: +-- job.input — the input table passed to jobs.create() +-- job.id — the job's UUID +-- job.progress() — persist progress (visible in the dashboard) +-- job.should_stop() — check for pause/cancel (cooperative) +-- job.wait(seconds) — sleep (0–3600s) +-- +-- Available APIs: db.*, http.*, xrpc.*, atproto.*, Record.*, env. +-- Return value becomes the job's result. + +function handle() + local input = job.input + + job.progress({ status = "working" }) + + if job.should_stop() then + return { partial = true } + end + + return { done = true } end ` diff --git a/web/tests/e2e/jobs.spec.ts b/web/tests/e2e/jobs.spec.ts new file mode 100644 --- /dev/null +++ b/web/tests/e2e/jobs.spec.ts @@ -0,0 +1,216 @@ +import { test, expect } from "@playwright/test" +import { randomUUID } from "crypto" +import pg from "pg" +import { loginAsTestAdmin } from "./auth-helper" + +const DB_URL = "postgres://happyview:happyview@localhost:5434/happyview_test" +const TEST_DID = "did:plc:e2e-test-admin" + +async function seedJob( + status: string, + jobType = "test.e2e.export", +): Promise { + const client = new pg.Client(DB_URL) + await client.connect() + try { + const id = randomUUID() + const now = new Date().toISOString() + await client.query( + `INSERT INTO happyview_jobs (id, job_type, status, input, progress, created_by, created_at) + VALUES ($1, $2, $3, $4, $5, $6, $7)`, + [ + id, + jobType, + status, + JSON.stringify({ source: "e2e-test" }), + JSON.stringify({}), + TEST_DID, + now, + ], + ) + return id + } finally { + await client.end() + } +} + +async function cleanupJobs(): Promise { + const client = new pg.Client(DB_URL) + await client.connect() + try { + await client.query( + "DELETE FROM happyview_jobs WHERE created_by = $1", + [TEST_DID], + ) + } finally { + await client.end() + } +} + +test.describe("Jobs Dashboard", () => { + test.beforeEach(async ({ page }) => { + await loginAsTestAdmin(page) + }) + + test.afterEach(async () => { + await cleanupJobs() + }) + + test("shows empty state when no jobs exist", async ({ page }) => { + await page.goto("/dashboard/jobs") + + await expect( + page.getByText("No background jobs yet"), + ).toBeVisible({ timeout: 5000 }) + }) + + test("lists seeded jobs in the table", async ({ page }) => { + await seedJob("pending", "test.e2e.alpha") + await seedJob("running", "test.e2e.beta") + + await page.goto("/dashboard/jobs") + + const rows = page.locator("table tbody tr") + await expect(rows).toHaveCount(2, { timeout: 5000 }) + + await expect(page.getByText("test.e2e.alpha")).toBeVisible() + await expect(page.getByText("test.e2e.beta")).toBeVisible() + }) + + test("filters jobs by status", async ({ page }) => { + await seedJob("pending", "test.e2e.pending-job") + await seedJob("completed", "test.e2e.completed-job") + + await page.goto("/dashboard/jobs") + + const rows = page.locator("table tbody tr") + await expect(rows).toHaveCount(2, { timeout: 5000 }) + + await page.getByRole("combobox").click() + await page.getByRole("option", { name: "Completed" }).click() + + await expect(rows).toHaveCount(1, { timeout: 5000 }) + await expect(page.getByText("test.e2e.completed-job")).toBeVisible() + await expect(page.getByText("test.e2e.pending-job")).not.toBeVisible() + }) + + test("opens detail sheet when clicking a job row", async ({ page }) => { + const id = await seedJob("pending", "test.e2e.detail") + + await page.goto("/dashboard/jobs") + + const row = page.locator("table tbody tr", { + hasText: "test.e2e.detail", + }) + await expect(row).toBeVisible({ timeout: 5000 }) + await row.click() + + const sheet = page.locator("[data-state='open'][role='dialog']") + await expect(sheet).toBeVisible({ timeout: 3000 }) + + await expect(sheet.getByText("Job Details")).toBeVisible() + await expect(sheet.getByText(id)).toBeVisible() + await expect(sheet.getByText("test.e2e.detail")).toBeVisible() + await expect(sheet.getByText("pending")).toBeVisible() + }) + + test("shows cancel button for running job and cancels it", async ({ + page, + }) => { + const id = await seedJob("running", "test.e2e.cancel") + + await page.goto("/dashboard/jobs") + + const row = page.locator("table tbody tr", { + hasText: "test.e2e.cancel", + }) + await expect(row).toBeVisible({ timeout: 5000 }) + await row.click() + + const sheet = page.locator("[data-state='open'][role='dialog']") + await expect(sheet).toBeVisible({ timeout: 3000 }) + + const cancelButton = sheet.getByRole("button", { name: "Cancel Job" }) + await expect(cancelButton).toBeVisible() + await cancelButton.click() + + await expect(page.getByText("Job cancelled")).toBeVisible({ + timeout: 5000, + }) + }) + + test("shows pause button for running job", async ({ page }) => { + await seedJob("running", "test.e2e.pause") + + await page.goto("/dashboard/jobs") + + const row = page.locator("table tbody tr", { + hasText: "test.e2e.pause", + }) + await expect(row).toBeVisible({ timeout: 5000 }) + await row.click() + + const sheet = page.locator("[data-state='open'][role='dialog']") + await expect(sheet).toBeVisible({ timeout: 3000 }) + + await expect( + sheet.getByRole("button", { name: "Pause Job" }), + ).toBeVisible() + }) + + test("shows resume button for paused job", async ({ page }) => { + await seedJob("paused", "test.e2e.resume") + + await page.goto("/dashboard/jobs") + + const row = page.locator("table tbody tr", { + hasText: "test.e2e.resume", + }) + await expect(row).toBeVisible({ timeout: 5000 }) + await row.click() + + const sheet = page.locator("[data-state='open'][role='dialog']") + await expect(sheet).toBeVisible({ timeout: 3000 }) + + await expect( + sheet.getByRole("button", { name: "Resume Job" }), + ).toBeVisible() + }) + + test("shows error section for failed job", async ({ page }) => { + const client = new pg.Client(DB_URL) + await client.connect() + try { + const id = randomUUID() + const now = new Date().toISOString() + await client.query( + `INSERT INTO happyview_jobs (id, job_type, status, input, progress, error, created_by, created_at, completed_at) + VALUES ($1, $2, 'failed', $3, $4, $5, $6, $7, $7)`, + [ + id, + "test.e2e.failed", + JSON.stringify({}), + JSON.stringify({}), + "something went wrong", + TEST_DID, + now, + ], + ) + } finally { + await client.end() + } + + await page.goto("/dashboard/jobs") + + const row = page.locator("table tbody tr", { + hasText: "test.e2e.failed", + }) + await expect(row).toBeVisible({ timeout: 5000 }) + await row.click() + + const sheet = page.locator("[data-state='open'][role='dialog']") + await expect(sheet).toBeVisible({ timeout: 3000 }) + + await expect(sheet.getByText("something went wrong")).toBeVisible() + }) +}) diff --git a/web/tests/e2e/script-job.spec.ts b/web/tests/e2e/script-job.spec.ts new file mode 100644 --- /dev/null +++ b/web/tests/e2e/script-job.spec.ts @@ -0,0 +1,90 @@ +import { test, expect } from "@playwright/test" +import { loginAsTestAdmin } from "./auth-helper" + +const JOB_TYPE = "test.e2e.myjob" +const TRIGGER_ID = `job.run:${JOB_TYPE}` + +async function cleanupScript( + request: import("@playwright/test").APIRequestContext, +) { + await request.delete(`/admin/scripts/${encodeURIComponent(TRIGGER_ID)}`) +} + +test.describe("Job Script Creation", () => { + test.beforeEach(async ({ page }) => { + await loginAsTestAdmin(page) + }) + + test.afterEach(async ({ page }) => { + await cleanupScript(page.request) + }) + + test("selecting Job source shows job type input and composes trigger id", async ({ + page, + }) => { + await page.goto("/dashboard/settings/scripts/new") + + const sourceSelect = page.locator("#source-pick") + await expect(sourceSelect).toBeVisible({ timeout: 5000 }) + + await sourceSelect.click() + await page.getByRole("option", { name: /Job/ }).click() + + const jobTypeInput = page.locator("#job-type-input") + await expect(jobTypeInput).toBeVisible() + + await expect(page.locator("#action-pick")).not.toBeVisible() + + await jobTypeInput.fill(JOB_TYPE) + + await expect(page.getByText(TRIGGER_ID)).toBeVisible() + }) + + test("creating a job script navigates to detail page", async ({ page }) => { + await page.goto("/dashboard/settings/scripts/new") + + await page.locator("#source-pick").click() + await page.getByRole("option", { name: /Job/ }).click() + + await page.locator("#job-type-input").fill(JOB_TYPE) + + await page.getByRole("button", { name: "Create script" }).click() + + await page.waitForURL( + `**/dashboard/settings/scripts/${encodeURIComponent(TRIGGER_ID)}`, + { timeout: 5000 }, + ) + + await expect(page.getByText("Job runner")).toBeVisible() + await expect(page.getByText(TRIGGER_ID)).toBeVisible() + }) + + test("job script has job-specific template body", async ({ page }) => { + await page.goto("/dashboard/settings/scripts/new") + + await page.locator("#source-pick").click() + await page.getByRole("option", { name: /Job/ }).click() + + await expect(page.getByText("job.input")).toBeVisible({ timeout: 3000 }) + await expect(page.getByText("job.should_stop")).toBeVisible() + }) + + test("job script appears in scripts list with Job runners family", async ({ + page, + }) => { + await page.request.post("/admin/scripts", { + data: { + id: TRIGGER_ID, + body: "function handle()\n return { ok = true }\nend", + }, + }) + + await page.goto("/dashboard/settings/scripts") + + const row = page.locator("table tbody tr", { hasText: JOB_TYPE }) + await expect(row).toBeVisible({ timeout: 5000 }) + + await expect(row.getByText("Job runner")).toBeVisible() + await expect(row.getByText("Job runners")).toBeVisible() + }) +}) diff --git a/packages/docs/content/blog/happyview-2.10.md b/packages/docs/content/blog/happyview-2.10.md --- a/packages/docs/content/blog/happyview-2.10.md +++ b/packages/docs/content/blog/happyview-2.10.md @@ -23,7 +23,7 @@ With a service identity in place, HappyView can act as a service proxy. A PDS sends a request with an `atproto-proxy` header pointing at your AppView, HappyView verifies the caller via service auth, runs your XRPC handler, and responds. This is how atproto apps are _supposed_ to work! Up to this point HappyView only supported direct connections via DPoP. -Full docs: [Service Identity](/docs/getting-started/service-identity). +Full docs: [Service Identity](/getting-started/service-identity). ## Permissioned spaces alignment @@ -63,7 +63,7 @@ - **Record operation log** - `listRepoOps` returns the oplog for sync - **Write notifications** - `registerNotify`, `notifyWrite`, `notifySpaceDeleted` -Full docs: [Permissioned Spaces](/docs/experimental/spaces/). +Full docs: [Permissioned Spaces](/experimental/spaces/). ## Blob utilities for Lua @@ -75,7 +75,7 @@ local new_blob_ref = uploaded.blob ``` -Full docs: [atproto API (`blob_download` / `blob_upload`)](/docs/api-reference/lua/atproto-api#atprotoblob_download). +Full docs: [atproto API (`blob_download` / `blob_upload`)](/api-reference/lua/atproto-api#atprotoblob_download). ## Prefixed database tables diff --git a/packages/docs/content/blog/happyview-2.9.md b/packages/docs/content/blog/happyview-2.9.md --- a/packages/docs/content/blog/happyview-2.9.md +++ b/packages/docs/content/blog/happyview-2.9.md @@ -28,11 +28,11 @@ For record events, the dispatcher tries the action-specific trigger first (e.g. `record.create:com.example.post`), then falls back to the wildcard `record.index:com.example.post`. This means you can have one general-purpose script that handles everything, or surgical scripts for specific actions — or both. -Scripts are managed through the dashboard under **Settings > Scripts**, or via the new [`/admin/scripts`](/docs/api-reference/admin/scripts) API endpoints. The lexicon detail page also shows which scripts target each lexicon, with links to create or edit them. +Scripts are managed through the dashboard under **Settings > Scripts**, or via the new [`/admin/scripts`](/api-reference/admin/scripts) API endpoints. The lexicon detail page also shows which scripts target each lexicon, with links to create or edit them. **If you're upgrading from v2.x to v2.9:** existing index hooks and lexicon scripts will be migrated to the new system automatically. -Full docs: [Record & Label Scripts](/docs/guides/label-scripts), [Lua Scripting](/docs/guides/lua-scripting), [Admin API — Scripts](/docs/api-reference/admin/scripts). +Full docs: [Record & Label Scripts](/guides/label-scripts), [Lua Scripting](/guides/lua-scripting), [Admin API — Scripts](/api-reference/admin/scripts). ## Backfill, but concurrent @@ -72,7 +72,7 @@ }) ``` -Full docs are in the [Database API reference](/docs/api-reference/lua/database-api). +Full docs are in the [Database API reference](/api-reference/lua/database-api). ## Auth fixes @@ -87,7 +87,7 @@ - `GET /oauth/sessions/{did}/devices` — list all active sessions - `DELETE /oauth/sessions/{did}/devices/{session_id}` — revoke a session -The existing `DELETE /oauth/sessions/{did}` endpoint still works: confidential clients revoke all device sessions for the user, and public clients revoke the session matching their DPoP key. Full details in the [Authentication guide](/docs/getting-started/authentication#6-managing-device-sessions). +The existing `DELETE /oauth/sessions/{did}` endpoint still works: confidential clients revoke all device sessions for the user, and public clients revoke the session matching their DPoP key. Full details in the [Authentication guide](/getting-started/authentication#6-managing-device-sessions). ## SDK fix diff --git a/web/src/app/dashboard/jobs/page.tsx b/web/src/app/dashboard/jobs/page.tsx new file mode 100644 --- /dev/null +++ b/web/src/app/dashboard/jobs/page.tsx @@ -0,0 +1,515 @@ +"use client"; + +import { useCallback, useEffect, useState } from "react"; +import { toast } from "sonner"; +import { + CheckCircle2, + ChevronDown, + Circle, + Loader2, + PauseCircle, + XCircle, +} from "lucide-react"; + +import { useCurrentUser } from "@/hooks/use-current-user"; +import { toastError } from "@/lib/format"; +import { + cancelJob, + getJobs, + pauseJob, + resumeJob, +} from "@/lib/api"; +import type { Job } from "@/types/jobs"; +import { SiteHeader } from "@/components/site-header"; +import { Badge } from "@/components/ui/badge"; +import { Button } from "@/components/ui/button"; +import { + Collapsible, + CollapsibleContent, + CollapsibleTrigger, +} from "@/components/ui/collapsible"; +import { + Select, + SelectContent, + SelectItem, + SelectTrigger, + SelectValue, +} from "@/components/ui/select"; +import { + Sheet, + SheetContent, + SheetFooter, + SheetHeader, + SheetTitle, +} from "@/components/ui/sheet"; +import { + Table, + TableBody, + TableCell, + TableHead, + TableHeader, + TableRow, +} from "@/components/ui/table"; + +const STATUS_OPTIONS = [ + { value: "all", label: "All statuses" }, + { value: "pending", label: "Pending" }, + { value: "running", label: "Running" }, + { value: "paused", label: "Paused" }, + { value: "completed", label: "Completed" }, + { value: "failed", label: "Failed" }, + { value: "cancelled", label: "Cancelled" }, +] as const; + +function statusBadge(status: string) { + switch (status) { + case "completed": + return ( + + completed + + ); + case "failed": + return failed; + case "cancelled": + return ( + + cancelled + + ); + case "cancelling": + return ( + + cancelling + + ); + case "pausing": + return ( + + pausing + + ); + case "paused": + return ( + + paused + + ); + case "running": + return ( + + running + + ); + case "pending": + return pending; + default: + return {status}; + } +} + +function statusIcon(status: string) { + switch (status) { + case "completed": + return ; + case "failed": + return ; + case "cancelled": + return ; + case "cancelling": + return ; + case "pausing": + return ; + case "paused": + return ; + case "running": + return ; + default: + return ; + } +} + +function hasContent(obj: Record | null | undefined): boolean { + if (!obj) return false; + return Object.keys(obj).length > 0; +} + +function relativeTime(dateStr: string): string { + const now = Date.now(); + const then = new Date(dateStr).getTime(); + const diff = now - then; + const seconds = Math.floor(diff / 1000); + if (seconds < 60) return "just now"; + const minutes = Math.floor(seconds / 60); + if (minutes < 60) return `${minutes}m ago`; + const hours = Math.floor(minutes / 60); + if (hours < 24) return `${hours}h ago`; + const days = Math.floor(hours / 24); + return `${days}d ago`; +} + +export default function JobsPage() { + const { hasPermission } = useCurrentUser(); + const [jobs, setJobs] = useState([]); + const [statusFilter, setStatusFilter] = useState("all"); + const [selectedJobId, setSelectedJobId] = useState(null); + const [loading, setLoading] = useState(true); + + const load = useCallback(() => { + const params = statusFilter !== "all" ? { status: statusFilter } : {}; + getJobs(params) + .then((resp) => { + setJobs(resp.jobs); + setLoading(false); + }) + .catch((e) => { + toastError("Failed to load jobs", e); + setLoading(false); + }); + }, [statusFilter]); + + useEffect(() => { + setLoading(true); + load(); + }, [load]); + + // Poll every 5 seconds for active jobs + const hasActiveJobs = jobs.some( + (j) => + j.status === "running" || + j.status === "pending" || + j.status === "cancelling" || + j.status === "pausing", + ); + + useEffect(() => { + const interval = setInterval(load, hasActiveJobs ? 3000 : 10000); + return () => clearInterval(interval); + }, [load, hasActiveJobs]); + + const selectedJob = jobs.find((j) => j.id === selectedJobId) ?? null; + const canManage = hasPermission("jobs:manage"); + + return ( + <> + +
+
+

Background Jobs

+ +
+ +
+ + + + + Type + Status + Created by + Created + + + + {loading && jobs.length === 0 && ( + + + + + + )} + {!loading && jobs.length === 0 && ( + + + {statusFilter !== "all" + ? `No ${statusFilter} jobs.` + : "No background jobs yet. Jobs are created by Lua scripts via jobs.create()."} + + + )} + {jobs.map((job) => ( + setSelectedJobId(job.id)} + onKeyDown={(e) => { + if (e.key === "Enter" || e.key === " ") { + e.preventDefault(); + setSelectedJobId(job.id); + } + }} + > + + {statusIcon(job.status)} + + + {job.job_type} + + {statusBadge(job.status)} + + {job.created_by} + + + {relativeTime(job.created_at)} + + + ))} + +
+
+ + { + if (!open) { + setSelectedJobId(null); + load(); + } + }} + > + + {selectedJob && ( + + )} + + +
+ + ); +} + +function JobDetail({ + job, + canManage, + onAction, +}: { + job: Job; + canManage: boolean; + onAction: () => void; +}) { + const [actionLoading, setActionLoading] = useState(null); + const isActive = + job.status === "running" || + job.status === "cancelling" || + job.status === "pausing"; + + async function handleCancel() { + setActionLoading("cancel"); + try { + await cancelJob(job.id); + toast.success("Job cancelled"); + onAction(); + } catch (e) { + toastError("Failed to cancel job", e); + } finally { + setActionLoading(null); + } + } + + async function handlePause() { + setActionLoading("pause"); + try { + await pauseJob(job.id); + toast.success("Job paused"); + onAction(); + } catch (e) { + toastError("Failed to pause job", e); + } finally { + setActionLoading(null); + } + } + + async function handleResume() { + setActionLoading("resume"); + try { + await resumeJob(job.id); + toast.success("Job resumed"); + onAction(); + } catch (e) { + toastError("Failed to resume job", e); + } finally { + setActionLoading(null); + } + } + + return ( + <> + + + Job Details + + +
+
+
+ Job ID +

{job.id}

+
+
+ Type +

{job.job_type}

+
+
+ Status +
{statusBadge(job.status)}
+
+
+ Created by +

{job.created_by}

+
+
+ Created +

+ {new Date(job.created_at).toLocaleString()} +

+
+ {job.started_at && ( +
+ Started +

+ {new Date(job.started_at).toLocaleString()} +

+
+ )} + {job.completed_at && ( +
+ Completed +

+ {new Date(job.completed_at).toLocaleString()} +

+
+ )} +
+ + {job.error && ( +
+ Error +
+ {job.error} +
+
+ )} + + + + {job.result && } +
+ + {canManage && ( + + {(job.status === "running" || job.status === "pausing") && ( + + )} + {job.status === "paused" && ( + + )} + {isActive && ( + + )} + {job.status === "paused" && ( + + )} + + )} + + ); +} + +function JsonSection({ + title, + data, + defaultOpen = false, +}: { + title: string; + data: Record | null; + defaultOpen?: boolean; +}) { + const [open, setOpen] = useState(defaultOpen); + const empty = !hasContent(data); + + if (empty) return null; + + return ( + + + + + +
+          {JSON.stringify(data, null, 2)}
+        
+
+
+ ); +} diff --git a/web/src/app/dashboard/settings/scripts/script-form.tsx b/web/src/app/dashboard/settings/scripts/script-form.tsx --- a/web/src/app/dashboard/settings/scripts/script-form.tsx +++ b/web/src/app/dashboard/settings/scripts/script-form.tsx @@ -4,6 +4,7 @@ import { MonacoEditor } from "@/components/monaco-editor"; import { Badge } from "@/components/ui/badge"; +import { Input } from "@/components/ui/input"; import { Label } from "@/components/ui/label"; import { Select, @@ -18,7 +19,12 @@ import { Textarea } from "@/components/ui/textarea"; import type { LexiconSummary } from "@/types/lexicons"; import type { TriggerKind } from "@/types/scripts"; -import { TRIGGER_KIND_LABELS, parseTriggerId } from "@/types/scripts"; +import { + DEFAULT_JOB_SCRIPT_BODY, + DEFAULT_SCRIPT_BODY, + TRIGGER_KIND_LABELS, + parseTriggerId, +} from "@/types/scripts"; /** * Sentinel suffix used when the operator picks "Actor" in the lexicon @@ -27,15 +33,27 @@ */ export const ACTOR_SUFFIX = "_actor"; +/** + * Sentinel value for the source selector when the operator picks "Job". + * The actual suffix is typed into a free-form input (the job type name). + */ +export const JOB_SOURCE = "_job"; + export interface ScriptFormState { /** Trigger kind selector value (e.g. `record.create`). */ kind: TriggerKind; /** * Suffix portion of the trigger id — usually an NSID (= a lexicon id), * or the literal `_actor` when `kind === "labeler.apply"` for - * actor-level labels. + * actor-level labels. For jobs, this is the user-typed job type name. */ suffix: string; + /** + * Which source-selector value was chosen. Usually identical to `suffix` + * (i.e. a lexicon NSID or `_actor`). For jobs this is `_job` while + * `suffix` holds the free-form job type name. + */ + source: string; description: string; body: string; } @@ -51,9 +69,15 @@ body: string; }): ScriptFormState { const parsed = parseTriggerId(args.id); + const kind = parsed?.kind ?? "record.index"; + const suffix = parsed?.suffix ?? ""; + let source = suffix; + if (kind === "job.run") source = JOB_SOURCE; + else if (suffix === ACTOR_SUFFIX) source = ACTOR_SUFFIX; return { - kind: parsed?.kind ?? "record.index", - suffix: parsed?.suffix ?? "", + kind, + suffix, + source, description: args.description ?? "", body: args.body, }; @@ -93,9 +117,23 @@ { kind: "xrpc.procedure", label: "Procedure handler" }, ]; -function actionsFor(suffix: string, lexicons: LexiconSummary[]): ActionOption[] { - if (suffix === ACTOR_SUFFIX) return ACTOR_ACTIONS; - const lex = lexicons.find((l) => l.id === suffix); +const JOB_ACTIONS: ActionOption[] = [{ kind: "job.run", label: "Job runner" }]; + +const JOB_TYPE_PATTERN = /^[a-z0-9][a-z0-9._-]*$/; + +export function isValidJobType(value: string): boolean { + return ( + value.length > 0 && value.length <= 128 && JOB_TYPE_PATTERN.test(value) + ); +} + +function actionsFor( + source: string, + lexicons: LexiconSummary[], +): ActionOption[] { + if (source === ACTOR_SUFFIX) return ACTOR_ACTIONS; + if (source === JOB_SOURCE) return JOB_ACTIONS; + const lex = lexicons.find((l) => l.id === source); if (!lex) return []; switch (lex.lexicon_type) { case "record": @@ -129,12 +167,15 @@ onChange, idLocked, lexicons, + lexiconsLoading, }: { state: ScriptFormState; onChange: (next: ScriptFormState) => void; idLocked?: boolean; /** Required when `idLocked` is false; ignored otherwise. */ lexicons?: LexiconSummary[]; + /** True while the lexicon list is being fetched. */ + lexiconsLoading?: boolean; }) { return (
@@ -145,6 +186,7 @@ state={state} onChange={onChange} lexicons={lexicons ?? []} + lexiconsLoading={lexiconsLoading} /> )} @@ -193,19 +235,23 @@ state, onChange, lexicons, + lexiconsLoading, }: { state: ScriptFormState; onChange: (next: ScriptFormState) => void; lexicons: LexiconSummary[]; + lexiconsLoading?: boolean; }) { const sortedLexicons = useMemo( () => [...lexicons].sort((a, b) => a.id.localeCompare(b.id)), [lexicons], ); const actions = useMemo( - () => actionsFor(state.suffix, lexicons), - [state.suffix, lexicons], + () => actionsFor(state.source, lexicons), + [state.source, lexicons], ); + + const isJob = state.source === JOB_SOURCE; const stateRef = useRef(state); stateRef.current = state; @@ -218,14 +264,34 @@ } }, [actions, onChange]); - function handleSuffixChange(next: string) { - // Pre-snap kind so the resolved trigger id badge updates immediately - // rather than flickering through an invalid state. + function handleSourceChange(next: string) { + const wasJob = state.source === JOB_SOURCE; + const isNowJob = next === JOB_SOURCE; + const bodyIsDefault = + state.body === DEFAULT_SCRIPT_BODY || + state.body === DEFAULT_JOB_SCRIPT_BODY; + + if (isNowJob) { + onChange({ + ...state, + source: JOB_SOURCE, + suffix: "", + kind: "job.run", + body: bodyIsDefault ? DEFAULT_JOB_SCRIPT_BODY : state.body, + }); + return; + } const nextActions = actionsFor(next, lexicons); const nextKind = nextActions.some((a) => a.kind === state.kind) ? state.kind : (nextActions[0]?.kind ?? state.kind); - onChange({ ...state, suffix: next, kind: nextKind }); + onChange({ + ...state, + source: next, + suffix: next, + kind: nextKind, + body: wasJob && bodyIsDefault ? DEFAULT_SCRIPT_BODY : state.body, + }); } const triggerPreview = @@ -235,12 +301,12 @@ <>
-
- - + {isJob ? ( + <> + + onChange({ ...state, suffix: e.target.value })} + placeholder="e.g. export, migrate, sync" + className="h-8 text-sm font-mono" + aria-invalid={ + state.suffix.length > 0 && !isValidJobType(state.suffix) + } + /> + {state.suffix.length > 0 && !isValidJobType(state.suffix) ? ( +

+ Lowercase letters, numbers, dots, hyphens, and underscores + only. +

+ ) : ( +

+ Must match the type passed to{" "} + + jobs.create() + {" "} + in the queuing script. +

+ )} + + ) : ( + <> + + + + )}
@@ -304,7 +420,9 @@ {triggerPreview} ) : ( - Pick a lexicon to compose the trigger id. + {isJob + ? "Enter a job type to compose the trigger id." + : "Pick a source to compose the trigger id."} )}

diff --git a/web/src/app/dashboard/settings/scripts/[id]/script-detail.tsx b/web/src/app/dashboard/settings/scripts/[id]/script-detail.tsx --- a/web/src/app/dashboard/settings/scripts/[id]/script-detail.tsx +++ b/web/src/app/dashboard/settings/scripts/[id]/script-detail.tsx @@ -66,13 +66,20 @@ ); }, [state, original]); - async function handleSave() { - if (!state || !script) return; + useEffect(() => { + if (!isDirty) return; + function onBeforeUnload(e: BeforeUnloadEvent) { + e.preventDefault(); + } + window.addEventListener("beforeunload", onBeforeUnload); + return () => window.removeEventListener("beforeunload", onBeforeUnload); + }, [isDirty]); + + const handleSave = useCallback(async () => { + if (!state || !script || !isDirty || saving) return; setSaving(true); setError(null); try { - // PATCH only the editable fields. Trigger id is the PK — to - // rename, delete and recreate. await patchScript(script.id, { body: state.body, description: state.description.trim() || null, @@ -83,7 +90,18 @@ } finally { setSaving(false); } - } + }, [state, script, isDirty, saving, load]); + + useEffect(() => { + function onKeyDown(e: KeyboardEvent) { + if ((e.metaKey || e.ctrlKey) && e.key === "Enter") { + e.preventDefault(); + handleSave(); + } + } + window.addEventListener("keydown", onKeyDown); + return () => window.removeEventListener("keydown", onKeyDown); + }, [handleSave]); async function handleDelete() { if (!script) return; @@ -126,6 +144,7 @@ record: "Record event", xrpc: "XRPC handler", labeler: "Label arrival", + job: "Job runner", }; return ( @@ -178,6 +197,9 @@ {canManage && ( )}
diff --git a/web/src/app/dashboard/settings/scripts/new/page.tsx b/web/src/app/dashboard/settings/scripts/new/page.tsx --- a/web/src/app/dashboard/settings/scripts/new/page.tsx +++ b/web/src/app/dashboard/settings/scripts/new/page.tsx @@ -1,20 +1,37 @@ "use client"; -import { Suspense, useEffect, useState } from "react"; +import { Suspense, useCallback, useEffect, useMemo, useState } from "react"; import { useRouter, useSearchParams } from "next/navigation"; import { useCurrentUser } from "@/hooks/use-current-user"; import { getLexicons, upsertScript } from "@/lib/api"; import type { LexiconSummary } from "@/types/lexicons"; import type { TriggerKind } from "@/types/scripts"; -import { DEFAULT_SCRIPT_BODY, parseTriggerId } from "@/types/scripts"; +import { + DEFAULT_JOB_SCRIPT_BODY, + DEFAULT_SCRIPT_BODY, + parseTriggerId, +} from "@/types/scripts"; import { SiteHeader } from "@/components/site-header"; +import { + AlertDialog, + AlertDialogAction, + AlertDialogCancel, + AlertDialogContent, + AlertDialogDescription, + AlertDialogFooter, + AlertDialogHeader, + AlertDialogTitle, + AlertDialogTrigger, +} from "@/components/ui/alert-dialog"; import { Button } from "@/components/ui/button"; import { + JOB_SOURCE, ScriptForm, type ScriptFormState, composeTriggerId, + isValidJobType, } from "../script-form"; function NewScriptInner() { @@ -28,6 +45,7 @@ // form even if the call fails (the operator can still pick "Actor" // and create a labeler.apply:_actor script). const [lexicons, setLexicons] = useState([]); + const [lexiconsLoading, setLexiconsLoading] = useState(true); const [saving, setSaving] = useState(false); const [error, setError] = useState(null); @@ -39,23 +57,36 @@ useEffect(() => { getLexicons() .then(setLexicons) - .catch(() => setLexicons([])); + .catch(() => setLexicons([])) + .finally(() => setLexiconsLoading(false)); }, []); - if (!hasPermission("scripts:manage")) { + const isDirty = useMemo(() => { + const defaultBody = + state.source === JOB_SOURCE ? DEFAULT_JOB_SCRIPT_BODY : DEFAULT_SCRIPT_BODY; return ( - <> - -
-

- You don't have permission to create scripts. -

-
- + state.suffix !== "" || + state.description !== "" || + state.body !== defaultBody ); - } + }, [state]); - async function handleSave() { + useEffect(() => { + if (!isDirty) return; + function onBeforeUnload(e: BeforeUnloadEvent) { + e.preventDefault(); + } + window.addEventListener("beforeunload", onBeforeUnload); + return () => window.removeEventListener("beforeunload", onBeforeUnload); + }, [isDirty]); + + const canSave = + !saving && + !!state.suffix && + !(state.source === JOB_SOURCE && !isValidJobType(state.suffix)); + + const handleSave = useCallback(async () => { + if (!canSave) return; setSaving(true); setError(null); try { @@ -70,6 +101,30 @@ setError(e instanceof Error ? e.message : String(e)); setSaving(false); } + }, [canSave, state, router]); + + useEffect(() => { + function onKeyDown(e: KeyboardEvent) { + if ((e.metaKey || e.ctrlKey) && e.key === "Enter") { + e.preventDefault(); + handleSave(); + } + } + window.addEventListener("keydown", onKeyDown); + return () => window.removeEventListener("keydown", onKeyDown); + }, [handleSave]); + + if (!hasPermission("scripts:manage")) { + return ( + <> + +
+

+ You don't have permission to create scripts. +

+
+ + ); } return ( @@ -78,11 +133,48 @@
{error &&

{error}

} - +
-
- + + + + Discard changes? + + You have unsaved changes that will be lost. + + + + Keep editing + router.push("/dashboard/settings/scripts")} + > + Discard + + + + + ) : ( + + )} +
@@ -92,7 +184,11 @@ export default function NewScriptPage() { return ( - + + } + > ); @@ -105,21 +201,26 @@ if (presetId) { const parsed = parseTriggerId(presetId); if (parsed) { + const isJob = parsed.kind === "job.run"; return { kind: parsed.kind, suffix: parsed.suffix, + source: isJob ? JOB_SOURCE : parsed.suffix, description: "", - body: DEFAULT_SCRIPT_BODY, + body: isJob ? DEFAULT_JOB_SCRIPT_BODY : DEFAULT_SCRIPT_BODY, }; } } // Fallbacks to a sensible default. Suffix starts empty so the form - // surfaces the "Pick a lexicon to compose the trigger id" hint. + // surfaces the "Pick a source to compose the trigger id" hint. const kind = (searchParams.get("kind") as TriggerKind | null) ?? "record.index"; + const source = searchParams.get("source") ?? searchParams.get("suffix") ?? ""; + const isJob = kind === "job.run" || source === JOB_SOURCE; return { - kind, - suffix: searchParams.get("suffix") ?? "", + kind: isJob ? "job.run" : kind, + suffix: isJob ? "" : (searchParams.get("suffix") ?? ""), + source: isJob ? JOB_SOURCE : source, description: "", - body: DEFAULT_SCRIPT_BODY, + body: isJob ? DEFAULT_JOB_SCRIPT_BODY : DEFAULT_SCRIPT_BODY, }; } -- tangled.sh