Skip to main content

API reference

Auto-generated from triggers.__all__ by scripts/gen_api_reference.py. Regenerate via make docs.

Public surface

Context (dataclass)

Context(flow_run_id: str, flow_id: str, flow_version: str, source_node_id: str, retry_count: int, idempotency_key: str, traceparent: str, log: Any, tracer: Any, ingress_target: str = '', ingress_source_type: str = '', raw_event: bytes = b'') -> None

Constructor fields

FieldRequiredTypeDefault
flow_run_idyesstr-
flow_idyesstr-
flow_versionyesstr-
source_node_idyesstr-
retry_countyesint-
idempotency_keyyesstr-
traceparentyesstr-
logyesAny-
traceryesAny-
ingress_targetnostr''
ingress_source_typenostr''
raw_eventnobytesb''

HTTP

Value: Transport(kind='http', addr='')

InputField (dataclass)

InputField(name: str, type: str, required: bool = False, description: str = '') -> None

Constructor fields

FieldRequiredTypeDefault
nameyesstr-
typeyesstr-
requirednoboolFalse
descriptionnostr''

Methods

  • to_dict(self) -> dict[str, Any]

KAFKA

Value: Transport(kind='kafka', addr='')

Lifecycle (enum)

Public lifecycle enum for service and action catalog metadata.

MemberValue
STABLE'stable'
BETA'beta'
DEPRECATED'deprecated'

Lookup (dataclass)

Lookup(name: str, params: tuple[InputField, ...] = (), result_fields: tuple[InputField, ...] = (), display_name: str = '', description: str = '') -> None

A synchronous, stateless request/response the engine may issue mid-flow on triggers.v2.service.<service>.lookup.<name>. Stateless means no state carried between executions, not side-effect-free: a handler may write.

Constructor fields

FieldRequiredTypeDefault
nameyesstr-
paramsnotuple[InputField, ...]()
result_fieldsnotuple[InputField, ...]()
display_namenostr''
descriptionnostr''

Methods

  • to_dict(self) -> dict[str, Any]

NATS

Value: Transport(kind='nats', addr='')

PermanentError (class)

Input is unrecoverable. Routes to DLQ kind="permanent"; acks; no retry; no event.

Constructor: PermanentError(code_or_msg: str, message: str | None = None) -> None

ResultEventDecl (dataclass)

ResultEventDecl(event_type: str, schema_id: str, schema_version: str = '', fields: tuple[InputField, ...] = ()) -> None

Catalog metadata for Studio chained-field resolution (v2 shape).

Constructor fields

FieldRequiredTypeDefault
event_typeyesstr-
schema_idyesstr-
schema_versionnostr''
fieldsnotuple[InputField, ...]()

Methods

  • to_dict(self) -> dict[str, Any]

Service (dataclass)

Service(name: str, display_name: str = '', description: str = '', lifecycle: Lifecycle = Lifecycle.STABLE, icon: str = '', settings: TriggersSettings | None = None, logger: Any = None) -> None

Constructor fields

| Field | Required | Type | Default | | -------------- | -------- | ----------------- | ------------------ | ------ | | name | yes | str | - | | display_name | no | str | '' | | description | no | str | '' | | lifecycle | no | Lifecycle | Lifecycle.STABLE | | icon | no | str | '' | | settings | no | TriggersSettings | None | None | | logger | no | Any | None |

Derived instance attributes

FieldTypeNotes
versionstrPopulated by the SDK after construction
drain_timeoutfloatPopulated by the SDK after construction
metrics_addrstrPopulated by the SDK after construction
http_addrstrPopulated by the SDK after construction
nats_urlstrPopulated by the SDK after construction
nats_creds_pathstrPopulated by the SDK after construction
nats_source_streamstrPopulated by the SDK after construction
kafka_brokerslist[str]Populated by the SDK after construction
redis_urlstrPopulated by the SDK after construction
max_in_flightintPopulated by the SDK after construction
dedup_window_secondsintPopulated by the SDK after construction
result_recovery_enabledboolPopulated by the SDK after construction
result_recovery_ttl_secondsfloatPopulated by the SDK after construction

Methods

  • action(self, name: str, *, timeout: float = 30.0, retry: Mapping[str, Any] | None = None, dlq_subject: str | None = None, description: str = '', lifecycle: Lifecycle = Lifecycle.STABLE, action_env: tuple[InputField, ...] | list[InputField] = (), result_event: ResultEventDecl | None = None, emit_result_event: bool = True, max_forwarding_age_seconds: float = 0, freshness_timestamp_field: str = '') -> Callable[[ActionHandler], ActionHandler]
  • build_manifest(self) -> Manifest
  • data_source(self, name: str, handler: DataSourceHandler, *, description: str = '', timeout_ms: int = 0) -> Service
  • data_source_subject(self, name: str) -> str
  • enable_dedup(self, *, prefix: str | None = None, window_seconds: int | None = None) -> Service
  • enable_rate_limit(self, *, prefix: str | None = None) -> Service
  • ingress(self, *transports: Transport) -> Service
  • lookup(self, name: str, *, display_name: str = '', description: str = '')
  • raise_if_invalid(self) -> None
  • run(self, *, bus: Any = None) -> None
  • wiring(self) -> ServiceWiring

ServiceError (class)

Inappropriate argument value (of correct type).

Constructor: not introspectable.

SkipForward (class)

Handler chose to no-op. Acks; no DLQ; no event.

Constructor: not introspectable.

Transport (dataclass)

Transport(kind: str, addr: str = '') -> None

Constructor fields

FieldRequiredTypeDefault
kindyesstr-
addrnostr''

Methods

  • at(self, addr: str) -> Transport

Service methods

Service.ingress(self, *transports: Transport) -> Service

Service.enable_dedup(self, *, prefix: str | None = None, window_seconds: int | None = None) -> Service

Service.enable_rate_limit(self, *, prefix: str | None = None) -> Service

Service.data_source(self, name: str, handler: DataSourceHandler, *, description: str = '', timeout_ms: int = 0) -> Service

Service.lookup(self, name: str, *, display_name: str = '', description: str = '')

Register a synchronous, stateless request/response.

Service.action(self, name: str, *, timeout: float = 30.0, retry: Mapping[str, Any] | None = None, dlq_subject: str | None = None, description: str = '', lifecycle: Lifecycle = Lifecycle.STABLE, action_env: tuple[InputField, ...] | list[InputField] = (), result_event: ResultEventDecl | None = None, emit_result_event: bool = True, max_forwarding_age_seconds: float = 0, freshness_timestamp_field: str = '') -> Callable[[ActionHandler], ActionHandler]

Service.wiring(self) -> ServiceWiring

Service.run(self, *, bus: Any = None) -> None