Skip to content

Latest commit

 

History

History
357 lines (280 loc) · 15.6 KB

File metadata and controls

357 lines (280 loc) · 15.6 KB

effect-machine

Type-safe state machines for Effect.

Commands

bun run gate          # typecheck + lint + test + build
bun test              # Run tests
bun run typecheck     # tsgo --noEmit, patched by @effect/tsgo
bun run lint          # type-aware oxlint; Effect diagnostics run through tsgo plugin
bun run fmt           # oxfmt

Conventions

  • Files: kebab-case (actor.ts, persistent-actor.ts)
  • States/Events: schema-first with State({...}) / Event({...}) - they ARE schemas
  • Empty structs: plain values - State.Idle (not callable)
  • Non-empty: State.Loading({ url }) - constructor requiring args
  • Machine creation: Machine.make({ state, event, initial }) - types inferred
  • Exports: all public API via src/index.ts
  • Namespace pattern: import { Machine } from "effect-machine" then Machine.make, etc.
  • v4 services use class X extends Context.Service<X, Shape>()("key") {}; keep serviceNotAsClass enabled.
  • effect stays peer-only at runtime; keep concrete Effect versions in dev deps for validation.

Fluent Builder

const machine = Machine.make({ state, event, initial })
  .on(State.Idle, Event.Start, () => State.Running)
  .when(
    State.Ready,
    Event.Submit,
    ({ state }) => state.valid,
    () => State.Saving,
  )
  .on([State.Draft, State.Review], Event.Cancel, () => State.Cancelled) // multi-state
  .onAny(Event.Reset, () => State.Idle) // wildcard (any state)
  .spawn(State.Running, () => Effect.flatMap(Poller, (poller) => poller.poll()))
  .timeout(State.Loading, { duration: Duration.seconds(30), event: Event.Timeout })
  .postpone(State.Connecting, Event.Data)
  .final(State.Done);
  • Builder methods mutate this, return this
  • Builder chain ends naturally — no terminal method needed
  • .onAny() fires when no specific .on() matches for that event

.task()

Async work that emits an event on completion:

// Explicit onSuccess mapping
.task(State.Loading, ({ state }) => fetchData(state.url), {
  onSuccess: (data) => Event.Loaded({ data }),
  onFailure: (error) => Event.Failed({ error: String(error) }),
})

// Shorthand — when task returns Event directly, onSuccess can be omitted
.task(State.Loading, ({ state }) => fetchData(state.url).pipe(Effect.map(d => Event.Loaded({ data: d }))), {
  onFailure: (error) => Event.Failed({ error: String(error) }),
})

// Multi-state
.task([State.Loading, State.Retrying], ({ state }) => fetchData(state.url), { onSuccess: ... })

State.with()

Construct state from existing source. Per-variant and union-level:

// Per-variant: target-specific, works cross-state
State.Active.with(state, { count: state.count + 1 });
State.Shipped.with(processingState, { trackingId: "TRACK-123" });
State.Idle.with(anyState); // → { _tag: "Idle" }

// Union-level: dispatches by _tag, preserves specific variant subtype
const updated = MyState.with(state, { queue: newQueue });
// If state is Streaming, returns Streaming (not union type)
// Partial keys not in target variant are silently dropped

Effect Services

Use standard Effect services in task, spawn, and background handlers:

class Api extends Context.Service<Api, { readonly fetch: (url: string) => Effect.Effect<Data> }>()(
  "app/Api",
) {}

const machine = Machine.make({ state, event, initial }).task(
  State.Loading,
  ({ state }) => Effect.flatMap(Api, (api) => api.fetch(state.url)),
  { onSuccess: (data) => Event.Loaded({ data }) },
);

const actor = yield * Machine.spawn(machine).pipe(Effect.provideService(Api, { fetch: Http.get }));
yield * actor.start;

Machine.spawn captures the current Effect context. The actor keeps that context when actor.start runs later.

Transition handlers and .when() predicates can use Effect services. Their error channels must be never.

Running Machines

Simple (no registry):

Machine.spawn returns an unstarted actor. Call actor.start to fork the event loop.

const actor = yield * Machine.spawn(machine);
yield * actor.start; // fork event loop, background effects, spawn effects
yield * actor.stop; // caller responsible

// Scope-aware — use Machine.scoped to bridge ActorScope from Scope:
yield *
  Effect.scoped(
    Machine.scoped(
      Effect.gen(function* () {
        const actor = yield* Machine.spawn(machine);
        yield* actor.start;
        // actor.stop called automatically when scope closes
      }),
    ),
  );

With registry:

system.spawn auto-starts — no actor.start needed.

const system = yield * ActorSystemService;
const actor = yield * system.spawn("my-id", machine);

ActorScope: Machine.spawn and system.spawn detect ActorScope via Effect.serviceOption — if present, attach stop finalizer; if absent, skip. Use Machine.scoped(effect) to bridge Scope.Scope → ActorScope. This is explicit opt-in — ambient Scope.Scope does NOT trigger auto-cleanup (prevents bugs where unrelated scopes tear down actors).

