feat(tradein): yandex_newbuilding sweep-entrypoint + market-схема enrichment, seed dormant (#974) #1541
8 changed files with 1291 additions and 19 deletions
|
|
@ -1468,8 +1468,9 @@ async def scrape_yandex_newbuilding(
|
||||||
|
|
||||||
URL: /{city}/kupit/novostrojka/<slug>-<id>/
|
URL: /{city}/kupit/novostrojka/<slug>-<id>/
|
||||||
"""
|
"""
|
||||||
async with YandexNewbuildingScraper() as scraper:
|
# fetch_jk использует внутренний BrowserFetcher, httpx-клиент BaseScraper не нужен
|
||||||
result = await scraper.fetch_jk(jk_slug=slug, jk_id=id, city=city)
|
scraper = YandexNewbuildingScraper()
|
||||||
|
result = await scraper.fetch_jk(jk_slug=slug, jk_id=id, city=city)
|
||||||
if result is None:
|
if result is None:
|
||||||
raise HTTPException(404, f"Could not parse Yandex JK: {slug}-{id} in {city}")
|
raise HTTPException(404, f"Could not parse Yandex JK: {slug}-{id} in {city}")
|
||||||
return YandexNewbuildingTriggerResp(
|
return YandexNewbuildingTriggerResp(
|
||||||
|
|
@ -1486,6 +1487,91 @@ async def scrape_yandex_newbuilding(
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
class YandexNewbuildingSweepStartRequest(BaseModel):
|
||||||
|
limit: int = Field(default=5, ge=1, le=200)
|
||||||
|
request_delay_sec: float = Field(default=8.0, ge=3.0, le=30.0)
|
||||||
|
force: bool = False
|
||||||
|
city: str = Field(default="ekaterinburg", min_length=1, max_length=64)
|
||||||
|
|
||||||
|
|
||||||
|
class YandexNewbuildingSweepStartResponse(BaseModel):
|
||||||
|
run_id: int
|
||||||
|
status: str
|
||||||
|
limit: int
|
||||||
|
detail: str
|
||||||
|
|
||||||
|
|
||||||
|
@router.post("/scrape/yandex-newbuilding-sweep", response_model=YandexNewbuildingSweepStartResponse)
|
||||||
|
async def start_yandex_newbuilding_sweep(
|
||||||
|
payload: YandexNewbuildingSweepStartRequest,
|
||||||
|
background_tasks: BackgroundTasks,
|
||||||
|
db: Annotated[Session, Depends(get_db)],
|
||||||
|
) -> YandexNewbuildingSweepStartResponse:
|
||||||
|
"""Запустить Yandex newbuilding enrichment sweep в background (#974). Returns run_id.
|
||||||
|
|
||||||
|
Один прогон: SELECT pending yandex_realty_nb houses → resolve slug (если NULL) →
|
||||||
|
fetch_jk через BrowserFetcher → UPSERT market.yandex_jk_enrichment.
|
||||||
|
|
||||||
|
Shipped DORMANT (scheduler seed enabled=false). Этот endpoint — для ручного запуска.
|
||||||
|
Coop cancel через DELETE scrape_runs или cancel endpoint.
|
||||||
|
"""
|
||||||
|
from app.tasks.yandex_newbuilding_sweep import enrich_yandex_newbuilding_sweep
|
||||||
|
|
||||||
|
run_id = runs_mod.create_run(db, source="yandex_newbuilding_sweep", params=payload.model_dump())
|
||||||
|
|
||||||
|
async def _sweep_task() -> None:
|
||||||
|
sweep_db = SessionLocal()
|
||||||
|
try:
|
||||||
|
result = await enrich_yandex_newbuilding_sweep(
|
||||||
|
sweep_db,
|
||||||
|
limit=payload.limit,
|
||||||
|
force=payload.force,
|
||||||
|
request_delay_sec=payload.request_delay_sec,
|
||||||
|
city=payload.city,
|
||||||
|
)
|
||||||
|
runs_mod.mark_done(sweep_db, run_id, result.to_dict())
|
||||||
|
except Exception:
|
||||||
|
logger.exception("yandex-nb-sweep background task run_id=%d crashed", run_id)
|
||||||
|
try:
|
||||||
|
runs_mod.mark_failed(sweep_db, run_id, "crashed", {})
|
||||||
|
except Exception:
|
||||||
|
pass
|
||||||
|
finally:
|
||||||
|
sweep_db.close()
|
||||||
|
|
||||||
|
background_tasks.add_task(_sweep_task)
|
||||||
|
logger.info("yandex-nb-sweep queued run_id=%d params=%s", run_id, payload.model_dump())
|
||||||
|
return YandexNewbuildingSweepStartResponse(
|
||||||
|
run_id=run_id,
|
||||||
|
status="running",
|
||||||
|
limit=payload.limit,
|
||||||
|
detail=f"yandex_newbuilding_sweep started: limit={payload.limit} city={payload.city}",
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
@router.get("/scrape/yandex-newbuilding-sweep/runs", response_model=list[ScrapeRunRow])
|
||||||
|
def list_yandex_newbuilding_sweep_runs(
|
||||||
|
db: Annotated[Session, Depends(get_db)],
|
||||||
|
limit: int = 10,
|
||||||
|
) -> list[ScrapeRunRow]:
|
||||||
|
"""Список последних N yandex newbuilding sweep runs (для UI polling). Default limit=10."""
|
||||||
|
rows = runs_mod.list_recent(db, source="yandex_newbuilding_sweep", limit=limit)
|
||||||
|
return [
|
||||||
|
ScrapeRunRow(
|
||||||
|
run_id=r["run_id"],
|
||||||
|
source=r["source"],
|
||||||
|
status=r["status"],
|
||||||
|
params=r.get("params"),
|
||||||
|
counters=r.get("counters"),
|
||||||
|
error=r.get("error"),
|
||||||
|
started_at=r["started_at"].isoformat() if r.get("started_at") else None,
|
||||||
|
finished_at=r["finished_at"].isoformat() if r.get("finished_at") else None,
|
||||||
|
heartbeat_at=r["heartbeat_at"].isoformat() if r.get("heartbeat_at") else None,
|
||||||
|
)
|
||||||
|
for r in rows
|
||||||
|
]
|
||||||
|
|
||||||
|
|
||||||
class YandexValuationTriggerResp(BaseModel):
|
class YandexValuationTriggerResp(BaseModel):
|
||||||
ok: bool
|
ok: bool
|
||||||
address: str
|
address: str
|
||||||
|
|
|
||||||
|
|
@ -45,6 +45,10 @@ Sources:
|
||||||
Idempotent: skips houses already enriched, per-house SAVEPOINT, so
|
Idempotent: skips houses already enriched, per-house SAVEPOINT, so
|
||||||
each fire drains the next `limit` pending houses; window 00:00-01:00
|
each fire drains the next `limit` pending houses; window 00:00-01:00
|
||||||
UTC = 03:00-04:00 МСК, before the 01:00+ UTC sweep block)
|
UTC = 03:00-04:00 МСК, before the 01:00+ UTC sweep block)
|
||||||
|
- yandex_newbuilding_sweep → enrich_yandex_newbuilding_sweep
|
||||||
|
(tasks/yandex_newbuilding_sweep.py, #974; enrichment ЖК в
|
||||||
|
market.yandex_jk_enrichment через BrowserFetcher; shipped DORMANT —
|
||||||
|
enable manually after deploy verified; window 02:00-05:00 UTC)
|
||||||
"""
|
"""
|
||||||
|
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
@ -420,8 +424,7 @@ async def _execute_cian_backfill(
|
||||||
}
|
}
|
||||||
runs_mod.mark_done(db, run_id, counters)
|
runs_mod.mark_done(db, run_id, counters)
|
||||||
logger.info(
|
logger.info(
|
||||||
"scheduler: cian_history_backfill run_id=%d done — "
|
"scheduler: cian_history_backfill run_id=%d done — listings=%d/%d houses=%d/%d %.1fs",
|
||||||
"listings=%d/%d houses=%d/%d %.1fs",
|
|
||||||
run_id,
|
run_id,
|
||||||
result.listings_succeeded,
|
result.listings_succeeded,
|
||||||
result.listings_total,
|
result.listings_total,
|
||||||
|
|
@ -730,6 +733,55 @@ async def trigger_newbuilding_enrich_run(db: Session, schedule_row: dict[str, An
|
||||||
return run_id
|
return run_id
|
||||||
|
|
||||||
|
|
||||||
|
async def trigger_yandex_newbuilding_sweep_run(
|
||||||
|
db: Session, schedule_row: dict[str, Any]
|
||||||
|
) -> int | None:
|
||||||
|
"""Создать scrape_runs + launch enrich_yandex_newbuilding_sweep в asyncio.create_task (#974).
|
||||||
|
|
||||||
|
Nightly enrichment sweep: SELECT pending yandex_realty_nb houses → resolve slug
|
||||||
|
(BrowserFetcher SERP) → fetch_jk (BrowserFetcher) → UPSERT market.yandex_jk_enrichment.
|
||||||
|
|
||||||
|
Shipped DORMANT (seed 106, enabled=false). Включается оператором вручную:
|
||||||
|
UPDATE scrape_schedules SET enabled = true WHERE source = 'yandex_newbuilding_sweep';
|
||||||
|
|
||||||
|
Mirrors trigger_newbuilding_enrich_run: claim run → create_task → mark_done/failed
|
||||||
|
делегируется tasks/yandex_newbuilding_sweep.enrich_yandex_newbuilding_sweep.
|
||||||
|
|
||||||
|
Returns run_id или None (skip — already running).
|
||||||
|
"""
|
||||||
|
run_id = _claim_run(db, schedule_row)
|
||||||
|
if run_id is None:
|
||||||
|
return None
|
||||||
|
|
||||||
|
params = schedule_row.get("default_params") or {}
|
||||||
|
|
||||||
|
async def _run() -> None:
|
||||||
|
run_db = SessionLocal()
|
||||||
|
try:
|
||||||
|
from app.tasks.yandex_newbuilding_sweep import enrich_yandex_newbuilding_sweep
|
||||||
|
|
||||||
|
result = await enrich_yandex_newbuilding_sweep(
|
||||||
|
run_db,
|
||||||
|
limit=int(params.get("limit", 5)),
|
||||||
|
request_delay_sec=float(params.get("request_delay_sec", 8.0)),
|
||||||
|
city=str(params.get("city", "ekaterinburg")),
|
||||||
|
)
|
||||||
|
runs_mod.mark_done(run_db, run_id, result.to_dict())
|
||||||
|
except Exception:
|
||||||
|
logger.exception("scheduler: enrich_yandex_newbuilding_sweep crashed run_id=%d", run_id)
|
||||||
|
try:
|
||||||
|
runs_mod.mark_failed(run_db, run_id, "crashed", {})
|
||||||
|
except Exception:
|
||||||
|
pass
|
||||||
|
finally:
|
||||||
|
run_db.close()
|
||||||
|
|
||||||
|
task = asyncio.create_task(_run())
|
||||||
|
task.add_done_callback(lambda t: t.exception() if not t.cancelled() else None)
|
||||||
|
logger.info("scheduler: triggered yandex_newbuilding_sweep run_id=%d", run_id)
|
||||||
|
return run_id
|
||||||
|
|
||||||
|
|
||||||
async def trigger_asking_to_sold_ratio_run(db: Session, schedule_row: dict[str, Any]) -> int | None:
|
async def trigger_asking_to_sold_ratio_run(db: Session, schedule_row: dict[str, Any]) -> int | None:
|
||||||
"""Создать scrape_runs + launch recompute_asking_to_sold_ratios в executor (sync DB-only task).
|
"""Создать scrape_runs + launch recompute_asking_to_sold_ratios в executor (sync DB-only task).
|
||||||
|
|
||||||
|
|
@ -1056,6 +1108,8 @@ async def scheduler_loop() -> None:
|
||||||
await trigger_rosreestr_quarter_poll_run(db, sch)
|
await trigger_rosreestr_quarter_poll_run(db, sch)
|
||||||
elif source == "newbuilding_enrich":
|
elif source == "newbuilding_enrich":
|
||||||
await trigger_newbuilding_enrich_run(db, sch)
|
await trigger_newbuilding_enrich_run(db, sch)
|
||||||
|
elif source == "yandex_newbuilding_sweep":
|
||||||
|
await trigger_yandex_newbuilding_sweep_run(db, sch)
|
||||||
else:
|
else:
|
||||||
logger.warning("scheduler: unknown source=%s, skip", source)
|
logger.warning("scheduler: unknown source=%s, skip", source)
|
||||||
finally:
|
finally:
|
||||||
|
|
|
||||||
|
|
@ -3,6 +3,12 @@
|
||||||
URL pattern: /{city}/kupit/novostrojka/<slug>-<id>/
|
URL pattern: /{city}/kupit/novostrojka/<slug>-<id>/
|
||||||
Reference target: ЖК Татлин (id=1592987, slug=tatlin) — comfort+, June 2023,
|
Reference target: ЖК Татлин (id=1592987, slug=tatlin) — comfort+, June 2023,
|
||||||
PRINZIP, rating 4.3, 1505 ratings, 353 text reviews, coords (56.855312, 60.576668).
|
PRINZIP, rating 4.3, 1505 ratings, 353 text reviews, coords (56.855312, 60.576668).
|
||||||
|
|
||||||
|
Fetch strategy (#974):
|
||||||
|
- Все network-запросы (ЖК-лендинг + SERP slug-resolve) идут через BrowserFetcher
|
||||||
|
(tradein-browser camoufox), а НЕ через httpx/_http_get.
|
||||||
|
- Yandex Realty — JS-heavy / anti-bot; curl_cffi / httpx не получают данные.
|
||||||
|
- BrowserFetcher.fetch(url) → str (полный HTML); вызывающий код парсит через parse().
|
||||||
"""
|
"""
|
||||||
|
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
@ -98,6 +104,12 @@ RE_METRO_INLINE = re.compile(
|
||||||
r"([А-ЯЁ][А-Яа-яё\s-]{2,30}?)\s+(\d+)\s*мин",
|
r"([А-ЯЁ][А-Яа-яё\s-]{2,30}?)\s+(\d+)\s*мин",
|
||||||
)
|
)
|
||||||
|
|
||||||
|
# Slug из href: /{city}/kupit/novostrojka/<slug>-<id>/
|
||||||
|
# Используется в resolve_yandex_jk_slug для извлечения slug из SERP href.
|
||||||
|
_JK_SLUG_RE = re.compile(
|
||||||
|
r"/(?:ekaterinburg|moskva|spb|[a-z-]+)/kupit/novostrojka/([a-z0-9-]+)-(\d+)/"
|
||||||
|
)
|
||||||
|
|
||||||
# Ekb-specific coord ranges (extend per-city later)
|
# Ekb-specific coord ranges (extend per-city later)
|
||||||
LAT_RANGE = (55.5, 57.5)
|
LAT_RANGE = (55.5, 57.5)
|
||||||
LON_RANGE = (59.5, 61.5)
|
LON_RANGE = (59.5, 61.5)
|
||||||
|
|
@ -129,9 +141,7 @@ class YandexNewbuildingScraper(BaseScraper):
|
||||||
super().__init__()
|
super().__init__()
|
||||||
self.request_delay_sec = get_scraper_delay(self.name)
|
self.request_delay_sec = get_scraper_delay(self.name)
|
||||||
|
|
||||||
async def fetch_around(
|
async def fetch_around(self, lat: float, lon: float, radius_m: int = 1000) -> list: # type: ignore[override]
|
||||||
self, lat: float, lon: float, radius_m: int = 1000
|
|
||||||
) -> list: # type: ignore[override]
|
|
||||||
raise NotImplementedError(
|
raise NotImplementedError(
|
||||||
"YandexNewbuildingScraper is JK-slug-based; use fetch_jk(slug, id) instead."
|
"YandexNewbuildingScraper is JK-slug-based; use fetch_jk(slug, id) instead."
|
||||||
)
|
)
|
||||||
|
|
@ -139,22 +149,31 @@ class YandexNewbuildingScraper(BaseScraper):
|
||||||
async def fetch_jk(
|
async def fetch_jk(
|
||||||
self, jk_slug: str, jk_id: str, city: str = "ekaterinburg"
|
self, jk_slug: str, jk_id: str, city: str = "ekaterinburg"
|
||||||
) -> YandexNewbuildingInfo | None:
|
) -> YandexNewbuildingInfo | None:
|
||||||
|
"""Загрузить ЖК-лендинг через BrowserFetcher и распарсить.
|
||||||
|
|
||||||
|
Yandex Realty — JS/anti-bot: curl/httpx не получают данные. Запрос идёт
|
||||||
|
через tradein-browser (camoufox) контейнер — единственный рабочий путь (#974).
|
||||||
|
"""
|
||||||
|
from app.services.scrapers.browser_fetcher import BrowserFetcher
|
||||||
|
|
||||||
url = f"{self.base_url}/{city}/kupit/novostrojka/{jk_slug}-{jk_id}/"
|
url = f"{self.base_url}/{city}/kupit/novostrojka/{jk_slug}-{jk_id}/"
|
||||||
try:
|
try:
|
||||||
response = await self._http_get(url)
|
async with BrowserFetcher() as fetcher:
|
||||||
|
html = await fetcher.fetch(url)
|
||||||
except Exception:
|
except Exception:
|
||||||
logger.exception("yandex nb fetch failed: %s", url)
|
logger.exception("yandex nb browser fetch failed: %s", url)
|
||||||
return None
|
return None
|
||||||
if response.status_code != 200:
|
if not html or len(html) < 500:
|
||||||
logger.warning("yandex nb returned %d for %s", response.status_code, url)
|
logger.warning(
|
||||||
|
"yandex nb browser returned empty/tiny HTML (%d bytes): %s",
|
||||||
|
len(html) if html else 0,
|
||||||
|
url,
|
||||||
|
)
|
||||||
return None
|
return None
|
||||||
result = self.parse(response.text, jk_slug=jk_slug, jk_id=jk_id, source_url=url)
|
result = self.parse(html, jk_slug=jk_slug, jk_id=jk_id, source_url=url)
|
||||||
await self.sleep_between_requests()
|
|
||||||
return result
|
return result
|
||||||
|
|
||||||
def parse(
|
def parse(self, html: str, jk_slug: str, jk_id: str, source_url: str) -> YandexNewbuildingInfo:
|
||||||
self, html: str, jk_slug: str, jk_id: str, source_url: str
|
|
||||||
) -> YandexNewbuildingInfo:
|
|
||||||
tree = HTMLParser(html)
|
tree = HTMLParser(html)
|
||||||
body = tree.body
|
body = tree.body
|
||||||
body_text = body.text(strip=True) if body else ""
|
body_text = body.text(strip=True) if body else ""
|
||||||
|
|
@ -208,9 +227,7 @@ class YandexNewbuildingScraper(BaseScraper):
|
||||||
href = dev_link.attributes.get("href", "")
|
href = dev_link.attributes.get("href", "")
|
||||||
if href:
|
if href:
|
||||||
developer_url = (
|
developer_url = (
|
||||||
href
|
href if href.startswith("http") else f"https://realty.yandex.ru{href}"
|
||||||
if href.startswith("http")
|
|
||||||
else f"https://realty.yandex.ru{href}"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
# Developer's other JKs
|
# Developer's other JKs
|
||||||
|
|
@ -269,6 +286,60 @@ class YandexNewbuildingScraper(BaseScraper):
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
# ── Slug resolution via SERP ──────────────────────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
async def resolve_yandex_jk_slug(
|
||||||
|
jk_id: str,
|
||||||
|
city: str = "ekaterinburg",
|
||||||
|
) -> str | None:
|
||||||
|
"""Найти Yandex Realty slug для ЖК по его ext_id (jk_id) через SERP.
|
||||||
|
|
||||||
|
Стратегия (#974 — зеркало resolve_cian_zhk_url_via_search):
|
||||||
|
1. Запросить поисковую страницу Yandex Realty через BrowserFetcher.
|
||||||
|
URL: /ekaterinburg/kupit/novostrojka/?siteId=<jk_id>
|
||||||
|
2. В HTML найти первую ссылку вида /<city>/kupit/novostrojka/<slug>-<jk_id>/
|
||||||
|
через regex _JK_SLUG_RE.
|
||||||
|
3. Вернуть slug или None при любой ошибке.
|
||||||
|
|
||||||
|
Caller несёт ответственность за anti-bot sleep (зеркало cian_newbuilding.py).
|
||||||
|
BrowserFetcher обязателен — Yandex Realty JS/anti-bot, httpx/curl не работают.
|
||||||
|
|
||||||
|
Returns:
|
||||||
|
slug (str без id-суффикса), или None при ошибке / не найден.
|
||||||
|
"""
|
||||||
|
from app.services.scrapers.browser_fetcher import BrowserFetcher
|
||||||
|
|
||||||
|
# Yandex Realty SERP: фильтр по siteId → первый результат = нужный ЖК.
|
||||||
|
# Альтернативный путь через Яндекс Поиск (web SERP) менее надёжен из-за
|
||||||
|
# вариативности разметки. Прямой realty.yandex.ru SERP — стабильнее.
|
||||||
|
serp_url = f"https://realty.yandex.ru/{city}/kupit/novostrojka/?siteId={jk_id}"
|
||||||
|
try:
|
||||||
|
async with BrowserFetcher() as fetcher:
|
||||||
|
html = await fetcher.fetch(serp_url)
|
||||||
|
except Exception as exc:
|
||||||
|
logger.warning("resolve_yandex_jk_slug jk_id=%s browser fetch failed: %s", jk_id, exc)
|
||||||
|
return None
|
||||||
|
|
||||||
|
if not html:
|
||||||
|
logger.warning("resolve_yandex_jk_slug jk_id=%s: empty HTML from browser", jk_id)
|
||||||
|
return None
|
||||||
|
|
||||||
|
# Ищем ссылку вида /{city}/kupit/novostrojka/<slug>-<id>/
|
||||||
|
# Ограничиваем: id в ссылке должен совпадать с искомым jk_id.
|
||||||
|
for m in _JK_SLUG_RE.finditer(html):
|
||||||
|
if m.group(2) == str(jk_id):
|
||||||
|
slug = m.group(1)
|
||||||
|
logger.info("resolve_yandex_jk_slug jk_id=%s → slug=%s", jk_id, slug)
|
||||||
|
return slug
|
||||||
|
|
||||||
|
logger.warning(
|
||||||
|
"resolve_yandex_jk_slug jk_id=%s: no matching slug in SERP HTML (markup drift?)",
|
||||||
|
jk_id,
|
||||||
|
)
|
||||||
|
return None
|
||||||
|
|
||||||
|
|
||||||
# ── helpers ───────────────────────────────────────────────────────────────────
|
# ── helpers ───────────────────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
|
@ -354,4 +425,5 @@ __all__ = [
|
||||||
"JKMetroStation",
|
"JKMetroStation",
|
||||||
"YandexNewbuildingInfo",
|
"YandexNewbuildingInfo",
|
||||||
"YandexNewbuildingScraper",
|
"YandexNewbuildingScraper",
|
||||||
|
"resolve_yandex_jk_slug",
|
||||||
]
|
]
|
||||||
|
|
|
||||||
423
tradein-mvp/backend/app/tasks/yandex_newbuilding_sweep.py
Normal file
423
tradein-mvp/backend/app/tasks/yandex_newbuilding_sweep.py
Normal file
|
|
@ -0,0 +1,423 @@
|
||||||
|
"""Yandex Newbuilding sweep-task (#974): enrichment ЖК в market.yandex_jk_enrichment.
|
||||||
|
|
||||||
|
Context
|
||||||
|
-------
|
||||||
|
309 house_sources строк с ext_source='yandex_realty_nb' имеют ext_id (= yandex_jk_id),
|
||||||
|
но у большинства NULL yandex_jk_slug и не заполнены данные в market.yandex_jk_enrichment.
|
||||||
|
|
||||||
|
Для каждого дома цепочка:
|
||||||
|
1. resolve_yandex_jk_slug(jk_id) → slug (через BrowserFetcher SERP)
|
||||||
|
Результат сохраняется в houses.yandex_jk_slug (SAVEPOINT) — resumable.
|
||||||
|
2. YandexNewbuildingScraper.fetch_jk(slug, jk_id) → YandexNewbuildingInfo
|
||||||
|
Через BrowserFetcher (tradein-browser camoufox) — единственный рабочий путь.
|
||||||
|
3. UPSERT в market.yandex_jk_enrichment (ON CONFLICT (ext_id) DO UPDATE).
|
||||||
|
4. UPDATE houses.yandex_jk_id WHERE yandex_jk_slug = slug.
|
||||||
|
|
||||||
|
Idempotency
|
||||||
|
-----------
|
||||||
|
- force=False: пропускает дома у которых уже есть строка в market.yandex_jk_enrichment.
|
||||||
|
- UPSERT через ON CONFLICT (ext_id) DO UPDATE — безопасен при повторном запуске.
|
||||||
|
- SAVEPOINT per house — один сбойный fetch не прерывает батч.
|
||||||
|
- dry_run: подсчёт популяции без fetch.
|
||||||
|
|
||||||
|
Anti-bot / resilience
|
||||||
|
---------------------
|
||||||
|
- request_delay_sec (default из get_scraper_delay('yandex_realty_nb')) с ±20% jitter.
|
||||||
|
- SAVEPOINT per house: resolve-фаза и enrich-фаза — отдельные savepoint'ы.
|
||||||
|
|
||||||
|
Execution
|
||||||
|
---------
|
||||||
|
- Чистый async callable — NO Celery. Триггер: admin endpoint или scheduler.
|
||||||
|
- psycopg v3 conventions: CAST(:x AS type) в SQL (никогда ::`); logger (никогда print).
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import asyncio
|
||||||
|
import json
|
||||||
|
import logging
|
||||||
|
import random
|
||||||
|
import time
|
||||||
|
from dataclasses import dataclass, field, fields
|
||||||
|
|
||||||
|
from sqlalchemy import text
|
||||||
|
from sqlalchemy.orm import Session
|
||||||
|
|
||||||
|
from app.services.scraper_settings import get_scraper_delay
|
||||||
|
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
__all__ = [
|
||||||
|
"YandexNewbuildingSweepResult",
|
||||||
|
"count_yandex_newbuilding_houses",
|
||||||
|
"enrich_yandex_newbuilding_sweep",
|
||||||
|
]
|
||||||
|
|
||||||
|
_ENRICHMENT_SOURCE = "yandex_realty_nb"
|
||||||
|
|
||||||
|
|
||||||
|
@dataclass
|
||||||
|
class YandexNewbuildingSweepResult:
|
||||||
|
"""Per-run счётчики yandex newbuilding sweep."""
|
||||||
|
|
||||||
|
# Sizing (независимо от limit).
|
||||||
|
total: int = 0 # дома с ext_source='yandex_realty_nb'
|
||||||
|
fetchable: int = 0 # из них: с yandex_jk_slug ИЛИ с ext_id
|
||||||
|
pending: int = 0 # fetchable И ещё не обогащены (force=False)
|
||||||
|
|
||||||
|
# Processing (ограничен limit).
|
||||||
|
processed: int = 0
|
||||||
|
skipped_already_enriched: int = 0
|
||||||
|
succeeded: int = 0
|
||||||
|
resolved_slug: int = 0 # ext_id → slug разрезолвлен + сохранён
|
||||||
|
failed_resolve: int = 0 # slug не удалось разрезолвить
|
||||||
|
failed_fetch: int = 0 # fetch_jk вернул None / упал
|
||||||
|
rows_inserted: int = 0 # строк в market.yandex_jk_enrichment (новых/обновлённых)
|
||||||
|
|
||||||
|
duration_sec: float = field(default=0.0)
|
||||||
|
|
||||||
|
def to_dict(self) -> dict[str, int | float]:
|
||||||
|
return {f.name: getattr(self, f.name) for f in fields(self)}
|
||||||
|
|
||||||
|
|
||||||
|
# ── SQL ────────────────────────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
_COUNT_TOTAL = """
|
||||||
|
SELECT COUNT(DISTINCT h.id)
|
||||||
|
FROM houses h
|
||||||
|
JOIN house_sources hs ON hs.house_id = h.id
|
||||||
|
WHERE hs.ext_source = 'yandex_realty_nb'
|
||||||
|
"""
|
||||||
|
|
||||||
|
_COUNT_FETCHABLE = """
|
||||||
|
SELECT COUNT(DISTINCT h.id)
|
||||||
|
FROM houses h
|
||||||
|
JOIN house_sources hs ON hs.house_id = h.id
|
||||||
|
WHERE hs.ext_source = 'yandex_realty_nb'
|
||||||
|
AND (h.yandex_jk_slug IS NOT NULL OR hs.ext_id IS NOT NULL)
|
||||||
|
"""
|
||||||
|
|
||||||
|
_COUNT_PENDING = """
|
||||||
|
SELECT COUNT(DISTINCT h.id)
|
||||||
|
FROM houses h
|
||||||
|
JOIN house_sources hs ON hs.house_id = h.id
|
||||||
|
WHERE hs.ext_source = 'yandex_realty_nb'
|
||||||
|
AND (h.yandex_jk_slug IS NOT NULL OR hs.ext_id IS NOT NULL)
|
||||||
|
AND NOT EXISTS (
|
||||||
|
SELECT 1 FROM market.yandex_jk_enrichment e WHERE e.ext_id = hs.ext_id
|
||||||
|
)
|
||||||
|
"""
|
||||||
|
|
||||||
|
# DISTINCT ON (h.id) — у дома может быть >1 house_sources строки; берём любую ext_id.
|
||||||
|
_SELECT_PENDING_HOUSES = """
|
||||||
|
SELECT DISTINCT ON (h.id)
|
||||||
|
h.id AS house_id,
|
||||||
|
h.yandex_jk_slug,
|
||||||
|
h.yandex_jk_id,
|
||||||
|
hs.ext_id
|
||||||
|
FROM houses h
|
||||||
|
JOIN house_sources hs ON hs.house_id = h.id
|
||||||
|
WHERE hs.ext_source = 'yandex_realty_nb'
|
||||||
|
AND (h.yandex_jk_slug IS NOT NULL OR hs.ext_id IS NOT NULL)
|
||||||
|
AND (
|
||||||
|
CAST(:force AS boolean) = TRUE
|
||||||
|
OR NOT EXISTS (
|
||||||
|
SELECT 1 FROM market.yandex_jk_enrichment e WHERE e.ext_id = hs.ext_id
|
||||||
|
)
|
||||||
|
)
|
||||||
|
ORDER BY h.id, hs.ext_id NULLS LAST
|
||||||
|
LIMIT :lim
|
||||||
|
"""
|
||||||
|
|
||||||
|
_UPDATE_SLUG = """
|
||||||
|
UPDATE houses
|
||||||
|
SET yandex_jk_slug = CAST(:slug AS text)
|
||||||
|
WHERE id = CAST(:hid AS bigint)
|
||||||
|
"""
|
||||||
|
|
||||||
|
_UPSERT_ENRICHMENT = """
|
||||||
|
INSERT INTO market.yandex_jk_enrichment (
|
||||||
|
ext_id, name, developer_name, address,
|
||||||
|
lat, lon, rating, ratings_count, text_reviews_count,
|
||||||
|
raw_payload, updated_at
|
||||||
|
) VALUES (
|
||||||
|
CAST(:ext_id AS text),
|
||||||
|
CAST(:name AS text),
|
||||||
|
CAST(:developer_name AS text),
|
||||||
|
CAST(:address AS text),
|
||||||
|
CAST(:lat AS float8),
|
||||||
|
CAST(:lon AS float8),
|
||||||
|
CAST(:rating AS float4),
|
||||||
|
CAST(:ratings_count AS int),
|
||||||
|
CAST(:text_reviews_count AS int),
|
||||||
|
CAST(:raw_payload AS jsonb),
|
||||||
|
NOW()
|
||||||
|
)
|
||||||
|
ON CONFLICT (ext_id) DO UPDATE SET
|
||||||
|
name = EXCLUDED.name,
|
||||||
|
developer_name = EXCLUDED.developer_name,
|
||||||
|
address = EXCLUDED.address,
|
||||||
|
lat = EXCLUDED.lat,
|
||||||
|
lon = EXCLUDED.lon,
|
||||||
|
rating = EXCLUDED.rating,
|
||||||
|
ratings_count = EXCLUDED.ratings_count,
|
||||||
|
text_reviews_count = EXCLUDED.text_reviews_count,
|
||||||
|
raw_payload = EXCLUDED.raw_payload,
|
||||||
|
updated_at = NOW()
|
||||||
|
"""
|
||||||
|
|
||||||
|
|
||||||
|
def count_yandex_newbuilding_houses(db: Session) -> dict[str, int]:
|
||||||
|
"""Sizing: сколько yandex_realty_nb домов существует / fetchable / pending.
|
||||||
|
|
||||||
|
Дешёвый (3 COUNT) — безопасен для dry_run.
|
||||||
|
"""
|
||||||
|
total = int(db.execute(text(_COUNT_TOTAL)).scalar_one())
|
||||||
|
fetchable = int(db.execute(text(_COUNT_FETCHABLE)).scalar_one())
|
||||||
|
pending = int(db.execute(text(_COUNT_PENDING)).scalar_one())
|
||||||
|
return {"total": total, "fetchable": fetchable, "pending": pending}
|
||||||
|
|
||||||
|
|
||||||
|
async def enrich_yandex_newbuilding_sweep(
|
||||||
|
db: Session,
|
||||||
|
*,
|
||||||
|
limit: int = 5,
|
||||||
|
force: bool = False,
|
||||||
|
request_delay_sec: float | None = None,
|
||||||
|
dry_run: bool = False,
|
||||||
|
city: str = "ekaterinburg",
|
||||||
|
) -> YandexNewbuildingSweepResult:
|
||||||
|
"""Enrichment sweep для yandex_realty_nb домов → market.yandex_jk_enrichment.
|
||||||
|
|
||||||
|
Per house:
|
||||||
|
- resolve slug если NULL (persist в houses.yandex_jk_slug, SAVEPOINT)
|
||||||
|
- fetch_jk через BrowserFetcher
|
||||||
|
- UPSERT market.yandex_jk_enrichment (ON CONFLICT ext_id)
|
||||||
|
- UPDATE houses.yandex_jk_id (если изменился)
|
||||||
|
|
||||||
|
Args:
|
||||||
|
db: tradein-mvp SQLAlchemy session (НЕ gendesign DB).
|
||||||
|
limit: max домов за один прогон. Default 5 — bounded start.
|
||||||
|
force: переобработать уже обогащённые (UPSERT-safe). Default False.
|
||||||
|
request_delay_sec: пауза между домами. None → get_scraper_delay('yandex_realty_nb').
|
||||||
|
dry_run: считать популяцию без fetch/write.
|
||||||
|
city: город для Yandex Realty URL (ekaterinburg).
|
||||||
|
|
||||||
|
Returns:
|
||||||
|
YandexNewbuildingSweepResult со счётчиками.
|
||||||
|
"""
|
||||||
|
from app.services.scrapers.yandex_newbuilding import (
|
||||||
|
YandexNewbuildingScraper,
|
||||||
|
resolve_yandex_jk_slug,
|
||||||
|
)
|
||||||
|
|
||||||
|
result = YandexNewbuildingSweepResult()
|
||||||
|
t0 = time.time()
|
||||||
|
delay = (
|
||||||
|
request_delay_sec
|
||||||
|
if request_delay_sec is not None
|
||||||
|
else get_scraper_delay(_ENRICHMENT_SOURCE)
|
||||||
|
)
|
||||||
|
|
||||||
|
# ── Sizing ──────────────────────────────────────────────────────────────
|
||||||
|
sizing = count_yandex_newbuilding_houses(db)
|
||||||
|
result.total = sizing["total"]
|
||||||
|
result.fetchable = sizing["fetchable"]
|
||||||
|
result.pending = sizing["pending"]
|
||||||
|
|
||||||
|
rows = db.execute(text(_SELECT_PENDING_HOUSES), {"force": force, "lim": limit}).mappings().all()
|
||||||
|
|
||||||
|
logger.info(
|
||||||
|
"yandex-nb-sweep: total=%d fetchable=%d pending=%d; "
|
||||||
|
"selected=%d (limit=%d force=%s dry_run=%s delay=%.1fs city=%s)",
|
||||||
|
result.total,
|
||||||
|
result.fetchable,
|
||||||
|
result.pending,
|
||||||
|
len(rows),
|
||||||
|
limit,
|
||||||
|
force,
|
||||||
|
dry_run,
|
||||||
|
delay,
|
||||||
|
city,
|
||||||
|
)
|
||||||
|
|
||||||
|
if dry_run:
|
||||||
|
logger.info(
|
||||||
|
"dry_run: would process %d yandex_realty_nb houses (ids=%s)",
|
||||||
|
len(rows),
|
||||||
|
[r["house_id"] for r in rows],
|
||||||
|
)
|
||||||
|
result.duration_sec = time.time() - t0
|
||||||
|
return result
|
||||||
|
|
||||||
|
for idx, row in enumerate(rows):
|
||||||
|
house_id: int = int(row["house_id"])
|
||||||
|
jk_slug: str | None = row["yandex_jk_slug"]
|
||||||
|
ext_id: str | None = row["ext_id"]
|
||||||
|
result.processed += 1
|
||||||
|
|
||||||
|
# Idempotency fast-path (belt-and-braces под concurrent writer)
|
||||||
|
if not force and jk_slug:
|
||||||
|
already = db.execute(
|
||||||
|
text(
|
||||||
|
"SELECT 1 FROM market.yandex_jk_enrichment WHERE ext_id = CAST(:ext_id AS text)"
|
||||||
|
),
|
||||||
|
{"ext_id": ext_id},
|
||||||
|
).fetchone()
|
||||||
|
if already is not None:
|
||||||
|
result.skipped_already_enriched += 1
|
||||||
|
logger.debug("skip house_id=%s — already enriched ext_id=%s", house_id, ext_id)
|
||||||
|
continue
|
||||||
|
|
||||||
|
# ── Resolve slug ─────────────────────────────────────────────────
|
||||||
|
if not jk_slug:
|
||||||
|
if not ext_id:
|
||||||
|
logger.warning("house_id=%s: нет yandex_jk_slug и нет ext_id — skip", house_id)
|
||||||
|
result.failed_resolve += 1
|
||||||
|
continue
|
||||||
|
|
||||||
|
try:
|
||||||
|
resolved = await resolve_yandex_jk_slug(ext_id, city=city)
|
||||||
|
except Exception as exc:
|
||||||
|
logger.warning(
|
||||||
|
"resolve_yandex_jk_slug house_id=%s ext_id=%s raised: %s",
|
||||||
|
house_id,
|
||||||
|
ext_id,
|
||||||
|
exc,
|
||||||
|
)
|
||||||
|
resolved = None
|
||||||
|
|
||||||
|
# Anti-bot sleep после resolve-fetch
|
||||||
|
await _sleep_with_jitter(delay, idx, len(rows), force=True)
|
||||||
|
|
||||||
|
if not resolved:
|
||||||
|
logger.warning("slug unresolved house_id=%s ext_id=%s — skip", house_id, ext_id)
|
||||||
|
result.failed_resolve += 1
|
||||||
|
continue
|
||||||
|
|
||||||
|
# Persist slug под SAVEPOINT
|
||||||
|
sp = db.begin_nested()
|
||||||
|
try:
|
||||||
|
db.execute(text(_UPDATE_SLUG), {"slug": resolved, "hid": house_id})
|
||||||
|
sp.commit()
|
||||||
|
db.commit()
|
||||||
|
except Exception as exc:
|
||||||
|
sp.rollback()
|
||||||
|
logger.warning(
|
||||||
|
"persist slug failed house_id=%s ext_id=%s: %s", house_id, ext_id, exc
|
||||||
|
)
|
||||||
|
result.failed_resolve += 1
|
||||||
|
continue
|
||||||
|
|
||||||
|
jk_slug = resolved
|
||||||
|
result.resolved_slug += 1
|
||||||
|
|
||||||
|
# ── Guard: ext_id обязателен для fetch_jk (формирует URL /{slug}-{id}/) ──
|
||||||
|
if not ext_id:
|
||||||
|
logger.warning(
|
||||||
|
"house_id=%s: jk_slug=%s известен, но ext_id пустой — "
|
||||||
|
"fetch_jk пропущен (пустой jk_id даёт битый URL)",
|
||||||
|
house_id,
|
||||||
|
jk_slug,
|
||||||
|
)
|
||||||
|
result.failed_fetch += 1
|
||||||
|
continue
|
||||||
|
|
||||||
|
# ── Fetch через BrowserFetcher ────────────────────────────────────
|
||||||
|
info = None
|
||||||
|
try:
|
||||||
|
scraper = YandexNewbuildingScraper()
|
||||||
|
info = await scraper.fetch_jk(jk_slug=jk_slug, jk_id=ext_id, city=city)
|
||||||
|
except Exception as exc:
|
||||||
|
logger.warning(
|
||||||
|
"fetch_jk failed house_id=%s jk_slug=%s ext_id=%s: %s",
|
||||||
|
house_id,
|
||||||
|
jk_slug,
|
||||||
|
ext_id,
|
||||||
|
exc,
|
||||||
|
)
|
||||||
|
result.failed_fetch += 1
|
||||||
|
await _sleep_with_jitter(delay, idx, len(rows))
|
||||||
|
continue
|
||||||
|
|
||||||
|
if info is None:
|
||||||
|
logger.warning(
|
||||||
|
"fetch_jk returned None house_id=%s jk_slug=%s (anti-bot / parse miss?)",
|
||||||
|
house_id,
|
||||||
|
jk_slug,
|
||||||
|
)
|
||||||
|
result.failed_fetch += 1
|
||||||
|
await _sleep_with_jitter(delay, idx, len(rows))
|
||||||
|
continue
|
||||||
|
|
||||||
|
# ── UPSERT market.yandex_jk_enrichment под SAVEPOINT ─────────────
|
||||||
|
sp = db.begin_nested()
|
||||||
|
try:
|
||||||
|
db.execute(
|
||||||
|
text(_UPSERT_ENRICHMENT),
|
||||||
|
{
|
||||||
|
"ext_id": info.ext_id,
|
||||||
|
"name": info.name,
|
||||||
|
"developer_name": info.developer_name,
|
||||||
|
"address": info.address,
|
||||||
|
"lat": info.lat,
|
||||||
|
"lon": info.lon,
|
||||||
|
"rating": info.rating,
|
||||||
|
"ratings_count": info.ratings_count,
|
||||||
|
"text_reviews_count": info.text_reviews_count,
|
||||||
|
"raw_payload": json.dumps(info.raw_payload or {}, ensure_ascii=False),
|
||||||
|
},
|
||||||
|
)
|
||||||
|
sp.commit()
|
||||||
|
db.commit()
|
||||||
|
result.rows_inserted += 1
|
||||||
|
result.succeeded += 1
|
||||||
|
logger.info(
|
||||||
|
"enriched house_id=%s ext_id=%s slug=%s name=%r",
|
||||||
|
house_id,
|
||||||
|
info.ext_id,
|
||||||
|
jk_slug,
|
||||||
|
info.name,
|
||||||
|
)
|
||||||
|
except Exception as exc:
|
||||||
|
sp.rollback()
|
||||||
|
logger.warning(
|
||||||
|
"UPSERT yandex_jk_enrichment failed house_id=%s ext_id=%s: %s",
|
||||||
|
house_id,
|
||||||
|
ext_id,
|
||||||
|
exc,
|
||||||
|
)
|
||||||
|
try:
|
||||||
|
db.rollback()
|
||||||
|
except Exception as rb_exc:
|
||||||
|
logger.warning("rollback failed house_id=%s: %s", house_id, rb_exc)
|
||||||
|
|
||||||
|
await _sleep_with_jitter(delay, idx, len(rows))
|
||||||
|
|
||||||
|
result.duration_sec = time.time() - t0
|
||||||
|
logger.info(
|
||||||
|
"yandex-nb-sweep done: processed=%d ok=%d skip=%d resolved=%d "
|
||||||
|
"resolve_fail=%d fetch_fail=%d rows_inserted=%d | %.1fs",
|
||||||
|
result.processed,
|
||||||
|
result.succeeded,
|
||||||
|
result.skipped_already_enriched,
|
||||||
|
result.resolved_slug,
|
||||||
|
result.failed_resolve,
|
||||||
|
result.failed_fetch,
|
||||||
|
result.rows_inserted,
|
||||||
|
result.duration_sec,
|
||||||
|
)
|
||||||
|
return result
|
||||||
|
|
||||||
|
|
||||||
|
async def _sleep_with_jitter(delay: float, idx: int, total: int, *, force: bool = False) -> None:
|
||||||
|
"""Polite anti-bot sleep с ±20% jitter.
|
||||||
|
|
||||||
|
Пропускается после последнего элемента (нет следующего fetch) — ЕСЛИ force=False.
|
||||||
|
force=True используется для resolve→enrich gap внутри одного дома.
|
||||||
|
"""
|
||||||
|
if delay <= 0:
|
||||||
|
return
|
||||||
|
if not force and idx >= total - 1:
|
||||||
|
return
|
||||||
|
await asyncio.sleep(delay * random.uniform(0.8, 1.2))
|
||||||
|
|
@ -0,0 +1,82 @@
|
||||||
|
-- 105_market_schema_yandex_enrichment.sql
|
||||||
|
-- Market-схема + таблица yandex_jk_enrichment + houses.yandex_jk_slug (#974).
|
||||||
|
--
|
||||||
|
-- Архитектурное решение (ADR):
|
||||||
|
-- Новые скрейперы пишут в схему market в СУЩЕСТВУЮЩЕМ tradein-postgres.
|
||||||
|
-- Существующие listings/houses не выносятся — market является дополнительным слоем
|
||||||
|
-- для cross-product ЖК-аналитики (Site Finder v2 / GG-форсайт).
|
||||||
|
--
|
||||||
|
-- Таблица market.yandex_jk_enrichment хранит parsed payload с Yandex Realty
|
||||||
|
-- ЖК-лендинга (slug, coords, rating, developer). Ключ — ext_id (= yandex_jk_id).
|
||||||
|
-- UPSERT через ON CONFLICT (ext_id) DO UPDATE.
|
||||||
|
--
|
||||||
|
-- houses.yandex_jk_slug: персистируется при resolve_yandex_jk_slug, используется
|
||||||
|
-- для возобновления sweep без повторного запроса к SERP.
|
||||||
|
--
|
||||||
|
-- ЗАВИСИМОСТИ:
|
||||||
|
-- 009_houses.sql (таблица houses), 031_houses_alter_yandex.sql (yandex_jk_id),
|
||||||
|
-- 101_gendesign_reader_role.sql (роль gendesign_reader).
|
||||||
|
-- Idempotent: IF NOT EXISTS + ADD COLUMN IF NOT EXISTS.
|
||||||
|
-- AUTO-APPLIED на прод при деплое через deploy-tradein.yml/_schema_migrations.
|
||||||
|
|
||||||
|
BEGIN;
|
||||||
|
|
||||||
|
-- ── Схема market ─────────────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
CREATE SCHEMA IF NOT EXISTS market;
|
||||||
|
|
||||||
|
COMMENT ON SCHEMA market IS
|
||||||
|
'Cross-product ЖК-аналитика: yandex/cian newbuilding enrichment, '
|
||||||
|
'Site Finder v2 / GG-форсайт (#974). '
|
||||||
|
'Читается gendesign_reader через GRANT ниже.';
|
||||||
|
|
||||||
|
-- ── market.yandex_jk_enrichment ──────────────────────────────────────────────
|
||||||
|
|
||||||
|
CREATE TABLE IF NOT EXISTS market.yandex_jk_enrichment (
|
||||||
|
id bigserial PRIMARY KEY,
|
||||||
|
ext_id text NOT NULL,
|
||||||
|
name text,
|
||||||
|
developer_name text,
|
||||||
|
address text,
|
||||||
|
lat float8,
|
||||||
|
lon float8,
|
||||||
|
rating float4,
|
||||||
|
ratings_count int,
|
||||||
|
text_reviews_count int,
|
||||||
|
raw_payload jsonb,
|
||||||
|
created_at timestamptz NOT NULL DEFAULT now(),
|
||||||
|
updated_at timestamptz NOT NULL DEFAULT now(),
|
||||||
|
CONSTRAINT uq_market_yandex_jk_ext_id UNIQUE (ext_id)
|
||||||
|
);
|
||||||
|
|
||||||
|
COMMENT ON TABLE market.yandex_jk_enrichment IS
|
||||||
|
'Parsed payload с Yandex Realty ЖК-лендинга (#974). '
|
||||||
|
'Ключ ext_id = yandex_jk_id (string). '
|
||||||
|
'UPSERT через ON CONFLICT (ext_id) DO UPDATE.';
|
||||||
|
|
||||||
|
CREATE INDEX IF NOT EXISTS idx_market_yandex_jk_ext_id
|
||||||
|
ON market.yandex_jk_enrichment (ext_id);
|
||||||
|
|
||||||
|
-- ── houses.yandex_jk_slug ────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
-- Примечание: yandex_jk_id уже существует из миграции 031_houses_alter_yandex.sql.
|
||||||
|
-- yandex_jk_slug — новый столбец, добавляем только если нет.
|
||||||
|
ALTER TABLE houses ADD COLUMN IF NOT EXISTS yandex_jk_slug text;
|
||||||
|
|
||||||
|
COMMENT ON COLUMN houses.yandex_jk_slug IS
|
||||||
|
'Yandex Realty ЖК slug из URL /{city}/kupit/novostrojka/<slug>-<id>/. '
|
||||||
|
'Персистируется sweep-таской при resolve_yandex_jk_slug (#974). '
|
||||||
|
'NULL = slug ещё не разрезолвлен.';
|
||||||
|
|
||||||
|
CREATE INDEX IF NOT EXISTS idx_houses_yandex_jk_slug
|
||||||
|
ON houses (yandex_jk_slug)
|
||||||
|
WHERE yandex_jk_slug IS NOT NULL;
|
||||||
|
|
||||||
|
-- ── Доступ для gendesign_reader ───────────────────────────────────────────────
|
||||||
|
-- Оба продукта (tradein + gendesign Site Finder v2) читают market-слой через роль
|
||||||
|
-- gendesign_reader (создана в 101_gendesign_reader_role.sql).
|
||||||
|
|
||||||
|
GRANT USAGE ON SCHEMA market TO gendesign_reader;
|
||||||
|
GRANT SELECT ON ALL TABLES IN SCHEMA market TO gendesign_reader;
|
||||||
|
|
||||||
|
COMMIT;
|
||||||
|
|
@ -0,0 +1,53 @@
|
||||||
|
-- 106_scrape_schedules_seed_yandex_newbuilding_sweep.sql
|
||||||
|
-- Seed row для Yandex newbuilding sweep (#974).
|
||||||
|
--
|
||||||
|
-- !!! DORMANT BY DESIGN !!! enabled = false.
|
||||||
|
-- Sweep обращается к Yandex Realty через tradein-browser (camoufox).
|
||||||
|
-- Засеян ВЫКЛЮЧЕННЫМ — включается оператором вручную:
|
||||||
|
-- UPDATE scrape_schedules SET enabled = true
|
||||||
|
-- WHERE source = 'yandex_newbuilding_sweep';
|
||||||
|
-- Зеркалит cian_city_sweep (103, тоже enabled = false, dormant).
|
||||||
|
--
|
||||||
|
-- default_params:
|
||||||
|
-- limit — дома за один прогон (старт консервативно 5).
|
||||||
|
-- request_delay_sec — пауза между ЖК (поверх browser wait).
|
||||||
|
-- city — город для Yandex Realty URL (ekaterinburg).
|
||||||
|
--
|
||||||
|
-- ЗАВИСИМОСТИ: 052_scrape_schedules.sql (таблица + UNIQUE(source)).
|
||||||
|
-- Idempotent: ON CONFLICT (source) DO NOTHING — безопасно запускать повторно.
|
||||||
|
--
|
||||||
|
-- yandex_newbuilding_sweep — окно 02:00-05:00 UTC (05:00-08:00 МСК).
|
||||||
|
|
||||||
|
BEGIN;
|
||||||
|
|
||||||
|
INSERT INTO scrape_schedules (
|
||||||
|
source,
|
||||||
|
enabled,
|
||||||
|
window_start_hour,
|
||||||
|
window_end_hour,
|
||||||
|
next_run_at,
|
||||||
|
default_params
|
||||||
|
)
|
||||||
|
VALUES
|
||||||
|
(
|
||||||
|
'yandex_newbuilding_sweep',
|
||||||
|
false, -- DORMANT (#974) — enable manually
|
||||||
|
2,
|
||||||
|
5,
|
||||||
|
((CURRENT_DATE + INTERVAL '1 day') + make_interval(hours => 2)) AT TIME ZONE 'UTC',
|
||||||
|
'{"limit": 5, "request_delay_sec": 8, "city": "ekaterinburg"}'::jsonb
|
||||||
|
)
|
||||||
|
ON CONFLICT (source) DO NOTHING;
|
||||||
|
|
||||||
|
COMMENT ON TABLE scrape_schedules IS
|
||||||
|
'In-app scheduler config (заменяет cron-script setup). '
|
||||||
|
'Sources: avito_city_sweep, yandex_city_sweep (dormant, #561), '
|
||||||
|
'cian_history_backfill, rosreestr_dkp_import, listing_source_snapshot (#570), '
|
||||||
|
'asking_to_sold_ratio_refresh (#648), refresh_search_matview (#769), '
|
||||||
|
'yandex_address_backfill (#855, EKB pilot), '
|
||||||
|
'sber_index_pull (#887, monthly СберИндекс city-level price index), '
|
||||||
|
'rosreestr_quarter_poll (#888, monthly Rosreestr new-quarter availability check), '
|
||||||
|
'cian_city_sweep (dormant, #973), '
|
||||||
|
'yandex_newbuilding_sweep (dormant, #974).';
|
||||||
|
|
||||||
|
COMMIT;
|
||||||
397
tradein-mvp/backend/tests/tasks/test_yandex_newbuilding_sweep.py
Normal file
397
tradein-mvp/backend/tests/tasks/test_yandex_newbuilding_sweep.py
Normal file
|
|
@ -0,0 +1,397 @@
|
||||||
|
"""Tests for yandex_newbuilding_sweep task (#974).
|
||||||
|
|
||||||
|
No real network, no real Postgres. BrowserFetcher, resolve_yandex_jk_slug and
|
||||||
|
fetch_jk are mocked. The DB session is an in-memory fake that models
|
||||||
|
market.yandex_jk_enrichment + houses.yandex_jk_slug persistence assertions.
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import os
|
||||||
|
import sys
|
||||||
|
from unittest.mock import AsyncMock, MagicMock, patch
|
||||||
|
|
||||||
|
# DATABASE_URL required by config before any app import.
|
||||||
|
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
|
||||||
|
|
||||||
|
# WeasyPrint stub — not installed in CI without GTK.
|
||||||
|
_wp_mock = MagicMock()
|
||||||
|
sys.modules.setdefault("weasyprint", _wp_mock)
|
||||||
|
|
||||||
|
import pytest # noqa: E402
|
||||||
|
|
||||||
|
from app.tasks.yandex_newbuilding_sweep import ( # noqa: E402
|
||||||
|
YandexNewbuildingSweepResult,
|
||||||
|
count_yandex_newbuilding_houses,
|
||||||
|
enrich_yandex_newbuilding_sweep,
|
||||||
|
)
|
||||||
|
|
||||||
|
# ── In-memory fake DB ─────────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
class _FakeResult:
|
||||||
|
def __init__(self, *, scalar=None, rows=None, fetchone=None):
|
||||||
|
self._scalar = scalar
|
||||||
|
self._rows = rows or []
|
||||||
|
self._fetchone = fetchone
|
||||||
|
|
||||||
|
def scalar_one(self):
|
||||||
|
return self._scalar
|
||||||
|
|
||||||
|
def mappings(self):
|
||||||
|
return self
|
||||||
|
|
||||||
|
def all(self):
|
||||||
|
return self._rows
|
||||||
|
|
||||||
|
def fetchone(self):
|
||||||
|
return self._fetchone
|
||||||
|
|
||||||
|
|
||||||
|
class FakeDB:
|
||||||
|
"""Minimal SQLAlchemy-Session fake for yandex_newbuilding_sweep.
|
||||||
|
|
||||||
|
Tracks:
|
||||||
|
- enrichment: set of ext_id strings (inserted via UPSERT)
|
||||||
|
- slug_updates: dict house_id → slug (persisted via UPDATE houses)
|
||||||
|
- savepoints: int (# of begin_nested calls)
|
||||||
|
"""
|
||||||
|
|
||||||
|
def __init__(self, pending_rows: list[dict]):
|
||||||
|
self._pending_rows = pending_rows
|
||||||
|
# market.yandex_jk_enrichment keyed by ext_id
|
||||||
|
self.enrichment: dict[str, dict] = {}
|
||||||
|
# houses.yandex_jk_slug written during resolve
|
||||||
|
self.slug_updates: dict[int, str] = {}
|
||||||
|
self.commits = 0
|
||||||
|
self.rollbacks = 0
|
||||||
|
self._sp_stack: list[_FakeSavepoint] = []
|
||||||
|
|
||||||
|
def execute(self, statement, params=None):
|
||||||
|
sql = str(statement)
|
||||||
|
params = params or {}
|
||||||
|
|
||||||
|
# COUNT queries
|
||||||
|
if "COUNT(DISTINCT h.id)" in sql and "yandex_realty_nb" in sql:
|
||||||
|
if "NOT EXISTS" in sql:
|
||||||
|
# _COUNT_PENDING
|
||||||
|
pending = sum(
|
||||||
|
1 for r in self._pending_rows if r.get("ext_id") not in self.enrichment
|
||||||
|
)
|
||||||
|
return _FakeResult(scalar=pending)
|
||||||
|
if "yandex_jk_slug IS NOT NULL" in sql:
|
||||||
|
# _COUNT_FETCHABLE
|
||||||
|
return _FakeResult(scalar=len(self._pending_rows))
|
||||||
|
# _COUNT_TOTAL
|
||||||
|
return _FakeResult(scalar=len(self._pending_rows))
|
||||||
|
|
||||||
|
# SELECT pending houses
|
||||||
|
if "SELECT DISTINCT ON (h.id)" in sql and "yandex_realty_nb" in sql:
|
||||||
|
lim = params.get("lim", len(self._pending_rows))
|
||||||
|
force = params.get("force", False)
|
||||||
|
out = []
|
||||||
|
for r in self._pending_rows[:lim]:
|
||||||
|
ext_id = r.get("ext_id")
|
||||||
|
if not force and ext_id in self.enrichment:
|
||||||
|
continue
|
||||||
|
# Reflect persisted slug
|
||||||
|
slug = self.slug_updates.get(r["house_id"]) or r.get("yandex_jk_slug")
|
||||||
|
out.append({**r, "yandex_jk_slug": slug})
|
||||||
|
return _FakeResult(rows=out)
|
||||||
|
|
||||||
|
# UPDATE houses SET yandex_jk_slug
|
||||||
|
if "UPDATE houses" in sql and "yandex_jk_slug" in sql:
|
||||||
|
self.slug_updates[params["hid"]] = params["slug"]
|
||||||
|
return _FakeResult()
|
||||||
|
|
||||||
|
# Idempotency check: SELECT 1 FROM market.yandex_jk_enrichment
|
||||||
|
if "SELECT 1 FROM market.yandex_jk_enrichment" in sql:
|
||||||
|
ext_id = params.get("ext_id")
|
||||||
|
if ext_id in self.enrichment:
|
||||||
|
return _FakeResult(fetchone=object())
|
||||||
|
return _FakeResult(fetchone=None)
|
||||||
|
|
||||||
|
# UPSERT market.yandex_jk_enrichment
|
||||||
|
if "INSERT INTO market.yandex_jk_enrichment" in sql:
|
||||||
|
self.enrichment[params["ext_id"]] = dict(params)
|
||||||
|
return _FakeResult()
|
||||||
|
|
||||||
|
return _FakeResult()
|
||||||
|
|
||||||
|
def begin_nested(self):
|
||||||
|
sp = _FakeSavepoint(self)
|
||||||
|
self._sp_stack.append(sp)
|
||||||
|
return sp
|
||||||
|
|
||||||
|
def commit(self):
|
||||||
|
self.commits += 1
|
||||||
|
|
||||||
|
def rollback(self):
|
||||||
|
self.rollbacks += 1
|
||||||
|
|
||||||
|
|
||||||
|
class _FakeSavepoint:
|
||||||
|
def __init__(self, db: FakeDB):
|
||||||
|
self._db = db
|
||||||
|
|
||||||
|
def commit(self):
|
||||||
|
pass
|
||||||
|
|
||||||
|
def rollback(self):
|
||||||
|
pass
|
||||||
|
|
||||||
|
|
||||||
|
# ── Tests: count_yandex_newbuilding_houses ───────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
def test_count_returns_sizing_dict():
|
||||||
|
db = FakeDB([{"house_id": 1, "ext_id": "100", "yandex_jk_slug": None}])
|
||||||
|
sizing = count_yandex_newbuilding_houses(db)
|
||||||
|
assert "total" in sizing
|
||||||
|
assert "fetchable" in sizing
|
||||||
|
assert "pending" in sizing
|
||||||
|
assert sizing["total"] >= 0
|
||||||
|
|
||||||
|
|
||||||
|
# ── Tests: enrich_yandex_newbuilding_sweep ───────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_dry_run_returns_counts_without_fetch():
|
||||||
|
"""dry_run: популяция заполнена, fetch не вызван."""
|
||||||
|
db = FakeDB(
|
||||||
|
[
|
||||||
|
{"house_id": 1, "ext_id": "111", "yandex_jk_slug": None, "yandex_jk_id": None},
|
||||||
|
]
|
||||||
|
)
|
||||||
|
result = await enrich_yandex_newbuilding_sweep(db, limit=5, dry_run=True)
|
||||||
|
assert isinstance(result, YandexNewbuildingSweepResult)
|
||||||
|
assert result.processed == 0 # dry_run = no processing
|
||||||
|
assert result.rows_inserted == 0
|
||||||
|
assert result.total >= 0
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_sweep_with_existing_slug_upserts_enrichment():
|
||||||
|
"""House с уже известным slug: fetch_jk вызван, UPSERT market + счётчик."""
|
||||||
|
db = FakeDB(
|
||||||
|
[
|
||||||
|
{
|
||||||
|
"house_id": 10,
|
||||||
|
"ext_id": "999",
|
||||||
|
"yandex_jk_slug": "tatlin",
|
||||||
|
"yandex_jk_id": "999",
|
||||||
|
}
|
||||||
|
]
|
||||||
|
)
|
||||||
|
|
||||||
|
from app.services.scrapers.yandex_newbuilding import YandexNewbuildingInfo
|
||||||
|
|
||||||
|
fake_info = YandexNewbuildingInfo(
|
||||||
|
ext_id="999",
|
||||||
|
ext_slug="tatlin",
|
||||||
|
source_url="https://realty.yandex.ru/ekaterinburg/kupit/novostrojka/tatlin-999/",
|
||||||
|
name="ЖК Татлин",
|
||||||
|
developer_name="ПРИНЦИПЗ",
|
||||||
|
lat=56.855,
|
||||||
|
lon=60.576,
|
||||||
|
rating=4.3,
|
||||||
|
ratings_count=100,
|
||||||
|
text_reviews_count=25,
|
||||||
|
)
|
||||||
|
|
||||||
|
with patch("app.services.scrapers.yandex_newbuilding.YandexNewbuildingScraper") as mock_scraper:
|
||||||
|
scraper_instance = MagicMock()
|
||||||
|
scraper_instance.fetch_jk = AsyncMock(return_value=fake_info)
|
||||||
|
mock_scraper.return_value = scraper_instance
|
||||||
|
|
||||||
|
result = await enrich_yandex_newbuilding_sweep(db, limit=5, request_delay_sec=0)
|
||||||
|
|
||||||
|
assert result.processed == 1
|
||||||
|
assert result.succeeded == 1
|
||||||
|
assert result.rows_inserted == 1
|
||||||
|
assert result.failed_fetch == 0
|
||||||
|
assert result.failed_resolve == 0
|
||||||
|
assert "999" in db.enrichment
|
||||||
|
assert db.enrichment["999"]["name"] == "ЖК Татлин"
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_sweep_resolves_slug_when_missing():
|
||||||
|
"""House без slug: resolve_yandex_jk_slug вызван, slug сохранён, fetch_jk вызван."""
|
||||||
|
db = FakeDB(
|
||||||
|
[
|
||||||
|
{
|
||||||
|
"house_id": 20,
|
||||||
|
"ext_id": "888",
|
||||||
|
"yandex_jk_slug": None,
|
||||||
|
"yandex_jk_id": None,
|
||||||
|
}
|
||||||
|
]
|
||||||
|
)
|
||||||
|
|
||||||
|
from app.services.scrapers.yandex_newbuilding import YandexNewbuildingInfo
|
||||||
|
|
||||||
|
fake_info = YandexNewbuildingInfo(
|
||||||
|
ext_id="888",
|
||||||
|
ext_slug="testjk",
|
||||||
|
source_url="https://realty.yandex.ru/ekaterinburg/kupit/novostrojka/testjk-888/",
|
||||||
|
name="ЖК Тест",
|
||||||
|
lat=56.8,
|
||||||
|
lon=60.5,
|
||||||
|
)
|
||||||
|
|
||||||
|
with (
|
||||||
|
patch(
|
||||||
|
"app.services.scrapers.yandex_newbuilding.resolve_yandex_jk_slug",
|
||||||
|
new_callable=AsyncMock,
|
||||||
|
return_value="testjk",
|
||||||
|
) as mock_resolve,
|
||||||
|
patch("app.services.scrapers.yandex_newbuilding.YandexNewbuildingScraper") as mock_scraper,
|
||||||
|
):
|
||||||
|
scraper_instance = MagicMock()
|
||||||
|
scraper_instance.fetch_jk = AsyncMock(return_value=fake_info)
|
||||||
|
mock_scraper.return_value = scraper_instance
|
||||||
|
|
||||||
|
result = await enrich_yandex_newbuilding_sweep(db, limit=5, request_delay_sec=0)
|
||||||
|
|
||||||
|
mock_resolve.assert_called_once_with("888", city="ekaterinburg")
|
||||||
|
assert result.resolved_slug == 1
|
||||||
|
assert result.succeeded == 1
|
||||||
|
assert result.rows_inserted == 1
|
||||||
|
# Slug persisted to houses
|
||||||
|
assert db.slug_updates.get(20) == "testjk"
|
||||||
|
assert "888" in db.enrichment
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_sweep_failed_resolve_counts():
|
||||||
|
"""resolve_yandex_jk_slug returns None → failed_resolve += 1, no upsert."""
|
||||||
|
db = FakeDB(
|
||||||
|
[
|
||||||
|
{
|
||||||
|
"house_id": 30,
|
||||||
|
"ext_id": "777",
|
||||||
|
"yandex_jk_slug": None,
|
||||||
|
"yandex_jk_id": None,
|
||||||
|
}
|
||||||
|
]
|
||||||
|
)
|
||||||
|
|
||||||
|
with patch(
|
||||||
|
"app.services.scrapers.yandex_newbuilding.resolve_yandex_jk_slug",
|
||||||
|
new_callable=AsyncMock,
|
||||||
|
return_value=None,
|
||||||
|
):
|
||||||
|
result = await enrich_yandex_newbuilding_sweep(db, limit=5, request_delay_sec=0)
|
||||||
|
|
||||||
|
assert result.failed_resolve == 1
|
||||||
|
assert result.succeeded == 0
|
||||||
|
assert result.rows_inserted == 0
|
||||||
|
assert "777" not in db.enrichment
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_sweep_failed_fetch_counts():
|
||||||
|
"""fetch_jk returns None → failed_fetch += 1, no upsert."""
|
||||||
|
db = FakeDB(
|
||||||
|
[
|
||||||
|
{
|
||||||
|
"house_id": 40,
|
||||||
|
"ext_id": "666",
|
||||||
|
"yandex_jk_slug": "existingslug",
|
||||||
|
"yandex_jk_id": "666",
|
||||||
|
}
|
||||||
|
]
|
||||||
|
)
|
||||||
|
|
||||||
|
with patch("app.services.scrapers.yandex_newbuilding.YandexNewbuildingScraper") as mock_scraper:
|
||||||
|
scraper_instance = MagicMock()
|
||||||
|
scraper_instance.fetch_jk = AsyncMock(return_value=None)
|
||||||
|
mock_scraper.return_value = scraper_instance
|
||||||
|
|
||||||
|
result = await enrich_yandex_newbuilding_sweep(db, limit=5, request_delay_sec=0)
|
||||||
|
|
||||||
|
assert result.failed_fetch == 1
|
||||||
|
assert result.succeeded == 0
|
||||||
|
assert "666" not in db.enrichment
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_sweep_idempotent_skip_already_enriched():
|
||||||
|
"""force=False: дом уже в enrichment → пропускается без fetch."""
|
||||||
|
db = FakeDB(
|
||||||
|
[
|
||||||
|
{
|
||||||
|
"house_id": 50,
|
||||||
|
"ext_id": "555",
|
||||||
|
"yandex_jk_slug": "existingslug",
|
||||||
|
"yandex_jk_id": "555",
|
||||||
|
}
|
||||||
|
]
|
||||||
|
)
|
||||||
|
# Pre-populate enrichment
|
||||||
|
db.enrichment["555"] = {"ext_id": "555", "name": "Already enriched"}
|
||||||
|
|
||||||
|
with patch("app.services.scrapers.yandex_newbuilding.YandexNewbuildingScraper") as mock_scraper:
|
||||||
|
scraper_instance = MagicMock()
|
||||||
|
scraper_instance.fetch_jk = AsyncMock(return_value=None)
|
||||||
|
mock_scraper.return_value = scraper_instance
|
||||||
|
|
||||||
|
result = await enrich_yandex_newbuilding_sweep(db, limit=5, request_delay_sec=0)
|
||||||
|
|
||||||
|
# The SELECT filters out already-enriched rows (force=False in SQL),
|
||||||
|
# so processed=0 (SELECT returns empty)
|
||||||
|
assert result.succeeded == 0
|
||||||
|
scraper_instance.fetch_jk.assert_not_called()
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_sweep_no_ext_id_fails_resolve():
|
||||||
|
"""House без ext_id и без slug → failed_resolve (нечего резолвить)."""
|
||||||
|
db = FakeDB(
|
||||||
|
[
|
||||||
|
{
|
||||||
|
"house_id": 60,
|
||||||
|
"ext_id": None,
|
||||||
|
"yandex_jk_slug": None,
|
||||||
|
"yandex_jk_id": None,
|
||||||
|
}
|
||||||
|
]
|
||||||
|
)
|
||||||
|
|
||||||
|
result = await enrich_yandex_newbuilding_sweep(db, limit=5, request_delay_sec=0)
|
||||||
|
|
||||||
|
assert result.failed_resolve == 1
|
||||||
|
assert result.succeeded == 0
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_sweep_slug_present_but_ext_id_null_skips_fetch():
|
||||||
|
"""House с jk_slug но ext_id=NULL → guard срабатывает, fetch_jk НЕ вызван.
|
||||||
|
|
||||||
|
M1 review: пустой jk_id даёт URL /{slug}-/, который является битым.
|
||||||
|
Ожидаем failed_fetch += 1 и scraper.fetch_jk не вызван.
|
||||||
|
"""
|
||||||
|
db = FakeDB(
|
||||||
|
[
|
||||||
|
{
|
||||||
|
"house_id": 70,
|
||||||
|
"ext_id": None,
|
||||||
|
"yandex_jk_slug": "tatlin",
|
||||||
|
"yandex_jk_id": None,
|
||||||
|
}
|
||||||
|
]
|
||||||
|
)
|
||||||
|
|
||||||
|
with patch("app.services.scrapers.yandex_newbuilding.YandexNewbuildingScraper") as mock_scraper:
|
||||||
|
scraper_instance = MagicMock()
|
||||||
|
scraper_instance.fetch_jk = AsyncMock()
|
||||||
|
mock_scraper.return_value = scraper_instance
|
||||||
|
|
||||||
|
result = await enrich_yandex_newbuilding_sweep(db, limit=5, request_delay_sec=0)
|
||||||
|
|
||||||
|
scraper_instance.fetch_jk.assert_not_called()
|
||||||
|
assert result.failed_fetch == 1
|
||||||
|
assert result.succeeded == 0
|
||||||
|
assert result.rows_inserted == 0
|
||||||
|
|
@ -1,6 +1,8 @@
|
||||||
"""Offline smoke for scheduler logic."""
|
"""Offline smoke for scheduler logic."""
|
||||||
|
|
||||||
|
import asyncio
|
||||||
import os
|
import os
|
||||||
|
from unittest.mock import MagicMock
|
||||||
|
|
||||||
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
|
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
|
||||||
|
|
||||||
|
|
@ -170,3 +172,106 @@ def test_claim_run_rechecks_running_under_lock(monkeypatch: pytest.MonkeyPatch)
|
||||||
assert _sched._claim_run(db, _SCHED_ROW) is None
|
assert _sched._claim_run(db, _SCHED_ROW) is None
|
||||||
assert calls["n"] == 0
|
assert calls["n"] == 0
|
||||||
assert db.rollbacks == 1 # лок освобождён
|
assert db.rollbacks == 1 # лок освобождён
|
||||||
|
|
||||||
|
|
||||||
|
# ── #974: trigger_yandex_newbuilding_sweep_run + dispatch ────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
_YANDEX_NB_SWEEP_ROW = {
|
||||||
|
"source": "yandex_newbuilding_sweep",
|
||||||
|
"window_start_hour": 2,
|
||||||
|
"window_end_hour": 5,
|
||||||
|
"default_params": {
|
||||||
|
"limit": 5,
|
||||||
|
"request_delay_sec": 8.0,
|
||||||
|
"city": "ekaterinburg",
|
||||||
|
},
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_trigger_yandex_newbuilding_sweep_creates_run(
|
||||||
|
monkeypatch: pytest.MonkeyPatch,
|
||||||
|
) -> None:
|
||||||
|
"""trigger_yandex_newbuilding_sweep_run claims run + launches create_task (#974)."""
|
||||||
|
db = _FakeSchedulerDB()
|
||||||
|
calls: dict[str, list] = {"create_run": [], "sweep": []}
|
||||||
|
|
||||||
|
def fake_create_run(_d, *, source, params):
|
||||||
|
calls["create_run"].append(source)
|
||||||
|
db.running[source] = True
|
||||||
|
return 301
|
||||||
|
|
||||||
|
monkeypatch.setattr(_sched.runs_mod, "create_run", fake_create_run)
|
||||||
|
|
||||||
|
# Patch enrich_yandex_newbuilding_sweep to a coroutine that records call
|
||||||
|
sweep_result = MagicMock()
|
||||||
|
sweep_result.to_dict.return_value = {}
|
||||||
|
|
||||||
|
async def fake_sweep(run_db, *, limit, **kwargs):
|
||||||
|
calls["sweep"].append(limit)
|
||||||
|
return sweep_result
|
||||||
|
|
||||||
|
# Patch inside the _run closure (lazy import)
|
||||||
|
monkeypatch.setattr(
|
||||||
|
"app.tasks.yandex_newbuilding_sweep.enrich_yandex_newbuilding_sweep",
|
||||||
|
fake_sweep,
|
||||||
|
)
|
||||||
|
monkeypatch.setattr(_sched, "SessionLocal", lambda: db)
|
||||||
|
|
||||||
|
# mark_done is called inside _run; mock it so it doesn't fail on fake db
|
||||||
|
monkeypatch.setattr(_sched.runs_mod, "mark_done", lambda *a, **kw: None)
|
||||||
|
monkeypatch.setattr(_sched.runs_mod, "mark_failed", lambda *a, **kw: None)
|
||||||
|
|
||||||
|
run_id = await _sched.trigger_yandex_newbuilding_sweep_run(db, _YANDEX_NB_SWEEP_ROW)
|
||||||
|
|
||||||
|
assert run_id == 301
|
||||||
|
assert calls["create_run"] == ["yandex_newbuilding_sweep"]
|
||||||
|
# Drain event-loop so the asyncio.create_task runs
|
||||||
|
await asyncio.sleep(0)
|
||||||
|
assert calls["sweep"] == [5]
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_trigger_yandex_newbuilding_sweep_skips_when_running(
|
||||||
|
monkeypatch: pytest.MonkeyPatch,
|
||||||
|
) -> None:
|
||||||
|
"""trigger_yandex_newbuilding_sweep_run returns None when source already running."""
|
||||||
|
db = _FakeSchedulerDB()
|
||||||
|
db.running["yandex_newbuilding_sweep"] = True
|
||||||
|
|
||||||
|
calls: dict[str, int] = {"n": 0}
|
||||||
|
|
||||||
|
def fake_create_run(_d, *, source, params):
|
||||||
|
calls["n"] += 1
|
||||||
|
return 999
|
||||||
|
|
||||||
|
monkeypatch.setattr(_sched.runs_mod, "create_run", fake_create_run)
|
||||||
|
|
||||||
|
result = await _sched.trigger_yandex_newbuilding_sweep_run(db, _YANDEX_NB_SWEEP_ROW)
|
||||||
|
|
||||||
|
assert result is None
|
||||||
|
assert calls["n"] == 0
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_scheduler_dispatch_routes_yandex_newbuilding_sweep(
|
||||||
|
monkeypatch: pytest.MonkeyPatch,
|
||||||
|
) -> None:
|
||||||
|
"""Dispatch block в scheduler_loop роутит yandex_newbuilding_sweep → trigger."""
|
||||||
|
triggered: list[str] = []
|
||||||
|
|
||||||
|
async def fake_trigger(db, sch):
|
||||||
|
triggered.append(sch["source"])
|
||||||
|
return 1
|
||||||
|
|
||||||
|
monkeypatch.setattr(_sched, "trigger_yandex_newbuilding_sweep_run", fake_trigger)
|
||||||
|
|
||||||
|
db = object()
|
||||||
|
sch = _YANDEX_NB_SWEEP_ROW
|
||||||
|
source = sch["source"]
|
||||||
|
# Replicate the elif chain from scheduler_loop for yandex_newbuilding_sweep
|
||||||
|
if source == "yandex_newbuilding_sweep":
|
||||||
|
await _sched.trigger_yandex_newbuilding_sweep_run(db, sch)
|
||||||
|
|
||||||
|
assert "yandex_newbuilding_sweep" in triggered
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue