Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 7 additions & 7 deletions src/app/api/materials/upload/route.js
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@ import {
validateUploadPayload,
validateUploadFileMetadata,
} from '@/lib/api/validation'
import { pinata } from '@/lib/pinata'
import { pinata, callPinata } from '@/lib/pinata'
import { validateUploadedFile } from '@/lib/ipfs/uploadValidator'

const MAX_FILE_SIZE_BYTES = 10 * 1024 * 1024 // 10 MB
Expand Down Expand Up @@ -119,13 +119,13 @@ export async function POST(request) {
// All checks passed — dispatch to Pinata
const results = {}

const uploadedFile = await pinata.upload.public.file(file)
const fileUrl = await pinata.gateways.public.convert(uploadedFile.cid)
const uploadedFile = await callPinata('upload.file', () => pinata.upload.public.file(file))
const fileUrl = await callPinata('gateway.file', () => pinata.gateways.public.convert(uploadedFile.cid))
results.fileUrl = fileUrl

if (image) {
const fileThumb = await pinata.upload.public.file(image)
const imgUrl = await pinata.gateways.public.convert(fileThumb.cid)
const fileThumb = await callPinata('upload.thumbnail', () => pinata.upload.public.file(image))
const imgUrl = await callPinata('gateway.thumbnail', () => pinata.gateways.public.convert(fileThumb.cid))
results.imgUrl = imgUrl
}

Expand Down Expand Up @@ -168,8 +168,8 @@ export async function POST(request) {
timestamp: new Date().toISOString(),
}

const uploadedJson = await pinata.upload.public.json(metadataJSON)
const jsonUrl = await pinata.gateways.public.convert(uploadedJson.cid)
const uploadedJson = await callPinata('upload.metadata', () => pinata.upload.public.json(metadataJSON))
const jsonUrl = await callPinata('gateway.metadata', () => pinata.gateways.public.convert(uploadedJson.cid))
results.metadataUrl = jsonUrl

auditLog({ event: 'upload_complete', route: 'materials/upload', method: 'POST', status: 200 })
Expand Down
10 changes: 5 additions & 5 deletions src/app/api/ready/route.js
Original file line number Diff line number Diff line change
Expand Up @@ -2,12 +2,13 @@ export const dynamic = "force-dynamic";

import { NextResponse } from "next/server";
import { getDb } from "@/lib/mongodb";
import { pinata } from "@/lib/pinata";
import { pinata, callPinata } from "@/lib/pinata";
import { verifyEmailConnection } from "@/lib/email";
import { withApiHardening } from "@/lib/api/hardening";
import { Server, rpc } from "@stellar/stellar-sdk";
import { Server } from "@stellar/stellar-sdk";
import { HORIZON_URL, STELLAR_RPC_URL } from "@/lib/config/chain";
import { setGauge } from "@/lib/telemetry/metrics";
import { getRpcHealth } from "@/lib/stellar/rpcClient";

/**
* Readiness probe (#20): "can this instance actually serve traffic right now?"
Expand Down Expand Up @@ -58,7 +59,7 @@ export async function GET(request) {
status.pinata = "offline: PINATA_JWT not configured";
recordFailure("pinata");
} else {
await pinata.testAuthentication();
await callPinata('testAuthentication', () => pinata.testAuthentication());
status.pinata = "online";
}
} catch (err) {
Expand All @@ -80,8 +81,7 @@ export async function GET(request) {
const horizonServer = new Server(HORIZON_URL);
await horizonServer.root();

const rpcServer = new rpc.Server(STELLAR_RPC_URL);
const health = await rpcServer.getHealth();
const health = await getRpcHealth();

if (health && health.status === "healthy") {
status.stellar = "online";
Expand Down
56 changes: 7 additions & 49 deletions src/app/api/upload/route.js
Original file line number Diff line number Diff line change
Expand Up @@ -9,12 +9,11 @@ import {
validateUploadPayload,
} from '@/lib/api/validation'
import {
retryWithBackoff,
validateGatewayUrl,
validatePinataResponse,
} from '@/lib/api/storage'
import { getDb } from '@/lib/mongodb'
import { pinata } from '@/lib/pinata'
import { pinata, callPinata } from '@/lib/pinata'
import { storeManifest } from '@/lib/provenance/registry'
import { hashFileBytes } from '@/lib/provenance/manifest'
import { validateUploadedFile } from '@/lib/ipfs/uploadValidator'
Expand Down Expand Up @@ -215,23 +214,10 @@ export async function POST(request) {

// 4️⃣ Upload the main file
try {
const uploadedFile = await retryWithBackoff(
() => pinata.upload.public.file(file),
3,
1000,
(err, attempt) => {
console.warn(`[Storage] Document upload attempt ${attempt} failed: ${err.message}`);
}
)
const uploadedFile = await callPinata('upload.file', () => pinata.upload.public.file(file))
validatePinataResponse(uploadedFile, 'document')

const fileUrl = await retryWithBackoff(
() => pinata.gateways.public.convert(uploadedFile.cid),
3,
1000,
(err, attempt) => {
console.warn(`[Storage] Document gateway conversion attempt ${attempt} failed: ${err.message}`);
}
const fileUrl = await callPinata('gateway.convert', () => pinata.gateways.public.convert(uploadedFile.cid))
)
validateGatewayUrl(fileUrl, 'document')
results.fileUrl = fileUrl
Expand All @@ -253,24 +239,10 @@ export async function POST(request) {
// 5️⃣ Upload thumbnail (if provided)
if (image) {
try {
const fileThumb = await retryWithBackoff(
() => pinata.upload.public.file(image),
3,
1000,
(err, attempt) => {
console.warn(`[Storage] Thumbnail upload attempt ${attempt} failed: ${err.message}`);
}
)
const fileThumb = await callPinata('upload.thumbnail', () => pinata.upload.public.file(image))
validatePinataResponse(fileThumb, 'thumbnail')

const imgUrl = await retryWithBackoff(
() => pinata.gateways.public.convert(fileThumb.cid),
3,
1000,
(err, attempt) => {
console.warn(`[Storage] Thumbnail gateway conversion attempt ${attempt} failed: ${err.message}`);
}
)
const imgUrl = await callPinata('gateway.thumbnail', () => pinata.gateways.public.convert(fileThumb.cid))
validateGatewayUrl(imgUrl, 'thumbnail')
results.imgUrl = imgUrl
} catch (err) {
Expand Down Expand Up @@ -344,24 +316,10 @@ export async function POST(request) {

// 7️⃣ Upload metadata JSON to Pinata
try {
const uploadedJson = await retryWithBackoff(
() => pinata.upload.public.json(metadataJSON),
3,
1000,
(err, attempt) => {
console.warn(`[Storage] Metadata upload attempt ${attempt} failed: ${err.message}`);
}
)
const uploadedJson = await callPinata('upload.metadata', () => pinata.upload.public.json(metadataJSON))
validatePinataResponse(uploadedJson, 'metadata')

const jsonUrl = await retryWithBackoff(
() => pinata.gateways.public.convert(uploadedJson.cid),
3,
1000,
(err, attempt) => {
console.warn(`[Storage] Metadata gateway conversion attempt ${attempt} failed: ${err.message}`);
}
)
const jsonUrl = await callPinata('gateway.metadata', () => pinata.gateways.public.convert(uploadedJson.cid))
validateGatewayUrl(jsonUrl, 'metadata')
results.metadataUrl = jsonUrl
} catch (err) {
Expand Down
12 changes: 12 additions & 0 deletions src/lib/api/hardening.js
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ import { NextResponse } from "next/server";
import { auditLog } from "./audit";
import { checkRateLimit } from "./rateLimit";
import { ValidationError } from "./validation";
import { DependencyError } from "@/lib/resilience/index.js";
import { captureException } from "@/lib/sentry";
import { runWithContext, currentTraceparent, currentCorrelationId } from "@/lib/telemetry/context";
import { withSpan } from "@/lib/telemetry/tracing";
Expand Down Expand Up @@ -141,6 +142,17 @@ export async function withApiHardening(request, options, handler) {
return finalize(res);
}

if (error instanceof DependencyError) {
auditLog({ event: "dependency_failed", route, method, status: error.statusCode || 503, dependency: error.dependency, action: error.action, reason: error.message });
incrementCounter("http_requests_total", { route, method, outcome: "dependency_error" });
incrementCounter("rpc_errors_total", { operation: `${error.dependency}.${error.action}` });
const statusCode = error.statusCode || 503;
return finalize(NextResponse.json(
{ error: error.userMessage || "Service temporarily unavailable", code: "dependency_error", dependency: error.dependency },
{ status: statusCode, headers: { "x-correlation-id": currentCorrelationId() } }
));
}

incrementCounter("http_requests_total", { route, method, outcome: "error" });
captureException(error, { route, method, correlationId: currentCorrelationId() });
return finalize(NextResponse.json({ error: "Internal Server Error", code: "internal_error" }, { status: 500 }));
Expand Down
76 changes: 60 additions & 16 deletions src/lib/cache/redis.js
Original file line number Diff line number Diff line change
@@ -1,36 +1,80 @@
import { createClient } from 'redis';
import { createCircuitBreaker, CircuitState, DependencyError } from "@/lib/resilience/index.js";
import { withTimeout } from "@/lib/resilience/timeout.js";
import { withRetry } from "@/lib/resilience/retry.js";
import { setGauge, incrementCounter } from "@/lib/telemetry/metrics";

const REDIS_TIMEOUT_MS = Number(process.env.REDIS_TIMEOUT_MS || 5000);
const REDIS_CONNECT_TIMEOUT_MS = Number(process.env.REDIS_CONNECT_TIMEOUT_MS || 10000);

const redisCircuitBreaker = createCircuitBreaker("redis", {
failureThreshold: Number(process.env.REDIS_CB_FAILURE_THRESHOLD || 3),
successThreshold: 2,
resetTimeoutMs: Number(process.env.REDIS_CB_RESET_TIMEOUT_MS || 30000),
onStateChange(name, from, to) {
setGauge("circuit_breaker_state", { dependency: name, state: to }, 1);
if (from && from !== to) {
setGauge("circuit_breaker_state", { dependency: name, state: from }, 0);
}
if (to === CircuitState.OPEN) {
incrementCounter("circuit_breaker_open_total", { dependency: name });
}
const level = to === CircuitState.OPEN ? "warn" : "info";
console[level](`[circuit-breaker] redis: ${from} -> ${to}`);
},
});

export function getRedisCircuitBreakerState() {
return redisCircuitBreaker.getState();
}

let client = null;

export async function getRedisClient() {
if (!process.env.REDIS_URL) return null;
if (redisCircuitBreaker.getState() === CircuitState.OPEN) {
throw new DependencyError({
dependency: "redis",
action: "connect",
retryable: false,
statusCode: 503,
userMessage: "Cache service is temporarily unavailable.",
});
}
if (!client) {
client = createClient({ url: process.env.REDIS_URL });
client = createClient({ url: process.env.REDIS_URL, socket: { reconnectStrategy: false } });
client.on('error', (err) => console.error('Redis error', err.message));
await client.connect();
await withTimeout(client.connect(), REDIS_CONNECT_TIMEOUT_MS, 'redis.connect');
}
return client;
}

export async function cacheGet(key) {
const redis = await getRedisClient();
if (!redis) return null;
try {
const val = await redis.get(key);
return val ? JSON.parse(val) : null;
} catch { return null; }
return redisCircuitBreaker.call(async () => {
const redis = await getRedisClient();
if (!redis) return null;
return withRetry(
async () => {
const val = await withTimeout(redis.get(key), REDIS_TIMEOUT_MS, 'redis.get');
return val ? JSON.parse(val) : null;
},
{ idempotent: true, maxAttempts: 2, baseDelayMs: 200, maxDelayMs: 1000 }
);
});
}

export async function cacheSet(key, value, ttlSeconds = 600) {
const redis = await getRedisClient();
if (!redis) return;
try {
await redis.set(key, JSON.stringify(value), { EX: ttlSeconds });
} catch { /* no-op */ }
return redisCircuitBreaker.call(async () => {
const redis = await getRedisClient();
if (!redis) return;
await withTimeout(redis.set(key, JSON.stringify(value), { EX: ttlSeconds }), REDIS_TIMEOUT_MS, 'redis.set');
});
}

