org.llm4s.agent.graph

Members list

Type members

Classlikes

final case class Checkpoint(formatVersion: Int, id: String, parent: Option[String], threadId: String, runId: String, status: CheckpointStatus, createdAt: Instant, snapshot: GraphSnapshot)

One durable point in a thread's execution - data only. Closures, codecs and update functions are rebound from the compiled graph on restore; every value inside carries the version of the codec that wrote it (VersionedJson).

One durable point in a thread's execution - data only. Closures, codecs and update functions are rebound from the compiled graph on restore; every value inside carries the version of the codec that wrote it (VersionedJson).

parent is the checkpoint this one supersedes. A checkpointer accepts a new checkpoint only if parent is the thread's latest, so two writers cannot both advance one thread.

Persist with Checkpoint.toJson and read with Checkpoint.fromJson, which migrates older formatVersions and refuses newer ones.

Attributes

Companion
object
Supertypes
trait Serializable
trait Product
trait Equals
class Object
trait Matchable
class Any
Show all
object Checkpoint

Attributes

Companion
class
Supertypes
trait Product
trait Mirror
class Object
trait Matchable
class Any
Self type
Checkpoint.type

Whether a thread's latest checkpoint is mid-execution or finished its run.

Whether a thread's latest checkpoint is mid-execution or finished its run.

Attributes

Supertypes
trait Enum
trait Serializable
trait Product
trait Equals
class Object
trait Matchable
class Any
Show all
trait Checkpointer

Durable storage for graph threads: each thread's latest checkpoint, the pending writes recorded against it, and its durable event log. One store owns all three so that a commit is atomic.

Durable storage for graph threads: each thread's latest checkpoint, the pending writes recorded against it, and its durable event log. One store owns all three so that a commit is atomic.

Implementations must be safe to call from several threads, and must:

  • apply a Commit entirely or not at all;
  • accept a new checkpoint only when its parent is the thread's latest checkpoint id (GraphError.CheckpointConflict otherwise), and only pending writes naming the resulting latest checkpoint (GraphError.InvalidCommit);
  • number a thread's events 1, 2, 3, ... in commit order, inside the commit, never reusing a number - including after compaction or a failed commit;
  • keep events until compactEvents removes them, and report the earliest sequence still available when asked to replay from before it (GraphError.ReplayUnavailable).

Fencing a commit with a run-claim token is Stage 2; the parent check is its precursor.

Attributes

Supertypes
class Object
trait Matchable
class Any
Known subtypes
final case class Command(update: StateUpdate, routes: List[Route])

A node's state update and routes. Static edges declared on the builder are additive: they are scheduled before the command's routes. A command with no routes and no outgoing edges ends that branch; the run completes when no branch has work left.

A node's state update and routes. Static edges declared on the builder are additive: they are scheduled before the command's routes. A command with no routes and no outgoing edges ends that branch; the run completes when no branch has work left.

Attributes

Companion
object
Supertypes
trait Serializable
trait Product
trait Equals
class Object
trait Matchable
class Any
Show all
object Command

Attributes

Companion
class
Supertypes
trait Product
trait Mirror
class Object
trait Matchable
class Any
Self type
Command.type
final case class Commit(checkpoint: Option[Checkpoint], pendingWrites: Vector[PendingWrite], events: Vector[EventDraft])

One atomic write to a thread. The checkpointer applies all of it or none of it:

One atomic write to a thread. The checkpointer applies all of it or none of it:

  • checkpoint, if present, becomes the thread's latest, provided its parent is the current latest; pending writes recorded against the old checkpoint are dropped.
  • pendingWrites are recorded against the (resulting) latest checkpoint, and must name it.
  • events are appended to the thread's durable log, each given the next sequence number.

Attributes

Supertypes
trait Serializable
trait Product
trait Equals
class Object
trait Matchable
class Any
Show all
final class CompiledGraph[I, O]

An immutable, validated graph, run in supersteps.

An immutable, validated graph, run in supersteps.

