,
]);
setRuntime(rt);
setIncidents((inc as unknown as { incidents?: PaymentIncident[] }).incidents ?? []);
setPayments((pay as unknown as { payments?: PaymentRecord[] }).payments ?? []);
+ setRefunds((refundPayload as unknown as { refunds?: RefundCase[] }).refunds ?? []);
setRisk(riskPayload);
} catch { /* */ }
setLoading(false);
@@ -328,6 +334,11 @@ export function PaymentsPageClient() {
{inc.detail || inc.reason || "—"}
+ {inc.refund_case_id ? (
+
+ 退款工单 #{inc.refund_case_id} · {inc.refund_status || "open"}
+
+ ) : null}
{compactMono(inc.user_id)}
@@ -355,6 +366,54 @@ export function PaymentsPageClient() {
+
+ 退款工单 ({refunds.length})
+
+ {refunds.length === 0 ? (
+ 暂无退款或售后工单
+ ) : (
+
+
+
+
+ | ID |
+ 状态 / 原因 |
+ 用户 / Intent |
+ Tx Hash |
+ 处理人 |
+ 更新时间 |
+
+
+
+ {refunds.map((refund) => (
+
+ | {refund.id} |
+
+ {refund.status || "open"}
+ {paymentReasonLabel(refund.reason)}
+ |
+
+ {compactMono(refund.user_id)}
+ {compactMono(refund.intent_id, 12, 6)}
+ |
+
+ {compactMono(refund.tx_hash)}
+ |
+
+ {refund.handled_by || refund.created_by || "—"}
+ |
+
+ {compactDate(refund.updated_at || refund.created_at)}
+ |
+
+ ))}
+
+
+
+ )}
+
+
+
成功支付记录 ({payments.length})
diff --git a/frontend/lib/ops-api.ts b/frontend/lib/ops-api.ts
index 474ab38d..d5d11c48 100644
--- a/frontend/lib/ops-api.ts
+++ b/frontend/lib/ops-api.ts
@@ -109,6 +109,37 @@ export const opsApi = {
if (reason) params.set("reason", reason);
return opsFetch>(`/api/ops/payments/incidents?${params}`);
},
+ auditLog(limit = 100, action?: string) {
+ const params = new URLSearchParams({ limit: String(limit) });
+ if (action) params.set("action", action);
+ return opsFetch>(`/api/ops/audit-log?${params}`);
+ },
+ refunds(limit = 50, status?: string) {
+ const params = new URLSearchParams({ limit: String(limit) });
+ if (status) params.set("status", status);
+ return opsFetch>(`/api/ops/refunds?${params}`);
+ },
+ createRefund(input: {
+ reason: string;
+ intent_id?: string;
+ tx_hash?: string;
+ user_id?: string;
+ amount_usdc?: string;
+ note?: string;
+ }) {
+ return opsFetch>("/api/ops/refunds", {
+ method: "POST",
+ headers: { "Content-Type": "application/json" },
+ body: JSON.stringify(input),
+ });
+ },
+ updateRefund(caseId: string | number, input: { status: string; note?: string }) {
+ return opsFetch>(`/api/ops/refunds/${caseId}`, {
+ method: "PATCH",
+ headers: { "Content-Type": "application/json" },
+ body: JSON.stringify(input),
+ });
+ },
feedback(limit = 100, status?: string) {
const params = new URLSearchParams({ limit: String(limit) });
if (status) params.set("status", status);
diff --git a/frontend/types/ops.ts b/frontend/types/ops.ts
index 9a67e8de..34989857 100644
--- a/frontend/types/ops.ts
+++ b/frontend/types/ops.ts
@@ -155,6 +155,8 @@ export type PaymentIncident = {
intent_id?: string;
user_id?: string;
tx_hash?: string;
+ refund_case_id?: number | string | null;
+ refund_status?: string;
payload_json?: string;
created_at?: string;
resolved?: boolean;
@@ -166,6 +168,42 @@ export type PaymentIncident = {
last_seen_at?: string;
};
+export type RefundCase = {
+ id: number;
+ status?: string;
+ reason?: string;
+ intent_id?: string;
+ tx_hash?: string;
+ user_id?: string;
+ amount_usdc?: string;
+ created_by?: string;
+ handled_by?: string;
+ notes?: Array<{ note?: string; by?: string; at?: string }>;
+ created_at?: string;
+ updated_at?: string;
+};
+
+export type RefundCasesPayload = {
+ refunds?: RefundCase[];
+};
+
+export type OpsAuditEvent = {
+ id: number;
+ action?: string;
+ actor_email?: string;
+ target_user_id?: string;
+ target_email?: string;
+ target_type?: string;
+ target_id?: string;
+ payload?: Record;
+ created_at?: string;
+};
+
+export type OpsAuditPayload = {
+ events?: OpsAuditEvent[];
+ total?: number;
+};
+
export type IncidentsPayload = {
incidents?: PaymentIncident[];
total?: number;
diff --git a/src/database/db_manager.py b/src/database/db_manager.py
index 8597ed30..768a38c0 100644
--- a/src/database/db_manager.py
+++ b/src/database/db_manager.py
@@ -596,6 +596,75 @@ class DBManager:
conn.execute(
"CREATE INDEX IF NOT EXISTS idx_payment_audit_events_created_at ON payment_audit_events(created_at DESC)"
)
+ conn.execute("""
+ CREATE TABLE IF NOT EXISTS ops_audit_events (
+ id INTEGER PRIMARY KEY AUTOINCREMENT,
+ action TEXT NOT NULL,
+ actor_email TEXT NOT NULL DEFAULT '',
+ target_user_id TEXT NOT NULL DEFAULT '',
+ target_email TEXT NOT NULL DEFAULT '',
+ target_type TEXT NOT NULL DEFAULT '',
+ target_id TEXT NOT NULL DEFAULT '',
+ payload_json TEXT NOT NULL,
+ created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
+ )
+ """)
+ conn.execute(
+ """
+ CREATE INDEX IF NOT EXISTS idx_ops_audit_events_created_at
+ ON ops_audit_events(created_at DESC)
+ """
+ )
+ conn.execute(
+ """
+ CREATE INDEX IF NOT EXISTS idx_ops_audit_events_action_created_at
+ ON ops_audit_events(action, created_at DESC)
+ """
+ )
+ conn.execute("""
+ CREATE TABLE IF NOT EXISTS points_ledger (
+ id INTEGER PRIMARY KEY AUTOINCREMENT,
+ telegram_id INTEGER,
+ supabase_user_id TEXT NOT NULL DEFAULT '',
+ supabase_email TEXT NOT NULL DEFAULT '',
+ source TEXT NOT NULL,
+ delta_points INTEGER NOT NULL,
+ balance_after INTEGER NOT NULL,
+ actor_email TEXT NOT NULL DEFAULT '',
+ reference_type TEXT NOT NULL DEFAULT '',
+ reference_id TEXT NOT NULL DEFAULT '',
+ metadata_json TEXT NOT NULL,
+ created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
+ )
+ """)
+ conn.execute(
+ """
+ CREATE INDEX IF NOT EXISTS idx_points_ledger_user_created_at
+ ON points_ledger(supabase_user_id, supabase_email, created_at DESC)
+ """
+ )
+ conn.execute("""
+ CREATE TABLE IF NOT EXISTS payment_refund_cases (
+ id INTEGER PRIMARY KEY AUTOINCREMENT,
+ status TEXT NOT NULL DEFAULT 'open',
+ reason TEXT NOT NULL,
+ intent_id TEXT NOT NULL DEFAULT '',
+ tx_hash TEXT NOT NULL DEFAULT '',
+ user_id TEXT NOT NULL DEFAULT '',
+ amount_usdc TEXT NOT NULL DEFAULT '',
+ created_by TEXT NOT NULL DEFAULT '',
+ handled_by TEXT NOT NULL DEFAULT '',
+ notes_json TEXT NOT NULL,
+ created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
+ updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
+ )
+ """)
+ conn.execute(
+ """
+ CREATE INDEX IF NOT EXISTS idx_payment_refund_cases_status_created_at
+ ON payment_refund_cases(status, created_at DESC)
+ """
+ )
conn.execute("""
CREATE TABLE IF NOT EXISTS observation_patch_events (
revision INTEGER PRIMARY KEY AUTOINCREMENT,
@@ -1673,6 +1742,479 @@ class DBManager:
)
conn.commit()
+ def append_ops_audit_event(
+ self,
+ *,
+ action: str,
+ actor_email: str = "",
+ target_user_id: str = "",
+ target_email: str = "",
+ target_type: str = "",
+ target_id: str = "",
+ payload: Optional[Dict[str, Any]] = None,
+ ) -> Dict[str, Any]:
+ normalized_action = str(action or "").strip().lower()
+ if not normalized_action:
+ return {"ok": False, "reason": "invalid_action"}
+ body = payload if isinstance(payload, dict) else {}
+ now = datetime.now().isoformat()
+ with self._get_connection() as conn:
+ conn.row_factory = sqlite3.Row
+ cursor = conn.execute(
+ """
+ INSERT INTO ops_audit_events (
+ action,
+ actor_email,
+ target_user_id,
+ target_email,
+ target_type,
+ target_id,
+ payload_json,
+ created_at
+ )
+ VALUES (?, ?, ?, ?, ?, ?, ?, ?)
+ """,
+ (
+ normalized_action,
+ str(actor_email or "").strip().lower(),
+ str(target_user_id or "").strip().lower(),
+ str(target_email or "").strip().lower(),
+ str(target_type or "").strip().lower(),
+ str(target_id or "").strip(),
+ json.dumps(body, ensure_ascii=False, default=str),
+ now,
+ ),
+ )
+ event_id = int(cursor.lastrowid)
+ conn.commit()
+ return {
+ "id": event_id,
+ "action": normalized_action,
+ "actor_email": str(actor_email or "").strip().lower(),
+ "target_user_id": str(target_user_id or "").strip().lower(),
+ "target_email": str(target_email or "").strip().lower(),
+ "target_type": str(target_type or "").strip().lower(),
+ "target_id": str(target_id or "").strip(),
+ "payload": body,
+ "created_at": now,
+ }
+
+ def list_ops_audit_events(
+ self,
+ *,
+ limit: int = 100,
+ action: str = "",
+ actor_email: str = "",
+ target_user_id: str = "",
+ ) -> List[Dict[str, Any]]:
+ safe_limit = max(1, min(int(limit or 100), 500))
+ clauses: List[str] = []
+ params: List[Any] = []
+ normalized_action = str(action or "").strip().lower()
+ normalized_actor = str(actor_email or "").strip().lower()
+ normalized_target_user = str(target_user_id or "").strip().lower()
+ if normalized_action:
+ clauses.append("action = ?")
+ params.append(normalized_action)
+ if normalized_actor:
+ clauses.append("actor_email = ?")
+ params.append(normalized_actor)
+ if normalized_target_user:
+ clauses.append("target_user_id = ?")
+ params.append(normalized_target_user)
+ where_sql = f"WHERE {' AND '.join(clauses)}" if clauses else ""
+ params.append(safe_limit)
+ with self._get_connection() as conn:
+ conn.row_factory = sqlite3.Row
+ rows = conn.execute(
+ f"""
+ SELECT id, action, actor_email, target_user_id, target_email,
+ target_type, target_id, payload_json, created_at
+ FROM ops_audit_events
+ {where_sql}
+ ORDER BY id DESC
+ LIMIT ?
+ """,
+ tuple(params),
+ ).fetchall()
+ events: List[Dict[str, Any]] = []
+ for row in rows:
+ try:
+ payload = json.loads(str(row["payload_json"] or "{}"))
+ except Exception:
+ payload = {}
+ events.append(
+ {
+ "id": int(row["id"]),
+ "action": str(row["action"] or ""),
+ "actor_email": str(row["actor_email"] or ""),
+ "target_user_id": str(row["target_user_id"] or ""),
+ "target_email": str(row["target_email"] or ""),
+ "target_type": str(row["target_type"] or ""),
+ "target_id": str(row["target_id"] or ""),
+ "payload": payload if isinstance(payload, dict) else {},
+ "created_at": row["created_at"],
+ }
+ )
+ return events
+
+ def _append_points_ledger_entry_conn(
+ self,
+ conn: sqlite3.Connection,
+ *,
+ telegram_id: Optional[int],
+ supabase_user_id: str = "",
+ supabase_email: str = "",
+ source: str,
+ delta_points: int,
+ balance_after: int,
+ actor_email: str = "",
+ reference_type: str = "",
+ reference_id: str = "",
+ metadata: Optional[Dict[str, Any]] = None,
+ ) -> None:
+ normalized_source = str(source or "").strip().lower()
+ if not normalized_source or int(delta_points or 0) == 0:
+ return
+ conn.execute(
+ """
+ INSERT INTO points_ledger (
+ telegram_id,
+ supabase_user_id,
+ supabase_email,
+ source,
+ delta_points,
+ balance_after,
+ actor_email,
+ reference_type,
+ reference_id,
+ metadata_json,
+ created_at
+ )
+ VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
+ """,
+ (
+ int(telegram_id) if telegram_id is not None else None,
+ str(supabase_user_id or "").strip().lower(),
+ str(supabase_email or "").strip().lower(),
+ normalized_source,
+ int(delta_points),
+ int(balance_after),
+ str(actor_email or "").strip().lower(),
+ str(reference_type or "").strip().lower(),
+ str(reference_id or "").strip(),
+ json.dumps(metadata if isinstance(metadata, dict) else {}, ensure_ascii=False, default=str),
+ datetime.now().isoformat(),
+ ),
+ )
+
+ def append_points_ledger_entry(
+ self,
+ *,
+ telegram_id: Optional[int] = None,
+ supabase_user_id: str = "",
+ supabase_email: str = "",
+ source: str,
+ delta_points: int,
+ balance_after: int,
+ actor_email: str = "",
+ reference_type: str = "",
+ reference_id: str = "",
+ metadata: Optional[Dict[str, Any]] = None,
+ ) -> None:
+ with self._get_connection() as conn:
+ self._append_points_ledger_entry_conn(
+ conn,
+ telegram_id=telegram_id,
+ supabase_user_id=supabase_user_id,
+ supabase_email=supabase_email,
+ source=source,
+ delta_points=delta_points,
+ balance_after=balance_after,
+ actor_email=actor_email,
+ reference_type=reference_type,
+ reference_id=reference_id,
+ metadata=metadata,
+ )
+ conn.commit()
+
+ def list_points_ledger_entries(
+ self,
+ *,
+ limit: int = 20,
+ supabase_user_id: str = "",
+ supabase_email: str = "",
+ ) -> List[Dict[str, Any]]:
+ safe_limit = max(1, min(int(limit or 20), 200))
+ normalized_user_id = str(supabase_user_id or "").strip().lower()
+ normalized_email = str(supabase_email or "").strip().lower()
+ if not normalized_user_id and not normalized_email:
+ return []
+ clauses: List[str] = []
+ params: List[Any] = []
+ if normalized_user_id:
+ clauses.append("supabase_user_id = ?")
+ params.append(normalized_user_id)
+ if normalized_email:
+ clauses.append("supabase_email = ?")
+ params.append(normalized_email)
+ where_sql = f"WHERE {' OR '.join(clauses)}" if clauses else ""
+ params.append(safe_limit)
+ with self._get_connection() as conn:
+ conn.row_factory = sqlite3.Row
+ rows = conn.execute(
+ f"""
+ SELECT id, telegram_id, supabase_user_id, supabase_email, source,
+ delta_points, balance_after, actor_email, reference_type,
+ reference_id, metadata_json, created_at
+ FROM points_ledger
+ {where_sql}
+ ORDER BY id DESC
+ LIMIT ?
+ """,
+ tuple(params),
+ ).fetchall()
+ out: List[Dict[str, Any]] = []
+ for row in rows:
+ try:
+ metadata = json.loads(str(row["metadata_json"] or "{}"))
+ except Exception:
+ metadata = {}
+ out.append(
+ {
+ "id": int(row["id"]),
+ "telegram_id": row["telegram_id"],
+ "supabase_user_id": str(row["supabase_user_id"] or ""),
+ "supabase_email": str(row["supabase_email"] or ""),
+ "source": str(row["source"] or ""),
+ "delta_points": int(row["delta_points"] or 0),
+ "balance_after": int(row["balance_after"] or 0),
+ "actor_email": str(row["actor_email"] or ""),
+ "reference_type": str(row["reference_type"] or ""),
+ "reference_id": str(row["reference_id"] or ""),
+ "metadata": metadata if isinstance(metadata, dict) else {},
+ "created_at": row["created_at"],
+ }
+ )
+ return out
+
+ def get_points_ledger_summary(
+ self,
+ *,
+ supabase_user_id: str = "",
+ supabase_email: str = "",
+ limit: int = 20,
+ ) -> Dict[str, Any]:
+ recent = self.list_points_ledger_entries(
+ limit=limit,
+ supabase_user_id=supabase_user_id,
+ supabase_email=supabase_email,
+ )
+ by_source: Dict[str, Dict[str, int]] = {}
+ for row in recent:
+ source = str(row.get("source") or "unknown")
+ bucket = by_source.setdefault(source, {"points": 0, "count": 0})
+ bucket["points"] += int(row.get("delta_points") or 0)
+ bucket["count"] += 1
+ balance = int(recent[0]["balance_after"]) if recent else (
+ self.get_points_by_supabase_user_id(supabase_user_id)
+ if supabase_user_id
+ else self.get_points_by_supabase_email(supabase_email)
+ )
+ return {
+ "balance": max(0, balance),
+ "recent": recent,
+ "by_source": by_source,
+ }
+
+ @staticmethod
+ def _refund_case_row_to_dict(row: sqlite3.Row) -> Dict[str, Any]:
+ try:
+ notes = json.loads(str(row["notes_json"] or "[]"))
+ except Exception:
+ notes = []
+ return {
+ "id": int(row["id"]),
+ "status": str(row["status"] or ""),
+ "reason": str(row["reason"] or ""),
+ "intent_id": str(row["intent_id"] or ""),
+ "tx_hash": str(row["tx_hash"] or ""),
+ "user_id": str(row["user_id"] or ""),
+ "amount_usdc": str(row["amount_usdc"] or ""),
+ "created_by": str(row["created_by"] or ""),
+ "handled_by": str(row["handled_by"] or ""),
+ "notes": notes if isinstance(notes, list) else [],
+ "created_at": row["created_at"],
+ "updated_at": row["updated_at"],
+ }
+
+ def create_refund_case(
+ self,
+ *,
+ reason: str,
+ intent_id: str = "",
+ tx_hash: str = "",
+ user_id: str = "",
+ amount_usdc: str = "",
+ created_by: str = "",
+ note: str = "",
+ ) -> Dict[str, Any]:
+ normalized_reason = str(reason or "").strip().lower()
+ if not normalized_reason:
+ return {"ok": False, "reason": "invalid_refund_reason"}
+ now = datetime.now().isoformat()
+ notes = []
+ note_text = str(note or "").strip()
+ if note_text:
+ notes.append(
+ {
+ "note": note_text,
+ "by": str(created_by or "").strip().lower(),
+ "at": now,
+ }
+ )
+ with self._get_connection() as conn:
+ conn.row_factory = sqlite3.Row
+ cursor = conn.execute(
+ """
+ INSERT INTO payment_refund_cases (
+ status,
+ reason,
+ intent_id,
+ tx_hash,
+ user_id,
+ amount_usdc,
+ created_by,
+ handled_by,
+ notes_json,
+ created_at,
+ updated_at
+ )
+ VALUES ('open', ?, ?, ?, ?, ?, ?, '', ?, ?, ?)
+ """,
+ (
+ normalized_reason,
+ str(intent_id or "").strip(),
+ str(tx_hash or "").strip().lower(),
+ str(user_id or "").strip().lower(),
+ str(amount_usdc or "").strip(),
+ str(created_by or "").strip().lower(),
+ json.dumps(notes, ensure_ascii=False, default=str),
+ now,
+ now,
+ ),
+ )
+ case_id = int(cursor.lastrowid)
+ row = conn.execute(
+ """
+ SELECT id, status, reason, intent_id, tx_hash, user_id,
+ amount_usdc, created_by, handled_by, notes_json,
+ created_at, updated_at
+ FROM payment_refund_cases
+ WHERE id = ?
+ """,
+ (case_id,),
+ ).fetchone()
+ conn.commit()
+ return self._refund_case_row_to_dict(row)
+
+ def update_refund_case(
+ self,
+ case_id: int,
+ *,
+ status: str,
+ handled_by: str = "",
+ note: str = "",
+ ) -> Optional[Dict[str, Any]]:
+ safe_id = int(case_id or 0)
+ normalized_status = str(status or "").strip().lower()
+ allowed = {"open", "processing", "refunded", "rejected", "closed"}
+ if safe_id <= 0 or normalized_status not in allowed:
+ return None
+ now = datetime.now().isoformat()
+ actor = str(handled_by or "").strip().lower()
+ with self._get_connection() as conn:
+ conn.row_factory = sqlite3.Row
+ row = conn.execute(
+ """
+ SELECT notes_json
+ FROM payment_refund_cases
+ WHERE id = ?
+ LIMIT 1
+ """,
+ (safe_id,),
+ ).fetchone()
+ if not row:
+ return None
+ try:
+ notes = json.loads(str(row["notes_json"] or "[]"))
+ except Exception:
+ notes = []
+ if not isinstance(notes, list):
+ notes = []
+ note_text = str(note or "").strip()
+ if note_text:
+ notes.append({"note": note_text, "by": actor, "at": now})
+ conn.execute(
+ """
+ UPDATE payment_refund_cases
+ SET status = ?,
+ handled_by = ?,
+ notes_json = ?,
+ updated_at = ?
+ WHERE id = ?
+ """,
+ (
+ normalized_status,
+ actor,
+ json.dumps(notes, ensure_ascii=False, default=str),
+ now,
+ safe_id,
+ ),
+ )
+ updated = conn.execute(
+ """
+ SELECT id, status, reason, intent_id, tx_hash, user_id,
+ amount_usdc, created_by, handled_by, notes_json,
+ created_at, updated_at
+ FROM payment_refund_cases
+ WHERE id = ?
+ """,
+ (safe_id,),
+ ).fetchone()
+ conn.commit()
+ return self._refund_case_row_to_dict(updated)
+
+ def list_refund_cases(
+ self,
+ *,
+ limit: int = 50,
+ status: Optional[str] = None,
+ ) -> List[Dict[str, Any]]:
+ safe_limit = max(1, min(int(limit or 50), 200))
+ normalized_status = str(status or "").strip().lower()
+ params: List[Any] = []
+ where_sql = ""
+ if normalized_status:
+ where_sql = "WHERE status = ?"
+ params.append(normalized_status)
+ params.append(safe_limit)
+ with self._get_connection() as conn:
+ conn.row_factory = sqlite3.Row
+ rows = conn.execute(
+ f"""
+ SELECT id, status, reason, intent_id, tx_hash, user_id,
+ amount_usdc, created_by, handled_by, notes_json,
+ created_at, updated_at
+ FROM payment_refund_cases
+ {where_sql}
+ ORDER BY id DESC
+ LIMIT ?
+ """,
+ tuple(params),
+ ).fetchall()
+ return [self._refund_case_row_to_dict(row) for row in rows]
+
def append_app_analytics_event(
self,
event_type: str,
@@ -1981,6 +2523,7 @@ class DBManager:
*,
points: int,
reason: str = "",
+ actor_email: str = "",
) -> Dict[str, Any]:
safe_points = int(points or 0)
if safe_points <= 0:
@@ -2017,7 +2560,7 @@ class DBManager:
user_row = conn.execute(
"""
- SELECT telegram_id, username, points, supabase_email
+ SELECT telegram_id, username, points, supabase_email, supabase_user_id
FROM users
WHERE lower(trim(COALESCE(supabase_email, ''))) = ?
LIMIT 1
@@ -2027,7 +2570,8 @@ class DBManager:
if not user_row:
user_row = conn.execute(
"""
- SELECT u.telegram_id, u.username, u.points, b.supabase_email
+ SELECT u.telegram_id, u.username, u.points, b.supabase_email,
+ b.supabase_user_id
FROM users u
JOIN supabase_bindings b ON b.telegram_id = u.telegram_id
WHERE lower(trim(COALESCE(b.supabase_email, ''))) = ?
@@ -2074,6 +2618,19 @@ class DBManager:
""",
(safe_points, normalized_reason, now, now, int(feedback_id)),
)
+ self._append_points_ledger_entry_conn(
+ conn,
+ telegram_id=telegram_id,
+ supabase_user_id=str(user_row["supabase_user_id"] or "").strip().lower(),
+ supabase_email=str(user_row["supabase_email"] or email),
+ source="feedback_reward",
+ delta_points=safe_points,
+ balance_after=after,
+ actor_email=actor_email,
+ reference_type="feedback",
+ reference_id=str(feedback_id),
+ metadata={"reason": normalized_reason},
+ )
updated_feedback_row = conn.execute(
"""
SELECT id, category, message, source, status, contact, user_id,
@@ -2764,6 +3321,12 @@ class DBManager:
self,
supabase_email: str,
amount: int,
+ *,
+ source: str = "manual_adjustment",
+ actor_email: str = "",
+ reference_type: str = "",
+ reference_id: str = "",
+ metadata: Optional[Dict[str, Any]] = None,
) -> Dict[str, Any]:
email = str(supabase_email or "").strip().lower()
points = int(amount or 0)
@@ -2776,7 +3339,7 @@ class DBManager:
conn.row_factory = sqlite3.Row
row = conn.execute(
"""
- SELECT telegram_id, username, points, supabase_email
+ SELECT telegram_id, username, points, supabase_email, supabase_user_id
FROM users
WHERE lower(trim(COALESCE(supabase_email, ''))) = ?
LIMIT 1
@@ -2797,6 +3360,19 @@ class DBManager:
""",
(after, telegram_id),
)
+ self._append_points_ledger_entry_conn(
+ conn,
+ telegram_id=telegram_id,
+ supabase_user_id=str(row["supabase_user_id"] or "").strip().lower(),
+ supabase_email=str(row["supabase_email"] or email),
+ source=source,
+ delta_points=points,
+ balance_after=after,
+ actor_email=actor_email,
+ reference_type=reference_type,
+ reference_id=reference_id,
+ metadata=metadata,
+ )
conn.commit()
self._sync_points_to_supabase_user_metadata(telegram_id, force=True)
return {
@@ -2813,6 +3389,12 @@ class DBManager:
self,
supabase_user_id: str,
amount: int,
+ *,
+ source: str = "manual_adjustment",
+ actor_email: str = "",
+ reference_type: str = "",
+ reference_id: str = "",
+ metadata: Optional[Dict[str, Any]] = None,
) -> Dict[str, Any]:
key = str(supabase_user_id or "").strip().lower()
points = int(amount or 0)
@@ -2849,6 +3431,19 @@ class DBManager:
""",
(after, telegram_id),
)
+ self._append_points_ledger_entry_conn(
+ conn,
+ telegram_id=telegram_id,
+ supabase_user_id=key,
+ supabase_email=str(row["supabase_email"] or ""),
+ source=source,
+ delta_points=points,
+ balance_after=after,
+ actor_email=actor_email,
+ reference_type=reference_type,
+ reference_id=reference_id,
+ metadata=metadata,
+ )
conn.commit()
self._sync_points_to_supabase_user_metadata(telegram_id, force=True)
return {
@@ -2866,6 +3461,12 @@ class DBManager:
self,
supabase_email: str,
amount: int,
+ *,
+ source: str = "points_redemption",
+ actor_email: str = "",
+ reference_type: str = "",
+ reference_id: str = "",
+ metadata: Optional[Dict[str, Any]] = None,
) -> Dict[str, Any]:
email = str(supabase_email or "").strip().lower()
points = int(amount or 0)
@@ -2878,7 +3479,7 @@ class DBManager:
conn.row_factory = sqlite3.Row
row = conn.execute(
"""
- SELECT telegram_id, username, points, supabase_email
+ SELECT telegram_id, username, points, supabase_email, supabase_user_id
FROM users
WHERE lower(trim(COALESCE(supabase_email, ''))) = ?
LIMIT 1
@@ -2902,6 +3503,19 @@ class DBManager:
"UPDATE users SET points = ? WHERE telegram_id = ?",
(after, telegram_id),
)
+ self._append_points_ledger_entry_conn(
+ conn,
+ telegram_id=telegram_id,
+ supabase_user_id=str(row["supabase_user_id"] or "").strip().lower(),
+ supabase_email=str(row["supabase_email"] or email),
+ source=source,
+ delta_points=-points,
+ balance_after=after,
+ actor_email=actor_email,
+ reference_type=reference_type,
+ reference_id=reference_id,
+ metadata=metadata,
+ )
conn.commit()
self._sync_points_to_supabase_user_metadata(telegram_id, force=True)
return {
diff --git a/src/payments/contract_checkout.py b/src/payments/contract_checkout.py
index 4f8d09c9..a484235d 100644
--- a/src/payments/contract_checkout.py
+++ b/src/payments/contract_checkout.py
@@ -2201,6 +2201,20 @@ class PaymentContractCheckoutService:
prefer="resolution=merge-duplicates,return=minimal",
allowed_status=[200, 201],
)
+ if str(status or "").strip().lower() == "refund_required":
+ self._db.append_payment_audit_event(
+ "payment_refund_required",
+ {
+ "reason": "refund_required",
+ "detail": detail,
+ "intent_id": intent.intent_id,
+ "user_id": getattr(intent, "user_id", "") or "",
+ "tx_hash": tx_hash_text,
+ "chain_id": int(intent.chain_id),
+ "from_address": _normalize_address(from_address) or "",
+ "receiver_expected": intent.receiver_address,
+ },
+ )
return {}
except Exception:
return {}
diff --git a/tests/test_ops_operational_closure.py b/tests/test_ops_operational_closure.py
new file mode 100644
index 00000000..029c8963
--- /dev/null
+++ b/tests/test_ops_operational_closure.py
@@ -0,0 +1,354 @@
+from src.database.db_manager import DBManager
+from web.schemas.auth import GrantPointsRequest
+import web.services.auth_api as auth_api
+import web.services.ops.payments as ops_payments
+import web.services.ops.users as ops_users
+import web.services.system_api as system_api
+
+
+def test_ops_audit_log_records_admin_action_roundtrip(tmp_path):
+ db = DBManager(str(tmp_path / "ops-audit.db"))
+
+ event = db.append_ops_audit_event(
+ action="manual_points_grant",
+ actor_email="ops@example.com",
+ target_user_id="user-1",
+ target_email="pilot@example.com",
+ payload={"points": 300},
+ )
+
+ rows = db.list_ops_audit_events(limit=10)
+
+ assert event["id"] > 0
+ assert rows[0]["action"] == "manual_points_grant"
+ assert rows[0]["actor_email"] == "ops@example.com"
+ assert rows[0]["target_user_id"] == "user-1"
+ assert rows[0]["target_email"] == "pilot@example.com"
+ assert rows[0]["payload"]["points"] == 300
+
+
+def test_points_ledger_explains_manual_and_feedback_sources(tmp_path):
+ db = DBManager(str(tmp_path / "points-ledger.db"))
+ db.upsert_user(1001, "pilot")
+ with db._get_connection() as conn: # noqa: SLF001
+ conn.execute(
+ """
+ UPDATE users
+ SET points = ?, supabase_user_id = ?, supabase_email = ?
+ WHERE telegram_id = ?
+ """,
+ (50, "user-1", "pilot@example.com", 1001),
+ )
+ conn.commit()
+
+ grant = db.grant_points_by_supabase_email(
+ "pilot@example.com",
+ 300,
+ source="ops_manual_grant",
+ actor_email="ops@example.com",
+ reference_type="ops_audit",
+ reference_id="audit-1",
+ )
+ feedback = db.append_user_feedback(
+ category="data",
+ message="METAR was stale.",
+ user_id="user-1",
+ user_email="pilot@example.com",
+ )
+ reward = db.grant_feedback_reward(
+ feedback["id"],
+ points=500,
+ reason="valid stale-data report",
+ actor_email="ops@example.com",
+ )
+
+ summary = db.get_points_ledger_summary(
+ supabase_user_id="user-1",
+ supabase_email="pilot@example.com",
+ limit=10,
+ )
+
+ assert grant["ok"] is True
+ assert reward["ok"] is True
+ assert summary["balance"] == 850
+ assert summary["by_source"]["ops_manual_grant"]["points"] == 300
+ assert summary["by_source"]["feedback_reward"]["points"] == 500
+ assert summary["recent"][0]["source"] == "feedback_reward"
+ assert summary["recent"][0]["metadata"]["reason"] == "valid stale-data report"
+
+
+def test_points_ledger_requires_user_identity(tmp_path):
+ db = DBManager(str(tmp_path / "points-ledger-identity.db"))
+ db.append_points_ledger_entry(
+ supabase_user_id="user-1",
+ supabase_email="pilot@example.com",
+ source="ops_manual_grant",
+ delta_points=100,
+ balance_after=100,
+ )
+
+ summary = db.get_points_ledger_summary(limit=10)
+
+ assert summary["balance"] == 0
+ assert summary["recent"] == []
+ assert summary["by_source"] == {}
+
+
+def test_refund_case_state_machine_roundtrip(tmp_path):
+ db = DBManager(str(tmp_path / "refund-cases.db"))
+
+ created = db.create_refund_case(
+ reason="duplicate_payment",
+ intent_id="intent-1",
+ tx_hash="0x" + "1" * 64,
+ user_id="user-1",
+ amount_usdc="29.9",
+ created_by="ops@example.com",
+ note="User submitted tx after order already paid.",
+ )
+ updated = db.update_refund_case(
+ created["id"],
+ status="processing",
+ handled_by="ops2@example.com",
+ note="Refund tx prepared.",
+ )
+ rows = db.list_refund_cases(limit=10)
+
+ assert created["status"] == "open"
+ assert updated["status"] == "processing"
+ assert updated["handled_by"] == "ops2@example.com"
+ assert updated["notes"][-1]["note"] == "Refund tx prepared."
+ assert rows[0]["id"] == created["id"]
+ assert rows[0]["reason"] == "duplicate_payment"
+
+
+def test_payment_incidents_include_refund_required_cases(monkeypatch):
+ class FakeDB:
+ def list_payment_audit_events(self, limit=50, event_type=None):
+ if event_type == "payment_intent_failed":
+ return [
+ {
+ "id": 1,
+ "event_type": "payment_intent_failed",
+ "payload": {
+ "reason": "receiver_mismatch",
+ "intent_id": "intent-1",
+ "user_id": "user-1",
+ "tx_hash": "0x" + "1" * 64,
+ },
+ "created_at": "2026-06-01T00:00:00",
+ }
+ ]
+ if event_type == "payment_refund_required":
+ return [
+ {
+ "id": 2,
+ "event_type": "payment_refund_required",
+ "payload": {
+ "reason": "refund_required",
+ "intent_id": "intent-2",
+ "user_id": "user-2",
+ "tx_hash": "0x" + "2" * 64,
+ },
+ "created_at": "2026-06-01T00:01:00",
+ }
+ ]
+ return []
+
+ def list_refund_cases(self, limit=50, status=None):
+ return [
+ {
+ "id": 10,
+ "status": "open",
+ "reason": "duplicate_payment",
+ "intent_id": "intent-3",
+ "user_id": "user-3",
+ "tx_hash": "0x" + "3" * 64,
+ "created_at": "2026-06-01T00:02:00",
+ }
+ ]
+
+ monkeypatch.setattr(ops_payments, "_require_ops", lambda request: {"email": "ops@example.com"})
+ monkeypatch.setattr(ops_payments, "_get_db", lambda: FakeDB())
+
+ payload = ops_payments.list_ops_payment_incidents(object(), limit=20)
+ reasons = {item["reason"] for item in payload["incidents"]}
+
+ assert {"receiver_mismatch", "refund_required", "duplicate_payment"}.issubset(
+ reasons
+ )
+ assert any(item.get("refund_case_id") == 10 for item in payload["incidents"])
+
+
+def test_ops_points_grant_writes_audit_and_ledger(monkeypatch):
+ calls = []
+
+ class FakeDB:
+ def grant_points_by_supabase_email(self, email, amount, **kwargs):
+ calls.append(("grant", email, amount, kwargs))
+ return {
+ "ok": True,
+ "supabase_user_id": "user-1",
+ "supabase_email": email,
+ "points_added": amount,
+ "points_after": 350,
+ }
+
+ def append_ops_audit_event(self, **kwargs):
+ calls.append(("audit", kwargs))
+ return {"id": 9, **kwargs}
+
+ monkeypatch.setattr(ops_users, "_require_ops", lambda request: {"email": "ops@example.com"})
+ monkeypatch.setattr(ops_users, "_get_db", lambda: FakeDB())
+
+ payload = ops_users.grant_ops_points(
+ object(),
+ GrantPointsRequest(email="pilot@example.com", points=300),
+ )
+
+ assert payload["ok"] is True
+ assert calls[0][0] == "grant"
+ assert calls[0][3]["source"] == "ops_manual_grant"
+ assert calls[0][3]["actor_email"] == "ops@example.com"
+ assert calls[1][0] == "audit"
+ assert calls[1][1]["action"] == "manual_points_grant"
+ assert payload["audit_event_id"] == 9
+
+
+def test_auth_me_payload_includes_points_ledger_summary(monkeypatch):
+ class FakeEntitlement:
+ enabled = True
+ require_subscription = False
+
+ @staticmethod
+ def ensure_signup_trial(user_id, email):
+ return None
+
+ @staticmethod
+ def get_subscription_window(*args, **kwargs):
+ return {"current": None, "rows": [], "total_expires_at": None}
+
+ @staticmethod
+ def get_latest_subscription_any_status(user_id):
+ return None
+
+ @staticmethod
+ def get_referral_summary(user_id):
+ return {"reward_points": 3500}
+
+ class FakeDB:
+ def get_points_ledger_summary(self, **kwargs):
+ assert kwargs["supabase_user_id"] == "user-1"
+ assert kwargs["supabase_email"] == "pilot@example.com"
+ return {
+ "balance": 850,
+ "by_source": {"feedback_reward": {"points": 500, "count": 1}},
+ "recent": [
+ {
+ "source": "feedback_reward",
+ "delta_points": 500,
+ "balance_after": 850,
+ }
+ ],
+ }
+
+ class RequestStub:
+ query_params = {}
+
+ def __init__(self):
+ self.state = type(
+ "State",
+ (),
+ {
+ "auth_user_id": "user-1",
+ "auth_email": "pilot@example.com",
+ },
+ )()
+
+ monkeypatch.setattr(auth_api.legacy_routes, "_assert_entitlement", lambda request: None)
+ monkeypatch.setattr(auth_api.legacy_routes, "_bind_optional_supabase_identity", lambda request: None)
+ monkeypatch.setattr(auth_api.legacy_routes, "_resolve_auth_points", lambda request: 850)
+ monkeypatch.setattr(
+ auth_api.legacy_routes,
+ "_resolve_weekly_profile",
+ lambda request: {"weekly_points": 0, "weekly_rank": None},
+ )
+ monkeypatch.setattr(auth_api.legacy_routes, "SUPABASE_ENTITLEMENT", FakeEntitlement())
+ monkeypatch.setattr(auth_api.legacy_routes, "_SUPABASE_AUTH_REQUIRED", False, raising=False)
+ monkeypatch.setattr(auth_api, "DBManager", lambda: FakeDB())
+ monkeypatch.setattr(auth_api.TelegramGroupPricing, "configured", False, raising=False)
+
+ payload = auth_api.get_auth_me_payload(RequestStub())
+
+ assert payload["points"] == 850
+ assert payload["points_ledger"]["balance"] == 850
+ assert payload["points_ledger"]["by_source"]["feedback_reward"]["points"] == 500
+
+
+def test_prometheus_exports_operational_closure_metrics(monkeypatch, tmp_path):
+ db_path = tmp_path / "metrics.db"
+ db_path.write_bytes(b"x" * 2048)
+ db_path_text = str(db_path)
+
+ class FakeDB:
+ db_path = db_path_text
+
+ def list_payment_audit_events(self, limit=500, event_type=None):
+ if event_type == "payment_intent_failed":
+ return [
+ {
+ "id": 1,
+ "event_type": "payment_intent_failed",
+ "payload": {"reason": "receiver_mismatch"},
+ "created_at": "2026-06-01T00:00:00",
+ },
+ {
+ "id": 2,
+ "event_type": "payment_intent_failed",
+ "payload": {
+ "reason": "event_mismatch",
+ "resolved_at": "2026-06-01T00:05:00",
+ },
+ "created_at": "2026-06-01T00:01:00",
+ },
+ ]
+ if event_type == "payment_refund_required":
+ return [
+ {
+ "id": 3,
+ "event_type": "payment_refund_required",
+ "payload": {"reason": "refund_required"},
+ "created_at": "2026-06-01T00:02:00",
+ }
+ ]
+ return []
+
+ def list_refund_cases(self, limit=500, status=None):
+ return [
+ {"id": 10, "status": "open", "reason": "duplicate_payment"},
+ {"id": 11, "status": "refunded", "reason": "receiver_mismatch"},
+ ]
+
+ monkeypatch.setattr(system_api.legacy_routes, "_require_ops_admin", lambda request: {"email": "ops@example.com"})
+ monkeypatch.setattr(system_api, "DBManager", lambda: FakeDB())
+ monkeypatch.setattr(
+ system_api,
+ "_realtime_status_payload",
+ lambda: {
+ "store": "degraded_sqlite",
+ "latest_revision": 42,
+ "sse_connections": 3,
+ "degraded_from": "redis",
+ },
+ )
+
+ response = system_api.get_prometheus_metrics_response(object())
+ text = response.body.decode("utf-8")
+
+ assert "polyweather_payment_incidents_open 2" in text
+ assert 'polyweather_payment_incidents_by_reason{reason="receiver_mismatch"} 1' in text
+ assert "polyweather_refund_cases_open 1" in text
+ assert "polyweather_sse_connections 3" in text
+ assert "polyweather_realtime_latest_revision 42" in text
+ assert "polyweather_realtime_redis_fallback 1" in text
+ assert "polyweather_sqlite_db_size_bytes 2048" in text
diff --git a/web/routers/ops.py b/web/routers/ops.py
index 69bc3370..97d5f4ca 100644
--- a/web/routers/ops.py
+++ b/web/routers/ops.py
@@ -4,6 +4,7 @@ from fastapi import APIRouter, Request, Response
from web.core import FeedbackRewardRequest, GrantPointsRequest
from web.services.ops_api import (
+ create_ops_refund_case,
extend_ops_subscription,
get_ops_analytics_funnel,
get_ops_billing_risk,
@@ -17,7 +18,9 @@ from web.services.ops_api import (
get_ops_source_health,
get_ops_truth_history,
get_ops_weekly_leaderboard,
+ list_ops_audit_log,
list_ops_feedback,
+ list_ops_refund_cases,
get_ops_user_subscriptions,
grant_ops_points,
transfer_ops_points,
@@ -25,6 +28,7 @@ from web.services.ops_api import (
list_ops_memberships,
list_ops_payment_incidents,
resolve_ops_payment_incident,
+ update_ops_refund_case,
list_ops_payments,
search_ops_users,
update_ops_config,
@@ -75,6 +79,23 @@ async def ops_weekly_leaderboard(request: Request, limit: int = 20):
return get_ops_weekly_leaderboard(request, limit=limit)
+@router.get("/api/ops/audit-log")
+async def ops_audit_log(
+ request: Request,
+ limit: int = 100,
+ action: str = "",
+ actor_email: str = "",
+ target_user_id: str = "",
+):
+ return list_ops_audit_log(
+ request,
+ limit=limit,
+ action=action,
+ actor_email=actor_email,
+ target_user_id=target_user_id,
+ )
+
+
@router.get("/api/ops/feedback")
async def ops_feedback(request: Request, limit: int = 100, status: str = ""):
return list_ops_feedback(request, limit=limit, status=status)
@@ -147,6 +168,40 @@ async def ops_payments(request: Request, limit: int = 50):
return list_ops_payments(request, limit=limit)
+@router.get("/api/ops/refunds")
+async def ops_refunds(request: Request, limit: int = 50, status: str = ""):
+ return list_ops_refund_cases(request, limit=limit, status=status)
+
+
+@router.post("/api/ops/refunds")
+async def ops_refund_create(request: Request):
+ import json as _json
+ body_bytes = await request.body()
+ body = _json.loads(body_bytes.decode("utf-8") or "{}")
+ return create_ops_refund_case(
+ request,
+ reason=str(body.get("reason") or "").strip(),
+ intent_id=str(body.get("intent_id") or "").strip(),
+ tx_hash=str(body.get("tx_hash") or "").strip(),
+ user_id=str(body.get("user_id") or "").strip(),
+ amount_usdc=str(body.get("amount_usdc") or "").strip(),
+ note=str(body.get("note") or "").strip(),
+ )
+
+
+@router.patch("/api/ops/refunds/{case_id}")
+async def ops_refund_update(request: Request, case_id: int):
+ import json as _json
+ body_bytes = await request.body()
+ body = _json.loads(body_bytes.decode("utf-8") or "{}")
+ return update_ops_refund_case(
+ request,
+ case_id=case_id,
+ status=str(body.get("status") or "").strip(),
+ note=str(body.get("note") or "").strip(),
+ )
+
+
@router.get("/api/ops/billing-risk")
async def ops_billing_risk(request: Request, days: int = 30, limit: int = 80):
return get_ops_billing_risk(request, days=days, limit=limit)
diff --git a/web/services/auth_api.py b/web/services/auth_api.py
index e0f29c5b..8e1099f9 100644
--- a/web/services/auth_api.py
+++ b/web/services/auth_api.py
@@ -381,6 +381,27 @@ def get_auth_me_payload(request: Request) -> Dict[str, Any]:
"weekly_profile",
lambda: legacy_routes._resolve_weekly_profile(request),
)
+ points_ledger = {
+ "balance": points,
+ "recent": [],
+ "by_source": {},
+ }
+ if user_id and not entitlement_scope:
+ try:
+ points_ledger = timer.measure(
+ "points_ledger",
+ lambda: DBManager().get_points_ledger_summary(
+ supabase_user_id=str(user_id or ""),
+ supabase_email=str(email or ""),
+ limit=8,
+ ),
+ )
+ except Exception:
+ points_ledger = {
+ "balance": points,
+ "recent": [],
+ "by_source": {},
+ }
def resolve_telegram_pricing() -> Any:
if not user_id:
@@ -409,6 +430,7 @@ def get_auth_me_payload(request: Request) -> Dict[str, Any]:
"user_id": user_id,
"email": email,
"points": points,
+ "points_ledger": points_ledger,
"weekly_points": weekly_profile["weekly_points"],
"weekly_rank": weekly_profile["weekly_rank"],
"entitlement_mode": (
diff --git a/web/services/ops/config.py b/web/services/ops/config.py
index 403070f0..1fbb7238 100644
--- a/web/services/ops/config.py
+++ b/web/services/ops/config.py
@@ -274,7 +274,8 @@ def grant_ops_subscription(
days: int = 30,
deduct_points: int = 0,
) -> dict[str, Any]:
- _require_ops(request)
+ admin = _require_ops(request) or {}
+ actor_email = str(admin.get("email") or "").strip().lower()
from datetime import datetime
import web.routes as legacy_routes # lazy – avoid circular import
@@ -344,11 +345,34 @@ def grant_ops_subscription(
if safe_deduct > 0:
db = _get_db()
deduct_result = db.deduct_points_by_supabase_email(
- normalized_email, safe_deduct
+ normalized_email,
+ safe_deduct,
+ source="ops_subscription_deduction",
+ actor_email=actor_email,
+ reference_type="subscription",
+ reference_id=user_id,
+ metadata={"plan_code": plan_code, "days": safe_days},
)
result["points_deducted"] = safe_deduct
result["points_result"] = deduct_result
+ try:
+ _get_db().append_ops_audit_event(
+ action="subscription_manual_grant",
+ actor_email=actor_email,
+ target_user_id=user_id,
+ target_email=normalized_email,
+ target_type="subscription",
+ payload={
+ "plan_code": plan_code,
+ "days": safe_days,
+ "expires_at": expires_at,
+ "deduct_points": safe_deduct,
+ },
+ )
+ except Exception:
+ pass
+
return result
@@ -357,7 +381,8 @@ def extend_ops_subscription(
email: str,
additional_days: int = 30,
) -> dict[str, Any]:
- _require_ops(request)
+ admin = _require_ops(request) or {}
+ actor_email = str(admin.get("email") or "").strip().lower()
from datetime import datetime
import web.routes as legacy_routes # lazy – avoid circular import
@@ -419,6 +444,22 @@ def extend_ops_subscription(
)
if patch_resp.ok:
legacy_routes.SUPABASE_ENTITLEMENT.invalidate_subscription_cache(user_id)
+ try:
+ _get_db().append_ops_audit_event(
+ action="subscription_manual_extend",
+ actor_email=actor_email,
+ target_user_id=user_id,
+ target_email=normalized_email,
+ target_type="subscription",
+ target_id=str(sub.get("id") or ""),
+ payload={
+ "additional_days": safe_days,
+ "previous_expires_at": current_expiry,
+ "new_expires_at": new_expiry,
+ },
+ )
+ except Exception:
+ pass
return {
"ok": True,
"email": normalized_email,
diff --git a/web/services/ops/payments.py b/web/services/ops/payments.py
index e0597b9f..db889889 100644
--- a/web/services/ops/payments.py
+++ b/web/services/ops/payments.py
@@ -146,6 +146,8 @@ def _normalize_payment_incident(item: Dict[str, Any]) -> Dict[str, Any]:
or confirm_failure.get("tx_hash")
or ""
).strip(),
+ "refund_case_id": payload.get("refund_case_id"),
+ "refund_status": str(payload.get("refund_status") or "").strip(),
"resolved": bool(resolved_at),
"resolved_at": resolved_at,
"resolved_by": str(payload.get("resolved_by") or "").strip(),
@@ -235,10 +237,47 @@ def list_ops_payment_incidents(
_require_ops(request)
db = _get_db()
safe_limit = max(1, min(int(limit or 50), 200))
- incidents = db.list_payment_audit_events(
- limit=max(safe_limit, 500),
- event_type="payment_intent_failed",
- )
+ incidents: List[Dict[str, Any]] = []
+ for event_type in ("payment_intent_failed", "payment_refund_required"):
+ try:
+ rows = db.list_payment_audit_events(
+ limit=max(safe_limit, 500),
+ event_type=event_type,
+ )
+ incidents.extend(
+ row for row in rows
+ if str(row.get("event_type") or "").strip().lower() == event_type
+ )
+ except Exception:
+ continue
+ terminal_refund_statuses = {"refunded", "rejected", "closed"}
+ list_refund_cases = getattr(db, "list_refund_cases", None)
+ if callable(list_refund_cases):
+ try:
+ refund_cases = list_refund_cases(limit=max(safe_limit, 500))
+ except Exception:
+ refund_cases = []
+ for case in refund_cases:
+ if not isinstance(case, dict):
+ continue
+ status = str(case.get("status") or "").strip().lower()
+ if not include_resolved and status in terminal_refund_statuses:
+ continue
+ incidents.append(
+ {
+ "id": int(case.get("id") or 0),
+ "event_type": "payment_refund_case",
+ "payload": {
+ "reason": str(case.get("reason") or "refund_required"),
+ "intent_id": case.get("intent_id"),
+ "user_id": case.get("user_id"),
+ "tx_hash": case.get("tx_hash"),
+ "refund_case_id": case.get("id"),
+ "refund_status": status,
+ },
+ "created_at": case.get("created_at"),
+ }
+ )
grouped = _group_payment_incidents(
incidents,
reason=reason,
@@ -250,6 +289,106 @@ def list_ops_payment_incidents(
}
+def list_ops_refund_cases(
+ request: Request,
+ limit: int = 50,
+ status: str = "",
+) -> Dict[str, Any]:
+ _require_ops(request)
+ db = _get_db()
+ safe_limit = max(1, min(int(limit or 50), 200))
+ return {
+ "refunds": db.list_refund_cases(
+ limit=safe_limit,
+ status=str(status or "").strip().lower() or None,
+ )
+ }
+
+
+def create_ops_refund_case(
+ request: Request,
+ *,
+ reason: str,
+ intent_id: str = "",
+ tx_hash: str = "",
+ user_id: str = "",
+ amount_usdc: str = "",
+ note: str = "",
+) -> Dict[str, Any]:
+ admin = _require_ops(request) or {}
+ actor_email = str(admin.get("email") or "").strip().lower()
+ db = _get_db()
+ created = db.create_refund_case(
+ reason=reason,
+ intent_id=intent_id,
+ tx_hash=tx_hash,
+ user_id=user_id,
+ amount_usdc=amount_usdc,
+ created_by=actor_email,
+ note=note,
+ )
+ if not created or created.get("ok") is False:
+ raise HTTPException(status_code=400, detail=created or "refund_case_failed")
+ db.append_ops_audit_event(
+ action="refund_case_create",
+ actor_email=actor_email,
+ target_user_id=str(user_id or ""),
+ target_type="refund_case",
+ target_id=str(created.get("id") or ""),
+ payload={
+ "reason": reason,
+ "intent_id": intent_id,
+ "tx_hash": tx_hash,
+ "amount_usdc": amount_usdc,
+ },
+ )
+ db.append_payment_audit_event(
+ "payment_refund_required",
+ {
+ "reason": str(reason or "refund_required").strip().lower(),
+ "intent_id": intent_id,
+ "user_id": user_id,
+ "tx_hash": tx_hash,
+ "refund_case_id": created.get("id"),
+ },
+ )
+ return {"ok": True, "refund": created}
+
+
+def update_ops_refund_case(
+ request: Request,
+ *,
+ case_id: int,
+ status: str,
+ note: str = "",
+) -> Dict[str, Any]:
+ admin = _require_ops(request) or {}
+ actor_email = str(admin.get("email") or "").strip().lower()
+ db = _get_db()
+ updated = db.update_refund_case(
+ case_id,
+ status=status,
+ handled_by=actor_email,
+ note=note,
+ )
+ if not updated:
+ raise HTTPException(status_code=404, detail="refund_case_not_found")
+ db.append_ops_audit_event(
+ action="refund_case_update",
+ actor_email=actor_email,
+ target_user_id=str(updated.get("user_id") or ""),
+ target_type="refund_case",
+ target_id=str(case_id),
+ payload={
+ "status": status,
+ "note": note,
+ "intent_id": updated.get("intent_id"),
+ "tx_hash": updated.get("tx_hash"),
+ },
+ )
+ return {"ok": True, "refund": updated}
+
+
def resolve_ops_payment_incident(request: Request, event_id: int) -> Dict[str, Any]:
admin = _require_ops(request) or {}
db = _get_db()
diff --git a/web/services/ops/users.py b/web/services/ops/users.py
index a75cbbde..b828d853 100644
--- a/web/services/ops/users.py
+++ b/web/services/ops/users.py
@@ -63,6 +63,25 @@ def search_ops_users(request: Request, q: str = "", limit: int = 20) -> Dict[str
return {"users": db.search_users(q, limit=limit)}
+def list_ops_audit_log(
+ request: Request,
+ *,
+ limit: int = 100,
+ action: str = "",
+ actor_email: str = "",
+ target_user_id: str = "",
+) -> Dict[str, Any]:
+ _require_ops(request)
+ db = _get_db()
+ rows = db.list_ops_audit_events(
+ limit=limit,
+ action=action,
+ actor_email=actor_email,
+ target_user_id=target_user_id,
+ )
+ return {"events": rows, "total": len(rows)}
+
+
def get_ops_weekly_leaderboard(request: Request, limit: int = 20) -> Dict[str, Any]:
_require_ops(request)
db = _get_db()
@@ -76,12 +95,34 @@ def get_ops_weekly_leaderboard(request: Request, limit: int = 20) -> Dict[str, A
def grant_ops_points(request: Request, body: GrantPointsRequest) -> Dict[str, Any]:
admin = _require_ops(request) or {}
db = _get_db()
- result = db.grant_points_by_supabase_email(body.email, body.points)
- result["operator_email"] = admin.get("email")
+ actor_email = str(admin.get("email") or "").strip().lower()
+ result = db.grant_points_by_supabase_email(
+ body.email,
+ body.points,
+ source="ops_manual_grant",
+ actor_email=actor_email,
+ reference_type="ops_action",
+ metadata={"action": "manual_points_grant"},
+ )
+ result["operator_email"] = actor_email
if not result.get("ok"):
reason = str(result.get("reason") or "grant_points_failed")
status_code = 404 if reason == "user_not_found" else 400
raise HTTPException(status_code=status_code, detail=result)
+ append_audit = getattr(db, "append_ops_audit_event", None)
+ if callable(append_audit):
+ audit = append_audit(
+ action="manual_points_grant",
+ actor_email=actor_email,
+ target_user_id=str(result.get("supabase_user_id") or ""),
+ target_email=str(result.get("supabase_email") or body.email),
+ target_type="user",
+ payload={
+ "points_added": int(result.get("points_added") or body.points),
+ "points_after": int(result.get("points_after") or 0),
+ },
+ )
+ result["audit_event_id"] = audit.get("id")
return result
@@ -170,12 +211,23 @@ def grant_ops_feedback_reward(
reason: str = "",
) -> Dict[str, Any]:
admin = _require_ops(request) or {}
+ actor_email = str(admin.get("email") or "").strip().lower()
db = _get_db()
- result = db.grant_feedback_reward(
- feedback_id,
- points=points,
- reason=reason,
- )
+ try:
+ result = db.grant_feedback_reward(
+ feedback_id,
+ points=points,
+ reason=reason,
+ actor_email=actor_email,
+ )
+ except TypeError as exc:
+ if "actor_email" not in str(exc):
+ raise
+ result = db.grant_feedback_reward(
+ feedback_id,
+ points=points,
+ reason=reason,
+ )
if not result.get("ok") and str(result.get("reason") or "") == "user_not_found":
feedback = result.get("feedback") if isinstance(result.get("feedback"), dict) else {}
reward_status = str(feedback.get("reward_status") or "").strip().lower()
@@ -203,11 +255,32 @@ def grant_ops_feedback_reward(
"supabase_user_id": supabase_user_id,
"feedback": updated_feedback,
}
- result["operator_email"] = admin.get("email")
+ result["operator_email"] = actor_email
if not result.get("ok"):
reason_code = str(result.get("reason") or "feedback_reward_failed")
status_code = 404 if reason_code in {"feedback_not_found", "user_not_found"} else 400
if reason_code == "already_rewarded":
status_code = 409
raise HTTPException(status_code=status_code, detail=result)
+ feedback = result.get("feedback") if isinstance(result.get("feedback"), dict) else {}
+ append_audit = getattr(db, "append_ops_audit_event", None)
+ if callable(append_audit):
+ audit = append_audit(
+ action="feedback_reward_grant",
+ actor_email=actor_email,
+ target_user_id=str(
+ result.get("supabase_user_id")
+ or feedback.get("user_id")
+ or ""
+ ),
+ target_email=str(result.get("supabase_email") or feedback.get("user_email") or ""),
+ target_type="feedback",
+ target_id=str(feedback_id),
+ payload={
+ "points_added": int(points or 0),
+ "reason": str(reason or ""),
+ "points_after": int(result.get("points_after") or 0),
+ },
+ )
+ result["audit_event_id"] = audit.get("id")
return result
diff --git a/web/services/ops_api.py b/web/services/ops_api.py
index 353cae2a..4aeb4605 100644
--- a/web/services/ops_api.py
+++ b/web/services/ops_api.py
@@ -35,6 +35,7 @@ from web.services.ops.users import ( # noqa: E402, F401
get_ops_weekly_leaderboard,
grant_ops_feedback_reward,
grant_ops_points,
+ list_ops_audit_log,
list_ops_feedback,
search_ops_users,
transfer_ops_points,
@@ -45,13 +46,16 @@ from web.services.ops.users import ( # noqa: E402, F401
# Payments / Billing / Memberships
# ---------------------------------------------------------------------------
from web.services.ops.payments import ( # noqa: E402, F401
+ create_ops_refund_case,
get_ops_billing_risk,
get_ops_memberships_growth,
get_ops_memberships_overview,
+ list_ops_refund_cases,
list_ops_memberships,
list_ops_payment_incidents,
list_ops_payments,
resolve_ops_payment_incident,
+ update_ops_refund_case,
)
# ---------------------------------------------------------------------------
diff --git a/web/services/system_api.py b/web/services/system_api.py
index 1afe9d76..eee60932 100644
--- a/web/services/system_api.py
+++ b/web/services/system_api.py
@@ -3,6 +3,8 @@
from __future__ import annotations
import time
+import os
+from collections import Counter
from typing import Any, Dict, Optional
from fastapi import BackgroundTasks, Request
@@ -10,7 +12,8 @@ from fastapi.concurrency import run_in_threadpool
from fastapi.responses import PlainTextResponse
from loguru import logger
-from src.utils.metrics import export_prometheus_metrics
+from src.database.db_manager import DBManager
+from src.utils.metrics import export_prometheus_metrics, gauge_set
from web.core import build_health_payload, build_system_status_payload
import web.routes as legacy_routes
@@ -128,7 +131,83 @@ def run_system_priority_warm(
def get_prometheus_metrics_response(request: Request) -> PlainTextResponse:
legacy_routes._require_ops_admin(request)
+ _refresh_operational_metrics()
return PlainTextResponse(
export_prometheus_metrics(),
media_type="text/plain; version=0.0.4; charset=utf-8",
)
+
+
+def _payment_event_reason(row: Dict[str, Any]) -> str:
+ payload = row.get("payload") if isinstance(row, dict) else {}
+ payload = payload if isinstance(payload, dict) else {}
+ confirm_failure = (
+ payload.get("confirm_failure")
+ if isinstance(payload.get("confirm_failure"), dict)
+ else {}
+ )
+ return str(
+ payload.get("reason")
+ or confirm_failure.get("reason")
+ or payload.get("error")
+ or "unknown"
+ ).strip().lower() or "unknown"
+
+
+def _payment_event_is_resolved(row: Dict[str, Any]) -> bool:
+ payload = row.get("payload") if isinstance(row, dict) else {}
+ payload = payload if isinstance(payload, dict) else {}
+ return bool(str(payload.get("resolved_at") or "").strip())
+
+
+def _refresh_operational_metrics() -> None:
+ try:
+ db = DBManager()
+ except Exception:
+ return
+
+ payment_events = []
+ for event_type in ("payment_intent_failed", "payment_refund_required"):
+ try:
+ payment_events.extend(
+ db.list_payment_audit_events(limit=500, event_type=event_type)
+ )
+ except Exception:
+ continue
+
+ open_events = [row for row in payment_events if not _payment_event_is_resolved(row)]
+ gauge_set("polyweather_payment_incidents_open", len(open_events))
+ for reason, count in Counter(_payment_event_reason(row) for row in open_events).items():
+ gauge_set("polyweather_payment_incidents_by_reason", count, reason=reason)
+
+ try:
+ refund_cases = db.list_refund_cases(limit=500)
+ except Exception:
+ refund_cases = []
+ terminal_statuses = {"refunded", "rejected", "closed"}
+ open_refunds = [
+ case
+ for case in refund_cases
+ if str(case.get("status") or "").strip().lower() not in terminal_statuses
+ ]
+ gauge_set("polyweather_refund_cases_open", len(open_refunds))
+
+ realtime = _realtime_status_payload()
+ gauge_set("polyweather_sse_connections", int(realtime.get("sse_connections") or 0))
+ gauge_set(
+ "polyweather_realtime_latest_revision",
+ int(realtime.get("latest_revision") or 0),
+ )
+ degraded_from = str(realtime.get("degraded_from") or "").strip().lower()
+ store = str(realtime.get("store") or "").strip().lower()
+ gauge_set(
+ "polyweather_realtime_redis_fallback",
+ 1 if degraded_from == "redis" or store == "degraded_sqlite" else 0,
+ )
+
+ db_path = str(getattr(db, "db_path", "") or "").strip()
+ if db_path:
+ try:
+ gauge_set("polyweather_sqlite_db_size_bytes", os.path.getsize(db_path))
+ except OSError:
+ pass
|