""" PliceCogs → MONSTERWIKI receiver server. Listens on 0.0.0.0:5002/plice for POST requests from the PliceCogs Chrome extension whenever a record is priced. Writes: - release_market — live market snapshot (have/want/prices/suggestions) - release_review — full review text + ratings (unique per release/user/date) - collection_pricing — the pricing action itself (condition, price, SKU, folder) Run: python3 plice_server.py (keep running in background — e.g. screen/tmux on stupendo) """ import json import psycopg2 import psycopg2.extras import socket from datetime import datetime from http.server import BaseHTTPRequestHandler, HTTPServer DB_NAME = "discogs_full" PORT = 5002 # Detect if we are running on 'ultra' HOSTNAME = socket.gethostname() IS_ULTRA = (HOSTNAME == 'ultra.local' or HOSTNAME == 'ultra') ULTRA_TAILSCALE_IP = '100.91.239.7' def get_conn(): if IS_ULTRA: # Local connection on ultra return psycopg2.connect(dbname=DB_NAME) else: # Remote connection to ultra via Tailscale return psycopg2.connect( dbname=DB_NAME, host=ULTRA_TAILSCALE_IP, user='johnking' # Assuming johnking as per standard local setup ) def ensure_tables(conn): with conn.cursor() as cur: cur.execute(""" -- Full-res cover image on the release row (not in XML dump) ALTER TABLE release ADD COLUMN IF NOT EXISTS cover_image TEXT; ALTER TABLE release ADD COLUMN IF NOT EXISTS images JSONB; -- Live market snapshot per release (one row per release, upserted) CREATE TABLE IF NOT EXISTS release_market ( release_id INTEGER PRIMARY KEY, have INTEGER, want INTEGER, for_sale INTEGER, low_price NUMERIC(10,2), median_price NUMERIC(10,2), high_price NUMERIC(10,2), last_sold TEXT, price_suggestions JSONB, fetched_at TIMESTAMP DEFAULT NOW() ); -- Community reviews (one row per review, unique by release+user+date) CREATE TABLE IF NOT EXISTS release_review ( id SERIAL PRIMARY KEY, release_id INTEGER NOT NULL, username TEXT, review_date TEXT, rating INTEGER, review_text TEXT, helpful_count INTEGER, replies JSONB, fetched_at TIMESTAMP DEFAULT NOW(), UNIQUE (release_id, username, review_date) ); -- Every pricing action recorded (append-only log) CREATE TABLE IF NOT EXISTS collection_pricing ( id SERIAL PRIMARY KEY, release_id INTEGER NOT NULL, artist TEXT, title TEXT, label TEXT, year INTEGER, media_condition TEXT, sleeve_condition TEXT, price TEXT, sku TEXT, folder TEXT, comment TEXT, priced_at TIMESTAMP DEFAULT NOW() ); -- Reviews from master pages (all versions aggregated) CREATE TABLE IF NOT EXISTS master_review ( id SERIAL PRIMARY KEY, master_id INTEGER NOT NULL, username TEXT, review_date TEXT, rating INTEGER, review_text TEXT, helpful INTEGER, version_ref TEXT, -- "referencing X LP ZEN88" description string version_id INTEGER, -- Discogs release ID if available replies JSONB, fetched_at TIMESTAMP DEFAULT NOW(), UNIQUE (master_id, username, review_date) ); -- Individual sales history rows (append-only; de-duped by release+date+price) CREATE TABLE IF NOT EXISTS release_sale ( id SERIAL PRIMARY KEY, release_id INTEGER NOT NULL, sale_date TEXT, condition TEXT, sleeve TEXT, price NUMERIC(10,2), price_orig TEXT, currency TEXT, comment TEXT, fetched_at TIMESTAMP DEFAULT NOW(), UNIQUE (release_id, sale_date, price) ); -- Snapshot of current marketplace listings at time of capture CREATE TABLE IF NOT EXISTS release_listing_snapshot ( id SERIAL PRIMARY KEY, release_id INTEGER NOT NULL, condition TEXT, sleeve TEXT, price TEXT, price_value NUMERIC(10,2), currency TEXT, seller TEXT, seller_rating NUMERIC(3,1), seller_percent NUMERIC(5,2), location TEXT, shipping TEXT, notes TEXT, snapshot_at TIMESTAMP DEFAULT NOW() ); -- Historical record of market price changes over time (append-only) -- A new row is inserted only when have/want/for_sale or any price changes. CREATE TABLE IF NOT EXISTS release_market_history ( id SERIAL PRIMARY KEY, release_id INTEGER NOT NULL, have INTEGER, want INTEGER, for_sale INTEGER, low_price NUMERIC(10,2), median_price NUMERIC(10,2), high_price NUMERIC(10,2), last_sold TEXT, recorded_at TIMESTAMP DEFAULT NOW() ); CREATE INDEX IF NOT EXISTS idx_rmh_release_recorded ON release_market_history (release_id, recorded_at DESC); CREATE TABLE IF NOT EXISTS release_gemini ( id SERIAL PRIMARY KEY, release_id INTEGER NOT NULL, prompt_text TEXT, input_json JSONB, output_text TEXT, created_at TIMESTAMP DEFAULT NOW() ); CREATE INDEX IF NOT EXISTS idx_release_gemini_release_created ON release_gemini (release_id, created_at DESC); ALTER TABLE release_gemini ADD COLUMN IF NOT EXISTS output_json JSONB; """) conn.commit() class Handler(BaseHTTPRequestHandler): def log_message(self, fmt, *args): pass # silence default access log def _cors(self): self.send_header('Access-Control-Allow-Origin', '*') self.send_header('Access-Control-Allow-Headers', 'Content-Type') self.send_header('Access-Control-Allow-Methods', 'POST, GET, OPTIONS') # Required for Chrome Private Network Access (Tailscale IPs are "private") self.send_header('Access-Control-Allow-Private-Network', 'true') def do_OPTIONS(self): self.send_response(200) self._cors() self.end_headers() def do_GET(self): if self.path.startswith('/plice/inventory'): import urllib.parse query = urllib.parse.urlparse(self.path).query params = urllib.parse.parse_qs(query) release_id = params.get('release_id', [None])[0] if not release_id: self.send_response(400) self._cors() self.end_headers() self.wfile.write(json.dumps({'ok': False, 'error': 'Missing release_id'}).encode()) return try: import subprocess # SSH command to query the remote VPS database # Table: wp_rmp_disc_inventory, Fields: release_id, sku, price # We sort by sku (which is a timestamp) to get the "oldest" first as requested ssh_cmd = [ 'ssh', 'mrpadmin@100.123.123.64', f"mysql -N -e \"SELECT sku, price FROM wp_rmp_disc_inventory WHERE release_id = {int(release_id)} ORDER BY sku ASC;\"" ] result = subprocess.run(ssh_cmd, capture_output=True, text=True, timeout=10) if result.returncode != 0: raise Exception(f"SSH/MySQL error: {result.stderr.strip()}") lines = result.stdout.strip().split('\n') inventory = [] for line in lines: if not line.strip(): continue parts = line.split('\t') if len(parts) >= 2: inventory.append({'sku': parts[0], 'price': parts[1]}) resp_body = json.dumps({'ok': True, 'inventory': inventory}).encode() self.send_response(200) self._cors() self.send_header('Content-Type', 'application/json') self.send_header('Content-Length', len(resp_body)) self.end_headers() self.wfile.write(resp_body) except Exception as e: print(f" ERROR (inventory): {e}") self.send_response(500) self._cors() self.end_headers() self.wfile.write(json.dumps({'ok': False, 'error': str(e)}).encode()) return if self.path in ('/plice/health', '/plice', '/plice/gemini'): body = json.dumps({'ok': True, 'server': 'plice_server'}).encode() self.send_response(200) self._cors() self.send_header('Content-Type', 'application/json') self.send_header('Content-Length', len(body)) self.end_headers() self.wfile.write(body) else: self.send_response(404) self.end_headers() def do_POST(self): if self.path == '/plice/gemini': length = int(self.headers.get('Content-Length', 0)) body = self.rfile.read(length) try: payload = json.loads(body) rows_written = self._store_gemini(payload) resp_body = json.dumps({'ok': True, 'ai_rows': rows_written}).encode() self.send_response(200) self._cors() self.send_header('Content-Type', 'application/json') self.send_header('Content-Length', len(resp_body)) self.end_headers() self.wfile.write(resp_body) ts = datetime.now().strftime('%H:%M:%S') print(f"[{ts}] [gemini] release {payload.get('release_id') or (payload.get('release') or {}).get('id')}") except Exception as e: print(f" ERROR (gemini): {e}") self.send_response(500) self._cors() self.end_headers() self.wfile.write(json.dumps({'ok': False, 'error': str(e)}).encode()) return if self.path == '/plice/master': length = int(self.headers.get('Content-Length', 0)) body = self.rfile.read(length) try: payload = json.loads(body) rows_written = self._store_master(payload) resp_body = json.dumps({'ok': True, 'review_rows': rows_written}).encode() self.send_response(200) self._cors() self.send_header('Content-Type', 'application/json') self.send_header('Content-Length', len(resp_body)) self.end_headers() self.wfile.write(resp_body) ts = datetime.now().strftime('%H:%M:%S') print(f"[{ts}] [master] {payload.get('url','?')} " f"{payload.get('artist','')} — {payload.get('title','')} " f"→ {rows_written} reviews") except Exception as e: print(f" ERROR (master): {e}") self.send_response(500) self._cors() self.end_headers() self.wfile.write(json.dumps({'ok': False, 'error': str(e)}).encode()) return if self.path != '/plice': self.send_response(404) self.end_headers() return length = int(self.headers.get('Content-Length', 0)) body = self.rfile.read(length) try: payload = json.loads(body) market_rows, review_rows, pricing_rows, sale_rows, listing_rows = self._store(payload) resp_body = json.dumps({ 'ok': True, 'market_rows': market_rows, 'review_rows': review_rows, 'pricing_rows': pricing_rows, 'sale_rows': sale_rows, 'listing_rows': listing_rows, }).encode() self.send_response(200) self._cors() self.send_header('Content-Type', 'application/json') self.send_header('Content-Length', len(resp_body)) self.end_headers() self.wfile.write(resp_body) ts = datetime.now().strftime('%H:%M:%S') release = payload.get('release', {}) pricing = payload.get('pricing', {}) print( f"[{ts}] release {release.get('id')} " f"{release.get('artist')} — {release.get('title')} " f"| price={pricing.get('price') or '(no price)'} " f"| reviews={review_rows} sales={sale_rows} listings={listing_rows} market_upsert={market_rows}" ) except Exception as e: print(f" ERROR: {e}") self.send_response(500) self._cors() self.end_headers() self.wfile.write(json.dumps({'ok': False, 'error': str(e)}).encode()) def _store_master(self, payload): master_id = payload.get('master_id') if not master_id: return 0 rows = 0 conn = get_conn() try: with conn.cursor() as cur: for r in payload.get('reviews', []): try: cur.execute(""" INSERT INTO master_review (master_id, username, review_date, rating, review_text, helpful, version_ref, version_id, replies, fetched_at) VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,NOW()) ON CONFLICT (master_id, username, review_date) DO UPDATE SET rating = EXCLUDED.rating, review_text = EXCLUDED.review_text, helpful = EXCLUDED.helpful, version_ref = EXCLUDED.version_ref, version_id = EXCLUDED.version_id, replies = EXCLUDED.replies, fetched_at = NOW() """, ( master_id, r.get('username'), r.get('date'), r.get('rating'), r.get('text'), r.get('helpful') or 0, r.get('version_ref'), r.get('version_id'), json.dumps(r.get('replies') or []), )) rows += 1 except Exception: conn.rollback() continue conn.commit() finally: conn.close() return rows def _store_gemini(self, payload): release_id = payload.get('release_id') or (payload.get('release') or {}).get('id') or payload.get('id') if not release_id: return 0 output_text = payload.get('output_text') or '' output_json = payload.get('output_json') if output_json is None and output_text: start = output_text.find('{') end = output_text.rfind('}') if start != -1 and end != -1 and end > start: try: output_json = json.loads(output_text[start:end + 1]) except Exception: output_json = None conn = get_conn() try: with conn.cursor() as cur: cur.execute(""" INSERT INTO release_gemini (release_id, prompt_text, input_json, output_text, output_json, created_at) VALUES (%s, %s, %s::jsonb, %s, %s::jsonb, NOW()) """, ( release_id, payload.get('prompt'), json.dumps(payload.get('input_json')) if payload.get('input_json') is not None else None, output_text, json.dumps(output_json) if output_json is not None else None, )) conn.commit() finally: conn.close() return 1 def _store(self, payload): release = payload.get('release', {}) market = payload.get('market', {}) community = payload.get('community', {}) pricing = payload.get('pricing', {}) sales_history = payload.get('sales_history') or {} curr_listings = payload.get('current_listings') or [] release_id = release.get('id') or payload.get('id') if not release_id: return 0, 0, 0 conn = get_conn() market_rows = review_rows = pricing_rows = 0 sale_rows = listing_rows = 0 try: with conn.cursor() as cur: # ── release.cover_image + images (not in XML dump) ──────────── cover_image = release.get('cover_image') images_raw = release.get('images') # list of URIs from API if cover_image or images_raw: cur.execute(""" UPDATE release SET cover_image = COALESCE(%s, cover_image), images = COALESCE(%s::jsonb, images) WHERE id = %s """, ( cover_image, json.dumps(images_raw) if images_raw else None, release_id, )) # ── release_market upsert ────────────────────────────────────── prices = market.get('prices', {}) def to_numeric(v): if v is None: return None if isinstance(v, (int, float)): return v # Strip currency symbols like "$12.50" import re m = re.search(r'[\d.]+', str(v)) return float(m.group()) if m else None new_have = market.get('have') new_want = market.get('want') new_for_sale= market.get('for_sale') new_low = to_numeric(prices.get('low')) new_median = to_numeric(prices.get('median')) new_high = to_numeric(prices.get('high')) new_sold = market.get('last_sold') new_suggestions = json.dumps(market.get('price_suggestions') or {}) # Read previous values to detect changes cur.execute(""" SELECT have, want, for_sale, low_price, median_price, high_price, last_sold FROM release_market WHERE release_id = %s """, (release_id,)) prev = cur.fetchone() cur.execute(""" INSERT INTO release_market (release_id, have, want, for_sale, low_price, median_price, high_price, last_sold, price_suggestions, fetched_at) VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, NOW()) ON CONFLICT (release_id) DO UPDATE SET have = EXCLUDED.have, want = EXCLUDED.want, for_sale = EXCLUDED.for_sale, low_price = EXCLUDED.low_price, median_price = EXCLUDED.median_price, high_price = EXCLUDED.high_price, last_sold = EXCLUDED.last_sold, price_suggestions = EXCLUDED.price_suggestions, fetched_at = NOW() """, ( release_id, new_have, new_want, new_for_sale, new_low, new_median, new_high, new_sold, new_suggestions, )) market_rows = 1 # Append to history only when something actually changed def _changed(a, b): if a is None and b is None: return False return str(a) != str(b) is_first_scan = prev is None if is_first_scan or any(_changed(*p) for p in [ (new_have, prev[0]), (new_want, prev[1]), (new_for_sale, prev[2]), (new_low, prev[3]), (new_median, prev[4]), (new_high, prev[5]), (new_sold, prev[6]), ]): cur.execute(""" INSERT INTO release_market_history (release_id, have, want, for_sale, low_price, median_price, high_price, last_sold, recorded_at) VALUES (%s, %s, %s, %s, %s, %s, %s, %s, NOW()) """, ( release_id, new_have, new_want, new_for_sale, new_low, new_median, new_high, new_sold, )) # ── release_review upsert ────────────────────────────────────── for rev in community.get('reviews', []): username = rev.get('username') or '' rev_date = rev.get('date') or '' rating_raw = rev.get('rating') rating = int(rating_raw) if rating_raw is not None else None try: cur.execute(""" INSERT INTO release_review (release_id, username, review_date, rating, review_text, helpful_count, replies, fetched_at) VALUES (%s, %s, %s, %s, %s, %s, %s, NOW()) ON CONFLICT (release_id, username, review_date) DO UPDATE SET rating = EXCLUDED.rating, review_text = EXCLUDED.review_text, helpful_count = EXCLUDED.helpful_count, replies = EXCLUDED.replies, fetched_at = NOW() """, ( release_id, username, rev_date, rating, rev.get('text') or '', rev.get('helpful') or 0, json.dumps(rev.get('replies') or []), )) review_rows += 1 except Exception as re: conn.rollback() print(f" review insert error: {re}") continue # ── release_sale upsert ─────────────────────────────────────── currency = sales_history.get('currency') for sale in sales_history.get('sales', []): try: cur.execute(""" INSERT INTO release_sale (release_id, sale_date, condition, sleeve, price, price_orig, currency, comment, fetched_at) VALUES (%s, %s, %s, %s, %s, %s, %s, %s, NOW()) ON CONFLICT (release_id, sale_date, price) DO NOTHING """, ( release_id, sale.get('date'), sale.get('condition'), sale.get('sleeve'), sale.get('price'), sale.get('price_orig'), currency, sale.get('comment'), )) sale_rows += 1 except Exception: conn.rollback() continue # ── release_listing_snapshot — append only if listings changed ── if curr_listings: # Compare seller set against most recent snapshot to avoid duplicate entries cur.execute(""" SELECT seller, price, condition FROM release_listing_snapshot WHERE release_id = %s AND snapshot_at = ( SELECT MAX(snapshot_at) FROM release_listing_snapshot WHERE release_id = %s ) ORDER BY seller, price """, (release_id, release_id)) prev_snap = {(r[0], r[1], r[2]) for r in cur.fetchall()} new_snap = {(l.get('seller'), l.get('price'), l.get('condition')) for l in curr_listings} listings_changed = (prev_snap != new_snap) if listings_changed: for lst in curr_listings: try: cur.execute(""" INSERT INTO release_listing_snapshot (release_id, condition, sleeve, price, price_value, currency, seller, seller_rating, seller_percent, location, shipping, notes, snapshot_at) VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,NOW()) """, ( release_id, lst.get('condition'), lst.get('sleeve'), lst.get('price'), lst.get('price_value'), lst.get('currency'), lst.get('seller'), lst.get('seller_rating'), lst.get('seller_percent'), lst.get('location'), lst.get('shipping'), lst.get('notes'), )) listing_rows += 1 except Exception: conn.rollback() continue # ── collection_pricing append ────────────────────────────────── if pricing: cur.execute(""" INSERT INTO collection_pricing (release_id, artist, title, label, year, media_condition, sleeve_condition, price, sku, folder, comment, priced_at) VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, NOW()) """, ( release_id, release.get('artist'), release.get('title'), release.get('label'), release.get('year'), pricing.get('media_condition'), pricing.get('sleeve_condition'), pricing.get('price'), pricing.get('sku'), pricing.get('folder'), pricing.get('comment'), )) pricing_rows = 1 conn.commit() finally: conn.close() return market_rows, review_rows, pricing_rows, sale_rows, listing_rows if __name__ == '__main__': conn = get_conn() ensure_tables(conn) conn.close() print(f"PliceCogs server listening on 0.0.0.0:{PORT}") if IS_ULTRA: print(f"Running LOCALLY on {HOSTNAME}. Database connection is local.") print(f"Extension should use http://localhost:{PORT}/plice") else: print(f"Running REMOTELY on {HOSTNAME}. Database connection → {ULTRA_TAILSCALE_IP}") print(f"Extension should use http://{ULTRA_TAILSCALE_IP}:{PORT}/plice") HTTPServer(('0.0.0.0', PORT), Handler).serve_forever()