Module adcp.reporting.ledger.provisional

Durable scheduling identities and immutable provisional observation metadata.

These records are private persistence metadata, not new reporting wire fields. A source execution's semantic request is frozen; its deadline is an attempt budget.

Functions

def utc(value: datetime) ‑> datetime.datetime
Expand source code
def utc(value: datetime) -> datetime:
    if value.tzinfo is None or value.utcoffset() is None:
        raise ValueError("provisional observation timestamps must be timezone-aware")
    return value.astimezone(timezone.utc)

Classes

class PendingProvisionalAcquisitionStore (*args, **kwargs)
Expand source code
@runtime_checkable
class PendingProvisionalAcquisitionStore(Protocol):
    """Optional lookup of a frozen acquisition before selecting its readiness window."""

    async def get_provisional_acquisition(
        self, *, account_id: str, reporting_obligation_id: str, ordinal: int
    ) -> ProvisionalAcquisition | None:
        pass

Optional lookup of a frozen acquisition before selecting its readiness window.

Ancestors

  • typing.Protocol
  • typing.Generic

Methods

async def get_provisional_acquisition(self, *, account_id: str, reporting_obligation_id: str, ordinal: int) ‑> ProvisionalAcquisition | None
Expand source code
async def get_provisional_acquisition(
    self, *, account_id: str, reporting_obligation_id: str, ordinal: int
) -> ProvisionalAcquisition | None:
    pass
class ProvisionalAcquisition (request_json: str,
ordinal: int,
policy: ProvisionalPolicy,
predecessor_revision_id: str | None = None,
previous_boundary: datetime | None = None)
Expand source code
@dataclass(frozen=True)
class ProvisionalAcquisition:
    request_json: str
    ordinal: int
    policy: ProvisionalPolicy
    predecessor_revision_id: str | None = None
    previous_boundary: datetime | None = None

    def __post_init__(self) -> None:
        if type(self.ordinal) is not int or self.ordinal < 0:
            raise ValueError("observation ordinal must be a nonnegative integer")
        self.request()
        if self.previous_boundary is not None:
            utc(self.previous_boundary)

    def request(self, *, deadline_at: datetime | None = None) -> ReportingSourceSliceRequestV1:
        # Parse on every access: callers cannot mutate the durable request's lists.
        request = ReportingSourceSliceRequestV1.model_validate_json(self.request_json)
        if deadline_at is not None:
            request = request.model_copy(update={"deadline_at": utc(deadline_at)})
        return request

    @property
    def account_id(self) -> str:
        return self.request().identity.account_id

    @property
    def obligation_id(self) -> str:
        return self.request().identity.reporting_obligation_id

    @property
    def execution_key(self) -> str:
        return self.request().identity.source_execution_key

    def binds(self, obligation: ReportingObligationRecord) -> bool:
        request = self.request()
        identity = request.identity
        return (
            identity.account_id == obligation.account_id
            and identity.reporting_obligation_id == obligation.reporting_obligation_id
            and identity.delivery_config_id == obligation.delivery_config_id
            and identity.delivery_config_version == obligation.delivery_config_version
            and identity.report_definition_id == obligation.report_definition_id
            and request.period.start == obligation.period.start
            and request.period.end == obligation.period.end
            and request.currency == obligation.currency
        )

    def to_wire(self) -> dict[str, Any]:
        return {
            "request": json.loads(self.request_json),
            "ordinal": self.ordinal,
            "policy": self.policy.to_wire(),
            "predecessor_revision_id": self.predecessor_revision_id,
            "previous_boundary": (
                utc(self.previous_boundary).isoformat()
                if self.previous_boundary is not None
                else None
            ),
        }

    @classmethod
    def from_wire(cls, value: dict[str, Any]) -> ProvisionalAcquisition:
        return cls(
            json.dumps(value["request"], sort_keys=True, separators=(",", ":")),
            value["ordinal"],
            ProvisionalPolicy.from_wire(value["policy"]),
            value["predecessor_revision_id"],
            (
                datetime.fromisoformat(value["previous_boundary"])
                if value["previous_boundary"] is not None
                else None
            ),
        )

ProvisionalAcquisition(request_json: 'str', ordinal: 'int', policy: 'ProvisionalPolicy', predecessor_revision_id: 'str | None' = None, previous_boundary: 'datetime | None' = None)

Static methods

def from_wire(value: dict[str, Any]) ‑> ProvisionalAcquisition

Instance variables

prop account_id : str
Expand source code
@property
def account_id(self) -> str:
    return self.request().identity.account_id
