fix(webapp,redis-worker): stop logging raw metadata, alert payloads, and job items (#4403)

This commit is contained in:
Chris Arderne
2026-07-29 17:59:36 +01:00
committed by GitHub
parent a09817169f
commit 8ebc8a41af
5 changed files with 74 additions and 27 deletions
+1 -1
View File
@@ -50,7 +50,7 @@ function flattenArgs(args: Array<Record<string, unknown> | undefined>) {
export const logger = new Logger(
"webapp",
(process.env.APP_LOG_LEVEL ?? "info") as LogLevel,
["examples", "output", "connectionString", "payload"],
["examples", "output", "connectionString", "payload", "metadata", "seedMetadata"],
sensitiveDataReplacer,
() => {
const fields = currentFieldsStore.getStore();
@@ -6,7 +6,11 @@ import type {
import { applyMetadataOperations, parsePacket } from "@trigger.dev/core/v3";
import type { PrismaClientOrTransaction } from "~/db.server";
import type { AuthenticatedEnvironment } from "~/services/apiAuth.server";
import { handleMetadataPacket, MetadataTooLargeError } from "~/utils/packets";
import {
handleMetadataPacket,
handleMetadataPacketWithByteLength,
MetadataTooLargeError,
} from "~/utils/packets";
import { ServiceValidationError } from "~/v3/services/common.server";
import { Effect, Schedule, Duration, Fiber } from "effect";
import { type RuntimeFiber } from "effect/Fiber";
@@ -91,9 +95,14 @@ export class UpdateMetadataService {
this._bufferedOperations.clear();
yield* Effect.sync(() => {
if (this.flushLoggingEnabled) {
if (this.flushLoggingEnabled && currentOperations.size > 0) {
const operationCount = Array.from(currentOperations.values()).reduce(
(sum, ops) => sum + ops.length,
0
);
this.logger.debug(`[UpdateMetadataService] Flushing operations`, {
operations: Object.fromEntries(currentOperations),
runCount: currentOperations.size,
operationCount,
});
}
});
@@ -520,9 +529,9 @@ export class UpdateMetadataService {
if (this.flushLoggingEnabled) {
this.logger.debug(`[updateRunMetadataWithOperations] Updated metadata for run`, {
metadata: applyResults.newMetadata,
operations: operations,
runId,
metadataKeyCount: Object.keys(applyResults.newMetadata).length,
operationCount: operations.length,
});
}
@@ -549,16 +558,18 @@ export class UpdateMetadataService {
body: UpdateMetadataRequestBody,
existingMetadata: IOPacket
): Promise<{ metadata: Record<string, unknown> | undefined; updatedAtMs?: number }> {
const metadataPacket = handleMetadataPacket(
const metadataPacketWithByteLength = handleMetadataPacketWithByteLength(
body.metadata,
"application/json",
this.maximumSize
);
if (!metadataPacket) {
if (!metadataPacketWithByteLength) {
return { metadata: {} };
}
const { packet: metadataPacket, byteLength: metadataSizeBytes } = metadataPacketWithByteLength;
let updatedAtMs: number | undefined;
if (
@@ -567,8 +578,8 @@ export class UpdateMetadataService {
) {
if (this.flushLoggingEnabled) {
this.logger.debug(`[updateRunMetadataDirectly] Updating metadata directly for run`, {
metadata: metadataPacket.data,
runId,
metadataSizeBytes,
});
}
@@ -607,7 +618,7 @@ export class UpdateMetadataService {
if (this.flushLoggingEnabled) {
this.logger.debug(`[ingestRunOperations] Ingesting operations for run`, {
runId,
bufferedOperations,
operationCount: bufferedOperations.length,
});
}
+9 -1
View File
@@ -13,6 +13,14 @@ export function handleMetadataPacket(
metadataType: string,
maximumSize: number
): IOPacket | undefined {
return handleMetadataPacketWithByteLength(metadata, metadataType, maximumSize)?.packet;
}
export function handleMetadataPacketWithByteLength(
metadata: any,
metadataType: string,
maximumSize: number
): { packet: IOPacket; byteLength: number } | undefined {
let metadataPacket: IOPacket | undefined = undefined;
if (typeof metadata === "string") {
@@ -33,5 +41,5 @@ export function handleMetadataPacket(
throw new MetadataTooLargeError(`Metadata exceeds maximum size of ${maximumSize} bytes`);
}
return metadataPacket;
return { packet: metadataPacket, byteLength };
}
@@ -455,7 +455,10 @@ export class DeliverAlertService extends BaseService {
error,
};
await this.#deliverWebhook(payload, webhookProperties.data);
await this.#deliverWebhook(payload, webhookProperties.data, {
webhookId: alert.channel.id,
runId: alert.taskRun.friendlyId,
});
break;
}
case "v2": {
@@ -516,7 +519,10 @@ export class DeliverAlertService extends BaseService {
},
};
await this.#deliverWebhook(payload, webhookProperties.data);
await this.#deliverWebhook(payload, webhookProperties.data, {
webhookId: alert.channel.id,
runId: alert.taskRun.friendlyId,
});
break;
}
@@ -577,7 +583,9 @@ export class DeliverAlertService extends BaseService {
vercel: this.#buildWebhookVercelObject(deploymentMeta.vercelDeploymentUrl),
};
await this.#deliverWebhook(payload, webhookProperties.data);
await this.#deliverWebhook(payload, webhookProperties.data, {
webhookId: alert.channel.id,
});
break;
}
case "v2": {
@@ -616,7 +624,9 @@ export class DeliverAlertService extends BaseService {
},
};
await this.#deliverWebhook(payload, webhookProperties.data);
await this.#deliverWebhook(payload, webhookProperties.data, {
webhookId: alert.channel.id,
});
break;
}
@@ -671,7 +681,9 @@ export class DeliverAlertService extends BaseService {
vercel: this.#buildWebhookVercelObject(deploymentMeta.vercelDeploymentUrl),
};
await this.#deliverWebhook(payload, webhookProperties.data);
await this.#deliverWebhook(payload, webhookProperties.data, {
webhookId: alert.channel.id,
});
break;
}
case "v2": {
@@ -716,7 +728,9 @@ export class DeliverAlertService extends BaseService {
},
};
await this.#deliverWebhook(payload, webhookProperties.data);
await this.#deliverWebhook(payload, webhookProperties.data, {
webhookId: alert.channel.id,
});
break;
}
@@ -1017,7 +1031,11 @@ export class DeliverAlertService extends BaseService {
}
}
async #deliverWebhook<T>(payload: T, webhook: ProjectAlertWebhookProperties) {
async #deliverWebhook<T>(
payload: T,
webhook: ProjectAlertWebhookProperties,
context: { webhookId: string; runId?: string }
) {
const rawPayload = JSON.stringify(payload);
const hashPayload = Buffer.from(rawPayload, "utf-8");
@@ -1046,15 +1064,17 @@ export class DeliverAlertService extends BaseService {
});
if (!response.ok) {
// Never log the request/response body here: it is customer-controlled alert
// content and may include stack traces or other application data.
logger.info("[DeliverAlert] Failed to send alert webhook", {
status: response.status,
statusText: response.statusText,
url: webhook.url,
body: payload,
signature,
urlHost: safeUrlHost(webhook.url),
webhookId: context.webhookId,
runId: context.runId,
});
throw new Error(`Failed to send alert webhook to ${webhook.url}`);
throw new Error(`Failed to send alert webhook to ${safeUrlHost(webhook.url)}`);
}
}
@@ -1435,3 +1455,11 @@ function isWebAPIHTTPError(error: unknown): error is WebAPIHTTPError {
function isWebAPIRateLimitedError(error: unknown): error is WebAPIRateLimitedError {
return (error as WebAPIRateLimitedError).code === ErrorCode.RateLimitedError;
}
function safeUrlHost(url: string): string {
try {
return new URL(url).host;
} catch {
return "unknown";
}
}
+5 -5
View File
@@ -140,7 +140,7 @@ class Worker<TCatalog extends WorkerCatalog> {
> = new Map();
constructor(private options: WorkerOptions<TCatalog>) {
this.logger = options.logger ?? new Logger("Worker", "debug");
this.logger = options.logger ?? new Logger("Worker", "debug", ["item"]);
this.tracer = options.tracer ?? trace.getTracer(options.name);
this.meter = options.meter ?? metrics.getMeter(options.name);
@@ -608,7 +608,8 @@ class Worker<TCatalog extends WorkerCatalog> {
this.logger.error("Unhandled error in processItem:", {
error: err,
workerId,
item,
id: queueItem.id,
job: queueItem.job,
});
}
);
@@ -933,11 +934,12 @@ class Worker<TCatalog extends WorkerCatalog> {
const errorLogLevel =
error && typeof error === "object" && "logLevel" in error ? error.logLevel : undefined;
// Never include the raw item/payload here: it is job data that may be
// customer-controlled. It is retrievable via `getJob(id)` if needed for triage.
const logAttributes = {
name: this.options.name,
id,
job,
item,
visibilityTimeoutMs,
error,
errorMessage,
@@ -994,7 +996,6 @@ class Worker<TCatalog extends WorkerCatalog> {
name: this.options.name,
id,
job,
item,
retryDate,
retryDelay,
visibilityTimeoutMs,
@@ -1015,7 +1016,6 @@ class Worker<TCatalog extends WorkerCatalog> {
name: this.options.name,
id,
job,
item,
visibilityTimeoutMs,
error: requeueError,
}