Imp.Clients.ReqLLMBatch (Imp v0.5.0)

Copy Markdown View Source

Provider-neutral, bounded, resumable execution for durable LM batches.

Requests have stable caller-supplied IDs and JSON-safe payloads. A dispatcher receives each normalized request plus attempt metadata and must return one of:

  • {:ok, output}
  • {:transient, reason}
  • {:terminal, reason}
  • {:malformed, reason}

The checkpoint is an auditable JSON document. Every dispatch intent and outcome is appended to its event history before execution continues. Writes use a synced temporary file followed by an atomic rename. On resume, a request that was dispatched but has no committed outcome becomes :ambiguous and is never replayed automatically.

Summary

Functions

Builds a dispatcher backed by an existing ReqLLM client.

Resumes a batch from its checkpoint.

Starts a new durable batch.

Types

context()

@type context() :: %{request_id: String.t(), attempt: pos_integer()}

dispatcher()

@type dispatcher() :: (request(), context() -> outcome())

outcome()

@type outcome() ::
  {:ok, term()}
  | {:transient, term()}
  | {:terminal, term()}
  | {:malformed, term()}

request()

@type request() :: %{id: String.t(), payload: term()}

summary()

@type summary() :: %{
  checkpoint: Path.t(),
  complete?: boolean(),
  counts: %{required(atom()) => non_neg_integer()},
  requests: [map()]
}

Functions

req_llm_dispatcher(client, call_opts \\ [])

@spec req_llm_dispatcher(Imp.Clients.ReqLLM.t(), keyword()) :: dispatcher()

Builds a dispatcher backed by an existing ReqLLM client.

Request payloads may be a message list or %{"messages" => messages}. The adapter is provider-neutral because provider and model selection remain in the client ("openai:...", "anthropic:...", or "gemini:...").

resume(checkpoint, dispatcher, opts \\ [])

@spec resume(Path.t(), dispatcher(), keyword()) :: {:ok, summary()} | {:error, term()}

Resumes a batch from its checkpoint.

Persisted requests, attempts, and retry policy are authoritative. Resume accepts only runtime options: :num_threads, :timeout, and :validate_output.

run(requests, dispatcher, opts)

@spec run([map()], dispatcher(), keyword()) :: {:ok, summary()} | {:error, term()}

Starts a new durable batch.

Required options:

  • :checkpoint - destination for the atomic JSON checkpoint

Runtime options are :num_threads (default 4), :max_attempts (default 3), :timeout (default 30_000), and :validate_output, an optional arity-one callback returning :ok or {:error, reason}.