from __future__ import annotations

from contextlib import contextmanager
from datetime import datetime, timezone
import re
from types import SimpleNamespace

import pytest
from fastapi.testclient import TestClient
from sqlalchemy import create_engine, event, select, text
from sqlalchemy.orm import Session, sessionmaker
from sqlalchemy.pool import StaticPool

from db.base import Base
from db.models import ChannelConnection, MagentoConnection, ShopifyConnection
from db.models_v2 import (
    V2Attribute,
    V2Category,
    V2ChangeEvent,
    V2ChannelPublishState,
    V2ConfigurableLink,
    V2MediaAsset,
    V2Product,
    V2ProductCategory,
    V2ProductCollection,
    V2QueueMedia,
    V2QueueProductData,
    V2QueueRelations,
    V2ShoppingAttributeOptionMap,
    V2SyncJob,
)
from v2_api import get_db


LEGACY_TABLES = [
    ChannelConnection.__table__,
    MagentoConnection.__table__,
    ShopifyConnection.__table__,
]


def _create_v2_test_schema(session: Session) -> None:
    ddl = [
        """
        CREATE TABLE v2_category (
            id INTEGER PRIMARY KEY AUTOINCREMENT,
            name TEXT NOT NULL,
            parent_id INTEGER,
            level INTEGER,
            node_type TEXT NOT NULL DEFAULT 'category',
            created_by TEXT,
            merged_into_id INTEGER,
            normalized_name TEXT NOT NULL DEFAULT '',
            dedupe_scope INTEGER NOT NULL DEFAULT 0,
            created_at TEXT DEFAULT CURRENT_TIMESTAMP,
            updated_at TEXT DEFAULT CURRENT_TIMESTAMP
        )
        """,
        """
        CREATE TABLE v2_attribute (
            id INTEGER PRIMARY KEY AUTOINCREMENT,
            code TEXT NOT NULL UNIQUE,
            label TEXT NOT NULL,
            data_type TEXT NOT NULL,
            is_variant_axis BOOLEAN NOT NULL DEFAULT 0,
            is_channel_required BOOLEAN NOT NULL DEFAULT 1,
            applies_to_product_types JSON NOT NULL DEFAULT '[]',
            created_at TEXT DEFAULT CURRENT_TIMESTAMP
        )
        """,
        """
        CREATE TABLE v2_product (
            id INTEGER PRIMARY KEY AUTOINCREMENT,
            sku TEXT NOT NULL UNIQUE,
            product_type TEXT NOT NULL,
            lifecycle_status TEXT NOT NULL,
            title TEXT NOT NULL,
            slug TEXT,
            seo_title TEXT,
            seo_description TEXT,
            description TEXT,
            primary_category_id INTEGER,
            shopping_l1_category_id INTEGER,
            shopping_l2_category_id INTEGER,
            capabilities JSON NOT NULL DEFAULT '{}',
            attributes JSON NOT NULL DEFAULT '{}',
            source TEXT NOT NULL DEFAULT 'manual',
            created_at TEXT DEFAULT CURRENT_TIMESTAMP,
            updated_at TEXT DEFAULT CURRENT_TIMESTAMP,
            FOREIGN KEY(primary_category_id) REFERENCES v2_category(id),
            FOREIGN KEY(shopping_l1_category_id) REFERENCES v2_category(id),
            FOREIGN KEY(shopping_l2_category_id) REFERENCES v2_category(id)
        )
        """,
        """
        CREATE TABLE v2_product_category (
            product_id INTEGER NOT NULL,
            category_id INTEGER NOT NULL,
            PRIMARY KEY (product_id, category_id),
            FOREIGN KEY(product_id) REFERENCES v2_product(id) ON DELETE CASCADE,
            FOREIGN KEY(category_id) REFERENCES v2_category(id)
        )
        """,
        """
        CREATE TABLE v2_product_collection (
            product_id INTEGER NOT NULL,
            collection_id INTEGER NOT NULL,
            is_primary BOOLEAN NOT NULL DEFAULT 1,
            PRIMARY KEY (product_id, collection_id),
            FOREIGN KEY(product_id) REFERENCES v2_product(id) ON DELETE CASCADE,
            FOREIGN KEY(collection_id) REFERENCES v2_category(id)
        )
        """,
        """
        CREATE TABLE v2_channel_attribute_mapping (
            id INTEGER PRIMARY KEY AUTOINCREMENT,
            channel TEXT NOT NULL,
            canonical_attribute_id INTEGER NOT NULL,
            channel_attribute_code TEXT,
            channel_attribute_id TEXT,
            sync_mode TEXT,
            is_optional BOOLEAN NOT NULL DEFAULT 0,
            created_at TEXT DEFAULT CURRENT_TIMESTAMP,
            UNIQUE(channel, canonical_attribute_id),
            FOREIGN KEY(canonical_attribute_id) REFERENCES v2_attribute(id) ON DELETE CASCADE
        )
        """,
        """
        CREATE TABLE v2_media_asset (
            id INTEGER PRIMARY KEY AUTOINCREMENT,
            product_id INTEGER NOT NULL,
            r2_key TEXT NOT NULL,
            public_url TEXT,
            sha256 TEXT NOT NULL,
            mime_type TEXT NOT NULL,
            width INTEGER,
            height INTEGER,
            alt_text TEXT,
            sort_order INTEGER NOT NULL DEFAULT 0,
            role_flags JSON NOT NULL DEFAULT '[]',
            is_active BOOLEAN NOT NULL DEFAULT 1,
            created_at TEXT DEFAULT CURRENT_TIMESTAMP,
            updated_at TEXT DEFAULT CURRENT_TIMESTAMP,
            UNIQUE(product_id, r2_key),
            FOREIGN KEY(product_id) REFERENCES v2_product(id) ON DELETE CASCADE
        )
        """,
        """
        CREATE TABLE v2_configurable_link (
            parent_id INTEGER NOT NULL,
            child_id INTEGER NOT NULL,
            axis_attributes JSON NOT NULL DEFAULT '{}',
            PRIMARY KEY (parent_id, child_id),
            FOREIGN KEY(parent_id) REFERENCES v2_product(id) ON DELETE CASCADE,
            FOREIGN KEY(child_id) REFERENCES v2_product(id) ON DELETE CASCADE
        )
        """,
        """
        CREATE TABLE v2_shopping_attribute_option_map (
            id INTEGER PRIMARY KEY AUTOINCREMENT,
            channel TEXT NOT NULL,
            shopping_attribute TEXT NOT NULL,
            category_id INTEGER NOT NULL,
            channel_option_id TEXT NOT NULL,
            status TEXT NOT NULL DEFAULT 'active',
            created_at TEXT DEFAULT CURRENT_TIMESTAMP,
            updated_at TEXT DEFAULT CURRENT_TIMESTAMP,
            UNIQUE(channel, shopping_attribute, category_id),
            FOREIGN KEY(category_id) REFERENCES v2_category(id)
        )
        """,
        """
        CREATE TABLE v2_change_event (
            id INTEGER PRIMARY KEY AUTOINCREMENT,
            product_id INTEGER NOT NULL,
            domain TEXT NOT NULL,
            change_type TEXT NOT NULL,
            payload JSON NOT NULL,
            batch_id TEXT,
            triggered_by TEXT NOT NULL,
            created_at TEXT DEFAULT CURRENT_TIMESTAMP,
            FOREIGN KEY(product_id) REFERENCES v2_product(id) ON DELETE CASCADE
        )
        """,
        """
        CREATE TABLE v2_sync_job (
            id TEXT PRIMARY KEY,
            source TEXT NOT NULL,
            channel TEXT,
            domain TEXT,
            batch_id TEXT,
            priority_lane TEXT NOT NULL DEFAULT 'normal',
            status TEXT NOT NULL DEFAULT 'pending',
            triggered_by TEXT NOT NULL,
            item_count INTEGER NOT NULL DEFAULT 0,
            success_count INTEGER NOT NULL DEFAULT 0,
            failure_count INTEGER NOT NULL DEFAULT 0,
            blocked_count INTEGER NOT NULL DEFAULT 0,
            started_at TEXT,
            completed_at TEXT,
            created_at TEXT DEFAULT CURRENT_TIMESTAMP,
            metadata JSON NOT NULL DEFAULT '{}'
        )
        """,
        """
        CREATE TABLE v2_queue_product_data (
            id INTEGER PRIMARY KEY AUTOINCREMENT,
            sync_job_id TEXT,
            product_id INTEGER NOT NULL,
            channel TEXT NOT NULL,
            priority INTEGER NOT NULL DEFAULT 5,
            priority_lane TEXT NOT NULL DEFAULT 'normal',
            status TEXT NOT NULL DEFAULT 'pending',
            block_reason TEXT,
            payload JSON NOT NULL,
            attempt_count INTEGER NOT NULL DEFAULT 0,
            max_attempts INTEGER NOT NULL DEFAULT 5,
            last_error TEXT,
            last_response JSON,
            locked_by TEXT,
            locked_at TEXT,
            completed_at TEXT,
            coalesce_key TEXT NOT NULL,
            created_at TEXT DEFAULT CURRENT_TIMESTAMP,
            scheduled_at TEXT DEFAULT CURRENT_TIMESTAMP,
            FOREIGN KEY(sync_job_id) REFERENCES v2_sync_job(id) ON DELETE SET NULL,
            FOREIGN KEY(product_id) REFERENCES v2_product(id) ON DELETE CASCADE
        )
        """,
        """
        CREATE UNIQUE INDEX uq_v2_queue_product_data_pending_coalesce
        ON v2_queue_product_data (coalesce_key) WHERE status = 'pending'
        """,
        """
        CREATE TABLE v2_queue_media (
            id INTEGER PRIMARY KEY AUTOINCREMENT,
            sync_job_id TEXT,
            product_id INTEGER NOT NULL,
            channel TEXT NOT NULL,
            priority INTEGER NOT NULL DEFAULT 5,
            priority_lane TEXT NOT NULL DEFAULT 'normal',
            status TEXT NOT NULL DEFAULT 'pending',
            block_reason TEXT,
            payload JSON NOT NULL,
            attempt_count INTEGER NOT NULL DEFAULT 0,
            max_attempts INTEGER NOT NULL DEFAULT 5,
            last_error TEXT,
            last_response JSON,
            locked_by TEXT,
            locked_at TEXT,
            completed_at TEXT,
            coalesce_key TEXT NOT NULL,
            created_at TEXT DEFAULT CURRENT_TIMESTAMP,
            scheduled_at TEXT DEFAULT CURRENT_TIMESTAMP,
            FOREIGN KEY(sync_job_id) REFERENCES v2_sync_job(id) ON DELETE SET NULL,
            FOREIGN KEY(product_id) REFERENCES v2_product(id) ON DELETE CASCADE
        )
        """,
        """
        CREATE UNIQUE INDEX uq_v2_queue_media_pending_coalesce
        ON v2_queue_media (coalesce_key) WHERE status = 'pending'
        """,
        """
        CREATE TABLE v2_queue_relations (
            id INTEGER PRIMARY KEY AUTOINCREMENT,
            sync_job_id TEXT,
            product_id INTEGER NOT NULL,
            channel TEXT NOT NULL,
            priority INTEGER NOT NULL DEFAULT 5,
            priority_lane TEXT NOT NULL DEFAULT 'normal',
            status TEXT NOT NULL DEFAULT 'pending',
            block_reason TEXT,
            payload JSON NOT NULL,
            attempt_count INTEGER NOT NULL DEFAULT 0,
            max_attempts INTEGER NOT NULL DEFAULT 5,
            last_error TEXT,
            last_response JSON,
            locked_by TEXT,
            locked_at TEXT,
            completed_at TEXT,
            coalesce_key TEXT NOT NULL,
            created_at TEXT DEFAULT CURRENT_TIMESTAMP,
            scheduled_at TEXT DEFAULT CURRENT_TIMESTAMP,
            FOREIGN KEY(sync_job_id) REFERENCES v2_sync_job(id) ON DELETE SET NULL,
            FOREIGN KEY(product_id) REFERENCES v2_product(id) ON DELETE CASCADE
        )
        """,
        """
        CREATE UNIQUE INDEX uq_v2_queue_relations_pending_coalesce
        ON v2_queue_relations (coalesce_key) WHERE status = 'pending'
        """,
        """
        CREATE TABLE v2_channel_publish_state (
            product_id INTEGER NOT NULL,
            channel TEXT NOT NULL,
            domain TEXT NOT NULL,
            status TEXT NOT NULL DEFAULT 'pending',
            block_reason TEXT,
            channel_ref_id TEXT,
            last_event_id INTEGER,
            last_queue_item_id INTEGER,
            last_pushed_at TEXT,
            last_verified_at TEXT,
            verification_mode TEXT,
            verification_result JSON,
            last_error TEXT,
            updated_at TEXT DEFAULT CURRENT_TIMESTAMP,
            PRIMARY KEY (product_id, channel, domain),
            FOREIGN KEY(product_id) REFERENCES v2_product(id) ON DELETE CASCADE,
            FOREIGN KEY(last_event_id) REFERENCES v2_change_event(id) ON DELETE SET NULL
        )
        """,
    ]
    for statement in ddl:
        session.execute(text(statement))
    session.commit()


@pytest.fixture()
def v2_session() -> Session:
    engine = create_engine(
        "sqlite://",
        connect_args={"check_same_thread": False},
        poolclass=StaticPool,
    )

    @event.listens_for(engine, "connect")
    def _sqlite_pragma(dbapi_connection, connection_record):
        cursor = dbapi_connection.cursor()
        cursor.execute("PRAGMA foreign_keys=ON")
        cursor.close()
        dbapi_connection.create_function(
            "regexp_replace",
            4,
            lambda value, pattern, replacement, flags: re.sub(pattern, replacement, value or ""),
            deterministic=True,
        )
        dbapi_connection.create_function("now", 0, lambda: datetime.now(timezone.utc).isoformat())

    Base.metadata.create_all(engine, tables=LEGACY_TABLES)
    session = sessionmaker(bind=engine, expire_on_commit=False)()
    _create_v2_test_schema(session)
    try:
        yield session
    finally:
        session.close()
        engine.dispose()


@pytest.fixture()
def v2_client(v2_session: Session):
    from main_v2 import app

    def _override_db():
        yield v2_session

    app.dependency_overrides[get_db] = _override_db
    try:
        yield TestClient(app)
    finally:
        app.dependency_overrides.clear()


def _add_active_connection(session: Session, channel_code: str = "dash-magento") -> None:
    session.add(
        ChannelConnection(
            channel_type="magento",
            channel_code=channel_code,
            display_name="Dashboard Magento",
            environment="production",
            status="active",
        )
    )
    session.commit()


