fix(tradein/matching): listing_sources датируется построчно, а не стартом транзакции (#2731) #2743

Merged
bot-backend merged 1 commit from fix/2731-listing-sources-clock into main 2026-08-06 16:35:30 +00:00
2 changed files with 147 additions and 5 deletions

View file

@ -246,7 +246,23 @@ def _upsert_listing_source(
source_url: str | None,
source_data: dict | None,
) -> None:
"""Insert or refresh listing_sources row for this source+ext_id."""
"""Insert or refresh listing_sources row for this source+ext_id.
Отметки времени statement_timestamp(), НЕ NOW() (#2731). Этот upsert вызывается
ПОСТРОЧНО из save_listings (hook `_link_listing_to_house`), а транзакция batch'а
коммитится один раз в конце, поэтому NOW() (== transaction_timestamp) давал одну
метку на весь вызов: прод-замер 2026-08-06 219 строк на 1 метку в 11:00,
235/1 в 10:00, 297/1 в 09:00, и так каждый час.
Чинится вместе с listings.scraped_at/last_seen_at, а не отдельно: сегодня
listings.last_seen_at = listing_sources.last_seen_at у 100% пар (2407 из 2407 за
сутки) именно потому, что обе колонки берут одну транзакционную метку. Почини
только одну вторая осталась бы замороженной на старте batch'а, и расхождение
выросло бы с миллисекунд (честная разница двух записей) до длительности прогона.
Все три колонки пишутся ОДНИМ statement'ом, поэтому statement_timestamp() даёт им
одинаковое значение; clock_timestamp() развёл бы их на микросекунды.
"""
raw = json.dumps(source_data) if source_data is not None else None
db.execute(
text("""
@ -257,15 +273,15 @@ def _upsert_listing_source(
price_rub, area_m2, floor, rooms_count, raw_payload
) VALUES (
CAST(:lid AS bigint), :s, :e,
CAST(:c AS real), :m, NOW(), NOW(),
:url, NOW(),
CAST(:c AS real), :m, statement_timestamp(), statement_timestamp(),
:url, statement_timestamp(),
CAST(:p AS bigint), CAST(:a AS numeric), :fl, :rc,
CAST(:raw AS jsonb)
)
ON CONFLICT (ext_source, ext_id) DO UPDATE SET
confidence = GREATEST(EXCLUDED.confidence, listing_sources.confidence),
last_seen_at = NOW(),
last_scraped_at = NOW(),
last_seen_at = statement_timestamp(),
last_scraped_at = statement_timestamp(),
price_rub = COALESCE(EXCLUDED.price_rub, listing_sources.price_rub),
area_m2 = COALESCE(EXCLUDED.area_m2, listing_sources.area_m2),
floor = COALESCE(EXCLUDED.floor, listing_sources.floor),

View file

@ -0,0 +1,126 @@
"""#2731 (вторая половина): listing_sources датируется построчно, как и listings.
`_upsert_listing_source` вызывается ПОСТРОЧНО из save_listings (hook
`_link_listing_to_house`), а транзакция batch'а коммитится один раз в конце. С `NOW()`
(== `transaction_timestamp()`) все строки вызова получали одну метку прод-замер
2026-08-06 по часам: 219 строк / 1 метка, 235 / 1, 297 / 1, 150 / 1
Почему это чинится ВМЕСТЕ с listings, а не отдельно: сегодня
`listings.last_seen_at = listing_sources.last_seen_at` у 100% пар (прод: 2407 из 2407 за
сутки) ровно потому, что обе колонки берут одну транзакционную метку. Если починить
только listings, вторая колонка осталась бы замороженной на старте batch'а, и расхождение
выросло бы с миллисекунд (честная разница двух соседних записей) до длительности прогона.
Фальсификация: на старом коде (`NOW()`) оба теста дают одну метку на все строки.
"""
from __future__ import annotations
import os
import re
from typing import Any
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost/test_db")
from app.services.matching.listings import _upsert_listing_source
# Колонки времени этого писателя. Пишутся ОДНИМ statement'ом → обязаны совпадать.
_COLS = ("matched_at", "last_seen_at", "last_scraped_at")
_STEP_SEC = 0.6 # три строки → 1.8 с работы одной транзакции
class _FakePg:
"""Модель PostgreSQL на две функции времени (см. test_2731_observation_timestamps).
`now()` замерзает на старте транзакции, `statement_timestamp()` на старте запроса;
`clock_timestamp()` вычисляется на каждый вызов, поэтому три колонки одного
statement'а разъехались бы.
"""
def __init__(self) -> None:
self.wall = 0.0
self.tx_start: float | None = None
self.stmt_start = 0.0
self._clock_calls = 0
self.rows: list[dict[str, float]] = []
def _value(self, func: str) -> float:
if func == "now":
assert self.tx_start is not None
return self.tx_start
if func == "statement_timestamp":
return self.stmt_start
if func == "clock_timestamp":
self._clock_calls += 1
return self.stmt_start + self._clock_calls * 1e-6
raise AssertionError(f"неизвестная функция времени: {func}")
def execute(self, stmt: Any, params: Any = None) -> None:
if self.tx_start is None:
self.tx_start = self.wall
self.stmt_start = self.wall
self.wall += _STEP_SEC
sql = re.sub(r"--[^\n]*", "", str(stmt))
if "INSERT INTO listing_sources" not in sql:
return None
# VALUES: отметки позиционные; ON CONFLICT: именованные. Берём обе формы.
cols = re.search(
r"INSERT INTO listing_sources\s*\((.*?)\)\s*VALUES\s*\((.*?)\)\s*\n", sql, re.S
)
row: dict[str, float] = {}
if cols is not None:
names = [c.strip() for c in cols.group(1).replace("\n", " ").split(",")]
vals = [v.strip() for v in cols.group(2).replace("\n", " ").split(",")]
for name, val in zip(names, vals, strict=False):
f = re.fullmatch(r"(\w+)\s*\(\s*\)", val)
if name in _COLS and f is not None:
row[name] = self._value(f.group(1).lower())
for col in _COLS: # ON CONFLICT DO UPDATE SET <col> = <func>()
m = re.search(rf"\b{col}\s*=\s*(\w+)\s*\(\s*\)", sql)
if m is not None:
row.setdefault(col, self._value(m.group(1).lower()))
self.rows.append(row)
return None
def _write(db: _FakePg, n: int = 3) -> None:
for i in range(n):
_upsert_listing_source(
db,
listing_id=i,
ext_source="cian",
ext_id=str(i),
method="source_link",
confidence=1.0,
price_rub=5_000_000,
area_m2=42.0,
floor=3,
rooms_count=1,
source_url=f"https://ekb.cian.ru/sale/flat/{i}/",
source_data=None,
)
def test_listing_source_marks_differ_per_row() -> None:
"""Писатель, отработавший дольше секунды, оставляет РАЗЛИЧНЫЕ метки у строк."""
db = _FakePg()
_write(db)
marks = [row["last_seen_at"] for row in db.rows]
assert len(marks) == 3
assert max(marks) - min(marks) > 1.0, "писатель отработал дольше секунды"
assert len(set(marks)) == len(marks), "у строк одной транзакции обязаны быть разные метки"
def test_all_three_marks_of_one_statement_agree() -> None:
"""matched_at / last_seen_at / last_scraped_at пишутся одним statement'ом — и совпадают.
`clock_timestamp()` развёл бы их на микросекунды: модель БД это воспроизводит,
поэтому такая «починка» провалит тест.
"""
db = _FakePg()
_write(db)
for row in db.rows:
assert set(row) == set(_COLS), "все три отметки обязаны быть записаны"
assert len(set(row.values())) == 1, "отметки одного statement'а обязаны совпадать"