From 6243ae30bb3454f116c17cf7ec2d5dc19fbc3e64 Mon Sep 17 00:00:00 2001 From: Matt Aitken Date: Fri, 24 May 2024 18:05:37 +0100 Subject: [PATCH 1/6] =?UTF-8?q?Send=20a=20=E2=80=9Csign-up=E2=80=9D=20even?= =?UTF-8?q?t=20to=20Loops=20(#1129)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- apps/webapp/app/env.server.ts | 2 + apps/webapp/app/services/loops.server.ts | 75 ++++++++++++++++++++ apps/webapp/app/services/telemetry.server.ts | 32 +++++---- 3 files changed, 97 insertions(+), 12 deletions(-) create mode 100644 apps/webapp/app/services/loops.server.ts diff --git a/apps/webapp/app/env.server.ts b/apps/webapp/app/env.server.ts index 71edf6a6c..07035ffb7 100644 --- a/apps/webapp/app/env.server.ts +++ b/apps/webapp/app/env.server.ts @@ -164,6 +164,8 @@ const EnvironmentSchema = z.object({ ALERT_RESEND_API_KEY: z.string().optional(), MAX_SEQUENTIAL_INDEX_FAILURE_COUNT: z.coerce.number().default(96), + + LOOPS_API_KEY: z.string().optional(), }); export type Environment = z.infer; diff --git a/apps/webapp/app/services/loops.server.ts b/apps/webapp/app/services/loops.server.ts new file mode 100644 index 000000000..75f72d027 --- /dev/null +++ b/apps/webapp/app/services/loops.server.ts @@ -0,0 +1,75 @@ +import { env } from "~/env.server"; +import { logger } from "./logger.server"; + +class LoopsClient { + constructor(private readonly apiKey: string) {} + + async userCreated({ + userId, + email, + name, + }: { + userId: string; + email: string; + name: string | null; + }) { + logger.info(`Loops send "sign-up" event`, { userId, email, name }); + return this.#sendEvent({ + email, + userId, + firstName: name?.split(" ").at(0), + eventName: "sign-up", + }); + } + + async #sendEvent({ + email, + userId, + firstName, + eventName, + eventProperties, + }: { + email: string; + userId: string; + firstName?: string; + eventName: string; + eventProperties?: Record; + }) { + const options = { + method: "POST", + headers: { Authorization: `Bearer ${this.apiKey}`, "Content-Type": "application/json" }, + body: JSON.stringify({ + email, + userId, + firstName, + eventName, + eventProperties, + }), + }; + + try { + const response = await fetch("https://app.loops.so/api/v1/events/send", options); + + if (!response.ok) { + logger.error(`Loops sendEvent ${eventName} bad status`, { status: response.status }); + return false; + } + + const responseBody = (await response.json()) as any; + + if (!responseBody.success) { + logger.error(`Loops sendEvent ${eventName} failed response`, { + message: responseBody.message, + }); + return false; + } + + return true; + } catch (error) { + logger.error(`Loops sendEvent ${eventName} failed`, { error }); + return false; + } + } +} + +export const loopsClient = env.LOOPS_API_KEY ? new LoopsClient(env.LOOPS_API_KEY) : null; diff --git a/apps/webapp/app/services/telemetry.server.ts b/apps/webapp/app/services/telemetry.server.ts index d6d83b23d..341a8094a 100644 --- a/apps/webapp/app/services/telemetry.server.ts +++ b/apps/webapp/app/services/telemetry.server.ts @@ -7,6 +7,7 @@ import type { Organization } from "~/models/organization.server"; import type { Project } from "~/models/project.server"; import type { User } from "~/models/user.server"; import { singleton } from "~/utils/singleton"; +import { loopsClient } from "./loops.server"; type Options = { postHogApiKey?: string; @@ -39,18 +40,19 @@ class Telemetry { user = { identify: ({ user, isNewUser }: { user: User; isNewUser: boolean }) => { - if (this.#posthogClient === undefined) return; - this.#posthogClient.identify({ - distinctId: user.id, - properties: { - email: user.email, - name: user.name, - authenticationMethod: user.authenticationMethod, - admin: user.admin, - createdAt: user.createdAt, - isNewUser, - }, - }); + if (this.#posthogClient) { + this.#posthogClient.identify({ + distinctId: user.id, + properties: { + email: user.email, + name: user.name, + authenticationMethod: user.authenticationMethod, + admin: user.admin, + createdAt: user.createdAt, + isNewUser, + }, + }); + } if (isNewUser) { this.#capture({ userId: user.id, @@ -64,6 +66,12 @@ class Telemetry { }, }); + loopsClient?.userCreated({ + userId: user.id, + email: user.email, + name: user.name, + }); + this.#triggerClient?.sendEvent({ name: "user.created", payload: { From ece6ca678a37ede9c0506aaf9c1586304fc0b1ed Mon Sep 17 00:00:00 2001 From: Eric Allam Date: Sat, 25 May 2024 21:43:16 +0100 Subject: [PATCH 2/6] Fix issue when using SDK in non-node environments by scoping the stream import with node: --- .changeset/new-pants-beg.md | 5 +++++ packages/core/src/v3/zodfetch.ts | 2 +- packages/core/tsup.config.ts | 1 + 3 files changed, 7 insertions(+), 1 deletion(-) create mode 100644 .changeset/new-pants-beg.md diff --git a/.changeset/new-pants-beg.md b/.changeset/new-pants-beg.md new file mode 100644 index 000000000..e484742c9 --- /dev/null +++ b/.changeset/new-pants-beg.md @@ -0,0 +1,5 @@ +--- +"@trigger.dev/core": patch +--- + +Fix issue when using SDK in non-node environments by scoping the stream import with node: diff --git a/packages/core/src/v3/zodfetch.ts b/packages/core/src/v3/zodfetch.ts index e77b27f9d..ac15f38e3 100644 --- a/packages/core/src/v3/zodfetch.ts +++ b/packages/core/src/v3/zodfetch.ts @@ -4,7 +4,7 @@ import { APIConnectionError, APIError } from "./apiErrors"; import { RetryOptions } from "./schemas"; import { calculateNextRetryDelay } from "./utils/retries"; import { FormDataEncoder } from "form-data-encoder"; -import { Readable } from "stream"; +import { Readable } from "node:stream"; export const defaultRetryOptions = { maxAttempts: 3, diff --git a/packages/core/tsup.config.ts b/packages/core/tsup.config.ts index edb33fc2d..f920f0982 100644 --- a/packages/core/tsup.config.ts +++ b/packages/core/tsup.config.ts @@ -16,4 +16,5 @@ export default defineConfig({ "./src/v3/prod/index.ts", "./src/v3/workers/index.ts", ], + external: ["node:stream"], }); From a56f9af9fe08746d378d22fbf391ef19ee23fcb6 Mon Sep 17 00:00:00 2001 From: Eric Allam Date: Sun, 26 May 2024 20:06:10 +0100 Subject: [PATCH 3/6] Safeguard against out of control v2 run executions --- apps/webapp/app/consts.ts | 1 + .../runs/performRunExecutionV3.server.ts | 30 +++++++++++++++++++ 2 files changed, 31 insertions(+) diff --git a/apps/webapp/app/consts.ts b/apps/webapp/app/consts.ts index 30dba9c7f..8b27c3c95 100644 --- a/apps/webapp/app/consts.ts +++ b/apps/webapp/app/consts.ts @@ -13,3 +13,4 @@ export const VERCEL_RESPONSE_TIMEOUT_STATUS_CODES = [408, 504]; export const MAX_BATCH_TRIGGER_ITEMS = 100; export const MAX_TASK_RUN_ATTEMPTS = 250; export const BULK_ACTION_RUN_LIMIT = 250; +export const MAX_JOB_RUN_EXECUTION_COUNT = 250; diff --git a/apps/webapp/app/services/runs/performRunExecutionV3.server.ts b/apps/webapp/app/services/runs/performRunExecutionV3.server.ts index 6f0ecd56c..de503e835 100644 --- a/apps/webapp/app/services/runs/performRunExecutionV3.server.ts +++ b/apps/webapp/app/services/runs/performRunExecutionV3.server.ts @@ -25,6 +25,7 @@ import { import { generateErrorMessage } from "zod-error"; import { eventRecordToApiJson } from "~/api.server"; import { + MAX_JOB_RUN_EXECUTION_COUNT, MAX_RUN_CHUNK_EXECUTION_LIMIT, MAX_RUN_YIELDED_EXECUTIONS, RUN_CHUNK_EXECUTION_BUFFER, @@ -141,6 +142,35 @@ export class PerformRunExecutionV3Service { }); } + // If the execution duration is greater than the maximum execution time, we need to fail the run + if (run.executionDuration >= run.organization.maximumExecutionTimePerRunInMs) { + await this.#failRunExecution( + this.#prismaClient, + run, + { + message: `Execution timed out after ${ + run.organization.maximumExecutionTimePerRunInMs / 1000 + } seconds`, + }, + "TIMED_OUT", + 0 + ); + return; + } + + if (run.executionCount >= MAX_JOB_RUN_EXECUTION_COUNT) { + await this.#failRunExecution( + this.#prismaClient, + run, + { + message: `Execution timed out after ${run.executionCount} executions`, + }, + "TIMED_OUT", + 0 + ); + return; + } + const client = new EndpointApi(run.environment.apiKey, run.endpoint.url); const event = eventRecordToApiJson(run.event); From d9ad72446e30877f98509949f5870dae3f5e4a5f Mon Sep 17 00:00:00 2001 From: Eric Allam Date: Sun, 26 May 2024 21:30:05 +0100 Subject: [PATCH 4/6] =?UTF-8?q?Abort=20v2=20runs=20when=20the=20job=20vers?= =?UTF-8?q?ion=20they=E2=80=99re=20associated=20with=20is=20disabled?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../runs/performRunExecutionV3.server.ts | 11 +++++++++ references/job-catalog/src/stressTest.ts | 23 +++++++++++++++++++ 2 files changed, 34 insertions(+) diff --git a/apps/webapp/app/services/runs/performRunExecutionV3.server.ts b/apps/webapp/app/services/runs/performRunExecutionV3.server.ts index de503e835..8d26f8888 100644 --- a/apps/webapp/app/services/runs/performRunExecutionV3.server.ts +++ b/apps/webapp/app/services/runs/performRunExecutionV3.server.ts @@ -142,6 +142,17 @@ export class PerformRunExecutionV3Service { }); } + if (run.version.status === "DISABLED") { + return await this.#failRunExecution( + this.#prismaClient, + run, + { + message: `Job version ${run.version.version} is disabled, aborting run.`, + }, + "ABORTED" + ); + } + // If the execution duration is greater than the maximum execution time, we need to fail the run if (run.executionDuration >= run.organization.maximumExecutionTimePerRunInMs) { await this.#failRunExecution( diff --git a/references/job-catalog/src/stressTest.ts b/references/job-catalog/src/stressTest.ts index 46689bf26..99603e90a 100644 --- a/references/job-catalog/src/stressTest.ts +++ b/references/job-catalog/src/stressTest.ts @@ -37,6 +37,29 @@ client.defineJob({ }, }); +client.defineJob({ + id: "stress-test-disabled", + name: "Stress Test Disabled", + version: "1.0.0", + trigger: eventTrigger({ + name: "stress.test.disabled", + }), + enabled: false, + run: async (payload, io, ctx) => { + await io.wait("wait-1", 20); + + await io.runTask( + `task-1`, + async (task) => { + await new Promise((resolve) => setTimeout(resolve, 10000)); + }, + { name: `Task 1` } + ); + + await io.wait("wait-2", 5); + }, +}); + client.defineJob({ id: "stress-test-2", name: "Stress Test 2", From ee3619bbb1d86ab726b204d7229e54aac6361d3d Mon Sep 17 00:00:00 2001 From: Matt Aitken Date: Mon, 27 May 2024 11:04:12 +0100 Subject: [PATCH 5/6] Added env var to enable/disable the v2 SqsEventConsumer --- apps/webapp/app/env.server.ts | 4 ++++ apps/webapp/app/services/events/sqsEventConsumer.ts | 6 ++++++ 2 files changed, 10 insertions(+) diff --git a/apps/webapp/app/env.server.ts b/apps/webapp/app/env.server.ts index 07035ffb7..ce14837b4 100644 --- a/apps/webapp/app/env.server.ts +++ b/apps/webapp/app/env.server.ts @@ -67,6 +67,10 @@ const EnvironmentSchema = z.object({ AWS_SQS_QUEUE_URL: z.string().optional(), AWS_SQS_BATCH_SIZE: z.coerce.number().int().optional().default(1), AWS_SQS_WAIT_TIME_MS: z.coerce.number().int().optional().default(100), + + /** v2 SQS Event Consumer Enabled */ + V2_SQS_EVENT_CONSUMER_ENABLED: z.string().default("false"), + DISABLE_SSE: z.string().optional(), OPENAI_API_KEY: z.string().optional(), diff --git a/apps/webapp/app/services/events/sqsEventConsumer.ts b/apps/webapp/app/services/events/sqsEventConsumer.ts index ceffee269..e351a75e9 100644 --- a/apps/webapp/app/services/events/sqsEventConsumer.ts +++ b/apps/webapp/app/services/events/sqsEventConsumer.ts @@ -134,6 +134,12 @@ export function getSharedSqsEventConsumer() { env.AWS_SQS_ACCESS_KEY_ID && env.AWS_SQS_SECRET_ACCESS_KEY ) { + if (env.V2_SQS_EVENT_CONSUMER_ENABLED !== "true") { + logger.info("V2_SQS_EVENT_CONSUMER_ENABLED isn't true, not starting consumer"); + return; + } + + logger.info("Starting SQS Event Consumer"); const consumer = new SqsEventConsumer(undefined, { queueUrl: env.AWS_SQS_QUEUE_URL, batchSize: env.AWS_SQS_BATCH_SIZE, From a5a5d3ae219b243ef0746f82ad8ac7635abab969 Mon Sep 17 00:00:00 2001 From: Matt Aitken Date: Mon, 27 May 2024 11:13:55 +0100 Subject: [PATCH 6/6] We could already disable the queue by not setting AWS_SQS_QUEUE_URL This reverts commit ee3619bbb1d86ab726b204d7229e54aac6361d3d. --- apps/webapp/app/env.server.ts | 4 ---- apps/webapp/app/services/events/sqsEventConsumer.ts | 6 ------ 2 files changed, 10 deletions(-) diff --git a/apps/webapp/app/env.server.ts b/apps/webapp/app/env.server.ts index ce14837b4..07035ffb7 100644 --- a/apps/webapp/app/env.server.ts +++ b/apps/webapp/app/env.server.ts @@ -67,10 +67,6 @@ const EnvironmentSchema = z.object({ AWS_SQS_QUEUE_URL: z.string().optional(), AWS_SQS_BATCH_SIZE: z.coerce.number().int().optional().default(1), AWS_SQS_WAIT_TIME_MS: z.coerce.number().int().optional().default(100), - - /** v2 SQS Event Consumer Enabled */ - V2_SQS_EVENT_CONSUMER_ENABLED: z.string().default("false"), - DISABLE_SSE: z.string().optional(), OPENAI_API_KEY: z.string().optional(), diff --git a/apps/webapp/app/services/events/sqsEventConsumer.ts b/apps/webapp/app/services/events/sqsEventConsumer.ts index e351a75e9..ceffee269 100644 --- a/apps/webapp/app/services/events/sqsEventConsumer.ts +++ b/apps/webapp/app/services/events/sqsEventConsumer.ts @@ -134,12 +134,6 @@ export function getSharedSqsEventConsumer() { env.AWS_SQS_ACCESS_KEY_ID && env.AWS_SQS_SECRET_ACCESS_KEY ) { - if (env.V2_SQS_EVENT_CONSUMER_ENABLED !== "true") { - logger.info("V2_SQS_EVENT_CONSUMER_ENABLED isn't true, not starting consumer"); - return; - } - - logger.info("Starting SQS Event Consumer"); const consumer = new SqsEventConsumer(undefined, { queueUrl: env.AWS_SQS_QUEUE_URL, batchSize: env.AWS_SQS_BATCH_SIZE,