Module adcp.reporting.projection.notifications

Isolated v2 status queues using the original SDK fanout/delivery transactions.

Classes

class PgReportingProjectionOutbox (*, pool: AsyncConnectionPool, clock: Callable[[], datetime] | None = None)
Expand source code
class PgReportingProjectionOutbox(PgReportingStatusOutbox):
    @asynccontextmanager
    async def _connection(self) -> AsyncIterator[Any]:
        async with self._pool.connection() as connection:
            yield _ProjectionQueueConnection(connection)

    async def create_schema(self) -> None:
        from adcp.reporting.projection.pg import PgReportingProjectionStore

        await PgReportingProjectionStore(pool=self._pool, notifications=True).create_schema()

Separate C events, expansion, delivery and activity, with generic HTTP logic.

Ancestors

Methods

async def create_schema(self) ‑> None
Expand source code
async def create_schema(self) -> None:
    from adcp.reporting.projection.pg import PgReportingProjectionStore

    await PgReportingProjectionStore(pool=self._pool, notifications=True).create_schema()

Inherited members