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
2 changes: 1 addition & 1 deletion apps/webapp/app/services/logger.server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down
29 changes: 20 additions & 9 deletions apps/webapp/app/services/metadata/updateMetadata.server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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,
});
}
});
Expand Down Expand Up @@ -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,
});
}

Expand All @@ -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 (
Expand All @@ -567,8 +578,8 @@ export class UpdateMetadataService {
) {
if (this.flushLoggingEnabled) {
this.logger.debug(`[updateRunMetadataDirectly] Updating metadata directly for run`, {
metadata: metadataPacket.data,
runId,
metadataSizeBytes,
});
}

Expand Down Expand Up @@ -607,7 +618,7 @@ export class UpdateMetadataService {
if (this.flushLoggingEnabled) {
this.logger.debug(`[ingestRunOperations] Ingesting operations for run`, {
runId,
bufferedOperations,
operationCount: bufferedOperations.length,
});
}

Expand Down
10 changes: 9 additions & 1 deletion apps/webapp/app/utils/packets.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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") {
Expand All @@ -33,5 +41,5 @@ export function handleMetadataPacket(
throw new MetadataTooLargeError(`Metadata exceeds maximum size of ${maximumSize} bytes`);
}

return metadataPacket;
return { packet: metadataPacket, byteLength };
}
50 changes: 39 additions & 11 deletions apps/webapp/app/v3/services/alerts/deliverAlert.server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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": {
Expand Down Expand Up @@ -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;
}
Expand Down Expand Up @@ -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": {
Expand Down Expand Up @@ -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;
}
Expand Down Expand Up @@ -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": {
Expand Down Expand Up @@ -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;
}
Expand Down Expand Up @@ -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");

Expand Down Expand Up @@ -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,
Comment thread
coderabbitai[bot] marked this conversation as resolved.
});

throw new Error(`Failed to send alert webhook to ${webhook.url}`);
throw new Error(`Failed to send alert webhook to ${safeUrlHost(webhook.url)}`);
}
}

Expand Down Expand Up @@ -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";
}
}
10 changes: 5 additions & 5 deletions packages/redis-worker/src/worker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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"]);
Comment thread
carderne marked this conversation as resolved.
this.tracer = options.tracer ?? trace.getTracer(options.name);
this.meter = options.meter ?? metrics.getMeter(options.name);

Expand Down Expand Up @@ -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,
});
}
);
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -994,7 +996,6 @@ class Worker<TCatalog extends WorkerCatalog> {
name: this.options.name,
id,
job,
item,
retryDate,
retryDelay,
visibilityTimeoutMs,
Expand All @@ -1015,7 +1016,6 @@ class Worker<TCatalog extends WorkerCatalog> {
name: this.options.name,
id,
job,
item,
visibilityTimeoutMs,
error: requeueError,
}
Expand Down