Python API reference (unreleased)
This reference is generated from source annotations and docstrings. Only implemented public exports are listed in the rendered page. The API contract describes the full first-release scope; a planned name there is not an implemented API.
Reading the reference
Use the arrows in On this page to browse a section, module, then class or function. Search finds interface names and descriptions; direct links open the matching outline branch. The usage guide provides worked examples.
On GitHub, this file contains generation directives. Read the complete online reference or follow the local build and deployment guide.
Package entrypoints
rpkiparrot
RPKI payload management and validation, without import-time I/O.
__version__ = '0.1.0rc1'
module-attribute
Client
Online client bound to one host-owned AnyIO loop and task tree.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
config
|
ClientConfig
|
Complete explicit configuration, allowing initially empty groups. |
required |
transports
|
Mapping[str, TransportFactory] | None
|
RTR source ID to fresh-stream factory. Bindings are copied. |
None
|
readers
|
Mapping[str, JsonReader] | None
|
JSON source ID to explicitly registered JsonReader; custom bindings require reader_profile_id and never alter global formats. |
None
|
persistence
|
PersistenceBackend | None
|
Optional synchronous complete-state backend owned by this Client's dedicated worker thread. Mutually exclusive with configured persistence. All methods, including close, execute on that thread. |
None
|
Constructing does no I/O. Enter the async context to validate dependencies
and prepare resources, then use await task_group.start(client.run).
Cancellation of a wait or subscription does not stop source synchronization.
__init__(config: ClientConfig, *, transports: Mapping[str, TransportFactory] | None = None, readers: Mapping[str, JsonReader] | None = None, persistence: PersistenceBackend | None = None) -> None
__aenter__() -> Client
async
__aexit__(exc_type: type[BaseException] | None, exc: BaseException | None, tb: TracebackType | None) -> None
async
run(*, task_status: TaskStatus[None] = anyio.TASK_STATUS_IGNORED) -> None
async
Run source and expiry tasks within the caller's task group, once.
started signals operational infrastructure, not payload readiness. Unexpected implementation failures propagate to the host task group; expected individual source failures retry without stopping other sources.
apply_config(config: ClientConfig, *, expected_revision: ConfigRevision, transports: Mapping[str, TransportFactory] | None = None, readers: Mapping[str, JsonReader] | None = None) -> ConfigReceipt
async
Atomically replace running configuration after a revision comparison.
Added/replaced sources start unloaded; this call does not wait for them to synchronize. Old source incarnations cannot publish late work. A cancellation after the commit may lose the receipt but does not roll back the committed revision. Concurrent candidate preparation is bounded to one; conflicting, invalid or over-budget candidates leave state intact. Database identity and clock changes require constructing a new Client.
get_snapshot() -> Snapshot
async
Acquire the current coherent view without waiting for initial readiness.
Expired views are rebuilt under managed admission. Returned snapshots remain fixed and independently reject queries after their own boundary. Pure time work does not cancel admitted inputs. When all build slots are occupied, this call waits for a slot, including slots held by same-group source builds and their waiters. Cancelling the wait does not abandon an already running worker or extend data validity.
flush(*, timeout: float | None = None) -> PersistReceipt
async
Wait for the current complete state or a covering successor to persist.
Requires run and configured/injected persistence, otherwise raises ConfigurationError. PersistenceError reports a failed attempt or timeout; cancelling this waiter never abandons an active database transaction. The returned receipt identifies the actual durable version, which may be newer than the version captured by this call. An omitted timeout captures the current PersistenceConfig.flush_timeout, or 600 seconds for an injected backend. Hot changes apply only to subsequent calls. The wait budget starts after obtaining the target snapshot; it is not a hard wall-clock bound on this method or on managed transaction cleanup.
wait_ready(*, required: frozenset[PayloadKind], timeout: float = 30.0, aspa_afi: Afi | None = None) -> Snapshot
async
Wait for usable required data without stopping synchronization on timeout.
aspa_afi selects one ASPA family; None requires both to be usable, even when their authorizations differ. It requires ASPA in required.
watch() -> AsyncIterator[Subscription]
async
Atomically observe an initial snapshot and register its subsequent events.
No duplicate initial event is yielded. A slow subscriber receives sticky ResyncRequiredError and must close this context and observe again. initial_status captures current connection/persistence observations at registration without rewriting the fixed initial snapshot's payload.
get_status() -> StatusReport
Sample safe lifecycle/data diagnostics without performing network or SQL I/O.
get_metrics() -> MetricsReport
Sample completed operations and current usage for this publisher epoch.
Sync/connection and commit durations are accumulated monotonic seconds;
snapshot_build_seconds retains the latest completed build measurement;
persistence_commit_seconds retains the latest commit attempt duration.
A reconnect is a completed connection attempt
after the first, whether successful or failed. Cancellation is excluded.
Source record gauges use records/source/
aclose() -> None
async
Stop new work, cancel sources, and await every owned resource cleanup.
Idempotent; cleanup exceeding cleanup_timeout remains closing and appears overdue in get_status. Cancellation never abandons a worker or connection. A decided source revocation finishes under the same build budget before final persistence, even if its protocol worker reports after close begins. This necessary cleanup does not extend the original final flush deadline.
MemoryStore
Single-owner synchronous store with complete-source atomic replacement.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
context
|
EvaluationContext | None
|
Defaults to online. Offline requires an explicit reference time. |
None
|
limits
|
Limits | None
|
Finite admission and query budgets fixed on returned snapshots. |
None
|
Write methods and snapshot acquisition run on the constructing thread. Returned immutable snapshots can be queried by other threads. Construction performs no network/file I/O, and there is no implicit event loop or worker.
__init__(*, context: EvaluationContext | None = None, limits: Limits | None = None) -> None
replace_source(source_id: str, dataset: ParsedDataset, *, freshness: FreshnessPolicy) -> Snapshot
Validate and atomically replace a complete source, returning its snapshot.
Raises:
| Type | Description |
|---|---|
InputError
|
Incomplete dataset, missing original generation for max_age, future generation, or invalid offline reference time. |
ResourceLimitError
|
Source or total support budget would be exceeded. |
ConfigurationError
|
Wrong owner thread or freshness configuration. |
Failure preserves the previous state. Re-reading old data retains its original deadline; max_age never starts at this call's time.
remove_source(source_id: str) -> Snapshot
Atomically remove one source; unknown IDs are idempotent.
snapshot() -> Snapshot
Acquire the current view, rebuilding at support or source boundaries.
The returned view may be unavailable; readiness is not implied. Existing views remain unchanged and reject queries after their own deadlines.
Models
rpkiparrot.store
Atomic complete-source publication and immutable indexed snapshots.
MemoryStore
Single-owner synchronous store with complete-source atomic replacement.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
context
|
EvaluationContext | None
|
Defaults to online. Offline requires an explicit reference time. |
None
|
limits
|
Limits | None
|
Finite admission and query budgets fixed on returned snapshots. |
None
|
Write methods and snapshot acquisition run on the constructing thread. Returned immutable snapshots can be queried by other threads. Construction performs no network/file I/O, and there is no implicit event loop or worker.
__init__(*, context: EvaluationContext | None = None, limits: Limits | None = None) -> None
replace_source(source_id: str, dataset: ParsedDataset, *, freshness: FreshnessPolicy) -> Snapshot
Validate and atomically replace a complete source, returning its snapshot.
Raises:
| Type | Description |
|---|---|
InputError
|
Incomplete dataset, missing original generation for max_age, future generation, or invalid offline reference time. |
ResourceLimitError
|
Source or total support budget would be exceeded. |
ConfigurationError
|
Wrong owner thread or freshness configuration. |
Failure preserves the previous state. Re-reading old data retains its original deadline; max_age never starts at this call's time.
remove_source(source_id: str) -> Snapshot
Atomically remove one source; unknown IDs are idempotent.
snapshot() -> Snapshot
Acquire the current view, rebuilding at support or source boundaries.
The returned view may be unavailable; readiness is not implied. Existing views remain unchanged and reject queries after their own deadlines.
rpkiparrot.models
Public immutable payloads, evaluation contexts, results and snapshot identities.
Fields and validation behavior are generated from their implementation docstrings. See the data contract for atomicity, time windows and configured-source trust.
Network = IPv4Network | IPv6Network
module-attribute
Afi
Bases: IntEnum
IANA address-family identifiers for separate ASPA authorization scopes.
Values 1 and 2 differ from IP network versions 4 and 6. An explicit family never inherits authorization from the other family's ASPA records.
IPV4 = 1
class-attribute
instance-attribute
IPv4 authorization scope (IANA AFI 1).
IPV6 = 2
class-attribute
instance-attribute
IPv6 authorization scope (IANA AFI 2).
Aspa
dataclass
One original ASPA assertion, before cross-record/source union.
The customer is nonzero. Providers are nonempty, unique, exclude the customer, and may contain zero only as the sole no-provider marker. afi=None asserts the same authorization for both IPv4 and IPv6; it never means unknown scope. An explicit Afi applies only to that address family.
customer: int
instance-attribute
Nonzero unsigned 32-bit customer ASN.
providers: frozenset[int]
instance-attribute
Nonempty unique provider ASNs; {0} alone means no authorized providers.
afi: Afi | None = None
class-attribute
instance-attribute
Explicit authorization scope; None asserts the same authorization for both families.
__post_init__() -> None
__init__(*, customer: int, providers: frozenset[int], afi: Afi | None = None) -> None
AspaContext
dataclass
Required local and neighbor ASNs and explicitly supplied receiving role.
local_asn: int
instance-attribute
Nonzero ASN of the local receiving AS.
neighbor_asn: int
instance-attribute
Nonzero ASN of the neighbor supplying the path.
relationship: NeighborRelationship
instance-attribute
Neighbor role as seen by the local receiving AS.
__post_init__() -> None
__init__(local_asn: int, neighbor_asn: int, relationship: NeighborRelationship) -> None
AspaInput
dataclass
Batch path/context with correlation ID and an explicit ASPA family scope.
afi=None requests a common view valid for both families; an evaluator must reject this request if their effective authorizations or usability differ.
path: AsPath
instance-attribute
Reconstructed external AS path to evaluate.
context: AspaContext
instance-attribute
Explicit receiving AS, neighbor and relationship context.
input_id: str | None = None
class-attribute
instance-attribute
Optional caller correlation string preserved in the corresponding batch item.
afi: Afi | None = None
class-attribute
instance-attribute
Family to evaluate; None requires identical usable authorization views for both.
__post_init__() -> None
__init__(path: AsPath, context: AspaContext, input_id: str | None = None, afi: Afi | None = None) -> None
AspaProviders
dataclass
Effective customer union with per-provider configured-source provenance.
present=False means no attestation in an available data view. Provider
maps are copied and recursively immutable. If nonzero providers exist the
effective union suppresses AS0; the original source assertions retain it.
afi identifies the evaluated family; None denotes an explicit common view
whose authorization and usability agree for both address families.
customer: int
instance-attribute
Customer ASN queried in the available ASPA view.
present: bool
instance-attribute
Whether an attestation exists; False is permitted only in a usable data view.
providers: Mapping[int, tuple[RecordSupport, ...]]
instance-attribute
Effective provider ASN to immutable original supports; nonzero providers suppress AS0.
afi: Afi | None = None
class-attribute
instance-attribute
Evaluated family, or None for an explicitly verified identical view in both families.
__post_init__() -> None
__init__(*, customer: int, present: bool, providers: Mapping[int, tuple[RecordSupport, ...]], afi: Afi | None = None) -> None
AspaRecord
dataclass
Bases: _RecordTimes
One ASPA object's interval; aggregate JSONExt entries use conservative bounds.
valid_from: datetime | None = None
class-attribute
instance-attribute
Optional inclusive original support start, as an aware UTC datetime.
expires_at: datetime | None = None
class-attribute
instance-attribute
Optional exclusive original support end; loading may impose an earlier source deadline.
upstream_label: str | None = None
class-attribute
instance-attribute
Optional producer label; it is not a configured source identity or trust decision.
aspa: Aspa
instance-attribute
Original customer-provider assertion, including its authorization scope.
__post_init__() -> None
__init__(*, valid_from: datetime | None = None, expires_at: datetime | None = None, upstream_label: str | None = None, aspa: Aspa) -> None
AspaResult
dataclass
Fixed-draft outcome with normalized path, evaluated afi and ramp evidence.
afi=None denotes an explicitly verified common view for both families.
status: AspaStatus
instance-attribute
Normal outcome of the fixed ASPA path algorithm.
path: AsPath
instance-attribute
Original reconstructed path supplied by the caller.
context: AspaContext
instance-attribute
Receiving context used by the algorithm.
normalized_path: tuple[int, ...]
instance-attribute
ASN sequence after the algorithm removes consecutive prepends.
reason_codes: tuple[str, ...]
instance-attribute
Stable reasons supporting this outcome.
snapshot_id: SnapshotId
instance-attribute
Fixed data version used for this result.
evaluated_at: datetime
instance-attribute
UTC time actually used for this path evaluation.
mode: EvaluationMode
instance-attribute
Online current-time or explicitly fixed offline evaluation.
algorithm_version: str = 'draft-ietf-sidrops-aspa-verification-28'
class-attribute
instance-attribute
Exact fixed ASPA verification specification identifier.
evidence: Mapping[str, object] | None = None
class-attribute
instance-attribute
Optional bounded ramp and authorization explanation.
afi: Afi | None = None
class-attribute
instance-attribute
Evaluated family, or None for an explicitly verified common view.
__post_init__() -> None
__init__(*, status: AspaStatus, path: AsPath, context: AspaContext, normalized_path: tuple[int, ...], reason_codes: tuple[str, ...], snapshot_id: SnapshotId, evaluated_at: datetime, mode: EvaluationMode, algorithm_version: str = 'draft-ietf-sidrops-aspa-verification-28', evidence: Mapping[str, object] | None = None, afi: Afi | None = None) -> None
AspaStatus
Bases: StrEnum
Fixed ASPA draft outcomes; unknown is missing attestation, not missing data.
VALID = 'valid'
class-attribute
instance-attribute
The fixed ASPA algorithm accepts the path in the supplied context.
INVALID = 'invalid'
class-attribute
instance-attribute
The fixed ASPA algorithm rejects the path in the supplied context.
UNKNOWN = 'unknown'
class-attribute
instance-attribute
Available data lacks the attestation needed for a conclusive path result.
AsPath
dataclass
Reconstructed external AS path, neighbor first and origin last.
segments: tuple[PathSegment, ...]
instance-attribute
Reconstructed external path segments, neighbor first and origin last.
__post_init__() -> None
sequence(asns: Iterable[int]) -> AsPath
classmethod
Construct an ordinary path; an empty path has zero segments.
__init__(segments: tuple[PathSegment, ...]) -> None
Availability
Bases: StrEnum
Data usability independent from transport connection state.
UNSUPPORTED = 'unsupported'
class-attribute
instance-attribute
The source or protocol explicitly does not supply this payload.
NOT_LOADED = 'not_loaded'
class-attribute
instance-attribute
No complete successful update has supplied this payload yet.
READY = 'ready'
class-attribute
instance-attribute
A complete usable view contains effective records.
READY_EMPTY = 'ready_empty'
class-attribute
instance-attribute
A complete usable view authoritatively contains no effective records.
EXPIRED = 'expired'
class-attribute
instance-attribute
Previously synchronized data can no longer supply a valid view.
UNAVAILABLE = 'unavailable'
class-attribute
instance-attribute
Availability was explicitly revoked, for example by a fatal protocol error.
BatchItem
dataclass
Bases: Generic[T]
One ordered result or safe error, never both; input_id is caller-owned.
index: int
instance-attribute
Zero-based position in the caller input sequence.
input_id: str | None
instance-attribute
Caller correlation string, or None when absent or itself invalid.
result: T | None = None
class-attribute
instance-attribute
Evaluated result, including normal invalid/unknown outcomes; None on execution error.
error: ErrorInfo | None = None
class-attribute
instance-attribute
Safe per-item input or availability error; exactly one of result and error is present.
__post_init__() -> None
__init__(*, index: int, input_id: str | None, result: T | None = None, error: ErrorInfo | None = None) -> None
BatchResult
dataclass
Bases: Generic[T]
One fixed snapshot's ordered batch; complete does not mean all valid.
snapshot_id: SnapshotId
instance-attribute
One fixed data version used for the entire batch.
items: tuple[BatchItem[T], ...]
instance-attribute
Ordered per-input results or safe errors.
complete: bool = True
class-attribute
instance-attribute
Whether every input was processed; this does not mean every route was valid.
success_count: int
property
Number of evaluated items, including normal invalid/unknown outcomes.
error_count: int
property
Number of per-item execution or input failures.
__post_init__() -> None
__init__(*, snapshot_id: SnapshotId, items: tuple[BatchItem[T], ...], complete: bool = True) -> None
BmpAnalysis
dataclass
Conservative scenario analysis with evaluated afi, never standard 'valid'.
afi=None denotes an explicitly verified common view for both families.
assessment: BmpAssessment
instance-attribute
Conservative BMP assessment, distinct from a standard ASPA status.
path: AsPath
instance-attribute
Original reconstructed path supplied for analysis.
context: BmpContext
instance-attribute
Known observation facts; missing facts are not silently inferred.
missing_context: tuple[str, ...]
instance-attribute
Names of facts absent or unsuitable for a conclusive assessment.
scenarios: tuple[AspaResult, ...]
instance-attribute
Standard ASPA results for the explicitly checked receiving scenarios.
snapshot_id: SnapshotId
instance-attribute
Fixed data version used for the scenario analysis.
evaluated_at: datetime
instance-attribute
UTC time used for the analysis.
mode: EvaluationMode
instance-attribute
Online current-time or explicitly fixed offline evaluation.
afi: Afi | None = None
class-attribute
instance-attribute
Evaluated family, or None for an explicitly verified common view.
__post_init__() -> None
__init__(*, assessment: BmpAssessment, path: AsPath, context: BmpContext, missing_context: tuple[str, ...], scenarios: tuple[AspaResult, ...], snapshot_id: SnapshotId, evaluated_at: datetime, mode: EvaluationMode, afi: Afi | None = None) -> None
BmpAssessment
Bases: StrEnum
Conservative BMP analysis, deliberately distinct from standard validation.
INVALID = 'invalid'
class-attribute
instance-attribute
Every applicable checked scenario is invalid with sufficient observation context.
NO_INVALID_EVIDENCE = 'no_invalid_evidence'
class-attribute
instance-attribute
All checked scenarios are valid; this is not a standard ASPA valid result.
INDETERMINATE = 'indeterminate'
class-attribute
instance-attribute
Missing context or inconclusive scenarios prevent a definite assessment.
BmpContext
dataclass
Known BMP facts; unknown direction or reconstruction prevents certainty.
observation_stage: str | None = None
class-attribute
instance-attribute
Known pre_policy or post_policy observation stage; None means unknown.
path_reconstructed: bool | None = None
class-attribute
instance-attribute
Whether the host reconstructed the external path, including AS4 information.
path_direction: str | None = None
class-attribute
instance-attribute
Known neighbor_to_origin ordering; unknown or other values prevent certainty.
local_asn: int | None = None
class-attribute
instance-attribute
Local receiving ASN, if known from the observation context.
neighbor_asn: int | None = None
class-attribute
instance-attribute
ASN of the supplying neighbor, if known.
relationship: NeighborRelationship | None = None
class-attribute
instance-attribute
Known neighbor role; None requests analysis across applicable scenarios.
__post_init__() -> None
__init__(*, observation_stage: str | None = None, path_reconstructed: bool | None = None, path_direction: str | None = None, local_asn: int | None = None, neighbor_asn: int | None = None, relationship: NeighborRelationship | None = None) -> None
ConfigRevision
dataclass
Configuration CAS identity, independent from snapshot generation.
store_id: str
instance-attribute
Configuration owner identifier used in compare-and-swap checks.
epoch: str
instance-attribute
Configuration owner lifetime; not necessarily the observing snapshot epoch.
revision: int
instance-attribute
Positive configuration sequence within this owner and epoch.
__post_init__() -> None
__init__(store_id: str, epoch: str, revision: int) -> None
ConnectionState
Bases: StrEnum
Transport state; retrying sources can still have usable data.
IDLE = 'idle'
class-attribute
instance-attribute
No connection attempt is currently running.
CONNECTING = 'connecting'
class-attribute
instance-attribute
A connection attempt or handshake is in progress.
CONNECTED = 'connected'
class-attribute
instance-attribute
The transport is connected; payload readiness is a separate property.
RETRYING = 'retrying'
class-attribute
instance-attribute
Synchronization failed and may retry; old data keeps its original lifetime.
STOPPED = 'stopped'
class-attribute
instance-attribute
The source worker has stopped; this does not extend or revoke data lifetime.
Diagnostic
dataclass
Input observation with stable code, human message and optional field path.
code: str
instance-attribute
Stable identifier for the parsing or completeness observation.
message: str
instance-attribute
Human-readable observation text, not a stable matching interface.
field: str | None = None
class-attribute
instance-attribute
Optional producer field path associated with the observation.
__post_init__() -> None
__init__(*, code: str, message: str, field: str | None = None) -> None
EvaluationContext
dataclass
Evaluation mode and optional aware offline reference time.
Use online() or offline(at=...). Offline evaluation assesses only the
supplied complete dataset; it cannot reconstruct historical global state.
mode: EvaluationMode
instance-attribute
Online lifetime checks or fixed offline evaluation.
reference_time: datetime | None = None
class-attribute
instance-attribute
Aware offline reference time, normalized to UTC; must be None online.
__post_init__() -> None
online() -> EvaluationContext
classmethod
Use current UTC and monotonic time without extending data lifetime.
offline(*, at: datetime) -> EvaluationContext
classmethod
Evaluate at an explicit timezone-aware reference time.
__init__(*, mode: EvaluationMode, reference_time: datetime | None = None) -> None
EvaluationMode
Bases: StrEnum
Online current-time evaluation or explicitly fixed offline evaluation.
ONLINE = 'online'
class-attribute
instance-attribute
Check the current online lifetime on every operation.
OFFLINE = 'offline'
class-attribute
instance-attribute
Evaluate at the explicit reference time, without claiming current validity.
NeighborRelationship
Bases: StrEnum
Neighbor's role as seen by the local receiving AS.
CUSTOMER = 'customer'
class-attribute
instance-attribute
The receiving AS sees the neighbor as its customer.
PEER = 'peer'
class-attribute
instance-attribute
The receiving AS sees the neighbor as its peer.
PROVIDER = 'provider'
class-attribute
instance-attribute
The receiving AS sees the neighbor as its provider.
ROUTE_SERVER = 'route_server'
class-attribute
instance-attribute
The neighbor is a route server; the fixed algorithm permits its ASN omission.
ROUTE_SERVER_CLIENT = 'route_server_client'
class-attribute
instance-attribute
The neighbor is a route-server client as seen by the local route server.
OriginInput
dataclass
Batch route input; scalar validation occurs per item during evaluation.
prefix: Network | str
instance-attribute
Route network or strict CIDR text; evaluated and checked per batch item.
asn: int | None
instance-attribute
Claimed origin ASN, or None when no origin is available.
input_id: str | None = None
class-attribute
instance-attribute
Optional caller correlation string preserved in the corresponding batch item.
__init__(prefix: Network | str, asn: int | None, input_id: str | None = None) -> None
OriginResult
dataclass
RFC 6811 outcome, fixed snapshot identity and per-route evaluation time.
status: OriginStatus
instance-attribute
Normal RFC 6811 validation outcome on available data.
prefix: Network
instance-attribute
Canonical route network that was evaluated.
asn: int | None
instance-attribute
Evaluated origin ASN, or None for a missing origin.
reason_codes: tuple[str, ...]
instance-attribute
Stable reasons supporting this outcome.
snapshot_id: SnapshotId
instance-attribute
Fixed data version used for this result.
evaluated_at: datetime
instance-attribute
UTC time actually used for the route evaluation.
mode: EvaluationMode
instance-attribute
Online current-time or explicitly fixed offline evaluation.
evidence: Mapping[str, object] | None = None
class-attribute
instance-attribute
Optional bounded explanation; None when explanation was not requested.
__post_init__() -> None
__init__(*, status: OriginStatus, prefix: Network, asn: int | None, reason_codes: tuple[str, ...], snapshot_id: SnapshotId, evaluated_at: datetime, mode: EvaluationMode, evidence: Mapping[str, object] | None = None) -> None
OriginStatus
Bases: StrEnum
RFC 6811 origin validation outcomes on available data.
VALID = 'valid'
class-attribute
instance-attribute
At least one covering VRP authorizes both the origin ASN and prefix length.
INVALID = 'invalid'
class-attribute
instance-attribute
VRPs cover the route, but none authorizes both origin and length.
NOTFOUND = 'notfound'
class-attribute
instance-attribute
The available VRP view contains no covering assertion.
ParsedDataset
dataclass
Complete normalized source candidate, from a reader or host application.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
present
|
frozenset[PayloadKind]
|
Explicitly supplied payload families; missing differs from empty. |
required |
vrps
|
tuple[VrpRecord, ...]
|
Original VRP support records, not a mutable index. |
()
|
aspas
|
tuple[AspaRecord, ...]
|
Original ASPA supports, before provider union. |
()
|
router_keys
|
tuple[RouterKeyRecord, ...]
|
Internal keys retained when supplied by a producer. |
()
|
generated_at
|
datetime | None
|
Original export time; never download time or file mtime. |
None
|
format_id
|
str
|
Producer format or |
'programmatic'
|
diagnostics
|
tuple[Diagnostic, ...]
|
Completeness and parsing observations. |
()
|
complete
|
bool
|
Host/adapter assertion of full required scope. False candidates can be inspected but cannot enter a MemoryStore, even offline. |
True
|
aspa_families
|
frozenset[Afi] | None
|
Explicitly covered ASPA address families, including empty results. None defaults to both families when ASPA is present, else neither. Every original record's scope must lie within this set; a family-independent record requires both. Frozen after construction. |
None
|
present: frozenset[PayloadKind]
instance-attribute
vrps: tuple[VrpRecord, ...] = ()
class-attribute
instance-attribute
aspas: tuple[AspaRecord, ...] = ()
class-attribute
instance-attribute
router_keys: tuple[RouterKeyRecord, ...] = ()
class-attribute
instance-attribute
generated_at: datetime | None = None
class-attribute
instance-attribute
format_id: str = 'programmatic'
class-attribute
instance-attribute
diagnostics: tuple[Diagnostic, ...] = ()
class-attribute
instance-attribute
complete: bool = True
class-attribute
instance-attribute
metadata: Mapping[str, object] = field(default_factory=dict)
class-attribute
instance-attribute
Copied immutable JSON-compatible producer metadata; it is not verified trust evidence.
aspa_families: frozenset[Afi] | None = None
class-attribute
instance-attribute
__post_init__() -> None
__init__(*, present: frozenset[PayloadKind], vrps: tuple[VrpRecord, ...] = (), aspas: tuple[AspaRecord, ...] = (), router_keys: tuple[RouterKeyRecord, ...] = (), generated_at: datetime | None = None, format_id: str = 'programmatic', diagnostics: tuple[Diagnostic, ...] = (), complete: bool = True, metadata: Mapping[str, object] = dict(), aspa_families: frozenset[Afi] | None = None) -> None
PathSegment
dataclass
Nonempty BGP path segment of unsigned, nonzero ASNs.
kind: PathSegmentKind
instance-attribute
BGP segment category; validation keeps unsupported categories explicit.
asns: tuple[int, ...]
instance-attribute
Nonempty ordered tuple of unsigned nonzero ASNs as supplied by the host.
__post_init__() -> None
__init__(kind: PathSegmentKind, asns: tuple[int, ...]) -> None
PathSegmentKind
Bases: StrEnum
BGP path segment kinds; confederations require external reconstruction.
SEQUENCE = 'sequence'
class-attribute
instance-attribute
An ordered AS_SEQUENCE segment.
SET = 'set'
class-attribute
instance-attribute
An unordered AS_SET segment, retained for explicit validation rejection.
CONFED_SEQUENCE = 'confed_sequence'
class-attribute
instance-attribute
A confederation sequence requiring external reconstruction.
CONFED_SET = 'confed_set'
class-attribute
instance-attribute
A confederation set requiring external reconstruction.
PayloadKind
Bases: StrEnum
Payload families; Router Keys are retained but have no dedicated query API.
VRP = 'vrp'
class-attribute
instance-attribute
Validated prefix and origin assertions used for ROV.
ASPA = 'aspa'
class-attribute
instance-attribute
Customer-provider authorizations with explicit address-family scope.
ROUTER_KEY = 'router_key'
class-attribute
instance-attribute
Keys retained for RTR, export and persistence; no dedicated query API.
RecordSupport
dataclass
Configured-source provenance and original bounded support interval.
source_id: str
instance-attribute
Configured source that supplies this original support.
expires_at: datetime
instance-attribute
Exclusive effective UTC end, bounded by original record and source deadlines.
generated_at: datetime | None = None
class-attribute
instance-attribute
Original producer or successful RTR update time, when known.
upstream_label: str | None = None
class-attribute
instance-attribute
Optional original producer label, distinct from source_id.
valid_from: datetime | None = None
class-attribute
instance-attribute
Inclusive original support start, when supplied; None adds no lower bound.
__post_init__() -> None
__init__(*, source_id: str, expires_at: datetime, generated_at: datetime | None = None, upstream_label: str | None = None, valid_from: datetime | None = None) -> None
SnapshotId
dataclass
Publisher identity; generation is ordered only within store_id and epoch.
store_id: str
instance-attribute
Publisher identifier; generation values from different publishers are incomparable.
epoch: str
instance-attribute
Publisher lifetime identifier, changed when an online publisher restarts.
generation: int
instance-attribute
Positive snapshot sequence within this store_id and epoch.
__post_init__() -> None
__init__(store_id: str, epoch: str, generation: int) -> None
SourceInfo
dataclass
Nonsecret source identity, connection, capability and freshness metadata.
aspa_capabilities declares both Afi keys. None copies the aggregate ASPA status to both, defaulting missing ASPA to unsupported. Explicit per-family statuses determine aggregate ASPA availability: both must be ready/empty; otherwise unavailable, expired, not_loaded, unsupported take that priority. A single ready family can never imply that both are available.
id: str
instance-attribute
Configured source identifier, distinct from producer labels.
identity: str
instance-attribute
Nonsecret digest of endpoint, interpretation and trust identity used for restoration.
kind: str
instance-attribute
Source category, such as rtr or json; not the producer format.
format_id: str | None
instance-attribute
Format of the last complete dataset; None before a successful commit.
capabilities: Mapping[PayloadKind, Availability]
instance-attribute
Availability of each payload family, independent of transport connection.
connection: ConnectionState = ConnectionState.IDLE
class-attribute
instance-attribute
Current transport or worker state; retrying data may remain usable.
reason_code: str | None = None
class-attribute
instance-attribute
Stable reason for the latest reported source state, when available.
last_sync_at: datetime | None = None
class-attribute
instance-attribute
UTC time of the last accepted complete synchronization, when known.
expires_at: datetime | None = None
class-attribute
instance-attribute
Exclusive source deadline; reconnects and restoration do not extend it.
generated_at: datetime | None = None
class-attribute
instance-attribute
Original dataset generation time, when supplied or established by RTR EOD.
protocol_version: int | None = None
class-attribute
instance-attribute
Negotiated RTR version of the complete dataset, or None for other inputs.
session_id: int | None = None
class-attribute
instance-attribute
RTR cache session identifier, when a complete session is available.
serial: int | None = None
class-attribute
instance-attribute
RTR serial of the complete dataset; interpreted only within its session.
last_error: ErrorInfo | None = None
class-attribute
instance-attribute
Most recent safe error information, when available.
aspa_capabilities: Mapping[Afi, Availability] | None = None
class-attribute
instance-attribute
Both AFI capabilities; explicit values determine aggregate ASPA availability.
__post_init__() -> None
__init__(*, id: str, identity: str, kind: str, format_id: str | None, capabilities: Mapping[PayloadKind, Availability], connection: ConnectionState = ConnectionState.IDLE, reason_code: str | None = None, last_sync_at: datetime | None = None, expires_at: datetime | None = None, generated_at: datetime | None = None, protocol_version: int | None = None, session_id: int | None = None, serial: int | None = None, last_error: ErrorInfo | None = None, aspa_capabilities: Mapping[Afi, Availability] | None = None) -> None
Vrp
dataclass
Validated prefix assertion, including AS0 assertions.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
prefix
|
Network
|
Strict IPv4/IPv6 network; host bits and interface zones are rejected.
Mapped IPv6 stays IPv6; JSON output uses |
required |
max_length
|
int
|
Maximum authorized length within this address family. |
required |
asn
|
int
|
Unsigned 32-bit ASN; bool is not accepted. |
required |
prefix: Network
instance-attribute
max_length: int
instance-attribute
asn: int
instance-attribute
__post_init__() -> None
__init__(*, prefix: Network, max_length: int, asn: int) -> None
VrpMatch
dataclass
One effective VRP and all its current source supports.
vrp: Vrp
instance-attribute
One deduplicated effective prefix assertion.
supports: tuple[RecordSupport, ...]
instance-attribute
All currently effective original supports for this assertion.
__post_init__() -> None
__init__(*, vrp: Vrp, supports: tuple[RecordSupport, ...]) -> None
VrpRecord
dataclass
Bases: _RecordTimes
VRP support interval and optional upstream label, without source identity.
valid_from and expires_at constrain the original object. The store
additionally applies the source deadline. Duplicate payloads may have
independent intervals and must not be collapsed into a longer lifetime.
valid_from: datetime | None = None
class-attribute
instance-attribute
Optional inclusive original support start, as an aware UTC datetime.
expires_at: datetime | None = None
class-attribute
instance-attribute
Optional exclusive original support end; loading may impose an earlier source deadline.
upstream_label: str | None = None
class-attribute
instance-attribute
Optional producer label; it is not a configured source identity or trust decision.
vrp: Vrp
instance-attribute
Original prefix assertion carried by this independently timed support.
__post_init__() -> None
__init__(*, valid_from: datetime | None = None, expires_at: datetime | None = None, upstream_label: str | None = None, vrp: Vrp) -> None
ConfigReceipt
dataclass
Atomic configuration commit result; readiness and persistence are separate.
A cancelled caller may miss a successful receipt; config_revision and watch events identify the committed result. changed=False leaves all IDs intact. Source ID tuples are deterministic and contain no injected objects/secrets.
changed: bool
instance-attribute
Whether a new configuration and snapshot were atomically committed.
config_revision: ConfigRevision
instance-attribute
Configuration owner revision that defines this view.
snapshot_id: SnapshotId
instance-attribute
Snapshot committed with this configuration; readiness and persistence may follow later.
added_source_ids: tuple[str, ...] = ()
class-attribute
instance-attribute
Sorted identifiers newly present in the configuration.
removed_source_ids: tuple[str, ...] = ()
class-attribute
instance-attribute
Sorted identifiers removed from the configuration.
restarted_source_ids: tuple[str, ...] = ()
class-attribute
instance-attribute
Sorted retained identifiers whose source identity or injected binding was replaced.
__post_init__() -> None
__init__(*, changed: bool, config_revision: ConfigRevision, snapshot_id: SnapshotId, added_source_ids: tuple[str, ...] = (), removed_source_ids: tuple[str, ...] = (), restarted_source_ids: tuple[str, ...] = ()) -> None
EventId
dataclass
Publisher epoch and event sequence, separate from snapshot generation.
store_id: str
instance-attribute
Publisher identifier for this event stream.
epoch: str
instance-attribute
Event-stream lifetime; restarting a publisher invalidates older cursors.
sequence: int
instance-attribute
Positive sequence within this publisher and epoch, including status-only events.
__post_init__() -> None
__init__(store_id: str, epoch: str, sequence: int) -> None
MetricsReport
dataclass
Explicit process-local metrics; counters restart with the publisher epoch.
Keys use fixed metric names and bounded source_id/kind labels. Counters count completed operations, gauges sample current usage, and durations are seconds. Pure validation functions do not mutate this or any hidden global collector.
sampled_at: datetime
instance-attribute
UTC time at which immutable metric values were sampled.
store_id: str
instance-attribute
Publisher or explicit collector identifier for these process-local metrics.
epoch: str
instance-attribute
Metric lifetime; counters do not continue across publisher restarts.
counters: Mapping[str, int] = field(default_factory=dict)
class-attribute
instance-attribute
Cumulative integer event counts using fixed metric names and bounded labels.
gauges: Mapping[str, float] = field(default_factory=dict)
class-attribute
instance-attribute
Sampled resource and state quantities; units depend on the documented metric name.
durations: Mapping[str, float] = field(default_factory=dict)
class-attribute
instance-attribute
Durations in seconds; each metric defines cumulative or latest-attempt semantics.
__post_init__() -> None
__init__(*, sampled_at: datetime, store_id: str, epoch: str, counters: Mapping[str, int] = dict(), gauges: Mapping[str, float] = dict(), durations: Mapping[str, float] = dict()) -> None
SnapshotEvent
dataclass
Committed snapshot transition and conservative route revalidation scope.
changed_vrps includes support-only changes. affected_prefixes cover every contained route, independent of ASN and max length. A large diff sets requires_full_revalidation and affected_kinds instead of silently truncating. Coalesced transitions identify both endpoints, never an invented intermediate generation. Sources contains changed SourceInfo entries or None for removal.
event_id: EventId
instance-attribute
Position of this transition in the bounded publisher event stream.
before_id: SnapshotId
instance-attribute
Snapshot immediately before this transition or coalesced interval.
after_id: SnapshotId
instance-attribute
Snapshot published by this transition or at the coalesced interval end.
before_config_revision: ConfigRevision
instance-attribute
Configuration revision attached to before_id.
after_config_revision: ConfigRevision
instance-attribute
Configuration revision attached to after_id.
reason: str
instance-attribute
Stable publication cause, such as synchronized or configuration_changed.
changed_vrps: tuple[Vrp, ...] = ()
class-attribute
instance-attribute
VRP values affected by payload or support changes, subject to event limits.
changed_customers: tuple[int, ...] = ()
class-attribute
instance-attribute
ASPA customer ASNs affected by authorization or support changes.
sources: Mapping[str, SourceInfo | None] = field(default_factory=dict)
class-attribute
instance-attribute
Changed source metadata keyed by source ID; None denotes removal.
affected_prefixes: tuple[Network, ...] = ()
class-attribute
instance-attribute
Conservative route scope: every route contained by these prefixes may need revalidation.
affected_customers: tuple[int, ...] = ()
class-attribute
instance-attribute
Conservative ASPA customer scope requiring path revalidation.
affected_kinds: frozenset[PayloadKind] = frozenset()
class-attribute
instance-attribute
Payload families affected, including when a bounded detailed diff is unavailable.
requires_full_revalidation: bool = False
class-attribute
instance-attribute
Whether consumers must revalidate all relevant routes for affected_kinds.
coalesced: bool = False
class-attribute
instance-attribute
Whether this event spans multiple actual commits; no intermediate IDs are invented.
__post_init__() -> None
__init__(*, event_id: EventId, before_id: SnapshotId, after_id: SnapshotId, before_config_revision: ConfigRevision, after_config_revision: ConfigRevision, reason: str, changed_vrps: tuple[Vrp, ...] = (), changed_customers: tuple[int, ...] = (), sources: Mapping[str, SourceInfo | None] = dict(), affected_prefixes: tuple[Network, ...] = (), affected_customers: tuple[int, ...] = (), affected_kinds: frozenset[PayloadKind] = frozenset(), requires_full_revalidation: bool = False, coalesced: bool = False) -> None
SourceUpdate
dataclass
One complete source synchronization offered to a SourceSink atomically.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
source
|
SourceInfo
|
Nonsecret identity, protocol/session metadata and original last_sync_at/expires_at. Both times are mandatory for a completed update; expires_at does not start at sink acceptance time. |
required |
dataset
|
ParsedDataset
|
Complete normalized original payloads and support windows. For RTR the format is rtr-v1/rtr-v2 and metadata['rtr'] retains the EOD refresh_interval/retry_interval/expire_interval for recovery. |
required |
reason
|
str
|
Stable publication cause, normally synchronized or incremental. |
'synchronized'
|
A sink rechecks source incarnation, completeness, resource limits and time. Constructing this object is not publication or a grant of online trust.
source: SourceInfo
instance-attribute
dataset: ParsedDataset
instance-attribute
reason: str = 'synchronized'
class-attribute
instance-attribute
__post_init__() -> None
__init__(*, source: SourceInfo, dataset: ParsedDataset, reason: str = 'synchronized') -> None
StatusEvent
dataclass
Connection or persistence observation, without a new data generation.
Persistence notifications carry the last acknowledged IDs and any current error. An empty sources mapping is valid for such notifications.
event_id: EventId
instance-attribute
Position of this observation in the publisher event stream.
observed_snapshot_id: SnapshotId
instance-attribute
Data snapshot current when this status observation was emitted.
sources: Mapping[str, SourceInfo]
instance-attribute
Changed source connection or diagnostic state; may be empty for persistence events.
persisted_snapshot_id: SnapshotId | None = None
class-attribute
instance-attribute
Last acknowledged durable snapshot identity, or None before any successful commit.
persisted_config_revision: ConfigRevision | None = None
class-attribute
instance-attribute
Configuration revision in the last acknowledged durable snapshot.
persistence_error: ErrorInfo | None = None
class-attribute
instance-attribute
Current safe persistence failure, if any; valid memory state remains independent.
__post_init__() -> None
__init__(*, event_id: EventId, observed_snapshot_id: SnapshotId, sources: Mapping[str, SourceInfo], persisted_snapshot_id: SnapshotId | None = None, persisted_config_revision: ConfigRevision | None = None, persistence_error: ErrorInfo | None = None) -> None
StatusReport
dataclass
Safe point-in-time lifecycle/data/persistence report, without doing I/O.
Lifecycle liveness and required data readiness are separate. Cleanup that exceeds its reporting threshold remains closing until all owned work exits; cleanup_resources and retiring_source_ids identify bounded pending cleanup. Persisted IDs describe an actual acknowledged state, not an intended commit. aspa_availability, aspa_usable_until and aspa_active_groups expose each address family independently; aggregate ASPA readiness requires both.
sampled_at: datetime
instance-attribute
UTC time at which this report was sampled.
lifecycle: Literal['created', 'open', 'running', 'closing', 'closed']
instance-attribute
Resource lifecycle state; closing persists until all managed work has exited.
snapshot_id: SnapshotId
instance-attribute
Fixed snapshot publication identity associated with this value.
config_revision: ConfigRevision
instance-attribute
Configuration owner revision that defines this view.
availability: Mapping[PayloadKind, Availability]
instance-attribute
Observed payload availability; online queries still perform their own lifetime checks.
usable_until: Mapping[PayloadKind, datetime | None]
instance-attribute
Exclusive next-recomputation boundary per payload, not an extension of source lifetime.
active_groups: Mapping[PayloadKind, str | None]
instance-attribute
Selected trust group per payload; aggregate ASPA is None if family groups differ.
sources: Mapping[str, SourceInfo]
instance-attribute
Current safe source observations, including changes after the last snapshot commit.
configured_limits: Limits
instance-attribute
Limits currently committed by the owning client.
algorithm_versions: Mapping[str, str]
instance-attribute
Exact fixed protocol and validation specification identifiers.
ready: bool
instance-attribute
Whether the owning client satisfies its configured required payloads.
aspa_availability: Mapping[Afi, Availability] = field(default_factory=dict)
class-attribute
instance-attribute
Observed availability independently for IPv4 and IPv6.
aspa_usable_until: Mapping[Afi, datetime | None] = field(default_factory=dict)
class-attribute
instance-attribute
Exclusive next-recomputation boundary for each ASPA family, or None if unavailable.
aspa_active_groups: Mapping[Afi, str | None] = field(default_factory=dict)
class-attribute
instance-attribute
Selected trust group per ASPA family, or None when no group is usable.
persisted_snapshot_id: SnapshotId | None = None
class-attribute
instance-attribute
Last acknowledged durable snapshot identity, or None before any successful commit.
persisted_config_revision: ConfigRevision | None = None
class-attribute
instance-attribute
Configuration revision in the last acknowledged durable snapshot.
recovered_from: SnapshotId | None = None
class-attribute
instance-attribute
Disk snapshot used during restoration, or None if no complete state was recovered.
persistence_error: ErrorInfo | None = None
class-attribute
instance-attribute
Current safe persistence failure, if any; valid memory state remains independent.
reconfiguring: bool = False
class-attribute
instance-attribute
Whether a configuration candidate is being prepared; it may not yet be committed.
retiring_source_ids: tuple[str, ...] = ()
class-attribute
instance-attribute
Removed or replaced source IDs whose owned cleanup has not finished.
cleanup_overdue: bool = False
class-attribute
instance-attribute
Whether cleanup exceeded its reporting threshold while work remains managed.
cleanup_elapsed: float = 0.0
class-attribute
instance-attribute
Seconds elapsed in the currently observed cleanup.
cleanup_resources: tuple[str, ...] = ()
class-attribute
instance-attribute
Bounded descriptions of owned resources still awaiting cleanup.
__post_init__() -> None
__init__(*, sampled_at: datetime, lifecycle: Literal['created', 'open', 'running', 'closing', 'closed'], snapshot_id: SnapshotId, config_revision: ConfigRevision, availability: Mapping[PayloadKind, Availability], usable_until: Mapping[PayloadKind, datetime | None], active_groups: Mapping[PayloadKind, str | None], sources: Mapping[str, SourceInfo], configured_limits: Limits, algorithm_versions: Mapping[str, str], ready: bool, aspa_availability: Mapping[Afi, Availability] = dict(), aspa_usable_until: Mapping[Afi, datetime | None] = dict(), aspa_active_groups: Mapping[Afi, str | None] = dict(), persisted_snapshot_id: SnapshotId | None = None, persisted_config_revision: ConfigRevision | None = None, recovered_from: SnapshotId | None = None, persistence_error: ErrorInfo | None = None, reconfiguring: bool = False, retiring_source_ids: tuple[str, ...] = (), cleanup_overdue: bool = False, cleanup_elapsed: float = 0.0, cleanup_resources: tuple[str, ...] = ()) -> None
Snapshot
dataclass
An immutable, fixed-generation source and effective payload view.
Public fields describe identity, configuration, original times, capabilities and selected groups. Queries check the relevant usable_until on every call; holding this object never preserves an expired attestation. Obtain a new snapshot from its owner after SnapshotExpiredError. The internal indexes, original source states and clock are not serialization or extension APIs. ASPA capabilities, deadlines and selected groups are also exposed per AFI. Aggregate ASPA readiness requires both families. Its aggregate active group is None when the two families select different groups.
id: SnapshotId
instance-attribute
Fixed publisher identity and generation; holding this object does not preserve lifetime.
config_revision: ConfigRevision
instance-attribute
Configuration owner revision that defines this view.
published_at: datetime
instance-attribute
UTC publication time, separate from original generation and support deadlines.
context: EvaluationContext
instance-attribute
Online or explicitly fixed offline evaluation context.
sources: Mapping[str, SourceInfo]
instance-attribute
Source metadata captured atomically with this snapshot; later status changes are separate.
active_groups: Mapping[PayloadKind, str | None]
instance-attribute
Selected trust group per payload; aggregate ASPA is None if family groups differ.
capabilities: Mapping[PayloadKind, Availability]
instance-attribute
Payload availability in this fixed view; empty and unavailable remain distinct.
usable_until: Mapping[PayloadKind, datetime | None]
instance-attribute
Exclusive next-recomputation boundary per payload, not an extension of source lifetime.
aspa_capabilities: Mapping[Afi, Availability] = field(default_factory=dict)
class-attribute
instance-attribute
Independent availability for IPv4 and IPv6 authorization scopes.
aspa_usable_until: Mapping[Afi, datetime | None] = field(default_factory=dict)
class-attribute
instance-attribute
Exclusive next-recomputation boundary for each ASPA family, or None if unavailable.
aspa_active_groups: Mapping[Afi, str | None] = field(default_factory=dict)
class-attribute
instance-attribute
Selected trust group per ASPA family, or None when no group is usable.
upstream_id: SnapshotId | None = None
class-attribute
instance-attribute
Original publisher snapshot identity retained by an observer or import, when available.
algorithm_versions: Mapping[str, str] = field(default_factory=lambda: {'rov': 'rfc6811+rfc8481', 'aspa': 'draft-ietf-sidrops-aspa-verification-28', 'aspa_profile': 'draft-ietf-sidrops-aspa-profile-29', 'rtr_v1': 'rfc8210', 'rtr_v2': 'draft-ietf-sidrops-8210bis-27', 'rtr_v2_legacy': 'draft-ietf-sidrops-8210bis-10', 'aspa_profile_legacy': 'draft-ietf-sidrops-aspa-profile-07', 'aspa_afi_semantics': 'rpkiparrot-afi-1', 'rtr_v2_legacy13': 'draft-ietf-sidrops-8210bis-13', 'aspa_profile_legacy13': 'draft-ietf-sidrops-aspa-profile-18'})
class-attribute
instance-attribute
Exact fixed protocol and validation specification identifiers.
policy_id: str = 'none'
class-attribute
instance-attribute
Applied policy identity; none explicitly denotes no local policy in this release.
__post_init__() -> None
__init__(*, id: SnapshotId, config_revision: ConfigRevision, published_at: datetime, context: EvaluationContext, sources: Mapping[str, SourceInfo], active_groups: Mapping[PayloadKind, str | None], capabilities: Mapping[PayloadKind, Availability], usable_until: Mapping[PayloadKind, datetime | None], aspa_capabilities: Mapping[Afi, Availability] = dict(), aspa_usable_until: Mapping[Afi, datetime | None] = dict(), aspa_active_groups: Mapping[Afi, str | None] = dict(), upstream_id: SnapshotId | None = None, algorithm_versions: Mapping[str, str] = (lambda: {'rov': 'rfc6811+rfc8481', 'aspa': 'draft-ietf-sidrops-aspa-verification-28', 'aspa_profile': 'draft-ietf-sidrops-aspa-profile-29', 'rtr_v1': 'rfc8210', 'rtr_v2': 'draft-ietf-sidrops-8210bis-27', 'rtr_v2_legacy': 'draft-ietf-sidrops-8210bis-10', 'aspa_profile_legacy': 'draft-ietf-sidrops-aspa-profile-07', 'aspa_afi_semantics': 'rpkiparrot-afi-1', 'rtr_v2_legacy13': 'draft-ietf-sidrops-8210bis-13', 'aspa_profile_legacy13': 'draft-ietf-sidrops-aspa-profile-18'})(), policy_id: str = 'none', _limits: Limits, _clock: Clock, _states: Mapping[str, SourceState], _vrps: tuple[VrpMatch, ...], _prefixes: Mapping[tuple[int, int, int], tuple[VrpMatch, ...]], _aspas: Mapping[tuple[Afi, int], AspaProviders], _aspa_sources: Mapping[tuple[Afi, int], frozenset[str]] = dict(), _groups: Mapping[str, tuple[int, tuple[str, ...]]] = dict(), _temporal: object | None = None) -> None
PayloadDiff
dataclass
Complete original-source difference for one payload kind.
Only-side entries use SupportDifference with the other side empty. Common entries retain both sides even when payloads match but expiration differs.
kind: PayloadKind
instance-attribute
Payload category compared using the original source records.
only_left: tuple[SupportDifference, ...]
instance-attribute
Payloads asserted only by the left source, retaining original support intervals.
only_right: tuple[SupportDifference, ...]
instance-attribute
Payloads asserted only by the right source, retaining original support intervals.
common: tuple[SupportDifference, ...]
instance-attribute
Shared payloads retaining both sources supports, including differing deadlines.
__post_init__() -> None
__init__(*, kind: PayloadKind, only_left: tuple[SupportDifference, ...], only_right: tuple[SupportDifference, ...], common: tuple[SupportDifference, ...]) -> None
SourceDiff
dataclass
Bounded full source comparison, with status and time metadata intact.
This result never chooses a trusted winner or silently truncates a side. Limits are inherited from the fixed snapshot used for comparison.
snapshot_id: SnapshotId
instance-attribute
Fixed snapshot whose original source states were compared.
left: SourceInfo
instance-attribute
Left source identity, capability and time metadata.
right: SourceInfo
instance-attribute
Right source identity, capability and time metadata.
payloads: Mapping[PayloadKind, PayloadDiff]
instance-attribute
Complete bounded differences keyed by payload kind; no silent truncation.
__post_init__() -> None
__init__(*, snapshot_id: SnapshotId, left: SourceInfo, right: SourceInfo, payloads: Mapping[PayloadKind, PayloadDiff]) -> None
SupportDifference
dataclass
One common payload with both original sources' bounded support intervals.
ASPA comparisons use one customer-provider edge per payload. Expired and suppressed AS0 assertions remain visible here; this is diagnostic data.
payload: Vrp | Aspa | RouterKey
instance-attribute
Original common or one-sided payload; ASPA uses one customer-provider edge with AFI.
left: tuple[RecordSupport, ...]
instance-attribute
Original bounded supports from the left source; empty when present only on the right.
right: tuple[RecordSupport, ...]
instance-attribute
Original bounded supports from the right source; empty when present only on the left.
__post_init__() -> None
__init__(*, payload: Vrp | Aspa | RouterKey, left: tuple[RecordSupport, ...], right: tuple[RecordSupport, ...]) -> None
Errors
rpkiparrot.errors
Stable, serializable failures; domain validation outcomes are not exceptions.
ErrorInfo
dataclass
Safe failure information, without traceback or credentials.
Attributes:
| Name | Type | Description |
|---|---|---|
code |
str
|
Stable machine-readable error category. |
message |
str
|
Human-readable explanation; not a matching interface. |
retryable |
bool
|
Whether retrying the same operation may succeed. |
source_id |
str | None
|
Affected configured source, when known. |
retry_after |
float | None
|
Minimum delay in seconds, when supplied. |
details |
Mapping[str, object]
|
Copied, recursively immutable safe diagnostic fields. |
code: str
instance-attribute
message: str
instance-attribute
retryable: bool = False
class-attribute
instance-attribute
source_id: str | None = None
class-attribute
instance-attribute
retry_after: float | None = None
class-attribute
instance-attribute
details: Mapping[str, object] = field(default_factory=dict)
class-attribute
instance-attribute
__init__(*, code: str, message: str, retryable: bool = False, source_id: str | None = None, retry_after: float | None = None, details: Mapping[str, object] = dict()) -> None
__post_init__() -> None
RpkiparrotError
Bases: Exception
Base failure with a stable code and safe, immutable diagnostic details.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
message
|
str
|
Explanation without secret input values. |
required |
code
|
str | None
|
Override the subclass's default machine-readable code. |
None
|
source_id
|
str | None
|
Configured source identifier. |
None
|
retryable
|
bool
|
Whether retry may succeed without configuration changes. |
False
|
retry_after
|
float | None
|
Minimum retry delay in seconds, if known. |
None
|
details
|
Mapping[str, object] | None
|
Safe diagnostic fields, copied on construction. |
None
|
Attributes:
| Name | Type | Description |
|---|---|---|
info |
Immutable, serializable failure details used by all subclasses. |
default_code = 'rpkiparrot_error'
class-attribute
instance-attribute
Default stable error category used unless the constructor supplies an explicit code.
info = ErrorInfo(code=code or self.default_code, message=message, source_id=source_id, retryable=retryable, retry_after=retry_after, details=details or {})
instance-attribute
code: str
property
Stable error code.
source_id: str | None
property
Source affected by the operation, if known.
retryable: bool
property
Whether retrying may succeed.
retry_after: float | None
property
Minimum retry delay in seconds, if known.
details: Mapping[str, object]
property
Immutable safe fields explaining the failure.
__init__(message: str, *, code: str | None = None, source_id: str | None = None, retryable: bool = False, retry_after: float | None = None, details: Mapping[str, object] | None = None) -> None
ConfigurationError
Bases: RpkiparrotError
Invalid configuration; correct it before retrying.
default_code = 'invalid_configuration'
class-attribute
instance-attribute
Default stable error category used unless the constructor supplies an explicit code.
info = ErrorInfo(code=code or self.default_code, message=message, source_id=source_id, retryable=retryable, retry_after=retry_after, details=details or {})
instance-attribute
code: str
property
Stable error code.
source_id: str | None
property
Source affected by the operation, if known.
retryable: bool
property
Whether retrying may succeed.
retry_after: float | None
property
Minimum retry delay in seconds, if known.
details: Mapping[str, object]
property
Immutable safe fields explaining the failure.
__init__(message: str, *, code: str | None = None, source_id: str | None = None, retryable: bool = False, retry_after: float | None = None, details: Mapping[str, object] | None = None) -> None
ConfigurationConflictError
Bases: ConfigurationError
Configuration CAS failed; obtain the current revision.
default_code = 'configuration_conflict'
class-attribute
instance-attribute
Default stable error category used unless the constructor supplies an explicit code.
info = ErrorInfo(code=code or self.default_code, message=message, source_id=source_id, retryable=retryable, retry_after=retry_after, details=details or {})
instance-attribute
code: str
property
Stable error code.
source_id: str | None
property
Source affected by the operation, if known.
retryable: bool
property
Whether retrying may succeed.
retry_after: float | None
property
Minimum retry delay in seconds, if known.
details: Mapping[str, object]
property
Immutable safe fields explaining the failure.
__init__(message: str, *, code: str | None = None, source_id: str | None = None, retryable: bool = False, retry_after: float | None = None, details: Mapping[str, object] | None = None) -> None
MissingExtraError
Bases: ConfigurationError
The explicitly requested optional component is not installed.
default_code = 'missing_extra'
class-attribute
instance-attribute
Default stable error category used unless the constructor supplies an explicit code.
info = ErrorInfo(code=code or self.default_code, message=message, source_id=source_id, retryable=retryable, retry_after=retry_after, details=details or {})
instance-attribute
code: str
property
Stable error code.
source_id: str | None
property
Source affected by the operation, if known.
retryable: bool
property
Whether retrying may succeed.
retry_after: float | None
property
Minimum retry delay in seconds, if known.
details: Mapping[str, object]
property
Immutable safe fields explaining the failure.
__init__(message: str, *, code: str | None = None, source_id: str | None = None, retryable: bool = False, retry_after: float | None = None, details: Mapping[str, object] | None = None) -> None
InputError
Bases: RpkiparrotError
Malformed, incomplete, or unsupported input; no partial commit occurred.
default_code = 'invalid_input'
class-attribute
instance-attribute
Default stable error category used unless the constructor supplies an explicit code.
info = ErrorInfo(code=code or self.default_code, message=message, source_id=source_id, retryable=retryable, retry_after=retry_after, details=details or {})
instance-attribute
code: str
property
Stable error code.
source_id: str | None
property
Source affected by the operation, if known.
retryable: bool
property
Whether retrying may succeed.
retry_after: float | None
property
Minimum retry delay in seconds, if known.
details: Mapping[str, object]
property
Immutable safe fields explaining the failure.
__init__(message: str, *, code: str | None = None, source_id: str | None = None, retryable: bool = False, retry_after: float | None = None, details: Mapping[str, object] | None = None) -> None
DataUnavailableError
Bases: RpkiparrotError
Required data is not usable; this is not a validation outcome.
default_code = 'not_loaded'
class-attribute
instance-attribute
Default stable error category used unless the constructor supplies an explicit code.
info = ErrorInfo(code=code or self.default_code, message=message, source_id=source_id, retryable=retryable, retry_after=retry_after, details=details or {})
instance-attribute
code: str
property
Stable error code.
source_id: str | None
property
Source affected by the operation, if known.
retryable: bool
property
Whether retrying may succeed.
retry_after: float | None
property
Minimum retry delay in seconds, if known.
details: Mapping[str, object]
property
Immutable safe fields explaining the failure.
__init__(message: str, *, code: str | None = None, source_id: str | None = None, retryable: bool = False, retry_after: float | None = None, details: Mapping[str, object] | None = None) -> None
SnapshotExpiredError
Bases: DataUnavailableError
The fixed view reached a change boundary; acquire a new snapshot.
default_code = 'snapshot_expired'
class-attribute
instance-attribute
Default stable error category used unless the constructor supplies an explicit code.
info = ErrorInfo(code=code or self.default_code, message=message, source_id=source_id, retryable=retryable, retry_after=retry_after, details=details or {})
instance-attribute
code: str
property
Stable error code.
source_id: str | None
property
Source affected by the operation, if known.
retryable: bool
property
Whether retrying may succeed.
retry_after: float | None
property
Minimum retry delay in seconds, if known.
details: Mapping[str, object]
property
Immutable safe fields explaining the failure.
__init__(message: str, *, code: str | None = None, source_id: str | None = None, retryable: bool = False, retry_after: float | None = None, details: Mapping[str, object] | None = None) -> None
ReadyTimeoutError
Bases: DataUnavailableError
Readiness wait expired without stopping background synchronization.
default_code = 'ready_timeout'
class-attribute
instance-attribute
Default stable error category used unless the constructor supplies an explicit code.
info = ErrorInfo(code=code or self.default_code, message=message, source_id=source_id, retryable=retryable, retry_after=retry_after, details=details or {})
instance-attribute
code: str
property
Stable error code.
source_id: str | None
property
Source affected by the operation, if known.
retryable: bool
property
Whether retrying may succeed.
retry_after: float | None
property
Minimum retry delay in seconds, if known.
details: Mapping[str, object]
property
Immutable safe fields explaining the failure.
__init__(message: str, *, code: str | None = None, source_id: str | None = None, retryable: bool = False, retry_after: float | None = None, details: Mapping[str, object] | None = None) -> None
TransportError
Bases: RpkiparrotError
Transport connection or I/O failed.
default_code = 'transport_error'
class-attribute
instance-attribute
Default stable error category used unless the constructor supplies an explicit code.
info = ErrorInfo(code=code or self.default_code, message=message, source_id=source_id, retryable=retryable, retry_after=retry_after, details=details or {})
instance-attribute
code: str
property
Stable error code.
source_id: str | None
property
Source affected by the operation, if known.
retryable: bool
property
Whether retrying may succeed.
retry_after: float | None
property
Minimum retry delay in seconds, if known.
details: Mapping[str, object]
property
Immutable safe fields explaining the failure.
__init__(message: str, *, code: str | None = None, source_id: str | None = None, retryable: bool = False, retry_after: float | None = None, details: Mapping[str, object] | None = None) -> None
ProtocolError
Bases: RpkiparrotError
RTR protocol failure; details may include a separate wire_error_code.
default_code = 'protocol_error'
class-attribute
instance-attribute
Default stable error category used unless the constructor supplies an explicit code.
info = ErrorInfo(code=code or self.default_code, message=message, source_id=source_id, retryable=retryable, retry_after=retry_after, details=details or {})
instance-attribute
code: str
property
Stable error code.
source_id: str | None
property
Source affected by the operation, if known.
retryable: bool
property
Whether retrying may succeed.
retry_after: float | None
property
Minimum retry delay in seconds, if known.
details: Mapping[str, object]
property
Immutable safe fields explaining the failure.
__init__(message: str, *, code: str | None = None, source_id: str | None = None, retryable: bool = False, retry_after: float | None = None, details: Mapping[str, object] | None = None) -> None
ResourceLimitError
Bases: RpkiparrotError
Operation exceeded an explicit limit; no truncated success is returned.
default_code = 'resource_limit'
class-attribute
instance-attribute
Default stable error category used unless the constructor supplies an explicit code.
info = ErrorInfo(code=code or self.default_code, message=message, source_id=source_id, retryable=retryable, retry_after=retry_after, details=details or {})
instance-attribute
code: str
property
Stable error code.
source_id: str | None
property
Source affected by the operation, if known.
retryable: bool
property
Whether retrying may succeed.
retry_after: float | None
property
Minimum retry delay in seconds, if known.
details: Mapping[str, object]
property
Immutable safe fields explaining the failure.
__init__(message: str, *, code: str | None = None, source_id: str | None = None, retryable: bool = False, retry_after: float | None = None, details: Mapping[str, object] | None = None) -> None
PersistenceError
Bases: RpkiparrotError
Persistence failed independently of the in-memory committed view.
default_code = 'persistence_failed'
class-attribute
instance-attribute
Default stable error category used unless the constructor supplies an explicit code.
info = ErrorInfo(code=code or self.default_code, message=message, source_id=source_id, retryable=retryable, retry_after=retry_after, details=details or {})
instance-attribute
code: str
property
Stable error code.
source_id: str | None
property
Source affected by the operation, if known.
retryable: bool
property
Whether retrying may succeed.
retry_after: float | None
property
Minimum retry delay in seconds, if known.
details: Mapping[str, object]
property
Immutable safe fields explaining the failure.
__init__(message: str, *, code: str | None = None, source_id: str | None = None, retryable: bool = False, retry_after: float | None = None, details: Mapping[str, object] | None = None) -> None
ResyncRequiredError
Bases: RpkiparrotError
Event continuity was lost; close the old subscription and observe anew.
default_code = 'resync_required'
class-attribute
instance-attribute
Default stable error category used unless the constructor supplies an explicit code.
info = ErrorInfo(code=code or self.default_code, message=message, source_id=source_id, retryable=retryable, retry_after=retry_after, details=details or {})
instance-attribute
code: str
property
Stable error code.
source_id: str | None
property
Source affected by the operation, if known.
retryable: bool
property
Whether retrying may succeed.
retry_after: float | None
property
Minimum retry delay in seconds, if known.
details: Mapping[str, object]
property
Immutable safe fields explaining the failure.
__init__(message: str, *, code: str | None = None, source_id: str | None = None, retryable: bool = False, retry_after: float | None = None, details: Mapping[str, object] | None = None) -> None
SnapshotGoneError
Bases: RpkiparrotError
A remote snapshot was evicted; restart the entire paginated operation.
default_code = 'snapshot_gone'
class-attribute
instance-attribute
Default stable error category used unless the constructor supplies an explicit code.
info = ErrorInfo(code=code or self.default_code, message=message, source_id=source_id, retryable=retryable, retry_after=retry_after, details=details or {})
instance-attribute
code: str
property
Stable error code.
source_id: str | None
property
Source affected by the operation, if known.
retryable: bool
property
Whether retrying may succeed.
retry_after: float | None
property
Minimum retry delay in seconds, if known.
details: Mapping[str, object]
property
Immutable safe fields explaining the failure.
__init__(message: str, *, code: str | None = None, source_id: str | None = None, retryable: bool = False, retry_after: float | None = None, details: Mapping[str, object] | None = None) -> None
ClosedError
Bases: RpkiparrotError
The stateful object cannot be used after closing.
default_code = 'closed'
class-attribute
instance-attribute
Default stable error category used unless the constructor supplies an explicit code.
info = ErrorInfo(code=code or self.default_code, message=message, source_id=source_id, retryable=retryable, retry_after=retry_after, details=details or {})
instance-attribute
code: str
property
Stable error code.
source_id: str | None
property
Source affected by the operation, if known.
retryable: bool
property
Whether retrying may succeed.
retry_after: float | None
property
Minimum retry delay in seconds, if known.
details: Mapping[str, object]
property
Immutable safe fields explaining the failure.
__init__(message: str, *, code: str | None = None, source_id: str | None = None, retryable: bool = False, retry_after: float | None = None, details: Mapping[str, object] | None = None) -> None
Configuration
rpkiparrot.config
Explicit immutable configuration; constructing it performs no I/O.
SourceConfig: TypeAlias = RtrSourceConfig | JsonFileSourceConfig | HttpJsonSourceConfig
module-attribute
A supported complete-source configuration; source IDs are globally unique.
ServiceConfig
dataclass
Bounded asyncio service configuration; constructing it performs no I/O.
host/port describe the listener the host must actually configure. workers must be one. Non-loopback listeners require bearer authentication and TLS. Direct TLS requires certificate/key paths and HTTPS ASGI requests. Proxy TLS requires bearer authentication, explicit peer IP/CIDR allowlists and one https forwarded protocol header; disable the ASGI server's proxy-header rewriting. Unauthenticated loopback mode also checks the actual request peer, rejecting missing/nonlocal peers.
token_provider is called once in a managed worker during lifespan startup. It supplies a nonempty bearer secret; the provider and its result are never included in status, snapshots or error diagnostics. Credential rotation requires a new service lifespan. The application configures no root logger.
Ordinary requests have finite execution and waiting slots. Long polls have separate execution slots and no queue. request_timeout covers admission, body receipt, processing, serialization and sending; managed workers retain their slot until they really finish after cancellation. Counts bound owned work, not Python heap usage. Limits in ClientConfig separately bound retained snapshots, subscriptions, event history, batches and exports. The default request deadline is 120 seconds; the SDK allows 180 seconds because its deadline also includes response parsing and verification.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
host
|
str
|
Actual listener IP or DNS name configured by the host. Only a literal loopback IP qualifies for unauthenticated local access. |
'127.0.0.1'
|
port
|
int
|
Actual listener port, an integer from 1 through 65535. |
8323
|
workers
|
int
|
Source-owning worker count; must be exactly one. |
1
|
auth_mode
|
Literal['none', 'bearer']
|
Either local-only |
'none'
|
token_provider
|
Callable[[], str] | None
|
Synchronous callable supplying the bearer secret once during startup; required exactly when auth_mode is bearer. Cancellation waits for its managed worker to return rather than abandoning it. |
None
|
tls_mode
|
Literal['none', 'direct', 'proxy']
|
|
'none'
|
tls_cert_file
|
str | Path | None
|
Certificate path for direct TLS; the host also passes it to the ASGI server. Configuration construction does not read it. |
None
|
tls_key_file
|
str | Path | None
|
Private-key path for direct TLS; required with the certificate path and never included in source snapshots or diagnostics. |
None
|
trusted_proxy_ips
|
tuple[str, ...]
|
Explicit IP/CIDR allowlist for proxy TLS, checked against the original TCP peer with proxy-header rewriting disabled. |
()
|
max_body_bytes
|
int
|
Buffered request-body byte limit; excess bodies return 413. |
4 * 1024 * 1024
|
max_response_bytes
|
int
|
Complete ordinary-response byte limit; snapshot export instead uses ClientConfig.limits.max_export_bytes. |
16 * 1024 * 1024
|
max_concurrent_requests
|
int
|
Execution slots for ordinary domain requests; cancelled managed work occupies its slot until it actually finishes. |
4
|
max_pending_requests
|
int
|
FIFO waiting slots for ordinary requests; zero disables waiting. Saturation returns 429 before full body intake. |
16
|
request_timeout
|
float
|
Total seconds from request entry through admission, body receipt, computation, serialization and response sending. |
120.0
|
max_concurrent_polls
|
int
|
Independent long-poll slots, without a waiting queue. |
32
|
long_poll_timeout
|
float
|
Maximum event-wait seconds per poll, shorter than request_timeout. A request may select a shorter wait, including zero. |
25.0
|
watch_idle_timeout
|
float
|
Seconds without an active poll before reclaiming a watch; an in-progress poll is not reclaimed as idle. |
60.0
|
page_size
|
int
|
Default number of items per page when the caller omits a limit. |
1000
|
max_page_size
|
int
|
Maximum explicit page size; must be at least page_size. |
10000
|
host: str = '127.0.0.1'
class-attribute
instance-attribute
port: int = 8323
class-attribute
instance-attribute
workers: int = 1
class-attribute
instance-attribute
auth_mode: Literal['none', 'bearer'] = 'none'
class-attribute
instance-attribute
token_provider: Callable[[], str] | None = field(default=None, repr=False, compare=False)
class-attribute
instance-attribute
tls_mode: Literal['none', 'direct', 'proxy'] = 'none'
class-attribute
instance-attribute
tls_cert_file: str | Path | None = None
class-attribute
instance-attribute
tls_key_file: str | Path | None = None
class-attribute
instance-attribute
trusted_proxy_ips: tuple[str, ...] = ()
class-attribute
instance-attribute
max_body_bytes: int = 4 * 1024 * 1024
class-attribute
instance-attribute
max_response_bytes: int = 16 * 1024 * 1024
class-attribute
instance-attribute
max_concurrent_requests: int = 4
class-attribute
instance-attribute
max_pending_requests: int = 16
class-attribute
instance-attribute
request_timeout: float = 120.0
class-attribute
instance-attribute
max_concurrent_polls: int = 32
class-attribute
instance-attribute
long_poll_timeout: float = 25.0
class-attribute
instance-attribute
watch_idle_timeout: float = 60.0
class-attribute
instance-attribute
page_size: int = 1000
class-attribute
instance-attribute
max_page_size: int = 10000
class-attribute
instance-attribute
__post_init__() -> None
__init__(*, host: str = '127.0.0.1', port: int = 8323, workers: int = 1, auth_mode: Literal['none', 'bearer'] = 'none', token_provider: Callable[[], str] | None = None, tls_mode: Literal['none', 'direct', 'proxy'] = 'none', tls_cert_file: str | Path | None = None, tls_key_file: str | Path | None = None, trusted_proxy_ips: tuple[str, ...] = (), max_body_bytes: int = 4 * 1024 * 1024, max_response_bytes: int = 16 * 1024 * 1024, max_concurrent_requests: int = 4, max_pending_requests: int = 16, request_timeout: float = 120.0, max_concurrent_polls: int = 32, long_poll_timeout: float = 25.0, watch_idle_timeout: float = 60.0, page_size: int = 1000, max_page_size: int = 10000) -> None
Endpoint
dataclass
Connection coordinates without credentials or TLS context.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
transport
|
Literal['tcp', 'tls', 'ssh']
|
|
required |
host
|
str
|
DNS name or unbracketed IPv4/IPv6 address to connect to. |
required |
port
|
int | None
|
TCP port; defaults to 323 for TCP, 324 for TLS, or 22 for SSH. |
None
|
local_address
|
str | None
|
Optional literal source IP, never a DNS/interface name. |
None
|
server_name
|
str | None
|
TLS DNS reference identifier. Defaults to a DNS host; an IP host requires this argument for TLS. CN and IP-ID authentication are not used. International names are passed to AnyIO's IDNA 2008 implementation without IDNA 2003 conversion. |
None
|
Construction normalizes values and raises ConfigurationError for invalid combinations. It performs no DNS lookup, file read, or socket operation.
transport: Literal['tcp', 'tls', 'ssh']
instance-attribute
host: str
instance-attribute
port: int | None = None
class-attribute
instance-attribute
local_address: str | None = None
class-attribute
instance-attribute
server_name: str | None = None
class-attribute
instance-attribute
__post_init__() -> None
__init__(*, transport: Literal['tcp', 'tls', 'ssh'], host: str, port: int | None = None, local_address: str | None = None, server_name: str | None = None) -> None
RtrSourceConfig
dataclass
Immutable RTR source and transport configuration, without startup I/O.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
id
|
str
|
Stable source identifier within the client configuration. |
required |
host
|
str
|
Connection DNS name or unbracketed IP address. |
required |
port
|
int | None
|
Defaults to 323 (TCP), 324 (TLS), or 22 (SSH). |
None
|
transport
|
Literal['tcp', 'tls', 'ssh']
|
|
'tcp'
|
trust_profile_id
|
str | None
|
Nonsecret identity for the host's trust policy. Change this when credentials or trust conditions change. |
None
|
ssl_context
|
SSLContext | None
|
Host-owned context, never mutated by the library. The host must load a client certificate containing SAN iPAddress and enable CERT_REQUIRED, check_hostname and SAN-only verification. Mutually exclusive with all three certificate file arguments. |
None
|
local_address
|
str | None
|
Optional literal source IP to bind before connecting. |
None
|
min_version
|
int
|
Lowest acceptable RTR version, 1 or 2. |
1
|
max_version
|
int
|
Highest version to request. Experimental v2 is opt-in. |
1
|
v2_profile
|
Literal['8210bis-27', '8210bis-10', '8210bis-13']
|
Explicit experimental draft, 8210bis-27 by default or 8210bis-10 for historical AFI-specific ASPA or 8210bis-13 for counted, dual-family ASPA. Never inferred from ambiguous received bytes. Changes replace the source. |
'8210bis-27'
|
connect_timeout
|
float
|
Positive total seconds for DNS, TCP and TLS or SSH authentication and opening the rpki-rtr subsystem. |
10.0
|
query_timeout
|
float
|
Positive total seconds from query through End of Data. |
120.0
|
check_order
|
bool
|
Check the optional v2 payload ordering rule by default. |
True
|
server_name
|
str | None
|
Cache DNS identity; required when TLS host is an IP. |
None
|
ca_file
|
str | Path | None
|
Optional CA PEM used instead of default trust anchors. Omit to use Python/OpenSSL default trust, which may be affected by SSL_CERT_FILE and SSL_CERT_DIR. Built-in contexts also follow create_default_context's SSLKEYLOGFILE diagnostics. A host-owned ssl_context controls its own trust and diagnostics instead. |
None
|
client_cert_file
|
str | Path | None
|
Required for built-in TLS; PEM certificate chain. The host supplies the protocol-required SAN iPAddress identity. |
None
|
client_key_file
|
str | Path | None
|
Optional private key PEM; omit for a combined PEM. Encrypted keys require a host-prepared context; the built-in loader never prompts for a password. |
None
|
ssh
|
SshConfig | None
|
Explicit SSH credentials and host trust, only for SSH transport. Required at entry for the built-in adapter; injected factories own their authentication. Credentials are never included in Endpoint. |
None
|
File paths become absolute without reading them. Context-entry preparation performs certificate loading in a managed worker. Credentials and context are absent from endpoint and repr. ConfigurationError indicates invalid arguments; this type does not authenticate or connect to a cache.
id: str
instance-attribute
host: str
instance-attribute
port: int | None = None
class-attribute
instance-attribute
transport: Literal['tcp', 'tls', 'ssh'] = 'tcp'
class-attribute
instance-attribute
trust_profile_id: str | None = None
class-attribute
instance-attribute
ssl_context: ssl.SSLContext | None = field(default=None, repr=False)
class-attribute
instance-attribute
local_address: str | None = None
class-attribute
instance-attribute
min_version: int = 1
class-attribute
instance-attribute
max_version: int = 1
class-attribute
instance-attribute
v2_profile: Literal['8210bis-27', '8210bis-10', '8210bis-13'] = '8210bis-27'
class-attribute
instance-attribute
connect_timeout: float = 10.0
class-attribute
instance-attribute
query_timeout: float = 120.0
class-attribute
instance-attribute
check_order: bool = True
class-attribute
instance-attribute
server_name: str | None = None
class-attribute
instance-attribute
ca_file: str | Path | None = field(default=None, repr=False)
class-attribute
instance-attribute
client_cert_file: str | Path | None = field(default=None, repr=False)
class-attribute
instance-attribute
client_key_file: str | Path | None = field(default=None, repr=False)
class-attribute
instance-attribute
ssh: SshConfig | None = field(default=None, repr=False)
class-attribute
instance-attribute
endpoint: Endpoint
property
Normalized connection coordinates containing no credential fields.
__post_init__() -> None
__init__(*, id: str, host: str, port: int | None = None, transport: Literal['tcp', 'tls', 'ssh'] = 'tcp', trust_profile_id: str | None = None, ssl_context: ssl.SSLContext | None = None, local_address: str | None = None, min_version: int = 1, max_version: int = 1, v2_profile: Literal['8210bis-27', '8210bis-10', '8210bis-13'] = '8210bis-27', connect_timeout: float = 10.0, query_timeout: float = 120.0, check_order: bool = True, server_name: str | None = None, ca_file: str | Path | None = None, client_cert_file: str | Path | None = None, client_key_file: str | Path | None = None, ssh: SshConfig | None = None) -> None
SshConfig
dataclass
Explicit SSH trust and authentication for the rpki-rtr subsystem.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
username
|
str
|
Nonempty SSH user name, without control characters. |
required |
known_hosts
|
str | Path
|
Required OpenSSH known-hosts file. Unknown, changed and revoked host keys are rejected; there is no trust-on-first-use mode. |
required |
client_keys
|
tuple[str | Path, ...]
|
Private key files to try for public-key authentication. Omit only when password_provider is supplied. Default key files, SSH agents and user SSH configuration are never consulted. |
()
|
password_provider
|
Callable[[], str] | None
|
Optional synchronous callback returning a nonempty password. Called once per prepared source in a managed worker; never called during configuration construction or on every retry. |
None
|
passphrase_provider
|
Callable[[], str] | None
|
Optional synchronous callback returning a nonempty private-key passphrase, applied to the configured key files. Encrypted keys never prompt interactively. |
None
|
Requires the ssh extra and asyncio when the built-in transport enters
its context. Construction validates types and makes paths absolute without
file access. Providers and credential paths are excluded from repr; do not
serialize this object to expose provider internals. Changing credentials or
trust requires a new source configuration and trust_profile_id; modifying a
file or callback in place does not reload already prepared material.
ConfigurationError reports invalid arguments or startup preparation failure.
username: str
instance-attribute
known_hosts: str | Path = field(repr=False)
class-attribute
instance-attribute
client_keys: tuple[str | Path, ...] = field(default=(), repr=False)
class-attribute
instance-attribute
password_provider: Callable[[], str] | None = field(default=None, repr=False)
class-attribute
instance-attribute
passphrase_provider: Callable[[], str] | None = field(default=None, repr=False)
class-attribute
instance-attribute
__post_init__() -> None
__eq__(other: object) -> bool
Compare immutable fields and provider identity without calling providers.
__hash__() -> int
Hash the same non-invoking in-process identity used for equality.
__init__(*, username: str, known_hosts: str | Path, client_keys: tuple[str | Path, ...] = (), password_provider: Callable[[], str] | None = None, passphrase_provider: Callable[[], str] | None = None) -> None
Limits
dataclass
Finite resource admission limits, not a bound on Python heap usage.
Counts charge each source's support separately and ASPA provider edges individually. Byte limits are decoded bytes. Exceeding a limit fails the complete operation. Values are initial engineering budgets subject to measured release acceptance; see the configuration and performance guides. The service retains at most four locatable snapshots for up to 600 seconds; retention never extends their source validity and capacity can evict earlier.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
max_sources
|
int
|
Configured source count; also bounds separately retained retiring source workers during reconfiguration. |
16
|
max_concurrent_builds
|
int
|
Admitted builds, including admitted source updates waiting for their group's full-build lock. |
2
|
max_pending_builds
|
int
|
Waiting requests for shared build admission. Only one configuration candidate may wait or execute at a time. |
16
|
max_total_staging_wire_bytes
|
int
|
Combined charged input bytes of all admitted sources, including work still cleaning up after removal. |
512 * 1024 * 1024
|
max_json_bytes
|
int
|
Maximum decoded JSON input bytes, after decompression. |
256 * 1024 * 1024
|
max_json_depth
|
int
|
Maximum nesting depth, including host-provided mappings. |
128
|
max_records_per_source
|
int
|
Normalized support units per source; each ASPA provider edge counts separately, including the AS0 marker. |
5000000
|
max_total_records
|
int
|
Combined normalized support units of retained sources; identical payloads in different sources consume separate units. |
10000000
|
max_staging_wire_bytes
|
int
|
Input byte budget for one source transaction. |
512 * 1024 * 1024
|
max_pdu_bytes_v1
|
int
|
Local RFC 8210 PDU byte limit, including its header; this is an engineering guard, not a protocol-defined maximum. |
1024 * 1024
|
max_pdu_bytes_v2_legacy
|
int
|
Local PDU byte limit for draft -10/-13 profiles. Draft -27 instead has its fixed 65,535-byte protocol maximum. |
1024 * 1024
|
max_batch_items
|
int
|
Items in one fixed-snapshot batch; excess input is rejected. |
10000
|
max_path_asns
|
int
|
ASNs in an input path; excess paths are never truncated. |
16384
|
max_subscribers
|
int
|
Simultaneous subscriptions owned by one manager. |
128
|
subscription_queue_size
|
int
|
Pending events per subscriber; overflow requires explicit resynchronization rather than silently dropping an event. |
64
|
event_window_size
|
int
|
Maximum retained event count, also constrained by bytes. |
1024
|
max_event_window_bytes
|
int
|
Maximum canonical event bytes retained in history. |
16 * 1024 * 1024
|
max_event_changes
|
int
|
Detailed changes before an event requests full revalidation of the affected payload families. |
10000
|
max_explanation_records
|
int
|
Detailed evidence records before an explanation explicitly reports truncation and supplies further query criteria. |
100
|
max_source_diff_records
|
int
|
Output payload and support differences allowed in one SourceDiff; overflow fails the whole comparison. |
10000
|
max_source_diff_bytes
|
int
|
Canonical UTF-8 JSON bytes allowed for a SourceDiff. |
16 * 1024 * 1024
|
retained_snapshots
|
int
|
Server-locatable snapshot count; caller-held local snapshots remain the caller's memory responsibility. |
4
|
retained_snapshot_seconds
|
float
|
Maximum server retention in seconds, measured from the original pin; neither access nor pinning renews validity. |
600.0
|
max_export_bytes
|
int
|
Complete snapshot exchange byte budget for export/import; an over-budget operation cannot return truncated success. |
512 * 1024 * 1024
|
max_recording_bytes
|
int
|
Recording output byte budget; exhaustion marks the recording incomplete without stopping source synchronization. |
256 * 1024 * 1024
|
recording_queue_size
|
int
|
Pending recording items before recording reports a gap. |
128
|
max_recording_queue_bytes
|
int
|
Charged recording queue bytes, including base64 allowance and per-item overhead; not a Python heap bound. |
8 * 1024 * 1024
|
max_sources: int = 16
class-attribute
instance-attribute
max_concurrent_builds: int = 2
class-attribute
instance-attribute
max_pending_builds: int = 16
class-attribute
instance-attribute
max_total_staging_wire_bytes: int = 512 * 1024 * 1024
class-attribute
instance-attribute
max_json_bytes: int = 256 * 1024 * 1024
class-attribute
instance-attribute
max_json_depth: int = 128
class-attribute
instance-attribute
max_records_per_source: int = 5000000
class-attribute
instance-attribute
max_total_records: int = 10000000
class-attribute
instance-attribute
max_staging_wire_bytes: int = 512 * 1024 * 1024
class-attribute
instance-attribute
max_pdu_bytes_v1: int = 1024 * 1024
class-attribute
instance-attribute
max_pdu_bytes_v2_legacy: int = 1024 * 1024
class-attribute
instance-attribute
max_batch_items: int = 10000
class-attribute
instance-attribute
max_path_asns: int = 16384
class-attribute
instance-attribute
max_subscribers: int = 128
class-attribute
instance-attribute
subscription_queue_size: int = 64
class-attribute
instance-attribute
event_window_size: int = 1024
class-attribute
instance-attribute
max_event_window_bytes: int = 16 * 1024 * 1024
class-attribute
instance-attribute
max_event_changes: int = 10000
class-attribute
instance-attribute
max_explanation_records: int = 100
class-attribute
instance-attribute
max_source_diff_records: int = 10000
class-attribute
instance-attribute
max_source_diff_bytes: int = 16 * 1024 * 1024
class-attribute
instance-attribute
retained_snapshots: int = 4
class-attribute
instance-attribute
retained_snapshot_seconds: float = 600.0
class-attribute
instance-attribute
max_export_bytes: int = 512 * 1024 * 1024
class-attribute
instance-attribute
max_recording_bytes: int = 256 * 1024 * 1024
class-attribute
instance-attribute
recording_queue_size: int = 128
class-attribute
instance-attribute
max_recording_queue_bytes: int = 8 * 1024 * 1024
class-attribute
instance-attribute
__init__(*, max_sources: int = 16, max_concurrent_builds: int = 2, max_pending_builds: int = 16, max_total_staging_wire_bytes: int = 512 * 1024 * 1024, max_json_bytes: int = 256 * 1024 * 1024, max_json_depth: int = 128, max_records_per_source: int = 5000000, max_total_records: int = 10000000, max_staging_wire_bytes: int = 512 * 1024 * 1024, max_pdu_bytes_v1: int = 1024 * 1024, max_pdu_bytes_v2_legacy: int = 1024 * 1024, max_batch_items: int = 10000, max_path_asns: int = 16384, max_subscribers: int = 128, subscription_queue_size: int = 64, event_window_size: int = 1024, max_event_window_bytes: int = 16 * 1024 * 1024, max_event_changes: int = 10000, max_explanation_records: int = 100, max_source_diff_records: int = 10000, max_source_diff_bytes: int = 16 * 1024 * 1024, retained_snapshots: int = 4, retained_snapshot_seconds: float = 600.0, max_export_bytes: int = 512 * 1024 * 1024, max_recording_bytes: int = 256 * 1024 * 1024, recording_queue_size: int = 128, max_recording_queue_bytes: int = 8 * 1024 * 1024) -> None
__post_init__() -> None
FreshnessPolicy
dataclass
Explicit upper bounds on source freshness; never use download time.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
max_age
|
float | None
|
Positive seconds measured from the original generated_at. |
None
|
valid_until
|
datetime | None
|
Absolute timezone-aware deadline, normalized to UTC. At least one argument is required; when both exist, use the earlier. |
None
|
max_age: float | None = None
class-attribute
instance-attribute
valid_until: datetime | None = None
class-attribute
instance-attribute
__init__(*, max_age: float | None = None, valid_until: datetime | None = None) -> None
__post_init__() -> None
JsonFileSourceConfig
dataclass
A periodically read complete JSON file, with an explicit freshness bound.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
id
|
str
|
Unique nonsecret source identifier. |
required |
path
|
str | Path
|
File path fixed against the construction-time working directory. |
required |
freshness
|
FreshnessPolicy
|
Original-time bounds; re-reading the file never renews them. |
required |
format
|
str
|
Built-in format name or an explicitly injected reader format. |
'auto'
|
reader_profile_id
|
str | None
|
Nonsecret parsing identity required for custom readers. |
None
|
poll_interval
|
float
|
Positive seconds between operations. A hot change takes effect for the next operation, without resetting an in-flight timer. |
30.0
|
id: str
instance-attribute
path: str | Path
instance-attribute
freshness: FreshnessPolicy
instance-attribute
format: str = 'auto'
class-attribute
instance-attribute
reader_profile_id: str | None = None
class-attribute
instance-attribute
poll_interval: float = 30.0
class-attribute
instance-attribute
__init__(*, id: str, path: str | Path, freshness: FreshnessPolicy, format: str = 'auto', reader_profile_id: str | None = None, poll_interval: float = 30.0) -> None
__post_init__() -> None
HttpJsonSourceConfig
dataclass
Optional HTTP JSON input with explicit credentials and TLS trust identity.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
id
|
str
|
Unique nonsecret source identifier. |
required |
url
|
str
|
Absolute HTTP(S) URL without userinfo or fragments. Credentials go in headers; URLs and headers are never copied to source diagnostics. |
required |
freshness
|
FreshnessPolicy
|
Original data bounds, also retained for HTTP 304 responses. |
required |
format
|
str
|
Built-in or explicitly injected reader format. |
'auto'
|
reader_profile_id
|
str | None
|
Nonsecret identity for custom parsing semantics. |
None
|
poll_interval
|
float
|
Positive seconds between completed requests. |
3600.0
|
request_timeout
|
float
|
Total request deadline in seconds, default 1200, including admission, connection, body streaming, parsing and publication. Managed cleanup may outlast this deadline. |
1200.0
|
headers
|
Mapping[str, str]
|
Explicit request headers, copied and excluded from repr. May include authentication; do not mutate configuration after binding. |
dict()
|
ssl_context
|
SSLContext | None
|
Optional host-owned HTTPS verification context. Certificate verification and hostname checks must remain enabled. The default uses HTTPX's certifi trust bundle, ignoring certificate environment variables and additional operating-system trust stores. |
None
|
trust_profile_id
|
str
|
Nonsecret trust/credential identity for recovery. Change it when the authentication or trust conditions change. |
'system-trust'
|
HTTP connections require the http extra. Construction validates only
values; importing core configuration never imports an HTTP implementation.
id: str
instance-attribute
url: str = field(repr=False)
class-attribute
instance-attribute
freshness: FreshnessPolicy
instance-attribute
format: str = 'auto'
class-attribute
instance-attribute
reader_profile_id: str | None = None
class-attribute
instance-attribute
poll_interval: float = 3600.0
class-attribute
instance-attribute
request_timeout: float = 1200.0
class-attribute
instance-attribute
headers: Mapping[str, str] = field(default_factory=dict, repr=False)
class-attribute
instance-attribute
ssl_context: SSLContext | None = field(default=None, repr=False)
class-attribute
instance-attribute
trust_profile_id: str = 'system-trust'
class-attribute
instance-attribute
__init__(*, id: str, url: str, freshness: FreshnessPolicy, format: str = 'auto', reader_profile_id: str | None = None, poll_interval: float = 3600.0, request_timeout: float = 1200.0, headers: Mapping[str, str] = dict(), ssl_context: SSLContext | None = None, trust_profile_id: str = 'system-trust') -> None
__post_init__() -> None
SourceGroupConfig
dataclass
Sources merged within one priority group, selected separately by payload.
Lower priority numbers are preferred. Priorities and IDs must be unique in ClientConfig. A recovered higher-priority group waits failback_delay seconds before replacing a usable backup; zero requests immediate recovery.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
id
|
str
|
Nonempty group identifier, unique within the Client configuration. |
required |
priority
|
int
|
Unique integer preference; smaller values are preferred. |
required |
sources
|
tuple[SourceConfig, ...]
|
Complete source configurations in this trust group. The tuple may be empty; each source ID must be unique across all groups. |
()
|
failback_delay
|
float
|
Seconds of healthy recovery before replacing a usable backup; zero permits immediate failback. ASPA selection is per AFI. |
30.0
|
id: str
instance-attribute
priority: int
instance-attribute
sources: tuple[SourceConfig, ...] = ()
class-attribute
instance-attribute
failback_delay: float = 30.0
class-attribute
instance-attribute
__init__(*, id: str, priority: int, sources: tuple[SourceConfig, ...] = (), failback_delay: float = 30.0) -> None
__post_init__() -> None
PersistenceConfig
dataclass
Explicit single-owner database configuration; opening is a lifecycle step.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
backend
|
Literal['sqlite', 'duckdb']
|
|
required |
path
|
str | Path
|
Database path fixed at construction; only owner opening creates it. |
required |
mode
|
Literal['owner']
|
Only |
'owner'
|
flush_timeout
|
float
|
Positive wait budget in seconds, default 600. Final close captures the same budget when shutdown begins. Neither waiting deadline permits abandoning an active database transaction or its worker thread; this is not a hard bound on total close time. |
_DEFAULT_FLUSH_TIMEOUT
|
backend: Literal['sqlite', 'duckdb']
instance-attribute
path: str | Path
instance-attribute
mode: Literal['owner'] = 'owner'
class-attribute
instance-attribute
flush_timeout: float = _DEFAULT_FLUSH_TIMEOUT
class-attribute
instance-attribute
__init__(*, backend: Literal['sqlite', 'duckdb'], path: str | Path, mode: Literal['owner'] = 'owner', flush_timeout: float = _DEFAULT_FLUSH_TIMEOUT) -> None
__post_init__() -> None
ClientConfig
dataclass
Complete immutable configuration for one owner Client.
Empty groups are valid and can acquire sources through apply_config. The required set controls default readiness and health; each explicit readiness call has its own requirements and deadline. startup_timeout opens a backup preparation window without cancelling a still-running primary. Cleanup timeout reports overdue cleanup and never abandons owned work.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
groups
|
tuple[SourceGroupConfig, ...]
|
Complete ordered-group configuration; priorities determine selection independently of tuple order. Empty configurations are valid. |
()
|
limits
|
Limits
|
Finite admission, input, query, event and retained-view budgets. |
Limits()
|
persistence
|
PersistenceConfig | None
|
Optional fixed owner database configuration; None keeps state in memory unless a backend is explicitly injected into Client. |
None
|
required
|
frozenset[PayloadKind]
|
Nonempty set of VRP and/or ASPA used for default readiness and health checks, independently of individual query requirements. |
frozenset({VRP})
|
startup_timeout
|
float
|
Seconds in each group's initial preparation window before preparing a backup for an unready capability. Changing this value does not move deadlines of already activated groups. |
30.0
|
cleanup_timeout
|
float
|
Seconds before reporting cleanup_overdue; owned work remains managed until it actually finishes. |
5.0
|
clock_tolerance
|
float
|
Allowed wall/monotonic drift and future generation tolerance in seconds; changing it requires a new Client. |
5.0
|
groups: tuple[SourceGroupConfig, ...] = ()
class-attribute
instance-attribute
limits: Limits = field(default_factory=Limits)
class-attribute
instance-attribute
persistence: PersistenceConfig | None = None
class-attribute
instance-attribute
required: frozenset[PayloadKind] = frozenset({PayloadKind.VRP})
class-attribute
instance-attribute
startup_timeout: float = 30.0
class-attribute
instance-attribute
cleanup_timeout: float = 5.0
class-attribute
instance-attribute
clock_tolerance: float = 5.0
class-attribute
instance-attribute
__init__(*, groups: tuple[SourceGroupConfig, ...] = (), limits: Limits = Limits(), persistence: PersistenceConfig | None = None, required: frozenset[PayloadKind] = frozenset({PayloadKind.VRP}), startup_timeout: float = 30.0, cleanup_timeout: float = 5.0, clock_tolerance: float = 5.0) -> None
__post_init__() -> None
SharedReaderConfig
dataclass
Read-only SQLite observer configuration, with no source connections.
The path must already exist when opened. Readers independently enforce original deadlines even after the writer exits. Neither the backend nor the clock tolerance can be hot-switched on an existing observer.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
path
|
str | Path
|
Existing SQLite database file, resolved against the construction working directory; this observer does not create or write it. |
required |
poll_interval
|
float
|
Seconds between checks for a new committed database head. |
1.0
|
limits
|
Limits
|
Local input, rebuild, query and event budgets of this observer. |
Limits()
|
clock_tolerance
|
float
|
Wall/monotonic drift and future-generation tolerance in seconds, fixed for this observer's lifetime. |
5.0
|
path: str | Path
instance-attribute
poll_interval: float = 1.0
class-attribute
instance-attribute
limits: Limits = field(default_factory=Limits)
class-attribute
instance-attribute
clock_tolerance: float = 5.0
class-attribute
instance-attribute
__init__(*, path: str | Path, poll_interval: float = 1.0, limits: Limits = Limits(), clock_tolerance: float = 5.0) -> None
__post_init__() -> None
JSON readers and snapshot loading
rpkiparrot.readers
Synchronous, bounded readers for supported producer JSON documents.
Parsing preserves original time bounds and completeness diagnostics. It does not
grant trust or freshness: submit a complete dataset to MemoryStore with an
explicit freshness policy before querying it. No network or file I/O occurs in
parse_json or a format adapter.
JsonFormatAdapter
Bases: Protocol
Explicit producer adapter, registered on one :class:JsonReader.
format_id must be a unique, nonempty string other than auto or a
built-in ID. parse is synchronous and performs no I/O. Return normalized
records with original time bounds, explicit capabilities and completeness;
the reader applies common size and dataset checks to the returned value.
Adapter exceptions are propagated and no source is partially committed.
format_id: str
property
Stable format identifier chosen explicitly by the caller.
parse(document: Mapping[str, object]) -> ParsedDataset
Parse an already checked JSON object without performing I/O.
JsonReader
Instance-local built-in and explicitly injected producer adapters.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
adapters
|
Sequence[JsonFormatAdapter]
|
Additional adapters. IDs cannot shadow built-ins. The input sequence is copied and calls on this reader are serialized, so a stateful adapter can safely be shared by multiple source tasks. |
()
|
Raises:
| Type | Description |
|---|---|
ConfigurationError
|
An adapter has an invalid or duplicate ID. |
supported_formats: tuple[str, ...]
property
Explicit supported IDs, including built-ins; auto uses only built-ins.
__init__(adapters: Sequence[JsonFormatAdapter] = ()) -> None
parse_json(data: bytes | str | dict[str, object], *, format: str = 'auto', limits: Limits | None = None) -> ParsedDataset
Parse one complete producer document without reading files or clocks.
str always means JSON text. Byte input must be UTF-8. Dict input is
checked and byte-budgeted but cannot recover duplicate keys; a diagnostic
records that limitation. Missing payload arrays remain unsupported, and
explicit empty arrays preserve the producer's capability. Auto detection
requires an explicit format for Routinator documents whose payloads are
all empty or absent: JSON and JSONExt share the same metadata. Incomplete
producer counts return complete=False for field research.
Raises:
| Type | Description |
|---|---|
InputError
|
Syntax, field types, producer semantics or format fail. |
ResourceLimitError
|
Bytes, nesting, supports or staged payload exceed their configured budgets. No partial dataset is returned. |
read_json(path: str | PathLike[str], *, format: str = 'auto', limits: Limits | None = None) -> ParsedDataset
Read a bounded local file and parse it; this method blocks its thread.
A changed file identity, size or timestamps rejects the candidate with
InputError. Online callers can retry without replacing their previous
dataset. The file is always closed, including on failure. File-system
OSError exceptions retain their usual Python meaning.
parse_json(data: bytes | str | dict[str, object], *, format: str = 'auto', limits: Limits | None = None) -> ParsedDataset
Parse producer JSON using built-ins; see :meth:JsonReader.parse_json.
Use an explicit JsonReader instance for custom adapters. This function
never changes a global registry and treats strings solely as JSON text.
read_json(path: str | PathLike[str], *, format: str = 'auto', limits: Limits | None = None) -> ParsedDataset
Read a local file using built-ins; see :meth:JsonReader.read_json.
load_snapshot(data: bytes | str, *, context: EvaluationContext | None = None, limits: Limits | None = None) -> Snapshot
Load a complete export selected by the host, preserving its original times.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
data
|
bytes | str
|
UTF-8 bytes or JSON text from |
required |
context
|
EvaluationContext | None
|
Defaults to online. Offline exports require an explicit offline context, whose reference cannot precede the original generation. |
None
|
limits
|
Limits | None
|
The export byte budget is independent from producer JSON limits. |
None
|
Returns:
| Type | Description |
|---|---|
Snapshot
|
A snapshot with a new local identity and |
Snapshot
|
exported identity. It may contain expired or unavailable capabilities. |
Snapshot
|
After a view-change boundary, reload the original bytes to re-evaluate |
Snapshot
|
remaining support; parsing and indexing do not extend original deadlines. |
Snapshot
|
Online monotonic time advances from the greater of the current clock |
Snapshot
|
floor and each source's recorded evaluation time, even after restart. |
Snapshot
|
Group definitions and unloaded sources are preserved. Each payload |
Snapshot
|
retains its previously selected group while that group remains usable; |
Snapshot
|
otherwise it selects the highest-priority usable group independently. |
Raises:
| Type | Description |
|---|---|
InputError
|
Schema, digest, counts, capabilities, times or mode are invalid. |
ResourceLimitError
|
The full export or normalized data exceeds a budget. |
Validation
rpkiparrot.validation
Pure ROV, ASPA, and conservative BMP evaluation on one immutable snapshot.
ROV follows RFC 6811 section 2 and RFC 8481. ASPA follows the explicitly experimental draft-ietf-sidrops-aspa-verification-28 sections 5.1--5.6. These functions neither choose routing policy nor collect global metrics.
validate_origin(snapshot: Snapshot, prefix: Network | str, asn: int | None, *, explain: bool = False) -> OriginResult
Evaluate one origin against every covering VRP (RFC 6811 section 2).
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
snapshot
|
Snapshot
|
Fixed source and policy view. Its VRP lifetime is checked at the start of this evaluation; holding it does not extend freshness. |
required |
prefix
|
Network | str
|
A strict CIDR or network object; host bits and interface zones are rejected. Mapped IPv6 is evaluated as IPv6, never as IPv4. |
required |
asn
|
int | None
|
Unsigned 32-bit origin, or None for RFC 6811's NONE. AS0 and NONE cannot match a VRP, but still produce notfound when none covers. |
required |
explain
|
bool
|
Include bounded covering records and matching predicates. |
False
|
Returns:
| Type | Description |
|---|---|
OriginResult
|
A valid, invalid, or notfound result with evaluation time and identity. |
OriginResult
|
Detailed and brief modes have identical status and reason codes. |
Raises:
| Type | Description |
|---|---|
InputError
|
Input types, ASN, or network are malformed. |
DataUnavailableError
|
VRP data is unavailable, including an expired fixed snapshot. An empty ready dataset can produce notfound. |
validate_aspa(snapshot: Snapshot, path: AsPath, *, context: AspaContext, afi: Afi | None = None, explain: bool = False) -> AspaResult
Apply ASPA verification draft -28 to an external neighbor-to-origin path.
afi selects the authorization family. None requires identical complete family views and never unions different authorizations. The caller reconstructs AS_PATH/AS4_PATH and supplies the local ASN, neighbor ASN, and relationship from the receiver's perspective. Only consecutive prepends are compressed. The local ASN is never inserted. Customer, peer, and route-server roles use upstream verification; provider uses downstream verification. Only ROUTE_SERVER omits the first-AS check.
Malformed paths raise InputError; confederation segments use its unsupported_path_context code. Empty paths, AS_SET, and a mismatched neighbor are ordinary invalid results once ASPA data is usable. No data raises DataUnavailableError even for paths requiring no authorization. Evidence is optional and bounded; its query field locates full records in the same snapshot. The result names the fixed experimental draft.
analyze_bmp_path(snapshot: Snapshot, path: AsPath, *, context: BmpContext, afi: Afi | None = None, explain: bool = False) -> BmpAnalysis
Conservatively analyze BMP observations without claiming standard valid.
afi selects the authorization family; None requires identical family views. Only pre_policy or post_policy inbound observations with explicitly reconstructed, neighbor_to_origin external paths are modeled. Missing direction, reconstruction, stage, or ASNs yields indeterminate. With an unknown relationship, all five supported relationships are considered; none is inferred from ASN values. All possible scenarios must be invalid to yield invalid, and all must be valid to yield no_invalid_evidence. A disagreement, unknown result, or unmodeled context is indeterminate.
Unavailable ASPA data raises DataUnavailableError even when context is incomplete. Malformed inputs raise InputError; unsupported confederation context is retained as indeterminate analysis. Returned scenarios share a snapshot and evaluation time. This function does not parse BMP packets.
validate_origins(snapshot: Snapshot, routes: Sequence[OriginInput], *, explain: bool = False) -> BatchResult[OriginResult]
Evaluate a bounded sequence in order, fixing data but not online time.
Each item rechecks VRP freshness. Input and data errors are saved per item, including an initially unavailable dataset or expiration during a batch. A malformed container or batch-size excess fails the whole call. Internal failures and cancellation propagate. complete means every item was processed, not that its route was valid. An empty sequence is complete.
validate_aspa_batch(snapshot: Snapshot, routes: Sequence[AspaInput], *, explain: bool = False) -> BatchResult[AspaResult]
Evaluate ordered ASPA inputs with the same batch semantics as ROV.
Inputs contain an AsPath and complete AspaContext. Each item samples time again and preserves input_id; resource, input, and availability errors are individual BatchItem errors. Oversized or malformed batch containers fail before evaluating any input. Cancellation is never an item error.
Queries and snapshot exports
rpkiparrot.query
Deterministic queries over one fixed immutable snapshot.
covering_vrps(snapshot: Snapshot, prefix: Network | str) -> tuple[VrpMatch, ...]
Return every same-family covering VRP, in public deterministic order.
Raises InputError for a non-network CIDR, DataUnavailableError for unavailable VRP data and SnapshotExpiredError at this view's next change boundary.
iter_vrps(snapshot: Snapshot, *, asn: int | None = None, source_id: str | None = None, family: int | None = None) -> Iterator[VrpMatch]
Traverse effective VRPs, checking the fixed view before each next item.
Source filters select records with that effective source support, preserving all effective supports. Dropping the iterator releases its reference; there are no hidden file handles or threads. Expiry never silently switches views. Like other Python generators, argument and initial availability checks run on the first next(), and freshness is checked again on subsequent next().
aspa_providers(snapshot: Snapshot, customer: int, *, afi: Afi | None = None) -> AspaProviders
Return the complete effective provider union and configured-source supports.
afi=None requires identical IPv4/IPv6 views; otherwise pass an explicit Afi. A missing customer has present=False only when ASPA data is usable. AS0 is retained only when no nonzero effective provider exists for that customer.
iter_aspas(snapshot: Snapshot, *, source_id: str | None = None, afi: Afi | None = None) -> Iterator[AspaProviders]
Yield each effective customer in the selected family in ascending ASN order.
None is accepted only when both family views have identical authorization, provenance, availability and lifetime. Otherwise specify Afi.IPV4 or IPV6. A source filter preserves the complete provider union in the selected family. Each next checks the fixed snapshot; errors start on the first next().
compare_sources(snapshot: Snapshot, left: str, right: str) -> SourceDiff
Compare two original sources, preserving expired and suppressed supports.
Each returned payload and support consumes max_source_diff_records; the complete canonical JSON result consumes max_source_diff_bytes. Exactly the limit is allowed. Exceeding either raises ResourceLimitError instead of returning a partially complete difference. Unknown source IDs are InputError. This diagnostic query does not grant data usability or choose trust policy.
export_snapshot(snapshot: Snapshot) -> bytes
Export complete original state with its original intervals and SHA-256.
The digest proves file integrity, not source authentication. Expired original records remain diagnostic data. ResourceLimitError prevents partial success above max_export_bytes, including the envelope; reload using readers.load_snapshot, not parse_json. Canonical record ordering preserves original supports, including equal-key ties; see the snapshot exchange format for the ordering and encoding rules. Mapped IPv6 uses fixed mixed notation across supported Python versions. Bounded record chunks preserve the complete canonical bytes and byte budget. Safe group definitions and unloaded sources are included; an empty record collection never substitutes for the source's explicit availability.
RTR codecs
rpkiparrot.rtr.codec
Pure RFC 8210 and explicit experimental 8210bis-27/-10/-13 wire codecs.
This module validates individual PDUs. It never establishes a session, performs I/O, applies a payload, or validates the direction/transaction of a message. Malformed constructor inputs raise InputError. Malformed received bytes raise ProtocolError, whose details include wire_error_code, received_type and reply_allowed. A false reply_allowed must never be turned into another Error Report. Local byte-budget failures instead raise ResourceLimitError.
Pdu: TypeAlias = SerialNotify | SerialQuery | ResetQuery | CacheResponse | Ipv4Prefix | Ipv6Prefix | EndOfData | CacheReset | RouterKey | ErrorReport | AspaPdu
module-attribute
SerialNotify
dataclass
Bases: _Serial
Cache update hint with explicit version, 16-bit session_id and 32-bit serial.
Session negotiation decides whether a received hint is usable. This codec accepts only versions 1 and 2, including for Serial Notify.
pdu_type: int = 0
class-attribute
Fixed wire PDU type number for this concrete message.
version: int
instance-attribute
RTR protocol version, restricted to 1 or 2; ASPA messages require version 2.
length: int
property
Total encoded byte length, including the eight-byte header.
session_id: int
instance-attribute
Unsigned 16-bit cache session identifier.
serial: int
instance-attribute
Unsigned 32-bit serial, ordered only within the cache session using serial arithmetic.
__init__(*, version: int, session_id: int, serial: int) -> None
__post_init__() -> None
SerialQuery
dataclass
Bases: _Serial
Incremental request identified by version, 16-bit session_id and 32-bit serial.
pdu_type: int = 1
class-attribute
Fixed wire PDU type number for this concrete message.
version: int
instance-attribute
RTR protocol version, restricted to 1 or 2; ASPA messages require version 2.
length: int
property
Total encoded byte length, including the eight-byte header.
session_id: int
instance-attribute
Unsigned 16-bit cache session identifier.
serial: int
instance-attribute
Unsigned 32-bit serial, ordered only within the cache session using serial arithmetic.
__init__(*, version: int, session_id: int, serial: int) -> None
__post_init__() -> None
ResetQuery
dataclass
Bases: _Versioned
Full-state request with an explicit protocol version; reserved bytes encode zero.
pdu_type: int = 2
class-attribute
Fixed wire PDU type number for this concrete message.
version: int
instance-attribute
RTR protocol version, restricted to 1 or 2; ASPA messages require version 2.
length: int
property
Total encoded byte length, including the eight-byte header.
__init__(*, version: int) -> None
__post_init__() -> None
CacheResponse
dataclass
Bases: _Session
Transaction start with explicit version and unsigned 16-bit session_id.
pdu_type: int = 3
class-attribute
Fixed wire PDU type number for this concrete message.
version: int
instance-attribute
RTR protocol version, restricted to 1 or 2; ASPA messages require version 2.
length: int
property
Total encoded byte length, including the eight-byte header.
session_id: int
instance-attribute
Unsigned 16-bit cache session identifier.
__init__(*, version: int, session_id: int) -> None
__post_init__() -> None
Ipv4Prefix
dataclass
Bases: _Versioned
IPv4 announcement or withdrawal, before any session applies it.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
version
|
int
|
Protocol 1 or 2. |
required |
announce
|
bool
|
True for announcement, False for exact-record withdrawal. |
required |
prefix
|
IPv4Network
|
Canonical IPv4 network. Nonzero host bits are rejected for both versions; this is a library acceptance policy for version 1. |
required |
max_length
|
int
|
Integer between prefix.prefixlen and 32 inclusive. |
required |
asn
|
int
|
Unsigned 32-bit origin ASN, including AS0; bool is rejected. |
required |
announce: bool
instance-attribute
prefix: IPv4Network
instance-attribute
max_length: int
instance-attribute
asn: int
instance-attribute
pdu_type: int = 4
class-attribute
Fixed wire PDU type number for this concrete message.
version: int
instance-attribute
RTR protocol version, restricted to 1 or 2; ASPA messages require version 2.
length: int
property
Total encoded byte length, including the eight-byte header.
__init__(*, version: int, announce: bool, prefix: IPv4Network, max_length: int, asn: int) -> None
__post_init__() -> None
Ipv6Prefix
dataclass
Bases: _Versioned
IPv6 announcement or withdrawal; the Ipv4Prefix rules apply with width 128.
announce: bool
instance-attribute
True announces a prefix; False withdraws the exact prefix, maximum length and ASN.
prefix: IPv6Network
instance-attribute
Canonical IPv6 network; host bits and zone identifiers are rejected.
max_length: int
instance-attribute
Maximum authorized length, between prefix.prefixlen and 128 inclusive.
asn: int
instance-attribute
Unsigned 32-bit origin ASN, including AS0; bool is rejected.
pdu_type: int = 6
class-attribute
Fixed wire PDU type number for this concrete message.
version: int
instance-attribute
RTR protocol version, restricted to 1 or 2; ASPA messages require version 2.
length: int
property
Total encoded byte length, including the eight-byte header.
__init__(*, version: int, announce: bool, prefix: IPv6Network, max_length: int, asn: int) -> None
__post_init__() -> None
EndOfData
dataclass
Bases: _Serial
Transaction end with three unsigned 32-bit intervals in seconds.
version, session_id and serial identify this transaction. refresh_interval, retry_interval and expire_interval here accept their entire wire ranges; RFC section 6 ranges and the requirement that expire exceed the other two intervals are checked by the session before it commits a transaction.
refresh_interval: int
instance-attribute
Wire refresh interval in seconds; session-level timing constraints apply at commit.
retry_interval: int
instance-attribute
Wire retry interval in seconds; session-level timing constraints apply at commit.
expire_interval: int
instance-attribute
Wire expiry interval in seconds; a successful EOD establishes the source lifetime.
pdu_type: int = 7
class-attribute
Fixed wire PDU type number for this concrete message.
version: int
instance-attribute
RTR protocol version, restricted to 1 or 2; ASPA messages require version 2.
length: int
property
Total encoded byte length, including the eight-byte header.
session_id: int
instance-attribute
Unsigned 16-bit cache session identifier.
serial: int
instance-attribute
Unsigned 32-bit serial, ordered only within the cache session using serial arithmetic.
__init__(*, version: int, session_id: int, serial: int, refresh_interval: int, retry_interval: int, expire_interval: int) -> None
__post_init__() -> None
CacheReset
dataclass
Bases: _Versioned
Cache response requiring full synchronization, with explicit protocol version.
pdu_type: int = 8
class-attribute
Fixed wire PDU type number for this concrete message.
version: int
instance-attribute
RTR protocol version, restricted to 1 or 2; ASPA messages require version 2.
length: int
property
Total encoded byte length, including the eight-byte header.
__init__(*, version: int) -> None
__post_init__() -> None
RouterKey
dataclass
Bases: _Versioned
Wire key announcement/withdrawal, distinct from any public key query API.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
version
|
int
|
Protocol 1 or 2. |
required |
announce
|
bool
|
True for announcement, False for withdrawal. |
required |
asn
|
int
|
Unsigned 32-bit ASN. |
required |
ski
|
bytes
|
Exactly 20 immutable bytes. |
required |
spki
|
bytes
|
Nonempty, immutable SPKI bytes within the uint32 length bound. Byte acceptance does not prove valid ASN.1, signature or trust. |
required |
Construction validates a profile-neutral candidate. encode_pdu additionally checks the selected profile's length bound, including the default -27 cap.
announce: bool
instance-attribute
asn: int
instance-attribute
ski: bytes
instance-attribute
spki: bytes
instance-attribute
pdu_type: int = 9
class-attribute
Fixed wire PDU type number for this concrete message.
length: int
property
Header, SKI, ASN and full SPKI size in bytes.
version: int
instance-attribute
RTR protocol version, restricted to 1 or 2; ASPA messages require version 2.
__init__(*, version: int, announce: bool, asn: int, ski: bytes, spki: bytes) -> None
__post_init__() -> None
ErrorReport
dataclass
Bases: _Versioned
A version-specific error with opaque offending bytes and UTF-8 text.
error_code is 0-8 for version 1 and 0-13 for version 2. The default encapsulated_pdu and text are empty. Construction validates a profile-neutral candidate; encode_pdu additionally checks -10/-13's error-code range 0-8 or -27's length bound and minimum four bytes for nonempty encapsulation. Encapsulation may be truncated and is never recursively decoded. The session must never send this PDU in response to an ErrorReport.
error_code: int
instance-attribute
Version-specific wire error number; profile-specific bounds also apply during encoding.
encapsulated_pdu: bytes = b''
class-attribute
instance-attribute
Opaque triggering PDU bytes, possibly truncated; never recursively decoded.
text: str = ''
class-attribute
instance-attribute
UTF-8 diagnostic text; an empty string is permitted.
pdu_type: int = 10
class-attribute
Fixed wire PDU type number for this concrete message.
length: int
property
Total byte size, including both length fields and UTF-8 encoded text.
version: int
instance-attribute
RTR protocol version, restricted to 1 or 2; ASPA messages require version 2.
__init__(*, version: int, error_code: int, encapsulated_pdu: bytes = b'', text: str = '') -> None
__post_init__() -> None
AspaPdu
dataclass
Bases: _Versioned
Experimental version-2 ASPA announcement, replacement or withdrawal.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
version
|
int
|
Must explicitly be 2. |
required |
announce
|
bool
|
True for announcement/replacement; False for customer withdrawal. |
required |
customer
|
int
|
Positive unsigned 32-bit ASN. |
required |
providers
|
tuple[int, ...]
|
Unique ASN tuple. Announcements require at least one provider, exclude customer, and allow AS0 only alone. Withdrawals require an empty tuple. An input list is copied, never sorted or deduplicated. -27 requires strictly increasing order and at most 16,380 providers; -10/-13 permit any order and at most 65,535. |
required |
afi
|
Afi | None
|
None for -27/-13's family-independent payload, or Afi.IPV4/Afi.IPV6 for -10. encode_pdu requires an explicitly matching v2_profile. |
None
|
v2_profile
|
str | None
|
Explicit wire representation. None preserves convenient construction by choosing -27 for afi=None and -10 for concrete afi. Runtime value is the chosen profile string. -13 must be explicit because it has a 16-byte header and no AFI. No wire guessing occurs. |
None
|
announce: bool
instance-attribute
customer: int
instance-attribute
providers: tuple[int, ...]
instance-attribute
afi: Afi | None = None
class-attribute
instance-attribute
v2_profile: str | None = None
class-attribute
instance-attribute
pdu_type: int = 11
class-attribute
Fixed wire PDU type number for this concrete message.
length: int
property
Header, customer and all provider ASN bytes.
version: int
instance-attribute
RTR protocol version, restricted to 1 or 2; ASPA messages require version 2.
__init__(*, version: int, announce: bool, customer: int, providers: tuple[int, ...], afi: Afi | None = None, v2_profile: str | None = None) -> None
__post_init__() -> None
PduDecoder
Incremental synchronous decoder for one transport's ordered byte stream.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
limits
|
Limits | None
|
Local v1 and v2 -10 frame bounds. -27 always caps at 65,535 bytes. |
None
|
v2_profile
|
str
|
"8210bis-27" (default), "8210bis-10" or "8210bis-13"; fixed here, invalid values raise InputError before any bytes are accepted. |
'8210bis-27'
|
Only one incomplete frame is buffered. A header is checked before its body is copied, even if feed receives a large concatenated input. Return values may contain multiple complete PDUs. No frame from a failed feed is returned. Any failure makes the decoder unusable; create a new decoder for a new stream. This object is owned by its calling task/thread and is not thread-safe.
__init__(*, limits: Limits | None = None, v2_profile: str = '8210bis-27') -> None
feed(data: bytes) -> tuple[Pdu, ...]
Consume bytes and return all complete frames, or fail the entire call.
Raises InputError for non-bytes, ProtocolError for bad frames, and ResourceLimitError for local bounds. Subsequent calls after failure or finish raise ClosedError. An empty input is allowed while open.
finish() -> None
End the stream, failing on a residual frame; successful repeats are safe.
The decoder closes even when this raises ProtocolError. A decoder that failed earlier instead raises ClosedError, preserving that failure's original no-reply and wire-error information at the original call.
encode_pdu(pdu: Pdu, *, limits: Limits | None = None, v2_profile: str = '8210bis-27') -> bytes
Encode one validated PDU; zero all reserved fields and unused flag bits.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
pdu
|
Pdu
|
One of the public frozen PDU dataclasses. InputError reports invalid constructor fields or a value that is not a supported PDU type. |
required |
limits
|
Limits | None
|
Optional local bounds for v1 and historical v2 profiles. Exceeding them raises ResourceLimitError. The -27 protocol cap is always 65,535 bytes. |
None
|
v2_profile
|
str
|
Explicit "8210bis-27" (default), "8210bis-10" or "8210bis-13". Invalid profile or a candidate incompatible with it raises InputError. |
'8210bis-27'
|
Returns:
| Type | Description |
|---|---|
bytes
|
Exactly one complete frame in network byte order. No I/O occurs. |
decode_pdu(data: bytes, *, limits: Limits | None = None, v2_profile: str = '8210bis-27') -> Pdu
Decode exactly one frame, enforcing structure and local admission limits.
Extra trailing bytes and truncated input fail, even if a valid PDU precedes them. Use PduDecoder for a fragmented or concatenated stream. This function ignores reserved receive fields, but never normalizes malformed prefixes or -27 ASPA provider ordering. -10/-13 preserve unordered unique announcements and discards a withdrawal's ignored count/list. Error encapsulation stays opaque. v2_profile explicitly selects "8210bis-27" (default), "8210bis-10" or "8210bis-13"; no automatic layout detection or family-scope conversion occurs.
Raises:
| Type | Description |
|---|---|
InputError
|
data is not bytes or v2_profile is unsupported. |
ProtocolError
|
Malformed frame; details specify wire_error_code and reply_allowed. Error Report input always prohibits an error response. |
ResourceLimitError
|
A v1 or historical v2 frame exceeds its configured byte bound. |
Transports
rpkiparrot.transports
AnyIO TCP/TLS/SSH transports and explicit ownership of injected streams.
Applications own the event loop and task tree. Factories create no detached tasks and preserve cancellation. A successful connect transfers one stream lease to its caller; the factory cleans up attempts that fail before handoff. Builtin TCP reads are pull-based, without an unbounded background receive queue. DNS, address attempts and TLS handshakes remain inside the connection deadline.
TransportFactory
Bases: Protocol
Create a fresh AnyIO byte stream for each connection attempt.
Implementations must clean up partial resources when connect fails or is
cancelled. After return, the caller closes the stream. No inheritance or
reconnectable attribute is required; ordinary factories can reconnect.
Custom transports are responsible for their own authentication and must
not log credentials. Endpoint contains no secret material.
The host retains ordinary factory lifecycle ownership: a method named
prepare or aclose alone never opts a custom factory into session management.
connect(endpoint: Endpoint) -> ByteStream
async
Connect to endpoint and transfer ownership on successful return.
May raise TransportError or a transport-specific I/O exception. Cancellation propagates after cleaning untransferred resources.
ConnectedTransport
Bases: _ManagedTransportFactory
Use one existing stream, with explicit close ownership and no reconnect.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
stream
|
ByteStream
|
Already connected AnyIO ByteStream supplied by the host. |
required |
owns_stream
|
bool
|
If true, this wrapper owns closing the original stream, including when closed before connect. If false, returned leases close independently and the host retains the original stream. |
required |
connect may be called once. It returns a wrapper whose aclose respects owns_stream; cancellation before handoff leaves this object responsible for its owned stream. The host/session should call aclose when finished. Further connection attempts raise ClosedError. Resource close is idempotent and protected from outer cancellation. A custom stream's close must eventually finish; this object does not abandon its cleanup.
reconnectable = False
class-attribute
instance-attribute
False: this binding can hand off its existing stream only once.
owns_stream: bool
property
Whether closing this transport also closes the original stream.
__init__(stream: ByteStream, owns_stream: bool) -> None
prepare() -> None
async
Check that the one-shot lease is available, without handing it out.
connect(endpoint: Endpoint) -> ByteStream
async
Hand out the stream once; endpoint does not open another connection.
Raises ClosedError after a prior handoff or aclose; invalid endpoint types raise ConfigurationError. Cancellation before handoff does not consume the one use, and aclose still releases any owned stream.
aclose() -> None
async
Release this transport and close the original only when owned.
Online lifecycle and subscriptions
rpkiparrot.client
Application-owned online source management and atomic reconfiguration.
Client
Online client bound to one host-owned AnyIO loop and task tree.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
config
|
ClientConfig
|
Complete explicit configuration, allowing initially empty groups. |
required |
transports
|
Mapping[str, TransportFactory] | None
|
RTR source ID to fresh-stream factory. Bindings are copied. |
None
|
readers
|
Mapping[str, JsonReader] | None
|
JSON source ID to explicitly registered JsonReader; custom bindings require reader_profile_id and never alter global formats. |
None
|
persistence
|
PersistenceBackend | None
|
Optional synchronous complete-state backend owned by this Client's dedicated worker thread. Mutually exclusive with configured persistence. All methods, including close, execute on that thread. |
None
|
Constructing does no I/O. Enter the async context to validate dependencies
and prepare resources, then use await task_group.start(client.run).
Cancellation of a wait or subscription does not stop source synchronization.
__init__(config: ClientConfig, *, transports: Mapping[str, TransportFactory] | None = None, readers: Mapping[str, JsonReader] | None = None, persistence: PersistenceBackend | None = None) -> None
__aenter__() -> Client
async
__aexit__(exc_type: type[BaseException] | None, exc: BaseException | None, tb: TracebackType | None) -> None
async
run(*, task_status: TaskStatus[None] = anyio.TASK_STATUS_IGNORED) -> None
async
Run source and expiry tasks within the caller's task group, once.
started signals operational infrastructure, not payload readiness. Unexpected implementation failures propagate to the host task group; expected individual source failures retry without stopping other sources.
apply_config(config: ClientConfig, *, expected_revision: ConfigRevision, transports: Mapping[str, TransportFactory] | None = None, readers: Mapping[str, JsonReader] | None = None) -> ConfigReceipt
async
Atomically replace running configuration after a revision comparison.
Added/replaced sources start unloaded; this call does not wait for them to synchronize. Old source incarnations cannot publish late work. A cancellation after the commit may lose the receipt but does not roll back the committed revision. Concurrent candidate preparation is bounded to one; conflicting, invalid or over-budget candidates leave state intact. Database identity and clock changes require constructing a new Client.
get_snapshot() -> Snapshot
async
Acquire the current coherent view without waiting for initial readiness.
Expired views are rebuilt under managed admission. Returned snapshots remain fixed and independently reject queries after their own boundary. Pure time work does not cancel admitted inputs. When all build slots are occupied, this call waits for a slot, including slots held by same-group source builds and their waiters. Cancelling the wait does not abandon an already running worker or extend data validity.
flush(*, timeout: float | None = None) -> PersistReceipt
async
Wait for the current complete state or a covering successor to persist.
Requires run and configured/injected persistence, otherwise raises ConfigurationError. PersistenceError reports a failed attempt or timeout; cancelling this waiter never abandons an active database transaction. The returned receipt identifies the actual durable version, which may be newer than the version captured by this call. An omitted timeout captures the current PersistenceConfig.flush_timeout, or 600 seconds for an injected backend. Hot changes apply only to subsequent calls. The wait budget starts after obtaining the target snapshot; it is not a hard wall-clock bound on this method or on managed transaction cleanup.
wait_ready(*, required: frozenset[PayloadKind], timeout: float = 30.0, aspa_afi: Afi | None = None) -> Snapshot
async
Wait for usable required data without stopping synchronization on timeout.
aspa_afi selects one ASPA family; None requires both to be usable, even when their authorizations differ. It requires ASPA in required.
watch() -> AsyncIterator[Subscription]
async
Atomically observe an initial snapshot and register its subsequent events.
No duplicate initial event is yielded. A slow subscriber receives sticky ResyncRequiredError and must close this context and observe again. initial_status captures current connection/persistence observations at registration without rewriting the fixed initial snapshot's payload.
get_status() -> StatusReport
Sample safe lifecycle/data diagnostics without performing network or SQL I/O.
get_metrics() -> MetricsReport
Sample completed operations and current usage for this publisher epoch.
Sync/connection and commit durations are accumulated monotonic seconds;
snapshot_build_seconds retains the latest completed build measurement;
persistence_commit_seconds retains the latest commit attempt duration.
A reconnect is a completed connection attempt
after the first, whether successful or failed. Cancellation is excluded.
Source record gauges use records/source/
aclose() -> None
async
Stop new work, cancel sources, and await every owned resource cleanup.
Idempotent; cleanup exceeding cleanup_timeout remains closing and appears overdue in get_status. Cancellation never abandons a worker or connection. A decided source revocation finishes under the same build budget before final persistence, even if its protocol worker reports after close begins. This necessary cleanup does not extend the original final flush deadline.
rpkiparrot.events
Bounded broadcast subscriptions for atomic snapshot publication.
Only Subscription is public. The owner serializes internal hub publication and registration without awaiting between snapshot exchange and event publication. Expensive snapshot comparisons can run in a managed worker before that commit.
Subscription
One independent bounded event stream and its atomic initial snapshot.
Obtain a subscription through Client.watch, not by constructing it. The initial snapshot belongs to the observation point before this stream's first event; it is not repeated by iteration. Retaining it does not extend its data validity. Use one async consumer per subscription. initial_status supplies current diagnostics at that same registration point, including connection observations since the fixed snapshot was published.
Queue/window overflow, publisher discontinuity or queue reconfiguration raises ResyncRequiredError on every subsequent next, including after aclose. Close that subscription, enter a new watch, and revalidate from its new initial snapshot. Normal close stops immediately and discards pending events rather than draining them. Cancelled next does not consume an event or unsubscribe; context exit/aclose releases this observer's resources.
initial: Snapshot
property
Immutable snapshot observed atomically when this watch registered.
initial_status: StatusReport | None
property
Current diagnostics sampled atomically with initial and registration.
Client and SharedSnapshotClient always supply this report. Snapshot source metadata remains fixed at its publication; connection-only and persistence observations since then are reflected here and followed by the event stream. None is reserved for internal standalone publishers.
__init__(initial: Snapshot, *, _hub: _EventHub, initial_status: StatusReport | None = None) -> None
__aenter__() -> Self
async
Use async context exit to release the observer on any outcome.
__aexit__(*exc: object) -> None
async
__aiter__() -> Self
__anext__() -> _Event
async
aclose() -> None
async
Unregister immediately; preserve any prior terminal resync signal.
Idempotent, with no I/O or cancellation checkpoint. Closing one observer neither closes the publisher nor interrupts another observer.
rpkiparrot.rtr.client
Application-owned AnyIO RTR sessions with synchronous transaction semantics.
SourceSink
Bases: Protocol
Atomically accept complete sources and observe transport/data status.
Implementations are called serially within the session's managed run task. They must never expose half a candidate on failure or cancellation. A status with UNAVAILABLE or EXPIRED capabilities revokes that source atomically before report returns; a retrying connection alone does not revoke still-valid data. Sink errors propagate unless a ResourceLimitError explicitly rejects admission before commit. publish returning successfully is the acknowledgement/serial advancement point. Cancellation or another exception before acknowledgement ends this single-use run, even if the sink had already committed a complete update. A completed or rejected query releases its candidate reservation before report; the controller admits required revocation as maintenance work. report may be entered with cancellation already pending after a managed protocol check. A standalone sink must protect its required trust-loss work; Client shields that work and awaits it during its own shutdown. For an already invalidated source, Client observes repeated unavailable connection reports without preempting another source's candidate.
publish(update: SourceUpdate) -> None
async
Commit one complete EOD candidate, preserving original times.
report(status: SourceInfo) -> None
async
Observe status; capability revocation has atomic data consequences.
RtrSession
One low-level RTR source, driven by a host-owned asyncio/Trio task.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
config
|
RtrSourceConfig
|
Validated source/transport configuration. v2 is opt-in; v2_profile explicitly selects -27, AFI-aware -10, or historical -13. |
required |
transport_factory
|
TransportFactory | None
|
A fresh-connection factory, or ConnectedTransport for one preconnected stream with explicit ownership. The built-in TCP/TLS/SSH factory is used when omitted. |
None
|
limits
|
Limits | None
|
Finite per-PDU and candidate budgets, captured per query. |
None
|
restored
|
SourceUpdate | None
|
Previously acknowledged complete SourceUpdate. Matching, unexpired state permits Serial Query; unusable recovery is reported as recovery_rejected and falls back to Reset Query. It never publishes or extends restored data before a fresh EOD. |
None
|
Enter the async context before task_group.start(session.run, sink).
Startup signals task ownership, not data readiness; publish occurs only
after valid EOD. Ordinary network failures retry while retaining original
deadlines. Fatal protocol errors report unavailable data for revocation.
A one-shot transport ends after its connection; it is never reused.
Standalone sessions stamp EOD with UTC wall time and use AnyIO's monotonic clock for live timers. When owned by Client, EOD and recovery instead use that owner's nondecreasing effective UTC authority. A new synchronization receives its complete protocol lifetime after a wall clock rollback; previously committed timestamps and expiry are never adjusted.
Construction performs no I/O. Entry prepares TLS/SSH dependencies/credentials. run is single-use and bound to the entered event loop. Cancellation clears staging and closes owned streams in a shielded cleanup; custom stream close must eventually finish. No loop or unmanaged background task is created.
__init__(config: RtrSourceConfig, *, transport_factory: TransportFactory | None = None, limits: Limits | None = None, restored: SourceUpdate | None = None) -> None
__aenter__() -> RtrSession
async
__aexit__(exc_type: type[BaseException] | None, exc: BaseException | None, traceback: TracebackType | None) -> None
async
run(sink: SourceSink, *, task_status: TaskStatus[None] = anyio.TASK_STATUS_IGNORED) -> None
async
Run once in the host task tree; propagate cancellation and sink bugs.
task_status.started is called after ownership is established, before networking. Protocol/network failures are reported then retried; source readiness is observed through publish, never inferred from startup.
aclose() -> None
async
Idempotently cancel run and await owned-resource cleanup.
Borrowed ConnectedTransport streams remain owned by the host. Closing
before run also closes an owned, not-yet-handed-out existing stream.
An owned stream or managed factory close failure raises TransportError
with code close_failed and retains its handle for a later attempt.
Complete-state persistence
rpkiparrot.persistence
Immutable complete-state persistence contracts and optional SQL backends.
Constructors do no I/O. Open each backend in its owning worker thread and keep all calls on that thread. Persistence never authenticates a source or renews its original deadlines; a host must re-evaluate identity, mode and time before publishing recovered data. SQLite and DuckDB drivers are imported on open.
BackendCapabilities
dataclass
Backend guarantees; shared_read means concurrent local SQLite observers.
backend_name: str
instance-attribute
Stable backend implementation name, such as sqlite or duckdb.
shared_read: bool
instance-attribute
Whether concurrent local shared readers are supported.
schema_version: int = 2
class-attribute
instance-attribute
Complete-state logical schema implemented by this backend; current version is 2.
single_owner: bool = True
class-attribute
instance-attribute
Whether the backend enforces one writer owner for the state.
atomic_commit: bool = True
class-attribute
instance-attribute
Whether a commit publishes one complete state atomically.
__init__(*, backend_name: str, shared_read: bool, schema_version: int = 2, single_owner: bool = True, atomic_commit: bool = True) -> None
PersistReceipt
dataclass
Actual durable head, including its original successful commit observation.
state_digest covers the canonical exchange payload, excluding committed_at. Neither that digest nor the backend's additional metadata checksum proves source authenticity. An idempotent retry returns the original receipt.
snapshot_id: SnapshotId
instance-attribute
Snapshot identity actually committed, not a requested or pending version.
config_revision: ConfigRevision
instance-attribute
Configuration revision durably committed with that snapshot.
committed_at: datetime
instance-attribute
Original successful UTC commit observation, retained on idempotent retries.
schema_version: int
instance-attribute
Logical schema of the acknowledged state; current receipts require version 2.
state_digest: str
instance-attribute
Lowercase SHA-256 of the canonical complete exchange payload; not an authenticity proof.
__init__(*, snapshot_id: SnapshotId, config_revision: ConfigRevision, committed_at: datetime, schema_version: int, state_digest: str) -> None
__post_init__() -> None
PersistedSource
dataclass
Original immutable records and source state, without runtime clock anchors.
dataset and info share existing frozen models. group_id is the original trust scope; evaluated_at is the fixed known evaluation floor. Unloaded sources are retained as precise empty incomplete placeholders.
dataset: ParsedDataset
instance-attribute
Complete original records, coverage and producer metadata, or an unloaded placeholder.
info: SourceInfo
instance-attribute
Nonsecret source identity, capability, protocol and original deadline metadata.
group_id: str
instance-attribute
Original trust group containing this source.
evaluated_at: datetime
instance-attribute
Fixed UTC evaluation floor already reached by this source state.
invalidated: bool = False
class-attribute
instance-attribute
Whether the source was explicitly revoked; retained records do not grant availability.
__init__(*, dataset: ParsedDataset, info: SourceInfo, group_id: str, evaluated_at: datetime, invalidated: bool = False) -> None
__post_init__() -> None
PersistedState
dataclass
One complete logical state, sharing record objects rather than JSON copies.
groups maps group ID to (integer priority, source ID tuple). active_groups selects one group independently per payload family. receipt is None before commit and populated by read_state. observed_at is a nondecreasing observed clock floor used for recovery; it and receipt are outside the stable payload digest. All payload times remain unchanged across retries and restarts. aspa_capabilities, aspa_active_groups and aspa_usable_until preserve each address family's independent view. None infers legacy both-family fields; current snapshots always supply the explicit maps. Schema 2 is required.
snapshot_id: SnapshotId
instance-attribute
Publisher snapshot identity of this complete state.
config_revision: ConfigRevision
instance-attribute
Configuration revision committed with the same complete state.
published_at: datetime
instance-attribute
Original UTC snapshot publication time.
context: EvaluationContext
instance-attribute
Online or fixed offline evaluation context of the state.
sources: tuple[PersistedSource, ...]
instance-attribute
Every configured original source, including precise unloaded placeholders.
groups: Mapping[str, tuple[int, tuple[str, ...]]]
instance-attribute
Group ID to priority and source-ID tuple; groups are separate trust scopes.
active_groups: Mapping[PayloadKind, str | None]
instance-attribute
Selected group per payload; aggregate ASPA may be None for differing family groups.
upstream_id: SnapshotId | None = None
class-attribute
instance-attribute
Original upstream publisher identity, when retained by an observer or import.
algorithm_versions: Mapping[str, str] = field(default_factory=lambda: dict(_VERSIONS))
class-attribute
instance-attribute
Exact fixed protocol and validation specification identifiers.
policy_id: str = 'none'
class-attribute
instance-attribute
Applied policy identifier; none explicitly denotes no local policy.
receipt: PersistReceipt | None = None
class-attribute
instance-attribute
Actual durable receipt after read_state; None on a state awaiting commit.
observed_at: datetime | None = None
class-attribute
instance-attribute
Nondecreasing UTC observation floor for recovery; excluded from the payload digest.
aspa_capabilities: Mapping[Afi, Availability] | None = None
class-attribute
instance-attribute
Independent IPv4/IPv6 availability; None requests legacy field inference.
aspa_active_groups: Mapping[Afi, str | None] | None = None
class-attribute
instance-attribute
Selected trust group per ASPA family; None requests legacy field inference.
aspa_usable_until: Mapping[Afi, datetime | None] | None = None
class-attribute
instance-attribute
Original next-recomputation boundary per family; None requests legacy inference.
__init__(*, snapshot_id: SnapshotId, config_revision: ConfigRevision, published_at: datetime, context: EvaluationContext, sources: tuple[PersistedSource, ...], groups: Mapping[str, tuple[int, tuple[str, ...]]], active_groups: Mapping[PayloadKind, str | None], upstream_id: SnapshotId | None = None, algorithm_versions: Mapping[str, str] = (lambda: dict(_VERSIONS))(), policy_id: str = 'none', receipt: PersistReceipt | None = None, observed_at: datetime | None = None, aspa_capabilities: Mapping[Afi, Availability] | None = None, aspa_active_groups: Mapping[Afi, str | None] | None = None, aspa_usable_until: Mapping[Afi, datetime | None] | None = None) -> None
__post_init__() -> None
PersistenceBackend
Bases: Protocol
Synchronous single-thread backend contract for complete atomic commits.
open binds the owning thread. read_state reads one consistent transaction; commit compares expected with the actual head (None means empty), verifies complete state and returns only after durable commit. close is idempotent. A host must await completion or rollback before releasing a worker.
capabilities: BackendCapabilities
property
Return supported backend guarantees without opening resources.
open() -> None
Acquire ownership, open the connection and validate/create schema.
read_state() -> PersistedState | None
Return a verified complete head with its receipt, or an empty database.
commit(state: PersistedState, *, expected: SnapshotId | None) -> PersistReceipt
Atomically compare and replace the head, or return an identical receipt.
close() -> None
Close connection and release ownership on the opening thread.
SQLiteBackend
Single-owner SQLite WAL backend requiring the sqlite extra.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
path
|
str | Path
|
Local database path, normalized before taking the owner lock. |
required |
limits
|
Limits | None
|
Complete-state load/commit budgets. |
None
|
busy_timeout
|
float
|
Finite seconds for SQLite lock waits; ownership conflicts fail immediately. WAL checkpoints are passive, not hard disk bounds. |
5.0
|
Construct without I/O; call open, read_state, commit and close on one worker thread. A second owner fails before opening the database. Shared observers are supplied by SharedSnapshotClient, never by another owner instance.
capabilities: BackendCapabilities
property
SQLite permits local short-transaction read-only observers.
__init__(path: str | Path, *, limits: Limits | None = None, busy_timeout: float = 5.0) -> None
open() -> None
Bind the opening thread, acquire ownership and check schema.
read_state() -> PersistedState | None
Read and verify a coherent head in one short read transaction.
Recognized transient driver/I/O failures raise PersistenceError with code persistence_failed and retryable=True. Detected schema/value corruption and library checksum mismatches remain persistence_corrupt. DuckDB's generic IOException can also mean block checksum damage: it raises persistence_failed without retryable=True, because the driver exposes no reliable subcode. No failed read publishes data or rewrites the file. Internal assertions propagate after rollback. The owner decides when to retry; the backend does not retry this call.
commit(state: PersistedState, *, expected: SnapshotId | None) -> PersistReceipt
Validate and atomically replace all tables with expected-head comparison.
Direct calls validate the complete state under this backend's limits. Managed builtin writes may reuse validation performed in the same call for the identical immutable state and limits. A changed backend budget still requires full validation before the transaction begins.
close() -> None
Close the connection before releasing ownership; idempotent.
DuckDBBackend
Bases: SQLiteBackend
Single-owner DuckDB backend requiring the duckdb extra.
All calls use one owning thread and an explicitly owned connection. This backend does not offer cross-process shared file readers; use the owner's service. Constructors are inert and the driver is imported only in open.
capabilities: BackendCapabilities
property
DuckDB allows one owner and no active shared file readers.
__init__(path: str | Path, *, limits: Limits | None = None) -> None
open() -> None
Bind the opening thread, acquire ownership and check schema.
read_state() -> PersistedState | None
Read and verify a coherent head in one short read transaction.
Recognized transient driver/I/O failures raise PersistenceError with code persistence_failed and retryable=True. Detected schema/value corruption and library checksum mismatches remain persistence_corrupt. DuckDB's generic IOException can also mean block checksum damage: it raises persistence_failed without retryable=True, because the driver exposes no reliable subcode. No failed read publishes data or rewrites the file. Internal assertions propagate after rollback. The owner decides when to retry; the backend does not retry this call.
commit(state: PersistedState, *, expected: SnapshotId | None) -> PersistReceipt
Validate and atomically replace all tables with expected-head comparison.
Direct calls validate the complete state under this backend's limits. Managed builtin writes may reuse validation performed in the same call for the identical immutable state and limits. A changed backend budget still requires full validation before the transaction begins.
close() -> None
Close the connection before releasing ownership; idempotent.
state_from_snapshot(snapshot: Snapshot) -> PersistedState
Capture original state in O(sources), sharing all immutable record tuples.
Runtime indexes and anchors are not serialized. A single clock observation supplies an additional receipt floor without mutating the fixed payload or changing the digest associated with an existing SnapshotId.
state_digest(state: PersistedState, *, limits: Limits | None = None) -> str
Hash canonical schema-two payload incrementally, without a full JSON copy.
The result equals the exchange payload SHA-256, including duplicate supports and fixed source evaluation times. Sorting holds record references, not a second dictionary per record. max_export_bytes limits canonical byte count. Small native records use bounded encoding fragments; large records and free metadata keep scalar streaming. This does not bound process memory. Mapped IPv6 prefixes use fixed mixed notation across Python versions. Raises InputError for inconsistent state and ResourceLimitError for limits.
SQLite shared observers
rpkiparrot.shared
Application-owned read-only observers of a local SQLite snapshot database.
SharedSnapshotClient
Observe complete SQLite commits with independent local time enforcement.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
config
|
SharedReaderConfig
|
Existing local SQLite path, polling interval, clock tolerance and finite resource limits. The read-only observer uses standard library sqlite3; the sqlite extra supplies locking for writers. |
required |
Construction does no I/O and starts no thread. Enter the async context to open a read-only connection and recover one coherent state, then start run in the host's AnyIO task group. All SQL uses one explicitly joined owning thread; no schema, source connections or writer locks are created.
Snapshot IDs belong to this reader. upstream_id and config_revision identify the observed writer state. Polling may coalesce writer generations; writer epoch changes require every existing watch to resynchronize. Read failures preserve previous data only until its original deadlines. A five-second cleanup reporting threshold never abandons a blocked operation or connection.
__init__(config: SharedReaderConfig) -> None
__aenter__() -> Self
async
Open and validate the existing database, then recover original deadlines.
__aexit__(*exc: object) -> None
async
Close owned resources while preserving any original body exception.
run(*, task_status: TaskStatus[None] = anyio.TASK_STATUS_IGNORED) -> None
async
Poll and expire data in the host task tree; started does not imply ready.
get_snapshot() -> Snapshot
async
Return the current view, locally re-evaluating time without any SQL read.
Head reconstruction and local time work share max_concurrent_builds. With one slot this call may wait for an admitted head worker; a caller can cancel its wait, while actual managed work retains its slot until completion. Previously obtained snapshots still reject expired data.
wait_ready(*, required: frozenset[PayloadKind], timeout: float = 30.0, aspa_afi: Afi | None = None) -> Snapshot
async
Wait for requested availability; timeout does not stop polling.
aspa_afi selects one family when required includes ASPA. None requires both families to be usable; it does not require identical providers. Specifying a family without ASPA in required is a configuration error.
watch() -> AsyncIterator[Subscription]
async
Atomically observe initial and subscribe; overflow requires a new watch.
get_status() -> StatusReport
Sample lifecycle, original deadlines and last read error without SQL I/O.
ready means VRP is usable; use wait_ready for another explicit required set. Persisted IDs and config_revision refer to the actual writer receipt. This diagnostic call is also available before entry and after close.
get_metrics() -> MetricsReport
Sample local reader operations and current bounded usage without I/O.
Source sync, reconnect and commit counters are zero: this observer does not run writer protocols or write the database. database_load_total counts accepted new durable heads, and coalesced_updates_total counts accepted jumps over intermediate generations within one writer epoch. Source gauges count original records with percent-encoded ID labels. Durations accumulate seconds except the latest snapshot_build_seconds.
aclose() -> None
async
Cancel polling and join all reads/builds before releasing the connection.
Idempotent after success. A failed backend close remains closing and may be retried; cancellation cannot leave an unowned SQLite worker running.
Remote HTTP SDK
rpkiparrot.remote
Optional HTTP SDK with typed domain results and fixed-snapshot iteration.
Importing this module needs only the core package. Entering RemoteClient needs
the http extra. asyncio and Trio are supported; the application owns its
event loop. No network or implicit background task is started by construction.
SnapshotRef
dataclass
Opaque reference to one retained server snapshot, with original deadlines.
The token is specific to its issuing service; do not parse it or compare its order. Retention does not extend usable_until. A reference carries no data or local trust: every request asks the server to check the original lifetime. Eviction raises SnapshotGoneError; expiration raises SnapshotExpiredError. ASPA capabilities, groups and deadlines remain separate for both families.
snapshot_id: SnapshotId
instance-attribute
Fixed snapshot publication identity associated with this value.
config_revision: ConfigRevision
instance-attribute
Configuration owner revision that defines this view.
snapshot_token: str
instance-attribute
Opaque service-issued token; it is scoped to that service and must not be parsed.
capabilities: Mapping[PayloadKind, Availability]
instance-attribute
Payload availability in this fixed view; empty and unavailable remain distinct.
usable_until: Mapping[PayloadKind, datetime | None]
instance-attribute
Exclusive next-recomputation boundary per payload, not an extension of source lifetime.
active_groups: Mapping[PayloadKind, str | None]
instance-attribute
Selected trust group per payload; aggregate ASPA is None if family groups differ.
aspa_capabilities: Mapping[Afi, Availability]
instance-attribute
Independent availability for IPv4 and IPv6 authorization scopes.
aspa_usable_until: Mapping[Afi, datetime | None]
instance-attribute
Exclusive next-recomputation boundary for each ASPA family, or None if unavailable.
aspa_active_groups: Mapping[Afi, str | None]
instance-attribute
Selected trust group per ASPA family, or None when no group is usable.
retention_deadline: datetime
instance-attribute
UTC retention deadline of this pin; reads do not renew it or any payload lifetime.
algorithm_versions: Mapping[str, str]
instance-attribute
Exact fixed protocol and validation specification identifiers.
mode: EvaluationMode
instance-attribute
Evaluation mode of the referenced server snapshot.
evaluated_at: datetime
instance-attribute
Server UTC evaluation time associated with this reference.
upstream_id: SnapshotId | None = None
class-attribute
instance-attribute
Original publisher snapshot identity retained by an observer or import, when available.
__post_init__() -> None
__init__(*, snapshot_id: SnapshotId, config_revision: ConfigRevision, snapshot_token: str, capabilities: Mapping[PayloadKind, Availability], usable_until: Mapping[PayloadKind, datetime | None], active_groups: Mapping[PayloadKind, str | None], aspa_capabilities: Mapping[Afi, Availability], aspa_usable_until: Mapping[Afi, datetime | None], aspa_active_groups: Mapping[Afi, str | None], retention_deadline: datetime, algorithm_versions: Mapping[str, str], mode: EvaluationMode, evaluated_at: datetime, upstream_id: SnapshotId | None = None) -> None
RemoteClient
HTTP client for one rpkiparrot /v1 service, used with async with.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
base_url
|
str
|
HTTP(S) service root, without credentials/query/fragment. |
required |
timeout
|
float
|
Positive total seconds per request, including body consumption, managed JSON parsing and typed model decoding. Defaults to 180 seconds; each page is a request, not a renewal of snapshot retention. |
180.0
|
token
|
str | None
|
Optional bearer credential, never copied into errors or results. |
None
|
transport
|
AsyncBaseTransport | None
|
Optional HTTPX async transport. This client owns and closes it; callers must not share its lifetime with another client. |
None
|
limits
|
Limits | None
|
Input batch/export limits; defaults to the core Limits. |
None
|
max_response_bytes
|
int
|
Decoded response cap for ordinary endpoints. Export uses limits.max_export_bytes instead. Oversize results fail whole. |
16 * 1024 * 1024
|
Connections are pooled, redirects and environment proxies are disabled. Default HTTPS verification uses HTTPX's certifi bundle, ignoring certificate environment variables and additional operating-system trust stores. Each request has one absolute deadline. Failures are typed domain errors or TransportError; unknown error codes survive as RpkiparrotError.code. Unknown schema/outcome enums raise ProtocolError(code='incompatible_response'). No query, page or subscription silently switches an explicit snapshot. Enter/use/close belong to one event loop. Pending watch registrations count against max_subscribers before network I/O. Closing rejects new work, cancels and joins admitted operations, then closes owned resources; managed decoder threads cannot be abandoned even after their request deadline. Entering prepares default TLS trust in a managed thread; concurrent close joins that preparation and disposes a late transport before returning. The default transport owns TCP attempts and incomplete TLS handshakes through cancellation. An injected transport retains its own I/O semantics. Default requests run in managed, joined workers: native caller cancellation requests structured cancellation and waits for protocol cleanup, including repeated cancellation during header/body reads. Normal EOF, response-limit and decoding failures all close the owned response and decoding iterators before returning. Default I/O close runs in a joined shielded task; iterators remain in their consuming worker. Cleanup preserves control-flow cancellation and does not replace an original failure with secondary ordinary errors. Injected streams retain their task ownership and close their own resources; their response closure need not close a shared connection. Raw cleanup for an injected stream can take up to one extra second before finalization.
__init__(base_url: str, *, timeout: float = 180.0, token: str | None = None, transport: httpx.AsyncBaseTransport | None = None, limits: Limits | None = None, max_response_bytes: int = 16 * 1024 * 1024) -> None
__aenter__() -> RemoteClient
async
__aexit__(exc_type: type[BaseException] | None, exc: BaseException | None, tb: TracebackType | None) -> None
async
aclose() -> None
async
Release subscriptions and owned HTTP resources, shielding final cleanup.
Local closure succeeds even when the server cannot acknowledge DELETE; unreachable remote watches expire under the server's idle timeout. Closing an already closed instance is harmless; it cannot be reopened. Concurrent close callers join the same cleanup. Pending registrations and request workers are joined, and subscription DELETEs run in parallel.
get_snapshot(*, required: Sequence[PayloadKind] = (), aspa_afi: Afi | None = None) -> SnapshotRef
async
Pin the current view, optionally requiring usable capabilities now.
No required capabilities means unavailable/empty initial state can be observed. aspa_afi selects one family's readiness; None requires both.
wait_ready(*, required: Sequence[PayloadKind] = (PayloadKind.VRP,), aspa_afi: Afi | None = None, timeout: float = 30.0) -> SnapshotRef
async
Poll readiness within one finite total deadline, returning a fixed ref.
Data unavailability is retried, while configuration, compatibility and transport failures propagate. Timeout never cancels server ownership.
validate_origin(prefix: Network | str, asn: int | None, *, snapshot: SnapshotRef | None = None, explain: bool = False) -> OriginResult
async
Apply RFC 6811 using one server view and its current freshness check.
validate_origins(routes: Sequence[OriginInput], *, snapshot: SnapshotRef | None = None, explain: bool = False) -> BatchResult[OriginResult]
async
Return a complete fixed-view batch with ordered per-item errors.
validate_aspa(path: AsPath, *, context: AspaContext, snapshot: SnapshotRef | None = None, afi: Afi | None = None, explain: bool = False) -> AspaResult
async
Validate a reconstructed path in the explicit family (or common view).
validate_aspa_batch(routes: Sequence[AspaInput], *, snapshot: SnapshotRef | None = None, explain: bool = False) -> BatchResult[AspaResult]
async
Evaluate family-specific ASPA inputs in order against one fixed view.
analyze_bmp_path(path: AsPath, *, context: BmpContext, snapshot: SnapshotRef | None = None, afi: Afi | None = None, explain: bool = False) -> BmpAnalysis
async
Return conservative BMP assessment, separate from ASPA validation.
iter_vrps(*, snapshot: SnapshotRef | None = None, asn: int | None = None, source_id: str | None = None, family: int | None = None, page_size: int | None = None) -> AsyncGenerator[VrpMatch, None]
async
Iterate deterministic VRPs; all pages pin one view and all filters.
Expiry/eviction fails without silently restarting. Abandoning iteration leaves no background task; server retention bounds the pinned token.
covering_vrps(prefix: Network | str, *, snapshot: SnapshotRef | None = None, page_size: int | None = None) -> tuple[VrpMatch, ...]
async
Collect every covering VRP from one fixed token, respecting batch cap.
iter_aspas(*, snapshot: SnapshotRef | None = None, source_id: str | None = None, afi: Afi | None = None, page_size: int | None = None) -> AsyncGenerator[AspaProviders, None]
async
Iterate customer unions within one fixed family/view, preserving support.
aspa_providers(customer: int, *, snapshot: SnapshotRef | None = None, afi: Afi | None = None) -> AspaProviders
async
Query one customer's complete provider union for the selected family.
compare_sources(left: str, right: str, *, snapshot: SnapshotRef | None = None) -> SourceDiff
async
Compare original source assertions without selecting a trusted winner.
export_snapshot(*, snapshot: SnapshotRef | None = None) -> bytes
async
Return complete bounded schema-2 export bytes, distinct from envelopes.
No local trust or recovery is created. Use load_snapshot explicitly if the caller chooses this server/export as an input trust source.
get_status() -> StatusReport
async
Read typed source/capability and acknowledged persistence status.
get_metrics() -> MetricsReport
async
Read process-local metrics; counters restart with the server epoch.
watch() -> AsyncIterator[RemoteSubscription]
async
Atomically observe an initial reference and its subsequent events.
The context owns DELETE cleanup. ResyncRequiredError requires a fresh watch; watch_conflict leaves continuity intact and can be retried. Pending creation reserves capacity; cancellation cleans up a known returned watch ID before releasing that reservation.
RemoteSubscription
One bounded server event stream, created by RemoteClient.watch().
initial is fixed at atomic registration. Exactly one task may consume
this object at a time. A request cancellation retains the last acknowledged
cursor; repeating it may replay events, deduplicated by publisher EventId.
The server may expire an idle watch; no implicit new initial point is made.
ResyncRequiredError is terminal for this observation. Concurrent close
callers wait for the same finite DELETE and active poll cleanup.
Attributes:
| Name | Type | Description |
|---|---|---|
initial |
Fixed snapshot reference captured atomically when the watch registered. |
|
initial_status |
Matching source and persistence observations at that registration point. |
initial = initial
instance-attribute
initial_status = initial_status
instance-attribute
__init__(client: RemoteClient, *, initial: SnapshotRef, initial_status: StatusReport, watch_id: str, cursor: str) -> None
__aiter__() -> RemoteSubscription
__anext__() -> _Event
async
aclose() -> None
async
Idempotently DELETE the watch; remote transport failure leaves idle GC.
Local buffers are released even when the server is unreachable. Closure shields the finite DELETE attempt from ambient cancellation.
HTTP service
rpkiparrot.service
Optional asyncio HTTP service; importing it does not require service extras.
ServiceConfig
dataclass
Bounded asyncio service configuration; constructing it performs no I/O.
host/port describe the listener the host must actually configure. workers must be one. Non-loopback listeners require bearer authentication and TLS. Direct TLS requires certificate/key paths and HTTPS ASGI requests. Proxy TLS requires bearer authentication, explicit peer IP/CIDR allowlists and one https forwarded protocol header; disable the ASGI server's proxy-header rewriting. Unauthenticated loopback mode also checks the actual request peer, rejecting missing/nonlocal peers.
token_provider is called once in a managed worker during lifespan startup. It supplies a nonempty bearer secret; the provider and its result are never included in status, snapshots or error diagnostics. Credential rotation requires a new service lifespan. The application configures no root logger.
Ordinary requests have finite execution and waiting slots. Long polls have separate execution slots and no queue. request_timeout covers admission, body receipt, processing, serialization and sending; managed workers retain their slot until they really finish after cancellation. Counts bound owned work, not Python heap usage. Limits in ClientConfig separately bound retained snapshots, subscriptions, event history, batches and exports. The default request deadline is 120 seconds; the SDK allows 180 seconds because its deadline also includes response parsing and verification.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
host
|
str
|
Actual listener IP or DNS name configured by the host. Only a literal loopback IP qualifies for unauthenticated local access. |
'127.0.0.1'
|
port
|
int
|
Actual listener port, an integer from 1 through 65535. |
8323
|
workers
|
int
|
Source-owning worker count; must be exactly one. |
1
|
auth_mode
|
Literal['none', 'bearer']
|
Either local-only |
'none'
|
token_provider
|
Callable[[], str] | None
|
Synchronous callable supplying the bearer secret once during startup; required exactly when auth_mode is bearer. Cancellation waits for its managed worker to return rather than abandoning it. |
None
|
tls_mode
|
Literal['none', 'direct', 'proxy']
|
|
'none'
|
tls_cert_file
|
str | Path | None
|
Certificate path for direct TLS; the host also passes it to the ASGI server. Configuration construction does not read it. |
None
|
tls_key_file
|
str | Path | None
|
Private-key path for direct TLS; required with the certificate path and never included in source snapshots or diagnostics. |
None
|
trusted_proxy_ips
|
tuple[str, ...]
|
Explicit IP/CIDR allowlist for proxy TLS, checked against the original TCP peer with proxy-header rewriting disabled. |
()
|
max_body_bytes
|
int
|
Buffered request-body byte limit; excess bodies return 413. |
4 * 1024 * 1024
|
max_response_bytes
|
int
|
Complete ordinary-response byte limit; snapshot export instead uses ClientConfig.limits.max_export_bytes. |
16 * 1024 * 1024
|
max_concurrent_requests
|
int
|
Execution slots for ordinary domain requests; cancelled managed work occupies its slot until it actually finishes. |
4
|
max_pending_requests
|
int
|
FIFO waiting slots for ordinary requests; zero disables waiting. Saturation returns 429 before full body intake. |
16
|
request_timeout
|
float
|
Total seconds from request entry through admission, body receipt, computation, serialization and response sending. |
120.0
|
max_concurrent_polls
|
int
|
Independent long-poll slots, without a waiting queue. |
32
|
long_poll_timeout
|
float
|
Maximum event-wait seconds per poll, shorter than request_timeout. A request may select a shorter wait, including zero. |
25.0
|
watch_idle_timeout
|
float
|
Seconds without an active poll before reclaiming a watch; an in-progress poll is not reclaimed as idle. |
60.0
|
page_size
|
int
|
Default number of items per page when the caller omits a limit. |
1000
|
max_page_size
|
int
|
Maximum explicit page size; must be at least page_size. |
10000
|
host: str = '127.0.0.1'
class-attribute
instance-attribute
port: int = 8323
class-attribute
instance-attribute
workers: int = 1
class-attribute
instance-attribute
auth_mode: Literal['none', 'bearer'] = 'none'
class-attribute
instance-attribute
token_provider: Callable[[], str] | None = field(default=None, repr=False, compare=False)
class-attribute
instance-attribute
tls_mode: Literal['none', 'direct', 'proxy'] = 'none'
class-attribute
instance-attribute
tls_cert_file: str | Path | None = None
class-attribute
instance-attribute
tls_key_file: str | Path | None = None
class-attribute
instance-attribute
trusted_proxy_ips: tuple[str, ...] = ()
class-attribute
instance-attribute
max_body_bytes: int = 4 * 1024 * 1024
class-attribute
instance-attribute
max_response_bytes: int = 16 * 1024 * 1024
class-attribute
instance-attribute
max_concurrent_requests: int = 4
class-attribute
instance-attribute
max_pending_requests: int = 16
class-attribute
instance-attribute
request_timeout: float = 120.0
class-attribute
instance-attribute
max_concurrent_polls: int = 32
class-attribute
instance-attribute
long_poll_timeout: float = 25.0
class-attribute
instance-attribute
watch_idle_timeout: float = 60.0
class-attribute
instance-attribute
page_size: int = 1000
class-attribute
instance-attribute
max_page_size: int = 10000
class-attribute
instance-attribute
__post_init__() -> None
__init__(*, host: str = '127.0.0.1', port: int = 8323, workers: int = 1, auth_mode: Literal['none', 'bearer'] = 'none', token_provider: Callable[[], str] | None = None, tls_mode: Literal['none', 'direct', 'proxy'] = 'none', tls_cert_file: str | Path | None = None, tls_key_file: str | Path | None = None, trusted_proxy_ips: tuple[str, ...] = (), max_body_bytes: int = 4 * 1024 * 1024, max_response_bytes: int = 16 * 1024 * 1024, max_concurrent_requests: int = 4, max_pending_requests: int = 16, request_timeout: float = 120.0, max_concurrent_polls: int = 32, long_poll_timeout: float = 25.0, watch_idle_timeout: float = 60.0, page_size: int = 1000, max_page_size: int = 10000) -> None
create_app(service_config: ServiceConfig, *, client_config: ClientConfig, transports: Mapping[str, TransportFactory] | None = None, readers: Mapping[str, JsonReader] | None = None, persistence: PersistenceBackend | None = None) -> FastAPI
Create an inert single-owner ASGI application with generated OpenAPI.
Requires the service extra. Host it on asyncio with one Uvicorn worker; configure the actual listener/TLS from ServiceConfig and disable automatic proxy-header rewriting. The ASGI lifespan owns Client, its source tasks, persistence and bounded watch pumps. Startup does not wait for source data.
Domain computation runs on managed threads; request cancellation never cancels background synchronization or abandons a still-running computation. app.state.client is available during lifespan for host apply_config calls. Transports/readers/persistence follow Client's explicit ownership contract. Construction creates no network connection, database or event loop. OpenAPI describes AFI values as integers and AFI object keys as strings, matching the JSON representation in nested status and event responses.
Protocol recording and offline replay
rpkiparrot.diagnostics
Explicit bounded RTR recording and synchronous offline protocol replay.
RecorderStatus
dataclass
A sampled recording status, independent of source synchronization.
complete becomes true only after successful finalization without gaps. bytes_written includes JSON framing; frames counts application byte chunks, not decoded PDUs. queued_bytes charges conservative encoded sizes, including the in-flight write. last_error contains safe, stable failure information.
lifecycle: str
instance-attribute
Recorder resource state, independent of RTR synchronization state.
complete: bool
instance-attribute
Whether finalization succeeded without gaps or lost recording events.
bytes_written: int
instance-attribute
Encoded bytes successfully written, including JSON framing.
frames: int
instance-attribute
Recorded application byte chunks; this is not a decoded PDU count.
queued_bytes: int
instance-attribute
Conservatively charged encoded bytes, including the in-flight write.
dropped: int
instance-attribute
Count of recording items lost to bounded queue or recording failures.
last_error: ErrorInfo | None
instance-attribute
Most recent safe recorder failure, or None when none has been reported.
__init__(*, lifecycle: str, complete: bool, bytes_written: int, frames: int, queued_bytes: int, dropped: int, last_error: ErrorInfo | None) -> None
ReplayDifference
dataclass
One bounded replay discrepancy, with at most 64 bytes of evidence.
sequence identifies the recording event, or None for a file-wide failure. code is stable; message is explanatory. expected/actual contain safe raw protocol prefixes, not transport credentials or a recoverable snapshot. reason identifies a file integrity failure (including an intact recorded gap versus a corrupt trailer) when code is recording_incomplete.
sequence: int | None
instance-attribute
Recording event number, or None for a file-wide discrepancy.
code: str
instance-attribute
Stable discrepancy category.
message: str
instance-attribute
Human-readable explanation, not a stable matching interface.
expected: bytes | None = None
class-attribute
instance-attribute
At most 64 expected protocol bytes, when applicable.
actual: bytes | None = None
class-attribute
instance-attribute
At most 64 observed protocol bytes, when applicable.
reason: str | None = None
class-attribute
instance-attribute
Specific integrity or gap reason for an incomplete recording, when applicable.
__init__(*, sequence: int | None, code: str, message: str, expected: bytes | None = None, actual: bytes | None = None, reason: str | None = None) -> None
ReplayEvent
dataclass
A virtual state transition, commit, error or unfinished transaction.
offset is seconds after the recording origin, mapped to reference_time. data is immutable JSON-compatible diagnostic content; commit payloads keep ASPA AFI but have no online source identity or restoration authority.
sequence: int
instance-attribute
Position of the corresponding event in the recording.
offset: float
instance-attribute
Seconds after the recording origin, mapped onto the chosen reference time.
kind: str
instance-attribute
Diagnostic event category, such as a transition, commit, error or unfinished transaction.
phase: str
instance-attribute
Virtual protocol phase associated with the event.
data: Mapping[str, object] = field(default_factory=dict)
class-attribute
instance-attribute
Immutable diagnostic fields; commit content is not trusted online state.
__post_init__() -> None
__init__(*, sequence: int, offset: float, kind: str, phase: str, data: Mapping[str, object] = dict()) -> None
ReplayReport
dataclass
Finite offline evidence, never an online snapshot or restoration input.
recording_complete means the input passed framing/hash/sequence checks. complete additionally requires reaching its end within report budgets; matched requires complete and no send/state discrepancies. events describe virtual immediate sink acknowledgements, not actual application commits. input_sha256 hashes the entire file. started_at and reference_time make the original/virtual clock mapping explicit. unverified lists unobserved host behavior even when a complete byte recording matches the protocol core.
schema_version: int
instance-attribute
Replay report schema version for consumers of diagnostic output.
protocol_versions: Mapping[str, str]
instance-attribute
Exact fixed protocol and profile identifiers used for replay.
input_sha256: str
instance-attribute
SHA-256 of the entire input recording file.
started_at: datetime
instance-attribute
Original recording start time used in the virtual-clock mapping.
reference_time: datetime
instance-attribute
Explicit UTC reference time onto which recording offsets are mapped.
recording_complete: bool
instance-attribute
Whether framing, sequence and integrity checks establish a complete input.
complete: bool
instance-attribute
Whether replay reached the input end within its report and resource budgets.
matched: bool
instance-attribute
Whether replay was complete and found no send or state discrepancies.
events: tuple[ReplayEvent, ...]
instance-attribute
Bounded virtual events; sink acknowledgements model immediate acceptance only.
differences: tuple[ReplayDifference, ...]
instance-attribute
Bounded discrepancies discovered during parsing or protocol replay.
unverified: tuple[str, ...]
instance-attribute
Host behaviors not established by byte replay, even when all recorded bytes match.
__post_init__() -> None
__init__(*, schema_version: int, protocol_versions: Mapping[str, str], input_sha256: str, started_at: datetime, reference_time: datetime, recording_complete: bool, complete: bool, matched: bool, events: tuple[ReplayEvent, ...], differences: tuple[ReplayDifference, ...], unverified: tuple[str, ...]) -> None
Recorder
Record one RTR session and its sequential reconnects without disk waits.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
path
|
str | Path
|
New JSONL file, created exclusively with owner-only permissions where supported. Existing files are never overwritten. |
required |
limits
|
Limits | None
|
Queue item/byte and total recording budgets. The in-flight write remains charged. Overflow stops sampling at an explicit gap. |
None
|
rtr_config
|
RtrSourceConfig | None
|
Fixed session controls for replay. Only source ID, version bounds, explicit v2 profile, ordering and query timeout are saved. Omission selects a fresh default v1 session; it never infers a draft. |
None
|
Construction and wrap perform no I/O. Enter the async context before using a wrapped factory. One fixed, non-daemon worker owns the file; aclose joins it even during cancellation. I/O errors are safe get_status diagnostics and do not break protocol traffic. Close all session streams before this object to obtain a complete recording. Same-instance concurrent sessions are rejected. Transport credentials and endpoints are never serialized. Raw application bytes can contain upstream diagnostic text, as explicitly selected by the caller. A complete file does not prove actual sink acknowledgement. Session entry and reconfiguration reject any disagreement with recorded session controls or protocol budgets; recording never silently follows a different query timeout or limit than its fixed header declares.
__init__(path: str | Path, *, limits: Limits | None = None, rtr_config: RtrSourceConfig | None = None) -> None
__aenter__() -> Recorder
async
__aexit__(exc_type: type[BaseException] | None, exc: BaseException | None, traceback: TracebackType | None) -> None
async
wrap(factory: TransportFactory | None = None) -> TransportFactory
Decorate a factory without changing stream ownership or reconnectability.
Omit factory to use the built-in TCP/TLS/SSH transport from an explicitly supplied rtr_config. Without that configuration an explicit factory is required. Neither form connects or reads credentials during wrap. The returned factory may be constructed before entry, but connect is valid only while this recorder is open. Its explicit prepare/aclose delegate lifecycle hooks; the recorder closes owned one-shot factories on exit, including streams never handed out. Returned stream close failures remain retryable and never discard the underlying handle.
get_status() -> RecorderStatus
Return immutable integrity/counter/error information without file I/O.
aclose() -> None
async
Stop sampling, drain accepted work, finalize integrity and join writer.
Cancellation never abandons a write or its thread. This does not close handed-out streams; their session still owns them. A live stream at recorder shutdown marks an incomplete trace. Calling again is safe. All unhanded managed factories receive a close attempt even if another fails; a cleanup failure leaves this recorder closing for retry.
replay(path: str | Path, *, reference_time: datetime, allow_incomplete: bool = False, limits: Limits | None = None, max_events: int = 10000, max_report_bytes: int = 16 * 1024 * 1024) -> ReplayReport
Validate JSONL and replay its protocol bytes on a synchronous virtual clock.
Only the selected input file is opened; there is no networking, sleeping, online database or event-loop creation. reference_time maps recording offset zero to an explicit UTC evaluation time. The header selects a fixed profile; legacy AFI is never guessed or merged. Recorded protocol limits may only be tightened by limits. Input bytes are bounded by max_recording_bytes.
max_events and max_report_bytes bound combined events/differences, reserving 1024 bytes for the final truncation notice. Exceeding either stops protocol evaluation with complete=False/report_limit while integrity validation still streams to EOF. Commits include counts and AFI-preserving ASPA diagnostics, not full VRP/key exports or usable SourceUpdate/Snapshot objects.
Missing trailers, gaps, hash/sequence/framing failures raise InputError unless allow_incomplete explicitly permits incomplete prefix diagnostics. Unknown header versions/session semantics remain errors. I/O failures raise InputError; oversized input raises ResourceLimitError. Ordinary send differences appear in the report. Complete matched transport evidence does not prove actual sink acknowledgement, initial recovery or external configuration changes. Pending protocol sends must precede subsequent receives. Unexplained local closure of an outstanding query is a difference; a quiescent local close is an observable disconnect. Incomplete reports retain a stable integrity reason, and a gap's trailer is verified before trusting that finite prefix.
Explicit validation metrics
rpkiparrot.metrics
Explicit bounded instrumentation for synchronous validation, without globals.
MetricsCollector
Instrument explicit calls while returning the original results and errors.
The five validation methods mirror :mod:rpkiparrot.validation, including
fixed-snapshot freshness, AFI selection, explanation limits and batch errors.
Pure functions and other collectors remain unaffected. Construction performs
no file/network I/O and creates no worker, event loop or background task.
Instances may be shared by managed worker threads: only constant-size metric
updates and sampling take a lock; domain computation runs outside that lock.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
timing
|
bool
|
Enable cumulative |
False
|
Counters
validation_total counts completed wrapper calls, including calls that
raise Exception; one batch is one call. validation_error_total counts
those calls with an exception, per-item error or incomplete batch. Normal
invalid/notfound/unknown/BMP indeterminate results are not execution errors.
validation_item_total and validation_item_error_total count known
item outcomes: a single call contributes one, a returned batch contributes
its items, and a batch that raises contributes no invented item outcomes.
Non-Exception BaseException interruptions increment only
validation_interrupted_total; they do not count as completed calls,
errors, items or timed observations. The fixed operation keys are
validation_{origin,origins,aspa,aspa_batch,bmp}_total and
validation_{origin,origins,aspa,aspa_batch,bmp}_error_total.
Durations
With timing enabled, validation_duration_seconds and the five
validation_<operation>_duration_seconds accumulate elapsed seconds
for completed calls, including execution errors. Concurrent calls each
contribute their own elapsed duration. Divide by corresponding call
totals for mean call latency; item totals separately describe throughput.
No arbitrary labels, source identities, inputs, result objects, exceptions or snapshots are retained. Report store_id/epoch identify this collector, not a data snapshot or manager. Create a new collector to start a new counting epoch; there is no reset, close or event-loop ownership requirement.
Raises:
| Type | Description |
|---|---|
InputError
|
timing is not a boolean. |
__init__(*, timing: bool = False) -> None
validate_origin(snapshot: Snapshot, prefix: Network | str, asn: int | None, *, explain: bool = False) -> OriginResult
Validate one origin, counting a completed call and its known item outcome.
Arguments, fixed-snapshot time checks, result identity and exceptions are those of validation.validate_origin. Invalid/notfound results are normal outcomes; only an execution exception increments the error counters.
validate_origins(snapshot: Snapshot, routes: Sequence[OriginInput], *, explain: bool = False) -> BatchResult[OriginResult]
Validate one bounded origin batch on the supplied fixed snapshot.
Calls validation.validate_origins unchanged. One returned batch adds one call, len(items) item outcomes and its per-item error count. A whole-call exception adds an erroneous call and no presumed item outcomes; an empty completed batch adds one successful call and zero items.
validate_aspa(snapshot: Snapshot, path: AsPath, *, context: AspaContext, afi: Afi | None = None, explain: bool = False) -> AspaResult
Validate an ASPA path with explicit context and optional address family.
Calls validation.validate_aspa unchanged, preserving source freshness and the requirement for equivalent family views when afi is None. Returns the same AspaResult; unknown/invalid results are not execution errors. Raised exceptions retain their original type, identity and diagnostic details.
validate_aspa_batch(snapshot: Snapshot, routes: Sequence[AspaInput], *, explain: bool = False) -> BatchResult[AspaResult]
Validate an ordered ASPA batch, preserving each input's AFI and context.
Calls validation.validate_aspa_batch unchanged. The batch counts as one completed call, including per-item errors; returned items separately count outcomes and errors. Whole-call exceptions do not invent item completions. The snapshot and online deadlines remain those of the original function.
analyze_bmp_path(snapshot: Snapshot, path: AsPath, *, context: BmpContext, afi: Afi | None = None, explain: bool = False) -> BmpAnalysis
Analyze one BMP observation without treating uncertainty as an error.
Calls validation.analyze_bmp_path with the original snapshot, context, family and explanation choice. The returned BmpAnalysis is not converted to standard ASPA valid; no_invalid_evidence and indeterminate are normal analysis outcomes. Execution exceptions retain their original identity.
get_metrics() -> MetricsReport
Atomically copy current counters and cumulative durations without I/O.
Returns an immutable MetricsReport sampled in UTC. All fixed counter keys start at zero; timing=False returns an empty durations mapping and gauges are always empty. Running calls have not contributed completed outcomes. Previously returned reports never change. store_id/epoch are this collector's stable UUID identity, independent of any measured snapshot.