diff --git a/data/world-inventory.json b/data/world-inventory.json index 07bb7f8..1c8773e 100644 --- a/data/world-inventory.json +++ b/data/world-inventory.json @@ -51,5 +51,42 @@ "events": 49394, "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" + } ] } diff --git a/generator/src/clickstream_generator/cli.py b/generator/src/clickstream_generator/cli.py index 6dd719a..dc2187b 100644 --- a/generator/src/clickstream_generator/cli.py +++ b/generator/src/clickstream_generator/cli.py @@ -1,4 +1,4 @@ -"""Интерфейс запуска: `python -m clickstream_generator batch|live`. +"""Интерфейс запуска: `python -m clickstream_generator batch|live|snapshot`. Зовущий — контейнер (спека генератора, раздел 9), и интерфейс сделан под него: параметры приходят аргументами или переменными окружения, логи идут в @@ -17,6 +17,11 @@ различие за числом, а команда называет его словом. Ограниченная пачка живому дню не полагается: ждать там нечего — ожидание снимает ускорение. +**Третья команда — `snapshot`** — второй источник стенда: слепок заказов в свой +топик. У неё свой смысл `--day`: это день прогона, а уезжает слепок дня D−1 +(правило сдвига живёт в проигрывателе). Прогон дня 0 не отправляет ничего — +на старте оси вчера нет. + **Приёмник выбирается тем, что для него назвали**: `--file` или `--brokers`. Оба сразу — ошибка, ни одного — тоже: молча выбранный по умолчанию приёмник однажды напишет в файл то, чего ждали в Kafka. @@ -56,14 +61,19 @@ def main(argv: Sequence[str] | None = None) -> int: _check(parser, options) try: with closing(_sink(parser, options)) as sink: - player.play( - sink, - seed=options.seed, - first_day=options.day, - days=options.days, - limit=options.limit, - speed=options.speed, - ) + if options.command == "snapshot": + player.play_snapshots( + sink, seed=options.seed, first_day=options.day, days=options.days + ) + else: + player.play( + sink, + seed=options.seed, + first_day=options.day, + days=options.days, + limit=options.limit, + speed=options.speed, + ) except (OSError, RuntimeError) as failure: logging.error("прогон не удался: %s", failure) return 1 @@ -76,9 +86,9 @@ def _parser() -> argparse.ArgumentParser: prog="python -m clickstream_generator", 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="пачкой, без пауз: заливка снимка и переигровка дня" ) _common(batch) @@ -91,7 +101,7 @@ def _parser() -> argparse.ArgumentParser: help="взять не больше N событий на весь прогон, а не день целиком", ) - live = modes.add_parser("live", help="живой день: темп модельного времени") + live = commands.add_parser("live", help="живой день: темп модельного времени") _common(live) live.set_defaults(limit=None) live.add_argument( @@ -101,11 +111,16 @@ def _parser() -> argparse.ArgumentParser: metavar="X", help=f"ускорение модельного времени (по умолчанию ×{DEFAULT_SPEED:.0f})", ) + + snapshot = commands.add_parser( + "snapshot", help="слепок заказов: прогон дня D отправляет слепок дня D−1" + ) + _common(snapshot) return parser def _common(parser: argparse.ArgumentParser) -> None: - """Что спрашивают у обоих режимов: какой мир, какие дни и куда.""" + """Что спрашивают у всех команд: какой мир, какие дни и куда.""" parser.add_argument( "--seed", type=int, @@ -118,8 +133,8 @@ def _common(parser: argparse.ArgumentParser) -> None: type=int, default=_env_int("GENERATOR_DAY", None), metavar="D", - help="номер дня на оси мира (D0 — первый); умолчания нет —" - " позицию ведёт зовущий", + help="номер дня на оси мира (D0 — первый); у snapshot это день прогона," + " а уезжает слепок дня D−1; умолчания нет — позицию ведёт зовущий", ) parser.add_argument( "--days", @@ -133,7 +148,7 @@ def _common(parser: argparse.ArgumentParser) -> None: type=Path, default=_env("GENERATOR_FILE", None, Path), metavar="ПУТЬ", - help="приёмник — файл: одно событие в строке", + help="приёмник — файл: одна запись в строке", ) parser.add_argument( "--brokers", diff --git a/generator/src/clickstream_generator/commerce.py b/generator/src/clickstream_generator/commerce.py index 751bab5..6eee879 100644 --- a/generator/src/clickstream_generator/commerce.py +++ b/generator/src/clickstream_generator/commerce.py @@ -151,6 +151,10 @@ class Purchases: # та самая, что уехала в событие как `purchaseRevenue`. Бэкенд назовёт # её `items_total` — у двух источников свои имена одному числу. revenue: NDArray[np.int64] + # Когда покупка подтверждена — та же абсолютная метка, что уехала в + # событие. У заказа она станет секундой `created_at`: строка в базе + # источника создаётся синхронно с покупкой. + moment: NDArray[np.datetime64] def __len__(self) -> int: return self.person_id.size @@ -185,6 +189,7 @@ _NO_PURCHASES = Purchases( quantity=(), coupon=(), 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), coupon=tuple(coupons[mine].tolist()), revenue=revenue[mine], + moment=rows["UTCEventTime"][here], ) diff --git a/generator/src/clickstream_generator/inventory.py b/generator/src/clickstream_generator/inventory.py index 0291acd..9e3aba7 100644 --- a/generator/src/clickstream_generator/inventory.py +++ b/generator/src/clickstream_generator/inventory.py @@ -26,6 +26,11 @@ пересчитывается он и обычным `sha256sum` по сыгранному в файл дню (как именно — в README репозитория). +Слепок заказов — второй артефакт мира, и хешами дней он не покрыт ни при каком +раскладе подпотоков, поэтому у каждого отправленного слепка своя строка. Их на +один меньше, чем дней: прогон дня 0 отправлять ещё нечего. Счёта строк у +слепка нет — опись описывает мир, а не доставку. + Собирается опись из каталога генератора целью `make inventory`, а свежесть её сторожит тест — как и у «описания выгрузки». """ @@ -39,6 +44,7 @@ from pathlib import Path from typing import Any 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.catalog import CATALOG_PATH from clickstream_generator.seeds import CANONICAL_SEED @@ -50,12 +56,25 @@ STARTING_DAYS = 8 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 { "seed": CANONICAL_SEED, "generator_version": version("clickstream-generator"), "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" -def _day(number: int) -> dict[str, Any]: +def _day(number: int, today: day_module.Day) -> dict[str, Any]: """Строка описи: номер дня, его дата, число событий и хеш байтов. Дата считается от D0 арифметикой, а не берётся из событий: ось модельного @@ -72,15 +91,34 @@ def _day(number: int) -> dict[str, Any]: стенда обрамляют счёт в `ods.event`. Соври она — подневная сверка это и покажет, каждый день сразу. """ - payloads = serialize.events(day_module.stream(CANONICAL_SEED, number)) + payloads = serialize.events(today) return { "day": number, - "date": (world.ORIGIN + timedelta(days=number)).isoformat(), + "date": _date(number), "events": len(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: return hashlib.sha256(payload).hexdigest() diff --git a/generator/src/clickstream_generator/orders.py b/generator/src/clickstream_generator/orders.py index cd8855d..800e734 100644 --- a/generator/src/clickstream_generator/orders.py +++ b/generator/src/clickstream_generator/orders.py @@ -29,6 +29,11 @@ его место занимает −1, а не ноль, — иначе «оплатили в секунду рождения» было бы не отличить от «не оплатили вовсе». +**Состояние на границе суток** — чтение готовой судьбы, а не накопление: +слепок дня D учитывает моменты не позже границы D|D+1 и по ним называет +статус. Отсюда «дыхание» окна — заказ, оплаченный назавтра, стоит в сегодняшнем +слепке как `created`, а завтра как `paid`. + **Дельта суммы — вычеркнутая позиция**: товара не оказалось в наличии, и заказ приезжает на строку короче клиентской корзины, а `items_total` меньше ровно на её полную стоимость. Момента у дельты нет — склад собрал заказ до @@ -79,6 +84,10 @@ class Orders: order_id: tuple[str, ...] # Пользователь магазина: тот же человек, что стоит за купившей кукой. user_id: NDArray[np.uint64] + # Когда строка заказа создана в базе источника: секунда той самой покупки + # — в модели строка создаётся синхронно с ней — и миллисекунда часов базы. + # От этого момента отсчитываются и моменты судьбы. + created_at: NDArray[np.datetime64] # Позиции заказа: номера товаров каталога и штуки, ячейка на заказ. У # заказа с дельтой позиций на одну меньше, чем в корзине клиента. product: tuple[NDArray[np.int64], ...] @@ -104,11 +113,13 @@ def of_day(seed: int, day: int, purchases: Purchases) -> Orders: delivery = _delivery(rng, len(purchases)) outcome, paid_after, cancelled_after = _fate(rng, len(purchases)) product, quantity, items_total = _delta(rng, purchases) + created_at = _created_at(rng, purchases) discount = _discount(purchases) return Orders( day=day, order_id=purchases.order_id, user_id=purchases.person_id, + created_at=created_at, product=product, quantity=quantity, 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]: """Стоимость доставки каждого заказа дня — броском по таблице весов.""" return _DELIVERY_PRICE[pick(rng, _DELIVERY_CUMULATIVE, orders)] @@ -147,6 +201,22 @@ def _fate( 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]: """Момент внутри окна: час по таблице весов и равномерная секунда в нём. diff --git a/generator/src/clickstream_generator/player.py b/generator/src/clickstream_generator/player.py index 7453cbb..61957c3 100644 --- a/generator/src/clickstream_generator/player.py +++ b/generator/src/clickstream_generator/player.py @@ -16,6 +16,11 @@ (раздел 5): это разные машины разной природы, и сложенные в одно число они перестают что-либо говорить. Сериализация считается частью генерации: она рождает те самые байты, которые сторожит опись. + +**Слепки заказов идут третьим ходом того же проигрывателя** и живут по правилу +сдвига: прогон дня D отправляет слепок дня D−1 — ночная выгрузка бэкенда за +вчера. Правило записано здесь одно и целиком, потому что оно и есть разница +между двумя источниками: трекер шлёт сегодняшний день, бэкенд — вчерашний. """ import logging @@ -26,6 +31,7 @@ import numpy as np from numpy.typing import NDArray from clickstream_generator import day as day_module +from clickstream_generator import orders as orders_module from clickstream_generator import serialize 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 отправляет слепок дня D−1: содержимое слепка — чистая функция + зерна и дня, от момента отправки оно не зависит, а сдвиг делает живой день + обычным. На старте оси вчера нет, поэтому прогон дня 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: """Пакетно: подряд и без пауз.""" for payload in payloads: diff --git a/generator/src/clickstream_generator/serialize.py b/generator/src/clickstream_generator/serialize.py index c2c893d..17dec5a 100644 --- a/generator/src/clickstream_generator/serialize.py +++ b/generator/src/clickstream_generator/serialize.py @@ -27,16 +27,25 @@ делает это одним вызовом на колонку, и на дне в полсотни тысяч событий разница заметна. Обратная сторона — день лежит в памяти дважды; проигрыватель поэтому и берёт его днями, а не горизонтом целиком. + +**Второй контракт провода — слепок заказов** (мастер-спека, раздел 2). Он не +похож на событие: одиннадцать ключей вместо сорока семи, деньги строками, а +не числами, времена с миллисекундами. Общее у них одно, зато главное: байты +рождаются здесь и только здесь. Запись слепка — один словарь с вложенным +списком и один `orjson.dumps`. """ +from collections.abc import Sequence +from datetime import timedelta from typing import Any import numpy as np import orjson 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.orders import Orders, at_boundary _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]: """Колонка питоновскими значениями — в той записи, в какой уедет на провод. diff --git a/generator/tests/conftest.py b/generator/tests/conftest.py new file mode 100644 index 0000000..097d561 --- /dev/null +++ b/generator/tests/conftest.py @@ -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) diff --git a/generator/tests/test_inventory.py b/generator/tests/test_inventory.py index 32f2e47..3e9b088 100644 --- a/generator/tests/test_inventory.py +++ b/generator/tests/test_inventory.py @@ -10,18 +10,53 @@ поднятом стенде, где расхождение счётчиков выглядит поломкой хранилища. """ +import hashlib import json 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" +# Прогон, слепок которого сверяется с описью байт в байт. День любой из +# отправляемых; этот дешевле прочих — его окно короче. +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")) - assert stored == build(), ( + assert stored == inventory, ( "опись мира отстала от кода — пересоберите: make inventory." " Разошлись хеши дней и каталога — правили 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 diff --git a/generator/tests/test_player.py b/generator/tests/test_player.py index 273ce11..9381e60 100644 --- a/generator/tests/test_player.py +++ b/generator/tests/test_player.py @@ -22,27 +22,6 @@ from clickstream_generator.sinks import FileSink 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): """Проиграть в файл и вернуть итог прогона.""" diff --git a/generator/tests/test_snapshot.py b/generator/tests/test_snapshot.py new file mode 100644 index 0000000..918b45c --- /dev/null +++ b/generator/tests/test_snapshot.py @@ -0,0 +1,337 @@ +"""Слепок заказов: окно, граница суток, запись на проводе и байты прогона. + +Слепок — второй артефакт мира, и сторожится он тем же, чем день событий: +формой записи и побайтовым повтором. Проверяются пять обещаний +(docs/architecture/orders/snapshot.md): + +1. **Окно.** Слепок дня D несёт заказы, рождённые в дни D−6…D, и у начала оси + усекается сам. +2. **Граница суток.** Состояние — чтение готовой судьбы на границе D|D+1: + моменты позже границы не учитываются, поэтому заказ «дышит» — `created` в + одном слепке, `paid` в следующем. +3. **Запись на проводе.** Одиннадцать ключей в порядке контракта, деньги + строками с двумя знаками, времена RFC 3339 с миллисекундами. +4. **Сдвиг отправки.** Прогон дня D отправляет слепок дня D−1; прогон дня 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, и окно усекается у начала оси. + + Сдвиг «ночная выгрузка за вчера» живёт в проигрывателе, поэтому и + спрашивается с него: прогоны дней 4…7 обязаны отправить слепки дней 3…6 — + теми же байтами, какие даёт сериализатор, по строке на заказ. + """ + 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()