Recovery + Durability

Lifecycle hooks for persistence. Replace the old PersistConfig.

const actor =
  yield *
  Machine.spawn(machine, {
    lifecycle: {
      recovery: {
        resolve: ({ actorId, generation, machineInitial }) =>
          storage.get(actorId).pipe(Effect.map(Option.fromNullable)),
      },
      durability: {
        save: ({ actorId, nextState, event }) => storage.set(actorId, nextState),
        shouldSave: (state, prev) => state._tag !== prev._tag,
      },
    },
  });
yield * actor.start;
  • Recovery — resolves initial state per generation. Runs during actor.start. Returns Option<S>: Some overrides initial, None uses machine.initial. generation is 0 for cold start, 1+ for supervision restarts.
  • Durability — saves state after committed transitions. Receives DurabilityCommit with actorId, generation, previousState, nextState, event. shouldSave is a sync predicate to skip uninteresting transitions.
  • hydrate overrides recovery — Machine.spawn(machine, { hydrate: state }) skips recovery.resolve entirely.

Supervision

Actors can automatically restart on defect with Supervision.restart():

import { Supervision } from "effect-machine";

const actor =
  yield *
  Machine.spawn(machine, {
    supervision: Supervision.restart({ maxRestarts: 3, within: "1 minute" }),
  });
yield * actor.start;

// Via system (auto-starts)
const actor =
  yield *
  system.spawn("id", machine, {
    supervision: Supervision.restart(),
  });

// Observe exit reason
const exit = yield * actor.awaitExit; // ActorExit<S>
const otherExit = yield * other.awaitExit; // ActorExit<OtherState>
  • Restart from machine.initial — always clean slate, never last-state
  • Actor ID survives — same identity across restarts
  • Pending requests fail — call/ask/sendWait behind crash get ActorStoppedError
  • Children die — both scopes close; children come back only if restart re-runs spawn/background
  • stop/drain are terminal — no restart
  • Final state = no restart — awaitExit resolves with ActorExit.Final
  • Budget — Schedule controls timing/count; exhaustion = terminal ActorExit.Defect
  • Classifier — shouldRestart optionally skips restart for specific defect types
  • Entity-machine: cluster-supervised via defectRetryPolicy, NOT local supervision

Child Actors

Spawn children from .spawn()/.background() handlers via self.spawn(id, childMachine):

machine.spawn(State.Active, ({ self }) =>
  Effect.gen(function* () {
    const child = yield* self.spawn("worker-1", workerMachine).pipe(Effect.orDie);
    yield* child.send(WorkerEvent.Start);
    // child auto-stopped when parent exits Active state
  }),
);
  • Children spawned in .spawn() handlers are state-scoped — auto-stopped on state exit
  • Children spawned in .background() handlers live for machine lifetime
  • self.spawn returns Effect<ActorRef, DuplicateActorError, R> — use Effect.orDie in handlers
  • Every ActorRef has actor.system for child access: actor.system.get("worker-1")

Lazy State-Owned Actors

Use ActorHost.make({ identity, spawn }) when consumers must request a child without owning its lifetime.

  • Construct the host in the service layer. spawn captures those services.
  • Run host.host(input) in the parent state's .spawn handler. It registers a generation and waits for a consumer.
  • Consumers call host.acquire(input). Identity values match with Object.is. The first matching consumer supplies the factory input.
  • For parent-owned startup data, type the second spawn(request, hostInput) argument and call host(input, hostInput). Consumers still call acquire(input). Host data belongs to that generation; pass immutable values. One-argument factories need no host data.
  • Concurrent consumers share startup and its result. Cancelling one consumer does not cancel startup. ActorHost starts actors from either Machine.spawn or system.spawn before it returns them.
  • The factory receives the host generation's Scope and ActorScope. State exit closes the child. Closing the host service also closes the current generation and fails pending consumers.
  • A second active host fails with ActorHostOccupiedError. A closed generation fails pending acquisition with ActorHostClosedError.
  • Keep session validity, authorization, and completion events in application services. ActorHost only owns actor creation and lifetime.
  • A failed factory stays failed until the hosting scope closes. Handle expected factory errors in the spawn handler with a transition or event. Effect.orDie is appropriate only for invariant failures; it defects the parent and cannot support a later reentry.
  • Once registered, host is interrupted when its generation closes. Consumers get ActorHostClosedError before actor cleanup starts. An acquisition belongs to the generation it observed; call acquire again after reentry. Pre-registration calls to a closed host service fail with ActorHostClosedError.

ActorRef API

