Skip to content

Commit 8ebc8a4

Browse files
authored
fix(webapp,redis-worker): stop logging raw metadata, alert payloads, and job items (#4403)
1 parent a098171 commit 8ebc8a4

5 files changed

Lines changed: 74 additions & 27 deletions

File tree

apps/webapp/app/services/logger.server.ts

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -50,7 +50,7 @@ function flattenArgs(args: Array<Record<string, unknown> | undefined>) {
5050
export const logger = new Logger(
5151
"webapp",
5252
(process.env.APP_LOG_LEVEL ?? "info") as LogLevel,
53-
["examples", "output", "connectionString", "payload"],
53+
["examples", "output", "connectionString", "payload", "metadata", "seedMetadata"],
5454
sensitiveDataReplacer,
5555
() => {
5656
const fields = currentFieldsStore.getStore();

apps/webapp/app/services/metadata/updateMetadata.server.ts

Lines changed: 20 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,11 @@ import type {
66
import { applyMetadataOperations, parsePacket } from "@trigger.dev/core/v3";
77
import type { PrismaClientOrTransaction } from "~/db.server";
88
import type { AuthenticatedEnvironment } from "~/services/apiAuth.server";
9-
import { handleMetadataPacket, MetadataTooLargeError } from "~/utils/packets";
9+
import {
10+
handleMetadataPacket,
11+
handleMetadataPacketWithByteLength,
12+
MetadataTooLargeError,
13+
} from "~/utils/packets";
1014
import { ServiceValidationError } from "~/v3/services/common.server";
1115
import { Effect, Schedule, Duration, Fiber } from "effect";
1216
import { type RuntimeFiber } from "effect/Fiber";
@@ -91,9 +95,14 @@ export class UpdateMetadataService {
9195
this._bufferedOperations.clear();
9296

9397
yield* Effect.sync(() => {
94-
if (this.flushLoggingEnabled) {
98+
if (this.flushLoggingEnabled && currentOperations.size > 0) {
99+
const operationCount = Array.from(currentOperations.values()).reduce(
100+
(sum, ops) => sum + ops.length,
101+
0
102+
);
95103
this.logger.debug(`[UpdateMetadataService] Flushing operations`, {
96-
operations: Object.fromEntries(currentOperations),
104+
runCount: currentOperations.size,
105+
operationCount,
97106
});
98107
}
99108
});
@@ -520,9 +529,9 @@ export class UpdateMetadataService {
520529

521530
if (this.flushLoggingEnabled) {
522531
this.logger.debug(`[updateRunMetadataWithOperations] Updated metadata for run`, {
523-
metadata: applyResults.newMetadata,
524-
operations: operations,
525532
runId,
533+
metadataKeyCount: Object.keys(applyResults.newMetadata).length,
534+
operationCount: operations.length,
526535
});
527536
}
528537

@@ -549,16 +558,18 @@ export class UpdateMetadataService {
549558
body: UpdateMetadataRequestBody,
550559
existingMetadata: IOPacket
551560
): Promise<{ metadata: Record<string, unknown> | undefined; updatedAtMs?: number }> {
552-
const metadataPacket = handleMetadataPacket(
561+
const metadataPacketWithByteLength = handleMetadataPacketWithByteLength(
553562
body.metadata,
554563
"application/json",
555564
this.maximumSize
556565
);
557566

558-
if (!metadataPacket) {
567+
if (!metadataPacketWithByteLength) {
559568
return { metadata: {} };
560569
}
561570

571+
const { packet: metadataPacket, byteLength: metadataSizeBytes } = metadataPacketWithByteLength;
572+
562573
let updatedAtMs: number | undefined;
563574

564575
if (
@@ -567,8 +578,8 @@ export class UpdateMetadataService {
567578
) {
568579
if (this.flushLoggingEnabled) {
569580
this.logger.debug(`[updateRunMetadataDirectly] Updating metadata directly for run`, {
570-
metadata: metadataPacket.data,
571581
runId,
582+
metadataSizeBytes,
572583
});
573584
}
574585

@@ -607,7 +618,7 @@ export class UpdateMetadataService {
607618
if (this.flushLoggingEnabled) {
608619
this.logger.debug(`[ingestRunOperations] Ingesting operations for run`, {
609620
runId,
610-
bufferedOperations,
621+
operationCount: bufferedOperations.length,
611622
});
612623
}
613624

apps/webapp/app/utils/packets.ts

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,14 @@ export function handleMetadataPacket(
1313
metadataType: string,
1414
maximumSize: number
1515
): IOPacket | undefined {
16+
return handleMetadataPacketWithByteLength(metadata, metadataType, maximumSize)?.packet;
17+
}
18+
19+
export function handleMetadataPacketWithByteLength(
20+
metadata: any,
21+
metadataType: string,
22+
maximumSize: number
23+
): { packet: IOPacket; byteLength: number } | undefined {
1624
let metadataPacket: IOPacket | undefined = undefined;
1725

1826
if (typeof metadata === "string") {
@@ -33,5 +41,5 @@ export function handleMetadataPacket(
3341
throw new MetadataTooLargeError(`Metadata exceeds maximum size of ${maximumSize} bytes`);
3442
}
3543

36-
return metadataPacket;
44+
return { packet: metadataPacket, byteLength };
3745
}

apps/webapp/app/v3/services/alerts/deliverAlert.server.ts

Lines changed: 39 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -455,7 +455,10 @@ export class DeliverAlertService extends BaseService {
455455
error,
456456
};
457457

458-
await this.#deliverWebhook(payload, webhookProperties.data);
458+
await this.#deliverWebhook(payload, webhookProperties.data, {
459+
webhookId: alert.channel.id,
460+
runId: alert.taskRun.friendlyId,
461+
});
459462
break;
460463
}
461464
case "v2": {
@@ -516,7 +519,10 @@ export class DeliverAlertService extends BaseService {
516519
},
517520
};
518521

519-
await this.#deliverWebhook(payload, webhookProperties.data);
522+
await this.#deliverWebhook(payload, webhookProperties.data, {
523+
webhookId: alert.channel.id,
524+
runId: alert.taskRun.friendlyId,
525+
});
520526

521527
break;
522528
}
@@ -577,7 +583,9 @@ export class DeliverAlertService extends BaseService {
577583
vercel: this.#buildWebhookVercelObject(deploymentMeta.vercelDeploymentUrl),
578584
};
579585

580-
await this.#deliverWebhook(payload, webhookProperties.data);
586+
await this.#deliverWebhook(payload, webhookProperties.data, {
587+
webhookId: alert.channel.id,
588+
});
581589
break;
582590
}
583591
case "v2": {
@@ -616,7 +624,9 @@ export class DeliverAlertService extends BaseService {
616624
},
617625
};
618626

619-
await this.#deliverWebhook(payload, webhookProperties.data);
627+
await this.#deliverWebhook(payload, webhookProperties.data, {
628+
webhookId: alert.channel.id,
629+
});
620630

621631
break;
622632
}
@@ -671,7 +681,9 @@ export class DeliverAlertService extends BaseService {
671681
vercel: this.#buildWebhookVercelObject(deploymentMeta.vercelDeploymentUrl),
672682
};
673683

674-
await this.#deliverWebhook(payload, webhookProperties.data);
684+
await this.#deliverWebhook(payload, webhookProperties.data, {
685+
webhookId: alert.channel.id,
686+
});
675687
break;
676688
}
677689
case "v2": {
@@ -716,7 +728,9 @@ export class DeliverAlertService extends BaseService {
716728
},
717729
};
718730

