← back to Omega Watches 2

database/schema.sql

596 lines

-- ============================================================================
-- OMEGA WATCHES 2.0 — ENTERPRISE DUAL-LEDGER SCHEMA
-- ============================================================================
-- Ledger 1: Primary Market MSRP (daily immutable snapshots)
-- Ledger 2: Secondary Market Events (auction, marketplace, dealer, forum)
-- ============================================================================

-- Create database (run separately as superuser)
-- CREATE DATABASE omega_watches_v2 OWNER omega_admin;

-- Extensions
CREATE EXTENSION IF NOT EXISTS "uuid-ossp";
CREATE EXTENSION IF NOT EXISTS "pg_trgm";
CREATE EXTENSION IF NOT EXISTS "btree_gin";

-- ============================================================================
-- ENUMS
-- ============================================================================

CREATE TYPE event_type_enum AS ENUM (
    'auction_realized',
    'marketplace_sold',
    'dealer_sold',
    'dealer_ask',
    'forum_listing',
    'price_guide'
);

CREATE TYPE seller_type_enum AS ENUM (
    'auction_house', 'dealer', 'private', 'marketplace', 'unknown'
);

CREATE TYPE fee_semantics_enum AS ENUM (
    'hammer_only', 'price_realised', 'total_to_buyer', 'unknown'
);

CREATE TYPE condition_grade_enum AS ENUM (
    'new_unworn', 'excellent', 'very_good', 'good', 'fair', 'poor', 'parts_only', 'unknown'
);

CREATE TYPE access_method_enum AS ENUM (
    'api', 'scrape', 'manual', 'export'
);

CREATE TYPE job_status_enum AS ENUM (
    'pending', 'running', 'completed', 'failed', 'cancelled'
);

CREATE TYPE severity_enum AS ENUM (
    'low', 'medium', 'high', 'critical'
);

-- ============================================================================
-- WATCH REFERENCE (Master Entity)
-- ============================================================================

CREATE TABLE watch_reference (
    id              UUID PRIMARY KEY DEFAULT uuid_generate_v4(),
    brand           TEXT NOT NULL DEFAULT 'Omega',
    collection      TEXT NOT NULL,
    model_name      TEXT NOT NULL,
    reference_number TEXT NOT NULL,
    calibre         TEXT,
    case_material   TEXT,
    case_diameter_mm NUMERIC(5,2),
    dial_color      TEXT,
    year_introduced  INTEGER,
    year_discontinued INTEGER,
    is_limited_edition BOOLEAN DEFAULT false,
    limited_edition_count INTEGER,
    notes           TEXT,
    search_vector   tsvector,
    created_at      TIMESTAMPTZ NOT NULL DEFAULT NOW(),
    updated_at      TIMESTAMPTZ NOT NULL DEFAULT NOW(),

    CONSTRAINT uq_reference UNIQUE (brand, reference_number),
    CONSTRAINT ck_year_range CHECK (year_introduced IS NULL OR (year_introduced >= 1848 AND year_introduced <= 2030)),
    CONSTRAINT ck_limited CHECK (
        (is_limited_edition = false AND limited_edition_count IS NULL)
        OR (is_limited_edition = true AND limited_edition_count > 0)
    )
);

CREATE INDEX idx_ref_collection ON watch_reference(collection);
CREATE INDEX idx_ref_brand ON watch_reference(brand);
CREATE INDEX idx_ref_reference ON watch_reference(reference_number);
CREATE INDEX idx_ref_year ON watch_reference(year_introduced);
CREATE INDEX idx_ref_search ON watch_reference USING GIN(search_vector);
CREATE INDEX idx_ref_model_trgm ON watch_reference USING GIN(model_name gin_trgm_ops);
CREATE INDEX idx_ref_refnum_trgm ON watch_reference USING GIN(reference_number gin_trgm_ops);

-- Auto-update search vector
CREATE OR REPLACE FUNCTION update_ref_search_vector() RETURNS TRIGGER AS $$
BEGIN
    NEW.search_vector :=
        setweight(to_tsvector('english', COALESCE(NEW.model_name, '')), 'A') ||
        setweight(to_tsvector('english', COALESCE(NEW.collection, '')), 'A') ||
        setweight(to_tsvector('english', COALESCE(NEW.reference_number, '')), 'B') ||
        setweight(to_tsvector('english', COALESCE(NEW.calibre, '')), 'C') ||
        setweight(to_tsvector('english', COALESCE(NEW.notes, '')), 'D');
    RETURN NEW;