Each superstep runs every ready task against the same committed ThreadState, then commits their updates and schedules the next frontier together:

  1. Tasks run in frontier order (possibly concurrently); none sees another's writes.
  2. If any task fails, the superstep commits nothing and the run fails with the first failure in frontier order.
  3. Each command is checked: every update is to a key in the node's declared write set, every route targets a node and join of this graph, and a task fans out to a join at most once.
  4. Updates apply task by task in frontier order, and in emission order within a task.
  5. The next frontier is, for each task in order, its static edges (declaration order) then its routes (route order, fan-out payloads in item order); then the targets of joins this commit released - static joins in declaration order, then dynamic activations in the order they opened.

The run completes when the frontier is empty and no join is waiting; a join still waiting at that point can never release and fails the run with GraphError.UnsatisfiedJoin.

Attributes

Supertypes
class Object
trait Matchable
class Any
enum Durability

When a run's checkpoints and events become durable. In every mode a durable event is delivered only after the commit that numbered it, so no subscriber sees a sequence number that a crash could later reuse. Suspensions (#1269) will be persisted synchronously in every mode before a suspended result is returned.

When a run's checkpoints and events become durable. In every mode a durable event is delivered only after the commit that numbered it, so no subscriber sees a sequence number that a crash could later reuse. Suspensions (#1269) will be persisted synchronously in every mode before a suspended result is returned.

Attributes

Supertypes
trait Enum
trait Serializable
trait Product
trait Equals
class Object
trait Matchable
class Any
Show all
final class DynamicJoin

A barrier attached to a fan-out. A Route.FanOut opens an activation recording the task id of every child it schedules; target is scheduled once all of them have completed. An empty fan-out releases immediately.

A barrier attached to a fan-out. A Route.FanOut opens an activation recording the task id of every child it schedules; target is scheduled once all of them have completed. An empty fan-out releases immediately.

Attributes

Supertypes
class Object
trait Matchable
class Any

A state operation as data: updates are encoded with the key's update codec.

A state operation as data: updates are encoded with the key's update codec.

Attributes

Supertypes
trait Enum
trait Serializable
trait Product
trait Equals
class Object
trait Matchable
class Any
Show all

A route as data: payloads are encoded with the target node's input codec.

A route as data: payloads are encoded with the target node's input codec.

Attributes

Supertypes
trait Enum
trait Serializable
trait Product
trait Equals
class Object
trait Matchable
class Any
Show all
final case class EventDraft(runId: String, checkpointId: Option[String], taskId: Option[String], nodeId: Option[String], timestamp: Instant, event: RunEvent)

A durable event before its commit assigns it a sequence number.

A durable event before its commit assigns it a sequence number.

Attributes

Supertypes
trait Serializable
trait Product
trait Equals
class Object
trait Matchable
class Any
Show all
final case class EventRecord(threadId: String, seq: Long, runId: String, checkpointId: Option[String], taskId: Option[String], nodeId: Option[String], timestamp: Instant, event: RunEvent)

A committed durable event. seq is per thread, allocated in the commit that stored the event, contiguous and never reused; an event is delivered to subscribers only after that commit.

A committed durable event. seq is per thread, allocated in the commit that stored the event, contiguous and never reused; an event is delivered to subscribers only after that commit.

Attributes

Companion
object
Supertypes
trait Serializable
trait Product
trait Equals
class Object
trait Matchable
class Any
Show all
object EventRecord

Attributes

Companion
class
Supertypes
trait Product
trait Mirror
class Object
trait Matchable
class Any
Self type
final class Execution

A graph run paused at a superstep boundary: the committed state, the ready frontier and the open join activations. Advance it with CompiledGraph.step; persist it with CompiledGraph.snapshot.

A graph run paused at a superstep boundary: the committed state, the ready frontier and the open join activations. Advance it with CompiledGraph.step; persist it with CompiledGraph.snapshot.

Attributes

Supertypes
class Object
trait Matchable
class Any
final class GraphBuilder

Builds a typed graph. The builder issues every handle - nodes, joins - and validates the whole graph in compile, reporting every problem at once.

Builds a typed graph. The builder issues every handle - nodes, joins - and validates the whole graph in compile, reporting every problem at once.

Nodes that route to each other in a cycle are declared first and implemented afterwards:

val b      = GraphBuilder("counter", "v1")
val count  = StateKey.replace[Int]("count", 0)
val tick   = b.declare[Unit]("tick")
b.implement(tick, writes = Set(count)) { (_, state, _) =>
 NodeResult.fromResult(state.get(count).map { n =>
   val next = Command.empty.update(count, n + 1)
   if n + 1 < 3 then next.goto(tick) else next
 })
}
val graph = b.compile(tick)(_.get(count))

The builder is mutable and meant to be used from one thread while the graph is assembled; the CompiledGraph it produces is immutable.

Attributes

Companion
object
Supertypes
class Object
trait Matchable
class Any
object GraphBuilder

Attributes

Companion
class
Supertypes
class Object
trait Matchable
class Any
Self type
sealed trait GraphError extends NonRecoverableError

Errors raised while building, running or restoring a typed graph.

Errors raised while building, running or restoring a typed graph.

Attributes

Companion
object
Supertypes
trait LLMError
trait Serializable
trait Product
trait Equals
class Object
trait Matchable
class Any
Show all
Known subtypes
object GraphError

Attributes

Companion
trait
Supertypes
trait Sum
trait Mirror
class Object
trait Matchable
class Any
Self type
GraphError.type
trait GraphNode[I]

A node's behaviour. It reads the superstep's committed snapshot - never another task's uncommitted writes - and returns its effects as data.

A node's behaviour. It reads the superstep's committed snapshot - never another task's uncommitted writes - and returns its effects as data.

Attributes

Supertypes
class Object
trait Matchable
class Any
final class GraphRuntime(checkpointer: Checkpointer, clock: Clock)

Runs compiled graphs on durable threads.

Runs compiled graphs on durable threads.

start begins a run: on a new thread from the graph's entry, on a thread whose latest run completed by applying the input to its committed state. It refuses a thread with an incomplete execution (GraphError.IncompleteRun). recover continues an incomplete execution from its latest checkpoint without new input: tasks with a pending write are not run again, and failed or unstarted tasks run once more (per-node retry policy is Stage 1).

subscribe replays a thread's committed events after a sequence number and then delivers new ones as their commits succeed, in ascending order with no gaps or duplicates, followed by live progress as it happens. Events are delivered on the committing thread; a dedicated ordered dispatcher with bounded queues is part of the run API (#1271).

Attributes

Supertypes
class Object
trait Matchable
class Any
final case class GraphSnapshot(graphId: String, graphVersion: String, fingerprint: String, superstep: Int, state: Map[String, VersionedJson], frontier: Vector[PendingTask], staticJoins: Vector[StaticArrivals], dynamicJoins: Vector[Activation])

A serializable picture of an Execution - data only, no closures or codecs. Node inputs and state values are encoded with the codecs of the graph that wrote them, and decoded and checked against the graph that restores them.

A serializable picture of an Execution - data only, no closures or codecs. Node inputs and state values are encoded with the codecs of the graph that wrote them, and decoded and checked against the graph that restores them.

Each value records its codec's version and is migrated on restore; the snapshot itself is versioned by the Checkpoint that carries it.

Attributes

Companion
object
Supertypes
trait Serializable
trait Product
trait Equals
class Object
trait Matchable
class Any
Show all
object GraphSnapshot

Attributes

Companion
class
Supertypes
trait Product
trait Mirror
class Object
trait Matchable
class Any
Self type
final class InMemoryCheckpointer extends Checkpointer

A Checkpointer in memory. Checkpoints and pending writes are stored as JSON and decoded on read, exactly as a database-backed store would, so nothing executable survives a round trip.

A Checkpointer in memory. Checkpoints and pending writes are stored as JSON and decoded on read, exactly as a database-backed store would, so nothing executable survives a round trip.

Attributes

Supertypes
trait Checkpointer
class Object
trait Matchable
class Any
object JoinId

Attributes

Supertypes
class Object
trait Matchable
class Any
Self type
JoinId.type
final class NodeContext

The running task's identity, for attribution and idempotency keys, and its event channels.

The running task's identity, for attribution and idempotency keys, and its event channels.

emit records a durable custom event: it is committed with this task's result, given a per-thread sequence number in that commit, delivered only after the commit and replayed to later subscribers. It is discarded if the task fails. progress is live-only, for token deltas and similar high-volume progress: delivered at once to current subscribers, never persisted or replayed. Outside a GraphRuntime both are no-ops.