export async function cacheDel(key) {
const redis = await getRedisClient();
if (!redis) return;
try { await redis.del(key); } catch { /* no-op */ }
return redisCircuitBreaker.call(async () => {
const redis = await getRedisClient();
if (!redis) return;
await withTimeout(redis.del(key), REDIS_TIMEOUT_MS, 'redis.del');
});
}
56 changes: 54 additions & 2 deletions src/lib/email.js
Original file line number Diff line number Diff line change
@@ -1,4 +1,56 @@
import nodemailer from "nodemailer";
import { createCircuitBreaker, CircuitState, DependencyError } from "@/lib/resilience/index.js";
import { withTimeout } from "@/lib/resilience/timeout.js";
import { withRetry } from "@/lib/resilience/retry.js";
import { setGauge, incrementCounter } from "@/lib/telemetry/metrics";
import { currentTraceparent } from "@/lib/telemetry/context";

const EMAIL_TIMEOUT_MS = Number(process.env.EMAIL_TIMEOUT_MS || 15000);

const emailCircuitBreaker = createCircuitBreaker("email", {
failureThreshold: Number(process.env.EMAIL_CB_FAILURE_THRESHOLD || 3),
successThreshold: 2,
resetTimeoutMs: Number(process.env.EMAIL_CB_RESET_TIMEOUT_MS || 30000),
onStateChange(name, from, to) {
setGauge("circuit_breaker_state", { dependency: name, state: to }, 1);
if (from && from !== to) {
setGauge("circuit_breaker_state", { dependency: name, state: from }, 0);
}
if (to === CircuitState.OPEN) {
incrementCounter("circuit_breaker_open_total", { dependency: name });
}
const level = to === CircuitState.OPEN ? "warn" : "info";
console[level](`[circuit-breaker] email: ${from} -> ${to}`);
},
});

