{"_id":"@aleju03/kewa","name":"@aleju03/kewa","dist-tags":{"latest":"0.1.0"},"versions":{"0.1.0":{"name":"@aleju03/kewa","version":"0.1.0","description":"A tiny durable job queue on SQLite: enqueue/claim/dedupe/backoff, worker lanes, and load shedding.","type":"module","license":"MIT","author":{"name":"aleju03"},"repository":{"type":"git","url":"git+https://github.com/aleju03/kewa.git"},"bugs":{"url":"https://github.com/aleju03/kewa/issues"},"homepage":"https://github.com/aleju03/kewa#readme","keywords":["job-queue","queue","sqlite","libsql","worker","background-jobs","durable","backoff","dedupe"],"main":"./dist/index.js","types":"./dist/index.d.ts","exports":{".":{"types":"./dist/index.d.ts","import":"./dist/index.js"}},"scripts":{"build":"tsc -p tsconfig.json","test":"vitest run","typecheck":"tsc -p tsconfig.json --noEmit","verify":"vitest run && tsc -p tsconfig.json --noEmit","example":"node --import tsx examples/demo.ts","prepublishOnly":"npm run build"},"dependencies":{"@libsql/client":"^0.17.3"},"devDependencies":{"@types/node":"^22.10.2","tsx":"^4.20.6","typescript":"^5.7.2","vitest":"^3.0.5"},"gitHead":"d01fd4eee0f2b70ba5555c8cd8d2017f69d7f915","_id":"@aleju03/kewa@0.1.0","_nodeVersion":"24.4.0","_npmVersion":"11.10.0","dist":{"integrity":"sha512-OIbbb4jqlZdv0Z+9Z/wKsx/5GD3gT8NaENJ7poW+K/oVxhyxtYq0BedHXLaHTyX0ZVpIy/hM8ke/7evGv/Hzxg==","shasum":"8e4c762c254d483e088ffbf2a273bb8634fe408f","tarball":"https://registry.npmjs.org/@aleju03/kewa/-/kewa-0.1.0.tgz","fileCount":11,"unpackedSize":50438,"signatures":[{"keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U","sig":"MEUCIEuXkt1x4uMPUHDPRwYIizniWzBH4oRYen+FokQVb37FAiEA35pL5dEvjh2RYHH5lsd79r00o63i92BQ0OH7TiO3XQ8="}]},"_npmUser":{"name":"aleju03","email":"alejimenezu@gmail.com"},"directories":{},"maintainers":[{"name":"aleju03","email":"alejimenezu@gmail.com"}],"_npmOperationalInternal":{"host":"s3://npm-registry-packages-npm-production","tmp":"tmp/kewa_0.1.0_1784838391373_0.20256146231785088"},"_hasShrinkwrap":false}},"time":{"created":"2026-07-23T20:26:31.110Z","0.1.0":"2026-07-23T20:26:31.541Z","modified":"2026-07-23T20:26:31.822Z"},"maintainers":[{"name":"aleju03","email":"alejimenezu@gmail.com"}],"description":"A tiny durable job queue on SQLite: enqueue/claim/dedupe/backoff, worker lanes, and load shedding.","homepage":"https://github.com/aleju03/kewa#readme","keywords":["job-queue","queue","sqlite","libsql","worker","background-jobs","durable","backoff","dedupe"],"repository":{"type":"git","url":"git+https://github.com/aleju03/kewa.git"},"author":{"name":"aleju03"},"bugs":{"url":"https://github.com/aleju03/kewa/issues"},"license":"MIT","readme":"<p align=\"center\">\n  <img src=\"assets/mascot.png\" alt=\"kewa mascot\" height=\"160\">\n</p>\n\n# kewa\n\nA tiny **durable job queue on SQLite**. Enqueue with dedupe keys, claim with locks, retry with backoff, and shed load when the backlog spikes, all in one table, no broker to run.\n\n`kewa` (queue) is the job-queue core pulled out of an always-on service that runs thousands of background jobs a day on a single SQLite file. It is deliberately unopinionated: you bring the job types and their handlers, and kewa owns the mechanics: durable state, dedupe-merge semantics, per-type retry backoff, worker lanes with a watchdog, and pressure-based load shedding.\n\n## Why\n\n- **One SQLite file, no broker.** Jobs live in a `jobs` table in your own database (a `file:` path, `:memory:`, or a remote libsql/Turso URL). No Redis, no RabbitMQ, no separate service to keep alive. Every claim is a local indexed update.\n- **Durable and idempotent.** A claimed job is locked with a lease; if a worker dies mid-job the lease expires and the job is reclaimed. Handlers are meant to be idempotent, so a redelivery is safe.\n- **Dedupe and debounce built in.** Every job has a unique `dedupe_key`. Re-enqueuing merges instead of duplicating: priority takes the max, and `run_after` normally moves earlier (an urgent request pulls a scheduled job forward) or, with `debounce`, later (a burst of events collapses into one delayed run).\n- **Backpressure that does not starve.** When the runnable backlog grows past a target depth, low-value jobs park as `deferred_pressure` and reactivate once there is headroom. Reserved lanes keep a small always-runnable reserve so no type is starved to zero under sustained load.\n\n## Install\n\n```bash\nnpm install @aleju03/kewa\n```\n\nRequires Node 18+. Storage is `@libsql/client`, so a `file:` path, `:memory:`, or a remote Turso URL all work.\n\n## Quick start\n\n```ts\nimport { JobQueue, WorkerRunner, createDb } from \"@aleju03/kewa\";\n\nconst db = await createDb({ url: \"file:jobs.db\" });\nconst queue = new JobQueue(db);\nawait queue.ensureSchema();\n\nconst runner = new WorkerRunner(queue, {\n  lanes: [\n    { name: \"mail\",  jobTypes: [\"send_email\"],   claimLimit: 4, intervalMs: 250 },\n    { name: \"media\", jobTypes: [\"resize_image\"], claimLimit: 2, intervalMs: 500 },\n  ],\n});\n\n// You bring the job types. kewa dispatches by job.type to what you register.\nrunner.registerHandler<{ to: string }>(\"send_email\", async (job, { signal }) => {\n  await sendEmail(job.payload.to, { signal });\n});\nrunner.registerHandler<{ id: number }>(\"resize_image\", async (job) => {\n  await resize(job.payload.id);\n});\n\nconst stop = runner.start();          // background lane timers begin claiming\n\n// Enqueue from anywhere that has the queue:\nawait queue.enqueue(\"send_email\", `welcome:${userId}`, { to: user.email }, { priority: 10 });\n\n// later: stop();  and  db.close();\n```\n\n`enqueue(type, dedupeKey, payload, options?)` inserts or merges a job. `registerHandler(type, fn)` binds a handler; a job whose `type` has no handler fails with a clear error. Handlers receive the `Job` and a `{ signal }` context (the watchdog aborts it on timeout).\n\n## The injection story: you bring the job types\n\nkewa ships zero domain knowledge. There is no built-in list of job names, no hardcoded handlers. Two things stay yours:\n\n1. **The job types and their payloads.** You pick the `type` strings and `registerHandler` a function for each. The payload is any JSON-serializable value, typed through the handler generic.\n2. **The tuning.** Depth thresholds, which types are sheddable, per-type caps, and reserved lanes are all constructor options with generic defaults. Nothing is assumed about your workload.\n\n```ts\nconst queue = new JobQueue(db, {\n  targetDepth: 200,                       // trim the runnable pool back to this under pressure\n  softPressureDepth: 160,                 // at/above this, sheddable types defer on enqueue\n  sheddableTypes: [\"rebuild_search_index\"],\n  typeCaps: { call_external_api: 20 },    // keep at most N of this type runnable at target depth\n  reservedLanes: { nightly_report: 1 },   // a small always-runnable reserve, invisible to shared depth\n});\n```\n\n## Lanes\n\nA `WorkerRunner` runs one or more **lanes**. Each lane polls on its own `intervalMs`, claims up to `claimLimit` jobs (optionally filtered to a set of `jobTypes`), and runs them concurrently. Separate lanes keep slow work from starving fast work: a wide interactive lane and a narrow `claimLimit: 1` backfill lane share the same queue without fighting.\n\n```ts\nrunner.pause();          // stop starting new jobs (in-flight jobs finish)\nrunner.resume();\nrunner.status();         // { paused, stopped, workerId, lanes: [{ name, activeJobs, ... }] }\nawait runner.runOnce();  // claim + run one batch synchronously (handy in tests)\n```\n\nEach job invocation is guarded by a **watchdog** (`DEFAULT_JOB_WATCHDOG_MS`, 10 minutes, or per-lane `jobTimeoutMs`). If a handler hangs past the ceiling, its lane is released so it keeps ticking, the job takes the normal fail-with-backoff path, and the handler's `AbortSignal` fires so cooperative handlers can stop at their next boundary.\n\n## Backoff\n\nA handler that throws marks its job `failed` and stamps `run_after` in the future, so the job is retried after a delay, not immediately. The default is exponential: base 30s, doubling per attempt, capped at 60 minutes. Override it globally, per handler, or per error:\n\n```ts\nconst runner = new WorkerRunner(queue, {\n  retryDelayMs: 5_000,                                   // flat override for every type\n  retryBackoff: { baseMs: 1_000, factor: 3, maxMs: 60_000 }, // or tune the exponential curve\n});\n\nrunner.registerHandler(\"call_external_api\", handler, { retryDelayMs: 60_000 }); // per handler\n\nconst runner2 = new WorkerRunner(queue, {\n  // Classify an error into fail (default), defer (retry, no failure recorded),\n  // or complete (give up and mark done), with an optional per-occurrence delay.\n  errorClassifier: (error) => {\n    if (error instanceof RateLimited) return { retryDelayMs: error.retryAfterMs };\n    if (error instanceof NotFound)    return { action: \"complete\", reason: \"gone\" };\n    return { action: \"fail\" };\n  },\n});\n```\n\nObserve the lifecycle with `onJobEvent`, which fires on `started` / `done` / `failed` / `deferred`:\n\n```ts\nnew WorkerRunner(queue, {\n  onJobEvent: (e) => log.info(e.status, { job: e.jobId, type: e.type, lane: e.lane, attempts: e.attempts }),\n});\n```\n\n## Pressure and shedding\n\n`await queue.pressure()` returns a snapshot: current runnable `depth`, `deferred` count, the configured thresholds, and whether the queue is `shedding`. Jobs scheduled for the future (retry backoffs, self-rescheduling work) count as appointments, not backlog, so they never trip shedding on their own.\n\nShedding runs on every `enqueue` and reserved-lane refills run on every `claim`; a long-lived deployment should also call `await queue.shedPressure()` on a timer (say every 30-60s) so parked work reactivates even when nothing is being enqueued. Reserved-lane types are held at their reserve independently of the shared pool, guaranteeing a trickle of progress for every type under any pressure.\n\nOther queue methods: `complete(id)`, `fail(id, error, retryDelayMs?)`, `defer(id, retryDelayMs?)`, `depth()`, `summary()` (counts by status/type with the newest error), `clearFailed(type?)`, and `hasRunnableOutsideTypes(types)`.\n\n## Run the demo\n\n```bash\nnpm install\nnpm run example   # builds a queue, registers handlers, runs two lanes, sheds load, retries a flaky job\nnpm test          # queue merge/backoff/shedding + runner end-to-end tests\n```\n\n## License\n\nMIT\n","readmeFilename":"README.md","_rev":"1-f3664b8567593389045c7bccf4654782"}