END;
$$ LANGUAGE plpgsql;

CREATE TRIGGER trg_ref_search BEFORE INSERT OR UPDATE ON watch_reference
    FOR EACH ROW EXECUTE FUNCTION update_ref_search_vector();

-- ============================================================================
-- LEDGER 1: MSRP SNAPSHOTS (Immutable Daily Observations)
-- ============================================================================

CREATE TABLE msrp_snapshot (
    msrp_id         BIGSERIAL,
    reference_id    UUID NOT NULL REFERENCES watch_reference(id),
    region_code     TEXT NOT NULL,          -- US, CH, UK, EU, JP, HK, SG
    currency        CHAR(3) NOT NULL,
    msrp_amount     NUMERIC(14,2) NOT NULL CHECK (msrp_amount > 0),
    captured_date   DATE NOT NULL,
    source_url      TEXT NOT NULL,
    page_hash       TEXT,
    parser_version  TEXT NOT NULL,
    created_at      TIMESTAMPTZ NOT NULL DEFAULT NOW(),

    PRIMARY KEY (msrp_id, captured_date),
    CONSTRAINT uq_msrp_daily UNIQUE (reference_id, region_code, captured_date)
) PARTITION BY RANGE (captured_date);

-- Partitions by year
CREATE TABLE msrp_snapshot_2024 PARTITION OF msrp_snapshot
    FOR VALUES FROM ('2024-01-01') TO ('2025-01-01');
CREATE TABLE msrp_snapshot_2025 PARTITION OF msrp_snapshot
    FOR VALUES FROM ('2025-01-01') TO ('2026-01-01');
CREATE TABLE msrp_snapshot_2026 PARTITION OF msrp_snapshot
    FOR VALUES FROM ('2026-01-01') TO ('2027-01-01');
CREATE TABLE msrp_snapshot_2027 PARTITION OF msrp_snapshot
    FOR VALUES FROM ('2027-01-01') TO ('2028-01-01');
CREATE TABLE msrp_snapshot_2028 PARTITION OF msrp_snapshot
    FOR VALUES FROM ('2028-01-01') TO ('2029-01-01');

CREATE INDEX idx_msrp_ref ON msrp_snapshot(reference_id);
CREATE INDEX idx_msrp_date ON msrp_snapshot(captured_date DESC);
CREATE INDEX idx_msrp_ref_region ON msrp_snapshot(reference_id, region_code);
CREATE INDEX idx_msrp_ref_date ON msrp_snapshot(reference_id, captured_date DESC);

-- MSRP change detection view
CREATE OR REPLACE VIEW msrp_changes AS
SELECT
    b.reference_id,
    r.reference_number,
    r.model_name,
    b.region_code,
    b.currency,
    a.msrp_amount AS old_price,
    b.msrp_amount AS new_price,
    b.captured_date AS change_date,
    ROUND(((b.msrp_amount - a.msrp_amount) / a.msrp_amount) * 100, 2) AS pct_change
FROM msrp_snapshot a
JOIN msrp_snapshot b
    ON a.reference_id = b.reference_id
    AND a.region_code = b.region_code
    AND a.captured_date = b.captured_date - INTERVAL '1 day'
JOIN watch_reference r ON b.reference_id = r.id
WHERE a.msrp_amount <> b.msrp_amount;

-- ============================================================================
-- LEDGER 2: SECONDARY MARKET EVENTS
-- ============================================================================