prop execution_key : str
Expand source code
@property
def execution_key(self) -> str:
    return self.request().identity.source_execution_key
prop obligation_id : str
Expand source code
@property
def obligation_id(self) -> str:
    return self.request().identity.reporting_obligation_id
var ordinal : int
var policy : ProvisionalPolicy
var predecessor_revision_id : str | None
var previous_boundary : datetime.datetime | None
var request_json : str

Methods

def binds(self, obligation: ReportingObligationRecord) ‑> bool
Expand source code
def binds(self, obligation: ReportingObligationRecord) -> bool:
    request = self.request()
    identity = request.identity
    return (
        identity.account_id == obligation.account_id
        and identity.reporting_obligation_id == obligation.reporting_obligation_id
        and identity.delivery_config_id == obligation.delivery_config_id
        and identity.delivery_config_version == obligation.delivery_config_version
        and identity.report_definition_id == obligation.report_definition_id
        and request.period.start == obligation.period.start
        and request.period.end == obligation.period.end
        and request.currency == obligation.currency
    )
def request(self, *, deadline_at: datetime | None = None) ‑> ReportingSourceSliceRequestV1
Expand source code
def request(self, *, deadline_at: datetime | None = None) -> ReportingSourceSliceRequestV1:
    # Parse on every access: callers cannot mutate the durable request's lists.
    request = ReportingSourceSliceRequestV1.model_validate_json(self.request_json)
    if deadline_at is not None:
        request = request.model_copy(update={"deadline_at": utc(deadline_at)})
    return request
def to_wire(self) ‑> dict[str, typing.Any]
Expand source code
def to_wire(self) -> dict[str, Any]:
    return {
        "request": json.loads(self.request_json),
        "ordinal": self.ordinal,
        "policy": self.policy.to_wire(),
        "predecessor_revision_id": self.predecessor_revision_id,
        "previous_boundary": (
            utc(self.previous_boundary).isoformat()
            if self.previous_boundary is not None
            else None
        ),
    }
class ProvisionalObservation (acquisition: ProvisionalAcquisition,
revision_id: str,
checked_at: datetime,
provisional_until: datetime,
next_due_at: datetime | None,
manifest_json: str)
Expand source code
@dataclass(frozen=True)
class ProvisionalObservation:
    acquisition: ProvisionalAcquisition
    revision_id: str
    checked_at: datetime
    provisional_until: datetime
    next_due_at: datetime | None
    manifest_json: str

    def __post_init__(self) -> None:
        utc(self.checked_at)
        utc(self.provisional_until)
        if self.next_due_at is not None:
            utc(self.next_due_at)

    def to_wire(self) -> dict[str, Any]:
        return {
            "acquisition": self.acquisition.to_wire(),
            "revision_id": self.revision_id,
            "checked_at": utc(self.checked_at).isoformat(),
            "provisional_until": utc(self.provisional_until).isoformat(),
            "next_due_at": (
                utc(self.next_due_at).isoformat() if self.next_due_at is not None else None
            ),
            "manifest": json.loads(self.manifest_json),
        }

    @classmethod
    def from_wire(cls, value: dict[str, Any]) -> ProvisionalObservation:
        return cls(
            ProvisionalAcquisition.from_wire(value["acquisition"]),
            value["revision_id"],
            datetime.fromisoformat(value["checked_at"]),
            datetime.fromisoformat(value["provisional_until"]),
            (
                datetime.fromisoformat(value["next_due_at"])
                if value["next_due_at"] is not None
                else None
            ),
            json.dumps(value["manifest"], sort_keys=True, separators=(",", ":")),
        )

ProvisionalObservation(acquisition: 'ProvisionalAcquisition', revision_id: 'str', checked_at: 'datetime', provisional_until: 'datetime', next_due_at: 'datetime | None', manifest_json: 'str')

Static methods

def from_wire(value: dict[str, Any]) ‑> ProvisionalObservation

Instance variables

var acquisition : ProvisionalAcquisition
var checked_at : datetime.datetime
var manifest_json : str
var next_due_at : datetime.datetime | None
var provisional_until : datetime.datetime
var revision_id : str

Methods

def to_wire(self) ‑> dict[str, typing.Any]
Expand source code
def to_wire(self) -> dict[str, Any]:
    return {
        "acquisition": self.acquisition.to_wire(),
        "revision_id": self.revision_id,
        "checked_at": utc(self.checked_at).isoformat(),
        "provisional_until": utc(self.provisional_until).isoformat(),
        "next_due_at": (
            utc(self.next_due_at).isoformat() if self.next_due_at is not None else None
        ),
        "manifest": json.loads(self.manifest_json),
    }
