diff --git a/docs/architecture.md b/docs/architecture.md index 18efd0c45..29751b959 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -32,11 +32,11 @@ flowchart TD storage, object storage, screenshots, and mail. - `packages/core`, `packages/db`, `packages/rpc`, and the remaining service packages own portable application and domain behavior. -- `@orbian/node` implements those contracts for a single Node deployment. - `selfhost/entry` composes the Community application. +- Trusted host services stay inside `selfhost/entry`, which composes the + Community application. -Release lockfiles pin `@orbian/sdk` and `@orbian/node` to immutable Orbian -commit artifacts. Contributors working in adjacent checkouts can switch to the +Release lockfiles pin `@orbian/sdk` to an immutable Orbian commit artifact. +Contributors working in adjacent checkouts can switch to the sibling workspace with `pnpm orbian:source workspace`; maintainers prepare a standalone release with `pnpm orbian:source `. diff --git a/packages/agent/package.json b/packages/agent/package.json index 99ac1005f..9abff5d0d 100644 --- a/packages/agent/package.json +++ b/packages/agent/package.json @@ -36,7 +36,7 @@ "typebox": "1.1.38" }, "devDependencies": { - "@orbian/node": "https://pkg.voidha.sh/orbian-node/fc444bad5db7da1302f6fa0b1a9bd7cadd800526", + "@effect/sql-pg": "catalog:", "@voidhash/tsconfig": "workspace:*", "typescript": "catalog:", "vite-plus": "catalog:", diff --git a/packages/agent/tests/AgentSessionCore.test.ts b/packages/agent/tests/AgentSessionCore.test.ts index 995aeaf8f..741a01250 100644 --- a/packages/agent/tests/AgentSessionCore.test.ts +++ b/packages/agent/tests/AgentSessionCore.test.ts @@ -5,8 +5,8 @@ import { type Context as PiContext, type Model, } from "@earendil-works/pi-ai"; -import { makeMemoryDurableEntityHost } from "@orbian/node/MemoryDurableEntity"; -import { makeNodeDurableEntitySession } from "@orbian/node/NodeDurableEntitySession"; +import { makeMemoryDurableEntityHost } from "./runtime/MemoryDurableEntity.ts"; +import { makeNodeDurableEntitySession } from "./runtime/NodeDurableEntitySession.ts"; import { Effect } from "effect"; import { describe, expect, it } from "vitest"; diff --git a/packages/agent/tests/AgentSessionPg.integration.test.ts b/packages/agent/tests/AgentSessionPg.integration.test.ts index 27d3304bd..98c7a2fbb 100644 --- a/packages/agent/tests/AgentSessionPg.integration.test.ts +++ b/packages/agent/tests/AgentSessionPg.integration.test.ts @@ -5,11 +5,11 @@ import { type Model, } from "@earendil-works/pi-ai"; import { DurableEntityHost } from "@orbian/sdk/DurableEntity"; -import { makeNodeDurableEntitySession } from "@orbian/node/NodeDurableEntitySession"; +import { makeNodeDurableEntitySession } from "./runtime/NodeDurableEntitySession.ts"; import { PgDurableEntityHostLive, type PgDurableEntityConfig, -} from "@orbian/node/DurableEntity"; +} from "./runtime/DurableEntity.ts"; import { Effect, ManagedRuntime, Redacted } from "effect"; import { describe, expect, it } from "vitest"; diff --git a/packages/agent/tests/SessionLog.test.ts b/packages/agent/tests/SessionLog.test.ts index 5e8e4beb2..e003defcb 100644 --- a/packages/agent/tests/SessionLog.test.ts +++ b/packages/agent/tests/SessionLog.test.ts @@ -1,4 +1,4 @@ -import { makeMemoryDurableEntityHost } from "@orbian/node/MemoryDurableEntity"; +import { makeMemoryDurableEntityHost } from "./runtime/MemoryDurableEntity.ts"; import { Effect } from "effect"; import { describe, expect, it } from "vitest"; diff --git a/packages/agent/tests/runtime/DurableEntity.ts b/packages/agent/tests/runtime/DurableEntity.ts new file mode 100644 index 000000000..6dad91537 --- /dev/null +++ b/packages/agent/tests/runtime/DurableEntity.ts @@ -0,0 +1,252 @@ +import { + DurableEntityHost, + type DurableEntityAddress, + type DurableEntityContext, + type DurableEntityHostShape, + type DurableEntitySession, +} from "@orbian/sdk/DurableEntity"; +import { Context, Effect, Layer, Semaphore } from "effect"; +import { SqlClient } from "effect/unstable/sql"; +import { createHash } from "node:crypto"; + +import { PgPlatformClientLive, type PgPlatformConfig } from "./Postgres.js"; + +/** Postgres connection parameters for the single-node durable entity host. */ +export type PgDurableEntityConfig = PgPlatformConfig; + +/** A persisted entity alarm ready to be dispatched by the Node scheduler. */ +export interface DueDurableEntityAlarm { + readonly address: DurableEntityAddress; + readonly scheduledTime: number; +} + +/** Adapter control plane used by the single-node alarm scheduler. */ +export interface NodeDurableEntityControlShape { + readonly listDueAlarms: ( + now: number, + limit: number, + ) => Effect.Effect>; +} + +/** Exposes persisted alarms to the single-node scheduler. */ +export class NodeDurableEntityControl extends Context.Service< + NodeDurableEntityControl, + NodeDurableEntityControlShape +>()("@voidhash/agent-test/NodeDurableEntityControl") {} + +interface EntityRuntimeState { + readonly lock: Semaphore.Semaphore; + readonly sessions: Map; +} + +interface KeyValueRow { + readonly value: unknown; +} + +interface AlarmRow { + readonly type: string; + readonly id: string; + readonly scheduledTime: number | string; +} + +const runtimeKey = (address: DurableEntityAddress): string => `${address.type}\u0000${address.id}`; + +const schemaName = (address: DurableEntityAddress): string => + `entity_${createHash("sha256").update(runtimeKey(address)).digest("hex").slice(0, 32)}`; + +const encodeJson = (value: unknown): string => { + const encoded = JSON.stringify(value); + if (encoded === undefined) { + throw new TypeError("Durable entity values must be JSON-serializable"); + } + return encoded; +}; + +const ensureTables = (sql: SqlClient.SqlClient) => + sql.withTransaction( + Effect.gen(function* () { + // PostgreSQL's IF NOT EXISTS DDL can still race in its catalog, so host + // processes serialize this tiny bootstrap migration with an advisory lock. + yield* sql`SELECT pg_advisory_xact_lock(hashtext('orbian_entity_schema_v1'))`; + yield* sql` + CREATE TABLE IF NOT EXISTS platform_entity_kv ( + entity_type TEXT NOT NULL, + entity_id TEXT NOT NULL, + key TEXT NOT NULL, + value_json JSONB NOT NULL, + PRIMARY KEY (entity_type, entity_id, key) + ) + `; + yield* sql` + CREATE TABLE IF NOT EXISTS platform_entity_alarms ( + entity_type TEXT NOT NULL, + entity_id TEXT NOT NULL, + scheduled_time BIGINT NOT NULL, + PRIMARY KEY (entity_type, entity_id) + ) + `; + yield* sql` + CREATE INDEX IF NOT EXISTS platform_entity_alarms_due_idx + ON platform_entity_alarms (scheduled_time) + `; + }), + ); + +const makePgHost = (sql: SqlClient.SqlClient): DurableEntityHostShape => { + const runtimeStates = new Map(); + const runDb = (effect: Effect.Effect): Effect.Effect => + effect.pipe(Effect.orDie); + + const stateFor = (address: DurableEntityAddress): EntityRuntimeState => { + const key = runtimeKey(address); + let state = runtimeStates.get(key); + if (!state) { + state = { lock: Semaphore.makeUnsafe(1), sessions: new Map() }; + runtimeStates.set(key, state); + } + return state; + }; + + const contextFor = ( + address: DurableEntityAddress, + state: EntityRuntimeState, + ): DurableEntityContext => ({ + address, + keyValue: { + get: (key) => + runDb( + sql` + SELECT value_json AS "value" + FROM platform_entity_kv + WHERE entity_type = ${address.type} + AND entity_id = ${address.id} + AND key = ${key} + `.pipe(Effect.map((rows) => rows[0]?.value)), + ), + put: (key, value) => + runDb( + Effect.suspend(() => { + const encoded = encodeJson(value); + return sql` + INSERT INTO platform_entity_kv (entity_type, entity_id, key, value_json) + VALUES (${address.type}, ${address.id}, ${key}, ${encoded}::jsonb) + ON CONFLICT (entity_type, entity_id, key) + DO UPDATE SET value_json = EXCLUDED.value_json + `.pipe(Effect.asVoid); + }), + ), + delete: (key) => + runDb( + sql` + DELETE FROM platform_entity_kv + WHERE entity_type = ${address.type} + AND entity_id = ${address.id} + AND key = ${key} + `.pipe(Effect.asVoid), + ), + }, + sql: { + execute: >>( + statement: string, + bindings: ReadonlyArray = [], + ) => { + const schema = schemaName(address); + return runDb( + sql.withTransaction( + Effect.gen(function* () { + yield* sql.unsafe(`CREATE SCHEMA IF NOT EXISTS ${schema}`); + yield* sql.unsafe(`SET LOCAL search_path TO ${schema}, public`); + return yield* sql.unsafe(statement, bindings); + }), + ), + ); + }, + }, + alarm: { + get: runDb( + sql` + SELECT entity_type AS "type", entity_id AS "id", scheduled_time AS "scheduledTime" + FROM platform_entity_alarms + WHERE entity_type = ${address.type} AND entity_id = ${address.id} + `.pipe(Effect.map((rows) => (rows[0] ? Number(rows[0].scheduledTime) : undefined))), + ), + set: (scheduledTime) => + runDb( + sql` + INSERT INTO platform_entity_alarms (entity_type, entity_id, scheduled_time) + VALUES (${address.type}, ${address.id}, ${scheduledTime}) + ON CONFLICT (entity_type, entity_id) + DO UPDATE SET scheduled_time = EXCLUDED.scheduled_time + `.pipe(Effect.asVoid), + ), + delete: runDb( + sql` + DELETE FROM platform_entity_alarms + WHERE entity_type = ${address.type} AND entity_id = ${address.id} + `.pipe(Effect.asVoid), + ), + }, + sessions: { + get: (sessionId) => Effect.sync(() => state.sessions.get(sessionId)), + list: Effect.sync(() => [...state.sessions.values()]), + attach: (session) => Effect.sync(() => void state.sessions.set(session.id, session)), + remove: (sessionId) => Effect.sync(() => void state.sessions.delete(sessionId)), + }, + }); + + return DurableEntityHost.of({ + run: (address, operation) => + Effect.suspend(() => { + const state = stateFor(address); + return state.lock.withPermit(Effect.suspend(() => operation(contextFor(address, state)))); + }), + }); +}; + +const makeControl = (sql: SqlClient.SqlClient): NodeDurableEntityControlShape => ({ + listDueAlarms: (now, limit) => + sql` + SELECT entity_type AS "type", entity_id AS "id", scheduled_time AS "scheduledTime" + FROM platform_entity_alarms + WHERE scheduled_time <= ${now} + ORDER BY scheduled_time ASC, entity_type ASC, entity_id ASC + LIMIT ${Math.max(0, Math.floor(limit))} + `.pipe( + Effect.map((rows) => + rows.map((row) => ({ + address: { type: row.type, id: row.id }, + scheduledTime: Number(row.scheduledTime), + })), + ), + Effect.orDie, + ), +}); + +/** + * Postgres-backed single-node entity layer. Database state and alarms survive + * process restarts; execution locks and active WebSocket sessions are local to + * the one Node process. + */ +export const PgDurableEntityHostLive = ( + config: PgDurableEntityConfig, +): Layer.Layer => + Layer.effect( + DurableEntityHost, + Effect.gen(function* () { + const sql = yield* SqlClient.SqlClient; + yield* ensureTables(sql); + return makePgHost(sql); + }), + ).pipe( + Layer.merge( + Layer.effect( + NodeDurableEntityControl, + Effect.gen(function* () { + const sql = yield* SqlClient.SqlClient; + return makeControl(sql); + }), + ), + ), + Layer.provide(PgPlatformClientLive(config)), + Layer.orDie, + ); diff --git a/packages/agent/tests/runtime/MemoryDurableEntity.ts b/packages/agent/tests/runtime/MemoryDurableEntity.ts new file mode 100644 index 000000000..0e420ba50 --- /dev/null +++ b/packages/agent/tests/runtime/MemoryDurableEntity.ts @@ -0,0 +1,72 @@ +import { + DurableEntityHost, + type DurableEntityContext, + type DurableEntityHostShape, + type DurableEntitySession, +} from "@orbian/sdk/DurableEntity"; +import { Effect, Layer, Semaphore } from "effect"; + +interface MemoryEntityState { + readonly lock: Semaphore.Semaphore; + readonly values: Map; + readonly sessions: Map; + alarm: number | undefined; +} + +const entityKey = (type: string, id: string): string => `${type}\u0000${id}`; + +/** + * Builds an isolated in-memory durable entity host. Operations for one address + * are FIFO-serialized; different addresses may run concurrently. + */ +export const makeMemoryDurableEntityHost = (): DurableEntityHostShape => { + const states = new Map(); + + const stateFor = (type: string, id: string): MemoryEntityState => { + const key = entityKey(type, id); + let state = states.get(key); + if (!state) { + state = { + lock: Semaphore.makeUnsafe(1), + values: new Map(), + sessions: new Map(), + alarm: undefined, + }; + states.set(key, state); + } + return state; + }; + + return DurableEntityHost.of({ + run: (address, operation) => + Effect.suspend(() => { + const state = stateFor(address.type, address.id); + const context: DurableEntityContext = { + address, + keyValue: { + get: (key) => Effect.sync(() => state.values.get(key)), + put: (key, value) => Effect.sync(() => void state.values.set(key, value)), + delete: (key) => Effect.sync(() => void state.values.delete(key)), + }, + alarm: { + get: Effect.sync(() => state.alarm), + set: (scheduledTime) => Effect.sync(() => void (state.alarm = scheduledTime)), + delete: Effect.sync(() => void (state.alarm = undefined)), + }, + sessions: { + get: (sessionId) => Effect.sync(() => state.sessions.get(sessionId)), + list: Effect.sync(() => [...state.sessions.values()]), + attach: (session) => Effect.sync(() => void state.sessions.set(session.id, session)), + remove: (sessionId) => Effect.sync(() => void state.sessions.delete(sessionId)), + }, + }; + return state.lock.withPermit(Effect.suspend(() => operation(context))); + }), + }); +}; + +/** In-memory entity host layer for tests and ephemeral local development. */ +export const MemoryDurableEntityHostLive: Layer.Layer = Layer.sync( + DurableEntityHost, + makeMemoryDurableEntityHost, +); diff --git a/packages/agent/tests/runtime/NodeDurableEntitySession.ts b/packages/agent/tests/runtime/NodeDurableEntitySession.ts new file mode 100644 index 000000000..856707a3a --- /dev/null +++ b/packages/agent/tests/runtime/NodeDurableEntitySession.ts @@ -0,0 +1,36 @@ +import type { DurableEntitySession } from "@orbian/sdk/DurableEntity"; +import { Effect } from "effect"; + +/** Minimal server-side WebSocket surface needed by the entity session adapter. */ +export interface NodeWebSocketLike { + readonly send: (message: string | Uint8Array) => unknown; + readonly close: (code?: number, reason?: string) => unknown; +} + +/** + * Wraps a Node WebSocket connection as a runtime-neutral durable entity + * session. Attachments remain in memory for the connection lifetime. + */ +export const makeNodeDurableEntitySession = ( + id: string, + socket: NodeWebSocketLike, + initialAttachment?: unknown, +): DurableEntitySession => { + let attachment = initialAttachment; + return { + id, + send: (message) => + Effect.sync(() => { + socket.send(message); + }), + close: (code, reason) => + Effect.sync(() => { + socket.close(code, reason); + }), + getAttachment: Effect.sync(() => attachment), + setAttachment: (nextAttachment) => + Effect.sync(() => { + attachment = nextAttachment; + }), + }; +}; diff --git a/packages/agent/tests/runtime/Postgres.ts b/packages/agent/tests/runtime/Postgres.ts new file mode 100644 index 000000000..3c1c84b34 --- /dev/null +++ b/packages/agent/tests/runtime/Postgres.ts @@ -0,0 +1,24 @@ +import * as PgClient from "@effect/sql-pg/PgClient"; +import type { Redacted } from "effect"; +import type { ConnectionOptions } from "node:tls"; + +/** Postgres connection parameters shared by single-node Orbian adapters. */ +export interface PgPlatformConfig { + readonly host: string; + readonly port: number; + readonly database: string; + readonly username: string; + readonly password: Redacted.Redacted; + readonly ssl?: boolean | ConnectionOptions; +} + +/** Builds a Postgres client layer for a single-node Orbian adapter. */ +export const PgPlatformClientLive = (config: PgPlatformConfig) => + PgClient.layer({ + host: config.host, + port: config.port, + database: config.database, + username: config.username, + password: config.password, + ssl: config.ssl, + }); diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 37f9df1c5..b9ded97c6 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -9,6 +9,9 @@ catalogs: '@better-auth/api-key': specifier: ^1.6.23 version: 1.6.23 + '@effect-aws/client-s3': + specifier: 2.0.0-beta.4 + version: 2.0.0-beta.4 '@effect/language-service': specifier: ^0.86.2 version: 0.86.6 @@ -60,6 +63,9 @@ catalogs: '@tanstack/zod-adapter': specifier: 1.167.0 version: 1.167.0 + '@types/nodemailer': + specifier: 8.0.1 + version: 8.0.1 '@types/react': specifier: ^19.1.0 version: 19.1.17 @@ -93,6 +99,12 @@ catalogs: lucide-react: specifier: ^0.555.0 version: 0.555.0 + nodemailer: + specifier: 9.0.3 + version: 9.0.3 + playwright-core: + specifier: 1.61.1 + version: 1.61.1 react: specifier: 19.2.7 version: 19.2.7 @@ -1304,9 +1316,9 @@ importers: specifier: 1.1.38 version: 1.1.38 devDependencies: - '@orbian/node': - specifier: https://pkg.voidha.sh/orbian-node/fc444bad5db7da1302f6fa0b1a9bd7cadd800526 - version: https://pkg.voidha.sh/orbian-node/fc444bad5db7da1302f6fa0b1a9bd7cadd800526(effect@4.0.0-beta.84) + '@effect/sql-pg': + specifier: 'catalog:' + version: 4.0.0-beta.84(effect@4.0.0-beta.84) '@voidhash/tsconfig': specifier: workspace:* version: link:../tsconfig @@ -1984,15 +1996,15 @@ importers: selfhost/entry: dependencies: + '@effect-aws/client-s3': + specifier: 'catalog:' + version: 2.0.0-beta.4(effect@4.0.0-beta.84) '@effect/platform-node': specifier: 'catalog:' version: 4.0.0-beta.84(effect@4.0.0-beta.84)(ioredis@5.9.1) '@effect/sql-pg': specifier: 'catalog:' version: 4.0.0-beta.84(effect@4.0.0-beta.84) - '@orbian/node': - specifier: https://pkg.voidha.sh/orbian-node/fc444bad5db7da1302f6fa0b1a9bd7cadd800526 - version: https://pkg.voidha.sh/orbian-node/fc444bad5db7da1302f6fa0b1a9bd7cadd800526(effect@4.0.0-beta.84) '@orbian/sdk': specifier: https://pkg.voidha.sh/orbian-sdk/fc444bad5db7da1302f6fa0b1a9bd7cadd800526 version: https://pkg.voidha.sh/orbian-sdk/fc444bad5db7da1302f6fa0b1a9bd7cadd800526(effect@4.0.0-beta.84) @@ -2044,6 +2056,12 @@ importers: jose: specifier: 'catalog:' version: 6.1.3 + nodemailer: + specifier: 'catalog:' + version: 9.0.3 + playwright-core: + specifier: 'catalog:' + version: 1.61.1 preact: specifier: ^10.25.4 version: 10.29.2 @@ -2066,6 +2084,9 @@ importers: '@types/node': specifier: ^24.0.12 version: 24.10.4 + '@types/nodemailer': + specifier: 'catalog:' + version: 8.0.1 '@types/ws': specifier: ^8.18.1 version: 8.18.1 @@ -5008,26 +5029,8 @@ packages: resolution: {integrity: sha512-a61ljmRVVyG5MC/698C8/FfFDw5a8LOIvyOLW5fztgUXqUpc1jOfQzOitSCbge657OgXXThmY3Tk8fpiDb4UcA==} engines: {node: '>= 20.0.0'} - '@orbian/core@https://pkg.voidha.sh/orbian-core/fc444ba': - resolution: {integrity: sha512-YxCHTGwpot5BvvILMBmBa1M5WbkJgIm8AySmFUAyVLUQwEKOjzyYSLDPmCSxKKoutWxrOw3m7HtRx2Q9Nc1WQw==, tarball: https://pkg.voidha.sh/orbian-core/fc444ba} - version: 0.0.0-gfc444ba - peerDependencies: - effect: 4.0.0-beta.84 - - '@orbian/node@https://pkg.voidha.sh/orbian-node/fc444bad5db7da1302f6fa0b1a9bd7cadd800526': - resolution: {integrity: sha512-2jPjjl/oXbyi+2zai62MJxcf3ZMXQnIw2dgTnP4k7cxp9vQBl5UFfM2THY3SR5y/sgzGvVmtR8GNJiD3IrkXEQ==, tarball: https://pkg.voidha.sh/orbian-node/fc444bad5db7da1302f6fa0b1a9bd7cadd800526} - version: 0.0.0-gfc444ba - peerDependencies: - effect: 4.0.0-beta.84 - - '@orbian/sdk@https://pkg.voidha.sh/orbian-sdk/fc444ba': - resolution: {integrity: sha512-3gM8r6ZIaZGWiLAGLtxWwbse4gmfmr/dX/t6BhYAqV13cxUsZR9M4hJoA0rSDXv/Kdy5oD2ynZCSPF/R7Bsc+w==, tarball: https://pkg.voidha.sh/orbian-sdk/fc444ba} - version: 0.0.0-gfc444ba - peerDependencies: - effect: 4.0.0-beta.84 - '@orbian/sdk@https://pkg.voidha.sh/orbian-sdk/fc444bad5db7da1302f6fa0b1a9bd7cadd800526': - resolution: {integrity: sha512-3gM8r6ZIaZGWiLAGLtxWwbse4gmfmr/dX/t6BhYAqV13cxUsZR9M4hJoA0rSDXv/Kdy5oD2ynZCSPF/R7Bsc+w==, tarball: https://pkg.voidha.sh/orbian-sdk/fc444bad5db7da1302f6fa0b1a9bd7cadd800526} + resolution: {tarball: https://pkg.voidha.sh/orbian-sdk/fc444bad5db7da1302f6fa0b1a9bd7cadd800526} version: 0.0.0-gfc444ba peerDependencies: effect: 4.0.0-beta.84 @@ -7966,6 +7969,9 @@ packages: '@types/node@24.10.4': resolution: {integrity: sha512-vnDVpYPMzs4wunl27jHrfmwojOGKya0xyM3sH+UE5iv5uPS6vX7UIoh6m+vQc5LGBq52HBKPIn/zcSZVzeDEZg==} + '@types/nodemailer@8.0.1': + resolution: {integrity: sha512-PxpaInm8V1JQDd4j0ds5HfvWQk8JupS1C0Picb96QJsrrRDjBH+DlK7L4ZdNSqNULhiZRQHc40nLVShaGxXAMw==} + '@types/offscreencanvas@2019.7.3': resolution: {integrity: sha512-ieXiYmgSRXUDeOntE1InxjWyvEelZGP63M+cGuquuRLuIKKT1osnkXjxev9B7d1nXSug5vpunx+gNlbVxMlC9A==} @@ -19441,27 +19447,6 @@ snapshots: '@orama/orama@3.1.18': {} - '@orbian/core@https://pkg.voidha.sh/orbian-core/fc444ba(effect@4.0.0-beta.84)': - dependencies: - '@orbian/sdk': https://pkg.voidha.sh/orbian-sdk/fc444ba(effect@4.0.0-beta.84) - effect: 4.0.0-beta.84 - - '@orbian/node@https://pkg.voidha.sh/orbian-node/fc444bad5db7da1302f6fa0b1a9bd7cadd800526(effect@4.0.0-beta.84)': - dependencies: - '@effect-aws/client-s3': 2.0.0-beta.4(effect@4.0.0-beta.84) - '@effect/sql-pg': 4.0.0-beta.84(effect@4.0.0-beta.84) - '@orbian/core': https://pkg.voidha.sh/orbian-core/fc444ba(effect@4.0.0-beta.84) - '@orbian/sdk': https://pkg.voidha.sh/orbian-sdk/fc444ba(effect@4.0.0-beta.84) - effect: 4.0.0-beta.84 - nodemailer: 9.0.3 - playwright-core: 1.61.1 - transitivePeerDependencies: - - pg-native - - '@orbian/sdk@https://pkg.voidha.sh/orbian-sdk/fc444ba(effect@4.0.0-beta.84)': - dependencies: - effect: 4.0.0-beta.84 - '@orbian/sdk@https://pkg.voidha.sh/orbian-sdk/fc444bad5db7da1302f6fa0b1a9bd7cadd800526(effect@4.0.0-beta.84)': dependencies: effect: 4.0.0-beta.84 @@ -23387,6 +23372,10 @@ snapshots: dependencies: undici-types: 7.16.0 + '@types/nodemailer@8.0.1': + dependencies: + '@types/node': 24.10.4 + '@types/offscreencanvas@2019.7.3': {} '@types/pg@8.20.0': diff --git a/scripts/set-orbian-source.mjs b/scripts/set-orbian-source.mjs index addb1b533..87832c67f 100644 --- a/scripts/set-orbian-source.mjs +++ b/scripts/set-orbian-source.mjs @@ -17,12 +17,10 @@ const packageOrigin = (process.env.ORBIAN_PACKAGE_ORIGIN ?? "https://pkg.voidha. ); const packageProjects = new Map([ ["@orbian/core", "orbian-core"], - ["@orbian/node", "orbian-node"], ["@orbian/sdk", "orbian-sdk"], ]); const workspacePackages = [ "../orbian/packages/core", - "../orbian/packages/node", "../orbian/packages/sdk", ]; diff --git a/selfhost/entry/package.json b/selfhost/entry/package.json index 7b24c4e33..bdb0a9fcb 100644 --- a/selfhost/entry/package.json +++ b/selfhost/entry/package.json @@ -20,6 +20,7 @@ "test": "vp test run -c vitest.mts" }, "dependencies": { + "@effect-aws/client-s3": "catalog:", "@effect/platform-node": "catalog:", "@effect/sql-pg": "catalog:", "@voidhash/agent": "workspace:*", @@ -36,10 +37,11 @@ "@voidhash/paywall-renderer-web-core": "workspace:*", "@voidhash/paywalls": "workspace:*", "@orbian/sdk": "https://pkg.voidha.sh/orbian-sdk/fc444bad5db7da1302f6fa0b1a9bd7cadd800526", - "@orbian/node": "https://pkg.voidha.sh/orbian-node/fc444bad5db7da1302f6fa0b1a9bd7cadd800526", "effect": "catalog:", "esbuild": "^0.25.10", "jose": "catalog:", + "nodemailer": "catalog:", + "playwright-core": "catalog:", "preact": "^10.25.4", "react": "catalog:", "react-dom": "catalog:", @@ -49,6 +51,7 @@ }, "devDependencies": { "@types/node": "^24.0.12", + "@types/nodemailer": "catalog:", "@types/ws": "^8.18.1", "@voidhash/mimic-core": "workspace:*", "@voidhash/tsconfig": "workspace:*", diff --git a/selfhost/entry/src/DurableEntityAlarms.ts b/selfhost/entry/src/DurableEntityAlarms.ts index 5859f238f..222b34933 100644 --- a/selfhost/entry/src/DurableEntityAlarms.ts +++ b/selfhost/entry/src/DurableEntityAlarms.ts @@ -1,5 +1,5 @@ import type { DurableEntityAddress } from "@orbian/sdk/DurableEntity"; -import type { NodeDurableEntityControlShape } from "@orbian/node/DurableEntity"; +import type { NodeDurableEntityControlShape } from "./runtime/DurableEntity.ts"; import { Effect } from "effect"; /** Handler for one durable-entity alarm type. */ diff --git a/selfhost/entry/src/agent/AgentNodeWebSocket.ts b/selfhost/entry/src/agent/AgentNodeWebSocket.ts index 9dc9e00eb..41513257f 100644 --- a/selfhost/entry/src/agent/AgentNodeWebSocket.ts +++ b/selfhost/entry/src/agent/AgentNodeWebSocket.ts @@ -20,7 +20,7 @@ import { AgentSessionIndexService, LocalUserSessionService, Workos } from "@void import type { AuthTokenVerifier } from "@voidhash/core/services/auth/AuthTokenVerifier"; import { Db } from "@voidhash/db"; import type { DurableEntityHostShape } from "@orbian/sdk/DurableEntity"; -import { makeNodeDurableEntitySession } from "@orbian/node/NodeDurableEntitySession"; +import { makeNodeDurableEntitySession } from "../runtime/NodeDurableEntitySession.ts"; import { Context, Effect, Redacted } from "effect"; import * as HttpHeaders from "effect/unstable/http/Headers"; import { WebSocketServer, type RawData } from "ws"; diff --git a/selfhost/entry/src/backend/Analytics.ts b/selfhost/entry/src/backend/Analytics.ts index afa108347..180a96965 100644 --- a/selfhost/entry/src/backend/Analytics.ts +++ b/selfhost/entry/src/backend/Analytics.ts @@ -25,9 +25,9 @@ import { Db } from "@voidhash/db"; import { KeyValueStore } from "@orbian/sdk/KeyValueStore"; import { PlatformRuntime } from "@orbian/sdk/PlatformRuntime"; import { QueueDriver } from "@orbian/sdk/Queue"; -import { PgKeyValueStoreLive } from "@orbian/node/KeyValueStore"; -import { NodePlatformRuntimeLive } from "@orbian/node/PlatformRuntime"; -import { PgQueueLive } from "@orbian/node/Queue"; +import { PgKeyValueStoreLive } from "../runtime/KeyValueStore.ts"; +import { NodePlatformRuntimeLive } from "../runtime/PlatformRuntime.ts"; +import { PgQueueLive } from "../runtime/Queue.ts"; import { Context, Effect, Layer, Redacted } from "effect"; import type { SelfhostRuntimeConfig } from "../config.ts"; diff --git a/selfhost/entry/src/backend/Backend.ts b/selfhost/entry/src/backend/Backend.ts index 9a9038d8b..6b6c6b2fc 100644 --- a/selfhost/entry/src/backend/Backend.ts +++ b/selfhost/entry/src/backend/Backend.ts @@ -11,7 +11,7 @@ import type { PublicFileStore } from "@voidhash/core/services/storage/PublicFile import { PaywallAssetConfig } from "@voidhash/core/services/paywallLocations/PaywallAssetConfig"; import { Db } from "@voidhash/db"; import { HostServiceTag } from "@voidhash/mimic-db/app/hostService"; -import { NodePlatformRuntimeLive } from "@orbian/node/PlatformRuntime"; +import { NodePlatformRuntimeLive } from "../runtime/PlatformRuntime.ts"; import { Effect, Layer } from "effect"; import type { SelfhostRuntimeConfig, SelfhostWorkosConfig } from "../config.ts"; diff --git a/selfhost/entry/src/backend/Background.ts b/selfhost/entry/src/backend/Background.ts index 5a6046b77..5a7d6baac 100644 --- a/selfhost/entry/src/backend/Background.ts +++ b/selfhost/entry/src/backend/Background.ts @@ -11,7 +11,7 @@ import { CronScheduler, } from "@orbian/sdk/CronScheduler"; import { PlatformRuntime } from "@orbian/sdk/PlatformRuntime"; -import { PgCronSchedulerLive } from "@orbian/node/CronScheduler"; +import { PgCronSchedulerLive } from "../runtime/CronScheduler.ts"; import { Context, Effect, Layer, Redacted } from "effect"; import type { SelfhostRuntimeConfig } from "../config.ts"; diff --git a/selfhost/entry/src/backend/ObjectStores.ts b/selfhost/entry/src/backend/ObjectStores.ts index 0b38c4f97..03257b8e3 100644 --- a/selfhost/entry/src/backend/ObjectStores.ts +++ b/selfhost/entry/src/backend/ObjectStores.ts @@ -11,7 +11,7 @@ import { PlatformRuntime } from "@orbian/sdk/PlatformRuntime"; import { S3ObjectStoreLive, type S3ObjectStoreConfig, -} from "@orbian/node/ObjectStore"; +} from "../runtime/ObjectStore.ts"; import { Effect, Layer, Option } from "effect"; const objectStoreCause = (cause: unknown): string => diff --git a/selfhost/entry/src/backend/Thumbnails.ts b/selfhost/entry/src/backend/Thumbnails.ts index e080e6a3b..03f0df788 100644 --- a/selfhost/entry/src/backend/Thumbnails.ts +++ b/selfhost/entry/src/backend/Thumbnails.ts @@ -22,8 +22,8 @@ import { Screenshot } from "@orbian/sdk/Screenshot"; import { ChromiumScreenshotLive, type ChromiumScreenshotConfig, -} from "@orbian/node/Screenshot"; -import { NodePlatformRuntimeLive } from "@orbian/node/PlatformRuntime"; +} from "../runtime/Screenshot.ts"; +import { NodePlatformRuntimeLive } from "../runtime/PlatformRuntime.ts"; import { Cause, Effect, Layer } from "effect"; import { mimicDocumentIdleQueueName } from "../mimic/MimicDocumentIdleQueue.ts"; diff --git a/selfhost/entry/src/backend/WorkflowPorts.ts b/selfhost/entry/src/backend/WorkflowPorts.ts index fbb850d83..afb984449 100644 --- a/selfhost/entry/src/backend/WorkflowPorts.ts +++ b/selfhost/entry/src/backend/WorkflowPorts.ts @@ -8,8 +8,8 @@ import { WebhookDeliveryWorkflow } from "@voidhash/core/services/webhookDispatch import { Db } from "@voidhash/db"; import { PlatformRuntime } from "@orbian/sdk/PlatformRuntime"; import { WorkflowRunner } from "@orbian/sdk/Workflow"; -import { NodePlatformRuntimeLive } from "@orbian/node/PlatformRuntime"; -import { PgWorkflowRunnerLive } from "@orbian/node/Workflow"; +import { NodePlatformRuntimeLive } from "../runtime/PlatformRuntime.ts"; +import { PgWorkflowRunnerLive } from "../runtime/Workflow.ts"; import { Effect, Layer, Redacted } from "effect"; import type { SelfhostRuntimeConfig } from "../config.ts"; diff --git a/selfhost/entry/src/config.ts b/selfhost/entry/src/config.ts index ee18bedb2..c241ab049 100644 --- a/selfhost/entry/src/config.ts +++ b/selfhost/entry/src/config.ts @@ -1,6 +1,6 @@ import type { DbConfig } from "@voidhash/db/db"; -import type { SmtpMailerConfig } from "@orbian/node/Mailer"; -import type { S3ObjectStoreConfig } from "@orbian/node/ObjectStore"; +import type { SmtpMailerConfig } from "./runtime/Mailer.ts"; +import type { S3ObjectStoreConfig } from "./runtime/ObjectStore.ts"; import { Redacted } from "effect"; const positiveIntegerFromEnv = (name: string, fallback: number): number => { diff --git a/selfhost/entry/src/main.ts b/selfhost/entry/src/main.ts index a2c40c729..f915b8c4d 100644 --- a/selfhost/entry/src/main.ts +++ b/selfhost/entry/src/main.ts @@ -17,8 +17,8 @@ import { HostServiceTag } from "@voidhash/mimic-db/app/hostService"; import { getConfig as getMimicConfig } from "@voidhash/mimic-db/config"; import { makeRoutesLive } from "@voidhash/mimic-db/http/rpc-app"; import { DurableEntityHost } from "@orbian/sdk/DurableEntity"; -import { NodeDurableEntityControl } from "@orbian/node/DurableEntity"; -import { SmtpMailerLive } from "@orbian/node/Mailer"; +import { NodeDurableEntityControl } from "./runtime/DurableEntity.ts"; +import { SmtpMailerLive } from "./runtime/Mailer.ts"; import { Context, Effect, Layer } from "effect"; import { HttpRouter } from "effect/unstable/http"; import { HttpApiBuilder } from "effect/unstable/httpapi"; diff --git a/selfhost/entry/src/mimic/MimicNode.ts b/selfhost/entry/src/mimic/MimicNode.ts index 6eed75ac7..ce17672bf 100644 --- a/selfhost/entry/src/mimic/MimicNode.ts +++ b/selfhost/entry/src/mimic/MimicNode.ts @@ -16,7 +16,7 @@ import { NodeDurableEntityControl, type PgDurableEntityConfig, PgDurableEntityHostLive, -} from "@orbian/node/DurableEntity"; +} from "../runtime/DurableEntity.ts"; import { Effect, Layer } from "effect"; import { PgControlStoreLive } from "./PgControlStore.ts"; diff --git a/selfhost/entry/src/mimic/MimicNodeWebSocket.ts b/selfhost/entry/src/mimic/MimicNodeWebSocket.ts index 60150652b..cfd14e5af 100644 --- a/selfhost/entry/src/mimic/MimicNodeWebSocket.ts +++ b/selfhost/entry/src/mimic/MimicNodeWebSocket.ts @@ -20,8 +20,8 @@ import { type DurableEntityContext, makeDurableEntityAddress, } from "@orbian/sdk/DurableEntity"; -import type { NodeDurableEntityControlShape } from "@orbian/node/DurableEntity"; -import { makeNodeDurableEntitySession } from "@orbian/node/NodeDurableEntitySession"; +import type { NodeDurableEntityControlShape } from "../runtime/DurableEntity.ts"; +import { makeNodeDurableEntitySession } from "../runtime/NodeDurableEntitySession.ts"; import { Duration, Effect, Fiber, Semaphore } from "effect"; import WebSocket, { WebSocketServer, type RawData } from "ws"; diff --git a/selfhost/entry/src/mimic/PgControlStore.ts b/selfhost/entry/src/mimic/PgControlStore.ts index 423b16088..e2c369975 100644 --- a/selfhost/entry/src/mimic/PgControlStore.ts +++ b/selfhost/entry/src/mimic/PgControlStore.ts @@ -10,7 +10,7 @@ import type { UserRecord, } from "@voidhash/mimic-db/core/store"; import { ControlStore } from "@voidhash/mimic-db/core/store"; -import type { PgDurableEntityConfig } from "@orbian/node/DurableEntity"; +import type { PgDurableEntityConfig } from "../runtime/DurableEntity.ts"; import { Effect, Layer } from "effect"; import { SqlClient } from "effect/unstable/sql"; diff --git a/selfhost/entry/src/mimic/config.ts b/selfhost/entry/src/mimic/config.ts index b9e414b18..7cc957016 100644 --- a/selfhost/entry/src/mimic/config.ts +++ b/selfhost/entry/src/mimic/config.ts @@ -1,5 +1,5 @@ import { makePgDocumentConfig } from "@voidhash/mimic-db/core/pg-store"; -import type { PgDurableEntityConfig } from "@orbian/node/DurableEntity"; +import type { PgDurableEntityConfig } from "../runtime/DurableEntity.ts"; import { Redacted } from "effect"; import type { MimicNodeConfig } from "./MimicNode.ts"; diff --git a/selfhost/entry/src/mimic/main.ts b/selfhost/entry/src/mimic/main.ts index 96bc16d4e..173959680 100644 --- a/selfhost/entry/src/mimic/main.ts +++ b/selfhost/entry/src/mimic/main.ts @@ -7,9 +7,9 @@ import { makeRoutesLive } from "@voidhash/mimic-db/http/rpc-app"; import { DurableEntityHost } from "@orbian/sdk/DurableEntity"; import { NodeDurableEntityControl, -} from "@orbian/node/DurableEntity"; -import { NodePlatformRuntimeLive } from "@orbian/node/PlatformRuntime"; -import { PgQueueLive } from "@orbian/node/Queue"; +} from "../runtime/DurableEntity.ts"; +import { NodePlatformRuntimeLive } from "../runtime/PlatformRuntime.ts"; +import { PgQueueLive } from "../runtime/Queue.ts"; import { Context, Effect, Layer } from "effect"; import { HttpRouter } from "effect/unstable/http"; diff --git a/selfhost/entry/src/runtime/CronScheduler.ts b/selfhost/entry/src/runtime/CronScheduler.ts new file mode 100644 index 000000000..400d0574c --- /dev/null +++ b/selfhost/entry/src/runtime/CronScheduler.ts @@ -0,0 +1,240 @@ +import { + type CronJob, + type CronRunOptions, + CronScheduler, + CronSchedulerError, + type CronSchedulerShape, +} from "@orbian/sdk/CronScheduler"; +import { PlatformRuntime } from "@orbian/sdk/PlatformRuntime"; +import { Cause, Cron, Duration, Effect, Layer, Result } from "effect"; +import { SqlClient } from "effect/unstable/sql"; + +import { PgPlatformClientLive, type PgPlatformConfig } from "./Postgres.js"; + +interface ClaimedRunRow { + readonly leaseToken: string; + readonly scheduledTime: number | string; +} + +interface ChangedRow { + readonly jobName: string; +} + +const ensureTable = (sql: SqlClient.SqlClient) => + sql.withTransaction( + Effect.gen(function* () { + yield* sql`SELECT pg_advisory_xact_lock(hashtext('orbian_cron_schema_v1'))`; + yield* sql` + CREATE TABLE IF NOT EXISTS platform_cron_state ( + job_name TEXT PRIMARY KEY, + expression TEXT NOT NULL, + time_zone TEXT, + last_scheduled_at_ms BIGINT, + next_scheduled_at_ms BIGINT NOT NULL, + lease_token TEXT, + leased_until_ms BIGINT, + updated_at_ms BIGINT NOT NULL + ) + `; + yield* sql` + CREATE INDEX IF NOT EXISTS platform_cron_state_due_idx + ON platform_cron_state (next_scheduled_at_ms) + `; + }), + ); + +const schedulerError = (jobName: string, operation: string, cause: unknown) => + new CronSchedulerError({ jobName, operation, cause: String(cause) }); + +const parseCron = (job: CronJob) => { + const parsed = Cron.parse(job.expression, job.timeZone); + return Result.isSuccess(parsed) + ? Effect.succeed(parsed.success) + : Effect.fail(schedulerError(job.name, "parse", parsed.failure.message)); +}; + +const leaseMillis = (job: CronJob): number => { + const value = job.leaseMillis; + return typeof value === "number" && Number.isFinite(value) && value > 0 + ? Math.floor(value) + : 300_000; +}; + +const claimRun = ( + sql: SqlClient.SqlClient, + job: CronJob, + cron: Cron.Cron, + now: number, +) => { + const initialNext = Cron.next(cron, new Date(now - 1)).getTime(); + const token = crypto.randomUUID(); + const timeZone = job.timeZone ?? null; + return sql.withTransaction( + Effect.gen(function* () { + yield* sql` + INSERT INTO platform_cron_state ( + job_name, expression, time_zone, next_scheduled_at_ms, updated_at_ms + ) VALUES ( + ${job.name}, ${job.expression}, ${timeZone}, ${initialNext}, ${now} + ) + ON CONFLICT (job_name) + DO UPDATE SET + next_scheduled_at_ms = EXCLUDED.next_scheduled_at_ms, + expression = EXCLUDED.expression, + time_zone = EXCLUDED.time_zone, + lease_token = NULL, + leased_until_ms = NULL, + updated_at_ms = EXCLUDED.updated_at_ms + WHERE platform_cron_state.expression IS DISTINCT FROM EXCLUDED.expression + OR platform_cron_state.time_zone IS DISTINCT FROM EXCLUDED.time_zone + `; + return yield* sql` + WITH candidate AS ( + SELECT job_name + FROM platform_cron_state + WHERE job_name = ${job.name} + AND next_scheduled_at_ms <= ${now} + AND (leased_until_ms IS NULL OR leased_until_ms <= ${now}) + FOR UPDATE SKIP LOCKED + ) + UPDATE platform_cron_state AS state + SET lease_token = ${token}, + leased_until_ms = ${now + leaseMillis(job)}, + updated_at_ms = ${now} + FROM candidate + WHERE state.job_name = candidate.job_name + RETURNING state.next_scheduled_at_ms AS "scheduledTime", state.lease_token AS "leaseToken" + `; + }), + ); +}; + +const releaseRun = (sql: SqlClient.SqlClient, jobName: string, leaseToken: string, now: number) => + sql` + UPDATE platform_cron_state + SET lease_token = NULL, leased_until_ms = NULL, updated_at_ms = ${now} + WHERE job_name = ${jobName} AND lease_token = ${leaseToken} + `; + +const completeRun = ( + sql: SqlClient.SqlClient, + jobName: string, + leaseToken: string, + scheduledTime: number, + nextScheduledTime: number, + now: number, +) => + sql` + UPDATE platform_cron_state + SET last_scheduled_at_ms = ${scheduledTime}, + next_scheduled_at_ms = ${nextScheduledTime}, + lease_token = NULL, + leased_until_ms = NULL, + updated_at_ms = ${now} + WHERE job_name = ${jobName} AND lease_token = ${leaseToken} + RETURNING job_name AS "jobName" + `; + +const tick = ( + sql: SqlClient.SqlClient, + job: CronJob, + inputNow: Date | undefined, +): Effect.Effect => + PlatformRuntime.pipe( + Effect.andThen(parseCron(job)), + Effect.flatMap((cron) => + Effect.suspend(() => { + const now = inputNow?.getTime() ?? Date.now(); + return claimRun(sql, job, cron, now).pipe( + Effect.mapError((cause) => schedulerError(job.name, "claim", cause)), + Effect.flatMap((rows) => { + const claimed = rows[0]; + if (!claimed) return Effect.succeed(false); + const scheduledTime = Number(claimed.scheduledTime); + const nextScheduledTime = Cron.next(cron, new Date(scheduledTime)).getTime(); + return job + .run({ + scheduledTime: new Date(scheduledTime), + catchUp: scheduledTime < now, + }) + .pipe( + Effect.matchCauseEffect({ + onFailure: (cause) => + releaseRun(sql, job.name, claimed.leaseToken, Date.now()).pipe( + Effect.mapError((releaseCause) => + schedulerError(job.name, "release", releaseCause), + ), + Effect.andThen( + Effect.fail(schedulerError(job.name, "run", Cause.pretty(cause))), + ), + ), + onSuccess: () => + completeRun( + sql, + job.name, + claimed.leaseToken, + scheduledTime, + nextScheduledTime, + Date.now(), + ).pipe( + Effect.flatMap((changed) => + changed.length === 1 + ? Effect.succeed(true) + : Effect.fail( + schedulerError( + job.name, + "complete", + "cron lease expired before completion", + ), + ), + ), + Effect.mapError((cause) => + cause instanceof CronSchedulerError + ? cause + : schedulerError(job.name, "complete", cause), + ), + ), + }), + ); + }), + ); + }), + ), + ); + +const makeScheduler = (sql: SqlClient.SqlClient): CronSchedulerShape => ({ + tick: (job, now) => tick(sql, job, now), + run: (job, options: CronRunOptions = {}) => { + const pollInterval = + typeof options.pollIntervalMillis === "number" && + Number.isFinite(options.pollIntervalMillis) && + options.pollIntervalMillis > 0 + ? Math.floor(options.pollIntervalMillis) + : 1_000; + return Effect.forever( + tick(sql, job, undefined).pipe( + Effect.catch((error) => + Effect.logError("scheduled job tick failed", { + jobName: job.name, + operation: error.operation, + cause: error.cause, + }).pipe(Effect.as(false)), + ), + Effect.flatMap((ran) => + ran ? Effect.yieldNow : Effect.sleep(Duration.millis(pollInterval)), + ), + ), + ); + }, +}); + +/** Postgres-backed cron scheduler with persisted catch-up and execution leases. */ +export const PgCronSchedulerLive = (config: PgPlatformConfig): Layer.Layer => + Layer.effect( + CronScheduler, + Effect.gen(function* () { + const sql = yield* SqlClient.SqlClient; + yield* ensureTable(sql); + return makeScheduler(sql); + }), + ).pipe(Layer.provide(PgPlatformClientLive(config)), Layer.orDie); diff --git a/selfhost/entry/src/runtime/DurableEntity.ts b/selfhost/entry/src/runtime/DurableEntity.ts new file mode 100644 index 000000000..4f96082c3 --- /dev/null +++ b/selfhost/entry/src/runtime/DurableEntity.ts @@ -0,0 +1,252 @@ +import { + DurableEntityHost, + type DurableEntityAddress, + type DurableEntityContext, + type DurableEntityHostShape, + type DurableEntitySession, +} from "@orbian/sdk/DurableEntity"; +import { Context, Effect, Layer, Semaphore } from "effect"; +import { SqlClient } from "effect/unstable/sql"; +import { createHash } from "node:crypto"; + +import { PgPlatformClientLive, type PgPlatformConfig } from "./Postgres.js"; + +/** Postgres connection parameters for the single-node durable entity host. */ +export type PgDurableEntityConfig = PgPlatformConfig; + +/** A persisted entity alarm ready to be dispatched by the Node scheduler. */ +export interface DueDurableEntityAlarm { + readonly address: DurableEntityAddress; + readonly scheduledTime: number; +} + +/** Adapter control plane used by the single-node alarm scheduler. */ +export interface NodeDurableEntityControlShape { + readonly listDueAlarms: ( + now: number, + limit: number, + ) => Effect.Effect>; +} + +/** Exposes persisted alarms to the single-node scheduler. */ +export class NodeDurableEntityControl extends Context.Service< + NodeDurableEntityControl, + NodeDurableEntityControlShape +>()("@voidhash/selfhost-entry/NodeDurableEntityControl") {} + +interface EntityRuntimeState { + readonly lock: Semaphore.Semaphore; + readonly sessions: Map; +} + +interface KeyValueRow { + readonly value: unknown; +} + +interface AlarmRow { + readonly type: string; + readonly id: string; + readonly scheduledTime: number | string; +} + +const runtimeKey = (address: DurableEntityAddress): string => `${address.type}\u0000${address.id}`; + +const schemaName = (address: DurableEntityAddress): string => + `entity_${createHash("sha256").update(runtimeKey(address)).digest("hex").slice(0, 32)}`; + +const encodeJson = (value: unknown): string => { + const encoded = JSON.stringify(value); + if (encoded === undefined) { + throw new TypeError("Durable entity values must be JSON-serializable"); + } + return encoded; +}; + +const ensureTables = (sql: SqlClient.SqlClient) => + sql.withTransaction( + Effect.gen(function* () { + // PostgreSQL's IF NOT EXISTS DDL can still race in its catalog, so host + // processes serialize this tiny bootstrap migration with an advisory lock. + yield* sql`SELECT pg_advisory_xact_lock(hashtext('orbian_entity_schema_v1'))`; + yield* sql` + CREATE TABLE IF NOT EXISTS platform_entity_kv ( + entity_type TEXT NOT NULL, + entity_id TEXT NOT NULL, + key TEXT NOT NULL, + value_json JSONB NOT NULL, + PRIMARY KEY (entity_type, entity_id, key) + ) + `; + yield* sql` + CREATE TABLE IF NOT EXISTS platform_entity_alarms ( + entity_type TEXT NOT NULL, + entity_id TEXT NOT NULL, + scheduled_time BIGINT NOT NULL, + PRIMARY KEY (entity_type, entity_id) + ) + `; + yield* sql` + CREATE INDEX IF NOT EXISTS platform_entity_alarms_due_idx + ON platform_entity_alarms (scheduled_time) + `; + }), + ); + +const makePgHost = (sql: SqlClient.SqlClient): DurableEntityHostShape => { + const runtimeStates = new Map(); + const runDb = (effect: Effect.Effect): Effect.Effect => + effect.pipe(Effect.orDie); + + const stateFor = (address: DurableEntityAddress): EntityRuntimeState => { + const key = runtimeKey(address); + let state = runtimeStates.get(key); + if (!state) { + state = { lock: Semaphore.makeUnsafe(1), sessions: new Map() }; + runtimeStates.set(key, state); + } + return state; + }; + + const contextFor = ( + address: DurableEntityAddress, + state: EntityRuntimeState, + ): DurableEntityContext => ({ + address, + keyValue: { + get: (key) => + runDb( + sql` + SELECT value_json AS "value" + FROM platform_entity_kv + WHERE entity_type = ${address.type} + AND entity_id = ${address.id} + AND key = ${key} + `.pipe(Effect.map((rows) => rows[0]?.value)), + ), + put: (key, value) => + runDb( + Effect.suspend(() => { + const encoded = encodeJson(value); + return sql` + INSERT INTO platform_entity_kv (entity_type, entity_id, key, value_json) + VALUES (${address.type}, ${address.id}, ${key}, ${encoded}::jsonb) + ON CONFLICT (entity_type, entity_id, key) + DO UPDATE SET value_json = EXCLUDED.value_json + `.pipe(Effect.asVoid); + }), + ), + delete: (key) => + runDb( + sql` + DELETE FROM platform_entity_kv + WHERE entity_type = ${address.type} + AND entity_id = ${address.id} + AND key = ${key} + `.pipe(Effect.asVoid), + ), + }, + sql: { + execute: >>( + statement: string, + bindings: ReadonlyArray = [], + ) => { + const schema = schemaName(address); + return runDb( + sql.withTransaction( + Effect.gen(function* () { + yield* sql.unsafe(`CREATE SCHEMA IF NOT EXISTS ${schema}`); + yield* sql.unsafe(`SET LOCAL search_path TO ${schema}, public`); + return yield* sql.unsafe(statement, bindings); + }), + ), + ); + }, + }, + alarm: { + get: runDb( + sql` + SELECT entity_type AS "type", entity_id AS "id", scheduled_time AS "scheduledTime" + FROM platform_entity_alarms + WHERE entity_type = ${address.type} AND entity_id = ${address.id} + `.pipe(Effect.map((rows) => (rows[0] ? Number(rows[0].scheduledTime) : undefined))), + ), + set: (scheduledTime) => + runDb( + sql` + INSERT INTO platform_entity_alarms (entity_type, entity_id, scheduled_time) + VALUES (${address.type}, ${address.id}, ${scheduledTime}) + ON CONFLICT (entity_type, entity_id) + DO UPDATE SET scheduled_time = EXCLUDED.scheduled_time + `.pipe(Effect.asVoid), + ), + delete: runDb( + sql` + DELETE FROM platform_entity_alarms + WHERE entity_type = ${address.type} AND entity_id = ${address.id} + `.pipe(Effect.asVoid), + ), + }, + sessions: { + get: (sessionId) => Effect.sync(() => state.sessions.get(sessionId)), + list: Effect.sync(() => [...state.sessions.values()]), + attach: (session) => Effect.sync(() => void state.sessions.set(session.id, session)), + remove: (sessionId) => Effect.sync(() => void state.sessions.delete(sessionId)), + }, + }); + + return DurableEntityHost.of({ + run: (address, operation) => + Effect.suspend(() => { + const state = stateFor(address); + return state.lock.withPermit(Effect.suspend(() => operation(contextFor(address, state)))); + }), + }); +}; + +const makeControl = (sql: SqlClient.SqlClient): NodeDurableEntityControlShape => ({ + listDueAlarms: (now, limit) => + sql` + SELECT entity_type AS "type", entity_id AS "id", scheduled_time AS "scheduledTime" + FROM platform_entity_alarms + WHERE scheduled_time <= ${now} + ORDER BY scheduled_time ASC, entity_type ASC, entity_id ASC + LIMIT ${Math.max(0, Math.floor(limit))} + `.pipe( + Effect.map((rows) => + rows.map((row) => ({ + address: { type: row.type, id: row.id }, + scheduledTime: Number(row.scheduledTime), + })), + ), + Effect.orDie, + ), +}); + +/** + * Postgres-backed single-node entity layer. Database state and alarms survive + * process restarts; execution locks and active WebSocket sessions are local to + * the one Node process. + */ +export const PgDurableEntityHostLive = ( + config: PgDurableEntityConfig, +): Layer.Layer => + Layer.effect( + DurableEntityHost, + Effect.gen(function* () { + const sql = yield* SqlClient.SqlClient; + yield* ensureTables(sql); + return makePgHost(sql); + }), + ).pipe( + Layer.merge( + Layer.effect( + NodeDurableEntityControl, + Effect.gen(function* () { + const sql = yield* SqlClient.SqlClient; + return makeControl(sql); + }), + ), + ), + Layer.provide(PgPlatformClientLive(config)), + Layer.orDie, + ); diff --git a/selfhost/entry/src/runtime/KeyValueStore.ts b/selfhost/entry/src/runtime/KeyValueStore.ts new file mode 100644 index 000000000..755d2724b --- /dev/null +++ b/selfhost/entry/src/runtime/KeyValueStore.ts @@ -0,0 +1,250 @@ +import { + type KeyValuePutOptions, + KeyValueStore, + KeyValueStoreError, + type KeyValueStoreShape, +} from "@orbian/sdk/KeyValueStore"; +import { PlatformRuntime } from "@orbian/sdk/PlatformRuntime"; +import { Effect, Layer, Option, SchemaParser } from "effect"; +import { SqlClient } from "effect/unstable/sql"; + +import { PgPlatformClientLive, type PgPlatformConfig } from "./Postgres.js"; + +interface ValueRow { + readonly value: unknown; +} + +interface KeyRow { + readonly key: string; +} + +const ensureTable = (sql: SqlClient.SqlClient) => + sql.withTransaction( + Effect.gen(function* () { + yield* sql`SELECT pg_advisory_xact_lock(hashtext('orbian_kv_schema_v1'))`; + yield* sql` + CREATE TABLE IF NOT EXISTS platform_key_value ( + namespace TEXT NOT NULL, + key TEXT NOT NULL, + value_json JSONB NOT NULL, + expires_at_ms BIGINT, + updated_at_ms BIGINT NOT NULL, + PRIMARY KEY (namespace, key) + ) + `; + yield* sql` + CREATE INDEX IF NOT EXISTS platform_key_value_expiry_idx + ON platform_key_value (expires_at_ms) + WHERE expires_at_ms IS NOT NULL + `; + }), + ); + +const expiry = (now: number, options: KeyValuePutOptions | undefined): number | null => { + const ttl = options?.ttlMillis; + return typeof ttl === "number" && Number.isFinite(ttl) && ttl > 0 ? now + Math.floor(ttl) : null; +}; + +const encodeJson = (value: unknown): Effect.Effect => + Effect.try({ + try: () => { + const encoded = JSON.stringify(value); + if (encoded === undefined) { + throw new TypeError("Key-value entries must be JSON-serializable"); + } + return encoded; + }, + catch: (cause) => cause, + }); + +const storeError = (namespace: string, operation: string, cause: unknown) => + new KeyValueStoreError({ namespace, operation, cause: String(cause) }); + +const makeStore = (sql: SqlClient.SqlClient): KeyValueStoreShape => ({ + get: (namespace, key, schema) => + PlatformRuntime.pipe( + Effect.andThen( + Effect.suspend(() => { + const now = Date.now(); + return sql` + SELECT value_json AS "value" + FROM platform_key_value + WHERE namespace = ${namespace} + AND key = ${key} + AND (expires_at_ms IS NULL OR expires_at_ms > ${now}) + `; + }), + ), + Effect.flatMap((rows) => + rows[0] + ? SchemaParser.decodeUnknownEffect(schema)(rows[0].value).pipe(Effect.map(Option.some)) + : Effect.succeedNone, + ), + Effect.mapError((cause) => storeError(namespace, "get", cause)), + ), + put: (namespace, key, value, schema, options) => + PlatformRuntime.pipe( + Effect.andThen(SchemaParser.encodeUnknownEffect(schema)(value)), + Effect.flatMap(encodeJson), + Effect.flatMap((encoded) => { + const now = Date.now(); + const expiresAt = expiry(now, options); + return sql` + INSERT INTO platform_key_value ( + namespace, key, value_json, expires_at_ms, updated_at_ms + ) VALUES ( + ${namespace}, ${key}, ${encoded}::jsonb, ${expiresAt}, ${now} + ) + ON CONFLICT (namespace, key) + DO UPDATE SET value_json = EXCLUDED.value_json, + expires_at_ms = EXCLUDED.expires_at_ms, + updated_at_ms = EXCLUDED.updated_at_ms + `; + }), + Effect.asVoid, + Effect.mapError((cause) => storeError(namespace, "put", cause)), + ), + putMany: (namespace, entries, schema, options) => + PlatformRuntime.pipe( + Effect.andThen( + Effect.forEach(entries, ({ key, value }) => + SchemaParser.encodeUnknownEffect(schema)(value).pipe( + Effect.flatMap(encodeJson), + Effect.map((encoded) => ({ key, encoded })), + ), + ), + ), + Effect.flatMap((encodedEntries) => { + const now = Date.now(); + const expiresAt = expiry(now, options); + return sql.withTransaction( + Effect.forEach( + encodedEntries, + ({ key, encoded }) => + sql` + INSERT INTO platform_key_value ( + namespace, key, value_json, expires_at_ms, updated_at_ms + ) VALUES ( + ${namespace}, ${key}, ${encoded}::jsonb, ${expiresAt}, ${now} + ) + ON CONFLICT (namespace, key) + DO UPDATE SET value_json = EXCLUDED.value_json, + expires_at_ms = EXCLUDED.expires_at_ms, + updated_at_ms = EXCLUDED.updated_at_ms + `, + { discard: true }, + ), + ); + }), + Effect.mapError((cause) => storeError(namespace, "putMany", cause)), + ), + existingKeys: (namespace, keys) => { + if (keys.length === 0) return Effect.succeed(new Set()); + return PlatformRuntime.pipe( + Effect.andThen( + Effect.suspend(() => { + const now = Date.now(); + return sql` + SELECT key + FROM platform_key_value + WHERE namespace = ${namespace} + AND ${sql.in("key", keys)} + AND (expires_at_ms IS NULL OR expires_at_ms > ${now}) + `; + }), + ), + Effect.map((rows) => new Set(rows.map(({ key }) => key))), + Effect.mapError((cause) => storeError(namespace, "existingKeys", cause)), + ); + }, + delete: (namespace, key) => + PlatformRuntime.pipe( + Effect.andThen( + sql`DELETE FROM platform_key_value WHERE namespace = ${namespace} AND key = ${key}`, + ), + Effect.asVoid, + Effect.mapError((cause) => storeError(namespace, "delete", cause)), + ), + deleteMany: (namespace, keys) => { + if (keys.length === 0) return Effect.void; + return PlatformRuntime.pipe( + Effect.andThen( + sql` + DELETE FROM platform_key_value + WHERE namespace = ${namespace} AND ${sql.in("key", keys)} + `, + ), + Effect.asVoid, + Effect.mapError((cause) => storeError(namespace, "deleteMany", cause)), + ); + }, + increment: (namespace, key, options) => + PlatformRuntime.pipe( + Effect.andThen( + Effect.suspend(() => { + const now = Date.now(); + const expiresAt = expiry(now, options); + return sql` + INSERT INTO platform_key_value ( + namespace, key, value_json, expires_at_ms, updated_at_ms + ) VALUES ( + ${namespace}, ${key}, '1'::jsonb, ${expiresAt}, ${now} + ) + ON CONFLICT (namespace, key) + DO UPDATE SET + value_json = to_jsonb( + CASE + WHEN platform_key_value.expires_at_ms IS NOT NULL + AND platform_key_value.expires_at_ms <= ${now} + THEN 1 + WHEN jsonb_typeof(platform_key_value.value_json) = 'number' + THEN (platform_key_value.value_json #>> '{}')::BIGINT + 1 + ELSE 1 + END + ), + expires_at_ms = EXCLUDED.expires_at_ms, + updated_at_ms = EXCLUDED.updated_at_ms + RETURNING value_json AS "value" + `; + }), + ), + Effect.map((rows) => Number(rows[0]?.value ?? 0)), + Effect.mapError((cause) => storeError(namespace, "increment", cause)), + ), + pruneExpired: (limit) => + PlatformRuntime.pipe( + Effect.andThen( + Effect.suspend(() => { + const now = Date.now(); + const rowLimit = Number.isFinite(limit) && limit > 0 ? Math.floor(limit) : 0; + return sql` + WITH expired AS ( + SELECT namespace, key + FROM platform_key_value + WHERE expires_at_ms IS NOT NULL AND expires_at_ms <= ${now} + ORDER BY expires_at_ms ASC + LIMIT ${rowLimit} + FOR UPDATE SKIP LOCKED + ) + DELETE FROM platform_key_value AS value + USING expired + WHERE value.namespace = expired.namespace AND value.key = expired.key + RETURNING value.key + `; + }), + ), + Effect.map((rows) => rows.length), + Effect.mapError((cause) => storeError("*", "pruneExpired", cause)), + ), +}); + +/** Postgres-backed typed key-value store with TTL and atomic counters. */ +export const PgKeyValueStoreLive = (config: PgPlatformConfig): Layer.Layer => + Layer.effect( + KeyValueStore, + Effect.gen(function* () { + const sql = yield* SqlClient.SqlClient; + yield* ensureTable(sql); + return makeStore(sql); + }), + ).pipe(Layer.provide(PgPlatformClientLive(config)), Layer.orDie); diff --git a/selfhost/entry/src/runtime/Mailer.ts b/selfhost/entry/src/runtime/Mailer.ts new file mode 100644 index 000000000..efd795121 --- /dev/null +++ b/selfhost/entry/src/runtime/Mailer.ts @@ -0,0 +1,137 @@ +import { + type MailAddress, + type MailDeliveryResult, + Mailer, + MailerError, + type MailerShape, + type MailMessage, +} from "@orbian/sdk/Mailer"; +import { PlatformRuntime } from "@orbian/sdk/PlatformRuntime"; +import { Effect, Layer, Redacted } from "effect"; +import nodemailer, { type Transporter } from "nodemailer"; + +/** SMTP connection, authentication, and sender defaults. */ +export interface SmtpMailerConfig { + readonly host: string; + readonly port: number; + readonly secure?: boolean; + readonly requireTls?: boolean; + readonly username?: string; + readonly password?: Redacted.Redacted; + readonly defaultFrom?: MailAddress; + readonly connectionTimeoutMillis?: number; + readonly greetingTimeoutMillis?: number; + readonly tlsRejectUnauthorized?: boolean; + readonly verifyOnStart?: boolean; +} + +const mailerError = (operation: string, cause: unknown) => + new MailerError({ operation, cause: String(cause) }); + +const mailbox = ({ address, name }: MailAddress) => ({ address, name }); + +const resultAddress = (value: unknown): string => { + if (typeof value === "string") return value; + if (typeof value === "object" && value !== null && "address" in value) { + const address = (value as { readonly address?: unknown }).address; + if (typeof address === "string") return address; + } + return String(value); +}; + +const validateMessage = ( + config: SmtpMailerConfig, + message: MailMessage, +): Effect.Effect => { + const from = message.from ?? config.defaultFrom; + if (!from) return Effect.fail(mailerError("validate", "sender is required")); + if (message.to.length === 0) { + return Effect.fail(mailerError("validate", "at least one recipient is required")); + } + if (message.text === undefined && message.html === undefined) { + return Effect.fail(mailerError("validate", "text or html content is required")); + } + return Effect.succeed(from); +}; + +const send = ( + transporter: Transporter, + config: SmtpMailerConfig, + message: MailMessage, +): Effect.Effect => + PlatformRuntime.pipe( + Effect.andThen(validateMessage(config, message)), + Effect.flatMap((from) => + Effect.tryPromise({ + try: () => + transporter.sendMail({ + from: mailbox(from), + to: message.to.map(mailbox), + cc: message.cc?.map(mailbox), + bcc: message.bcc?.map(mailbox), + replyTo: message.replyTo ? mailbox(message.replyTo) : undefined, + subject: message.subject, + text: message.text, + html: message.html, + headers: message.headers, + }), + catch: (cause) => mailerError("send", cause), + }), + ), + Effect.map((info) => ({ + messageId: info.messageId, + accepted: (info.accepted as ReadonlyArray).map(resultAddress), + rejected: (info.rejected as ReadonlyArray).map(resultAddress), + })), + ); + +const makeMailer = (transporter: Transporter, config: SmtpMailerConfig): MailerShape => ({ + send: (message) => send(transporter, config, message), +}); + +const makeTransporter = (config: SmtpMailerConfig): Effect.Effect => { + if ((config.username === undefined) !== (config.password === undefined)) { + return Effect.fail( + mailerError("configure", "SMTP username and password must be provided together"), + ); + } + return Effect.succeed( + nodemailer.createTransport({ + host: config.host, + port: config.port, + secure: config.secure ?? false, + requireTLS: config.requireTls ?? false, + auth: + config.username && config.password + ? { + user: config.username, + pass: Redacted.value(config.password), + } + : undefined, + connectionTimeout: config.connectionTimeoutMillis ?? 10_000, + greetingTimeout: config.greetingTimeoutMillis ?? 10_000, + tls: { + rejectUnauthorized: config.tlsRejectUnauthorized ?? true, + }, + }), + ); +}; + +/** SMTP-backed mailer with optional authenticated TLS and startup verification. */ +export const SmtpMailerLive = (config: SmtpMailerConfig): Layer.Layer => + Layer.effect( + Mailer, + Effect.acquireRelease( + makeTransporter(config).pipe( + Effect.tap((transporter) => + config.verifyOnStart + ? Effect.tryPromise({ + try: () => transporter.verify(), + catch: (cause) => mailerError("verify", cause), + }).pipe(Effect.asVoid) + : Effect.void, + ), + ), + (transporter) => Effect.sync(() => transporter.close()), + ).pipe(Effect.map((transporter) => makeMailer(transporter, config))), + ); diff --git a/selfhost/entry/src/runtime/MemoryDurableEntity.ts b/selfhost/entry/src/runtime/MemoryDurableEntity.ts new file mode 100644 index 000000000..0e420ba50 --- /dev/null +++ b/selfhost/entry/src/runtime/MemoryDurableEntity.ts @@ -0,0 +1,72 @@ +import { + DurableEntityHost, + type DurableEntityContext, + type DurableEntityHostShape, + type DurableEntitySession, +} from "@orbian/sdk/DurableEntity"; +import { Effect, Layer, Semaphore } from "effect"; + +interface MemoryEntityState { + readonly lock: Semaphore.Semaphore; + readonly values: Map; + readonly sessions: Map; + alarm: number | undefined; +} + +const entityKey = (type: string, id: string): string => `${type}\u0000${id}`; + +/** + * Builds an isolated in-memory durable entity host. Operations for one address + * are FIFO-serialized; different addresses may run concurrently. + */ +export const makeMemoryDurableEntityHost = (): DurableEntityHostShape => { + const states = new Map(); + + const stateFor = (type: string, id: string): MemoryEntityState => { + const key = entityKey(type, id); + let state = states.get(key); + if (!state) { + state = { + lock: Semaphore.makeUnsafe(1), + values: new Map(), + sessions: new Map(), + alarm: undefined, + }; + states.set(key, state); + } + return state; + }; + + return DurableEntityHost.of({ + run: (address, operation) => + Effect.suspend(() => { + const state = stateFor(address.type, address.id); + const context: DurableEntityContext = { + address, + keyValue: { + get: (key) => Effect.sync(() => state.values.get(key)), + put: (key, value) => Effect.sync(() => void state.values.set(key, value)), + delete: (key) => Effect.sync(() => void state.values.delete(key)), + }, + alarm: { + get: Effect.sync(() => state.alarm), + set: (scheduledTime) => Effect.sync(() => void (state.alarm = scheduledTime)), + delete: Effect.sync(() => void (state.alarm = undefined)), + }, + sessions: { + get: (sessionId) => Effect.sync(() => state.sessions.get(sessionId)), + list: Effect.sync(() => [...state.sessions.values()]), + attach: (session) => Effect.sync(() => void state.sessions.set(session.id, session)), + remove: (sessionId) => Effect.sync(() => void state.sessions.delete(sessionId)), + }, + }; + return state.lock.withPermit(Effect.suspend(() => operation(context))); + }), + }); +}; + +/** In-memory entity host layer for tests and ephemeral local development. */ +export const MemoryDurableEntityHostLive: Layer.Layer = Layer.sync( + DurableEntityHost, + makeMemoryDurableEntityHost, +); diff --git a/selfhost/entry/src/runtime/NodeDurableEntitySession.ts b/selfhost/entry/src/runtime/NodeDurableEntitySession.ts new file mode 100644 index 000000000..856707a3a --- /dev/null +++ b/selfhost/entry/src/runtime/NodeDurableEntitySession.ts @@ -0,0 +1,36 @@ +import type { DurableEntitySession } from "@orbian/sdk/DurableEntity"; +import { Effect } from "effect"; + +/** Minimal server-side WebSocket surface needed by the entity session adapter. */ +export interface NodeWebSocketLike { + readonly send: (message: string | Uint8Array) => unknown; + readonly close: (code?: number, reason?: string) => unknown; +} + +/** + * Wraps a Node WebSocket connection as a runtime-neutral durable entity + * session. Attachments remain in memory for the connection lifetime. + */ +export const makeNodeDurableEntitySession = ( + id: string, + socket: NodeWebSocketLike, + initialAttachment?: unknown, +): DurableEntitySession => { + let attachment = initialAttachment; + return { + id, + send: (message) => + Effect.sync(() => { + socket.send(message); + }), + close: (code, reason) => + Effect.sync(() => { + socket.close(code, reason); + }), + getAttachment: Effect.sync(() => attachment), + setAttachment: (nextAttachment) => + Effect.sync(() => { + attachment = nextAttachment; + }), + }; +}; diff --git a/selfhost/entry/src/runtime/ObjectStore.ts b/selfhost/entry/src/runtime/ObjectStore.ts new file mode 100644 index 000000000..ec8a76b00 --- /dev/null +++ b/selfhost/entry/src/runtime/ObjectStore.ts @@ -0,0 +1,126 @@ +import { S3 } from "@effect-aws/client-s3"; +import { ObjectStore, ObjectStoreError, type ObjectStoreShape } from "@orbian/sdk/ObjectStore"; +import { PlatformRuntime } from "@orbian/sdk/PlatformRuntime"; +import { Effect, Layer, Option, Redacted } from "effect"; + +/** S3-compatible connection and bucket parameters. */ +export interface S3ObjectStoreConfig { + readonly bucketName: string; + readonly region: string; + readonly endpoint?: string; + readonly accessKeyId: string; + readonly secretAccessKey: Redacted.Redacted; + readonly forcePathStyle?: boolean; +} + +const storeError = (config: S3ObjectStoreConfig, key: string, operation: string, cause: unknown) => + new ObjectStoreError({ + bucketName: config.bucketName, + key, + operation, + cause: String(cause), + }); + +const isNotFound = (error: unknown): boolean => { + if (typeof error !== "object" || error === null) return false; + const tagged = error as { readonly _tag?: unknown }; + return tagged._tag === "NoSuchKey" || tagged._tag === "NotFound"; +}; + +const makeStore = (config: S3ObjectStoreConfig, client: S3.Type): ObjectStoreShape => ({ + bucketName: config.bucketName, + put: ({ key, body, contentType, cacheControl }) => + PlatformRuntime.pipe( + Effect.andThen( + client.putObject({ + Bucket: config.bucketName, + Key: key, + Body: body, + ContentType: contentType, + CacheControl: cacheControl, + }), + ), + Effect.asVoid, + Effect.mapError((cause) => storeError(config, key, "put", cause)), + ), + get: (key) => + PlatformRuntime.pipe( + Effect.andThen(client.getObject({ Bucket: config.bucketName, Key: key })), + Effect.flatMap((output) => { + const stream = output.Body; + if (!stream) { + return Effect.fail(storeError(config, key, "get", "response body is missing")); + } + return Effect.tryPromise({ + try: () => stream.transformToByteArray(), + catch: (cause) => storeError(config, key, "get", cause), + }).pipe( + Effect.map((body) => + Option.some({ + body, + contentType: output.ContentType ?? null, + etag: output.ETag ?? null, + size: output.ContentLength ?? body.byteLength, + }), + ), + ); + }), + Effect.catch((cause) => + isNotFound(cause) + ? Effect.succeedNone + : Effect.fail( + cause instanceof ObjectStoreError ? cause : storeError(config, key, "get", cause), + ), + ), + ), + head: (key) => + PlatformRuntime.pipe( + Effect.andThen(client.headObject({ Bucket: config.bucketName, Key: key })), + Effect.map((output) => + Option.some({ + contentType: output.ContentType ?? null, + etag: output.ETag ?? null, + size: output.ContentLength ?? 0, + }), + ), + Effect.catch((cause) => + isNotFound(cause) + ? Effect.succeedNone + : Effect.fail(storeError(config, key, "head", cause)), + ), + ), + delete: (key) => + PlatformRuntime.pipe( + Effect.andThen(client.deleteObject({ Bucket: config.bucketName, Key: key })), + Effect.asVoid, + Effect.mapError((cause) => storeError(config, key, "delete", cause)), + ), +}); + +/** S3-compatible object store layer for AWS S3, MinIO, Garage, or R2. */ +export const S3ObjectStoreLive = (config: S3ObjectStoreConfig): Layer.Layer => + Layer.effect( + ObjectStore, + Effect.map(S3, (client) => makeStore(config, client)), + ).pipe( + Layer.provide( + S3.layer({ + region: config.region, + endpoint: config.endpoint, + forcePathStyle: config.forcePathStyle ?? config.endpoint !== undefined, + // AWS SDK v3 enables CRC32 checksums by default. Some S3-compatible + // stores reject those optional headers, so custom endpoints use the + // compatibility mode while native AWS S3 retains its stronger default. + ...(config.endpoint === undefined + ? {} + : { + requestChecksumCalculation: "WHEN_REQUIRED" as const, + responseChecksumValidation: "WHEN_REQUIRED" as const, + }), + credentials: { + accessKeyId: config.accessKeyId, + secretAccessKey: Redacted.value(config.secretAccessKey), + }, + }), + ), + ); diff --git a/selfhost/entry/src/runtime/PgWorkflowEngine.ts b/selfhost/entry/src/runtime/PgWorkflowEngine.ts new file mode 100644 index 000000000..fdcbce856 --- /dev/null +++ b/selfhost/entry/src/runtime/PgWorkflowEngine.ts @@ -0,0 +1,698 @@ +import { Duration, Effect, Exit, Fiber, Layer, Option, Schema, Scope } from "effect"; +import { SqlClient } from "effect/unstable/sql"; +import { Workflow, WorkflowEngine } from "effect/unstable/workflow"; + +import { PgPlatformClientLive, type PgPlatformConfig } from "./Postgres.js"; + +interface RegisteredWorkflow { + readonly workflow: Workflow.Any; + readonly execute: ( + payload: object, + executionId: string, + ) => Effect.Effect< + unknown, + unknown, + WorkflowEngine.WorkflowInstance | WorkflowEngine.WorkflowEngine + >; + readonly scope: Scope.Scope; +} + +interface ClaimedExecutionRow { + readonly executionId: string; + readonly interrupted: boolean; + readonly leaseToken: string; + readonly parentExecutionId: string | null; + readonly payload: unknown; + readonly workflowName: string; +} + +interface ExecutionResultRow { + readonly result: unknown; +} + +interface ExecutionIdRow { + readonly executionId: string; +} + +interface ActivityResultRow { + readonly result: unknown; +} + +interface DeferredResultRow { + readonly exit: unknown; +} + +interface DueClockRow { + readonly deferredName: string; + readonly executionId: string; + readonly exit: unknown; + readonly workflowName: string; +} + +interface ActiveExecution { + readonly instance: WorkflowEngine.WorkflowInstance["Service"]; + interrupt: Effect.Effect; +} + +const activityResultCodec = Schema.toCodecJson( + Workflow.Result({ success: Schema.Any, error: Schema.Any }), +); + +type JsonCodec = Schema.Codec; + +const jsonCodec = (schema: Schema.Top): JsonCodec => + Schema.toCodecJson(schema) as unknown as JsonCodec; + +const encodeJson = (value: unknown): Effect.Effect => + Effect.try({ + try: () => { + const encoded = JSON.stringify(value); + if (encoded === undefined) throw new TypeError("Workflow state must be JSON-serializable"); + return encoded; + }, + catch: (cause) => cause, + }); + +const ensureTables = (sql: SqlClient.SqlClient) => + sql.withTransaction( + Effect.gen(function* () { + yield* sql`SELECT pg_advisory_xact_lock(hashtext('orbian_workflow_schema_v1'))`; + yield* sql` + CREATE TABLE IF NOT EXISTS platform_workflow_execution ( + execution_id TEXT PRIMARY KEY, + workflow_name TEXT NOT NULL, + payload_json JSONB NOT NULL, + parent_execution_id TEXT, + status TEXT NOT NULL, + result_json JSONB, + interrupted BOOLEAN NOT NULL DEFAULT FALSE, + resume_requested BOOLEAN NOT NULL DEFAULT FALSE, + lease_token TEXT, + leased_until_ms BIGINT, + created_at_ms BIGINT NOT NULL, + updated_at_ms BIGINT NOT NULL + ) + `; + yield* sql` + ALTER TABLE platform_workflow_execution + ADD COLUMN IF NOT EXISTS resume_requested BOOLEAN NOT NULL DEFAULT FALSE + `; + yield* sql` + CREATE INDEX IF NOT EXISTS platform_workflow_execution_claim_idx + ON platform_workflow_execution (status, leased_until_ms, updated_at_ms) + `; + yield* sql` + CREATE TABLE IF NOT EXISTS platform_workflow_activity ( + execution_id TEXT NOT NULL, + activity_name TEXT NOT NULL, + attempt INTEGER NOT NULL, + result_json JSONB NOT NULL, + updated_at_ms BIGINT NOT NULL, + PRIMARY KEY (execution_id, activity_name, attempt) + ) + `; + yield* sql` + CREATE TABLE IF NOT EXISTS platform_workflow_deferred ( + execution_id TEXT NOT NULL, + deferred_name TEXT NOT NULL, + exit_json JSONB NOT NULL, + updated_at_ms BIGINT NOT NULL, + PRIMARY KEY (execution_id, deferred_name) + ) + `; + yield* sql` + CREATE TABLE IF NOT EXISTS platform_workflow_clock ( + execution_id TEXT NOT NULL, + workflow_name TEXT NOT NULL, + deferred_name TEXT NOT NULL, + due_at_ms BIGINT NOT NULL, + exit_json JSONB NOT NULL, + created_at_ms BIGINT NOT NULL, + PRIMARY KEY (execution_id, deferred_name) + ) + `; + yield* sql` + CREATE INDEX IF NOT EXISTS platform_workflow_clock_due_idx + ON platform_workflow_clock (due_at_ms) + `; + }), + ); + +const workflowResultCodec = (workflow: Workflow.Any) => + Schema.toCodecJson( + Workflow.Result({ + success: workflow.successSchema, + error: workflow.errorSchema, + }), + ) as unknown as JsonCodec; + +const encodePayload = (workflow: Workflow.Any, payload: object) => + Schema.encodeEffect(jsonCodec(workflow.payloadSchema))(payload).pipe(Effect.flatMap(encodeJson)); + +const decodePayload = (workflow: Workflow.Any, payload: unknown) => + Schema.decodeEffect(jsonCodec(workflow.payloadSchema))(payload as Schema.Json); + +const encodeWorkflowResult = (workflow: Workflow.Any, result: Workflow.Result) => + Schema.encodeEffect(workflowResultCodec(workflow))(result).pipe(Effect.flatMap(encodeJson)); + +const decodeWorkflowResult = (workflow: Workflow.Any, result: unknown) => + Schema.decodeEffect(workflowResultCodec(workflow))(result as Schema.Json) as Effect.Effect< + Workflow.Result + >; + +const encodeActivityResult = (result: Workflow.Result) => + Schema.encodeEffect(activityResultCodec)(result).pipe(Effect.flatMap(encodeJson)); + +const decodeActivityResult = (result: unknown) => + Schema.decodeEffect(activityResultCodec)(result as Schema.Json); + +const makePgEngine = (sql: SqlClient.SqlClient) => + Effect.gen(function* () { + const scope = yield* Effect.scope; + const registered = new Map(); + const active = new Map(); + const leaseMillis = 30_000; + + let engine: WorkflowEngine.WorkflowEngine["Service"]; + + const claimExecution = (executionId: string) => + Effect.suspend(() => { + const now = Date.now(); + const token = crypto.randomUUID(); + return sql` + WITH candidate AS ( + SELECT execution_id + FROM platform_workflow_execution + WHERE execution_id = ${executionId} + AND status IN ('pending', 'running') + AND (leased_until_ms IS NULL OR leased_until_ms <= ${now}) + FOR UPDATE SKIP LOCKED + ) + UPDATE platform_workflow_execution AS execution + SET status = 'running', + result_json = NULL, + resume_requested = FALSE, + lease_token = ${token}, + leased_until_ms = ${now + leaseMillis}, + updated_at_ms = ${now} + FROM candidate + WHERE execution.execution_id = candidate.execution_id + RETURNING execution.execution_id AS "executionId", + execution.workflow_name AS "workflowName", + execution.payload_json AS "payload", + execution.parent_execution_id AS "parentExecutionId", + execution.interrupted, + execution.lease_token AS "leaseToken" + `; + }); + + const releaseUnregistered = (executionId: string, leaseToken: string) => + Effect.suspend(() => { + const now = Date.now(); + return sql` + UPDATE platform_workflow_execution + SET status = 'pending', lease_token = NULL, leased_until_ms = NULL, updated_at_ms = ${now} + WHERE execution_id = ${executionId} AND lease_token = ${leaseToken} + `; + }); + + const releaseFailedExecution = (executionId: string, leaseToken: string) => + Effect.suspend(() => { + const now = Date.now(); + return sql` + UPDATE platform_workflow_execution + SET status = 'pending', lease_token = NULL, leased_until_ms = NULL, updated_at_ms = ${now} + WHERE execution_id = ${executionId} AND lease_token = ${leaseToken} + `; + }); + + const persistExecutionResult = ( + row: ClaimedExecutionRow, + workflow: Workflow.Any, + result: Workflow.Result, + ) => + Effect.gen(function* () { + const encoded = yield* encodeWorkflowResult(workflow, result); + const now = Date.now(); + const suspended = result._tag === "Suspended"; + const changed = yield* sql` + UPDATE platform_workflow_execution + SET status = CASE + WHEN ${suspended} AND resume_requested THEN 'pending' + ELSE ${suspended ? "suspended" : "complete"} + END, + result_json = CASE + WHEN ${suspended} AND resume_requested THEN NULL + ELSE ${encoded}::jsonb + END, + resume_requested = FALSE, + lease_token = NULL, + leased_until_ms = NULL, + updated_at_ms = ${now} + WHERE execution_id = ${row.executionId} AND lease_token = ${row.leaseToken} + RETURNING execution_id AS "executionId" + `; + return changed.length === 1; + }); + + const heartbeatLease = (row: ClaimedExecutionRow) => + Effect.forever( + Effect.sleep(leaseMillis / 3).pipe( + Effect.andThen( + Effect.suspend(() => { + const now = Date.now(); + return sql` + UPDATE platform_workflow_execution + SET leased_until_ms = ${now + leaseMillis}, updated_at_ms = ${now} + WHERE execution_id = ${row.executionId} + AND lease_token = ${row.leaseToken} + AND status = 'running' + `; + }), + ), + Effect.catchCause((cause) => + Effect.logError("workflow execution lease heartbeat failed", { + executionId: row.executionId, + workflowName: row.workflowName, + cause: String(cause), + }), + ), + ), + ); + + const runClaimedExecution = ( + row: ClaimedExecutionRow, + entry: RegisteredWorkflow, + instance: WorkflowEngine.WorkflowInstance["Service"], + ) => + Effect.scoped( + Effect.gen(function* () { + yield* heartbeatLease(row).pipe(Effect.forkScoped); + const payload = yield* decodePayload(entry.workflow, row.payload); + const execution = instance.interrupted + ? Effect.interrupt + : entry.execute(payload as object, row.executionId); + const result = yield* execution.pipe( + Workflow.intoResult, + Effect.provideService(WorkflowEngine.WorkflowInstance, instance), + Effect.provideService(WorkflowEngine.WorkflowEngine, engine), + ); + const persisted = yield* persistExecutionResult(row, entry.workflow, result); + if (persisted && result._tag === "Complete" && row.parentExecutionId) { + yield* sql` + UPDATE platform_workflow_execution + SET resume_requested = TRUE, + status = CASE WHEN status = 'suspended' THEN 'pending' ELSE status END, + result_json = CASE WHEN status = 'suspended' THEN NULL ELSE result_json END, + lease_token = CASE WHEN status = 'suspended' THEN NULL ELSE lease_token END, + leased_until_ms = CASE + WHEN status = 'suspended' THEN NULL + ELSE leased_until_ms + END, + updated_at_ms = ${Date.now()} + WHERE execution_id = ${row.parentExecutionId} + AND status IN ('running', 'suspended') + `; + yield* startExecution(row.parentExecutionId); + } + }), + ).pipe( + Effect.onExit((exit) => + Exit.isFailure(exit) + ? releaseFailedExecution(row.executionId, row.leaseToken).pipe( + Effect.tap(() => + Effect.logError("workflow execution failed before producing a result", { + executionId: row.executionId, + workflowName: row.workflowName, + cause: String(exit.cause), + }), + ), + ) + : Effect.void, + ), + ); + + function startExecution(executionId: string): Effect.Effect { + if (active.has(executionId)) return Effect.void; + return Effect.gen(function* () { + if (active.has(executionId)) return; + const rows = yield* claimExecution(executionId); + const row = rows[0]; + if (!row) return; + const entry = registered.get(row.workflowName); + if (!entry) { + yield* releaseUnregistered(executionId, row.leaseToken); + return; + } + + const instance = WorkflowEngine.WorkflowInstance.initial(entry.workflow, executionId); + instance.interrupted = row.interrupted; + const activeExecution: ActiveExecution = { instance, interrupt: Effect.void }; + active.set(executionId, activeExecution); + const fiber = yield* runClaimedExecution(row, entry, instance).pipe( + Effect.ensuring( + Effect.sync(() => { + active.delete(executionId); + }), + ), + Effect.forkIn(entry.scope), + ); + activeExecution.interrupt = Fiber.interrupt(fiber); + }); + } + + const pollResult = (workflow: Workflow.Any, executionId: string) => + sql` + SELECT result_json AS "result" + FROM platform_workflow_execution + WHERE execution_id = ${executionId} AND result_json IS NOT NULL + `.pipe( + Effect.flatMap((rows) => + rows[0] + ? decodeWorkflowResult(workflow, rows[0].result).pipe(Effect.map(Option.some)) + : Effect.succeedNone, + ), + ); + + function waitForResult( + workflow: Workflow.Any, + executionId: string, + ): Effect.Effect, unknown> { + return pollResult(workflow, executionId).pipe( + Effect.flatMap((result) => + Option.isSome(result) + ? Effect.succeed(result.value) + : Effect.sleep("20 millis").pipe(Effect.andThen(waitForResult(workflow, executionId))), + ), + ); + } + + const ensureExecution = ( + workflow: Workflow.Any, + executionId: string, + payload: object, + parent: WorkflowEngine.WorkflowInstance["Service"] | undefined, + ) => + Effect.gen(function* () { + const encoded = yield* encodePayload(workflow, payload); + const now = Date.now(); + yield* sql` + INSERT INTO platform_workflow_execution ( + execution_id, workflow_name, payload_json, parent_execution_id, + status, created_at_ms, updated_at_ms + ) VALUES ( + ${executionId}, ${workflow._tag}, ${encoded}::jsonb, + ${parent?.executionId ?? null}, 'pending', ${now}, ${now} + ) + ON CONFLICT (execution_id) DO NOTHING + `; + }); + + const processDueClocks = Effect.gen(function* () { + const now = Date.now(); + const rows = yield* sql.withTransaction( + Effect.gen(function* () { + const due = yield* sql` + SELECT execution_id AS "executionId", workflow_name AS "workflowName", + deferred_name AS "deferredName", exit_json AS "exit" + FROM platform_workflow_clock + WHERE due_at_ms <= ${now} + ORDER BY due_at_ms ASC + LIMIT 50 + FOR UPDATE SKIP LOCKED + `; + yield* Effect.forEach( + due, + (clock) => + Effect.gen(function* () { + const encoded = yield* encodeJson(clock.exit); + yield* sql` + INSERT INTO platform_workflow_deferred ( + execution_id, deferred_name, exit_json, updated_at_ms + ) VALUES ( + ${clock.executionId}, ${clock.deferredName}, ${encoded}::jsonb, ${now} + ) + ON CONFLICT (execution_id, deferred_name) DO NOTHING + `; + yield* sql` + UPDATE platform_workflow_execution + SET resume_requested = TRUE, + status = CASE WHEN status = 'suspended' THEN 'pending' ELSE status END, + result_json = CASE + WHEN status = 'suspended' THEN NULL + ELSE result_json + END, + lease_token = CASE + WHEN status = 'suspended' THEN NULL + ELSE lease_token + END, + leased_until_ms = CASE + WHEN status = 'suspended' THEN NULL + ELSE leased_until_ms + END, + updated_at_ms = ${now} + WHERE execution_id = ${clock.executionId} + AND status IN ('running', 'suspended') + `; + yield* sql` + DELETE FROM platform_workflow_clock + WHERE execution_id = ${clock.executionId} + AND deferred_name = ${clock.deferredName} + `; + }), + { discard: true }, + ); + return due; + }), + ); + yield* Effect.forEach(rows, (row) => startExecution(row.executionId), { + discard: true, + }); + }); + + const recoverExecutions = Effect.gen(function* () { + const workflowNames = [...registered.keys()]; + if (workflowNames.length === 0) return; + const now = Date.now(); + const rows = yield* sql` + SELECT execution_id AS "executionId" + FROM platform_workflow_execution + WHERE ${sql.in("workflow_name", workflowNames)} + AND status IN ('pending', 'running') + AND (leased_until_ms IS NULL OR leased_until_ms <= ${now}) + ORDER BY updated_at_ms ASC + LIMIT 50 + `; + yield* Effect.forEach(rows, (row) => startExecution(row.executionId), { + discard: true, + }); + }); + + const tick = processDueClocks.pipe(Effect.andThen(recoverExecutions)); + + const executeWorkflow = ( + workflow: Workflow.Any, + options: { + readonly executionId: string; + readonly payload: object; + readonly discard: Discard; + readonly parent?: WorkflowEngine.WorkflowInstance["Service"] | undefined; + }, + ): Effect.Effect> => + Effect.gen(function* () { + yield* ensureExecution(workflow, options.executionId, options.payload, options.parent); + yield* startExecution(options.executionId); + if (options.discard) return; + return yield* waitForResult(workflow, options.executionId); + }).pipe(Effect.orDie) as Effect.Effect< + Discard extends true ? void : Workflow.Result + >; + + engine = WorkflowEngine.makeUnsafe({ + register: (workflow, execute) => + Effect.gen(function* () { + registered.set(workflow._tag, { + workflow, + execute, + scope: yield* Effect.scope, + }); + yield* recoverExecutions; + }).pipe(Effect.orDie), + execute: executeWorkflow, + poll: (workflow, executionId) => pollResult(workflow, executionId).pipe(Effect.orDie), + interrupt: (_workflow, executionId) => + Effect.gen(function* () { + yield* sql` + UPDATE platform_workflow_execution + SET interrupted = TRUE, resume_requested = FALSE, + status = 'pending', result_json = NULL, + lease_token = NULL, leased_until_ms = NULL, updated_at_ms = ${Date.now()} + WHERE execution_id = ${executionId} AND status <> 'complete' + `; + const running = active.get(executionId); + if (running) { + running.instance.interrupted = true; + yield* running.interrupt; + } else { + yield* startExecution(executionId); + } + }).pipe(Effect.orDie), + interruptUnsafe: (_workflow, executionId) => + Effect.gen(function* () { + yield* sql` + UPDATE platform_workflow_execution + SET interrupted = TRUE, resume_requested = FALSE, + status = 'pending', result_json = NULL, + lease_token = NULL, leased_until_ms = NULL, updated_at_ms = ${Date.now()} + WHERE execution_id = ${executionId} AND status <> 'complete' + `; + const running = active.get(executionId); + if (running) { + running.instance.interrupted = true; + yield* running.interrupt; + } else { + yield* startExecution(executionId); + } + }).pipe(Effect.orDie), + resume: (_workflow, executionId) => + Effect.gen(function* () { + yield* sql` + UPDATE platform_workflow_execution + SET interrupted = FALSE, resume_requested = FALSE, + status = 'pending', result_json = NULL, + lease_token = NULL, leased_until_ms = NULL, updated_at_ms = ${Date.now()} + WHERE execution_id = ${executionId} AND status <> 'complete' + `; + yield* startExecution(executionId); + }).pipe(Effect.orDie), + activityExecute: (activity, attempt) => + Effect.gen(function* () { + const parent = yield* WorkflowEngine.WorkflowInstance; + const rows = yield* sql` + SELECT result_json AS "result" + FROM platform_workflow_activity + WHERE execution_id = ${parent.executionId} + AND activity_name = ${activity.name} + AND attempt = ${attempt} + `; + if (rows[0]) { + const stored = yield* decodeActivityResult(rows[0].result); + if (stored._tag === "Complete") return stored; + yield* sql` + DELETE FROM platform_workflow_activity + WHERE execution_id = ${parent.executionId} + AND activity_name = ${activity.name} + AND attempt = ${attempt} + `; + } + + const instance = WorkflowEngine.WorkflowInstance.initial( + parent.workflow, + parent.executionId, + ); + instance.interrupted = parent.interrupted; + const result = yield* activity.executeEncoded.pipe( + Workflow.intoResult, + Effect.provideService(WorkflowEngine.WorkflowInstance, instance), + ); + const encoded = yield* encodeActivityResult(result); + yield* sql` + INSERT INTO platform_workflow_activity ( + execution_id, activity_name, attempt, result_json, updated_at_ms + ) VALUES ( + ${parent.executionId}, ${activity.name}, ${attempt}, ${encoded}::jsonb, + ${Date.now()} + ) + ON CONFLICT (execution_id, activity_name, attempt) + DO NOTHING + `; + return result; + }).pipe(Effect.orDie), + deferredResult: (deferred) => + Effect.gen(function* () { + const instance = yield* WorkflowEngine.WorkflowInstance; + const rows = yield* sql` + SELECT exit_json AS "exit" + FROM platform_workflow_deferred + WHERE execution_id = ${instance.executionId} + AND deferred_name = ${deferred.name} + `; + return rows[0] ? Option.some(rows[0].exit as Exit.Exit) : Option.none(); + }).pipe(Effect.orDie), + deferredDone: (options) => + Effect.gen(function* () { + const encoded = yield* encodeJson(options.exit); + const now = Date.now(); + yield* sql` + INSERT INTO platform_workflow_deferred ( + execution_id, deferred_name, exit_json, updated_at_ms + ) VALUES ( + ${options.executionId}, ${options.deferredName}, ${encoded}::jsonb, ${now} + ) + ON CONFLICT (execution_id, deferred_name) DO NOTHING + `; + yield* sql` + UPDATE platform_workflow_execution + SET resume_requested = TRUE, + status = CASE WHEN status = 'suspended' THEN 'pending' ELSE status END, + result_json = CASE WHEN status = 'suspended' THEN NULL ELSE result_json END, + lease_token = CASE WHEN status = 'suspended' THEN NULL ELSE lease_token END, + leased_until_ms = CASE + WHEN status = 'suspended' THEN NULL + ELSE leased_until_ms + END, + updated_at_ms = ${now} + WHERE execution_id = ${options.executionId} + AND status IN ('running', 'suspended') + `; + yield* startExecution(options.executionId); + }).pipe(Effect.orDie), + scheduleClock: (workflow, options) => + Effect.gen(function* () { + const exitCodec = options.clock.deferred.exitSchema as unknown as Schema.Codec< + Exit.Exit, + Schema.Json, + never + >; + const encodedExit = yield* Schema.encodeEffect(exitCodec)(Exit.void).pipe( + Effect.flatMap(encodeJson), + ); + const now = Date.now(); + const dueAt = now + Duration.toMillis(options.clock.duration); + yield* sql` + INSERT INTO platform_workflow_clock ( + execution_id, workflow_name, deferred_name, due_at_ms, + exit_json, created_at_ms + ) VALUES ( + ${options.executionId}, ${workflow._tag}, ${options.clock.deferred.name}, + ${dueAt}, ${encodedExit}::jsonb, ${now} + ) + ON CONFLICT (execution_id, deferred_name) DO NOTHING + `; + }).pipe(Effect.orDie), + }); + + yield* Effect.forever( + tick.pipe( + Effect.catchCause((cause) => + Effect.logError("workflow runner tick failed", { cause: String(cause) }), + ), + Effect.andThen(Effect.sleep("100 millis")), + ), + ).pipe(Effect.forkIn(scope)); + + return engine; + }); + +/** Postgres-backed Effect workflow engine with durable activities and clocks. */ +export const PgWorkflowEngineLive = ( + config: PgPlatformConfig, +): Layer.Layer => + Layer.effect( + WorkflowEngine.WorkflowEngine, + Effect.gen(function* () { + const sql = yield* SqlClient.SqlClient; + yield* ensureTables(sql); + return yield* makePgEngine(sql); + }), + ).pipe(Layer.provide(PgPlatformClientLive(config)), Layer.orDie); diff --git a/selfhost/entry/src/runtime/PlatformRuntime.ts b/selfhost/entry/src/runtime/PlatformRuntime.ts new file mode 100644 index 000000000..a0ef730ac --- /dev/null +++ b/selfhost/entry/src/runtime/PlatformRuntime.ts @@ -0,0 +1,6 @@ +import { Layer } from "effect"; + +import { PlatformRuntime } from "@orbian/sdk/PlatformRuntime"; + +/** Platform runtime marker for the single-process Node composition. */ +export const NodePlatformRuntimeLive = Layer.succeed(PlatformRuntime, PlatformRuntime.of({})); diff --git a/selfhost/entry/src/runtime/Postgres.ts b/selfhost/entry/src/runtime/Postgres.ts new file mode 100644 index 000000000..3c1c84b34 --- /dev/null +++ b/selfhost/entry/src/runtime/Postgres.ts @@ -0,0 +1,24 @@ +import * as PgClient from "@effect/sql-pg/PgClient"; +import type { Redacted } from "effect"; +import type { ConnectionOptions } from "node:tls"; + +/** Postgres connection parameters shared by single-node Orbian adapters. */ +export interface PgPlatformConfig { + readonly host: string; + readonly port: number; + readonly database: string; + readonly username: string; + readonly password: Redacted.Redacted; + readonly ssl?: boolean | ConnectionOptions; +} + +/** Builds a Postgres client layer for a single-node Orbian adapter. */ +export const PgPlatformClientLive = (config: PgPlatformConfig) => + PgClient.layer({ + host: config.host, + port: config.port, + database: config.database, + username: config.username, + password: config.password, + ssl: config.ssl, + }); diff --git a/selfhost/entry/src/runtime/Queue.ts b/selfhost/entry/src/runtime/Queue.ts new file mode 100644 index 000000000..8766c72ec --- /dev/null +++ b/selfhost/entry/src/runtime/Queue.ts @@ -0,0 +1,339 @@ +import { + QueueConsumerError, + type QueueConsumerOptions, + QueueDriver, + type QueueDriverShape, + QueueProducerError, + type QueueProducer, +} from "@orbian/sdk/Queue"; +import { PlatformRuntime } from "@orbian/sdk/PlatformRuntime"; +import { Cause, Duration, Effect, Layer, Schema, SchemaParser } from "effect"; +import { SqlClient } from "effect/unstable/sql"; + +import { PgPlatformClientLive, type PgPlatformConfig } from "./Postgres.js"; + +interface QueueRow { + readonly id: string; + readonly body: unknown; + readonly attempt: number | string; + readonly sequence: number | string; +} + +interface DecodedQueueRow { + readonly row: QueueRow; + readonly message: A; +} + +const defaultOptions = { + batchSize: 10, + maxRetries: 3, + retryDelayMillis: 1_000, + visibilityTimeoutMillis: 30_000, + pollIntervalMillis: 250, +} as const; + +const positiveInteger = (value: number | undefined, fallback: number): number => + value === undefined || !Number.isFinite(value) || value <= 0 ? fallback : Math.floor(value); + +const nonNegativeInteger = (value: number | undefined, fallback: number): number => + value === undefined || !Number.isFinite(value) || value < 0 ? fallback : Math.floor(value); + +const resolvedOptions = (options: QueueConsumerOptions | undefined) => ({ + batchSize: positiveInteger(options?.batchSize, defaultOptions.batchSize), + maxRetries: nonNegativeInteger(options?.maxRetries, defaultOptions.maxRetries), + retryDelayMillis: nonNegativeInteger(options?.retryDelayMillis, defaultOptions.retryDelayMillis), + visibilityTimeoutMillis: positiveInteger( + options?.visibilityTimeoutMillis, + defaultOptions.visibilityTimeoutMillis, + ), + pollIntervalMillis: positiveInteger( + options?.pollIntervalMillis, + defaultOptions.pollIntervalMillis, + ), + deadLetterQueue: options?.deadLetterQueue, +}); + +const encodeJson = (value: unknown): Effect.Effect => + Effect.try({ + try: () => { + const encoded = JSON.stringify(value); + if (encoded === undefined) throw new TypeError("Queue messages must be JSON-serializable"); + return encoded; + }, + catch: (cause) => cause, + }); + +const ensureTable = (sql: SqlClient.SqlClient) => + sql.withTransaction( + Effect.gen(function* () { + yield* sql`SELECT pg_advisory_xact_lock(hashtext('orbian_queue_schema_v1'))`; + yield* sql` + CREATE TABLE IF NOT EXISTS platform_queue_messages ( + id TEXT PRIMARY KEY, + queue_name TEXT NOT NULL, + body_json JSONB NOT NULL, + attempt INTEGER NOT NULL DEFAULT 0, + available_at_ms BIGINT NOT NULL, + leased_until_ms BIGINT, + created_at_ms BIGINT NOT NULL, + last_error TEXT, + sequence BIGSERIAL NOT NULL + ) + `; + yield* sql` + ALTER TABLE platform_queue_messages + ADD COLUMN IF NOT EXISTS sequence BIGSERIAL + `; + yield* sql` + CREATE INDEX IF NOT EXISTS platform_queue_messages_claim_v2_idx + ON platform_queue_messages (queue_name, available_at_ms, sequence) + `; + }), + ); + +const producerError = (queueName: string, cause: unknown) => + new QueueProducerError({ queueName, cause: String(cause) }); + +const consumerError = (queueName: string, cause: unknown) => + new QueueConsumerError({ queueName, cause: String(cause) }); + +const makeProducer = ( + sql: SqlClient.SqlClient, + queueName: string, + schema: Schema.Codec, +): QueueProducer => { + const encode = SchemaParser.encodeUnknownEffect(schema); + + const publishEncoded = (encoded: unknown) => + Effect.gen(function* () { + const body = yield* encodeJson(encoded); + const now = Date.now(); + yield* sql` + INSERT INTO platform_queue_messages ( + id, queue_name, body_json, available_at_ms, created_at_ms + ) VALUES ( + ${crypto.randomUUID()}, ${queueName}, ${body}::jsonb, ${now}, ${now} + ) + `; + }); + + const publish = (message: A) => + PlatformRuntime.pipe( + Effect.andThen(encode(message)), + Effect.flatMap(publishEncoded), + Effect.mapError((cause) => producerError(queueName, cause)), + ); + + const publishBatch = (messages: ReadonlyArray) => + PlatformRuntime.pipe( + Effect.andThen(Effect.forEach(messages, (message) => encode(message))), + Effect.flatMap((encoded) => + sql.withTransaction(Effect.forEach(encoded, publishEncoded, { discard: true })), + ), + Effect.mapError((cause) => producerError(queueName, cause)), + ); + + return { publish, publishBatch }; +}; + +const claimBatch = ( + sql: SqlClient.SqlClient, + queueName: string, + options: ReturnType, +) => + Effect.suspend(() => { + const now = Date.now(); + return sql` + WITH claimed AS ( + SELECT id + FROM platform_queue_messages + WHERE queue_name = ${queueName} + AND available_at_ms <= ${now} + AND (leased_until_ms IS NULL OR leased_until_ms <= ${now}) + ORDER BY sequence ASC + FOR UPDATE SKIP LOCKED + LIMIT ${options.batchSize} + ) + UPDATE platform_queue_messages AS message + SET attempt = message.attempt + 1, + leased_until_ms = ${now + options.visibilityTimeoutMillis} + FROM claimed + WHERE message.id = claimed.id + RETURNING message.id, message.body_json AS "body", message.attempt, message.sequence + `.pipe( + Effect.map((rows) => + [...rows].sort((left, right) => Number(left.sequence) - Number(right.sequence)), + ), + ); + }); + +const deleteRows = (sql: SqlClient.SqlClient, rows: ReadonlyArray) => + sql.withTransaction( + Effect.forEach( + rows, + (row) => + sql` + DELETE FROM platform_queue_messages + WHERE id = ${row.id} + `, + { discard: true }, + ), + ); + +const retryRows = ( + sql: SqlClient.SqlClient, + rows: ReadonlyArray, + queueName: string, + options: ReturnType, + cause: Cause.Cause, +) => { + const now = Date.now(); + const maxAttempts = options.maxRetries + 1; + const error = Cause.pretty(cause); + return sql + .withTransaction( + Effect.forEach( + rows, + (row) => { + if (Number(row.attempt) < maxAttempts) { + return sql` + UPDATE platform_queue_messages + SET available_at_ms = ${now + options.retryDelayMillis}, + leased_until_ms = NULL, + last_error = ${error} + WHERE id = ${row.id} + `; + } + if (options.deadLetterQueue) { + return Effect.gen(function* () { + const body = yield* encodeJson(row.body); + yield* sql` + INSERT INTO platform_queue_messages ( + id, queue_name, body_json, available_at_ms, created_at_ms, last_error + ) VALUES ( + ${crypto.randomUUID()}, ${options.deadLetterQueue}, ${body}::jsonb, + ${now}, ${now}, ${error} + ) + `; + yield* sql` + DELETE FROM platform_queue_messages + WHERE id = ${row.id} + `; + }); + } + return sql` + DELETE FROM platform_queue_messages + WHERE id = ${row.id} + `; + }, + { discard: true }, + ), + ) + .pipe( + Effect.tap(() => + Effect.logWarning("queue batch failed; delivery state updated", { + queueName, + messageCount: rows.length, + cause: error, + }), + ), + ); +}; + +const processBatch = ( + sql: SqlClient.SqlClient, + queueName: string, + schema: Schema.Codec, + handleBatch: (messages: ReadonlyArray) => Effect.Effect, + inputOptions: QueueConsumerOptions | undefined, +): Effect.Effect => { + const options = resolvedOptions(inputOptions); + const decode = SchemaParser.decodeUnknownEffect(schema); + + return PlatformRuntime.pipe( + Effect.andThen(claimBatch(sql, queueName, options)), + Effect.flatMap((rows) => + Effect.gen(function* () { + if (rows.length === 0) return 0; + + const decoded: Array> = []; + const poison: Array = []; + yield* Effect.forEach( + rows, + (row) => + decode(row.body).pipe( + Effect.matchCauseEffect({ + onFailure: (cause) => + Effect.logWarning("queue payload decode failed; acking poison message", { + queueName, + messageId: row.id, + cause: Cause.pretty(cause), + }).pipe( + Effect.tap(() => + Effect.sync(() => { + poison.push(row); + }), + ), + ), + onSuccess: (message) => + Effect.sync(() => { + decoded.push({ row, message }); + }), + }), + ), + { discard: true }, + ); + + if (poison.length > 0) yield* deleteRows(sql, poison); + if (decoded.length === 0) return rows.length; + + yield* handleBatch(decoded.map(({ message }) => message)).pipe( + Effect.matchCauseEffect({ + onFailure: (cause) => + retryRows( + sql, + decoded.map(({ row }) => row), + queueName, + options, + cause, + ), + onSuccess: () => + deleteRows( + sql, + decoded.map(({ row }) => row), + ), + }), + ); + return rows.length; + }), + ), + Effect.mapError((cause) => consumerError(queueName, cause)), + ); +}; + +const makeQueueDriver = (sql: SqlClient.SqlClient): QueueDriverShape => ({ + producer: (queueName, schema) => makeProducer(sql, queueName, schema), + processBatch: (queueName, schema, handleBatch, options) => + processBatch(sql, queueName, schema, handleBatch, options), + consumeBatch: (queueName, schema, handleBatch, options) => { + const resolved = resolvedOptions(options); + return Effect.forever( + processBatch(sql, queueName, schema, handleBatch, options).pipe( + Effect.flatMap((count) => + count === 0 ? Effect.sleep(Duration.millis(resolved.pollIntervalMillis)) : Effect.void, + ), + ), + ); + }, +}); + +/** Postgres-backed queue driver with durable leases, retries, and dead letters. */ +export const PgQueueLive = (config: PgPlatformConfig): Layer.Layer => + Layer.effect( + QueueDriver, + Effect.gen(function* () { + const sql = yield* SqlClient.SqlClient; + yield* ensureTable(sql); + return makeQueueDriver(sql); + }), + ).pipe(Layer.provide(PgPlatformClientLive(config)), Layer.orDie); diff --git a/selfhost/entry/src/runtime/Screenshot.ts b/selfhost/entry/src/runtime/Screenshot.ts new file mode 100644 index 000000000..32650161e --- /dev/null +++ b/selfhost/entry/src/runtime/Screenshot.ts @@ -0,0 +1,143 @@ +import { + Screenshot, + ScreenshotError, + type ScreenshotOptions, + type ScreenshotShape, +} from "@orbian/sdk/Screenshot"; +import { PlatformRuntime } from "@orbian/sdk/PlatformRuntime"; +import { Effect, Layer } from "effect"; +import { chromium, type Browser } from "playwright-core"; + +/** Headless Chromium launch and resource limits. */ +export interface ChromiumScreenshotConfig { + readonly executablePath?: string; + readonly disableSandbox?: boolean; + readonly timeoutMillis?: number; + readonly maxWidth?: number; + readonly maxHeight?: number; + readonly maxDeviceScaleFactor?: number; + readonly maxHtmlBytes?: number; + readonly maxRenderedPixels?: number; +} + +const screenshotError = (operation: string, cause: unknown) => + new ScreenshotError({ operation, cause: String(cause) }); + +const positiveInteger = (value: number, maximum: number): boolean => + Number.isInteger(value) && value > 0 && value <= maximum; + +/** Validates screenshot memory and viewport budgets before Chromium is invoked. */ +export const validateChromiumScreenshotOptions = ( + options: ScreenshotOptions, + config: ChromiumScreenshotConfig, +): Effect.Effect => { + const maxWidth = config.maxWidth ?? 4_096; + const maxHeight = config.maxHeight ?? 4_096; + const maxScale = config.maxDeviceScaleFactor ?? 4; + const maxHtmlBytes = config.maxHtmlBytes ?? 4 * 1_024 * 1_024; + const maxRenderedPixels = config.maxRenderedPixels ?? 16_777_216; + if (!positiveInteger(options.width, maxWidth)) { + return Effect.fail(screenshotError("validate", `width must be between 1 and ${maxWidth}`)); + } + if (!positiveInteger(options.height, maxHeight)) { + return Effect.fail(screenshotError("validate", `height must be between 1 and ${maxHeight}`)); + } + if ( + !Number.isFinite(options.deviceScaleFactor) || + options.deviceScaleFactor < 1 || + options.deviceScaleFactor > maxScale + ) { + return Effect.fail( + screenshotError("validate", `deviceScaleFactor must be between 1 and ${maxScale}`), + ); + } + if (new TextEncoder().encode(options.html).byteLength > maxHtmlBytes) { + return Effect.fail(screenshotError("validate", `html must be at most ${maxHtmlBytes} bytes`)); + } + const renderedPixels = + options.width * options.height * options.deviceScaleFactor * options.deviceScaleFactor; + if (renderedPixels > maxRenderedPixels) { + return Effect.fail( + screenshotError("validate", `rendered image must be at most ${maxRenderedPixels} pixels`), + ); + } + return Effect.void; +}; + +const render = (browser: Browser, config: ChromiumScreenshotConfig, options: ScreenshotOptions) => + validateChromiumScreenshotOptions(options, config).pipe( + Effect.andThen( + Effect.acquireUseRelease( + Effect.tryPromise({ + try: () => + browser.newContext({ + viewport: { width: Math.floor(options.width), height: Math.floor(options.height) }, + deviceScaleFactor: options.deviceScaleFactor, + javaScriptEnabled: false, + serviceWorkers: "block", + }), + catch: (cause) => screenshotError("openContext", cause), + }), + (context) => + Effect.tryPromise({ + try: async () => { + await context.setOffline(true); + const page = await context.newPage(); + await page.route("**/*", (route) => route.abort("blockedbyclient")); + await page.setContent(options.html, { + waitUntil: "load", + timeout: config.timeoutMillis ?? 15_000, + }); + return new Uint8Array( + await page.screenshot({ + type: "png", + fullPage: false, + animations: "disabled", + timeout: config.timeoutMillis ?? 15_000, + }), + ); + }, + catch: (cause) => screenshotError("render", cause), + }), + (context) => + Effect.promise(() => context.close()).pipe( + Effect.catchCause((cause) => + Effect.logWarning("failed to close screenshot browser context", { + cause: String(cause), + }), + ), + ), + ), + ), + ); + +const makeRenderer = (browser: Browser, config: ChromiumScreenshotConfig): ScreenshotShape => ({ + renderPng: (options) => PlatformRuntime.pipe(Effect.andThen(render(browser, config, options))), +}); + +/** Chromium-backed PNG screenshot layer with network and JavaScript disabled. */ +export const ChromiumScreenshotLive = ( + config: ChromiumScreenshotConfig = {}, +): Layer.Layer => + Layer.effect( + Screenshot, + Effect.acquireRelease( + Effect.tryPromise({ + try: () => + chromium.launch({ + executablePath: config.executablePath ?? process.env.CHROMIUM_EXECUTABLE_PATH, + headless: true, + args: ["--disable-dev-shm-usage", ...(config.disableSandbox ? ["--no-sandbox"] : [])], + }), + catch: (cause) => screenshotError("launch", cause), + }), + (browser) => + Effect.promise(() => browser.close()).pipe( + Effect.catchCause((cause) => + Effect.logWarning("failed to close screenshot browser", { + cause: String(cause), + }), + ), + ), + ).pipe(Effect.map((browser) => makeRenderer(browser, config))), + ); diff --git a/selfhost/entry/src/runtime/Workflow.ts b/selfhost/entry/src/runtime/Workflow.ts new file mode 100644 index 000000000..1cb9ef6cb --- /dev/null +++ b/selfhost/entry/src/runtime/Workflow.ts @@ -0,0 +1,214 @@ +import { + type WorkflowDefinition, + type WorkflowExecutionResult, + type WorkflowHandlerContext, + WorkflowRunner, + WorkflowRunnerError, + type WorkflowRunnerShape, + type WorkflowStepOptions, +} from "@orbian/sdk/Workflow"; +import { PlatformRuntime } from "@orbian/sdk/PlatformRuntime"; +import { Cause, Effect, Exit, Layer, Option, Schema } from "effect"; +import { Activity, DurableClock, Workflow, WorkflowEngine } from "effect/unstable/workflow"; + +import { PgWorkflowEngineLive } from "./PgWorkflowEngine.js"; +import type { PgPlatformConfig } from "./Postgres.js"; + +const runnerError = (workflowName: string, operation: string, cause: unknown) => + new WorkflowRunnerError({ workflowName, operation, cause: String(cause) }); + +const catchRunnerCause = ( + effect: Effect.Effect, + workflowName: string, + operation: string, +): Effect.Effect => + effect.pipe( + Effect.catchCause((cause) => { + const squashed = Cause.squash(cause); + return Effect.fail( + squashed instanceof WorkflowRunnerError + ? squashed + : runnerError(workflowName, operation, Cause.pretty(cause)), + ); + }), + ); + +const toNativeWorkflow = < + const Name extends string, + Payload extends Schema.Struct.Fields, + Success extends Schema.Top, +>( + workflow: WorkflowDefinition, +) => + Workflow.make(workflow.name, { + payload: workflow.payload, + success: workflow.success, + error: WorkflowRunnerError, + idempotencyKey: workflow.idempotencyKey, + }); + +const executionResult = ( + result: Workflow.Result, +): Effect.Effect> => { + if (result._tag === "Suspended") { + return Effect.succeed({ status: "suspended" }); + } + return Effect.succeed( + Exit.match(result.exit, { + onFailure: (cause) => { + if (Cause.hasInterrupts(cause)) { + return { status: "interrupted" as const }; + } + const error = Cause.squash(cause); + return { + status: "failed" as const, + error: + error instanceof WorkflowRunnerError ? error : runnerError("unknown", "poll", error), + }; + }, + onSuccess: (value) => ({ status: "succeeded" as const, value }), + }), + ); +}; + +const makeStep = ( + workflowName: string, + options: WorkflowStepOptions, +): Effect.Effect => + Activity.make({ + name: options.name, + success: options.success, + error: WorkflowRunnerError, + execute: catchRunnerCause( + PlatformRuntime.pipe(Effect.andThen(options.execute)), + workflowName, + `step:${options.name}`, + ), + }) as unknown as Effect.Effect; + +const sleepUntil = ( + workflowName: string, + name: string, + scheduledTime: Date, +): Effect.Effect => + catchRunnerCause( + PlatformRuntime.pipe( + Effect.andThen( + Effect.suspend(() => { + const delay = scheduledTime.getTime() - Date.now(); + return delay <= 0 + ? Effect.void + : DurableClock.sleep({ + name, + duration: delay, + inMemoryThreshold: 1, + }); + }), + ), + ), + workflowName, + `sleep:${name}`, + ) as unknown as Effect.Effect; + +type AnyWorkflowDefinition = WorkflowDefinition; + +type AnyWorkflowHandler = ( + payload: Schema.Struct.Type, + context: WorkflowHandlerContext, +) => Effect.Effect; + +const makeRunner = (engine: WorkflowEngine.WorkflowEngine["Service"]): WorkflowRunnerShape => + ({ + register: (workflow: AnyWorkflowDefinition, handler: AnyWorkflowHandler) => { + const native = toNativeWorkflow(workflow); + return catchRunnerCause( + engine.register(native, (payload, executionId) => { + const context: WorkflowHandlerContext = { + executionId, + step: (options) => makeStep(workflow.name, options), + sleepUntil: (name, scheduledTime) => sleepUntil(workflow.name, name, scheduledTime), + }; + return catchRunnerCause( + PlatformRuntime.pipe(Effect.andThen(handler(payload, context))), + workflow.name, + "run", + ); + }), + workflow.name, + "register", + ); + }, + dispatch: ( + workflow: AnyWorkflowDefinition, + payload: Schema.Struct.Type, + ) => { + const native = toNativeWorkflow(workflow); + return catchRunnerCause( + PlatformRuntime.pipe( + Effect.andThen(native.execute(payload, { discard: true })), + Effect.provideService(WorkflowEngine.WorkflowEngine, engine), + ), + workflow.name, + "dispatch", + ); + }, + execute: ( + workflow: AnyWorkflowDefinition, + payload: Schema.Struct.Type, + ) => { + const native = toNativeWorkflow(workflow); + return catchRunnerCause( + PlatformRuntime.pipe( + Effect.andThen(native.execute(payload)), + Effect.provideService(WorkflowEngine.WorkflowEngine, engine), + ), + workflow.name, + "execute", + ); + }, + poll: (workflow: AnyWorkflowDefinition, executionId: string) => { + const native = toNativeWorkflow(workflow); + return catchRunnerCause( + PlatformRuntime.pipe( + Effect.andThen(native.poll(executionId)), + Effect.provideService(WorkflowEngine.WorkflowEngine, engine), + Effect.flatMap( + Option.match({ + onNone: () => Effect.succeedNone, + onSome: (result) => executionResult(result).pipe(Effect.map(Option.some)), + }), + ), + ), + workflow.name, + "poll", + ); + }, + resume: (workflow: AnyWorkflowDefinition, executionId: string) => { + const native = toNativeWorkflow(workflow); + return catchRunnerCause( + PlatformRuntime.pipe( + Effect.andThen(native.resume(executionId)), + Effect.provideService(WorkflowEngine.WorkflowEngine, engine), + ), + workflow.name, + "resume", + ); + }, + interrupt: (workflow: AnyWorkflowDefinition, executionId: string) => { + const native = toNativeWorkflow(workflow); + return catchRunnerCause( + PlatformRuntime.pipe( + Effect.andThen(native.interrupt(executionId)), + Effect.provideService(WorkflowEngine.WorkflowEngine, engine), + ), + workflow.name, + "interrupt", + ); + }, + }) as unknown as WorkflowRunnerShape; + +/** Postgres-backed provider-neutral workflow runner. */ +export const PgWorkflowRunnerLive = (config: PgPlatformConfig): Layer.Layer => + Layer.effect(WorkflowRunner, Effect.map(WorkflowEngine.WorkflowEngine, makeRunner)).pipe( + Layer.provide(PgWorkflowEngineLive(config)), + ); diff --git a/selfhost/entry/tests/AgentNodeWebSocket.integration.test.ts b/selfhost/entry/tests/AgentNodeWebSocket.integration.test.ts index 454a884f5..5ff1b7c64 100644 --- a/selfhost/entry/tests/AgentNodeWebSocket.integration.test.ts +++ b/selfhost/entry/tests/AgentNodeWebSocket.integration.test.ts @@ -8,7 +8,7 @@ import { Workos, } from "@voidhash/core/services"; import { Db } from "@voidhash/db"; -import { makeMemoryDurableEntityHost } from "@orbian/node/MemoryDurableEntity"; +import { makeMemoryDurableEntityHost } from "../src/runtime/MemoryDurableEntity.ts"; import { Context, Effect, Redacted } from "effect"; import { WebSocket } from "ws"; import { afterEach, describe, expect, it } from "vite-plus/test"; diff --git a/selfhost/entry/tests/MimicDocumentIdle.test.ts b/selfhost/entry/tests/MimicDocumentIdle.test.ts index e11b14d6f..ba2eb7642 100644 --- a/selfhost/entry/tests/MimicDocumentIdle.test.ts +++ b/selfhost/entry/tests/MimicDocumentIdle.test.ts @@ -4,8 +4,8 @@ import { type MimicDocumentIdleMessageType, } from "@voidhash/mimic-db/ws/idle-notify"; import { makeDurableEntityAddress } from "@orbian/sdk/DurableEntity"; -import type { NodeDurableEntityControlShape } from "@orbian/node/DurableEntity"; -import { makeMemoryDurableEntityHost } from "@orbian/node/MemoryDurableEntity"; +import type { NodeDurableEntityControlShape } from "../src/runtime/DurableEntity.ts"; +import { makeMemoryDurableEntityHost } from "../src/runtime/MemoryDurableEntity.ts"; import { Effect } from "effect"; import { describe, expect, it, vi } from "vitest"; diff --git a/selfhost/entry/tests/MimicNode.integration.test.ts b/selfhost/entry/tests/MimicNode.integration.test.ts index 06d7f7f5b..86da1b6e4 100644 --- a/selfhost/entry/tests/MimicNode.integration.test.ts +++ b/selfhost/entry/tests/MimicNode.integration.test.ts @@ -7,7 +7,7 @@ import { DurableEntityHost, makeDurableEntityAddress, } from "@orbian/sdk/DurableEntity"; -import { NodeDurableEntityControl } from "@orbian/node/DurableEntity"; +import { NodeDurableEntityControl } from "../src/runtime/DurableEntity.ts"; import { Effect, ManagedRuntime, Redacted } from "effect"; import { describe, expect, it } from "vitest"; import WebSocket from "ws";