From b88dd4d6e716d7ed6b5e8f1b3d088fda22f76aa4 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Fri, 14 Aug 2026 11:11:17 +0000 Subject: [PATCH] feat(server): add SQLite KV store exposed as $kv in script sandbox Co-authored-by: Nasyarobby Putra --- packages/server/kv-store.js | 325 ++++++++++++++++++ .../migrations/20260814130000_script_state.js | 24 ++ packages/server/script-sandbox.js | 3 + packages/server/test/kv-smoke.js | 52 +++ 4 files changed, 404 insertions(+) create mode 100644 packages/server/kv-store.js create mode 100644 packages/server/migrations/20260814130000_script_state.js create mode 100644 packages/server/test/kv-smoke.js diff --git a/packages/server/kv-store.js b/packages/server/kv-store.js new file mode 100644 index 0000000..57cd5cf --- /dev/null +++ b/packages/server/kv-store.js @@ -0,0 +1,325 @@ +import { db } from "./db.js"; + +const MAX_KEY_LENGTH = 512; +const MAX_NAMESPACE_LENGTH = 512; +const MAX_VALUE_BYTES = 256 * 1024; +const DEFAULT_LIST_LIMIT = 100; +const MAX_LIST_LIMIT = 500; + +function nowIso() { + return new Date().toISOString(); +} + +/** + * @param {unknown} value + */ +function valuesEqual(a, b) { + if (a === b) return true; + if (a == null || b == null) return a === b; + try { + return JSON.stringify(a) === JSON.stringify(b); + } catch { + return false; + } +} + +/** + * @param {string} label + * @param {unknown} value + */ +function assertString(label, value) { + if (typeof value !== "string" || value.length === 0) { + throw new Error(`${label} must be a non-empty string`); + } +} + +/** + * @param {string} label + * @param {string} value + * @param {number} max + */ +function assertMaxLength(label, value, max) { + if (value.length > max) { + throw new Error(`${label} must be at most ${max} characters`); + } +} + +/** + * @param {string} namespace + */ +function assertNamespace(namespace) { + assertString("namespace", namespace); + assertMaxLength("namespace", namespace, MAX_NAMESPACE_LENGTH); +} + +/** + * @param {string} key + */ +function assertKey(key) { + assertString("key", key); + assertMaxLength("key", key, MAX_KEY_LENGTH); +} + +/** + * @param {unknown} value + * @returns {string} + */ +export function serializeKvValue(value) { + let json; + try { + json = JSON.stringify(value); + } catch { + throw new Error("value must be JSON-serializable"); + } + if (Buffer.byteLength(json, "utf8") > MAX_VALUE_BYTES) { + throw new Error(`value exceeds ${MAX_VALUE_BYTES} byte limit`); + } + return json; +} + +/** + * @param {string | null | undefined} value + * @returns {unknown} + */ +function deserializeKvValue(value) { + if (value == null) return null; + try { + return JSON.parse(value); + } catch { + return value; + } +} + +/** + * @param {{ expires_at?: string | null }} row + */ +function isExpired(row) { + if (!row.expires_at) return false; + return Date.parse(row.expires_at) <= Date.now(); +} + +/** + * @param {string} namespace + * @param {string} key + * @returns {Promise} + */ +export async function kvGet(namespace, key) { + assertNamespace(namespace); + assertKey(key); + + const row = await db("script_state").where({ namespace, key }).first(); + if (!row) return null; + + if (isExpired(row)) { + await db("script_state").where({ namespace, key }).del(); + return null; + } + + return deserializeKvValue(row.value); +} + +/** + * @param {string} namespace + * @param {string} key + * @param {unknown} value + * @param {{ expiresAt?: string | Date | null }} [opts] + */ +export async function kvSet(namespace, key, value, opts = {}) { + assertNamespace(namespace); + assertKey(key); + + const json = serializeKvValue(value); + const updated_at = nowIso(); + let expires_at = null; + if (opts.expiresAt != null) { + expires_at = + opts.expiresAt instanceof Date + ? opts.expiresAt.toISOString() + : String(opts.expiresAt); + } + + await db("script_state") + .insert({ + namespace, + key, + value: json, + updated_at, + expires_at, + }) + .onConflict(["namespace", "key"]) + .merge({ + value: json, + updated_at, + expires_at, + }); +} + +/** + * @param {string} namespace + * @param {string} key + * @returns {Promise} + */ +export async function kvDelete(namespace, key) { + assertNamespace(namespace); + assertKey(key); + const deleted = await db("script_state").where({ namespace, key }).del(); + return deleted > 0; +} + +/** + * @param {string} namespace + * @param {string} key + * @param {unknown} expected + * @param {unknown} next + * @param {{ expiresAt?: string | Date | null }} [opts] + * @returns {Promise<{ ok: boolean, previous: unknown }>} + */ +export async function kvCompareAndSet(namespace, key, expected, next, opts = {}) { + assertNamespace(namespace); + assertKey(key); + + return db.transaction(async (trx) => { + const row = await trx("script_state").where({ namespace, key }).first(); + + if (row && isExpired(row)) { + await trx("script_state").where({ namespace, key }).del(); + } + + const currentRow = + row && !isExpired(row) + ? row + : await trx("script_state").where({ namespace, key }).first(); + const previous = currentRow ? deserializeKvValue(currentRow.value) : null; + + if (!valuesEqual(previous, expected)) { + return { ok: false, previous }; + } + + const json = serializeKvValue(next); + const updated_at = nowIso(); + let expires_at = null; + if (opts.expiresAt != null) { + expires_at = + opts.expiresAt instanceof Date + ? opts.expiresAt.toISOString() + : String(opts.expiresAt); + } + + if (currentRow) { + await trx("script_state").where({ namespace, key }).update({ + value: json, + updated_at, + expires_at, + }); + } else { + await trx("script_state").insert({ + namespace, + key, + value: json, + updated_at, + expires_at, + }); + } + + return { ok: true, previous }; + }); +} + +/** + * @param {string} namespace + * @param {{ limit?: number }} [opts] + */ +export async function kvList(namespace, opts = {}) { + assertNamespace(namespace); + const limit = Math.min( + Math.max(opts.limit ?? DEFAULT_LIST_LIMIT, 1), + MAX_LIST_LIMIT, + ); + + const rows = await db("script_state") + .where({ namespace }) + .orderBy("updated_at", "desc") + .limit(limit); + + const items = []; + for (const row of rows) { + if (isExpired(row)) { + await db("script_state").where({ namespace, key: row.key }).del(); + continue; + } + items.push({ + key: row.key, + value: deserializeKvValue(row.value), + updatedAt: row.updated_at, + expiresAt: row.expires_at ?? null, + }); + } + return items; +} + +/** + * @param {string} defaultNamespace + */ +export function createKvApi(defaultNamespace) { + assertNamespace(defaultNamespace); + + /** + * @param {{ namespace?: string }} [opts] + */ + function resolveNamespace(opts = {}) { + const namespace = opts.namespace ?? defaultNamespace; + assertNamespace(namespace); + return namespace; + } + + return { + namespace: defaultNamespace, + + /** + * @param {string} key + * @param {{ namespace?: string }} [opts] + */ + get(key, opts) { + return kvGet(resolveNamespace(opts), key); + }, + + /** + * @param {string} key + * @param {unknown} value + * @param {{ namespace?: string, expiresAt?: string | Date | null }} [opts] + */ + set(key, value, opts) { + const { namespace, expiresAt } = opts ?? {}; + return kvSet(resolveNamespace(opts), key, value, { expiresAt }); + }, + + /** + * @param {string} key + * @param {{ namespace?: string }} [opts] + */ + delete(key, opts) { + return kvDelete(resolveNamespace(opts), key); + }, + + /** + * @param {string} key + * @param {unknown} expected + * @param {unknown} next + * @param {{ namespace?: string, expiresAt?: string | Date | null }} [opts] + */ + compareAndSet(key, expected, next, opts) { + const { expiresAt } = opts ?? {}; + return kvCompareAndSet(resolveNamespace(opts), key, expected, next, { + expiresAt, + }); + }, + + /** + * @param {{ namespace?: string, limit?: number }} [opts] + */ + list(opts) { + const { namespace, limit } = opts ?? {}; + return kvList(resolveNamespace(opts), { limit }); + }, + }; +} diff --git a/packages/server/migrations/20260814130000_script_state.js b/packages/server/migrations/20260814130000_script_state.js new file mode 100644 index 0000000..065eef0 --- /dev/null +++ b/packages/server/migrations/20260814130000_script_state.js @@ -0,0 +1,24 @@ +/** + * @param {import("knex").Knex} knex + */ +export async function up(knex) { + await knex.schema.createTable("script_state", (t) => { + t.text("namespace").notNullable(); + t.text("key").notNullable(); + t.text("value").notNullable(); + t.text("updated_at").notNullable(); + t.text("expires_at"); + t.primary(["namespace", "key"]); + }); + + await knex.schema.raw( + "CREATE INDEX script_state_namespace_updated_at_idx ON script_state (namespace, updated_at DESC)", + ); +} + +/** + * @param {import("knex").Knex} knex + */ +export async function down(knex) { + await knex.schema.dropTableIfExists("script_state"); +} diff --git a/packages/server/script-sandbox.js b/packages/server/script-sandbox.js index f01ffa7..0af779d 100644 --- a/packages/server/script-sandbox.js +++ b/packages/server/script-sandbox.js @@ -3,6 +3,7 @@ import path from "path"; import vm from "node:vm"; import { createRequire } from "node:module"; import axios from "axios"; +import { createKvApi } from "./kv-store.js"; import { SCRIPTS_DIR } from "./paths.js"; const hostRequire = createRequire(import.meta.url); @@ -275,11 +276,13 @@ function createRestrictedRequire(screenedAxios) { function createScriptSandbox({ log, script, workflowName }) { const scriptLog = log.child({ workflow: workflowName, script }); const $axios = createScreenedAxios(scriptLog); + const $kv = createKvApi(workflowName); const sandbox = { ...pickBuiltins(), log: scriptLog, console: createConsole(scriptLog), $axios, + $kv, require: createRestrictedRequire($axios), }; diff --git a/packages/server/test/kv-smoke.js b/packages/server/test/kv-smoke.js new file mode 100644 index 0000000..931b157 --- /dev/null +++ b/packages/server/test/kv-smoke.js @@ -0,0 +1,52 @@ +import { migrate } from "../db.js"; +import { createKvApi, kvDelete } from "../kv-store.js"; +import { runScriptSource } from "../script-sandbox.js"; +import { log } from "../logger.js"; + +await migrate(); + +const ns = "test/kv-smoke"; +await kvDelete(ns, "counter"); +await kvDelete(ns, "cas"); + +const api = createKvApi(ns); + +await api.set("counter", { n: 1 }); +const v1 = await api.get("counter"); +if (v1?.n !== 1) throw new Error(`expected n=1, got ${JSON.stringify(v1)}`); + +const casFail = await api.compareAndSet("cas", null, { first: true }); +if (!casFail.ok || casFail.previous !== null) { + throw new Error(`first CAS failed: ${JSON.stringify(casFail)}`); +} + +const casFail2 = await api.compareAndSet("cas", null, { second: true }); +if (casFail2.ok) throw new Error("second CAS with null expected should fail"); + +const casOk = await api.compareAndSet("cas", { first: true }, { second: true }); +if (!casOk.ok) throw new Error("CAS update should succeed"); + +const listed = await api.list(); +if (!listed.some((item) => item.key === "cas")) { + throw new Error(`list missing cas: ${JSON.stringify(listed)}`); +} + +const script = ` +export default async function () { + const prev = await $kv.get("script-key"); + await $kv.set("script-key", { prev, now: Date.now() }); + return { prev, namespace: $kv.namespace }; +} +`; + +const out = await runScriptSource("kv-smoke.js", script, { data: {} }, { + log, + workflowName: ns, +}); +if (out.namespace !== ns) throw new Error(`wrong namespace: ${out.namespace}`); + +await kvDelete(ns, "counter"); +await kvDelete(ns, "cas"); +await kvDelete(ns, "script-key"); + +console.log("kv smoke test passed");