pliceclogs-og/rfid-daemon/plice_server.py
type-two e5b8504bc6 Import pliceclogs as-is — the original Discogs seller extension
Snapshot of the working tree exactly as it stood, no edits. This is the predecessor
PRICEGOD was rewritten from ("kept intact, untouched" per pricegod/README.md), and it
is still the only place the DYMO scale is actually implemented —
rfid-daemon/index.js:1068-1265: HID discovery, parseScaleReport, one-shot read, and a
streaming /weight + /scale/start|stop session whose JSON shape PRICEGOD's daemon.js
already speaks.

Preserved verbatim on purpose (hence -og), including the known bug: DYMO_PIDS at
index.js:1082 is [0x8003, 0x8004], so it cannot see the bench M25 (0x8009). Fix that
in whatever daemon inherits the scale, not here.

node_modules stays ignored (22M of the 24M tree). No credentials in the import: the
two PEM markers in utils.js/sheets.js only strip headers off a key read from settings,
and mrpadmin / johnking are an SSH and a Postgres username, both key/trust auth.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-24 13:47:09 +10:00

689 lines
30 KiB
Python

"""
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()