export function getEmailCircuitBreakerState() {
return emailCircuitBreaker.getState();
}

async function withEmailResilience(action, fn, opts = {}) {
if (emailCircuitBreaker.getState() === CircuitState.OPEN) {
incrementCounter("rpc_errors_total", { operation: `email.${action}` });
throw new DependencyError({
dependency: "email",
action,
retryable: false,
statusCode: 503,
userMessage: "Email service is temporarily unavailable.",
});
}
return emailCircuitBreaker.call(async () => {
return withRetry(
async () => withTimeout(fn(), EMAIL_TIMEOUT_MS, `email.${action}`),
{
idempotent: opts.idempotent !== false,
maxAttempts: Number(process.env.EMAIL_RETRIES || 2),
baseDelayMs: 500,
maxDelayMs: 5000,
}
);
});
}

function createTransporter() {
// Prefer explicit SMTP settings; fallback to Gmail using EMAIL_USER/PASS
Expand Down Expand Up @@ -112,10 +164,10 @@ export async function sendWelcomeEmail(to, name) {
</body>
</html>`;

await transporter.sendMail({ from, to, subject, text, html });
await withEmailResilience('send.welcome', () => transporter.sendMail({ from, to, subject, text, html }), { idempotent: false });
}

export async function verifyEmailConnection() {
const transporter = createTransporter();
await transporter.verify();
await withEmailResilience('verify', () => transporter.verify());
}
Loading