CREATE TABLE market_event (
    event_id            BIGSERIAL PRIMARY KEY,
    -- Source identity (hard dedup key)
    source_name         TEXT NOT NULL,
    source_listing_id   TEXT NOT NULL,
    -- Event classification
    event_type          event_type_enum NOT NULL,
    -- Watch identity (may be null if unresolved)
    reference_id        UUID REFERENCES watch_reference(id),
    reference_number_raw TEXT,
    serial_number       TEXT,
    year_text           TEXT,
    case_material       TEXT,
    movement            TEXT,
    -- Condition
    condition_raw       TEXT,
    condition_normalized condition_grade_enum NOT NULL DEFAULT 'unknown',
    condition_score     INTEGER CHECK (condition_score IS NULL OR (condition_score >= 0 AND condition_score <= 100)),
    -- Provenance
    provenance_raw      TEXT,
    -- Sale metadata
    sale_date           DATE,
    location_text       TEXT,
    seller_type         seller_type_enum NOT NULL DEFAULT 'unknown',
    -- Price fields (explicit fee semantics)
    price_amount        NUMERIC(14,2),
    currency            CHAR(3) NOT NULL DEFAULT 'USD',
    hammer_price        NUMERIC(14,2),
    buyers_premium      NUMERIC(14,2),
    total_to_buyer      NUMERIC(14,2),
    fee_semantics       fee_semantics_enum NOT NULL DEFAULT 'unknown',
    -- FX normalization
    price_usd           NUMERIC(14,2),
    fx_rate             NUMERIC(12,6),
    fx_source           TEXT,
    -- Inflation adjustment
    price_usd_real      NUMERIC(14,2),
    cpi_base_year       INTEGER,
    -- Images and source
    images_json         JSONB,
    source_url          TEXT NOT NULL,
    scraped_at          TIMESTAMPTZ NOT NULL DEFAULT NOW(),
    page_hash           TEXT,
    parser_version      TEXT NOT NULL,
    -- Quality controls
    quality_score       INTEGER CHECK (quality_score IS NULL OR (quality_score >= 0 AND quality_score <= 100)),
    is_quarantined      BOOLEAN NOT NULL DEFAULT false,
    quarantine_reason   TEXT,
    -- Dedup
    dedup_cluster_id    UUID,
    -- Timestamps
    created_at          TIMESTAMPTZ NOT NULL DEFAULT NOW(),

    CONSTRAINT uq_source_listing UNIQUE (source_name, source_listing_id)
);

CREATE INDEX idx_me_source ON market_event(source_name);
CREATE INDEX idx_me_type ON market_event(event_type);
CREATE INDEX idx_me_ref ON market_event(reference_id);
CREATE INDEX idx_me_refnum_raw ON market_event(reference_number_raw);
CREATE INDEX idx_me_serial ON market_event(serial_number) WHERE serial_number IS NOT NULL;
CREATE INDEX idx_me_sale_date ON market_event(sale_date DESC) WHERE sale_date IS NOT NULL;
CREATE INDEX idx_me_price_usd ON market_event(price_usd DESC) WHERE price_usd IS NOT NULL;
CREATE INDEX idx_me_quarantine ON market_event(is_quarantined) WHERE is_quarantined = true;
CREATE INDEX idx_me_dedup ON market_event(dedup_cluster_id) WHERE dedup_cluster_id IS NOT NULL;
CREATE INDEX idx_me_scraped ON market_event(scraped_at DESC);
CREATE INDEX idx_me_condition ON market_event(condition_normalized);
CREATE INDEX idx_me_seller ON market_event(seller_type);
CREATE INDEX idx_me_refnum_trgm ON market_event USING GIN(reference_number_raw gin_trgm_ops);

-- ============================================================================
-- DATA SOURCES
-- ============================================================================

CREATE TABLE data_source (
    id              UUID PRIMARY KEY DEFAULT uuid_generate_v4(),
    name            TEXT UNIQUE NOT NULL,       -- sothebys, christies, ebay, chrono24, etc.
    display_name    TEXT NOT NULL,
    source_type     TEXT NOT NULL,              -- auction, marketplace, dealer, forum, manufacturer, index
    access_method   access_method_enum NOT NULL,
    base_url        TEXT,
    robots_txt_compliant BOOLEAN NOT NULL DEFAULT true,
    tos_reviewed    BOOLEAN NOT NULL DEFAULT false,
    tos_allows_collection BOOLEAN,
    tos_notes       TEXT,
    reliability_score INTEGER NOT NULL DEFAULT 3 CHECK (reliability_score >= 1 AND reliability_score <= 5),
    coverage_years_from INTEGER,
    coverage_years_to INTEGER,
    rate_limit_rpm  INTEGER,                    -- requests per minute
    notes           TEXT,
    is_active       BOOLEAN NOT NULL DEFAULT true,
    created_at      TIMESTAMPTZ NOT NULL DEFAULT NOW(),
    updated_at      TIMESTAMPTZ NOT NULL DEFAULT NOW()
);

-- ============================================================================
-- COLLECTOR RUNS (Job History)
-- ============================================================================

