diff --git a/apps/webapp/app/models/workflowRun.server.ts b/apps/webapp/app/models/workflowRun.server.ts index 260c65fda..ba60ce6a8 100644 --- a/apps/webapp/app/models/workflowRun.server.ts +++ b/apps/webapp/app/models/workflowRun.server.ts @@ -82,7 +82,8 @@ export async function failWorkflowRun( export async function completeWorkflowRun( output: string, runId: string, - apiKey: string + apiKey: string, + timestamp: string ) { const workflowRun = await findWorkflowRunScopedToApiKey(runId, apiKey); @@ -118,6 +119,7 @@ export async function completeWorkflowRun( context: {}, startedAt: new Date(), finishedAt: new Date(), + ts: timestamp, }, update: {}, }); @@ -128,7 +130,8 @@ export async function triggerEventInRun( key: string, event: z.infer, runId: string, - apiKey: string + apiKey: string, + timestamp: string ) { const workflowRun = await findWorkflowRunScopedToApiKey(runId, apiKey); @@ -146,6 +149,7 @@ export async function triggerEventInRun( context: {}, startedAt: new Date(), finishedAt: new Date(), + ts: timestamp, }); if (step.status === "EXISTING") { @@ -172,7 +176,8 @@ export async function logMessageInRun( key: string, log: z.infer, runId: string, - apiKey: string + apiKey: string, + timestamp: string ) { const workflowRun = await findWorkflowRunScopedToApiKey(runId, apiKey); @@ -189,6 +194,7 @@ export async function logMessageInRun( status: "SUCCESS", startedAt: new Date(), finishedAt: new Date(), + ts: timestamp, }); } diff --git a/apps/webapp/app/models/workflowRunPresenter.server.ts b/apps/webapp/app/models/workflowRunPresenter.server.ts index 5ab4eaa5f..65fbcd36c 100644 --- a/apps/webapp/app/models/workflowRunPresenter.server.ts +++ b/apps/webapp/app/models/workflowRunPresenter.server.ts @@ -196,7 +196,7 @@ function getWorkflowRun(prismaClient: PrismaClient, id: string) { }, }, }, - orderBy: { startedAt: "asc" }, + orderBy: { ts: "asc" }, }, }, }); diff --git a/apps/webapp/app/services/delays/initiateDelay.server.ts b/apps/webapp/app/services/delays/initiateDelay.server.ts index a2691881d..4339375ae 100644 --- a/apps/webapp/app/services/delays/initiateDelay.server.ts +++ b/apps/webapp/app/services/delays/initiateDelay.server.ts @@ -16,7 +16,11 @@ export class InitiateDelay { this.#prismaClient = prismaClient; } - async call(runId: string, delay: { key: string; wait: Wait }) { + async call( + runId: string, + timestamp: string, + delay: { key: string; wait: Wait } + ) { const delayUntil = this.#calculateDelayUntil(delay.wait); // Make sure the delay is not more than 1 year in the future @@ -32,6 +36,7 @@ export class InitiateDelay { context: { delayUntil: delayUntil.toISOString() }, status: "RUNNING", startedAt: new Date(), + ts: timestamp, }); if (idempotentStep.status === "EXISTING") { diff --git a/apps/webapp/app/services/messageBroker.server.ts b/apps/webapp/app/services/messageBroker.server.ts index ad738f157..dc72dd6fc 100644 --- a/apps/webapp/app/services/messageBroker.server.ts +++ b/apps/webapp/app/services/messageBroker.server.ts @@ -143,7 +143,8 @@ async function createTriggerSubscriber() { data.key, data.log, properties["x-workflow-run-id"], - properties["x-api-key"] + properties["x-api-key"], + properties["x-timestamp"] ); return true; @@ -166,7 +167,8 @@ async function createTriggerSubscriber() { await completeWorkflowRun( data.output, properties["x-workflow-run-id"], - properties["x-api-key"] + properties["x-api-key"], + properties["x-timestamp"] ); return true; @@ -185,6 +187,7 @@ async function createTriggerSubscriber() { data.key, properties["x-workflow-run-id"], properties["x-api-key"], + properties["x-timestamp"], data.request ); @@ -195,7 +198,8 @@ async function createTriggerSubscriber() { data.key, data.event, properties["x-workflow-run-id"], - properties["x-api-key"] + properties["x-api-key"], + properties["x-timestamp"] ); return true; @@ -203,10 +207,14 @@ async function createTriggerSubscriber() { INITIALIZE_DELAY: async (id, data, properties) => { const service = new InitiateDelay(); - await service.call(properties["x-workflow-run-id"], { - key: data.key, - wait: data.wait, - }); + await service.call( + properties["x-workflow-run-id"], + properties["x-timestamp"], + { + key: data.key, + wait: data.wait, + } + ); return true; }, diff --git a/apps/webapp/app/services/requests/createIntegrationRequest.server.ts b/apps/webapp/app/services/requests/createIntegrationRequest.server.ts index 5559a3767..2694e6480 100644 --- a/apps/webapp/app/services/requests/createIntegrationRequest.server.ts +++ b/apps/webapp/app/services/requests/createIntegrationRequest.server.ts @@ -16,6 +16,7 @@ export class CreateIntegrationRequest { key: string, runId: string, apiKey: string, + timestamp: string, data: { service: string; endpoint: string; @@ -61,6 +62,7 @@ export class CreateIntegrationRequest { endpoint: data.endpoint, }, status: "PENDING", + ts: timestamp, }); if (idempotentStep.status === "EXISTING") { diff --git a/apps/webapp/prisma/migrations/20230105173643_add_timestamp_to_steps/migration.sql b/apps/webapp/prisma/migrations/20230105173643_add_timestamp_to_steps/migration.sql new file mode 100644 index 000000000..5a767a48f --- /dev/null +++ b/apps/webapp/prisma/migrations/20230105173643_add_timestamp_to_steps/migration.sql @@ -0,0 +1,14 @@ +/* + Warnings: + + - Added the required column `timestamp` to the `WorkflowRunStep` table without a default value. This is not possible if the table is not empty. + +*/ +-- AlterTable +ALTER TABLE "WorkflowRunStep" ADD COLUMN "ts" INTEGER NULL; + +-- Add timestamps to existing steps based on the step createdAt (converting to unix timestamp since timestamp is an Integer) +UPDATE "WorkflowRunStep" SET ts = extract(epoch from "createdAt") * 1000; + +-- Make timestamp required +ALTER TABLE "WorkflowRunStep" ALTER COLUMN "ts" SET NOT NULL; diff --git a/apps/webapp/prisma/migrations/20230105175909_make_ts_a_string/migration.sql b/apps/webapp/prisma/migrations/20230105175909_make_ts_a_string/migration.sql new file mode 100644 index 000000000..f539653c5 --- /dev/null +++ b/apps/webapp/prisma/migrations/20230105175909_make_ts_a_string/migration.sql @@ -0,0 +1,2 @@ +-- AlterTable +ALTER TABLE "WorkflowRunStep" ALTER COLUMN "ts" SET DATA TYPE TEXT; diff --git a/apps/webapp/prisma/schema.prisma b/apps/webapp/prisma/schema.prisma index d7fb192a7..232eaf159 100644 --- a/apps/webapp/prisma/schema.prisma +++ b/apps/webapp/prisma/schema.prisma @@ -386,6 +386,7 @@ model WorkflowRunStep { runId String idempotencyKey String + ts String type WorkflowRunStepType input Json? diff --git a/apps/wss/src/runController.ts b/apps/wss/src/runController.ts index 5c01ebea0..ffb29fe87 100644 --- a/apps/wss/src/runController.ts +++ b/apps/wss/src/runController.ts @@ -154,25 +154,27 @@ export class WorkflowRunController { async close() { await this.#subscriber.close(); - await this.#publisher.publish( - "WORKFLOW_RUN_DISCONNECTED", - { - id: this.#runId, - }, - this.#publishProperties, - { partitionKey: this.#runId } - ); + await this.publish("WORKFLOW_RUN_DISCONNECTED", { + id: this.#runId, + }); this.#logger.debug("Workflow run closed"); } async publish( eventName: TEventName, - data: z.infer + data: z.infer, + timestamp: number = Date.now() ) { this.#logger.debug(`Publishing event ${eventName} with data`, data); - return this.#publisher.publish(eventName, data, this.#publishProperties, { + const properties = { + ...this.#publishProperties, + "x-timestamp": String(timestamp), + }; + + return this.#publisher.publish(eventName, data, properties, { + orderingKey: this.#runId, partitionKey: this.#runId, }); } diff --git a/examples/send-to-slack/src/index.ts b/examples/send-to-slack/src/index.ts index 4f2c9f1c2..374c22f78 100644 --- a/examples/send-to-slack/src/index.ts +++ b/examples/send-to-slack/src/index.ts @@ -23,12 +23,6 @@ const trigger = new Trigger({ await ctx.waitFor("initial-wait", { minutes: 1 }); - await ctx.logger.error("Error message!", { event }); - - await ctx.waitUntil("initial-wait-until", new Date(Date.now() + 1000 * 60)); - - await ctx.logger.info("Info message"); - const response = await slack.postMessage("send-to-slack", { channel: "test-integrations", text: `New domain created: ${event.domain} by customer ${event.customerId}`, @@ -36,8 +30,6 @@ const trigger = new Trigger({ await ctx.logger.debug("Debug message"); - await ctx.logger.warn("Warning message!"); - return response.message; }, }); diff --git a/packages/internal-bridge/src/logger.ts b/packages/internal-bridge/src/logger.ts index e62494ddb..b3ccd59e1 100644 --- a/packages/internal-bridge/src/logger.ts +++ b/packages/internal-bridge/src/logger.ts @@ -10,7 +10,9 @@ export class Logger { constructor(name: string, level: LogLevel = "info") { this.#name = name; - this.#level = logLevels.indexOf(level); + this.#level = logLevels.indexOf( + (process.env.TRIGGER_LOG_LEVEL ?? level) as LogLevel + ); } log(...args: any[]) { diff --git a/packages/internal-bridge/src/schemas/server.ts b/packages/internal-bridge/src/schemas/server.ts index be985347b..08e52851b 100644 --- a/packages/internal-bridge/src/schemas/server.ts +++ b/packages/internal-bridge/src/schemas/server.ts @@ -11,6 +11,7 @@ export const ServerRPCSchema = { runId: z.string(), key: z.string(), wait: WaitSchema, + timestamp: z.string(), }), response: z.boolean(), }, @@ -23,6 +24,7 @@ export const ServerRPCSchema = { endpoint: z.string(), params: z.any(), }), + timestamp: z.string(), }), response: z.boolean(), }, @@ -35,6 +37,7 @@ export const ServerRPCSchema = { level: z.enum(["DEBUG", "INFO", "WARN", "ERROR"]), properties: z.string().optional(), }), + timestamp: z.string(), }), response: z.boolean(), }, @@ -43,6 +46,7 @@ export const ServerRPCSchema = { runId: z.string(), key: z.string(), event: CustomEventSchema, + timestamp: z.string(), }), response: z.boolean(), }, @@ -70,6 +74,7 @@ export const ServerRPCSchema = { START_WORKFLOW_RUN: { request: z.object({ runId: z.string(), + timestamp: z.string(), }), response: z.boolean(), }, @@ -77,6 +82,7 @@ export const ServerRPCSchema = { request: z.object({ runId: z.string(), output: z.string(), + timestamp: z.string(), }), response: z.boolean(), }, @@ -88,6 +94,7 @@ export const ServerRPCSchema = { message: z.string(), stackTrace: z.string().optional(), }), + timestamp: z.string(), }), response: z.boolean(), }, diff --git a/packages/internal-platform/src/logger.ts b/packages/internal-platform/src/logger.ts index e62494ddb..b3ccd59e1 100644 --- a/packages/internal-platform/src/logger.ts +++ b/packages/internal-platform/src/logger.ts @@ -10,7 +10,9 @@ export class Logger { constructor(name: string, level: LogLevel = "info") { this.#name = name; - this.#level = logLevels.indexOf(level); + this.#level = logLevels.indexOf( + (process.env.TRIGGER_LOG_LEVEL ?? level) as LogLevel + ); } log(...args: any[]) { diff --git a/packages/internal-platform/src/messages/schemas/customEvents.ts b/packages/internal-platform/src/messages/schemas/customEvents.ts index fc7b3b010..8c6aa0779 100644 --- a/packages/internal-platform/src/messages/schemas/customEvents.ts +++ b/packages/internal-platform/src/messages/schemas/customEvents.ts @@ -1,6 +1,6 @@ -import { z } from "zod"; import { CustomEventSchema } from "@trigger.dev/common-schemas"; -import { WorkflowRunEventPropertiesSchema } from "../sharedSchemas"; +import { z } from "zod"; +import { WorkflowSendRunEventPropertiesSchema } from "../sharedSchemas"; export const wss = { TRIGGER_CUSTOM_EVENT: { @@ -8,6 +8,6 @@ export const wss = { key: z.string(), event: CustomEventSchema, }), - properties: WorkflowRunEventPropertiesSchema, + properties: WorkflowSendRunEventPropertiesSchema, }, }; diff --git a/packages/internal-platform/src/messages/sharedSchemas.ts b/packages/internal-platform/src/messages/sharedSchemas.ts index cf2426596..01c3324bd 100644 --- a/packages/internal-platform/src/messages/sharedSchemas.ts +++ b/packages/internal-platform/src/messages/sharedSchemas.ts @@ -20,4 +20,5 @@ export const WorkflowSendEventPropertiesSchema = z.object({ export const WorkflowSendRunEventPropertiesSchema = WorkflowSendEventPropertiesSchema.extend({ "x-workflow-run-id": z.string(), + "x-timestamp": z.string(), }); diff --git a/packages/internal-platform/src/messages/zodPublisher.ts b/packages/internal-platform/src/messages/zodPublisher.ts index 064e5023b..a06bb3e17 100644 --- a/packages/internal-platform/src/messages/zodPublisher.ts +++ b/packages/internal-platform/src/messages/zodPublisher.ts @@ -13,6 +13,7 @@ export type PublishOptions = { deliverAfter?: number; deliverAt?: number; partitionKey?: string; + orderingKey?: string; }; export type ZodPublisherOptions = @@ -100,6 +101,7 @@ export class ZodPublisher { const id = ulid(); this.#logger.debug("Publishing message", { + topic: this.#config.topic, type, data, properties, @@ -123,6 +125,7 @@ export class ZodPublisher { deliverAfter: options?.deliverAfter, deliverAt: options?.deliverAt, partitionKey: options?.partitionKey, + orderingKey: options?.orderingKey, }); return response.toString(); diff --git a/packages/internal-platform/src/messages/zodSubscriber.ts b/packages/internal-platform/src/messages/zodSubscriber.ts index b0eb13715..074e08f8d 100644 --- a/packages/internal-platform/src/messages/zodSubscriber.ts +++ b/packages/internal-platform/src/messages/zodSubscriber.ts @@ -135,15 +135,14 @@ export class ZodSubscriber { throw new Error(`Unknown message type: ${rawMessage.type}`); } - this.#logger.info( - `Handling message of type ${rawMessage.type}, parsing data and properties`, - rawMessage.data, - rawProperties - ); - const message = messageSchema.data.parse(rawMessage.data); const properties = messageSchema.properties.parse(rawProperties); + this.#logger.debug("Received message, calling handler", { + message, + properties, + }); + const handler = this.#handlers[typeName]; const returnValue = await handler(rawMessage.id, message, properties); diff --git a/packages/trigger-sdk/src/client.ts b/packages/trigger-sdk/src/client.ts index dfda068b9..e5de5ef61 100644 --- a/packages/trigger-sdk/src/client.ts +++ b/packages/trigger-sdk/src/client.ts @@ -246,6 +246,7 @@ export class TriggerClient { message, properties: JSON.stringify(properties ?? {}), }, + timestamp: String(highPrecisionTimestamp()), }); }), fireEvent: async (key, event) => { @@ -253,6 +254,7 @@ export class TriggerClient { runId: data.id, key, event: JSON.parse(JSON.stringify(event)), + timestamp: String(highPrecisionTimestamp()), }); }, waitFor: async (key, options) => { @@ -273,6 +275,7 @@ export class TriggerClient { hours: options.hours, days: options.days, }, + timestamp: String(highPrecisionTimestamp()), }); await result; @@ -294,6 +297,7 @@ export class TriggerClient { type: "SCHEDULE_FOR", scheduledFor: date.toISOString(), }, + timestamp: String(highPrecisionTimestamp()), }); await result; @@ -327,6 +331,7 @@ export class TriggerClient { endpoint: options.endpoint, params: options.params, }, + timestamp: String(highPrecisionTimestamp()), }); const output = await result; @@ -340,6 +345,7 @@ export class TriggerClient { serverRPC .send("START_WORKFLOW_RUN", { runId: data.id, + timestamp: String(highPrecisionTimestamp()), }) .then(() => { return this.#trigger.options @@ -348,6 +354,7 @@ export class TriggerClient { return serverRPC.send("COMPLETE_WORKFLOW_RUN", { runId: data.id, output: JSON.stringify(output), + timestamp: String(highPrecisionTimestamp()), }); }) .catch((anyError) => { @@ -377,6 +384,7 @@ export class TriggerClient { return serverRPC.send("SEND_WORKFLOW_ERROR", { runId: data.id, error, + timestamp: String(highPrecisionTimestamp()), }); }); }) @@ -384,6 +392,7 @@ export class TriggerClient { return serverRPC.send("SEND_WORKFLOW_ERROR", { runId: data.id, error: anyError, + timestamp: String(highPrecisionTimestamp()), }); }); } @@ -461,3 +470,9 @@ export const sleep = (ms: number) => new Promise((resolve) => setTimeout(resolve, ms)); const messageKey = (runId: string, key: string) => `${runId}:${key}`; + +function highPrecisionTimestamp() { + const [seconds, nanoseconds] = process.hrtime(); + + return seconds * 1e9 + nanoseconds; +}