Persistence

Event sourcing model

Writes append after the transition succeeds. Recovery loads an optional snapshot, then replays later events (transition actions run; entry/producing do not).

{
  enum OrderState derives Finite:
    case Pending, Paid, Shipped

  enum OrderEvent derives Finite:
    case Pay, Ship

  import OrderState.*, OrderEvent.*

  val machine = Machine(
    assembly[OrderState, OrderEvent](
      Pending via Pay to Paid,
      Paid via Ship to Shipped,
    )
  )

  val orderId: OrderId = "order-persist-1"

  ZIO
    .scoped {
      for
        fsm   <- FSMRuntime(orderId, machine, Pending)
        _     <- fsm.send(Pay)
        _     <- fsm.send(Ship)
        _     <- fsm.saveSnapshot
        state <- fsm.currentState
        seq   <- fsm.lastSequenceNr
      yield (state.toString, seq)
    }
    .provide(
      InMemoryEventStore.layer[OrderId, OrderState, OrderEvent],
      TimeoutStrategy.fiber[OrderId],
      LockingStrategy.optimistic[OrderId],
    )
    .asDoc
}
(Shipped,2)

Recover after restart

There is no separate recover API: construct FSMRuntime again with the same id against the same EventStore. Session one writes history; session two resumes at Shipped:

{
  enum OrderState derives Finite:
    case Pending, Paid, Shipped

  enum OrderEvent derives Finite:
    case Pay, Ship

  import OrderState.*, OrderEvent.*

  val machine = Machine(
    assembly[OrderState, OrderEvent](
      Pending via Pay to Paid,
      Paid via Ship to Shipped,
    )
  )

  val orderId: OrderId = "order-recover-1"

  ZIO.scoped {
    for
      store <- InMemoryEventStore.make[OrderId, OrderState, OrderEvent]()
      _     <- ZIO
        .scoped {
          FSMRuntime(orderId, machine, Pending).flatMap { fsm =>
            fsm.send(Pay) *> fsm.send(Ship) *> fsm.saveSnapshot
          }
        }
        .provide(
          ZLayer.succeed(store),
          TimeoutStrategy.fiber[OrderId],
          LockingStrategy.optimistic[OrderId],
        )
      recovered <- ZIO
        .scoped {
          FSMRuntime(orderId, machine, Pending).flatMap(_.currentState)
        }
        .provide(
          ZLayer.succeed(store),
          TimeoutStrategy.fiber[OrderId],
          LockingStrategy.optimistic[OrderId],
        )
    yield recovered
  }.asDoc
}
Shipped

EventStore and codecs

Implement EventStore[Id, S, E] for your backend (append, loadEvents, snapshots, …). append must use optimistic locking: atomically check expectedSeqNr, then increment.

PostgreSQL ships as mechanoid-postgres. Derive JSON codecs with import mechanoid.postgres.* (finiteJsonCodec from Finite) and initialize schema via PostgresSchema.initialize (see examples/heartbeat).

Optimistic locking

Concurrent writers that lose the race see SequenceConflictError. Reload and retry, or move up to Distributed Coordination to prevent conflicts upfront.

Next: Durable Timeouts.