CREATE TABLE collector_run (
    run_id              BIGSERIAL PRIMARY KEY,
    source_id           UUID NOT NULL REFERENCES data_source(id),
    job_type            TEXT NOT NULL,           -- daily_msrp, backfill, incremental, full_crawl
    status              job_status_enum NOT NULL DEFAULT 'pending',
    started_at          TIMESTAMPTZ,
    completed_at        TIMESTAMPTZ,
    records_fetched     INTEGER NOT NULL DEFAULT 0,
    records_parsed      INTEGER NOT NULL DEFAULT 0,
    records_inserted    INTEGER NOT NULL DEFAULT 0,
    records_updated     INTEGER NOT NULL DEFAULT 0,
    records_quarantined INTEGER NOT NULL DEFAULT 0,
    records_deduplicated INTEGER NOT NULL DEFAULT 0,
    error_message       TEXT,
    parser_version      TEXT NOT NULL,
    duration_ms         INTEGER,
    metadata            JSONB,
    created_at          TIMESTAMPTZ NOT NULL DEFAULT NOW()
);

CREATE INDEX idx_cr_source ON collector_run(source_id);
CREATE INDEX idx_cr_status ON collector_run(status);
CREATE INDEX idx_cr_started ON collector_run(started_at DESC);
CREATE INDEX idx_cr_type ON collector_run(job_type);

-- ============================================================================
-- PARSER VERSIONS
-- ============================================================================

CREATE TABLE parser_version (
    id              UUID PRIMARY KEY DEFAULT uuid_generate_v4(),
    source_id       UUID NOT NULL REFERENCES data_source(id),
    version         TEXT NOT NULL,
    selectors_json  JSONB NOT NULL,
    dom_hash        TEXT,
    is_current      BOOLEAN NOT NULL DEFAULT false,
    tested_at       TIMESTAMPTZ,
    test_pass_rate  NUMERIC(5,2),
    notes           TEXT,
    created_at      TIMESTAMPTZ NOT NULL DEFAULT NOW(),

    CONSTRAINT uq_parser_version UNIQUE (source_id, version)
);

CREATE INDEX idx_pv_source ON parser_version(source_id);
CREATE INDEX idx_pv_current ON parser_version(is_current) WHERE is_current = true;

-- ============================================================================
-- QUARANTINE
-- ============================================================================

CREATE TABLE quarantine_record (
    id              BIGSERIAL PRIMARY KEY,
    event_id        BIGINT REFERENCES market_event(event_id),
    msrp_id         BIGINT,
    reason          TEXT NOT NULL,
    severity        severity_enum NOT NULL DEFAULT 'medium',
    details         JSONB,
    resolved        BOOLEAN NOT NULL DEFAULT false,
    resolved_by     TEXT,
    resolved_at     TIMESTAMPTZ,
    resolution_notes TEXT,
    created_at      TIMESTAMPTZ NOT NULL DEFAULT NOW()
);

CREATE INDEX idx_qr_resolved ON quarantine_record(resolved) WHERE resolved = false;
CREATE INDEX idx_qr_severity ON quarantine_record(severity);
CREATE INDEX idx_qr_event ON quarantine_record(event_id);

-- ============================================================================
-- CONDITION MAPPINGS (per-venue normalization)
-- ============================================================================

CREATE TABLE condition_mapping (
    id              UUID PRIMARY KEY DEFAULT uuid_generate_v4(),
    source_name     TEXT NOT NULL,
    source_category TEXT NOT NULL,
    normalized_grade condition_grade_enum NOT NULL,
    score_0_100     INTEGER NOT NULL CHECK (score_0_100 >= 0 AND score_0_100 <= 100),
    version         INTEGER NOT NULL DEFAULT 1,
    created_at      TIMESTAMPTZ NOT NULL DEFAULT NOW(),

    CONSTRAINT uq_condition_map UNIQUE (source_name, source_category, version)
);

-- ============================================================================
-- FX RATES
-- ============================================================================

CREATE TABLE fx_rate (
    id              BIGSERIAL PRIMARY KEY,
    currency_from   CHAR(3) NOT NULL,
    currency_to     CHAR(3) NOT NULL DEFAULT 'USD',
    rate            NUMERIC(12,6) NOT NULL CHECK (rate > 0),
    rate_date       DATE NOT NULL,
    source          TEXT NOT NULL,           -- fed_h10, ecb, manual
    created_at      TIMESTAMPTZ NOT NULL DEFAULT NOW(),

    CONSTRAINT uq_fx_daily UNIQUE (currency_from, currency_to, rate_date, source)
);

