← 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;