@contextmanager
def _session_scope(session: Session):
    try:
        yield session
    finally:
        pass


def test_v2_product_create_enqueues_product_data(v2_client: TestClient, v2_session: Session):
    _add_active_connection(v2_session)

    response = v2_client.post(
        "/api/v2/products",
        json={
            "sku": "SKU-100",
            "title": "Sample Cabinet",
            "product_type": "simple",
            "lifecycle_status": "active",
            "attributes": {},
            "category_ids": [],
            "collection_ids": [],
        },
    )

    assert response.status_code == 200
    product = v2_session.execute(select(V2Product).where(V2Product.sku == "SKU-100")).scalar_one()
    event_rows = v2_session.execute(select(V2ChangeEvent)).scalars().all()
    queue_rows = v2_session.execute(select(V2QueueProductData)).scalars().all()
    state_rows = v2_session.execute(select(V2ChannelPublishState)).scalars().all()
    sync_jobs = v2_session.execute(select(V2SyncJob)).scalars().all()

    assert product.title == "Sample Cabinet"
    assert len(event_rows) == 1
    assert event_rows[0].domain == "product_data"
    assert len(queue_rows) == 1
    assert queue_rows[0].product_id == product.id
    assert queue_rows[0].channel == "dash-magento"
    assert queue_rows[0].status == "pending"
    assert len(state_rows) == 1
    assert state_rows[0].channel == "dash-magento"
    assert state_rows[0].domain == "product_data"
    assert state_rows[0].status == "pending"
    assert len(sync_jobs) == 1
    assert sync_jobs[0].item_count == 1
    assert sync_jobs[0].blocked_count == 0


