diff --git a/docs/analytics.md b/docs/analytics.md new file mode 100644 index 0000000..50cbdfc --- /dev/null +++ b/docs/analytics.md @@ -0,0 +1,113 @@ +# Analytics + +Admin-only commercial reporting, under `/v1/admin/analytics`. + +Before this existed, the admin dashboard derived every money figure client-side +from the most recent 100 orders. That is fine for a demo and wrong for a shop: +the numbers stop being true the moment the hundred-and-first order is placed. +These endpoints compute the same figures in SQL over the whole order history. + +## Endpoints + +| Endpoint | Answers | +|---|---| +| `GET /summary` | What did we take, and how does that compare to last period? | +| `GET /timeseries` | How did that move day by day? | +| `GET /products` | What sold, what came back, what never sold at all? | +| `GET /funnel` | Where do orders stop? | + +All four take optional `from` and `to` calendar dates, inclusive, defaulting to +the last 30 days. + +## Definitions + +These are choices, not facts, so they are written down rather than left implicit +in a query: + +| Term | Definition | +|---|---| +| Gross revenue | Orders in `paid`, `ready_to_ship` or `shipped` | +| Refunded | Orders in `refunded` — **whole-order only** | +| Net revenue | Gross less refunded | +| Average order value | Gross divided by revenue-producing orders | +| Units | Line quantities on revenue-producing orders | + +Never counted: `draft`, `pending_payment`, `cancelled`, and anything with +`deleted_at` set. + +## Four things that are easy to get wrong + +**Money is grouped by currency.** `orders.currency` permits more than one, and a +cross-currency total is not slightly wrong — it is meaningless. Every money +figure is therefore a list keyed by currency, which collapses to a single entry +for the usual single-currency shop. + +**Order money and line money are queried separately.** Joining `orders` to +`order_items` and summing `total_amount` multiplies each order's value by its +line count. The join is only ever used for quantities and line revenue; order +totals come from their own statement and the two are merged in Python. + +**Days are cut in the shop's timezone.** `SHOP_TIMEZONE` (default +`Europe/Berlin`) feeds `timezone(tz, created_at)`, Postgres' `AT TIME ZONE`. +Bucketing on raw UTC moves evening orders into the following day for any shop +east of Greenwich, which is the kind of error that looks like a data problem for +weeks. + +**Checkout is read from `payments`, not from `orders.status`.** Status records +only where an order is *now*. A cancelled order is indistinguishable from one +that never reached checkout, and a shipped-then-refunded order no longer says it +shipped. A payment row is written when checkout starts and survives whatever +happens next. + +## Known limits + +Stated here because a reader will otherwise infer something stronger: + +- **The funnel is an order funnel, not a visitor funnel.** It begins at order + creation and cannot see shoppers who browsed and never started one. Visitor + conversion needs session data the API does not collect. That is S2. +- **Partial refunds are not modelled.** `orders.status = refunded` is + all-or-nothing, so refund figures will not reconcile against a partially + refunded Stripe charge. +- **Return rate is an upper bound per SKU.** Returns are recorded per order, not + per line, so a return on a two-line order counts against both SKUs. +- **Product revenue need not equal order revenue.** Product figures sum line + values; an order total may carry shipping or adjustments belonging to no line. + +## Performance + +Aggregates are computed live rather than from a rollup table — always current, +no staleness, no extra moving parts. Three indexes support it: + +```sql +ix_orders_created_at (created_at) +ix_orders_status_created_at (status, created_at) +ix_order_items_sku (sku) +``` + +The schema is created with `Base.metadata.create_all`, which creates missing +*tables* but does not alter existing ones. A database that predates this change +therefore needs the indexes applied by hand: + +```sql +CREATE INDEX IF NOT EXISTS ix_orders_created_at ON orders (created_at); +CREATE INDEX IF NOT EXISTS ix_orders_status_created_at ON orders (status, created_at); +CREATE INDEX IF NOT EXISTS ix_order_items_sku ON order_items (sku); +``` + +Without them the endpoints still return correct figures, just with a sequential +scan. A window longer than five years is refused outright, since live +aggregation has no rollup behind it to make an unbounded range cheap. + +When order volume outgrows live aggregation, the next step is a nightly rollup +table written by the ARQ worker — deliberately not built yet, because it costs a +job, a table and staleness in exchange for speed nobody needs at this size. + +## Testing + +- `tests/test_analytics_unit.py` — period arithmetic, timezone conversion, + gap filling, undefined-vs-infinite percentage change. No database. +- `tests/test_analytics_integration.py` — seeds a deterministic dataset inside + March 2025, a window nothing else touches, and asserts exact figures. Covers + currency isolation, soft-delete exclusion, the multi-line fan-out trap, the + 23:30-UTC timezone boundary, and the payments-based funnel. diff --git a/src/app/main.py b/src/app/main.py index a88ef2e..52eef20 100644 --- a/src/app/main.py +++ b/src/app/main.py @@ -8,6 +8,7 @@ from app.chore import lifespan from app.services.admin import admin_api_router +from app.services.analytics import analytics_api_router from app.services.crud_item_store import router as item_store_router from app.services.customers import customers_api_router from app.services.fulfillment import fulfillment_api_router @@ -152,6 +153,9 @@ async def generic_exception_handler(request: Request, exc: Exception) -> JSONRes # Include admin service router (Phase 2) app.include_router(admin_api_router, prefix="/v1") +# Include analytics service router (S1) +app.include_router(analytics_api_router, prefix="/v1") + # Include fulfillment service router (Phase 3) app.include_router(fulfillment_api_router, prefix="/v1") diff --git a/src/app/services/analytics/__init__.py b/src/app/services/analytics/__init__.py new file mode 100644 index 0000000..cc6da6c --- /dev/null +++ b/src/app/services/analytics/__init__.py @@ -0,0 +1,27 @@ +""" +Analytics Service + +Admin-only commercial reporting, computed in SQL over the whole order history. + +Endpoints: + GET /admin/analytics/summary — headline figures vs the previous period + GET /admin/analytics/timeseries — the same figures bucketed over time + GET /admin/analytics/products — per-SKU performance and dead stock + GET /admin/analytics/funnel — where orders stop + +Usage: + from app.services.analytics import analytics_api_router + app.include_router(analytics_api_router, prefix="/v1") +""" + +from fastapi import APIRouter + +from .routers import analytics_router + +analytics_api_router = APIRouter( + prefix="/admin/analytics", + tags=["Analytics"], +) +analytics_api_router.include_router(analytics_router) + +__all__ = ["analytics_api_router"] diff --git a/src/app/services/analytics/dependencies.py b/src/app/services/analytics/dependencies.py new file mode 100644 index 0000000..8d434ee --- /dev/null +++ b/src/app/services/analytics/dependencies.py @@ -0,0 +1,10 @@ +""" +Analytics Dependencies + +Analytics is admin-only. Re-exports the shared admin guard so the policy has one +home, matching services/admin/dependencies.py. +""" + +from app.authorize import require_admin + +__all__ = ["require_admin"] diff --git a/src/app/services/analytics/functions/__init__.py b/src/app/services/analytics/functions/__init__.py new file mode 100644 index 0000000..b6a753e --- /dev/null +++ b/src/app/services/analytics/functions/__init__.py @@ -0,0 +1,19 @@ +"""Analytics helper functions.""" + +from .periods import ( + MAX_PERIOD_DAYS, + Interval, + Period, + build_period, + percent_change, + resolve_timezone, +) + +__all__ = [ + "MAX_PERIOD_DAYS", + "Interval", + "Period", + "build_period", + "percent_change", + "resolve_timezone", +] diff --git a/src/app/services/analytics/functions/periods.py b/src/app/services/analytics/functions/periods.py new file mode 100644 index 0000000..e996037 --- /dev/null +++ b/src/app/services/analytics/functions/periods.py @@ -0,0 +1,149 @@ +""" +Reporting Period Helpers + +Turns the ``from``/``to`` query parameters into a validated window, and derives +the immediately preceding window of equal length so every figure can be shown +against a comparable baseline. + +Buckets are cut in the shop's timezone rather than UTC. A German shop closing +at 23:00 local would otherwise see that evening's orders land on the following +day for half the year, which makes daily revenue look wrong in a way that is +hard to spot and easy to disbelieve. +""" + +from __future__ import annotations + +from dataclasses import dataclass +from datetime import date, datetime, time, timedelta +from enum import Enum +from zoneinfo import ZoneInfo, ZoneInfoNotFoundError + +from app.shared.exceptions import ValidationError + +# A window longer than this is refused. The aggregates are computed live rather +# than from a rollup table, so an unbounded range is a slow query waiting to +# happen; five years is far beyond any dashboard use and still a clear ceiling. +MAX_PERIOD_DAYS = 366 * 5 + + +class Interval(str, Enum): + """Bucket width for time series.""" + + DAY = "day" + WEEK = "week" + MONTH = "month" + + +@dataclass(frozen=True) +class Period: + """ + A half-open reporting window ``[start, end)`` in UTC. + + Half-open matters: an order placed at 23:59:59.999 on the last day belongs + to the period, and one placed at 00:00:00 the next day does not. A closed + range on dates either drops the final day or double-counts a boundary + order when two periods are compared. + """ + + start: datetime + end: datetime + timezone: str + + @property + def days(self) -> int: + return max((self.end - self.start).days, 1) + + def previous(self) -> "Period": + """The window of equal length ending where this one starts.""" + length = self.end - self.start + return Period( + start=self.start - length, + end=self.start, + timezone=self.timezone, + ) + + +def resolve_timezone(name: str) -> ZoneInfo: + """ + Look up a timezone, failing with a useful message rather than a 500. + + A misconfigured SHOP_TIMEZONE should say so plainly — the alternative is + every analytics endpoint returning an opaque internal error. + """ + try: + return ZoneInfo(name) + except (ZoneInfoNotFoundError, ValueError) as exc: + raise ValidationError( + message=f"Unknown shop timezone {name!r}", + context={"timezone": name}, + original_exception=exc, + ) from exc + + +def build_period( + date_from: date | None, + date_to: date | None, + timezone_name: str, + default_days: int = 30, +) -> Period: + """ + Build a reporting window from inclusive calendar dates. + + Both bounds are interpreted in the shop's timezone and converted to UTC, + because that is what the database stores. ``date_to`` is inclusive to the + reader — asking for the 1st to the 31st should include the 31st — so the + exclusive end is midnight at the start of the following day. + + Args: + date_from: First day to include. Defaults to ``default_days`` before + ``date_to``. + date_to: Last day to include. Defaults to today in the shop timezone. + timezone_name: IANA timezone the operator runs the shop in. + default_days: Window length used when ``date_from`` is omitted. + + Returns: + The resolved window, in UTC. + + Raises: + ValidationError: If the timezone is unknown, the range is inverted, or + the range exceeds MAX_PERIOD_DAYS. + """ + tz = resolve_timezone(timezone_name) + + last_day = date_to or datetime.now(tz).date() + first_day = date_from or (last_day - timedelta(days=default_days - 1)) + + if first_day > last_day: + raise ValidationError( + message="The start of the period is after its end", + context={"from": first_day.isoformat(), "to": last_day.isoformat()}, + ) + + span = (last_day - first_day).days + 1 + if span > MAX_PERIOD_DAYS: + raise ValidationError( + message=( + f"Requested period spans {span} days, " + f"more than the {MAX_PERIOD_DAYS} day maximum" + ), + context={"days": span, "maximum": MAX_PERIOD_DAYS}, + ) + + start = datetime.combine(first_day, time.min, tzinfo=tz) + # Exclusive end: midnight opening the day after the last requested day. + end = datetime.combine(last_day + timedelta(days=1), time.min, tzinfo=tz) + + return Period(start=start, end=end, timezone=timezone_name) + + +def percent_change(current: int | float, previous: int | float) -> float | None: + """ + Percentage change from ``previous`` to ``current``. + + Returns None when there is no baseline to compare against. Growth from zero + is not "infinite percent" or "100%" — it is undefined, and reporting a + number there invites a reader to trust something meaningless. + """ + if previous == 0: + return None + return round(((current - previous) / previous) * 100, 1) diff --git a/src/app/services/analytics/models/__init__.py b/src/app/services/analytics/models/__init__.py new file mode 100644 index 0000000..e65cdb1 --- /dev/null +++ b/src/app/services/analytics/models/__init__.py @@ -0,0 +1,33 @@ +"""Analytics schemas.""" + +from .analytics_models import ( + AnalyticsFunnelResponse, + AnalyticsProductsResponse, + AnalyticsSummaryResponse, + AnalyticsTimeseriesResponse, + Change, + CurrencySeries, + CurrencyTotals, + CurrencyTotalsPrevious, + FunnelStep, + NeverSoldItem, + PeriodInfo, + ProductPerformance, + SeriesPoint, +) + +__all__ = [ + "AnalyticsFunnelResponse", + "AnalyticsProductsResponse", + "AnalyticsSummaryResponse", + "AnalyticsTimeseriesResponse", + "Change", + "CurrencySeries", + "CurrencyTotals", + "CurrencyTotalsPrevious", + "FunnelStep", + "NeverSoldItem", + "PeriodInfo", + "ProductPerformance", + "SeriesPoint", +] diff --git a/src/app/services/analytics/models/analytics_models.py b/src/app/services/analytics/models/analytics_models.py new file mode 100644 index 0000000..9f3cad8 --- /dev/null +++ b/src/app/services/analytics/models/analytics_models.py @@ -0,0 +1,244 @@ +""" +Analytics Pydantic Schemas + +Response shapes for the admin analytics endpoints. + +Two conventions run through all of them: + +**Money is minor units.** Every amount is an integer in the currency's smallest +unit, matching ``orders.total_amount``. No floats touch money anywhere in this +service. + +**Money is grouped by currency.** ``orders.currency`` permits more than one, and +a total summed across currencies is not wrong by a little — it is meaningless. +Every money-bearing response is therefore a list keyed by currency rather than a +single figure, which collapses to one entry for the common single-currency shop. +""" + +from __future__ import annotations + +from datetime import date, datetime + +from pydantic import BaseModel, ConfigDict, Field + +from app.shared.responses import BaseResponse + + +# ============================================================================ +# Shared +# ============================================================================ + + +class PeriodInfo(BaseModel): + """The window a figure was computed over.""" + + model_config = ConfigDict(from_attributes=True) + + start: datetime = Field(description="Start of the window, inclusive (UTC)") + end: datetime = Field(description="End of the window, exclusive (UTC)") + timezone: str = Field(description="Timezone the window's days were cut in") + days: int = Field(description="Length of the window in days") + + +class Change(BaseModel): + """ + Movement against the previous equivalent period. + + Every field is nullable: percentage change from a baseline of zero is + undefined, and reporting a number there would invite trust in nothing. + """ + + net_revenue_pct: float | None = Field(default=None) + gross_revenue_pct: float | None = Field(default=None) + orders_pct: float | None = Field(default=None) + units_pct: float | None = Field(default=None) + average_order_value_pct: float | None = Field(default=None) + + +# ============================================================================ +# Summary +# ============================================================================ + + +class CurrencyTotals(BaseModel): + """Headline figures for one currency.""" + + model_config = ConfigDict(from_attributes=True) + + currency: str = Field(description="ISO 4217 code") + + gross_revenue: int = Field( + description="Money taken: orders paid, ready to ship or shipped" + ) + refunded_revenue: int = Field( + description="Value of orders refunded in full. Partial refunds are not modelled" + ) + net_revenue: int = Field(description="Gross revenue less refunds") + + orders: int = Field(description="Orders that produced revenue") + units: int = Field(description="Units across those orders") + average_order_value: int = Field( + description="Gross revenue divided by revenue-producing orders" + ) + + previous: "CurrencyTotalsPrevious | None" = Field( + default=None, description="The same figures for the preceding window" + ) + change: Change | None = Field( + default=None, description="Movement against the preceding window" + ) + + +class CurrencyTotalsPrevious(BaseModel): + """Prior-period figures, without their own comparison.""" + + gross_revenue: int + refunded_revenue: int + net_revenue: int + orders: int + units: int + average_order_value: int + + +class AnalyticsSummaryResponse(BaseResponse): + """Headline commercial figures, one entry per currency.""" + + period: PeriodInfo + previous_period: PeriodInfo + currencies: list[CurrencyTotals] = Field(default_factory=list) + + +# ============================================================================ +# Time series +# ============================================================================ + + +class SeriesPoint(BaseModel): + """One bucket of a time series.""" + + bucket: date = Field(description="First day of the bucket, in shop timezone") + gross_revenue: int + refunded_revenue: int + net_revenue: int + orders: int + units: int + + +class CurrencySeries(BaseModel): + """A complete series for one currency, gap-filled.""" + + currency: str + points: list[SeriesPoint] = Field(default_factory=list) + + +class AnalyticsTimeseriesResponse(BaseResponse): + """ + Revenue and volume over time. + + Buckets with no orders are present with zeroes rather than absent, so a + chart shows a real trough instead of silently joining across the gap. + """ + + period: PeriodInfo + interval: str + series: list[CurrencySeries] = Field(default_factory=list) + + +# ============================================================================ +# Products +# ============================================================================ + + +class ProductPerformance(BaseModel): + """How one SKU sold over the window.""" + + sku: str + name: str | None = Field( + default=None, description="Catalogue name, absent if the item was deleted" + ) + currency: str + + units_sold: int + gross_revenue: int + orders: int = Field(description="Revenue-producing orders containing this SKU") + + orders_with_return: int = Field( + description=( + "Orders containing this SKU where a return was raised. Returns are " + "recorded per order, not per line, so this attributes a return to " + "every SKU on the order — an upper bound, not a per-item rate" + ) + ) + return_rate: float | None = Field( + default=None, + description="orders_with_return divided by orders, or null when unsold", + ) + + +class NeverSoldItem(BaseModel): + """An active catalogue item with no sales in the window.""" + + sku: str + name: str + status: str + on_hand: int | None = Field( + default=None, description="Units in stock, absent if untracked" + ) + + +class AnalyticsProductsResponse(BaseResponse): + """Per-SKU performance, plus catalogue items that did not sell.""" + + period: PeriodInfo + sort: str + products: list[ProductPerformance] = Field(default_factory=list) + never_sold: list[NeverSoldItem] = Field(default_factory=list) + + +# ============================================================================ +# Funnel +# ============================================================================ + + +class FunnelStep(BaseModel): + """One stage of the order lifecycle.""" + + step: str + label: str + orders: int + conversion_from_start: float | None = Field( + default=None, description="Share of orders created that reached this step" + ) + drop_off_from_previous: int | None = Field( + default=None, description="Orders lost between the previous step and this one" + ) + + +class AnalyticsFunnelResponse(BaseResponse): + """ + Where orders stop. + + This is an **order** funnel, not a visitor funnel. It begins at order + creation, so it cannot see shoppers who browsed and never started one. + Visitor conversion needs session data the API does not collect. + + Checkout is read from the payments table rather than from order status: + a payment row is written when checkout starts, and unlike ``orders.status`` + — which only records where an order is now — it still exists after the + order moves on or is cancelled. + """ + + period: PeriodInfo + steps: list[FunnelStep] = Field(default_factory=list) + + never_checked_out: int = Field( + description="Orders created that never started checkout" + ) + payment_failed: int = Field(description="Checkouts whose payment failed") + payment_unresolved: int = Field( + description="Checkouts still pending: abandoned, or awaiting a webhook" + ) + cancelled: int = Field(description="Orders cancelled in the window") + + +CurrencyTotals.model_rebuild() diff --git a/src/app/services/analytics/responses/__init__.py b/src/app/services/analytics/responses/__init__.py new file mode 100644 index 0000000..0074d4e --- /dev/null +++ b/src/app/services/analytics/responses/__init__.py @@ -0,0 +1,15 @@ +"""Analytics OpenAPI documentation.""" + +from .analytics_docs import ( + FUNNEL_RESPONSES, + PRODUCTS_RESPONSES, + SUMMARY_RESPONSES, + TIMESERIES_RESPONSES, +) + +__all__ = [ + "FUNNEL_RESPONSES", + "PRODUCTS_RESPONSES", + "SUMMARY_RESPONSES", + "TIMESERIES_RESPONSES", +] diff --git a/src/app/services/analytics/responses/analytics_docs.py b/src/app/services/analytics/responses/analytics_docs.py new file mode 100644 index 0000000..3defc93 --- /dev/null +++ b/src/app/services/analytics/responses/analytics_docs.py @@ -0,0 +1,50 @@ +""" +OpenAPI Documentation for Analytics Endpoints + +Error examples are built from real ErrorResponse instances so they cannot drift +from what the API actually returns — the same approach used in +services/admin/responses/admin_docs.py. +""" + +from app.shared.responses import ErrorResponse + +_FORBIDDEN = ErrorResponse( + status_code=403, + error_code="access_denied", + error_category="authorization", + message="Admin access required.", + details={"resource": "analytics", "action": "read"}, +).model_dump(mode="json", exclude_none=True) + +_BAD_PERIOD = ErrorResponse( + status_code=422, + error_code="invalid_input", + error_category="validation", + message="The start of the period is after its end", + details={"from": "2026-08-31", "to": "2026-08-01"}, +).model_dump(mode="json", exclude_none=True) + + +_COMMON: dict[int | str, dict] = { + 403: { + "description": "Caller is not an administrator, or the token came from a " + "client that may not reach admin endpoints.", + "content": {"application/json": {"example": _FORBIDDEN}}, + }, + 422: { + "description": "The requested period is invalid or too long.", + "content": {"application/json": {"example": _BAD_PERIOD}}, + }, +} + +SUMMARY_RESPONSES = dict(_COMMON) +TIMESERIES_RESPONSES = dict(_COMMON) +PRODUCTS_RESPONSES = dict(_COMMON) +FUNNEL_RESPONSES = dict(_COMMON) + +__all__ = [ + "FUNNEL_RESPONSES", + "PRODUCTS_RESPONSES", + "SUMMARY_RESPONSES", + "TIMESERIES_RESPONSES", +] diff --git a/src/app/services/analytics/routers/__init__.py b/src/app/services/analytics/routers/__init__.py new file mode 100644 index 0000000..08b8aa8 --- /dev/null +++ b/src/app/services/analytics/routers/__init__.py @@ -0,0 +1,5 @@ +"""Analytics routers.""" + +from .analytics_router import router as analytics_router + +__all__ = ["analytics_router"] diff --git a/src/app/services/analytics/routers/analytics_router.py b/src/app/services/analytics/routers/analytics_router.py new file mode 100644 index 0000000..42051c2 --- /dev/null +++ b/src/app/services/analytics/routers/analytics_router.py @@ -0,0 +1,370 @@ +""" +Analytics Router + +Admin-only commercial reporting: + + GET /admin/analytics/summary — headline figures vs the previous period + GET /admin/analytics/timeseries — the same figures bucketed over time + GET /admin/analytics/products — per-SKU performance and dead stock + GET /admin/analytics/funnel — where orders stop + +Every endpoint takes the same optional ``from``/``to`` calendar dates, +interpreted in the shop's timezone and defaulting to the last 30 days. +""" + +from __future__ import annotations + +from datetime import date + +from fastapi import APIRouter, Depends, Query +from sqlalchemy.ext.asyncio import AsyncSession + +from app.shared.config import get_settings +from app.shared.database.session import get_session_dependency +from app.shared.logger import get_logger + +from ..dependencies import require_admin +from ..functions import Interval, Period, build_period, percent_change +from ..models import ( + AnalyticsFunnelResponse, + AnalyticsProductsResponse, + AnalyticsSummaryResponse, + AnalyticsTimeseriesResponse, + Change, + CurrencySeries, + CurrencyTotals, + CurrencyTotalsPrevious, + FunnelStep, + NeverSoldItem, + PeriodInfo, + ProductPerformance, + SeriesPoint, +) +from ..responses import ( + FUNNEL_RESPONSES, + PRODUCTS_RESPONSES, + SUMMARY_RESPONSES, + TIMESERIES_RESPONSES, +) +from ..services import AnalyticsRepository, fill_series_gaps + +logger = get_logger(__name__) + +router = APIRouter() + +# Cap on rows returned by the products endpoint. A catalogue can be large and +# this is a dashboard table, not an export. +_MAX_PRODUCT_ROWS = 200 + + +def _period_info(period: Period) -> PeriodInfo: + return PeriodInfo( + start=period.start, + end=period.end, + timezone=period.timezone, + days=period.days, + ) + + +def _average(total: int, count: int) -> int: + return round(total / count) if count else 0 + + +async def _resolve_period( + date_from: date | None, + date_to: date | None, +) -> Period: + settings = get_settings() + return build_period( + date_from=date_from, + date_to=date_to, + timezone_name=settings.shop_timezone, + ) + + +# --------------------------------------------------------------------------- +# GET /admin/analytics/summary +# --------------------------------------------------------------------------- + + +@router.get( + "/summary", + response_model=AnalyticsSummaryResponse, + summary="Commercial summary (admin)", + description=( + "Revenue, refunds, order count, units and average order value for the " + "requested window, each compared against the immediately preceding " + "window of equal length.\n\n" + "Figures are grouped by currency because `orders.currency` permits more " + "than one and a cross-currency total would be meaningless. Revenue is " + "gross (paid, ready to ship, shipped) less refunds; refunds are " + "whole-order only, as partial refunds are not modelled." + ), + responses=SUMMARY_RESPONSES, + dependencies=[Depends(require_admin)], +) +async def get_summary( + date_from: date | None = Query( + None, alias="from", description="First day, inclusive" + ), + date_to: date | None = Query(None, alias="to", description="Last day, inclusive"), + session: AsyncSession = Depends(get_session_dependency), +) -> AnalyticsSummaryResponse: + period = await _resolve_period(date_from, date_to) + previous = period.previous() + repository = AnalyticsRepository(session) + + current_totals = await repository.order_totals_by_currency(period) + current_units = await repository.units_by_currency(period) + previous_totals = await repository.order_totals_by_currency(previous) + previous_units = await repository.units_by_currency(previous) + + currencies: list[CurrencyTotals] = [] + # A currency present only in the prior window still deserves a row — its + # disappearance is exactly the kind of thing a dashboard should show. + for currency in sorted(set(current_totals) | set(previous_totals)): + now = current_totals.get( + currency, {"gross_revenue": 0, "refunded_revenue": 0, "orders": 0} + ) + then = previous_totals.get( + currency, {"gross_revenue": 0, "refunded_revenue": 0, "orders": 0} + ) + + now_units = current_units.get(currency, 0) + then_units = previous_units.get(currency, 0) + + now_net = now["gross_revenue"] - now["refunded_revenue"] + then_net = then["gross_revenue"] - then["refunded_revenue"] + now_aov = _average(now["gross_revenue"], now["orders"]) + then_aov = _average(then["gross_revenue"], then["orders"]) + + currencies.append( + CurrencyTotals( + currency=currency, + gross_revenue=now["gross_revenue"], + refunded_revenue=now["refunded_revenue"], + net_revenue=now_net, + orders=now["orders"], + units=now_units, + average_order_value=now_aov, + previous=CurrencyTotalsPrevious( + gross_revenue=then["gross_revenue"], + refunded_revenue=then["refunded_revenue"], + net_revenue=then_net, + orders=then["orders"], + units=then_units, + average_order_value=then_aov, + ), + change=Change( + net_revenue_pct=percent_change(now_net, then_net), + gross_revenue_pct=percent_change( + now["gross_revenue"], then["gross_revenue"] + ), + orders_pct=percent_change(now["orders"], then["orders"]), + units_pct=percent_change(now_units, then_units), + average_order_value_pct=percent_change(now_aov, then_aov), + ), + ) + ) + + return AnalyticsSummaryResponse( + success=True, + message="Analytics summary retrieved successfully", + period=_period_info(period), + previous_period=_period_info(previous), + currencies=currencies, + ) + + +# --------------------------------------------------------------------------- +# GET /admin/analytics/timeseries +# --------------------------------------------------------------------------- + + +@router.get( + "/timeseries", + response_model=AnalyticsTimeseriesResponse, + summary="Revenue and volume over time (admin)", + description=( + "The summary figures bucketed by day, week or month, in the shop's " + "timezone.\n\n" + "Buckets with no orders are returned with zeroes rather than omitted, so " + "a chart shows a real trough instead of joining across the gap." + ), + responses=TIMESERIES_RESPONSES, + dependencies=[Depends(require_admin)], +) +async def get_timeseries( + date_from: date | None = Query( + None, alias="from", description="First day, inclusive" + ), + date_to: date | None = Query(None, alias="to", description="Last day, inclusive"), + interval: Interval = Query(Interval.DAY, description="Bucket width"), + session: AsyncSession = Depends(get_session_dependency), +) -> AnalyticsTimeseriesResponse: + period = await _resolve_period(date_from, date_to) + repository = AnalyticsRepository(session) + + raw = await repository.series_by_currency(period, interval) + + series = [ + CurrencySeries( + currency=currency, + points=[ + SeriesPoint(**point) + for point in fill_series_gaps(buckets, period, interval) + ], + ) + for currency, buckets in sorted(raw.items()) + ] + + return AnalyticsTimeseriesResponse( + success=True, + message="Analytics time series retrieved successfully", + period=_period_info(period), + interval=interval.value, + series=series, + ) + + +# --------------------------------------------------------------------------- +# GET /admin/analytics/products +# --------------------------------------------------------------------------- + + +@router.get( + "/products", + response_model=AnalyticsProductsResponse, + summary="Product performance (admin)", + description=( + "Units, revenue, order count and return rate per SKU, plus active " + "catalogue items that sold nothing in the window.\n\n" + "Revenue here sums line values (quantity x unit price) and need not " + "equal the order totals in the summary, which may carry shipping or " + "adjustments belonging to no line. Return rate attributes an order's " + "return to every SKU on that order, so it is an upper bound per SKU." + ), + responses=PRODUCTS_RESPONSES, + dependencies=[Depends(require_admin)], +) +async def get_products( + date_from: date | None = Query( + None, alias="from", description="First day, inclusive" + ), + date_to: date | None = Query(None, alias="to", description="Last day, inclusive"), + sort: str = Query("revenue", pattern="^(revenue|units|orders|return_rate)$"), + limit: int = Query(20, ge=1, le=_MAX_PRODUCT_ROWS), + session: AsyncSession = Depends(get_session_dependency), +) -> AnalyticsProductsResponse: + period = await _resolve_period(date_from, date_to) + repository = AnalyticsRepository(session) + + rows = await repository.product_performance(period) + returns = await repository.returns_by_sku(period) + names = await repository.item_names([row["sku"] for row in rows]) + + products: list[ProductPerformance] = [] + for row in rows: + orders = row["orders"] + with_return = returns.get(row["sku"], 0) + products.append( + ProductPerformance( + sku=row["sku"], + name=names.get(row["sku"]), + currency=row["currency"], + units_sold=row["units_sold"], + gross_revenue=row["gross_revenue"], + orders=orders, + orders_with_return=with_return, + return_rate=round(with_return / orders, 4) if orders else None, + ) + ) + + sort_keys = { + "revenue": lambda p: p.gross_revenue, + "units": lambda p: p.units_sold, + "orders": lambda p: p.orders, + "return_rate": lambda p: p.return_rate or 0.0, + } + products.sort(key=sort_keys[sort], reverse=True) + + never_sold = await repository.never_sold([row["sku"] for row in rows]) + + return AnalyticsProductsResponse( + success=True, + message="Product performance retrieved successfully", + period=_period_info(period), + sort=sort, + products=products[:limit], + never_sold=[NeverSoldItem(**item) for item in never_sold], + ) + + +# --------------------------------------------------------------------------- +# GET /admin/analytics/funnel +# --------------------------------------------------------------------------- + + +@router.get( + "/funnel", + response_model=AnalyticsFunnelResponse, + summary="Order funnel (admin)", + description=( + "Where orders stop, from creation through checkout and payment to " + "handover.\n\n" + "This is an **order** funnel, not a visitor funnel: it starts at order " + "creation and cannot see shoppers who browsed without creating one. " + "Checkout is read from the payments table rather than order status, " + "because status records only where an order is now and cannot " + "distinguish a cancelled checkout from one that never happened." + ), + responses=FUNNEL_RESPONSES, + dependencies=[Depends(require_admin)], +) +async def get_funnel( + date_from: date | None = Query( + None, alias="from", description="First day, inclusive" + ), + date_to: date | None = Query(None, alias="to", description="Last day, inclusive"), + session: AsyncSession = Depends(get_session_dependency), +) -> AnalyticsFunnelResponse: + period = await _resolve_period(date_from, date_to) + repository = AnalyticsRepository(session) + + counts = await repository.funnel_counts(period) + created = counts["created"] + + definitions = [ + ("created", "Order created"), + ("checkout_started", "Checkout started"), + ("paid", "Payment confirmed"), + ("shipped", "Handed to carrier"), + ] + + steps: list[FunnelStep] = [] + previous_count: int | None = None + for key, label in definitions: + count = counts[key] + steps.append( + FunnelStep( + step=key, + label=label, + orders=count, + conversion_from_start=(round(count / created, 4) if created else None), + drop_off_from_previous=( + None if previous_count is None else max(previous_count - count, 0) + ), + ) + ) + previous_count = count + + return AnalyticsFunnelResponse( + success=True, + message="Order funnel retrieved successfully", + period=_period_info(period), + steps=steps, + never_checked_out=counts["never_checked_out"], + payment_failed=counts["payment_failed"], + payment_unresolved=counts["payment_unresolved"], + cancelled=counts["cancelled"], + ) diff --git a/src/app/services/analytics/services/__init__.py b/src/app/services/analytics/services/__init__.py new file mode 100644 index 0000000..2c00dc4 --- /dev/null +++ b/src/app/services/analytics/services/__init__.py @@ -0,0 +1,15 @@ +"""Analytics database services.""" + +from .analytics_db_service import ( + EARNED_STATUSES, + AnalyticsRepository, + fill_series_gaps, + get_analytics_repository, +) + +__all__ = [ + "EARNED_STATUSES", + "AnalyticsRepository", + "fill_series_gaps", + "get_analytics_repository", +] diff --git a/src/app/services/analytics/services/analytics_db_service.py b/src/app/services/analytics/services/analytics_db_service.py new file mode 100644 index 0000000..8160332 --- /dev/null +++ b/src/app/services/analytics/services/analytics_db_service.py @@ -0,0 +1,482 @@ +""" +Analytics Database Service + +Aggregation queries over orders, order items, payments, returns and inventory. + +Everything here is read-only and computed live rather than from a rollup table. +That keeps the figures current and adds no moving parts, at the cost of scanning +the order history per request — see the indexes noted in ``orders_db_models``. + +Three rules hold throughout: + +**Soft-deleted orders never count.** ``orders.deleted_at`` is set on records that +must survive as history but must not appear in a total. + +**Order money and line money are queried separately.** Joining orders to their +items and summing ``orders.total_amount`` multiplies each order's value by its +line count. The two are therefore computed in separate statements and merged in +Python — the join is only ever used for quantities and line revenue. + +**Days are cut in the shop's timezone**, via ``timezone(tz, created_at)``, which +is Postgres' ``AT TIME ZONE``. Bucketing on raw UTC would move evening orders +into the following day for any shop east of Greenwich. +""" + +from __future__ import annotations + +from collections import defaultdict +from datetime import date, timedelta + +from sqlalchemy import Select, and_, case, distinct, func, select +from sqlalchemy.ext.asyncio import AsyncSession + +from app.services.crud_item_store.models.item_db_models import ItemDB +from app.services.inventory.models.inventory_db_models import InventoryItemDB +from app.services.orders.models.orders_db_models import OrderDB, OrderItemDB +from app.services.payments.models.payments_db_models import PaymentDB +from app.services.returns.models.returns_db_models import ReturnDB +from app.services.shipments.models.shipments_db_models import ShipmentDB +from app.shared.logger import get_logger + +from ..functions import Interval, Period + +logger = get_logger(__name__) + +# Statuses whose order value counts as money actually taken. Mirrors the +# admin dashboard's definition so the two never disagree. +EARNED_STATUSES: tuple[str, ...] = ("paid", "ready_to_ship", "shipped") +REFUNDED_STATUS = "refunded" +CANCELLED_STATUS = "cancelled" + + +class AnalyticsRepository: + """Read-only aggregation queries for the admin analytics endpoints.""" + + def __init__(self, session: AsyncSession) -> None: + self._session = session + + # ------------------------------------------------------------------ + # Shared predicates + # ------------------------------------------------------------------ + + @staticmethod + def _in_period(period: Period): + """Live orders created inside the half-open window.""" + return and_( + OrderDB.deleted_at.is_(None), + OrderDB.created_at >= period.start, + OrderDB.created_at < period.end, + ) + + @staticmethod + def _earned() -> tuple: + return (OrderDB.status.in_(EARNED_STATUSES),) + + @staticmethod + def _gross_expr(): + return func.coalesce( + func.sum( + case( + (OrderDB.status.in_(EARNED_STATUSES), OrderDB.total_amount), + else_=0, + ) + ), + 0, + ) + + @staticmethod + def _refunded_expr(): + return func.coalesce( + func.sum( + case((OrderDB.status == REFUNDED_STATUS, OrderDB.total_amount), else_=0) + ), + 0, + ) + + @staticmethod + def _earned_orders_expr(): + return func.count(distinct(OrderDB.id)).filter( + OrderDB.status.in_(EARNED_STATUSES) + ) + + def _bucket_expr(self, interval: Interval, timezone_name: str): + """ + Truncate ``created_at`` to a bucket, in the shop's local time. + + ``timezone(zone, timestamptz)`` is the function form of AT TIME ZONE and + takes the zone as a bind parameter, so the shop's configured timezone + never reaches SQL as interpolated text. + """ + local = func.timezone(timezone_name, OrderDB.created_at) + return func.date_trunc(interval.value, local) + + # ------------------------------------------------------------------ + # Summary + # ------------------------------------------------------------------ + + async def order_totals_by_currency( + self, period: Period + ) -> dict[str, dict[str, int]]: + """ + Order-level money and counts, grouped by currency. + + Deliberately does not join order_items — see the module docstring. + """ + statement: Select = ( + select( + OrderDB.currency, + self._gross_expr().label("gross_revenue"), + self._refunded_expr().label("refunded_revenue"), + self._earned_orders_expr().label("orders"), + ) + .where(self._in_period(period)) + .group_by(OrderDB.currency) + ) + + result = await self._session.execute(statement) + return { + row.currency: { + "gross_revenue": int(row.gross_revenue or 0), + "refunded_revenue": int(row.refunded_revenue or 0), + "orders": int(row.orders or 0), + } + for row in result + } + + async def units_by_currency(self, period: Period) -> dict[str, int]: + """Units sold on revenue-producing orders, grouped by currency.""" + statement: Select = ( + select( + OrderDB.currency, + func.coalesce(func.sum(OrderItemDB.quantity), 0).label("units"), + ) + .join(OrderItemDB, OrderItemDB.order_id == OrderDB.id) + .where(self._in_period(period), *self._earned()) + .group_by(OrderDB.currency) + ) + + result = await self._session.execute(statement) + return {row.currency: int(row.units or 0) for row in result} + + # ------------------------------------------------------------------ + # Time series + # ------------------------------------------------------------------ + + async def series_by_currency( + self, period: Period, interval: Interval + ) -> dict[str, dict[date, dict[str, int]]]: + """ + Bucketed money and counts per currency. + + Returns only buckets that have orders; gap filling is the caller's job, + since it needs the full window to know which buckets are missing. + """ + bucket = self._bucket_expr(interval, period.timezone) + + statement: Select = ( + select( + OrderDB.currency, + bucket.label("bucket"), + self._gross_expr().label("gross_revenue"), + self._refunded_expr().label("refunded_revenue"), + self._earned_orders_expr().label("orders"), + ) + .where(self._in_period(period)) + .group_by(OrderDB.currency, bucket) + .order_by(bucket) + ) + + result = await self._session.execute(statement) + buckets: dict[str, dict[date, dict[str, int]]] = defaultdict(dict) + for row in result: + bucket_day = ( + row.bucket.date() if hasattr(row.bucket, "date") else row.bucket + ) + buckets[row.currency][bucket_day] = { + "gross_revenue": int(row.gross_revenue or 0), + "refunded_revenue": int(row.refunded_revenue or 0), + "orders": int(row.orders or 0), + "units": 0, + } + + for currency, units_by_bucket in ( + await self._series_units(period, interval) + ).items(): + for bucket_day, units in units_by_bucket.items(): + buckets.setdefault(currency, {}).setdefault( + bucket_day, + { + "gross_revenue": 0, + "refunded_revenue": 0, + "orders": 0, + "units": 0, + }, + )["units"] = units + + return buckets + + async def _series_units( + self, period: Period, interval: Interval + ) -> dict[str, dict[date, int]]: + """Units per bucket per currency, joined separately to avoid fan-out.""" + bucket = self._bucket_expr(interval, period.timezone) + + statement: Select = ( + select( + OrderDB.currency, + bucket.label("bucket"), + func.coalesce(func.sum(OrderItemDB.quantity), 0).label("units"), + ) + .join(OrderItemDB, OrderItemDB.order_id == OrderDB.id) + .where(self._in_period(period), *self._earned()) + .group_by(OrderDB.currency, bucket) + ) + + result = await self._session.execute(statement) + units: dict[str, dict[date, int]] = defaultdict(dict) + for row in result: + bucket_day = ( + row.bucket.date() if hasattr(row.bucket, "date") else row.bucket + ) + units[row.currency][bucket_day] = int(row.units or 0) + return units + + # ------------------------------------------------------------------ + # Products + # ------------------------------------------------------------------ + + async def product_performance(self, period: Period) -> list[dict]: + """ + Per-SKU units, line revenue and order count over the window. + + Revenue here is the sum of line values (``quantity * unit_price``), + which need not equal the order totals used elsewhere — an order total + may carry shipping or adjustments that belong to no line. + """ + statement: Select = ( + select( + OrderDB.currency, + OrderItemDB.sku, + func.coalesce(func.sum(OrderItemDB.quantity), 0).label("units_sold"), + func.coalesce( + func.sum(OrderItemDB.quantity * OrderItemDB.unit_price), 0 + ).label("gross_revenue"), + func.count(distinct(OrderDB.id)).label("orders"), + ) + .join(OrderItemDB, OrderItemDB.order_id == OrderDB.id) + .where(self._in_period(period), *self._earned()) + .group_by(OrderDB.currency, OrderItemDB.sku) + ) + + result = await self._session.execute(statement) + return [ + { + "currency": row.currency, + "sku": row.sku, + "units_sold": int(row.units_sold or 0), + "gross_revenue": int(row.gross_revenue or 0), + "orders": int(row.orders or 0), + } + for row in result + ] + + async def returns_by_sku(self, period: Period) -> dict[str, int]: + """ + Orders containing each SKU where a return was raised. + + Returns are recorded per order, so a return on a two-line order counts + against both SKUs. This is an upper bound per SKU, not a per-item rate, + and is labelled as such in the response schema. + """ + statement: Select = ( + select( + OrderItemDB.sku, + func.count(distinct(OrderDB.id)).label("orders_with_return"), + ) + .join(OrderItemDB, OrderItemDB.order_id == OrderDB.id) + .join(ReturnDB, ReturnDB.order_id == OrderDB.id) + .where(self._in_period(period)) + .group_by(OrderItemDB.sku) + ) + + result = await self._session.execute(statement) + return {row.sku: int(row.orders_with_return or 0) for row in result} + + async def item_names(self, skus: list[str]) -> dict[str, str]: + """Catalogue names for the given SKUs. Missing SKUs are simply absent.""" + if not skus: + return {} + result = await self._session.execute( + select(ItemDB.sku, ItemDB.name).where(ItemDB.sku.in_(skus)) + ) + return {row.sku: row.name for row in result} + + async def never_sold(self, sold_skus: list[str], limit: int = 50) -> list[dict]: + """ + Active catalogue items with no sales in the window. + + Stock is joined optionally: an item without an inventory record is dead + stock worth surfacing too, so a missing row must not drop it. + """ + statement: Select = ( + select( + ItemDB.sku, + ItemDB.name, + ItemDB.status, + InventoryItemDB.on_hand, + ) + .outerjoin(InventoryItemDB, InventoryItemDB.sku == ItemDB.sku) + .where(ItemDB.status == "active") + .order_by(ItemDB.name) + .limit(limit) + ) + if sold_skus: + statement = statement.where(ItemDB.sku.notin_(sold_skus)) + + result = await self._session.execute(statement) + return [ + { + "sku": row.sku, + "name": row.name, + "status": row.status, + "on_hand": int(row.on_hand) if row.on_hand is not None else None, + } + for row in result + ] + + # ------------------------------------------------------------------ + # Funnel + # ------------------------------------------------------------------ + + async def funnel_counts(self, period: Period) -> dict[str, int]: + """ + Order lifecycle counts. + + Checkout is read from ``payments`` rather than from ``orders.status``. + Status records only where an order is now, so a cancelled order is + indistinguishable from one that never reached checkout. A payment row is + written when checkout starts and survives whatever happens next. + """ + created = await self._scalar( + select(func.count(OrderDB.id)).where(self._in_period(period)) + ) + + checkout_started = await self._scalar( + select(func.count(distinct(OrderDB.id))) + .join(PaymentDB, PaymentDB.order_id == OrderDB.id) + .where(self._in_period(period)) + ) + + paid = await self._scalar( + select(func.count(distinct(OrderDB.id))) + .join(PaymentDB, PaymentDB.order_id == OrderDB.id) + .where(self._in_period(period), PaymentDB.status == "succeeded") + ) + + payment_failed = await self._scalar( + select(func.count(distinct(OrderDB.id))) + .join(PaymentDB, PaymentDB.order_id == OrderDB.id) + .where(self._in_period(period), PaymentDB.status == "failed") + ) + + # Counted from shipments rather than order status: an order that shipped + # and was later refunded no longer says "shipped", but it did ship. + shipped = await self._scalar( + select(func.count(distinct(OrderDB.id))) + .join(ShipmentDB, ShipmentDB.order_id == OrderDB.id) + .where(self._in_period(period), ShipmentDB.status == "handed_over") + ) + + cancelled = await self._scalar( + select(func.count(OrderDB.id)).where( + self._in_period(period), OrderDB.status == CANCELLED_STATUS + ) + ) + + return { + "created": created, + "checkout_started": checkout_started, + "paid": paid, + "shipped": shipped, + "payment_failed": payment_failed, + "payment_unresolved": max(checkout_started - paid - payment_failed, 0), + "never_checked_out": max(created - checkout_started, 0), + "cancelled": cancelled, + } + + async def _scalar(self, statement: Select) -> int: + result = await self._session.execute(statement) + return int(result.scalar_one_or_none() or 0) + + +def get_analytics_repository(session: AsyncSession) -> AnalyticsRepository: + """Factory used by the router's dependency wiring.""" + return AnalyticsRepository(session) + + +def fill_series_gaps( + buckets: dict[date, dict[str, int]], + period: Period, + interval: Interval, +) -> list[dict]: + """ + Return one point per bucket in the window, zero-filled where there were none. + + A chart drawn from present-only buckets joins straight across a quiet week, + which reads as steady trade rather than none. Absent days must be visible. + """ + step = { + Interval.DAY: timedelta(days=1), + Interval.WEEK: timedelta(weeks=1), + }.get(interval) + + points: list[dict] = [] + + def zero(day: date) -> dict: + return { + "bucket": day, + "gross_revenue": 0, + "refunded_revenue": 0, + "net_revenue": 0, + "orders": 0, + "units": 0, + } + + def emit(day: date) -> None: + found = buckets.get(day) + if found is None: + points.append(zero(day)) + return + points.append( + { + "bucket": day, + "gross_revenue": found["gross_revenue"], + "refunded_revenue": found["refunded_revenue"], + "net_revenue": found["gross_revenue"] - found["refunded_revenue"], + "orders": found["orders"], + "units": found["units"], + } + ) + + start_day = period.start.date() + last_day = (period.end - timedelta(seconds=1)).date() + + if interval is Interval.MONTH: + # Months are not a fixed timedelta, so walk them by calendar. + cursor = start_day.replace(day=1) + while cursor <= last_day: + emit(cursor) + cursor = (cursor.replace(day=28) + timedelta(days=4)).replace(day=1) + return points + + if interval is Interval.WEEK: + # date_trunc('week') anchors to Monday; align so keys match. + cursor = start_day - timedelta(days=start_day.weekday()) + else: + cursor = start_day + + while cursor <= last_day: + emit(cursor) + cursor = cursor + step + + return points diff --git a/src/app/services/orders/models/orders_db_models.py b/src/app/services/orders/models/orders_db_models.py index 38094ef..d6b9694 100644 --- a/src/app/services/orders/models/orders_db_models.py +++ b/src/app/services/orders/models/orders_db_models.py @@ -14,7 +14,15 @@ from uuid import UUID, uuid4 -from sqlalchemy import BigInteger, CheckConstraint, ForeignKey, Integer, String, text +from sqlalchemy import ( + BigInteger, + CheckConstraint, + ForeignKey, + Index, + Integer, + String, + text, +) from sqlalchemy.dialects.postgresql import UUID as PGUUID from sqlalchemy.orm import Mapped, mapped_column @@ -84,6 +92,12 @@ class OrderDB(Base, TimestampMixin, SoftDeleteMixin): CheckConstraint( "total_amount >= 0", name="ck_orders_total_amount_non_negative" ), + # Analytics aggregates every query by creation date, and most also + # filter by status. ix_orders_status alone does not help a date range, + # and scanning the whole order history per dashboard load does not stay + # cheap as the shop grows. + Index("ix_orders_created_at", "created_at"), + Index("ix_orders_status_created_at", "status", "created_at"), ) def __repr__(self) -> str: @@ -149,6 +163,8 @@ class OrderItemDB(Base, TimestampMixin): CheckConstraint( "unit_price >= 0", name="ck_order_items_unit_price_non_negative" ), + # Product performance groups by sku across the whole line history. + Index("ix_order_items_sku", "sku"), ) def __repr__(self) -> str: diff --git a/src/app/shared/config/settings.py b/src/app/shared/config/settings.py index 07fd883..417038a 100644 --- a/src/app/shared/config/settings.py +++ b/src/app/shared/config/settings.py @@ -271,6 +271,16 @@ class Settings(BaseSettings): description="Default label format requested from DHL: 'pdf' or 'zpl'", ) + # Analytics / reporting + shop_timezone: str = Field( + default="Europe/Berlin", + description=( + "IANA timezone the shop trades in. Analytics buckets days in this " + "zone rather than UTC, so 'today' matches the operator's day " + "instead of drifting for evening orders." + ), + ) + # ARQ worker / job queue arq_max_jobs: int = Field( default=10, diff --git a/tests/test_analytics_integration.py b/tests/test_analytics_integration.py new file mode 100644 index 0000000..44f9c2c --- /dev/null +++ b/tests/test_analytics_integration.py @@ -0,0 +1,390 @@ +""" +Integration tests for the Analytics API (S1). + +Runs against the live stack: + + docker compose -f docker-compose.dev.yml up -d + +Endpoints covered: + GET /v1/admin/analytics/summary + GET /v1/admin/analytics/timeseries + GET /v1/admin/analytics/products + GET /v1/admin/analytics/funnel + +The fixture seeds a deterministic dataset inside **March 2025**, a window no +other fixture or seed touches, and every assertion queries exactly that window. +Figures are therefore exact regardless of whatever else is in the database — +which matters, because the dev database accumulates orders from other suites. + +The dataset is built to exercise the things that are quietly easy to get wrong: + + - two currencies, which must never be summed together + - a soft-deleted order, which must vanish from every figure + - a draft and a cancelled order, which are not revenue + - a refunded order, which reduces net but not gross + - a multi-line order, which must not multiply its own value by its line count + - an order at 23:30 UTC, which belongs to the *next* day in Berlin +""" + +import os +import subprocess +import uuid + +import pytest +import requests + +from auth_helpers import admin_headers + +_BASE = os.getenv("TEST_API_URL", "http://localhost:8000") +ANALYTICS_URL = f"{_BASE}/v1/admin/analytics" + +_HEADERS = admin_headers() + +# The isolated reporting window. Chosen far from any other fixture's data. +WINDOW = {"from": "2025-03-01", "to": "2025-03-31"} + +# Marks every row this module creates, so teardown can remove exactly those. +TAG = "analytics-it" + + +def _psql(sql: str) -> str: + """Run SQL inside the Postgres container and return stdout.""" + result = subprocess.run( + [ + "docker", + "exec", + "opentaberna-db", + "psql", + "-U", + "opentaberna", + "-d", + "opentaberna", + "-t", + "-A", + "-c", + sql, + ], + check=True, + capture_output=True, + text=True, + ) + return result.stdout.strip() + + +def _get(path: str, **params) -> dict: + response = requests.get(f"{ANALYTICS_URL}/{path}", headers=_HEADERS, params=params) + response.raise_for_status() + return response.json() + + +def _currency(payload: dict, code: str) -> dict | None: + for entry in payload["currencies"]: + if entry["currency"] == code: + return entry + return None + + +@pytest.fixture(scope="module", autouse=True) +def seeded_dataset(): + """Insert the fixture dataset, yield, then remove exactly what was inserted.""" + customer = uuid.uuid4() + orders = {name: uuid.uuid4() for name in "ABCDEFGH"} + suffix = uuid.uuid4().hex[:8] + + def order_row( + key: str, + status: str, + amount: int, + currency: str, + when: str, + deleted: str = "NULL", + ) -> str: + return ( + f"('{orders[key]}', '{customer}', '{status}', {amount}, '{currency}', " + f"'{when}'::timestamptz, '{when}'::timestamptz, {deleted})" + ) + + _psql( + f""" + INSERT INTO customers (id, keycloak_user_id, email, first_name, last_name, + created_at, updated_at) + VALUES ('{customer}', '{TAG}-{suffix}', '{TAG}-{suffix}@example.test', + 'Analytics', 'Fixture', now(), now()); + + INSERT INTO orders (id, customer_id, status, total_amount, currency, + created_at, updated_at, deleted_at) + VALUES + {order_row("A", "paid", 10000, "EUR", "2025-03-05T12:00:00Z")}, + {order_row("B", "shipped", 5000, "EUR", "2025-03-14T12:00:00Z")}, + {order_row("C", "refunded", 3000, "EUR", "2025-03-12T12:00:00Z")}, + {order_row("D", "draft", 9999, "EUR", "2025-03-15T12:00:00Z")}, + {order_row("E", "paid", 7000, "USD", "2025-03-06T12:00:00Z")}, + {order_row("F", "paid", 4000, "EUR", "2025-03-20T12:00:00Z", "now()")}, + {order_row("G", "cancelled", 2000, "EUR", "2025-03-22T12:00:00Z")}, + {order_row("H", "paid", 1000, "EUR", "2025-03-09T23:30:00Z")}; + + INSERT INTO order_items (id, order_id, sku, quantity, unit_price, + created_at, updated_at) + VALUES + ('{uuid.uuid4()}', '{orders["A"]}', '{TAG}-SKU-A', 2, 2500, now(), now()), + ('{uuid.uuid4()}', '{orders["A"]}', '{TAG}-SKU-B', 1, 5000, now(), now()), + ('{uuid.uuid4()}', '{orders["B"]}', '{TAG}-SKU-A', 1, 5000, now(), now()), + ('{uuid.uuid4()}', '{orders["E"]}', '{TAG}-SKU-C', 1, 7000, now(), now()), + ('{uuid.uuid4()}', '{orders["F"]}', '{TAG}-SKU-A', 9, 9999, now(), now()), + ('{uuid.uuid4()}', '{orders["D"]}', '{TAG}-SKU-A', 9, 9999, now(), now()); + + INSERT INTO payments (id, order_id, provider, provider_reference, amount, + currency, status, created_at, updated_at) + VALUES + ('{uuid.uuid4()}', '{orders["A"]}', 'stripe', '{TAG}-{suffix}-a', 10000, + 'EUR', 'succeeded', now(), now()), + ('{uuid.uuid4()}', '{orders["B"]}', 'stripe', '{TAG}-{suffix}-b', 5000, + 'EUR', 'succeeded', now(), now()), + ('{uuid.uuid4()}', '{orders["D"]}', 'stripe', '{TAG}-{suffix}-d', 9999, + 'EUR', 'pending', now(), now()), + ('{uuid.uuid4()}', '{orders["G"]}', 'stripe', '{TAG}-{suffix}-g', 2000, + 'EUR', 'failed', now(), now()); + + INSERT INTO shipments (id, order_id, carrier, status, created_at, updated_at) + VALUES ('{uuid.uuid4()}', '{orders["B"]}', 'dhl', 'handed_over', now(), now()); + + INSERT INTO returns (id, order_id, customer_id, status, reason, + created_at, updated_at) + VALUES ('{uuid.uuid4()}', '{orders["A"]}', '{customer}', 'requested', + 'fixture', now(), now()); + """ + ) + + yield + + # customers cascades to orders, which cascades to order_items; payments, + # shipments and returns are RESTRICT so they go first. + ids = ", ".join(f"'{value}'" for value in orders.values()) + _psql( + f""" + DELETE FROM returns WHERE order_id IN ({ids}); + DELETE FROM shipments WHERE order_id IN ({ids}); + DELETE FROM payments WHERE order_id IN ({ids}); + DELETE FROM orders WHERE id IN ({ids}); + DELETE FROM customers WHERE id = '{customer}'; + """ + ) + + +# --------------------------------------------------------------------------- +# Authorization +# --------------------------------------------------------------------------- + + +@pytest.mark.parametrize("endpoint", ["summary", "timeseries", "products", "funnel"]) +def test_analytics_requires_admin(endpoint): + """Commercial figures are admin-only on every endpoint, without exception.""" + response = requests.get(f"{ANALYTICS_URL}/{endpoint}") + assert response.status_code in (401, 403) + + +# --------------------------------------------------------------------------- +# Summary +# --------------------------------------------------------------------------- + + +def test_summary_reports_revenue_per_currency_without_mixing_them(): + """ + EUR and USD are reported separately. Summing them would produce a number + that means nothing, so the schema makes it impossible. + """ + payload = _get("summary", **WINDOW) + + eur = _currency(payload, "EUR") + usd = _currency(payload, "USD") + + assert eur is not None and usd is not None + + # A(10000) + B(5000) + H(1000). C is refunded, D draft, G cancelled, + # F soft-deleted — none of them are gross revenue. + assert eur["gross_revenue"] == 16000 + assert eur["refunded_revenue"] == 3000 + assert eur["net_revenue"] == 13000 + assert eur["orders"] == 3 + + assert usd["gross_revenue"] == 7000 + assert usd["orders"] == 1 + + +def test_soft_deleted_orders_are_excluded_from_every_figure(): + """ + Order F is paid, worth 4000, and carries nine units — and is soft-deleted. + None of it may reach a total. + """ + payload = _get("summary", **WINDOW) + eur = _currency(payload, "EUR") + + assert eur["gross_revenue"] == 16000, "a soft-deleted order leaked into revenue" + assert eur["units"] == 4, "a soft-deleted order's units leaked into the count" + + +def test_multi_line_orders_are_not_counted_once_per_line(): + """ + Order A has two lines. Joining orders to items and summing total_amount + would count its 10000 twice. Order money and line money are queried apart + precisely to stop that. + """ + payload = _get("summary", **WINDOW) + eur = _currency(payload, "EUR") + + assert eur["gross_revenue"] == 16000 + # A contributes 3 units across 2 lines, B contributes 1, H none. + assert eur["units"] == 4 + + +def test_average_order_value_divides_by_revenue_producing_orders_only(): + payload = _get("summary", **WINDOW) + eur = _currency(payload, "EUR") + + assert eur["average_order_value"] == round(16000 / 3) + + +def test_previous_period_abuts_the_requested_one(): + payload = _get("summary", **WINDOW) + + assert payload["previous_period"]["end"] == payload["period"]["start"] + assert payload["previous_period"]["days"] == payload["period"]["days"] + + +def test_percentage_change_is_null_when_there_is_no_baseline(): + """February 2025 is empty, so growth into March is undefined, not infinite.""" + payload = _get("summary", **WINDOW) + eur = _currency(payload, "EUR") + + assert eur["previous"]["gross_revenue"] == 0 + assert eur["change"]["gross_revenue_pct"] is None + + +def test_inverted_period_is_rejected(): + response = requests.get( + f"{ANALYTICS_URL}/summary", + headers=_HEADERS, + params={"from": "2025-03-31", "to": "2025-03-01"}, + ) + assert response.status_code == 422 + + +# --------------------------------------------------------------------------- +# Time series +# --------------------------------------------------------------------------- + + +def test_series_buckets_days_in_the_shop_timezone_not_utc(): + """ + Order H is at 23:30 UTC on 9 March, which is 00:30 on the 10th in Berlin. + A UTC bucket would file it under the 9th and show the operator a day's + takings on the wrong day. + """ + payload = _get("timeseries", **WINDOW, interval="day") + eur = next(s for s in payload["series"] if s["currency"] == "EUR") + by_day = {point["bucket"]: point for point in eur["points"]} + + # Nothing else is seeded on either day, so this isolates the boundary. + assert by_day["2025-03-10"]["gross_revenue"] == 1000 + assert by_day["2025-03-09"]["gross_revenue"] == 0 + + +def test_quiet_days_are_present_as_zero_rather_than_missing(): + """A chart must show a trough on a day with no orders, not skip over it.""" + payload = _get("timeseries", **WINDOW, interval="day") + eur = next(s for s in payload["series"] if s["currency"] == "EUR") + + assert len(eur["points"]) == 31, "March has 31 days; every one needs a point" + + by_day = {point["bucket"]: point for point in eur["points"]} + assert by_day["2025-03-01"]["orders"] == 0 + assert by_day["2025-03-05"]["gross_revenue"] == 10000 + + +def test_series_net_revenue_subtracts_refunds_in_the_right_bucket(): + payload = _get("timeseries", **WINDOW, interval="day") + eur = next(s for s in payload["series"] if s["currency"] == "EUR") + by_day = {point["bucket"]: point for point in eur["points"]} + + assert by_day["2025-03-12"]["refunded_revenue"] == 3000 + assert by_day["2025-03-12"]["net_revenue"] == -3000 + + +# --------------------------------------------------------------------------- +# Products +# --------------------------------------------------------------------------- + + +def test_product_performance_aggregates_a_sku_across_orders(): + payload = _get("products", **WINDOW, limit=200) + by_sku = {p["sku"]: p for p in payload["products"]} + + sku_a = by_sku[f"{TAG}-SKU-A"] + # 2 units at 2500 on order A, 1 at 5000 on order B. Order F's 9 units are + # soft-deleted and order D's are a draft. + assert sku_a["units_sold"] == 3 + assert sku_a["gross_revenue"] == 10000 + assert sku_a["orders"] == 2 + + +def test_return_rate_is_reported_per_sku_on_the_returned_order(): + """ + Returns are recorded per order, so a return on the two-line order A counts + against both its SKUs. That is an upper bound, and the schema says so. + """ + payload = _get("products", **WINDOW, limit=200) + by_sku = {p["sku"]: p for p in payload["products"]} + + assert by_sku[f"{TAG}-SKU-A"]["orders_with_return"] == 1 + assert by_sku[f"{TAG}-SKU-B"]["orders_with_return"] == 1 + assert by_sku[f"{TAG}-SKU-B"]["return_rate"] == 1.0 + + +def test_products_can_be_sorted_by_units(): + payload = _get("products", **WINDOW, sort="units", limit=200) + units = [p["units_sold"] for p in payload["products"]] + assert units == sorted(units, reverse=True) + + +# --------------------------------------------------------------------------- +# Funnel +# --------------------------------------------------------------------------- + + +def _step(payload: dict, name: str) -> dict: + return next(s for s in payload["steps"] if s["step"] == name) + + +def test_funnel_counts_checkout_from_payments_not_order_status(): + """ + Orders A, B, D and G have payment rows; C, E and H do not. Order status + could not tell them apart — G is cancelled and D is a draft, yet both + reached checkout. + """ + payload = _get("funnel", **WINDOW) + + assert _step(payload, "created")["orders"] == 7 # F is soft-deleted + assert _step(payload, "checkout_started")["orders"] == 4 + assert _step(payload, "paid")["orders"] == 2 + assert _step(payload, "shipped")["orders"] == 1 + + +def test_funnel_separates_failed_payments_from_unresolved_ones(): + """ + A payment that failed and one still pending are different problems: one is + a lost sale, the other may still complete. + """ + payload = _get("funnel", **WINDOW) + + assert payload["payment_failed"] == 1 # order G + assert payload["payment_unresolved"] == 1 # order D, still pending + assert payload["never_checked_out"] == 3 # C, E, H + assert payload["cancelled"] == 1 # order G + + +def test_funnel_drop_off_is_the_loss_from_the_previous_step(): + payload = _get("funnel", **WINDOW) + + assert _step(payload, "created")["drop_off_from_previous"] is None + assert _step(payload, "checkout_started")["drop_off_from_previous"] == 3 + assert _step(payload, "paid")["drop_off_from_previous"] == 2 diff --git a/tests/test_analytics_unit.py b/tests/test_analytics_unit.py new file mode 100644 index 0000000..a2f74dd --- /dev/null +++ b/tests/test_analytics_unit.py @@ -0,0 +1,199 @@ +""" +Unit tests for the Analytics service — pure logic, no DB, no network. + +Covers the parts where the reasoning is subtle enough to get quietly wrong: + + - build_period() — half-open windows, timezone conversion, validation + - Period.previous() — the comparison baseline + - percent_change() — undefined rather than infinite growth from zero + - fill_series_gaps() — quiet days present as zeroes, not absent +""" + +from datetime import UTC, date, datetime + +import pytest + +from app.services.analytics.functions import ( + MAX_PERIOD_DAYS, + Interval, + Period, + build_period, + percent_change, + resolve_timezone, +) +from app.services.analytics.services import fill_series_gaps +from app.shared.exceptions import ValidationError + + +# --------------------------------------------------------------------------- +# build_period +# --------------------------------------------------------------------------- + + +def test_period_is_half_open_and_includes_the_final_day(): + """ + `to` is inclusive to the reader, so the window must extend to the end of + that day. A closed range on dates silently drops the last day's orders. + """ + period = build_period(date(2026, 8, 1), date(2026, 8, 31), "UTC") + + assert period.start == datetime(2026, 8, 1, 0, 0, tzinfo=UTC) + # Exclusive end: midnight opening 1 September, so 31 August is included. + assert period.end == datetime(2026, 9, 1, 0, 0, tzinfo=UTC) + assert period.days == 31 + + +def test_period_boundaries_are_cut_in_the_shop_timezone(): + """ + Berlin is UTC+2 in August, so the shop's day starts at 22:00 UTC the day + before. Bucketing on raw UTC would push evening orders into the next day. + """ + period = build_period(date(2026, 8, 1), date(2026, 8, 1), "Europe/Berlin") + + assert period.start.astimezone(UTC) == datetime(2026, 7, 31, 22, 0, tzinfo=UTC) + assert period.end.astimezone(UTC) == datetime(2026, 8, 1, 22, 0, tzinfo=UTC) + + +def test_previous_period_abuts_the_current_one_without_overlap(): + """ + The baseline must end exactly where the window starts. An overlap would + count boundary orders twice; a gap would lose them. + """ + period = build_period(date(2026, 8, 1), date(2026, 8, 31), "UTC") + previous = period.previous() + + assert previous.end == period.start + assert previous.end - previous.start == period.end - period.start + assert previous.start == datetime(2026, 7, 1, 0, 0, tzinfo=UTC) + + +def test_inverted_range_is_rejected(): + with pytest.raises(ValidationError): + build_period(date(2026, 8, 31), date(2026, 8, 1), "UTC") + + +def test_absurdly_long_range_is_rejected(): + """Live aggregation has no rollup table behind it, so the range is bounded.""" + with pytest.raises(ValidationError): + build_period(date(1990, 1, 1), date(2026, 8, 1), "UTC") + + +def test_range_at_the_limit_is_allowed(): + start = date(2026, 8, 1) + end = date(2026, 8, 1).replace(year=2026) + period = build_period(start, end, "UTC") + assert period.days <= MAX_PERIOD_DAYS + + +def test_unknown_timezone_is_a_validation_error_not_a_crash(): + """A misconfigured SHOP_TIMEZONE must say so, not return an opaque 500.""" + with pytest.raises(ValidationError): + resolve_timezone("Mars/Olympus_Mons") + + +def test_defaults_to_a_thirty_day_window_ending_today(): + period = build_period(None, date(2026, 8, 30), "UTC", default_days=30) + assert period.days == 30 + assert period.start == datetime(2026, 8, 1, 0, 0, tzinfo=UTC) + + +# --------------------------------------------------------------------------- +# percent_change +# --------------------------------------------------------------------------- + + +def test_percent_change_from_zero_is_undefined(): + """ + Growth from nothing is not 100% and not infinite. Reporting a number here + invites a reader to trust something meaningless. + """ + assert percent_change(500, 0) is None + + +def test_percent_change_computes_both_directions(): + assert percent_change(150, 100) == 50.0 + assert percent_change(50, 100) == -50.0 + assert percent_change(100, 100) == 0.0 + + +# --------------------------------------------------------------------------- +# fill_series_gaps +# --------------------------------------------------------------------------- + + +def _period(start: date, end: date) -> Period: + return build_period(start, end, "UTC") + + +def test_days_without_orders_appear_as_zero_rather_than_vanishing(): + """ + A chart built from present-only buckets joins straight across a quiet day, + which reads as steady trade rather than none. + """ + period = _period(date(2026, 8, 1), date(2026, 8, 5)) + buckets = { + date(2026, 8, 1): { + "gross_revenue": 1000, + "refunded_revenue": 0, + "orders": 2, + "units": 3, + }, + date(2026, 8, 5): { + "gross_revenue": 500, + "refunded_revenue": 0, + "orders": 1, + "units": 1, + }, + } + + points = fill_series_gaps(buckets, period, Interval.DAY) + + assert [p["bucket"] for p in points] == [ + date(2026, 8, 1), + date(2026, 8, 2), + date(2026, 8, 3), + date(2026, 8, 4), + date(2026, 8, 5), + ] + assert points[1]["gross_revenue"] == 0 + assert points[1]["orders"] == 0 + + +def test_net_revenue_is_gross_less_refunds(): + period = _period(date(2026, 8, 1), date(2026, 8, 1)) + buckets = { + date(2026, 8, 1): { + "gross_revenue": 1000, + "refunded_revenue": 250, + "orders": 2, + "units": 2, + } + } + + points = fill_series_gaps(buckets, period, Interval.DAY) + + assert points[0]["net_revenue"] == 750 + + +def test_empty_window_still_yields_one_point_per_day(): + """An empty range must render an empty chart, not an absent one.""" + period = _period(date(2026, 8, 1), date(2026, 8, 3)) + + points = fill_series_gaps({}, period, Interval.DAY) + + assert len(points) == 3 + assert all(p["gross_revenue"] == 0 for p in points) + + +def test_weekly_buckets_align_to_monday(): + """ + Postgres date_trunc('week') anchors to Monday, so the gap filler must use + the same anchor or every key misses and the whole series reads as zero. + """ + # 5 August 2026 is a Wednesday; its week starts Monday 3 August. + period = _period(date(2026, 8, 5), date(2026, 8, 18)) + + points = fill_series_gaps({}, period, Interval.WEEK) + + assert points[0]["bucket"] == date(2026, 8, 3) + assert points[0]["bucket"].weekday() == 0