diff --git a/src/admin/jobs.rs b/src/admin/jobs.rs index 27bca7c..46d3b1c 100644 --- a/src/admin/jobs.rs +++ b/src/admin/jobs.rs @@ -23,7 +23,7 @@ pub async fn list_jobs( ) -> Result, AppError> { auth.require(Permission::JobsRead).await?; - let limit = query.limit.unwrap_or(50).min(100); + let limit = query.limit.unwrap_or(50).clamp(1, 100); let (jobs_list, cursor) = jobs::db::list_jobs( &state, query.status.as_deref(), @@ -72,7 +72,7 @@ pub async fn cancel_job( jobs::db::set_status(&state, &id, "cancelled").await?; Ok(Json(serde_json::json!({ "status": "cancelled" }))) } - _ => Err(AppError::BadRequest(format!( + _ => Err(AppError::Conflict(format!( "cannot cancel job with status: {}", job.status ))), @@ -91,7 +91,7 @@ pub async fn pause_job( .ok_or_else(|| AppError::NotFound("job not found".into()))?; if job.status != "running" { - return Err(AppError::BadRequest(format!( + return Err(AppError::Conflict(format!( "cannot pause job with status: {}", job.status ))); @@ -113,7 +113,7 @@ pub async fn resume_job( .ok_or_else(|| AppError::NotFound("job not found".into()))?; if job.status != "paused" { - return Err(AppError::BadRequest(format!( + return Err(AppError::Conflict(format!( "cannot resume job with status: {}", job.status ))); diff --git a/src/jobs/db.rs b/src/jobs/db.rs index 10120c8..6d00d6a 100644 --- a/src/jobs/db.rs +++ b/src/jobs/db.rs @@ -2,7 +2,7 @@ use serde_json::Value; use uuid::Uuid; use crate::AppState; -use crate::db::{adapt_sql, now_rfc3339}; +use crate::db::{DatabaseBackend, adapt_sql, now_rfc3339}; use crate::error::AppError; use super::Job; @@ -284,10 +284,16 @@ pub async fn find_interrupted_jobs(state: &AppState) -> Vec { 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 sql = match state.db_backend { + DatabaseBackend::Postgres => 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 FOR UPDATE SKIP LOCKED) RETURNING *", + state.db_backend, + ), + DatabaseBackend::Sqlite => 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) AND status = 'pending' RETURNING *", + state.db_backend, + ), + }; let row: Option = sqlx::query_as(&sql) .bind(&now) diff --git a/src/lua/jobs_api.rs b/src/lua/jobs_api.rs index 959628d..f50b397 100644 --- a/src/lua/jobs_api.rs +++ b/src/lua/jobs_api.rs @@ -1,9 +1,13 @@ use mlua::{Lua, LuaSerdeExt, Result as LuaResult}; -use std::sync::Arc; +use regex::Regex; +use std::sync::{Arc, LazyLock}; use crate::AppState; use crate::jobs; +static JOB_TYPE_PATTERN: LazyLock = + LazyLock::new(|| Regex::new(r"^[a-z0-9][a-z0-9._-]*$").unwrap()); + /// Register the `jobs` table for queuing jobs from scripts. /// Available in all script contexts (procedure, query, record-event). pub fn register_jobs_api( @@ -31,6 +35,15 @@ pub fn register_jobs_api( .unwrap_or(false); async move { + if job_type.is_empty() + || job_type.len() > 128 + || !JOB_TYPE_PATTERN.is_match(&job_type) + { + return Err(mlua::Error::runtime( + "job_type must be 1-128 characters matching /^[a-z0-9][a-z0-9._-]*$/", + )); + } + let caller = caller_did.as_deref().ok_or_else(|| { mlua::Error::runtime("jobs.create requires an authenticated caller") })?; @@ -294,4 +307,24 @@ mod tests { .unwrap(); assert!(result); } + + #[tokio::test] + async fn jobs_create_rejects_invalid_job_type() { + 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(); + + for bad in [ + "", + "UPPER", + "has space", + "has:colon", + "-leading-dash", + ".leading-dot", + ] { + let script = format!(r#"return jobs.create("{bad}", {{}})"#); + let result: mlua::Result = lua.load(&script).eval_async().await; + assert!(result.is_err(), "expected error for job_type={bad:?}"); + } + } } diff --git a/src/lua/scripts.rs b/src/lua/scripts.rs index df4ae59..439c7c0 100644 --- a/src/lua/scripts.rs +++ b/src/lua/scripts.rs @@ -32,9 +32,13 @@ //! [`super::execute::execute_procedure_script`] / //! [`super::execute::execute_query_script`] directly. +use regex::Regex; use serde::{Deserialize, Serialize}; use serde_json::Value; -use std::sync::Arc; +use std::sync::{Arc, LazyLock}; + +static JOB_TYPE_RE: LazyLock = + LazyLock::new(|| Regex::new(r"^[a-z0-9][a-z0-9._-]*$").unwrap()); use crate::AppState; use crate::db::{DatabaseBackend, adapt_sql, now_rfc3339}; @@ -60,6 +64,7 @@ pub enum TriggerKind { XrpcQuery, XrpcProcedure, LabelerApply, + JobRun, } /// A trigger id parsed into `(kind, suffix)`. The suffix is either an NSID @@ -82,6 +87,7 @@ impl ParsedTrigger { TriggerKind::XrpcQuery => format!("xrpc.query:{}", self.suffix), TriggerKind::XrpcProcedure => format!("xrpc.procedure:{}", self.suffix), TriggerKind::LabelerApply => format!("labeler.apply:{}", self.suffix), + TriggerKind::JobRun => format!("job.run:{}", self.suffix), } } @@ -92,7 +98,8 @@ impl ParsedTrigger { format!( "trigger id '{id}' must contain a ':' separator; \ valid prefixes: record.{{index,create,update,delete}}:, \ - xrpc.{{query,procedure}}:, labeler.apply:" + xrpc.{{query,procedure}}:, labeler.apply:, \ + job.run:" ) })?; @@ -108,18 +115,21 @@ impl ParsedTrigger { "xrpc.query" => TriggerKind::XrpcQuery, "xrpc.procedure" => TriggerKind::XrpcProcedure, "labeler.apply" => TriggerKind::LabelerApply, + "job.run" => TriggerKind::JobRun, other => { return Err(format!( "unknown trigger prefix '{other}'; valid prefixes: \ record.{{index,create,update,delete}}, xrpc.{{query,procedure}}, \ - labeler.apply" + labeler.apply, job.run" )); } }; - // Suffix validation: NSID for everything except `labeler.apply:_actor`. - match (kind, suffix) { - (TriggerKind::LabelerApply, "_actor") => {} + // Suffix validation: NSID for most triggers, but `labeler.apply:_actor` + // and `job.run:` have their own formats. + match kind { + TriggerKind::JobRun => validate_job_type(suffix)?, + TriggerKind::LabelerApply if suffix == "_actor" => {} _ => validate_nsid(suffix)?, } @@ -130,6 +140,20 @@ impl ParsedTrigger { } } +fn validate_job_type(job_type: &str) -> Result<(), String> { + if job_type.is_empty() || job_type.len() > 128 { + return Err(format!( + "invalid job type '{job_type}': must be 1–128 characters" + )); + } + if !JOB_TYPE_RE.is_match(job_type) { + return Err(format!( + "invalid job type '{job_type}': must match /^[a-z0-9][a-z0-9._-]*$/" + )); + } + Ok(()) +} + /// Minimal NSID validation: at least two dot-separated segments, each /// non-empty and matching `[a-zA-Z][a-zA-Z0-9-]*`. Mirrors the AT Protocol /// spec's character class for everyday use; full Unicode strictness lives @@ -900,6 +924,21 @@ mod tests { assert!(err.contains("valid prefixes")); } + #[test] + fn parse_job_run_trigger() { + let t = ParsedTrigger::parse("job.run:test.export").unwrap(); + assert_eq!(t.kind, TriggerKind::JobRun); + assert_eq!(t.suffix, "test.export"); + assert_eq!(t.id(), "job.run:test.export"); + } + + #[test] + fn rejects_bad_job_type() { + assert!(ParsedTrigger::parse("job.run:UPPER").is_err()); + assert!(ParsedTrigger::parse("job.run:has space").is_err()); + assert!(ParsedTrigger::parse("job.run:").is_err()); + } + #[test] fn rejects_unknown_prefix() { let err = ParsedTrigger::parse("garbage:com.example.thing").unwrap_err(); diff --git a/tests/e2e_jobs.rs b/tests/e2e_jobs.rs index 1ab363c..eead318 100644 --- a/tests/e2e_jobs.rs +++ b/tests/e2e_jobs.rs @@ -272,7 +272,7 @@ async fn cancel_paused_job_sets_cancelled() { #[tokio::test] #[serial] -async fn cancel_completed_job_returns_400() { +async fn cancel_completed_job_returns_409() { common::require_db!(); let app = TestApp::new().await; @@ -289,7 +289,7 @@ async fn cancel_completed_job_returns_400() { .await .unwrap(); - assert_eq!(resp.status(), StatusCode::BAD_REQUEST); + assert_eq!(resp.status(), StatusCode::CONFLICT); } // --------------------------------------------------------------------------- @@ -322,7 +322,7 @@ async fn pause_running_job_sets_pausing() { #[tokio::test] #[serial] -async fn pause_pending_job_returns_400() { +async fn pause_pending_job_returns_409() { common::require_db!(); let app = TestApp::new().await; @@ -339,7 +339,7 @@ async fn pause_pending_job_returns_400() { .await .unwrap(); - assert_eq!(resp.status(), StatusCode::BAD_REQUEST); + assert_eq!(resp.status(), StatusCode::CONFLICT); } // --------------------------------------------------------------------------- @@ -372,7 +372,7 @@ async fn resume_paused_job_sets_pending() { #[tokio::test] #[serial] -async fn resume_running_job_returns_400() { +async fn resume_running_job_returns_409() { common::require_db!(); let app = TestApp::new().await; @@ -389,7 +389,7 @@ async fn resume_running_job_returns_400() { .await .unwrap(); - assert_eq!(resp.status(), StatusCode::BAD_REQUEST); + assert_eq!(resp.status(), StatusCode::CONFLICT); } // --------------------------------------------------------------------------- diff --git a/web/src/app/dashboard/jobs/page.tsx b/web/src/app/dashboard/jobs/page.tsx index 72a09f6..43a0026 100644 --- a/web/src/app/dashboard/jobs/page.tsx +++ b/web/src/app/dashboard/jobs/page.tsx @@ -129,9 +129,14 @@ function statusIcon(status: string) { } } -function hasContent(obj: Record | null | undefined): boolean { - if (!obj) return false; - return Object.keys(obj).length > 0; +function hasContent(value: unknown): boolean { + if (value == null) return false; + if (typeof value === "object") { + return Array.isArray(value) + ? value.length > 0 + : Object.keys(value as Record).length > 0; + } + return true; } function relativeTime(dateStr: string): string { @@ -421,7 +426,7 @@ function JobDetail({ - {job.result && } + {hasContent(job.result) && } {canManage && ( @@ -484,7 +489,7 @@ function JsonSection({ defaultOpen = false, }: { title: string; - data: Record | null; + data: unknown; defaultOpen?: boolean; }) { const [open, setOpen] = useState(defaultOpen); diff --git a/web/src/app/dashboard/settings/scripts/[id]/script-detail.tsx b/web/src/app/dashboard/settings/scripts/[id]/script-detail.tsx index 1624515..b2cac0d 100644 --- a/web/src/app/dashboard/settings/scripts/[id]/script-detail.tsx +++ b/web/src/app/dashboard/settings/scripts/[id]/script-detail.tsx @@ -70,6 +70,7 @@ export default function ScriptDetail() { if (!isDirty) return; function onBeforeUnload(e: BeforeUnloadEvent) { e.preventDefault(); + e.returnValue = ""; } window.addEventListener("beforeunload", onBeforeUnload); return () => window.removeEventListener("beforeunload", onBeforeUnload); diff --git a/web/src/app/dashboard/settings/scripts/new/page.tsx b/web/src/app/dashboard/settings/scripts/new/page.tsx index dda5e0f..8ed4071 100644 --- a/web/src/app/dashboard/settings/scripts/new/page.tsx +++ b/web/src/app/dashboard/settings/scripts/new/page.tsx @@ -75,6 +75,7 @@ function NewScriptInner() { if (!isDirty) return; function onBeforeUnload(e: BeforeUnloadEvent) { e.preventDefault(); + e.returnValue = ""; } window.addEventListener("beforeunload", onBeforeUnload); return () => window.removeEventListener("beforeunload", onBeforeUnload); diff --git a/web/src/types/jobs.ts b/web/src/types/jobs.ts index c457980..74dcec1 100644 --- a/web/src/types/jobs.ts +++ b/web/src/types/jobs.ts @@ -1,18 +1,19 @@ 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; + id: string + job_type: string + status: string + input: unknown + progress: unknown + result: unknown | null + error: string | null + created_by: string + inherit_auth: boolean + started_at: string | null + completed_at: string | null + created_at: string } export interface JobsListResponse { - jobs: Job[]; - cursor: string | null; + jobs: Job[] + cursor: string | null } diff --git a/web/tests/e2e/script-job.spec.ts b/web/tests/e2e/script-job.spec.ts index adc1448..d10629c 100644 --- a/web/tests/e2e/script-job.spec.ts +++ b/web/tests/e2e/script-job.spec.ts @@ -4,6 +4,23 @@ import { loginAsTestAdmin } from "./auth-helper" const JOB_TYPE = "test.e2e.myjob" const TRIGGER_ID = `job.run:${JOB_TYPE}` +async function seedScript( + request: import("@playwright/test").APIRequestContext, +) { + const resp = await request.post("/admin/scripts", { + data: { + id: TRIGGER_ID, + body: "function handle()\n return { ok = true }\nend", + }, + }) + if (!resp.ok()) { + const text = await resp.text() + if (!text.includes("already exists")) { + throw new Error(`Failed to seed script: ${resp.status()} ${text}`) + } + } +} + async function cleanupScript( request: import("@playwright/test").APIRequestContext, ) { @@ -48,15 +65,21 @@ test.describe("Job Script Creation", () => { await page.locator("#job-type-input").fill(JOB_TYPE) - await page.getByRole("button", { name: "Create script" }).click() + const createButton = page.getByRole("button", { name: "Create script" }) + await expect(createButton).toBeEnabled({ timeout: 3000 }) + await createButton.click() await page.waitForURL( `**/dashboard/settings/scripts/${encodeURIComponent(TRIGGER_ID)}`, - { timeout: 5000 }, + { timeout: 10000 }, ) - await expect(page.getByText("Job runner")).toBeVisible() - await expect(page.getByText(TRIGGER_ID)).toBeVisible() + await expect( + page.getByText("Job runner", { exact: true }), + ).toBeVisible() + await expect( + page.getByText(TRIGGER_ID, { exact: true }), + ).toBeVisible() }) test("job script has job-specific template body", async ({ page }) => { @@ -65,26 +88,23 @@ test.describe("Job Script Creation", () => { 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() + await expect(page.getByText("job.input").first()).toBeVisible({ + timeout: 3000, + }) + await expect(page.getByText("job.should_stop").first()).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 seedScript(page.request) 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() + await expect(row.getByText("Job runner", { exact: true })).toBeVisible() + await expect(row.getByText("Job runners", { exact: true })).toBeVisible() }) })