diff --git a/apps/webapp/app/services/logger.server.ts b/apps/webapp/app/services/logger.server.ts index 15b248f0d94..b03802afefd 100644 --- a/apps/webapp/app/services/logger.server.ts +++ b/apps/webapp/app/services/logger.server.ts @@ -50,7 +50,7 @@ function flattenArgs(args: Array | 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(); diff --git a/apps/webapp/app/services/metadata/updateMetadata.server.ts b/apps/webapp/app/services/metadata/updateMetadata.server.ts index c44dacf65c5..07dc236c976 100644 --- a/apps/webapp/app/services/metadata/updateMetadata.server.ts +++ b/apps/webapp/app/services/metadata/updateMetadata.server.ts @@ -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 | 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, }); } diff --git a/apps/webapp/app/utils/packets.ts b/apps/webapp/app/utils/packets.ts index 7a522d6f7af..27217c7f2ae 100644 --- a/apps/webapp/app/utils/packets.ts +++ b/apps/webapp/app/utils/packets.ts @@ -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 }; } diff --git a/apps/webapp/app/v3/services/alerts/deliverAlert.server.ts b/apps/webapp/app/v3/services/alerts/deliverAlert.server.ts index 0905f7c768f..070be65f7ad 100644 --- a/apps/webapp/app/v3/services/alerts/deliverAlert.server.ts +++ b/apps/webapp/app/v3/services/alerts/deliverAlert.server.ts @@ -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(payload: T, webhook: ProjectAlertWebhookProperties) { + async #deliverWebhook( + 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"; + } +} diff --git a/packages/redis-worker/src/worker.ts b/packages/redis-worker/src/worker.ts index 58988b81896..b9240acaa63 100644 --- a/packages/redis-worker/src/worker.ts +++ b/packages/redis-worker/src/worker.ts @@ -140,7 +140,7 @@ class Worker { > = new Map(); constructor(private options: WorkerOptions) { - 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 { this.logger.error("Unhandled error in processItem:", { error: err, workerId, - item, + id: queueItem.id, + job: queueItem.job, }); } ); @@ -933,11 +934,12 @@ class Worker { 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 { name: this.options.name, id, job, - item, retryDate, retryDelay, visibilityTimeoutMs, @@ -1015,7 +1016,6 @@ class Worker { name: this.options.name, id, job, - item, visibilityTimeoutMs, error: requeueError, }