Skip to content

API Reference

BaseRepository

metaorm.repositories.BaseRepository

Source code in metaorm/repositories.py
class BaseRepository:
    def __init_subclass__(
        cls,
        table: type[BaseTable] | None = None,
        filter_: type[BaseFilter] | None = None,
        dto: type[BaseModel] | None = None,
        **kwargs,
    ):
        super().__init_subclass__(**kwargs)

        if (table := table or getattr(cls, "_table_type", None)) is None:
            raise TypeError(
                f"{cls.__name__} must specify 'table' keyword argument",
            )
        if (filter_ := filter_ or getattr(cls, "_filter_type", None)) is None:
            raise TypeError(
                f"{cls.__name__} must specify 'filter_' keyword argument",
            )

        cls._table_type = table
        cls._filter_type = filter_
        cls._dto_type = dto or getattr(cls, "_dto_type", None)

    def __init__(
        self,
        settings: RepositorySettings | None = None,
        container: RepositoriesContainer | None = None,
    ):
        if container is not None:
            self._container = container
        elif settings is not None:
            self._container = RepositoriesContainer(settings=settings)
        else:
            raise TypeError("Either 'container' or 'settings' must be provided")

    async def get_items_count(
        self,
        filter_: BaseFilter | None = None,
    ) -> int:
        table = self.get_table_type()
        statement = select(func.count()).select_from(table)
        statement = append_to_statement(
            statement=statement,
            model=table,
            filter_=filter_,
        )

        async with self.transaction():
            result = await self.session.exec(statement)
            items_count = result.first()

        return items_count

    async def get_item(
        self,
        filter_: BaseFilter | None = None,
        sort: BaseSort | None = None,
    ) -> BaseModel | None:
        async for item in self.get_items(filter_=filter_, sort=sort):
            return item

    async def get_items(
        self,
        filter_: BaseFilter | None = None,
        pagination: BasePagination | None = None,
        sort: BaseSort | None = None,
        options: Sequence[Any] | None = None,
    ) -> AsyncGenerator[BaseModel]:
        table = self.get_table_type()
        statement = select(table)
        statement = append_to_statement(
            statement=statement,
            model=table,
            filter_=filter_,
            pagination=pagination,
            sort=sort,
        )
        if options:
            statement = statement.options(*options)

        async with self.transaction():
            result = await self.session.exec(statement)
            result = result.yield_per(100)
            for db_item in result:
                yield self._convert_from_table(db_item)

    async def create_item(self, item: BaseModel) -> BaseModel:
        table = self.get_table_type()
        values = self._convert_to_table(item).to_values()
        statement = insert(table).values(values).returning(table)

        async with self.transaction():
            try:
                result = await self.session.exec(statement)
            except IntegrityError as error:
                error_message = str(error.orig).lower()
                if "unique" in error_message or "duplicate" in error_message:
                    raise AlreadyExistsError(
                        f"Record for model '{table.__name__}' already exists",
                    ) from error
                raise DatabaseException(
                    f"Integrity error for model '{table.__name__}'",
                ) from error
            db_item = result.scalar_one()

        return self._convert_from_table(db_item)

    async def update_items(
        self,
        filter_: BaseFilter | None = None,
        options: Sequence[Any] | None = None,
        **values,
    ) -> AsyncGenerator[BaseModel]:
        table = self.get_table_type()
        statement = update(table)
        statement = append_to_statement(
            statement=statement,
            model=table,
            filter_=filter_,
        )
        statement = statement.values(**values).returning(table)
        if options:
            statement = statement.options(*options)

        async with self.transaction():
            result = await self.session.exec(statement)
            result = result.yield_per(100)
            for db_item in result.scalars():
                yield self._convert_from_table(db_item)

    async def delete_items(
        self,
        filter_: BaseFilter | None = None,
    ) -> None:
        table = self.get_table_type()
        statement = delete(table)
        statement = append_to_statement(
            statement=statement,
            model=table,
            filter_=filter_,
        )

        async with self.transaction():
            await self.session.exec(statement)

    @property
    def session(self) -> AsyncSession | None:
        return self._container.session

    @property
    def transaction(self):
        return self._container.transaction

    @property
    def nested_transaction(self):
        return self._container.nested_transaction

    async def create_tables(self) -> None:
        table = self.get_table_type()
        async with self._container.engine.begin() as connection:
            await connection.run_sync(
                table.metadata.create_all,
                tables=[table.__table__],
            )

    def _convert_to_table(self, item: BaseModel) -> BaseModel:
        table = self.get_table_type()
        if isinstance(item, table):
            return item
        return table.from_item(item=item)

    def _convert_from_table(self, table: BaseTable) -> BaseModel:
        if self.get_dto_type() is None:
            return table
        return table.to_item()

    @classmethod
    def get_table_type(cls) -> type[BaseTable] | None:
        return cls._table_type

    @classmethod
    def get_filter_type(cls) -> type[BaseFilter] | None:
        return cls._filter_type

    @classmethod
    def get_dto_type(cls) -> type[BaseModel] | None:
        return cls._dto_type

