Skip to content
Open
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
1 change: 1 addition & 0 deletions rivetkit-typescript/packages/rivetkit-napi/index.d.ts
Original file line number Diff line number Diff line change
Expand Up @@ -372,6 +372,7 @@ export declare class ActorContext {
sleep(): void
destroy(): void
destroyRequested(): boolean
generation(): number | null
waitForDestroyCompletion(): Promise<void>
setPreventSleep(preventSleep: boolean): void
preventSleep(): boolean
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -42,7 +42,7 @@
ThreadsafeFunction<DisconnectPredicatePayload, ErrorStrategy::CalleeHandled>;
type RunRestartHook = Arc<dyn Fn(Option<(i64, u64)>) -> anyhow::Result<()> + Send + Sync + 'static>;
pub(crate) type RegisteredTask = Pin<Box<dyn Future<Output = ()> + Send + 'static>>;

Check warning on line 45 in rivetkit-typescript/packages/rivetkit-napi/src/actor_context.rs

View workflow job for this annotation

GitHub Actions / Rustfmt

Diff in /home/runner/work/rivet/rivet/rivetkit-typescript/packages/rivetkit-napi/src/actor_context.rs
// Keyed by generation so a newer generation of the same actor never shares, resets, or clears the
// JS runtime state of an older generation that is still shutting down on this host.
static ACTOR_CONTEXT_SHARED: LazyLock<
Expand Down Expand Up @@ -607,6 +607,11 @@
self.inner.is_destroy_requested()
}

#[napi]
pub fn generation(&self) -> Option<u32> {
self.inner.generation()
}

#[napi]
pub async fn wait_for_destroy_completion(&self) {
self.inner.wait_for_destroy_completion_public().await;
Expand Down
5 changes: 5 additions & 0 deletions rivetkit-typescript/packages/rivetkit-wasm/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1281,6 +1281,11 @@ impl WasmActorContext {
self.inner.actor_id().to_owned()
}

#[wasm_bindgen(js_name = generation)]
pub fn generation(&self) -> Option<u32> {
self.inner.generation()
}

#[wasm_bindgen]
pub fn name(&self) -> String {
self.inner.name().to_owned()
Expand Down
2 changes: 2 additions & 0 deletions rivetkit-typescript/packages/rivetkit/src/actor/config.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1079,6 +1079,8 @@ export interface RunControl {
/** @experimental */
export interface RunInspectorFactoryContext {
actorId: string;
/** Actor generation the inspector belongs to, when the runtime reports one. */
actorGeneration?: number;
control: RunControl;
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -600,6 +600,10 @@ export class NapiCoreRuntime implements CoreRuntime {
return asNativeActorContext(ctx).actorId();
}

actorGeneration(ctx: ActorContextHandle): number | undefined {
return asNativeActorContext(ctx).generation() ?? undefined;
}

runWithActorInvocationContext<T>(ctx: ActorContextHandle, run: () => T): T {
return this.#runAs(asNativeActorContext(ctx), run);
}
Expand Down
34 changes: 29 additions & 5 deletions rivetkit-typescript/packages/rivetkit/src/registry/native.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2963,6 +2963,10 @@ export class ActorContextHandleAdapter {
return callNativeSync(() => this.#runtime.actorId(this.#ctx));
}

get actorGeneration(): number | undefined {
return callNativeSync(() => this.#runtime.actorGeneration(this.#ctx));
}

get name(): string {
return callNativeSync(() => this.#runtime.actorName(this.#ctx));
}
Expand Down Expand Up @@ -3658,6 +3662,14 @@ class NativeWorkflowRuntimeAdapter {
#onMessagesReceived: () => void;
#completions = new Map<string, (response?: unknown) => Promise<void>>();

/**
* Actor generation, used to keep workflow registrations of different generations apart. Read
* lazily because most contexts never run a workflow.
*/
get generation(): number | undefined {
return this.#ctx.actorGeneration;
}

readonly id: string;
readonly driver: {
kvBatchGet: (
Expand Down Expand Up @@ -4104,9 +4116,14 @@ export function buildNativeFactory(
const actorId = callNativeSync(() => runtime.actorId(ctx));
const restart = () =>
callNativeSync(() => runtime.actorRestartRunHandler(ctx));
return (runHandlerCoordinator?.getInspector(actorId, restart)
?.workflow ??
getRunInspectorConfig(config.run, actorId)?.workflow) as
const actorGeneration = callNativeSync(() =>
runtime.actorGeneration(ctx),
);
return (runHandlerCoordinator?.getInspector(
actorId,
restart,
actorGeneration,
)?.workflow ?? getRunInspectorConfig(config.run, actorId)?.workflow) as
| NativeWorkflowInspectorConfig
| undefined;
};
Expand Down Expand Up @@ -4824,7 +4841,10 @@ export function buildNativeFactory(
saveActorState,
);
} finally {
runHandlerCoordinator?.destroy(actorId);
runHandlerCoordinator?.destroy(
actorId,
callNativeSync(() => runtime.actorGeneration(ctx)),
);
disposeRunInspector(config.run, actorId);
await actorCtx.dispose();
}
Expand All @@ -4838,7 +4858,10 @@ export function buildNativeFactory(
const actorId = callNativeSync(() => runtime.actorId(ctx));
// Close run control before user cleanup so replay cannot race actor
// destruction. Recreating this actor id receives a fresh controller.
runHandlerCoordinator?.destroy(actorId);
runHandlerCoordinator?.destroy(
actorId,
callNativeSync(() => runtime.actorGeneration(ctx)),
);
disposeRunInspector(config.run, actorId);
try {
if (typeof config.onDestroy === "function") {
Expand Down Expand Up @@ -5349,6 +5372,7 @@ export function buildNativeFactory(
runtime.actorRestartRunHandler(ctx),
),
executeRun,
callNativeSync(() => runtime.actorGeneration(ctx)),
);
// Only a start suppressed behind a failed replay must preserve a
// consumed durable wake. Closed actors must remain closed.
Expand Down
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
import { describe, expect, test, vi } from "vitest";
import { defineRunHandler, type RunControl } from "@/actor/config";
import { workflow } from "@/workflow/mod";
import { RunHandlerCoordinator } from "./run-handler-coordinator";

class Deferred<T = void> {
Expand Down Expand Up @@ -51,6 +52,107 @@ function createCoordinator() {
}

describe("RunHandlerCoordinator", () => {
test("workflow inspector registrations stay per generation", () => {
const coordinator = new RunHandlerCoordinator(workflow(async () => {}));
const restart = vi.fn();

const oldInspector = coordinator.getInspector(
"actor-1",
restart,
1,
)?.workflow;
const newInspector = coordinator.getInspector(
"actor-1",
restart,
2,
)?.workflow;
expect(oldInspector).toBeDefined();
expect(newInspector).toBeDefined();
expect(newInspector).not.toBe(oldInspector);

// A superseded generation cannot re-initialize or rebind the current generation's
// workflow registration.
expect(coordinator.getInspector("actor-1", restart, 1)).toBeUndefined();
expect(coordinator.getInspector("actor-1", restart, 2)?.workflow).toBe(
newInspector,
);
});

test("a superseded generation's run is rejected before it touches the inspector", async () => {
const coordinator = new RunHandlerCoordinator(workflow(async () => {}));
const restart = vi.fn();
const newInspector = coordinator.getInspector(
"actor-1",
restart,
2,
)?.workflow;

await expect(
coordinator.run("actor-1", restart, async () => {}, 1),
).resolves.toBe("closed");
expect(coordinator.getInspector("actor-1", restart, 2)?.workflow).toBe(
newInspector,
);
});

test("a newer generation does not wait on an older generation's stuck run", async () => {
const coordinator = new RunHandlerCoordinator(
defineRunHandler(async () => {}, {
inspectorKind: "workflow",
createInspector: () => ({
inspector: {
workflow: {
getHistory: () => null,
getState: async () => null,
onHistoryUpdated: () => () => {},
replayFromStep: async () => null,
},
},
}),
}),
);
const oldRun = new Deferred();
const oldOutcome = coordinator.run(
"actor-1",
() => {},
() => oldRun.promise,
1,
);

// The lost generation's JS callback cannot be cancelled, so it stays active.
await expect(
coordinator.run(
"actor-1",
() => {},
async () => {},
2,
),
).resolves.toBe("ran");

// The superseded generation cannot start new runs, and its late destroy leaves the
// newer generation's state in place.
await expect(
coordinator.run(
"actor-1",
() => {},
async () => {},
1,
),
).resolves.toBe("closed");
coordinator.destroy("actor-1", 1);
await expect(
coordinator.run(
"actor-1",
() => {},
async () => {},
2,
),
).resolves.toBe("ran");

oldRun.resolve();
await expect(oldOutcome).resolves.toBe("ran");
});

test("fails explicitly when a JavaScript factory omits required workflow controls", async () => {
const coordinator = new RunHandlerCoordinator(
defineRunHandler(async () => {}, {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,8 @@ import { RivetError } from "@/actor/errors";
type RunHandler = ((...args: any[]) => any) | { run: (...args: any[]) => any };

interface ActorRunState {
/** Actor generation that owns this state, when the runtime reports one. */
actorGeneration: number | undefined;
active: boolean;
closed: boolean;
exclusive: boolean;
Expand Down Expand Up @@ -46,6 +48,26 @@ function notifyStateChanged(state: ActorRunState): void {
}
}

function createRunState(actorGeneration: number | undefined): ActorRunState {
return {
actorGeneration,
active: false,
closed: false,
exclusive: false,
exclusiveGeneration: 0,
inspectorInitialized: false,
queued: 0,
waiters: new Set(),
};
}

function closeRunState(state: ActorRunState): void {
state.closed = true;
state.inspector?.dispose?.();
state.inspector = undefined;
notifyStateChanged(state);
}

function waitForStateChange(state: ActorRunState): Promise<void> {
return new Promise((resolve) => {
state.waiters.add(resolve);
Expand All @@ -69,8 +91,14 @@ export class RunHandlerCoordinator {
actorId: string,
restart: () => void | Promise<void>,
callback: () => void | Promise<void>,
actorGeneration?: number,
): Promise<RunHandlerOutcome> {
const state = this.#getOrCreate(actorId);
const state = this.#getOrCreate(actorId, actorGeneration);
// Closed state belongs to a superseded generation. It must not initialize an inspector,
// which would rebind the workflow registration of the current generation.
if (state.closed) {
return "closed";
}
state.restart = restart;
this.#initializeInspector(actorId, state);
state.queued += 1;
Expand Down Expand Up @@ -110,21 +138,33 @@ export class RunHandlerCoordinator {
getInspector(
actorId: string,
restart: () => void | Promise<void>,
actorGeneration?: number,
): RunInspectorConfig | undefined {
const state = this.#getOrCreate(actorId);
const state = this.#getOrCreate(actorId, actorGeneration);
if (state.closed) {
return undefined;
}
state.restart = restart;
this.#initializeInspector(actorId, state);
return state.inspector?.inspector;
}

destroy(actorId: string): void {
/**
* Closes the run state for `actorId`. With `actorGeneration`, only state owned by that
* generation is closed, so an older generation's late cleanup cannot close a newer one.
*/
destroy(actorId: string, actorGeneration?: number): void {
const state = this.#states.get(actorId);
if (!state) return;
if (
actorGeneration !== undefined &&
state.actorGeneration !== undefined &&
state.actorGeneration !== actorGeneration
) {
return;
}

state.closed = true;
state.inspector?.dispose?.();
state.inspector = undefined;
notifyStateChanged(state);
closeRunState(state);
this.#states.delete(actorId);
}

Expand Down Expand Up @@ -173,18 +213,32 @@ export class RunHandlerCoordinator {
};
}

#getOrCreate(actorId: string): ActorRunState {
#getOrCreate(
actorId: string,
actorGeneration: number | undefined,
): ActorRunState {
let state = this.#states.get(actorId);
if (
state &&
actorGeneration !== undefined &&
state.actorGeneration !== undefined &&
state.actorGeneration !== actorGeneration
) {
if (actorGeneration < state.actorGeneration) {
// A superseded generation gets closed, detached state so it cannot run or
// block the current generation.
const detached = createRunState(actorGeneration);
detached.closed = true;
return detached;
}
// A newer generation replaces state left by an older one, which may still be
// marked active because its JS callback cannot be cancelled.
closeRunState(state);
this.#states.delete(actorId);
state = undefined;
}
if (!state) {
state = {
active: false,
closed: false,
exclusive: false,
exclusiveGeneration: 0,
inspectorInitialized: false,
queued: 0,
waiters: new Set(),
};
state = createRunState(actorGeneration);
this.#states.set(actorId, state);
}
return state;
Expand All @@ -194,6 +248,7 @@ export class RunHandlerCoordinator {
if (state.inspectorInitialized) return;
const inspector = createRunInspector(this.#run, {
actorId,
actorGeneration: state.actorGeneration,
control: this.#control(actorId, state),
});
if (this.inspectorKind === "workflow") {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -583,6 +583,7 @@ export interface CoreRuntime {
writes: RuntimeWorkflowKvWrite[],
): Promise<void>;
actorId(ctx: ActorContextHandle): string;
actorGeneration(ctx: ActorContextHandle): number | undefined;
/**
* Runs one actor callback with `ctx` as the current invocation: operations
* on retained handles for the same actor resolve to it, and its Core span
Expand Down
Loading
Loading