Imp.Run (Imp v0.5.0)

Copy Markdown View Source

An addressable execution of an Imp program with ordered semantic events.

Imp.call/2 is the minimal program boundary. Imp.Run is the optional runtime boundary for hosts that need to observe a composed program while it is running, cancel its in-flight effects, or explicitly authorize validated ReActV2 and RLM tool effects. Observation and cancellation use the owned run context; authorization decisions are carried explicitly in Imp.Execution. Events describe Imp execution only: they contain no ACP, MCP, UI or transport concepts. One run-owned delivery process invokes the event sink serially, so a slow observer preserves event order without delaying cancellation or owner cleanup. A sink should still hand work off promptly, because a blocked sink holds up its own later events and barriers.

The sink's return value is ignored. When a sink raises, throws or exits, the run's owner is sent {:imp_run_event_sink_failed, run_id, %{sequence: sequence, kind: kind, reason: {class, reason}}}, where class is :error, :throw or :exit, and delivery goes on with the next event. Imp does not know whether the sink stored the event before it failed (a store call that timed out may have landed), so it does not mark a hole itself; the host, which knows its store, decides what to record. The event is still in the snapshot.

Stopping or cancelling a run ends delivery. stop/1 first waits up to five seconds for the sink to finish what it has; cancel_with_events/3 does not wait. The event the sink was holding is then reported as a sink failure with reason :in_sink_when_stopped, because it may have been stored, and each event after it, which the sink never received, is reported as {:imp_run_event_undelivered, run_id, %{sequence: sequence, kind: kind}}. kind is nil for an event the snapshot no longer holds. These reports are in the owner's mailbox when stop/1 or cancel_with_events/3 returns, after any report the sink's own failures produced, in sequence order, and each event is reported at most once. The same reports are sent when the sink's process dies outright, for example because it was linked to a process that crashed; the run's control then ends too, and so does the run: an effect in flight has its cancellation called with {:run_control_ended, reason} and the task is killed.

events/1 reads the retained native sequence independently of sink progress. cancel_with_events/3 snapshots that sequence before cleanup, including one owner-recorded cancellation outcome. The snapshot is in-memory evidence, not a durable effects log: node or owner death loses it, and cancellation says nothing about whether an unfinished remote write landed. Persist authorization before dispatch when that guarantee is needed. Imp.Run.Event.to_map/1 serializes a redacted event. Model request and response observations cover Imp.LM.request/2; ReActV2 and RLM emit the semantic tool call and result events, and each :tool_result carries metadata.outcome, the call's Imp.Tool.outcome/1.

A :model_request carries the messages as its input and the rest of the request as metadata: :options, the request options with the tool definitions removed, and :tools_hash, a SHA-256 of those definitions (or nil when the request sent none). The definitions themselves are emitted once per distinct hash per run, as a :tools_sent event whose input is the tool list as sent, so a run's record holds every request whole without repeating a roster that does not change.

Capture is bounded by start/3's :max_event_bytes, :max_events and :max_snapshot_bytes. An oversized event payload becomes a digest and size marker before sink delivery. A failed event keeps a small error marker in its place, with the validated HTTP status, provider code and retryability when present and when the marker fits, never the error message or the request and response bodies. Snapshot eviction adds a :capture_gap marker. A sink receives every bounded event; the snapshot is a bounded recent window.

Imp.Run.Event.kinds/0 lists every event kind.

Summary

Functions

Sends {:imp_run_barrier, tag} to receiver once every event the run emitted before this call has been handed to its :event_sink.

Cancels registered effects before terminating the outer supervised task.

Cancels a run and returns its terminal event snapshot before releasing control.

Emits one ordered, redacted event when called inside an Imp.Run.

Returns the ordered redacted events retained by a running or completed run before stop.

Returns a new random identifier, prefix followed by _ and 16 URL-safe characters, of the kind Imp gives runs, model calls and tool calls.

Registers fun to be called with the cancel reason when the current run is cancelled, and returns a reference for unregister_cancellable/1.

Starts an unlinked supervised program run owned by the calling process.

Releases the event/cancellation control process after a run completes.

Removes a registration register_cancellable/1 made; nil is accepted and ignored.

Types

t()

@type t() :: %Imp.Run{control: pid(), id: String.t(), task: Task.t()}

Functions

barrier(run, receiver, tag)

@spec barrier(t(), pid(), term()) :: :ok

Sends {:imp_run_barrier, tag} to receiver once every event the run emitted before this call has been handed to its :event_sink.

A host that needs its sink to have seen a run's events before it acts, such as one that stores them and then reports the turn finished, waits for this message. The call exits when the run's control process is gone, and then the message never comes.

cancel(run, reason \\ :cancelled, timeout \\ 5000)

@spec cancel(t(), term(), timeout()) :: :ok

Cancels registered effects before terminating the outer supervised task.

