API reference

The public extension surface of the periplo package: what a product built on Periplo implements and passes to periplo.bootstrap.create_app. See Extending Periplo for how the pieces fit together.

Extensions

Extension points a product built on Periplo may bring to create_app.

A field left None keeps the open-core, single-tenant default. Nothing here is wired to entry points or any other implicit discovery: what runs is only ever what the embedding product passes in explicitly.

class periplo.extensions.Extensions(*, authenticator=None, tenants=None, authorizer=None, audit=None, orchestrators=None, credentials=None, data_planes=None)[source]

Bases: object

What an embedding product brings to the composition root.

Parameters:
authenticator

Identity port. Default: AnonymousAuthenticator.

Type:

loom.rest.auth.abc.Authenticator | None

tenants

Tenant port. Default: SingleTenant.

Type:

periplo.tenancy.TenantResolver | None

authorizer

Authorization port. Default: SwitchAuthorizer.

Type:

periplo.access.Authorizer | None

audit

Audit port. Default: LogAuditSink.

Type:

periplo.access.AuditSink | None

orchestrators

Orchestrator port. Default: SingleOrchestrator.

Type:

periplo.etl.provider.OrchestratorProvider | None

credentials

Storage credentials port. Default: ProcessCredentials.

Type:

periplo.credentials.CredentialsProvider | None

data_planes

Data plane port. Default: SinglePlane, over the sources and storage the settings describe.

Type:

periplo.data_plane.DataPlanes | None

Identity and tenancy

Per-request identity and tenant.

