Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
17 changes: 14 additions & 3 deletions docs/reliable-reporting-service.md
Original file line number Diff line number Diff line change
Expand Up @@ -136,9 +136,20 @@ both managed delivery and receipts. Capability output is derived from the
components actually installed.

Unexpected failures are isolated by configuration or extension, logged, and
included on `ReliableReportingTurn`. A configured `worker_error_handler` can
page or emit telemetry; the lifecycle worker continues with its next turn
instead of silently stopping.
included on `ReliableReportingTurn.configuration_errors` or `extension_errors`.
Other configurations and extensions still run, and the background scheduler
retries on its next turn. These failures do not change service readiness.
`worker_error_handler(component, error)` receives the original exception and
the component name: `configuration:{account_id}:{delivery_config_id}@{version}`,
`materialization`, or `notification`. SDK logs contain only the component kind;
adopter callbacks control any further diagnostics. A callback failure is logged
without interrupting the remaining work.

Failures outside an individual configuration or extension turn stop the
background scheduler, notify the callback as `service` with the original error,
and withdraw readiness. Startup, scheduler, and resource cleanup failures are
terminal: after admitted work settles, `wait()` raises a sanitized
`ReliableReportingServiceError`. Restart those services with a new instance.

`run_worker()` is also available for an external scheduler. Calls within one
service process are serialized. Ledger writes are convergent, but deployments
Expand Down
28 changes: 28 additions & 0 deletions src/adcp/reporting/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -92,6 +92,21 @@
from adcp.reporting.service import (
ReportingAdapter as ReportingAdapter,
)
from adcp.reporting.service_lifecycle import (
ReliableReportingServiceError as ReliableReportingServiceError,
)
from adcp.reporting.service_lifecycle import (
ReliableReportingShutdownTimeoutError as ReliableReportingShutdownTimeoutError,
)
from adcp.reporting.service_lifecycle import (
ReliableReportingState as ReliableReportingState,
)
from adcp.reporting.service_lifecycle import (
ReliableReportingUnavailableError as ReliableReportingUnavailableError,
)
from adcp.reporting.service_lifecycle import (
ReportingServiceResource as ReportingServiceResource,
)
from adcp.reporting.testing import (
DeterministicReportingClock as DeterministicReportingClock,
)
Expand Down Expand Up @@ -121,6 +136,14 @@
_LAZY_EXPORTS = {
"ReliableReportingConfigurationError": ("service", "ReliableReportingConfigurationError"),
"ReliableReportingService": ("service", "ReliableReportingService"),
"ReliableReportingServiceError": ("service_lifecycle", "ReliableReportingServiceError"),
"ReliableReportingShutdownTimeoutError": (
"service_lifecycle",
"ReliableReportingShutdownTimeoutError",
),
"ReliableReportingState": ("service_lifecycle", "ReliableReportingState"),
"ReliableReportingUnavailableError": ("service_lifecycle", "ReliableReportingUnavailableError"),
"ReportingServiceResource": ("service_lifecycle", "ReportingServiceResource"),
"ReportingAccountContext": ("service", "ReportingAccountContext"),
"ReportingAdapter": ("service", "ReportingAdapter"),
"DeterministicReportingClock": ("testing", "DeterministicReportingClock"),
Expand Down Expand Up @@ -161,6 +184,11 @@ def __dir__() -> list[str]:
"ReportingTier",
"ReliableReportingConfigurationError",
"ReliableReportingService",
"ReliableReportingServiceError",
"ReliableReportingShutdownTimeoutError",
"ReliableReportingState",
"ReliableReportingUnavailableError",
"ReportingServiceResource",
"ReportingAccountContext",
"ReportingAdapter",
"DeterministicReportingClock",
Expand Down
47 changes: 47 additions & 0 deletions src/adcp/reporting/_settlement.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,47 @@
"""Settlement of work Python cannot safely interrupt."""

from __future__ import annotations

import asyncio
from typing import Any, TypeVar

_T = TypeVar("_T")


async def settle_task(task: asyncio.Task[_T]) -> _T:
"""Defer caller cancellation until an owned task has actually finished.

Shield alone only keeps the child alive: the parent still returns early on
cancellation. Retain the child and consume its outcome, including under
repeated cancellation, before propagating the caller's cancellation. This
barrier deliberately has no timeout; a supervisor can time out its *wait*
while retaining ownership of the still-running operation.
"""
cancelled: asyncio.CancelledError | None = None
while not task.done():
try:
await asyncio.shield(task)
except asyncio.CancelledError as error:
cancelled = error
except Exception:
break # result() below retrieves the child's exception exactly once
if cancelled is not None:
if not task.cancelled():
task.exception()
raise cancelled
return task.result()


async def cancel_and_settle(task: asyncio.Task[Any]) -> None:
"""Request cancellation, then join without losing repeated caller cancels.

Joining through a task that consumes the child's cancellation lets us
distinguish its expected CancelledError from a new cancellation of the
caller. The latter is propagated only after the child has settled.
"""
task.cancel()

async def join() -> None:
await asyncio.gather(task, return_exceptions=True)

await settle_task(asyncio.create_task(join()))
15 changes: 12 additions & 3 deletions src/adcp/reporting/inline_source.py
Original file line number Diff line number Diff line change
Expand Up @@ -73,6 +73,7 @@

from pydantic import TypeAdapter, ValidationError

from adcp.reporting._settlement import settle_task
from adcp.reporting.currency import (
ReportingCurrencyError,
validate_currency,
Expand Down Expand Up @@ -470,7 +471,7 @@ async def stage(
digest = hashlib.sha256(payload).hexdigest()
object_ref = f"{source_execution_key}.{ordinal}"
target = self._path(account_id, object_ref, digest)
await asyncio.to_thread(self._write, target, payload)
await settle_task(asyncio.create_task(asyncio.to_thread(self._write, target, payload)))
return object_ref, digest

@staticmethod
Expand Down Expand Up @@ -499,7 +500,9 @@ async def read(
cancel: asyncio.Event,
) -> bytes:
target = self._path(account_id, object_ref, object_generation)
payload: bytes = await asyncio.to_thread(target.read_bytes)
payload: bytes = await settle_task(
asyncio.create_task(asyncio.to_thread(target.read_bytes))
)
if hashlib.sha256(payload).hexdigest() != object_generation:
raise OSError("staged object bytes no longer match their pinned generation")
return payload
Expand Down Expand Up @@ -724,7 +727,13 @@ async def _invoke(
if _is_async_callable(self._fetch):
answer = await self._fetch(request)
else:
answer = await asyncio.shield(asyncio.to_thread(self._fetch, request))
# Keep both the thread and a dynamically returned awaitable owned
# until they settle. Shield alone abandons the wait on cancellation.
async def run_sync() -> Any:
result = await asyncio.to_thread(self._fetch, request)
return await result if isinstance(result, Awaitable) else result

answer = await settle_task(asyncio.create_task(run_sync()))
if isinstance(answer, Awaitable):
# A plain callable that hands back a coroutine; await it here
# rather than letting it reach the coercion step unfinished.
Expand Down
12 changes: 11 additions & 1 deletion src/adcp/reporting/ledger/producer.py
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,7 @@
from datetime import datetime, timedelta, timezone
from typing import TYPE_CHECKING, Any, TypeAlias

from adcp.reporting._settlement import cancel_and_settle
from adcp.reporting.canonical_json import canonical_json_utf8_v1
from adcp.reporting.currency import (
ReportingCurrencyError,
Expand Down Expand Up @@ -832,13 +833,22 @@ async def acquire_obligation(
constituents=constituents,
)
cancel = asyncio.Event()
execution = asyncio.create_task(self._source.execute(request, cancel=cancel))
try:
result = await asyncio.wait_for(
self._source.execute(request, cancel=cancel),
asyncio.shield(execution),
timeout=self._offerings.slice_timeout.total_seconds(),
)
except asyncio.CancelledError:
cancel.set()
await cancel_and_settle(execution)
raise
except asyncio.TimeoutError:
cancel.set()
# wait_for's own cancellation join can be interrupted by a second
# cancellation of this producer (notably on Python 3.10). Retain
# the execution explicitly until even synchronous work has settled.
await cancel_and_settle(execution)
turn.slices_failed.append(obligation.reporting_obligation_id)
self._note_escalation(obligation, turn, now=now)
return None
Expand Down
Loading
Loading