719-
await this.#deliverWebhook(payload, webhookProperties.data);
731+
await this.#deliverWebhook(payload, webhookProperties.data, {
732+
webhookId: alert.channel.id,
733+
});
720734

721735
break;
722736
}
@@ -1017,7 +1031,11 @@ export class DeliverAlertService extends BaseService {
10171031
}
10181032
}
10191033

1020-
async #deliverWebhook<T>(payload: T, webhook: ProjectAlertWebhookProperties) {
1034+
async #deliverWebhook<T>(
1035+
payload: T,
1036+
webhook: ProjectAlertWebhookProperties,
1037+
context: { webhookId: string; runId?: string }
1038+
) {
10211039
const rawPayload = JSON.stringify(payload);
10221040
const hashPayload = Buffer.from(rawPayload, "utf-8");
10231041

@@ -1046,15 +1064,17 @@ export class DeliverAlertService extends BaseService {
10461064
});
10471065

10481066
if (!response.ok) {
1067+
// Never log the request/response body here: it is customer-controlled alert
1068+
// content and may include stack traces or other application data.
10491069
logger.info("[DeliverAlert] Failed to send alert webhook", {
10501070
status: response.status,
10511071
statusText: response.statusText,
1052-
url: webhook.url,
1053-
body: payload,
1054-
signature,
1072+
urlHost: safeUrlHost(webhook.url),
1073+
webhookId: context.webhookId,
1074+
runId: context.runId,
10551075
});
10561076

1057-
throw new Error(`Failed to send alert webhook to ${webhook.url}`);
1077+
throw new Error(`Failed to send alert webhook to ${safeUrlHost(webhook.url)}`);
10581078
}
10591079
}
10601080