Methods:

__init_subclass__(table=None, filter_=None, dto=None, **kwargs)

Source code in metaorm/repositories.py
def __init_subclass__(
    cls,
    table: type[BaseTable] | None = None,
    filter_: type[BaseFilter] | None = None,
    dto: type[BaseModel] | None = None,
    **kwargs,
):
    super().__init_subclass__(**kwargs)

    if (table := table or getattr(cls, "_table_type", None)) is None:
        raise TypeError(
            f"{cls.__name__} must specify 'table' keyword argument",
        )
    if (filter_ := filter_ or getattr(cls, "_filter_type", None)) is None:
        raise TypeError(
            f"{cls.__name__} must specify 'filter_' keyword argument",
        )

    cls._table_type = table
    cls._filter_type = filter_
    cls._dto_type = dto or getattr(cls, "_dto_type", None)

__init__(settings=None, container=None)

Source code in metaorm/repositories.py
def __init__(
    self,
    settings: RepositorySettings | None = None,
    container: RepositoriesContainer | None = None,
):
    if container is not None:
        self._container = container
    elif settings is not None:
        self._container = RepositoriesContainer(settings=settings)
    else:
        raise TypeError("Either 'container' or 'settings' must be provided")

get_items_count(filter_=None) async

Source code in metaorm/repositories.py
async def get_items_count(
    self,
    filter_: BaseFilter | None = None,
) -> int:
    table = self.get_table_type()
    statement = select(func.count()).select_from(table)
    statement = append_to_statement(
        statement=statement,
        model=table,
        filter_=filter_,
    )

    async with self.transaction():
        result = await self.session.exec(statement)
        items_count = result.first()

    return items_count

get_item(filter_=None, sort=None) async

Source code in metaorm/repositories.py
async def get_item(
    self,
    filter_: BaseFilter | None = None,
    sort: BaseSort | None = None,
) -> BaseModel | None:
    async for item in self.get_items(filter_=filter_, sort=sort):
        return item

get_items(filter_=None, pagination=None, sort=None, options=None) async

Source code in metaorm/repositories.py
async def get_items(
    self,
    filter_: BaseFilter | None = None,
    pagination: BasePagination | None = None,
    sort: BaseSort | None = None,
    options: Sequence[Any] | None = None,
) -> AsyncGenerator[BaseModel]:
    table = self.get_table_type()
    statement = select(table)
    statement = append_to_statement(
        statement=statement,
        model=table,
        filter_=filter_,
        pagination=pagination,
        sort=sort,
    )
    if options:
        statement = statement.options(*options)

    async with self.transaction():
        result = await self.session.exec(statement)
        result = result.yield_per(100)
        for db_item in result:
            yield self._convert_from_table(db_item)

create_item(item) async

Source code in metaorm/repositories.py
async def create_item(self, item: BaseModel) -> BaseModel:
    table = self.get_table_type()
    values = self._convert_to_table(item).to_values()
    statement = insert(table).values(values).returning(table)

    async with self.transaction():
        try:
            result = await self.session.exec(statement)
        except IntegrityError as error:
            error_message = str(error.orig).lower()
            if "unique" in error_message or "duplicate" in error_message:
                raise AlreadyExistsError(
                    f"Record for model '{table.__name__}' already exists",
                ) from error
            raise DatabaseException(
                f"Integrity error for model '{table.__name__}'",
            ) from error
        db_item = result.scalar_one()

    return self._convert_from_table(db_item)

update_items(filter_=None, options=None, **values) async

Source code in metaorm/repositories.py
async def update_items(
    self,
    filter_: BaseFilter | None = None,
    options: Sequence[Any] | None = None,
    **values,
) -> AsyncGenerator[BaseModel]:
    table = self.get_table_type()
    statement = update(table)
    statement = append_to_statement(
        statement=statement,
        model=table,
        filter_=filter_,
    )
    statement = statement.values(**values).returning(table)
    if options:
        statement = statement.options(*options)

    async with self.transaction():
        result = await self.session.exec(statement)
        result = result.yield_per(100)
        for db_item in result.scalars():
            yield self._convert_from_table(db_item)

delete_items(filter_=None) async

Source code in metaorm/repositories.py
async def delete_items(
    self,
    filter_: BaseFilter | None = None,
) -> None:
    table = self.get_table_type()
    statement = delete(table)
    statement = append_to_statement(
        statement=statement,
        model=table,
        filter_=filter_,
    )

    async with self.transaction():
        await self.session.exec(statement)

create_tables() async

Source code in metaorm/repositories.py
async def create_tables(self) -> None:
    table = self.get_table_type()
    async with self._container.engine.begin() as connection:
        await connection.run_sync(
            table.metadata.create_all,
            tables=[table.__table__],
        )