def test_v2_missing_attribute_mapping_blocks_product_queue(v2_client: TestClient, v2_session: Session):
    _add_active_connection(v2_session, channel_code="shopify-live")
    v2_session.add(
        V2Attribute(
            code="finish",
            label="Finish",
            data_type="select",
            is_channel_required=True,
            applies_to_product_types=["simple"],
        )
    )
    v2_session.commit()

    response = v2_client.post(
        "/api/v2/products",
        json={
            "sku": "SKU-BLOCK",
            "title": "Blocked Cabinet",
            "product_type": "simple",
            "lifecycle_status": "active",
            "attributes": {"finish": "Oak"},
            "category_ids": [],
            "collection_ids": [],
        },
    )

    assert response.status_code == 200
    queue_rows = v2_session.execute(select(V2QueueProductData)).scalars().all()
    state_rows = v2_session.execute(select(V2ChannelPublishState)).scalars().all()
    sync_jobs = v2_session.execute(select(V2SyncJob)).scalars().all()

    assert queue_rows == []
    assert len(state_rows) == 1
    assert state_rows[0].status == "blocked"
    assert state_rows[0].block_reason == "missing_channel_mapping:finish"
    assert len(sync_jobs) == 1
    assert sync_jobs[0].status == "failed"
    assert sync_jobs[0].item_count == 1
    assert sync_jobs[0].blocked_count == 1


def test_v2_media_and_variation_writes_enqueue_expected_domains(v2_client: TestClient, v2_session: Session):
    _add_active_connection(v2_session, channel_code="magento-main")

    parent = V2Product(sku="PARENT-1", title="Parent", product_type="simple", lifecycle_status="active", source="manual")
    child_a = V2Product(sku="CHILD-1A", title="Child A", product_type="simple", lifecycle_status="active", source="manual")
    child_b = V2Product(sku="CHILD-1B", title="Child B", product_type="simple", lifecycle_status="active", source="manual")
    v2_session.add_all([parent, child_a, child_b])
    v2_session.commit()

    media_response = v2_client.put(
        "/api/v2/products/PARENT-1/media",
        json={
            "assets": [
                {
                    "r2_key": "products/PARENT-1/front.jpg",
                    "public_url": "https://cdn.example.com/front.jpg",
                    "sha256": "abc123",
                    "mime_type": "image/jpeg",
                    "sort_order": 0,
                    "role_flags": ["base"],
                    "is_active": True,
                }
            ]
        },
    )
    assert media_response.status_code == 200

    variation_response = v2_client.put(
        "/api/v2/variations/PARENT-1",
        json={
            "parent_title": "Parent Configurable",
            "children": [
                {"sku": "CHILD-1A", "axis_attributes": {"finish": "Oak"}},
                {"sku": "CHILD-1B", "axis_attributes": {"finish": "Walnut"}},
            ],
        },
    )
    assert variation_response.status_code == 200

    media_rows = v2_session.execute(select(V2QueueMedia)).scalars().all()
    relation_rows = v2_session.execute(select(V2QueueRelations)).scalars().all()
    product_rows = v2_session.execute(select(V2QueueProductData)).scalars().all()
    states = v2_session.execute(
        select(V2ChannelPublishState).where(V2ChannelPublishState.product_id == parent.id)
    ).scalars().all()

    assert len(media_rows) == 1
    assert media_rows[0].channel == "magento-main"
    assert len(relation_rows) == 1
    assert relation_rows[0].channel == "magento-main"
    assert len(product_rows) == 1
    assert product_rows[0].channel == "magento-main"
    assert sorted((row.domain, row.status) for row in states) == [
        ("media", "pending"),
        ("product_data", "pending"),
        ("relations", "pending"),
    ]


