Skip to content

API reference

Auto-generated from docstrings. See the guides for narrative documentation; this page is for looking up exact signatures.

Agent

sofias_sdk_lite.AgentBuilder

Builder for constructing Agent instances declaratively.

The builder pattern allows configuring all aspects of an agent in a fluent, method-chaining style. Validation happens at build() time, catching configuration errors before runtime.

Example

agent = ( AgentBuilder("inbox_assist", version="1.0.0") .with_description("Agent that processes incoming emails") .with_sdk_config(sdk_config) .with_settings_class(InboxAssistSettings) .with_contract(input_schema=EmailInput, output_schema=EmailResponse) .add_llm_node("classifier", classifier_config, classifier_contract) .add_llm_node("drafter", drafter_config, drafter_contract) .add_function_node("merger", merger_contract, process_fn=merge_outputs) .set_entry_node("classifier") .add_route("classifier", FieldValueStrategy( field="category", mapping={"draft": "drafter", "review": "reviewer"}, )) .set_terminal("drafter") .set_terminal("reviewer") .with_retry_policy(RetryPolicy(max_retries=3)) .with_circuit_breaker(CircuitBreakerConfig(failure_threshold=5)) .build() )

__init__(name, *, version='1.0.0')

Initialize the builder with agent identity.

Parameters:

Name Type Description Default
name str

Unique name identifying this agent.

required
version str

Semantic version of the agent (e.g., "1.0.0").

'1.0.0'

add_aggregator_node(name, config, contract, retry_handler=None)

Register an AggregatorNode in the agent.

AggregatorNode collects responses from previously dispatched delegations using configurable resolution policies (all, any, majority). It requires a DelegationTransport to be configured via with_delegation_transport().

Parameters:

Name Type Description Default
name str

Unique name for this node within the agent.

required
config AggregatorNodeConfig

Configuration for the aggregator node.

required
contract NodeContract

Input/output contract for the node.

required
retry_handler Callable[[list[str]], Awaitable[list[str]]] | None

Async callable for retry_missing policy. Receives list of missing correlation IDs, returns new IDs to wait for. Required when config.on_timeout='retry_missing'.

None

Returns:

Type Description
AgentBuilder

Self for method chaining.

add_delegation_node(name, config, contract, input_mappers=None)

Register a DelegationNode in the agent.

DelegationNode delegates tasks to other agents and waits for their responses. It requires a DelegationTransport to be configured via with_delegation_transport().

Parameters:

Name Type Description Default
name str

Unique name for this node within the agent.

required
config DelegationNodeConfig

Configuration for the delegation node.

required
contract NodeContract

Input/output contract for the node.

required
input_mappers dict[str, Callable[[dict], dict]] | None

Optional dict of mapper functions. DelegationTargets reference these by name via input_mapping. Keys are mapper names, values are functions that transform input dicts.

None

Returns:

Type Description
AgentBuilder

Self for method chaining.

Example

builder.with_delegation_transport(my_transport).add_delegation_node( name="delegate_analysis", config=DelegationNodeConfig( name="delegate_analysis", targets=[ DelegationTarget( agent_name="analyzer", input_mapping="transform_for_analyzer", ), ], ), contract=analysis_contract, input_mappers={ "transform_for_analyzer": lambda d: {"text": d["content"]}, }, )

add_edge(from_node, to_node)

Add a static edge between two nodes.

The source node will always route to the target node, regardless of its output. For conditional routing, use add_route().

Parameters:

Name Type Description Default
from_node str

Source node name.

required
to_node str

Target node name.

required

Returns:

Type Description
AgentBuilder

Self for method chaining.

add_fan_out(from_node, targets, join_node, on_error='fail_all', timeout_seconds=300)

Declare a fan-out from a node to multiple parallel branches.

After from_node executes, all targets run concurrently. When every branch completes, execution continues at join_node which receives the original data plus a parallel_results dict keyed by target node name.

Parameters:

Name Type Description Default
from_node str

Node whose output triggers the fan-out.

required
targets list[str]

Nodes to execute in parallel (min 2).

required
join_node str

Node that collects the parallel results.

required
on_error str

"fail_all" (default) cancels remaining branches on first failure; "continue_partial" waits for all.

'fail_all'
timeout_seconds int

Maximum time for all branches.

300

Returns:

Type Description
AgentBuilder

Self for method chaining.

add_function_node(name, contract, process_fn=None, config=None, node=None)

Register a FunctionNode in the agent.

There are two ways to register a function node: 1. Provide a process_fn callable (for simple transformations) 2. Provide a FunctionNode instance (for subclassed nodes with complex logic)

Parameters:

Name Type Description Default
name str

Unique name for this node within the agent.

required
contract NodeContract

Input/output contract for the node.

required
process_fn Callable[[dict[str, Any], dict[str, Any] | None], dict[str, Any]] | None

Optional callable that executes the node logic.

None
config FunctionNodeConfig | None

Optional configuration for the node.

None
node FunctionNode | None

Optional pre-built FunctionNode instance (for subclassed nodes).

None

Returns:

Type Description
AgentBuilder

Self for method chaining.

add_llm_node(name, node_config=None, contract=None, tools=None, node=None)

Register an LLM node in the agent.

There are two ways to register an LLM node: 1. Provide node_config and contract (for standard LLMNode instances) 2. Provide a pre-built LLMNode instance (for subclassed nodes)

Parameters:

Name Type Description Default
name str

Unique name for this node within the agent.

required
node_config LLMNodeConfig | None

Configuration for the LLM node (required if node not provided).

None
contract NodeContract | None

Input/output contract for the node (required if node not provided).

None
tools list[BaseTool] | None

Optional list of tools available to this node.

None
node LLMNode | None

Optional pre-built LLMNode instance (for subclassed nodes).

None

Returns:

Type Description
AgentBuilder

Self for method chaining.

add_planner_node(name, contract, llm=None, system_prompt=None, description=None, available_nodes=None)

Register a PlannerNode in the agent.

The PlannerNode uses an LLM to generate an ExecutionPlan (DAG) based on the input and the list of available nodes. The available nodes are collected automatically from all registered nodes at build time, unless explicitly specified.

Parameters:

Name Type Description Default
name str

Unique name for this node within the agent.

required
contract NodeContract

Input/output contract for the planner node.

required
llm LLMCallable | None

Optional LLM callable. Falls back to the agent-level LLM.

None
system_prompt str | None

Custom system prompt for plan generation.

None
description str | None

Human-readable description.

None
available_nodes list[str] | None

Optional list of node names the planner can use. If None, all non-planner nodes are available.

None

Returns:

Type Description
AgentBuilder

Self for method chaining.

add_route(from_node, strategy)

Register a routing strategy for a node.

Parameters:

Name Type Description Default
from_node str

Name of the source node.

required
strategy RoutingStrategy

Routing strategy to determine the next node.

required

Returns:

Type Description
AgentBuilder

Self for method chaining.

build()

Build the Agent instance.

Validates all configuration and constructs the Agent. If validation fails, raises AgentBuildError with ALL errors found (not just the first one).

Returns:

Type Description
Agent

Configured Agent instance ready for execution.

Raises:

Type Description
AgentBuildError

If validation fails (contains all errors).

set_entry_node(name)

Set the entry point of the graph.

Parameters:

Name Type Description Default
name str

Name of the node where execution starts.

required

Returns:

Type Description
AgentBuilder

Self for method chaining.

set_plan_only(name)

Mark a node as plan-only.

Plan-only nodes exist solely for PlanExecutor to invoke. They do not participate in the main graph routing and are excluded from the validation rule 'every non-terminal node must have a route'.