@@ -1435,3 +1455,11 @@ function isWebAPIHTTPError(error: unknown): error is WebAPIHTTPError {
14351455
function isWebAPIRateLimitedError(error: unknown): error is WebAPIRateLimitedError {
14361456
return (error as WebAPIRateLimitedError).code === ErrorCode.RateLimitedError;
14371457
}
1458+
1459+
function safeUrlHost(url: string): string {
1460+
try {
1461+
return new URL(url).host;
1462+
} catch {
1463+
return "unknown";
1464+
}
1465+
}

packages/redis-worker/src/worker.ts

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -140,7 +140,7 @@ class Worker<TCatalog extends WorkerCatalog> {
140140
> = new Map();
141141

142142
constructor(private options: WorkerOptions<TCatalog>) {
143-
this.logger = options.logger ?? new Logger("Worker", "debug");
143+
this.logger = options.logger ?? new Logger("Worker", "debug", ["item"]);
144144
this.tracer = options.tracer ?? trace.getTracer(options.name);
145145
this.meter = options.meter ?? metrics.getMeter(options.name);
146146

@@ -608,7 +608,8 @@ class Worker<TCatalog extends WorkerCatalog> {
608608
this.logger.error("Unhandled error in processItem:", {
609609
error: err,
610610
workerId,
611-
item,
611+
id: queueItem.id,
612+
job: queueItem.job,
612613
});
613614
}
614615
);
@@ -933,11 +934,12 @@ class Worker<TCatalog extends WorkerCatalog> {
933934
const errorLogLevel =
934935
error && typeof error === "object" && "logLevel" in error ? error.logLevel : undefined;
935936

937+
// Never include the raw item/payload here: it is job data that may be
938+
// customer-controlled. It is retrievable via `getJob(id)` if needed for triage.
936939
const logAttributes = {
937940
name: this.options.name,
938941
id,
939942
job,
940-
item,
941943
visibilityTimeoutMs,
942944
error,
943945
errorMessage,
@@ -994,7 +996,6 @@ class Worker<TCatalog extends WorkerCatalog> {
994996
name: this.options.name,
995997
id,
996998
job,
997-
item,
998999
retryDate,
9991000
retryDelay,
10001001
visibilityTimeoutMs,
@@ -1015,7 +1016,6 @@ class Worker<TCatalog extends WorkerCatalog> {
10151016
name: this.options.name,
10161017
id,
10171018
job,
1018-
item,
10191019
visibilityTimeoutMs,
10201020
error: requeueError,
10211021
}

0 commit comments

Comments
 (0)