The cancellations are given timeout between them, and the task another timeout to end before it is killed. A cancellation still running after its timeout is abandoned, so one that never returns delays the cancel by timeout rather than holding it.

cancel_with_events(run, reason \\ :cancelled, timeout \\ 5000)

Cancels a run and returns its terminal event snapshot before releasing control.

This snapshot remains available even when an asynchronous event sink blocks. Cancellation does not imply an unfinished external write did not happen.

emit(kind, attrs \\ [])

@spec emit(atom(), keyword() | map()) :: :ok

Emits one ordered, redacted event when called inside an Imp.Run.

events(run)

Returns the ordered redacted events retained by a running or completed run before stop.

new_event_id(prefix \\ "event")

@spec new_event_id(String.t()) :: String.t()

Returns a new random identifier, prefix followed by _ and 16 URL-safe characters, of the kind Imp gives runs, model calls and tool calls.

A host recording its own events beside a run's uses it so its identifiers have the same shape.

register_cancellable(fun)

@spec register_cancellable((term() -> term())) :: reference() | nil

Registers fun to be called with the cancel reason when the current run is cancelled, and returns a reference for unregister_cancellable/1.

A host running work of its own inside a run, such as an external request or a process it started, registers how to stop it, so cancel/3 stops that work too. Work registered after the run was cancelled is stopped at once, and the result is nil; so is work registered after the run's control has ended, whose fun is called with {:run_control_ended, reason}. Outside a run nothing is registered and the result is nil.

start(program, inputs, opts \\ [])

@spec start(struct(), map() | keyword(), keyword()) :: {:ok, t()} | {:error, term()}

Starts an unlinked supervised program run owned by the calling process.

By default a run takes a place in the machine-wide pool that all Imp tasks share, bounded by the :async_max_workers setting, and start/3 waits for a place when the pool is full.

Pass admission: {pool, limit} to count the run in a pool the host names instead, such as one per agent: at most limit runs hold a place in pool at once, and when it is full start/3 returns {:error, :busy} straight away, having stopped the control process it started for the run and started no task. The host keeps its own queue and starts the next run when one of its runs ends. The limit is read on each start, so a host that changes its setting passes the new one. A run in a named pool does not count against the machine-wide pool. Tasks it starts inside itself take places in the machine-wide pool as usual, except a stream of Imp tasks the run enumerates itself, which runs one item at a time on the run's own place, as it does for any run.

A run inherits the caller's Imp.Deadline, and deadline: 30_000 bounds it further. Requests made through Imp.Clients.ReqLLM inside the run are capped to the time left, Imp.Predict.ReActV2 makes no further request once it has passed, and tasks the run starts carry the same bound; a custom Imp.LM decides for itself whether to read it. The deadline does not bound the wait for a place in the pool: start/3 waits as it does without one, and returns {:error, :deadline_exceeded}, having started no work, when the deadline passed while it waited. Nor does it interrupt a tool that is already running; use cancel/2 for that.

An option this function does not know raises ArgumentError, so a misspelled :authorize cannot start a run whose tool calls nobody is asked about.

Options

  • :event_sink (function of arity 1) - Called with each Imp.Run.Event, in order, from one delivery process. Its return value is ignored; see the module documentation for what a sink that fails causes. When absent, events are only kept in the snapshot.

  • :authorize - Decides each ReActV2 and RLM tool call before it runs: a function of an Imp.Execution.Authorization returning :allow, {:deny, reason} or {:cancel, reason}. When absent, tool calls are not asked about. The default value is nil.

  • :authorization_timeout (pos_integer/0) - Milliseconds :authorize has for one decision; a decision that takes longer is {:deny, :authorization_timeout}. The default value is 30000.

  • :admission - {pool, limit}: count the run in the host's pool pool, where at most limit runs hold a place at once, instead of the machine-wide pool. A full pool returns {:error, :busy} at once.

  • :id (String.t/0) - The run's id. When absent, a random one is made.

  • :deadline - Bounds the run with Imp.Deadline: milliseconds counted from the call to start/3, :infinity, or {:deadline, absolute} in monotonic milliseconds. It is capped by the caller's own deadline. When absent, the run inherits the caller's deadline, if any.

  • :max_event_bytes - Largest event, in bytes of its external term format, delivered and kept whole; a larger one becomes a digest and size marker. :infinity never measures an event. The default value is 65536.

  • :max_events - Events the snapshot keeps; older ones are evicted. :infinity evicts none. The default value is 512.

  • :max_snapshot_bytes - Bytes the snapshot keeps; older events are evicted. :infinity evicts none. A host that must keep a complete record of a run sets all three bounds to :infinity. The default value is 4194304.

stop(run)

@spec stop(t()) :: :ok

Releases the event/cancellation control process after a run completes.

A run still going when its control is released is ended with it, as when its control ends for any other reason.

unregister_cancellable(ref)

@spec unregister_cancellable(reference() | nil) :: :ok

Removes a registration register_cancellable/1 made; nil is accepted and ignored.