feat(generator): добавлена команда слепков заказов #100

Merged
ddmitry merged 1 commits from feat/92-order-snapshot into main 2026-08-18 17:17:18 +03:00
11 changed files with 740 additions and 46 deletions
+37
View File
@@ -51,5 +51,42 @@
"events": 49394, "events": 49394,
"sha256": "286ccc2451847ac8de215b5e902c328371021ad5502bd4afca0494b6f313aef7" "sha256": "286ccc2451847ac8de215b5e902c328371021ad5502bd4afca0494b6f313aef7"
} }
],
"snapshots": [
{
"day": 0,
"date": "2026-06-01",
"sha256": "98736f45970f65e04256cd3c108a07332a7276769b7a5dbf8b5617b75116a3c3"
},
{
"day": 1,
"date": "2026-06-02",
"sha256": "141c593f939eb695b8b639d9f47cc92efb5de1aa2aeffbc5f3abac63fa603628"
},
{
"day": 2,
"date": "2026-06-03",
"sha256": "c6f594f03bf1bba119bf591d5232fe6ce8d3793effab64d2171208f267ff8cac"
},
{
"day": 3,
"date": "2026-06-04",
"sha256": "382b20303450bbb7928e1f118cf36cc005c6d6747a215715c5d5c8dce0822f95"
},
{
"day": 4,
"date": "2026-06-05",
"sha256": "ef743bad81a87725a767be45bfb18dfc74c20173ba0da5c41ccf254640c5af76"
},
{
"day": 5,
"date": "2026-06-06",
"sha256": "1618cf0f9a6816a1e302e0518a067cd4c7ca295dd1b7f4660c73f77da5182dfe"
},
{
"day": 6,
"date": "2026-06-07",
"sha256": "4e32a4cc9a8cc997fbb4bdf4727c22eb3a763e0d262262eccfe6a3094a14d142"
}
] ]
} }
+31 -16
View File
@@ -1,4 +1,4 @@
"""Интерфейс запуска: `python -m clickstream_generator batch|live`. """Интерфейс запуска: `python -m clickstream_generator batch|live|snapshot`.
Зовущий — контейнер (спека генератора, раздел 9), и интерфейс сделан под него: Зовущий — контейнер (спека генератора, раздел 9), и интерфейс сделан под него:
параметры приходят аргументами или переменными окружения, логи идут в параметры приходят аргументами или переменными окружения, логи идут в
@@ -17,6 +17,11 @@
различие за числом, а команда называет его словом. Ограниченная пачка живому различие за числом, а команда называет его словом. Ограниченная пачка живому
дню не полагается: ждать там нечего — ожидание снимает ускорение. дню не полагается: ждать там нечего — ожидание снимает ускорение.
**Третья команда — `snapshot`** — второй источник стенда: слепок заказов в свой
топик. У неё свой смысл `--day`: это день прогона, а уезжает слепок дня D−1
(правило сдвига живёт в проигрывателе). Прогон дня 0 не отправляет ничего —
на старте оси вчера нет.
**Приёмник выбирается тем, что для него назвали**: `--file` или `--brokers`. **Приёмник выбирается тем, что для него назвали**: `--file` или `--brokers`.
Оба сразу — ошибка, ни одного — тоже: молча выбранный по умолчанию приёмник Оба сразу — ошибка, ни одного — тоже: молча выбранный по умолчанию приёмник
однажды напишет в файл то, чего ждали в Kafka. однажды напишет в файл то, чего ждали в Kafka.
@@ -56,14 +61,19 @@ def main(argv: Sequence[str] | None = None) -> int:
_check(parser, options) _check(parser, options)
try: try:
with closing(_sink(parser, options)) as sink: with closing(_sink(parser, options)) as sink:
player.play( if options.command == "snapshot":
sink, player.play_snapshots(
seed=options.seed, sink, seed=options.seed, first_day=options.day, days=options.days
first_day=options.day, )
days=options.days, else:
limit=options.limit, player.play(
speed=options.speed, sink,
) seed=options.seed,
first_day=options.day,
days=options.days,
limit=options.limit,
speed=options.speed,
)
except (OSError, RuntimeError) as failure: except (OSError, RuntimeError) as failure:
logging.error("прогон не удался: %s", failure) logging.error("прогон не удался: %s", failure)
return 1 return 1
@@ -76,9 +86,9 @@ def _parser() -> argparse.ArgumentParser:
prog="python -m clickstream_generator", prog="python -m clickstream_generator",
description="Проигрыватель модельных дней кликстрима в файл или Kafka.", description="Проигрыватель модельных дней кликстрима в файл или Kafka.",
) )
modes = parser.add_subparsers(dest="mode", required=True) commands = parser.add_subparsers(dest="command", required=True)
batch = modes.add_parser( batch = commands.add_parser(
"batch", help="пачкой, без пауз: заливка снимка и переигровка дня" "batch", help="пачкой, без пауз: заливка снимка и переигровка дня"
) )
_common(batch) _common(batch)
@@ -91,7 +101,7 @@ def _parser() -> argparse.ArgumentParser:
help="взять не больше N событий на весь прогон, а не день целиком", help="взять не больше N событий на весь прогон, а не день целиком",
) )
live = modes.add_parser("live", help="живой день: темп модельного времени") live = commands.add_parser("live", help="живой день: темп модельного времени")
_common(live) _common(live)
live.set_defaults(limit=None) live.set_defaults(limit=None)
live.add_argument( live.add_argument(
@@ -101,11 +111,16 @@ def _parser() -> argparse.ArgumentParser:
metavar="X", metavar="X",
help=f"ускорение модельного времени (по умолчанию ×{DEFAULT_SPEED:.0f})", help=f"ускорение модельного времени (по умолчанию ×{DEFAULT_SPEED:.0f})",
) )
snapshot = commands.add_parser(
"snapshot", help="слепок заказов: прогон дня D отправляет слепок дня D−1"
)
_common(snapshot)
return parser return parser
def _common(parser: argparse.ArgumentParser) -> None: def _common(parser: argparse.ArgumentParser) -> None:
"""Что спрашивают у обоих режимов: какой мир, какие дни и куда.""" """Что спрашивают у всех команд: какой мир, какие дни и куда."""
parser.add_argument( parser.add_argument(
"--seed", "--seed",
type=int, type=int,
@@ -118,8 +133,8 @@ def _common(parser: argparse.ArgumentParser) -> None:
type=int, type=int,
default=_env_int("GENERATOR_DAY", None), default=_env_int("GENERATOR_DAY", None),
metavar="D", metavar="D",
help="номер дня на оси мира (D0 — первый); умолчания нет —" help="номер дня на оси мира (D0 — первый); у snapshot это день прогона,"
" позицию ведёт зовущий", " а уезжает слепок дня D−1; умолчания нет — позицию ведёт зовущий",
) )
parser.add_argument( parser.add_argument(
"--days", "--days",
@@ -133,7 +148,7 @@ def _common(parser: argparse.ArgumentParser) -> None:
type=Path, type=Path,
default=_env("GENERATOR_FILE", None, Path), default=_env("GENERATOR_FILE", None, Path),
metavar="ПУТЬ", metavar="ПУТЬ",
help="приёмник — файл: одно событие в строке", help="приёмник — файл: одна запись в строке",
) )
parser.add_argument( parser.add_argument(
"--brokers", "--brokers",
@@ -151,6 +151,10 @@ class Purchases:
# та самая, что уехала в событие как `purchaseRevenue`. Бэкенд назовёт # та самая, что уехала в событие как `purchaseRevenue`. Бэкенд назовёт
# её `items_total` — у двух источников свои имена одному числу. # её `items_total` — у двух источников свои имена одному числу.
revenue: NDArray[np.int64] revenue: NDArray[np.int64]
# Когда покупка подтверждена — та же абсолютная метка, что уехала в
# событие. У заказа она станет секундой `created_at`: строка в базе
# источника создаётся синхронно с покупкой.
moment: NDArray[np.datetime64]
def __len__(self) -> int: def __len__(self) -> int:
return self.person_id.size return self.person_id.size
@@ -185,6 +189,7 @@ _NO_PURCHASES = Purchases(
quantity=(), quantity=(),
coupon=(), coupon=(),
revenue=np.empty(0, dtype=np.int64), revenue=np.empty(0, dtype=np.int64),
moment=np.empty(0, dtype="datetime64[s]"),
) )
@@ -666,6 +671,7 @@ def _purchases(
quantity=tuple(draws.quantity[group] for group in bought), quantity=tuple(draws.quantity[group] for group in bought),
coupon=tuple(coupons[mine].tolist()), coupon=tuple(coupons[mine].tolist()),
revenue=revenue[mine], revenue=revenue[mine],
moment=rows["UTCEventTime"][here],
) )
@@ -26,6 +26,11 @@
пересчитывается он и обычным `sha256sum` по сыгранному в файл дню (как пересчитывается он и обычным `sha256sum` по сыгранному в файл дню (как
именно — в README репозитория). именно — в README репозитория).
Слепок заказов — второй артефакт мира, и хешами дней он не покрыт ни при каком
раскладе подпотоков, поэтому у каждого отправленного слепка своя строка. Их на
один меньше, чем дней: прогон дня 0 отправлять ещё нечего. Счёта строк у
слепка нет — опись описывает мир, а не доставку.
Собирается опись из каталога генератора целью `make inventory`, а свежесть её Собирается опись из каталога генератора целью `make inventory`, а свежесть её
сторожит тест — как и у «описания выгрузки». сторожит тест — как и у «описания выгрузки».
""" """
@@ -39,6 +44,7 @@ from pathlib import Path
from typing import Any from typing import Any
from clickstream_generator import day as day_module from clickstream_generator import day as day_module
from clickstream_generator import orders as orders_module
from clickstream_generator import serialize, world from clickstream_generator import serialize, world
from clickstream_generator.catalog import CATALOG_PATH from clickstream_generator.catalog import CATALOG_PATH
from clickstream_generator.seeds import CANONICAL_SEED from clickstream_generator.seeds import CANONICAL_SEED
@@ -50,12 +56,25 @@ STARTING_DAYS = 8
def build() -> dict[str, Any]: def build() -> dict[str, Any]:
"""Опись целиком: паспорт мира и по строке на каждый его день.""" """Опись целиком: паспорт мира, строка на день и строка на слепок.
Дни играются по одному и отпускаются: заказы дня остаются, потому что из
них собираются слепки, а полсотни тысяч событий восьми дней сразу в память
не нужны.
"""
days = []
orders = []
for number in range(STARTING_DAYS):
today = day_module.stream(CANONICAL_SEED, number)
days.append(_day(number, today))
orders.append(today.orders)
return { return {
"seed": CANONICAL_SEED, "seed": CANONICAL_SEED,
"generator_version": version("clickstream-generator"), "generator_version": version("clickstream-generator"),
"catalog_sha256": _digest(CATALOG_PATH.read_bytes()), "catalog_sha256": _digest(CATALOG_PATH.read_bytes()),
"days": [_day(number) for number in range(STARTING_DAYS)], "days": days,
"snapshots": [_snapshot(number, orders) for number in range(STARTING_DAYS - 1)],
} }
@@ -64,7 +83,7 @@ def render() -> str:
return json.dumps(build(), ensure_ascii=False, indent=2) + "\n" return json.dumps(build(), ensure_ascii=False, indent=2) + "\n"
def _day(number: int) -> dict[str, Any]: def _day(number: int, today: day_module.Day) -> dict[str, Any]:
"""Строка описи: номер дня, его дата, число событий и хеш байтов. """Строка описи: номер дня, его дата, число событий и хеш байтов.
Дата считается от D0 арифметикой, а не берётся из событий: ось модельного Дата считается от D0 арифметикой, а не берётся из событий: ось модельного
@@ -72,15 +91,34 @@ def _day(number: int) -> dict[str, Any]:
стенда обрамляют счёт в `ods.event`. Соври она — подневная сверка это и стенда обрамляют счёт в `ods.event`. Соври она — подневная сверка это и
покажет, каждый день сразу. покажет, каждый день сразу.
""" """
payloads = serialize.events(day_module.stream(CANONICAL_SEED, number)) payloads = serialize.events(today)
return { return {
"day": number, "day": number,
"date": (world.ORIGIN + timedelta(days=number)).isoformat(), "date": _date(number),
"events": len(payloads), "events": len(payloads),
"sha256": _digest(b"".join(payload + b"\n" for payload in payloads)), "sha256": _digest(b"".join(payload + b"\n" for payload in payloads)),
} }
def _snapshot(number: int, orders: list[orders_module.Orders]) -> dict[str, Any]:
"""Строка описи слепка: какой день снят, его дата и хеш отправленных байтов.
Байты те же, что уезжают в топик `orders`, с переводом строки после
каждого заказа: пересъёмка слепка обязана дать их снова.
"""
window = [orders[born] for born in orders_module.window(number)]
payloads = serialize.orders(window, number)
return {
"day": number,
"date": _date(number),
"sha256": _digest(b"".join(payload + b"\n" for payload in payloads)),
}
def _date(number: int) -> str:
return (world.ORIGIN + timedelta(days=number)).isoformat()
def _digest(payload: bytes) -> str: def _digest(payload: bytes) -> str:
return hashlib.sha256(payload).hexdigest() return hashlib.sha256(payload).hexdigest()
@@ -29,6 +29,11 @@
его место занимает −1, а не ноль, — иначе «оплатили в секунду рождения» было его место занимает −1, а не ноль, — иначе «оплатили в секунду рождения» было
бы не отличить от «не оплатили вовсе». бы не отличить от «не оплатили вовсе».
**Состояние на границе суток** — чтение готовой судьбы, а не накопление:
слепок дня D учитывает моменты не позже границы D|D+1 и по ним называет
статус. Отсюда «дыхание» окна — заказ, оплаченный назавтра, стоит в сегодняшнем
слепке как `created`, а завтра как `paid`.
**Дельта суммы — вычеркнутая позиция**: товара не оказалось в наличии, и **Дельта суммы — вычеркнутая позиция**: товара не оказалось в наличии, и
заказ приезжает на строку короче клиентской корзины, а `items_total` меньше заказ приезжает на строку короче клиентской корзины, а `items_total` меньше
ровно на её полную стоимость. Момента у дельты нет — склад собрал заказ до ровно на её полную стоимость. Момента у дельты нет — склад собрал заказ до
@@ -79,6 +84,10 @@ class Orders:
order_id: tuple[str, ...] order_id: tuple[str, ...]
# Пользователь магазина: тот же человек, что стоит за купившей кукой. # Пользователь магазина: тот же человек, что стоит за купившей кукой.
user_id: NDArray[np.uint64] user_id: NDArray[np.uint64]
# Когда строка заказа создана в базе источника: секунда той самой покупки
# — в модели строка создаётся синхронно с ней — и миллисекунда часов базы.
# От этого момента отсчитываются и моменты судьбы.
created_at: NDArray[np.datetime64]
# Позиции заказа: номера товаров каталога и штуки, ячейка на заказ. У # Позиции заказа: номера товаров каталога и штуки, ячейка на заказ. У
# заказа с дельтой позиций на одну меньше, чем в корзине клиента. # заказа с дельтой позиций на одну меньше, чем в корзине клиента.
product: tuple[NDArray[np.int64], ...] product: tuple[NDArray[np.int64], ...]
@@ -104,11 +113,13 @@ def of_day(seed: int, day: int, purchases: Purchases) -> Orders:
delivery = _delivery(rng, len(purchases)) delivery = _delivery(rng, len(purchases))
outcome, paid_after, cancelled_after = _fate(rng, len(purchases)) outcome, paid_after, cancelled_after = _fate(rng, len(purchases))
product, quantity, items_total = _delta(rng, purchases) product, quantity, items_total = _delta(rng, purchases)
created_at = _created_at(rng, purchases)
discount = _discount(purchases) discount = _discount(purchases)
return Orders( return Orders(
day=day, day=day,
order_id=purchases.order_id, order_id=purchases.order_id,
user_id=purchases.person_id, user_id=purchases.person_id,
created_at=created_at,
product=product, product=product,
quantity=quantity, quantity=quantity,
items_total=items_total, items_total=items_total,
@@ -121,6 +132,49 @@ def of_day(seed: int, day: int, purchases: Purchases) -> Orders:
) )
def window(day: int) -> range:
"""Дни рождения, чьи заказы несёт слепок дня `day`.
Окно изменяемости — константа мира; у начала оси оно усекается само, а не
сторожем: дней до D0 попросту нет.
"""
return range(max(0, day - world.ORDER_WINDOW_DAYS + 1), day + 1)
def at_boundary(rows: Orders, day: int) -> tuple[list[str], NDArray[np.datetime64]]:
"""Статус заказов и момент их последнего изменения на границе `day`|`day+1`.
Судьба решена при рождении, поэтому слепок её только читает: момент позже
границы для него ещё не случился. Отмена перевешивает оплату — у дороги
«оплачен и отменён» она поздняя, и заказ на границе уже отменён.
`updated_at` — поздний учтённый момент, а без единого заказ показывает своё
рождение: строку с тех пор никто не трогал.
"""
edge = _boundary(day)
paid = rows.created_at + rows.paid_after.astype("timedelta64[s]")
cancelled = rows.created_at + rows.cancelled_after.astype("timedelta64[s]")
# Момента, которого у исхода нет, в данных нет вовсе: там −1, и без маски
# он прикинулся бы моментом за секунду до рождения.
got_paid = (rows.paid_after >= 0) & (paid <= edge)
got_cancelled = (rows.cancelled_after >= 0) & (cancelled <= edge)
status = np.where(got_cancelled, "cancelled", np.where(got_paid, "paid", "created"))
updated = np.where(
got_cancelled, cancelled, np.where(got_paid, paid, rows.created_at)
)
return status.tolist(), updated
def _boundary(day: int) -> np.datetime64:
"""Граница суток `day`|`day+1` абсолютной меткой: полночь пояса счётчика.
Модельные сутки считаются в поясе счётчика, а моменты заказа — метки UTC:
между ними ровно смещение пояса.
"""
midnight = np.datetime64(world.ORIGIN, "s") + np.timedelta64(day + 1, "D")
return midnight - np.timedelta64(world.COUNTER_TIMEZONE_MINUTES, "m")
def _delivery(rng: np.random.Generator, orders: int) -> NDArray[np.int64]: def _delivery(rng: np.random.Generator, orders: int) -> NDArray[np.int64]:
"""Стоимость доставки каждого заказа дня — броском по таблице весов.""" """Стоимость доставки каждого заказа дня — броском по таблице весов."""
return _DELIVERY_PRICE[pick(rng, _DELIVERY_CUMULATIVE, orders)] return _DELIVERY_PRICE[pick(rng, _DELIVERY_CUMULATIVE, orders)]
@@ -147,6 +201,22 @@ def _fate(
return outcome, paid, cancelled return outcome, paid, cancelled
def _created_at(
rng: np.random.Generator, purchases: Purchases
) -> NDArray[np.datetime64]:
"""Момент создания строки заказа: секунда покупки и миллисекунда часов базы.
Секунда приходит из события — строка создаётся синхронно с покупкой, и
сдвигать её значило бы подделывать аудит источника. Миллисекунду трекер не
видит вовсе: у него своё разрешение, у базы своё, и три дописанных нуля
выдали бы секундную модель за миллисекундную (исследование формата слепка).
"""
millisecond = rng.integers(0, 1000, len(purchases))
return purchases.moment.astype("datetime64[ms]") + millisecond.astype(
"timedelta64[ms]"
)
def _moment(rng: np.random.Generator, orders: int) -> NDArray[np.int64]: def _moment(rng: np.random.Generator, orders: int) -> NDArray[np.int64]:
"""Момент внутри окна: час по таблице весов и равномерная секунда в нём. """Момент внутри окна: час по таблице весов и равномерная секунда в нём.
@@ -16,6 +16,11 @@
(раздел 5): это разные машины разной природы, и сложенные в одно число они (раздел 5): это разные машины разной природы, и сложенные в одно число они
перестают что-либо говорить. Сериализация считается частью генерации: она перестают что-либо говорить. Сериализация считается частью генерации: она
рождает те самые байты, которые сторожит опись. рождает те самые байты, которые сторожит опись.
**Слепки заказов идут третьим ходом того же проигрывателя** и живут по правилу
сдвига: прогон дня D отправляет слепок дня D1 ночная выгрузка бэкенда за
вчера. Правило записано здесь одно и целиком, потому что оно и есть разница
между двумя источниками: трекер шлёт сегодняшний день, бэкенд вчерашний.
""" """
import logging import logging
@@ -26,6 +31,7 @@ import numpy as np
from numpy.typing import NDArray from numpy.typing import NDArray
from clickstream_generator import day as day_module from clickstream_generator import day as day_module
from clickstream_generator import orders as orders_module
from clickstream_generator import serialize from clickstream_generator import serialize
from clickstream_generator.sinks import Sink from clickstream_generator.sinks import Sink
@@ -103,6 +109,70 @@ def play(
) )
def play_snapshots(sink: Sink, seed: int, first_day: int, days: int = 1) -> None:
"""Отправить слепки прогонов `days` дней подряд начиная с `first_day`.
Прогон дня D отправляет слепок дня D1: содержимое слепка чистая функция
зерна и дня, от момента отправки оно не зависит, а сдвиг делает живой день
обычным. На старте оси вчера нет, поэтому прогон дня 0 не отправляет
ничего.
Слепок собирается переигровкой дней своего окна: заказы дня производная
всей воронки, дешёвого пути к ним нет. Окна соседних слепков перекрываются
почти целиком, и внутри одного прогона день играется один раз иначе
стартовый диапазон стоил бы полусотни проигрышей вместо восьми. Между
прогонами не остаётся ничего: кэш на томе был бы состоянием, которого у
проигрывателя нет.
"""
played: dict[int, orders_module.Orders] = {}
sent = 0
generated = 0.0
delivered = 0.0
for number in range(first_day, first_day + days):
taken = number - 1
if taken < 0:
continue
clock = time.monotonic()
window = [
_orders_of(played, seed, born) for born in orders_module.window(taken)
]
payloads = serialize.orders(window, taken)
spent = time.monotonic() - clock
generated += spent
log.info(
"слепок дня %d: заказов %d, генерация %.1f с", taken, len(payloads), spent
)
clock = time.monotonic()
_send(sink, payloads)
sink.flush()
spent = time.monotonic() - clock
delivered += spent
sent += len(payloads)
log.info(
"слепок дня %d: отправлено %d, доставка %.1f с", taken, len(payloads), spent
)
log.info(
"итого отправлено %d заказов: генерация %.1f с, доставка %.1f с",
sent,
generated,
delivered,
)
def _orders_of(
played: dict[int, orders_module.Orders], seed: int, number: int
) -> orders_module.Orders:
"""Заказы дня `number`, сыгранного один раз на весь прогон."""
if number not in played:
played[number] = day_module.stream(seed, number).orders
return played[number]
def _send(sink: Sink, payloads: list[bytes]) -> None: def _send(sink: Sink, payloads: list[bytes]) -> None:
"""Пакетно: подряд и без пауз.""" """Пакетно: подряд и без пауз."""
for payload in payloads: for payload in payloads:
@@ -27,16 +27,25 @@
делает это одним вызовом на колонку, и на дне в полсотни тысяч событий разница делает это одним вызовом на колонку, и на дне в полсотни тысяч событий разница
заметна. Обратная сторона день лежит в памяти дважды; проигрыватель поэтому заметна. Обратная сторона день лежит в памяти дважды; проигрыватель поэтому
и берёт его днями, а не горизонтом целиком. и берёт его днями, а не горизонтом целиком.
**Второй контракт провода слепок заказов** (мастер-спека, раздел 2). Он не
похож на событие: одиннадцать ключей вместо сорока семи, деньги строками, а
не числами, времена с миллисекундами. Общее у них одно, зато главное: байты
рождаются здесь и только здесь. Запись слепка один словарь с вложенным
списком и один `orjson.dumps`.
""" """
from collections.abc import Sequence
from datetime import timedelta
from typing import Any from typing import Any
import numpy as np import numpy as np
import orjson import orjson
from numpy.typing import NDArray from numpy.typing import NDArray
from clickstream_generator import schema from clickstream_generator import catalog, schema, world
from clickstream_generator.day import Day from clickstream_generator.day import Day
from clickstream_generator.orders import Orders, at_boundary
_ARRAY_PREFIX = "Array(" _ARRAY_PREFIX = "Array("
@@ -60,6 +69,78 @@ def events(day: Day, limit: int | None = None) -> list[bytes]:
] ]
def orders(window: Sequence[Orders], day: int) -> list[bytes]:
"""Канонические байты слепка дня `day`: по документу JSON на заказ.
`window` заказы дней окна, от раннего дня к позднему: слепок несёт их
подряд, и порядок строк выходит порядком рождения заказов, он же
возрастание `order_id`. Какие это дни, решает `orders.window`.
Деньги уезжают строками с ровно двумя знаками, а не числами: у заказа они
станут `Decimal`, и дробь двоичного числа была бы потерей точности до
всякого разбора. Времена метки UTC с миллисекундами; `snapshot_date`
одинакова во всей выгрузке это дата дня, состояние которого снято.
"""
goods = catalog.catalog()
sku = goods.sku.tolist()
prices = [_money(price) for price in goods.price.tolist()]
snapshot_date = (world.ORIGIN + timedelta(days=day)).isoformat()
payloads = []
for rows in window:
status, updated = at_boundary(rows, day)
created_at = _moments(rows.created_at)
updated_at = _moments(updated)
user_id = rows.user_id.tolist()
items_total = rows.items_total.tolist()
discount = rows.discount.tolist()
delivery = rows.delivery.tolist()
total = rows.total.tolist()
for number, order_id in enumerate(rows.order_id):
payloads.append(
orjson.dumps(
{
"order_id": order_id,
"user_id": user_id[number],
"status": status[number],
"created_at": created_at[number],
"updated_at": updated_at[number],
"items_total": _money(items_total[number]),
"discount": _money(discount[number]),
"delivery": _money(delivery[number]),
"total": _money(total[number]),
"items": [
{"sku": sku[item], "qty": count, "price": prices[item]}
for item, count in zip(
rows.product[number].tolist(),
rows.quantity[number].tolist(),
strict=True,
)
],
"snapshot_date": snapshot_date,
}
)
)
return payloads
def _money(kopecks: int) -> str:
"""Копейки — строкой с ровно двумя знаками: `129990` → `1299.90`."""
return f"{kopecks // 100}.{kopecks % 100:02d}"
def _moments(values: NDArray[np.datetime64]) -> list[str]:
"""Метки времени — строками RFC 3339 в UTC с миллисекундами.
Три знака стоят всегда, в том числе `.000`: одинаковая длина дробной части
и одинаковая зона дают хронологическую сортировку простым сравнением строк,
а разбор в хранилище идёт по точному шаблону.
"""
ms: NDArray[np.datetime64] = values.astype("datetime64[ms]")
return np.datetime_as_string(ms, unit="ms", timezone="UTC").tolist()
def _values(column: schema.Column, values: NDArray[Any]) -> list[Any]: def _values(column: schema.Column, values: NDArray[Any]) -> list[Any]:
"""Колонка питоновскими значениями — в той записи, в какой уедет на провод. """Колонка питоновскими значениями — в той записи, в какой уедет на провод.
+26
View File
@@ -0,0 +1,26 @@
"""Общее для всех проверок: чистое окружение прогона.
Переменные запуска живут на этой машине по-настоящему стенд их экспортирует.
Тест, читающий их у машины, зелен у одного и красен у другого, поэтому интерфейс
запуска проверяется только тем, что ему передали аргументами.
"""
import pytest
LAUNCH_VARIABLES = (
"GENERATOR_SEED",
"GENERATOR_DAY",
"GENERATOR_DAYS",
"GENERATOR_LIMIT",
"GENERATOR_SPEED",
"GENERATOR_FILE",
"KAFKA_BOOTSTRAP_SERVERS",
"KAFKA_TOPIC",
)
@pytest.fixture(autouse=True)
def bare_environment(monkeypatch):
"""Прогон тестов не зависит от того, что задано в окружении машины."""
for name in LAUNCH_VARIABLES:
monkeypatch.delenv(name, raising=False)
+38 -3
View File
@@ -10,18 +10,53 @@
поднятом стенде, где расхождение счётчиков выглядит поломкой хранилища. поднятом стенде, где расхождение счётчиков выглядит поломкой хранилища.
""" """
import hashlib
import json import json
from pathlib import Path from pathlib import Path
from clickstream_generator.inventory import build import pytest
from clickstream_generator import cli
from clickstream_generator.inventory import STARTING_DAYS, build
INVENTORY_PATH = Path(__file__).resolve().parents[2] / "data" / "world-inventory.json" INVENTORY_PATH = Path(__file__).resolve().parents[2] / "data" / "world-inventory.json"
# Прогон, слепок которого сверяется с описью байт в байт. День любой из
# отправляемых; этот дешевле прочих — его окно короче.
RUN_DAY = 2
def test_inventory_is_up_to_date():
@pytest.fixture(scope="module")
def inventory() -> dict:
"""Опись, собранная из кода: мир пересчитывается один раз на весь модуль."""
return build()
def test_inventory_is_up_to_date(inventory: dict):
stored = json.loads(INVENTORY_PATH.read_text(encoding="utf-8")) stored = json.loads(INVENTORY_PATH.read_text(encoding="utf-8"))
assert stored == build(), ( assert stored == inventory, (
"опись мира отстала от кода — пересоберите: make inventory." "опись мира отстала от кода — пересоберите: make inventory."
" Разошлись хеши дней и каталога — правили data/catalog/products.csv;" " Разошлись хеши дней и каталога — правили data/catalog/products.csv;"
" разошлись только дни — правили генератор" " разошлись только дни — правили генератор"
) )
def test_the_inventory_holds_the_hash_of_every_sent_snapshot(inventory: dict, tmp_path):
"""У каждого отправленного слепка — своя строка с хешем его байтов.
Байты слепка не покрыты хешами дней ни при каком раскладе подпотоков: это
второй артефакт мира, и побайтовое обещание сторожит опись. Слепков на день
меньше, чем дней: прогон дня 0 не отправляет ничего.
Хеш сверяется с настоящей выгрузкой, а не с самим собой: в файл уезжают те
же байты, что и в Kafka, по строке на заказ.
"""
days = [row["day"] for row in inventory["snapshots"]]
assert days == list(range(STARTING_DAYS - 1))
path = tmp_path / "snapshot.jsonl"
assert cli.main(["snapshot", "--day", str(RUN_DAY), "--file", str(path)]) == 0
sent = hashlib.sha256(path.read_bytes()).hexdigest()
row = next(row for row in inventory["snapshots"] if row["day"] == RUN_DAY - 1)
assert row["sha256"] == sent
-21
View File
@@ -22,27 +22,6 @@ from clickstream_generator.sinks import FileSink
DAY = 2 DAY = 2
# Переменные, которыми зовущий задаёт прогон. Тест, читающий их из окружения
# машины, зелен у одного и красен у другого — а на этой машине они как раз и
# живут: стенд их экспортирует.
LAUNCH_VARIABLES = (
"GENERATOR_SEED",
"GENERATOR_DAY",
"GENERATOR_DAYS",
"GENERATOR_LIMIT",
"GENERATOR_SPEED",
"GENERATOR_FILE",
"KAFKA_BOOTSTRAP_SERVERS",
"KAFKA_TOPIC",
)
@pytest.fixture(autouse=True)
def bare_environment(monkeypatch):
"""Прогон тестов не зависит от того, что задано в окружении машины."""
for name in LAUNCH_VARIABLES:
monkeypatch.delenv(name, raising=False)
def _play(path, **options): def _play(path, **options):
"""Проиграть в файл и вернуть итог прогона.""" """Проиграть в файл и вернуть итог прогона."""
+337
View File
@@ -0,0 +1,337 @@
"""Слепок заказов: окно, граница суток, запись на проводе и байты прогона.
Слепок второй артефакт мира, и сторожится он тем же, чем день событий:
формой записи и побайтовым повтором. Проверяются пять обещаний
(docs/architecture/orders/snapshot.md):
1. **Окно.** Слепок дня D несёт заказы, рождённые в дни D6D, и у начала оси
усекается сам.
2. **Граница суток.** Состояние чтение готовой судьбы на границе D|D+1:
моменты позже границы не учитываются, поэтому заказ «дышит» `created` в
одном слепке, `paid` в следующем.
3. **Запись на проводе.** Одиннадцать ключей в порядке контракта, деньги
строками с двумя знаками, времена RFC 3339 с миллисекундами.
4. **Сдвиг отправки.** Прогон дня D отправляет слепок дня D1; прогон дня 0
ничего.
5. **Повтор.** Пересъёмка слепка даёт те же байты.
День рождения заказа берётся из его номера, а не из `created_at`: номер
считается в поясе счётчика, а `created_at` абсолютная метка, и у ночной
покупки их даты расходятся.
"""
import hashlib
import json
import re
import subprocess
import sys
from contextlib import closing
from datetime import UTC, date, datetime, time, timedelta
from decimal import Decimal
import numpy as np
import pytest
from clickstream_generator import catalog, cli, commerce, player, serialize, world
from clickstream_generator import day as day_module
from clickstream_generator.seeds import CANONICAL_SEED
from clickstream_generator.sinks import FileSink
# Первый слепок с полным окном и слепок у начала оси, где окно усечено.
FULL_WINDOW_DAY = world.ORDER_WINDOW_DAYS - 1
SHORT_WINDOW_DAY = 3
# Порядок ключей записи — порядок полей контракта (мастер-спека, раздел 2).
CONTRACT = (
"order_id",
"user_id",
"status",
"created_at",
"updated_at",
"items_total",
"discount",
"delivery",
"total",
"items",
"snapshot_date",
)
ITEM = ("sku", "qty", "price")
STATUSES = {"created", "paid", "cancelled"}
MONEY = re.compile(r"\d+\.\d{2}")
MOMENT = re.compile(r"\d{4}-\d{2}-\d{2}T\d{2}:\d{2}:\d{2}\.\d{3}Z")
@pytest.fixture(scope="module")
def days() -> list[day_module.Day]:
"""Дни окна: слепок дня 6 собирается переигровкой дней 0…6."""
return [
day_module.stream(CANONICAL_SEED, number)
for number in range(world.ORDER_WINDOW_DAYS)
]
def snapshot_records(days: list[day_module.Day], number: int) -> list[dict]:
"""Слепок дня `number` разобранным обратно из канонических байтов."""
window = [
today.orders for today in days[max(0, number - FULL_WINDOW_DAY) : number + 1]
]
return [json.loads(payload) for payload in serialize.orders(window, number)]
def test_the_record_matches_the_wire_contract(days: list[day_module.Day]):
"""Запись слепка — ровно то, что примет строгий приём хранилища.
Набор и порядок ключей, деньги строками с двумя знаками, времена RFC 3339
в UTC с миллисекундами, `items` обычным массивом всё это контракт
провода, а не вкус: разбор в ODS сверяет набор ключей и форму значений, и
строка мимо формы уходит в брак целиком (спецификация приёма заказов).
Деньги сверяются с позициями и между собой: строка на проводе обязана
отвечать тому же миру, что и заказ в памяти. Цена позиции цена каталога,
а `created_at` время той самой покупки, что уехала в трекер: строка
заказа создаётся синхронно с ней.
"""
goods = catalog.catalog()
price = dict(zip(goods.sku.tolist(), goods.price.tolist(), strict=True))
bought = _purchase_times(days[SHORT_WINDOW_DAY])
records = snapshot_records(days, SHORT_WINDOW_DAY)
assert records
for record in records:
assert list(record) == list(CONTRACT)
assert isinstance(record["user_id"], int)
assert record["status"] in STATUSES
assert MOMENT.fullmatch(record["created_at"])
assert MOMENT.fullmatch(record["updated_at"])
assert record["snapshot_date"] == _date(SHORT_WINDOW_DAY).isoformat()
assert isinstance(record["items"], list)
assert record["items"]
lines = Decimal(0)
for item in record["items"]:
assert list(item) == list(ITEM)
assert MONEY.fullmatch(item["price"])
assert isinstance(item["qty"], int)
assert item["qty"] > 0
assert Decimal(item["price"]) * 100 == price[item["sku"]]
lines += Decimal(item["price"]) * item["qty"]
money = {name: record[name] for name in CONTRACT[5:9]}
for value in money.values():
assert MONEY.fullmatch(value)
assert lines == Decimal(money["items_total"])
assert Decimal(money["total"]) == (
Decimal(money["items_total"])
- Decimal(money["discount"])
+ Decimal(money["delivery"])
)
if _born(record) == _date(SHORT_WINDOW_DAY):
assert record["created_at"][:19] + "Z" == bought[record["order_id"]]
def test_the_moments_carry_real_milliseconds(days: list[day_module.Day]):
"""Миллисекунды — часы базы источника, а не три дописанных нуля.
Три знака дробной части контракт требует всегда, и `.000` формально им
отвечают но тогда миллисекундная точность была бы обещанием без модели за
ним, а разбор в `DateTime64(3)` показывал бы менти ровную секундную сетку
там, где у источника её нет (исследование формата слепка).
Фаза у обоих моментов одна: часы базы ставят её строке при создании, а
судьба заказа отмеряется от неё целыми секундами.
"""
records = snapshot_records(days, SHORT_WINDOW_DAY)
assert records
fractions = {record["created_at"][20:23] for record in records}
assert fractions - {"000"}
for record in records:
assert record["updated_at"][20:23] == record["created_at"][20:23]
def test_the_rows_go_in_the_order_of_birth(days: list[day_module.Day]):
"""Порядок строк слепка — порядок рождения заказов, он же рост номера.
Хешу слепка в описи нужен именно названный порядок: детерминизм даёт его
даром, но обещание побайтового повтора держится на нём, а не на удаче.
"""
numbers = [record["order_id"] for record in snapshot_records(days, FULL_WINDOW_DAY)]
assert numbers
assert numbers == sorted(numbers)
assert len(set(numbers)) == len(numbers)
def test_the_state_is_read_at_the_boundary_of_the_day(days: list[day_module.Day]):
"""Дыхание окна: оплаченный назавтра заказ в сегодняшнем слепке `created`.
Это и есть состояние на границе суток: учитываются моменты не позже
границы, поэтому оплата следующего дня в сегодняшний слепок не попадает.
`updated_at` поздний учтённый момент, а без единого `created_at`:
заказ, с которым на границе ещё ничего не случилось, показывает своё
рождение.
Граница модельных суток полночь пояса счётчика: по нему считается
модельный день, а `created_at` и `updated_at` уезжают абсолютной меткой.
"""
before = {
record["order_id"]: record
for record in snapshot_records(days, SHORT_WINDOW_DAY)
}
after = {
record["order_id"]: record
for record in snapshot_records(days, SHORT_WINDOW_DAY + 1)
}
breathed = [
number
for number, record in before.items()
if record["status"] == "created" and after[number]["status"] == "paid"
]
assert breathed
for number in breathed:
assert before[number]["updated_at"] == before[number]["created_at"]
assert after[number]["created_at"] == before[number]["created_at"]
assert after[number]["updated_at"] > after[number]["created_at"]
edge = _boundary(SHORT_WINDOW_DAY)
for record in before.values():
assert record["updated_at"] <= edge
def test_the_run_sends_the_snapshot_of_the_day_before(
days: list[day_module.Day], tmp_path
):
"""Прогон дня D отправляет слепок дня D−1, и окно усекается у начала оси.
Сдвиг «ночная выгрузка за вчера» живёт в проигрывателе, поэтому и
спрашивается с него: прогоны дней 47 обязаны отправить слепки дней 36
теми же байтами, какие даёт сериализатор, по строке на заказ.
"""
path = tmp_path / "snapshots.jsonl"
with closing(FileSink(path)) as sink:
player.play_snapshots(
sink,
seed=CANONICAL_SEED,
first_day=SHORT_WINDOW_DAY + 1,
days=FULL_WINDOW_DAY - SHORT_WINDOW_DAY + 1,
)
assert path.read_bytes() == b"".join(
payload + b"\n"
for number in range(SHORT_WINDOW_DAY, FULL_WINDOW_DAY + 1)
for payload in serialize.orders(
[today.orders for today in days[: number + 1]], number
)
)
sent = [json.loads(line) for line in path.read_bytes().splitlines()]
for number in (SHORT_WINDOW_DAY, FULL_WINDOW_DAY):
window = {
_born(record)
for record in sent
if record["snapshot_date"] == _date(number).isoformat()
}
assert window == {
_date(born) for born in range(max(0, number - FULL_WINDOW_DAY), number + 1)
}
def test_the_range_plays_every_day_of_its_windows_once(monkeypatch, tmp_path):
"""Диапазон одним прогоном переигрывает каждый день окна один раз.
Окна соседних слепков перекрываются почти целиком, и без переиспользования
стартовый диапазон стоил бы полусотни проигрышей вместо восьми. День здесь
настоящий: считается не содержимое, а то, сколько раз его спросили.
"""
asked: list[int] = []
honest = day_module.stream
def spy(seed: int, number: int) -> day_module.Day:
asked.append(number)
return honest(seed, number)
monkeypatch.setattr(player.day_module, "stream", spy)
with closing(FileSink(tmp_path / "range.jsonl")) as sink:
player.play_snapshots(sink, seed=CANONICAL_SEED, first_day=1, days=3)
assert sorted(asked) == [0, 1, 2]
def test_the_first_run_sends_nothing(tmp_path):
"""На старте оси дня −1 нет: прогону дня 0 отправлять нечего.
Не ошибка и не пустой слепок, а отсутствие выгрузки: первый слепок дня 0
уезжает прогоном дня 1.
"""
path = tmp_path / "day-zero.jsonl"
assert cli.main(["snapshot", "--day", "0", "--file", str(path)]) == 0
assert path.read_bytes() == b""
def test_two_runs_give_the_same_snapshot(tmp_path):
"""Пересъёмка слепка даёт те же байты — тем и лечится пропущенный день.
Прогоны идут разными процессами по тому же доводу, что и у дня событий:
внутри одного интерпретатора общее зерно хеширования спрятало бы
зависимость канона от порядка обхода множества.
"""
first, second = tmp_path / "first.jsonl", tmp_path / "second.jsonl"
_run_apart(first)
_run_apart(second)
assert first.read_bytes()
assert _digest(first) == _digest(second)
def _run_apart(path) -> None:
"""Снять слепок отдельным процессом; окружение он берёт от нас."""
finished = subprocess.run(
[
sys.executable,
"-m",
"clickstream_generator",
"snapshot",
"--day",
"2",
"--file",
str(path),
],
capture_output=True,
text=True,
timeout=300,
check=False,
)
assert finished.returncode == 0, finished.stderr
def _purchase_times(today: day_module.Day) -> dict[str, str]:
"""Когда покупка уехала в трекер — по номеру заказа, абсолютной меткой."""
here = today.columns["EventType"] == commerce.PURCHASE
numbers = [cell[0] for cell in today.columns["purchaseID"][here]]
times = np.datetime_as_string(
today.columns["UTCEventTime"][here], unit="s", timezone="UTC"
)
return dict(zip(numbers, times.tolist(), strict=True))
def _born(record: dict) -> date:
"""День рождения заказа: его номер начинается датой дня покупки."""
return datetime.strptime(record["order_id"].split("-")[0], "%Y%m%d").date()
def _date(number: int) -> date:
return world.ORIGIN + timedelta(days=number)
def _boundary(number: int) -> str:
"""Граница суток `number`|`number+1` — полночь пояса счётчика в UTC."""
midnight = datetime.combine(_date(number + 1), time.min, UTC)
edge = midnight - timedelta(minutes=world.COUNTER_TIMEZONE_MINUTES)
return edge.strftime("%Y-%m-%dT%H:%M:%S.") + f"{edge.microsecond // 1000:03d}Z"
def _digest(path) -> str:
return hashlib.sha256(path.read_bytes()).hexdigest()