def test_v2_product_create_applies_slug_seo_and_color_defaults(v2_client: TestClient, v2_session: Session):
    response = v2_client.post(
        "/api/v2/products",
        json={
            "sku": "SKU-NORM",
            "title": "Kitchen Cabinets / Quest Metro Mist",
            "product_type": "simple",
            "lifecycle_status": "active",
            "attributes": {
                "category_l1": "kitchen cabinet",
                "collection": "Kitchen Cabinets / Quest Metro Mist",
                "color": "White",
            },
            "category_ids": [],
            "collection_ids": [],
        },
    )

    assert response.status_code == 200
    product = v2_session.execute(select(V2Product).where(V2Product.sku == "SKU-NORM")).scalar_one()
    assert product.slug == "kitchen-cabinets-quest-metro-mist"
    assert product.seo_title == "Kitchen Cabinets / Quest Metro Mist"
    assert product.seo_description is None
    assert product.attributes["category_l1"] == "Kitchen Cabinets"
    assert product.attributes["collection"] == "Quest Metro Mist"
    assert product.attributes["color"] == "White"
    assert product.attributes["color_finish"] == "White"


def test_v2_collection_and_taxonomy_creation_canonicalize_names(v2_client: TestClient, v2_session: Session):
    hub_response = v2_client.post(
        "/api/v2/taxonomy/nodes",
        json={
            "name": "kitchen cabinet",
            "node_type": "hub",
            "created_by": "test",
        },
    )
    assert hub_response.status_code == 200
    hub_id = hub_response.json()["category_id"]

    collection_response = v2_client.post(
        "/api/v2/collections",
        json={
            "name": "Kitchen Cabinets / Anna Snow White",
            "parent_id": hub_id,
            "created_by": "test",
        },
    )

    assert collection_response.status_code == 200
    body = collection_response.json()
    assert body["collection"]["name"] == "Anna Snow White"
    assert body["collection"]["path_slug"] == "kitchen-cabinets/anna-snow-white"


def test_v2_worker_processes_pending_item_and_verifies_publish_state(
    v2_client: TestClient,
    v2_session: Session,
    monkeypatch: pytest.MonkeyPatch,
):
    from app.jobs import v2_sync_worker

    _add_active_connection(v2_session, channel_code="magento-main")
    create_response = v2_client.post(
        "/api/v2/products",
        json={
            "sku": "SKU-WORKER-1",
            "title": "Worker Product",
            "product_type": "simple",
            "lifecycle_status": "active",
            "attributes": {},
            "category_ids": [],
            "collection_ids": [],
        },
    )
    assert create_response.status_code == 200

    monkeypatch.setattr(v2_sync_worker, "get_session", lambda: _session_scope(v2_session))
    monkeypatch.setenv("V2_SYNC_EXECUTION_MODE", "simulate")

    processed = v2_sync_worker.run_one(worker_id="test-worker")

    assert processed is True
    queue_row = v2_session.execute(select(V2QueueProductData)).scalar_one()
    publish_state = v2_session.execute(select(V2ChannelPublishState)).scalar_one()
    sync_job = v2_session.execute(select(V2SyncJob)).scalar_one()

    assert queue_row.status == "done"
    assert queue_row.locked_at is None
    assert queue_row.locked_by is None
    assert queue_row.last_response["status"] == "simulated"
    assert publish_state.status == "verified"
    assert publish_state.verification_result["status"] == "simulated"
    assert publish_state.last_verified_at is not None
    assert sync_job.status == "completed"
    assert sync_job.success_count == 1
    assert sync_job.failure_count == 0