get_table_type() classmethod

Source code in metaorm/repositories.py
@classmethod
def get_table_type(cls) -> type[BaseTable] | None:
    return cls._table_type

get_filter_type() classmethod

Source code in metaorm/repositories.py
@classmethod
def get_filter_type(cls) -> type[BaseFilter] | None:
    return cls._filter_type

get_dto_type() classmethod

Source code in metaorm/repositories.py
@classmethod
def get_dto_type(cls) -> type[BaseModel] | None:
    return cls._dto_type

BaseTable

metaorm.tables.BaseTable

Bases: SQLModel

Source code in metaorm/tables.py
class BaseTable[ItemType: BaseModel](SQLModel):
    @classmethod
    def from_item(cls, item: ItemType) -> Self:
        raise NotImplementedError

    def to_item(self) -> ItemType:
        raise NotImplementedError

    def to_values(self) -> dict[str, Any]:
        return {
            column.name: getattr(self, column.name) for column in self.__table__.columns
        }

RepositoriesContainer

metaorm.container.RepositoriesContainer

Source code in metaorm/container.py
class RepositoriesContainer:
    def __init__(self, settings: RepositorySettings):
        engine_parameters = {
            "url": settings.dsn,
            "pool_recycle": settings.pool_recycle,
        }
        if not settings.dsn.startswith("sqlite"):
            engine_parameters["pool_timeout"] = settings.pool_timeout
            engine_parameters["pool_size"] = settings.pool_size

        self._engine = create_async_engine(**engine_parameters)
        self._session_context = contextvars.ContextVar(
            "session_context",
            default=None,
        )

    @property
    def engine(self) -> AsyncEngine:
        return self._engine

    @property
    def session(self) -> AsyncSession | None:
        return self._session_context.get(None)

    async def create_tables(self, *repository_classes: type[RepositoryType]) -> None:
        repositories = [
            await self.get_repository(repository_class=repository_class)
            for repository_class in repository_classes
        ]
        for repository in repositories:
            await repository.create_tables()

    @asynccontextmanager
    async def transaction(self) -> AsyncGenerator[AsyncSession]:
        existing_session = self._session_context.get(None)
        if existing_session is not None:
            yield existing_session
            return

        session_parameters = {
            "bind": self._engine,
            "expire_on_commit": False,
        }
        async with AsyncSession(**session_parameters) as session:
            token = self._session_context.set(session)
            try:
                async with session.begin():
                    yield session
            finally:
                self._session_context.reset(token)

    @asynccontextmanager
    async def nested_transaction(self) -> AsyncGenerator[AsyncSession]:
        existing_session = self._session_context.get(None)
        if existing_session is not None:
            async with existing_session.begin_nested():
                yield existing_session
            return

        session_parameters = {
            "bind": self._engine,
            "expire_on_commit": False,
        }
        async with AsyncSession(**session_parameters) as session:
            token = self._session_context.set(session)
            try:
                async with session.begin_nested():
                    yield session
            finally:
                self._session_context.reset(token)

    def get_repository(self, repository_class: type[RepositoryType]) -> RepositoryType:
        return repository_class(container=self)

RepositorySettings

metaorm.settings.RepositorySettings

Bases: BaseModel

Source code in metaorm/settings.py
4
5
6
7
8
class RepositorySettings(BaseModel):
    dsn: str = Field(default="sqlite+aiosqlite:///db.sqlite3", pattern=r"^.+://")
    pool_size: int = Field(default=5, ge=1)
    pool_recycle: int = Field(default=60, ge=1)  # in seconds: 1 minute
    pool_timeout: int = Field(default=60, ge=1)  # in seconds: 1 minute

Exceptions

metaorm.exceptions.DatabaseException

Bases: Exception

Source code in metaorm/exceptions.py
class DatabaseException(Exception):
    pass

metaorm.exceptions.NotFoundError

Bases: DatabaseException

Source code in metaorm/exceptions.py
5
6
7
class NotFoundError(DatabaseException):
    def __init__(self, detail: str = "Not found"):
        super().__init__(detail)

metaorm.exceptions.HaveNoSessionError

Bases: DatabaseException

Source code in metaorm/exceptions.py
class HaveNoSessionError(DatabaseException):
    def __init__(self):
        super().__init__("Have no actual session")

metaorm.exceptions.AlreadyExistsError

Bases: DatabaseException

Source code in metaorm/exceptions.py
class AlreadyExistsError(DatabaseException):
    def __init__(self, detail: str = "Record already exists"):
        super().__init__(detail)

Re-exports

The following symbols are re-exported from metaorm for convenience:

  • BaseFilter, BasePagination, BaseSort, OffsetPagination, PagePagination — from pydantic-filters
  • Field, Relationship — from sqlmodel