Two questions every /api/* request answers before it does anything else: who is calling (the Authenticator port of Loom) and which tenant owns the data (TenantResolver, defined here). Both answers are published in a single RequestContext, read back with current_context().

The open-core defaults (AnonymousAuthenticator, SingleTenant) implement single-tenant behaviour: every caller is anonymous and every request belongs to DEFAULT_TENANT.

class periplo.tenancy.Tenant(id)[source]

Bases: object

The owner of the data and ETL a request refers to.

Parameters:

id (str)

periplo.tenancy.DEFAULT_TENANT: Final = Tenant(id='default')

The single tenant of an open-core, single-client installation.

class periplo.tenancy.RequestContext(tenant, identity)[source]

Bases: object

Who is calling, and on behalf of which tenant, for one request.

Parameters:
  • tenant (Tenant)

  • identity (Identity)

periplo.tenancy.current_context()[source]

Return the tenant and identity of the /api/* request being served.

Only periplo.http_context.RequestContextMiddleware sets this, and only for /api/* requests. Outside of one — at start-up, in a health probe, or in a background task — there is no context, and this raises LookupError rather than defaulting to the single tenant.

A task started from a request but meant to outlive it (for example background discovery) MUST run in an empty contextvars.Context() so it does not inherit the request’s tenant; see test_discovery_task_does_not_inherit_the_request_context.

Raises:

LookupError – Outside of a request handled by the middleware.

Return type:

RequestContext

periplo.tenancy.set_context(context)[source]

Publish context for the current request; undo it with reset_context().

Parameters:

context (RequestContext)

Return type:

Token

periplo.tenancy.reset_context(token)[source]

Restore whatever context was published before the matching set_context().

Parameters:

token (Token)

Return type:

None

class periplo.tenancy.TenantResolver(*args, **kwargs)[source]

Bases: Protocol

async resolve(identity, credentials)[source]

Return the tenant that owns credentials.

A header sent more than once reaches credentials.headers as its last value only: a resolver must never rely on seeing every repetition.

Raises:
  • TenantUnresolved – When the request cannot be attributed to any tenant. The caller answers 401 without revealing which tenants exist.

  • loom Unauthenticated or Forbidden – Answered 401 or 403.

Parameters:
  • identity (Identity)

  • credentials (RequestCredentials)

Return type:

Tenant

class periplo.tenancy.SingleTenant[source]

Bases: object

Default TenantResolver: every request is DEFAULT_TENANT.

exception periplo.tenancy.TenantUnresolved[source]

Bases: Exception

Raised by a TenantResolver that cannot attribute a request to a tenant.

exception periplo.tenancy.ForeignTenant(tenant_id)[source]

Bases: Exception

A single-tenant default received a tenant other than its own.

Raised by defaults that are a tenant boundary (SwitchAuthorizer, SingleOrchestrator, SinglePlane, require_default_tenant()) when composed with a resolver that reports more than one tenant. It is a composition error, not a client error: it always ends in 500.

Parameters:

tenant_id (str)

Return type:

None

periplo.tenancy.require_default_tenant(context)[source]

Guard a single-tenant default, which only ever serves DEFAULT_TENANT.

Raises:

ForeignTenant – When context belongs to a tenant other than DEFAULT_TENANT.

Parameters:

context (RequestContext)

Return type:

None

class periplo.tenancy.AnonymousAuthenticator[source]

Bases: object

Default identity port: every caller is ANONYMOUS.

Implements Loom’s Authenticator protocol structurally; it never refuses a request, matching the open source’s lack of authentication.

Any other authenticator sees a header sent more than once as its last value only (RequestCredentials.headers is a plain mapping). It refuses by returning None or raising loom Unauthenticated (401), or by raising loom Forbidden (403).

Authorization and audit

Authorization and audit.

Two questions every state-affecting call answers, once identity and tenant are known (periplo.tenancy): may this happen (the Authorizer port) and who should know it happened (the AuditSink port). Both are composed behind Access, the single object routers call.

The open-core defaults (SwitchAuthorizer, LogAuditSink) implement single-tenant behaviour: PERIPLO_ETL_ALLOW_OPERATE gates operating ETLs, and audit events go to the structured log.

Every route asks in the same shape:

  • a listing: require(action, WHOLE), then one Access.visible() over its items;

  • a single item: Access.reveal(). With a plain Authorizer that is require on the item, and a refusal answers 403. With a FilteringAuthorizer it is require(action, WHOLE), then visible about the item alone, and a hidden item raises the route’s own not-found error, byte-identical to one that does not exist.

A 404 only ever comes from filtering: a refusal is never remapped to one.

A product names its own actions with a dotted prefix of its own ("shop.export"), so their values never collide with a core Action; Access refuses one that does.

periplo.access.AUDIT_WRITE_TIMEOUT_SECONDS: Final = 5.0

How long any audit write may take; see Access.operate and Access._record.

periplo.access.Target

What an action is about, as an Authorizer receives it.

The convention, shared by every route and by any private Authorizer:

  • () — the resource as a whole (list ETLs, read the catalog, run a query);

  • ("etl", name) — one deployment (etl_target());

  • ("run", run_id, etl) — one run and the deployment it belongs to, with etl empty when the orchestrator cannot tell or the deployment is gone (run_target());

  • ("table", database, table) — one catalog table (table_target());

  • ("query", query_id) — one admitted query (query_target()).

alias of tuple[str, …]

periplo.access.etl_of(target)[source]

The deployment an etl or run target is about; None for any other target.

Parameters:

target (tuple[str, ...])

Return type:

str | None

class periplo.access.Action(*values)[source]

Bases: StrEnum

The actions of the core. A product adds its own as another StrEnum.

ADMIN_CATALOG = 'admin_catalog'

The sources, their discovery reports and starting a discovery.

periplo.access.STATE_CHANGING: Final = frozenset({Action.OPERATE_ETL})

Core actions a denied audit event is worth recording for from require.

Access.operate records every denial it meets, whatever the action.

exception periplo.access.Denied(message='You are not allowed to do this', *, code='forbidden')[source]

Bases: Forbidden

Raised by an Authorizer to refuse an action, with a caller-chosen code.

code must map to 403 in HttpErrorMapper: either it already does (a code shared with another Denied), or this registers it. A code already mapped to a different status is a programming error, not a request to relax that other mapping. Registration is process-global and permanent.

Parameters:
Return type:

None

periplo.access.as_403(error)[source]

error itself when its code answers 403, else a Denied that does.

A Forbidden whose code Loom’s mapper does not know would end in 500 on the catalog, and one whose code maps elsewhere (say not_found) would answer that status: every surface must answer a refusal the same way, with 403. An unregistered code is kept (Denied registers it, process-global and permanent); one already mapped to another status is replaced by the generic forbidden.

Parameters:

error (Forbidden)

Return type:

Forbidden

periplo.access.as_401(error)[source]

error itself when its code answers 401, else one that does.

The counterpart of as_403(): an unregistered code is registered as 401 (process-global and permanent, like Denied) and kept; one already mapped to another status is replaced by the generic unauthenticated.

Parameters:

error (Unauthenticated)

Return type:

Unauthenticated

class periplo.access.Justification(*, grant)[source]

Bases: object

What an Authorizer names as the reason it allowed an action.

Parameters:

grant (str)

grant: str

The grant that allowed it, recorded as grant in the detail of the audit events of Access.operate.

class periplo.access.Authorizer(*args, **kwargs)[source]

Bases: Protocol

async authorize(context, action, target)[source]

Return to allow; raise loom Forbidden (or Denied) to refuse.

Anything else fails the request, never a silent allow. action is an Action or a product’s own StrEnum. target follows the Target convention; () asks whether the action is allowed at all. A run target carries its deployment, or "" when it is not known: deny "" unless runs without a deployment are granted explicitly. Keeping the run within the tenant is the OrchestratorProvider’s job.

Returns:

The Justification of the action, or None.

Parameters:
Return type:

Justification | None

class periplo.access.FilteringAuthorizer(*args, **kwargs)[source]

Bases: Authorizer, Protocol

An Authorizer that can also tell which items of a collection are visible.

async visible(context, action, targets)[source]

The targets the caller may see, asked once for the whole collection.

Only asked once authorize(context, action, WHOLE) has allowed the action. Raise like authorize to refuse the collection as a whole. Anything not among targets in the answer is ignored.

Parameters:
Return type:

Collection[tuple[str, …]]

class periplo.access.SwitchAuthorizer(*, allow_operate)[source]

Bases: object

Default Authorizer: the single-tenant allow_operate switch.

Parameters:

allow_operate (bool)

class periplo.access.AuditEvent(*, at, tenant, subject, action, target, outcome, detail=<factory>)[source]

Bases: Struct

One fact worth keeping about a state-changing or denied action.

Parameters:
action: str

The Action or product StrEnum member itself, encoded as its value.

detail: dict[str, str | int]

run_id | normalized sql, query_id, state, rows, bytes, code, and grant, from the Justification the authorizer named.

Never row data.

class periplo.access.AuditSink(*args, **kwargs)[source]

Bases: Protocol

async record(event)[source]

Raise when the event could not be stored.

Parameters:

event (AuditEvent)

Return type:

None

class periplo.access.LogAuditSink[source]

Bases: object

Default AuditSink: a structured log line, for any tenant.

Never logs a full SQL statement: only its sql_sha256 and the first 512 normalized characters. AuditEvent.detail itself keeps the full normalized SQL, for private sinks that need it.

class periplo.access.Access(authorizer, audit)[source]

Bases: object

The one object routers call to authorize an action and audit it.

Parameters:
async require(action, target=WHOLE)[source]

Authorize action for the current request, or raise a denial.

Any loom Forbidden is re-raised through as_403() and any Unauthenticated through as_401(), so each answers the same status on every surface.

Parameters:
Return type:

RequestContext

async visible(action, targets)[source]

The targets the current request may see, in their order.

Without a FilteringAuthorizer every target, and the authorizer is not asked: the listing’s own require(action, WHOLE) already decided. With one, a single visible call; its answer is kept only where it names one of targets. Refusals are mapped as in require(); any other exception propagates.

Parameters:
Return type:

list[tuple[str, …]]

async reveal(action, target, *, hidden)[source]

Authorize action on one item, or raise: a denial, or hidden() when hidden.

Without a FilteringAuthorizer this is require() on target. With one, require() on the whole resource, then visible() on target.

Returns:

The request’s context.

Parameters:
Return type:

RequestContext

async allows(action)[source]

Whether action would be authorized on the whole resource, without raising.

A Forbidden/Unauthenticated denial answers False; any other exception propagates: a broken authorizer is an error, not a False.

Parameters:

action (StrEnum)

Return type:

bool

async operate(action, target, work, *, describe=lambda _: ...)[source]

Authorize, audit and perform a state-changing action, of the core or a product.

  1. Authorization as in require; a denial is audited (best-effort) for any action and re-raised.

  2. The requested event, not captured and time-bounded: a sink that fails or runs out of time stops the action before work ever runs: an action that cannot be audited must not happen.

  3. work().

  4. A best-effort, shielded and time-bounded outcome event: succeeded; failed with the failure’s code when it has one; or cancelled when the request was cancelled mid-action. A describe that breaks is logged and the action still counts as succeeded: it has already happened.

Every event carries the authorizer’s grant, when it named one, as grant.

Parameters:
Return type:

_T

async operate_revealed(view, action, target, work, *, hidden, describe=lambda _: ...)[source]

operate() on an item the caller must also be allowed to view.

Without a FilteringAuthorizer this is operate(). With one, view is authorized on the whole resource and target must be visible to it first: a refusal is audited as a denied action and re-raised, and a hidden item is audited as a denied action with the code hidden before hidden() is raised.

Parameters:
Return type:

_T

async record_best_effort(event)[source]

Record event for a bounded time, logging (never raising) when the sink fails.

Unshielded: a request that is already cancelled may cut the write short.

Parameters:

event (AuditEvent)

Return type:

None

async record_shielded(event)[source]

Record event best-effort, even from a cancelled request, for a bounded time.

Shielded so a disconnect or cancellation still leaves the event; bounded by AUDIT_WRITE_TIMEOUT_SECONDS so a hung sink never holds the request (or its query slot) forever.

Parameters:

event (AuditEvent)

Return type:

None

Storage credentials

Storage credentials, per tenant.

Every read of storage opens it with the ReadCredentials a CredentialsProvider returns for the tenant it reads for: discovery, a table’s log, its statistics and history, and a query’s data files. Every route asks for them through one CredentialsGate, which bounds the wait and, for a product’s provider, refuses anything but explicit keys (check_credentials()).

The open-core default, ProcessCredentials, answers PROCESS_CREDENTIALS for every tenant: no storage options, so delta-rs and pyarrow use the process’s own credential chain.

periplo.credentials.ACCESS_KEY_ID_KEYS: Final = frozenset({'access_key_id', 'aws_access_key_id'})

The names object_store reads an access key id from, compared without case.

periplo.credentials.CONNECTION_KEYS: Final = frozenset({'aws_allow_http', 'aws_endpoint', 'aws_endpoint_url', 'aws_region', 'aws_virtual_hosted_style_request'})

Where and how to reach storage; none of them says where credentials come from.

periplo.credentials.REDIRECTING_ENVIRONMENT: Final = frozenset({'allow_http', 'aws_allow_http', 'aws_endpoint', 'aws_endpoint_url', 'aws_endpoint_url_s3', 'aws_s3_allow_unsafe_rename', 'aws_s3_express', 'aws_skip_signature', 'aws_unsigned_payload', 'aws_virtual_hosted_style_request', 'endpoint', 'endpoint_url', 's3_express', 'skip_signature', 'unsigned_payload', 'virtual_hosted_style_request'})

Environment variables delta-rs reads, in any case, that change where a tenant read goes or how it is signed, whatever keys it carries.

periplo.credentials.ENDPOINT_KEYS: Final = frozenset({'aws_endpoint', 'aws_endpoint_url'})

The names of the endpoint; at most one of them may be given.

periplo.credentials.EXPIRY_MARGIN: Final = datetime.timedelta(seconds=60)

How long credentials must still work when a provider hands them over.

periplo.credentials.PROCESS_CREDENTIALS: Final = ReadCredentials(<0 storage options>, expires_at=None)

the process’s own credential chain, as with no provider at all.

Type:

No storage options

periplo.credentials.SECRET_ACCESS_KEY_KEYS: Final = frozenset({'aws_secret_access_key', 'secret_access_key'})

The names object_store reads a secret access key from, compared without case.

periplo.credentials.SESSION_TOKEN_KEYS: Final = frozenset({'aws_session_token', 'aws_token', 'session_token', 'token'})

The names object_store reads a session token from, compared without case.

class periplo.credentials.CredentialsGate(provider, *, timeout, checked, clock=lambda : ...)[source]

Bases: object

How every route asks for a tenant’s credentials.

The wait is bounded by timeout (504); a provider that breaks is answered 500 and logged by the type of its error only, since its message may carry a secret. With checked, what the provider answers goes through check_credentials() before anything reads with it.

Parameters:
async for_tenant(tenant)[source]

Raises CredentialsUnavailable, CredentialsTimeout or CredentialsFailed.

Parameters:

tenant (Tenant)

Return type:

ReadCredentials

class periplo.credentials.CredentialsProvider(*args, **kwargs)[source]

Bases: Protocol

async for_tenant(tenant)[source]

The credentials every read for tenant uses now.

Asked on the event loop, once per request that reads storage and once per discovery, through a CredentialsGate. Answer explicit keys only (see check_credentials()).

Raises:

CredentialsUnavailable – When there are none to give; answered 503. Anything else is answered 500, and logged by its type only. Neither is ever replaced by the process’s own credentials.

Parameters:

tenant (Tenant)

Return type:

ReadCredentials

exception periplo.credentials.CredentialsUnavailable(message='Storage credentials are not available right now')[source]

Bases: LoomError

A tenant’s storage credentials are missing or unusable: retry shortly (503).

Raised by a CredentialsProvider that has none to give, and by Periplo for credentials it will not read with. The message must already be safe to show.

Parameters:

message (str)

Return type:

None

class periplo.credentials.ProcessCredentials[source]

Bases: object

Default CredentialsProvider: PROCESS_CREDENTIALS for every tenant.

class periplo.credentials.ReadCredentials(storage_options, expires_at=None)[source]

Bases: object

What storage is opened with, for one tenant, until expires_at.

Not a dataclass: neither repr nor dataclasses.asdict ever shows the options.

Parameters:
  • storage_options (Mapping[str, str])

  • expires_at (datetime | None)

property storage_options: Mapping[str, str]

A read-only view of the object_store configuration keys.

property uses_process_chain: bool

Whether storage is reached with the process’s own credential chain.

same_storage(other)[source]

Whether other reaches storage with exactly the same options.

Parameters:

other (ReadCredentials)

Return type:

bool

storage_key()[source]

A hashable identity of the storage options. It holds the options themselves.

Return type:

tuple[tuple[str, str], …]

periplo.credentials.check_credentials(credentials, *, now)[source]

Refuse credentials that could let storage fall back to any other identity.

Accepted: exactly one access key id, one secret access key and one session token, temporary credentials, under the names object_store reads (ACCESS_KEY_ID_KEYS, SECRET_ACCESS_KEY_KEYS, SESSION_TOKEN_KEYS), plus any of CONNECTION_KEYS with at most one of ENDPOINT_KEYS; names compared without case, no value empty; and an expires_at later than now plus EXPIRY_MARGIN. Anything else, a profile, a role or keys without a token among them, would let delta-rs complete the options from the process’s environment.

Raises:

CredentialsUnavailable – Naming neither a key nor a value.

Parameters:
Return type:

None

periplo.credentials.check_environment(environ)[source]

Refuse to serve a credentials provider from an environment that redirects its reads.

The process’s own credentials are allowed: every answer of a provider carries its keys, so no tenant read is signed with them. What is refused is any of REDIRECTING_ENVIRONMENT, which delta-rs would apply to every read.

Raises:

ConfigError – Naming the variables, never their values.

Parameters:

environ (Mapping[str, str])

Return type:

None

periplo.credentials.check_storage_options(options)[source]

The shape half of check_credentials(): explicit keys, and nothing else.

Raises:

CredentialsUnavailable – Naming neither a key nor a value.

Parameters:

options (Mapping[str, str])

Return type:

None

periplo.credentials.option(options, names)[source]

The value of whichever of names is in options, compared without case.

Parameters:
Return type:

str | None

Data planes

Data planes: everything that reads one tenant’s data.

A DataPlane holds a tenant’s sources and catalog, its metadata readers and its query engine, bound to that tenant. Every catalog and query route resolves the plane of the request’s tenant with plane_for() before it authorizes anything, and reads nothing outside it. Nothing in one plane is shared with another but the byte budget of the metadata cache, where each plane keys its own values.

The open-core default, SinglePlane, serves one plane to the single DEFAULT_TENANT. periplo.planes.build_data_plane builds a plane the way the default one is built.

class periplo.data_plane.DataPlane(*, tenant, configuration, state, metadata, stats, history, engine, on_close)[source]

Bases: object

One tenant’s sources, catalog, metadata readers and query engine.

Parameters:
  • tenant (Tenant)

  • configuration (Configuration)

  • state (CatalogState)

  • metadata (TableMetadataReader)

  • stats (TableStatsReader)

  • history (TableHistoryReader)

  • engine (QueryEngine)

  • on_close (Callable[[], Awaitable[None]])

tenant: Tenant

The only tenant this plane ever serves.

configuration: Configuration

The sources and the labels the catalog is grouped and described by.

on_close: Callable[[], Awaitable[None]]

Waits for the plane’s background reads; see aclose().

async aclose()[source]

Release the plane on shutdown, once its background reads are done.

Return type:

None

exception periplo.data_plane.NoDataPlane(tenant_id)[source]

Bases: Exception

Raised by a DataPlanes for a tenant it has no plane for.

A composition error, not a client error: it always ends in 500, and is never answered with another tenant’s plane.

Parameters:

tenant_id (str)

Return type:

None

exception periplo.data_plane.WrongDataPlane[source]

Bases: Exception

A DataPlanes answered a plane bound to another tenant (500).

Return type:

None

class periplo.data_plane.DataPlanes(*args, **kwargs)[source]

Bases: Protocol

async for_tenant(tenant)[source]

The data plane every catalog and query read for tenant uses.

Raises:

NoDataPlane – When tenant has none (500). Anything else fails the request as well; neither is ever answered with another plane.

Parameters:

tenant (Tenant)

Return type:

DataPlane

async aclose()[source]

Release every plane on shutdown, with DataPlane.aclose().

Return type:

None

async periplo.data_plane.plane_for(planes, tenant)[source]

The plane planes resolves for tenant, only when it is bound to tenant.

Raises:

WrongDataPlane – When the plane answered is bound to another tenant.

Parameters:
Return type:

DataPlane

class periplo.data_plane.SinglePlane(plane)[source]

Bases: object

Default DataPlanes: one plane, for DEFAULT_TENANT.

A tenant boundary like SwitchAuthorizer: any other tenant raises ForeignTenant.

Parameters:

plane (DataPlane)

periplo.planes.build_data_plane(tenant, configuration, storage, settings, *, cache, opener=open_table)[source]

The data plane of tenant, built the way the default one is.

Parameters:
  • tenant (Tenant) – The only tenant the plane serves; its cached metadata is keyed by it, and the default tenant’s by nothing, as a single plane’s always was.

  • configuration (Configuration) – The plane’s sources and labels.

  • storage (Callable[[ReadCredentials], Storage]) – The storage that lists and probes with the given credentials, such as periplo.catalog.adapters.s3_lister.S3Storage.

  • settings (Settings) – Discovery, snapshot and cache limits.

  • cache (CacheGateway) – The metadata cache, from metadata_cache(), shared by every plane.

  • opener (Callable[[str, ReadCredentials], DeltaTable]) – How a table is opened; the default is the only one outside tests.

Return type:

DataPlane

periplo.planes.metadata_cache(settings)[source]

A new metadata cache, bounded by settings.metadata_cache_bytes.

Each call builds a cache of its own and leaves every cache handed out before as it is: a gateway keeps the cache it was built on. Call it once for all your planes.

Parameters:

settings (Settings)

Return type:

CacheGateway

class periplo.catalog.adapters.s3_lister.S3Storage(*, max_listers=MAX_LISTERS, clock=lambda : ...)[source]

Bases: object

StorageFor over S3: the lister of each read’s credentials.

The process’s own credentials share one lister for good. Any others get one of their own, kept until they expire, so each bucket’s filesystem is built once per set of credentials rather than once per probe.

Parameters:
  • max_listers (int)

  • clock (Callable[[], datetime])

ETL orchestrators

The ETL orchestrator seen per tenant.

Which orchestrator answers a request is itself a question with a port, OrchestratorProvider, and a single-tenant open-core default, SingleOrchestrator, that wraps whatever this installation was started with (or nothing, when the integration is off).

class periplo.etl.provider.OrchestratorProvider(*args, **kwargs)[source]

Bases: Protocol

async for_tenant(tenant)[source]

The orchestrator that serves tenant, or None when ETL is off for it.

None answers 404 etl_not_configured; anything raised here is a composition error, not a request outcome, and ends in 500: it never falls back to the default.

A provider serving more than one tenant is what keeps each tenant within its own deployments and runs (run routes are authorized by run id alone), so every orchestrator it hands out must be bounded by that tenant’s tags: an adapter built with require_tags=True, which refuses to exist without them.

Parameters:

tenant (Tenant)

Return type:

Orchestrator | None

async aclose()[source]

Release whatever orchestrators this provider holds, on shutdown.

Return type:

None

class periplo.etl.provider.SingleOrchestrator(orchestrator)[source]

Bases: object

Default OrchestratorProvider: the one orchestrator of a single-tenant install.

Any tenant other than DEFAULT_TENANT is a composition error for this default: it is a tenant boundary, like SwitchAuthorizer.

Parameters:

orchestrator (Orchestrator | None)

class periplo.etl.ports.Orchestrator(*args, **kwargs)[source]

Bases: Protocol

A workflow orchestrator seen through the configured tags, which bound every query.

The adapter is the tenant boundary: it only ever returns deployments and runs within its tags, and drops (and logs) anything else the orchestrator answers. With no tags it sees the whole workspace, which only a single-tenant install may want: an adapter handed out by a multi-tenant OrchestratorProvider must refuse to be built without them (the adapter’s require_tags=True).

Every method raises an EtlError subclass on failure (Unknown, Ambiguous, Upstream, Rejected); nothing else escapes.

async list_deployments(only=None)[source]

Every visible deployment, ordered by name then flow_name, with the summary.

With only, just the deployments named in it: the summary, running runs, upcoming runs and histories count those alone.

Parameters:

only (Collection[str] | None)

Return type:

EtlList

async list_runs(name, limit)[source]

The newest limit runs of the deployment name, most recently started first.

Scheduled runs, which have no start_at, come last.

Parameters:
Return type:

list[FlowRun]

async get_run(run_id)[source]

The run run_id, if it exists within the configured tags.

Parameters:

run_id (str)

Return type:

RunDetail

async get_logs(run_id, after, limit, task_runs=None, q=None, min_level=None)[source]

At most limit (never above 200) log lines of the run run_id.

Without after: the last limit lines. With after (a next from a previous page, opaque): the first limit lines whose timestamp is at or after it; the boundary line is repeated and the client deduplicates by id. Without task_runs: flow-level lines only. With it: those task runs’ lines (several make up a process’s scope: its marker plus its steps). q searches the orchestrator’s own log text (at most 200 characters); min_level is a floor.

Parameters:
  • run_id (str)

  • after (str | None)

  • limit (int)

  • task_runs (list[str] | None)

  • q (str | None)

  • min_level (int | None)

Return type:

LogPage

async get_tasks(run_id)[source]

The run’s attempts, grouped into processes and steps.

Parameters:

run_id (str)

Return type:

RunTasks

async get_grid(name, limit)[source]

The limit most recent non-SCHEDULED runs of name as a process grid.

Parameters:
Return type:

RunGrid

async get_step(run_id, task_run_id)[source]

One step of the run: its facts and the first page of its own logs.

Parameters:
  • run_id (str)

  • task_run_id (str)

Return type:

StepDetail

async create_run(name, parameters)[source]

Launch a run of the deployment name with parameters, or its own when None.

Parameters:
Return type:

RunDetail

async set_schedule(name, active)[source]

Activate or deactivate every schedule of the deployment name; re-read it.

Resuming a deployment that is paused also un-pauses it. Raises Rejected when the deployment has no schedule at all.

Parameters:
Return type:

Deployment

async aclose()[source]

Release the underlying connections.

Return type:

None

class periplo.etl.ports.RunResolver(*args, **kwargs)[source]

Bases: Protocol

An Orchestrator that can tell which deployment a run belongs to.

With it, a run is authorized with its deployment in the target (("run", run_id, etl)); without it, by the run id alone.

async run_etl(run_id)[source]

The name of the deployment of run_id; None when that deployment is gone.

Raises Unknown for a run that does not exist within the configured tags.

Parameters:

run_id (str)

Return type:

str | None