def test_v2_worker_resets_stale_processing_back_to_pending(v2_session: Session):
    from db.v2_jobs import reset_stale_processing_rows

    job_id = "11111111-1111-1111-1111-111111111111"
    product = V2Product(
        sku="SKU-STALE",
        title="Stale",
        product_type="simple",
        lifecycle_status="active",
        source="manual",
    )
    v2_session.add(product)
    v2_session.flush()

    queue_row = V2QueueProductData(
        sync_job_id=job_id,
        product_id=int(product.id),
        channel="stale-ch",
        priority=5,
        priority_lane="normal",
        status="processing",
        payload={"sku": "SKU-STALE"},
        attempt_count=1,
        max_attempts=5,
        coalesce_key=f"stale-ch:{int(product.id)}:product_data",
        locked_by="worker-old",
        locked_at=datetime(2026, 8, 1, tzinfo=timezone.utc),
    )
    state = V2ChannelPublishState(
        product_id=int(product.id),
        channel="stale-ch",
        domain="product_data",
        status="processing",
    )
    job = V2SyncJob(
        id=job_id,
        source="test",
        channel="stale-ch",
        domain="product_data",
        priority_lane="normal",
        status="running",
        triggered_by="test",
        item_count=1,
    )
    v2_session.add(job)
    v2_session.flush()
    v2_session.add(queue_row)
    v2_session.add(state)
    v2_session.commit()

    reset_count = reset_stale_processing_rows(v2_session)

    assert reset_count == 1
    refreshed = v2_session.get(V2QueueProductData, int(queue_row.id))
    refreshed_state = v2_session.get(
        V2ChannelPublishState,
        {"product_id": int(product.id), "channel": "stale-ch", "domain": "product_data"},
    )
    assert refreshed.status == "pending"
    assert refreshed.locked_at is None
    assert refreshed.locked_by is None
    assert refreshed.last_error == "Reset after stale processing lock"
    assert refreshed_state.status == "pending"
    assert refreshed_state.verification_result["status"] == "stale_reset"


def test_v2_worker_marks_failure_for_invalid_execution_mode(
    v2_client: TestClient,
    v2_session: Session,
    monkeypatch: pytest.MonkeyPatch,
):
    from app.jobs import v2_sync_worker

    _add_active_connection(v2_session, channel_code="magento-main")
    create_response = v2_client.post(
        "/api/v2/products",
        json={
            "sku": "SKU-WORKER-FAIL",
            "title": "Worker Failure Product",
            "product_type": "simple",
            "lifecycle_status": "active",
            "attributes": {},
            "category_ids": [],
            "collection_ids": [],
        },
    )
    assert create_response.status_code == 200

    monkeypatch.setattr(v2_sync_worker, "get_session", lambda: _session_scope(v2_session))
    monkeypatch.setenv("V2_SYNC_EXECUTION_MODE", "bogus")

    processed = v2_sync_worker.run_one(worker_id="test-worker-fail")

    assert processed is True
    queue_row = v2_session.execute(select(V2QueueProductData)).scalar_one()
    publish_state = v2_session.execute(select(V2ChannelPublishState)).scalar_one()
    sync_job = v2_session.execute(select(V2SyncJob)).scalar_one()

    assert queue_row.status == "failed"
    assert "Unsupported V2 sync execution mode" in (queue_row.last_error or "")
    assert publish_state.status == "failed"
    assert "Unsupported V2 sync execution mode" in (publish_state.last_error or "")
    assert sync_job.status == "failed"
    assert sync_job.success_count == 0
    assert sync_job.failure_count == 1


def test_v2_worker_live_shopify_product_bridge(
    v2_session: Session,
    monkeypatch: pytest.MonkeyPatch,
):
    from app.jobs import v2_sync_worker
    from db.models import ShopifyConnection

    captured: dict[str, object] = {}

    def _fake_push_products(session, **kwargs):
        captured.update(kwargs)
        return {"status": "ok", "pushed": 1, "failed": 0}

    v2_session.add(
        ShopifyConnection(
            shop_code="shop-live",
            shop_domain="shop-live.example.com",
            api_version="2026-04",
            admin_access_token="token",
            status="active",
            environment="production",
        )
    )
    v2_session.commit()

    monkeypatch.setattr(v2_sync_worker, "get_session", lambda: _session_scope(v2_session))
    monkeypatch.setenv("V2_SYNC_EXECUTION_MODE", "live")
    monkeypatch.setattr("shopify.product_sync.push_products", _fake_push_products)

    result = v2_sync_worker.process_claimed_item(
        domain="product_data",
        row=SimpleNamespace(id=91, channel="shop-live", payload={"sku": "SHOP-SKU-1"}),
    )

    assert result["status"] == "synced"
    assert result["backend"] == "shopify"
    assert captured["dry_run"] is False
    assert captured["force"] is True
    assert captured["skus"] == ["SHOP-SKU-1"]
    assert captured["shop_code"] == "shop-live"
    assert captured["connection_id"] == 1000001