CREATE INDEX idx_fx_pair_date ON fx_rate(currency_from, currency_to, rate_date DESC);

-- ============================================================================
-- CPI DATA (for inflation adjustment)
-- ============================================================================

CREATE TABLE cpi_data (
    id              BIGSERIAL PRIMARY KEY,
    country_code    CHAR(2) NOT NULL,       -- US, UK, CH
    year            INTEGER NOT NULL,
    month           INTEGER CHECK (month >= 1 AND month <= 12),
    cpi_value       NUMERIC(10,4) NOT NULL,
    base_year       INTEGER NOT NULL,
    source          TEXT NOT NULL,           -- bls, ons
    created_at      TIMESTAMPTZ NOT NULL DEFAULT NOW(),

    CONSTRAINT uq_cpi UNIQUE (country_code, year, month, source)
);

-- ============================================================================
-- AUDIT LOG
-- ============================================================================

CREATE TABLE audit_log (
    id              BIGSERIAL,
    table_name      TEXT NOT NULL,
    record_id       TEXT NOT NULL,
    operation       TEXT NOT NULL CHECK (operation IN ('INSERT', 'UPDATE', 'DELETE')),
    old_values      JSONB,
    new_values      JSONB,
    changed_by      TEXT NOT NULL DEFAULT 'system',
    ip_address      INET,
    changed_at      TIMESTAMPTZ NOT NULL DEFAULT NOW(),

    PRIMARY KEY (id, changed_at)
) PARTITION BY RANGE (changed_at);

CREATE TABLE audit_log_2025 PARTITION OF audit_log
    FOR VALUES FROM ('2025-01-01') TO ('2026-01-01');
CREATE TABLE audit_log_2026 PARTITION OF audit_log
    FOR VALUES FROM ('2026-01-01') TO ('2027-01-01');
CREATE TABLE audit_log_2027 PARTITION OF audit_log
    FOR VALUES FROM ('2027-01-01') TO ('2028-01-01');

CREATE INDEX idx_audit_table ON audit_log(table_name, record_id);
CREATE INDEX idx_audit_time ON audit_log(changed_at DESC);

-- ============================================================================
-- RAW SNAPSHOTS (HTML/PDF archive for auditability)
-- ============================================================================

CREATE TABLE raw_snapshot (
    id              BIGSERIAL PRIMARY KEY,
    source_name     TEXT NOT NULL,
    url             TEXT NOT NULL,
    content_hash    TEXT NOT NULL,
    content_type    TEXT NOT NULL,           -- text/html, application/pdf
    content_size_bytes INTEGER,
    storage_path    TEXT,                    -- local filesystem path
    scraped_at      TIMESTAMPTZ NOT NULL DEFAULT NOW(),
    parser_version  TEXT,

    CONSTRAINT uq_snapshot UNIQUE (source_name, url, content_hash)
);

CREATE INDEX idx_rs_source ON raw_snapshot(source_name);
CREATE INDEX idx_rs_scraped ON raw_snapshot(scraped_at DESC);

-- ============================================================================
-- TRIGGERS
-- ============================================================================

-- Auto-update updated_at
CREATE OR REPLACE FUNCTION update_timestamp() RETURNS TRIGGER AS $$
BEGIN
    NEW.updated_at = NOW();
    RETURN NEW;
END;
$$ LANGUAGE plpgsql;

CREATE TRIGGER trg_ref_updated BEFORE UPDATE ON watch_reference
    FOR EACH ROW EXECUTE FUNCTION update_timestamp();
CREATE TRIGGER trg_source_updated BEFORE UPDATE ON data_source
    FOR EACH ROW EXECUTE FUNCTION update_timestamp();

-- Audit trigger for market_event
CREATE OR REPLACE FUNCTION audit_market_event() RETURNS TRIGGER AS $$
BEGIN
    IF TG_OP = 'INSERT' THEN
        INSERT INTO audit_log(table_name, record_id, operation, new_values)
        VALUES ('market_event', NEW.event_id::TEXT, 'INSERT', row_to_json(NEW)::jsonb);
        RETURN NEW;
    ELSIF TG_OP = 'UPDATE' THEN
        INSERT INTO audit_log(table_name, record_id, operation, old_values, new_values)
        VALUES ('market_event', NEW.event_id::TEXT, 'UPDATE', row_to_json(OLD)::jsonb, row_to_json(NEW)::jsonb);
        RETURN NEW;
    ELSIF TG_OP = 'DELETE' THEN
        INSERT INTO audit_log(table_name, record_id, operation, old_values)
        VALUES ('market_event', OLD.event_id::TEXT, 'DELETE', row_to_json(OLD)::jsonb);
        RETURN OLD;
    END IF;
END;
$$ LANGUAGE plpgsql;

CREATE TRIGGER trg_audit_market_event
    AFTER INSERT OR UPDATE OR DELETE ON market_event
    FOR EACH ROW EXECUTE FUNCTION audit_market_event();

-- ============================================================================
-- SEED DATA SOURCES (from spec)
-- ============================================================================

INSERT INTO data_source (name, display_name, source_type, access_method, base_url, reliability_score, tos_reviewed, tos_allows_collection, tos_notes, coverage_years_from, coverage_years_to, rate_limit_rpm) VALUES
    ('sothebys',    'Sotheby''s',           'auction',      'scrape',  'https://www.sothebys.com',       5, true, NULL, 'Login-gated results; treat as permissioned environment', 1960, 2026, 10),
    ('christies',   'Christie''s',          'auction',      'scrape',  'https://www.christies.com',      5, true, NULL, 'Price realised includes hammer + buyer''s premium', 1960, 2026, 10),
    ('phillips',    'Phillips',             'auction',      'scrape',  'https://www.phillips.com',       5, true, NULL, 'Strong watch focus; buyer''s premium schedules published', 1970, 2026, 10),
    ('antiquorum',  'Antiquorum',           'auction',      'scrape',  'https://www.antiquorum.swiss',   5, true, NULL, 'Graded condition breakdowns; movement/case numbers', 1970, 2026, 10),
    ('bonhams',     'Bonhams',              'auction',      'scrape',  'https://www.bonhams.com',        4, true, NULL, 'Dedicated watches department', 1980, 2026, 10),
    ('ebay',        'eBay',                 'marketplace',  'api',     'https://www.ebay.com',           4, true, false, 'ToS prohibits robots/scrapers without permission; Feb 20 2026 updated agreement', 2000, 2026, 30),
    ('chrono24',    'Chrono24',             'marketplace',  'scrape',  'https://www.chrono24.com',       4, true, NULL, 'robots.txt disallows some agents including Scrapy; structured listing attributes', 2005, 2026, 5),
    ('the1916co',   'The 1916 Company',     'dealer',       'scrape',  'https://www.the1916company.com', 3, false, NULL, 'Standard ToS; selection bias', 2010, 2026, 10),
    ('hodinkee',    'Hodinkee',             'dealer',       'manual',  'https://shop.hodinkee.com',      4, true, false, 'ToS prohibits robots/spiders/data mining/extraction tools', 2015, 2026, NULL),
    ('watchuseek',  'Watchuseek',           'forum',        'scrape',  'https://www.watchuseek.com',     3, false, NULL, 'Rules require leaving asking price after sold', 2000, 2026, 5),
    ('omegaforums', 'Omega Forums',         'forum',        'scrape',  'https://omegaforums.net',        3, false, NULL, 'Rule: do not delete listing or price when sold; requires price/currency', 2005, 2026, 5),
    ('watchcharts', 'WatchCharts',          'index',        'api',     'https://watchcharts.com',        3, true, NULL, 'Paid API with Data Credits; business pricing from $5k/yr; best for triangulation', 2018, 2026, 60),
    ('omega_official', 'Omega Official',    'manufacturer', 'scrape',  'https://www.omegawatches.com',   5, false, NULL, 'Official MSRP source; may need data-sharing agreement for long-term compliance', 1957, 2026, 5);

-- ============================================================================
-- INITIAL CONDITION MAPPINGS
-- ============================================================================

INSERT INTO condition_mapping (source_name, source_category, normalized_grade, score_0_100) VALUES
    -- Chrono24 mappings
    ('chrono24', 'New/Unworn',          'new_unworn',   100),
    ('chrono24', 'Very good',           'very_good',    75),
    ('chrono24', 'Good',                'good',         60),
    ('chrono24', 'Fair',                'fair',         40),
    ('chrono24', 'Incomplete',          'parts_only',   10),
    -- Auction mappings
    ('auction_generic', 'Mint',         'new_unworn',   95),
    ('auction_generic', 'Excellent',    'excellent',     85),
    ('auction_generic', 'Very Good',    'very_good',    75),
    ('auction_generic', 'Good',         'good',         60),
    ('auction_generic', 'Fair',         'fair',         40),
    ('auction_generic', 'Poor',         'poor',         20),
    -- eBay mappings
    ('ebay', 'New with tags',           'new_unworn',   100),
    ('ebay', 'New without tags',        'new_unworn',   95),
    ('ebay', 'Pre-owned',              'good',          60),
    ('ebay', 'For parts or not working', 'parts_only',  10);