Parameters:

Name Type Description Default
name str

Name of the node to mark as plan-only.

required

Returns:

Type Description
AgentBuilder

Self for method chaining.

set_terminal(name)

Mark a node as terminal (end of graph).

Parameters:

Name Type Description Default
name str

Name of the terminal node.

required

Returns:

Type Description
AgentBuilder

Self for method chaining.

with_agent_llm_config(config)

Set agent-level LLM configuration.

Parameters:

Name Type Description Default
config AgentLLMConfig

Agent-level LLM settings.

required

Returns:

Type Description
AgentBuilder

Self for method chaining.

with_circuit_breaker(config)

Set the circuit breaker configuration.

Parameters:

Name Type Description Default
config CircuitBreakerConfig

Circuit breaker configuration.

required

Returns:

Type Description
AgentBuilder

Self for method chaining.

with_context(ctx)

Set the static context shared with all nodes.

The context is read-only during execution. It contains information that doesn't change (e.g., user metadata, conversation history).

Parameters:

Name Type Description Default
ctx dict[str, Any]

Context dictionary.

required

Returns:

Type Description
AgentBuilder

Self for method chaining.

with_contract(input_schema, output_schema)

Set the agent's input/output contract.

Parameters:

Name Type Description Default
input_schema type[InputContract]

Pydantic model class for agent input validation.

required
output_schema type[OutputContract]

Pydantic model class for agent output validation.

required

Returns:

Type Description
AgentBuilder

Self for method chaining.

with_conversation_state(state)

Set the conversation state store for multi-turn flow persistence.

The ConversationState is injected into all FunctionNodes at build time, enabling them to persist and retrieve temporary state across multiple agent executions within the same conversation.

Parameters:

Name Type Description Default
state Any

ConversationState adapter wrapping a ConversationStateProvider implementation.

required

Returns:

Type Description
AgentBuilder

Self for method chaining.

Example

from sofias_sdk_lite.state import ( ConversationState, InMemoryStateProvider, )

builder.with_conversation_state( ConversationState(InMemoryStateProvider()) )

with_default_route_to_terminal()

Auto-route unrouted nodes to the terminal node.

When enabled, nodes without an explicit route (via add_route or add_edge) will automatically route to the terminal node at build time. This eliminates boilerplate for graphs where most nodes converge to a single terminal.

Returns:

Type Description
AgentBuilder

Self for method chaining.

with_delegation_transport(transport)

Set the transport for delegation nodes.

The transport is required if the agent contains any DelegationNodes. It provides the send/receive capabilities for delegation communication.

Parameters:

Name Type Description Default
transport DelegationTransport

Transport implementing the DelegationTransport protocol.

required

Returns:

Type Description
AgentBuilder

Self for method chaining.

with_description(description)

Set the agent's description.

Parameters:

Name Type Description Default
description str

Human-readable description of what this agent does.

required

Returns:

Type Description
AgentBuilder

Self for method chaining.

with_fallback(node, fallback_node)

Register a fallback node for when a node fails.

Parameters:

Name Type Description Default
node str

Name of the node that may fail.

required
fallback_node str

Name of the node to execute on failure.

required

Returns:

Type Description
AgentBuilder

Self for method chaining.

with_llm(llm)

Set the LLM callable for all LLM nodes.

Optional. When omitted, build() uses the LLM installed by sofias_sdk_lite.llm.default_llm (what AgentRunner does with the per-task settings) or, failing that, one built from the SOFIAS_LLM_* environment variables via sofias_sdk_lite.llm.create_llm.

Parameters:

Name Type Description Default
llm LLMCallable

LLM implementation for making completions.

required

Returns:

Type Description
AgentBuilder

Self for method chaining.

with_middleware(middleware)

Add a middleware to the agent execution pipeline.

Middleware hooks are called before and after each node execution. Multiple middleware can be added — they execute in registration order.

Parameters:

Name Type Description Default
middleware Any

Object implementing before_node/after_node/on_error.

required

Returns:

Type Description
AgentBuilder

Self for method chaining.

with_plan_failure_policy(default_policy='fail_plan', overrides=None, timeout_seconds=600)

Configure failure handling for plan execution.

Parameters:

Name Type Description Default
default_policy Literal['fail_plan', 'skip_dependents', 'continue_partial']

Default policy when a plan step fails. One of "fail_plan", "skip_dependents", "continue_partial".

'fail_plan'
overrides dict[str, str] | None

Per-node policy overrides (node_name -> policy).

None
timeout_seconds int

Global timeout for plan execution.

600

Returns:

Type Description
AgentBuilder

Self for method chaining.

with_response_workflow(workflow)

Set the response workflow for delivering agent output.

The workflow is invoked after execution completes (success or error) to deliver the response through the appropriate channel (chat stream, reply queue, webhook, etc.).

Parameters:

Name Type Description Default
workflow ResponseWorkflow

Implementation of the ResponseWorkflow protocol.

required

Returns:

Type Description
AgentBuilder

Self for method chaining.

with_retry_policy(policy)

Set the default retry policy for all nodes.

Parameters:

Name Type Description Default
policy RetryPolicy

Default retry policy.

required

Returns:

Type Description
AgentBuilder

Self for method chaining.

with_retry_policy_for_node(name, policy)

Set a node-specific retry policy.

Parameters:

Name Type Description Default
name str

Name of the node.

required
policy RetryPolicy

Retry policy for this specific node.

required

Returns:

Type Description
AgentBuilder

Self for method chaining.

with_sdk_config(config)

Set the SDK configuration.

Parameters:

Name Type Description Default
config SDKConfig

Global SDK configuration.

required

Returns:

Type Description
AgentBuilder

Self for method chaining.

with_settings_class(cls)

Set the runtime settings class.

Parameters:

Name Type Description Default
cls type[AgentSettings]

The AgentSettings subclass for runtime configuration.

required

Returns:

Type Description
AgentBuilder

Self for method chaining.

with_system_prompt(prompt)

Set an agent-level system prompt propagated to all LLM nodes.

Node-level system prompts (set via LLMNodeConfig.system_prompt) take precedence over this agent-level prompt.

Parameters:

Name Type Description Default
prompt str

System prompt to inject into every LLM call.

required

Returns:

Type Description
AgentBuilder

Self for method chaining.

with_tool_registry(registry)

Set the tool registry for dynamic tool discovery.

The ToolRegistry is injected into LLM nodes built from configs, enabling them to discover and execute tools from external providers in addition to static BaseTool instances.

Pre-built nodes (registered via node=) are not affected.

Parameters:

Name Type Description Default
registry ToolRegistry

ToolRegistry wrapping a ToolProvider implementation.

required

Returns:

Type Description
AgentBuilder

Self for method chaining.

sofias_sdk_lite.Agent

The complete agent graph ready for execution.

The Agent is a pure orchestrator. It does not contain business logic. Its responsibilities are: 1. Receive input and resolve runtime configuration 2. Traverse the graph executing nodes according to routing decisions 3. Manage errors through the error handler 4. Produce validated output

The agent does NOT execute LLM calls or tools directly. It delegates to Node instances which handle their own execution.

Example

Agent is typically built via AgentBuilder

agent = ( AgentBuilder("my_agent", version="1.0.0") .with_settings_class(MySettings) .with_contract(input_schema=MyInput, output_schema=MyOutput) .add_node("start", node_config, contract) .set_entry_node("start") .set_terminal("start") .build() )

Execute the agent