def test_v2_worker_live_magento_media_bridge(
    v2_session: Session,
    monkeypatch: pytest.MonkeyPatch,
):
    from app.jobs import v2_sync_worker
    from db.magento_repositories import SqlAlchemyMagentoSyncQueueRepository
    from db.models import MagentoConnection

    captured: dict[str, object] = {}

    def _fake_enqueue(self, connection_id, label, options=None, supersede_queued=False):
        captured["enqueue_connection_id"] = connection_id
        captured["enqueue_label"] = label
        captured["enqueue_options"] = dict(options or {})
        return 77

    def _fake_get_status(self, queue_id):
        captured["status_queue_id"] = queue_id
        return {"status": "done", "progress": {"success_count": 1, "error_count": 0}}

    def _fake_run_one(queue_id=None):
        captured["run_one_queue_id"] = queue_id
        return True

    v2_session.add(
        MagentoConnection(
            tenant_id="tenant-a",
            store_base_url="https://magento.example.com",
            magento_api_base_url="https://magento.example.com/rest/V1",
            environment="production",
            consumer_key="key",
            consumer_secret="secret",
            access_token="token",
            access_token_secret="token-secret",
            store_code="magento-live",
            status="active",
        )
    )
    v2_session.commit()

    monkeypatch.setattr(v2_sync_worker, "get_session", lambda: _session_scope(v2_session))
    monkeypatch.setenv("V2_SYNC_EXECUTION_MODE", "live")
    monkeypatch.setattr(SqlAlchemyMagentoSyncQueueRepository, "enqueue", _fake_enqueue)
    monkeypatch.setattr(SqlAlchemyMagentoSyncQueueRepository, "get_status", _fake_get_status)
    monkeypatch.setattr("app.jobs.magento_sync_worker.run_one", _fake_run_one)

    result = v2_sync_worker.process_claimed_item(
        domain="media",
        row=SimpleNamespace(id=55, channel="magento-live", payload={"sku": "MAG-SKU-1"}),
    )

    assert result["status"] == "synced"
    assert result["backend"] == "magento_sync_queue"
    assert result["queue_id"] == 77
    assert captured["enqueue_connection_id"] == 1
    assert captured["run_one_queue_id"] == 77
    assert captured["status_queue_id"] == 77
    assert captured["enqueue_options"]["images_only"] is True
    assert captured["enqueue_options"]["force_images"] is True
    assert captured["enqueue_options"]["limit_skus"] == ["MAG-SKU-1"]


def test_v2_queue_item_list_detail_and_retry_endpoints(v2_client: TestClient, v2_session: Session):
    _add_active_connection(v2_session, channel_code="dash-magento")
    create_response = v2_client.post(
        "/api/v2/products",
        json={
            "sku": "SKU-QUEUE-1",
            "title": "Queue Endpoint Product",
            "product_type": "simple",
            "lifecycle_status": "active",
            "attributes": {},
            "category_ids": [],
            "collection_ids": [],
        },
    )
    assert create_response.status_code == 200

    queue_row = v2_session.execute(select(V2QueueProductData)).scalar_one()
    queue_row.status = "failed"
    queue_row.last_error = "Boom"
    v2_session.commit()

    list_response = v2_client.get("/api/v2/queues/product_data/items", params={"status": "failed"})
    assert list_response.status_code == 200
    list_body = list_response.json()
    assert list_body["total"] == 1
    assert list_body["items"][0]["sku"] == "SKU-QUEUE-1"
    assert list_body["items"][0]["status"] == "failed"
    assert list_body["items"][0]["last_error"] == "Boom"

    detail_response = v2_client.get(f"/api/v2/queues/product_data/items/{int(queue_row.id)}")
    assert detail_response.status_code == 200
    detail_body = detail_response.json()
    assert detail_body["id"] == int(queue_row.id)
    assert detail_body["sku"] == "SKU-QUEUE-1"
    assert detail_body["payload"]["sku"] == "SKU-QUEUE-1"

    retry_response = v2_client.post(f"/api/v2/queues/product_data/items/{int(queue_row.id)}/retry")
    assert retry_response.status_code == 200
    assert retry_response.json()["queue_status"] == "pending"

    refreshed = v2_session.get(V2QueueProductData, int(queue_row.id))
    refreshed_state = v2_session.get(
        V2ChannelPublishState,
        {"product_id": int(refreshed.product_id), "channel": "dash-magento", "domain": "product_data"},
    )
    assert refreshed.status == "pending"
    assert refreshed_state.status == "pending"


def test_v2_queue_item_replay_alias_and_retry_guard(v2_client: TestClient, v2_session: Session):
    _add_active_connection(v2_session, channel_code="dash-magento")
    create_response = v2_client.post(
        "/api/v2/products",
        json={
            "sku": "SKU-QUEUE-2",
            "title": "Queue Replay Product",
            "product_type": "simple",
            "lifecycle_status": "active",
            "attributes": {},
            "category_ids": [],
            "collection_ids": [],
        },
    )
    assert create_response.status_code == 200

    queue_row = v2_session.execute(select(V2QueueProductData)).scalar_one()

    guard_response = v2_client.post(f"/api/v2/queues/product_data/items/{int(queue_row.id)}/retry")
    assert guard_response.status_code == 409
    assert "cannot be retried" in guard_response.json()["detail"]

    queue_row.status = "dead_letter"
    queue_row.last_error = "Too many attempts"
    v2_session.commit()

    replay_response = v2_client.post(f"/api/v2/queues/product_data/items/{int(queue_row.id)}/replay")
    assert replay_response.status_code == 200
    assert replay_response.json()["queue_status"] == "pending"


