fix: error-only Sentry telemetry, storm-proof capture, and crash fixes (#476)

Removes Sentry tracing entirely (BullMQ idle polling burned 4.8M transactions in 2 days at the baked 0.1 rate), decouples PostHog sampling, and replaces the type-only error scrub with a vetted-field sanitizer plus SafeError/ToolInputError contracts. One classified capture path with per-signature throttles and a per-process ceiling makes storms impossible (NODE-1E was 4,541 events from one 30s loop). Browser errors move to a dedicated web Sentry project with their own source maps. Adds the SNAPOTTER_TELEMETRY runtime kill switch and silences test fleets.

Crash fixes: remote 204/304 SSRF process kill (NODE-20), conversion-preset boot crash loop (NODE-21), Redis version preflight + unhandled subscribe rejection (NODE-1T), Sign PDF on plain-http origins (NODE-1K/1M), wavesurfer/pdf.js teardown rejections (NODE-1P/1N), bundle-import ZlibError to 400 (NODE-1Z), chart-maker input errors declassified (NODE-1H/1J), asset requests skip the session DB lookup (NODE-1D).
This commit is contained in:
SnapOtter
2026-07-10 21:41:49 +08:00
committed by GitHub
parent 3d1744aec8
commit ae6a4c8b7c
75 changed files with 2198 additions and 260 deletions
+20 -4
View File
@@ -5,21 +5,22 @@ import cors from "@fastify/cors";
import rateLimit from "@fastify/rate-limit";
import { trace } from "@opentelemetry/api";
import { getDispatcherStatus, initDispatcher, isGpuAvailable } from "@snapotter/ai";
import { APP_VERSION } from "@snapotter/shared";
import { APP_VERSION, SafeError } from "@snapotter/shared";
import { eq, sql } from "drizzle-orm";
import Fastify from "fastify";
import { env } from "./config.js";
import { closeDb, db, schema } from "./db/index.js";
import { runMigrations } from "./db/migrate.js";
import { startCancelListener, stopCancelListener } from "./jobs/cancel.js";
import { closeRedis, pingRedis } from "./jobs/connection.js";
import { assertRedisCompatible, closeRedis, pingRedis } from "./jobs/connection.js";
import { closeFlowProducer, closeQueueEvents, warmQueueEvents } from "./jobs/enqueue.js";
import { closeQueues, perPoolHealth, queueCounts } from "./jobs/queues.js";
import { enqueueSystemJob, SYSTEM_JOBS, scheduleSystemJobs } from "./jobs/system-jobs.js";
import { closeWorkers, startWorkers } from "./jobs/worker.js";
import { captureException, initAnalytics, shutdownAnalytics } from "./lib/analytics.js";
import { initAnalytics, shutdownAnalytics } from "./lib/analytics.js";
import { shouldRunStartupCleanup } from "./lib/cleanup.js";
import { buildCsp } from "./lib/csp.js";
import { reportError } from "./lib/error-report.js";
import { stripInternalPaths } from "./lib/errors.js";
import { ensureAiDirs, recoverInterruptedInstalls } from "./lib/feature-status.js";
import { logger } from "./lib/logger.js";
@@ -89,6 +90,16 @@ try {
console.error(err);
process.exit(1);
}
// BullMQ v5 requires Redis >= 6.2. Fail fast with an actionable message instead
// of crash-looping later on ReplyErrors from an incompatible server.
try {
await assertRedisCompatible();
} catch (err) {
const detected = err instanceof SafeError && err.code ? ` (detected ${err.code})` : "";
console.error(`FATAL: ${(err as Error).message}${detected}`);
process.exit(1);
}
console.log("Redis connected");
// Verify the local storage directories are writable before serving. A non-root
@@ -276,7 +287,12 @@ app.setErrorHandler((error: Error & { statusCode?: number }, request, reply) =>
{ err: error, url: request.url, method: request.method },
"Unhandled request error",
);
captureException(error);
void reportError(error, {
source: "http",
route: request.routeOptions?.url ?? undefined,
method: request.method,
statusCode,
});
} else {
request.log.warn({ err: error, url: request.url, method: request.method }, "Request error");
}
+26 -46
View File
@@ -1,19 +1,25 @@
import { existsSync } from "node:fs";
import { ANALYTICS_BAKED } from "@snapotter/shared";
import { analyticsEnabled, gatePrimed } from "./lib/analytics-gate.js";
import { analyticsEnabled, gatePrimed, telemetryEnvKilled } from "./lib/analytics-gate.js";
import { buildBeforeSend } from "./lib/sentry-scrub.js";
// Sentry inits at process load, before the gate cache is primed. Until the
// first successful read, stay silent rather than emit on the default-ON cache,
// so an opted-out instance never reports even a boot-window crash.
const sentryActive = () => gatePrimed() && analyticsEnabled();
// Collapse any absolute path in a stack frame filename to its basename, so
// even our own source paths never carry a workspace or job directory.
function basename(p: string): string {
const i = Math.max(p.lastIndexOf("/"), p.lastIndexOf("\\"));
return i >= 0 ? p.slice(i + 1) : p;
// All-in-one detection: docker/entrypoint.sh exports EMBEDDED_MODE=1 before
// exec'ing s6-overlay, and the snapotter service run script is with-contenv,
// so the marker reaches this process. URL absence is not a usable signal:
// embedded mode sets loopback DATABASE_URL/REDIS_URL before boot, and native
// dev commonly leaves DATABASE_URL unset (config.ts defaults it).
function deployMode(): string {
if (process.env.EMBEDDED_MODE) return "embedded";
if (existsSync("/.dockerenv")) return "external";
return "native";
}
if (ANALYTICS_BAKED.sentryDsn) {
if (ANALYTICS_BAKED.sentryDsn && !telemetryEnvKilled()) {
try {
const Sentry = await import("@sentry/node");
const { APP_VERSION } = await import("@snapotter/shared");
@@ -21,53 +27,27 @@ if (ANALYTICS_BAKED.sentryDsn) {
// attribute to a build; falls back to APP_VERSION for non-image runs.
const release = process.env.SENTRY_RELEASE || APP_VERSION;
// buildBeforeSend is typed on loose Record shapes so sentry-scrub.ts never
// imports @sentry/node; cast at this one boundary to the SDK callback type.
type SentryOptions = NonNullable<Parameters<typeof Sentry.init>[0]>;
Sentry.init({
dsn: ANALYTICS_BAKED.sentryDsn,
release,
environment: process.env.NODE_ENV || "production",
tracesSampleRate: ANALYTICS_BAKED.sampleRate,
environment: process.env.SNAPOTTER_ENV || "production",
sendDefaultPii: false,
// Release-health request-sessions and client-report envelopes are sent outside
// beforeSend/beforeSendTransaction, so the runtime opt-out below would not stop
// them. Disable both so an opted-out instance truly stops phoning home.
// Errors only. No traces options are set at all, so the SDK never
// starts traces and BullMQ/pg idle polling can't become transactions
// again (the July 2026 quota incident).
integrations: [Sentry.httpIntegration({ trackIncomingRequestsAsSessions: false })],
sendClientReports: false,
// Runtime opt-out: drop the whole transaction when analytics is off.
tracesSampler: () => (sentryActive() ? ANALYTICS_BAKED.sampleRate : 0),
beforeSend(event) {
if (!sentryActive()) return null; // kill switch (covers auto-captured errors)
// Allow-list: emit only error type + a basename-collapsed stack.
event.message = undefined;
event.logentry = undefined; // structured twin of message (captureMessage path)
event.server_name = undefined; // hostname is not anonymous
event.request = undefined;
event.extra = undefined;
event.contexts = undefined;
event.breadcrumbs = undefined;
event.user = undefined;
if (event.exception?.values) {
for (const ex of event.exception.values) {
ex.value = ex.type; // never the raw message body
if (ex.stacktrace?.frames) {
for (const frame of ex.stacktrace.frames) {
if (frame.filename) frame.filename = basename(frame.filename);
frame.abs_path = undefined;
frame.vars = undefined;
}
}
}
}
return event;
},
beforeBreadcrumb() {
return null; // breadcrumbs can carry URLs/messages with content; drop them
},
beforeSendTransaction(event) {
return sentryActive() ? event : null;
},
maxBreadcrumbs: 0,
beforeBreadcrumb: () => null,
initialScope: { tags: { deploy_mode: deployMode() } },
beforeSend: buildBeforeSend(sentryActive) as unknown as SentryOptions["beforeSend"],
});
console.log("[sentry] initialized, release:", release);
console.log("[sentry] initialized (errors only), release:", release);
} catch {
// @sentry/node not available
}
+18
View File
@@ -16,7 +16,25 @@ import { env } from "../config.js";
import { db, schema } from "../db/index.js";
import { getSettingString } from "../lib/settings-helpers.js";
/**
* In-memory license check, no DB access. This evaluator runs every 60s on
* every instance and previously queried settings even when unlicensed, where
* no supported path can create alert destinations. Gated for licensing parity
* with the SIEM job's NODE-1E fix (that job's unguarded settings read is what
* stormed Sentry, not this one).
*/
async function alertsLicensed(): Promise<boolean> {
try {
const { isFeatureEnabled } = await import("@snapotter/enterprise");
return isFeatureEnabled("admin_alerts");
} catch {
return false;
}
}
export async function evaluateAlerts(): Promise<void> {
if (!(await alertsLicensed())) return;
// 1. Read webhook destinations from settings
const destJson = await getSettingString("webhook_destinations", "[]");
let destinations: { url: string; authHeader: string; enabled: boolean; type: string }[];
+28
View File
@@ -4,6 +4,8 @@
* Uses ioredis with settings compatible with BullMQ's requirements
* (maxRetriesPerRequest: null for blocking commands).
*/
import { SafeError } from "@snapotter/shared";
import type { ConnectionOptions } from "bullmq";
import Redis from "ioredis";
import { env } from "../config.js";
@@ -64,3 +66,29 @@ export async function closeRedis(): Promise<void> {
_shared = null;
}
}
/** Pure check: returns a SafeError for known-incompatible versions, else null. */
export function checkRedisInfoCompatible(info: string): SafeError | null {
const m = info.match(/redis_version:(\d+)\.(\d+)/);
if (!m) return null; // managed Redis may hide INFO details; do not block boot
const major = Number(m[1]);
const minor = Number(m[2]);
if (major > 6 || (major === 6 && minor >= 2)) return null;
return new SafeError("Redis 6.2 or newer is required. Point REDIS_URL at Redis 8.", {
kind: "operational",
code: `redis-${major}.${minor}`,
});
}
/** Boot preflight: BullMQ v5 needs Redis >= 6.2. Fails fast with a clear message. */
export async function assertRedisCompatible(): Promise<void> {
let info: string;
try {
info = await sharedRedis().info("server");
} catch {
console.warn("[redis] INFO not permitted; skipping version preflight");
return;
}
const err = checkRedisInfoCompatible(info);
if (err) throw err;
}
+17
View File
@@ -35,7 +35,24 @@ async function readSettingValue(key: string): Promise<string | null> {
return row?.value ?? null;
}
/**
* In-memory license check, no DB access. This job fires every 30s on every
* instance; before this gate existed, the readSiemConfig settings read below
* ran on unlicensed instances too and turned every DB outage into a Sentry
* event storm (NODE-1E, July 2026).
*/
async function siemLicensed(): Promise<boolean> {
try {
const { isFeatureEnabled } = await import("@snapotter/enterprise");
return isFeatureEnabled("siem_forwarding");
} catch {
return false;
}
}
export async function runSiemForward(): Promise<{ forwarded: number } | undefined> {
if (!(await siemLicensed())) return;
// 1. Read SIEM config
const config = await readSiemConfig();
if (!config?.enabled || !config.webhookUrl) {
+18 -13
View File
@@ -25,14 +25,15 @@ import { mkdir, readFile, rm } from "node:fs/promises";
import { tmpdir } from "node:os";
import { join } from "node:path";
import { context, propagation, ROOT_CONTEXT, SpanStatusCode, trace } from "@opentelemetry/api";
import { ANALYTICS_EVENTS, getBundleForTool, TOOLS } from "@snapotter/shared";
import { ANALYTICS_EVENTS, getBundleForTool, isToolInputError, TOOLS } from "@snapotter/shared";
import { type Job, UnrecoverableError, Worker } from "bullmq";
import { eq } from "drizzle-orm";
import { env } from "../config.js";
import { db, schema } from "../db/index.js";
import { captureException, trackEvent } from "../lib/analytics.js";
import { trackEvent } from "../lib/analytics.js";
import { analyticsEnabled } from "../lib/analytics-gate.js";
import { resolveConcurrency } from "../lib/env.js";
import { reportError } from "../lib/error-report.js";
import { friendlyError } from "../lib/errors.js";
import { logger } from "../lib/logger.js";
import { jobDuration, jobsTotal } from "../lib/metrics.js";
@@ -375,7 +376,8 @@ async function processToolJob(job: Job<ToolJobData>): Promise<ToolJobResult> {
// friendlyError(finalError)). Expected validation rejections -- bad user
// input, not a server fault -- would otherwise flood error logs, so skip
// them here; they still reach the OTel span recorded below.
const isValidationError = err instanceof Error && err.name === "InputValidationError";
const isValidationError =
err instanceof Error && (err.name === "InputValidationError" || isToolInputError(err));
if (!isCanceled && !isTimeout && !isValidationError) {
logger.error({ err, jobId, toolId: data.toolId }, "tool job failed");
}
@@ -383,7 +385,7 @@ async function processToolJob(job: Job<ToolJobData>): Promise<ToolJobResult> {
// Record error on the OTel span
if (span) {
span.setStatus({ code: SpanStatusCode.ERROR, message: finalError });
span.recordException(err instanceof Error ? err : new Error(String(err)));
span.recordException(err instanceof Error ? err : String(err));
span.addEvent("job.failed");
}
@@ -447,9 +449,6 @@ async function processToolJob(job: Job<ToolJobData>): Promise<ToolJobResult> {
},
data.analyticsDistinctId,
);
if (!isCanceled && !isTimeout) {
void captureException(err instanceof Error ? err : new Error(String(err)));
}
}
if (isCanceled) throw new UnrecoverableError("Canceled");
@@ -888,9 +887,12 @@ export function startWorkers(): void {
});
worker.on("failed", (job, err) => {
if (analyticsEnabled() && job) {
void captureException(err instanceof Error ? err : new Error(String(err)));
}
if (!job) return;
void reportError(err, {
source: "worker",
pool,
toolId: (job.data as ToolJobData | undefined)?.toolId,
});
});
workers.push(worker);
@@ -916,9 +918,12 @@ export function startWorkers(): void {
});
worker.on("failed", (job, err) => {
if (analyticsEnabled() && job) {
void captureException(err instanceof Error ? err : new Error(String(err)));
}
if (!job) return;
void reportError(err, {
source: "worker",
pool,
toolId: (job.data as ToolJobData | undefined)?.toolId,
});
});
workers.push(worker);
+7
View File
@@ -25,8 +25,15 @@ async function defaultReader(): Promise<boolean | undefined> {
return rows[0].value !== "false";
}
/** Runtime kill switch honored in ALL builds: SNAPOTTER_TELEMETRY=0|false|off. */
export function telemetryEnvKilled(): boolean {
const v = process.env.SNAPOTTER_TELEMETRY;
return v === "0" || v === "false" || v === "off";
}
/** Compile-time bake, with a NON-PRODUCTION-only override so tests can force it on. */
export function bakedEnabled(): boolean {
if (telemetryEnvKilled()) return false;
if (process.env.NODE_ENV !== "production") {
const o = process.env.ANALYTICS_BAKED_OVERRIDE;
if (o === "on") return true;
+11 -9
View File
@@ -59,13 +59,10 @@ export async function initAnalytics(): Promise<void> {
}
export async function captureException(error: unknown): Promise<void> {
try {
if (!analyticsEnabled()) return;
const Sentry = await import("@sentry/node");
Sentry.captureException(error);
} catch {
// analytics must never throw
}
// Deprecated shim: route through the classified path. New code calls
// reportError directly with a source.
const { reportError } = await import("./error-report.js");
await reportError(error, { source: "boot" });
}
export async function shutdownAnalytics(): Promise<void> {
@@ -90,8 +87,13 @@ export async function trackEvent(
): Promise<void> {
try {
if (!analyticsEnabled() || !posthogClient) return;
if (ANALYTICS_BAKED.sampleRate < 1.0) {
if (ANALYTICS_BAKED.sampleRate <= 0.0 || Math.random() >= ANALYTICS_BAKED.sampleRate) return;
if (ANALYTICS_BAKED.posthogSampleRate < 1.0) {
if (
ANALYTICS_BAKED.posthogSampleRate <= 0.0 ||
Math.random() >= ANALYTICS_BAKED.posthogSampleRate
) {
return;
}
}
posthogClient.capture({
distinctId: distinctId ?? (await getInstanceId()),
+123
View File
@@ -0,0 +1,123 @@
/**
* The single deliberate Sentry capture path for the API.
*
* Classes:
* - expected: user input / client aborts / cancels. Never sent.
* - operational: someone's environment is broken (db down, disk full).
* Sent once per signature per hour, level=warning, fingerprinted per class.
* - bug: our fault. Sent up to 10 per signature per hour.
*
* State is per-process, so a crash-looping instance always reports its first
* event after each restart. The beforeSend ceiling (sentry-scrub.ts) is the
* final backstop and also covers SDK-captured uncaught exceptions.
*/
import {
connectivityClass,
isClientAbort,
isSafeMessageError,
isToolInputError,
} from "@snapotter/shared";
import { analyticsEnabled } from "./analytics-gate.js";
export type ErrorClass = "expected" | "operational" | "bug";
const HOUR_MS = 3600_000;
const LIMITS: Record<Exclude<ErrorClass, "expected">, number> = { operational: 1, bug: 10 };
const OPERATIONAL_CODES = new Set(["ENOSPC", "EACCES", "EROFS", "EMFILE", "ENFILE"]);
export interface ReportContext {
source: "http" | "worker" | "cron" | "boot";
toolId?: string;
pool?: string;
route?: string;
method?: string;
statusCode?: number;
subsystem?: string;
}
export function classifyError(err: unknown, source?: ReportContext["source"]): ErrorClass {
if (isToolInputError(err)) return "expected";
const e = err as { name?: string; message?: string; code?: string } | null;
if (e && typeof e.message === "string" && /^(Canceled$|Timed out after )/.test(e.message)) {
return "expected";
}
// The next two shortcuts only make sense at the HTTP boundary (undefined
// keeps the http-ish default for direct calls). Off the request path a bare
// ECONNRESET is an upstream socket loss, not a client abort, and a ZodError
// means schema drift: settings were already validated at the boundary, so a
// worker-side parse failure is our bug.
if (source === "http" || source === undefined) {
if (isClientAbort(err)) return "expected";
// ZodError = settings validation; InputValidationError = upload validation
// (apps/api/src/modality/contract.ts). Both are user-input problems.
if (e?.name === "ZodError" || e?.name === "InputValidationError") return "expected";
}
if (isSafeMessageError(err)) return err.kind === "bug" ? "bug" : "operational";
if (connectivityClass(err)) return "operational";
if (e?.code && OPERATIONAL_CODES.has(e.code)) return "operational";
return "bug";
}
const seen = new Map<string, { count: number; windowStart: number }>();
export function shouldReport(
cls: Exclude<ErrorClass, "expected">,
signature: string,
now = Date.now(),
): boolean {
const key = `${cls}:${signature}`;
const entry = seen.get(key);
if (!entry || now - entry.windowStart > HOUR_MS) {
seen.set(key, { count: 1, windowStart: now });
return true;
}
entry.count++;
return entry.count <= LIMITS[cls];
}
export function resetThrottleForTests(): void {
seen.clear();
}
export function errorSignature(err: unknown): string {
const e = err as { name?: string; code?: string; stack?: string } | null;
const name = e?.name ?? "Unknown";
const code = e?.code ?? "-";
let frame = "-";
if (typeof e?.stack === "string") {
const line = e.stack.split("\n").find((l) => l.includes("/apps/") || l.includes("/packages/"));
const m = line?.match(/([^/\\]+\.[cm]?[jt]sx?):(\d+)/);
if (m) frame = `${m[1]}:${m[2]}`;
}
return `${name}:${code}:${frame}`;
}
/** Fire-and-forget; never throws, never blocks. */
export async function reportError(err: unknown, ctx: ReportContext): Promise<void> {
try {
if (!analyticsEnabled()) return;
const cls = classifyError(err, ctx.source);
if (cls === "expected") return;
if (!shouldReport(cls, errorSignature(err))) return;
const Sentry = await import("@sentry/node");
const net = connectivityClass(err);
Sentry.withScope((scope) => {
scope.setLevel(cls === "operational" ? "warning" : "error");
scope.setTag("source", ctx.source);
scope.setTag("error_class", cls);
const code = (err as { code?: string } | null)?.code;
if (code) scope.setTag("error_code", code);
if (ctx.toolId) scope.setTag("tool_id", ctx.toolId);
if (ctx.pool) scope.setTag("pool", ctx.pool);
if (ctx.route) scope.setTag("route", ctx.route);
if (ctx.method) scope.setTag("method", ctx.method);
if (ctx.statusCode) scope.setTag("status_code", String(ctx.statusCode));
if (ctx.subsystem) scope.setTag("subsystem", ctx.subsystem);
if (net) scope.setFingerprint(["connectivity", net]);
Sentry.captureException(err instanceof Error ? err : new Error(String(err)));
});
} catch {
// telemetry must never throw
}
}
+18
View File
@@ -741,6 +741,24 @@ export async function importBundleArchive(
extractor.on("finish", () => res());
extractor.on("error", rej);
stream.on("error", rej);
}).catch((err: unknown) => {
// Malformed uploads (non-gzip data, corrupt/truncated gzip, garbage
// tar) surface as ZlibError or tar parse errors here. Map them to a
// 400-able validation error instead of letting them escape as a 500
// (Sentry NODE-1Z). Fatal node-tar parse errors always carry tarCode
// (TAR_ABORT, TAR_BAD_ARCHIVE); recoverable ones never reach "error".
if (err instanceof ImportValidationError) throw err;
const name = (err as Error | null)?.name ?? "";
const msg = String((err as Error | null)?.message ?? "");
const tarCode = (err as { tarCode?: unknown } | null)?.tarCode;
if (
name === "ZlibError" ||
typeof tarCode === "string" ||
/unexpected end of (file|data)|invalid tar|incorrect header check|zlib/i.test(msg)
) {
throw new ImportValidationError("Not a valid bundle archive");
}
throw err;
});
// Read and validate bundle.json
+22 -14
View File
@@ -4,6 +4,7 @@ import { mkdir, readFile, statfs, unlink, writeFile } from "node:fs/promises";
import { extname, join } from "node:path";
import type { Readable } from "node:stream";
import type { S3StorageModule } from "@snapotter/enterprise";
import { SafeError } from "@snapotter/shared";
import { env } from "../config.js";
const MIN_FREE_BYTES = 100 * 1024 * 1024;
@@ -13,9 +14,12 @@ async function assertDiskSpace(dir: string): Promise<void> {
const stats = await statfs(dir);
const freeBytes = stats.bfree * stats.bsize;
if (freeBytes < MIN_FREE_BYTES) {
const err = new Error("Insufficient disk space") as Error & { statusCode: number };
err.statusCode = 507;
throw err;
// ENOSPC here is synthesized from the free-space floor check, not a syscall errno.
throw new SafeError("Insufficient disk space", {
kind: "operational",
code: "ENOSPC",
statusCode: 507,
});
}
} catch (e) {
if (e instanceof Error && (e as Error & { statusCode?: number }).statusCode === 507) throw e;
@@ -102,11 +106,11 @@ export async function ensureStorageDir(): Promise<void> {
await mkdir(env.FILES_STORAGE_PATH, { recursive: true });
} catch (e) {
if (e instanceof Error && (e as NodeJS.ErrnoException).code === "EACCES") {
const err = new Error("Storage directory is not writable") as Error & {
statusCode: number;
};
err.statusCode = 503;
throw err;
throw new SafeError("Storage directory is not writable", {
kind: "operational",
code: (e as NodeJS.ErrnoException).code,
statusCode: 503,
});
}
throw e;
}
@@ -126,11 +130,11 @@ export async function saveFile(buffer: Buffer, originalName: string): Promise<st
await writeFile(join(env.FILES_STORAGE_PATH, storedName), buffer);
} catch (e) {
if (e instanceof Error && (e as NodeJS.ErrnoException).code === "EACCES") {
const err = new Error("Storage directory is not writable") as Error & {
statusCode: number;
};
err.statusCode = 503;
throw err;
throw new SafeError("Storage directory is not writable", {
kind: "operational",
code: (e as NodeJS.ErrnoException).code,
statusCode: 503,
});
}
throw e;
}
@@ -185,7 +189,11 @@ async function ensureThumbDir(): Promise<void> {
await mkdir(join(env.FILES_STORAGE_PATH, THUMB_DIR), { recursive: true });
} catch (err: unknown) {
if ((err as NodeJS.ErrnoException).code === "EACCES") {
throw Object.assign(new Error("Storage directory is not writable"), { statusCode: 503 });
throw new SafeError("Storage directory is not writable", {
kind: "operational",
code: (err as NodeJS.ErrnoException).code,
statusCode: 503,
});
}
throw err;
}
+100
View File
@@ -0,0 +1,100 @@
/**
* Sentry beforeSend for the API: allowlist-first scrubbing plus a per-process
* event ceiling. Kept pure (factory + injected gate) so it is unit-testable
* without initializing the SDK. See the telemetry overhaul spec for the rules.
*/
import { rebuildErrorValue } from "@snapotter/shared";
const CEILING_PER_HOUR = 20;
const HOUR_MS = 3600_000;
const TAG_ALLOWLIST = new Set([
"source",
"tool_id",
"pool",
"route",
"method",
"error_class",
"error_code",
"deploy_mode",
"subsystem",
"status_code",
]);
function basename(p: string): string {
const i = Math.max(p.lastIndexOf("/"), p.lastIndexOf("\\"));
return i >= 0 ? p.slice(i + 1) : p;
}
// Sentry event/hint are typed loosely on purpose: this module must not import
// @sentry/node (instrument.ts loads the SDK lazily and passes events through).
type AnyEvent = Record<string, unknown>;
type AnyHint = { originalException?: unknown };
/** Narrow to a plain mutable object, or null for anything else (fail-closed). */
function asObj(value: unknown): AnyEvent | null {
return value !== null && typeof value === "object" && !Array.isArray(value)
? (value as AnyEvent)
: null;
}
export function buildBeforeSend(isActive: () => boolean) {
let windowStart = 0;
let sentInWindow = 0;
return function beforeSend(event: AnyEvent, hint: AnyHint): AnyEvent | null {
if (!isActive()) return null;
const now = Date.now();
if (now - windowStart > HOUR_MS) {
windowStart = now;
sentInWindow = 0;
}
if (++sentInWindow > CEILING_PER_HOUR) return null;
event.message = undefined;
event.logentry = undefined;
event.server_name = undefined;
event.request = undefined;
event.extra = undefined;
event.breadcrumbs = undefined;
event.user = undefined;
const ctx = asObj(event.contexts);
const keep: AnyEvent = {};
const os = asObj(ctx?.os);
if (os?.name) keep.os = { name: os.name, version: os.version };
const runtime = asObj(ctx?.runtime);
if (runtime?.name) keep.runtime = { name: runtime.name, version: runtime.version };
event.contexts = Object.keys(keep).length ? keep : undefined;
const tags = asObj(event.tags);
if (tags) {
for (const key of Object.keys(tags)) {
if (!TAG_ALLOWLIST.has(key)) delete tags[key];
}
}
const rebuilt = rebuildErrorValue(hint?.originalException);
const values = asObj(event.exception)?.values;
if (Array.isArray(values)) {
for (let i = 0; i < values.length; i++) {
const ex = asObj(values[i]);
if (!ex) continue;
// The last entry is the original error; linked/outer wrappers get type-only.
ex.value = i === values.length - 1 && rebuilt ? rebuilt : ex.type;
const frames = asObj(ex.stacktrace)?.frames;
if (Array.isArray(frames)) {
for (const entry of frames) {
const frame = asObj(entry);
if (!frame) continue;
if (typeof frame.filename === "string") frame.filename = basename(frame.filename);
frame.abs_path = undefined;
frame.vars = undefined;
}
}
}
}
return event;
};
}
+35 -7
View File
@@ -172,6 +172,29 @@ function normalizeSafeFetchOptions(options?: AbortSignal | SafeFetchOptions): Sa
return options;
}
// Statuses that undici's Response constructor rejects a body for. An empty
// Buffer still counts as a body, so these must pass null explicitly. A remote
// server controls this status; before this guard, a 204/304 reply crashed the
// whole process from inside the 'end' event handler (Sentry NODE-20).
const NULL_BODY_STATUSES = new Set([101, 204, 205, 304]);
export function toFetchResponse(
// ArrayBuffer-backed (what Buffer.concat/from/alloc return); the Response
// constructor's BodyInit does not accept SharedArrayBuffer-backed views.
body: Buffer<ArrayBuffer>,
statusCode: number | undefined,
statusText: string,
headers: Headers,
): Response {
const status =
statusCode !== undefined && statusCode >= 200 && statusCode <= 599 ? statusCode : 502;
return new Response(NULL_BODY_STATUSES.has(status) ? null : body, {
status,
statusText,
headers,
});
}
function withResponseSizeLimit(response: Response, maxBytes?: number): Response {
if (maxBytes === undefined || !response.body) return response;
@@ -275,13 +298,18 @@ export async function safeFetch(
for (const v of vals) headers.append(key, v);
}
}
resolve(
new Response(body, {
status: incomingMessage.statusCode ?? 500,
statusText: incomingMessage.statusMessage ?? "",
headers,
}),
);
try {
resolve(
toFetchResponse(
body,
incomingMessage.statusCode,
incomingMessage.statusMessage ?? "",
headers,
),
);
} catch (err) {
reject(err);
}
});
incomingMessage.on("error", (err) => {
if (settled) return;
+3 -1
View File
@@ -6414,7 +6414,9 @@ paths:
type: string
sentryDsn:
type: string
sampleRate:
sentryDsnWeb:
type: string
posthogSampleRate:
type: number
instanceId:
type: string
+11
View File
@@ -1100,8 +1100,19 @@ function isPublicRoute(url: string): boolean {
return PUBLIC_PATHS.some((path) => url.startsWith(path));
}
/** SPA bundle assets are public by definition (the login page needs them). */
export function isStaticAssetRequest(method: string, url: string): boolean {
return (
(method === "GET" || method === "HEAD") && url.startsWith("/assets/") && !url.includes("..")
);
}
export async function authMiddleware(app: FastifyInstance): Promise<void> {
app.addHook("preHandler", async (request: FastifyRequest, reply: FastifyReply) => {
// Skip the session DB lookup for bundle assets: they are served on every
// page load and an unreachable DB must not 500 them (Sentry NODE-1D).
if (isStaticAssetRequest(request.method, request.url)) return;
if (!env.AUTH_ENABLED) {
(request as FastifyRequest & { user?: AuthUser }).user = {
id: "anonymous",
+4 -2
View File
@@ -13,7 +13,8 @@ export async function analyticsRoutes(app: FastifyInstance): Promise<void> {
posthogApiKey: "",
posthogHost: "",
sentryDsn: "",
sampleRate: 0,
sentryDsnWeb: "",
posthogSampleRate: 0,
instanceId: "",
};
}
@@ -28,7 +29,8 @@ export async function analyticsRoutes(app: FastifyInstance): Promise<void> {
posthogApiKey: ANALYTICS_BAKED.posthogApiKey,
posthogHost: ANALYTICS_BAKED.posthogHost,
sentryDsn: ANALYTICS_BAKED.sentryDsn,
sampleRate: ANALYTICS_BAKED.sampleRate,
sentryDsnWeb: ANALYTICS_BAKED.sentryDsnWeb,
posthogSampleRate: ANALYTICS_BAKED.posthogSampleRate,
instanceId: row?.value ?? "",
};
});
+3 -1
View File
@@ -238,7 +238,9 @@ function ensureSubscriber(): void {
sseSubscriber.on("error", (err) => {
console.error("SSE progress subscriber error", err);
});
void sseSubscriber.subscribe(progressChannel());
void sseSubscriber.subscribe(progressChannel()).catch((err) => {
console.error("SSE progress subscribe failed", err);
});
sseSubscriber.on("message", (_channel: string, message: string) => {
try {
const parsed = JSON.parse(message) as { jobId?: string };
+6 -5
View File
@@ -1,3 +1,4 @@
import { ToolInputError } from "@snapotter/shared";
import type { FastifyInstance } from "fastify";
import Papa from "papaparse";
import sharp from "sharp";
@@ -207,26 +208,26 @@ export function registerChartMaker(app: FastifyInstance) {
try {
data = parseInput(input.buffer, input.filename);
} catch (err) {
throw new Error(err instanceof Error ? err.message : "Failed to parse input");
throw new ToolInputError(err instanceof Error ? err.message : "Failed to parse input");
}
if (data.length === 0) {
throw new Error("No data points found in input");
throw new ToolInputError("No data points found in input");
}
if (data.length > 100) {
throw new Error("Too many data points (max 100)");
throw new ToolInputError("Too many data points (max 100)");
}
// Validate numeric values
for (const point of data) {
if (Number.isNaN(point.value)) {
throw new Error("Column 2 must be numeric");
throw new ToolInputError("Column 2 must be numeric");
}
}
// Negative values render as invalid/degenerate SVG (negative bar heights,
// backward pie arcs that Sharp silently drops); reject with a clear message.
if (data.some((point) => point.value < 0)) {
throw new Error("Chart values must be zero or greater");
throw new ToolInputError("Chart values must be zero or greater");
}
let svg: string;
@@ -35,6 +35,7 @@ function presetSchema(base: string) {
* own parameterized registrars.
*/
export function registerConversionPresets(app: FastifyInstance): number {
let skipped = 0;
for (const preset of CONVERSION_PRESETS) {
const cfg = BASE_CONFIG[preset.base];
if (cfg.group === "image-to-pdf") {
@@ -53,9 +54,16 @@ export function registerConversionPresets(app: FastifyInstance): number {
// group "registry": delegate to the base tool's processV2 with locked settings merged in.
const baseConfig = getToolConfig(preset.base);
if (!baseConfig) {
throw new Error(
`Preset "${preset.id}" base "${preset.base}" is not registered in the tool registry`,
// A missing base must degrade to a disabled preset, not a boot crash:
// one bad environment crash-looped an instance 287 times (Sentry
// NODE-21). CI still fails hard via the tool-route drift test, so the
// official image can't ship with a hole.
app.log.error(
{ presetId: preset.id, base: preset.base },
"conversion preset base tool missing; preset disabled",
);
skipped++;
continue;
}
createToolRoute(app, {
toolId: preset.id,
@@ -73,5 +81,5 @@ export function registerConversionPresets(app: FastifyInstance): number {
});
}
return CONVERSION_PRESETS.length;
return CONVERSION_PRESETS.length - skipped;
}