From 51315fc3c8afa7b24530fca013fc9e0e216bd86a Mon Sep 17 00:00:00 2001 From: nicktrn <55853254+nicktrn@users.noreply.github.com> Date: Tue, 2 Apr 2024 10:22:52 +0100 Subject: [PATCH] v3: shared queue pool (#994) * configurable retry delay checkpoint threshold * shared queue consumer pool * add more attributes to marqs spans * inject trace context into queued messages --- apps/coordinator/src/index.ts | 32 ++++++- apps/webapp/app/env.server.ts | 3 + apps/webapp/app/v3/handleSocketIo.server.ts | 1 + apps/webapp/app/v3/marqs/index.server.ts | 83 +++++++++++++++++-- .../v3/marqs/sharedQueueConsumer.server.ts | 43 +++++++--- apps/webapp/app/v3/sharedSocketConnection.ts | 69 +++++++++++++-- 6 files changed, 203 insertions(+), 28 deletions(-) diff --git a/apps/coordinator/src/index.ts b/apps/coordinator/src/index.ts index 176bc5087..29e8f4fee 100644 --- a/apps/coordinator/src/index.ts +++ b/apps/coordinator/src/index.ts @@ -18,6 +18,8 @@ collectDefaultMetrics(); const HTTP_SERVER_PORT = Number(process.env.HTTP_SERVER_PORT || 8020); const NODE_NAME = process.env.NODE_NAME || "coordinator"; +const DEFAULT_RETRY_DELAY_THRESHOLD_IN_MS = 30_000; + const REGISTRY_HOST = process.env.REGISTRY_HOST || "localhost:5000"; const CHECKPOINT_PATH = process.env.CHECKPOINT_PATH || "/checkpoints"; const REGISTRY_TLS_VERIFY = process.env.REGISTRY_TLS_VERIFY === "false" ? "false" : "true"; @@ -334,12 +336,19 @@ class TaskCoordinator { { resolve: (value: void) => void; reject: (err?: any) => void } >(); + #delayThresholdInMs: number; + constructor( private port: number, private host = "0.0.0.0" ) { this.#httpServer = this.#createHttpServer(); this.#checkpointer.initialize(); + this.#delayThresholdInMs = this.#getDelayThreshold(); + + if (process.env.DELAY_THRESHOLD_IN_MS) { + this.#delayThresholdInMs = this.#getDelayThreshold(); + } const io = new Server(this.#httpServer); this.#prodWorkerNamespace = this.#createProdWorkerNamespace(io); @@ -356,6 +365,27 @@ class TaskCoordinator { register.registerMetric(connectedTasksTotal); } + #getDelayThreshold() { + if (!process.env.RETRY_DELAY_THRESHOLD_IN_MS) { + return DEFAULT_RETRY_DELAY_THRESHOLD_IN_MS; + } + + const threshold = parseInt(process.env.RETRY_DELAY_THRESHOLD_IN_MS); + + if (isNaN(threshold)) { + logger.log( + "RETRY_DELAY_THRESHOLD_IN_MS parses as NaN, must supply integer. Will use default instead.", + { + RETRY_DELAY_THRESHOLD_IN_MS: process.env.RETRY_DELAY_THRESHOLD_IN_MS, + DEFAULT_DELAY_THRESHOLD_IN_MS: DEFAULT_RETRY_DELAY_THRESHOLD_IN_MS, + } + ); + return DEFAULT_RETRY_DELAY_THRESHOLD_IN_MS; + } + + return threshold; + } + #createPlatformSocket() { if (!PLATFORM_ENABLED) { console.log("INFO: platform connection disabled"); @@ -629,7 +659,7 @@ class TaskCoordinator { return; } - if (completion.retry.delay < 10_000) { + if (completion.retry.delay < this.#delayThresholdInMs) { completeWithoutCheckpoint(false); return; } diff --git a/apps/webapp/app/env.server.ts b/apps/webapp/app/env.server.ts index 84272dbe9..4996a5091 100644 --- a/apps/webapp/app/env.server.ts +++ b/apps/webapp/app/env.server.ts @@ -102,6 +102,9 @@ const EnvironmentSchema = z.object({ EVENTS_BATCH_SIZE: z.coerce.number().int().default(100), EVENTS_BATCH_INTERVAL: z.coerce.number().int().default(1000), EVENTS_DEFAULT_LOG_RETENTION: z.coerce.number().int().default(7), + SHARED_QUEUE_CONSUMER_POOL_SIZE: z.coerce.number().int().default(10), + SHARED_QUEUE_CONSUMER_INTERVAL_MS: z.coerce.number().int().default(100), + SHARED_QUEUE_CONSUMER_NEXT_TICK_INTERVAL_MS: z.coerce.number().int().default(100), // Development OTEL environment variables DEV_OTEL_EXPORTER_OTLP_ENDPOINT: z.string().optional(), diff --git a/apps/webapp/app/v3/handleSocketIo.server.ts b/apps/webapp/app/v3/handleSocketIo.server.ts index 6cb5b226c..c6353d9e2 100644 --- a/apps/webapp/app/v3/handleSocketIo.server.ts +++ b/apps/webapp/app/v3/handleSocketIo.server.ts @@ -168,6 +168,7 @@ function createSharedQueueConsumerNamespace(io: Server) { namespace: sharedQueue.namespace, socket, logger, + poolSize: env.SHARED_QUEUE_CONSUMER_POOL_SIZE, }); sharedSocketConnection.onClose.attach((closeEvent) => { diff --git a/apps/webapp/app/v3/marqs/index.server.ts b/apps/webapp/app/v3/marqs/index.server.ts index 8afcb2520..a2dcd1fe6 100644 --- a/apps/webapp/app/v3/marqs/index.server.ts +++ b/apps/webapp/app/v3/marqs/index.server.ts @@ -1,4 +1,9 @@ -import { Span, SpanKind, SpanOptions, trace } from "@opentelemetry/api"; +import { Span, SpanKind, SpanOptions, context, propagation, trace } from "@opentelemetry/api"; +import { + SEMATTRS_MESSAGE_ID, + SEMATTRS_MESSAGING_OPERATION, + SEMATTRS_MESSAGING_SYSTEM, +} from "@opentelemetry/semantic-conventions"; import { flattenAttributes } from "@trigger.dev/core/v3"; import Redis, { type Callback, type RedisOptions, type Result } from "ioredis"; import { env } from "~/env.server"; @@ -98,6 +103,8 @@ export class MarQS { const parentQueue = this.keys.envSharedQueueKey(env); + propagation.inject(context.active(), messageData); + const messagePayload: MessagePayload = { version: "1", data: messageData, @@ -117,7 +124,15 @@ export class MarQS { await this.#callEnqueueMessage(messagePayload); }, - { kind: SpanKind.PRODUCER, attributes: { ...attributesFromAuthenticatedEnv(env) } } + { + kind: SpanKind.PRODUCER, + attributes: { + [SEMATTRS_MESSAGING_OPERATION]: "publish", + [SEMATTRS_MESSAGE_ID]: messageId, + [SEMATTRS_MESSAGING_SYSTEM]: "marqs", + ...attributesFromAuthenticatedEnv(env), + }, + } ); } @@ -160,6 +175,7 @@ export class MarQS { if (message) { span.setAttributes({ + [SEMATTRS_MESSAGE_ID]: message.messageId, [SemanticAttributes.QUEUE]: message.queue, [SemanticAttributes.MESSAGE_ID]: message.messageId, [SemanticAttributes.CONCURRENCY_KEY]: message.concurrencyKey, @@ -171,7 +187,14 @@ export class MarQS { return message; }, - { kind: SpanKind.CONSUMER, attributes: { ...attributesFromAuthenticatedEnv(env) } } + { + kind: SpanKind.CONSUMER, + attributes: { + [SEMATTRS_MESSAGING_OPERATION]: "receive", + [SEMATTRS_MESSAGING_SYSTEM]: "marqs", + ...attributesFromAuthenticatedEnv(env), + }, + } ); } @@ -218,6 +241,7 @@ export class MarQS { if (message) { span.setAttributes({ + [SEMATTRS_MESSAGE_ID]: message.messageId, [SemanticAttributes.QUEUE]: message.queue, [SemanticAttributes.MESSAGE_ID]: message.messageId, [SemanticAttributes.CONCURRENCY_KEY]: message.concurrencyKey, @@ -229,7 +253,13 @@ export class MarQS { return message; }, - { kind: SpanKind.CONSUMER } + { + kind: SpanKind.CONSUMER, + attributes: { + [SEMATTRS_MESSAGING_OPERATION]: "receive", + [SEMATTRS_MESSAGING_SYSTEM]: "marqs", + }, + } ); } @@ -259,7 +289,14 @@ export class MarQS { messageId, }); }, - { kind: SpanKind.CONSUMER } + { + kind: SpanKind.CONSUMER, + attributes: { + [SEMATTRS_MESSAGING_OPERATION]: "ack", + [SEMATTRS_MESSAGE_ID]: messageId, + [SEMATTRS_MESSAGING_SYSTEM]: "marqs", + }, + } ); } @@ -305,7 +342,14 @@ export class MarQS { await this.#callEnqueueMessage(newMessage); }, - { kind: SpanKind.CONSUMER } + { + kind: SpanKind.CONSUMER, + attributes: { + [SEMATTRS_MESSAGING_OPERATION]: "replace", + [SEMATTRS_MESSAGE_ID]: messageId, + [SEMATTRS_MESSAGING_SYSTEM]: "marqs", + }, + } ); } @@ -370,7 +414,14 @@ export class MarQS { messageScore: retryAt, }); }, - { kind: SpanKind.CONSUMER } + { + kind: SpanKind.CONSUMER, + attributes: { + [SEMATTRS_MESSAGING_OPERATION]: "nack", + [SEMATTRS_MESSAGE_ID]: messageId, + [SEMATTRS_MESSAGING_SYSTEM]: "marqs", + }, + } ); } @@ -411,7 +462,14 @@ export class MarQS { return message.data; }, - { attributes: { [SemanticAttributes.MESSAGE_ID]: messageId } } + { + attributes: { + [SEMATTRS_MESSAGING_OPERATION]: "receive", + [SEMATTRS_MESSAGE_ID]: messageId, + [SEMATTRS_MESSAGING_SYSTEM]: "marqs", + [SemanticAttributes.MESSAGE_ID]: messageId, + }, + } ); } @@ -455,7 +513,14 @@ export class MarQS { return choice; }, - { kind: SpanKind.CONSUMER, attributes: { [SemanticAttributes.PARENT_QUEUE]: parentQueue } } + { + kind: SpanKind.CONSUMER, + attributes: { + [SEMATTRS_MESSAGING_OPERATION]: "receive", + [SEMATTRS_MESSAGING_SYSTEM]: "marqs", + [SemanticAttributes.PARENT_QUEUE]: parentQueue, + }, + } ); } diff --git a/apps/webapp/app/v3/marqs/sharedQueueConsumer.server.ts b/apps/webapp/app/v3/marqs/sharedQueueConsumer.server.ts index c787fa8eb..78f714743 100644 --- a/apps/webapp/app/v3/marqs/sharedQueueConsumer.server.ts +++ b/apps/webapp/app/v3/marqs/sharedQueueConsumer.server.ts @@ -1,4 +1,12 @@ -import { Context, ROOT_CONTEXT, Span, SpanKind, context, trace } from "@opentelemetry/api"; +import { + Context, + ROOT_CONTEXT, + Span, + SpanKind, + context, + propagation, + trace, +} from "@opentelemetry/api"; import { Machine, ProdTaskRunExecution, @@ -31,19 +39,24 @@ import { findCurrentWorkerDeployment } from "../models/workerDeployment.server"; const tracer = trace.getTracer("sharedQueueConsumer"); +const WithTraceContext = z.object({ + traceparent: z.string().optional(), + tracestate: z.string().optional(), +}); + const MessageBody = z.discriminatedUnion("type", [ - z.object({ + WithTraceContext.extend({ type: z.literal("EXECUTE"), taskIdentifier: z.string(), checkpointEventId: z.string().optional(), }), - z.object({ + WithTraceContext.extend({ type: z.literal("RESUME"), completedAttemptIds: z.string().array(), resumableAttemptId: z.string(), checkpointEventId: z.string().optional(), }), - z.object({ + WithTraceContext.extend({ type: z.literal("RESUME_AFTER_DURATION"), resumableAttemptId: z.string(), checkpointEventId: z.string(), @@ -57,6 +70,7 @@ export type SharedQueueConsumerOptions = { traceTimeoutSeconds?: number; nextTickInterval?: number; interval?: number; + parentContext?: Context; }; export class SharedQueueConsumer { @@ -83,6 +97,7 @@ export class SharedQueueConsumer { traceTimeoutSeconds: options.traceTimeoutSeconds ?? 60, // 60 seconds nextTickInterval: options.nextTickInterval ?? 1000, // 1 second interval: options.interval ?? 100, // 100ms + parentContext: options.parentContext ?? ROOT_CONTEXT, }; } @@ -187,8 +202,17 @@ export class SharedQueueConsumer { this.#doWork().finally(() => {}); } + #endCurrentSpan() { + if (this._currentSpan) { + this._currentSpan.setAttribute("tasks.period.failures", this._taskFailures); + this._currentSpan.setAttribute("tasks.period.successes", this._taskSuccesses); + this._currentSpan.end(); + } + } + async #doWork() { if (!this._enabled) { + this.#endCurrentSpan(); return; } @@ -199,12 +223,9 @@ export class SharedQueueConsumer { this._currentSpanContext === undefined || this._endSpanInNextIteration ) { - if (this._currentSpan) { - this._currentSpan.setAttribute("tasks.period.failures", this._taskFailures); - this._currentSpan.setAttribute("tasks.period.successes", this._taskSuccesses); + this.#endCurrentSpan(); - this._currentSpan.end(); - } + const parentContext = this._options.parentContext ?? ROOT_CONTEXT; // Create a new trace this._currentSpan = tracer.startSpan( @@ -212,11 +233,11 @@ export class SharedQueueConsumer { { kind: SpanKind.CONSUMER, }, - ROOT_CONTEXT + parentContext ); // Get the span trace context - this._currentSpanContext = trace.setSpan(ROOT_CONTEXT, this._currentSpan); + this._currentSpanContext = trace.setSpan(parentContext, this._currentSpan); this._perTraceCountdown = this._options.maximumItemsPerTrace; this._lastNewTrace = new Date(); diff --git a/apps/webapp/app/v3/sharedSocketConnection.ts b/apps/webapp/app/v3/sharedSocketConnection.ts index ac6d0fb76..d889f065c 100644 --- a/apps/webapp/app/v3/sharedSocketConnection.ts +++ b/apps/webapp/app/v3/sharedSocketConnection.ts @@ -10,7 +10,56 @@ import { Evt } from "evt"; import { randomUUID } from "node:crypto"; import { logger } from "~/services/logger.server"; import { SharedQueueConsumer } from "./marqs/sharedQueueConsumer.server"; -import { DisconnectReason, Namespace, Socket } from "socket.io"; +import type { DisconnectReason, Namespace, Socket } from "socket.io"; +import { ROOT_CONTEXT, Span, SpanKind, trace } from "@opentelemetry/api"; +import { env } from "~/env.server"; + +const tracer = trace.getTracer("sharedQueueConsumerPool"); + +interface SharedQueueConsumerPoolOptions { + sender: ZodMessageSender; + poolSize: number; +} + +class SharedQueueConsumerPool { + #consumers: SharedQueueConsumer[]; + #span: Span; + + constructor(opts: SharedQueueConsumerPoolOptions) { + this.#span = tracer.startSpan( + "SharedQueueConsumerPool()", + { + kind: SpanKind.CONSUMER, + attributes: { + "pool.size": opts.poolSize, + }, + }, + ROOT_CONTEXT + ); + + const spanContext = trace.setSpan(ROOT_CONTEXT, this.#span); + + this.#consumers = Array(opts.poolSize) + .fill(null) + .map( + () => + new SharedQueueConsumer(opts.sender, { + interval: env.SHARED_QUEUE_CONSUMER_INTERVAL_MS, + nextTickInterval: env.SHARED_QUEUE_CONSUMER_NEXT_TICK_INTERVAL_MS, + parentContext: spanContext, + }) + ); + } + + async start() { + await Promise.allSettled(this.#consumers.map((consumer) => consumer.start())); + } + + async stop() { + await Promise.allSettled(this.#consumers.map((consumer) => consumer.stop())); + this.#span.end(); + } +} interface SharedSocketConnectionOptions { namespace: Namespace< @@ -22,6 +71,7 @@ interface SharedSocketConnectionOptions { MessageCatalogToSocketIoEvents >; logger?: StructuredLogger; + poolSize?: number; } export class SharedSocketConnection { @@ -29,8 +79,9 @@ export class SharedSocketConnection { public onClose: Evt = new Evt(); private _sender: ZodMessageSender; - private _sharedConsumer: SharedQueueConsumer; + private _sharedQueueConsumerPool: SharedQueueConsumerPool; private _messageHandler: ZodMessageHandler; + private _defaultPoolSize = 10; constructor(opts: SharedSocketConnectionOptions) { this.id = randomUUID(); @@ -50,9 +101,13 @@ export class SharedSocketConnection { }, }); - this._sharedConsumer = new SharedQueueConsumer(this._sender, { - interval: 100, - nextTickInterval: 1000, + logger.log("Starting SharedQueueConsumer pool", { + poolSize: opts.poolSize ?? this._defaultPoolSize, + }); + + this._sharedQueueConsumerPool = new SharedQueueConsumerPool({ + poolSize: opts.poolSize ?? this._defaultPoolSize, + sender: this._sender, }); opts.socket.on("disconnect", this.#handleClose.bind(this)); @@ -62,7 +117,7 @@ export class SharedSocketConnection { schema: clientWebsocketMessages, messages: { READY_FOR_TASKS: async (payload) => { - this._sharedConsumer.start(); + this._sharedQueueConsumerPool.start(); }, BACKGROUND_WORKER_DEPRECATED: async (payload) => { // await this._sharedConsumer.deprecateBackgroundWorker(payload.backgroundWorkerId); @@ -89,7 +144,7 @@ export class SharedSocketConnection { } async #handleClose(ev: DisconnectReason) { - await this._sharedConsumer.stop(); + await this._sharedQueueConsumerPool.stop(); this.onClose.post(ev); }