def test_v2_sync_job_detail_includes_queue_items(v2_client: TestClient, v2_session: Session):
    _add_active_connection(v2_session, channel_code="dash-magento")
    create_response = v2_client.post(
        "/api/v2/products",
        json={
            "sku": "SKU-JOB-1",
            "title": "Sync Job Detail Product",
            "product_type": "simple",
            "lifecycle_status": "active",
            "attributes": {},
            "category_ids": [],
            "collection_ids": [],
        },
    )
    assert create_response.status_code == 200

    sync_job = v2_session.execute(select(V2SyncJob)).scalar_one()
    queue_row = v2_session.execute(select(V2QueueProductData)).scalar_one()
    queue_row.status = "done"
    queue_row.last_response = {"status": "synced", "backend": "shopify"}
    sync_job.status = "completed"
    sync_job.success_count = 1
    v2_session.commit()

    response = v2_client.get(f"/api/v2/sync-jobs/{sync_job.id}")

    assert response.status_code == 200
    body = response.json()
    assert body["id"] == sync_job.id
    assert body["status"] == "completed"
    assert body["queue_counts"]["product_data"] == 1
    assert body["queue_counts"]["media"] == 0
    assert body["queue_items"][0]["id"] == int(queue_row.id)
    assert body["queue_items"][0]["sku"] == "SKU-JOB-1"
    assert body["queue_items"][0]["last_response"]["backend"] == "shopify"


def test_v2_publish_state_list_filters(v2_client: TestClient, v2_session: Session):
    _add_active_connection(v2_session, channel_code="dash-magento")
    first = v2_client.post(
        "/api/v2/products",
        json={
            "sku": "SKU-STATE-1",
            "title": "Publish State One",
            "product_type": "simple",
            "lifecycle_status": "active",
            "attributes": {},
            "category_ids": [],
            "collection_ids": [],
        },
    )
    second = v2_client.post(
        "/api/v2/products",
        json={
            "sku": "SKU-STATE-2",
            "title": "Publish State Two",
            "product_type": "simple",
            "lifecycle_status": "active",
            "attributes": {},
            "category_ids": [],
            "collection_ids": [],
        },
    )
    assert first.status_code == 200
    assert second.status_code == 200

    states = v2_session.execute(select(V2ChannelPublishState).order_by(V2ChannelPublishState.product_id.asc())).scalars().all()
    states[0].status = "verified"
    states[0].verification_result = {"status": "synced"}
    states[1].status = "failed"
    states[1].last_error = "No remote product"
    v2_session.commit()

    filtered = v2_client.get("/api/v2/publish-states", params={"status": "failed", "channel": "dash-magento"})
    assert filtered.status_code == 200
    body = filtered.json()
    assert body["total"] == 1
    assert body["items"][0]["sku"] == "SKU-STATE-2"
    assert body["items"][0]["status"] == "failed"
    assert body["items"][0]["last_error"] == "No remote product"

    by_sku = v2_client.get("/api/v2/publish-states", params={"sku": "SKU-STATE-1"})
    assert by_sku.status_code == 200
    by_sku_body = by_sku.json()
    assert by_sku_body["total"] == 1
    assert by_sku_body["items"][0]["sku"] == "SKU-STATE-1"
    assert by_sku_body["items"][0]["status"] == "verified"


def test_v2_runtime_status_reports_mode_connections_and_queue_depths(
    v2_client: TestClient,
    v2_session: Session,
    monkeypatch: pytest.MonkeyPatch,
):
    _add_active_connection(v2_session, channel_code="dash-magento")
    create_response = v2_client.post(
        "/api/v2/products",
        json={
            "sku": "SKU-RUNTIME-1",
            "title": "Runtime Product",
            "product_type": "simple",
            "lifecycle_status": "active",
            "attributes": {},
            "category_ids": [],
            "collection_ids": [],
        },
    )
    assert create_response.status_code == 200

    monkeypatch.setenv("V2_SYNC_EXECUTION_MODE", "live")
    monkeypatch.setenv("V2_SYNC_WORKER_ENABLED", "true")
    monkeypatch.setenv("V2_SYNC_STALE_MINUTES", "90")

    response = v2_client.get("/api/v2/ops/runtime")

    assert response.status_code == 200
    body = response.json()
    assert body["ok"] is True
    assert body["service"] == "plytixmage-v2"
    assert body["execution_mode"] == "live"
    assert body["embedded_worker_enabled"] is True
    assert body["stale_minutes"] == 90
    assert body["active_connections"] == 1
    assert body["queue_depths"]["product_data"]["pending"] == 1