response = await agent.execute( AgentMessage( content=MyInput(query="Hello"), runtime_config={"model": "gpt-4o"}, ) ) assert response.status == ResponseStatus.SUCCESS

context property

The read-only context shared with all nodes.

description property

The agent's description.

entry_node property

The name of the entry node.

name property

The agent's name.

version property

The agent's version.

__init__(*, name, version, description, nodes, router, error_handler, contract, entry_node, settings_resolver, context=None, response_workflow=None, plan_executor=None, middleware=None)

Initialize the agent.

This constructor is typically called by AgentBuilder.build(), not directly by users.

Parameters:

Name Type Description Default
name str

Unique name identifying this agent.

required
version str

Semantic version of the agent.

required
description str | None

Human-readable description.

required
nodes dict[str, BaseNode]

Dictionary mapping node names to Node instances.

required
router Router

Router managing graph edges and routing decisions.

required
error_handler ErrorHandler

ErrorHandler for retry, fallback, circuit breaker.

required
contract AgentContract

Contract defining agent input/output schemas.

required
entry_node str

Name of the node where execution starts.

required
settings_resolver SettingsResolver[AgentSettings]

Resolver for runtime settings.

required
context dict[str, Any] | None

Read-only context shared with all nodes.

None
response_workflow ResponseWorkflow | None

Optional workflow for delivering the response.

None
plan_executor PlanExecutor | None

Optional executor for plans produced by PlannerNodes.

None
middleware list[Any] | None

Optional list of AgentMiddleware for before/after node hooks.

None

__repr__()

Return a string representation of the agent.

execute(message) async

Execute the agent graph with the given input message.

This is the main entry point for agent execution. The flow is: 1. Resolve runtime config from message and set in context 2. Validate message content against agent contract 3. Execute graph nodes following routing decisions 4. Validate final output against agent contract 5. Build AgentResponse with execution metadata 6. Clean up context and return result

For controlled business errors (validation failures, routing errors, node execution errors handled by the error handler), returns an AgentResponse with status=ERROR and error_message describing the problem. Infrastructure failures (unexpected exceptions not derived from AgentSDKError) are re-raised.

Parameters:

Name Type Description Default
message AgentMessage

Typed input envelope containing domain data and metadata.

required

Returns:

Type Description
AgentResponse

AgentResponse with validated output and execution metadata.

Raises:

Type Description
Exception

Only for infrastructure failures not covered by AgentSDKError (unexpected crashes, connectivity issues, etc.).

get_graph_info()

Get information about the graph structure.

Useful for debugging, logging, and observability.

Returns:

Type Description
dict[str, Any]

Dictionary containing:

dict[str, Any]
  • name: Agent name
dict[str, Any]
  • version: Agent version
dict[str, Any]
  • description: Agent description
dict[str, Any]
  • entry_node: Name of the entry node
dict[str, Any]
  • nodes: List of all node names
dict[str, Any]
  • terminals: Set of terminal node names
dict[str, Any]
  • routes: Dictionary of source node to target nodes
dict[str, Any]
  • fallbacks: Dictionary of node to fallback node mappings

stream_execute(message) async

Execute the agent graph with streaming.

Same flow as execute() but yields StreamEvent instances in real time. For LLMNodes, text chunks are yielded as they arrive from the LLM. For non-streaming nodes, only NodeStart/NodeComplete events are emitted.

If a StreamingResponseWorkflow is configured, on_stream_event() is called for each event. send_response() is called at the end.

Parameters:

Name Type Description Default
message AgentMessage

Typed input envelope containing domain data and metadata.

required

Yields:

Type Description
AsyncGenerator[StreamEvent, None]

StreamEvent instances for real-time visibility into execution.

sofias_sdk_lite.AgentMessage

Bases: StrictContract, Generic[T]

Typed input envelope for agent execution.

Generic over the agent's InputContract subclass, so each agent can declare its expected input type:

AgentMessage[RealEstateInput]
AgentMessage[SupportTicketInput]

The content field is validated against the concrete InputContract, while the envelope fields carry transport-level metadata.

Attributes:

Name Type Description
content T

Domain input data, validated against the agent's InputContract.

conversation_id str | None

Conversation thread identifier.

trace_id str | None

Distributed trace identifier for observability.

runtime_config dict[str, Any] | None

Runtime configuration overrides (model, temperature, etc.).

metadata dict[str, Any]

Additional transport-level metadata.

sofias_sdk_lite.AgentResponse

Bases: StrictContract

Typed output envelope from agent execution.

Wraps the domain output (validated against AgentContract.output_schema) with execution metadata. Returned by Agent.execute().

For controlled business errors, status is ERROR and error_message describes the problem. Infrastructure failures (graph compilation, unexpected node explosions) still raise exceptions.

Attributes:

Name Type Description
content dict[str, Any]

Domain output data, validated against the agent's OutputContract.

status ResponseStatus

Whether the execution succeeded or hit a business error.

agent_name str

Name of the agent that produced this response.

agent_version str

Version of the agent that produced this response.

execution_path list[str]

Ordered list of nodes executed.

error_message str | None

Error description when status is ERROR.

execution_time_ms float | None

Wall-clock execution time in milliseconds.

metadata dict[str, Any]

Response metadata (run_id, trace_id, timing, etc.).

Contracts

sofias_sdk_lite.InputContract

Bases: StrictContract

Base class for node input contracts.

Agents should inherit from this class to define their input schemas. Inherits strict validation from StrictContract.

sofias_sdk_lite.OutputContract

Bases: StrictContract

Base class for node output contracts.

Agents should inherit from this class to define their output schemas. Inherits strict validation from StrictContract.

sofias_sdk_lite.NodeContract

Bases: Generic[InputT, OutputT]

Encapsulates the input/output contract pair for a node.

Provides validation methods to verify data against the defined schemas.

__init__(input_schema, output_schema)

Initialize the node contract.

Parameters:

Name Type Description Default
input_schema type[InputT]

The Pydantic model class for input validation.

required
output_schema type[OutputT]

The Pydantic model class for output validation.

required

validate_input(data)

Validate input data against the input schema.

Parameters:

Name Type Description Default
data dict[str, Any]

The data dictionary to validate.

required

Returns:

Type Description
InputT

The validated input model instance.

Raises:

Type Description
InputValidationError

If validation fails.

validate_output(data)

Validate output data against the output schema.

Parameters:

Name Type Description Default
data dict[str, Any]

The data dictionary to validate.

required

Returns:

Type Description
OutputT

The validated output model instance.

Raises:

Type Description
OutputValidationError

If validation fails.

sofias_sdk_lite.StrictContract

Bases: BaseModel

Base class for all strictly-validated contracts in the SDK.

This class provides the common validation configuration used by both: - Node contracts (InputContract, OutputContract) for intra-graph data flow - Delegation contracts (DelegationRequest, DelegationResponse) for inter-agent protocol

Configuration
  • extra="ignore": Silently drop any fields not defined in the schema (not "forbid" — unknown fields do not raise; they are discarded)
  • validate_default=True: Validate default values
  • strict=True: Use strict type coercion (no automatic conversions)

Nodes

sofias_sdk_lite.BaseNode

Bases: ABC

Abstract base class for all node types in the Agent SDK.

BaseNode defines the common interface that all nodes must implement. It provides a template method pattern for execution: the public execute() method handles input/output validation, while subclasses implement the _run() method with their specific logic.

This design allows the router, error handler, and other components to work with any node type uniformly through the execute() interface.

Subclasses
  • LLMNode: Nodes that invoke an LLM with optional tool loop.
  • FunctionNode: Nodes that execute deterministic Python logic.