actor.send(event); // fire-and-forget
actor.call(event); // request-reply, returns ProcessEventResult
actor.ask(event); // typed reply (event must have Event.reply())
actor.waitFor(State.X); // wait for state (constructor or predicate)
actor.sendAndWait(ev, X); // send + wait for state
actor.awaitFinal; // wait for final state
actor.awaitExit; // completes when actor stops
actor.drain; // process remaining queue, then stop
actor.subscribe(fn); // sync callback, returns unsubscribe
actor.system; // ActorSystem
actor.children; // ReadonlyMap<string, ActorRef>

// JavaScript client for code outside Effect
actor.client.send(event);
actor.client.stop();
actor.client.getSnapshot();
actor.client.matches(tag);
actor.client.canSync(event); // Boolean predicates only
await actor.client.can(event); // all predicates

ask / reply

Events declare reply schemas via Event.reply(). Handlers use Machine.reply():

const MyEvent = Event({
  GetCount: Event.reply({}, Schema.Number),
  Reset: {},
});

.on(State.Active, Event.GetCount, ({ state }) =>
  Machine.reply(state, state.count),
)

const count = yield* actor.ask(Event.GetCount);  // number

Handler Type Constraints

Method Allowed R Why
.on() / .reenter() Any R State transition Effects
.when() predicate Any R Boolean or Effect condition
.spawn() / .background() Scope Finalizers allowed
.task() Any R Async work can use Effect services
  • Transition handlers and predicates can require services
  • Task, spawn, and background handlers can require services
  • Handlers cannot produce errors — error channel fixed to never
  • Handlers must return machine's state schema — wrong states rejected at compile time

Gotchas

  • waitFor owns its listener with Effect.acquireUseRelease. Keep the post-subscription recheck inside that lifetime. Cancellation and predicate defects must remove the listener.

  • actor.stop cancels pending startup and waits for recovery cleanup. Stop callers can cancel their own wait without cancelling shutdown. A stop from recovery marks that startup interrupted. Shutdown waits for its protected regions and finalizers to finish. Recovery cleanup defects reach every stop caller.

  • Protect both shutdown owner creation and its cache publication from interruption. Protecting only the cached body can cache cancellation when its protected region ends.

  • Recovery self-stop must also identify the supervisor fiber. A protected supervisor must not join an owner that waits for that same supervisor. Preserve supervised recovery cleanup defects on stop.

  • A supervised generation can fail during runtime.start. Let the loop read its recorded exit and apply the restart policy. Do not let that start failure end the supervisor before it completes the actor exit.

  • The runtime event loop must be ready before startup completes. Host client sends must still commit synchronous transitions before returning.

  • Publish Active only while the same generation is still Starting. A terminal lifecycle must never return to Active.

  • actor.start shares one startup result across concurrent and repeated calls, including failure or interruption. Stop the actor and spawn a new actor to retry initialization. It cannot restart a terminal actor.

  • ActorScope cleanup must match the registered actor identity before removing an ID. A later actor can reuse that ID.

  • Machine.spawn returns an unstarted actor — must call yield* actor.start. system.spawn auto-starts.

  • Never throw in Effect.gen — use yield* Effect.fail()

  • yield* Effect.yieldNow after send() to let effects run

  • simulate()/createTestHarness() don't run spawn effects

  • Same-state transitions skip spawn/finalizers — use .reenter() to force

  • Empty structs: State.Idle not State.Idle()

  • .onAny() only fires when no specific .on() matches

  • self.spawn errors with DuplicateActorError — wrap with Effect.orDie

  • Non-Effect code uses actor.client

  • Pending call/ask Deferreds settled with ActorStoppedError on stop

  • ask() only accepts events with Event.reply() — non-reply events are a type error

  • Reply decode failures (schema mismatch) are defects

Cluster / Entity Machines

Wire machines to @effect/cluster for distributed actors:

import { toEntity, EntityMachine } from "effect-machine/cluster";

const OrderEntity = toEntity(orderMachine, { type: "Order" });
const OrderEntityLayer = EntityMachine.layer(OrderEntity, orderMachine, {
  initializeState: (entityId) => OrderState.Pending({ orderId: entityId }),
  persistence: { strategy: "journal" },
});
  • toEntity generates Entity with Send/Ask/GetState/WatchState RPCs
  • EntityMachine.layer wires machine to cluster via shared runtime kernel
  • EntityActorRef: typed client wrapper (send/ask/snapshot/watch/waitFor)

Entity Persistence

Opt-in via EntityMachineOptions.persistence:

  • Snapshot strategy (default): background scheduler + deactivation finalizer
  • Journal strategy: inline event append on each RPC, replay on reactivation
  • PersistenceAdapter service tag resolved from context
  • Journal append failures defect entity — cluster retry restarts from snapshot
  • Hydration: snapshot → journal replay → initializeState → machine.initial

Cluster Gotchas

  • Entity tests use Entity.makeTestClient + ShardingConfig.layer + Effect.scoped
  • EntityMachine.layer accepts raw Machine
  • Entity RPCs use .tag field (not ._tag) to distinguish request types
  • WatchState test skipped due to effect beta Queue bug

Documentation

  • SKILL.md — AI agent quick reference