-- ============================================================================
-- MATERIALIZED VIEWS
-- ============================================================================

-- Reference price summary
CREATE MATERIALIZED VIEW ref_price_summary AS
SELECT
    r.id AS reference_id,
    r.reference_number,
    r.model_name,
    r.collection,
    COUNT(DISTINCT me.event_id) AS event_count,
    COUNT(DISTINCT me.event_id) FILTER (WHERE me.event_type = 'auction_realized') AS auction_count,
    COUNT(DISTINCT me.event_id) FILTER (WHERE me.event_type = 'marketplace_sold') AS marketplace_count,
    MIN(me.price_usd) FILTER (WHERE me.price_usd > 0) AS min_price_usd,
    MAX(me.price_usd) AS max_price_usd,
    AVG(me.price_usd)::NUMERIC(14,2) AS avg_price_usd,
    PERCENTILE_CONT(0.5) WITHIN GROUP (ORDER BY me.price_usd) AS median_price_usd,
    MIN(me.sale_date) AS earliest_sale,
    MAX(me.sale_date) AS latest_sale,
    (SELECT ms.msrp_amount FROM msrp_snapshot ms
     WHERE ms.reference_id = r.id AND ms.region_code = 'US'
     ORDER BY ms.captured_date DESC LIMIT 1) AS current_msrp_usd,
    NOW() AS calculated_at
FROM watch_reference r
LEFT JOIN market_event me ON me.reference_id = r.id AND NOT me.is_quarantined
GROUP BY r.id, r.reference_number, r.model_name, r.collection;

CREATE UNIQUE INDEX idx_rps_ref ON ref_price_summary(reference_id);
CREATE INDEX idx_rps_collection ON ref_price_summary(collection);

-- Source health summary
CREATE MATERIALIZED VIEW source_health AS
SELECT
    ds.id AS source_id,
    ds.name,
    ds.display_name,
    ds.is_active,
    COUNT(cr.run_id) AS total_runs,
    COUNT(cr.run_id) FILTER (WHERE cr.status = 'completed') AS completed_runs,
    COUNT(cr.run_id) FILTER (WHERE cr.status = 'failed') AS failed_runs,
    MAX(cr.completed_at) AS last_completed,
    MAX(cr.started_at) AS last_started,
    AVG(cr.duration_ms)::INTEGER AS avg_duration_ms,
    SUM(cr.records_inserted) AS total_records_inserted,
    SUM(cr.records_quarantined) AS total_records_quarantined,
    CASE
        WHEN COUNT(cr.run_id) > 0 THEN
            ROUND(COUNT(cr.run_id) FILTER (WHERE cr.status = 'completed')::NUMERIC / COUNT(cr.run_id) * 100, 1)
        ELSE 0
    END AS success_rate_pct,
    NOW() AS calculated_at
FROM data_source ds
LEFT JOIN collector_run cr ON cr.source_id = ds.id
GROUP BY ds.id, ds.name, ds.display_name, ds.is_active;

CREATE UNIQUE INDEX idx_sh_source ON source_health(source_id);

-- Refresh function
CREATE OR REPLACE FUNCTION refresh_materialized_views() RETURNS void AS $$
BEGIN
    REFRESH MATERIALIZED VIEW CONCURRENTLY ref_price_summary;
    REFRESH MATERIALIZED VIEW CONCURRENTLY source_health;
END;
$$ LANGUAGE plpgsql;

-- ============================================================================
-- GRANTS
-- ============================================================================

GRANT ALL ON ALL TABLES IN SCHEMA public TO omega_admin;
GRANT ALL ON ALL SEQUENCES IN SCHEMA public TO omega_admin;
GRANT EXECUTE ON ALL FUNCTIONS IN SCHEMA public TO omega_admin;

-- Done
SELECT 'Omega Watches 2.0 schema created successfully' AS status;