Example

class MyCustomNode(BaseNode): async def _run(self, input_data: dict, context: dict | None = None) -> dict: # Custom logic here return {"result": "processed"}

contract property

The node's input/output contract.

description property

Human-readable description of the node.

name property

The node's unique name.

node_type property

The type of this node. Subclasses override to declare their type.

__init__(name, contract, description=None)

Initialize the base node.

Parameters:

Name Type Description Default
name str

Unique name identifying this node within the agent.

required
contract NodeContract

Input/output contract for validation.

required
description str | None

Human-readable description of what this node does.

None

execute(input_data, context=None) async

Execute the node with the given input.

This is the main entry point for node execution. It implements a template method pattern: 1. Validates input against the contract 2. Calls _run() (implemented by subclass) 3. Validates output against the contract 4. Returns the validated result

Parameters:

Name Type Description Default
input_data dict[str, Any]

Input dictionary to process.

required
context dict[str, Any] | None

Optional context dictionary with additional information.

None

Returns:

Type Description
dict[str, Any]

Output dictionary conforming to the output contract.

Raises:

Type Description
InputValidationError

If input validation fails.

OutputValidationError

If output validation fails.

NodeExecutionError

If execution fails in _run().

stream_execute(input_data, context=None) async

Execute the node with streaming.

For LLMNodes that implement _stream_run(), yields streaming events as they occur. For other node types, falls back to _run() and yields only NodeStartEvent and NodeCompleteEvent.

Parameters:

Name Type Description Default
input_data dict[str, Any]

Input dictionary to process.

required
context dict[str, Any] | None

Optional context dictionary with additional information.

None

Yields:

Type Description
AsyncGenerator[StreamEvent, None]

StreamEvent instances for each notable event during execution.

sofias_sdk_lite.FunctionNode

Bases: BaseNode

A node that executes deterministic Python logic without an LLM.

FunctionNode is designed for pure data transformations, filtering, validation, merging outputs, or any logic that doesn't require an LLM. It inherits input/output validation from BaseNode.

There are two ways to use FunctionNode:

  1. Pass a callable to the constructor (for simple transformations):

    def format_data(input_data: dict, context: dict | None = None) -> dict: return {"formatted": input_data["raw"].upper()}

    node = FunctionNode( name="formatter", contract=contract, process_fn=format_data, )

  2. Subclass for complex logic:

    class DataMerger(FunctionNode): def init(self, name: str, contract: NodeContract): super().init(name=name, contract=contract) self._cache = {}

    async def process(self, input_data: dict, context: dict | None = None) -> dict:
        # Complex logic with state
        return {"merged": ...}
    
Example

node = FunctionNode( name="transformer", contract=NodeContract(input_schema=RawInput, output_schema=FormattedOutput), process_fn=lambda data, ctx: {"value": data["input"] * 2}, ) result = await node.execute({"input": 5})

result: {"value": 10}

__init__(name, contract, process_fn=None, description=None)

Initialize the function node.

Parameters:

Name Type Description Default
name str

Unique name identifying this node within the agent.

required
contract NodeContract

Input/output contract for validation.

required
process_fn ProcessFn | None

Optional function that executes the node logic. If not provided, subclass must override process().

None
description str | None

Human-readable description of what this node does.

None

delete_state(key) async

Delete persisted state for the current conversation.

Reads the conversation_id automatically from ExecutionContext. Requires that AgentBuilder.with_conversation_state() was called.

Parameters:

Name Type Description Default
key str

State key to delete.

required

Raises:

Type Description
NodeExecutionError

If no conversation state is configured.

get_state(key) async

Retrieve persisted state for the current conversation.

Reads the conversation_id automatically from ExecutionContext. Requires that AgentBuilder.with_conversation_state() was called.

Parameters:

Name Type Description Default
key str

State key (e.g. "validation:doc123").

required

Returns:

Type Description
dict[str, Any] | None

The stored dict, or None if not found.

Raises:

Type Description
NodeExecutionError

If no conversation state is configured.

process(input_data, context=None) async

Execute the node's processing logic.

Override this method in subclasses to implement custom logic. This method is called when no process_fn is provided to the constructor.

Parameters:

Name Type Description Default
input_data dict[str, Any]

Validated input dictionary.

required
context dict[str, Any] | None

Optional context dictionary.

None

Returns:

Type Description
dict[str, Any]

Output dictionary.

Raises:

Type Description
NotImplementedError

If neither process_fn was provided nor this method was overridden in a subclass.

set_state(key, value) async

Persist state for the current conversation.

Reads the conversation_id automatically from ExecutionContext. Requires that AgentBuilder.with_conversation_state() was called.

Parameters:

Name Type Description Default
key str

State key (e.g. "validation:doc123").

required
value dict[str, Any]

Dict to persist. Overwrites any existing value.

required

Raises:

Type Description
NodeExecutionError

If no conversation state is configured.

sofias_sdk_lite.LLMNode

Bases: BaseNode

A node that uses an LLM for execution.

LLMNode encapsulates: - Prompt assembly with template resolution - LLM invocation with optional tool loop - Configurable hooks for observability and intervention

The node inherits input/output validation from BaseNode. Internally, it may execute multiple LLM <-> tool cycles before producing the final result.

Example

node = LLMNode( config=node_config, contract=NodeContract(input_schema=MyInput, output_schema=MyOutput), llm=my_llm_callable, tools=[tool1, tool2], ) result = await node.execute({"query": "Hello"})

config property

The node's configuration.

__init__(config, contract, llm, tools=None, prompt_assembler=None, *, system_prompt=None, settings_resolver=None, tool_registry=None, on_tool_call=None, on_tool_result=None, on_llm_response=None, should_continue=None, on_chunk=None)

Initialize the LLM node.

Parameters:

Name Type Description Default
config LLMNodeConfig

Node configuration (name, prompt, LLM settings, etc.).

required
contract NodeContract

Input/output contract for validation.

required
llm LLMCallable

LLM callable for making completions.

required
tools list[BaseTool] | None

Optional list of BaseTool instances available to this node.

None
prompt_assembler PromptAssembler | None

Custom prompt assembler. If None, one is created from config.

None
settings_resolver SettingsResolver[AgentSettings] | None

Optional resolver for runtime settings access.

None
tool_registry ToolRegistry | None

Optional registry for dynamic tool provisioning. Any backend may implement this protocol; the node has no knowledge of or dependency on a specific tool provider.

None
on_tool_call OnToolCallHook | None

Hook called before each tool execution.

None
on_tool_result OnToolResultHook | None

Hook called after each tool execution.

None
on_llm_response OnLLMResponseHook | None

Hook called after each LLM response.

None
should_continue ShouldContinueHook | None

Hook to control loop continuation.

None
on_chunk OnChunkHook | None

Hook called for each text chunk during streaming.

None

sofias_sdk_lite.DelegationNode

Bases: BaseNode

Node that delegates tasks to other agents.

DelegationNode sends tasks to one or more target agents and waits for their responses. It supports three execution modes:

  • Sequential: Delegates to targets one at a time, in order.
  • Parallel: Delegates to all targets simultaneously and waits for all responses.
  • DAG: Executes steps respecting dependencies, with configurable failure policies.

The node uses a DelegationTransport for communication, which handles the underlying message passing (RabbitMQ, HTTP, etc.).

Example (static targets): config = DelegationNodeConfig( name="delegate_to_summarizer", targets=[DelegationTarget(agent_name="summarizer")], )

node = DelegationNode(
    config=config,
    contract=my_contract,
    transport=my_transport,
)

