job_context
Worker group types for structured concurrent worker execution.
- class pipecat.pipeline.job_context.JobParams(*, name: str | None = None, payload: dict | None = None, timeout: float | None = None, label: str | None = None, cancellable: bool = True)[source]
Bases:
BaseModelConfiguration for a job sent to a single worker.
- Parameters:
name – Optional job name, routing the request to the matching
@job(name=...)handler on the worker.payload – Optional structured data describing the work.
timeout – Optional timeout in seconds, covering both the wait for the worker to become ready and the job itself.
label – Optional human-readable description of the work, e.g.
"Research: Radiohead". AUIWorkertitles the client’s progress card with it.cancellable – Whether an external requester, such as the client UI, may ask for the job to be cancelled. Cancellation the worker initiates itself (shutdown, timeout,
cancel_on_error) proceeds either way.
- class pipecat.pipeline.job_context.JobGroupParams(*, name: str | None = None, payload: dict | None = None, timeout: float | None = None, label: str | None = None, cancellable: bool = True, cancel_on_error: bool = True)[source]
Bases:
JobParamsConfiguration for a job sent to several workers at once.
Carries every
JobParamsfield, which apply to the group as a whole, plus how the group reacts to one of its workers failing.- Parameters:
cancel_on_error – Whether a worker responding with an error status cancels the rest of the group.
- class pipecat.pipeline.job_context.JobParamsT
Either params class, so
resolve_job_paramshands back what it was asked for.alias of TypeVar(‘JobParamsT’, bound=
JobParams)
- pipecat.pipeline.job_context.resolve_job_params(params: JobParamsT | None, params_class: type[JobParamsT], **deprecated) JobParamsT[source]
Fold the deprecated per-argument job spelling into a params object.
Dispatch methods accept a params object and, until 2.0.0, the individual arguments it replaced. Call this once at the entry point a caller reaches, then pass the result down as
paramsso nothing warns twice.Deprecated since version 1.8.0: No replacement. Will be removed in 2.0.0, along with the individual arguments it exists to absorb.
- Parameters:
params – The params object the caller passed, if any.
params_class – The class to build when only individual arguments came in.
**deprecated – The individual arguments. Each defaults to None, so a non-None value means the caller passed it explicitly.
- Returns:
The params to dispatch with.
- Raises:
TypeError – If both
paramsand individual arguments were passed.
- class pipecat.pipeline.job_context.JobStatus(*values)[source]
Bases:
StrEnumStatus of a completed worker.
Inherits from
strso values compare naturally with plain strings and serialize without extra handling.- Parameters:
COMPLETED – The worker finished successfully.
CANCELLED – The worker was cancelled by the requester.
FAILED – The worker failed due to a logical or business error.
ERROR – The worker encountered an unexpected runtime error.
- COMPLETED = 'completed'
- CANCELLED = 'cancelled'
- FAILED = 'failed'
- ERROR = 'error'
- exception pipecat.pipeline.job_context.JobError[source]
Bases:
ExceptionRaised when a worker is cancelled due to a worker error or timeout.
- exception pipecat.pipeline.job_context.JobGroupError[source]
Bases:
ExceptionRaised when a worker group is cancelled due to a worker error or timeout.
- class pipecat.pipeline.job_context.JobGroupResponse(job_id: str, responses: dict[str, dict])[source]
Bases:
objectCollected results from a completed job group.
- Parameters:
job_id – The shared job identifier.
responses – Collected responses keyed by worker name.
- class pipecat.pipeline.job_context.JobEvent(type: str, data: dict | None = None)[source]
Bases:
objectAn event received from a worker during a single-worker job.
- Parameters:
type – The event type.
data – Optional event payload.
- class pipecat.pipeline.job_context.JobGroupEvent(type: str, worker_name: str, data: dict | None = None)[source]
Bases:
objectAn event received from a worker during job group execution.
- Parameters:
type – The event type.
worker_name – The name of the worker that sent the event.
data – Optional event payload.
- class pipecat.pipeline.job_context.JobGroup(job_id: str, worker_names: list[str], responses: dict[str, dict] = <factory>, timeout_task: ~_asyncio.Task | None = None, cancel_on_error: bool = True, label: str | None = None, cancellable: bool = True, terminated: set[str] = <factory>, event_queue: ~asyncio.queues.Queue | None = None, _done: ~asyncio.locks.Event = <factory>, _error: str | None = None)[source]
Bases:
objectTracks a group of workers launched together.
- Parameters:
job_id – Shared identifier for all workers in this group.
worker_names – Names of the workers in the group, in dispatch order.
responses – Collected responses keyed by worker name.
timeout_task – Optional asyncio worker that cancels the group on timeout.
cancel_on_error – Whether to cancel the group if a worker errors.
label – Optional human-readable description of the work, from
JobGroupParams.label.cancellable – Whether an external requester may ask for the group to be cancelled, from
JobGroupParams.cancellable.terminated – Names of the workers that have reached a terminal state, whether by responding or by ending their stream.
event_queue – Optional queue for streaming events to a
JobGroupContextasync iterator.
- async wait() None[source]
Wait for all workers in the group to respond.
- Raises:
JobGroupError – If the group was cancelled due to error or timeout.
- class pipecat.pipeline.job_context.JobGroupContext(worker: BaseWorker, worker_names: tuple[str, ...], *, params: JobGroupParams | None = None)[source]
Bases:
objectAsync context manager and iterator for structured job group execution.
Sends job requests on enter, waits for all responses on exit. Supports
async forto receive intermediate events (updates and streaming data) from workers while waiting for completion.On normal completion, results are available via
responses. On worker error (withcancel_on_error=True) or timeout, raisesJobGroupError. If theasync withblock raises, remaining jobs are cancelled.Example:
async with self.job_group( "w1", "w2", params=JobGroupParams(payload=data) ) as tg: async for event in tg: print(f"{event.worker_name} [{event.type}]: {event.data}") for name, result in tg.responses.items(): print(name, result)
- __init__(worker: BaseWorker, worker_names: tuple[str, ...], *, params: JobGroupParams | None = None)[source]
Initialize the JobGroupContext.
- Parameters:
worker – The parent BaseWorker that owns this job group.
worker_names – Names of the workers to send the job to.
params – How to run the group. Defaults to
JobGroupParamsdefaults.
- class pipecat.pipeline.job_context.JobContext(worker: BaseWorker, worker_name: str, *, params: JobParams | None = None)[source]
Bases:
objectAsync context manager and iterator for a single-worker job.
Sends a job request on enter, waits for the response on exit. Supports
async forto receive intermediate events (updates and streaming data) from the worker while waiting for completion.On normal completion, the result is available via
response. On worker error or timeout, raisesJobError. If theasync withblock raises, the job is cancelled.Example:
async with self.job("worker", params=JobParams(payload=data)) as t: async for event in t: print(f"[{event.type}]: {event.data}") print(t.response)