Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 6 additions & 7 deletions .github/workflows/publish.yml
Original file line number Diff line number Diff line change
Expand Up @@ -34,22 +34,20 @@ jobs:
- uses: actions/setup-node@v4
with:
node-version: 24
cache: pnpm
cache-dependency-path: pnpm-lock.yaml
registry-url: https://registry.npmjs.org
- name: Install OIDC-capable npm
run: npm install --global npm@11.16.0
- run: pnpm install --frozen-lockfile
- id: release
run: >-
node scripts/resolve-release.mjs
--version=${{ inputs.version }}
--tag=${{ inputs.dist_tag }}
--branch=${{ github.ref_name }}
- name: Verify npm authentication and unused version
env:
NODE_AUTH_TOKEN: ${{ secrets.NPM_TOKEN }}
- name: Verify npm registry and unused version
run: |
set -euo pipefail
if [ -n "${NODE_AUTH_TOKEN:-}" ]; then npm whoami; else npm ping; fi
npm ping
if npm view "@rivet-dev/workflows@${{ steps.release.outputs.version }}" version >/dev/null 2>&1; then
echo "@rivet-dev/workflows@${{ steps.release.outputs.version }} already exists" >&2
exit 1
Expand All @@ -63,7 +61,8 @@ jobs:
npm pack ./packages/workflows --silent --pack-destination .pack
- name: Publish
env:
NODE_AUTH_TOKEN: ${{ secrets.NPM_TOKEN }}
# setup-node provides a dummy token; clear it so npm uses OIDC.
NODE_AUTH_TOKEN: ""
run: >-
npm publish
.pack/rivet-dev-workflows-${{ steps.release.outputs.version }}.tgz
Expand Down
7 changes: 3 additions & 4 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -9,15 +9,14 @@ pnpm add @rivet-dev/workflows rivetkit
```

```ts
import { actor } from "rivetkit";
import { workflow } from "@rivet-dev/workflows";