result = await node.execute({"text": "Long document..."})
# result = {"summarizer": {"summary": "..."}}

Example (DAG with dynamic targets): config = DelegationNodeConfig( name="dag_executor", execution_mode="dag", dynamic_targets=True, )

plan = {
    "steps": [
        {"id": "s1", "agent_name": "search", "input": {"q": "A"}},
        {"id": "s2", "agent_name": "search", "input": {"q": "B"}},
        {"id": "s3", "agent_name": "synth", "input": {}, "depends_on": ["s1", "s2"]},
    ]
}

result = await node.execute(plan)
# result = {"results": {...}, "failed": {...}, "skipped": [...], "all_completed": bool}

config property

The node's configuration.

__init__(config, contract, transport, input_mappers=None, source_agent_name=None, *, settings_resolver=None)

Initialize the DelegationNode.

Parameters:

Name Type Description Default
config DelegationNodeConfig

Configuration for the delegation node.

required
contract NodeContract

Input/output contract for validation.

required
transport DelegationTransport

Transport for sending requests and receiving responses.

required
input_mappers dict[str, InputMapper] | None

Optional dict mapping mapper names to functions. DelegationTargets reference these by name via input_mapping.

None
source_agent_name str | None

Name of the agent that owns this node. Used in metadata for tracing.

None
settings_resolver SettingsResolver | None

Optional resolver for runtime settings. When present, agent_timeout from AgentSettings overrides per-target timeout_seconds.

None

sofias_sdk_lite.AggregatorNode

Bases: BaseNode

Node that collects responses from previously dispatched delegations.

AggregatorNode implements a fan-in pattern: given a list of correlation IDs (from prior DelegationNode sends or any other source), it waits for responses via the DelegationTransport and resolves according to a configurable policy.

Resolution policies: - 'all': Wait until every expected response arrives. - 'any': Resolve as soon as at least one response arrives. - 'majority': Resolve when more than floor(expected * threshold) responses arrive.

Retry support: When on_timeout='retry_missing', the node calls the retry_handler with the list of unresponsive correlation IDs. The handler re-dispatches work and returns the new IDs to wait for. Each retry gets its own timeout window controlled by retry_timeout_seconds (or timeout_seconds if not set).

Example

config = AggregatorNodeConfig( name="collect_results", resolution_policy="all", timeout_seconds=60.0, )

node = AggregatorNode( config=config, contract=my_contract, transport=my_transport, )

result = await node.execute({"pending_ids": ["corr-1", "corr-2"]})

result = {

"results": {"corr-1": {...}, "corr-2": {...}},

"pending": [],

"is_partial": False,

"total_expected": 2,

"total_received": 2,

}

config property

The node's configuration.

__init__(config, contract, transport, retry_handler=None)

Initialize the AggregatorNode.

Parameters:

Name Type Description Default
config AggregatorNodeConfig

Configuration for the aggregator node.

required
contract NodeContract

Input/output contract for validation.

required
transport DelegationTransport

Transport for receiving responses.

required
retry_handler RetryHandler | None

Async callable for retry_missing policy. Receives list of missing correlation IDs, returns new IDs to wait for. Required when config.on_timeout='retry_missing'.

None

sofias_sdk_lite.PlannerNode

Bases: BaseNode

Node that generates an ExecutionPlan via LLM.

The PlannerNode: 1. Receives input data (task description, context, or re-plan info) 2. Calls an LLM with a prompt describing available nodes 3. Parses the LLM response as an ExecutionPlan 4. Returns the plan as its output

The Agent detects this node by type (node_type == PLANNER) and routes the output to PlanExecutor for execution.

Example

planner = PlannerNode( name="planner", contract=planner_contract, llm=my_llm, available_nodes=[ NodeDescriptor(name="search", node_type=NodeType.FUNCTION, description="Search docs"), NodeDescriptor(name="summarize", node_type=NodeType.LLM, description="Summarize text"), ], )

__init__(name, contract, llm, available_nodes, system_prompt=None, description=None)

Initialize the PlannerNode.

Parameters:

Name Type Description Default
name str

Unique name identifying this node.

required
contract NodeContract

Input/output contract for validation.

required
llm LLMCallable

LLM callable for generating plans.

required
available_nodes list[NodeDescriptor]

List of nodes the planner can include in plans.

required
system_prompt str | None

Custom system prompt. If None, a default is built.

None
description str | None

Human-readable description.

None

Errors

sofias_sdk_lite.ErrorHandler

Manages error handling for graph node execution.

Wraps node execution with: 1. Circuit breaker: prevents repeatedly calling a failing node 2. Retry: attempts to recover from transient failures 3. Fallback: routes to alternative nodes when recovery fails

The handler does NOT execute nodes directly. It wraps the execution callable provided by the agent.

Example

handler = ErrorHandler(config=config, store=store) result = await handler.handle_node_execution( node_name="processor", execute_fn=node.execute, input_data={"query": "hello"}, )

handle_node_execution(node_name, execute_fn, input_data) async

Execute a node with full error handling.

Flow: 1. Check circuit breaker state (OPEN -> fallback directly) 2. Execute with retry logic 3. On success: record success, return result 4. On failure after retries: record failure, try fallback, or raise

register_fallback_executor(node_name, executor)

Register a fallback executor for a node.

The fallback executor receives the original input plus error context.

sofias_sdk_lite.RetryPolicy

Bases: BaseModel

Configuration for retry behavior.

backoff_strategy = BackoffStrategy.EXPONENTIAL class-attribute instance-attribute

Strategy for calculating delay between retries.

delay = 1.0 class-attribute instance-attribute

Base delay in seconds between retries.

max_delay = 30.0 class-attribute instance-attribute

Maximum delay in seconds (caps exponential/linear growth).

max_retries = 3 class-attribute instance-attribute

Maximum number of retry attempts.

retryable_exceptions = Field(default=(NodeExecutionError, ToolExecutionError, OutputValidationError)) class-attribute instance-attribute

Exception types that should trigger retries.

OutputValidationError is retried by default because LLM providers may return error payloads (e.g. content-filter messages) inside a successful HTTP response; the node parses them as valid JSON but they fail output-contract validation. Retrying gives the LLM another chance to produce a conforming response.

sofias_sdk_lite.CircuitBreakerConfig

Bases: BaseModel

Configuration for circuit breaker behavior.

failure_threshold = 5 class-attribute instance-attribute

Number of consecutive failures to trip the circuit.

recovery_timeout = 30.0 class-attribute instance-attribute

Seconds to wait in OPEN state before trying HALF_OPEN.

Config

sofias_sdk_lite.BaseAgentSettings

Bases: AgentSettings

Runtime settings shared by all agents.

Carries the LLM gateway keys the Sofias platform injects per task (model_name, router_url, router_api_key). AgentRunner reads them through sofias_sdk_lite.llm.create_llm and installs the resulting client as the default for AgentBuilder.build(), so agent code only declares LLM nodes and never constructs a client.

All three are optional here so function-only agents validate with an empty payload; has_llm_gateway tells whether a model was supplied.

Subclass to add agent-specific fields — feature flags, downstream service URLs, or anything else a given agent depends on.

CREDENTIAL_FIELDS = frozenset({'router_api_key'}) class-attribute

Names of fields that hold credentials (API keys, tokens, secrets).

Extend on your subclass (BaseAgentSettings.CREDENTIAL_FIELDS | {...}), typing the fields as pydantic.SecretStr so repr/model_dump never leak the raw value. The SDK excludes these fields wherever settings are dumped into a dict that reaches an LLM prompt (LLMNode template variables), so a {api_key} placeholder can never inline a secret. Read them at the point of use with .get_secret_value().

Example::

class MySettings(BaseAgentSettings):
    CREDENTIAL_FIELDS = BaseAgentSettings.CREDENTIAL_FIELDS | {"search_api_key"}

    search_api_key: SecretStr = SecretStr("")

has_llm_gateway property

Whether these settings name a model (and therefore can build an LLM).

model_name = '' class-attribute instance-attribute

Model identifier or gateway alias. Empty means "no LLM configured here".

router_api_key = SecretStr('') class-attribute instance-attribute

Gateway bearer token. Empty sends unauthenticated requests.

router_provider = '' class-attribute instance-attribute

Adapter flavour: "sofias" (default when empty) or an OpenAI-compatible label.

router_url = '' class-attribute instance-attribute

OpenAI-compatible gateway root, injected by the platform. Required to build an LLM.

sofias_sdk_lite.ConfigSource

Bases: Protocol

Supplies runtime configuration for an agent, keyed by agent name and conversation.

Implementations may fetch configuration from anywhere: a database, a remote service, a local file, environment variables, or an in-memory dict. The only contract is the fetch coroutine's signature and return type.

Example

class DatabaseConfigSource: async def fetch( self, agent_name: str, conversation_id: str | None = None ) -> dict[str, Any]: return await my_db.get_agent_config(agent_name)

fetch(agent_name, conversation_id=None) async

Fetch the runtime configuration payload for an agent.

Parameters:

Name Type Description Default
agent_name str

The name/identity of the agent requesting configuration.

required
conversation_id str | None

Optional conversation identifier, for sources that vary configuration per-conversation (e.g. A/B tests, per-tenant overrides).

None

Returns:

Type Description
dict[str, Any]

A raw configuration dictionary, typically passed to a

dict[str, Any]

sofias_sdk_lite.config.runtime_config.SettingsResolver for

dict[str, Any]

validation into an AgentSettings subclass.

sofias_sdk_lite.StaticConfigSource

A ConfigSource that always returns the same static dict.

Useful for local development, tests, and any deployment where a single fixed configuration applies to every agent and conversation.

Example

source = StaticConfigSource({"model_name": "gpt-4o-mini", "api_key": "..."}) config = await source.fetch("my_agent")

__init__(config)

Initialize with the static configuration to always return.

Parameters:

Name Type Description Default
config dict[str, Any]

The configuration dictionary to return from every fetch call. Stored by reference; a shallow copy is returned on each call so callers cannot mutate the stored configuration through the returned dict.

required

fetch(agent_name, conversation_id=None) async

Return a copy of the static configuration.

Parameters:

Name Type Description Default
agent_name str

Unused; accepted to satisfy the ConfigSource protocol.

required
conversation_id str | None

Unused; accepted to satisfy the ConfigSource protocol.

None

Returns:

Type Description
dict[str, Any]

A shallow copy of the configured dict.

sofias_sdk_lite.EnvConfigSource

A ConfigSource that reads a JSON blob from a single environment variable.

Useful for container-based deployments that inject the full configuration payload as one environment variable containing a JSON object, e.g.:

export SOFIAS_AGENT_CONFIG='{"model_name": "gpt-4o-mini", "api_key": "sk-..."}'
Example

source = EnvConfigSource() # reads SOFIAS_AGENT_CONFIG config = await source.fetch("my_agent")

__init__(env_var='SOFIAS_AGENT_CONFIG')

Initialize with the name of the environment variable to read.

Parameters:

Name Type Description Default
env_var str

Name of the environment variable holding a JSON object. Defaults to "SOFIAS_AGENT_CONFIG".

'SOFIAS_AGENT_CONFIG'

fetch(agent_name, conversation_id=None) async

Read and parse the configured environment variable.

Parameters:

Name Type Description Default
agent_name str

Unused for lookup, but included in error context if the environment variable is missing or invalid.

required
conversation_id str | None

Unused; accepted to satisfy the ConfigSource protocol.

None

Returns:

Type Description
dict[str, Any]

The parsed JSON object as a dict.

Raises:

Type Description
ConfigSourceError

If the environment variable is unset, empty, not valid JSON, or does not decode to a JSON object (dict).

Messaging & RabbitMQ

sofias_sdk_lite.AgentTaskMessage

Bases: _BaseTaskMessage

Canonical schema for the message a conversational agent receives via RabbitMQ.

Used by both the publisher (platform -> agent) and the delegation transport (agent -> agent). The receiver does not need to distinguish the origin.

sofias_sdk_lite.StreamFragment

Bases: BaseModel

A single fragment published to a RabbitMQ stream during streaming.

Every agent that streams responses to the frontend must publish messages conforming to this schema so the frontend can consume them uniformly regardless of the source agent.

sofias_sdk_lite.RabbitMQConfig dataclass

Configuration for RabbitMQ connections (AMQP and Streams).

Only connection-level settings belong here. Topology parameters (queue_name, exchange_name, routing_key, stream_name) belong in the constructors of Consumer / Publisher.

amqp_url property

Build an AMQP URL from the individual fields.

sofias_sdk_lite.RabbitMQClient

Manages the connection and channel lifecycle for RabbitMQ.

Accepts a :class:RabbitMQConfig and supports usage as an async context manager::

async with RabbitMQClient(config) as client:
    ...

config property

The configuration used to create this client.

Runner & workflows

sofias_sdk_lite.AgentRunner

Bases: ABC

Base class for running an agent against RabbitMQ.

Handles the RabbitMQ connection, consumer loop, per-message config resolution, conversation history, response delivery, and graceful shutdown on SIGTERM/SIGINT. Subclasses only need to describe how to build the agent and how to shape its input.

Each turn runs inside an agent.handle span (when OpenTelemetry is installed) and the inbound W3C traceparent is echoed on every response fragment, so a chat turn is one trace end to end.

History handed to prepare_input is windowed by max_messages and max_context_tokens; a SummarizingHistoryProvider gets maybe_summarize called in the background after each response.

auto_summarize = True class-attribute instance-attribute

After each delivered response, call maybe_summarize on the history provider (if it implements SummarizingHistoryProvider) as a background task. Never blocks the next message; drained on graceful shutdown.

max_context_tokens = 65536 class-attribute instance-attribute

Estimated token cap for the history handed to prepare_input; oldest turns are dropped until the remainder fits (the latest turn always stays).

max_messages = 100 class-attribute instance-attribute

Most recent turns handed to prepare_input (oldest dropped first).

after_execution(task, response) async

Hook called after successful agent execution. Override for custom logic.

build_agent(settings, workflow) abstractmethod

Build the Agent instance for a single message.

Parameters:

Name Type Description Default
settings BaseAgentSettings

Resolved settings for this invocation (from ConfigSource.fetch() merged with task.agent_configuration, validated against settings_class).

required
workflow Any

The ResponseWorkflow to attach to the agent (a ChatResponseWorkflow for conversational tasks, or a BaseDelegationResponseWorkflow when task.reply_to is set).

required

create_llm(settings)

Build (or reuse) the LLM client for this task's settings.

The default reads model_name / router_url / router_api_key from settings through sofias_sdk_lite.llm.create_llm and caches one client per distinct gateway configuration. Returning None means "no LLM from settings": AgentBuilder.build() then falls back to the SOFIAS_LLM_* environment variables.

Override to plug in a different client (a fake in tests, a provider the bundled adapters do not cover, per-tenant routing, ...).

create_response_workflow(client, task, *, traceparent=None)

Create the response workflow for a single message.

