Compare commits
No commits in common. "5d51eae7057737c95f6d2ed92a99158cbbc42cff" and "27ba08b821a1d02ade801d2fca29e6005cde8e03" have entirely different histories.
5d51eae705
...
27ba08b821
5 changed files with 17 additions and 104 deletions
|
|
@ -28,27 +28,9 @@ from pydantic import BaseModel, Field, field_validator
|
||||||
# (endpoint не прокинут через config, в отличие от SERP-классов и BaseScraper.__aenter__
|
# (endpoint не прокинут через config, в отличие от SERP-классов и BaseScraper.__aenter__
|
||||||
# соседних провайдеров) — вне scope этой задачи (см. "только consume, не трогать
|
# соседних провайдеров) — вне scope этой задачи (см. "только consume, не трогать
|
||||||
# scraper_kit provider-логику"), фиксируем как follow-up.
|
# scraper_kit provider-логику"), фиксируем как follow-up.
|
||||||
#
|
|
||||||
# scraper_kit-миграция (#2397 slice A, эпик #2277): 5 debug city-sweep/full-load
|
|
||||||
# роутов (avito-city-sweep, cian-city-sweep, cian-full-load, yandex-city-sweep,
|
|
||||||
# yandex-full-load) переключены с app.services.scrape_pipeline (legacy) на
|
|
||||||
# scraper_kit.orchestration.pipeline — тот же DI-паттерн, что и app.scheduler_main
|
|
||||||
# ._run_kit_scheduler / scraper_kit.orchestration.scheduler._job_*: config/matcher/
|
|
||||||
# enrichment/proxy_provider инжектируются явно вместо module-level импортов app.*
|
|
||||||
# внутри scrape_pipeline.py. shutdown_requested НЕ прокидывается (нет SIGTERM-drain
|
|
||||||
# семантики у admin BackgroundTasks — эквивалент дефолту lambda: False, поведение не
|
|
||||||
# меняется). Боевой scheduler.py (production sweep-cron) остаётся на legacy — это
|
|
||||||
# отдельный, ещё больший шаг (см. vault Event_Legacy_Scrapers_Partial_Deletion_Jul04).
|
|
||||||
from scraper_kit.base import save_listings
|
from scraper_kit.base import save_listings
|
||||||
from scraper_kit.browser_fetcher import BrowserFetcher
|
from scraper_kit.browser_fetcher import BrowserFetcher
|
||||||
from scraper_kit.orchestration.pipeline import (
|
from scraper_kit.orchestration.pipeline import DEFAULT_REGION_CODE
|
||||||
DEFAULT_REGION_CODE,
|
|
||||||
run_avito_city_sweep,
|
|
||||||
run_cian_city_sweep,
|
|
||||||
run_cian_full_load,
|
|
||||||
run_yandex_city_sweep,
|
|
||||||
run_yandex_full_load,
|
|
||||||
)
|
|
||||||
from scraper_kit.providers.avito.detail import fetch_detail, save_detail_enrichment
|
from scraper_kit.providers.avito.detail import fetch_detail, save_detail_enrichment
|
||||||
from scraper_kit.providers.avito.houses import fetch_house_catalog, save_house_catalog_enrichment
|
from scraper_kit.providers.avito.houses import fetch_house_catalog, save_house_catalog_enrichment
|
||||||
from scraper_kit.providers.avito.imv import (
|
from scraper_kit.providers.avito.imv import (
|
||||||
|
|
@ -72,12 +54,14 @@ from app.services import cian_session as cian_session_svc
|
||||||
from app.services import scrape_runs as runs_mod
|
from app.services import scrape_runs as runs_mod
|
||||||
from app.services.geocoder import geocode
|
from app.services.geocoder import geocode
|
||||||
from app.services.scheduler import has_running_run
|
from app.services.scheduler import has_running_run
|
||||||
from app.services.scraper_adapters import (
|
from app.services.scrape_pipeline import (
|
||||||
RealEnrichmentJobs,
|
run_avito_city_sweep,
|
||||||
RealMatcherAdapter,
|
run_cian_city_sweep,
|
||||||
RealProxyProvider,
|
run_cian_full_load,
|
||||||
RealScraperConfig,
|
run_yandex_city_sweep,
|
||||||
|
run_yandex_full_load,
|
||||||
)
|
)
|
||||||
|
from app.services.scraper_adapters import RealMatcherAdapter, RealProxyProvider, RealScraperConfig
|
||||||
from app.services.scraper_settings import get_scraper_delay, invalidate_cache
|
from app.services.scraper_settings import get_scraper_delay, invalidate_cache
|
||||||
from app.services.scrapers.yandex_newbuilding import YandexNewbuildingScraper
|
from app.services.scrapers.yandex_newbuilding import YandexNewbuildingScraper
|
||||||
from app.tasks.avito_detail_backfill import run_avito_detail_backfill
|
from app.tasks.avito_detail_backfill import run_avito_detail_backfill
|
||||||
|
|
@ -856,10 +840,6 @@ async def start_avito_city_sweep(
|
||||||
await run_avito_city_sweep(
|
await run_avito_city_sweep(
|
||||||
sweep_db,
|
sweep_db,
|
||||||
run_id=run_id,
|
run_id=run_id,
|
||||||
config=RealScraperConfig(),
|
|
||||||
matcher=RealMatcherAdapter(),
|
|
||||||
enrichment=RealEnrichmentJobs(),
|
|
||||||
proxy_provider=_kit_proxy_provider(),
|
|
||||||
radius_m=payload.radius_m,
|
radius_m=payload.radius_m,
|
||||||
pages_per_anchor=payload.pages_per_anchor,
|
pages_per_anchor=payload.pages_per_anchor,
|
||||||
enrich_houses=payload.enrich_houses,
|
enrich_houses=payload.enrich_houses,
|
||||||
|
|
@ -950,9 +930,6 @@ async def start_cian_city_sweep(
|
||||||
await run_cian_city_sweep(
|
await run_cian_city_sweep(
|
||||||
sweep_db,
|
sweep_db,
|
||||||
run_id=run_id,
|
run_id=run_id,
|
||||||
config=RealScraperConfig(),
|
|
||||||
matcher=RealMatcherAdapter(),
|
|
||||||
proxy_provider=_kit_proxy_provider(),
|
|
||||||
radius_m=payload.radius_m,
|
radius_m=payload.radius_m,
|
||||||
pages_per_anchor=payload.pages_per_anchor,
|
pages_per_anchor=payload.pages_per_anchor,
|
||||||
request_delay_sec=payload.request_delay_sec,
|
request_delay_sec=payload.request_delay_sec,
|
||||||
|
|
@ -1082,9 +1059,6 @@ async def start_cian_full_load(
|
||||||
await run_cian_full_load(
|
await run_cian_full_load(
|
||||||
task_db,
|
task_db,
|
||||||
run_id=run_id,
|
run_id=run_id,
|
||||||
config=RealScraperConfig(),
|
|
||||||
matcher=RealMatcherAdapter(),
|
|
||||||
proxy_provider=_kit_proxy_provider(),
|
|
||||||
price_cap_per_bucket=payload.price_cap_per_bucket,
|
price_cap_per_bucket=payload.price_cap_per_bucket,
|
||||||
request_delay_sec=payload.request_delay_sec,
|
request_delay_sec=payload.request_delay_sec,
|
||||||
concurrency=payload.concurrency,
|
concurrency=payload.concurrency,
|
||||||
|
|
@ -1186,10 +1160,6 @@ async def start_yandex_full_load(
|
||||||
await run_yandex_full_load(
|
await run_yandex_full_load(
|
||||||
task_db,
|
task_db,
|
||||||
run_id=run_id,
|
run_id=run_id,
|
||||||
config=RealScraperConfig(),
|
|
||||||
matcher=RealMatcherAdapter(),
|
|
||||||
enrichment=RealEnrichmentJobs(),
|
|
||||||
proxy_provider=_kit_proxy_provider(),
|
|
||||||
price_cap_per_bucket=payload.price_cap_per_bucket,
|
price_cap_per_bucket=payload.price_cap_per_bucket,
|
||||||
request_delay_sec=payload.request_delay_sec,
|
request_delay_sec=payload.request_delay_sec,
|
||||||
concurrency=payload.concurrency,
|
concurrency=payload.concurrency,
|
||||||
|
|
@ -1254,10 +1224,6 @@ async def start_yandex_city_sweep(
|
||||||
await run_yandex_city_sweep(
|
await run_yandex_city_sweep(
|
||||||
sweep_db,
|
sweep_db,
|
||||||
run_id=run_id,
|
run_id=run_id,
|
||||||
config=RealScraperConfig(),
|
|
||||||
matcher=RealMatcherAdapter(),
|
|
||||||
enrichment=RealEnrichmentJobs(),
|
|
||||||
proxy_provider=_kit_proxy_provider(),
|
|
||||||
radius_m=payload.radius_m,
|
radius_m=payload.radius_m,
|
||||||
pages_per_anchor=payload.pages_per_anchor,
|
pages_per_anchor=payload.pages_per_anchor,
|
||||||
request_delay_sec=payload.request_delay_sec,
|
request_delay_sec=payload.request_delay_sec,
|
||||||
|
|
|
||||||
|
|
@ -717,14 +717,12 @@ def test_cian_start_endpoint_validates_delay_too_low(app_with_admin) -> None:
|
||||||
|
|
||||||
|
|
||||||
def test_cian_start_endpoint_ok(app_with_admin) -> None:
|
def test_cian_start_endpoint_ok(app_with_admin) -> None:
|
||||||
"""Valid request → 200, run_id; kit sweep вызван с DI (config/matcher, БЕЗ enrichment)."""
|
"""Valid request → 200, run_id в ответе."""
|
||||||
from unittest.mock import AsyncMock, patch
|
from unittest.mock import AsyncMock, patch
|
||||||
|
|
||||||
from app.services.scraper_adapters import RealMatcherAdapter, RealScraperConfig
|
|
||||||
|
|
||||||
with (
|
with (
|
||||||
patch("app.services.scrape_runs.create_run", return_value=42),
|
patch("app.services.scrape_runs.create_run", return_value=42),
|
||||||
patch("app.api.v1.admin.run_cian_city_sweep", new_callable=AsyncMock) as mock_sweep,
|
patch("app.services.scrape_pipeline.run_cian_city_sweep", new_callable=AsyncMock),
|
||||||
):
|
):
|
||||||
r = app_with_admin.post(
|
r = app_with_admin.post(
|
||||||
"/api/v1/admin/scrape/cian-city-sweep",
|
"/api/v1/admin/scrape/cian-city-sweep",
|
||||||
|
|
@ -737,13 +735,6 @@ def test_cian_start_endpoint_ok(app_with_admin) -> None:
|
||||||
assert body["pages_per_anchor"] == 2
|
assert body["pages_per_anchor"] == 2
|
||||||
assert body["detail_top_n"] == 5
|
assert body["detail_top_n"] == 5
|
||||||
|
|
||||||
mock_sweep.assert_called_once()
|
|
||||||
kw = mock_sweep.call_args.kwargs
|
|
||||||
assert isinstance(kw["config"], RealScraperConfig)
|
|
||||||
assert isinstance(kw["matcher"], RealMatcherAdapter)
|
|
||||||
assert "proxy_provider" in kw
|
|
||||||
assert "enrichment" not in kw
|
|
||||||
|
|
||||||
|
|
||||||
def test_cian_cancel_endpoint(app_with_admin) -> None:
|
def test_cian_cancel_endpoint(app_with_admin) -> None:
|
||||||
"""Cancel endpoint возвращает run_id + cancelled=True."""
|
"""Cancel endpoint возвращает run_id + cancelled=True."""
|
||||||
|
|
|
||||||
|
|
@ -313,17 +313,11 @@ def test_start_endpoint_validates_request_delay_too_low(app_with_admin) -> None:
|
||||||
|
|
||||||
|
|
||||||
def test_start_endpoint_ok(app_with_admin) -> None:
|
def test_start_endpoint_ok(app_with_admin) -> None:
|
||||||
"""Valid request → 200, run_id в ответе; kit sweep вызван с DI (config/matcher/enrichment)."""
|
"""Valid request → 200, run_id в ответе."""
|
||||||
from app.services.scraper_adapters import (
|
|
||||||
RealEnrichmentJobs,
|
|
||||||
RealMatcherAdapter,
|
|
||||||
RealScraperConfig,
|
|
||||||
)
|
|
||||||
|
|
||||||
with (
|
with (
|
||||||
patch("app.api.v1.admin.has_running_run", return_value=False),
|
patch("app.api.v1.admin.has_running_run", return_value=False),
|
||||||
patch("app.services.scrape_runs.create_run", return_value=7),
|
patch("app.services.scrape_runs.create_run", return_value=7),
|
||||||
patch("app.api.v1.admin.run_avito_city_sweep", new_callable=AsyncMock) as mock_sweep,
|
patch("app.services.scrape_pipeline.run_avito_city_sweep", new_callable=AsyncMock),
|
||||||
):
|
):
|
||||||
r = app_with_admin.post(
|
r = app_with_admin.post(
|
||||||
"/api/v1/admin/scrape/avito-city-sweep",
|
"/api/v1/admin/scrape/avito-city-sweep",
|
||||||
|
|
@ -335,13 +329,6 @@ def test_start_endpoint_ok(app_with_admin) -> None:
|
||||||
assert body["status"] == "running"
|
assert body["status"] == "running"
|
||||||
assert body["pages_per_anchor"] == 2
|
assert body["pages_per_anchor"] == 2
|
||||||
|
|
||||||
mock_sweep.assert_called_once()
|
|
||||||
kw = mock_sweep.call_args.kwargs
|
|
||||||
assert isinstance(kw["config"], RealScraperConfig)
|
|
||||||
assert isinstance(kw["matcher"], RealMatcherAdapter)
|
|
||||||
assert isinstance(kw["enrichment"], RealEnrichmentJobs)
|
|
||||||
assert "proxy_provider" in kw
|
|
||||||
|
|
||||||
|
|
||||||
def test_cancel_endpoint_exists(app_with_admin) -> None:
|
def test_cancel_endpoint_exists(app_with_admin) -> None:
|
||||||
"""Cancel endpoint returns run_id + cancelled flag."""
|
"""Cancel endpoint returns run_id + cancelled flag."""
|
||||||
|
|
|
||||||
|
|
@ -453,24 +453,14 @@ def test_city_sweep_start_request_enrich_imv_false() -> None:
|
||||||
|
|
||||||
|
|
||||||
def test_sweep_endpoint_passes_enrich_imv(app_with_admin) -> None:
|
def test_sweep_endpoint_passes_enrich_imv(app_with_admin) -> None:
|
||||||
"""POST /scrape/avito-city-sweep с enrich_imv=False → sweep вызывается без IMV.
|
"""POST /scrape/avito-city-sweep с enrich_imv=False → sweep вызывается без IMV."""
|
||||||
|
|
||||||
kit sweep вызван с DI (config/matcher/enrichment) — mock должен перехватывать
|
|
||||||
имя, забинженное в app.api.v1.admin (kit-импорт), не устаревший scrape_pipeline.
|
|
||||||
"""
|
|
||||||
from app.services.scraper_adapters import (
|
|
||||||
RealEnrichmentJobs,
|
|
||||||
RealMatcherAdapter,
|
|
||||||
RealScraperConfig,
|
|
||||||
)
|
|
||||||
|
|
||||||
with (
|
with (
|
||||||
patch("app.api.v1.admin.has_running_run", return_value=False),
|
patch("app.api.v1.admin.has_running_run", return_value=False),
|
||||||
patch("app.services.scrape_runs.create_run", return_value=42),
|
patch("app.services.scrape_runs.create_run", return_value=42),
|
||||||
patch(
|
patch(
|
||||||
"app.api.v1.admin.run_avito_city_sweep",
|
"app.services.scrape_pipeline.run_avito_city_sweep",
|
||||||
new_callable=AsyncMock,
|
new_callable=AsyncMock,
|
||||||
) as mock_sweep,
|
),
|
||||||
):
|
):
|
||||||
r = app_with_admin.post(
|
r = app_with_admin.post(
|
||||||
"/api/v1/admin/scrape/avito-city-sweep",
|
"/api/v1/admin/scrape/avito-city-sweep",
|
||||||
|
|
@ -480,14 +470,6 @@ def test_sweep_endpoint_passes_enrich_imv(app_with_admin) -> None:
|
||||||
body = r.json()
|
body = r.json()
|
||||||
assert body["run_id"] == 42
|
assert body["run_id"] == 42
|
||||||
|
|
||||||
mock_sweep.assert_called_once()
|
|
||||||
kw = mock_sweep.call_args.kwargs
|
|
||||||
assert isinstance(kw["config"], RealScraperConfig)
|
|
||||||
assert isinstance(kw["matcher"], RealMatcherAdapter)
|
|
||||||
assert isinstance(kw["enrichment"], RealEnrichmentJobs)
|
|
||||||
assert "proxy_provider" in kw
|
|
||||||
assert kw["enrich_imv"] is False
|
|
||||||
|
|
||||||
|
|
||||||
@pytest.fixture
|
@pytest.fixture
|
||||||
def app_with_admin():
|
def app_with_admin():
|
||||||
|
|
|
||||||
|
|
@ -623,18 +623,12 @@ def test_yandex_start_endpoint_validates_delay_too_low(app_with_admin) -> None:
|
||||||
|
|
||||||
|
|
||||||
def test_yandex_start_endpoint_ok(app_with_admin) -> None:
|
def test_yandex_start_endpoint_ok(app_with_admin) -> None:
|
||||||
"""Valid request → 200, run_id в ответе; kit sweep вызван с DI (config/matcher/enrichment)."""
|
"""Valid request → 200, run_id в ответе."""
|
||||||
from unittest.mock import patch
|
from unittest.mock import patch
|
||||||
|
|
||||||
from app.services.scraper_adapters import (
|
|
||||||
RealEnrichmentJobs,
|
|
||||||
RealMatcherAdapter,
|
|
||||||
RealScraperConfig,
|
|
||||||
)
|
|
||||||
|
|
||||||
with (
|
with (
|
||||||
patch("app.services.scrape_runs.create_run", return_value=55),
|
patch("app.services.scrape_runs.create_run", return_value=55),
|
||||||
patch("app.api.v1.admin.run_yandex_city_sweep", new_callable=AsyncMock) as mock_sweep,
|
patch("app.services.scrape_pipeline.run_yandex_city_sweep", new_callable=AsyncMock),
|
||||||
):
|
):
|
||||||
r = app_with_admin.post(
|
r = app_with_admin.post(
|
||||||
"/api/v1/admin/scrape/yandex-city-sweep",
|
"/api/v1/admin/scrape/yandex-city-sweep",
|
||||||
|
|
@ -647,13 +641,6 @@ def test_yandex_start_endpoint_ok(app_with_admin) -> None:
|
||||||
assert body["pages_per_anchor"] == 2
|
assert body["pages_per_anchor"] == 2
|
||||||
assert body["detail_top_n"] == 0
|
assert body["detail_top_n"] == 0
|
||||||
|
|
||||||
mock_sweep.assert_called_once()
|
|
||||||
kw = mock_sweep.call_args.kwargs
|
|
||||||
assert isinstance(kw["config"], RealScraperConfig)
|
|
||||||
assert isinstance(kw["matcher"], RealMatcherAdapter)
|
|
||||||
assert isinstance(kw["enrichment"], RealEnrichmentJobs)
|
|
||||||
assert "proxy_provider" in kw
|
|
||||||
|
|
||||||
|
|
||||||
def test_yandex_cancel_endpoint(app_with_admin) -> None:
|
def test_yandex_cancel_endpoint(app_with_admin) -> None:
|
||||||
"""Cancel endpoint возвращает run_id + cancelled=True."""
|
"""Cancel endpoint возвращает run_id + cancelled=True."""
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue