feat(server): add SQLite KV store exposed as $kv in script sandbox
Co-authored-by: Nasyarobby Putra <nasyarobby@gmail.com>
This commit is contained in:
@@ -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<unknown>}
|
||||
*/
|
||||
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<boolean>}
|
||||
*/
|
||||
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 });
|
||||
},
|
||||
};
|
||||
}
|
||||
@@ -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");
|
||||
}
|
||||
@@ -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),
|
||||
};
|
||||
|
||||
|
||||
@@ -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");
|
||||
Reference in New Issue
Block a user