Returns a ChatResponseWorkflow publishing to RunnerConfig.response_stream. Override to use a different transport or workflow (e.g. a BaseStreamingResponseWorkflow subclass for token-by-token streaming).

Parameters:

Name Type Description Default
client RabbitMQClient

The connected RabbitMQ client.

required
task AgentTaskMessage

The incoming task message.

required
traceparent str | None

W3C traceparent received on the inbound message. Echoing it onto the response fragments keeps a chat turn a single trace across the whole pipeline. Forward it to your workflow when overriding.

None

prepare_input(task, history, role) abstractmethod

Build the AgentMessage to pass to Agent.execute().

Parameters:

Name Type Description Default
task AgentTaskMessage

The incoming task message.

required
history list[Message]

Prior messages for this conversation, oldest first.

required
role str

The resolved role for this request (see resolve_role).

required

resolve_role(payload)

Map an incoming user_role through RunnerConfig.role_map.

Override for custom resolution logic. Unmapped roles pass through unchanged rather than raising or defaulting silently.

run()

Entry point: connect, consume, run until shutdown.

sofias_sdk_lite.RunnerConfig

Bases: BaseModel

Configuration for running an agent against RabbitMQ.

Every field that determines wire-level behavior (queue name, response stream, role mapping) is explicit here — the runner has no built-in naming convention linking an agent's identity to its queue or stream.

agent_name instance-attribute

The agent's identity, independent of the queue it consumes from.

Used for StreamFragment.agent and passed to ConfigSource.fetch().

queue instance-attribute

The RabbitMQ queue this runner consumes tasks from.

rabbitmq instance-attribute

Connection settings for the RabbitMQ broker.

response_stream = 'agent.responses' class-attribute instance-attribute

Stream name for ChatResponseWorkflow to publish to.

This is a neutral default, not a required naming convention — override it to match whatever your deployment uses.

role_map = Field(default_factory=dict) class-attribute instance-attribute

Optional mapping from an incoming user_role value to your agent's own role vocabulary. Unmapped roles pass through unchanged.

sofias_sdk_lite.HistoryProvider

Bases: Protocol

Protocol defining the interface for conversation history backends.

Unlike MemoryProvider (long-term summaries/facts/preferences) or ConversationStateProvider (ephemeral flow-scoped data), a HistoryProvider deals only with the ordered list of turns exchanged in a conversation.

Example implementation::

class RedisHistory:
    def __init__(self, redis_client):
        self._redis = redis_client

    async def get(self, conversation_id: str) -> list[Message]:
        raw = await self._redis.lrange(f"history:{conversation_id}", 0, -1)
        return [Message.model_validate_json(item) for item in raw]

    async def append(self, conversation_id: str, message: Message) -> None:
        await self._redis.rpush(
            f"history:{conversation_id}", message.model_dump_json()
        )

append(conversation_id, message) async

Append a message to a conversation's history.

Parameters:

Name Type Description Default
conversation_id str

Unique identifier for the conversation.

required
message Message

The message to record.

required

get(conversation_id) async

Retrieve the full message history for a conversation.

Parameters:

Name Type Description Default
conversation_id str

Unique identifier for the conversation.

required

Returns:

Type Description
list[Message]

List of messages in chronological order. Empty list if the

list[Message]

conversation has no recorded history.

sofias_sdk_lite.SummarizingHistoryProvider

Bases: HistoryProvider, Protocol

A HistoryProvider that can fold old turns into a summary.

AgentRunner calls maybe_summarize after the response has been delivered, as a background task that never blocks the next message and is drained on graceful shutdown. Implementations decide whether a summary is due (e.g. by message count or token estimate) and how to store it; the runner only guarantees the call happens once per turn.

Example implementation::

class SummarizingRedisHistory(RedisHistory):
    async def maybe_summarize(self, conversation_id: str) -> None:
        messages = await self.get(conversation_id)
        if len(messages) < 40:
            return
        summary = await self._llm_summarize(messages[:-10])
        await self._replace(conversation_id, [summary, *messages[-10:]])

maybe_summarize(conversation_id) async

Summarise older turns if the conversation needs it.

Called fire-and-forget after each delivered response. Must be safe to call when nothing needs summarising.

sofias_sdk_lite.ChatResponseWorkflow

Publishes agent responses as StreamFragment chunks to a RabbitMQ stream.

Implements the ResponseWorkflow protocol (send_response). Also provides send_error and close for infrastructure-level error handling and cleanup.

Parameters:

Name Type Description Default
agent_name str

Name used in StreamFragment.agent.

required
client RabbitMQClient

RabbitMQ client for publishing.

required
stream_name str

Target RabbitMQ stream name. No default naming convention is assumed — pass whatever stream your deployment uses (see RunnerConfig.response_stream).

required
conversation_id str

Current conversation ID.

required
message_id str | int

Current message ID.

required
chunk_size int

Maximum characters per fragment (default 100).

_DEFAULT_CHUNK_SIZE
traceparent str | None

W3C traceparent header from the inbound message. When set, it is echoed on every fragment's AMQP application properties so the streamed response continues the producer's trace end to end. If OpenTelemetry is installed and a span is active, that span's context wins over the inbound header.

None

close() async

Close the underlying stream producer.

send_error(error) async

Publish an error fragment.

send_response(response, context) async

Publish the agent's answer as chunked StreamFragments.

Extracts the answer from response.content (tries "answer" then "content" keys) and chunks it into fragments. An empty answer is replaced by EMPTY_RESPONSE_FALLBACK so a failed turn is never delivered as silence. The turn's token usage (response.usage) is attached to the terminal fragment's state.usage.

sofias_sdk_lite.BaseStreamingResponseWorkflow

Publishes agent responses as StreamFragment chunks, incrementally.

Implements the StreamingResponseWorkflow protocol. Subclass it and override the hook methods for agent-specific business logic; the wire mechanics (_publish_fragment, chunking, traceparent echo, usage) are inherited as-is.

Parameters:

Name Type Description Default
agent_name str

Used in StreamFragment.agent.

required
client RabbitMQClient | None

RabbitMQ client. Pass exactly one of client/publisher; with client the workflow builds (and owns/closes) its own RabbitMQPublisher.

None
publisher RabbitMQPublisher | None

An already-built, externally-owned RabbitMQPublisher (e.g. shared across many conversations for connection reuse). close() is a no-op when the publisher was injected this way.

None
stream_name str

Target RabbitMQ stream.

required
conversation_id str

Current conversation ID.

required
message_id str | int

Current message ID (echoed by the agent, matched back by the consumer — accepts int | str).

required
chunk_size int

Maximum characters per fragment.

_DEFAULT_CHUNK_SIZE
traceparent str | None

W3C traceparent from the inbound message, echoed on every fragment unless an active OpenTelemetry span overrides it.

None
stream_node_names set[str] | None

