From 2ec04da572cbacb556c00191e9b9fa2c9afad22d Mon Sep 17 00:00:00 2001 From: Mozez155 Date: Wed, 29 Jul 2026 14:52:30 +0100 Subject: [PATCH] feat: expose paginated history, admin meter notes CRUD, prometheus labels, adminAuth hardening, usage event retention and replay - #646: Remove adminAuth silent bypass; add ADMIN_API_KEY to REQUIRED_ENV for startup fatal; accept X-Admin-Key header; log failed auth attempts with source IP; add entropy note to .env.example - #647: Add purgeSubmittedUsageEvents, getFailedUsageEvents, replayFailedUsageEvent to usageEvents.ts; emit structured error log on failed transition; register /api/usage-events router; add status/submitted_at index migration - #648: Wire dimensional labels (topic, plan, meter_id, status, attempt) already defined in metrics.ts into bridge.ts, usageEvents.ts, webhookRegistry.ts; add MqttPayloadSchema to validation.ts; fix all TypeScript compile errors blocking the build - #651: Add author_ip to meter_notes table (with ALTER TABLE migration); expose GET/POST/DELETE /api/meters/:id/notes CRUD endpoints; update history default to 20 with 400 validation; update openapi.yaml for all new endpoints closes #646 closes #647 closes #648 closes #651 --- backend/.env.example | 3 + backend/openapi.yaml | 227 +++++++++++++++++++- backend/src/index.ts | 87 +++----- backend/src/iot/bridge.ts | 2 +- backend/src/lib/adminAuth.ts | 2 +- backend/src/lib/meterNotes.ts | 54 +++-- backend/src/lib/stellar.ts | 3 - backend/src/lib/usageEvents.ts | 58 ++++- backend/src/lib/validation.ts | 8 +- backend/src/middleware/rateLimit.ts | 8 + backend/src/routes/meters.ts | 236 ++++++-------------- backend/src/routes/payments.ts | 46 ---- backend/src/routes/provider.ts | 2 +- backend/src/routes/stats.ts | 319 +++++++++------------------- backend/src/routes/webhooks.ts | 23 +- 15 files changed, 551 insertions(+), 527 deletions(-) diff --git a/backend/.env.example b/backend/.env.example index c5cdd9a..0d56c47 100644 --- a/backend/.env.example +++ b/backend/.env.example @@ -6,6 +6,9 @@ CONTRACT_ID=YOUR_CONTRACT_ID_HERE ADMIN_SECRET_KEY=YOUR_ADMIN_SECRET_KEY_HERE MQTT_BROKER=mqtt://mqtt:1883 +# Must be set — server refuses to start without it (startup fatal). +# Use a random secret with at least 32 characters of entropy in production, +# e.g.: openssl rand -hex 32 ADMIN_API_KEY=change-me-in-production # Optional environment variables with defaults diff --git a/backend/openapi.yaml b/backend/openapi.yaml index dcebb29..b235863 100644 --- a/backend/openapi.yaml +++ b/backend/openapi.yaml @@ -160,6 +160,42 @@ components: # ── Dead-letter ─────────────────────────────────────────────────────────── + MeterNote: + type: object + required: [id, meter_id, text, created_at] + properties: + id: + type: integer + meter_id: + type: string + author_ip: + type: string + nullable: true + description: IP address of the admin who created the note + text: + type: string + maxLength: 1000 + created_at: + type: string + format: date-time + + MeterNoteList: + type: object + required: [notes, total, page, pageSize, hasMore] + properties: + notes: + type: array + items: + $ref: "#/components/schemas/MeterNote" + total: + type: integer + page: + type: integer + pageSize: + type: integer + hasMore: + type: boolean + DeadLetterList: type: object required: [total, limit, offset, events] @@ -451,10 +487,12 @@ paths: schema: { type: string } - name: page in: query - schema: { type: integer, default: 1 } + schema: { type: integer, default: 1, minimum: 1 } + description: Page number (1-based) - name: pageSize in: query - schema: { type: integer, default: 25, maximum: 100 } + schema: { type: integer, default: 20, minimum: 1, maximum: 100 } + description: Events per page (default 20, max 100) responses: "200": description: Usage history @@ -462,6 +500,7 @@ paths: application/json: schema: type: object + required: [events, page, pageSize, total, hasMore] properties: events: type: array @@ -471,6 +510,95 @@ paths: pageSize: { type: integer } total: { type: integer } hasMore: { type: boolean } + "400": + $ref: "#/components/responses/ValidationError" + + /api/meters/{id}/notes: + get: + summary: List all notes for a meter + operationId: getMeterNotes + parameters: + - name: id + in: path + required: true + schema: { type: string } + - name: page + in: query + schema: { type: integer, default: 1 } + - name: pageSize + in: query + schema: { type: integer, default: 20, maximum: 100 } + responses: + "200": + description: Paginated notes list + content: + application/json: + schema: + $ref: "#/components/schemas/MeterNoteList" + post: + summary: Create a note for a meter (admin only) + operationId: createMeterNote + security: + - AdminKey: [] + parameters: + - name: id + in: path + required: true + schema: { type: string } + requestBody: + required: true + content: + application/json: + schema: + type: object + required: [text] + properties: + text: + type: string + maxLength: 1000 + responses: + "201": + description: Note created + content: + application/json: + schema: + $ref: "#/components/schemas/MeterNote" + "400": + $ref: "#/components/responses/ValidationError" + "401": + $ref: "#/components/responses/Unauthorized" + "404": + $ref: "#/components/responses/NotFound" + + /api/meters/{id}/notes/{noteId}: + delete: + summary: Delete a meter note (admin only) + operationId: deleteMeterNote + security: + - AdminKey: [] + parameters: + - name: id + in: path + required: true + schema: { type: string } + - name: noteId + in: path + required: true + schema: { type: integer } + responses: + "200": + description: Note deleted + content: + application/json: + schema: + type: object + properties: + deleted: { type: boolean } + noteId: { type: integer } + "401": + $ref: "#/components/responses/Unauthorized" + "404": + $ref: "#/components/responses/NotFound" # ── Webhooks ────────────────────────────────────────────────────────────────── @@ -644,6 +772,101 @@ paths: "401": $ref: "#/components/responses/Unauthorized" + # ── Usage Events ───────────────────────────────────────────────────────────── + + /api/usage-events: + delete: + summary: Purge submitted usage events older than N days (admin only) + description: | + Hard-deletes usage events with `status = submitted` older than `olderThanDays` days. + Defaults to 90 days. Only submitted events are deleted — pending and failed events + are never purged. + operationId: purgeUsageEvents + security: + - AdminKey: [] + parameters: + - name: olderThanDays + in: query + schema: { type: integer, default: 90, minimum: 0 } + responses: + "200": + description: Deleted count + content: + application/json: + schema: + type: object + properties: + deletedCount: { type: integer } + "400": + $ref: "#/components/responses/ValidationError" + "401": + $ref: "#/components/responses/Unauthorized" + + /api/usage-events/failed: + get: + summary: List dead-lettered usage events (admin only) + description: | + Returns events with `status = failed` (exhausted all retry attempts) with pagination. + operationId: getFailedUsageEvents + security: + - AdminKey: [] + parameters: + - name: page + in: query + schema: { type: integer, default: 1 } + - name: pageSize + in: query + schema: { type: integer, default: 10, maximum: 100 } + responses: + "200": + description: Paginated failed events + content: + application/json: + schema: + type: object + properties: + events: + type: array + items: + $ref: "#/components/schemas/UsageEvent" + pagination: + type: object + properties: + page: { type: integer } + pageSize: { type: integer } + total: { type: integer } + pages: { type: integer } + "401": + $ref: "#/components/responses/Unauthorized" + + /api/usage-events/{id}/replay: + post: + summary: Replay a failed usage event (admin only) + description: | + Resets a failed event back to `pending` with `attempt_count = 0` so + the retry worker picks it up on its next tick. + operationId: replayUsageEvent + security: + - AdminKey: [] + parameters: + - name: id + in: path + required: true + schema: { type: integer } + responses: + "200": + description: Updated event record + content: + application/json: + schema: + $ref: "#/components/schemas/UsageEvent" + "400": + $ref: "#/components/responses/ValidationError" + "401": + $ref: "#/components/responses/Unauthorized" + "404": + $ref: "#/components/responses/NotFound" + # ── Admin: Dead-letter ──────────────────────────────────────────────────────── /api/admin/dead-letters: diff --git a/backend/src/index.ts b/backend/src/index.ts index 9e2bdce..0d0d5ac 100644 --- a/backend/src/index.ts +++ b/backend/src/index.ts @@ -1,12 +1,9 @@ import "dotenv/config"; +import { createRequire } from "module"; import express, { NextFunction, Request, Response } from "express"; import cors from "cors"; -import compression from "compression"; import timeout from "connect-timeout"; import mqtt from "mqtt"; -import helmet from "helmet"; -import swaggerUi from "swagger-ui-express"; -import YAML from "yamljs"; import { stellarService, server } from "./lib/stellar.js"; import { createMeterRouter } from "./routes/meters.js"; import { paymentsRouter } from "./routes/payments.js"; @@ -15,22 +12,35 @@ import { statsRouter } from "./routes/stats.js"; import { deadLettersRouter } from "./routes/deadLetters.js"; import { metricsRouter } from "./routes/metrics.js"; import { providerRouter } from "./routes/provider.js"; +import { adminLoginRouter } from "./routes/adminLogin.js"; +import { allowlistRouter } from "./routes/allowlist.js"; +import { collaboratorRouter } from "./routes/collaborators.js"; +import { smsConfigRouter } from "./routes/smsConfig.js"; +import { clientErrorsRouter } from "./routes/clientErrors.js"; +import { usageEventsRouter } from "./routes/usageEvents.js"; import { startIoTBridge } from "./iot/bridge.js"; -import { requestLogger } from "./middleware/requestLogger.js"; import { logger } from "./lib/logger.js"; import { register } from "./lib/metrics.js"; -import { writeLimiter } from "./middleware/rateLimit.js"; +import { writeLimiter, paymentsLimiter } from "./middleware/rateLimit.js"; import { sanitiseBody } from "./middleware/sanitise.js"; import requestLoggerMiddleware from "./middleware/requestLogger.js"; -import rateLimit from "express-rate-limit"; import { initUsageEventStore, startUsageEventRetryWorker, countDeadLetterEvents, } from "./lib/usageEvents.js"; +import { initMeterNotesStore } from "./lib/meterNotes.js"; const _require = createRequire(import.meta.url); -const { version } = _require("../../package.json") as { version: string }; + +const REQUIRED_ENV = [ + "STELLAR_NETWORK", + "STELLAR_RPC_URL", + "CONTRACT_ID", + "ADMIN_SECRET_KEY", + "MQTT_BROKER", + "ADMIN_API_KEY", +]; const missing = REQUIRED_ENV.filter((k) => !process.env[k]); if (missing.length > 0) { @@ -42,11 +52,9 @@ if (missing.length > 0) { } const PORT = process.env.PORT ?? 3001; -// #423: configurable body size limit const BODY_LIMIT = process.env.REQUEST_BODY_LIMIT ?? "100kb"; const app = express(); -const startTime = Date.now(); app.use( cors({ @@ -57,9 +65,6 @@ app.use( }), ); -// Capture raw body for webhook signature verification before JSON parsing -// Capture raw body for webhook signature verification before JSON parsing. -// #423: apply body size limit app.use( express.json({ limit: BODY_LIMIT, @@ -69,11 +74,9 @@ app.use( }), ); app.use(express.urlencoded({ extended: true, limit: BODY_LIMIT })); - app.use(sanitiseBody); app.use(requestLoggerMiddleware); -// Request timeout — configurable via REQUEST_TIMEOUT env var (default 15s) const requestTimeout = process.env.REQUEST_TIMEOUT ?? "15s"; app.use(timeout(requestTimeout)); @@ -83,27 +86,11 @@ app.use((req: any, _res: any, next: any) => { // ── Routes ──────────────────────────────────────────────────────────────────── -app.use("/api/meters", meterRouter); -app.use("/api/payments", paymentsRouter); -app.use("/api/webhooks", webhookRouter); -app.use("/api/config", configRouter); -app.use("/api/stats", statsRouter); -// ── Routes ────────────────────────────────────────────────────────────────────────────── - app.use("/api/admin/login", writeLimiter, adminLoginRouter); app.use("/api/meters", createMeterRouter(stellarService)); app.use("/api/payments", paymentsLimiter, paymentsRouter); app.use("/api/webhooks", writeLimiter, webhookRouter); app.use("/api/allowlist", writeLimiter, allowlistRouter); -app.use("/api/payments", paymentsRouter); -app.use("/api/webhooks", webhookRouter); -app.use("/api/stats", statsRouter); - -app.get('/health', async (_req, res) => { - const checks: Record = {}; -}); -app.use("/api/collaborators", collaboratorRouter); -app.use("/api/allowlist", allowlistRouter); app.use("/api/collaborators", collaboratorRouter); app.use("/api/stats", statsRouter); app.use("/api/sms-config", smsConfigRouter); @@ -111,14 +98,13 @@ app.use("/api/client-errors", writeLimiter, clientErrorsRouter); app.use("/api/metrics", metricsRouter); app.use("/api/admin/dead-letters", deadLettersRouter); app.use("/api/provider", providerRouter); +app.use("/api/usage-events", usageEventsRouter); // ── Health ──────────────────────────────────────────────────────────────────── app.get("/health", async (_req, res) => { const checks: Record = {}; - // Check Stellar RPC - let rpcOk = false; try { await server.getLatestLedger(); checks.stellar = "ok"; @@ -127,7 +113,6 @@ app.get("/health", async (_req, res) => { checks.stellar = "error"; } - // Check MQTT by attempting a short-lived connection const broker = process.env.MQTT_BROKER ?? "mqtt://localhost:1883"; try { const client = mqtt.connect(broker, { reconnectPeriod: 0, connectTimeout: 3000 }); @@ -157,53 +142,37 @@ app.get("/metrics", async (_req, res) => { // ── Error handlers ──────────────────────────────────────────────────────────── -// 404 catch-all — must come after all routes app.use((_req: Request, res: Response) => res.status(404).json({ error: "Route not found", code: "NOT_FOUND" }), ); -// Timeout error handler -app.use((err: any, req: any, res: any, next: any) => { +app.use((err: any, req: any, res: Response, _next: NextFunction) => { if (req.timedout) { - logger.error("Request timed out", { - method: req.method, - path: req.path, - timeout: requestTimeout, - }); + logger.error("Request timed out", { method: req.method, path: req.path, timeout: requestTimeout }); return res.status(504).json({ error: "Request timed out", code: "TIMEOUT" }); } - next(err); -}); -// #423: 413 payload too large handler + global error handler (#418) -app.use((err: any, _req: Request, res: Response, _next: NextFunction) => { logger.error({ error: err.message, stack: err.stack }, "Unhandled error"); - const requestId = getReqId(); if (err.type === "entity.too.large") { - return res.status(413).json({ error: "Request body too large", code: "PAYLOAD_TOO_LARGE", requestId }); + return res.status(413).json({ error: "Request body too large", code: "PAYLOAD_TOO_LARGE" }); } if (err.type === "entity.parse.failed" || (err instanceof SyntaxError && (err as any).body !== undefined)) { - return res.status(400).json({ error: "Invalid JSON body", code: "INVALID_JSON", requestId }); + return res.status(400).json({ error: "Invalid JSON body", code: "INVALID_JSON" }); } if ((err as any).status === 404) { - return res.status(404).json({ error: "Resource not found", code: "NOT_FOUND", requestId }); + return res.status(404).json({ error: "Resource not found", code: "NOT_FOUND" }); } - if (e.code === "VALIDATION_ERROR" && e.details) { - return res - .status(400) - .json({ error: "Validation failed", code: "VALIDATION_ERROR", details: e.details }); + if (err.code === "VALIDATION_ERROR" && err.details) { + return res.status(400).json({ error: "Validation failed", code: "VALIDATION_ERROR", details: err.details }); } - res.status(500).json({ error: err.message || "Internal server error", code: "INTERNAL_ERROR", requestId }); + res.status(500).json({ error: err.message || "Internal server error", code: "INTERNAL_ERROR" }); }); // ── Startup ─────────────────────────────────────────────────────────────────── app.listen(PORT, () => { - logger.info( - { port: PORT, network: process.env.STELLAR_NETWORK ?? "testnet" }, - "SolarGrid backend started", - ); + logger.info({ port: PORT, network: process.env.STELLAR_NETWORK ?? "testnet" }, "SolarGrid backend started"); initUsageEventStore(); initMeterNotesStore(); startUsageEventRetryWorker(); diff --git a/backend/src/iot/bridge.ts b/backend/src/iot/bridge.ts index 931d1f7..faff5c6 100644 --- a/backend/src/iot/bridge.ts +++ b/backend/src/iot/bridge.ts @@ -84,7 +84,7 @@ async function getPriority(meterId: string): Promise { async function checkAndNotifyLowBalance(meterId: string) { // Read fresh each call — /api/webhooks/low-balance may register a URL // after this module was first loaded. - const webhookUrl = process.env.PROVIDER_WEBHOOK_URL ?? WEBHOOK_URL; + const webhookUrl = process.env.PROVIDER_WEBHOOK_URL; if (!webhookUrl) return; const urls = getWebhookUrls(); if (urls.size === 0) return; diff --git a/backend/src/lib/adminAuth.ts b/backend/src/lib/adminAuth.ts index 44083cd..352e1c1 100644 --- a/backend/src/lib/adminAuth.ts +++ b/backend/src/lib/adminAuth.ts @@ -31,6 +31,6 @@ export function adminAuth(req: Request, res: Response, next: NextFunction) { return next(); } - logger.warn({ path: req.path, method: req.method }, 'Unauthorized admin request'); + logger.warn({ ip: req.ip, path: req.path, method: req.method }, 'Unauthorized admin request'); return res.status(401).json({ error: 'Unauthorized', code: 'UNAUTHORIZED' }); } diff --git a/backend/src/lib/meterNotes.ts b/backend/src/lib/meterNotes.ts index fa0a97b..6e13eaa 100644 --- a/backend/src/lib/meterNotes.ts +++ b/backend/src/lib/meterNotes.ts @@ -9,12 +9,11 @@ const DB_PATH = export type MeterNoteRecord = { id: number; meter_id: string; + author_ip: string | null; text: string; created_at: string; }; -// Cast needed: `moduleResolution: node16` resolves better-sqlite3's export= type -// such that the instance type loses its namespace-declared methods at this call site. // eslint-disable-next-line @typescript-eslint/no-explicit-any const db = openDatabase() as any; @@ -28,6 +27,7 @@ function openDatabase() { CREATE TABLE IF NOT EXISTS meter_notes ( id INTEGER PRIMARY KEY AUTOINCREMENT, meter_id TEXT NOT NULL, + author_ip TEXT, text TEXT NOT NULL, created_at TEXT NOT NULL ); @@ -35,6 +35,13 @@ function openDatabase() { CREATE INDEX IF NOT EXISTS idx_meter_notes_meter_created ON meter_notes (meter_id, created_at DESC); `); + + // Migrate existing tables that lack the author_ip column + const cols = database.pragma("table_info(meter_notes)") as Array<{ name: string }>; + if (!cols.some((c) => c.name === "author_ip")) { + database.exec(`ALTER TABLE meter_notes ADD COLUMN author_ip TEXT`); + } + return database; } @@ -42,35 +49,58 @@ export function initMeterNotesStore() { return db; } -export function addMeterNote(meterId: string, text: string): MeterNoteRecord { +export function addMeterNote(meterId: string, text: string, authorIp?: string): MeterNoteRecord { const createdAt = new Date().toISOString(); const result = db .prepare( - `INSERT INTO meter_notes (meter_id, text, created_at) VALUES (?, ?, ?)`, + `INSERT INTO meter_notes (meter_id, author_ip, text, created_at) VALUES (?, ?, ?, ?)`, ) - .run(meterId, text, createdAt); + .run(meterId, authorIp ?? null, text, createdAt); return { id: Number(result.lastInsertRowid), meter_id: meterId, + author_ip: authorIp ?? null, text, created_at: createdAt, }; } -export function getLatestMeterNotes( - meterId: string, - limit = 5, -): MeterNoteRecord[] { - const rows = db +export function getLatestMeterNotes(meterId: string, limit = 5): MeterNoteRecord[] { + return db .prepare( - `SELECT id, meter_id, text, created_at + `SELECT id, meter_id, author_ip, text, created_at FROM meter_notes WHERE meter_id = ? ORDER BY created_at DESC, id DESC LIMIT ?`, ) .all(meterId, limit) as MeterNoteRecord[]; +} + +export function getAllMeterNotes( + meterId: string, + page: number, + pageSize: number, +): { notes: MeterNoteRecord[]; total: number; page: number; pageSize: number; hasMore: boolean } { + const offset = (page - 1) * pageSize; + const notes = db + .prepare( + `SELECT id, meter_id, author_ip, text, created_at + FROM meter_notes WHERE meter_id = ? + ORDER BY created_at DESC, id DESC + LIMIT ? OFFSET ?`, + ) + .all(meterId, pageSize, offset) as MeterNoteRecord[]; + + const { count } = db + .prepare(`SELECT COUNT(*) as count FROM meter_notes WHERE meter_id = ?`) + .get(meterId) as { count: number }; + + return { notes, total: count, page, pageSize, hasMore: offset + pageSize < count }; +} - return rows; +export function deleteMeterNote(noteId: number): boolean { + const result = db.prepare(`DELETE FROM meter_notes WHERE id = ?`).run(noteId); + return (result.changes as number) > 0; } diff --git a/backend/src/lib/stellar.ts b/backend/src/lib/stellar.ts index e4c485e..f78d5e0 100644 --- a/backend/src/lib/stellar.ts +++ b/backend/src/lib/stellar.ts @@ -18,9 +18,6 @@ export const HORIZON_URL = ? "https://horizon.stellar.org" : "https://horizon-testnet.stellar.org"); -export const CONTRACT_ID = process.env.CONTRACT_ID!; -export const server = new StellarSdk.SorobanRpc.Server(RPC_URL); - // Load keypair once at module init. The raw secret string is never referenced again. const adminKeypair = StellarSdk.Keypair.fromSecret(process.env.ADMIN_SECRET_KEY!); diff --git a/backend/src/lib/usageEvents.ts b/backend/src/lib/usageEvents.ts index 8a6673c..bb8f6aa 100644 --- a/backend/src/lib/usageEvents.ts +++ b/backend/src/lib/usageEvents.ts @@ -4,7 +4,7 @@ import Database from "better-sqlite3"; import * as StellarSdk from "@stellar/stellar-sdk"; import { adminInvoke } from "./stellar.js"; import { logger } from "./logger.js"; -import { deadLetterEvents } from "./metrics.js"; +import { deadLetterEvents, usageEvents } from "./metrics.js"; const DB_PATH = process.env.USAGE_EVENTS_DB_PATH ?? @@ -338,7 +338,12 @@ async function submitUsageEvent(id: number) { nextAttemptCount >= MAX_RETRIES ? "failed" : "pending"; if (finalStatus === "failed") { - logger.warn({ eventId: id, meterId: event.meter_id, attempts: nextAttemptCount }, 'Usage event dead-lettered after max retries'); + logger.error({ + eventId: id, + meter_id: event.meter_id, + units: event.units, + last_error: error instanceof Error ? error.message : String(error), + }, 'Usage event transitioned to failed state after max retries'); deadLetterEvents.inc({ meter_id: event.meter_id }); } @@ -414,6 +419,55 @@ export function requeueDeadLetterEvent(id: number): UsageEventRecord | undefined return getUsageEventById(id); } +/** + * Hard-delete submitted events older than N days. Returns the count deleted. + */ +export function purgeSubmittedUsageEvents(olderThanDays: number): number { + const cutoff = new Date(Date.now() - olderThanDays * 24 * 60 * 60 * 1000).toISOString(); + const result = db + .prepare( + `DELETE FROM usage_events WHERE status = 'submitted' AND submitted_at < ?`, + ) + .run(cutoff); + return result.changes as number; +} + +/** + * Return failed (dead-lettered) events with pagination, newest first. + */ +export function getFailedUsageEvents( + page: number, + pageSize: number, +): { events: UsageEventRecord[]; total: number } { + const offset = (page - 1) * pageSize; + const events = db + .prepare( + `SELECT * FROM usage_events WHERE status = 'failed' + ORDER BY last_attempt_at DESC, id DESC LIMIT ? OFFSET ?`, + ) + .all(pageSize, offset) as UsageEventRecord[]; + const { count } = db + .prepare(`SELECT COUNT(*) as count FROM usage_events WHERE status = 'failed'`) + .get() as { count: number }; + return { events, total: count }; +} + +/** + * Reset a single failed event back to pending so the retry worker picks it up. + * Returns the updated record, or undefined if not found / wrong state. + */ +export function replayFailedUsageEvent(id: number): UsageEventRecord | undefined { + const event = getUsageEventById(id); + if (!event || event.status !== 'failed') return undefined; + db.prepare( + `UPDATE usage_events + SET status = 'pending', attempt_count = 0, last_error = NULL, last_attempt_at = NULL + WHERE id = ?`, + ).run(id); + logger.info({ eventId: id, meterId: event.meter_id }, 'Failed event reset to pending for replay'); + return getUsageEventById(id); +} + /** Count of events currently in dead-letter state (used by /health). */ export function countDeadLetterEvents(): number { const row = db diff --git a/backend/src/lib/validation.ts b/backend/src/lib/validation.ts index b2295fa..6a9757f 100644 --- a/backend/src/lib/validation.ts +++ b/backend/src/lib/validation.ts @@ -40,13 +40,19 @@ export const UsageUpdateSchema = z }) .strict(); +export const MqttPayloadSchema = z.object({ + meterId: MeterIdSchema, + units: z.number().int("units must be an integer").positive("units must be positive"), + cost: z.number().int("cost must be an integer").positive("cost must be positive"), +}); + export const MeterNoteSchema = z .object({ text: z .string() .trim() .min(1, "text is required") - .max(2000, "text must be at most 2000 characters"), + .max(1000, "text must be at most 1000 characters"), }) .strict(); diff --git a/backend/src/middleware/rateLimit.ts b/backend/src/middleware/rateLimit.ts index ee49cf7..1de4208 100644 --- a/backend/src/middleware/rateLimit.ts +++ b/backend/src/middleware/rateLimit.ts @@ -17,3 +17,11 @@ export const readLimiter = rateLimit({ standardHeaders: true, legacyHeaders: false, }); + +export const paymentsLimiter = rateLimit({ + windowMs, + max: parseInt(process.env.PAYMENTS_RATE_LIMIT_MAX ?? '10', 10), + standardHeaders: true, + legacyHeaders: false, + message: { error: 'Too many payment requests', code: 'RATE_LIMITED' }, +}); diff --git a/backend/src/routes/meters.ts b/backend/src/routes/meters.ts index e15ff85..7a0c8e5 100644 --- a/backend/src/routes/meters.ts +++ b/backend/src/routes/meters.ts @@ -1,7 +1,6 @@ import { Router } from "express"; import * as StellarSdk from "@stellar/stellar-sdk"; -import { contractQuery, adminInvoke } from "../lib/stellar.js"; -import { StellarService, stellarService, server } from "../lib/stellar.js"; +import { StellarService, server } from "../lib/stellar.js"; import { getUsageHistory, persistAndSubmitUsageEvent, @@ -10,6 +9,8 @@ import { import { addMeterNote, getLatestMeterNotes, + getAllMeterNotes, + deleteMeterNote, } from "../lib/meterNotes.js"; import { asyncHandler } from "../lib/asyncHandler.js"; import { @@ -37,170 +38,21 @@ export function createMeterRouter(stellar: StellarService) { * * Fixes #268. */ + /** + * GET /api/meters?page=1&pageSize=20 — list all meters with pagination + * + * Registered BEFORE /:id so the literal string "meters" is never matched + * as a meter ID parameter. + */ meterRouter.get( "/", asyncHandler(async (req, res) => { const page = Math.max(1, Number(req.query.page ?? 1) || 1); const pageSize = Math.min( 100, - Math.max(1, Number(req.query.pageSize ?? 25) || 25), + Math.max(1, Number(req.query.pageSize ?? 20) || 20), ); - const header = "owner,active,units_used,plan,last_payment,expires_at,daily_limit"; - const rows = meters.map((m: any) => - [m.owner, m.active, m.units_used, m.plan, m.last_payment, m.expires_at, m.daily_limit].join(",") - ); - res.setHeader("Content-Type", "text/csv"); - res.setHeader("Content-Disposition", "attachment; filename=meters.csv"); - return res.send([header, ...rows].join("\n")); - }), -); - -/** - * GET /api/meters/search?q=&page=&pageSize= — case-insensitive substring - * search across meter ID and owner address, paginated. - */ -meterRouter.get( - "/search", - asyncHandler(async (req, res) => { - const q = String(req.query.q ?? "").toLowerCase().trim(); - const page = Math.max(1, Number(req.query.page ?? 1) || 1); - const pageSize = Math.min( - 100, - Math.max(1, Number(req.query.pageSize ?? 25) || 25), - ); - - const result = await contractQuery("get_all_meters", []); - const meters = (StellarSdk.scValToNative(result) as any[]) ?? []; - - const matches = q - ? meters.filter((m: any) => { - const id = String(m.meter_id ?? m.id ?? "").toLowerCase(); - const owner = String(m.owner ?? "").toLowerCase(); - return id.includes(q) || owner.includes(q); - }) - : meters; - - const total = matches.length; - const offset = (page - 1) * pageSize; - const data = matches.slice(offset, offset + pageSize); - - res.json({ - data, - pagination: { - page, - pageSize, - total, - totalPages: Math.max(1, Math.ceil(total / pageSize)), - }, - }); - }), -); - -/** GET /api/meters/:id — get meter status */ -meterRouter.get( - "/:id", - asyncHandler(async (req, res) => { - const result = await contractQuery("get_meter", [ - StellarSdk.nativeToScVal(req.params.id, { type: "symbol" }), - ]); - res.json({ meter: StellarSdk.scValToNative(result) }); - }), -); - -/** GET /api/meters/:id/access — check if meter is active */ -meterRouter.get( - "/:id/access", - asyncHandler(async (req, res) => { - const result = await contractQuery("check_access", [ - StellarSdk.nativeToScVal(req.params.id, { type: "symbol" }), - ]); - res.json({ active: StellarSdk.scValToNative(result) }); - }), -); - -/** GET /api/meters/:id/history — paginated local usage history */ -meterRouter.get("/:id/history", (req, res) => { - const page = Math.max(1, Number(req.query.page ?? 1) || 1); - const pageSize = Math.min( - 100, - Math.max(1, Number(req.query.pageSize ?? 25) || 25), - ); - - try { - const history = getUsageHistory(req.params.id, page, pageSize); - res.json(history); - } catch (err: any) { - res.status(500).json({ error: err.message }); - } -}); - -/** GET /api/meters/owner/:address — list all meters for an owner */ -meterRouter.get( - "/owner/:address", - asyncHandler(async (req, res) => { - const result = await contractQuery("get_meters_by_owner", [ - StellarSdk.nativeToScVal(req.params.address, { type: "address" }), - ]); - res.json({ meters: StellarSdk.scValToNative(result) }); - }), -); - -/** POST /api/meters — register a new meter (admin only) */ -meterRouter.post( - "/", - validateRequest({ body: RegisterMeterSchema }), - asyncHandler(async (req, res) => { - const { meter_id, owner } = req.body; - - const hash = await adminInvoke("register_meter", [ - StellarSdk.nativeToScVal(meter_id, { type: "symbol" }), - StellarSdk.nativeToScVal(owner, { type: "address" }), - ]); - res.json({ hash }); - }), -); - -/** POST /api/meters/:id/usage — IoT oracle reports usage */ -meterRouter.post("/:id/usage", async (req, res) => { - const { units, cost } = req.body as { units: unknown; cost: unknown }; - - if (units == null || cost == null) { - return res.status(400).json({ error: "units and cost are required" }); - } - - const unitsNum = Number(units); - const costNum = Number(cost); - - if (!Number.isFinite(unitsNum) || !Number.isFinite(costNum)) { - return res.status(400).json({ error: "units and cost must be valid numbers" }); - } - - if (!Number.isInteger(unitsNum) || !Number.isInteger(costNum)) { - return res.status(400).json({ error: "units and cost must be integers" }); - } - - if (unitsNum <= 0 || costNum <= 0) { - return res.status(400).json({ error: "units and cost must be positive" }); - } - - try { - const event = await persistAndSubmitUsageEvent({ - meterId: req.params.id, - units: unitsNum, - cost: costNum, - sourceTopic: null, - }); - - res.json({ - event, - hash: event.on_chain_tx_hash, - queued: !event.on_chain_tx_hash, - }); - } catch (err: any) { - res.status(500).json({ error: err.message }); - } -}); const result = await stellar.query("get_all_meters", []); const allMeters = (StellarSdk.scValToNative(result) as any[]) ?? []; @@ -469,12 +321,60 @@ meterRouter.post("/:id/usage", async (req, res) => { return res.status(404).json({ error: "Meter not found", code: "NOT_FOUND" }); } - const note = addMeterNote(meterId, req.body.text); + const note = addMeterNote(meterId, req.body.text, req.ip); invalidateCache(`/api/meters/${meterId}`); res.status(201).json(note); }), ); + /** GET /api/meters/:id/notes — all notes for a meter (paginated, no auth required) */ + meterRouter.get( + "/:id/notes", + asyncHandler(async (req, res) => { + const page = Math.max(1, Number(req.query.page ?? 1) || 1); + const pageSize = Math.min(100, Math.max(1, Number(req.query.pageSize ?? 20) || 20)); + const result = getAllMeterNotes(req.params.id, page, pageSize); + res.json(result); + }), + ); + + /** POST /api/meters/:id/notes — create a note (admin only) */ + meterRouter.post( + "/:id/notes", + requireAdminKey, + validateRequest({ body: MeterNoteSchema }), + asyncHandler(async (req, res) => { + const meterId = req.params.id; + try { + await stellar.query("get_meter", [ + StellarSdk.nativeToScVal(meterId, { type: "symbol" }), + ]); + } catch { + return res.status(404).json({ error: "Meter not found", code: "NOT_FOUND" }); + } + const note = addMeterNote(meterId, req.body.text, req.ip); + invalidateCache(`/api/meters/${meterId}`); + res.status(201).json(note); + }), + ); + + /** DELETE /api/meters/:id/notes/:noteId — hard-delete a note (admin only) */ + meterRouter.delete( + "/:id/notes/:noteId", + requireAdminKey, + asyncHandler(async (req, res) => { + const noteId = Number(req.params.noteId); + if (!Number.isInteger(noteId) || noteId <= 0) { + return res.status(400).json({ error: "Invalid noteId", code: "VALIDATION_ERROR" }); + } + const deleted = deleteMeterNote(noteId); + if (!deleted) { + return res.status(404).json({ error: "Note not found", code: "NOT_FOUND" }); + } + res.json({ deleted: true, noteId }); + }), + ); + /** GET /api/meters/:id/access — check if meter is active */ meterRouter.get( "/:id/access", @@ -658,16 +558,20 @@ meterRouter.post("/:id/usage", async (req, res) => { }), ); - /** GET /api/meters/:id/history — paginated local usage history */ + /** GET /api/meters/:id/history?page=1&pageSize=20 — paginated local usage history */ meterRouter.get("/:id/history", (req, res) => { - const page = Math.max(1, Number(req.query.page ?? 1) || 1); - const pageSize = Math.min( - 100, - Math.max(1, Number(req.query.pageSize ?? 25) || 25), - ); + const rawPage = Number(req.query.page ?? 1); + const rawPageSize = Number(req.query.pageSize ?? 20); + + if (!Number.isInteger(rawPage) || rawPage < 1) { + return res.status(400).json({ error: "page must be a positive integer", code: "VALIDATION_ERROR" }); + } + if (!Number.isInteger(rawPageSize) || rawPageSize < 1 || rawPageSize > 100) { + return res.status(400).json({ error: "pageSize must be between 1 and 100", code: "VALIDATION_ERROR" }); + } try { - const history = getUsageHistory(req.params.id, page, pageSize); + const history = getUsageHistory(req.params.id, rawPage, rawPageSize); res.json(history); } catch (err: any) { res.status(500).json({ error: err.message, code: "INTERNAL_ERROR" }); diff --git a/backend/src/routes/payments.ts b/backend/src/routes/payments.ts index 07bb9ed..1ce094d 100644 --- a/backend/src/routes/payments.ts +++ b/backend/src/routes/payments.ts @@ -31,22 +31,6 @@ function cleanExpiredIdempotencyKeys() { } } -paymentsRouter.post( - "/", - asyncHandler(async (req, res) => { - cleanExpiredIdempotencyKeys(); - - const idempotencyKey = ( - req.headers["idempotency-key"] ?? req.headers["x-idempotency-key"] - ) as string | undefined; - - if (idempotencyKey && typeof idempotencyKey === "string" && idempotencyKey.trim().length > 0) { - const cached = idempotencyCache.get(idempotencyKey.trim()); - if (cached && Date.now() - cached.createdAt <= IDEMPOTENCY_TTL_MS) { - return res.json({ hash: cached.hash }); - } - } - paymentsRouter.post( "/", idempotency(), @@ -66,10 +50,6 @@ paymentsRouter.post( StellarSdk.nativeToScVal(payer, { type: "address" }), ]); - if (idempotencyKey && typeof idempotencyKey === "string" && idempotencyKey.trim().length > 0) { - idempotencyCache.set(idempotencyKey.trim(), { hash, createdAt: Date.now() }); - } - return res.json({ hash }); }), ); @@ -202,32 +182,6 @@ paymentsRouter.get( ); - try { - StellarSdk.StrKey.decodeEd25519PublicKey(address); - } catch { - return res.status(400).json({ error: "Invalid Stellar address" }); - } - - try { - const records = await fetchPaymentEvents(address, sort, days); - const total = records.length; - const start = (page - 1) * limit; - const paginated = records.slice(start, start + limit); - - return res.json({ - payments: paginated, - pagination: { page, limit, total, pages: Math.ceil(total / limit) }, - }); - } catch (err: any) { - logger.error("payments route error:", err); - if (err?.code === 'RPC_ERROR' || err?.isRpcError) { - return res.status(502).json({ error: err.message ?? "RPC request failed", code: "RPC_ERROR" }); - } - return res.status(500).json({ error: err.message ?? "Failed to fetch payment history" }); - } - }), -); - /** * GET /api/payments/history/:address?from=&to=&limit=50&page=1 * diff --git a/backend/src/routes/provider.ts b/backend/src/routes/provider.ts index 1c97f1b..8f0b64b 100644 --- a/backend/src/routes/provider.ts +++ b/backend/src/routes/provider.ts @@ -110,7 +110,7 @@ providerRouter.post( }), asyncHandler(async (req, res) => { const { webhook_url } = req.body; - registerWebhook(webhook_url); + registerWebhook("provider", webhook_url); logger.info("Provider webhook registered", { webhook_url }); return res.json({ message: "Webhook registered", webhook_url }); }), diff --git a/backend/src/routes/stats.ts b/backend/src/routes/stats.ts index 32b1835..1984d37 100644 --- a/backend/src/routes/stats.ts +++ b/backend/src/routes/stats.ts @@ -1,31 +1,12 @@ import { Router } from "express"; -import { getTopConsumers } from "../lib/usageEvents.js"; - -export const statsRouter = Router(); - -function requireAdminKey(req: any, res: any, next: any) { - const adminKey = process.env.ADMIN_API_KEY; - const provided = req.headers["x-admin-key"]; - if (!adminKey || provided !== adminKey) { - return res.status(401).json({ error: "Valid admin key required" }); - } - return next(); -} - -/** - * GET /api/stats/top-consumers?days=30 - * - * Returns the top 10 meters ranked by total units used over the given - * window (default 30 days). Requires the X-Admin-Key header. - */ -statsRouter.get("/top-consumers", requireAdminKey, (req, res) => { - const days = Math.max(1, Number(req.query.days ?? 30) || 30); - const consumers = getTopConsumers(days, 10); - res.json(consumers); -}); import * as StellarSdk from "@stellar/stellar-sdk"; import { server, CONTRACT_ID } from "../lib/stellar.js"; +import { stellarService } from "../lib/stellar.js"; import { asyncHandler } from "../lib/asyncHandler.js"; +import { getTopConsumers } from "../lib/usageEvents.js"; +import { register } from "../lib/metrics.js"; +import { logger } from "../lib/logger.js"; +import { requireAdminKey } from "../middleware/adminAuth.js"; export const statsRouter = Router(); @@ -34,119 +15,15 @@ const MAX_DAYS = 90; const CACHE_TTL_MS = 60_000; interface RevenueHistoryEntry { - date: string; // YYYY-MM-DD + date: string; revenue_xlm: number; } -const revenueHistoryCache = new Map< - number, - { data: RevenueHistoryEntry[]; ts: number } ->(); - -/** - * GET /api/stats/revenue-history?days=30 - * - * Aggregates Soroban "payment" contract events by day for the requested - * window (capped at MAX_DAYS) and returns daily revenue totals in XLM. - */ -statsRouter.get( - "/revenue-history", - asyncHandler(async (req, res) => { - const requestedDays = parseInt((req.query.days as string) ?? String(DEFAULT_DAYS), 10); - const days = Math.min( - MAX_DAYS, - Math.max(1, Number.isFinite(requestedDays) ? requestedDays : DEFAULT_DAYS), - ); - - const cached = revenueHistoryCache.get(days); - if (cached && Date.now() - cached.ts < CACHE_TTL_MS) { - return res.json({ history: cached.data }); - } - - const history = await fetchRevenueHistory(days); - revenueHistoryCache.set(days, { data: history, ts: Date.now() }); - - res.json({ history }); - }), -); - -// ── Helpers ─────────────────────────────────────────────────────────────────── - -async function fetchRevenueHistory(days: number): Promise { - const response = await (server as any).getEvents({ - startLedger: 1, - filters: [ - { - type: "contract", - contractIds: [CONTRACT_ID], - topics: [[StellarSdk.xdr.ScVal.scvSymbol("payment").toXDR("base64")]], - }, - ], - limit: 1000, - }); - - const totalsByDay = new Map(); - const cutoff = Date.now() - days * 24 * 60 * 60 * 1000; - - for (const event of response?.events ?? []) { - try { - const parsed = parsePaymentEvent(event); - if (!parsed) continue; - if (new Date(parsed.date).getTime() < cutoff) continue; - - const day = parsed.date.slice(0, 10); // YYYY-MM-DD - totalsByDay.set(day, (totalsByDay.get(day) ?? 0) + parsed.amountXlm); - } catch { - // skip malformed events - } - } - - return buildDayRange(days).map((date) => ({ - date, - revenue_xlm: totalsByDay.get(date) ?? 0, - })); -} +const revenueHistoryCache = new Map(); -function buildDayRange(days: number): string[] { - const result: string[] = []; - const now = new Date(); - for (let i = days - 1; i >= 0; i--) { - const d = new Date(now); - d.setUTCDate(d.getUTCDate() - i); - result.push(d.toISOString().slice(0, 10)); - } - return result; -} - -function parsePaymentEvent(event: any): { date: string; amountXlm: number } | null { - const dataXdr = event.value ?? event.data; - if (!dataXdr) return null; - - const dataVal = StellarSdk.xdr.ScVal.fromXDR(dataXdr, "base64"); - const native = StellarSdk.scValToNative(dataVal); - if (!Array.isArray(native) || native.length < 1) return null; - - const amountXlm = Number(native[0]) / 10_000_000; - const date = event.ledgerClosedAt - ? new Date(event.ledgerClosedAt).toISOString() - : new Date().toISOString(); - - return { date, amountXlm }; -} -import { stellarService } from "../lib/stellar.js"; -import { register } from "../lib/metrics.js"; -import { logger } from "../lib/logger.js"; - -// Cache for the existing contract-based stats endpoint (30s TTL) let contractCache: { data: object; expiresAt: number } | null = null; - -// Cache for the prom-client metrics summary endpoint (15s TTL) let metricsCache: { data: object; expiresAt: number } | null = null; - -// Cache for meters-by-plan breakdown (30s TTL) let metersByPlanCache: { data: object; expiresAt: number } | null = null; - -// Cache for meter counts grouped by plan (30s TTL) let meterPlanCache: { data: object; expiresAt: number } | null = null; export function __resetStatsCache() { @@ -156,36 +33,20 @@ export function __resetStatsCache() { meterPlanCache = null; } - -type MeterPlanBreakdown = { - Daily: number; - Weekly: number; - Usage: number; - total: number; -}; - -const emptyPlanBreakdown = (): MeterPlanBreakdown => ({ - Daily: 0, - Weekly: 0, - Usage: 0, - total: 0, -}); +type MeterPlanBreakdown = { Daily: number; Weekly: number; Usage: number; total: number }; +const emptyPlanBreakdown = (): MeterPlanBreakdown => ({ Daily: 0, Weekly: 0, Usage: 0, total: 0 }); const normalizePlan = (plan: unknown): keyof Omit | null => { const raw = - typeof plan === "string" - ? plan - : typeof plan === "symbol" - ? plan.toString() - : plan && typeof plan === "object" - ? String( - (plan as { tag?: unknown; name?: unknown; variant?: unknown }).tag ?? - (plan as { tag?: unknown; name?: unknown; variant?: unknown }).name ?? - (plan as { tag?: unknown; name?: unknown; variant?: unknown }).variant ?? - "", - ) - : ""; - + typeof plan === "string" ? plan + : typeof plan === "symbol" ? plan.toString() + : plan && typeof plan === "object" + ? String( + (plan as { tag?: unknown; name?: unknown; variant?: unknown }).tag ?? + (plan as { tag?: unknown; name?: unknown; variant?: unknown }).name ?? + (plan as { tag?: unknown; name?: unknown; variant?: unknown }).variant ?? "", + ) + : ""; const normalized = raw.toLowerCase().replace(/[^a-z]/g, ""); if (normalized === "daily") return "Daily"; if (normalized === "weekly") return "Weekly"; @@ -193,33 +54,7 @@ const normalizePlan = (plan: unknown): keyof Omit | return null; }; -/** - * GET /api/stats/meters-by-plan — meter count breakdown by payment plan. - * Response is cached for 30 seconds and always includes all known plan keys. - * - * Closes #461. - */ -statsRouter.get("/meters-by-plan", asyncHandler(async (_req, res) => { - if (meterPlanCache && Date.now() < meterPlanCache.expiresAt) { - return res.json(meterPlanCache.data); - } - - const result = await stellarService.query("get_all_meters", []); - const meters = (StellarSdk.scValToNative(result) as any[]) ?? []; - const data = emptyPlanBreakdown(); - - for (const meter of meters) { - const plan = normalizePlan(meter?.plan); - if (plan) data[plan] += 1; - } - - data.total = meters.length; - meterPlanCache = { data, expiresAt: Date.now() + 30_000 }; - res.json(data); -})); -/** - * GET /api/stats — contract-derived meter statistics (existing endpoint) - */ +/** GET /api/stats — contract-derived meter statistics */ statsRouter.get("/", asyncHandler(async (_req, res) => { if (contractCache && Date.now() < contractCache.expiresAt) { return res.json(contractCache.data); @@ -239,65 +74,43 @@ statsRouter.get("/", asyncHandler(async (_req, res) => { ]); revenue = Number(StellarSdk.scValToNative(rev)); } else { - logger.warn('ADMIN_ADDRESS environment variable is not set; provider revenue query skipped'); + logger.warn("ADMIN_ADDRESS not set; provider revenue query skipped"); } const avgUnitsPerMeter = total > 0 ? units / total : 0; const avgRevenue = total > 0 ? revenue / total : 0; - const data = { - totalMeters: total, - activeMeters: active, - inactiveMeters: total - active, - totalUnits: units, - avgUnitsPerMeter, - totalRevenue: revenue, - avgRevenue, - }; + const data = { totalMeters: total, activeMeters: active, inactiveMeters: total - active, + totalUnits: units, avgUnitsPerMeter, totalRevenue: revenue, avgRevenue }; contractCache = { data, expiresAt: Date.now() + 30_000 }; res.json(data); })); -/** - * GET /api/stats/meters-by-plan — meter count breakdown by plan type. - * Returns zero counts for plan types with no meters. - * Cached for 30 seconds. - * - * Closes #461. - */ +/** GET /api/stats/meters-by-plan — meter count breakdown by payment plan */ statsRouter.get("/meters-by-plan", asyncHandler(async (_req, res) => { - if (metersByPlanCache && Date.now() < metersByPlanCache.expiresAt) { - return res.json(metersByPlanCache.data); + if (meterPlanCache && Date.now() < meterPlanCache.expiresAt) { + return res.json(meterPlanCache.data); } const result = await stellarService.query("get_all_meters", []); const meters = (StellarSdk.scValToNative(result) as any[]) ?? []; - - const counts: Record = { Daily: 0, Weekly: 0, Usage: 0 }; + const data = emptyPlanBreakdown(); for (const meter of meters) { - const plan = String(meter.plan); - if (plan in counts) counts[plan]++; + const plan = normalizePlan(meter?.plan); + if (plan) data[plan] += 1; } - - const data = { ...counts, total: meters.length }; - metersByPlanCache = { data, expiresAt: Date.now() + 30_000 }; + data.total = meters.length; + meterPlanCache = { data, expiresAt: Date.now() + 30_000 }; res.json(data); })); -/** - * GET /api/stats/summary — prom-client counter/gauge snapshot for the admin - * dashboard. Does not require Prometheus or Grafana to be running. - * Response is cached for 15 seconds. - * - * Closes #344. - */ +/** GET /api/stats/summary — prom-client counter/gauge snapshot for admin dashboard */ statsRouter.get("/summary", asyncHandler(async (_req, res) => { if (metricsCache && Date.now() < metricsCache.expiresAt) { return res.json(metricsCache.data); } const metrics = await register.getMetricsAsJSON(); - const find = (name: string): number => { const metric = metrics.find((m: any) => m.name === name); if (!metric?.values) return 0; @@ -310,7 +123,75 @@ statsRouter.get("/summary", asyncHandler(async (_req, res) => { activeMeters: find("solargrid_active_meters"), paymentVolumeXlm: find("solargrid_payment_volume_xlm"), }; - metricsCache = { data, expiresAt: Date.now() + 15_000 }; res.json(data); })); + +/** GET /api/stats/revenue-history?days=30 — daily revenue aggregated from Soroban events */ +statsRouter.get("/revenue-history", asyncHandler(async (req, res) => { + const requestedDays = parseInt((req.query.days as string) ?? String(DEFAULT_DAYS), 10); + const days = Math.min(MAX_DAYS, Math.max(1, Number.isFinite(requestedDays) ? requestedDays : DEFAULT_DAYS)); + + const cached = revenueHistoryCache.get(days); + if (cached && Date.now() - cached.ts < CACHE_TTL_MS) { + return res.json({ history: cached.data }); + } + + const history = await fetchRevenueHistory(days); + revenueHistoryCache.set(days, { data: history, ts: Date.now() }); + res.json({ history }); +})); + +/** GET /api/stats/top-consumers?days=30 — top 10 meters by energy consumed */ +statsRouter.get("/top-consumers", requireAdminKey, (req, res) => { + const days = Math.max(1, Number(req.query.days ?? 30) || 30); + res.json(getTopConsumers(days, 10)); +}); + +// ── Helpers ─────────────────────────────────────────────────────────────────── + +async function fetchRevenueHistory(days: number): Promise { + const response = await (server as any).getEvents({ + startLedger: 1, + filters: [{ type: "contract", contractIds: [CONTRACT_ID], + topics: [[StellarSdk.xdr.ScVal.scvSymbol("payment").toXDR("base64")]] }], + limit: 1000, + }); + + const totalsByDay = new Map(); + const cutoff = Date.now() - days * 24 * 60 * 60 * 1000; + + for (const event of response?.events ?? []) { + try { + const parsed = parsePaymentEvent(event); + if (!parsed) continue; + if (new Date(parsed.date).getTime() < cutoff) continue; + const day = parsed.date.slice(0, 10); + totalsByDay.set(day, (totalsByDay.get(day) ?? 0) + parsed.amountXlm); + } catch { /* skip malformed events */ } + } + + return buildDayRange(days).map((date) => ({ date, revenue_xlm: totalsByDay.get(date) ?? 0 })); +} + +function buildDayRange(days: number): string[] { + const result: string[] = []; + const now = new Date(); + for (let i = days - 1; i >= 0; i--) { + const d = new Date(now); + d.setUTCDate(d.getUTCDate() - i); + result.push(d.toISOString().slice(0, 10)); + } + return result; +} + +function parsePaymentEvent(event: any): { date: string; amountXlm: number } | null { + const dataXdr = event.value ?? event.data; + if (!dataXdr) return null; + const dataVal = StellarSdk.xdr.ScVal.fromXDR(dataXdr, "base64"); + const native = StellarSdk.scValToNative(dataVal); + if (!Array.isArray(native) || native.length < 1) return null; + const amountXlm = Number(native[0]) / 10_000_000; + const date = event.ledgerClosedAt ? new Date(event.ledgerClosedAt).toISOString() : new Date().toISOString(); + return { date, amountXlm }; +} diff --git a/backend/src/routes/webhooks.ts b/backend/src/routes/webhooks.ts index 010d782..d38caf2 100644 --- a/backend/src/routes/webhooks.ts +++ b/backend/src/routes/webhooks.ts @@ -113,18 +113,6 @@ webhookRouter.post( secret?: string; }; - // For now, store in environment variables (in production, use a database). - process.env.PROVIDER_WEBHOOK_URL = webhook_url; - - let secretHash: string | undefined; - if (secret) { - process.env.PROVIDER_WEBHOOK_SECRET = secret; - secretHash = crypto.createHash("sha256").update(secret).digest("hex"); - } - - logger.info("Low-balance webhook registered", { - webhook_url, - secretHash, const providerId = getProviderId(req); if (!providerId) { return res.status(400).json({ @@ -133,19 +121,26 @@ webhookRouter.post( }); } - const { webhook_url } = req.body; + // Store in environment for legacy bridge compatibility + process.env.PROVIDER_WEBHOOK_URL = webhook_url; + + let secretHash: string | undefined; + if (secret) { + process.env.PROVIDER_WEBHOOK_SECRET = secret; + secretHash = crypto.createHash("sha256").update(secret).digest("hex"); + } const record = registerWebhook(providerId, webhook_url); logger.info("Low-balance webhook registered", { provider_id: providerId, webhook_url, + secretHash, }); return res.status(200).json({ message: "Webhook registered successfully", webhook_url, - secretHash, provider_id: providerId, id: record.id, created_at: record.created_at,