Skip to content

Commit e2e9ab4

Browse files
Merge pull request #253 from appdevforall/feat/ADFA-4838-rest-job-engine
ADFA-4838: durable REST job engine for live content (kiwix)
2 parents e8c30a9 + 9c6b7dc commit e2e9ab4

8 files changed

Lines changed: 637 additions & 2 deletions

File tree

‎controller/docs/ADR-4832-live-content-channel.md‎

Lines changed: 29 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
# ADR — Content management on the live system (app ↔ in-server channel)
22

3-
Status: Accepted. Phase 1 (socket.io bridge) in progress; Phase 2 (REST + durable jobs) planned. Ticket: ADFA-4832.
3+
Status: Accepted. Phase 1 (socket.io bridge) shipped (ADFA-4832). Phase 2 (REST + durable jobs, full scope) accepted and in progress — see "Phase 2 (accepted) — design of record" below. Tickets: ADFA-4832 (epic thread), plus the Phase 2 tickets listed in that section.
44

55
## Context
66

@@ -29,6 +29,34 @@ Status: Accepted. Phase 1 (socket.io bridge) in progress; Phase 2 (REST + durabl
2929
- Maps/books on a live system still need the same migration (follow-up). Phase 1 skips the maps proot on the live path to avoid re-introducing a collision.
3030
- If the app process is fully killed, the job dies (server kills on disconnect). Phase 2 durable jobs remove this.
3131

32+
## Phase 2 (accepted) — design of record
33+
34+
Scope decision: **full Phase 2 in one rootfs rebuild** — kiwix + maps + books on the REST engine, IPv4/aria2 parity, and the dashboard packaging (build + supervision). Iteration decision: a **dev "push static/" path** copies the updated dashboard into the installed rootfs on-device and restarts `dash-node`, so we don't pay the ~2h rebuild per test; the full rebuild happens once at the end.
35+
36+
### Server (in-server dashboard)
37+
38+
1. **Job manager, module-scoped (not per-socket).** The download+index job is owned by the dashboard process, not by any connection. A client `disconnect` no longer kills anything (Phase 1 killed on disconnect). This is the core durability change.
39+
2. **Durable journal via `better-sqlite3`** (already a dependency). Table `jobs(id, type, target, phase, percent, speed, error, created, updated)`. On dashboard boot, reconcile: incomplete download jobs resume (aria2 `--continue` already resumes partials), indexing re-runs if it hadn't finished.
40+
3. **Progress parsed server-side into structured fields.** aria2/index stdout is parsed on the server into `{phase: downloading|indexing|done|error, percent, speedBytesPerSec, error}` and stored on the job. Clients read structured JSON — no more client-side terminal-scraping.
41+
4. **REST contract (localhost:8085 → :4000), short calls:**
42+
- `POST /api/{kiwix|maps|books}/download {ids:[...]}` → `{jobId}`
43+
- `GET /api/{...}/jobs/:id` → the structured job row
44+
- `POST /api/{...}/jobs/:id/cancel`
45+
- `GET /api/{...}/catalog`, `DELETE /api/{...}/item/:id`
46+
The socket.io handlers stay only as a thin compat shell (or are removed) — the job manager is the single source of truth.
47+
5. **aria2 IPv4 profiler parity.** Port the app's IPv4-preference behavior to the server aria2 spawn (closes a divergence noted below).
48+
6. **Packaging.** Compile TS→JS (drop `ts-node` dev mode), run under `pdsm` supervision with a health check; `dash-node` restarts cleanly on failure. This — not a separate cmdsrv — is the answer to "Node is fragile".
49+
50+
### Android
51+
52+
7. **REST content client replaces the socket.io `LiveContentClient`.** Short localhost calls: `POST` to start, **poll `GET` every ~1s** for structured progress, `POST cancel`. Owned by the foreground `InstallService`. There is no long-lived connection to lose; the job survives UI/config churn and even a dashboard/socket restart (durable job + `--continue`). Wire into `InstallService.downloadAndIndexKiwix` and the maps/books paths.
53+
54+
### Tickets
55+
56+
- Dashboard: durable REST job engine + endpoints + sqlite + aria2 IPv4 parity (kiwix/maps/books).
57+
- Dashboard packaging: TS→JS build + pdsm supervision/health + dev push-static path.
58+
- Android: REST content client (foreground poll) replacing socket.io Phase 1.
59+
3260
## References
3361

3462
`PRootEngine.java`, `InstallService.downloadAndIndexKiwix`, `ServerController`, `static/dashboard/server.ts` + `sockets/kiwix.socket.ts`, `LiveContentClient.java`. Ticket ADFA-4832.

‎static/dashboard/dash-node-nginx.conf‎

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,21 @@ location ^~ /dashboard/ {
1515
proxy_set_header X-Forwarded-Proto $scheme;
1616
}
1717

18+
# Proxy for the REST content API (ADFA-4838, localhost only)
19+
location ^~ /api/ {
20+
allow 127.0.0.1;
21+
allow ::1;
22+
deny all;
23+
24+
error_page 403 =404 /404.html;
25+
26+
proxy_pass http://127.0.0.1:4000/api/;
27+
proxy_set_header Host $http_host;
28+
proxy_set_header X-Real-IP $remote_addr;
29+
proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for;
30+
proxy_set_header X-Forwarded-Proto $scheme;
31+
}
32+
1833
# Proxy for WebSockets (socket.io uses it by default)
1934
location ^~ /socket.io/ {
2035
# Basic security (localhost only)

‎static/dashboard/routes.ts‎

Lines changed: 67 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,67 @@
1+
// routes.ts — ADFA-4838
2+
//
3+
// REST surface over the durable job engine. Short, stateless calls the app (and the
4+
// web UI) use instead of a long-lived socket: start a job, poll its structured status,
5+
// cancel it. The job itself lives in the dashboard process (see sockets/jobs.ts), so
6+
// none of these calls hold state — a client can drop and re-attach by polling the id.
7+
import express, { Router, Request, Response } from 'express';
8+
import { jobs, Job, JobType } from './sockets/jobs';
9+
10+
const VALID_TYPES: JobType[] = ['kiwix', 'maps', 'books'];
11+
function isType(t: string): t is JobType {
12+
return (VALID_TYPES as string[]).includes(t);
13+
}
14+
15+
/** Public shape returned to clients: ids as an array, no raw JSON column. */
16+
function toApi(job: Job) {
17+
let ids: string[] = [];
18+
try { ids = JSON.parse(job.target) as string[]; } catch { ids = []; }
19+
return {
20+
id: job.id,
21+
type: job.type,
22+
ids,
23+
phase: job.phase,
24+
percent: job.percent,
25+
speed: job.speed,
26+
detail: job.detail,
27+
error: job.error,
28+
updated: job.updated,
29+
};
30+
}
31+
32+
export const apiRouter: Router = express.Router();
33+
34+
// Start a content job → 202 { ...job }
35+
apiRouter.post('/:type/download', (req: Request, res: Response): void => {
36+
const type = String(req.params.type);
37+
if (!isType(type)) { res.status(404).json({ error: 'unknown type' }); return; }
38+
// kiwix sends { ids: ["file.zim"] }; maps/books send { items: [ {...} ] }.
39+
const body = req.body as { ids?: unknown; items?: unknown };
40+
const items: unknown[] = Array.isArray(body?.items)
41+
? body.items
42+
: Array.isArray(body?.ids) ? body.ids : [];
43+
if (items.length === 0) { res.status(400).json({ error: 'items (or ids) required' }); return; }
44+
res.status(202).json(toApi(jobs.create(type, items)));
45+
});
46+
47+
// Poll one job's structured status.
48+
apiRouter.get('/:type/jobs/:id', (req: Request, res: Response): void => {
49+
const job = jobs.get(String(req.params.id));
50+
if (!job || job.type !== String(req.params.type)) { res.status(404).json({ error: 'not found' }); return; }
51+
res.json(toApi(job));
52+
});
53+
54+
// List a type's jobs (most recent first).
55+
apiRouter.get('/:type/jobs', (req: Request, res: Response): void => {
56+
const type = String(req.params.type);
57+
if (!isType(type)) { res.status(404).json({ error: 'unknown type' }); return; }
58+
res.json(jobs.list(type).map(toApi));
59+
});
60+
61+
// Cancel a running job.
62+
apiRouter.post('/:type/jobs/:id/cancel', (req: Request, res: Response): void => {
63+
const job = jobs.get(String(req.params.id));
64+
if (!job || job.type !== String(req.params.type)) { res.status(404).json({ error: 'not found' }); return; }
65+
jobs.cancel(job.id);
66+
res.json({ ok: true });
67+
});

‎static/dashboard/server.ts‎

Lines changed: 15 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,14 @@ import { handleKiwixEvents } from './sockets/kiwix.socket';
1010
import { handleHomeEvents } from './sockets/home.socket';
1111
import { handleBooksEvents } from './sockets/books.socket';
1212

13+
// ADFA-4838: durable REST job engine. Importing the runner modules registers them with
14+
// the shared JobManager; the router exposes them over /api.
15+
import { jobs } from './sockets/jobs';
16+
import './sockets/kiwix.exec';
17+
import './sockets/maps.exec';
18+
import './sockets/books.exec';
19+
import { apiRouter } from './routes';
20+
1321
const app = express();
1422
const server = http.createServer(app);
1523
const io = new Server(server);
@@ -28,9 +36,13 @@ app.set('view engine', 'ejs');
2836
app.set('views', path.join(__dirname, 'views'));
2937
app.use(express.static(path.join(__dirname, 'public')));
3038

39+
// ADFA-4838: JSON body parsing + the REST content API.
40+
app.use(express.json());
41+
app.use('/api', apiRouter);
42+
3143
// Main route
3244
app.get('/', (req, res) => {
33-
res.render('index');
45+
res.render('index');
3446
});
3547

3648
// Main Socket Connection Handler
@@ -53,6 +65,8 @@ server.listen(PORT, () => {
5365
console.log(`===========================================`);
5466
console.log(`K2Go Dashboard active on port ${PORT}`);
5567
console.log(`===========================================`);
68+
// ADFA-4838: resume any content jobs that were mid-flight before a restart.
69+
try { jobs.reconcileOnBoot(); } catch (e) { console.error('[jobs] reconcile failed', e); }
5670
});
5771

5872
// ==========================================
Lines changed: 120 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,120 @@
1+
// sockets/books.exec.ts — ADFA-4838
2+
//
3+
// Books runner for the durable job engine: for each requested Gutenberg book, download
4+
// the EPUB and upload it to the local Calibre-Web. Ported from the Phase 1 books.socket
5+
// download_books_batch handler (auth/CSRF + fetch + upload), made durable and reporting
6+
// structured per-book progress. A job item is { id, title, url }.
7+
import { jobs, RunnerContext } from './jobs';
8+
import fs from 'fs';
9+
import path from 'path';
10+
11+
const CALIBRE_WEB_LOCAL_URL = 'http://127.0.0.1:8083';
12+
const TMP_DIR = '/tmp/books_downloader/';
13+
const SYSTEM_USER_AGENT = 'K2Go Dashboard/1.0 (https://github.com/appdevforall/KnowledgeToGo)';
14+
// Defaults match the Phase 1 handler; a credential override can be added later if needed.
15+
const CALIBRE_WEB_USER = 'Admin';
16+
const CALIBRE_WEB_PASS = 'changeme';
17+
18+
interface BookItem { id?: string; title?: string; url?: string; }
19+
20+
/** Authenticate against Calibre-Web and return a usable cookie + fresh CSRF token. */
21+
async function getCalibreSession(): Promise<{ cookie: string; csrfToken: string }> {
22+
const loginPageRes = await fetch(`${CALIBRE_WEB_LOCAL_URL}/login`);
23+
const initialCookies = loginPageRes.headers.getSetCookie().map((c) => c.split(';')[0]).join('; ');
24+
const loginHtml = await loginPageRes.text();
25+
26+
const csrfMatch = loginHtml.match(/name="csrf_token" value="(.*?)"/);
27+
if (!csrfMatch) throw new Error('Could not find CSRF token on login page');
28+
const csrfToken = csrfMatch[1];
29+
30+
const loginData = new URLSearchParams();
31+
loginData.append('csrf_token', csrfToken);
32+
loginData.append('username', CALIBRE_WEB_USER);
33+
loginData.append('password', CALIBRE_WEB_PASS);
34+
35+
const authRes = await fetch(`${CALIBRE_WEB_LOCAL_URL}/login`, {
36+
method: 'POST',
37+
headers: {
38+
Cookie: initialCookies,
39+
'Content-Type': 'application/x-www-form-urlencoded',
40+
Referer: `${CALIBRE_WEB_LOCAL_URL}/login`,
41+
},
42+
body: loginData,
43+
redirect: 'manual',
44+
});
45+
46+
if (authRes.status !== 302 && authRes.status !== 303) {
47+
throw new Error('Invalid Calibre-Web credentials');
48+
}
49+
50+
const authCookieString = authRes.headers.getSetCookie().map((c) => c.split(';')[0]).join('; ');
51+
52+
const homePageRes = await fetch(`${CALIBRE_WEB_LOCAL_URL}/`, { headers: { Cookie: authCookieString } });
53+
const homeHtml = await homePageRes.text();
54+
const finalCsrfMatch =
55+
homeHtml.match(/name="csrf_token"\s+value="([^"]+)"/i) ||
56+
homeHtml.match(/value="([^"]+)"\s+name="csrf_token"/i);
57+
const finalCsrfToken = finalCsrfMatch ? finalCsrfMatch[1] : csrfToken;
58+
59+
return { cookie: authCookieString, csrfToken: finalCsrfToken };
60+
}
61+
62+
const booksRunner: (ctx: RunnerContext) => Promise<void> = async (ctx) => {
63+
if (!fs.existsSync(TMP_DIR)) fs.mkdirSync(TMP_DIR, { recursive: true });
64+
65+
const books = ctx.items
66+
.map((x) => x as BookItem)
67+
.filter((b) => b && typeof b.url === 'string' && typeof b.id === 'string');
68+
if (books.length === 0) throw new Error('no books requested');
69+
70+
ctx.update({ phase: 'processing', percent: 0 });
71+
72+
let session: { cookie: string; csrfToken: string };
73+
try {
74+
session = await getCalibreSession();
75+
} catch (e) {
76+
throw new Error('Calibre-Web authentication failed');
77+
}
78+
79+
let done = 0;
80+
for (const book of books) {
81+
ctx.throwIfCanceled();
82+
const id = String(book.id);
83+
const title = String(book.title ?? id);
84+
const url = String(book.url);
85+
ctx.update({ phase: 'processing', detail: title, percent: Math.round((done / books.length) * 100) });
86+
87+
const tmp = path.join(TMP_DIR, `pg_${id}.epub`);
88+
try {
89+
const response = await fetch(url, {
90+
headers: { 'User-Agent': SYSTEM_USER_AGENT, Accept: 'application/epub+zip' },
91+
});
92+
if (!response.ok) throw new Error(`HTTP ${response.status} from Gutenberg`);
93+
94+
const fileBuffer = await response.arrayBuffer();
95+
fs.writeFileSync(tmp, Buffer.from(fileBuffer));
96+
97+
const form = new FormData();
98+
form.append('csrf_token', session.csrfToken);
99+
form.append('btn-upload', new Blob([fileBuffer], { type: 'application/epub+zip' }), `${title}.epub`);
100+
101+
const uploadRes = await fetch(`${CALIBRE_WEB_LOCAL_URL}/upload`, {
102+
method: 'POST',
103+
headers: { Cookie: session.cookie, Referer: `${CALIBRE_WEB_LOCAL_URL}/` },
104+
body: form,
105+
});
106+
if (!uploadRes.ok) throw new Error(`Calibre-Web rejected upload: ${uploadRes.status}`);
107+
} finally {
108+
if (fs.existsSync(tmp)) fs.unlinkSync(tmp);
109+
}
110+
111+
done++;
112+
ctx.update({ phase: 'processing', percent: Math.round((done / books.length) * 100) });
113+
}
114+
115+
ctx.update({ phase: 'done', percent: 100 });
116+
};
117+
118+
jobs.registerRunner('books', booksRunner);
119+
120+
export { booksRunner };

0 commit comments

Comments
 (0)