Skip to content

Commit 9c6b7dc

Browse files
ADFA-4838: maps + books runners (migrate to the REST engine like kiwix)
Port the remaining live-content types off the Phase 1 socket handlers onto the durable job engine, so all three run the same way. - Generalize the job payload: items are now arbitrary (strings for kiwix ZIMs; objects for maps/books). RunnerContext exposes items (raw) + ids (string convenience); routes accept { items } or { ids }. - maps.exec.ts: tile-extract.py extract, reusing the socket handler's box validation; structured progress (parses a % when the script prints one). - books.exec.ts: Calibre-Web auth/CSRF + Gutenberg fetch + upload, per-book progress. Deletion (maps/books) stays a separate op for now. Typechecks (tsc --noEmit).
1 parent 30ac091 commit 9c6b7dc

5 files changed

Lines changed: 195 additions & 10 deletions

File tree

‎static/dashboard/routes.ts‎

Lines changed: 7 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -35,10 +35,13 @@ export const apiRouter: Router = express.Router();
3535
apiRouter.post('/:type/download', (req: Request, res: Response): void => {
3636
const type = String(req.params.type);
3737
if (!isType(type)) { res.status(404).json({ error: 'unknown type' }); return; }
38-
const body = req.body as { ids?: unknown };
39-
const ids = Array.isArray(body?.ids) ? body.ids.map((x) => String(x)) : [];
40-
if (ids.length === 0) { res.status(400).json({ error: 'ids required' }); return; }
41-
res.status(202).json(toApi(jobs.create(type, ids)));
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)));
4245
});
4346

4447
// Poll one job's structured status.

‎static/dashboard/server.ts‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -14,6 +14,8 @@ import { handleBooksEvents } from './sockets/books.socket';
1414
// the shared JobManager; the router exposes them over /api.
1515
import { jobs } from './sockets/jobs';
1616
import './sockets/kiwix.exec';
17+
import './sockets/maps.exec';
18+
import './sockets/books.exec';
1719
import { apiRouter } from './routes';
1820

1921
const app = express();
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 };

‎static/dashboard/sockets/jobs.ts‎

Lines changed: 11 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -42,6 +42,9 @@ export interface JobUpdate {
4242
/** Passed to a runner: report progress, check cancellation, spawn tracked children. */
4343
export interface RunnerContext {
4444
readonly job: Job;
45+
/** Raw request items (strings for kiwix; objects for maps/books). */
46+
readonly items: unknown[];
47+
/** Convenience view: just the string items (kiwix ZIM filenames). */
4548
readonly ids: string[];
4649
update(patch: JobUpdate): void;
4750
isCanceled(): boolean;
@@ -88,12 +91,13 @@ class JobManager {
8891
this.runners.set(type, runner);
8992
}
9093

91-
/** Create + start a job. Returns the persisted row immediately (phase 'queued'). */
92-
create(type: JobType, ids: string[]): Job {
94+
/** Create + start a job. Returns the persisted row immediately (phase 'queued').
95+
* `items` are strings for kiwix (ZIM filenames) or objects for maps/books. */
96+
create(type: JobType, items: unknown[]): Job {
9397
const now = Date.now();
9498
const id = `${type}-${now}-${Math.random().toString(36).slice(2, 8)}`;
9599
const job: Job = {
96-
id, type, target: JSON.stringify(ids), phase: 'queued',
100+
id, type, target: JSON.stringify(items), phase: 'queued',
97101
percent: -1, speed: 0, detail: null, error: null, created: now, updated: now,
98102
};
99103
this.db.prepare(
@@ -159,11 +163,12 @@ class JobManager {
159163
const rt = { canceled: false, procs: new Set<ChildProcess>() };
160164
this.runtime.set(job.id, rt);
161165

162-
let ids: string[] = [];
163-
try { ids = JSON.parse(job.target) as string[]; } catch { ids = []; }
166+
let items: unknown[] = [];
167+
try { const parsed = JSON.parse(job.target); items = Array.isArray(parsed) ? parsed : []; } catch { items = []; }
168+
const ids = items.filter((x): x is string => typeof x === 'string');
164169

165170
const ctx: RunnerContext = {
166-
job, ids,
171+
job, items, ids,
167172
update: (p) => this.patch(job.id, p),
168173
isCanceled: () => rt.canceled,
169174
throwIfCanceled: () => { if (rt.canceled) throw new CanceledError(); },
Lines changed: 55 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,55 @@
1+
// sockets/maps.exec.ts — ADFA-4838
2+
//
3+
// Maps runner for the durable job engine: extract a tile region via tile-extract.py.
4+
// Ported from the Phase 1 maps.socket handler; reuses its box validation. A job item
5+
// is { name, box, noninteractive? }; deletion stays a separate op (not a long job).
6+
import { jobs, RunnerContext, CanceledError } from './jobs';
7+
import { parseBox } from './maps.socket';
8+
import path from 'path';
9+
10+
const SCRIPTS_DIR = '/opt/iiab/maps/tile-extract/';
11+
const EXTRACT_SCRIPT = path.join(SCRIPTS_DIR, 'tile-extract.py');
12+
// name: letters/digits/hyphen/underscore, 1..34 (same rule as the socket handler).
13+
const NAME_RE = /^[A-Za-z0-9_-]{1,34}$/;
14+
15+
interface MapItem { name?: string; box?: string; noninteractive?: boolean; }
16+
17+
const mapsRunner: (ctx: RunnerContext) => Promise<void> = async (ctx) => {
18+
const item = (ctx.items[0] ?? {}) as MapItem;
19+
const name = String(item.name ?? '');
20+
const rawBox = String(item.box ?? '');
21+
if (!NAME_RE.test(name)) throw new Error('invalid region name (A-Z a-z 0-9 _ -, length 1-34)');
22+
const parsed = parseBox(rawBox);
23+
if (!parsed.ok) throw new Error(parsed.error);
24+
25+
ctx.update({ phase: 'processing', percent: -1, detail: name });
26+
const args = [EXTRACT_SCRIPT, 'extract', name, parsed.box, 'noninteractive'];
27+
28+
await new Promise<void>((resolve, reject) => {
29+
const p = ctx.spawn('sudo', args, { env: { ...process.env, PYTHONUNBUFFERED: '1' } });
30+
const onData = (buf: Buffer) => {
31+
const text = buf.toString();
32+
ctx.log(text.trim());
33+
// tile-extract prints progress lines; surface a % if present, else stay indeterminate.
34+
const re = /(\d+)\s*%/g;
35+
let m: RegExpExecArray | null;
36+
let last = -1;
37+
while ((m = re.exec(text)) !== null) last = parseInt(m[1], 10);
38+
if (last >= 0) ctx.update({ phase: 'processing', percent: last });
39+
};
40+
p.stdout?.on('data', onData);
41+
p.stderr?.on('data', onData);
42+
p.on('error', reject);
43+
p.on('exit', (code, signal) => {
44+
if (signal === 'SIGKILL' || ctx.isCanceled()) return reject(new CanceledError());
45+
if (code === 0) resolve();
46+
else reject(new Error(`tile-extract exited with code ${code}`));
47+
});
48+
});
49+
50+
ctx.update({ phase: 'done', percent: 100 });
51+
};
52+
53+
jobs.registerRunner('maps', mapsRunner);
54+
55+
export { mapsRunner };

0 commit comments

Comments
 (0)