Skip to content

Custom task

Run a SalaryInfo custom task

main.py
"""Run a salary_info custom task.

This method running a task created on the basis of a quantum loop.
Effectiveness running task depends on the number of processor threads.
"""

import anyio
from datetime import datetime
from typing import Annotated, Any
from zoneinfo import ZoneInfo

from pydantic import EmailStr, Field
from pydantic_extra_types.phone_numbers import PhoneNumber, PhoneNumberValidator

from scruby import CustomTask, Scruby, ScrubyModel
from scruby.aggregation import Average, Max, Min, Sum
from pprint import pprint as pp


class Salesman(ScrubyModel):
    """Salesman model."""

    username: str
    first_name: str
    last_name: str
    birthday: datetime
    email: EmailStr
    phone: Annotated[PhoneNumber, PhoneNumberValidator(number_format="E164"), Field(strict=False)]
    salary: int
    # key is always at bottom
    key: Annotated[
        str,
        Field(
            frozen=True,
            default_factory=lambda data: data["phone"],
        ),
    ]


class SalaryInfo(CustomTask):
    """Custom task.

    Get information about sales salaries.

    The result should be the fields:
    `max_salary`, `min_salary`, `average_salary`, `sum_salaries`, `count_sellers`, and `salesman_list`.
    """

    def __init__(self) -> None:
        """Initializing the task."""
        self.stop_signal = False
        self.max_salary = Max()
        self.min_salary = Min()
        self.average_salary = Average()
        self.sum_salaries = Sum()
        self.salesman_list: list[Any] = []
        self.salary_info: dict[str, Any] = {}

    def accept(self, doc: Any) -> None:
        """Operation with a document."""
        self.max_salary.set(doc.salary)
        self.min_salary.set(doc.salary)
        self.average_salary.set(doc.salary)
        self.sum_salaries.set(doc.salary)
        self.salesman_list.append(doc)

    def result(self) -> Any | None:
        """Return result."""
        # Add data to result
        count_sellers = len(self.salesman_list)
        if count_sellers > 0:
            self.salary_info["max_salary"] = self.max_salary.get()
            self.salary_info["min_salary"] = self.min_salary.get()
            self.salary_info["average_salary"] = float(self.average_salary.get())
            self.salary_info["sum_salaries"] = int(self.sum_salaries.get())
            self.salary_info["count_sellers"] = count_sellers
            self.salary_info["salesman_list"] = self.salesman_list
        # Return
        return self.salary_info or None


async def main() -> None:
    """Example."""
    # Activate database
    Scruby.run()

    # Get collection `Salesman`
    salesman_coll = Scruby(Salesman)

    # Create sellers
    for num in range(1, 10):
        salesman = Salesman(
            username=f"salesman_{num}",
            first_name="John",
            last_name="Smith",
            birthday=datetime(1970, 1, num, tzinfo=ZoneInfo("UTC")),
            email=f"John_Smith_{num}@gmail.com",
            phone=f"+44798612345{num}",
        )
        await salesman_coll.add_doc(salesman)

    # Get salary information for sellers named John
    result: dict[str, Any] | None = salesman_coll.run_custom_task(
        custom_task=SalaryInfo(),
        filter_fn=lambda doc: doc.first_name == "John",
    )
    # Print to console
    if result is not None
        pp(result)
    else:
        print("No sellers named John")

    # Full database deletion.
    # Hint: The main purpose is tests.
    Scruby.napalm()


if __name__ == "__main__":
    anyio.run(main)

Run a SalaryInfoAsJson custom task

main.py
"""Run a salary_info_as_json custom task.

This method running a task created on the basis of a quantum loop.
Effectiveness running task depends on the number of processor threads.
"""

import anyio
from datetime import datetime
from typing import Annotated, Any
from zoneinfo import ZoneInfo

import orjson
from pydantic import EmailStr, Field
from pydantic_extra_types.phone_numbers import PhoneNumber, PhoneNumberValidator

from scruby import CustomTask, Scruby, ScrubyModel
from scruby.aggregation import Average, Max, Min, Sum
from pprint import pprint as pp


class Salesman(ScrubyModel):
    """Salesman model."""

    username: str
    first_name: str
    last_name: str
    birthday: datetime
    email: EmailStr
    phone: Annotated[PhoneNumber, PhoneNumberValidator(number_format="E164"), Field(strict=False)]
    salary: int
    # key is always at bottom
    key: Annotated[
        str,
        Field(
            frozen=True,
            default_factory=lambda data: data["phone"],
        ),
    ]


class SalaryInfoAsJson(CustomTask):
    """Custom task.

    Get information about sales salaries in json format.

    The result should be the fields:
    `max_salary`, `min_salary`, `average_salary`, `sum_salaries`, `count_sellers`, and `salesman_list`.
    """

    def __init__(self, **kwargs) -> None:
        """Initializing the task."""
        self.stop_signal = False
        self.max_salary = Max()
        self.min_salary = Min()
        self.average_salary = Average()
        self.sum_salaries = Sum()
        self.salesman_list: list[Any] = []
        self.salary_info: dict[str, Any] = {}

    def accept(self, doc: Any) -> None:
        """Operation with a document."""
        self.max_salary.set(doc.salary)
        self.min_salary.set(doc.salary)
        self.average_salary.set(doc.salary)
        self.sum_salaries.set(doc.salary)
        self.salesman_list.append(doc)

    def result(self) -> Any | None:
        """Return result."""
        # Add data to result
        count_sellers = len(self.salesman_list)
        if count_sellers > 0:
            self.salary_info["max_salary"] = self.max_salary.get()
            self.salary_info["min_salary"] = self.min_salary.get()
            self.salary_info["average_salary"] = float(self.average_salary.get())
            self.salary_info["sum_salaries"] = int(self.sum_salaries.get())
            self.salary_info["count_sellers"] = count_sellers
            self.salary_info["salesman_list"] = [doc.model_dump() for doc in self.salesman_list]
        # Convert to JSON-string
        result_json: str = orjson.dumps(self.salary_info).decode("utf-8")
        # Return
        return result_json if count_sellers > 0 else None


async def main() -> None:
    """Example."""
    # Activate database
    Scruby.run()

    # Get collection `Salesman`
    salesman_coll = Scruby(Salesman)

    # Create sellers
    for num in range(1, 10):
        salesman = Salesman(
            username=f"salesman_{num}",
            first_name="John",
            last_name="Smith",
            birthday=datetime(1970, 1, num, tzinfo=ZoneInfo("UTC")),
            email=f"John_Smith_{num}@gmail.com",
            phone=f"+44798612345{num}",
        )
        await salesman_coll.add_doc(salesman)

    # Get salary information for sellers named John
    result_json: str | None = await salesman_coll.run_custom_task(
        custom_task=SalaryInfoAsJson(),
        filter_fn=lambda doc: doc.first_name == "John",
    )
    # Print to console
    if result is not None
        result = orjson.loads(result_json)
        pp(result)
    else:
        print("No sellers named John")

    # Full database deletion.
    # Hint: The main purpose is tests.
    Scruby.napalm()


if __name__ == "__main__":
    anyio.run(main)