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: BaseModel

Configuration 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". A UIWorker titles 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: JobParams

Configuration for a job sent to several workers at once.

Carries every JobParams field, 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_params hands 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 params so 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 params and individual arguments were passed.

class pipecat.pipeline.job_context.JobStatus(*values)[source]

Bases: StrEnum

Status of a completed worker.

Inherits from str so 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: Exception

Raised when a worker is cancelled due to a worker error or timeout.

exception pipecat.pipeline.job_context.JobGroupError[source]

Bases: Exception

Raised 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: object

Collected 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: object

An event received from a worker during a single-worker job.

Parameters:
  • type – The event type.

  • data – Optional event payload.

UPDATE: ClassVar[str] = 'update'
STREAM_START: ClassVar[str] = 'stream_start'
STREAM_DATA: ClassVar[str] = 'stream_data'
STREAM_END: ClassVar[str] = 'stream_end'
data: dict | None = None
class pipecat.pipeline.job_context.JobGroupEvent(type: str, worker_name: str, data: dict | None = None)[source]

Bases: object

An 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.

UPDATE: ClassVar[str] = 'update'
STREAM_START: ClassVar[str] = 'stream_start'
STREAM_DATA: ClassVar[str] = 'stream_data'
STREAM_END: ClassVar[str] = 'stream_end'
data: dict | None = None
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: object

Tracks 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 JobGroupContext async iterator.

timeout_task: Task | None = None
cancel_on_error: bool = True
label: str | None = None
cancellable: bool = True
event_queue: Queue | None = None
property is_done: bool

Whether the group has completed or failed.

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.

complete() → None[source]

Signal that all workers have responded.

fail(reason: str | None = None) → None[source]

Signal that the group was cancelled.

Parameters:

reason – Human-readable reason for the failure.

class pipecat.pipeline.job_context.JobGroupContext(worker: BaseWorker, worker_names: tuple[str, ...], *, params: JobGroupParams | None = None)[source]

Bases: object

Async context manager and iterator for structured job group execution.

Sends job requests on enter, waits for all responses on exit. Supports async for to receive intermediate events (updates and streaming data) from workers while waiting for completion.

On normal completion, results are available via responses. On worker error (with cancel_on_error=True) or timeout, raises JobGroupError. If the async with block 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 JobGroupParams defaults.

property job_id: str

The shared job identifier for this group.

property responses: dict[str, dict]

Collected responses keyed by worker name.

class pipecat.pipeline.job_context.JobContext(worker: BaseWorker, worker_name: str, *, params: JobParams | None = None)[source]

Bases: object

Async context manager and iterator for a single-worker job.

Sends a job request on enter, waits for the response on exit. Supports async for to 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, raises JobError. If the async with block 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)
__init__(worker: BaseWorker, worker_name: str, *, params: JobParams | None = None)[source]

Initialize the JobContext.

Parameters:
  • worker – The parent BaseWorker that owns this job.

  • worker_name – Name of the worker to send the job to.

  • params – How to run the job. Defaults to JobParams defaults.

property job_id: str

The job identifier.

property response: dict

The worker’s response payload.