diff --git a/bun.lock b/bun.lock
--- a/bun.lock
+++ b/bun.lock
@@ -26,7 +26,7 @@
"fumadocs-mdx": "^15.0.4",
"fumadocs-ui": "^16.8.10",
"lucide-react": "^1.14.0",
- "mermaid": "^11.6.0",
+ "mermaid": "^11.16.0",
"next": "^16.1.6",
"next-themes": "^0.4.6",
"react": "^19.2.0",
@@ -219,15 +219,7 @@
"@braintree/sanitize-url": ["@braintree/sanitize-url@7.1.2", "", {}, "sha512-jigsZK+sMF/cuiB7sERuo9V7N9jx+dhmHHnQyDSVdpZwVutaBu7WvNYqMDLSgFgfB30n452TP3vjDAvFC973mA=="],
- "@chevrotain/cst-dts-gen": ["@chevrotain/cst-dts-gen@12.0.0", "", { "dependencies": { "@chevrotain/gast": "12.0.0", "@chevrotain/types": "12.0.0" } }, "sha512-fSL4KXjTl7cDgf0B5Rip9Q05BOrYvkJV/RrBTE/bKDN096E4hN/ySpcBK5B24T76dlQ2i32Zc3PAE27jFnFrKg=="],
-
- "@chevrotain/gast": ["@chevrotain/gast@12.0.0", "", { "dependencies": { "@chevrotain/types": "12.0.0" } }, "sha512-1ne/m3XsIT8aEdrvT33so0GUC+wkctpUPK6zU9IlOyJLUbR0rg4G7ZiApiJbggpgPir9ERy3FRjT6T7lpgetnQ=="],
-
- "@chevrotain/regexp-to-ast": ["@chevrotain/regexp-to-ast@12.0.0", "", {}, "sha512-p+EW9MaJwgaHguhoqwOtx/FwuGr+DnNn857sXWOi/mClXIkPGl3rn7hGNWvo31HA3vyeQxjqe+H36yZJwYU8cA=="],
-
- "@chevrotain/types": ["@chevrotain/types@12.0.0", "", {}, "sha512-S+04vjFQKeuYw0/eW3U52LkAHQsB1ASxsPGsLPUyQgrZ2iNNibQrsidruDzjEX2JYfespXMG0eZmXlhA6z7nWA=="],
-
- "@chevrotain/utils": ["@chevrotain/utils@12.0.0", "", {}, "sha512-lB59uJoaGIfOOL9knQqQRfhl9g7x8/wqFkp13zTdkRu1huG9kg6IJs1O8hqj9rs6h7orGxHJUKb+mX3rPbWGhA=="],
+ "@chevrotain/types": ["@chevrotain/types@11.1.2", "", {}, "sha512-U+HFai5+zmJCkK86QsaJtoITlboZHBqrVketcO2ROv865xfCMSFpELQoz1GkX5GzME8pTa+3kbKrZHQtI0gdbw=="],
"@clack/core": ["@clack/core@1.4.2", "", { "dependencies": { "fast-wrap-ansi": "^0.2.0", "sisteransi": "^1.0.5" } }, "sha512-0Ty/1Gfm+Kb07sXcuESjyKfwEhSy4Ns1AgeEisHb/bDY5fWme0tTeTkU14T1Gmcs17YIjB/teiDe4uaCghbYqQ=="],
@@ -379,7 +371,7 @@
"@mdx-js/mdx": ["@mdx-js/mdx@3.1.1", "", { "dependencies": { "@types/estree": "^1.0.0", "@types/estree-jsx": "^1.0.0", "@types/hast": "^3.0.0", "@types/mdx": "^2.0.0", "acorn": "^8.0.0", "collapse-white-space": "^2.0.0", "devlop": "^1.0.0", "estree-util-is-identifier-name": "^3.0.0", "estree-util-scope": "^1.0.0", "estree-walker": "^3.0.0", "hast-util-to-jsx-runtime": "^2.0.0", "markdown-extensions": "^2.0.0", "recma-build-jsx": "^1.0.0", "recma-jsx": "^1.0.0", "recma-stringify": "^1.0.0", "rehype-recma": "^1.0.0", "remark-mdx": "^3.0.0", "remark-parse": "^11.0.0", "remark-rehype": "^11.0.0", "source-map": "^0.7.0", "unified": "^11.0.0", "unist-util-position-from-estree": "^2.0.0", "unist-util-stringify-position": "^4.0.0", "unist-util-visit": "^5.0.0", "vfile": "^6.0.0" } }, "sha512-f6ZO2ifpwAQIpzGWaBQT2TXxPv6z3RBzQKpVftEWN78Vl/YweF1uwussDx8ECAXVtr3Rs89fKyG9YlzUs9DyGQ=="],
- "@mermaid-js/parser": ["@mermaid-js/parser@1.1.0", "", { "dependencies": { "langium": "^4.0.0" } }, "sha512-gxK9ZX2+Fex5zu8LhRQoMeMPEHbc73UKZ0FQ54YrQtUxE1VVhMwzeNtKRPAu5aXks4FasbMe4xB4bWrmq6Jlxw=="],
+ "@mermaid-js/parser": ["@mermaid-js/parser@1.2.0", "", { "dependencies": { "@chevrotain/types": "~11.1.2" } }, "sha512-oYPyv8A4As1yH5Bx+04iQEQxXuIQDe0GKCNSRgao6z8AM9jixXIfP0vsppRLvGf+nKIOb9/LdpWA4YuJiVvESA=="],
"@next/env": ["@next/env@16.2.6", "", {}, "sha512-gd8HoHN4ufj73WmR3JmVolrpJR47ILK6LouP5xElPglaVxir6e1a7VzvTvDWkOoPXT9rkkTzyCxBu4yeZfZwcw=="],
@@ -801,10 +793,6 @@
"character-reference-invalid": ["character-reference-invalid@2.0.1", "", {}, "sha512-iBZ4F4wRbyORVsu0jPV7gXkOsGYjGHPmAyv+HiHG8gi5PtC9KI2j1+v8/tlibRvjoWX027ypmG/n0HtO5t7unw=="],
- "chevrotain": ["chevrotain@12.0.0", "", { "dependencies": { "@chevrotain/cst-dts-gen": "12.0.0", "@chevrotain/gast": "12.0.0", "@chevrotain/regexp-to-ast": "12.0.0", "@chevrotain/types": "12.0.0", "@chevrotain/utils": "12.0.0" } }, "sha512-csJvb+6kEiQaqo1woTdSAuOWdN0WTLIydkKrBnS+V5gZz0oqBrp4kQ35519QgK6TpBThiG3V1vNSHlIkv4AglQ=="],
-
- "chevrotain-allstar": ["chevrotain-allstar@0.4.1", "", { "dependencies": { "lodash-es": "^4.17.21" }, "peerDependencies": { "chevrotain": "^12.0.0" } }, "sha512-PvVJm3oGqrveUVW2Vt/eZGeiAIsJszYweUcYwcskg9e+IubNYKKD+rHHem7A6XVO22eDAL+inxNIGAzZ/VIWlA=="],
-
"chokidar": ["chokidar@5.0.0", "", { "dependencies": { "readdirp": "^5.0.0" } }, "sha512-TQMmc3w+5AxjpL8iIiwebF73dRDF4fBIieAqGn9RGCWaEVwQ6Fb2cGe31Yns0RRIzii5goJ1Y7xbMwo1TxMplw=="],
"class-variance-authority": ["class-variance-authority@0.7.1", "", { "dependencies": { "clsx": "^2.1.1" } }, "sha512-Ka+9Trutv7G8M6WT6SeiRWz792K5qEqIGEGzXKhAE6xOWAY6pPH8U+9IY3oCMv6kqTmLsv7Xh/2w2RigkePMsg=="],
@@ -871,7 +859,7 @@
"csstype": ["csstype@3.2.3", "", {}, "sha512-z1HGKcYy2xA8AGQfwrn0PAy+PB7X/GSj3UVJW9qKyn43xWa+gl5nXmU4qqLMRzWVLFC8KusUX8T/0kCiOYpAIQ=="],
- "cytoscape": ["cytoscape@3.33.2", "", {}, "sha512-sj4HXd3DokGhzZAdjDejGvTPLqlt84vNFN8m7bGsOzDY5DyVcxIb2ejIXat2Iy7HxWhdT/N1oKyheJ5YdpsGuw=="],
+ "cytoscape": ["cytoscape@3.34.0", "", {}, "sha512-62rNSrioXw93uliKFBwjukeQyeWwH2PqDrTac31r2P6464u3AUvTk0xS4LVvT251g7IgkFunrI48ZEZGjywSOg=="],
"cytoscape-cose-bilkent": ["cytoscape-cose-bilkent@4.1.0", "", { "dependencies": { "cose-base": "^1.0.0" }, "peerDependencies": { "cytoscape": "^3.2.0" } }, "sha512-wgQlVIUJF13Quxiv5e1gstZ08rnZj2XaLHGoFMYXz7SkNfCDOOteKBE6SYRfA9WxxI/iBc3ajfDoc6hb/MRAHQ=="],
@@ -996,6 +984,8 @@
"error-ex": ["error-ex@1.3.4", "", { "dependencies": { "is-arrayish": "^0.2.1" } }, "sha512-sqQamAnR14VgCr1A618A3sGrygcpK+HEbenA/HiEAkkUwcZIIB/tgWqHFxWgOyDh4nB4JCRimh79dR5Ywc9MDQ=="],
"es-errors": ["es-errors@1.3.0", "", {}, "sha512-Zf5H2Kxt2xjTvbJvP2ZWLEICxA6j+hAmMzIlypy4xcBg1vKVnx89Wy0GbS+kf5cwCVFFzdCFh2XSCFNULS6csw=="],
+
+ "es-toolkit": ["es-toolkit@1.49.0", "", {}, "sha512-G5iZ6Pc/FNRY/soKZHC+TxGDD83rHUDXxzaWhGCX44vAv/tMs56WMusnm/KMNK+luUPsgA9U28cGr4RDlSzL2g=="],
"esast-util-from-estree": ["esast-util-from-estree@2.0.0", "", { "dependencies": { "@types/estree-jsx": "^1.0.0", "devlop": "^1.0.0", "estree-util-visit": "^2.0.0", "unist-util-position-from-estree": "^2.0.0" } }, "sha512-4CyanoAudUSBAn5K13H4JhsMH6L9ZP7XbLVe/dKybkxMO7eDyLsT8UHl9TRNrU2Gr9nz+FovfSIjuXWJ81uVwQ=="],
@@ -1245,8 +1235,6 @@
"khroma": ["khroma@2.1.0", "", {}, "sha512-Ls993zuzfayK269Svk9hzpeGUKob/sIgZzyHYdjQoAdQetRKpOLj+k/QQQ/6Qi0Yz65mlROrfd+Ev+1+7dz9Kw=="],
- "langium": ["langium@4.2.2", "", { "dependencies": { "@chevrotain/regexp-to-ast": "~12.0.0", "chevrotain": "~12.0.0", "chevrotain-allstar": "~0.4.1", "vscode-languageserver": "~9.0.1", "vscode-languageserver-textdocument": "~1.0.11", "vscode-uri": "~3.1.0" } }, "sha512-JUshTRAfHI4/MF9dH2WupvjSXyn8JBuUEWazB8ZVJUtXutT0doDlAv1XKbZ1Pb5sMexa8FF4CFBc0iiul7gbUQ=="],
-
"layout-base": ["layout-base@1.0.2", "", {}, "sha512-8h2oVEZNktL4BH2JCOI90iD1yXwL6iNW7KcCKT2QZgQJR2vbqDsldCTPRU9NifTCqHZci57XvQQ15YTu+sTYPg=="],
"lightningcss": ["lightningcss@1.32.0", "", { "dependencies": { "detect-libc": "^2.0.3" }, "optionalDependencies": { "lightningcss-android-arm64": "1.32.0", "lightningcss-darwin-arm64": "1.32.0", "lightningcss-darwin-x64": "1.32.0", "lightningcss-freebsd-x64": "1.32.0", "lightningcss-linux-arm-gnueabihf": "1.32.0", "lightningcss-linux-arm64-gnu": "1.32.0", "lightningcss-linux-arm64-musl": "1.32.0", "lightningcss-linux-x64-gnu": "1.32.0", "lightningcss-linux-x64-musl": "1.32.0", "lightningcss-win32-arm64-msvc": "1.32.0", "lightningcss-win32-x64-msvc": "1.32.0" } }, "sha512-NXYBzinNrblfraPGyrbPoD19C1h9lfI/1mzgWYvXUTe414Gz/X1FD2XBZSZM7rRTrMA8JL3OtAaGifrIKhQ5yQ=="],
@@ -1351,7 +1339,7 @@
"merge2": ["merge2@1.4.1", "", {}, "sha512-8q7VEgMJW4J8tcfVPy8g09NcQwZdbwFEqhe/WZkoIzjn/3TGDwtOCYtXGxA3O8tPzpczCCDgv+P2P5y00ZJOOg=="],
- "mermaid": ["mermaid@11.14.0", "", { "dependencies": { "@braintree/sanitize-url": "^7.1.1", "@iconify/utils": "^3.0.2", "@mermaid-js/parser": "^1.1.0", "@types/d3": "^7.4.3", "@upsetjs/venn.js": "^2.0.0", "cytoscape": "^3.33.1", "cytoscape-cose-bilkent": "^4.1.0", "cytoscape-fcose": "^2.2.0", "d3": "^7.9.0", "d3-sankey": "^0.12.3", "dagre-d3-es": "7.0.14", "dayjs": "^1.11.19", "dompurify": "^3.3.1", "katex": "^0.16.25", "khroma": "^2.1.0", "lodash-es": "^4.17.23", "marked": "^16.3.0", "roughjs": "^4.6.6", "stylis": "^4.3.6", "ts-dedent": "^2.2.0", "uuid": "^11.1.0" } }, "sha512-GSGloRsBs+JINmmhl0JDwjpuezCsHB4WGI4NASHxL3fHo3o/BRXTxhDLKnln8/Q0lRFRyDdEjmk1/d5Sn1Xz8g=="],
+ "mermaid": ["mermaid@11.16.0", "", { "dependencies": { "@braintree/sanitize-url": "^7.1.2", "@iconify/utils": "^3.0.2", "@mermaid-js/parser": "^1.2.0", "@types/d3": "^7.4.3", "@upsetjs/venn.js": "^2.0.0", "cytoscape": "^3.33.3", "cytoscape-cose-bilkent": "^4.1.0", "cytoscape-fcose": "^2.2.0", "d3": "^7.9.0", "d3-sankey": "^0.12.3", "dagre-d3-es": "7.0.14", "dayjs": "^1.11.20", "dompurify": "^3.3.3", "es-toolkit": "^1.45.1", "katex": "^0.16.45", "khroma": "^2.1.0", "marked": "^16.3.0", "roughjs": "^4.6.6", "stylis": "^4.3.6", "ts-dedent": "^2.2.0", "uuid": "^11.1.0 || ^12 || ^13 || ^14.0.0" } }, "sha512-Zvm3kbstgdpvIJPPItlL7fppIZ3kibvc1oZIGxdvk9t6UFz6flv+Jw7FtRGKwfcI8OckmH04LqG6LlS6X4B1pA=="],
"micromark": ["micromark@4.0.2", "", { "dependencies": { "@types/debug": "^4.0.0", "debug": "^4.0.0", "decode-named-character-reference": "^1.0.0", "devlop": "^1.0.0", "micromark-core-commonmark": "^2.0.0", "micromark-factory-space": "^2.0.0", "micromark-util-character": "^2.0.0", "micromark-util-chunked": "^2.0.0", "micromark-util-combine-extensions": "^2.0.0", "micromark-util-decode-numeric-character-reference": "^2.0.0", "micromark-util-encode": "^2.0.0", "micromark-util-normalize-identifier": "^2.0.0", "micromark-util-resolve-all": "^2.0.0", "micromark-util-sanitize-uri": "^2.0.0", "micromark-util-subtokenize": "^2.0.0", "micromark-util-symbol": "^2.0.0", "micromark-util-types": "^2.0.0" } }, "sha512-zpe98Q6kvavpCr1NPVSCMebCKfD7CA2NqZ+rykeNhONIJBpc1tFKt9hucLGwha3jNTNI8lHpctWJWoimVF4PfA=="],
@@ -1872,18 +1860,6 @@
"vfile-location": ["vfile-location@5.0.3", "", { "dependencies": { "@types/unist": "^3.0.0", "vfile": "^6.0.0" } }, "sha512-5yXvWDEgqeiYiBe1lbxYF7UMAIm/IcopxMHrMQDq3nvKcjPKIhZklUKL+AE7J7uApI4kwe2snsK+eI6UTj9EHg=="],
"vfile-message": ["vfile-message@4.0.3", "", { "dependencies": { "@types/unist": "^3.0.0", "unist-util-stringify-position": "^4.0.0" } }, "sha512-QTHzsGd1EhbZs4AsQ20JX1rC3cOlt/IWJruk893DfLRr57lcnOeMaWG4K0JrRta4mIJZKth2Au3mM3u03/JWKw=="],
-
- "vscode-jsonrpc": ["vscode-jsonrpc@8.2.0", "", {}, "sha512-C+r0eKJUIfiDIfwJhria30+TYWPtuHJXHtI7J0YlOmKAo7ogxP20T0zxB7HZQIFhIyvoBPwWskjxrvAtfjyZfA=="],
-
- "vscode-languageserver": ["vscode-languageserver@9.0.1", "", { "dependencies": { "vscode-languageserver-protocol": "3.17.5" }, "bin": { "installServerIntoExtension": "bin/installServerIntoExtension" } }, "sha512-woByF3PDpkHFUreUa7Hos7+pUWdeWMXRd26+ZX2A8cFx6v/JPTtd4/uN0/jB6XQHYaOlHbio03NTHCqrgG5n7g=="],
-
- "vscode-languageserver-protocol": ["vscode-languageserver-protocol@3.17.5", "", { "dependencies": { "vscode-jsonrpc": "8.2.0", "vscode-languageserver-types": "3.17.5" } }, "sha512-mb1bvRJN8SVznADSGWM9u/b07H7Ecg0I3OgXDuLdn307rl/J3A9YD6/eYOssqhecL27hK1IPZAsaqh00i/Jljg=="],
-
- "vscode-languageserver-textdocument": ["vscode-languageserver-textdocument@1.0.12", "", {}, "sha512-cxWNPesCnQCcMPeenjKKsOCKQZ/L6Tv19DTRIGuLWe32lyzWhihGVJ/rcckZXJxfdKCFvRLS3fpBIsV/ZGX4zA=="],
-
- "vscode-languageserver-types": ["vscode-languageserver-types@3.17.5", "", {}, "sha512-Ld1VelNuX9pdF39h2Hgaeb5hEZM2Z3jUrrMgWQAu82jMtZp7p3vJT3BzToKtZI7NgQssZje5o0zryOrhQvzQAg=="],
-
- "vscode-uri": ["vscode-uri@3.1.0", "", {}, "sha512-/BpdSx+yCQGnCvecbyXdxHDkuk55/G3xwnC0GqY4gmQ3j+A+g8kzzgB4Nk/SINjqn6+waqw3EgbVF2QKExkRxQ=="],
"web-namespaces": ["web-namespaces@2.0.1", "", {}, "sha512-bKr1DkiNa2krS7qxNtdrtHAmzuYGFQLiQ13TsorsdT6ULTkPLKuu5+GsFpDlg6JFjUTwX2DyhMPG2be8uPrqsQ=="],
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_409() {
+ 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::CONFLICT);
+}
+
+// ---------------------------------------------------------------------------
+// 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_409() {
+ 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::CONFLICT);
+}
+
+// ---------------------------------------------------------------------------
+// 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_409() {
+ 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::CONFLICT);
+}
+
+// ---------------------------------------------------------------------------
+// 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/tests/spaces_db.rs b/tests/spaces_db.rs
--- a/tests/spaces_db.rs
+++ b/tests/spaces_db.rs
@@ -307,6 +307,7 @@
} else {
None
},
+ value: None,
created_at: now_rfc3339(),
};
oplog::append_op(&pool, backend, &entry)
@@ -356,6 +357,7 @@
rkey: format!("item-{i}"),
cid: Some(format!("bafy{i}")),
prev: None,
+ value: None,
created_at: now_rfc3339(),
};
oplog::append_op(&pool, backend, &entry)
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/packages/docs/package.json b/packages/docs/package.json
--- a/packages/docs/package.json
+++ b/packages/docs/package.json
@@ -20,7 +20,7 @@
"fumadocs-mdx": "^15.0.4",
"fumadocs-ui": "^16.8.10",
"lucide-react": "^1.14.0",
- "mermaid": "^11.6.0",
+ "mermaid": "^11.16.0",
"next": "^16.1.6",
"next-themes": "^0.4.6",
"react": "^19.2.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).clamp(1, 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::Conflict(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::Conflict(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::Conflict(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,305 @@
+use serde_json::Value;
+use uuid::Uuid;
+
+use crate::AppState;
+use crate::db::{DatabaseBackend, 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