Compare commits

...

44 Commits

Author SHA1 Message Date
Eric Allam 764df23d19 Release 3.0.0-beta.39 2024-06-19 12:22:19 +01:00
github-actions[bot] 4b961a6ae2 chore: Update version for release (beta) (#1169)
Co-authored-by: github-actions[bot] <github-actions[bot]@users.noreply.github.com>
2024-06-19 12:21:30 +01:00
Eric Allam 8757fdceef v3: [prod] force flush timeout should be 1s 2024-06-19 12:16:11 +01:00
Eric Allam 2404e88ac5 Add a IMPORTANT note to the ProdTaskRunExecution 2024-06-19 11:16:40 +01:00
Eric Allam 88b36f5090 Add a default on machine preset 2024-06-19 11:07:57 +01:00
Eric Allam b73ae3f927 Release 3.0.0-beta.38 2024-06-19 10:40:29 +01:00
github-actions[bot] b605b892ac chore: Update version for release (beta) (#1166)
Co-authored-by: github-actions[bot] <github-actions[bot]@users.noreply.github.com>
2024-06-19 10:39:24 +01:00
Eric Allam 233316f7e8 TaskRun.startedAt will now be lazily migrated 2024-06-18 15:49:14 +01:00
Eric Allam 1b90ffbb8c Add a changeset for the usage tracking PR (I forgot) 2024-06-18 14:56:25 +01:00
Matt Aitken b45ca4e146 If the timezone has changed then we need to reschedule the next scheduled run 2024-06-18 14:49:16 +01:00
Eric Allam 25d15578f7 Remove the update migration 2024-06-18 13:51:29 +01:00
Eric Allam fe865a0f49 Make startedAt backwards compat by giving it a default 2024-06-18 13:16:36 +01:00
Eric Allam 0ed93a748e v3: Remove aggressive otel flush timeouts in dev/prod 2024-06-18 11:48:35 +01:00
Michael Li e02320f65d fix: allow command login to read api url from cli args (#1168)
* fix: allow login to read api url from args

* Create strong-phones-smoke.md

---------

Co-authored-by: Eric Allam <eallam@icloud.com>
2024-06-18 09:43:42 +01:00
Eric Allam 85a543d8ec v3: usage tracking (#1163)
* Starting to measure wall time and cpu time in the workers, and reporting that via otel and to completed task run attempts

* Move usage tracking outside of the executor

* WIP prod usage tracking

* WIP

* WIP custom fetch to openmeter

* Create a usage client

* WIP

* WIP

* Implement new machine preset stuff and send usage reports to OpenMeter from webapp

* WIP

* Expose usage info to the client

* Add usage and cost to TaskEvent

* Add ability to globally configure the task machine preset

* Report start run usage

* Change the machine docs to use presets

* setExpirationTime to 24h

* Removed logs

* Update machines.mdx

* Removed console.logs

* Handle revalidating JWT tokens

* Couple tweaks

---------

Co-authored-by: Matt Aitken <matt@mattaitken.com>
2024-06-18 09:40:23 +01:00
Ryan Lee 10ceb85a92 fix: cloudfront timeouts not being caught (#1162) 2024-06-14 15:11:39 +01:00
Matt Aitken c405ae7117 Schedule limits and timezone support (#1165)
* Added maximumScheduleInstancesLimit column to Org, default to 20

* Docs on the schedule limits and improved soft-limit communication

* Added limit info to the schedules list page

* Created a task that creates schedules, useful for testing

* Make deduplicationKey required when creating/updating a schedule using the SDK

* New schedule button shows an alert if you’re over the limit

* Added timezone to the form and db

* WIP on the timezone dropdown for the create/edit schedule form

* Use the new filter search for timezones

* Made the timezone dropdown faster by fixing the virtualization

* The preview table is working and added a nice message about daylight savings

* Created a page where you can view the full list of timezones

The URL is included in the error message if you send an invalid time using the SDK

* Creating tasks with the timezone

* Added timezone support the the scheduler and the schedules list

* Added timezone support to more of the schedules UI

* The timezone comes through to scheduled runs with nice JSDocs

* Allow setting the timezone from the SDK

* Always have a timezone on a schedule

* Updated jsdocs

* Updated catalog example

* Changed the column to be a string, not null. Added the timezone across the SDK

* API endpoint for getting the timezones

* Added an SDK function to get the list of timezones

* Added timezones to the docs

* Changeset: Added timezone support to schedules

* Added support for testing timezone

* Tidied up imports

* Imports

* Imports

* Update limits.mdx

* Fixed a couple type issues and use the already exported zodfetch

---------

Co-authored-by: Eric Allam <eallam@icloud.com>
2024-06-14 13:31:27 +01:00
nicktrn 3687fcb61e Make pod cleaner interval configurable 2024-06-14 12:12:42 +01:00
Émile Ré d4ccdf7105 v3 CLI compiling E2E test suite (#1135)
* Boilerplate server-only use case

* wip: integration suite instrumentation setup

* Working poc testing compileProject

* Add pnpm script to run e2e tests only

* Use vitest globals

* Remove commented line

* Remove useless export

* Add modifier to test only one fixture project

* Handle package manager and log level choice

* Update server-only example

* Setup / teardown + split compile for package manager capabilities

* Ignore yarn files

* Fix issue with corepack, store version in engines field

* Rename test file

* Fix npm updates yarn.lock

* Move typecheking in a dedicated test

* Stop bundling the compile command to allow for more granular testing

* Put config resolving in separate test

* Add no-config test case and add test case expected errors configuration

* Add wantCompilationError option

* Add dependencies handling

* Use packageManager passed as option to resolve required deps

* Remove unused guard clauses

* Add postinstall & hash handling step

* Add worker start test

* Handle yarn.lock copy renaming on sigterm and sigkill

* Update vitest and use concurrent option

* Add a readme file

* Add CI workflow

* Fix handle cli deps

* Run cli v3 e2e tests on publish action

* Increase timeout on deps resolving step

* Add changeset

* Remove .pnp.cjs as we use yarn with nodeLinker node-modules

* Add missing .yarnrc.yml file

* No need to build CLI to run E2E tests

* Remove bun.lockb files

* Update beige-pears-explode.md

---------

Co-authored-by: Eric Allam <eallam@icloud.com>
2024-06-14 11:28:36 +01:00
nicktrn e08b4569e5 Improve checkpoint restore logging 2024-06-13 10:18:41 +01:00
nicktrn 79da0ca9b5 Update lockfile 2024-06-12 17:37:58 +01:00
github-actions[bot] 3aca603a33 chore: Update version for release (beta) (#1153)
Co-authored-by: github-actions[bot] <github-actions[bot]@users.noreply.github.com>
2024-06-12 17:34:47 +01:00
Eric Allam c9e97d6b78 Add debug logging when we skip importing a resource span 2024-06-12 17:27:12 +01:00
nicktrn 01633c9c03 Print error logs on dev worker failure 2024-06-12 13:08:50 +01:00
nicktrn 691990d79e Remove catalog example that was causing typecheck issues 2024-06-12 12:48:19 +01:00
nicktrn b2ba403dd3 Improve error logs cli message 2024-06-12 12:04:59 +01:00
nicktrn 1d47cab69f Ensure rollback triggers execution of tasks waiting for deploy 2024-06-12 11:51:51 +01:00
nicktrn e23047f9ad update pre-release script 2024-06-12 09:49:47 +01:00
nicktrn 68d32429b6 v3: checkpoint failover and misc fixes (#1157)
* configurable checkpoint registry namespace

* add missing task create await

* remove unused messages

* changeset

* update self-hosting docs

* capture and display stderr for failed deploys

* add missing lockfile changes

* stderr changeset

* fix cli stderr message

* update error logs label
2024-06-11 14:06:50 +01:00
Eric Allam 36ac79ac66 Make sure users cannot revoke invites for other orgs 2024-06-11 11:09:52 +01:00
Eric Allam ca94f0cac3 cli v3: fix deploys by using sub-path imports (fixes execa require ESM error) 2024-06-11 11:08:15 +01:00
nicktrn a5d8e453a5 v3: deploy rollbacks (#1154)
* add deployment rollbacks

* update test presenter

* update task list presenter

* fix test page for dev env

* read replica for the test presenter
2024-06-11 10:09:29 +01:00
Eric Allam c332519e72 v3: Add a much needed index on TaskEvent.spanId 2024-06-10 20:54:13 +01:00
Matt Aitken 52112c3bfc Docs: added v3 limits page 2024-06-10 16:07:08 +01:00
Eric Allam eae294a332 Add back in the v2 timeout task thing 2024-06-10 16:06:21 +01:00
Eric Allam 465cd0335c v2: No longer eagerly timeout runs when no tasks are created 2024-06-10 14:54:22 +01:00
nicktrn 35dbaedf69 v3: self-hosting (#1147)
* add amin email regex env var

* fix displayed init command for self-hosted setups

* shared env var to disable telemetry in cli and webapp

* pin sdk version during init

* if specified, add api url to dev command shown after init

* improve checkpoint support detection

* control forced checkpoint simulation via env var

* add public init to providers

* better checkpoint support check for coordinator

* add docker to coordinator image

* update docker provider containerfile

* bump remaining containers to node 20

* add infra image build to default publish workflow

* lockfile

* remove concurrency group from infra workflow

* add docker provider to build matrix

* fix var subst

* checkpoint test is docker specific

* enable v3 projects by default on self-hosted instances

* fix v3 setup command again

* add default posthog key

* self-hosting docs

* add latest tags to versioned infra and webapp builds

* some checkpoint errors should skip retrying

* add changeset

* shorten paragraph

* some docs updates

* update tunnelling section

* add registry setup section

* use correct cli push flag

* add checkout to v3 branch

* update the worker machine setup steps

* fix infra build

* small docs update

* remove unused feature function

* Revert "remove unused feature function"

This reverts commit cfe07887a12b6893dca8ce499964481a9b3dc9db.

* fix self-hosted v3 feature gate

* add note about missing arm support

* simplify helper script syntax
2024-06-10 14:13:04 +01:00
Eric Allam c11a77f50b cli v3: increase otel force flush timeout to 30s from 500ms 2024-06-07 20:46:03 +01:00
Eric Allam fb52b9efea Remove redundant log (you’re welcome baselime) 2024-06-07 20:41:34 +01:00
Matt Aitken 0896b9fffc Use the read replica more (#1152)
* Switch to read replica: getEvent API endpoint

* Switch to read replica: v2 run list presenter

* Switch to read replica: Job presenter

* Switch to read replica: Job list presenter

* Switch to read replica: billing client

* Switch to read replica: OrgUsagePresenter

* Switch to read replica: OrgBillingPlanPresenter

* Switch to read replica: ScheduleListPresenter

* Switch to read replica: EventRepository taskEvent.findMany
2024-06-07 15:47:54 +01:00
Matt Aitken 3a2dd983c5 Fix for sendEvent same id causing multiple runs (#1151)
* Proof of concept

* When ingesting events, if it’s already been delivered then don’t continue

* DeliverEvent: throw AlreadyDeliveredError and don’t retry if that’s thrown

* Test for duplicate event ids

* Return the original event so sendEvent doesn’t fail, don’t enqueue

* Add AlreadyDeliveredError to the logged out message

---------

Co-authored-by: Eric Allam <eallam@icloud.com>
2024-06-07 15:10:53 +01:00
Eric Allam a627ca67d1 Use a different logger key for the marqs version 2024-06-07 14:19:05 +01:00
Eric Allam afc180aa70 marqs: Concurrency monitor that runs periodically and vacuums completed runs (#1150) 2024-06-07 13:36:02 +01:00
Eric Allam 393af1b7c5 Add ability to enable v2 marqs requeuing (and remove deprecated marqs visibility requeuing) 2024-06-06 17:14:15 +01:00
288 changed files with 17017 additions and 1136 deletions
+5
View File
@@ -0,0 +1,5 @@
---
"trigger.dev": patch
---
Add an e2e suite to test compiling with v3 CLI.
+5
View File
@@ -0,0 +1,5 @@
---
"trigger.dev": patch
---
cli v3: increase otel force flush timeout to 30s from 500ms
+7
View File
@@ -0,0 +1,7 @@
---
"trigger.dev": patch
"@trigger.dev/core": patch
"@trigger.dev/sdk": patch
---
v3: Usage tracking
+5
View File
@@ -0,0 +1,5 @@
---
"@trigger.dev/core": patch
---
v3: Remove aggressive otel flush timeouts in dev/prod
+5
View File
@@ -0,0 +1,5 @@
---
"trigger.dev": patch
---
Output stderr logs on dev worker failure
+12
View File
@@ -46,10 +46,12 @@
"changesets": [
"afraid-sheep-joke",
"angry-eagles-trade",
"beige-pears-explode",
"beige-pens-dance",
"big-tomatoes-deliver",
"blue-pumas-whisper",
"breezy-gorillas-mate",
"brown-spies-burn",
"chilled-hornets-move",
"clean-pianos-listen",
"clever-apes-collect",
@@ -70,6 +72,7 @@
"green-bags-wink",
"hot-buckets-behave",
"hot-fishes-retire",
"hot-wasps-sin",
"itchy-chairs-itch",
"khaki-apricots-design",
"khaki-poems-lay",
@@ -86,6 +89,7 @@
"many-ligers-pump",
"mighty-camels-joke",
"mighty-flowers-train",
"mighty-parrots-sin",
"nasty-jars-pump",
"new-pants-beg",
"new-rivers-tell",
@@ -93,6 +97,7 @@
"ninety-pets-travel",
"odd-poets-own",
"pink-pumas-rhyme",
"plenty-ducks-beam",
"polite-ducks-switch",
"polite-rockets-matter",
"poor-flowers-cross",
@@ -103,22 +108,28 @@
"rich-kangaroos-unite",
"rotten-beers-refuse",
"rotten-dryers-exercise",
"rude-toys-compare",
"selfish-ducks-sort",
"serious-hats-rest",
"shaggy-spoons-taste",
"sharp-emus-compare",
"sharp-zebras-serve",
"shiny-coats-cry",
"silly-suits-switch",
"silver-doors-juggle",
"six-ligers-exist",
"sixty-insects-watch",
"slow-buses-own",
"slow-sloths-retire",
"smart-needles-move",
"smart-olives-eat",
"spicy-lamps-smoke",
"spicy-terms-bow",
"strange-ghosts-matter",
"strange-sheep-pull",
"strong-lemons-add",
"strong-owls-know",
"strong-phones-smoke",
"stupid-adults-sniff",
"stupid-bulldogs-applaud",
"sweet-lizards-press",
@@ -127,6 +138,7 @@
"tame-guests-know",
"tender-moose-tell",
"tender-oranges-rhyme",
"tender-turkeys-compete",
"thin-parents-heal",
"thirty-islands-kiss",
"tidy-balloons-suffer",
+7
View File
@@ -0,0 +1,7 @@
---
"@trigger.dev/core-apps": patch
"trigger.dev": patch
"@trigger.dev/core": patch
---
Capture and display stderr on index failures
+5
View File
@@ -0,0 +1,5 @@
---
"trigger.dev": patch
---
v3: [prod] force flush timeout should be 1s
+5
View File
@@ -0,0 +1,5 @@
---
"@trigger.dev/core": patch
---
Make deduplicationKey required when creating/updating a schedule
+7
View File
@@ -0,0 +1,7 @@
---
"@trigger.dev/core-apps": patch
"@trigger.dev/core": patch
---
- Fix uncaught provider exception
- Remove unused provider messages
+9
View File
@@ -0,0 +1,9 @@
---
"@trigger.dev/core-apps": patch
"trigger.dev": patch
---
- Fix init command SDK pinning
- Show --api-url / -a flag where needed
- CLI now also respects `TRIGGER_TELEMETRY_DISABLED`
- Dedicated docker checkpoint test function
+5
View File
@@ -0,0 +1,5 @@
---
"trigger.dev": patch
---
fix: allow command login to read api url from cli args
+6
View File
@@ -0,0 +1,6 @@
---
"@trigger.dev/sdk": patch
"@trigger.dev/core": patch
---
Added timezone support to schedules
+2
View File
@@ -25,6 +25,8 @@ DEV_OTEL_BATCH_PROCESSING_ENABLED="0"
# OPTIONAL VARIABLES
# This is used for validating emails that are allowed to log in. Every email that do not match this regex will be rejected.
# WHITELISTED_EMAILS="authorized@yahoo\.com|authorized@gmail\.com"
# Accounts with these emails will get global admin rights. This grants access to the admin UI.
# ADMIN_EMAILS="admin@example\.com|another-admin@example\.com"
# This is used for logging in via GitHub. You can leave these commented out if you don't want to use GitHub for authentication.
# AUTH_GITHUB_CLIENT_ID=
# AUTH_GITHUB_CLIENT_SECRET=
+47 -3
View File
@@ -1,9 +1,53 @@
name: "🧪 E2E Tests"
name: "E2E"
on:
workflow_call:
inputs:
package:
description: The identifier of the job to run
default: webapp
required: false
type: string
jobs:
e2e:
name: "🧪 E2E Tests"
cli-v3:
name: "🧪 CLI v3 tests"
if: inputs.package == 'cli-v3' || inputs.package == ''
runs-on: buildjet-8vcpu-ubuntu-2204
strategy:
fail-fast: false
matrix:
package-manager: ["npm", "pnpm", "yarn"]
steps:
- name: ⬇️ Checkout repo
uses: actions/checkout@v3
with:
fetch-depth: 0
- name: ⎔ Setup pnpm
uses: pnpm/action-setup@v2.2.4
with:
version: 8.15.5
- name: ⎔ Setup node
uses: buildjet/setup-node@v3
with:
node-version: 20.11.1
cache: "pnpm"
- name: 📥 Download deps
run: pnpm install --frozen-lockfile --filter trigger.dev...
- name: 🔧 Build v3 cli monorepo dependencies
run: pnpm run build --filter trigger.dev^...
- name: 🔧 Build worker template files
run: pnpm --filter trigger.dev run build:workers
- name: Run E2E Tests
run: |
PM=${{ matrix.package-manager }} pnpm --filter trigger.dev run test:e2e
webapp:
name: "🧪 Webapp tests"
if: inputs.package == 'webapp' || inputs.package == ''
runs-on: buildjet-16vcpu-ubuntu-2204
steps:
- name: 🐳 Login to Docker Hub
+2
View File
@@ -29,4 +29,6 @@ jobs:
# e2e:
# uses: ./.github/workflows/e2e.yml
# with:
# package: webapp
# secrets: inherit
+15 -2
View File
@@ -39,11 +39,25 @@ jobs:
exit 1
fi
echo "::set-output name=version::${IMAGE_TAG}"
- name: 🔢 Get the commit hash
id: get_commit
run: |
echo ::set-output name=sha_short::$(echo ${{ github.sha }} | cut -c1-7)
- name: 📛 Set the tags
id: set_tags
run: |
ref_without_tag=ghcr.io/triggerdotdev/trigger.dev
image_tags=$ref_without_tag:${{ steps.get_version.outputs.version }}
# if it's a versioned tag, also tag it as latest
if [[ "${{ github.ref_name }}" == v.docker.* ]]; then
image_tags=$image_tags,$ref_without_tag:latest
fi
echo "IMAGE_TAGS=${image_tags}" >> "$GITHUB_OUTPUT"
- name: 🐙 Login to GitHub Container Registry
uses: docker/login-action@v2
with:
@@ -56,6 +70,5 @@ jobs:
with:
file: ./docker/Dockerfile
platforms: linux/amd64,linux/arm64
tags: |
ghcr.io/triggerdotdev/trigger.dev:${{ steps.get_version.outputs.version }}
tags: ${{ steps.set_tags.outputs.IMAGE_TAGS }}
push: true
+38 -11
View File
@@ -1,6 +1,7 @@
name: "🚢 Publish Infra Images"
on:
workflow_call:
push:
tags:
- "infra-dev-*"
@@ -29,9 +30,6 @@ permissions:
packages: write
contents: read
concurrency:
group: ${{ github.workflow }}-${{ github.ref }}
env:
AWS_REGION: us-east-1
@@ -39,7 +37,7 @@ jobs:
build:
strategy:
matrix:
package: [coordinator, kubernetes-provider]
package: [coordinator, docker-provider, kubernetes-provider]
runs-on: buildjet-16vcpu-ubuntu-2204
env:
DOCKER_BUILDKIT: "1"
@@ -48,20 +46,40 @@ jobs:
- name: Generate image reference
id: prep
# WARNING: This step expects the workflow to have been triggered by a specific tag format of: infra-${env}-*
run: |
env=$(echo ${{ github.ref_name }} | cut -d- -f2)
sha=${GITHUB_SHA::7}
ts=$(date +%s)
# set image repo
if [[ "${{ matrix.package }}" == *-provider ]]; then
provider_type=$(echo ${{ matrix.package }} | cut -d- -f1)
provider_type=$(echo "${{ matrix.package }}" | cut -d- -f1)
repository=provider/${provider_type}
else
repository=${{ matrix.package }}
repository="${{ matrix.package }}"
fi
echo "IMAGE_TAG=${env}-${sha}-${ts}" >> "$GITHUB_OUTPUT"
echo "REPOSITORY=${repository}" >> "$GITHUB_OUTPUT"
# set image tag
if [[ "${{ github.ref_type }}" == "tag" ]]; then
if [[ "${{ github.ref_name }}" == infra-*-* ]]; then
env=$(echo ${{ github.ref_name }} | cut -d- -f2)
sha=$(echo ${{ github.sha }} | head -c7)
ts=$(date +%s)
image_tag=${env}-${sha}-${ts}
elif [[ "${{ github.ref_name }}" == v.docker.* ]]; then
version="${GITHUB_REF_NAME#v.docker.}"
image_tag="v${version}"
elif [[ "${{ github.ref_name }}" == build-* ]]; then
image_tag="${GITHUB_REF_NAME#build-}"
else
echo "Invalid tag: ${{ github.ref_name }}"
exit 1
fi
elif [[ "${{ github.ref_name }}" == "main" ]]; then
image_tag="main"
else
echo "Invalid reference: ${{ github.ref }}"
exit 1
fi
echo "IMAGE_TAG=${image_tag}" >> "$GITHUB_OUTPUT"
- name: Set up Docker Buildx
uses: docker/setup-buildx-action@v3
@@ -92,3 +110,12 @@ jobs:
REGISTRY: ghcr.io/triggerdotdev
REPOSITORY: ${{ steps.prep.outputs.REPOSITORY }}
IMAGE_TAG: ${{ steps.prep.outputs.IMAGE_TAG }}
- name: 🐙 Push 'latest' to GitHub Container Registry
if: startsWith(github.ref_name, 'v.docker.')
run: |
docker tag infra_image $REGISTRY/$REPOSITORY:latest
docker push $REGISTRY/$REPOSITORY:latest
env:
REGISTRY: ghcr.io/triggerdotdev
REPOSITORY: ${{ steps.prep.outputs.REPOSITORY }}
+10 -3
View File
@@ -49,11 +49,18 @@ jobs:
uses: ./.github/workflows/unit-tests.yml
secrets: inherit
# e2e:
# uses: ./.github/workflows/e2e.yml
# secrets: inherit
e2e:
uses: ./.github/workflows/e2e.yml
with:
package: cli-v3
secrets: inherit
publish:
needs: [typecheck, units]
uses: ./.github/workflows/publish-docker.yml
secrets: inherit
publish-infra:
needs: [typecheck, units]
uses: ./.github/workflows/publish-infra.yml
secrets: inherit
+4 -4
View File
@@ -1,19 +1,19 @@
# syntax=docker/dockerfile:labs
FROM node:18-bullseye-slim@sha256:a4edd54dcfdcacc8a4100fee71498e8671d99556a1acf5614539214a70092426 AS node-18
FROM node:20-bookworm-slim@sha256:72f2f046a5f8468db28730b990b37de63ce93fd1a72a40f531d6aa82afdf0d46 AS node-20
WORKDIR /app
FROM node-18 AS pruner
FROM node-20 AS pruner
COPY --chown=node:node . .
RUN npx -q turbo@1.10.9 prune --scope=coordinator --docker
RUN find . -name "node_modules" -type d -prune -exec rm -rf '{}' +
FROM node-18 AS base
FROM node-20 AS base
RUN apt-get update \
&& apt-get install -y buildah ca-certificates dumb-init \
&& apt-get install -y buildah ca-certificates dumb-init docker.io \
&& rm -rf /var/lib/apt/lists/*
COPY --chown=node:node .gitignore .gitignore
+49 -50
View File
@@ -12,7 +12,7 @@ import {
} from "@trigger.dev/core/v3";
import { ZodNamespace } from "@trigger.dev/core/v3/zodNamespace";
import { ZodSocketConnection } from "@trigger.dev/core/v3/zodSocket";
import { HttpReply, getTextBody, SimpleLogger } from "@trigger.dev/core-apps";
import { HttpReply, getTextBody, SimpleLogger, testDockerCheckpoint } from "@trigger.dev/core-apps";
import { ExponentialBackoff } from "./backoff";
import { collectDefaultMetrics, register, Gauge } from "prom-client";
@@ -43,6 +43,7 @@ const SIMULATE_CHECKPOINT_FAILURE_SECONDS = parseInt(
);
const REGISTRY_HOST = process.env.REGISTRY_HOST || "localhost:5000";
const REGISTRY_NAMESPACE = process.env.REGISTRY_NAMESPACE || "trigger";
const CHECKPOINT_PATH = process.env.CHECKPOINT_PATH || "/checkpoints";
const REGISTRY_TLS_VERIFY = process.env.REGISTRY_TLS_VERIFY === "false" ? "false" : "true";
@@ -72,7 +73,10 @@ type CheckpointAndPushOptions = {
type CheckpointAndPushResult =
| { success: true; checkpoint: CheckpointData }
| { success: false; reason?: "CANCELED" | "DISABLED" | "ERROR" | "IN_PROGRESS" | "NO_SUPPORT" };
| {
success: false;
reason?: "CANCELED" | "DISABLED" | "ERROR" | "IN_PROGRESS" | "NO_SUPPORT" | "SKIP_RETRYING";
};
type CheckpointData = {
location: string;
@@ -125,70 +129,58 @@ class Checkpointer {
constructor(private opts = { forceSimulate: false }) {}
async initialize(): Promise<CheckpointerInitializeReturn> {
async init(): Promise<CheckpointerInitializeReturn> {
if (this.#initialized) {
return this.#getInitializeReturn();
return this.#getInitReturn(this.#canCheckpoint);
}
this.#logger.log(`${this.#dockerMode ? "Docker" : "Kubernetes"} mode`);
if (this.#dockerMode) {
try {
await $`criu --version`;
} catch (error) {
this.#logger.error("No checkpoint support: Missing CRIU binary");
this.#logger.error("Will simulate instead");
this.#canCheckpoint = false;
this.#initialized = true;
const testCheckpoint = await testDockerCheckpoint();
return this.#getInitializeReturn();
if (testCheckpoint.ok) {
return this.#getInitReturn(true);
}
try {
await $`docker checkpoint`;
} catch (error) {
this.#logger.error(
"No checkpoint support: Docker needs to have experimental features enabled"
);
this.#logger.error("Will simulate instead");
this.#canCheckpoint = false;
this.#initialized = true;
return this.#getInitializeReturn();
}
this.#logger.error(testCheckpoint.message, testCheckpoint.error ?? "");
return this.#getInitReturn(false);
} else {
try {
await $`buildah login --get-login ${REGISTRY_HOST}`;
} catch (error) {
this.#logger.error(`No checkpoint support: Not logged in to registry ${REGISTRY_HOST}`);
this.#canCheckpoint = false;
this.#initialized = true;
return this.#getInitializeReturn();
return this.#getInitReturn(false);
}
}
this.#logger.log(
`Full checkpoint support${
this.#dockerMode && this.opts.forceSimulate ? " with forced simulation enabled." : "!"
}`
);
this.#initialized = true;
this.#canCheckpoint = true;
return this.#getInitializeReturn();
return this.#getInitReturn(true);
}
#getInitializeReturn(): CheckpointerInitializeReturn {
#getInitReturn(canCheckpoint: boolean): CheckpointerInitializeReturn {
this.#initialized = true;
this.#canCheckpoint = canCheckpoint;
if (canCheckpoint) {
this.#logger.log("Full checkpoint support!");
}
const willSimulate = this.#dockerMode && (!this.#canCheckpoint || this.opts.forceSimulate);
if (willSimulate) {
this.#logger.log("Simulation mode enabled. Containers will be paused, not checkpointed.", {
forceSimulate: this.opts.forceSimulate,
});
}
return {
canCheckpoint: this.#canCheckpoint,
willSimulate: this.#dockerMode && (!this.#canCheckpoint || this.opts.forceSimulate),
canCheckpoint,
willSimulate,
};
}
#getImageRef(projectRef: string, deploymentVersion: string, shortCode: string) {
return `${REGISTRY_HOST}/trigger/${projectRef}:${deploymentVersion}.prod-${shortCode}`;
return `${REGISTRY_HOST}/${REGISTRY_NAMESPACE}/${projectRef}:${deploymentVersion}.prod-${shortCode}`;
}
#getExportLocation(projectRef: string, deploymentVersion: string, shortCode: string) {
@@ -327,6 +319,11 @@ class Checkpointer {
return result;
}
if (result.reason === "SKIP_RETRYING") {
this.#logger.log("Skipping retrying", { runId });
return result;
}
continue;
} catch (error) {
this.#logger.error("Checkpoint error", {
@@ -355,7 +352,7 @@ class Checkpointer {
projectRef,
deploymentVersion,
}: CheckpointAndPushOptions): Promise<CheckpointAndPushResult> {
await this.initialize();
await this.init();
const options = {
runId,
@@ -473,7 +470,8 @@ class Checkpointer {
// Create checkpoint (CRI)
if (!this.#canCheckpoint) {
throw new Error("No checkpoint support in kubernetes mode.");
this.#logger.error("No checkpoint support in kubernetes mode.");
return { success: false, reason: "SKIP_RETRYING" };
}
const containerId = this.#logger.debug(
@@ -484,7 +482,8 @@ class Checkpointer {
);
if (!containerId.stdout) {
throw new Error("could not find container id");
this.#logger.error("could not find container id", { options, containterName });
return { success: false, reason: "SKIP_RETRYING" };
}
const start = performance.now();
@@ -617,7 +616,7 @@ class TaskCoordinator {
private host = "0.0.0.0"
) {
this.#httpServer = this.#createHttpServer();
this.#checkpointer.initialize();
this.#checkpointer.init();
this.#delayThresholdInMs = this.#getDelayThreshold();
if (process.env.DELAY_THRESHOLD_IN_MS) {
@@ -1034,7 +1033,7 @@ class TaskCoordinator {
return;
}
const { canCheckpoint, willSimulate } = await this.#checkpointer.initialize();
const { canCheckpoint, willSimulate } = await this.#checkpointer.init();
const willCheckpointAndRestore = canCheckpoint || willSimulate;
@@ -1131,7 +1130,7 @@ class TaskCoordinator {
return;
}
const { canCheckpoint, willSimulate } = await this.#checkpointer.initialize();
const { canCheckpoint, willSimulate } = await this.#checkpointer.init();
const willCheckpointAndRestore = canCheckpoint || willSimulate;
@@ -1185,7 +1184,7 @@ class TaskCoordinator {
socket.on("WAIT_FOR_TASK", async (message, callback) => {
logger.log("[WAIT_FOR_TASK]", message);
const { canCheckpoint, willSimulate } = await this.#checkpointer.initialize();
const { canCheckpoint, willSimulate } = await this.#checkpointer.init();
const willCheckpointAndRestore = canCheckpoint || willSimulate;
@@ -1227,7 +1226,7 @@ class TaskCoordinator {
socket.on("WAIT_FOR_BATCH", async (message, callback) => {
logger.log("[WAIT_FOR_BATCH]", message);
const { canCheckpoint, willSimulate } = await this.#checkpointer.initialize();
const { canCheckpoint, willSimulate } = await this.#checkpointer.init();
const willCheckpointAndRestore = canCheckpoint || willSimulate;
+41 -10
View File
@@ -1,16 +1,47 @@
# syntax=docker/dockerfile:labs
FROM node:18-slim AS base
RUN apt-get update \
&& apt-get install -y dumb-init
FROM base
FROM node:20-alpine@sha256:7a91aa397f2e2dfbfcdad2e2d72599f374e0b0172be1d86eeb73f1d33f36a4b2 AS node-20-alpine
WORKDIR /app
COPY --chown=node dist/index.mjs /app/
FROM node-20-alpine AS pruner
COPY --chown=node:node . .
RUN npx -q turbo@1.10.9 prune --scope=docker-provider --docker
RUN find . -name "node_modules" -type d -prune -exec rm -rf '{}' +
FROM node-20-alpine AS base
RUN apk add --no-cache dumb-init docker
COPY --chown=node:node .gitignore .gitignore
COPY --from=pruner --chown=node:node /app/out/json/ .
COPY --from=pruner --chown=node:node /app/out/pnpm-lock.yaml ./pnpm-lock.yaml
COPY --from=pruner --chown=node:node /app/out/pnpm-workspace.yaml ./pnpm-workspace.yaml
FROM base AS dev-deps
RUN corepack enable
ENV NODE_ENV development
RUN --mount=type=cache,id=pnpm,target=/root/.local/share/pnpm/store pnpm fetch --no-frozen-lockfile
RUN --mount=type=cache,id=pnpm,target=/root/.local/share/pnpm/store pnpm install --ignore-scripts --no-frozen-lockfile
FROM base AS builder
RUN corepack enable
COPY --from=pruner --chown=node:node /app/out/full/ .
COPY --from=dev-deps --chown=node:node /app/ .
COPY --chown=node:node turbo.json turbo.json
RUN pnpm run -r --filter docker-provider build:bundle
FROM base AS runner
RUN corepack enable
ENV NODE_ENV production
COPY --from=builder --chown=node:node /app/apps/docker-provider/dist/index.mjs ./index.mjs
EXPOSE 8000
ENTRYPOINT [ "/usr/bin/dumb-init", "--", "/usr/local/bin/node", "/app/index.mjs" ]
USER node
CMD [ "/usr/bin/dumb-init", "--", "/usr/local/bin/node", "./index.mjs" ]
+51 -75
View File
@@ -6,6 +6,8 @@ import {
TaskOperationsRestoreOptions,
TaskOperationsCreateOptions,
TaskOperationsIndexOptions,
isExecaChildProcess,
testDockerCheckpoint,
} from "@trigger.dev/core-apps";
import { setTimeout } from "node:timers/promises";
import { PostStartCauses, PreStopCauses } from "@trigger.dev/core/v3";
@@ -23,70 +25,58 @@ const FORCE_CHECKPOINT_SIMULATION = ["1", "true"].includes(
const logger = new SimpleLogger(`[${MACHINE_NAME}]`);
type InitializeReturn = {
type TaskOperationsInitReturn = {
canCheckpoint: boolean;
willSimulate: boolean;
};
function isExecaChildProcess(maybeExeca: unknown): maybeExeca is Awaited<ExecaChildProcess> {
return typeof maybeExeca === "object" && maybeExeca !== null && "escapedCommand" in maybeExeca;
}
class DockerTaskOperations implements TaskOperations {
#initialized = false;
#canCheckpoint = false;
constructor(private opts = { forceSimulate: false }) {}
async #initialize(): Promise<InitializeReturn> {
async init(): Promise<TaskOperationsInitReturn> {
if (this.#initialized) {
return this.#getInitializeReturn();
return this.#getInitReturn(this.#canCheckpoint);
}
logger.log("Initializing task operations");
if (this.opts.forceSimulate) {
logger.log("Forced simulation enabled. Will simulate regardless of checkpoint support.");
const testCheckpoint = await testDockerCheckpoint();
if (testCheckpoint.ok) {
return this.#getInitReturn(true);
}
try {
await $`criu --version`;
} catch (error) {
logger.error("No checkpoint support: Missing CRIU binary. Will simulate instead.");
this.#canCheckpoint = false;
this.#initialized = true;
return this.#getInitializeReturn();
}
try {
await $`docker checkpoint`;
} catch (error) {
logger.error("No checkpoint support: Docker needs to have experimental features enabled");
logger.error("Will simulate instead");
this.#canCheckpoint = false;
this.#initialized = true;
return this.#getInitializeReturn();
}
logger.log("Full checkpoint support!");
this.#initialized = true;
this.#canCheckpoint = true;
return this.#getInitializeReturn();
logger.error(testCheckpoint.message, testCheckpoint.error);
return this.#getInitReturn(false);
}
#getInitializeReturn(): InitializeReturn {
#getInitReturn(canCheckpoint: boolean): TaskOperationsInitReturn {
this.#initialized = true;
this.#canCheckpoint = canCheckpoint;
if (canCheckpoint) {
logger.log("Full checkpoint support!");
}
const willSimulate = !canCheckpoint || this.opts.forceSimulate;
if (willSimulate) {
logger.log("Simulation mode enabled. Containers will be paused, not checkpointed.", {
forceSimulate: this.opts.forceSimulate,
});
}
return {
canCheckpoint: this.#canCheckpoint,
willSimulate: !this.#canCheckpoint || this.opts.forceSimulate,
canCheckpoint,
willSimulate,
};
}
async index(opts: TaskOperationsIndexOptions) {
await this.#initialize();
await this.init();
const containerName = this.#getIndexContainerName(opts.shortCode);
@@ -95,41 +85,27 @@ class DockerTaskOperations implements TaskOperations {
port: COORDINATOR_PORT,
});
try {
logger.debug(
await execa("docker", [
"run",
"--network=host",
"--rm",
`--env=INDEX_TASKS=true`,
`--env=TRIGGER_SECRET_KEY=${opts.apiKey}`,
`--env=TRIGGER_API_URL=${opts.apiUrl}`,
`--env=TRIGGER_ENV_ID=${opts.envId}`,
`--env=OTEL_EXPORTER_OTLP_ENDPOINT=${OTEL_EXPORTER_OTLP_ENDPOINT}`,
`--env=POD_NAME=${containerName}`,
`--env=COORDINATOR_HOST=${COORDINATOR_HOST}`,
`--env=COORDINATOR_PORT=${COORDINATOR_PORT}`,
`--name=${containerName}`,
`${opts.imageRef}`,
])
);
} catch (error: any) {
if (!isExecaChildProcess(error)) {
throw error;
}
logger.error("Index failed:", {
opts,
exitCode: error.exitCode,
escapedCommand: error.escapedCommand,
stdout: error.stdout,
stderr: error.stderr,
});
}
logger.debug(
await execa("docker", [
"run",
"--network=host",
"--rm",
`--env=INDEX_TASKS=true`,
`--env=TRIGGER_SECRET_KEY=${opts.apiKey}`,
`--env=TRIGGER_API_URL=${opts.apiUrl}`,
`--env=TRIGGER_ENV_ID=${opts.envId}`,
`--env=OTEL_EXPORTER_OTLP_ENDPOINT=${OTEL_EXPORTER_OTLP_ENDPOINT}`,
`--env=POD_NAME=${containerName}`,
`--env=COORDINATOR_HOST=${COORDINATOR_HOST}`,
`--env=COORDINATOR_PORT=${COORDINATOR_PORT}`,
`--name=${containerName}`,
`${opts.imageRef}`,
])
);
}
async create(opts: TaskOperationsCreateOptions) {
await this.#initialize();
await this.init();
const containerName = this.#getRunContainerName(opts.runId);
@@ -165,7 +141,7 @@ class DockerTaskOperations implements TaskOperations {
}
async restore(opts: TaskOperationsRestoreOptions) {
await this.#initialize();
await this.init();
const containerName = this.#getRunContainerName(opts.runId);
@@ -194,7 +170,7 @@ class DockerTaskOperations implements TaskOperations {
}
async delete(opts: { runId: string }) {
await this.#initialize();
await this.init();
const containerName = this.#getRunContainerName(opts.runId);
await this.#sendPreStop(containerName);
@@ -203,7 +179,7 @@ class DockerTaskOperations implements TaskOperations {
}
async get(opts: { runId: string }) {
await this.#initialize();
await this.init();
logger.log("noop: get");
}
+3 -3
View File
@@ -1,14 +1,14 @@
FROM node:18-alpine@sha256:ca9f6cb0466f9638e59e0c249d335a07c867cd50c429b5c7830dda1bed584649 AS node-18-alpine
FROM node:20-alpine@sha256:7a91aa397f2e2dfbfcdad2e2d72599f374e0b0172be1d86eeb73f1d33f36a4b2 AS node-20-alpine
WORKDIR /app
FROM node-18-alpine AS pruner
FROM node-20-alpine AS pruner
COPY --chown=node:node . .
RUN npx -q turbo@1.10.9 prune --scope=kubernetes-provider --docker
RUN find . -name "node_modules" -type d -prune -exec rm -rf '{}' +
FROM node-18-alpine AS base
FROM node-20-alpine AS base
RUN apk add --no-cache dumb-init
+15 -5
View File
@@ -7,7 +7,12 @@ import {
TaskOperationsIndexOptions,
TaskOperationsRestoreOptions,
} from "@trigger.dev/core-apps";
import { Machine, PostStartCauses, PreStopCauses, EnvironmentType } from "@trigger.dev/core/v3";
import {
MachinePreset,
PostStartCauses,
PreStopCauses,
EnvironmentType,
} from "@trigger.dev/core/v3";
import { randomUUID } from "crypto";
import { TaskMonitor } from "./taskMonitor";
import { PodCleaner } from "./podCleaner";
@@ -16,6 +21,7 @@ const RUNTIME_ENV = process.env.KUBERNETES_PORT ? "kubernetes" : "local";
const NODE_NAME = process.env.NODE_NAME || "local";
const OTEL_EXPORTER_OTLP_ENDPOINT =
process.env.OTEL_EXPORTER_OTLP_ENDPOINT ?? "http://0.0.0.0:4318";
const POD_CLEANER_INTERVAL_SECONDS = Number(process.env.POD_CLEANER_INTERVAL_SECONDS || "300");
const logger = new SimpleLogger(`[${NODE_NAME}]`);
logger.log(`running in ${RUNTIME_ENV} mode`);
@@ -47,6 +53,10 @@ class KubernetesTaskOperations implements TaskOperations {
this.#k8sApi = this.#createK8sApi();
}
async init() {
// noop
}
async index(opts: TaskOperationsIndexOptions) {
await this.#createJob(
{
@@ -394,10 +404,10 @@ class KubernetesTaskOperations implements TaskOperations {
};
}
#getResourcesFromMachineConfig(config: Machine): ComputeResources {
#getResourcesFromMachineConfig(preset: MachinePreset): ComputeResources {
return {
cpu: `${config.cpu}`,
memory: `${config.memory}G`,
cpu: `${preset.cpu}`,
memory: `${preset.memory}G`,
};
}
@@ -551,7 +561,7 @@ taskMonitor.start();
const podCleaner = new PodCleaner({
runtimeEnv: RUNTIME_ENV,
namespace: "default",
intervalInSeconds: 300,
intervalInSeconds: POD_CLEANER_INTERVAL_SECONDS,
});
podCleaner.start();
+33 -3
View File
@@ -8,6 +8,7 @@ import {
} from "./primitives/ClientTabs";
import { ClipboardField } from "./primitives/ClipboardField";
import { Paragraph } from "./primitives/Paragraph";
import { useAppOrigin } from "~/hooks/useAppOrigin";
export function InitCommand({ appOrigin, apiKey }: { appOrigin: string; apiKey: string }) {
return (
@@ -133,9 +134,38 @@ export function TriggerDevStep({ extra }: { extra?: string }) {
// Trigger.dev version 3 setup commands
const v3PackageTag = "beta";
function getApiUrlArg() {
const appOrigin = useAppOrigin();
let apiUrl: string | undefined = undefined;
switch (appOrigin) {
case "https://cloud.trigger.dev":
// don't display the arg, use the CLI default
break;
case "https://test-cloud.trigger.dev":
apiUrl = "https://test-api.trigger.dev";
break;
case "https://internal.trigger.dev":
apiUrl = "https://internal-api.trigger.dev";
break;
default:
apiUrl = appOrigin;
break;
}
return apiUrl ? `-a ${apiUrl}` : undefined;
}
export function InitCommandV3() {
const project = useProject();
const projectRef = project.ref;
const apiUrlArg = getApiUrlArg();
const initCommandParts = [`trigger.dev@${v3PackageTag}`, "init", `-p ${projectRef}`, apiUrlArg];
const initCommand = initCommandParts.filter(Boolean).join(" ");
return (
<ClientTabs defaultValue="npm">
<ClientTabsList>
@@ -148,7 +178,7 @@ export function InitCommandV3() {
variant="primary/medium"
iconButton
className="mb-4"
value={`npx trigger.dev@${v3PackageTag} init -p ${projectRef}`}
value={`npx ${initCommand}`}
/>
</ClientTabsContent>
<ClientTabsContent value={"pnpm"}>
@@ -156,7 +186,7 @@ export function InitCommandV3() {
variant="primary/medium"
iconButton
className="mb-4"
value={`pnpm dlx trigger.dev@${v3PackageTag} init -p ${projectRef}`}
value={`pnpm dlx ${initCommand}`}
/>
</ClientTabsContent>
<ClientTabsContent value={"yarn"}>
@@ -164,7 +194,7 @@ export function InitCommandV3() {
variant="primary/medium"
iconButton
className="mb-4"
value={`yarn dlx trigger.dev@${v3PackageTag} init -p ${projectRef}`}
value={`yarn dlx ${initCommand}`}
/>
</ClientTabsContent>
</ClientTabs>
@@ -6,6 +6,7 @@ type DateTimeProps = {
timeZone?: string;
includeSeconds?: boolean;
includeTime?: boolean;
showTimezone?: boolean;
};
export const DateTime = ({
@@ -13,6 +14,7 @@ export const DateTime = ({
timeZone,
includeSeconds = true,
includeTime = true,
showTimezone = false,
}: DateTimeProps) => {
const locales = useLocales();
@@ -42,7 +44,12 @@ export const DateTime = ({
);
}, [locales, includeSeconds, realDate]);
return <Fragment>{formattedDateTime.replace(/\s/g, String.fromCharCode(32))}</Fragment>;
return (
<Fragment>
{formattedDateTime.replace(/\s/g, String.fromCharCode(32))}
{showTimezone ? ` (${timeZone ?? "UTC"})` : null}
</Fragment>
);
};
export function formatDateTime(
@@ -8,6 +8,7 @@ import { ShortcutDefinition, useShortcutKeys } from "~/hooks/useShortcutKeys";
import { cn } from "~/utils/cn";
import { ShortcutKey } from "./ShortcutKey";
import { ChevronDown } from "lucide-react";
import { MatchSorterOptions, matchSorter } from "match-sorter";
const sizes = {
small: {
@@ -75,7 +76,10 @@ export interface SelectProps<TValue extends string | string[], TItem>
showHeading?: boolean;
items?: TItem[] | Section<TItem>[];
empty?: React.ReactNode;
filter?: (item: ItemFromSection<TItem>, search: string, title?: string) => boolean;
filter?:
| boolean
| MatchSorterOptions<TItem>
| ((item: ItemFromSection<TItem>, search: string, title?: string) => boolean);
children:
| React.ReactNode
| ((
@@ -129,18 +133,44 @@ export function Select<TValue extends string | string[], TItem>({
if (!items) return [];
if (!searchValue || !filter) return items;
if (typeof filter === "function") {
if (isSection(items)) {
return items
.map((section) => ({
...section,
items: section.items.filter((item) =>
filter(item as ItemFromSection<TItem>, searchValue, section.title)
),
}))
.filter((section) => section.items.length > 0);
}
return items.filter((item) => filter(item as ItemFromSection<TItem>, searchValue));
}
if (typeof filter === "boolean" && filter) {
if (isSection(items)) {
return items
.map((section) => ({
...section,
items: matchSorter(section.items, searchValue),
}))
.filter((section) => section.items.length > 0);
}
return matchSorter(items, searchValue);
}
if (isSection(items)) {
return items
.map((section) => ({
...section,
items: section.items.filter((item) =>
filter(item as ItemFromSection<TItem>, searchValue, section.title)
),
items: matchSorter(section.items, searchValue, filter),
}))
.filter((section) => section.items.length > 0);
}
return items.filter((item) => filter(item as ItemFromSection<TItem>, searchValue));
return matchSorter(items, searchValue, filter);
}, [searchValue, items]);
const enableItemShortcuts = allowItemShortcuts && matches.length === items?.length;
@@ -20,6 +20,17 @@ export function DeploymentError({ errorData }: DeploymentErrorProps) {
maxLines={20}
/>
)}
{errorData.stderr && (
<>
<DeploymentErrorHeader title="Error logs:" />
<CodeBlock
showCopyButton={false}
showLineNumbers={false}
code={errorData.stderr}
maxLines={20}
/>
</>
)}
</div>
);
}
@@ -0,0 +1,55 @@
import { ArrowPathIcon } from "@heroicons/react/20/solid";
import { Form, useNavigation } from "@remix-run/react";
import { Button } from "~/components/primitives/Buttons";
import {
DialogContent,
DialogDescription,
DialogFooter,
DialogHeader,
} from "~/components/primitives/Dialog";
type RollbackDeploymentDialogProps = {
projectId: string;
deploymentShortCode: string;
redirectPath: string;
};
export function RollbackDeploymentDialog({
projectId,
deploymentShortCode,
redirectPath,
}: RollbackDeploymentDialogProps) {
const navigation = useNavigation();
const formAction = `/resources/${projectId}/deployments/${deploymentShortCode}/rollback`;
const isLoading = navigation.formAction === formAction;
return (
<DialogContent key="rollback">
<DialogHeader>Roll back to this deployment?</DialogHeader>
<DialogDescription>
This deployment will become the default for all future runs. Tasks triggered but not
included in this deploy will remain queued until you roll back to or create a new deployment
with these tasks included.
</DialogDescription>
<DialogFooter>
<Form
action={`/resources/${projectId}/deployments/${deploymentShortCode}/rollback`}
method="post"
>
<Button
type="submit"
name="redirectUrl"
value={redirectPath}
variant="primary/small"
LeadingIcon={isLoading ? "spinner-white" : ArrowPathIcon}
disabled={isLoading}
shortcut={{ modifiers: ["meta"], key: "enter" }}
>
{isLoading ? "Rolling back..." : "Roll back deployment"}
</Button>
</Form>
</DialogFooter>
</DialogContent>
);
}
@@ -0,0 +1,63 @@
import { useVirtualizer } from "@tanstack/react-virtual";
import { useRef } from "react";
import { SelectItem } from "../primitives/Select";
export function TimezoneList({ timezones }: { timezones: string[] }) {
const parentRef = useRef<HTMLDivElement>(null);
const rowVirtualizer = useVirtualizer({
count: timezones.length,
getScrollElement: () => parentRef.current,
estimateSize: () => 28,
});
return (
<div
ref={parentRef}
className="max-h-[calc(min(480px,var(--popover-available-height))-2.35rem)] overflow-y-auto overscroll-contain scrollbar-thin scrollbar-track-transparent scrollbar-thumb-charcoal-600"
>
<div
style={{
height: `${rowVirtualizer.getTotalSize()}px`,
width: "100%",
position: "relative",
}}
>
{rowVirtualizer.getVirtualItems().map((virtualItem) => (
<TimezoneCell
key={virtualItem.key}
size={virtualItem.size}
start={virtualItem.start}
timezone={timezones[virtualItem.index]}
/>
))}
</div>
</div>
);
}
function TimezoneCell({
timezone,
size,
start,
}: {
timezone: string;
size: number;
start: number;
}) {
return (
<SelectItem
value={timezone}
style={{
position: "absolute",
top: 0,
left: 0,
width: "100%",
height: `${size}px`,
transform: `translateY(${start}px)`,
}}
>
{timezone}
</SelectItem>
);
}
+21 -2
View File
@@ -1,5 +1,5 @@
import { z } from "zod";
import { SecretStoreOptionsSchema } from "./services/secrets/secretStoreOptionsSchema.server";
import { z } from "zod";
import { isValidRegex } from "./utils/regex";
import { isValidDatabaseUrl } from "./utils/db";
@@ -27,15 +27,17 @@ const EnvironmentSchema = z.object({
.string()
.refine(isValidRegex, "WHITELISTED_EMAILS must be a valid regex.")
.optional(),
ADMIN_EMAILS: z.string().refine(isValidRegex, "ADMIN_EMAILS must be a valid regex.").optional(),
REMIX_APP_PORT: z.string().optional(),
LOGIN_ORIGIN: z.string().default("http://localhost:3030"),
APP_ORIGIN: z.string().default("http://localhost:3030"),
APP_ENV: z.string().default(process.env.NODE_ENV),
SERVICE_NAME: z.string().default("trigger.dev webapp"),
SECRET_STORE: SecretStoreOptionsSchema.default("DATABASE"),
POSTHOG_PROJECT_KEY: z.string().optional(),
POSTHOG_PROJECT_KEY: z.string().default("phc_LFH7kJiGhdIlnO22hTAKgHpaKhpM8gkzWAFvHmf5vfS"),
TELEMETRY_TRIGGER_API_KEY: z.string().optional(),
TELEMETRY_TRIGGER_API_URL: z.string().optional(),
TRIGGER_TELEMETRY_DISABLED: z.string().optional(),
HIGHLIGHT_PROJECT_ID: z.string().optional(),
AUTH_GITHUB_CLIENT_ID: z.string().optional(),
AUTH_GITHUB_CLIENT_SECRET: z.string().optional(),
@@ -185,6 +187,23 @@ const EnvironmentSchema = z.object({
.default(60 * 1000 * 15),
V2_MARQS_DEFAULT_ENV_CONCURRENCY: z.coerce.number().int().default(100),
V2_MARQS_VERBOSE: z.string().default("0"),
V3_MARQS_CONCURRENCY_MONITOR_ENABLED: z.string().default("0"),
V2_MARQS_CONCURRENCY_MONITOR_ENABLED: z.string().default("0"),
/* Usage settings */
USAGE_EVENT_URL: z.string().optional(),
PROD_USAGE_HEARTBEAT_INTERVAL_MS: z.coerce.number().int().optional(),
CENTS_PER_HOUR_MICRO: z.coerce.number().default(0),
CENTS_PER_HOUR_SMALL_1X: z.coerce.number().default(0),
CENTS_PER_HOUR_SMALL_2X: z.coerce.number().default(0),
CENTS_PER_HOUR_MEDIUM_1X: z.coerce.number().default(0),
CENTS_PER_HOUR_MEDIUM_2X: z.coerce.number().default(0),
CENTS_PER_HOUR_LARGE_1X: z.coerce.number().default(0),
CENTS_PER_HOUR_LARGE_2X: z.coerce.number().default(0),
BASE_RUN_COST_IN_CENTS: z.coerce.number().default(0),
USAGE_OPEN_METER_API_KEY: z.string().optional(),
USAGE_OPEN_METER_BASE_URL: z.string().optional(),
});
export type Environment = z.infer<typeof EnvironmentSchema>;
+19 -11
View File
@@ -7,20 +7,28 @@ export type TriggerFeatures = {
alertsEnabled: boolean;
};
// If the request host is cloud.trigger.dev then we are on the managed cloud
// or if env.NODE_ENV is development
export function featuresForRequest(request: Request): TriggerFeatures {
const url = requestUrl(request);
const isManagedCloud =
url.host === "cloud.trigger.dev" ||
url.host === "test-cloud.trigger.dev" ||
url.host === "internal.trigger.dev" ||
process.env.CLOUD_ENV === "development";
function isManagedCloud(host: string): boolean {
return (
host === "cloud.trigger.dev" ||
host === "test-cloud.trigger.dev" ||
host === "internal.trigger.dev" ||
process.env.CLOUD_ENV === "development"
);
}
function featuresForHost(host: string): TriggerFeatures {
return {
isManagedCloud,
isManagedCloud: isManagedCloud(host),
v3Enabled: env.V3_ENABLED === "true",
alertsEnabled: env.ALERT_FROM_EMAIL !== undefined && env.ALERT_RESEND_API_KEY !== undefined,
};
}
export function featuresForRequest(request: Request): TriggerFeatures {
const url = requestUrl(request);
return featuresForUrl(url);
}
export function featuresForUrl(url: URL): TriggerFeatures {
return featuresForHost(url.host);
}
@@ -27,6 +27,7 @@ export function detectResponseIsTimeout(rawBody: string, response?: Response) {
return (
isResponseVercelTimeout(response) ||
isResponseCloudfrontTimeout(response) ||
isResponseDenoDeployTimeout(rawBody, response) ||
isResponseCloudflareTimeout(rawBody, response)
);
@@ -50,3 +51,7 @@ function isResponseVercelTimeout(response: Response) {
function isResponseDenoDeployTimeout(rawBody: string, response: Response) {
return response.status === 502 && rawBody.includes("TIME_LIMIT");
}
function isResponseCloudfrontTimeout(response: Response) {
return response.status === 504 && typeof response.headers.get("x-amz-cf-id") === "string";
}
+1
View File
@@ -245,6 +245,7 @@ export async function revokeInvite({
const invite = await prisma.orgMemberInvite.delete({
where: {
id: inviteId,
organizationId: org.id,
},
select: {
email: true,
@@ -8,10 +8,10 @@ import type {
import { customAlphabet } from "nanoid";
import slug from "slug";
import { prisma, PrismaClientOrTransaction } from "~/db.server";
import { createProject } from "./project.server";
import { generate } from "random-words";
import { createApiKeyForEnv, createPkApiKeyForEnv, envSlug } from "./api-key.server";
import { env } from "~/env.server";
import { featuresForUrl } from "~/features.server";
export type { Organization };
@@ -52,6 +52,8 @@ export async function createOrganization(
);
}
const features = featuresForUrl(new URL(env.APP_ORIGIN));
const organization = await prisma.organization.create({
data: {
title,
@@ -64,6 +66,7 @@ export async function createOrganization(
role: "ADMIN",
},
},
v3Enabled: features.v3Enabled && !features.isManagedCloud,
},
include: {
members: true,
+11 -2
View File
@@ -47,12 +47,21 @@ export async function findOrCreateMagicLinkUser(
},
});
const adminEmailRegex = env.ADMIN_EMAILS ? new RegExp(env.ADMIN_EMAILS) : undefined;
const makeAdmin = adminEmailRegex ? adminEmailRegex.test(input.email) : false;
const user = await prisma.user.upsert({
where: {
email: input.email,
},
update: { email: input.email },
create: { email: input.email, authenticationMethod: "MAGIC_LINK" },
update: {
email: input.email,
},
create: {
email: input.email,
authenticationMethod: "MAGIC_LINK",
admin: makeAdmin, // only on create, to prevent automatically removing existing admins
},
});
return {
+12 -10
View File
@@ -67,7 +67,11 @@ export type ZodTasks<TConsumerSchema extends MessageCatalogSchema> = {
maxAttempts?: number;
jobKeyMode?: "replace" | "preserve_run_at" | "unsafe_dedupe";
flags?: string[];
handler: (payload: z.infer<TConsumerSchema[K]>, job: GraphileJob) => Promise<void>;
handler: (
payload: z.infer<TConsumerSchema[K]>,
job: GraphileJob,
helpers: JobHelpers
) => Promise<void>;
};
};
@@ -80,7 +84,11 @@ export type ZodRecurringTasks = {
[key: string]: {
match: string;
options?: CronItemOptions;
handler: (payload: RecurringTaskPayload, job: GraphileJob) => Promise<void>;
handler: (
payload: RecurringTaskPayload,
job: GraphileJob,
helpers: JobHelpers
) => Promise<void>;
};
};
@@ -330,12 +338,6 @@ export class ZodWorker<TMessageCatalog extends MessageCatalogSchema> {
}
}
logger.debug("Enqueuing worker task", {
identifier,
payload,
spec,
});
const { job, durationInMs } = await this.#addJob(
identifier as string,
payload,
@@ -581,7 +583,7 @@ export class ZodWorker<TMessageCatalog extends MessageCatalogSchema> {
},
async (span) => {
try {
await task.handler(payload, job);
await task.handler(payload, job, helpers);
} catch (error) {
if (error instanceof Error) {
span.recordException(error);
@@ -662,7 +664,7 @@ export class ZodWorker<TMessageCatalog extends MessageCatalogSchema> {
},
async (span) => {
try {
await recurringTask.handler(payload._cron, job);
await recurringTask.handler(payload._cron, job, helpers);
} catch (error) {
if (error instanceof Error) {
span.recordException(error);
@@ -10,15 +10,12 @@ import { User } from "~/models/user.server";
import { z } from "zod";
import { projectPath } from "~/utils/pathBuilder";
import { JobRunStatus } from "@trigger.dev/database";
import { BasePresenter } from "./v3/basePresenter.server";
export type ProjectJob = Awaited<ReturnType<JobListPresenter["call"]>>[0];
export class JobListPresenter {
#prismaClient: PrismaClient;
constructor(prismaClient: PrismaClient = prisma) {
this.#prismaClient = prismaClient;
}
export class JobListPresenter extends BasePresenter {
public async call({
userId,
@@ -39,7 +36,7 @@ export class JobListPresenter {
? { some: { integration: { slug: integrationSlug } } }
: {};
const jobs = await this.#prismaClient.job.findMany({
const jobs = await this._replica.job.findMany({
select: {
id: true,
slug: true,
@@ -106,7 +103,7 @@ export class JobListPresenter {
}[];
if (jobs.length > 0) {
latestRuns = await this.#prismaClient.$queryRaw<
latestRuns = await this._replica.$queryRaw<
{
createdAt: Date;
status: JobRunStatus;
@@ -11,13 +11,10 @@ import { User } from "~/models/user.server";
import { z } from "zod";
import { projectPath } from "~/utils/pathBuilder";
import { Job } from "@trigger.dev/database";
import { BasePresenter } from "./v3/basePresenter.server";
export class JobPresenter {
#prismaClient: PrismaClient;
constructor(prismaClient: PrismaClient = prisma) {
this.#prismaClient = prismaClient;
}
export class JobPresenter extends BasePresenter {
public async call({
userId,
@@ -30,7 +27,7 @@ export class JobPresenter {
projectSlug: Project["slug"];
organizationSlug: Organization["slug"];
}) {
const job = await this.#prismaClient.job.findFirst({
const job = await this._replica.job.findFirst({
select: {
id: true,
slug: true,
@@ -1,14 +1,7 @@
import { PrismaClient, prisma } from "~/db.server";
import { logger } from "~/services/logger.server";
import { BillingService } from "../services/billing.server";
import { BasePresenter } from "./v3/basePresenter.server";
export class OrgBillingPlanPresenter {
#prismaClient: PrismaClient;
constructor(prismaClient: PrismaClient = prisma) {
this.#prismaClient = prismaClient;
}
export class OrgBillingPlanPresenter extends BasePresenter {
public async call({ slug, isManagedCloud }: { slug: string; isManagedCloud: boolean }) {
const billingPresenter = new BillingService(isManagedCloud);
const plans = await billingPresenter.getPlans();
@@ -17,7 +10,7 @@ export class OrgBillingPlanPresenter {
return;
}
const organization = await this.#prismaClient.organization.findFirst({
const organization = await this._replica.organization.findFirst({
where: {
slug,
},
@@ -27,7 +20,7 @@ export class OrgBillingPlanPresenter {
return;
}
const maxConcurrency = await this.#prismaClient.$queryRaw<
const maxConcurrency = await this._replica.$queryRaw<
{ organization_id: string; max_concurrent_runs: BigInt }[]
>`WITH events AS (
SELECT
@@ -1,17 +1,12 @@
import { estimate } from "@trigger.dev/billing";
import { sqlDatabaseSchema, PrismaClient, prisma } from "~/db.server";
import { sqlDatabaseSchema } from "~/db.server";
import { featuresForRequest } from "~/features.server";
import { BillingService } from "~/services/billing.server";
import { BasePresenter } from "./v3/basePresenter.server";
export class OrgUsagePresenter {
#prismaClient: PrismaClient;
constructor(prismaClient: PrismaClient = prisma) {
this.#prismaClient = prismaClient;
}
export class OrgUsagePresenter extends BasePresenter {
public async call({ userId, slug, request }: { userId: string; slug: string; request: Request }) {
const organization = await this.#prismaClient.organization.findFirst({
const organization = await this._replica.organization.findFirst({
where: {
slug,
members: {
@@ -27,7 +22,7 @@ export class OrgUsagePresenter {
}
// Get count of runs since the start of the current month
const runsCount = await this.#prismaClient.jobRun.count({
const runsCount = await this._replica.jobRun.count({
where: {
organizationId: organization.id,
createdAt: {
@@ -48,7 +43,7 @@ export class OrgUsagePresenter {
// ]
// This will be used to generate the chart on the usage page
// Use prisma queryRaw for this since prisma doesn't support grouping by month
const monthlyRunsDataRaw = await this.#prismaClient.$queryRaw<
const monthlyRunsDataRaw = await this._replica.$queryRaw<
{
month: string;
count: number;
@@ -64,7 +59,7 @@ export class OrgUsagePresenter {
const monthlyRunsDataDisplay = fillInMissingRunMonthlyData(monthlyRunsData, 6);
// Max concurrency each day over past 30 days
const concurrencyChartRawData = await this.#prismaClient.$queryRaw<
const concurrencyChartRawData = await this._replica.$queryRaw<
{ day: Date; max_concurrent_runs: BigInt }[]
>`
WITH time_boundaries AS (
@@ -115,7 +110,7 @@ export class OrgUsagePresenter {
concurrencyChartRawData
);
const dailyRunsRawData = await this.#prismaClient.$queryRaw<
const dailyRunsRawData = await this._replica.$queryRaw<
{ day: Date; runs: BigInt }[]
>`SELECT date_trunc('day', "createdAt") as day, COUNT(*) as runs FROM ${sqlDatabaseSchema}."JobRun" WHERE "organizationId" = ${organization.id} AND "createdAt" >= NOW() - INTERVAL '30 days' AND "internal" = FALSE GROUP BY day`;
@@ -7,6 +7,7 @@ import {
} from "~/components/runs/RunStatuses";
import { PrismaClient, prisma } from "~/db.server";
import { getUsername } from "~/utils/username";
import { BasePresenter } from "./v3/basePresenter.server";
type RunListOptions = {
userId: string;
@@ -27,12 +28,8 @@ const DEFAULT_PAGE_SIZE = 20;
export type RunList = Awaited<ReturnType<RunListPresenter["call"]>>;
export class RunListPresenter {
#prismaClient: PrismaClient;
constructor(prismaClient: PrismaClient = prisma) {
this.#prismaClient = prismaClient;
}
export class RunListPresenter extends BasePresenter {
public async call({
userId,
@@ -53,7 +50,7 @@ export class RunListPresenter {
const directionMultiplier = direction === "forward" ? 1 : -1;
// Find the organization that the user is a member of
const organization = await this.#prismaClient.organization.findFirstOrThrow({
const organization = await this._replica.organization.findFirstOrThrow({
select: {
id: true,
},
@@ -64,7 +61,7 @@ export class RunListPresenter {
});
// Find the project scoped to the organization
const project = await this.#prismaClient.project.findFirstOrThrow({
const project = await this._replica.project.findFirstOrThrow({
select: {
id: true,
},
@@ -75,7 +72,7 @@ export class RunListPresenter {
});
const job = jobSlug
? await this.#prismaClient.job.findFirstOrThrow({
? await this._replica.job.findFirstOrThrow({
where: {
slug: jobSlug,
projectId: project.id,
@@ -84,10 +81,10 @@ export class RunListPresenter {
: undefined;
const event = eventId
? await this.#prismaClient.eventRecord.findUnique({ where: { id: eventId } })
? await this._replica.eventRecord.findUnique({ where: { id: eventId } })
: undefined;
const runs = await this.#prismaClient.jobRun.findMany({
const runs = await this._replica.jobRun.findMany({
select: {
id: true,
number: true,
@@ -80,7 +80,7 @@ export class ApiRetrieveRunPresenter extends BasePresenter {
version: taskRun.lockedToVersion ? taskRun.lockedToVersion.version : undefined,
createdAt: taskRun.createdAt ?? undefined,
updatedAt: taskRun.updatedAt ?? undefined,
startedAt: taskRun.lockedAt ?? undefined,
startedAt: taskRun.startedAt ?? taskRun.lockedAt ?? undefined,
finishedAt: ApiRetrieveRunPresenter.isStatusFinished(apiStatus)
? taskRun.updatedAt
: undefined,
@@ -7,6 +7,9 @@ import { getUsername } from "~/utils/username";
const pageSize = 20;
export type DeploymentList = Awaited<ReturnType<DeploymentListPresenter["call"]>>;
export type DeploymentListItem = DeploymentList["deployments"][0];
export class DeploymentListPresenter {
#prismaClient: PrismaClient;
@@ -136,6 +139,8 @@ LIMIT ${pageSize} OFFSET ${pageSize * (page - 1)};`;
deployedAt: deployment.deployedAt,
tasksCount: deployment.tasksCount ? Number(deployment.tasksCount) : null,
label: label?.label,
isCurrent: label?.label === "current",
isDeployed: deployment.status === "DEPLOYED",
environment: {
id: environment.id,
type: environment.type,
@@ -17,6 +17,7 @@ export type ErrorData = {
name: string;
message: string;
stack?: string;
stderr?: string;
};
export class DeploymentPresenter {
@@ -177,17 +178,20 @@ export class DeploymentPresenter {
name: parsedErrorData.data.name,
message: parsedErrorData.data.message,
stack: createTaskMetadataFailedErrorStack(parsedError.data),
stderr: parsedErrorData.data.stderr,
};
} else {
return {
name: parsedErrorData.data.name,
message: parsedErrorData.data.message,
stderr: parsedErrorData.data.stderr,
};
}
} else {
return {
name: parsedErrorData.data.name,
message: parsedErrorData.data.message,
stderr: parsedErrorData.data.stderr,
};
}
}
@@ -196,6 +200,7 @@ export class DeploymentPresenter {
name: parsedErrorData.data.name,
message: parsedErrorData.data.message,
stack: parsedErrorData.data.stack,
stderr: parsedErrorData.data.stderr,
};
}
}
@@ -1,6 +1,7 @@
import { RuntimeEnvironmentType } from "@trigger.dev/database";
import { PrismaClient, prisma } from "~/db.server";
import { displayableEnvironment } from "~/models/runtimeEnvironment.server";
import { getTimezones } from "~/utils/timezones.server";
type EditScheduleOptions = {
userId: string;
@@ -74,6 +75,7 @@ export class EditSchedulePresenter {
return {
possibleTasks: possibleTasks.map((task) => task.slug),
possibleEnvironments,
possibleTimezones: getTimezones(),
schedule: await this.#getExistingSchedule(friendlyId, possibleEnvironments),
};
}
@@ -91,6 +93,7 @@ export class EditSchedulePresenter {
externalId: true,
deduplicationKey: true,
userProvidedDeduplicationKey: true,
timezone: true,
taskIdentifier: true,
instances: {
select: {
@@ -156,6 +156,7 @@ export class RunListPresenter extends BasePresenter {
runtimeEnvironmentId: string;
status: TaskRunStatus;
createdAt: Date;
startedAt: Date | null;
lockedAt: Date | null;
updatedAt: Date;
isTest: boolean;
@@ -172,6 +173,7 @@ export class RunListPresenter extends BasePresenter {
tr."runtimeEnvironmentId" AS "runtimeEnvironmentId",
tr.status AS status,
tr."createdAt" AS "createdAt",
tr."startedAt" AS "startedAt",
tr."lockedAt" AS "lockedAt",
tr."updatedAt" AS "updatedAt",
tr."isTest" AS "isTest",
@@ -272,13 +274,15 @@ export class RunListPresenter extends BasePresenter {
const hasFinished = FINISHED_STATUSES.includes(run.status);
const startedAt = run.startedAt ?? run.lockedAt;
return {
id: run.id,
friendlyId: run.runFriendlyId,
number: Number(run.number),
createdAt: run.createdAt.toISOString(),
updatedAt: run.updatedAt.toISOString(),
startedAt: run.lockedAt ? run.lockedAt.toISOString() : undefined,
startedAt: startedAt ? startedAt.toISOString() : undefined,
hasFinished,
finishedAt: hasFinished ? run.updatedAt.toISOString() : undefined,
isTest: run.isTest,
@@ -4,6 +4,7 @@ import { PrismaClient, prisma, sqlDatabaseSchema } from "~/db.server";
import { displayableEnvironment } from "~/models/runtimeEnvironment.server";
import { getUsername } from "~/utils/username";
import { calculateNextScheduledTimestamp } from "~/v3/utils/calculateNextSchedule.server";
import { BasePresenter } from "./basePresenter.server";
type ScheduleListOptions = {
projectId: string;
@@ -21,6 +22,7 @@ export type ScheduleListItem = {
userProvidedDeduplicationKey: boolean;
cron: string;
cronDescription: string;
timezone: string;
externalId: string | null;
nextRun: Date;
lastRun: Date | undefined;
@@ -34,13 +36,7 @@ export type ScheduleListItem = {
export type ScheduleList = Awaited<ReturnType<ScheduleListPresenter["call"]>>;
export type ScheduleListAppliedFilters = ScheduleList["filters"];
export class ScheduleListPresenter {
#prismaClient: PrismaClient;
constructor(prismaClient: PrismaClient = prisma) {
this.#prismaClient = prismaClient;
}
export class ScheduleListPresenter extends BasePresenter {
public async call({
userId,
projectId,
@@ -54,7 +50,7 @@ export class ScheduleListPresenter {
tasks !== undefined || environments !== undefined || (search !== undefined && search !== "");
// Find the project scoped to the organization
const project = await this.#prismaClient.project.findFirstOrThrow({
const project = await this._replica.project.findFirstOrThrow({
select: {
id: true,
environments: {
@@ -75,14 +71,25 @@ export class ScheduleListPresenter {
},
},
},
organization: {
select: {
maximumSchedulesLimit: true,
},
},
},
where: {
id: projectId,
},
});
const schedulesCount = await this._prisma.taskSchedule.count({
where: {
projectId,
},
});
//get all possible scheduled tasks
const possibleTasks = await this.#prismaClient.backgroundWorkerTask.findMany({
const possibleTasks = await this._replica.backgroundWorkerTask.findMany({
distinct: ["slug"],
where: {
projectId: project.id,
@@ -93,7 +100,7 @@ export class ScheduleListPresenter {
//do this here to protect against SQL injection
search = search && search !== "" ? `%${search}%` : undefined;
const totalCount = await this.#prismaClient.taskSchedule.count({
const totalCount = await this._replica.taskSchedule.count({
where: {
projectId: project.id,
taskIdentifier: tasks ? { in: tasks } : undefined,
@@ -135,7 +142,7 @@ export class ScheduleListPresenter {
},
});
const rawSchedules = await this.#prismaClient.taskSchedule.findMany({
const rawSchedules = await this._replica.taskSchedule.findMany({
select: {
id: true,
friendlyId: true,
@@ -144,6 +151,7 @@ export class ScheduleListPresenter {
userProvidedDeduplicationKey: true,
generatorExpression: true,
generatorDescription: true,
timezone: true,
externalId: true,
instances: {
select: {
@@ -199,7 +207,7 @@ export class ScheduleListPresenter {
const latestRuns =
rawSchedules.length > 0
? await this.#prismaClient.$queryRaw<{ scheduleId: string; createdAt: Date }[]>`
? await this._replica.$queryRaw<{ scheduleId: string; createdAt: Date }[]>`
SELECT t."scheduleId", t."createdAt"
FROM (
SELECT "scheduleId", MAX("createdAt") as "LatestRun"
@@ -222,10 +230,11 @@ export class ScheduleListPresenter {
userProvidedDeduplicationKey: schedule.userProvidedDeduplicationKey,
cron: schedule.generatorExpression,
cronDescription: schedule.generatorDescription,
timezone: schedule.timezone,
active: schedule.active,
externalId: schedule.externalId,
lastRun: latestRun?.createdAt,
nextRun: calculateNextScheduledTimestamp(schedule.generatorExpression),
nextRun: calculateNextScheduledTimestamp(schedule.generatorExpression, schedule.timezone),
environments: schedule.instances.map((instance) => {
const environment = project.environments.find((env) => env.id === instance.environmentId);
if (!environment) {
@@ -249,6 +258,10 @@ export class ScheduleListPresenter {
return displayableEnvironment(environment, userId);
}),
hasFilters,
limits: {
used: schedulesCount,
limit: project.organization.maximumSchedulesLimit,
},
filters: {
tasks,
environments,
@@ -10,10 +10,16 @@ import type { Organization } from "~/models/organization.server";
import type { Project } from "~/models/project.server";
import { displayableEnvironment } from "~/models/runtimeEnvironment.server";
import type { User } from "~/models/user.server";
import { filterOrphanedEnvironments, sortEnvironments } from "~/utils/environmentSort";
import {
filterOrphanedEnvironments,
onlyDevEnvironments,
exceptDevEnvironments,
sortEnvironments,
} from "~/utils/environmentSort";
import { logger } from "~/services/logger.server";
import { BasePresenter } from "./basePresenter.server";
import { TaskRunStatus } from "~/database-types";
import { CURRENT_DEPLOYMENT_LABEL } from "~/consts";
export type Task = {
slug: string;
@@ -72,6 +78,9 @@ export class TaskListPresenter extends BasePresenter {
},
});
const devEnvironments = onlyDevEnvironments(project.environments);
const nonDevEnvironments = exceptDevEnvironments(project.environments);
const tasks = await this._replica.$queryRaw<
{
id: string;
@@ -83,12 +92,21 @@ export class TaskListPresenter extends BasePresenter {
triggerSource: TaskTriggerSource;
}[]
>`
WITH workers AS (
WITH non_dev_workers AS (
SELECT wd."workerId" AS id
FROM ${sqlDatabaseSchema}."WorkerDeploymentPromotion" wdp
INNER JOIN ${sqlDatabaseSchema}."WorkerDeployment" wd
ON wd.id = wdp."deploymentId"
WHERE wdp."environmentId" IN (${Prisma.join(nonDevEnvironments.map((e) => e.id))})
AND wdp."label" = ${CURRENT_DEPLOYMENT_LABEL}
),
workers AS (
SELECT DISTINCT ON ("runtimeEnvironmentId") id, "runtimeEnvironmentId", version
FROM ${sqlDatabaseSchema}."BackgroundWorker"
WHERE "runtimeEnvironmentId" IN (${Prisma.join(
filterOrphanedEnvironments(project.environments).map((e) => e.id)
filterOrphanedEnvironments(devEnvironments).map((e) => e.id)
)})
OR id IN (SELECT id FROM non_dev_workers)
ORDER BY "runtimeEnvironmentId", "createdAt" DESC
)
SELECT tasks.id, slug, "filePath", "exportName", "triggerSource", tasks."runtimeEnvironmentId", tasks."createdAt"
@@ -293,7 +311,7 @@ export class TaskListPresenter extends BasePresenter {
>`
SELECT
tr."taskIdentifier",
AVG(EXTRACT(EPOCH FROM (tr."updatedAt" - tr."lockedAt"))) as duration
AVG(EXTRACT(EPOCH FROM (tr."updatedAt" - COALESCE(tr."startedAt", tr."lockedAt")))) as duration
FROM
${sqlDatabaseSchema}."TaskRun" as tr
WHERE
@@ -3,7 +3,8 @@ import { sqlDatabaseSchema, PrismaClient, prisma } from "~/db.server";
import { TestSearchParams } from "~/routes/_app.orgs.$organizationSlug.projects.v3.$projectParam.test/route";
import { sortEnvironments } from "~/utils/environmentSort";
import { createSearchParams } from "~/utils/searchParams";
import { getUsername } from "~/utils/username";
import { findCurrentWorkerDeployment } from "~/v3/models/workerDeployment.server";
import { BasePresenter } from "./basePresenter.server";
type TaskListOptions = {
userId: string;
@@ -15,16 +16,10 @@ export type TaskList = Awaited<ReturnType<TestPresenter["call"]>>;
export type TaskListItem = NonNullable<TaskList["tasks"]>[0];
export type SelectedEnvironment = NonNullable<TaskList["selectedEnvironment"]>;
export class TestPresenter {
#prismaClient: PrismaClient;
constructor(prismaClient: PrismaClient = prisma) {
this.#prismaClient = prismaClient;
}
export class TestPresenter extends BasePresenter {
public async call({ userId, projectSlug, url }: TaskListOptions) {
// Find the project scoped to the organization
const project = await this.#prismaClient.project.findFirstOrThrow({
const project = await this._replica.project.findFirstOrThrow({
select: {
id: true,
environments: {
@@ -85,31 +80,8 @@ export class TestPresenter {
};
}
//get all possible tasks
const tasks = await this.#prismaClient.$queryRaw<
{
id: string;
version: string;
taskIdentifier: string;
filePath: string;
exportName: string;
friendlyId: string;
triggerSource: TaskTriggerSource;
}[]
>`WITH workers AS (
SELECT
bw.*,
ROW_NUMBER() OVER(ORDER BY string_to_array(bw.version, '.')::int[] DESC) AS rn
FROM
${sqlDatabaseSchema}."BackgroundWorker" bw
WHERE "runtimeEnvironmentId" = ${matchingEnvironment.id}
),
latest_workers AS (SELECT * FROM workers WHERE rn = 1)
SELECT bwt.id, version, slug as "taskIdentifier", "filePath", "exportName", bwt."friendlyId", bwt."triggerSource"
FROM latest_workers
JOIN ${sqlDatabaseSchema}."BackgroundWorkerTask" bwt ON bwt."workerId" = latest_workers.id
ORDER BY bwt."exportName" ASC;
`;
const isDev = matchingEnvironment.type === "DEVELOPMENT";
const tasks = await this.#getTasks(matchingEnvironment.id, isDev);
return {
hasSelectedEnvironment: true as const,
@@ -118,8 +90,7 @@ export class TestPresenter {
tasks: tasks.map((task) => {
return {
id: task.id,
version: task.version,
taskIdentifier: task.taskIdentifier,
taskIdentifier: task.slug,
filePath: task.filePath,
exportName: task.exportName,
friendlyId: task.friendlyId,
@@ -128,4 +99,35 @@ export class TestPresenter {
}),
};
}
async #getTasks(envId: string, isDev: boolean) {
if (isDev) {
return await this._replica.$queryRaw<
{
id: string;
version: string;
slug: string;
filePath: string;
exportName: string;
friendlyId: string;
triggerSource: TaskTriggerSource;
}[]
>`WITH workers AS (
SELECT
bw.*,
ROW_NUMBER() OVER(ORDER BY string_to_array(bw.version, '.')::int[] DESC) AS rn
FROM
${sqlDatabaseSchema}."BackgroundWorker" bw
WHERE "runtimeEnvironmentId" = ${envId}
),
latest_workers AS (SELECT * FROM workers WHERE rn = 1)
SELECT bwt.id, version, slug, "filePath", "exportName", bwt."friendlyId", bwt."triggerSource"
FROM latest_workers
JOIN ${sqlDatabaseSchema}."BackgroundWorkerTask" bwt ON bwt."workerId" = latest_workers.id
ORDER BY bwt."exportName" ASC;`;
} else {
const currentDeployment = await findCurrentWorkerDeployment(envId);
return currentDeployment?.worker?.tasks ?? [];
}
}
}
@@ -6,6 +6,7 @@ import {
TaskTriggerSource,
} from "@trigger.dev/database";
import { sqlDatabaseSchema, PrismaClient, prisma } from "~/db.server";
import { getTimezones } from "~/utils/timezones.server";
import { getUsername } from "~/utils/username";
type TestTaskOptions = {
@@ -37,6 +38,7 @@ export type TestTask =
| {
triggerSource: "SCHEDULED";
task: Task;
possibleTimezones: string[];
runs: ScheduledRun[];
};
@@ -61,6 +63,7 @@ export type ScheduledRun = Omit<RawRun, "number" | "payload"> & {
timestamp: Date;
lastTimestamp?: Date;
externalId?: string;
timezone: string;
};
};
@@ -168,9 +171,11 @@ export class TestTaskPresenter {
),
};
case "SCHEDULED":
const possibleTimezones = getTimezones();
return {
triggerSource: "SCHEDULED",
task: taskWithEnvironment,
possibleTimezones,
runs: (
await Promise.all(
latestRuns.map(async (r) => {
@@ -195,6 +200,9 @@ export class TestTaskPresenter {
async function getScheduleTaskRunPayload(run: RawRun) {
const payload = await parsePacket({ data: run.payload, dataType: run.payloadType });
if (!payload.timezone) {
payload.timezone = "UTC";
}
const parsed = ScheduledTaskPayload.safeParse(payload);
return parsed;
}
@@ -1,8 +1,8 @@
import { ScheduleObject } from "@trigger.dev/core/v3";
import { PrismaClient, prisma } from "~/db.server";
import { displayableEnvironment } from "~/models/runtimeEnvironment.server";
import { nextScheduledTimestamps } from "~/v3/utils/calculateNextSchedule.server";
import { RunListPresenter } from "./RunListPresenter.server";
import { ScheduleObject } from "@trigger.dev/core/v3";
import { displayableEnvironment } from "~/models/runtimeEnvironment.server";
type ViewScheduleOptions = {
userId?: string;
@@ -24,6 +24,7 @@ export class ViewSchedulePresenter {
friendlyId: true,
generatorExpression: true,
generatorDescription: true,
timezone: true,
externalId: true,
deduplicationKey: true,
userProvidedDeduplicationKey: true,
@@ -68,7 +69,7 @@ export class ViewSchedulePresenter {
}
const nextRuns = schedule.active
? nextScheduledTimestamps(schedule.generatorExpression, new Date(), 5)
? nextScheduledTimestamps(schedule.generatorExpression, schedule.timezone, new Date(), 5)
: [];
const runPresenter = new RunListPresenter(this.#prismaClient);
@@ -82,6 +83,7 @@ export class ViewSchedulePresenter {
return {
schedule: {
...schedule,
timezone: schedule.timezone,
cron: schedule.generatorExpression,
cronDescription: schedule.generatorDescription,
nextRuns,
@@ -105,6 +107,7 @@ export class ViewSchedulePresenter {
expression: result.schedule.cron,
description: result.schedule.cronDescription,
},
timezone: result.schedule.timezone,
externalId: result.schedule.externalId ?? undefined,
deduplicationKey: result.schedule.userProvidedDeduplicationKey
? result.schedule.deduplicationKey ?? undefined
@@ -1,7 +1,6 @@
import { CommandLineIcon, ServerIcon } from "@heroicons/react/20/solid";
import { Outlet, useParams } from "@remix-run/react";
import { ArrowPathIcon, CommandLineIcon, ServerIcon } from "@heroicons/react/20/solid";
import { Outlet, useLocation, useParams } from "@remix-run/react";
import { LoaderFunctionArgs } from "@remix-run/server-runtime";
import { TerminalIcon, TerminalSquareIcon } from "lucide-react";
import { typedjson, useTypedLoaderData } from "remix-typedjson";
import { z } from "zod";
import { BlankstateInstructions } from "~/components/BlankstateInstructions";
@@ -9,8 +8,9 @@ import { UserAvatar } from "~/components/UserProfilePhoto";
import { EnvironmentLabel } from "~/components/environments/EnvironmentLabel";
import { MainCenteredContainer, PageBody, PageContainer } from "~/components/layout/AppLayout";
import { Badge } from "~/components/primitives/Badge";
import { LinkButton } from "~/components/primitives/Buttons";
import { Button, LinkButton } from "~/components/primitives/Buttons";
import { DateTime } from "~/components/primitives/DateTime";
import { Dialog, DialogTrigger } from "~/components/primitives/Dialog";
import { NavBar, PageTitle } from "~/components/primitives/PageHeader";
import { PaginationControls } from "~/components/primitives/Pagination";
import { Paragraph } from "~/components/primitives/Paragraph";
@@ -24,24 +24,26 @@ import {
TableBlankRow,
TableBody,
TableCell,
TableCellChevron,
TableCellMenu,
TableHeader,
TableHeaderCell,
TableRow,
} from "~/components/primitives/Table";
import { TextLink } from "~/components/primitives/TextLink";
import { DeploymentStatus } from "~/components/runs/v3/DeploymentStatus";
import { RollbackDeploymentDialog } from "~/components/runs/v3/RollbackDeploymentDialog";
import { useOrganization } from "~/hooks/useOrganizations";
import { useProject } from "~/hooks/useProject";
import { useUser } from "~/hooks/useUser";
import { DeploymentListPresenter } from "~/presenters/v3/DeploymentListPresenter.server";
import {
DeploymentListItem,
DeploymentListPresenter,
} from "~/presenters/v3/DeploymentListPresenter.server";
import { requireUserId } from "~/services/session.server";
import { cn } from "~/utils/cn";
import {
ProjectParamSchema,
docsPath,
v3DeploymentPath,
v3DeploymentsPath,
v3EnvironmentVariablesPath,
} from "~/utils/pathBuilder";
import { createSearchParams } from "~/utils/searchParams";
@@ -166,7 +168,7 @@ export default function Page() {
""
)}
</TableCell>
<TableCellChevron to={path} />
<DeploymentActionsCell deployment={deployment} path={path} />
</TableRow>
);
})
@@ -240,3 +242,35 @@ function CreateDeploymentInstructions() {
</MainCenteredContainer>
);
}
function DeploymentActionsCell({
deployment,
path,
}: {
deployment: DeploymentListItem;
path: string;
}) {
const location = useLocation();
const project = useProject();
if (deployment.isCurrent || !deployment.isDeployed) return <TableCell to={path}>{""}</TableCell>;
return (
<TableCellMenu isSticky>
{!deployment.isCurrent && deployment.isDeployed && (
<Dialog>
<DialogTrigger asChild>
<Button variant="small-menu-item" LeadingIcon={ArrowPathIcon}>
Rollback
</Button>
</DialogTrigger>
<RollbackDeploymentDialog
projectId={project.id}
deploymentShortCode={deployment.shortCode}
redirectPath={`${location.pathname}${location.search}`}
/>
</Dialog>
)}
</TableCellMenu>
);
}
@@ -185,7 +185,8 @@ export default function Page() {
const location = useLocation();
const organization = useOrganization();
const project = useProject();
const user = useUser();
const isUtc = schedule.timezone === "UTC";
return (
<div className="grid h-full max-h-full grid-rows-[2.5rem_1fr_3.25rem] overflow-hidden bg-background-bright">
@@ -210,6 +211,7 @@ export default function Page() {
<Paragraph variant="small">{schedule.cronDescription}</Paragraph>
</div>
</Property>
<Property label="Timezone">{schedule.timezone}</Property>
<Property label="Environments">
<EnvironmentLabels size="small" environments={schedule.environments} />
</Property>
@@ -245,19 +247,21 @@ export default function Page() {
<Table>
<TableHeader>
<TableRow>
{!isUtc && <TableHeaderCell>{schedule.timezone}</TableHeaderCell>}
<TableHeaderCell>UTC</TableHeaderCell>
<TableHeaderCell>Local time</TableHeaderCell>
</TableRow>
</TableHeader>
<TableBody>
{schedule.nextRuns.map((run, index) => (
<TableRow key={index}>
{!isUtc && (
<TableCell>
<DateTime date={run} timeZone={schedule.timezone} />
</TableCell>
)}
<TableCell>
<DateTime date={run} timeZone="UTC" />
</TableCell>
<TableCell>
<DateTime date={run} />
</TableCell>
</TableRow>
))}
</TableBody>
@@ -21,7 +21,7 @@ export const loader = async ({ request, params }: LoaderFunctionArgs) => {
};
export default function Page() {
const { schedule, possibleTasks, possibleEnvironments, showGenerateField } =
const { schedule, possibleTasks, possibleEnvironments, possibleTimezones, showGenerateField } =
useTypedLoaderData<typeof loader>();
return (
@@ -29,6 +29,7 @@ export default function Page() {
schedule={schedule}
possibleTasks={possibleTasks}
possibleEnvironments={possibleEnvironments}
possibleTimezones={possibleTimezones}
showGenerateField={showGenerateField}
/>
);
@@ -20,7 +20,7 @@ export const loader = async ({ request, params }: LoaderFunctionArgs) => {
};
export default function Page() {
const { schedule, possibleTasks, possibleEnvironments, showGenerateField } =
const { schedule, possibleTasks, possibleEnvironments, possibleTimezones, showGenerateField } =
useTypedLoaderData<typeof loader>();
return (
@@ -29,6 +29,7 @@ export default function Page() {
possibleTasks={possibleTasks}
possibleEnvironments={possibleEnvironments}
showGenerateField={showGenerateField}
possibleTimezones={possibleTimezones}
/>
);
}
@@ -4,12 +4,21 @@ import { Outlet, useLocation, useParams } from "@remix-run/react";
import { LoaderFunctionArgs } from "@remix-run/server-runtime";
import { typedjson, useTypedLoaderData } from "remix-typedjson";
import { BlankstateInstructions } from "~/components/BlankstateInstructions";
import { Feedback } from "~/components/Feedback";
import { AdminDebugTooltip } from "~/components/admin/debugTooltip";
import { InlineCode } from "~/components/code/InlineCode";
import { EnvironmentLabel, EnvironmentLabels } from "~/components/environments/EnvironmentLabel";
import { EnvironmentLabels } from "~/components/environments/EnvironmentLabel";
import { MainCenteredContainer, PageBody, PageContainer } from "~/components/layout/AppLayout";
import { LinkButton } from "~/components/primitives/Buttons";
import { Button, LinkButton } from "~/components/primitives/Buttons";
import { DateTime } from "~/components/primitives/DateTime";
import {
Dialog,
DialogContent,
DialogDescription,
DialogFooter,
DialogHeader,
DialogTrigger,
} from "~/components/primitives/Dialog";
import { NavBar, PageAccessories, PageTitle } from "~/components/primitives/PageHeader";
import { PaginationControls } from "~/components/primitives/Pagination";
import { Paragraph } from "~/components/primitives/Paragraph";
@@ -78,6 +87,7 @@ export default function Page() {
possibleEnvironments,
hasFilters,
filters,
limits,
currentPage,
totalPages,
} = useTypedLoaderData<typeof loader>();
@@ -107,15 +117,43 @@ export default function Page() {
</PropertyTable>
</AdminDebugTooltip>
<LinkButton
LeadingIcon={PlusIcon}
to={`${v3NewSchedulePath(organization, project)}${location.search}`}
variant="primary/small"
shortcut={{ key: "n" }}
disabled={possibleTasks.length === 0 || isShowingNewPane}
>
New schedule
</LinkButton>
{limits.used >= limits.limit ? (
<Dialog>
<DialogTrigger asChild>
<Button
LeadingIcon={PlusIcon}
variant="primary/small"
shortcut={{ key: "n" }}
disabled={possibleTasks.length === 0 || isShowingNewPane}
>
New schedule
</Button>
</DialogTrigger>
<DialogContent>
<DialogHeader>You've exceeded your limit</DialogHeader>
<DialogDescription>
You've used {limits.used}/{limits.limit} of your schedules. You can request more
schedules.
</DialogDescription>
<DialogFooter>
<Feedback
button={<Button variant="primary/medium">Request more</Button>}
defaultValue="help"
/>
</DialogFooter>
</DialogContent>
</Dialog>
) : (
<LinkButton
LeadingIcon={PlusIcon}
to={`${v3NewSchedulePath(organization, project)}${location.search}`}
variant="primary/small"
shortcut={{ key: "n" }}
disabled={possibleTasks.length === 0 || isShowingNewPane}
>
New schedule
</LinkButton>
)}
</PageAccessories>
</NavBar>
<PageBody scrollable={false}>
@@ -142,7 +180,21 @@ export default function Page() {
</div>
<SchedulesTable schedules={schedules} hasFilters={hasFilters} />
<div className="mt-2 justify-end">
<div className="mt-2 justify-between">
<Paragraph variant="extra-small" className="mt-3">
<span className={limits.used >= limits.limit ? "text-warning" : ""}>
You've used {limits.used}/{limits.limit} of your schedules.
</span>{" "}
<Feedback
button={
<button className=" text-secondary transition hover:text-indigo-400">
Request more
</button>
}
defaultValue="help"
/>
.
</Paragraph>
<PaginationControls currentPage={currentPage} totalPages={totalPages} />
</div>
</div>
@@ -236,12 +288,13 @@ function SchedulesTable({
<TableRow>
<TableHeaderCell>ID</TableHeaderCell>
<TableHeaderCell>Task ID</TableHeaderCell>
<TableHeaderCell>External ID</TableHeaderCell>
<TableHeaderCell>CRON</TableHeaderCell>
<TableHeaderCell hiddenLabel>CRON description</TableHeaderCell>
<TableHeaderCell>External ID</TableHeaderCell>
<TableHeaderCell>Timezone</TableHeaderCell>
<TableHeaderCell>Next run</TableHeaderCell>
<TableHeaderCell>Last run</TableHeaderCell>
<TableHeaderCell>Deduplication key</TableHeaderCell>
<TableHeaderCell>Next run (UTC)</TableHeaderCell>
<TableHeaderCell>Last run (UTC)</TableHeaderCell>
<TableHeaderCell>Environments</TableHeaderCell>
<TableHeaderCell>Enabled</TableHeaderCell>
</TableRow>
@@ -262,6 +315,9 @@ function SchedulesTable({
<TableCell to={path} className={cellClass}>
{schedule.taskIdentifier}
</TableCell>
<TableCell to={path} className={cellClass}>
{schedule.externalId ? schedule.externalId : ""}
</TableCell>
<TableCell to={path} className={cellClass}>
{schedule.cron}
</TableCell>
@@ -269,17 +325,21 @@ function SchedulesTable({
{schedule.cronDescription}
</TableCell>
<TableCell to={path} className={cellClass}>
{schedule.externalId ? schedule.externalId : ""}
{schedule.timezone}
</TableCell>
<TableCell to={path} className={cellClass}>
<DateTime date={schedule.nextRun} timeZone={schedule.timezone} />
</TableCell>
<TableCell to={path} className={cellClass}>
{schedule.lastRun ? (
<DateTime date={schedule.lastRun} timeZone={schedule.timezone} />
) : (
""
)}
</TableCell>
<TableCell to={path} className={cellClass}>
{schedule.userProvidedDeduplicationKey ? schedule.deduplicationKey : ""}
</TableCell>
<TableCell to={path} className={cellClass}>
<DateTime date={schedule.nextRun} timeZone="utc" />
</TableCell>
<TableCell to={path} className={cellClass}>
{schedule.lastRun ? <DateTime date={schedule.lastRun} timeZone="utc" /> : ""}
</TableCell>
<TableCell to={path} className={cellClass}>
<EnvironmentLabels environments={schedule.environments} size="small" />
</TableCell>
@@ -26,8 +26,10 @@ import {
ResizablePanel,
ResizablePanelGroup,
} from "~/components/primitives/Resizable";
import { Select } from "~/components/primitives/Select";
import { TextLink } from "~/components/primitives/TextLink";
import { TaskRunStatusCombo } from "~/components/runs/v3/TaskRunStatus";
import { TimezoneList } from "~/components/scheduled/timezones";
import { redirectBackWithErrorMessage, redirectWithSuccessMessage } from "~/models/message.server";
import {
ScheduledRun,
@@ -95,7 +97,13 @@ export default function Page() {
return <StandardTaskForm task={result.task} runs={result.runs} />;
}
case "SCHEDULED": {
return <ScheduledTaskForm task={result.task} runs={result.runs} />;
return (
<ScheduledTaskForm
task={result.task}
runs={result.runs}
possibleTimezones={result.possibleTimezones}
/>
);
}
}
}
@@ -215,12 +223,21 @@ function StandardTaskForm({ task, runs }: { task: TestTask["task"]; runs: Standa
);
}
function ScheduledTaskForm({ task, runs }: { task: TestTask["task"]; runs: ScheduledRun[] }) {
function ScheduledTaskForm({
task,
runs,
possibleTimezones,
}: {
task: TestTask["task"];
runs: ScheduledRun[];
possibleTimezones: string[];
}) {
const lastSubmission = useActionData();
const [selectedCodeSampleId, setSelectedCodeSampleId] = useState(runs.at(0)?.id);
const [timestampValue, setTimestampValue] = useState<Date | undefined>();
const [lastTimestampValue, setLastTimestampValue] = useState<Date | undefined>();
const [externalIdValue, setExternalIdValue] = useState<string | undefined>();
const [timezoneValue, setTimezoneValue] = useState<string>("UTC");
//set initial values
useEffect(() => {
@@ -233,11 +250,20 @@ function ScheduledTaskForm({ task, runs }: { task: TestTask["task"]; runs: Sched
setTimestampValue(initialRun.payload.timestamp);
setLastTimestampValue(initialRun.payload.lastTimestamp);
setExternalIdValue(initialRun.payload.externalId);
setTimezoneValue(initialRun.payload.timezone);
}, [selectedCodeSampleId]);
const [
form,
{ timestamp, lastTimestamp, externalId, triggerSource, taskIdentifier, environmentId },
{
timestamp,
lastTimestamp,
externalId,
triggerSource,
taskIdentifier,
environmentId,
timezone,
},
] = useForm({
id: "test-task-scheduled",
// TODO: type this
@@ -314,6 +340,30 @@ function ScheduledTaskForm({ task, runs }: { task: TestTask["task"]; runs: Sched
</Hint>
<FormError id={lastTimestamp.errorId}>{lastTimestamp.error}</FormError>
</InputGroup>
<InputGroup>
<Label htmlFor={timezone.id}>Timezone</Label>
<Select
{...conform.select(timezone)}
placeholder="Select a timezone"
defaultValue={timezoneValue}
value={timezoneValue}
setValue={(e) => {
if (Array.isArray(e)) return;
setTimezoneValue(e);
}}
items={possibleTimezones}
filter={{ keys: [(item) => item.replace(/\//g, " ").replace(/_/g, " ")] }}
dropdownIcon
variant="tertiary/medium"
>
{(matches) => <TimezoneList timezones={matches} />}
</Select>
<Hint>
The Timestamp and Last timestamp are in UTC so this just changes the timezone
string that comes through in the payload.
</Hint>
<FormError id={timezone.errorId}>{timezone.error}</FormError>
</InputGroup>
<InputGroup>
<Label required={false} htmlFor={externalId.id}>
External ID
@@ -0,0 +1,35 @@
import { Link } from "@remix-run/react";
import type { LoaderFunctionArgs } from "@remix-run/server-runtime";
import { typedjson, useTypedLoaderData } from "remix-typedjson";
import { LogoIcon } from "~/components/LogoIcon";
import { Header1 } from "~/components/primitives/Headers";
import { Paragraph } from "~/components/primitives/Paragraph";
import { getTimezones } from "~/utils/timezones.server";
export const loader = async ({ request }: LoaderFunctionArgs) => {
return typedjson({
timezones: getTimezones(),
});
};
export default function Page() {
const { timezones } = useTypedLoaderData<typeof loader>();
return (
<div className="grid grid-rows-[2.5rem,1fr]">
<div className="flex items-center border-b border-b-grid-dimmed px-3">
<Link to="/">
<LogoIcon className="relative -top-px mr-2 h-4 w-4 min-w-[1rem]" />
</Link>
</div>
<div className="overflow-y-auto p-8 scrollbar-thin scrollbar-track-transparent scrollbar-thumb-charcoal-600">
<Header1 spacing>Supported timezones</Header1>
<Paragraph spacing>We support these timezones when creating a schedule.</Paragraph>
<ul className="">
{timezones.map((timezone) => (
<li key={timezone}>{timezone}</li>
))}
</ul>
</div>
</div>
);
}
+9 -8
View File
@@ -1,5 +1,8 @@
import { ActionFunctionArgs, json } from "@remix-run/server-runtime";
import { InitializeDeploymentRequestBody, InitializeDeploymentResponseBody } from "@trigger.dev/core/v3";
import {
InitializeDeploymentRequestBody,
InitializeDeploymentResponseBody,
} from "@trigger.dev/core/v3";
import { env } from "~/env.server";
import { authenticateApiRequest } from "~/services/apiAuth.server";
import { logger } from "~/services/logger.server";
@@ -37,13 +40,11 @@ export async function action({ request, params }: ActionFunctionArgs) {
contentHash: deployment.contentHash,
shortCode: deployment.shortCode,
version: deployment.version,
externalBuildData: deployment.externalBuildData as InitializeDeploymentResponseBody["externalBuildData"],
externalBuildData:
deployment.externalBuildData as InitializeDeploymentResponseBody["externalBuildData"],
imageTag,
registryHost: env.DEPLOY_REGISTRY_HOST
}
registryHost: env.DEPLOY_REGISTRY_HOST,
};
return json(
responseBody,
{ status: 200 }
);
return json(responseBody, { status: 200 });
}
@@ -123,7 +123,7 @@ export async function loader({ params, request }: LoaderFunctionArgs) {
const repository = new EnvironmentVariablesRepository();
const variables = await repository.getEnvironment(environment.project.id, environment.id, true);
const variables = await repository.getEnvironment(environment.project.id, environment.id);
const environmentVariable = variables.find((v) => v.key === parsedParams.data.name);
@@ -80,7 +80,7 @@ export async function loader({ params, request }: LoaderFunctionArgs) {
const repository = new EnvironmentVariablesRepository();
const variables = await repository.getEnvironment(environment.project.id, environment.id, true);
const variables = await repository.getEnvironment(environment.project.id, environment.id);
return json(variables.map((variable) => ({ name: variable.key, value: variable.value })));
}
@@ -2,7 +2,7 @@ import { LoaderFunctionArgs, json } from "@remix-run/server-runtime";
import { z } from "zod";
import { prisma } from "~/db.server";
import { authenticateApiRequest } from "~/services/apiAuth.server";
import { EnvironmentVariablesRepository } from "~/v3/environmentVariables/environmentVariablesRepository.server";
import { resolveVariablesForEnvironment } from "~/v3/environmentVariables/environmentVariablesRepository.server";
const ParamsSchema = z.object({
projectRef: z.string(),
@@ -41,9 +41,7 @@ export async function loader({ request, params }: LoaderFunctionArgs) {
return json({ error: "Project not found" }, { status: 404 });
}
const repository = new EnvironmentVariablesRepository();
const variables = await repository.getEnvironmentVariables(project.id, authenticatedEnv.id);
const variables = await resolveVariablesForEnvironment(authenticatedEnv);
return json({
variables: variables.reduce((acc: Record<string, string>, variable) => {
@@ -79,9 +79,9 @@ export async function action({ request, params }: ActionFunctionArgs) {
friendlyId: parsedParams.data.scheduleId,
taskIdentifier: body.data.task,
cron: body.data.cron,
timezone: body.data.timezone,
environments: [authenticationResult.environment.id],
externalId: body.data.externalId,
deduplicationKey: body.data.deduplicationKey,
};
const schedule = await service.call(authenticationResult.environment.projectId, options);
@@ -95,6 +95,7 @@ export async function action({ request, params }: ActionFunctionArgs) {
expression: schedule.cron,
description: schedule.cronDescription,
},
timezone: schedule.timezone,
externalId: schedule.externalId ?? undefined,
deduplicationKey: schedule.deduplicationKey,
environments: schedule.environments,
@@ -43,6 +43,7 @@ export async function action({ request }: ActionFunctionArgs) {
environments: [authenticationResult.environment.id],
externalId: body.data.externalId,
deduplicationKey: body.data.deduplicationKey,
timezone: body.data.timezone,
};
const schedule = await service.call(authenticationResult.environment.projectId, options);
@@ -56,6 +57,7 @@ export async function action({ request }: ActionFunctionArgs) {
expression: schedule.cron,
description: schedule.cronDescription,
},
timezone: schedule.timezone,
externalId: schedule.externalId ?? undefined,
deduplicationKey: schedule.deduplicationKey,
environments: schedule.environments,
@@ -111,6 +113,7 @@ export async function loader({ request }: LoaderFunctionArgs) {
expression: schedule.cron,
description: schedule.cronDescription,
},
timezone: schedule.timezone,
deduplicationKey: schedule.userProvidedDeduplicationKey
? schedule.deduplicationKey
: undefined,
@@ -0,0 +1,28 @@
import type { LoaderFunctionArgs } from "@remix-run/server-runtime";
import { json } from "@remix-run/server-runtime";
import { z } from "zod";
import { apiCors } from "~/utils/apiCors";
import { getTimezones } from "~/utils/timezones.server";
const SearchParamsSchema = z.object({
excludeUtc: z.preprocess((value) => value === "true", z.boolean()).default(false),
});
export async function loader({ request }: LoaderFunctionArgs) {
if (request.method.toUpperCase() === "OPTIONS") {
return apiCors(request, json({}));
}
const rawSearchParams = new URL(request.url).searchParams;
const params = SearchParamsSchema.safeParse(Object.fromEntries(rawSearchParams.entries()));
if (!params.success) {
return apiCors(
request,
json({ error: "Invalid request parameters", issues: params.error.issues }, { status: 400 })
);
}
const timezones = getTimezones(!params.data.excludeUtc);
return apiCors(request, json({ timezones }));
}
@@ -0,0 +1,97 @@
import { ActionFunctionArgs } from "@remix-run/server-runtime";
import { MachinePresetName } from "@trigger.dev/core/v3";
import { z } from "zod";
import { prisma } from "~/db.server";
import { validateJWTTokenAndRenew } from "~/services/apiAuth.server";
import { logger } from "~/services/logger.server";
import { workerQueue } from "~/services/worker.server";
import { machinePresetFromName } from "~/v3/machinePresets.server";
import { reportUsageEvent } from "~/v3/openMeter.server";
const JWTPayloadSchema = z.object({
environment_id: z.string(),
org_id: z.string(),
project_id: z.string(),
run_id: z.string(),
machine_preset: z.string(),
});
const BodySchema = z.object({
durationMs: z.number(),
});
export async function action({ request }: ActionFunctionArgs) {
// Ensure this is a POST request
if (request.method.toUpperCase() !== "POST") {
return { status: 405, body: "Method Not Allowed" };
}
const jwtResult = await validateJWTTokenAndRenew(request, JWTPayloadSchema);
if (!jwtResult) {
return { status: 401, body: "Unauthorized" };
}
const rawJson = await request.json();
const json = BodySchema.safeParse(rawJson);
if (!json.success) {
logger.error("Failed to parse request body", { rawJson });
return { status: 400, body: "Bad Request" };
}
const preset = machinePresetFromName(jwtResult.payload.machine_preset as MachinePresetName);
logger.debug("[/api/v1/usage/ingest] Reporting usage", { jwtResult, json: json.data, preset });
if (json.data.durationMs > 0) {
const costInCents = json.data.durationMs * preset.centsPerMs;
await prisma.taskRun.update({
where: {
id: jwtResult.payload.run_id,
},
data: {
usageDurationMs: {
increment: json.data.durationMs,
},
costInCents: {
increment: json.data.durationMs * preset.centsPerMs,
},
},
});
try {
await reportUsageEvent({
source: "webapp",
type: "usage",
subject: jwtResult.payload.org_id,
data: {
durationMs: json.data.durationMs,
costInCents: String(costInCents),
},
});
} catch (e) {
logger.error("Failed to report usage event, enqueing v3.reportUsage", { error: e });
await workerQueue.enqueue("v3.reportUsage", {
orgId: jwtResult.payload.org_id,
data: {
costInCents: String(costInCents),
},
additionalData: {
durationMs: json.data.durationMs,
},
});
}
}
return new Response(null, {
status: 200,
headers: {
"x-trigger-jwt": jwtResult.jwt,
},
});
}
@@ -2,7 +2,7 @@ import type { LoaderFunctionArgs } from "@remix-run/server-runtime";
import { json } from "@remix-run/server-runtime";
import { GetEvent } from "@trigger.dev/core";
import { z } from "zod";
import { prisma } from "~/db.server";
import { $replica } from "~/db.server";
import { authenticateApiRequest } from "~/services/apiAuth.server";
import { apiCors } from "~/utils/apiCors";
@@ -59,7 +59,7 @@ function toJSON(eventRecord: FoundEventRecord): GetEvent {
type FoundEventRecord = NonNullable<Awaited<ReturnType<typeof findEventRecord>>>;
async function findEventRecord(eventId: string, environmentId: string) {
return await prisma.eventRecord.findUnique({
return await $replica.eventRecord.findUnique({
select: {
eventId: true,
name: true,
@@ -0,0 +1,93 @@
import { parse } from "@conform-to/zod";
import { ActionFunction, json } from "@remix-run/node";
import { z } from "zod";
import { prisma } from "~/db.server";
import { redirectWithErrorMessage, redirectWithSuccessMessage } from "~/models/message.server";
import { logger } from "~/services/logger.server";
import { requireUserId } from "~/services/session.server";
import { RollbackDeploymentService } from "~/v3/services/rollbackDeployment.server";
export const rollbackSchema = z.object({
redirectUrl: z.string(),
});
const ParamSchema = z.object({
projectId: z.string(),
deploymentShortCode: z.string(),
});
export const action: ActionFunction = async ({ request, params }) => {
const userId = await requireUserId(request);
const { projectId, deploymentShortCode } = ParamSchema.parse(params);
console.log("projectId", projectId);
console.log("deploymentShortCode", deploymentShortCode);
const formData = await request.formData();
const submission = parse(formData, { schema: rollbackSchema });
if (!submission.value) {
return json(submission);
}
try {
const project = await prisma.project.findUnique({
where: {
id: projectId,
organization: {
members: {
some: {
userId,
},
},
},
},
});
if (!project) {
return redirectWithErrorMessage(submission.value.redirectUrl, request, "Project not found");
}
const deployment = await prisma.workerDeployment.findUnique({
where: {
projectId_shortCode: {
projectId: project.id,
shortCode: deploymentShortCode,
},
},
});
if (!deployment) {
return redirectWithErrorMessage(
submission.value.redirectUrl,
request,
"Deployment not found"
);
}
const rollbackService = new RollbackDeploymentService();
await rollbackService.call(deployment);
return redirectWithSuccessMessage(
submission.value.redirectUrl,
request,
"Rolled back deployment"
);
} catch (error) {
if (error instanceof Error) {
logger.error("Failed to roll back deployment", {
error: {
name: error.name,
message: error.message,
stack: error.stack,
},
});
submission.error = { runParam: error.message };
return json(submission);
} else {
logger.error("Failed to roll back deployment", { error });
submission.error = { runParam: JSON.stringify(error) };
return json(submission);
}
}
};
@@ -3,9 +3,10 @@ import { parse } from "@conform-to/zod";
import { CheckIcon, XMarkIcon } from "@heroicons/react/20/solid";
import { Form, useActionData, useLocation, useNavigation } from "@remix-run/react";
import { ActionFunctionArgs, json } from "@remix-run/server-runtime";
import { useVirtualizer } from "@tanstack/react-virtual";
import { parseExpression } from "cron-parser";
import cronstrue from "cronstrue";
import { useState } from "react";
import { useRef, useState } from "react";
import {
environmentTextClassName,
environmentTitle,
@@ -42,6 +43,7 @@ import { ProjectParamSchema, docsPath, v3SchedulesPath } from "~/utils/pathBuild
import { CronPattern, UpsertSchedule } from "~/v3/schedules";
import { UpsertTaskScheduleService } from "~/v3/services/upsertTaskSchedule.server";
import { AIGeneratedCronField } from "../resources.orgs.$organizationSlug.projects.$projectParam.schedules.new.natural-language";
import { TimezoneList } from "~/components/scheduled/timezones";
const cronFormat = `* * * * *
┬ ┬ ┬ ┬ ┬
@@ -117,9 +119,12 @@ export function UpsertScheduleForm({
schedule,
possibleTasks,
possibleEnvironments,
possibleTimezones,
showGenerateField,
}: EditableScheduleElements & { showGenerateField: boolean }) {
const lastSubmission = useActionData();
const [selectedTimezone, setSelectedTimezone] = useState<string>(schedule?.timezone ?? "UTC");
const isUtc = selectedTimezone === "UTC";
const [cronPattern, setCronPattern] = useState<string>(schedule?.cron ?? "");
const navigation = useNavigation();
const isLoading = navigation.state !== "idle";
@@ -127,18 +132,20 @@ export function UpsertScheduleForm({
const project = useProject();
const location = useLocation();
const [form, { taskIdentifier, cron, externalId, environments, deduplicationKey }] = useForm({
id: "create-schedule",
// TODO: type this
lastSubmission: lastSubmission as any,
shouldRevalidate: "onSubmit",
onValidate({ formData }) {
return parse(formData, { schema: UpsertSchedule });
},
});
const [form, { taskIdentifier, cron, timezone, externalId, environments, deduplicationKey }] =
useForm({
id: "create-schedule",
// TODO: type this
lastSubmission: lastSubmission as any,
shouldRevalidate: "onSubmit",
onValidate({ formData }) {
return parse(formData, { schema: UpsertSchedule });
},
});
let cronPatternResult: CronPatternResult | undefined = undefined;
let nextRuns: Date[] | undefined = undefined;
if (cronPattern !== "") {
const result = CronPattern.safeParse(cronPattern);
@@ -149,7 +156,10 @@ export function UpsertScheduleForm({
};
} else {
try {
const expression = parseExpression(cronPattern, { utc: true });
const expression = parseExpression(
cronPattern,
isUtc ? { utc: true } : { tz: selectedTimezone }
);
cronPatternResult = {
isValid: true,
description: cronstrue.toString(cronPattern),
@@ -195,6 +205,7 @@ export function UpsertScheduleForm({
items={possibleTasks}
filter={(task, search) => task.toLowerCase().includes(search.toLowerCase())}
dropdownIcon
variant="tertiary/medium"
>
{(matches) => (
<>
@@ -241,25 +252,52 @@ export function UpsertScheduleForm({
<ValidCronMessage isValid={false} message={cronPatternResult.error} />
)}
</InputGroup>
<InputGroup>
<Label htmlFor={timezone.id}>Timezone</Label>
<Select
{...conform.select(timezone)}
placeholder="Select a timezone"
defaultValue={selectedTimezone}
value={selectedTimezone}
setValue={(e) => {
if (Array.isArray(e)) return;
setSelectedTimezone(e);
}}
items={possibleTimezones}
filter={{ keys: [(item) => item.replace(/\//g, " ").replace(/_/g, " ")] }}
dropdownIcon
variant="tertiary/medium"
>
{(matches) => <TimezoneList timezones={matches} />}
</Select>
<Hint>
{isUtc
? "UTC will not change with daylight savings time."
: "This will automatically adjust for daylight savings time."}
</Hint>
<FormError id={timezone.errorId}>{timezone.error}</FormError>
</InputGroup>
{nextRuns !== undefined && (
<div className="flex flex-col gap-1">
<Header3>Next 5 runs</Header3>
<Table>
<TableHeader>
<TableRow>
{!isUtc && <TableHeaderCell>{selectedTimezone}</TableHeaderCell>}
<TableHeaderCell>UTC</TableHeaderCell>
<TableHeaderCell>Local time</TableHeaderCell>
</TableRow>
</TableHeader>
<TableBody>
{nextRuns.map((run, index) => (
<TableRow key={index}>
{!isUtc && (
<TableCell>
<DateTime date={run} timeZone={selectedTimezone} />
</TableCell>
)}
<TableCell>
<DateTime date={run} timeZone="UTC" />
</TableCell>
<TableCell>
<DateTime date={run} />
</TableCell>
</TableRow>
))}
</TableBody>
+128
View File
@@ -1,6 +1,7 @@
import { Prettify } from "@trigger.dev/core";
import { z } from "zod";
import {
RuntimeEnvironment,
findEnvironmentByApiKey,
findEnvironmentByPublicApiKey,
} from "~/models/runtimeEnvironment.server";
@@ -12,6 +13,9 @@ import {
import { prisma } from "~/db.server";
import { json } from "@remix-run/server-runtime";
import { findProjectByRef } from "~/models/project.server";
import { SignJWT, jwtVerify, errors } from "jose";
import { env } from "~/env.server";
import { logger } from "./logger.server";
type Optional<T, K extends keyof T> = Prettify<Omit<T, K> & Partial<Pick<T, K>>>;
@@ -209,3 +213,127 @@ export async function authenticatedEnvironmentForAuthentication(
}
}
}
const JWT_SECRET = new TextEncoder().encode(env.SESSION_SECRET);
const JWT_ALGORITHM = "HS256";
const DEFAULT_JWT_EXPIRATION_IN_MS = 1000 * 60 * 60; // 1 hour
export async function generateJWTTokenForEnvironment(
environment: RuntimeEnvironment,
payload: Record<string, string>
) {
const jwt = await new SignJWT({
environment_id: environment.id,
org_id: environment.organizationId,
project_id: environment.projectId,
...payload,
})
.setProtectedHeader({ alg: JWT_ALGORITHM })
.setIssuedAt()
.setIssuer("https://id.trigger.dev")
.setAudience("https://api.trigger.dev")
.setExpirationTime(calculateJWTExpiration())
.sign(JWT_SECRET);
return jwt;
}
export async function validateJWTTokenAndRenew<T extends z.ZodTypeAny>(
request: Request,
payloadSchema: T
): Promise<{ payload: z.infer<T>; jwt: string } | undefined> {
try {
const jwt = request.headers.get("x-trigger-jwt");
if (!jwt) {
logger.debug("Missing JWT token in request", {
headers: Object.fromEntries(request.headers),
});
return;
}
const { payload: rawPayload } = await jwtVerify(jwt, JWT_SECRET, {
issuer: "https://id.trigger.dev",
audience: "https://api.trigger.dev",
});
const payload = payloadSchema.safeParse(rawPayload);
if (!payload.success) {
logger.error("Failed to validate JWT", { payload: rawPayload, issues: payload.error.issues });
return;
}
const renewedJwt = await renewJWTToken(payload.data);
return {
payload: payload.data,
jwt: renewedJwt,
};
} catch (error) {
if (error instanceof errors.JWTExpired) {
// Now we need to try and renew the token using the API key auth
const authenticatedEnv = await authenticateApiRequest(request);
if (!authenticatedEnv) {
logger.error("Failed to renew JWT token, missing or invalid Authorization header", {
error: error.message,
});
return;
}
const payload = payloadSchema.safeParse(error.payload);
if (!payload.success) {
logger.error("Failed to parse jwt payload after expired", {
payload: error.payload,
issues: payload.error.issues,
});
return;
}
const renewedJwt = await generateJWTTokenForEnvironment(authenticatedEnv.environment, {
...payload.data,
});
logger.debug("Renewed JWT token from Authorization header API Key", {
environment: authenticatedEnv.environment,
payload: payload.data,
});
return {
payload: payload.data,
jwt: renewedJwt,
};
}
logger.error("Failed to validate JWT token", { error });
}
}
async function renewJWTToken(payload: Record<string, string>) {
const jwt = await new SignJWT(payload)
.setProtectedHeader({ alg: JWT_ALGORITHM })
.setIssuedAt()
.setIssuer("https://id.trigger.dev")
.setAudience("https://api.trigger.dev")
.setExpirationTime(calculateJWTExpiration())
.sign(JWT_SECRET);
return jwt;
}
function calculateJWTExpiration() {
if (env.PROD_USAGE_HEARTBEAT_INTERVAL_MS) {
return (
(Date.now() + Math.max(DEFAULT_JWT_EXPIRATION_IN_MS, env.PROD_USAGE_HEARTBEAT_INTERVAL_MS)) /
1000
);
}
return (Date.now() + DEFAULT_JWT_EXPIRATION_IN_MS) / 1000;
}
@@ -145,11 +145,13 @@ export const apiRateLimiter = authorizationRateLimitMiddleware({
"/api/internal/stripe_webhooks",
"/api/v1/authorization-code",
"/api/v1/token",
"/api/v1/usage/ingest",
/^\/api\/v1\/tasks\/[^\/]+\/callback\/[^\/]+$/, // /api/v1/tasks/$id/callback/$secret
/^\/api\/v1\/runs\/[^\/]+\/tasks\/[^\/]+\/callback\/[^\/]+$/, // /api/v1/runs/$runId/tasks/$id/callback/$secret
/^\/api\/v1\/http-endpoints\/[^\/]+\/env\/[^\/]+\/[^\/]+$/, // /api/v1/http-endpoints/$httpEndpointId/env/$envType/$shortcode
/^\/api\/v1\/sources\/http\/[^\/]+$/, // /api/v1/sources/http/$id
/^\/api\/v1\/endpoints\/[^\/]+\/[^\/]+\/index\/[^\/]+$/, // /api/v1/endpoints/$environmentId/$endpointSlug/index/$indexHookIdentifier
"/api/v1/timezones",
],
log: {
rejections: env.API_RATE_LIMIT_REJECTION_LOGS_ENABLED === "1",
+5 -3
View File
@@ -1,5 +1,5 @@
import { BillingClient, SetPlanBody } from "@trigger.dev/billing";
import { PrismaClient, prisma } from "~/db.server";
import { $replica, PrismaClient, PrismaReplicaClient, prisma } from "~/db.server";
import { env } from "~/env.server";
import { logger } from "~/services/logger.server";
import { organizationBillingPath } from "~/utils/pathBuilder";
@@ -7,9 +7,11 @@ import { organizationBillingPath } from "~/utils/pathBuilder";
export class BillingService {
#billingClient: BillingClient | undefined;
#prismaClient: PrismaClient;
#replica: PrismaReplicaClient;
constructor(isManagedCloud: boolean, prismaClient: PrismaClient = prisma) {
constructor(isManagedCloud: boolean, prismaClient: PrismaClient = prisma, replica: PrismaReplicaClient = $replica) {
this.#prismaClient = prismaClient;
this.#replica = replica;
if (isManagedCloud && process.env.BILLING_API_URL && process.env.BILLING_API_KEY) {
this.#billingClient = new BillingClient({
url: process.env.BILLING_API_URL,
@@ -35,7 +37,7 @@ export class BillingService {
firstDayOfNextMonth.setMonth(firstDayOfNextMonth.getMonth() + 1);
firstDayOfNextMonth.setHours(0, 0, 0, 0);
const currentRunCount = await this.#prismaClient.jobRun.count({
const currentRunCount = await this.#replica.jobRun.count({
where: {
organizationId: orgId,
createdAt: {
@@ -6,6 +6,12 @@ import { $transaction, PrismaClientOrTransaction, prisma } from "~/db.server";
import { logger } from "~/services/logger.server";
import { workerQueue } from "../worker.server";
class AlreadyDeliveredError extends Error {
constructor() {
super("Event already delivered");
}
}
export class DeliverEventService {
#prismaClient: PrismaClientOrTransaction;
@@ -14,81 +20,111 @@ export class DeliverEventService {
}
public async call(id: string) {
await $transaction(
this.#prismaClient,
async (tx) => {
const eventRecord = await tx.eventRecord.findUniqueOrThrow({
where: {
id,
},
include: {
environment: {
include: {
organization: true,
project: true,
try {
await $transaction(
this.#prismaClient,
async (tx) => {
const eventRecord = await tx.eventRecord.findUniqueOrThrow({
where: {
id,
},
include: {
environment: {
include: {
organization: true,
project: true,
},
},
},
},
});
});
const possibleEventDispatchers = await tx.eventDispatcher.findMany({
where: {
environmentId: eventRecord.environmentId,
event: {
has: eventRecord.name,
if (eventRecord.deliveredAt) {
logger.debug("Event already delivered", {
eventRecord: eventRecord.id,
});
return;
}
const possibleEventDispatchers = await tx.eventDispatcher.findMany({
where: {
environmentId: eventRecord.environmentId,
event: {
has: eventRecord.name,
},
source: eventRecord.source,
enabled: true,
manual: false,
},
source: eventRecord.source,
enabled: true,
manual: false,
},
});
});
logger.debug("Found possible event dispatchers", {
possibleEventDispatchers,
eventRecord: eventRecord.id,
});
const matchingEventDispatchers = possibleEventDispatchers.filter((eventDispatcher) =>
this.#evaluateEventRule(eventDispatcher, eventRecord)
);
if (matchingEventDispatchers.length === 0) {
logger.debug("No matching event dispatchers", {
logger.debug("Found possible event dispatchers", {
possibleEventDispatchers,
eventRecord: eventRecord.id,
});
return;
}
const matchingEventDispatchers = possibleEventDispatchers.filter((eventDispatcher) =>
this.#evaluateEventRule(eventDispatcher, eventRecord)
);
logger.debug("Found matching event dispatchers", {
matchingEventDispatchers,
eventRecord: eventRecord.id,
});
if (matchingEventDispatchers.length === 0) {
logger.debug("No matching event dispatchers", {
eventRecord: eventRecord.id,
});
await Promise.all(
matchingEventDispatchers.map((eventDispatcher) =>
workerQueue.enqueue(
"events.invokeDispatcher",
{
id: eventDispatcher.id,
eventRecordId: eventRecord.id,
},
{ tx }
return;
}
logger.debug("Found matching event dispatchers", {
matchingEventDispatchers,
eventRecord: eventRecord.id,
});
await Promise.all(
matchingEventDispatchers.map((eventDispatcher) =>
workerQueue.enqueue(
"events.invokeDispatcher",
{
id: eventDispatcher.id,
eventRecordId: eventRecord.id,
},
{ tx }
)
)
)
);
);
await tx.eventRecord.update({
where: {
id: eventRecord.id,
},
data: {
deliveredAt: new Date(),
},
// Optimistically mark the event as delivered
const lockedRecord = await tx.eventRecord.updateMany({
where: {
id: eventRecord.id,
deliveredAt: null,
},
data: {
deliveredAt: new Date(),
},
});
if (lockedRecord.count === 0) {
//this means we've already delivered it, because there were no records with deliveredAt = null
//by throwing it will rollback the transaction, stopping the queue from processing the event again
throw new AlreadyDeliveredError();
}
},
{ timeout: 10000 }
);
}
catch (error) {
if (error instanceof AlreadyDeliveredError) {
logger.debug("Event already delivered, AlreadyDeliveredError", {
eventRecord: id,
});
},
{ timeout: 10000 }
);
//we swallow the error because we don't want to retry
return;
}
throw error;
}
}
#evaluateEventRule(dispatcher: EventDispatcher, eventRecord: EventRecord): boolean {
@@ -110,6 +110,14 @@ export class IngestSendEvent {
},
});
if (existingEventLog?.deliveredAt) {
logger.debug("Event already delivered", {
eventRecordId: existingEventLog.id,
deliveredAt: existingEventLog.deliveredAt,
});
return existingEventLog;
}
const eventLog = await (existingEventLog
? this.updateEvent({ tx, existingEventLog, reqEvent: event, deliverAt })
: this.createEvent({
@@ -127,6 +135,15 @@ export class IngestSendEvent {
if (!createdEvent) return;
if (createdEvent.deliveredAt) {
logger.debug("Event already delivered", {
eventRecordId: createdEvent.id,
deliveredAt: createdEvent.deliveredAt,
});
//return the event if it was already delivered, don't enqueue it again
return createdEvent;
}
//rate limit
const result = await rateLimiter?.limit(environment.organizationId);
if (result && !result.success) {
@@ -46,7 +46,7 @@ import { forceYieldCoordinator } from "./forceYieldCoordinator.server";
import { ResumeRunService } from "./resumeRun.server";
type FoundRun = NonNullable<Awaited<ReturnType<typeof findRun>>>;
type FoundTask = FoundRun["tasks"][number];
type FoundTask = NonNullable<Awaited<ReturnType<typeof getCompletedTasksForRun>>>[number];
// We need to limit the cached tasks to not be too large >3.5MB when serialized
const TOTAL_CACHED_TASK_BYTE_LIMIT = 3500000;
@@ -81,6 +81,14 @@ export class PerformRunExecutionV3Service {
}
public async call(input: PerformRunExecutionV3Input, driftInMs: number = 0) {
logger.debug("PerformRunExecutionV3Service.call", { input, driftInMs });
if (Array.isArray(input.id)) {
logger.error("PerformRunExecutionV3Service.call: input.id is an array", { input });
throw new Error("input.id must be a string");
}
const run = await findRun(this.#prismaClient, input.id);
if (!run) {
@@ -221,11 +229,14 @@ export class PerformRunExecutionV3Service {
});
}
const taskCount = await getTaskCountForRun(this.#prismaClient, run.id);
const tasks = await getCompletedTasksForRun(this.#prismaClient, run.id);
const sourceContext = RunSourceContextSchema.safeParse(run.event.sourceContext);
const executionBody = await this.#createExecutionBody(
run,
run.tasks,
tasks,
startedAt,
false,
connections.auth,
@@ -418,7 +429,8 @@ export class PerformRunExecutionV3Service {
this.#prismaClient,
run,
input,
durationInMs
durationInMs,
taskCount
);
} else {
return await this.#failRunExecutionWithRetry(
@@ -1065,7 +1077,8 @@ export class PerformRunExecutionV3Service {
prisma: PrismaClientOrTransaction,
run: FoundRun,
input: PerformRunExecutionV3Input,
durationInMs: number
durationInMs: number,
existingTaskCount: number
) {
await $transaction(prisma, async (tx) => {
const executionDuration = run.executionDuration + durationInMs;
@@ -1086,31 +1099,25 @@ export class PerformRunExecutionV3Service {
return;
}
const runWithLatestTask = await tx.jobRun.findUniqueOrThrow({
where: {
id: run.id,
},
select: {
tasks: {
select: {
id: true,
name: true,
status: true,
displayKey: true,
},
take: 1,
orderBy: { createdAt: "desc" },
},
_count: {
select: {
tasks: true,
},
},
},
});
const newTaskCount = await getTaskCountForRun(tx, run.id);
if (runWithLatestTask._count.tasks === run._count.tasks) {
const latestTask = runWithLatestTask.tasks[0];
if (newTaskCount === existingTaskCount) {
const latestTask = await tx.task.findFirst({
select: {
id: true,
name: true,
status: true,
displayKey: true,
},
where: {
runId: run.id,
status: "RUNNING",
},
orderBy: {
createdAt: "desc",
},
take: 1,
});
const cause =
latestTask?.status === "RUNNING"
@@ -1254,6 +1261,35 @@ function prepareNoOpTasksBloomFilter(possibleTasks: FoundTask[]): string {
return filter.serialize();
}
async function getTaskCountForRun(prisma: PrismaClientOrTransaction, runId: string) {
return await prisma.task.count({
where: {
runId,
},
});
}
async function getCompletedTasksForRun(prisma: PrismaClientOrTransaction, runId: string) {
return await prisma.task.findMany({
where: {
runId,
status: "COMPLETED",
},
select: {
id: true,
idempotencyKey: true,
status: true,
noop: true,
output: true,
outputIsUndefined: true,
parentId: true,
},
orderBy: {
id: "asc",
},
});
}
async function findRun(prisma: PrismaClientOrTransaction, id: string) {
return await prisma.jobRun.findUnique({
where: { id },
@@ -1278,25 +1314,6 @@ async function findRun(prisma: PrismaClientOrTransaction, id: string) {
},
},
},
tasks: {
where: {
status: {
in: ["COMPLETED"],
},
},
select: {
id: true,
idempotencyKey: true,
status: true,
noop: true,
output: true,
outputIsUndefined: true,
parentId: true,
},
orderBy: {
id: "asc",
},
},
event: true,
version: {
include: {
@@ -1309,11 +1326,6 @@ async function findRun(prisma: PrismaClientOrTransaction, id: string) {
recipientMethod: "ENDPOINT",
},
},
_count: {
select: {
tasks: true,
},
},
},
});
}
@@ -22,6 +22,11 @@ class Telemetry {
#triggerClient: TriggerClient | undefined = undefined;
constructor({ postHogApiKey, trigger }: Options) {
if (env.TRIGGER_TELEMETRY_DISABLED !== undefined) {
console.log("📉 Telemetry disabled");
return;
}
if (postHogApiKey) {
this.#posthogClient = new PostHog(postHogApiKey, { host: "https://eu.posthog.com" });
} else {
+37
View File
@@ -45,6 +45,8 @@ import { PerformTaskOperationService } from "./tasks/performTaskOperation.server
import { ProcessCallbackTimeoutService } from "./tasks/processCallbackTimeout.server";
import { ResumeTaskService } from "./tasks/resumeTask.server";
import { RequeueV2Message } from "~/v3/marqs/requeueV2Message.server";
import { MarqsConcurrencyMonitor } from "~/v3/marqs/concurrencyMonitor.server";
import { reportUsageEvent } from "~/v3/openMeter.server";
const workerCatalog = {
indexEndpoint: z.object({
@@ -168,6 +170,13 @@ const workerCatalog = {
"v2.requeueMessage": z.object({
runId: z.string(),
}),
"v3.reportUsage": z.object({
orgId: z.string(),
data: z.object({
costInCents: z.string(),
}),
additionalData: z.record(z.any()).optional(),
}),
};
const executionWorkerCatalog = {
@@ -298,6 +307,19 @@ function getWorkerQueue() {
await eventRepository.truncateEvents();
},
},
"marqs.v3.queueConcurrencyMonitor": {
// run every 5 minutes
match: "*/5 * * * *",
handler: async (payload, job, helpers) => {
await MarqsConcurrencyMonitor.initiateV3Monitoring(helpers.abortSignal);
},
},
"marqs.v2.queueConcurrencyMonitor": {
match: "*/5 * * * *", // run every 5 minutes
handler: async (payload, job, helpers) => {
await MarqsConcurrencyMonitor.initiateV2Monitoring(helpers.abortSignal);
},
},
},
tasks: {
"events.invokeDispatcher": {
@@ -635,6 +657,21 @@ function getWorkerQueue() {
await service.call(payload.runId);
},
},
"v3.reportUsage": {
priority: 0,
maxAttempts: 8,
handler: async (payload, job) => {
await reportUsageEvent({
source: "webapp",
type: "usage",
subject: payload.orgId,
data: {
costInCents: payload.data.costInCents,
...payload.additionalData,
},
});
},
},
},
});
}
+8
View File
@@ -58,3 +58,11 @@ export function filterOrphanedEnvironments<T extends FilterableEnvironment>(
return false;
});
}
export function onlyDevEnvironments<T extends FilterableEnvironment>(environments: T[]): T[] {
return environments.filter((e) => e.type === "DEVELOPMENT");
}
export function exceptDevEnvironments<T extends FilterableEnvironment>(environments: T[]): T[] {
return environments.filter((e) => e.type !== "DEVELOPMENT");
}
@@ -0,0 +1,7 @@
export function getTimezones(includeUtc = true) {
const possibleTimezones = Intl.supportedValuesOf("timeZone").sort();
if (includeUtc) {
possibleTimezones.unshift("UTC");
}
return possibleTimezones;
}
@@ -1,4 +1,9 @@
import { Prisma, PrismaClient, RuntimeEnvironmentType } from "@trigger.dev/database";
import {
Prisma,
PrismaClient,
RuntimeEnvironment,
RuntimeEnvironmentType,
} from "@trigger.dev/database";
import { z } from "zod";
import { environmentTitle } from "~/components/environments/EnvironmentLabel";
import { $transaction, prisma } from "~/db.server";
@@ -427,11 +432,7 @@ export class EnvironmentVariablesRepository implements Repository {
return results;
}
async getEnvironment(
projectId: string,
environmentId: string,
excludeInternalVariables?: boolean
): Promise<EnvironmentVariable[]> {
async getEnvironment(projectId: string, environmentId: string): Promise<EnvironmentVariable[]> {
const project = await this.prismaClient.project.findUnique({
where: {
id: projectId,
@@ -453,124 +454,7 @@ export class EnvironmentVariablesRepository implements Repository {
return [];
}
return this.getEnvironmentVariables(projectId, environmentId, excludeInternalVariables);
}
async #getTriggerEnvironmentVariables(environmentId: string): Promise<EnvironmentVariable[]> {
const environment = await this.prismaClient.runtimeEnvironment.findFirst({
where: {
id: environmentId,
},
});
if (!environment) {
return [];
}
if (environment.type === "DEVELOPMENT") {
return [
{
key: "OTEL_EXPORTER_OTLP_ENDPOINT",
value: env.DEV_OTEL_EXPORTER_OTLP_ENDPOINT ?? env.APP_ORIGIN,
},
].concat(
env.DEV_OTEL_BATCH_PROCESSING_ENABLED === "1"
? [
{
key: "OTEL_BATCH_PROCESSING_ENABLED",
value: "1",
},
{
key: "OTEL_SPAN_MAX_EXPORT_BATCH_SIZE",
value: env.DEV_OTEL_SPAN_MAX_EXPORT_BATCH_SIZE,
},
{
key: "OTEL_SPAN_SCHEDULED_DELAY_MILLIS",
value: env.DEV_OTEL_SPAN_SCHEDULED_DELAY_MILLIS,
},
{
key: "OTEL_SPAN_EXPORT_TIMEOUT_MILLIS",
value: env.DEV_OTEL_SPAN_EXPORT_TIMEOUT_MILLIS,
},
{
key: "OTEL_SPAN_MAX_QUEUE_SIZE",
value: env.DEV_OTEL_SPAN_MAX_QUEUE_SIZE,
},
{
key: "OTEL_LOG_MAX_EXPORT_BATCH_SIZE",
value: env.DEV_OTEL_LOG_MAX_EXPORT_BATCH_SIZE,
},
{
key: "OTEL_LOG_SCHEDULED_DELAY_MILLIS",
value: env.DEV_OTEL_LOG_SCHEDULED_DELAY_MILLIS,
},
{
key: "OTEL_LOG_EXPORT_TIMEOUT_MILLIS",
value: env.DEV_OTEL_LOG_EXPORT_TIMEOUT_MILLIS,
},
{
key: "OTEL_LOG_MAX_QUEUE_SIZE",
value: env.DEV_OTEL_LOG_MAX_QUEUE_SIZE,
},
]
: []
);
}
return [
{
key: "TRIGGER_SECRET_KEY",
value: environment.apiKey,
},
{
key: "TRIGGER_API_URL",
value: env.APP_ORIGIN,
},
{
key: "TRIGGER_RUNTIME_WAIT_THRESHOLD_IN_MS",
value: String(env.RUNTIME_WAIT_THRESHOLD_IN_MS),
},
...(env.PROD_OTEL_BATCH_PROCESSING_ENABLED === "1"
? [
{
key: "OTEL_BATCH_PROCESSING_ENABLED",
value: "1",
},
{
key: "OTEL_SPAN_MAX_EXPORT_BATCH_SIZE",
value: env.PROD_OTEL_SPAN_MAX_EXPORT_BATCH_SIZE,
},
{
key: "OTEL_SPAN_SCHEDULED_DELAY_MILLIS",
value: env.PROD_OTEL_SPAN_SCHEDULED_DELAY_MILLIS,
},
{
key: "OTEL_SPAN_EXPORT_TIMEOUT_MILLIS",
value: env.PROD_OTEL_SPAN_EXPORT_TIMEOUT_MILLIS,
},
{
key: "OTEL_SPAN_MAX_QUEUE_SIZE",
value: env.PROD_OTEL_SPAN_MAX_QUEUE_SIZE,
},
{
key: "OTEL_LOG_MAX_EXPORT_BATCH_SIZE",
value: env.PROD_OTEL_LOG_MAX_EXPORT_BATCH_SIZE,
},
{
key: "OTEL_LOG_SCHEDULED_DELAY_MILLIS",
value: env.PROD_OTEL_LOG_SCHEDULED_DELAY_MILLIS,
},
{
key: "OTEL_LOG_EXPORT_TIMEOUT_MILLIS",
value: env.PROD_OTEL_LOG_EXPORT_TIMEOUT_MILLIS,
},
{
key: "OTEL_LOG_MAX_QUEUE_SIZE",
value: env.PROD_OTEL_LOG_MAX_QUEUE_SIZE,
},
]
: []),
];
return this.getEnvironmentVariables(projectId, environmentId);
}
async #getSecretEnvironmentVariables(
@@ -597,18 +481,9 @@ export class EnvironmentVariablesRepository implements Repository {
async getEnvironmentVariables(
projectId: string,
environmentId: string,
excludeInternalVariables?: boolean
environmentId: string
): Promise<EnvironmentVariable[]> {
const secretEnvVars = await this.#getSecretEnvironmentVariables(projectId, environmentId);
if (excludeInternalVariables) {
return secretEnvVars;
}
const triggerEnvVars = await this.#getTriggerEnvironmentVariables(environmentId);
return [...secretEnvVars, ...triggerEnvVars];
return this.#getSecretEnvironmentVariables(projectId, environmentId);
}
async delete(projectId: string, options: DeleteEnvironmentVariable): Promise<Result> {
@@ -782,3 +657,158 @@ export class EnvironmentVariablesRepository implements Repository {
}
}
}
export const environmentVariablesRepository = new EnvironmentVariablesRepository();
export async function resolveVariablesForEnvironment(runtimeEnvironment: RuntimeEnvironment) {
const projectSecrets = await environmentVariablesRepository.getEnvironmentVariables(
runtimeEnvironment.projectId,
runtimeEnvironment.id
);
const builtInVariables =
runtimeEnvironment.type === "DEVELOPMENT"
? await resolveBuiltInDevVariables(runtimeEnvironment)
: await resolveBuiltInProdVariables(runtimeEnvironment);
return [...projectSecrets, ...builtInVariables];
}
async function resolveBuiltInDevVariables(runtimeEnvironment: RuntimeEnvironment) {
let result: Array<EnvironmentVariable> = [
{
key: "OTEL_EXPORTER_OTLP_ENDPOINT",
value: env.DEV_OTEL_EXPORTER_OTLP_ENDPOINT ?? env.APP_ORIGIN,
},
];
if (env.DEV_OTEL_BATCH_PROCESSING_ENABLED === "1") {
result = result.concat([
{
key: "OTEL_BATCH_PROCESSING_ENABLED",
value: "1",
},
{
key: "OTEL_SPAN_MAX_EXPORT_BATCH_SIZE",
value: env.DEV_OTEL_SPAN_MAX_EXPORT_BATCH_SIZE,
},
{
key: "OTEL_SPAN_SCHEDULED_DELAY_MILLIS",
value: env.DEV_OTEL_SPAN_SCHEDULED_DELAY_MILLIS,
},
{
key: "OTEL_SPAN_EXPORT_TIMEOUT_MILLIS",
value: env.DEV_OTEL_SPAN_EXPORT_TIMEOUT_MILLIS,
},
{
key: "OTEL_SPAN_MAX_QUEUE_SIZE",
value: env.DEV_OTEL_SPAN_MAX_QUEUE_SIZE,
},
{
key: "OTEL_LOG_MAX_EXPORT_BATCH_SIZE",
value: env.DEV_OTEL_LOG_MAX_EXPORT_BATCH_SIZE,
},
{
key: "OTEL_LOG_SCHEDULED_DELAY_MILLIS",
value: env.DEV_OTEL_LOG_SCHEDULED_DELAY_MILLIS,
},
{
key: "OTEL_LOG_EXPORT_TIMEOUT_MILLIS",
value: env.DEV_OTEL_LOG_EXPORT_TIMEOUT_MILLIS,
},
{
key: "OTEL_LOG_MAX_QUEUE_SIZE",
value: env.DEV_OTEL_LOG_MAX_QUEUE_SIZE,
},
]);
}
const commonVariables = await resolveCommonBuiltInVariables(runtimeEnvironment);
return [...result, ...commonVariables];
}
async function resolveBuiltInProdVariables(runtimeEnvironment: RuntimeEnvironment) {
let result: Array<EnvironmentVariable> = [
{
key: "TRIGGER_SECRET_KEY",
value: runtimeEnvironment.apiKey,
},
{
key: "TRIGGER_API_URL",
value: env.APP_ORIGIN,
},
{
key: "TRIGGER_RUNTIME_WAIT_THRESHOLD_IN_MS",
value: String(env.RUNTIME_WAIT_THRESHOLD_IN_MS),
},
{
key: "TRIGGER_ORG_ID",
value: runtimeEnvironment.organizationId,
},
];
if (env.PROD_OTEL_BATCH_PROCESSING_ENABLED === "1") {
result = result.concat([
{
key: "OTEL_BATCH_PROCESSING_ENABLED",
value: "1",
},
{
key: "OTEL_SPAN_MAX_EXPORT_BATCH_SIZE",
value: env.PROD_OTEL_SPAN_MAX_EXPORT_BATCH_SIZE,
},
{
key: "OTEL_SPAN_SCHEDULED_DELAY_MILLIS",
value: env.PROD_OTEL_SPAN_SCHEDULED_DELAY_MILLIS,
},
{
key: "OTEL_SPAN_EXPORT_TIMEOUT_MILLIS",
value: env.PROD_OTEL_SPAN_EXPORT_TIMEOUT_MILLIS,
},
{
key: "OTEL_SPAN_MAX_QUEUE_SIZE",
value: env.PROD_OTEL_SPAN_MAX_QUEUE_SIZE,
},
{
key: "OTEL_LOG_MAX_EXPORT_BATCH_SIZE",
value: env.PROD_OTEL_LOG_MAX_EXPORT_BATCH_SIZE,
},
{
key: "OTEL_LOG_SCHEDULED_DELAY_MILLIS",
value: env.PROD_OTEL_LOG_SCHEDULED_DELAY_MILLIS,
},
{
key: "OTEL_LOG_EXPORT_TIMEOUT_MILLIS",
value: env.PROD_OTEL_LOG_EXPORT_TIMEOUT_MILLIS,
},
{
key: "OTEL_LOG_MAX_QUEUE_SIZE",
value: env.PROD_OTEL_LOG_MAX_QUEUE_SIZE,
},
]);
}
if (env.PROD_USAGE_HEARTBEAT_INTERVAL_MS && env.USAGE_EVENT_URL) {
result = result.concat([
{
key: "USAGE_HEARTBEAT_INTERVAL_MS",
value: String(env.PROD_USAGE_HEARTBEAT_INTERVAL_MS),
},
{
key: "USAGE_EVENT_URL",
value: env.USAGE_EVENT_URL,
},
]);
}
const commonVariables = await resolveCommonBuiltInVariables(runtimeEnvironment);
return [...result, ...commonVariables];
}
async function resolveCommonBuiltInVariables(
runtimeEnvironment: RuntimeEnvironment
): Promise<Array<EnvironmentVariable>> {
return [];
}
@@ -76,16 +76,8 @@ export interface Repository {
create(projectId: string, options: CreateEnvironmentVariables): Promise<CreateResult>;
edit(projectId: string, options: EditEnvironmentVariable): Promise<Result>;
getProject(projectId: string): Promise<ProjectEnvironmentVariable[]>;
getEnvironment(
projectId: string,
environmentId: string,
excludeInternalVariables?: boolean
): Promise<EnvironmentVariable[]>;
getEnvironmentVariables(
projectId: string,
environmentId: string,
excludeInternalVariables?: boolean
): Promise<EnvironmentVariable[]>;
getEnvironment(projectId: string, environmentId: string): Promise<EnvironmentVariable[]>;
getEnvironmentVariables(projectId: string, environmentId: string): Promise<EnvironmentVariable[]>;
delete(projectId: string, options: DeleteEnvironmentVariable): Promise<Result>;
deleteValue(projectId: string, options: DeleteEnvironmentVariableValue): Promise<Result>;
}
+2 -2
View File
@@ -334,14 +334,14 @@ export class EventRepository {
}
async queryEvents(queryOptions: QueryOptions): Promise<TaskEventRecord[]> {
return await this.db.taskEvent.findMany({
return await this.readReplica.taskEvent.findMany({
where: queryOptions,
});
}
async queryIncompleteEvents(queryOptions: QueryOptions) {
// First we will find all the events that match the query options (selecting minimal data).
const taskEvents = await this.db.taskEvent.findMany({
const taskEvents = await this.readReplica.taskEvent.findMany({
where: queryOptions,
select: {
spanId: true,
@@ -0,0 +1,81 @@
import { MachineConfig, MachinePreset, MachinePresetName } from "@trigger.dev/core/v3";
import { env } from "~/env.server";
import { logger } from "~/services/logger.server";
export const presets = {
micro: {
cpu: 0.25,
memory: 0.25,
centsPerMs: env.CENTS_PER_HOUR_MICRO / 3_600_000,
},
"small-1x": {
cpu: 0.5,
memory: 0.5,
centsPerMs: env.CENTS_PER_HOUR_SMALL_1X / 3_600_000,
},
"small-2x": {
cpu: 1,
memory: 1,
centsPerMs: env.CENTS_PER_HOUR_SMALL_2X / 3_600_000,
},
"medium-1x": {
cpu: 1,
memory: 2,
centsPerMs: env.CENTS_PER_HOUR_MEDIUM_1X / 3_600_000,
},
"medium-2x": {
cpu: 2,
memory: 4,
centsPerMs: env.CENTS_PER_HOUR_MEDIUM_2X / 3_600_000,
},
"large-1x": {
cpu: 4,
memory: 8,
centsPerMs: env.CENTS_PER_HOUR_LARGE_1X / 3_600_000,
},
"large-2x": {
cpu: 8,
memory: 16,
centsPerMs: env.CENTS_PER_HOUR_LARGE_2X / 3_600_000,
},
};
export function machinePresetFromConfig(config: unknown): MachinePreset {
const parsedConfig = MachineConfig.safeParse(config);
if (!parsedConfig.success) {
logger.error("Failed to parse machine config", { config });
return machinePresetFromName("small-1x");
}
if (parsedConfig.data.preset) {
return machinePresetFromName(parsedConfig.data.preset);
}
if (parsedConfig.data.cpu && parsedConfig.data.memory) {
const name = derivePresetNameFromValues(parsedConfig.data.cpu, parsedConfig.data.memory);
return machinePresetFromName(name);
}
return machinePresetFromName("small-1x");
}
export function machinePresetFromName(name: MachinePresetName): MachinePreset {
return {
name,
...presets[name],
};
}
// Finds the smallest machine preset name that satisfies the given CPU and memory requirements
function derivePresetNameFromValues(cpu: number, memory: number): MachinePresetName {
for (const [name, preset] of Object.entries(presets)) {
if (preset.cpu >= cpu && preset.memory >= memory) {
return name as MachinePresetName;
}
}
return "small-1x";
}
@@ -0,0 +1,216 @@
import { Logger } from "@trigger.dev/core-backend";
import { Redis } from "ioredis";
import { prisma } from "~/db.server";
import { logger } from "~/services/logger.server";
import { MarQS, marqs as marqsv3 } from "./index.server";
import { env } from "~/env.server";
import { marqsv2 } from "./v2.server";
export type MarqsConcurrencyMonitorOptions = {
dryRun?: boolean;
abortSignal?: AbortSignal;
};
export interface MarqsConcurrencyResolveCompletedRunsCallback {
(candidateRunIds: string[]): Promise<Array<{ id: string }>>;
}
export class MarqsConcurrencyMonitor {
private _logger: Logger;
constructor(
private marqs: MarQS,
private callback: MarqsConcurrencyResolveCompletedRunsCallback,
private options: MarqsConcurrencyMonitorOptions = {}
) {
this._logger = logger.child({
component: "marqs",
operation: "concurrencyMonitor",
dryRun: this.dryRun,
marqs: marqs.name,
});
}
get dryRun() {
return typeof this.options.dryRun === "boolean" ? this.options.dryRun : false;
}
get keys() {
return this.marqs.keys;
}
get signal() {
return this.options.abortSignal;
}
public async call() {
this._logger.debug("[MarqsConcurrencyMonitor] Initiating monitoring");
const stats = {
streamCallbacks: 0,
processedKeys: 0,
};
const { stream, redis } = this.marqs.queueConcurrencyScanStream(10, () => {
this._logger.debug("[MarqsConcurrencyMonitor] stream closed", {
stats,
});
});
stream.on("data", async (keys) => {
stream.pause();
if (this.signal?.aborted) {
stream.destroy();
return;
}
stats.streamCallbacks++;
const uniqueKeys = Array.from(new Set<string>(keys));
if (uniqueKeys.length === 0) {
stream.resume();
return;
}
this._logger.debug("[MarqsConcurrencyMonitor] correcting queues concurrency", {
keys: uniqueKeys,
});
stats.processedKeys += uniqueKeys.length;
await Promise.all(uniqueKeys.map((key) => this.#processKey(key, redis))).finally(() => {
stream.resume();
});
});
}
async #processKey(key: string, redis: Redis) {
key = this.keys.stripKeyPrefix(key);
const orgKey = this.keys.orgCurrentConcurrencyKeyFromQueue(key);
const envKey = this.keys.envCurrentConcurrencyKeyFromQueue(key);
// Next, we need to get all the items from the key, and any parent keys (org, env, queue) using sunion.
const runIds = await redis.sunion(orgKey, envKey, key);
if (runIds.length === 0) {
return;
}
const perfNow = performance.now();
const completeRuns = await this.callback(runIds);
const durationMs = performance.now() - perfNow;
const completedRunIds = completeRuns.map((run) => run.id);
if (completedRunIds.length === 0) {
this._logger.debug("[MarqsConcurrencyMonitor] no completed runs found", {
key,
orgKey,
envKey,
runIds,
durationMs,
});
return;
}
this._logger.debug("[MarqsConcurrencyMonitor] removing completed runs from queue", {
key,
orgKey,
envKey,
completedRunIds,
durationMs,
});
if (this.dryRun) {
return;
}
const pipeline = redis.pipeline();
pipeline.srem(key, ...completedRunIds);
pipeline.srem(orgKey, ...completedRunIds);
pipeline.srem(envKey, ...completedRunIds);
try {
await pipeline.exec();
} catch (e) {
this._logger.error("[MarqsConcurrencyMonitor] error removing completed runs from queue", {
key,
orgKey,
envKey,
completedRunIds,
error: e,
});
}
}
static async initiateV3Monitoring(abortSignal?: AbortSignal) {
if (!marqsv3) {
return;
}
const instance = new MarqsConcurrencyMonitor(
marqsv3,
(runIds) =>
prisma.taskRun.findMany({
select: { id: true },
where: {
id: {
in: runIds,
},
status: {
in: [
"CANCELED",
"COMPLETED_SUCCESSFULLY",
"COMPLETED_WITH_ERRORS",
"CRASHED",
"SYSTEM_FAILURE",
"INTERRUPTED",
],
},
},
}),
{ dryRun: env.V3_MARQS_CONCURRENCY_MONITOR_ENABLED === "0", abortSignal }
);
await instance.call();
}
static async initiateV2Monitoring(abortSignal?: AbortSignal) {
if (!marqsv2) {
return;
}
const instance = new MarqsConcurrencyMonitor(
marqsv2,
(runIds) =>
prisma.jobRun.findMany({
select: { id: true },
where: {
id: {
in: runIds,
},
status: {
in: [
"CANCELED",
"SUCCESS",
"FAILURE",
"TIMED_OUT",
"ABORTED",
"CANCELED",
"INVALID_PAYLOAD",
],
},
},
}),
{ dryRun: env.V2_MARQS_CONCURRENCY_MONITOR_ENABLED === "0", abortSignal }
);
await instance.call();
}
}
@@ -15,7 +15,8 @@ import { createNewSession, disconnectSession } from "~/models/runtimeEnvironment
import { AuthenticatedEnvironment } from "~/services/apiAuth.server";
import { logger } from "~/services/logger.server";
import { marqs, sanitizeQueueName } from "~/v3/marqs/index.server";
import { EnvironmentVariablesRepository } from "../environmentVariables/environmentVariablesRepository.server";
import { resolveVariablesForEnvironment } from "../environmentVariables/environmentVariablesRepository.server";
import { FailedTaskRunService } from "../failedTaskRun.server";
import { CancelTaskRunService } from "../services/cancelTaskRun.server";
import { CompleteAttemptService } from "../services/completeAttempt.server";
import { CreateTaskRunAttemptService } from "../services/createTaskRunAttempt.server";
@@ -25,7 +26,6 @@ import {
tracer,
} from "../tracer.server";
import { DevSubscriber, devPubSub } from "./devPubSub.server";
import { FailedTaskRunService } from "../failedTaskRun.server";
const MessageBody = z.discriminatedUnion("type", [
z.object({
@@ -415,6 +415,7 @@ export class DevQueueConsumer {
lockedById: backgroundTask.id,
status: "EXECUTING",
lockedToVersionId: backgroundWorker.id,
startedAt: existingTaskRun.startedAt ?? new Date(),
},
include: {
attempts: {
@@ -473,11 +474,7 @@ export class DevQueueConsumer {
return;
}
const environmentRepository = new EnvironmentVariablesRepository();
const variables = await environmentRepository.getEnvironmentVariables(
this.env.project.id,
this.env.id
);
const variables = await resolveVariablesForEnvironment(this.env);
if (backgroundWorker.supportsLazyAttempts) {
const payload: TaskRunExecutionLazyAttemptPayload = {
@@ -524,6 +521,7 @@ export class DevQueueConsumer {
lockedAt: null,
lockedById: null,
status: "PENDING",
startedAt: existingTaskRun.startedAt,
},
}),
]);
@@ -581,6 +579,7 @@ export class DevQueueConsumer {
lockedAt: null,
lockedById: null,
status: "PENDING",
startedAt: existingTaskRun.startedAt,
},
}),
]);
+20 -55
View File
@@ -78,8 +78,6 @@ export class MarQS {
this.keys = options.keysProducer;
this.queuePriorityStrategy = options.queuePriorityStrategy;
// Spawn options.workers workers to requeue visible messages
this.#startRequeuingWorkers();
this.#startRebalanceWorkers();
this.#registerCommands();
}
@@ -790,64 +788,31 @@ export class MarQS {
}
}
#startRequeuingWorkers() {
// Start a new worker to requeue visible messages
for (let i = 0; i < this.options.workers; i++) {
const worker = new AsyncWorker(this.#requeueVisibleMessages.bind(this), 1000);
queueConcurrencyScanStream(count: number = 100, onEndCallback?: () => void) {
const pattern = this.keys.queueCurrentConcurrencyScanPattern();
this.#requeueingWorkers.push(worker);
logger.debug("Starting queue concurrency scan stream", {
pattern,
component: "marqs",
operation: "queueConcurrencyScanStream",
service: this.name,
count,
});
worker.start();
}
}
const redis = this.redis.duplicate();
async #requeueVisibleMessages() {
// Remove any of the messages from the timeoutQueue that have expired
const messages = await this.redis.zrangebyscore(
constants.MESSAGE_VISIBILITY_TIMEOUT_QUEUE,
0,
Date.now(),
"LIMIT",
0,
10
);
const stream = redis.scanStream({
match: pattern,
type: "set",
count,
});
if (messages.length === 0) {
return;
}
stream.on("end", () => {
onEndCallback?.();
redis.quit();
});
for (let i = 0; i < messages.length; i++) {
const message = messages[i];
const messageData = await this.redis.get(this.keys.messageKey(message));
if (!messageData) {
// The message has been removed for some reason (TTL, etc.), so we should remove it from the timeout queue
await this.redis.zrem(constants.MESSAGE_VISIBILITY_TIMEOUT_QUEUE, message);
continue;
}
const parsedMessage = MessagePayload.safeParse(JSON.parse(messageData));
if (!parsedMessage.success) {
await this.redis.zrem(constants.MESSAGE_VISIBILITY_TIMEOUT_QUEUE, message);
continue;
}
await this.#callNackMessage({
messageKey: this.keys.messageKey(message),
messageQueue: parsedMessage.data.queue,
parentQueue: parsedMessage.data.parentQueue,
concurrencyKey: this.keys.currentConcurrencyKeyFromQueue(parsedMessage.data.queue),
envConcurrencyKey: this.keys.envCurrentConcurrencyKeyFromQueue(parsedMessage.data.queue),
orgConcurrencyKey: this.keys.orgCurrentConcurrencyKeyFromQueue(parsedMessage.data.queue),
visibilityQueue: constants.MESSAGE_VISIBILITY_TIMEOUT_QUEUE,
messageId: parsedMessage.data.messageId,
messageScore: parsedMessage.data.timestamp,
});
}
return { stream, redis };
}
async #rebalanceParentQueues() {
@@ -13,12 +13,16 @@ const constants = {
} as const;
export class MarQSShortKeyProducer implements MarQSKeyProducer {
constructor(private _prefix: string) {}
constructor(private _prefix: string) { }
sharedQueueScanPattern() {
return `${this._prefix}*${constants.SHARED_QUEUE}`;
}
queueCurrentConcurrencyScanPattern() {
return `${this._prefix}${constants.ORG_PART}:*:${constants.ENV_PART}:*:queue:*:${constants.CURRENT_CONCURRENCY_PART}`;
}
stripKeyPrefix(key: string): string {
if (key.startsWith(this._prefix)) {
return key.slice(this._prefix.length);
@@ -1,6 +1,6 @@
import { Context, ROOT_CONTEXT, Span, SpanKind, context, trace } from "@opentelemetry/api";
import {
Machine,
MachinePreset,
ProdTaskRunExecution,
ProdTaskRunExecutionPayload,
TaskRunError,
@@ -15,15 +15,19 @@ import { ZodMessageSender } from "@trigger.dev/core/v3/zodMessageHandler";
import {
BackgroundWorker,
BackgroundWorkerTask,
RuntimeEnvironment,
TaskRun,
TaskRunAttemptStatus,
TaskRunStatus,
} from "@trigger.dev/database";
import { z } from "zod";
import { prisma } from "~/db.server";
import { findEnvironmentById } from "~/models/runtimeEnvironment.server";
import { logger } from "~/services/logger.server";
import { singleton } from "~/utils/singleton";
import { marqs, sanitizeQueueName } from "~/v3/marqs/index.server";
import { EnvironmentVariablesRepository } from "../environmentVariables/environmentVariablesRepository.server";
import { resolveVariablesForEnvironment } from "../environmentVariables/environmentVariablesRepository.server";
import { FailedTaskRunService } from "../failedTaskRun.server";
import { generateFriendlyId } from "../friendlyIdentifiers";
import { socketIo } from "../handleSocketIo.server";
import {
@@ -31,12 +35,14 @@ import {
getWorkerDeploymentFromWorker,
getWorkerDeploymentFromWorkerTask,
} from "../models/workerDeployment.server";
import { RestoreCheckpointService } from "../services/restoreCheckpoint.server";
import { SEMINTATTRS_FORCE_RECORDING, tracer } from "../tracer.server";
import { CrashTaskRunService } from "../services/crashTaskRun.server";
import { FailedTaskRunService } from "../failedTaskRun.server";
import { CreateTaskRunAttemptService } from "../services/createTaskRunAttempt.server";
import { findEnvironmentById } from "~/models/runtimeEnvironment.server";
import { RestoreCheckpointService } from "../services/restoreCheckpoint.server";
import { tracer } from "../tracer.server";
import { generateJWTTokenForEnvironment } from "~/services/apiAuth.server";
import { EnvironmentVariable } from "../environmentVariables/repository";
import { machinePresetFromConfig } from "../machinePresets.server";
import { env } from "~/env.server";
const WithTraceContext = z.object({
traceparent: z.string().optional(),
@@ -405,6 +411,9 @@ export class SharedQueueConsumer {
lockedAt: new Date(),
lockedById: backgroundTask.id,
lockedToVersionId: deployment.worker.id,
startedAt: existingTaskRun.startedAt ?? new Date(),
baseCostInCents: env.BASE_RUN_COST_IN_CENTS,
machinePreset: machinePresetFromConfig(backgroundTask.machineConfig ?? {}).name,
},
include: {
runtimeEnvironment: true,
@@ -505,18 +514,7 @@ export class SharedQueueConsumer {
});
} else {
const machineConfig = lockedTaskRun.lockedBy?.machineConfig;
const machine = Machine.safeParse(machineConfig ?? {});
if (!machine.success) {
logger.error("Failed to parse machine config", {
queueMessage: message.data,
messageId: message.messageId,
machineConfig,
});
await this.#ackAndDoMoreWork(message.messageId);
return;
}
const machine = machinePresetFromConfig(machineConfig ?? {});
await this._sender.send("BACKGROUND_WORKER_MESSAGE", {
backgroundWorkerId: deployment.worker.friendlyId,
@@ -524,7 +522,7 @@ export class SharedQueueConsumer {
type: "SCHEDULE_ATTEMPT",
image: deployment.imageReference,
version: deployment.version,
machine: machine.data,
machine,
// identifiers
id: "placeholder", // TODO: Remove this completely in a future release
envId: lockedTaskRun.runtimeEnvironment.id,
@@ -554,6 +552,7 @@ export class SharedQueueConsumer {
lockedAt: null,
lockedById: null,
status: lockedTaskRun.status,
startedAt: existingTaskRun.startedAt,
},
}),
]);
@@ -1008,6 +1007,8 @@ class SharedQueueTasks {
const { backgroundWorkerTask, taskRun, queue } = attempt;
const machinePreset = machinePresetFromConfig(backgroundWorkerTask.machineConfig ?? {});
const execution: ProdTaskRunExecution = {
task: {
id: backgroundWorkerTask.slug,
@@ -1028,9 +1029,13 @@ class SharedQueueTasks {
payloadType: taskRun.payloadType,
context: taskRun.context,
createdAt: taskRun.createdAt,
startedAt: taskRun.startedAt ?? taskRun.createdAt,
tags: taskRun.tags.map((tag) => tag.name),
isTest: taskRun.isTest,
idempotencyKey: taskRun.idempotencyKey ?? undefined,
durationMs: taskRun.usageDurationMs,
costInCents: taskRun.costInCents,
baseCostInCents: taskRun.baseCostInCents,
},
queue: {
id: queue.friendlyId,
@@ -1061,12 +1066,13 @@ class SharedQueueTasks {
contentHash: attempt.backgroundWorker.contentHash,
version: attempt.backgroundWorker.version,
},
machine: machinePreset,
};
const environmentRepository = new EnvironmentVariablesRepository();
const variables = await environmentRepository.getEnvironmentVariables(
attempt.runtimeEnvironment.projectId,
attempt.runtimeEnvironmentId
const variables = await this.#buildEnvironmentVariables(
attempt.runtimeEnvironment,
taskRun,
machinePreset
);
const payload: ProdTaskRunExecutionPayload = {
@@ -1126,6 +1132,9 @@ class SharedQueueTasks {
id: runId,
runtimeEnvironmentId: environment.id,
},
include: {
lockedBy: true,
},
});
if (!run) {
@@ -1133,11 +1142,9 @@ class SharedQueueTasks {
return;
}
const environmentRepository = new EnvironmentVariablesRepository();
const variables = await environmentRepository.getEnvironmentVariables(
environment.projectId,
environment.id
);
const machinePreset = machinePresetFromConfig(run.lockedBy?.machineConfig ?? {});
const variables = await this.#buildEnvironmentVariables(environment, run, machinePreset);
return {
traceContext: run.traceContext as Record<string, unknown>,
@@ -1178,6 +1185,31 @@ class SharedQueueTasks {
await service.call(completion.id, completion);
}
async #buildEnvironmentVariables(
environment: RuntimeEnvironment,
run: TaskRun,
machinePreset: MachinePreset
): Promise<Array<EnvironmentVariable>> {
const variables = await resolveVariablesForEnvironment(environment);
const jwt = await generateJWTTokenForEnvironment(environment, {
run_id: run.id,
machine_preset: machinePreset.name,
});
return [
...variables,
...[
{ key: "TRIGGER_JWT", value: jwt },
{ key: "TRIGGER_RUN_ID", value: run.id },
{
key: "TRIGGER_MACHINE_PRESET",
value: machinePreset.name,
},
],
];
}
}
export const sharedQueueTasks = singleton("sharedQueueTasks", () => new SharedQueueTasks());
+1
View File
@@ -29,6 +29,7 @@ export interface MarQSKeyProducer {
envSharedQueueKey(env: AuthenticatedEnvironment): string;
sharedQueueKey(): string;
sharedQueueScanPattern(): string;
queueCurrentConcurrencyScanPattern(): string;
concurrencyLimitKeyFromQueue(queue: string): string;
currentConcurrencyKeyFromQueue(queue: string): string;
currentConcurrencyKey(
+1 -1
View File
@@ -82,7 +82,7 @@ function getMarQSClient() {
defaultEnvConcurrency: env.V2_MARQS_DEFAULT_ENV_CONCURRENCY, // this is so we aren't limited by the environment concurrency
defaultOrgConcurrency: env.DEFAULT_ORG_EXECUTION_CONCURRENCY_LIMIT,
visibilityTimeoutInMs: env.V2_MARQS_VISIBILITY_TIMEOUT_MS, // 15 minutes
enableRebalancing: !env.MARQS_DISABLE_REBALANCING,
enableRebalancing: env.V2_MARQS_CONSUMER_POOL_ENABLED === "1",
});
}
+47
View File
@@ -0,0 +1,47 @@
import { randomUUID } from "node:crypto";
import { env } from "~/env.server";
import { logger } from "~/services/logger.server";
export type UsageEvent = {
source: string;
subject: string;
type: string;
id?: string;
time?: Date;
data?: Record<string, unknown>;
};
export async function reportUsageEvent(event: UsageEvent) {
if (!env.USAGE_OPEN_METER_BASE_URL || !env.USAGE_OPEN_METER_API_KEY) {
return;
}
const body = {
specversion: "1.0",
id: event.id ?? randomUUID(),
source: event.source,
type: event.type,
time: (event.time ?? new Date()).toISOString(),
subject: event.subject,
datacontenttype: "application/json",
data: event.data,
};
const url = `${env.USAGE_OPEN_METER_BASE_URL}/api/v1/events`;
logger.debug("Reporting usage event to OpenMeter", { url, body });
const response = await fetch(url, {
method: "POST",
body: JSON.stringify(body),
headers: {
"Content-Type": "application/cloudevents+json",
Authorization: `Bearer ${env.USAGE_OPEN_METER_API_KEY}`,
Accept: "application/json",
},
});
if (!response.ok) {
logger.error(`Failed to report usage event: ${response.status} ${response.statusText}`);
}
}
+49 -1
View File
@@ -121,7 +121,14 @@ class OTLPExporter {
(attribute) => attribute.key === SemanticInternalAttributes.TRIGGER
);
if (!triggerAttribute) return false;
if (!triggerAttribute) {
logger.debug("Skipping resource span without trigger attribute", {
attributes: resourceSpan.resource?.attributes,
spans: resourceSpan.scopeSpans.flatMap((scopeSpan) => scopeSpan.spans),
});
return;
}
return isBoolValue(triggerAttribute.value) ? triggerAttribute.value.boolValue : false;
});
@@ -300,6 +307,19 @@ function convertSpansToCreateableEvents(resourceSpan: ResourceSpans): Array<Crea
"."
)
) ?? resourceProperties.attemptNumber,
usageDurationMs:
extractDoubleAttribute(
span.attributes ?? [],
SemanticInternalAttributes.USAGE_DURATION_MS
) ??
extractNumberAttribute(
span.attributes ?? [],
SemanticInternalAttributes.USAGE_DURATION_MS
),
usageCostInCents: extractDoubleAttribute(
span.attributes ?? [],
SemanticInternalAttributes.USAGE_COST_IN_CENTS
),
};
})
.filter(Boolean);
@@ -353,6 +373,20 @@ function extractResourceProperties(attributes: KeyValue[]) {
queueName: extractStringAttribute(attributes, SemanticInternalAttributes.QUEUE_NAME),
batchId: extractStringAttribute(attributes, SemanticInternalAttributes.BATCH_ID),
idempotencyKey: extractStringAttribute(attributes, SemanticInternalAttributes.IDEMPOTENCY_KEY),
machinePreset: extractStringAttribute(
attributes,
SemanticInternalAttributes.MACHINE_PRESET_NAME
),
machinePresetCpu:
extractDoubleAttribute(attributes, SemanticInternalAttributes.MACHINE_PRESET_CPU) ??
extractNumberAttribute(attributes, SemanticInternalAttributes.MACHINE_PRESET_CPU),
machinePresetMemory:
extractDoubleAttribute(attributes, SemanticInternalAttributes.MACHINE_PRESET_MEMORY) ??
extractNumberAttribute(attributes, SemanticInternalAttributes.MACHINE_PRESET_MEMORY),
machinePresetCentsPerMs: extractDoubleAttribute(
attributes,
SemanticInternalAttributes.MACHINE_PRESET_CENTS_PER_MS
),
};
}
@@ -604,6 +638,20 @@ function extractNumberAttribute(
return isIntValue(attribute?.value) ? Number(attribute.value.intValue) : fallback;
}
function extractDoubleAttribute(attributes: KeyValue[], name: string): number | undefined;
function extractDoubleAttribute(attributes: KeyValue[], name: string, fallback: number): number;
function extractDoubleAttribute(
attributes: KeyValue[],
name: string,
fallback?: number
): number | undefined {
const attribute = attributes.find((attribute) => attribute.key === name);
if (!attribute) return fallback;
return isDoubleValue(attribute?.value) ? Number(attribute.value.doubleValue) : fallback;
}
function extractBooleanAttribute(attributes: KeyValue[], name: string): boolean | undefined;
function extractBooleanAttribute(attributes: KeyValue[], name: string, fallback: boolean): boolean;
function extractBooleanAttribute(
+1
View File
@@ -55,6 +55,7 @@ export const UpsertSchedule = z.object({
),
externalId: z.string().optional(),
deduplicationKey: z.string().optional(),
timezone: z.string().optional(),
});
export type UpsertSchedule = z.infer<typeof UpsertSchedule>;
@@ -88,6 +88,7 @@ export class CompleteAttemptService extends BaseService {
completedAt: new Date(),
output: completion.output,
outputType: completion.outputType,
usageDurationMs: completion.usage?.durationMs,
taskRun: {
update: {
data: {
@@ -138,6 +139,7 @@ export class CompleteAttemptService extends BaseService {
// We need to cancel the task run instead of fail it
const cancelService = new CancelAttemptService();
// TODO: handle usages
await cancelService.call(
taskRunAttempt.friendlyId,
taskRunAttempt.taskRunId,
@@ -157,6 +159,7 @@ export class CompleteAttemptService extends BaseService {
status: "FAILED",
completedAt: new Date(),
error: completion.error,
usageDurationMs: completion.usage?.durationMs,
},
});
@@ -5,18 +5,21 @@ import { logger } from "~/services/logger.server";
import { generateFriendlyId } from "../friendlyIdentifiers";
import { BaseService, ServiceValidationError } from "./baseService.server";
import { TaskRun, TaskRunAttempt } from "@trigger.dev/database";
import { machinePresetFromConfig } from "../machinePresets.server";
import { workerQueue } from "~/services/worker.server";
export class CreateTaskRunAttemptService extends BaseService {
public async call(
runId: string,
env?: AuthenticatedEnvironment,
authenticatedEnv?: AuthenticatedEnvironment,
setToExecuting = true
): Promise<{
execution: TaskRunExecution;
run: TaskRun;
attempt: TaskRunAttempt;
}> {
const environment = env ?? (await getAuthenticatedEnvironmentFromRun(runId, this._prisma));
const environment =
authenticatedEnv ?? (await getAuthenticatedEnvironmentFromRun(runId, this._prisma));
if (!environment) {
throw new ServiceValidationError("Environment not found", 404);
@@ -128,6 +131,20 @@ export class CreateTaskRunAttemptService extends BaseService {
throw new ServiceValidationError("Failed to create task run attempt", 500);
}
if (taskRunAttempt.number === 1 && taskRun.baseCostInCents > 0) {
await workerQueue.enqueue("v3.reportUsage", {
orgId: environment.organizationId,
data: {
costInCents: String(taskRun.baseCostInCents),
},
additionalData: {
runId: taskRun.id,
},
});
}
const machinePreset = machinePresetFromConfig(taskRun.lockedBy.machineConfig ?? {});
const execution: TaskRunExecution = {
task: {
id: taskRun.lockedBy.slug,
@@ -151,6 +168,10 @@ export class CreateTaskRunAttemptService extends BaseService {
tags: taskRun.tags.map((tag) => tag.name),
isTest: taskRun.isTest,
idempotencyKey: taskRun.idempotencyKey ?? undefined,
startedAt: taskRun.startedAt ?? taskRun.createdAt,
durationMs: taskRun.usageDurationMs,
costInCents: taskRun.costInCents,
baseCostInCents: taskRun.baseCostInCents,
},
queue: {
id: queue.friendlyId,
@@ -176,6 +197,7 @@ export class CreateTaskRunAttemptService extends BaseService {
taskRun.batchItems[0] && taskRun.batchItems[0].batchTaskRun
? { id: taskRun.batchItems[0].batchTaskRun.friendlyId }
: undefined,
machine: machinePreset,
};
return {
@@ -1,14 +1,28 @@
import { PerformDeploymentAlertsService } from "./alerts/performDeploymentAlerts.server";
import { BaseService } from "./baseService.server";
import { logger } from "~/services/logger.server";
import { WorkerDeploymentStatus } from "@trigger.dev/database";
const FINAL_DEPLOYMENT_STATUSES: WorkerDeploymentStatus[] = [
"CANCELED",
"DEPLOYED",
"FAILED",
"TIMED_OUT",
];
export class DeploymentIndexFailed extends BaseService {
public async call(
maybeFriendlyId: string,
error: { name: string; message: string; stack?: string }
error: {
name: string;
message: string;
stack?: string;
stderr?: string;
}
) {
const isFriendlyId = maybeFriendlyId.startsWith("deployment_");
const deployment = await this._prisma.workerDeployment.update({
const deployment = await this._prisma.workerDeployment.findUnique({
where: isFriendlyId
? {
friendlyId: maybeFriendlyId,
@@ -16,6 +30,25 @@ export class DeploymentIndexFailed extends BaseService {
: {
id: maybeFriendlyId,
},
});
if (!deployment) {
logger.error("Worker deployment not found", { maybeFriendlyId });
return;
}
if (FINAL_DEPLOYMENT_STATUSES.includes(deployment.status)) {
logger.error("Worker deployment already in final state", {
id: deployment.id,
status: deployment.status,
});
return;
}
const failedDeployment = await this._prisma.workerDeployment.update({
where: {
id: deployment.id,
},
data: {
status: "FAILED",
failedAt: new Date(),
@@ -23,8 +56,8 @@ export class DeploymentIndexFailed extends BaseService {
},
});
await PerformDeploymentAlertsService.enqueue(deployment.id, this._prisma);
await PerformDeploymentAlertsService.enqueue(failedDeployment.id, this._prisma);
return deployment;
return failedDeployment;
}
}
@@ -20,6 +20,7 @@ export class RegisterNextTaskScheduleInstanceService extends BaseService {
const nextScheduledTimestamp = calculateNextScheduledTimestamp(
instance.taskSchedule.generatorExpression,
instance.taskSchedule.timezone,
instance.lastScheduledTimestamp ?? new Date()
);
@@ -1,9 +1,9 @@
import { TaskRunStatus, type Checkpoint, TaskRunAttemptStatus } from "@trigger.dev/database";
import { TaskRunAttemptStatus, TaskRunStatus, type Checkpoint } from "@trigger.dev/database";
import { logger } from "~/services/logger.server";
import { socketIo } from "../handleSocketIo.server";
import { CreateCheckpointRestoreEventService } from "./createCheckpointRestoreEvent.server";
import { machinePresetFromConfig } from "../machinePresets.server";
import { BaseService } from "./baseService.server";
import { Machine } from "@trigger.dev/core/v3";
import { CreateCheckpointRestoreEventService } from "./createCheckpointRestoreEvent.server";
const RESTORABLE_RUN_STATUSES: TaskRunStatus[] = ["WAITING_TO_RESUME"];
const RESTORABLE_ATTEMPT_STATUSES: TaskRunAttemptStatus[] = ["PAUSED"];
@@ -45,7 +45,7 @@ export class RestoreCheckpointService extends BaseService {
});
if (!checkpointEvent) {
logger.error("Checkpoint event not found", params);
logger.error("Checkpoint event not found", { eventId: params.eventId });
return;
}
@@ -56,30 +56,26 @@ export class RestoreCheckpointService extends BaseService {
if (!runIsRestorable) {
logger.error("Run is unrestorable", {
id: checkpoint.runId,
status: checkpoint.run.status,
eventId: params.eventId,
runId: checkpoint.runId,
runStatus: checkpoint.run.status,
attemptId: checkpoint.attemptId,
});
return;
}
if (!attemptIsRestorable && !params.isRetry) {
logger.error("Attempt is unrestorable", {
id: checkpoint.attemptId,
status: checkpoint.attempt.status,
eventId: params.eventId,
runId: checkpoint.runId,
attemptId: checkpoint.attemptId,
attemptStatus: checkpoint.attempt.status,
});
return;
}
const { machineConfig } = checkpoint.attempt.backgroundWorkerTask;
const machine = Machine.safeParse(machineConfig ?? {});
if (!machine.success) {
logger.error("Failed to parse machine config", {
attemptId: checkpoint.attemptId,
machineConfig: checkpoint.attempt.backgroundWorkerTask.machineConfig,
});
return;
}
const machine = machinePresetFromConfig(machineConfig ?? {});
const restoreEvent = await this._prisma.checkpointRestoreEvent.findFirst({
where: {
@@ -90,6 +86,8 @@ export class RestoreCheckpointService extends BaseService {
if (restoreEvent) {
logger.error("Restore event already exists", {
runId: checkpoint.runId,
attemptId: checkpoint.attemptId,
checkpointId: checkpoint.id,
restoreEventId: restoreEvent.id,
});
@@ -106,7 +104,7 @@ export class RestoreCheckpointService extends BaseService {
location: checkpoint.location,
reason: checkpoint.reason ?? undefined,
imageRef: checkpoint.imageRef,
machine: machine.data,
machine,
// identifiers
checkpointId: checkpoint.id,
envId: checkpoint.runtimeEnvironment.id,

Some files were not shown because too many files have changed in this diff Show More