class ProvisionalObservationStore (*args, **kwargs)
Expand source code
@runtime_checkable
class ProvisionalObservationStore(Protocol):
    async def reserve_provisional_acquisition(
        self, acquisition: ProvisionalAcquisition
    ) -> ProvisionalAcquisition:
        pass

    async def get_provisional_observation(
        self, *, account_id: str, reporting_obligation_id: str
    ) -> ProvisionalObservation | None:
        pass

    async def commit_provisional_observation(
        self,
        observation: ProvisionalObservation,
        revision: ReportingRevisionRecord,
        rows: Sequence[dict[str, Any]],
    ) -> ReportingRevisionRecord:
        pass

Base class for protocol classes.

Protocol classes are defined as::

class Proto(Protocol):
    def meth(self) -> int:
        ...

Such classes are primarily used with static type checkers that recognize structural subtyping (static duck-typing).

For example::

class C:
    def meth(self) -> int:
        return 0

def func(x: Proto) -> int:
    return x.meth()

func(C())  # Passes static type check

See PEP 544 for details. Protocol classes decorated with @typing.runtime_checkable act as simple-minded runtime protocols that check only the presence of given attributes, ignoring their type signatures. Protocol classes can be generic, they are defined as::

class GenProto(Protocol[T]):
    def meth(self) -> T:
        ...

Ancestors

  • typing.Protocol
  • typing.Generic

Methods

async def commit_provisional_observation(self,
observation: ProvisionalObservation,
revision: ReportingRevisionRecord,
rows: Sequence[dict[str, Any]]) ‑> ReportingRevisionRecord
Expand source code
async def commit_provisional_observation(
    self,
    observation: ProvisionalObservation,
    revision: ReportingRevisionRecord,
    rows: Sequence[dict[str, Any]],
) -> ReportingRevisionRecord:
    pass
async def get_provisional_observation(self, *, account_id: str, reporting_obligation_id: str) ‑> ProvisionalObservation | None
Expand source code
async def get_provisional_observation(
    self, *, account_id: str, reporting_obligation_id: str
) -> ProvisionalObservation | None:
    pass
async def reserve_provisional_acquisition(self,
acquisition: ProvisionalAcquisition) ‑> ProvisionalAcquisition
Expand source code
async def reserve_provisional_acquisition(
    self, acquisition: ProvisionalAcquisition
) -> ProvisionalAcquisition:
    pass
class ProvisionalPolicy (window: timedelta,
cadence: timedelta,
official_close_lag: timedelta | None = None)
Expand source code
@dataclass(frozen=True)
class ProvisionalPolicy:
    window: timedelta
    cadence: timedelta
    official_close_lag: timedelta | None = None

    def __post_init__(self) -> None:
        if self.window < timedelta(0) or self.cadence <= timedelta(0):
            raise ValueError("provisional window must be nonnegative and cadence positive")
        if self.official_close_lag is not None and self.official_close_lag < timedelta(0):
            raise ValueError("official close lag must be nonnegative")

    def to_wire(self) -> dict[str, int | None]:
        return {
            "window_us": self.window // timedelta(microseconds=1),
            "cadence_us": self.cadence // timedelta(microseconds=1),
            "official_close_lag_us": (
                self.official_close_lag // timedelta(microseconds=1)
                if self.official_close_lag is not None
                else None
            ),
        }

    @classmethod
    def from_wire(cls, value: dict[str, Any]) -> ProvisionalPolicy:
        lag = value["official_close_lag_us"]
        return cls(
            timedelta(microseconds=value["window_us"]),
            timedelta(microseconds=value["cadence_us"]),
            timedelta(microseconds=lag) if lag is not None else None,
        )

ProvisionalPolicy(window: 'timedelta', cadence: 'timedelta', official_close_lag: 'timedelta | None' = None)

Static methods

def from_wire(value: dict[str, Any]) ‑> ProvisionalPolicy

Instance variables

var cadence : datetime.timedelta
var official_close_lag : datetime.timedelta | None
var window : datetime.timedelta

Methods

def to_wire(self) ‑> dict[str, int | None]
Expand source code
def to_wire(self) -> dict[str, int | None]:
    return {
        "window_us": self.window // timedelta(microseconds=1),
        "cadence_us": self.cadence // timedelta(microseconds=1),
        "official_close_lag_us": (
            self.official_close_lag // timedelta(microseconds=1)
            if self.official_close_lag is not None
            else None
        ),
    }