180 lines
7.0 KiB
Python
180 lines
7.0 KiB
Python
from __future__ import annotations
|
|
|
|
from sqlalchemy import text
|
|
from sqlalchemy.orm import Session
|
|
|
|
from iaai_sync_service.models import (
|
|
BODY_TYPE_ENUM_VALUES,
|
|
CAR_TABLE_NAME,
|
|
COUNTRY_ENUM_VALUES,
|
|
CURRENCY_ENUM_VALUES,
|
|
DRIVE_ENUM_VALUES,
|
|
GEARBOX_ENUM_VALUES,
|
|
IMAGE_TABLE_NAME,
|
|
ORIGIN_ENUM_VALUES,
|
|
SELLING_TYPE_ENUM_VALUES,
|
|
SYNC_RUN_TABLE_NAME,
|
|
STEERING_WHEEL_ENUM_VALUES,
|
|
)
|
|
|
|
SCHEMA_BOOTSTRAP_LOCK_KEY = 901_700_001
|
|
|
|
|
|
def _create_enum_if_missing(session: Session, type_name: str, values: tuple[str, ...]) -> None:
|
|
values_sql = ", ".join("'" + value.replace("'", "''") + "'" for value in values)
|
|
session.execute(
|
|
text(
|
|
f"""
|
|
DO $$
|
|
BEGIN
|
|
IF NOT EXISTS (
|
|
SELECT 1
|
|
FROM pg_type t
|
|
JOIN pg_namespace n ON n.oid = t.typnamespace
|
|
WHERE n.nspname = 'public'
|
|
AND t.typname = '{type_name}'
|
|
) THEN
|
|
CREATE TYPE public.{type_name} AS ENUM ({values_sql});
|
|
END IF;
|
|
END
|
|
$$;
|
|
"""
|
|
)
|
|
)
|
|
|
|
|
|
def _add_enum_value_if_missing(session: Session, type_name: str, value: str) -> None:
|
|
escaped = value.replace("'", "''")
|
|
session.execute(text(f"ALTER TYPE public.{type_name} ADD VALUE IF NOT EXISTS '{escaped}'"))
|
|
|
|
|
|
def _acquire_schema_bootstrap_lock(session: Session) -> None:
|
|
session.execute(text("SELECT pg_advisory_xact_lock(:key)"), {"key": SCHEMA_BOOTSTRAP_LOCK_KEY})
|
|
|
|
|
|
def ensure_schema(session: Session) -> None:
|
|
bind = session.get_bind()
|
|
if bind is None or bind.dialect.name != "postgresql":
|
|
return
|
|
|
|
_acquire_schema_bootstrap_lock(session)
|
|
|
|
_create_enum_if_missing(session, "currencyenum", CURRENCY_ENUM_VALUES)
|
|
_create_enum_if_missing(session, "driveenum", DRIVE_ENUM_VALUES)
|
|
_create_enum_if_missing(session, "gearboxenum", GEARBOX_ENUM_VALUES)
|
|
_create_enum_if_missing(session, "steeringwheelenum", STEERING_WHEEL_ENUM_VALUES)
|
|
_create_enum_if_missing(session, "bodytypeenum", BODY_TYPE_ENUM_VALUES)
|
|
_create_enum_if_missing(session, "countryenum", COUNTRY_ENUM_VALUES)
|
|
_create_enum_if_missing(session, "originenum", ORIGIN_ENUM_VALUES)
|
|
_create_enum_if_missing(session, "sellingtypeenum", SELLING_TYPE_ENUM_VALUES)
|
|
|
|
for value in CURRENCY_ENUM_VALUES:
|
|
_add_enum_value_if_missing(session, "currencyenum", value)
|
|
for value in DRIVE_ENUM_VALUES:
|
|
_add_enum_value_if_missing(session, "driveenum", value)
|
|
for value in GEARBOX_ENUM_VALUES:
|
|
_add_enum_value_if_missing(session, "gearboxenum", value)
|
|
for value in STEERING_WHEEL_ENUM_VALUES:
|
|
_add_enum_value_if_missing(session, "steeringwheelenum", value)
|
|
for value in BODY_TYPE_ENUM_VALUES:
|
|
_add_enum_value_if_missing(session, "bodytypeenum", value)
|
|
for value in COUNTRY_ENUM_VALUES:
|
|
_add_enum_value_if_missing(session, "countryenum", value)
|
|
for value in ORIGIN_ENUM_VALUES:
|
|
_add_enum_value_if_missing(session, "originenum", value)
|
|
for value in SELLING_TYPE_ENUM_VALUES:
|
|
_add_enum_value_if_missing(session, "sellingtypeenum", value)
|
|
|
|
session.execute(
|
|
text(
|
|
f"""
|
|
CREATE TABLE IF NOT EXISTS public.{CAR_TABLE_NAME} (
|
|
id INTEGER GENERATED BY DEFAULT AS IDENTITY PRIMARY KEY,
|
|
parser_id VARCHAR(50) NOT NULL UNIQUE,
|
|
brand VARCHAR(50) NOT NULL,
|
|
model VARCHAR(50) NOT NULL,
|
|
year INTEGER NULL,
|
|
price BIGINT NULL,
|
|
currency public.currencyenum NOT NULL DEFAULT 'USD',
|
|
mileage INTEGER NOT NULL DEFAULT 0,
|
|
country public.countryenum NOT NULL DEFAULT 'NA',
|
|
is_sold BOOLEAN NOT NULL DEFAULT FALSE,
|
|
color VARCHAR NOT NULL DEFAULT 'other',
|
|
drive public.driveenum NULL,
|
|
gearbox public.gearboxenum NULL,
|
|
steering_wheel public.steeringwheelenum NULL,
|
|
body_type public.bodytypeenum NOT NULL DEFAULT 'OTHER',
|
|
engine_volume INTEGER NULL,
|
|
selling_type public.sellingtypeenum NOT NULL DEFAULT 'NA',
|
|
one_owner BOOLEAN NOT NULL DEFAULT FALSE,
|
|
new_car BOOLEAN NOT NULL DEFAULT FALSE,
|
|
is_hidden BOOLEAN NOT NULL DEFAULT FALSE,
|
|
origin public.originenum NOT NULL DEFAULT 'NA',
|
|
origin_url VARCHAR NOT NULL,
|
|
origin_id VARCHAR NOT NULL UNIQUE,
|
|
is_damaged BOOLEAN NOT NULL DEFAULT FALSE,
|
|
evaluation VARCHAR NULL,
|
|
non_smoking BOOLEAN NOT NULL DEFAULT TRUE,
|
|
rental BOOLEAN NOT NULL DEFAULT FALSE,
|
|
repair_history BOOLEAN NOT NULL DEFAULT FALSE,
|
|
slug VARCHAR NOT NULL,
|
|
last_seen_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
|
|
)
|
|
"""
|
|
)
|
|
)
|
|
session.execute(
|
|
text(
|
|
f"""
|
|
CREATE TABLE IF NOT EXISTS public.{IMAGE_TABLE_NAME} (
|
|
id BIGINT GENERATED BY DEFAULT AS IDENTITY PRIMARY KEY,
|
|
fullres_image VARCHAR NOT NULL,
|
|
preview_image VARCHAR NOT NULL,
|
|
order_index INTEGER NOT NULL,
|
|
car_id INTEGER NOT NULL REFERENCES public.{CAR_TABLE_NAME}(id) ON DELETE CASCADE
|
|
)
|
|
"""
|
|
)
|
|
)
|
|
session.execute(
|
|
text(
|
|
f"""
|
|
CREATE TABLE IF NOT EXISTS public.{SYNC_RUN_TABLE_NAME} (
|
|
id BIGINT GENERATED BY DEFAULT AS IDENTITY PRIMARY KEY,
|
|
started_at TIMESTAMPTZ NOT NULL,
|
|
finished_at TIMESTAMPTZ NULL,
|
|
status TEXT NOT NULL,
|
|
lane TEXT NOT NULL,
|
|
ids_fetched INTEGER NOT NULL DEFAULT 0,
|
|
cars_upserted INTEGER NOT NULL DEFAULT 0,
|
|
cars_failed INTEGER NOT NULL DEFAULT 0,
|
|
images_upserted INTEGER NOT NULL DEFAULT 0,
|
|
error_summary TEXT NULL
|
|
)
|
|
"""
|
|
)
|
|
)
|
|
|
|
session.execute(text(f"CREATE INDEX IF NOT EXISTS ix_iaai_images_car_id ON public.{IMAGE_TABLE_NAME} (car_id)"))
|
|
session.execute(
|
|
text(f"CREATE INDEX IF NOT EXISTS ix_iaai_images_car_order ON public.{IMAGE_TABLE_NAME} (car_id, order_index)")
|
|
)
|
|
session.execute(
|
|
text(
|
|
f"CREATE INDEX IF NOT EXISTS ix_iaai_cars_last_seen_id ON public.{CAR_TABLE_NAME} "
|
|
"(last_seen_at DESC, id DESC)"
|
|
)
|
|
)
|
|
session.execute(
|
|
text(
|
|
f"CREATE INDEX IF NOT EXISTS ix_iaai_cars_active_feed ON public.{CAR_TABLE_NAME} "
|
|
"(last_seen_at DESC, id DESC) WHERE is_sold = FALSE AND is_hidden = FALSE"
|
|
)
|
|
)
|
|
session.execute(
|
|
text(f"CREATE INDEX IF NOT EXISTS ix_iaai_sync_runs_started_at ON public.{SYNC_RUN_TABLE_NAME} (started_at)")
|
|
)
|
|
session.execute(
|
|
text(f"CREATE INDEX IF NOT EXISTS ix_iaai_sync_runs_status ON public.{SYNC_RUN_TABLE_NAME} (status)")
|
|
)
|