If set, only TextChunkEvents from these graph nodes are streamed (multi-node agents that only want to expose one node's output).

None
stream_key str | None

If set, every fragment carries this key instead of defaulting to conversation_id — lets the consumer multiplex several independent streams onto the same conversation (e.g. a background job's progress bubble vs. the main reply).

None
first_flush_prefix str

One-shot text prepended to the very first published fragment (e.g. a "done" chip that replaces an already-open progress indicator the instant real content lands).

''

close() async

Close the publisher, if this workflow built its own.

send_notice(message) async

Publish message as a normal, completed reply — no "Error:" prefix, so a degraded-but-graceful message renders as an ordinary chat bubble rather than an error banner.

sofias_sdk_lite.NullWorkflow

A ResponseWorkflow that records calls instead of delivering anywhere.

Useful for testing agents and runners without a real transport.

close() async

Mark this workflow as closed.

send_error(error) async

Record the error instead of delivering it.

send_response(response, context) async

Record the response instead of delivering it.

sofias_sdk_lite.usage_state_from_list(usage)

Collapse an AgentResponse.usage list into a single StreamFragment state.usage object, or None when there is nothing to report.

Every response workflow that publishes a terminal fragment should call this and merge the result into state: it is the one place a billing consumer reads token usage from. Sums tokens across every LLM call in the turn (a turn can invoke several models: the main chat call plus e.g. an embedding call for RAG), and names the entry with the most completion tokens as model, which picks the actual generation call over near-zero-completion side calls.

LLM clients

sofias_sdk_lite.create_llm(source=None, /, *, model=None, base_url=None, api_key=None, provider=None, temperature=None, max_tokens=None, timeout=None, max_retries=None, reasoning_effort=None, **adapter_kwargs)

Build a ready-to-use LLM client.

Parameters:

Name Type Description Default
source LLMGatewayConfig | Any | None

An LLMGatewayConfig, an agent settings object exposing model_name / router_url / router_api_key, or None to fall back to keyword arguments and then the environment.

None
model str | None

Overrides the resolved model identifier or gateway alias.

None
base_url str | None

Overrides the resolved gateway URL.

None
api_key str | None

Overrides the resolved bearer token.

None
provider str | None

Overrides the adapter flavour ("sofias" or an OpenAI-compatible label).

None
temperature float | None

Overrides the sampling temperature.

None
max_tokens int | None

Overrides the output budget.

None
timeout float | None

Overrides the per-request timeout in seconds.

None
max_retries int | None

Overrides the number of retries on transient failures.

None
reasoning_effort str | None

Overrides the per-model-family default (Sofias gateway only).

None
**adapter_kwargs Any

Forwarded to the adapter constructor (e.g. transport= for tests, default_headers=).

{}

Raises:

Type Description
LLMConfigurationError

If no model or gateway URL can be resolved, or the provider is unknown.

sofias_sdk_lite.LLMGatewayConfig

Bases: BaseModel

Everything needed to build an LLM client.

model names the model; base_url must be set before create_llm can build a client (the platform injects it). provider defaults to "sofias", the Sofias gateway.

api_key = SecretStr('') class-attribute instance-attribute

Bearer token. Empty sends no Authorization header.

base_url = '' class-attribute instance-attribute

OpenAI-compatible API root. Injected by the platform; required to build a client.

max_retries = 3 class-attribute instance-attribute

Retries on transient failures.

max_tokens = None class-attribute instance-attribute

Output budget; None keeps the adapter default.

model instance-attribute

Model identifier or gateway alias (e.g. "default").

provider = PROVIDER_GATEWAY class-attribute instance-attribute

"sofias" or any OpenAI-compatible label ("openai", "litellm", ...).

reasoning_effort = None class-attribute instance-attribute

Sofias gateway only: overrides the per-model-family default.

temperature = None class-attribute instance-attribute

Sampling temperature; None keeps the adapter default.

timeout = 120.0 class-attribute instance-attribute

Per-request timeout in seconds.

with_overrides(**overrides)

Return a copy with the non-None overrides applied.

sofias_sdk_lite.GatewayLLM

Bases: OpenAICompatibleLLM

StreamableLLMCallable for the Sofias gateway.

Parameters:

Name Type Description Default
api_key str

Gateway API key (empty string sends no Authorization).

required
model str

Model identifier or gateway alias (e.g. "default").

required
base_url str

Gateway host or API root. Required; a bare host gets /v1 appended (see normalize_gateway_base_url).

required
temperature float

Sampling temperature.

0.1
max_retries int

Retries on transient failures.

3
timeout float

Per-request timeout in seconds.

120.0
reasoning_effort str | None

None applies the per-family default; a non-empty string is sent as-is on every request; "" disables injection so the model default applies.

None
max_tokens int | None

Output budget; None applies DEFAULT_GATEWAY_MAX_TOKENS.

None
default_headers dict[str, str] | None

Extra headers sent on every request.

None
transport AsyncBaseTransport | None

Optional httpx transport (for tests).

None

Raises:

Type Description
LLMConfigurationError

If base_url is empty.

sofias_sdk_lite.OpenAICompatibleLLM

StreamableLLMCallable over the OpenAI chat-completions HTTP API.

Parameters:

Name Type Description Default
api_key str

Bearer token. Empty string sends no Authorization header (for gateways that do not authenticate on an internal network).

required
base_url str

API root, e.g. https://api.openai.com/v1.

required
model str

Model identifier (or gateway alias) sent on every request.

required
provider str

Label reported in LLMResponse.provider and logs. Gateways listed in GATEWAY_PROVIDERS have the real upstream provider and model read back from the response's extra_fields.

'openai'
temperature float

Sampling temperature.

0.3
max_retries int

Retries on transient failures (transport errors, 408/409/429/5xx). 0 disables retrying.

3
timeout float

Per-request timeout in seconds.

120.0
max_tokens int | None

Output budget sent on every request; None omits it.

None
default_headers dict[str, str] | None

Extra headers sent on every request.

None
transport AsyncBaseTransport | None

Optional httpx transport (tests inject a MockTransport here).

None

model property

Model identifier sent on every request.

provider property

Provider label configured for this client.

aclose() async

Close the underlying HTTP client.

invoke(prompt, tools=None, messages=None, system_prompt=None, response_format=None, parallel_tool_calls=None, **_extra) async

Run one chat completion and return the parsed response.

stream_invoke(prompt, tools=None, messages=None, system_prompt=None, response_format=None, parallel_tool_calls=None, **_extra) async

Run one chat completion as a server-sent-events stream.

Yields text deltas as they arrive; tool calls are accumulated and yielded once complete, followed by a final usage event when the gateway reports one.

sofias_sdk_lite.default_llm(llm)

Make llm the default for AgentBuilder.build() calls inside the block.

Passing None leaves any outer default untouched.

LLM protocol

sofias_sdk_lite.LLMCallable

Bases: Protocol

Protocol defining how an LLM is invoked.

Implement this against any LLM client (HTTP-based, local, mocked, etc.) to make it usable by the Sofias agent SDK.

sofias_sdk_lite.StreamableLLMCallable

Bases: LLMCallable, Protocol

An LLMCallable that also supports streaming invocation.

sofias_sdk_lite.LLMResponse

Bases: BaseModel

Response from a single LLM invocation.

has_tool_calls property

Whether the LLM requested one or more tool calls.

is_final property

Whether this response is a terminal answer (no tool calls to run).

sofias_sdk_lite.ToolSpec

Bases: BaseModel

Specification for a tool made available to the LLM.

LLM errors

sofias_sdk_lite.LLMConfigurationError

Bases: AgentSDKError

Error when no usable LLM gateway configuration can be resolved.

Raised by sofias_sdk_lite.llm.create_llm when neither the explicit arguments, the agent settings, nor the environment provide enough information (at minimum a model identifier) to build an LLM client.

sofias_sdk_lite.LLMRequestError

Bases: AgentSDKError

Error when an LLM gateway request fails after exhausting retries.

Carries the HTTP status (None for transport-level failures) and a truncated response body so callers can log or branch on it.

sofias_sdk_lite.EmptyLLMResponseError

Bases: AgentSDKError

Error when the LLM returns neither text content nor tool calls.

LLMNode treats this like any other invocation failure and retries the call, so a transient empty completion does not end the turn.