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]
|
|
dict[str, Any]
|
|
dict[str, Any]
|
|
dict[str, Any]
|
|
dict[str, Any]
|
|
dict[str, Any]
|
|
dict[str, Any]
|
|
dict[str, Any]
|
|
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:
-
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, )
-
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]
|
|
dict[str, Any]
|
validation into an |
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
|
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 |
required |
conversation_id
|
str | None
|
Unused; accepted to satisfy the |
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'
|
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 |
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
|
required |
workflow
|
Any
|
The |
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 |
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 |
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 |
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 |
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 |
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 |
required |
client
|
RabbitMQClient | None
|
RabbitMQ client. Pass exactly one of |
None
|
publisher
|
RabbitMQPublisher | None
|
An already-built, externally-owned |
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 |
required |
chunk_size
|
int
|
Maximum characters per fragment. |
_DEFAULT_CHUNK_SIZE
|
traceparent
|
str | None
|
W3C |
None
|
stream_node_names
|
set[str] | None
|
If set, only |
None
|
stream_key
|
str | None
|
If set, every fragment carries this key instead of
defaulting to |
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 |
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 ( |
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.
|
{}
|
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 |
required |
model
|
str
|
Model identifier or gateway alias (e.g. |
required |
base_url
|
str
|
Gateway host or API root. Required; a bare host gets
|
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
|
max_tokens
|
int | None
|
Output budget; |
None
|
default_headers
|
dict[str, str] | None
|
Extra headers sent on every request. |
None
|
transport
|
AsyncBaseTransport | None
|
Optional |
None
|
Raises:
| Type | Description |
|---|---|
LLMConfigurationError
|
If |
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 |
required |
base_url
|
str
|
API root, e.g. |
required |
model
|
str
|
Model identifier (or gateway alias) sent on every request. |
required |
provider
|
str
|
Label reported in |
'openai'
|
temperature
|
float
|
Sampling temperature. |
0.3
|
max_retries
|
int
|
Retries on transient failures (transport errors,
408/409/429/5xx). |
3
|
timeout
|
float
|
Per-request timeout in seconds. |
120.0
|
max_tokens
|
int | None
|
Output budget sent on every request; |
None
|
default_headers
|
dict[str, str] | None
|
Extra headers sent on every request. |
None
|
transport
|
AsyncBaseTransport | None
|
Optional |
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
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.