llm
LLM worker package – LLMWorker, LLMContextWorker, BackendLLMWorker, and the @tool decorator.
- class pipecat.workers.llm.BackendLLMWorker(*, llm: LLMService[Any], context: LLMContext | None = None, name: str | None = None, transform_output: Callable[[BackendOutput], Awaitable[BackendOutput]] | None = None, user_params: LLMUserAggregatorParams | None = None, assistant_params: LLMAssistantAggregatorParams | None = None)[source]
Bases:
LLMContextWorkerA worker that runs an LLM service as the backend a frontend delegates to.
The worker owns the backend’s conversation: an
LLMContextplus the aggregator pair, so multi-step tool calling works as it does in any pipeline. Each delegation arrives as arunjob carrying the request text, which is appended to the context as one user message, and runs the LLM until it produces a final answer.Job contract (
@job(name="run"), one delegation at a time):request payload:
{"request": str}— the text to put to the backend, composed by the frontend.updates: a
BackendOutputpayload for every piece of output — reasoning summaries, what the backend says before calling tools, and its final answer.response:
{"text": str}— the final answer, or""if the delegation ended without one. A backend LLM failure answers the job withJobStatus.ERROR, so the frontend hears about it as soon as it happens.
The final answer arrives twice, as the last update and as the response, so a caller uses one or the other: a frontend relaying output as it arrives reads the updates, while one that needs a return value reads the response and skips updates marked
is_final.Example:
backend = BackendLLMWorker( llm=AnthropicLLMService(api_key=...), context=LLMContext( [{"role": "system", "content": BACKEND_INSTRUCTIONS}], tools=[get_weather], ), )
- __init__(*, llm: LLMService[Any], context: LLMContext | None = None, name: str | None = None, transform_output: Callable[[BackendOutput], Awaitable[BackendOutput]] | None = None, user_params: LLMUserAggregatorParams | None = None, assistant_params: LLMAssistantAggregatorParams | None = None)[source]
Initialize the backend worker.
- Parameters:
llm – The backend LLM service.
context – The backend’s context, typically carrying its tools. A fresh empty context when omitted.
name – Worker name; auto-generated when omitted.
transform_output – Called with each
BackendOutputbefore it is sent, to adjust its text or whether the user may hear it. Without one, only the final answer asks to be spoken: a frontend filling the wait is usually mid-sentence when progress arrives, and speaking it talks over them.user_params – Optional parameters for the user aggregator. Defaults to external turn strategies: the backend has no audio, so the default VAD and turn-analysis strategies (and the model the latter loads) have nothing to do here.
assistant_params – Optional parameters for the assistant aggregator.
- async run_delegation(message: BusJobRequestMessage)[source]
Run one delegation to completion, streaming what the backend produces as updates.
- Parameters:
message – The job request; see the class docstring for the payload.
- class pipecat.workers.llm.LLMWorker(name: str, *, llm: LLMService[Any], pipeline: Pipeline | None = None, active: bool = False, bridged: tuple[str, ...] | None = None, defer_tool_frames: bool = True)[source]
Bases:
PipelineWorkerWorker with an LLM pipeline and automatic tool registration.
Methods decorated with
@toolare registered as direct functions on the LLM and tracked so that frames queued during tool execution can be deferred until all tools complete.Example:
class MyTask(LLMWorker): @tool async def my_function(self, params, arg: str): ... worker = MyTask("worker", bus=bus, llm=OpenAILLMService(api_key="..."))
- __init__(name: str, *, llm: LLMService[Any], pipeline: Pipeline | None = None, active: bool = False, bridged: tuple[str, ...] | None = None, defer_tool_frames: bool = True)[source]
Initialize the LLMWorker.
- Parameters:
name – Unique name for this worker.
llm – The LLM service.
@tooldecorated methods are automatically registered on it.pipeline – Optional pipeline override. When
None, defaults toPipeline([llm]). Subclasses can pass a custom pipeline that wraps the LLM with additional processors.active – Whether the worker starts active. Defaults to False.
bridged – Bridge configuration forwarded to
PipelineWorker. Pass()to wrap the LLM pipeline with bus edge processors so it can exchange frames with another bridged worker.defer_tool_frames – Whether to defer frames queued during tool execution until all tools complete. Defaults to True.
- property llm: LLMService
The LLM service this worker wraps.
- async on_activated(args: dict | None) None[source]
Configure the LLM with tools and activation messages.
- Parameters:
args – Optional activation arguments with messages to append.
- async queue_frame(frame: Frame, direction: FrameDirection = FrameDirection.DOWNSTREAM) None[source]
Queue a frame, holding it if a tool handler queued it.
A frame queued from inside one of this worker’s
@toolhandlers, or from anything that handler awaits, is held and delivered once the last tool finishes. Frames from anywhere else are queued immediately: the worker’s own traffic and frames arriving over the bus, which run outside any handler’s context, and frames a handler on a different worker queues here, which this worker would never release.- Parameters:
frame – Any
Frameto deliver.direction – Direction the frame should travel. Defaults to
FrameDirection.DOWNSTREAM.
- build_tools() list[source]
Return the tools for this worker’s LLM.
By default, returns all methods decorated with
@tool. Override to provide additional or different tools.- Returns:
List of tool functions.
- async end(*, reason: str | None = None, messages: list | None = None, result_callback: Callable[[...], Any] | None = None) None[source]
Request a graceful end of the session.
When called from a
@toolhandler, deliver the function call result first withawait params.result_callback(result): the LLM output it triggers is delivered before the session ends.- Parameters:
reason – Optional human-readable reason for ending.
messages –
Optional LLM messages to inject and speak before ending. The LLM runs immediately so the output is delivered before the session terminates.
Deprecated since version 1.8.0: Call
params.result_callback(result)beforeend()instead. Will be removed in 2.0.0.result_callback –
The
result_callbackfrom FunctionCallParams.Deprecated since version 1.8.0: Call
params.result_callback(result)beforeend()instead. Will be removed in 2.0.0.
- async activate_worker(worker_name: str, *, args: WorkerActivationArgs | None = None, deactivate_self: bool = False, messages: list | None = None, result_callback: Callable[[...], Any] | None = None) None[source]
Activate another worker, draining this worker’s pipeline to hand over.
When called from a
@toolhandler, deliver the function call result first withawait params.result_callback(result): the output it triggers is delivered before the target is activated. The handover itself waits until the tool call asking for it has finished.- Parameters:
worker_name – The name of the worker to activate.
args – Optional
WorkerActivationArgsforwarded to the target worker’son_activatedhandler.deactivate_self – Whether to deactivate this worker before activating the target. Deactivating this worker drains its pipeline first; staying active does not. A worker that stays active and wants to drain anyway can call
flush_pipeline()before this.messages –
Optional LLM messages to inject and deliver before activating the target. The LLM runs immediately so the output is delivered before the transfer completes.
Deprecated since version 1.8.0: Call
params.result_callback(result)beforeactivate_worker()instead. Will be removed in 2.0.0.result_callback –
The
result_callbackfrom FunctionCallParams.Deprecated since version 1.8.0: Call
params.result_callback(result)beforeactivate_worker()instead. Will be removed in 2.0.0.
- async process_deferred_tool_frames(frames: list[tuple[Frame, FrameDirection]]) list[tuple[Frame, FrameDirection]][source]
Process deferred frames before they are flushed.
Called after all in-flight tools complete, before the deferred frames are queued into the pipeline. Override to inspect, modify, reorder, or filter the frames.
- Parameters:
frames – The deferred frames collected during tool execution.
- Returns:
The frames to queue. Return the list as-is for default behavior.
- class pipecat.workers.llm.LLMWorkerActivationArgs(metadata: dict | None = None, messages: list | None = None, run_llm: bool | None = None)[source]
Bases:
WorkerActivationArgsActivation arguments for LLM workers.
- Parameters:
messages – LLM context messages to inject on activation.
run_llm – Whether to run the LLM after appending messages. Defaults to True when
messagesis set.
- class pipecat.workers.llm.LLMContextWorker(name: str, *, llm: LLMService[Any], active: bool = False, bridged: tuple[str, ...] | None = None, defer_tool_frames: bool = True, context: LLMContext | None = None, user_params: LLMUserAggregatorParams | None = None, assistant_params: LLMAssistantAggregatorParams | None = None)[source]
Bases:
LLMWorkerLLM worker that owns an LLMContext and a context aggregator pair.
Useful for workers that need to track their own conversation history, typically workers that run their own LLM pipeline outside of a shared transport pipeline. Subclasses do not need to instantiate the context or aggregators themselves; the pipeline is built as
[user_aggregator, llm, assistant_aggregator]automatically.Example:
worker = LLMContextWorker( "worker", llm=OpenAILLMService(...), ) @worker.assistant_aggregator.event_handler("on_assistant_turn_stopped") async def _on_stopped(aggregator, message): ...
- __init__(name: str, *, llm: LLMService[Any], active: bool = False, bridged: tuple[str, ...] | None = None, defer_tool_frames: bool = True, context: LLMContext | None = None, user_params: LLMUserAggregatorParams | None = None, assistant_params: LLMAssistantAggregatorParams | None = None)[source]
Initialize the LLMContextWorker.
- Parameters:
name – Unique name for this worker.
llm – The LLM service.
active – Whether the worker starts active. Defaults to False.
bridged – Bridge configuration forwarded to
PipelineWorker. Pass()to wrap the pipeline with bus edges so it can exchange frames with another bridged worker.defer_tool_frames – Whether to defer frames queued during tool execution until all tools complete. Defaults to True.
context – Optional pre-built LLMContext. When omitted, a fresh empty context is created.
user_params – Optional parameters for the user aggregator.
assistant_params – Optional parameters for the assistant aggregator.
- property context: LLMContext
The LLMContext owned by this worker.
- property user_aggregator: LLMUserAggregator
The user-side context aggregator.
- property assistant_aggregator: LLMAssistantAggregator
The assistant-side context aggregator.
- class pipecat.workers.llm.BackendOutput(text: str, is_thought: bool = False, is_final: bool = False, prefers_spoken: bool = True)[source]
Bases:
objectOne piece of output from a backend, on its way to the frontend.
- Parameters:
text – The text the backend produced.
is_thought – Whether this is a reasoning summary rather than a response.
is_final – Whether this is the backend’s answer to the delegation, as opposed to progress on the way to it.
prefers_spoken – Whether the backend would like the user to hear this. This is just a hint to the frontend: it may choose to follow it or not.
OpenAILiveLLMService’s live model takes it into consideration, but ultimately decides what to speak (or not) based on the conversation.
- pipecat.workers.llm.tool(fn=None, *, cancel_on_interruption=True, timeout_secs=None, timeout=None)[source]
Mark a method as a tool.
On
LLMWorkersubclasses, decorated methods are automatically registered with the LLM viaregister_direct_functionand included inbuild_tools().This is the worker-flavored variant of
@tool_options: it attaches the samecancel_on_interruption/timeout_secscall options and additionally marks the method (with_pipecat_is_llm_tool) so the worker collects it from the MRO.Can be used with or without arguments:
@tool async def my_tool(self, params, arg: str): ... @tool(cancel_on_interruption=False, timeout_secs=60) async def my_tool(self, params, arg: str): ...
- Parameters:
fn – The function to decorate (when used without arguments).
cancel_on_interruption – Whether to cancel this tool call when an interruption occurs. Defaults to True. Only applies to
LLMWorkertools.timeout_secs – Optional timeout in seconds for this tool call. A call that runs past it is cancelled. Defaults to None (uses the LLM service default). Only applies to
LLMWorkertools.timeout –
Deprecated alias for
timeout_secs.Deprecated since version 1.4.0: Use
timeout_secsinstead. Will be removed in 2.0.0.