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
}ShippedEventStore 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.