diff --git a/package.json b/package.json index 0988f78..20e4b1e 100644 --- a/package.json +++ b/package.json @@ -1,6 +1,6 @@ { "name": "@ztimson/zim-utils", - "version": "0.2.0", + "version": "0.2.1", "description": "Native, dependency-light ZIM archive reader/searcher and Kiwix catalog downloader for Node.js", "author": "Zak Timson", "license": "MIT", diff --git a/src/catalog.js b/src/catalog.js index e57d496..bdd7508 100644 --- a/src/catalog.js +++ b/src/catalog.js @@ -1,11 +1,11 @@ import {fuzzyMatch} from './utils.js'; import {decodeHtml} from '@ztimson/utils'; -export const CATALOG_URL = 'https://library.kiwix.org/catalog/v2/entries'; +export const CATALOG_URL = 'https://library.kiwix.org'; const PAGE_SIZE = 100; -function parseEntries(xml, baseUrl = 'https://library.kiwix.org') { +function parseEntries(xml, catalog = CATALOG_URL) { const blocks = xml.match(/[\s\S]*?<\/entry>/g) || []; return blocks.map(b => { const grab = re => (b.match(re) || [])[1] || ''; @@ -14,29 +14,31 @@ function parseEntries(xml, baseUrl = 'https://library.kiwix.org') { || b.match(/]*type=["']image\/[^"']*["'][^>]*href=["']([^"']+)["']/); let tags = grab(/([^<]*)<\/tags>/); if(tags) tags = tags.split(';'); + const name = grab(/([^<]*)<\/name>/); return { id: grab(/([^<]*)<\/id>/), title: decodeHtml(grab(/([^<]*)<\/title>/)), updated: new Date(grab(/<updated>([^<]*)<\/updated>/)), summary: decodeHtml(grab(/<summary>([^<]*)<\/summary>/)), language: grab(/<language>([^<]*)<\/language>/), - name: grab(/<name>([^<]*)<\/name>/), + name, category: grab(/<category>([^<]*)<\/category>/), tags, mediaCount: Number(grab(/<mediaCount>([^<]*)<\/mediaCount>/)) || 0, author: grab(/<author>\s*<name>([^<]*)<\/name>\s*<\/author>/m), publisher: grab(/<publisher>\s*<name>([^<]*)<\/name>\s*<\/publisher>/m), articleCount: Number(grab(/<articleCount>([^<]*)<\/articleCount>/)) || 0, - sizeMb: linkMatch ? (Number((b.match(/length=["'](\d+)["']/) || [])[1] || 0) / 1024 / 1024).toFixed(1) : '?', - href: linkMatch ? linkMatch[1] : null, - icon: iconMatch ? new URL(iconMatch[1], baseUrl).href : null, + sizeMb: linkMatch ? +(Number((b.match(/length=["'](\d+)["']/) || [])[1] || 0) / 1024 / 1024).toFixed(1) : '?', + download: linkMatch ? linkMatch[1] : null, + icon: iconMatch ? new URL(iconMatch[1], catalog).href : null, + viewer: name ? new URL(`viewer#${name}`, catalog).href : null, }; }); } async function fetchEntries(term, lang, url = CATALOG_URL) { const params = new URLSearchParams({q: term, count: String(PAGE_SIZE), lang: lang || 'eng'}); - const res = await fetch(`${url}?${params}`); + const res = await fetch(`${url}/catalog/v2/entries?${params}`); if (!res.ok) throw new Error(`${res.status} ${res.statusText}`); return parseEntries(await res.text()); } @@ -44,7 +46,7 @@ async function fetchEntries(term, lang, url = CATALOG_URL) { /** Looks up a single catalog entry by exact `name` (used to check for available updates). */ export async function zimCatalogInfo(name, url = CATALOG_URL) { const params = new URLSearchParams({name, count: '5'}); - const res = await fetch(`${url}?${params}`); + const res = await fetch(`${url}/catalog/v2/entries?${params}`); if (!res.ok) return null; return parseEntries(await res.text()).find(e => e.name === name) || null; } diff --git a/src/decompress-worker.js b/src/decompress-worker.js new file mode 100644 index 0000000..bb91888 --- /dev/null +++ b/src/decompress-worker.js @@ -0,0 +1,32 @@ +'use strict'; + +import {parentPort} from 'node:worker_threads'; +import {decompress as lzmaDecompress} from 'lzma1'; +import {ZstdCodec} from 'zstd-codec'; + +let zstdStreamingPromise = null; + +function getZstd() { + if (!zstdStreamingPromise) { + zstdStreamingPromise = new Promise(resolve => ZstdCodec.run(zstd => resolve(new zstd.Streaming()))); + } + return zstdStreamingPromise; +} + +/** Rebuilds a full 13-byte "alone" LZMA header from libzim's truncated 5-byte one (unknown size). */ +function toLzmaAloneStream(body) { + const header = Buffer.concat([body.subarray(0, 5), Buffer.alloc(8, 0xff)]); + return Buffer.concat([header, body.subarray(5)]); +} + +parentPort.on('message', async ({id, compType, body}) => { + try { + const buf = Buffer.from(body); + const data = compType === 4 + ? Buffer.from(lzmaDecompress(toLzmaAloneStream(buf))) + : Buffer.from((await getZstd()).decompress(new Uint8Array(buf))); + parentPort.postMessage({id, data}, [data.buffer]); + } catch (e) { + parentPort.postMessage({id, error: e.message}); + } +}); diff --git a/src/decompress.js b/src/decompress.js new file mode 100644 index 0000000..9f577a8 --- /dev/null +++ b/src/decompress.js @@ -0,0 +1,50 @@ +'use strict'; + +import {Worker} from 'node:worker_threads'; +import path from 'node:path'; +import os from 'node:os'; +import {fileURLToPath} from 'node:url'; + +const __dirname = path.dirname(fileURLToPath(import.meta.url)); + +/** Round-robin pool of worker threads for off-main-thread cluster decompression. */ +class DecompressPool { + #workers = []; + #next = 0; + #pending = new Map(); // id -> {resolve, reject} + #nextId = 0; + #active = 0; + + constructor(size = Math.max(1, os.cpus().length - 1)) { + for (let i = 0; i < size; i++) this.#spawn(); + } + + #spawn() { + const worker = new Worker(path.join(__dirname, 'decompress-worker.js')); + worker.on('message', ({id, data, error}) => { + const task = this.#pending.get(id); + if (!task) return; + this.#pending.delete(id); + // No work left in flight -> safe to let the process exit again. + if (--this.#active === 0) for (const w of this.#workers) w.unref(); + error ? task.reject(new Error(error)) : task.resolve(Buffer.from(data)); + }); + worker.unref(); // idle workers must never keep a short script/server alive + this.#workers.push(worker); + } + + /** Decompresses `{compType, body}` on the next worker in rotation. `body` must be a Uint8Array/Buffer view. */ + run({compType, body}) { + const worker = this.#workers[this.#next]; + this.#next = (this.#next + 1) % this.#workers.length; + const id = this.#nextId++; + const owned = new Uint8Array(body); + if (this.#active++ === 0) for (const w of this.#workers) w.ref(); + return new Promise((resolve, reject) => { + this.#pending.set(id, {resolve, reject}); + worker.postMessage({id, compType, body: owned}, [owned.buffer]); + }); + } +} + +export const decompressPool = new DecompressPool(); diff --git a/src/manager.js b/src/manager.js index d9cf273..e06e505 100644 --- a/src/manager.js +++ b/src/manager.js @@ -4,16 +4,24 @@ import {pipeline} from 'node:stream/promises'; import {Readable} from 'node:stream'; import {ZimReader} from './reader.js'; import {zimCatalog, zimCatalogInfo, CATALOG_URL} from './catalog.js'; -import {fuzzyMatch, titleFromUrl} from './utils.js'; + +const DEFAULT_READER_TTL = 60_000; /** Manages a local directory of ZIM archives: listing, update checks, downloads, and reading. */ export class ZimManager { #catalog; #dir; + #readerTTL; + #clusterTTL; + #readers = new Map(); // file -> {reader, timer} + #opening = new Map(); // file -> Promise<ZimReader>, dedupes concurrent first-open races - constructor(dir, catalog = CATALOG_URL) { + /** @param {{catalog?: string, readerTTL?: number, clusterTTL?: number}} [opts] */ + constructor(dir, catalog = CATALOG_URL, {readerTTL = DEFAULT_READER_TTL, clusterTTL} = {}) { this.#catalog = catalog; this.#dir = dir; + this.#readerTTL = readerTTL; + this.#clusterTTL = clusterTTL; // undefined -> ZimReader's own default } async #download(url, destPath) { @@ -58,13 +66,13 @@ export class ZimManager { const localDate = localMatch?.meta?.updated ?? null; if (!force && localMatch && remoteDate && localDate && remoteDate <= localDate) return {name, status: 'skipped', reason: 'up to date'}; - if (!catalogEntry.href) return {name, status: 'skipped', reason: 'missing download link'}; + if (!catalogEntry.download) return {name, status: 'skipped', reason: 'missing download link'}; - const filename = path.basename(new URL(catalogEntry.href).pathname).replace(/\.meta4$/i, ''); + const filename = path.basename(new URL(catalogEntry.download).pathname).replace(/\.meta4$/i, ''); const destPath = path.join(this.#dir, filename); - await this.#download(catalogEntry.href, destPath); - if (localMatch && localMatch.file !== destPath) await new ZimReader(localMatch.file).delete().catch(() => {}); - return {name, status: 'updated', file: destPath}; + await this.#download(catalogEntry.download, destPath); + if (localMatch && localMatch.file !== filename) await this.#evict(path.join(this.#dir, localMatch.file)); + return {name, status: 'updated', file: filename}; } /** Resolves a file path or catalog `name` to a local file path. */ @@ -75,7 +83,16 @@ export class ZimManager { const local = await this.list(); const match = local.find(l => l.meta?.name === fileOrName || path.basename(l.file) === fileOrName); if (!match) throw new Error(`ZIM not found locally: ${fileOrName}`); - return match.file; + return path.join(this.#dir, match.file); + } + + /** Closes and drops a cached reader for `file`, if any (used before delete/replace). */ + async #evict(file) { + const entry = this.#readers.get(file); + if (!entry) return new ZimReader(file).delete().catch(() => {}); + clearTimeout(entry.timer); + this.#readers.delete(file); + await entry.reader.delete(); } catalog(search, opts) { @@ -84,8 +101,8 @@ export class ZimManager { async delete(fileOrName) { const file = await this.#resolveFile(fileOrName); - await new ZimReader(file).delete(); - return {file, status: 'deleted'}; + await this.#evict(file); + return {file: path.basename(file), status: 'deleted'}; } /** Checks whether a local ZIM has a newer version in the catalog, without downloading. */ @@ -105,7 +122,7 @@ export class ZimManager { const files = (await fs.promises.readdir(this.#dir)).filter(f => f.endsWith('.zim')); return Promise.all(files.map(async f => { const file = path.join(this.#dir, f); - return {file, meta: await this.#readMeta(file)}; + return {file: f, ...(await this.#readMeta(file))}; })); } @@ -118,7 +135,7 @@ export class ZimManager { const local = await this.list(); const localMatch = local.find(l => l.meta?.name === name) ?? null; - const catalogEntry = await zimCatalogInfo(name, this.#catalog) || {name, updated: null, href}; + const catalogEntry = await zimCatalogInfo(name, this.#catalog) || {name, updated: null, download: href}; return this.#update(name, catalogEntry, localMatch, force); } @@ -141,11 +158,38 @@ export class ZimManager { return results; } - /** Opens a `ZimReader` for a local ZIM, resolved by file path or catalog `name`. Caller must `.close()` it. */ + /** Opens a fresh `ZimReader` for a local ZIM, resolved by file path or catalog `name`. Caller must `.close()` it. */ async open(fileOrName) { await this.#ensureDir(); const file = await this.#resolveFile(fileOrName); - return new ZimReader(file).open(); + return new ZimReader(file, {clusterTTL: this.#clusterTTL}).open(); + } + + /** + * Returns a cached, persistently-open `ZimReader` for serving requests — avoids + * re-opening the file per request. Idle-evicted after `readerTTL` ms of no use. + */ + async getCached(fileOrName) { + await this.#ensureDir(); + const file = await this.#resolveFile(fileOrName); + let entry = this.#readers.get(file); + if (!entry) { + let pending = this.#opening.get(file); + if (!pending) { + pending = new ZimReader(file, {clusterTTL: this.#clusterTTL}).open(); + this.#opening.set(file, pending); + } + const reader = await pending; + this.#opening.delete(file); + entry = this.#readers.get(file) ?? {reader}; + this.#readers.set(file, entry); + } + clearTimeout(entry.timer); + entry.timer = setTimeout(() => { + this.#readers.delete(file); + entry.reader.close(); + }, this.#readerTTL).unref(); + return entry.reader; } /** Fuzzy-searches titles across every local ZIM in the library, merging & re-ranking hits by score. */ @@ -154,9 +198,9 @@ export class ZimManager { const perZim = await Promise.all(local.map(async ({file, meta}) => { let reader; try { - reader = await new ZimReader(file).open(); + reader = await new ZimReader(path.join(this.#dir, file)).open(); const hits = await reader.search(terms, {limit, htmlOnly}); - return hits.map(h => ({...h, file, name: meta?.name})); + return hits.map(h => ({...h, file})); } catch { return []; } finally { diff --git a/src/reader.js b/src/reader.js index 19960af..82ac83b 100644 --- a/src/reader.js +++ b/src/reader.js @@ -2,8 +2,7 @@ import fs from 'node:fs'; import path from 'node:path'; -import {decompress as lzmaDecompress} from 'lzma1'; -import {ZstdCodec} from 'zstd-codec'; +import {decompressPool} from './decompress.js'; import {fuzzyMatch, titleFromUrl} from './utils.js'; const INDEX_VERSION = 1; @@ -12,21 +11,8 @@ const NS_CONTENT = 'C'; const NS_METADATA = 'M'; const TITLE_SENTINEL = 0xffffffffffffffffn; // Indicator -> ZIM v6+ archives with no title -let zstdStreamingPromise = null; - -/** Lazily initialised, shared Zstd streaming decompressor (handles unknown-size frames). */ -function getZstd() { - if (!zstdStreamingPromise) { - zstdStreamingPromise = new Promise(resolve => ZstdCodec.run(zstd => resolve(new zstd.Streaming()))); - } - return zstdStreamingPromise; -} - -/** Rebuilds a full 13-byte "alone" LZMA header from libzim's truncated 5-byte one (unknown size). */ -function toLzmaAloneStream(body) { - const header = Buffer.concat([body.subarray(0, 5), Buffer.alloc(8, 0xff)]); - return Buffer.concat([header, body.subarray(5)]); -} +const DEFAULT_CLUSTER_CACHE_MAX = 32; +const DEFAULT_CLUSTER_TTL = 60_000; /** Native, dependency-light reader for .zim archives. Supports zstd & LZMA cluster compression. */ export class ZimReader { @@ -35,11 +21,19 @@ export class ZimReader { #mimeTypes = []; #hasTitleListing = false; #index; + #clusterCache = new Map(); // clusterNumber -> {data, extended, timer} + #pending = new Map(); // clusterNumber -> Promise, dedupes concurrent misses + #clusterCacheMax; + #clusterTTL; get articleCount() { return this.#header?.articleCount ?? 0; } + get mediaCount() { return this.#header?.clusterCount ?? 0; } - constructor(path) { + /** @param {{clusterCacheMax?: number, clusterTTL?: number}} [opts] clusterTTL in ms; 0/null disables idle eviction. */ + constructor(path, {clusterCacheMax = DEFAULT_CLUSTER_CACHE_MAX, clusterTTL = DEFAULT_CLUSTER_TTL} = {}) { this.path = path; + this.#clusterCacheMax = clusterCacheMax; + this.#clusterTTL = clusterTTL; } /** Full O(n) scan over the URL pointer list, used when there's no title index. */ @@ -66,7 +60,15 @@ export class ZimReader { return null; } - async #getBlob(clusterNumber, blobNumber) { + /** Resets a cluster's idle-eviction timer. No-op when TTL disabled. */ + #touch(clusterNumber, entry) { + if (!this.#clusterTTL) return; + clearTimeout(entry.timer); + entry.timer = setTimeout(() => this.#clusterCache.delete(clusterNumber), this.#clusterTTL).unref(); + } + + /** Fetches + decompresses a cluster exactly once, offloading decompression to the worker pool. */ + async #loadCluster(clusterNumber) { const start = await this.#ptr64(this.#header.clusterPtrPos, clusterNumber); const isLast = clusterNumber === this.#header.clusterCount - 1; const end = isLast @@ -79,14 +81,35 @@ export class ZimReader { const body = raw.subarray(1); let data; - if (compType <= 1) data = body; - else if (compType === 4) data = Buffer.from(lzmaDecompress(toLzmaAloneStream(body))); - else if (compType === 5) data = Buffer.from((await getZstd()).decompress(new Uint8Array(body))); + if (compType <= 1) data = Buffer.from(body); + else if (compType === 4 || compType === 5) data = await decompressPool.run({compType, body}); else throw new Error(`Unsupported cluster compression type: ${compType}`); if (!data) throw new Error(`Cluster ${clusterNumber} failed to decompress (compType ${compType})`); - const readPtr = i => extended ? Number(data.readBigUInt64LE(i * 8)) : data.readUInt32LE(i * 4); - return data.subarray(readPtr(blobNumber), readPtr(blobNumber + 1)); + return {data, extended}; + } + + async #getBlob(clusterNumber, blobNumber) { + let entry = this.#clusterCache.get(clusterNumber); + if (!entry) { + let pending = this.#pending.get(clusterNumber); + if (!pending) { + pending = this.#loadCluster(clusterNumber); + this.#pending.set(clusterNumber, pending); + } + entry = await pending; + this.#pending.delete(clusterNumber); + this.#clusterCache.set(clusterNumber, entry); + if (this.#clusterCache.size > this.#clusterCacheMax) { + const oldestKey = this.#clusterCache.keys().next().value; + clearTimeout(this.#clusterCache.get(oldestKey)?.timer); + this.#clusterCache.delete(oldestKey); + } + } + this.#touch(clusterNumber, entry); + + const readPtr = i => entry.extended ? Number(entry.data.readBigUInt64LE(i * 8)) : entry.data.readUInt32LE(i * 4); + return entry.data.subarray(readPtr(blobNumber), readPtr(blobNumber + 1)); } async #icon(size = 48) { @@ -227,6 +250,9 @@ export class ZimReader { async close() { if (this.#fd) await this.#fd.close(); this.#fd = null; + for (const entry of this.#clusterCache.values()) clearTimeout(entry.timer); + this.#clusterCache.clear(); + this.#pending.clear(); } /** Deletes the zim archive and its cached index (if any). Safe to call on unopened readers. */ @@ -249,7 +275,6 @@ export class ZimReader { ['Title', 'Creator', 'Publisher', 'Date', 'Description', 'Language', 'Name', 'Tags'].map(get) ); return { - id: name, title, updated: date ? new Date(date) : null, summary: description, @@ -257,12 +282,11 @@ export class ZimReader { name, category: tags ? tags.split(';')[0] || '' : '', tags: tags ? tags.split(';') : [], - mediaCount: 0, author: creator, publisher, articleCount: this.articleCount, - sizeMb: ((await fs.promises.stat(this.path)).size / 1024 / 1024).toFixed(1), - href: null, + mediaCount: this.mediaCount, + sizeMb: +((await fs.promises.stat(this.path)).size / 1024 / 1024).toFixed(1), icon: await this.#icon(), }; } @@ -325,7 +349,7 @@ export class ZimReader { const urlScore = fuzzyMatch(urlTitle, ...termList).max; scored.push({url: dirent.url, title: dirent.title.length > urlTitle.length ? dirent.title : urlTitle, namespace: NS_CONTENT, score: Math.max(titleScore, urlScore)}); } - const {id, summary, mediaCount, articleCount, sizeMb, href, ...meta} = await this.metadata(); + const {summary, mediaCount, articleCount, sizeMb, ...meta} = await this.metadata(); return scored.filter(a => a.score > 0).toSorted((a, b) => b.score - a.score).slice(0, limit).map(a => ({...meta, ...a})); } }