export const report = actor({
run: workflow(async (ctx) => {
export const report = workflow({
run: async (ctx) => {
await ctx.step("generate", async (step) => {
step.log.info("generating report");
});
}),
},
});
```

Expand Down
7 changes: 3 additions & 4 deletions packages/workflows/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -5,13 +5,12 @@ Durable, replayable workflows for Rivet Actors.
[Documentation](https://rivet.dev/workflows/docs)

```ts
import { actor } from "rivetkit";
import { workflow } from "@rivet-dev/workflows";

export const example = actor({
run: workflow(async (ctx) => {
export const example = workflow({
run: async (ctx) => {
await ctx.step("hello", async () => "world");
}),
},
});
```

Expand Down
7 changes: 6 additions & 1 deletion packages/workflows/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,11 @@
"version": "2.3.10",
"description": "Durable, replayable workflows for Rivet Actors",
"license": "Apache-2.0",
"repository": {
"type": "git",
"url": "https://github.com/rivet-dev/workflows.git",
"directory": "packages/workflows"
},
"keywords": [
"rivet",
"workflow",
Expand Down Expand Up @@ -67,7 +72,7 @@
"commander": "^12.0.0",
"legacy-rivetkit": "npm:rivetkit@2.3.7",
"legacy-workflow-engine": "npm:@rivetkit/workflow-engine@2.3.7",
"rivetkit": "0.0.0-feat-workflows-public-host-apis.0ff6164",
"rivetkit": "0.0.0-feat-workflows-public-host-apis.1550fe4",
"tsup": "^8.4.0",
"tsx": "^4.7.0",
"typescript": "^5.7.3",
Expand Down
194 changes: 178 additions & 16 deletions packages/workflows/src/rivetkit/driver.ts
Original file line number Diff line number Diff line change
@@ -1,8 +1,5 @@
import type { ActorQueue, ActorRun, RunContext } from "rivetkit";
import {
WORKFLOW_STORAGE_V1,
type WorkflowStorageHandle,
} from "rivetkit/storage";
import type { RawAccess } from "rivetkit/db";
import type {
EngineDriver,
KVEntry,
Expand All @@ -12,6 +9,14 @@ import type {
WorkflowMessageIdentity,
} from "../index.js";

const WORKFLOW_STORAGE_PREFIX = new Uint8Array([6, 1]);
const WORKFLOW_UPSERT_SQL =
"INSERT INTO _rivet_wf_kv (key, value) VALUES (?, ?) ON CONFLICT(key) DO UPDATE SET value = excluded.value";

const WORKFLOW_SQLITE_MAX_VALUE_BYTES = 256 * 1024;
const WORKFLOW_SQLITE_MAX_BATCH_ROWS = 128;
const WORKFLOW_SQLITE_MAX_BATCH_BYTES = 512 * 1024;

function track<T>(
runCtx: RunContext<any, any, any, any, any, any, any, any>,
promise: Promise<T>,
Expand All @@ -25,6 +30,169 @@ function track<T>(
return promise;
}

function prefixWorkflowKey(key: Uint8Array): Uint8Array {
const prefixed = new Uint8Array(
WORKFLOW_STORAGE_PREFIX.byteLength + key.byteLength,
);
prefixed.set(WORKFLOW_STORAGE_PREFIX);
prefixed.set(key, WORKFLOW_STORAGE_PREFIX.byteLength);
return prefixed;
}

function stripWorkflowKey(key: Uint8Array): Uint8Array {
if (
key.byteLength < WORKFLOW_STORAGE_PREFIX.byteLength ||
!WORKFLOW_STORAGE_PREFIX.every((byte, index) => key[index] === byte)
) {
throw new Error("workflow SQLite key escaped the [6, 1] namespace");
}
return key.slice(WORKFLOW_STORAGE_PREFIX.byteLength);
}

function computeUpperBound(prefix: Uint8Array): Uint8Array {
const upperBound = prefix.slice();
for (let index = upperBound.length - 1; index >= 0; index--) {
if (upperBound[index] !== 0xff) {
upperBound[index]++;
return upperBound.slice(0, index + 1);
}
}

// Every workflow key begins with 6, so a finite upper bound always exists.
throw new Error("workflow storage prefix has no upper bound");
}

function normalizeSqlBlob(value: unknown): Uint8Array {
if (value instanceof Uint8Array) {
return value;
}
if (value instanceof ArrayBuffer) {
return new Uint8Array(value);
}
if (ArrayBuffer.isView(value)) {
return new Uint8Array(value.buffer, value.byteOffset, value.byteLength);
}
if (Array.isArray(value)) {
const bytes = new Uint8Array(value.length);
for (const [index, byte] of value.entries()) {
if (!Number.isInteger(byte) || byte < 0 || byte > 255) {
throw new Error("workflow SQLite value was not a byte array");
}
bytes[index] = byte;
}
return bytes;
}
throw new Error("workflow SQLite value was not a blob");
}

function validateWrites(writes: KVWrite[]): void {
if (writes.length > WORKFLOW_SQLITE_MAX_BATCH_ROWS) {
throw new Error(
`Workflow batch contains ${writes.length} rows, exceeding the ${WORKFLOW_SQLITE_MAX_BATCH_ROWS} row limit`,
);
}

let batchBytes = 0;
for (const write of writes) {
if (write.value.byteLength > WORKFLOW_SQLITE_MAX_VALUE_BYTES) {
throw new Error(
`Workflow value is ${write.value.byteLength} bytes, exceeding the ${WORKFLOW_SQLITE_MAX_VALUE_BYTES} byte limit`,
);
}
batchBytes +=
WORKFLOW_STORAGE_PREFIX.byteLength +
write.key.byteLength +
write.value.byteLength;
}

if (batchBytes > WORKFLOW_SQLITE_MAX_BATCH_BYTES) {
throw new Error(
`Workflow batch is ${batchBytes} bytes, exceeding the ${WORKFLOW_SQLITE_MAX_BATCH_BYTES} byte limit`,
);
}
}

class WorkflowStorage {
#db: RawAccess;

constructor(db: RawAccess) {
this.#db = db;
}

async get(key: Uint8Array): Promise<Uint8Array | null> {
const rows = await this.#db.execute<{ value: unknown }>(
"SELECT value FROM _rivet_wf_kv WHERE key = ?",
prefixWorkflowKey(key),
);
const value = rows[0]?.value;
return value == null ? null : normalizeSqlBlob(value);
}

async set(key: Uint8Array, value: Uint8Array): Promise<void> {
await this.batch([{ key, value }], false);
}

async delete(key: Uint8Array): Promise<void> {
await this.#db.execute(
"DELETE FROM _rivet_wf_kv WHERE key = ?",
prefixWorkflowKey(key),
);
}

async deletePrefix(prefix: Uint8Array): Promise<void> {
const start = prefixWorkflowKey(prefix);
await this.#db.execute(
"DELETE FROM _rivet_wf_kv WHERE key >= ? AND key < ?",
start,
computeUpperBound(start),
);
}

async deleteRange(start: Uint8Array, end: Uint8Array): Promise<void> {
await this.#db.execute(
"DELETE FROM _rivet_wf_kv WHERE key >= ? AND key < ?",
prefixWorkflowKey(start),
prefixWorkflowKey(end),
);
}

async list(prefix: Uint8Array): Promise<KVEntry[]> {
const start = prefixWorkflowKey(prefix);
const rows = await this.#db.execute<{ key: unknown; value: unknown }>(
"SELECT key, value FROM _rivet_wf_kv WHERE key >= ? AND key < ? ORDER BY key ASC",
start,
computeUpperBound(start),
);
return rows.map((row) => ({
key: stripWorkflowKey(normalizeSqlBlob(row.key)),
value: normalizeSqlBlob(row.value),
}));
}

async batch(writes: KVWrite[], includeState: boolean): Promise<void> {
if (writes.length === 0) return;
validateWrites(writes);

const commit = async (tx: RawAccess) => {
for (const write of writes) {
await tx.execute(
WORKFLOW_UPSERT_SQL,
prefixWorkflowKey(write.key),
write.value,
);
}
};

if (includeState) {
await this.#db.transaction(commit, {
experimental: { includeState: true },
});
} else {
await this.#db.transaction(commit);
}
}
}

class ActorWorkflowMessageDriver implements WorkflowMessageDriver {
#runCtx: RunContext<any, any, any, any, any, any, any, any>;
#queue: ActorQueue;
Expand Down Expand Up @@ -95,15 +263,15 @@ export class ActorWorkflowDriver implements EngineDriver {
readonly workerPollInterval = 100;
readonly messageDriver: WorkflowMessageDriver;
#runCtx: RunContext<any, any, any, any, any, any, any, any>;
#storage: WorkflowStorageHandle;
#storage: WorkflowStorage;
#queue: ActorQueue;
#run: ActorRun;

constructor(runCtx: RunContext<any, any, any, any, any, any, any, any>) {
this.#runCtx = runCtx;
this.messageDriver = new ActorWorkflowMessageDriver(runCtx);
this.#queue = runCtx.queue;
this.#storage = runCtx.storage.open(WORKFLOW_STORAGE_V1);
this.#storage = new WorkflowStorage(runCtx.db);
this.#run = runCtx.run;
}

Expand Down Expand Up @@ -132,9 +300,7 @@ export class ActorWorkflowDriver implements EngineDriver {
}

async batch(writes: KVWrite[]): Promise<void> {
if (writes.length === 0) return;

await track(this.#runCtx, this.#storage.flushWithState(writes));
await track(this.#runCtx, this.#storage.batch(writes, true));
}

async setAlarm(_workflowId: string, wakeAt: number): Promise<void> {
Expand Down Expand Up @@ -184,11 +350,11 @@ export class ActorWorkflowControlDriver implements EngineDriver {
readonly workerPollInterval = 100;
readonly messageDriver: WorkflowMessageDriver =
new NoopWorkflowMessageDriver();
#storage: WorkflowStorageHandle;
#storage: WorkflowStorage;
#run: ActorRun;

constructor(runCtx: RunContext<any, any, any, any, any, any, any, any>) {
this.#storage = runCtx.storage.open(WORKFLOW_STORAGE_V1);
this.#storage = new WorkflowStorage(runCtx.db);
this.#run = runCtx.run;
}

Expand Down Expand Up @@ -217,11 +383,7 @@ export class ActorWorkflowControlDriver implements EngineDriver {
}

async batch(writes: KVWrite[]): Promise<void> {
if (writes.length === 0) {
return;
}

await this.#storage.batch(writes);
await this.#storage.batch(writes, false);
}

async setAlarm(_workflowId: string, wakeAt: number): Promise<void> {
Expand Down
4 changes: 2 additions & 2 deletions packages/workflows/src/rivetkit/inspector.ts
Original file line number Diff line number Diff line change
@@ -1,9 +1,9 @@
import * as transport from "rivetkit/inspector/workflow";
import * as transport from "rivetkit/experimental/inspector/workflow";
import {
encodeWorkflowHistoryTransport,
encodeWorkflowInspectorValue,
type WorkflowInspectorAdapter,
} from "rivetkit/inspector/workflow";
} from "rivetkit/experimental/inspector/workflow";
import type {
BranchStatus,
BranchStatusType,
Expand Down
Loading
Loading