Spaces:
Sleeping
Sleeping
Add persistent Postgres storage (dual SQLite/Postgres db.py + psycopg)
Browse files- rmn-backend/README.md +5 -1
- rmn-backend/db.py +145 -109
- rmn-backend/requirements.txt +2 -0
rmn-backend/README.md
CHANGED
|
@@ -78,7 +78,11 @@ Notes:
|
|
| 78 |
an affiliate feed (Amazon Product Advertising API, Impact, CJ, Rakuten). Because the rest
|
| 79 |
of the app reads from the catalog/DB, this is an isolated swap. Affiliate links here are
|
| 80 |
also where the revenue model lives.
|
| 81 |
-
6. **Deploy**
|
|
|
|
|
|
|
|
|
|
|
|
|
| 82 |
|
| 83 |
## Moving this into Claude Code
|
| 84 |
This folder is the seed. In Claude Code, point it at this directory and ask, step by step:
|
|
|
|
| 78 |
an affiliate feed (Amazon Product Advertising API, Impact, CJ, Rakuten). Because the rest
|
| 79 |
of the app reads from the catalog/DB, this is an isolated swap. Affiliate links here are
|
| 80 |
also where the revenue model lives.
|
| 81 |
+
6. ~~**Deploy**~~ ✅ Done — live on Hugging Face Spaces (Docker) at
|
| 82 |
+
https://horoburger-rmn.hf.space. The FastAPI app serves both the API and the static
|
| 83 |
+
frontend from one service (see `/Dockerfile` and `/README.md` at the repo root). Storage
|
| 84 |
+
is ephemeral there, so the remaining production step is a managed Postgres (Neon/Supabase
|
| 85 |
+
free tier) for durable accounts/collections.
|
| 86 |
|
| 87 |
## Moving this into Claude Code
|
| 88 |
This folder is the seed. In Claude Code, point it at this directory and ask, step by step:
|
rmn-backend/db.py
CHANGED
|
@@ -1,17 +1,26 @@
|
|
| 1 |
"""
|
| 2 |
-
RMN persistence layer
|
| 3 |
|
| 4 |
-
|
| 5 |
-
|
| 6 |
-
|
|
|
|
|
|
|
| 7 |
|
| 8 |
-
|
| 9 |
-
the function signatures below are the contract the rest of the app uses.
|
| 10 |
"""
|
| 11 |
-
import
|
| 12 |
from datetime import datetime, timezone
|
| 13 |
|
| 14 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 15 |
|
| 16 |
|
| 17 |
class DuplicateEmail(Exception):
|
|
@@ -19,67 +28,95 @@ class DuplicateEmail(Exception):
|
|
| 19 |
|
| 20 |
|
| 21 |
def _conn():
|
|
|
|
|
|
|
| 22 |
conn = sqlite3.connect(DB_PATH)
|
| 23 |
conn.row_factory = sqlite3.Row
|
| 24 |
conn.execute("PRAGMA foreign_keys = ON")
|
| 25 |
return conn
|
| 26 |
|
| 27 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 28 |
def _now():
|
| 29 |
return datetime.now(timezone.utc).isoformat()
|
| 30 |
|
| 31 |
|
| 32 |
def init_db():
|
| 33 |
with _conn() as c:
|
| 34 |
-
|
| 35 |
-
|
| 36 |
-
|
| 37 |
-
|
| 38 |
-
|
| 39 |
-
|
| 40 |
-
|
| 41 |
-
|
| 42 |
-
|
| 43 |
-
|
| 44 |
-
|
| 45 |
-
|
| 46 |
-
|
| 47 |
-
|
| 48 |
-
|
| 49 |
-
|
| 50 |
-
|
| 51 |
-
|
| 52 |
-
|
| 53 |
-
|
| 54 |
-
|
| 55 |
-
|
| 56 |
-
|
| 57 |
-
|
| 58 |
-
|
| 59 |
-
|
| 60 |
-
|
| 61 |
-
|
| 62 |
-
|
| 63 |
-
|
| 64 |
-
|
| 65 |
-
|
| 66 |
-
|
| 67 |
-
|
| 68 |
-
|
| 69 |
-
|
| 70 |
-
|
| 71 |
-
|
| 72 |
-
|
| 73 |
-
|
| 74 |
-
|
| 75 |
-
|
| 76 |
-
|
| 77 |
-
|
| 78 |
-
|
| 79 |
-
|
| 80 |
-
|
| 81 |
-
|
| 82 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 83 |
|
| 84 |
|
| 85 |
# ---------- passwords ----------
|
|
@@ -100,20 +137,26 @@ def _check_password(password, stored):
|
|
| 100 |
|
| 101 |
# ---------- users & sessions ----------
|
| 102 |
def create_user(email, password):
|
| 103 |
-
|
| 104 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 105 |
cur = c.execute(
|
| 106 |
-
"INSERT INTO users(email, pw_hash, created_at) VALUES(?,?,?)",
|
| 107 |
-
|
| 108 |
-
|
| 109 |
-
|
| 110 |
raise DuplicateEmail(email)
|
| 111 |
-
|
| 112 |
|
| 113 |
|
| 114 |
def verify_user(email, password):
|
| 115 |
with _conn() as c:
|
| 116 |
-
row = c.execute("SELECT id, pw_hash FROM users WHERE email=?", (email,)).fetchone()
|
| 117 |
if not row or not _check_password(password, row["pw_hash"]):
|
| 118 |
return None
|
| 119 |
return row["id"]
|
|
@@ -122,7 +165,7 @@ def verify_user(email, password):
|
|
| 122 |
def create_session(user_id):
|
| 123 |
token = secrets.token_urlsafe(32)
|
| 124 |
with _conn() as c:
|
| 125 |
-
c.execute("INSERT INTO sessions(token, user_id, created_at) VALUES(?,?,?)",
|
| 126 |
(token, user_id, _now()))
|
| 127 |
return token
|
| 128 |
|
|
@@ -131,43 +174,41 @@ def user_for_token(token):
|
|
| 131 |
if not token:
|
| 132 |
return None
|
| 133 |
with _conn() as c:
|
| 134 |
-
row = c.execute("SELECT user_id FROM sessions WHERE token=?", (token,)).fetchone()
|
| 135 |
return row["user_id"] if row else None
|
| 136 |
|
| 137 |
|
| 138 |
def delete_session(token):
|
| 139 |
with _conn() as c:
|
| 140 |
-
c.execute("DELETE FROM sessions WHERE token=?", (token,))
|
| 141 |
|
| 142 |
|
| 143 |
# ---------- collections ----------
|
| 144 |
def add_saved(user_id, collection, product_id):
|
| 145 |
with _conn() as c:
|
| 146 |
-
c.execute(
|
| 147 |
-
|
| 148 |
-
|
| 149 |
|
| 150 |
|
| 151 |
def remove_saved(user_id, collection, product_id):
|
| 152 |
with _conn() as c:
|
| 153 |
-
c.execute(
|
| 154 |
-
|
| 155 |
-
(user_id, collection, product_id))
|
| 156 |
|
| 157 |
|
| 158 |
def get_collections(user_id):
|
| 159 |
"""Return {collection_name: [product_id, ...]} ordered by when saved."""
|
| 160 |
with _conn() as c:
|
| 161 |
-
rows = c.execute(
|
| 162 |
-
|
| 163 |
-
"WHERE user_id=? ORDER BY created_at", (user_id,)).fetchall()
|
| 164 |
out = {}
|
| 165 |
for r in rows:
|
| 166 |
out.setdefault(r["collection"], []).append(r["product_id"])
|
| 167 |
return out
|
| 168 |
|
| 169 |
|
| 170 |
-
# ---------- product catalog (
|
| 171 |
_PRODUCT_COLS = ["id", "name", "brand", "category", "subcategory", "color",
|
| 172 |
"price", "prev_price", "on_sale", "rating", "reviews", "seller",
|
| 173 |
"in_stock", "emoji", "specs", "description", "text"]
|
|
@@ -188,7 +229,8 @@ def catalog_count():
|
|
| 188 |
|
| 189 |
|
| 190 |
def import_products(products):
|
| 191 |
-
|
|
|
|
| 192 |
now = _now()
|
| 193 |
with _conn() as c:
|
| 194 |
c.executemany(
|
|
@@ -200,24 +242,22 @@ def import_products(products):
|
|
| 200 |
p["price"], p["prev_price"], int(p["on_sale"]), p["rating"], p["reviews"],
|
| 201 |
p["seller"], int(p["in_stock"]), p["emoji"], json.dumps(p["specs"]),
|
| 202 |
p["description"], p["_text"]) for p in products])
|
| 203 |
-
# seed one price-history point per product so tracking has a baseline
|
| 204 |
c.executemany("INSERT INTO price_history(product_id, price, ts) VALUES(?,?,?)",
|
| 205 |
[(p["id"], p["price"], now) for p in products])
|
| 206 |
|
| 207 |
|
| 208 |
def ensure_catalog(json_path="catalog.json"):
|
| 209 |
-
|
| 210 |
-
if catalog_count() > 0:
|
| 211 |
return
|
| 212 |
if not os.path.exists(json_path):
|
| 213 |
-
raise FileNotFoundError(
|
| 214 |
-
f"{json_path} not found — run `python generate_catalog.py` first.")
|
| 215 |
with open(json_path) as f:
|
| 216 |
import_products(json.load(f))
|
| 217 |
|
| 218 |
|
| 219 |
def reseed_catalog_from_json(json_path="catalog.json"):
|
| 220 |
-
|
|
|
|
| 221 |
with _conn() as c:
|
| 222 |
c.execute("DELETE FROM products")
|
| 223 |
c.execute("DELETE FROM price_history")
|
|
@@ -226,13 +266,14 @@ def reseed_catalog_from_json(json_path="catalog.json"):
|
|
| 226 |
|
| 227 |
|
| 228 |
def load_catalog():
|
|
|
|
|
|
|
| 229 |
with _conn() as c:
|
| 230 |
rows = c.execute("SELECT * FROM products").fetchall()
|
| 231 |
return [_row_to_product(r) for r in rows]
|
| 232 |
|
| 233 |
|
| 234 |
def category_breakdown():
|
| 235 |
-
"""Counts per product type (subcategory), with its parent category."""
|
| 236 |
with _conn() as c:
|
| 237 |
rows = c.execute(
|
| 238 |
"SELECT category, subcategory, COUNT(*) AS n FROM products "
|
|
@@ -243,56 +284,52 @@ def category_breakdown():
|
|
| 243 |
|
| 244 |
def update_price(product_id, new_price, prev_price):
|
| 245 |
with _conn() as c:
|
| 246 |
-
c.execute("UPDATE products SET price=?, prev_price=?, on_sale=? WHERE id=?",
|
| 247 |
(new_price, prev_price, int(new_price < prev_price), product_id))
|
| 248 |
-
c.execute("INSERT INTO price_history(product_id, price, ts) VALUES(?,?,?)",
|
| 249 |
(product_id, new_price, _now()))
|
| 250 |
|
| 251 |
|
| 252 |
def price_history(product_id, limit=30):
|
| 253 |
with _conn() as c:
|
| 254 |
-
rows = c.execute(
|
| 255 |
-
|
| 256 |
-
"ORDER BY ts DESC LIMIT ?", (product_id, limit)).fetchall()
|
| 257 |
return [{"price": r["price"], "ts": r["ts"]} for r in reversed(rows)]
|
| 258 |
|
| 259 |
|
| 260 |
# ---------- price tracking ----------
|
| 261 |
def add_watch(user_id, product_id, target_price=None):
|
| 262 |
with _conn() as c:
|
| 263 |
-
c.execute(
|
| 264 |
-
|
| 265 |
-
|
| 266 |
-
|
| 267 |
-
(user_id, product_id, target_price, _now()))
|
| 268 |
|
| 269 |
|
| 270 |
def remove_watch(user_id, product_id):
|
| 271 |
with _conn() as c:
|
| 272 |
-
c.execute("DELETE FROM price_watches WHERE user_id=? AND product_id=?",
|
| 273 |
(user_id, product_id))
|
| 274 |
|
| 275 |
|
| 276 |
def get_watches(user_id):
|
| 277 |
with _conn() as c:
|
| 278 |
-
rows = c.execute(
|
| 279 |
-
|
| 280 |
-
"WHERE user_id=? ORDER BY created_at DESC", (user_id,)).fetchall()
|
| 281 |
return [{"product_id": r["product_id"], "target_price": r["target_price"],
|
| 282 |
"created_at": r["created_at"]} for r in rows]
|
| 283 |
|
| 284 |
|
| 285 |
def watchers_for(product_id):
|
| 286 |
with _conn() as c:
|
| 287 |
-
rows = c.execute(
|
| 288 |
-
|
| 289 |
-
(product_id,)).fetchall()
|
| 290 |
return [(r["user_id"], r["target_price"]) for r in rows]
|
| 291 |
|
| 292 |
|
| 293 |
def watched_product_ids(user_id):
|
| 294 |
with _conn() as c:
|
| 295 |
-
rows = c.execute("SELECT product_id FROM price_watches WHERE user_id=?",
|
| 296 |
(user_id,)).fetchall()
|
| 297 |
return [r["product_id"] for r in rows]
|
| 298 |
|
|
@@ -300,9 +337,8 @@ def watched_product_ids(user_id):
|
|
| 300 |
# ---------- alerts ----------
|
| 301 |
def add_alert(user_id, product_id, old_price, new_price):
|
| 302 |
with _conn() as c:
|
| 303 |
-
c.execute(
|
| 304 |
-
|
| 305 |
-
"VALUES(?,?,?,?,?)", (user_id, product_id, old_price, new_price, _now()))
|
| 306 |
|
| 307 |
|
| 308 |
def get_alerts(user_id, unseen_only=False):
|
|
@@ -311,10 +347,10 @@ def get_alerts(user_id, unseen_only=False):
|
|
| 311 |
q += " AND seen=0"
|
| 312 |
q += " ORDER BY ts DESC LIMIT 50"
|
| 313 |
with _conn() as c:
|
| 314 |
-
rows = c.execute(q, (user_id,)).fetchall()
|
| 315 |
return [dict(r) for r in rows]
|
| 316 |
|
| 317 |
|
| 318 |
def mark_alerts_seen(user_id):
|
| 319 |
with _conn() as c:
|
| 320 |
-
c.execute("UPDATE alerts SET seen=1 WHERE user_id=?", (user_id,))
|
|
|
|
| 1 |
"""
|
| 2 |
+
RMN persistence layer.
|
| 3 |
|
| 4 |
+
Works with two backends, chosen at runtime:
|
| 5 |
+
* SQLite (default) — a local file (rmn.db); zero setup, used for dev.
|
| 6 |
+
* PostgreSQL — used automatically when the DATABASE_URL env var is set
|
| 7 |
+
(e.g. a free Neon/Supabase database), so accounts/collections/tracking
|
| 8 |
+
persist permanently in production.
|
| 9 |
|
| 10 |
+
Same public function signatures either way — the rest of the app is unchanged.
|
|
|
|
| 11 |
"""
|
| 12 |
+
import hashlib, hmac, secrets, json, os
|
| 13 |
from datetime import datetime, timezone
|
| 14 |
|
| 15 |
+
DATABASE_URL = os.environ.get("DATABASE_URL") # set this -> use Postgres
|
| 16 |
+
PG = bool(DATABASE_URL)
|
| 17 |
+
DB_PATH = os.environ.get("RMN_DB", "rmn.db") # SQLite path (writable on hosts)
|
| 18 |
+
|
| 19 |
+
if PG:
|
| 20 |
+
import psycopg
|
| 21 |
+
from psycopg.rows import dict_row
|
| 22 |
+
else:
|
| 23 |
+
import sqlite3
|
| 24 |
|
| 25 |
|
| 26 |
class DuplicateEmail(Exception):
|
|
|
|
| 28 |
|
| 29 |
|
| 30 |
def _conn():
|
| 31 |
+
if PG:
|
| 32 |
+
return psycopg.connect(DATABASE_URL, row_factory=dict_row)
|
| 33 |
conn = sqlite3.connect(DB_PATH)
|
| 34 |
conn.row_factory = sqlite3.Row
|
| 35 |
conn.execute("PRAGMA foreign_keys = ON")
|
| 36 |
return conn
|
| 37 |
|
| 38 |
|
| 39 |
+
def _q(sql):
|
| 40 |
+
"""SQLite uses ? placeholders; Postgres uses %s."""
|
| 41 |
+
return sql.replace("?", "%s") if PG else sql
|
| 42 |
+
|
| 43 |
+
|
| 44 |
+
def _is_dupe(e):
|
| 45 |
+
return type(e).__name__ in ("IntegrityError", "UniqueViolation")
|
| 46 |
+
|
| 47 |
+
|
| 48 |
def _now():
|
| 49 |
return datetime.now(timezone.utc).isoformat()
|
| 50 |
|
| 51 |
|
| 52 |
def init_db():
|
| 53 |
with _conn() as c:
|
| 54 |
+
if PG:
|
| 55 |
+
stmts = [
|
| 56 |
+
"""CREATE TABLE IF NOT EXISTS users(
|
| 57 |
+
id SERIAL PRIMARY KEY, email TEXT UNIQUE NOT NULL,
|
| 58 |
+
pw_hash TEXT NOT NULL, created_at TEXT NOT NULL)""",
|
| 59 |
+
"""CREATE TABLE IF NOT EXISTS sessions(
|
| 60 |
+
token TEXT PRIMARY KEY,
|
| 61 |
+
user_id INTEGER NOT NULL REFERENCES users(id) ON DELETE CASCADE,
|
| 62 |
+
created_at TEXT NOT NULL)""",
|
| 63 |
+
"""CREATE TABLE IF NOT EXISTS saved_items(
|
| 64 |
+
user_id INTEGER NOT NULL REFERENCES users(id) ON DELETE CASCADE,
|
| 65 |
+
collection TEXT NOT NULL, product_id INTEGER NOT NULL,
|
| 66 |
+
created_at TEXT NOT NULL,
|
| 67 |
+
PRIMARY KEY(user_id, collection, product_id))""",
|
| 68 |
+
"""CREATE TABLE IF NOT EXISTS products(
|
| 69 |
+
id INTEGER PRIMARY KEY, name TEXT, brand TEXT, category TEXT,
|
| 70 |
+
subcategory TEXT, color TEXT, price DOUBLE PRECISION,
|
| 71 |
+
prev_price DOUBLE PRECISION, on_sale INTEGER, rating DOUBLE PRECISION,
|
| 72 |
+
reviews INTEGER, seller TEXT, in_stock INTEGER, emoji TEXT,
|
| 73 |
+
specs TEXT, description TEXT, text TEXT)""",
|
| 74 |
+
"""CREATE TABLE IF NOT EXISTS price_history(
|
| 75 |
+
product_id INTEGER NOT NULL, price DOUBLE PRECISION NOT NULL, ts TEXT NOT NULL)""",
|
| 76 |
+
"""CREATE TABLE IF NOT EXISTS price_watches(
|
| 77 |
+
user_id INTEGER NOT NULL REFERENCES users(id) ON DELETE CASCADE,
|
| 78 |
+
product_id INTEGER NOT NULL, target_price DOUBLE PRECISION,
|
| 79 |
+
created_at TEXT NOT NULL, PRIMARY KEY(user_id, product_id))""",
|
| 80 |
+
"""CREATE TABLE IF NOT EXISTS alerts(
|
| 81 |
+
id SERIAL PRIMARY KEY,
|
| 82 |
+
user_id INTEGER NOT NULL REFERENCES users(id) ON DELETE CASCADE,
|
| 83 |
+
product_id INTEGER NOT NULL, old_price DOUBLE PRECISION NOT NULL,
|
| 84 |
+
new_price DOUBLE PRECISION NOT NULL, ts TEXT NOT NULL,
|
| 85 |
+
seen INTEGER NOT NULL DEFAULT 0)""",
|
| 86 |
+
]
|
| 87 |
+
for s in stmts:
|
| 88 |
+
c.execute(s)
|
| 89 |
+
else:
|
| 90 |
+
c.executescript("""
|
| 91 |
+
CREATE TABLE IF NOT EXISTS users(
|
| 92 |
+
id INTEGER PRIMARY KEY AUTOINCREMENT, email TEXT UNIQUE NOT NULL,
|
| 93 |
+
pw_hash TEXT NOT NULL, created_at TEXT NOT NULL);
|
| 94 |
+
CREATE TABLE IF NOT EXISTS sessions(
|
| 95 |
+
token TEXT PRIMARY KEY,
|
| 96 |
+
user_id INTEGER NOT NULL REFERENCES users(id) ON DELETE CASCADE,
|
| 97 |
+
created_at TEXT NOT NULL);
|
| 98 |
+
CREATE TABLE IF NOT EXISTS saved_items(
|
| 99 |
+
user_id INTEGER NOT NULL REFERENCES users(id) ON DELETE CASCADE,
|
| 100 |
+
collection TEXT NOT NULL, product_id INTEGER NOT NULL,
|
| 101 |
+
created_at TEXT NOT NULL, PRIMARY KEY(user_id, collection, product_id));
|
| 102 |
+
CREATE TABLE IF NOT EXISTS products(
|
| 103 |
+
id INTEGER PRIMARY KEY, name TEXT, brand TEXT, category TEXT,
|
| 104 |
+
subcategory TEXT, color TEXT, price REAL, prev_price REAL, on_sale INTEGER,
|
| 105 |
+
rating REAL, reviews INTEGER, seller TEXT, in_stock INTEGER,
|
| 106 |
+
emoji TEXT, specs TEXT, description TEXT, text TEXT);
|
| 107 |
+
CREATE INDEX IF NOT EXISTS idx_products_sub ON products(subcategory);
|
| 108 |
+
CREATE TABLE IF NOT EXISTS price_history(
|
| 109 |
+
product_id INTEGER NOT NULL, price REAL NOT NULL, ts TEXT NOT NULL);
|
| 110 |
+
CREATE TABLE IF NOT EXISTS price_watches(
|
| 111 |
+
user_id INTEGER NOT NULL REFERENCES users(id) ON DELETE CASCADE,
|
| 112 |
+
product_id INTEGER NOT NULL, target_price REAL, created_at TEXT NOT NULL,
|
| 113 |
+
PRIMARY KEY(user_id, product_id));
|
| 114 |
+
CREATE TABLE IF NOT EXISTS alerts(
|
| 115 |
+
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
| 116 |
+
user_id INTEGER NOT NULL REFERENCES users(id) ON DELETE CASCADE,
|
| 117 |
+
product_id INTEGER NOT NULL, old_price REAL NOT NULL, new_price REAL NOT NULL,
|
| 118 |
+
ts TEXT NOT NULL, seen INTEGER NOT NULL DEFAULT 0);
|
| 119 |
+
""")
|
| 120 |
|
| 121 |
|
| 122 |
# ---------- passwords ----------
|
|
|
|
| 137 |
|
| 138 |
# ---------- users & sessions ----------
|
| 139 |
def create_user(email, password):
|
| 140 |
+
pw, now = _hash_password(password), _now()
|
| 141 |
+
try:
|
| 142 |
+
with _conn() as c:
|
| 143 |
+
if PG:
|
| 144 |
+
row = c.execute(
|
| 145 |
+
"INSERT INTO users(email, pw_hash, created_at) VALUES(%s,%s,%s) RETURNING id",
|
| 146 |
+
(email, pw, now)).fetchone()
|
| 147 |
+
return row["id"]
|
| 148 |
cur = c.execute(
|
| 149 |
+
"INSERT INTO users(email, pw_hash, created_at) VALUES(?,?,?)", (email, pw, now))
|
| 150 |
+
return cur.lastrowid
|
| 151 |
+
except Exception as e:
|
| 152 |
+
if _is_dupe(e):
|
| 153 |
raise DuplicateEmail(email)
|
| 154 |
+
raise
|
| 155 |
|
| 156 |
|
| 157 |
def verify_user(email, password):
|
| 158 |
with _conn() as c:
|
| 159 |
+
row = c.execute(_q("SELECT id, pw_hash FROM users WHERE email=?"), (email,)).fetchone()
|
| 160 |
if not row or not _check_password(password, row["pw_hash"]):
|
| 161 |
return None
|
| 162 |
return row["id"]
|
|
|
|
| 165 |
def create_session(user_id):
|
| 166 |
token = secrets.token_urlsafe(32)
|
| 167 |
with _conn() as c:
|
| 168 |
+
c.execute(_q("INSERT INTO sessions(token, user_id, created_at) VALUES(?,?,?)"),
|
| 169 |
(token, user_id, _now()))
|
| 170 |
return token
|
| 171 |
|
|
|
|
| 174 |
if not token:
|
| 175 |
return None
|
| 176 |
with _conn() as c:
|
| 177 |
+
row = c.execute(_q("SELECT user_id FROM sessions WHERE token=?"), (token,)).fetchone()
|
| 178 |
return row["user_id"] if row else None
|
| 179 |
|
| 180 |
|
| 181 |
def delete_session(token):
|
| 182 |
with _conn() as c:
|
| 183 |
+
c.execute(_q("DELETE FROM sessions WHERE token=?"), (token,))
|
| 184 |
|
| 185 |
|
| 186 |
# ---------- collections ----------
|
| 187 |
def add_saved(user_id, collection, product_id):
|
| 188 |
with _conn() as c:
|
| 189 |
+
c.execute(_q("INSERT INTO saved_items(user_id, collection, product_id, created_at) "
|
| 190 |
+
"VALUES(?,?,?,?) ON CONFLICT DO NOTHING"),
|
| 191 |
+
(user_id, collection, product_id, _now()))
|
| 192 |
|
| 193 |
|
| 194 |
def remove_saved(user_id, collection, product_id):
|
| 195 |
with _conn() as c:
|
| 196 |
+
c.execute(_q("DELETE FROM saved_items WHERE user_id=? AND collection=? AND product_id=?"),
|
| 197 |
+
(user_id, collection, product_id))
|
|
|
|
| 198 |
|
| 199 |
|
| 200 |
def get_collections(user_id):
|
| 201 |
"""Return {collection_name: [product_id, ...]} ordered by when saved."""
|
| 202 |
with _conn() as c:
|
| 203 |
+
rows = c.execute(_q("SELECT collection, product_id FROM saved_items "
|
| 204 |
+
"WHERE user_id=? ORDER BY created_at"), (user_id,)).fetchall()
|
|
|
|
| 205 |
out = {}
|
| 206 |
for r in rows:
|
| 207 |
out.setdefault(r["collection"], []).append(r["product_id"])
|
| 208 |
return out
|
| 209 |
|
| 210 |
|
| 211 |
+
# ---------- product catalog (SQLite-only; in prod the catalog lives in catalog.json) ----------
|
| 212 |
_PRODUCT_COLS = ["id", "name", "brand", "category", "subcategory", "color",
|
| 213 |
"price", "prev_price", "on_sale", "rating", "reviews", "seller",
|
| 214 |
"in_stock", "emoji", "specs", "description", "text"]
|
|
|
|
| 229 |
|
| 230 |
|
| 231 |
def import_products(products):
|
| 232 |
+
if PG:
|
| 233 |
+
return # catalog is served from catalog.json, not stored in Postgres
|
| 234 |
now = _now()
|
| 235 |
with _conn() as c:
|
| 236 |
c.executemany(
|
|
|
|
| 242 |
p["price"], p["prev_price"], int(p["on_sale"]), p["rating"], p["reviews"],
|
| 243 |
p["seller"], int(p["in_stock"]), p["emoji"], json.dumps(p["specs"]),
|
| 244 |
p["description"], p["_text"]) for p in products])
|
|
|
|
| 245 |
c.executemany("INSERT INTO price_history(product_id, price, ts) VALUES(?,?,?)",
|
| 246 |
[(p["id"], p["price"], now) for p in products])
|
| 247 |
|
| 248 |
|
| 249 |
def ensure_catalog(json_path="catalog.json"):
|
| 250 |
+
if PG or catalog_count() > 0:
|
|
|
|
| 251 |
return
|
| 252 |
if not os.path.exists(json_path):
|
| 253 |
+
raise FileNotFoundError(f"{json_path} not found — run `python generate_catalog.py` first.")
|
|
|
|
| 254 |
with open(json_path) as f:
|
| 255 |
import_products(json.load(f))
|
| 256 |
|
| 257 |
|
| 258 |
def reseed_catalog_from_json(json_path="catalog.json"):
|
| 259 |
+
if PG:
|
| 260 |
+
return
|
| 261 |
with _conn() as c:
|
| 262 |
c.execute("DELETE FROM products")
|
| 263 |
c.execute("DELETE FROM price_history")
|
|
|
|
| 266 |
|
| 267 |
|
| 268 |
def load_catalog():
|
| 269 |
+
if PG:
|
| 270 |
+
return []
|
| 271 |
with _conn() as c:
|
| 272 |
rows = c.execute("SELECT * FROM products").fetchall()
|
| 273 |
return [_row_to_product(r) for r in rows]
|
| 274 |
|
| 275 |
|
| 276 |
def category_breakdown():
|
|
|
|
| 277 |
with _conn() as c:
|
| 278 |
rows = c.execute(
|
| 279 |
"SELECT category, subcategory, COUNT(*) AS n FROM products "
|
|
|
|
| 284 |
|
| 285 |
def update_price(product_id, new_price, prev_price):
|
| 286 |
with _conn() as c:
|
| 287 |
+
c.execute(_q("UPDATE products SET price=?, prev_price=?, on_sale=? WHERE id=?"),
|
| 288 |
(new_price, prev_price, int(new_price < prev_price), product_id))
|
| 289 |
+
c.execute(_q("INSERT INTO price_history(product_id, price, ts) VALUES(?,?,?)"),
|
| 290 |
(product_id, new_price, _now()))
|
| 291 |
|
| 292 |
|
| 293 |
def price_history(product_id, limit=30):
|
| 294 |
with _conn() as c:
|
| 295 |
+
rows = c.execute(_q("SELECT price, ts FROM price_history WHERE product_id=? "
|
| 296 |
+
"ORDER BY ts DESC LIMIT ?"), (product_id, limit)).fetchall()
|
|
|
|
| 297 |
return [{"price": r["price"], "ts": r["ts"]} for r in reversed(rows)]
|
| 298 |
|
| 299 |
|
| 300 |
# ---------- price tracking ----------
|
| 301 |
def add_watch(user_id, product_id, target_price=None):
|
| 302 |
with _conn() as c:
|
| 303 |
+
c.execute(_q("INSERT INTO price_watches(user_id, product_id, target_price, created_at) "
|
| 304 |
+
"VALUES(?,?,?,?) "
|
| 305 |
+
"ON CONFLICT(user_id, product_id) DO UPDATE SET target_price=excluded.target_price"),
|
| 306 |
+
(user_id, product_id, target_price, _now()))
|
|
|
|
| 307 |
|
| 308 |
|
| 309 |
def remove_watch(user_id, product_id):
|
| 310 |
with _conn() as c:
|
| 311 |
+
c.execute(_q("DELETE FROM price_watches WHERE user_id=? AND product_id=?"),
|
| 312 |
(user_id, product_id))
|
| 313 |
|
| 314 |
|
| 315 |
def get_watches(user_id):
|
| 316 |
with _conn() as c:
|
| 317 |
+
rows = c.execute(_q("SELECT product_id, target_price, created_at FROM price_watches "
|
| 318 |
+
"WHERE user_id=? ORDER BY created_at DESC"), (user_id,)).fetchall()
|
|
|
|
| 319 |
return [{"product_id": r["product_id"], "target_price": r["target_price"],
|
| 320 |
"created_at": r["created_at"]} for r in rows]
|
| 321 |
|
| 322 |
|
| 323 |
def watchers_for(product_id):
|
| 324 |
with _conn() as c:
|
| 325 |
+
rows = c.execute(_q("SELECT user_id, target_price FROM price_watches WHERE product_id=?"),
|
| 326 |
+
(product_id,)).fetchall()
|
|
|
|
| 327 |
return [(r["user_id"], r["target_price"]) for r in rows]
|
| 328 |
|
| 329 |
|
| 330 |
def watched_product_ids(user_id):
|
| 331 |
with _conn() as c:
|
| 332 |
+
rows = c.execute(_q("SELECT product_id FROM price_watches WHERE user_id=?"),
|
| 333 |
(user_id,)).fetchall()
|
| 334 |
return [r["product_id"] for r in rows]
|
| 335 |
|
|
|
|
| 337 |
# ---------- alerts ----------
|
| 338 |
def add_alert(user_id, product_id, old_price, new_price):
|
| 339 |
with _conn() as c:
|
| 340 |
+
c.execute(_q("INSERT INTO alerts(user_id, product_id, old_price, new_price, ts) "
|
| 341 |
+
"VALUES(?,?,?,?,?)"), (user_id, product_id, old_price, new_price, _now()))
|
|
|
|
| 342 |
|
| 343 |
|
| 344 |
def get_alerts(user_id, unseen_only=False):
|
|
|
|
| 347 |
q += " AND seen=0"
|
| 348 |
q += " ORDER BY ts DESC LIMIT 50"
|
| 349 |
with _conn() as c:
|
| 350 |
+
rows = c.execute(_q(q), (user_id,)).fetchall()
|
| 351 |
return [dict(r) for r in rows]
|
| 352 |
|
| 353 |
|
| 354 |
def mark_alerts_seen(user_id):
|
| 355 |
with _conn() as c:
|
| 356 |
+
c.execute(_q("UPDATE alerts SET seen=1 WHERE user_id=?"), (user_id,))
|
rmn-backend/requirements.txt
CHANGED
|
@@ -3,3 +3,5 @@ uvicorn[standard]
|
|
| 3 |
# semantic search (level 2) — small static embeddings, no torch
|
| 4 |
model2vec
|
| 5 |
numpy
|
|
|
|
|
|
|
|
|
| 3 |
# semantic search (level 2) — small static embeddings, no torch
|
| 4 |
model2vec
|
| 5 |
numpy
|
| 6 |
+
# persistent storage in production (used only when DATABASE_URL is set)
|
| 7 |
+
psycopg[binary]
|