Attributes

Supertypes
class Object
trait Matchable
class Any
object NodeId

Attributes

Supertypes
class Object
trait Matchable
class Any
Self type
NodeId.type
final class NodeRef[I]

A handle to a node that consumes I, issued by a GraphBuilder.

A handle to a node that consumes I, issued by a GraphBuilder.

Handles are the only way to route: a node cannot name another by string. Goto takes a NodeRef[Unit] and Send a payload of the target's input type, so a mis-typed route does not compile. A handle from another builder is rejected when the graph compiles, or when a node returns it at run time.

Attributes

Supertypes
class Object
trait Matchable
class Any
enum NodeResult

What a node task produced.

What a node task produced.

Attributes

Companion
object
Supertypes
trait Enum
trait Serializable
trait Product
trait Equals
class Object
trait Matchable
class Any
Show all
object NodeResult

Attributes

Companion
enum
Supertypes
trait Sum
trait Mirror
class Object
trait Matchable
class Any
Self type
NodeResult.type
final case class PendingWrite(checkpointId: String, taskId: String, nodeId: String, operations: Vector[EncodedOperation], routes: Vector[EncodedRoute])

A completed task's result, recorded against the checkpoint whose frontier it ran in, before the superstep commits. On recovery the task is not run again: its command is decoded from here.

A completed task's result, recorded against the checkpoint whose frontier it ran in, before the superstep commits. On recovery the task is not run again: its command is decoded from here.

Attributes

Supertypes
trait Serializable
trait Product
trait Equals
class Object
trait Matchable
class Any
Show all
enum Route

A routing decision returned by a node; scheduled for the next superstep.

A routing decision returned by a node; scheduled for the next superstep.

Attributes

Supertypes
trait Enum
trait Serializable
trait Product
trait Equals
class Object
trait Matchable
class Any
Show all
Known subtypes
class Send[I]
class FanOut[I]
enum RunEvent

What happened in a run. Durable: persisted in the thread's event log and replayable.

What happened in a run. Durable: persisted in the thread's event log and replayable.

Attributes

Supertypes
trait Enum
trait Serializable
trait Product
trait Equals
class Object
trait Matchable
class Any
Show all
object RunId

Attributes

Supertypes
class Object
trait Matchable
class Any
Self type
RunId.type
enum RunResult[+O]

The outcome of a run.

The outcome of a run.

Attributes

Supertypes
trait Enum
trait Serializable
trait Product
trait Equals
class Object
trait Matchable
class Any
Show all
Known subtypes
class Completed[O]
class Failed
final class SchemaVersion

The current version of a persisted JSON shape, with the migrations that bring older versions up to it. A state key's value, a state key's update and a node's input each carry one; their stable codec id is the owner's id with this version (<keyId>@<version>), and every encoded value records the version it was written with (see VersionedJson).

The current version of a persisted JSON shape, with the migrations that bring older versions up to it. A state key's value, a state key's update and a node's input each carry one; their stable codec id is the owner's id with this version (<keyId>@<version>), and every encoded value records the version it was written with (see VersionedJson).

A step keyed n migrates JSON written at version n to version n + 1. Steps are plain functions of the compiled graph, never persisted.

Attributes

Companion
object
Supertypes
class Object
trait Matchable
class Any
object SchemaVersion

Attributes

Companion
class
Supertypes
class Object
trait Matchable
class Any
Self type
final class StateKey[A, U]

A typed slot in a graph's thread state.

A typed slot in a graph's thread state.

A is the stored value and U the update a node emits for it. Every update - including a single writer's - is applied to the committed value with applyUpdate, in deterministic task and emission order; replacement is the special case U = A (see StateKey.replace). The function belongs to the compiled graph, never to a checkpoint: checkpoints hold values encoded with stateCodec and pending updates encoded with updateCodec, each tagged with the codec's SchemaVersion so an older checkpoint is migrated before it is decoded.

Keys are compared by identity. A graph rejects two distinct keys with the same id.

Attributes

Companion
object
Supertypes
class Object
trait Matchable
class Any
object StateKey

Attributes

Companion
class
Supertypes
class Object
trait Matchable
class Any
Self type
StateKey.type
object StateKeyId

Attributes

Supertypes
class Object
trait Matchable
class Any
Self type
StateKeyId.type
final class StateUpdate

An ordered list of state operations emitted by one node task.

An ordered list of state operations emitted by one node task.

Operations apply in emission order. combine is sequential composition: a.combine(b) applies a's operations and then b's, so for a StateKey.replace key the later value wins and for an operation-valued key both operations apply. Updates from different tasks in one superstep are never combined; the scheduler applies them task by task in frontier order.

Attributes

Companion
object
Supertypes
class Object
trait Matchable
class Any
object StateUpdate

Attributes

Companion
class
Supertypes
class Object
trait Matchable
class Any
Self type
final class StaticJoin

A barrier that schedules target once every source node has completed a task since its last release. Each source counts once per activation, however many of its tasks complete.

A barrier that schedules target once every source node has completed a task since its last release. Each source counts once per activation, however many of its tasks complete.

Attributes

Supertypes
class Object
trait Matchable
class Any
enum Step[+O]

The outcome of one superstep.

The outcome of one superstep.

Attributes

Supertypes
trait Enum
trait Serializable
trait Product
trait Equals
class Object
trait Matchable
class Any
Show all
Known subtypes
class Next
class Done[O]
final case class StoredCheckpoint(checkpoint: Checkpoint, pendingWrites: Vector[PendingWrite])

A thread's latest checkpoint with the pending writes recorded against it.

A thread's latest checkpoint with the pending writes recorded against it.

Attributes

Supertypes
trait Serializable
trait Product
trait Equals
class Object
trait Matchable
class Any
Show all

What a subscriber receives.

What a subscriber receives.

Attributes

Supertypes
trait Enum
trait Serializable
trait Product
trait Equals
class Object
trait Matchable
class Any
Show all
trait Subscription

A registration with GraphRuntime.subscribe.

A registration with GraphRuntime.subscribe.

Attributes

Supertypes
class Object
trait Matchable
class Any
object TaskId

Attributes

Supertypes
class Object
trait Matchable
class Any
Self type
TaskId.type
object ThreadId

Attributes

Supertypes
class Object
trait Matchable
class Any
Self type
ThreadId.type
final class ThreadState

The committed state of a graph thread: a typed key/value store over the graph's registered StateKeys. It is immutable, so every task in a superstep reads the same snapshot.

The committed state of a graph thread: a typed key/value store over the graph's registered StateKeys. It is immutable, so every task in a superstep reads the same snapshot.

Reads are checked: a key the graph did not register - or a different key instance sharing a registered key's id - is an GraphError.UnknownStateKey. A registered key with no stored entry reads as its initial value. The erasure to Any stays private to this boundary.

Attributes

Supertypes
class Object
trait Matchable
class Any
final case class VersionedJson(version: Int, value: Value)

A persisted value and the version of the codec that wrote it.

A persisted value and the version of the codec that wrote it.

Attributes

Supertypes
trait Serializable
trait Product
trait Equals
class Object
trait Matchable
class Any
Show all

Types

opaque type JoinId

Stable identifier of a static or dynamic join barrier; persisted in snapshots.

Stable identifier of a static or dynamic join barrier; persisted in snapshots.

Attributes

opaque type NodeId

Stable identifier of a node in a compiled graph; persisted in snapshots.

Stable identifier of a node in a compiled graph; persisted in snapshots.

Attributes

opaque type RunId

One execution attempt on a thread; start, recover and (later) resume each begin a new run.

One execution attempt on a thread; start, recover and (later) resume each begin a new run.

Attributes

opaque type StateKeyId

Stable identifier of a StateKey; persisted in snapshots.

Stable identifier of a StateKey; persisted in snapshots.

Attributes

opaque type TaskId

Identifier of one scheduled node execution; unique within a thread's history.

Identifier of one scheduled node execution; unique within a thread's history.

Attributes

opaque type ThreadId

A conversation or workflow thread: the address of its checkpoints and event log.

A conversation or workflow thread: the address of its checkpoints and event log.

Attributes