Compare commits
159 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| b53a575464 | |||
| ad68a3cc05 | |||
| be58638904 | |||
| 8f43aecacc | |||
| 6fd7560d5c | |||
| eb3b7b6f9e | |||
| 7aed154854 | |||
| 50506dce9f | |||
| d5772e57bd | |||
| 14c2bdf89b | |||
| 7976d924fb | |||
| 1e667ec28f | |||
| 7e97dcb93e | |||
| b9b86c89a7 | |||
| cd5d2ae92b | |||
| 0e77e7ef7d | |||
| 76a5c6204f | |||
| b171fde483 | |||
| 5ae3da6b4e | |||
| f565829959 | |||
| d57dec6919 | |||
| 75ec4ac6a6 | |||
| 374b6b9c0c | |||
| d0d3a64bd6 | |||
| 568da01785 | |||
| c75e29a9a7 | |||
| e5d26bd12d | |||
| b6f31ab651 | |||
| 50d46a8513 | |||
| 4cc61ac0ec | |||
| a696359c3e | |||
| 52b6f48a94 | |||
| d22a460555 | |||
| 9ba2a217a4 | |||
| 39885a427f | |||
| ccb0bc510a | |||
| 56d66ee07c | |||
| 4ca8887972 | |||
| 89bffc066c | |||
| 34ca7667d3 | |||
| 3e327acc0f | |||
| 77ad4127cb | |||
| 8a5076aacf | |||
| 5399f6bfb7 | |||
| ecef199660 | |||
| 4acfb8f4bb | |||
| 2ef278db67 | |||
| c7a55804d9 | |||
| da6a66efff | |||
| 7c36a1a4b0 | |||
| 225effb599 | |||
| 98eb6ed4f9 | |||
| 3069ebf0d8 | |||
| e133e628ca | |||
| 098932ea96 | |||
| 65f960e883 | |||
| ccbeff47e6 | |||
| 7c8f2df105 | |||
| fd44dabfe0 | |||
| 5daed3f69d | |||
| 596bf78e55 | |||
| 6ca66b76f4 | |||
| 29ef0395ce | |||
| 55d1f8c677 | |||
| 9835f4ec55 | |||
| 7fae10db23 | |||
| 8cf1f0a37d | |||
| dba4313c5c | |||
| 506613dc92 | |||
| 764df23d19 | |||
| 4b961a6ae2 | |||
| 8757fdceef | |||
| 2404e88ac5 | |||
| 88b36f5090 | |||
| b73ae3f927 | |||
| b605b892ac | |||
| 233316f7e8 | |||
| 1b90ffbb8c | |||
| b45ca4e146 | |||
| 25d15578f7 | |||
| fe865a0f49 | |||
| 0ed93a748e | |||
| e02320f65d | |||
| 85a543d8ec | |||
| 10ceb85a92 | |||
| c405ae7117 | |||
| 3687fcb61e | |||
| d4ccdf7105 | |||
| e08b4569e5 | |||
| 79da0ca9b5 | |||
| 3aca603a33 | |||
| c9e97d6b78 | |||
| 01633c9c03 | |||
| 691990d79e | |||
| b2ba403dd3 | |||
| 1d47cab69f | |||
| e23047f9ad | |||
| 68d32429b6 | |||
| 36ac79ac66 | |||
| ca94f0cac3 | |||
| a5d8e453a5 | |||
| c332519e72 | |||
| 52112c3bfc | |||
| eae294a332 | |||
| 465cd0335c | |||
| 35dbaedf69 | |||
| c11a77f50b | |||
| fb52b9efea | |||
| 0896b9fffc | |||
| 3a2dd983c5 | |||
| a627ca67d1 | |||
| afc180aa70 | |||
| 393af1b7c5 | |||
| 6a91fb89b8 | |||
| df7d1de16d | |||
| 8e8ed4a3bf | |||
| 8fc8f57b39 | |||
| 9ebd91ccec | |||
| 665f7c9756 | |||
| 928a632e23 | |||
| 74db2de1bc | |||
| 93acca6c3c | |||
| ebe079d83c | |||
| d272996de3 | |||
| 531bd4970d | |||
| 5c9eb25b5a | |||
| c970e892a7 | |||
| a867b6e5ae | |||
| d44abbd0fc | |||
| 1cc680ac1e | |||
| 9b049bc480 | |||
| c24a23b551 | |||
| ee1ae1fca6 | |||
| 8e5ef176a4 | |||
| 58b6b1aa0d | |||
| 9c0ae1459f | |||
| 2f15a84320 | |||
| a49a0ff416 | |||
| b703ffed29 | |||
| b4f9b70ae2 | |||
| 51bb4c887a | |||
| ba71f959e2 | |||
| bc7bbd4576 | |||
| 5fe23e4b3f | |||
| 7b3b2e0d8e | |||
| 3900ddadce | |||
| ca9e827bd3 | |||
| 04e936b69b | |||
| 98ef170299 | |||
| e69ffd314a | |||
| 782d4f75ae | |||
| b6de469d07 | |||
| 0dd3447c31 | |||
| a5a5d3ae21 | |||
| ee3619bbb1 | |||
| d9ad72446e | |||
| a56f9af9fe | |||
| ece6ca678a | |||
| 6243ae30bb |
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
---
|
||||
|
||||
Add an e2e suite to test compiling with v3 CLI.
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
---
|
||||
|
||||
cli v3: increase otel force flush timeout to 30s from 500ms
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
---
|
||||
|
||||
Increase dev worker timeout
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
---
|
||||
|
||||
Add sox and audiowaveform binaries to worker images
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"@trigger.dev/core": patch
|
||||
---
|
||||
|
||||
Support triggering tasks with non-URL friendly characters in the ID
|
||||
@@ -0,0 +1,7 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
"@trigger.dev/core": patch
|
||||
"@trigger.dev/sdk": patch
|
||||
---
|
||||
|
||||
v3: Usage tracking
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"@trigger.dev/sdk": patch
|
||||
---
|
||||
|
||||
Fix for calling trigger and passing a custom queue
|
||||
@@ -0,0 +1,30 @@
|
||||
---
|
||||
"@trigger.dev/core-apps": patch
|
||||
"trigger.dev": patch
|
||||
"@trigger.dev/core": patch
|
||||
---
|
||||
|
||||
Tasks should now be much more robust and resilient to reconnects during crucial operations and other failure scenarios.
|
||||
|
||||
Task runs now have to signal checkpointable state prior to ALL checkpoints. This ensures flushing always happens.
|
||||
|
||||
All important socket.io RPCs will now be retried with backoff. Actions relying on checkpoints will be replayed if we haven't been checkpointed and restored as expected, e.g. after reconnect.
|
||||
|
||||
Other changes:
|
||||
|
||||
- Fix retry check in shared queue
|
||||
- Fix env var sync spinner
|
||||
- Heartbeat between retries
|
||||
- Fix retry prep
|
||||
- Fix prod worker no tasks detection
|
||||
- Fail runs above `MAX_TASK_RUN_ATTEMPTS`
|
||||
- Additional debug logs in all places
|
||||
- Prevent crashes due to failed socket schema parsing
|
||||
- Remove core-apps barrel
|
||||
- Upgrade socket.io-client to fix an ACK memleak
|
||||
- Additional index failure logs
|
||||
- Prevent message loss during reconnect
|
||||
- Prevent burst of heartbeats on reconnect
|
||||
- Prevent crash on failed cleanup
|
||||
- Handle at-least-once lazy execute message delivery
|
||||
- Handle uncaught entry point exceptions
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"@trigger.dev/core": patch
|
||||
---
|
||||
|
||||
v3: Remove aggressive otel flush timeouts in dev/prod
|
||||
@@ -0,0 +1,6 @@
|
||||
---
|
||||
"@trigger.dev/sdk": patch
|
||||
"@trigger.dev/core": patch
|
||||
---
|
||||
|
||||
v3: Trigger delayed runs and reschedule them
|
||||
@@ -0,0 +1,9 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
"@trigger.dev/core": patch
|
||||
---
|
||||
|
||||
- Improve non-zero exit code error messages
|
||||
- Detect OOM conditions within worker child processes
|
||||
- Internal errors can have optional stack traces
|
||||
- Docker provider can be set to enforce machine presets
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"@trigger.dev/core": patch
|
||||
---
|
||||
|
||||
Fix issue when using SDK in non-node environments by scoping the stream import with node:
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
---
|
||||
|
||||
Output stderr logs on dev worker failure
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"@trigger.dev/sdk": patch
|
||||
---
|
||||
|
||||
Use global setTimeout to ensure cross-runtime support
|
||||
@@ -0,0 +1,7 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
---
|
||||
|
||||
Add e2e fixtures corresponding to past issues
|
||||
Implement e2e suite parallelism
|
||||
Enhance log level for specific e2e suite messages
|
||||
+41
-1
@@ -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",
|
||||
@@ -63,12 +65,17 @@
|
||||
"eight-pumas-float",
|
||||
"eleven-paws-join",
|
||||
"famous-boats-tease",
|
||||
"fast-colts-relax",
|
||||
"few-students-share",
|
||||
"five-toes-destroy",
|
||||
"friendly-walls-repair",
|
||||
"funny-swans-destroy",
|
||||
"gorgeous-gorillas-compete",
|
||||
"green-bags-wink",
|
||||
"hot-buckets-behave",
|
||||
"hot-fishes-retire",
|
||||
"hot-wasps-sin",
|
||||
"itchy-chairs-itch",
|
||||
"khaki-apricots-design",
|
||||
"khaki-poems-lay",
|
||||
"late-icons-lie",
|
||||
@@ -83,45 +90,73 @@
|
||||
"lovely-drinks-flash",
|
||||
"many-ligers-pump",
|
||||
"mighty-camels-joke",
|
||||
"mighty-eggs-grab",
|
||||
"mighty-flowers-train",
|
||||
"mighty-parrots-sin",
|
||||
"modern-stingrays-end",
|
||||
"nasty-jars-pump",
|
||||
"nervous-planets-sparkle",
|
||||
"new-pants-beg",
|
||||
"new-rivers-tell",
|
||||
"nice-bulldogs-turn",
|
||||
"ninety-pets-travel",
|
||||
"odd-poets-own",
|
||||
"pink-pumas-rhyme",
|
||||
"plenty-ducks-beam",
|
||||
"polite-ducks-switch",
|
||||
"polite-pears-grow",
|
||||
"polite-pots-walk",
|
||||
"polite-rockets-matter",
|
||||
"poor-flowers-cross",
|
||||
"purple-garlics-shop",
|
||||
"purple-spiders-care",
|
||||
"rare-lamps-promise",
|
||||
"rare-roses-float",
|
||||
"real-planets-stare",
|
||||
"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-buses-obey",
|
||||
"silly-forks-kiss",
|
||||
"silly-suits-switch",
|
||||
"silver-doors-juggle",
|
||||
"six-ligers-exist",
|
||||
"six-rats-hunt",
|
||||
"sixty-insects-watch",
|
||||
"slow-buses-own",
|
||||
"slow-kiwis-hide",
|
||||
"slow-sloths-retire",
|
||||
"smart-needles-move",
|
||||
"smart-olives-eat",
|
||||
"sour-pugs-teach",
|
||||
"spicy-frogs-remain",
|
||||
"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-ducks-remember",
|
||||
"sweet-lizards-press",
|
||||
"swift-dragons-peel",
|
||||
"tall-bees-wave",
|
||||
"tall-masks-repeat",
|
||||
"tame-apricots-clap",
|
||||
"tame-guests-know",
|
||||
"tender-moose-tell",
|
||||
"tender-oranges-rhyme",
|
||||
"tender-turkeys-compete",
|
||||
"thick-carrots-sneeze",
|
||||
"thin-parents-heal",
|
||||
"thirty-islands-kiss",
|
||||
"tidy-balloons-suffer",
|
||||
@@ -130,8 +165,13 @@
|
||||
"tiny-doors-type",
|
||||
"tiny-elephants-scream",
|
||||
"tricky-bulldogs-heal",
|
||||
"tricky-keys-attack",
|
||||
"tricky-ladybugs-unite",
|
||||
"twelve-knives-notice",
|
||||
"two-pumas-wait",
|
||||
"warm-planes-taste"
|
||||
"violet-clocks-notice",
|
||||
"warm-olives-provide",
|
||||
"warm-planes-taste",
|
||||
"young-snails-sell"
|
||||
]
|
||||
}
|
||||
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
---
|
||||
|
||||
Await file watcher cleanup in dev
|
||||
@@ -0,0 +1,7 @@
|
||||
---
|
||||
"@trigger.dev/core-apps": patch
|
||||
"trigger.dev": patch
|
||||
"@trigger.dev/core": patch
|
||||
---
|
||||
|
||||
Capture and display stderr on index failures
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
---
|
||||
|
||||
v3: [prod] force flush timeout should be 1s
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"@trigger.dev/core": patch
|
||||
---
|
||||
|
||||
Add callback to checkpoint created message
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
---
|
||||
|
||||
v3: Copy over more of the project's package.json keys into the deployed package.json (support for custom config like zenstack)
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"@trigger.dev/core": patch
|
||||
---
|
||||
|
||||
Make deduplicationKey required when creating/updating a schedule
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
---
|
||||
|
||||
Improved ESM module require error detection logic
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
---
|
||||
|
||||
Set the deploy timeout to 3mins from 1min
|
||||
@@ -0,0 +1,6 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
"@trigger.dev/core": patch
|
||||
---
|
||||
|
||||
v3: vercel edge runtime support
|
||||
@@ -0,0 +1,7 @@
|
||||
---
|
||||
"@trigger.dev/core-apps": patch
|
||||
"@trigger.dev/core": patch
|
||||
---
|
||||
|
||||
- Fix uncaught provider exception
|
||||
- Remove unused provider messages
|
||||
@@ -0,0 +1,6 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
"@trigger.dev/core": patch
|
||||
---
|
||||
|
||||
v3: fix otel flushing causing CLEANUP ack timeout errors by always setting a forceFlushTimeoutMillis value
|
||||
@@ -0,0 +1,10 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
---
|
||||
|
||||
- Prevent downgrades during update check and advise to upgrade CLI
|
||||
- Detect bun and use npm instead
|
||||
- During init, fail early and advise if not a TypeScript project
|
||||
- During init, allow specifying custom package manager args
|
||||
- Add links to dev worker started message
|
||||
- Fix links in unsupported terminals
|
||||
@@ -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
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
---
|
||||
|
||||
fix: allow command login to read api url from cli args
|
||||
@@ -0,0 +1,6 @@
|
||||
---
|
||||
"@trigger.dev/sdk": patch
|
||||
"@trigger.dev/core": patch
|
||||
---
|
||||
|
||||
Management SDK overhaul and adding the runs.list API
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"@trigger.dev/sdk": patch
|
||||
---
|
||||
|
||||
v3: Adding SDK functions for triggering tasks in a typesafe way, without importing task file
|
||||
@@ -0,0 +1,9 @@
|
||||
---
|
||||
"@trigger.dev/core-apps": patch
|
||||
"trigger.dev": patch
|
||||
"@trigger.dev/core": patch
|
||||
---
|
||||
|
||||
- Fix artifact detection logs
|
||||
- Fix OOM detection and error messages
|
||||
- Add test link to cli deployment completion
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
---
|
||||
|
||||
v3: postInstall config option now replaces the postinstall script found in package.json
|
||||
@@ -0,0 +1,6 @@
|
||||
---
|
||||
"@trigger.dev/sdk": patch
|
||||
"@trigger.dev/core": patch
|
||||
---
|
||||
|
||||
Added timezone support to schedules
|
||||
@@ -0,0 +1,6 @@
|
||||
---
|
||||
"@trigger.dev/core": patch
|
||||
"@trigger.dev/sdk": patch
|
||||
---
|
||||
|
||||
v3: Include presigned urls for downloading large payloads and outputs when using runs.retrieve
|
||||
@@ -0,0 +1,14 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
"@trigger.dev/core": patch
|
||||
---
|
||||
|
||||
- Clear paused states before retry
|
||||
- Detect and handle unrecoverable worker errors
|
||||
- Remove checkpoints after successful push
|
||||
- Permanently switch to DO hosted busybox image
|
||||
- Fix IPC timeout issue, or at least handle it more gracefully
|
||||
- Handle checkpoint failures
|
||||
- Basic chaos monkey for checkpoint testing
|
||||
- Stack traces are back in the dashboard
|
||||
- Display final errors on root span
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"@trigger.dev/core": patch
|
||||
---
|
||||
|
||||
v3: fix missing init output in task run function when no middleware is defined
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
---
|
||||
|
||||
Fix jsonc-parser import
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"@trigger.dev/core": patch
|
||||
---
|
||||
|
||||
Improve handling of IPC timeouts and fix checkpoint cancellation after failures
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
---
|
||||
|
||||
Increase cleanup IPC timeout
|
||||
+3
-1
@@ -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=
|
||||
@@ -69,7 +71,7 @@ COORDINATOR_SECRET=coordinator-secret # generate the actual secret with `openssl
|
||||
# OBJECT_STORE_BASE_URL="https://{bucket}.{accountId}.r2.cloudflarestorage.com"
|
||||
# OBJECT_STORE_ACCESS_KEY_ID=
|
||||
# OBJECT_STORE_SECRET_ACCESS_KEY=
|
||||
# RUNTIME_WAIT_THRESHOLD_IN_MS=10000
|
||||
# CHECKPOINT_THRESHOLD_IN_MS=10000
|
||||
|
||||
# These control the server-side internal telemetry
|
||||
# INTERNAL_OTEL_TRACE_EXPORTER_URL=<URL to send traces to>
|
||||
|
||||
@@ -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@v4
|
||||
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
|
||||
@@ -19,7 +63,7 @@ jobs:
|
||||
submodules: recursive
|
||||
|
||||
- name: ⎔ Setup pnpm
|
||||
uses: pnpm/action-setup@v2.2.4
|
||||
uses: pnpm/action-setup@v4
|
||||
with:
|
||||
version: 8.15.5
|
||||
|
||||
|
||||
@@ -29,4 +29,6 @@ jobs:
|
||||
|
||||
# e2e:
|
||||
# uses: ./.github/workflows/e2e.yml
|
||||
# with:
|
||||
# package: webapp
|
||||
# secrets: inherit
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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 }}
|
||||
|
||||
@@ -51,9 +51,16 @@ jobs:
|
||||
|
||||
# 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
|
||||
|
||||
@@ -28,7 +28,7 @@ jobs:
|
||||
fetch-depth: 0
|
||||
|
||||
- name: ⎔ Setup pnpm
|
||||
uses: pnpm/action-setup@v2.2.4
|
||||
uses: pnpm/action-setup@v4
|
||||
with:
|
||||
version: 8.15.5
|
||||
|
||||
|
||||
@@ -12,7 +12,7 @@ jobs:
|
||||
fetch-depth: 0
|
||||
|
||||
- name: ⎔ Setup pnpm
|
||||
uses: pnpm/action-setup@v2.2.4
|
||||
uses: pnpm/action-setup@v4
|
||||
with:
|
||||
version: 8.15.5
|
||||
|
||||
|
||||
@@ -12,7 +12,7 @@ jobs:
|
||||
fetch-depth: 0
|
||||
|
||||
- name: ⎔ Setup pnpm
|
||||
uses: pnpm/action-setup@v2.2.4
|
||||
uses: pnpm/action-setup@v4
|
||||
with:
|
||||
version: 8.15.5
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -21,8 +21,7 @@
|
||||
"execa": "^8.0.1",
|
||||
"nanoid": "^5.0.6",
|
||||
"prom-client": "^15.1.0",
|
||||
"socket.io": "4.7.4",
|
||||
"socket.io-client": "4.7.4"
|
||||
"socket.io": "4.7.4"
|
||||
},
|
||||
"devDependencies": {
|
||||
"@types/node": "^18",
|
||||
|
||||
@@ -0,0 +1,95 @@
|
||||
import type { Execa$ } from "execa";
|
||||
import { setTimeout as timeout } from "node:timers/promises";
|
||||
|
||||
class ChaosMonkeyError extends Error {
|
||||
constructor(message: string) {
|
||||
super(message);
|
||||
this.name = "ChaosMonkeyError";
|
||||
}
|
||||
}
|
||||
|
||||
export class ChaosMonkey {
|
||||
private chaosEventRate = 0.2;
|
||||
private delayInSeconds = 45;
|
||||
|
||||
constructor(private enabled = false) {
|
||||
if (this.enabled) {
|
||||
console.log("🍌 Chaos monkey enabled");
|
||||
}
|
||||
}
|
||||
|
||||
static Error = ChaosMonkeyError;
|
||||
|
||||
enable() {
|
||||
this.enabled = true;
|
||||
console.log("🍌 Chaos monkey enabled");
|
||||
}
|
||||
|
||||
disable() {
|
||||
this.enabled = false;
|
||||
console.log("🍌 Chaos monkey disabled");
|
||||
}
|
||||
|
||||
async call({
|
||||
$,
|
||||
throwErrors = true,
|
||||
addDelays = true,
|
||||
}: {
|
||||
$?: Execa$<string>;
|
||||
throwErrors?: boolean;
|
||||
addDelays?: boolean;
|
||||
} = {}) {
|
||||
if (!this.enabled) {
|
||||
return;
|
||||
}
|
||||
|
||||
const random = Math.random();
|
||||
|
||||
if (random > this.chaosEventRate) {
|
||||
// Don't interfere with normal operation
|
||||
return;
|
||||
}
|
||||
|
||||
const chaosEvents: Array<() => Promise<any>> = [];
|
||||
|
||||
if (addDelays) {
|
||||
chaosEvents.push(async () => {
|
||||
console.log("🍌 Chaos monkey: Add delay");
|
||||
|
||||
if ($) {
|
||||
await $`sleep ${this.delayInSeconds}`;
|
||||
} else {
|
||||
await timeout(this.delayInSeconds * 1000);
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
if (throwErrors) {
|
||||
chaosEvents.push(async () => {
|
||||
console.log("🍌 Chaos monkey: Throw error");
|
||||
|
||||
if ($) {
|
||||
await $`false`;
|
||||
} else {
|
||||
throw new ChaosMonkey.Error("🍌 Chaos monkey: Throw error");
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
if (chaosEvents.length === 0) {
|
||||
console.error("🍌 Chaos monkey: No events selected");
|
||||
return;
|
||||
}
|
||||
|
||||
const randomIndex = Math.floor(Math.random() * chaosEvents.length);
|
||||
|
||||
const chaosEvent = chaosEvents[randomIndex];
|
||||
|
||||
if (!chaosEvent) {
|
||||
console.error("🍌 Chaos monkey: No event found");
|
||||
return;
|
||||
}
|
||||
|
||||
await chaosEvent();
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,587 @@
|
||||
import { ExponentialBackoff } from "@trigger.dev/core-apps/backoff";
|
||||
import { isExecaChildProcess, testDockerCheckpoint } from "@trigger.dev/core-apps/checkpoints";
|
||||
import { SimpleLogger } from "@trigger.dev/core-apps/logger";
|
||||
import { $ } from "execa";
|
||||
import { nanoid } from "nanoid";
|
||||
import fs from "node:fs/promises";
|
||||
import { ChaosMonkey } from "./chaosMonkey";
|
||||
|
||||
type CheckpointerInitializeReturn = {
|
||||
canCheckpoint: boolean;
|
||||
willSimulate: boolean;
|
||||
};
|
||||
|
||||
type CheckpointAndPushOptions = {
|
||||
runId: string;
|
||||
leaveRunning?: boolean;
|
||||
projectRef: string;
|
||||
deploymentVersion: string;
|
||||
shouldHeartbeat?: boolean;
|
||||
};
|
||||
|
||||
type CheckpointAndPushResult =
|
||||
| { success: true; checkpoint: CheckpointData }
|
||||
| {
|
||||
success: false;
|
||||
reason?: "CANCELED" | "DISABLED" | "ERROR" | "IN_PROGRESS" | "NO_SUPPORT" | "SKIP_RETRYING";
|
||||
};
|
||||
|
||||
type CheckpointData = {
|
||||
location: string;
|
||||
docker: boolean;
|
||||
};
|
||||
|
||||
type CheckpointerOptions = {
|
||||
dockerMode: boolean;
|
||||
forceSimulate: boolean;
|
||||
heartbeat: (runId: string) => void;
|
||||
registryHost?: string;
|
||||
registryNamespace?: string;
|
||||
registryTlsVerify?: boolean;
|
||||
disableCheckpointSupport?: boolean;
|
||||
checkpointPath?: string;
|
||||
simulateCheckpointFailure?: boolean;
|
||||
simulateCheckpointFailureSeconds?: number;
|
||||
simulatePushFailure?: boolean;
|
||||
simulatePushFailureSeconds?: number;
|
||||
chaosMonkey?: ChaosMonkey;
|
||||
};
|
||||
|
||||
async function getFileSize(filePath: string): Promise<number> {
|
||||
try {
|
||||
const stats = await fs.stat(filePath);
|
||||
return stats.size;
|
||||
} catch (error) {
|
||||
console.error("Error getting file size:", error);
|
||||
return -1;
|
||||
}
|
||||
}
|
||||
|
||||
async function getParsedFileSize(filePath: string) {
|
||||
const sizeInBytes = await getFileSize(filePath);
|
||||
|
||||
let message = `Size in bytes: ${sizeInBytes}`;
|
||||
|
||||
if (sizeInBytes > 1024 * 1024) {
|
||||
const sizeInMB = (sizeInBytes / 1024 / 1024).toFixed(2);
|
||||
message = `Size in MB (rounded): ${sizeInMB}`;
|
||||
} else if (sizeInBytes > 1024) {
|
||||
const sizeInKB = (sizeInBytes / 1024).toFixed(2);
|
||||
message = `Size in KB (rounded): ${sizeInKB}`;
|
||||
}
|
||||
|
||||
return {
|
||||
path: filePath,
|
||||
sizeInBytes,
|
||||
message,
|
||||
};
|
||||
}
|
||||
|
||||
export class Checkpointer {
|
||||
#initialized = false;
|
||||
#canCheckpoint = false;
|
||||
#dockerMode: boolean;
|
||||
|
||||
#logger = new SimpleLogger("[checkptr]");
|
||||
#abortControllers = new Map<string, AbortController>();
|
||||
#failedCheckpoints = new Map<string, unknown>();
|
||||
#waitingForRetry = new Set<string>();
|
||||
|
||||
private registryHost: string;
|
||||
private registryNamespace: string;
|
||||
private registryTlsVerify: boolean;
|
||||
|
||||
private disableCheckpointSupport: boolean;
|
||||
private checkpointPath: string;
|
||||
|
||||
private simulateCheckpointFailure: boolean;
|
||||
private simulateCheckpointFailureSeconds: number;
|
||||
private simulatePushFailure: boolean;
|
||||
private simulatePushFailureSeconds: number;
|
||||
|
||||
private chaosMonkey: ChaosMonkey;
|
||||
|
||||
constructor(private opts: CheckpointerOptions) {
|
||||
this.#dockerMode = opts.dockerMode;
|
||||
|
||||
this.registryHost = opts.registryHost ?? "localhost:5000";
|
||||
this.registryNamespace = opts.registryNamespace ?? "trigger";
|
||||
this.registryTlsVerify = opts.registryTlsVerify ?? true;
|
||||
|
||||
this.disableCheckpointSupport = opts.disableCheckpointSupport ?? false;
|
||||
this.checkpointPath = opts.checkpointPath ?? "/checkpoints";
|
||||
|
||||
this.simulateCheckpointFailure = opts.simulateCheckpointFailure ?? false;
|
||||
this.simulateCheckpointFailureSeconds = opts.simulateCheckpointFailureSeconds ?? 300;
|
||||
this.simulatePushFailure = opts.simulatePushFailure ?? false;
|
||||
this.simulatePushFailureSeconds = opts.simulatePushFailureSeconds ?? 300;
|
||||
|
||||
this.chaosMonkey = opts.chaosMonkey ?? new ChaosMonkey(!!process.env.CHAOS_MONKEY_ENABLED);
|
||||
}
|
||||
|
||||
async init(): Promise<CheckpointerInitializeReturn> {
|
||||
if (this.#initialized) {
|
||||
return this.#getInitReturn(this.#canCheckpoint);
|
||||
}
|
||||
|
||||
this.#logger.log(`${this.#dockerMode ? "Docker" : "Kubernetes"} mode`);
|
||||
|
||||
if (this.#dockerMode) {
|
||||
const testCheckpoint = await testDockerCheckpoint();
|
||||
|
||||
if (testCheckpoint.ok) {
|
||||
return this.#getInitReturn(true);
|
||||
}
|
||||
|
||||
this.#logger.error(testCheckpoint.message, testCheckpoint.error ?? "");
|
||||
return this.#getInitReturn(false);
|
||||
} else {
|
||||
try {
|
||||
await $`buildah login --get-login ${this.registryHost}`;
|
||||
} catch (error) {
|
||||
this.#logger.error(`No checkpoint support: Not logged in to registry ${this.registryHost}`);
|
||||
return this.#getInitReturn(false);
|
||||
}
|
||||
}
|
||||
|
||||
return this.#getInitReturn(true);
|
||||
}
|
||||
|
||||
#getInitReturn(canCheckpoint: boolean): CheckpointerInitializeReturn {
|
||||
this.#canCheckpoint = canCheckpoint;
|
||||
|
||||
if (canCheckpoint) {
|
||||
if (!this.#initialized) {
|
||||
this.#logger.log("Full checkpoint support!");
|
||||
}
|
||||
}
|
||||
|
||||
this.#initialized = true;
|
||||
|
||||
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,
|
||||
willSimulate,
|
||||
};
|
||||
}
|
||||
|
||||
#getImageRef(projectRef: string, deploymentVersion: string, shortCode: string) {
|
||||
return `${this.registryHost}/${this.registryNamespace}/${projectRef}:${deploymentVersion}.prod-${shortCode}`;
|
||||
}
|
||||
|
||||
#getExportLocation(projectRef: string, deploymentVersion: string, shortCode: string) {
|
||||
const basename = `${projectRef}-${deploymentVersion}-${shortCode}`;
|
||||
|
||||
if (this.#dockerMode) {
|
||||
return basename;
|
||||
} else {
|
||||
return `${this.checkpointPath}/${basename}.tar`;
|
||||
}
|
||||
}
|
||||
|
||||
async checkpointAndPush(opts: CheckpointAndPushOptions): Promise<CheckpointData | undefined> {
|
||||
const start = performance.now();
|
||||
this.#logger.log(`checkpointAndPush() start`, { start, opts });
|
||||
|
||||
let interval: NodeJS.Timer | undefined;
|
||||
|
||||
if (opts.shouldHeartbeat) {
|
||||
interval = setInterval(() => {
|
||||
this.#logger.log("Sending heartbeat", { runId: opts.runId });
|
||||
this.opts.heartbeat(opts.runId);
|
||||
}, 20_000);
|
||||
}
|
||||
|
||||
try {
|
||||
const result = await this.#checkpointAndPushWithBackoff(opts);
|
||||
|
||||
const end = performance.now();
|
||||
this.#logger.log(`checkpointAndPush() end`, {
|
||||
start,
|
||||
end,
|
||||
diff: end - start,
|
||||
opts,
|
||||
success: result.success,
|
||||
});
|
||||
|
||||
if (!result.success) {
|
||||
return;
|
||||
}
|
||||
|
||||
return result.checkpoint;
|
||||
} finally {
|
||||
if (opts.shouldHeartbeat) {
|
||||
clearInterval(interval);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
isCheckpointing(runId: string) {
|
||||
return this.#abortControllers.has(runId) || this.#waitingForRetry.has(runId);
|
||||
}
|
||||
|
||||
cancelCheckpoint(runId: string): boolean {
|
||||
// If the last checkpoint failed, pretend we canceled it
|
||||
// This ensures tasks don't wait for external resume messages to continue
|
||||
if (this.#hasFailedCheckpoint(runId)) {
|
||||
this.#clearFailedCheckpoint(runId);
|
||||
return true;
|
||||
}
|
||||
|
||||
if (this.#waitingForRetry.has(runId)) {
|
||||
this.#waitingForRetry.delete(runId);
|
||||
return true;
|
||||
}
|
||||
|
||||
const controller = this.#abortControllers.get(runId);
|
||||
|
||||
if (!controller) {
|
||||
this.#logger.debug("Nothing to cancel", { runId });
|
||||
return false;
|
||||
}
|
||||
|
||||
controller.abort("cancelCheckpointing()");
|
||||
this.#abortControllers.delete(runId);
|
||||
|
||||
return true;
|
||||
}
|
||||
|
||||
async #checkpointAndPushWithBackoff({
|
||||
runId,
|
||||
leaveRunning = true, // This mirrors kubernetes behaviour more accurately
|
||||
projectRef,
|
||||
deploymentVersion,
|
||||
}: CheckpointAndPushOptions): Promise<CheckpointAndPushResult> {
|
||||
this.#logger.log("Checkpointing with backoff", {
|
||||
runId,
|
||||
leaveRunning,
|
||||
projectRef,
|
||||
deploymentVersion,
|
||||
});
|
||||
|
||||
const backoff = new ExponentialBackoff()
|
||||
.type("EqualJitter")
|
||||
.base(3)
|
||||
.max(3 * 3600)
|
||||
.maxElapsed(48 * 3600);
|
||||
|
||||
for await (const { delay, retry } of backoff) {
|
||||
try {
|
||||
if (retry > 0) {
|
||||
this.#logger.error("Retrying checkpoint", {
|
||||
runId,
|
||||
retry,
|
||||
delay,
|
||||
});
|
||||
|
||||
this.#waitingForRetry.add(runId);
|
||||
await new Promise((resolve) => setTimeout(resolve, delay.milliseconds));
|
||||
|
||||
if (!this.#waitingForRetry.has(runId)) {
|
||||
this.#logger.log("Checkpoint canceled while waiting for retry", { runId });
|
||||
return { success: false, reason: "CANCELED" };
|
||||
} else {
|
||||
this.#waitingForRetry.delete(runId);
|
||||
}
|
||||
}
|
||||
|
||||
const result = await this.#checkpointAndPush({
|
||||
runId,
|
||||
leaveRunning,
|
||||
projectRef,
|
||||
deploymentVersion,
|
||||
});
|
||||
|
||||
if (result.success) {
|
||||
return result;
|
||||
}
|
||||
|
||||
if (result.reason === "CANCELED") {
|
||||
this.#logger.log("Checkpoint canceled, won't retry", { runId });
|
||||
// Don't fail the checkpoint, as it was canceled
|
||||
return result;
|
||||
}
|
||||
|
||||
if (result.reason === "IN_PROGRESS") {
|
||||
this.#logger.log("Checkpoint already in progress, won't retry", { runId });
|
||||
this.#failCheckpoint(runId, result.reason);
|
||||
return result;
|
||||
}
|
||||
|
||||
if (result.reason === "NO_SUPPORT") {
|
||||
this.#logger.log("No checkpoint support, won't retry", { runId });
|
||||
this.#failCheckpoint(runId, result.reason);
|
||||
return result;
|
||||
}
|
||||
|
||||
if (result.reason === "DISABLED") {
|
||||
this.#logger.log("Checkpoint support disabled, won't retry", { runId });
|
||||
this.#failCheckpoint(runId, result.reason);
|
||||
return result;
|
||||
}
|
||||
|
||||
if (result.reason === "SKIP_RETRYING") {
|
||||
this.#logger.log("Skipping retrying", { runId });
|
||||
return result;
|
||||
}
|
||||
|
||||
continue;
|
||||
} catch (error) {
|
||||
this.#logger.error("Checkpoint error", {
|
||||
retry,
|
||||
runId,
|
||||
delay,
|
||||
error: error instanceof Error ? error.message : error,
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
this.#logger.error(`Checkpoint failed after exponential backoff`, {
|
||||
runId,
|
||||
leaveRunning,
|
||||
projectRef,
|
||||
deploymentVersion,
|
||||
});
|
||||
this.#failCheckpoint(runId, "ERROR");
|
||||
|
||||
return { success: false, reason: "ERROR" };
|
||||
}
|
||||
|
||||
async #checkpointAndPush({
|
||||
runId,
|
||||
leaveRunning = true, // This mirrors kubernetes behaviour more accurately
|
||||
projectRef,
|
||||
deploymentVersion,
|
||||
}: CheckpointAndPushOptions): Promise<CheckpointAndPushResult> {
|
||||
await this.init();
|
||||
|
||||
const options = {
|
||||
runId,
|
||||
leaveRunning,
|
||||
projectRef,
|
||||
deploymentVersion,
|
||||
};
|
||||
|
||||
if (!this.#dockerMode && !this.#canCheckpoint) {
|
||||
this.#logger.error("No checkpoint support. Simulation requires docker.");
|
||||
return { success: false, reason: "NO_SUPPORT" };
|
||||
}
|
||||
|
||||
if (this.isCheckpointing(runId)) {
|
||||
this.#logger.error("Checkpoint procedure already in progress", { options });
|
||||
return { success: false, reason: "IN_PROGRESS" };
|
||||
}
|
||||
|
||||
// This is a new checkpoint, clear any last failure for this run
|
||||
this.#clearFailedCheckpoint(runId);
|
||||
|
||||
if (this.disableCheckpointSupport) {
|
||||
this.#logger.error("Checkpoint support disabled", { options });
|
||||
return { success: false, reason: "DISABLED" };
|
||||
}
|
||||
|
||||
const controller = new AbortController();
|
||||
this.#abortControllers.set(runId, controller);
|
||||
|
||||
const $$ = $({ signal: controller.signal });
|
||||
|
||||
const shortCode = nanoid(8);
|
||||
const imageRef = this.#getImageRef(projectRef, deploymentVersion, shortCode);
|
||||
const exportLocation = this.#getExportLocation(projectRef, deploymentVersion, shortCode);
|
||||
|
||||
const cleanup = async () => {
|
||||
if (this.#dockerMode) {
|
||||
return;
|
||||
}
|
||||
|
||||
try {
|
||||
await $`rm ${exportLocation}`;
|
||||
this.#logger.log("Deleted checkpoint archive", { exportLocation });
|
||||
|
||||
await $`buildah rmi ${imageRef}`;
|
||||
this.#logger.log("Deleted checkpoint image", { imageRef });
|
||||
} catch (error) {
|
||||
this.#logger.error("Failure during checkpoint cleanup", { exportLocation, error });
|
||||
}
|
||||
};
|
||||
|
||||
try {
|
||||
await this.chaosMonkey.call({ $: $$ });
|
||||
|
||||
this.#logger.log("Checkpointing:", { options });
|
||||
|
||||
const containterName = this.#getRunContainerName(runId);
|
||||
|
||||
// Create checkpoint (docker)
|
||||
if (this.#dockerMode) {
|
||||
try {
|
||||
if (this.opts.forceSimulate || !this.#canCheckpoint) {
|
||||
this.#logger.log("Simulating checkpoint");
|
||||
this.#logger.debug(await $$`docker pause ${containterName}`);
|
||||
} else {
|
||||
if (this.simulateCheckpointFailure) {
|
||||
if (performance.now() < this.simulateCheckpointFailureSeconds * 1000) {
|
||||
this.#logger.error("Simulating checkpoint failure", { options });
|
||||
throw new Error("SIMULATE_CHECKPOINT_FAILURE");
|
||||
}
|
||||
}
|
||||
|
||||
if (leaveRunning) {
|
||||
this.#logger.debug(
|
||||
await $$`docker checkpoint create --leave-running ${containterName} ${exportLocation}`
|
||||
);
|
||||
} else {
|
||||
this.#logger.debug(
|
||||
await $$`docker checkpoint create ${containterName} ${exportLocation}`
|
||||
);
|
||||
}
|
||||
}
|
||||
} catch (error) {
|
||||
this.#logger.error("Failed while creating docker checkpoint", { exportLocation });
|
||||
throw error;
|
||||
}
|
||||
|
||||
this.#logger.log("checkpoint created:", {
|
||||
runId,
|
||||
location: exportLocation,
|
||||
});
|
||||
|
||||
return {
|
||||
success: true,
|
||||
checkpoint: {
|
||||
location: exportLocation,
|
||||
docker: true,
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
// Create checkpoint (CRI)
|
||||
if (!this.#canCheckpoint) {
|
||||
this.#logger.error("No checkpoint support in kubernetes mode.");
|
||||
return { success: false, reason: "SKIP_RETRYING" };
|
||||
}
|
||||
|
||||
const containerId = this.#logger.debug(
|
||||
// @ts-expect-error
|
||||
await $$`crictl ps`
|
||||
.pipeStdout($$({ stdin: "pipe" })`grep ${containterName}`)
|
||||
.pipeStdout($$({ stdin: "pipe" })`cut -f1 ${"-d "}`)
|
||||
);
|
||||
|
||||
if (!containerId.stdout) {
|
||||
this.#logger.error("could not find container id", { options, containterName });
|
||||
return { success: false, reason: "SKIP_RETRYING" };
|
||||
}
|
||||
|
||||
const start = performance.now();
|
||||
|
||||
if (this.simulateCheckpointFailure) {
|
||||
if (performance.now() < this.simulateCheckpointFailureSeconds * 1000) {
|
||||
this.#logger.error("Simulating checkpoint failure", { options });
|
||||
throw new Error("SIMULATE_CHECKPOINT_FAILURE");
|
||||
}
|
||||
}
|
||||
|
||||
// Create checkpoint
|
||||
this.#logger.debug(await $$`crictl checkpoint --export=${exportLocation} ${containerId}`);
|
||||
const postCheckpoint = performance.now();
|
||||
|
||||
// Print checkpoint size
|
||||
const size = await getParsedFileSize(exportLocation);
|
||||
this.#logger.log("checkpoint archive created", { size, options });
|
||||
|
||||
// Create image from checkpoint
|
||||
const container = this.#logger.debug(await $$`buildah from scratch`);
|
||||
const postFrom = performance.now();
|
||||
|
||||
this.#logger.debug(await $$`buildah add ${container} ${exportLocation} /`);
|
||||
const postAdd = performance.now();
|
||||
|
||||
this.#logger.debug(
|
||||
await $$`buildah config --annotation=io.kubernetes.cri-o.annotations.checkpoint.name=counter ${container}`
|
||||
);
|
||||
const postConfig = performance.now();
|
||||
|
||||
this.#logger.debug(await $$`buildah commit ${container} ${imageRef}`);
|
||||
const postCommit = performance.now();
|
||||
|
||||
this.#logger.debug(await $$`buildah rm ${container}`);
|
||||
const postRm = performance.now();
|
||||
|
||||
if (this.simulatePushFailure) {
|
||||
if (performance.now() < this.simulatePushFailureSeconds * 1000) {
|
||||
this.#logger.error("Simulating push failure", { options });
|
||||
throw new Error("SIMULATE_PUSH_FAILURE");
|
||||
}
|
||||
}
|
||||
|
||||
// Push checkpoint image
|
||||
this.#logger.debug(
|
||||
await $$`buildah push --tls-verify=${String(this.registryTlsVerify)} ${imageRef}`
|
||||
);
|
||||
const postPush = performance.now();
|
||||
|
||||
const perf = {
|
||||
"crictl checkpoint": postCheckpoint - start,
|
||||
"buildah from": postFrom - postCheckpoint,
|
||||
"buildah add": postAdd - postFrom,
|
||||
"buildah config": postConfig - postAdd,
|
||||
"buildah commit": postCommit - postConfig,
|
||||
"buildah rm": postRm - postCommit,
|
||||
"buildah push": postPush - postRm,
|
||||
};
|
||||
|
||||
this.#logger.log("Checkpointed and pushed image to:", { location: imageRef, perf });
|
||||
|
||||
return {
|
||||
success: true,
|
||||
checkpoint: {
|
||||
location: imageRef,
|
||||
docker: false,
|
||||
},
|
||||
};
|
||||
} catch (error) {
|
||||
if (isExecaChildProcess(error)) {
|
||||
if (error.isCanceled) {
|
||||
this.#logger.error("Checkpoint canceled", { options, error });
|
||||
|
||||
return { success: false, reason: "CANCELED" };
|
||||
}
|
||||
|
||||
this.#logger.error("Checkpoint command error", { options, error });
|
||||
|
||||
return { success: false, reason: "ERROR" };
|
||||
}
|
||||
|
||||
this.#logger.error("Unhandled checkpoint error", { options, error });
|
||||
|
||||
return { success: false, reason: "ERROR" };
|
||||
} finally {
|
||||
this.#abortControllers.delete(runId);
|
||||
await cleanup();
|
||||
}
|
||||
}
|
||||
|
||||
#failCheckpoint(runId: string, error: unknown) {
|
||||
this.#failedCheckpoints.set(runId, error);
|
||||
}
|
||||
|
||||
#clearFailedCheckpoint(runId: string) {
|
||||
this.#failedCheckpoints.delete(runId);
|
||||
}
|
||||
|
||||
#hasFailedCheckpoint(runId: string) {
|
||||
return this.#failedCheckpoints.has(runId);
|
||||
}
|
||||
|
||||
#getRunContainerName(suffix: string) {
|
||||
return `task-run-${suffix}`;
|
||||
}
|
||||
}
|
||||
+396
-375
File diff suppressed because it is too large
Load Diff
@@ -4,6 +4,8 @@ PLATFORM_WS_PORT=3030
|
||||
PLATFORM_SECRET=provider-secret
|
||||
SECURE_CONNECTION=false
|
||||
|
||||
OTEL_EXPORTER_OTLP_ENDPOINT=http://0.0.0.0:3030/otel
|
||||
|
||||
# Use this if you are on macOS
|
||||
# COORDINATOR_HOST="host.docker.internal"
|
||||
# OTEL_EXPORTER_OTLP_ENDPOINT="http://host.docker.internal:4318"
|
||||
@@ -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" ]
|
||||
|
||||
@@ -18,8 +18,7 @@
|
||||
"dependencies": {
|
||||
"@trigger.dev/core": "workspace:*",
|
||||
"@trigger.dev/core-apps": "workspace:*",
|
||||
"execa": "^8.0.1",
|
||||
"socket.io-client": "^4.7.4"
|
||||
"execa": "^8.0.1"
|
||||
},
|
||||
"devDependencies": {
|
||||
"@types/node": "^18.19.8",
|
||||
|
||||
@@ -1,87 +1,84 @@
|
||||
import { $, type ExecaChildProcess, execa } from "execa";
|
||||
import {
|
||||
SimpleLogger,
|
||||
TaskOperations,
|
||||
ProviderShell,
|
||||
TaskOperationsRestoreOptions,
|
||||
TaskOperations,
|
||||
TaskOperationsCreateOptions,
|
||||
TaskOperationsIndexOptions,
|
||||
} from "@trigger.dev/core-apps";
|
||||
TaskOperationsRestoreOptions,
|
||||
} from "@trigger.dev/core-apps/provider";
|
||||
import { SimpleLogger } from "@trigger.dev/core-apps/logger";
|
||||
import { isExecaChildProcess, testDockerCheckpoint } from "@trigger.dev/core-apps/checkpoints";
|
||||
import { setTimeout } from "node:timers/promises";
|
||||
import { PostStartCauses, PreStopCauses } from "@trigger.dev/core/v3";
|
||||
|
||||
const MACHINE_NAME = process.env.MACHINE_NAME || "local";
|
||||
const COORDINATOR_PORT = process.env.COORDINATOR_PORT || 8020;
|
||||
const COORDINATOR_HOST = process.env.COORDINATOR_HOST || "127.0.0.1";
|
||||
|
||||
const OTEL_EXPORTER_OTLP_ENDPOINT =
|
||||
process.env.OTEL_EXPORTER_OTLP_ENDPOINT || "http://0.0.0.0:4318";
|
||||
|
||||
const FORCE_CHECKPOINT_SIMULATION = ["1", "true"].includes(
|
||||
process.env.FORCE_CHECKPOINT_SIMULATION ?? "true"
|
||||
);
|
||||
|
||||
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.#canCheckpoint = canCheckpoint;
|
||||
|
||||
if (canCheckpoint) {
|
||||
if (!this.#initialized) {
|
||||
logger.log("Full checkpoint support!");
|
||||
}
|
||||
}
|
||||
|
||||
this.#initialized = true;
|
||||
|
||||
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);
|
||||
|
||||
@@ -90,60 +87,51 @@ 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);
|
||||
|
||||
const runArgs = [
|
||||
"run",
|
||||
"--network=host",
|
||||
"--detach",
|
||||
`--env=TRIGGER_ENV_ID=${opts.envId}`,
|
||||
`--env=TRIGGER_RUN_ID=${opts.runId}`,
|
||||
`--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}`,
|
||||
];
|
||||
|
||||
if (process.env.ENFORCE_MACHINE_PRESETS) {
|
||||
runArgs.push(`--cpus=${opts.machine.cpu}`, `--memory=${opts.machine.memory}G`);
|
||||
}
|
||||
|
||||
runArgs.push(`${opts.image}`);
|
||||
|
||||
try {
|
||||
logger.debug(
|
||||
await execa("docker", [
|
||||
"run",
|
||||
"--network=host",
|
||||
"--detach",
|
||||
`--env=TRIGGER_ENV_ID=${opts.envId}`,
|
||||
`--env=TRIGGER_RUN_ID=${opts.runId}`,
|
||||
`--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.image}`,
|
||||
])
|
||||
);
|
||||
logger.debug(await execa("docker", runArgs));
|
||||
} catch (error) {
|
||||
if (!isExecaChildProcess(error)) {
|
||||
throw error;
|
||||
@@ -160,7 +148,7 @@ class DockerTaskOperations implements TaskOperations {
|
||||
}
|
||||
|
||||
async restore(opts: TaskOperationsRestoreOptions) {
|
||||
await this.#initialize();
|
||||
await this.init();
|
||||
|
||||
const containerName = this.#getRunContainerName(opts.runId);
|
||||
|
||||
@@ -189,7 +177,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);
|
||||
@@ -198,7 +186,7 @@ class DockerTaskOperations implements TaskOperations {
|
||||
}
|
||||
|
||||
async get(opts: { runId: string }) {
|
||||
await this.#initialize();
|
||||
await this.init();
|
||||
|
||||
logger.log("noop: get");
|
||||
}
|
||||
@@ -278,7 +266,7 @@ class DockerTaskOperations implements TaskOperations {
|
||||
}
|
||||
|
||||
const provider = new ProviderShell({
|
||||
tasks: new DockerTaskOperations({ forceSimulate: true }),
|
||||
tasks: new DockerTaskOperations({ forceSimulate: FORCE_CHECKPOINT_SIMULATION }),
|
||||
type: "docker",
|
||||
});
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -19,8 +19,7 @@
|
||||
"@kubernetes/client-node": "^0.20.0",
|
||||
"@trigger.dev/core": "workspace:*",
|
||||
"@trigger.dev/core-apps": "workspace:*",
|
||||
"p-queue": "^8.0.1",
|
||||
"socket.io-client": "^4.7.4"
|
||||
"p-queue": "^8.0.1"
|
||||
},
|
||||
"devDependencies": {
|
||||
"dotenv": "^16.4.2",
|
||||
|
||||
@@ -1,22 +1,36 @@
|
||||
import * as k8s from "@kubernetes/client-node";
|
||||
import {
|
||||
ProviderShell,
|
||||
SimpleLogger,
|
||||
TaskOperations,
|
||||
TaskOperationsCreateOptions,
|
||||
TaskOperationsIndexOptions,
|
||||
TaskOperationsRestoreOptions,
|
||||
} from "@trigger.dev/core-apps";
|
||||
import { Machine, PostStartCauses, PreStopCauses, EnvironmentType } from "@trigger.dev/core/v3";
|
||||
} from "@trigger.dev/core-apps/provider";
|
||||
import { SimpleLogger } from "@trigger.dev/core-apps/logger";
|
||||
import {
|
||||
MachinePreset,
|
||||
PostStartCauses,
|
||||
PreStopCauses,
|
||||
EnvironmentType,
|
||||
} from "@trigger.dev/core/v3";
|
||||
import { randomUUID } from "crypto";
|
||||
import { TaskMonitor } from "./taskMonitor";
|
||||
import { PodCleaner } from "./podCleaner";
|
||||
import { UptimeHeartbeat } from "./uptimeHeartbeat";
|
||||
|
||||
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 UPTIME_HEARTBEAT_URL = process.env.UPTIME_HEARTBEAT_URL;
|
||||
const UPTIME_INTERVAL_SECONDS = Number(process.env.UPTIME_INTERVAL_SECONDS || "60");
|
||||
const UPTIME_MAX_PENDING_RUNS = Number(process.env.UPTIME_MAX_PENDING_RUNS || "25");
|
||||
const UPTIME_MAX_PENDING_INDECES = Number(process.env.UPTIME_MAX_PENDING_INDECES || "10");
|
||||
const UPTIME_MAX_PENDING_ERRORS = Number(process.env.UPTIME_MAX_PENDING_ERRORS || "10");
|
||||
|
||||
const logger = new SimpleLogger(`[${NODE_NAME}]`);
|
||||
logger.log(`running in ${RUNTIME_ENV} mode`);
|
||||
|
||||
@@ -47,6 +61,10 @@ class KubernetesTaskOperations implements TaskOperations {
|
||||
this.#k8sApi = this.#createK8sApi();
|
||||
}
|
||||
|
||||
async init() {
|
||||
// noop
|
||||
}
|
||||
|
||||
async index(opts: TaskOperationsIndexOptions) {
|
||||
await this.#createJob(
|
||||
{
|
||||
@@ -212,7 +230,7 @@ class KubernetesTaskOperations implements TaskOperations {
|
||||
},
|
||||
{
|
||||
name: "populate-taskinfo",
|
||||
image: "docker.io/library/busybox",
|
||||
image: "registry.digitalocean.com/trigger/busybox",
|
||||
imagePullPolicy: "IfNotPresent",
|
||||
command: ["/bin/sh", "-c"],
|
||||
args: ["printenv COORDINATOR_HOST | tee /etc/taskinfo/coordinator-host"],
|
||||
@@ -316,6 +334,9 @@ class KubernetesTaskOperations implements TaskOperations {
|
||||
{
|
||||
name: "registry-trigger",
|
||||
},
|
||||
{
|
||||
name: "registry-trigger-failover",
|
||||
},
|
||||
],
|
||||
nodeSelector: {
|
||||
nodetype: "worker",
|
||||
@@ -391,10 +412,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`,
|
||||
};
|
||||
}
|
||||
|
||||
@@ -516,27 +537,28 @@ provider.listen();
|
||||
|
||||
const taskMonitor = new TaskMonitor({
|
||||
runtimeEnv: RUNTIME_ENV,
|
||||
onIndexFailure: async (deploymentId, failureInfo) => {
|
||||
logger.log("Indexing failed", { deploymentId, failureInfo });
|
||||
onIndexFailure: async (deploymentId, details) => {
|
||||
logger.log("Indexing failed", { deploymentId, details });
|
||||
|
||||
try {
|
||||
provider.platformSocket.send("INDEXING_FAILED", {
|
||||
deploymentId,
|
||||
error: {
|
||||
name: `Crashed with exit code ${failureInfo.exitCode}`,
|
||||
message: failureInfo.reason,
|
||||
stack: failureInfo.logs,
|
||||
name: `Crashed with exit code ${details.exitCode}`,
|
||||
message: details.reason,
|
||||
stack: details.logs,
|
||||
},
|
||||
overrideCompletion: details.overrideCompletion,
|
||||
});
|
||||
} catch (error) {
|
||||
logger.error(error);
|
||||
}
|
||||
},
|
||||
onRunFailure: async (runId, failureInfo) => {
|
||||
logger.log("Run failed:", { runId, failureInfo });
|
||||
onRunFailure: async (runId, details) => {
|
||||
logger.log("Run failed:", { runId, details });
|
||||
|
||||
try {
|
||||
provider.platformSocket.send("WORKER_CRASHED", { runId, ...failureInfo });
|
||||
provider.platformSocket.send("WORKER_CRASHED", { runId, ...details });
|
||||
} catch (error) {
|
||||
logger.error(error);
|
||||
}
|
||||
@@ -548,7 +570,23 @@ taskMonitor.start();
|
||||
const podCleaner = new PodCleaner({
|
||||
runtimeEnv: RUNTIME_ENV,
|
||||
namespace: "default",
|
||||
intervalInSeconds: 300,
|
||||
intervalInSeconds: POD_CLEANER_INTERVAL_SECONDS,
|
||||
});
|
||||
|
||||
podCleaner.start();
|
||||
|
||||
if (UPTIME_HEARTBEAT_URL) {
|
||||
const uptimeHeartbeat = new UptimeHeartbeat({
|
||||
runtimeEnv: RUNTIME_ENV,
|
||||
namespace: "default",
|
||||
intervalInSeconds: UPTIME_INTERVAL_SECONDS,
|
||||
pingUrl: UPTIME_HEARTBEAT_URL,
|
||||
maxPendingRuns: UPTIME_MAX_PENDING_RUNS,
|
||||
maxPendingIndeces: UPTIME_MAX_PENDING_INDECES,
|
||||
maxPendingErrors: UPTIME_MAX_PENDING_ERRORS,
|
||||
});
|
||||
|
||||
uptimeHeartbeat.start();
|
||||
} else {
|
||||
logger.log("Uptime heartbeat is disabled, set UPTIME_HEARTBEAT_URL to enable.");
|
||||
}
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
import * as k8s from "@kubernetes/client-node";
|
||||
import { SimpleLogger } from "@trigger.dev/core-apps";
|
||||
import { SimpleLogger } from "@trigger.dev/core-apps/logger";
|
||||
|
||||
type PodCleanerOptions = {
|
||||
runtimeEnv: "local" | "kubernetes";
|
||||
|
||||
@@ -1,25 +1,20 @@
|
||||
import * as k8s from "@kubernetes/client-node";
|
||||
import { SimpleLogger } from "@trigger.dev/core-apps";
|
||||
import { SimpleLogger } from "@trigger.dev/core-apps/logger";
|
||||
import { EXIT_CODE_ALREADY_HANDLED, EXIT_CODE_CHILD_NONZERO } from "@trigger.dev/core-apps/process";
|
||||
import { setTimeout } from "timers/promises";
|
||||
import PQueue from "p-queue";
|
||||
import type { Prettify } from "@trigger.dev/core/v3";
|
||||
|
||||
type IndexFailureHandler = (
|
||||
deploymentId: string,
|
||||
failureInfo: {
|
||||
exitCode: number;
|
||||
reason: string;
|
||||
logs: string;
|
||||
}
|
||||
) => Promise<any>;
|
||||
type FailureDetails = Prettify<{
|
||||
exitCode: number;
|
||||
reason: string;
|
||||
logs: string;
|
||||
overrideCompletion: boolean;
|
||||
}>;
|
||||
|
||||
type RunFailureHandler = (
|
||||
runId: string,
|
||||
failureInfo: {
|
||||
exitCode: number;
|
||||
reason: string;
|
||||
logs: string;
|
||||
}
|
||||
) => Promise<any>;
|
||||
type IndexFailureHandler = (deploymentId: string, details: FailureDetails) => Promise<any>;
|
||||
|
||||
type RunFailureHandler = (runId: string, details: FailureDetails) => Promise<any>;
|
||||
|
||||
type TaskMonitorOptions = {
|
||||
runtimeEnv: "local" | "kubernetes";
|
||||
@@ -144,8 +139,10 @@ export class TaskMonitor {
|
||||
const containerState = this.#getContainerStateSummary(containerStatus.state);
|
||||
const exitCode = containerState.exitCode ?? -1;
|
||||
|
||||
// We use this special exit code to signal any errors were already handled elsewhere
|
||||
if (exitCode === 111) {
|
||||
if (exitCode === EXIT_CODE_ALREADY_HANDLED) {
|
||||
this.#logger.debug("Ignoring pod failure, already handled by worker", {
|
||||
podName,
|
||||
});
|
||||
return;
|
||||
}
|
||||
|
||||
@@ -162,6 +159,7 @@ export class TaskMonitor {
|
||||
|
||||
let reason = rawReason || "Unknown error";
|
||||
let logs = rawLogs || "";
|
||||
let overrideCompletion = false;
|
||||
|
||||
switch (rawReason) {
|
||||
case "Error":
|
||||
@@ -181,7 +179,10 @@ export class TaskMonitor {
|
||||
}
|
||||
break;
|
||||
case "OOMKilled":
|
||||
reason = "Out of memory! Try increasing the memory on this task.";
|
||||
overrideCompletion = true;
|
||||
reason = `${
|
||||
exitCode === EXIT_CODE_CHILD_NONZERO ? "Child process" : "Parent process"
|
||||
} ran out of memory! Try choosing a machine preset with more memory for this task.`;
|
||||
break;
|
||||
default:
|
||||
break;
|
||||
@@ -191,7 +192,8 @@ export class TaskMonitor {
|
||||
exitCode,
|
||||
reason,
|
||||
logs,
|
||||
};
|
||||
overrideCompletion,
|
||||
} satisfies FailureDetails;
|
||||
|
||||
const app = pod.metadata?.labels?.app;
|
||||
|
||||
|
||||
@@ -0,0 +1,272 @@
|
||||
import * as k8s from "@kubernetes/client-node";
|
||||
import { SimpleLogger } from "@trigger.dev/core-apps/logger";
|
||||
|
||||
type UptimeHeartbeatOptions = {
|
||||
runtimeEnv: "local" | "kubernetes";
|
||||
pingUrl: string;
|
||||
namespace?: string;
|
||||
intervalInSeconds?: number;
|
||||
maxPendingRuns?: number;
|
||||
maxPendingIndeces?: number;
|
||||
maxPendingErrors?: number;
|
||||
leadingEdge?: boolean;
|
||||
};
|
||||
|
||||
export class UptimeHeartbeat {
|
||||
private enabled = false;
|
||||
private namespace: string;
|
||||
|
||||
private intervalInSeconds: number;
|
||||
private maxPendingRuns: number;
|
||||
private maxPendingIndeces: number;
|
||||
private maxPendingErrors: number;
|
||||
|
||||
private leadingEdge = true;
|
||||
|
||||
private logger = new SimpleLogger("[UptimeHeartbeat]");
|
||||
private k8sClient: {
|
||||
core: k8s.CoreV1Api;
|
||||
kubeConfig: k8s.KubeConfig;
|
||||
};
|
||||
|
||||
constructor(private opts: UptimeHeartbeatOptions) {
|
||||
this.namespace = opts.namespace ?? "default";
|
||||
|
||||
this.intervalInSeconds = opts.intervalInSeconds ?? 60;
|
||||
this.maxPendingRuns = opts.maxPendingRuns ?? 25;
|
||||
this.maxPendingIndeces = opts.maxPendingIndeces ?? 10;
|
||||
this.maxPendingErrors = opts.maxPendingErrors ?? 10;
|
||||
|
||||
this.k8sClient = this.#createK8sClient();
|
||||
}
|
||||
|
||||
#createK8sClient() {
|
||||
const kubeConfig = new k8s.KubeConfig();
|
||||
|
||||
if (this.opts.runtimeEnv === "local") {
|
||||
kubeConfig.loadFromDefault();
|
||||
} else if (this.opts.runtimeEnv === "kubernetes") {
|
||||
kubeConfig.loadFromCluster();
|
||||
} else {
|
||||
throw new Error(`Unsupported runtime environment: ${this.opts.runtimeEnv}`);
|
||||
}
|
||||
|
||||
return {
|
||||
core: kubeConfig.makeApiClient(k8s.CoreV1Api),
|
||||
kubeConfig: kubeConfig,
|
||||
};
|
||||
}
|
||||
|
||||
#isRecord(candidate: unknown): candidate is Record<string, unknown> {
|
||||
if (typeof candidate !== "object" || candidate === null) {
|
||||
return false;
|
||||
} else {
|
||||
return true;
|
||||
}
|
||||
}
|
||||
|
||||
#logK8sError(err: unknown, debugOnly = false) {
|
||||
if (debugOnly) {
|
||||
this.logger.debug("K8s API Error", err);
|
||||
} else {
|
||||
this.logger.error("K8s API Error", err);
|
||||
}
|
||||
}
|
||||
|
||||
#handleK8sError(err: unknown) {
|
||||
if (!this.#isRecord(err) || !this.#isRecord(err.body)) {
|
||||
this.#logK8sError(err);
|
||||
return;
|
||||
}
|
||||
|
||||
this.#logK8sError(err, true);
|
||||
|
||||
if (typeof err.body.message === "string") {
|
||||
this.#logK8sError({ message: err.body.message });
|
||||
return;
|
||||
}
|
||||
|
||||
this.#logK8sError({ body: err.body });
|
||||
}
|
||||
|
||||
async #getPods(opts: {
|
||||
namespace: string;
|
||||
fieldSelector?: string;
|
||||
labelSelector?: string;
|
||||
}): Promise<Array<k8s.V1Pod> | undefined> {
|
||||
const listReturn = await this.k8sClient.core
|
||||
.listNamespacedPod(
|
||||
opts.namespace,
|
||||
undefined, // pretty
|
||||
undefined, // allowWatchBookmarks
|
||||
undefined, // _continue
|
||||
opts.fieldSelector,
|
||||
opts.labelSelector,
|
||||
this.maxPendingRuns * 2, // limit
|
||||
undefined, // resourceVersion
|
||||
undefined, // resourceVersionMatch
|
||||
undefined, // sendInitialEvents
|
||||
this.intervalInSeconds, // timeoutSeconds,
|
||||
undefined // watch
|
||||
)
|
||||
.catch(this.#handleK8sError.bind(this));
|
||||
|
||||
return listReturn?.body.items;
|
||||
}
|
||||
|
||||
async #getPendingIndeces(): Promise<Array<k8s.V1Pod> | undefined> {
|
||||
return await this.#getPods({
|
||||
namespace: this.namespace,
|
||||
fieldSelector: "status.phase=Pending",
|
||||
labelSelector: "app=task-index",
|
||||
});
|
||||
}
|
||||
|
||||
async #getPendingTasks(): Promise<Array<k8s.V1Pod> | undefined> {
|
||||
return await this.#getPods({
|
||||
namespace: this.namespace,
|
||||
fieldSelector: "status.phase=Pending",
|
||||
labelSelector: "app=task-run",
|
||||
});
|
||||
}
|
||||
|
||||
#countPods(pods: Array<k8s.V1Pod>): number {
|
||||
return pods.length;
|
||||
}
|
||||
|
||||
#filterPendingPods(
|
||||
pods: Array<k8s.V1Pod>,
|
||||
waitingReason: "CreateContainerError" | "RunContainerError"
|
||||
): Array<k8s.V1Pod> {
|
||||
return pods.filter((pod) => {
|
||||
const containerStatus = pod.status?.containerStatuses?.[0];
|
||||
return containerStatus?.state?.waiting?.reason === waitingReason;
|
||||
});
|
||||
}
|
||||
|
||||
async #sendPing() {
|
||||
this.logger.log("Sending ping");
|
||||
|
||||
const start = Date.now();
|
||||
const controller = new AbortController();
|
||||
|
||||
const timeoutMs = (this.intervalInSeconds * 1000) / 2;
|
||||
|
||||
const fetchTimeout = setTimeout(() => {
|
||||
controller.abort();
|
||||
}, timeoutMs);
|
||||
|
||||
try {
|
||||
const response = await fetch(this.opts.pingUrl, {
|
||||
signal: controller.signal,
|
||||
});
|
||||
|
||||
if (!response.ok) {
|
||||
this.logger.error("Failed to send ping, response not OK", {
|
||||
status: response.status,
|
||||
});
|
||||
return;
|
||||
}
|
||||
|
||||
const elapsedMs = Date.now() - start;
|
||||
this.logger.log("Ping sent", { elapsedMs });
|
||||
} catch (error) {
|
||||
if (error instanceof DOMException && error.name === "AbortError") {
|
||||
this.logger.log("Ping timeout", { timeoutSeconds: timeoutMs });
|
||||
return;
|
||||
}
|
||||
|
||||
this.logger.error("Failed to send ping", error);
|
||||
} finally {
|
||||
clearTimeout(fetchTimeout);
|
||||
}
|
||||
}
|
||||
|
||||
async #heartbeat() {
|
||||
this.logger.log("Performing heartbeat");
|
||||
|
||||
const start = Date.now();
|
||||
|
||||
const pendingTasks = await this.#getPendingTasks();
|
||||
|
||||
if (!pendingTasks) {
|
||||
this.logger.error("Failed to get pending tasks");
|
||||
return;
|
||||
}
|
||||
|
||||
const totalPendingTasks = this.#countPods(pendingTasks);
|
||||
|
||||
const pendingIndeces = await this.#getPendingIndeces();
|
||||
|
||||
if (!pendingIndeces) {
|
||||
this.logger.error("Failed to get pending indeces");
|
||||
return;
|
||||
}
|
||||
|
||||
const totalPendingIndeces = this.#countPods(pendingIndeces);
|
||||
|
||||
const elapsedMs = Date.now() - start;
|
||||
|
||||
this.logger.log("Finished heartbeat checks", { elapsedMs });
|
||||
|
||||
if (totalPendingTasks > this.maxPendingRuns) {
|
||||
this.logger.log("Too many pending tasks, skipping heartbeat", { totalPendingTasks });
|
||||
return;
|
||||
}
|
||||
|
||||
if (totalPendingIndeces > this.maxPendingIndeces) {
|
||||
this.logger.log("Too many pending indeces, skipping heartbeat", { totalPendingIndeces });
|
||||
return;
|
||||
}
|
||||
|
||||
const totalCreateContainerErrors = this.#countPods(
|
||||
this.#filterPendingPods(pendingTasks, "CreateContainerError")
|
||||
);
|
||||
const totalRunContainerErrors = this.#countPods(
|
||||
this.#filterPendingPods(pendingTasks, "RunContainerError")
|
||||
);
|
||||
|
||||
if (totalCreateContainerErrors + totalRunContainerErrors > this.maxPendingErrors) {
|
||||
this.logger.log("Too many pending tasks with errors, skipping heartbeat", {
|
||||
totalRunContainerErrors,
|
||||
totalCreateContainerErrors,
|
||||
});
|
||||
return;
|
||||
}
|
||||
|
||||
await this.#sendPing();
|
||||
|
||||
this.logger.log("Heartbeat done", { totalPendingTasks, elapsedMs });
|
||||
}
|
||||
|
||||
async start() {
|
||||
this.enabled = true;
|
||||
this.logger.log("Starting");
|
||||
|
||||
if (this.leadingEdge) {
|
||||
await this.#heartbeat();
|
||||
}
|
||||
|
||||
const heartbeat = setInterval(async () => {
|
||||
if (!this.enabled) {
|
||||
clearInterval(heartbeat);
|
||||
return;
|
||||
}
|
||||
|
||||
try {
|
||||
await this.#heartbeat();
|
||||
} catch (error) {
|
||||
this.logger.error("Error while heartbeating", error);
|
||||
}
|
||||
}, this.intervalInSeconds * 1000);
|
||||
}
|
||||
|
||||
async stop() {
|
||||
if (!this.enabled) {
|
||||
return;
|
||||
}
|
||||
|
||||
this.enabled = false;
|
||||
this.logger.log("Shutting down..");
|
||||
}
|
||||
}
|
||||
@@ -7,9 +7,9 @@
|
||||
"dev": "wrangler dev"
|
||||
},
|
||||
"devDependencies": {
|
||||
"@cloudflare/workers-types": "^4.20230419.0",
|
||||
"@cloudflare/workers-types": "^4.20240512.0",
|
||||
"typescript": "^5.0.4",
|
||||
"wrangler": "^3.0.0"
|
||||
"wrangler": "^3.57.1"
|
||||
},
|
||||
"dependencies": {
|
||||
"@aws-sdk/client-sqs": "^3.445.0",
|
||||
|
||||
@@ -1,5 +1,7 @@
|
||||
import { queueEvent } from "./events/queueEvent";
|
||||
import { queueEvents } from "./events/queueEvents";
|
||||
import { applyRateLimit } from "./rateLimit";
|
||||
import { Ratelimit } from "./rateLimiter";
|
||||
|
||||
export interface Env {
|
||||
/** The hostname needs to be changed to allow requests to pass to the Trigger.dev platform */
|
||||
@@ -9,6 +11,8 @@ export interface Env {
|
||||
AWS_SQS_SECRET_ACCESS_KEY: string;
|
||||
AWS_SQS_QUEUE_URL: string;
|
||||
AWS_SQS_REGION: string;
|
||||
//rate limiter
|
||||
API_RATE_LIMITER: Ratelimit;
|
||||
}
|
||||
|
||||
export default {
|
||||
@@ -25,13 +29,13 @@ export default {
|
||||
switch (url.pathname) {
|
||||
case "/api/v1/events": {
|
||||
if (request.method === "POST") {
|
||||
return queueEvent(request, env);
|
||||
return applyRateLimit(request, env, () => queueEvent(request, env));
|
||||
}
|
||||
break;
|
||||
}
|
||||
case "/api/v1/events/bulk": {
|
||||
if (request.method === "POST") {
|
||||
return queueEvents(request, env);
|
||||
return applyRateLimit(request, env, () => queueEvents(request, env));
|
||||
}
|
||||
break;
|
||||
}
|
||||
|
||||
@@ -0,0 +1,46 @@
|
||||
import { Env } from "src";
|
||||
import { getApiKeyFromRequest } from "./apikey";
|
||||
import { json } from "./json";
|
||||
|
||||
export async function applyRateLimit(
|
||||
request: Request,
|
||||
env: Env,
|
||||
fn: () => Promise<Response>
|
||||
): Promise<Response> {
|
||||
const apiKey = getApiKeyFromRequest(request);
|
||||
if (apiKey) {
|
||||
const result = await env.API_RATE_LIMITER.limit({ key: `apikey-${apiKey.apiKey}` });
|
||||
const { success } = result;
|
||||
console.log(`Rate limiter`, {
|
||||
success,
|
||||
key: `${apiKey.apiKey.substring(0, 12)}...`,
|
||||
});
|
||||
if (!success) {
|
||||
//60s in the future
|
||||
const reset = Date.now() + 60 * 1000;
|
||||
const secondsUntilReset = Math.max(0, (reset - new Date().getTime()) / 1000);
|
||||
|
||||
return json(
|
||||
{
|
||||
title: "Rate Limit Exceeded",
|
||||
status: 429,
|
||||
type: "https://developer.mozilla.org/en-US/docs/Web/HTTP/Status/429",
|
||||
detail: `Rate limit exceeded. Retry in ${secondsUntilReset} seconds.`,
|
||||
error: `Rate limit exceeded. Retry in ${secondsUntilReset} seconds.`,
|
||||
reset,
|
||||
},
|
||||
{
|
||||
status: 429,
|
||||
headers: {
|
||||
"x-ratelimit-reset": reset.toString(),
|
||||
},
|
||||
}
|
||||
);
|
||||
}
|
||||
} else {
|
||||
console.log(`Rate limiter: no API key for request`);
|
||||
}
|
||||
|
||||
//call the original function
|
||||
return fn();
|
||||
}
|
||||
@@ -0,0 +1,23 @@
|
||||
export interface Ratelimit {
|
||||
/*
|
||||
* The ratelimit function
|
||||
* @param {RatelimitOptions} options
|
||||
* @returns {Promise<RatelimitResponse>}
|
||||
*/
|
||||
limit: (options: RatelimitOptions) => Promise<RatelimitResponse>;
|
||||
}
|
||||
|
||||
export interface RatelimitOptions {
|
||||
/*
|
||||
* The key to identify the user, can be an IP address, user ID, etc.
|
||||
*/
|
||||
key: string;
|
||||
}
|
||||
|
||||
export interface RatelimitResponse {
|
||||
/*
|
||||
* The ratelimit success status
|
||||
* @returns {boolean}
|
||||
*/
|
||||
success: boolean;
|
||||
}
|
||||
@@ -1,7 +1,33 @@
|
||||
name = "proxy"
|
||||
main = "src/index.ts"
|
||||
compatibility_date = "2023-10-30"
|
||||
compatibility_date = "2024-05-13"
|
||||
compatibility_flags = [ "nodejs_compat" ]
|
||||
|
||||
[env.staging]
|
||||
[env.prod]
|
||||
# The rate limiting API is in open beta.
|
||||
[[env.staging.unsafe.bindings]]
|
||||
name = "API_RATE_LIMITER"
|
||||
type = "ratelimit"
|
||||
# An identifier you define, that is unique to your Cloudflare account.
|
||||
# Must be an integer.
|
||||
namespace_id = "1"
|
||||
|
||||
# Limit: the number of tokens allowed within a given period in a single
|
||||
# Cloudflare location
|
||||
# Period: the duration of the period, in seconds. Must be either 10 or 60
|
||||
simple = { limit = 100, period = 60 }
|
||||
|
||||
|
||||
[env.prod]
|
||||
# The rate limiting API is in open beta.
|
||||
[[env.prod.unsafe.bindings]]
|
||||
name = "API_RATE_LIMITER"
|
||||
type = "ratelimit"
|
||||
# An identifier you define, that is unique to your Cloudflare account.
|
||||
# Must be an integer.
|
||||
namespace_id = "2"
|
||||
|
||||
# Limit: the number of tokens allowed within a given period in a single
|
||||
# Cloudflare location
|
||||
# Period: the duration of the period, in seconds. Must be either 10 or 60
|
||||
simple = { limit = 300, period = 60 }
|
||||
@@ -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;
|
||||
|
||||
@@ -528,7 +528,7 @@ export type Tree<TData> = {
|
||||
/** A tree but flattened so it can easily be used for DOM elements */
|
||||
export type FlatTreeItem<TData> = {
|
||||
id: string;
|
||||
parentId: string | undefined;
|
||||
parentId?: string | undefined;
|
||||
children: string[];
|
||||
hasChildren: boolean;
|
||||
/** The indentation level, the root is 0 */
|
||||
|
||||
@@ -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>
|
||||
);
|
||||
}
|
||||
@@ -146,7 +146,6 @@ function FilterMenu(props: RunFiltersProps) {
|
||||
|
||||
const filterTrigger = (
|
||||
<SelectTrigger
|
||||
autoFocus
|
||||
icon={
|
||||
<div className="flex size-4 items-center justify-center">
|
||||
<ListFilterIcon className="size-3.5" />
|
||||
|
||||
@@ -3,10 +3,12 @@ import {
|
||||
BoltSlashIcon,
|
||||
BugAntIcon,
|
||||
CheckCircleIcon,
|
||||
ClockIcon,
|
||||
FireIcon,
|
||||
NoSymbolIcon,
|
||||
PauseCircleIcon,
|
||||
RectangleStackIcon,
|
||||
TrashIcon,
|
||||
XCircleIcon,
|
||||
} from "@heroicons/react/20/solid";
|
||||
import { TaskRunStatus } from "@trigger.dev/database";
|
||||
@@ -16,6 +18,7 @@ import { Spinner } from "~/components/primitives/Spinner";
|
||||
import { cn } from "~/utils/cn";
|
||||
|
||||
export const allTaskRunStatuses = [
|
||||
"DELAYED",
|
||||
"WAITING_FOR_DEPLOY",
|
||||
"PENDING",
|
||||
"EXECUTING",
|
||||
@@ -28,10 +31,12 @@ export const allTaskRunStatuses = [
|
||||
"PAUSED",
|
||||
"INTERRUPTED",
|
||||
"SYSTEM_FAILURE",
|
||||
"EXPIRED",
|
||||
] as const satisfies Readonly<Array<TaskRunStatus>>;
|
||||
|
||||
export const filterableTaskRunStatuses = [
|
||||
"WAITING_FOR_DEPLOY",
|
||||
"DELAYED",
|
||||
"PENDING",
|
||||
"EXECUTING",
|
||||
"RETRYING_AFTER_FAILURE",
|
||||
@@ -42,9 +47,11 @@ export const filterableTaskRunStatuses = [
|
||||
"CRASHED",
|
||||
"INTERRUPTED",
|
||||
"SYSTEM_FAILURE",
|
||||
"EXPIRED",
|
||||
] as const satisfies Readonly<Array<TaskRunStatus>>;
|
||||
|
||||
const taskRunStatusDescriptions: Record<TaskRunStatus, string> = {
|
||||
DELAYED: "Task has been delayed and is waiting to be executed",
|
||||
PENDING: "Task is waiting to be executed",
|
||||
WAITING_FOR_DEPLOY: "Task needs to be deployed first to start executing",
|
||||
EXECUTING: "Task is currently being executed",
|
||||
@@ -57,9 +64,10 @@ const taskRunStatusDescriptions: Record<TaskRunStatus, string> = {
|
||||
SYSTEM_FAILURE: "Task has failed due to a system failure",
|
||||
PAUSED: "Task has been paused by the user",
|
||||
CRASHED: "Task has crashed and won't be retried",
|
||||
EXPIRED: "Task has surpassed its ttl and won't be executed",
|
||||
};
|
||||
|
||||
export const QUEUED_STATUSES: TaskRunStatus[] = ["PENDING", "WAITING_FOR_DEPLOY"];
|
||||
export const QUEUED_STATUSES: TaskRunStatus[] = ["PENDING", "WAITING_FOR_DEPLOY", "DELAYED"];
|
||||
|
||||
export const RUNNING_STATUSES: TaskRunStatus[] = [
|
||||
"EXECUTING",
|
||||
@@ -74,6 +82,7 @@ export const FINISHED_STATUSES: TaskRunStatus[] = [
|
||||
"INTERRUPTED",
|
||||
"SYSTEM_FAILURE",
|
||||
"CRASHED",
|
||||
"EXPIRED",
|
||||
];
|
||||
|
||||
export function descriptionForTaskRunStatus(status: TaskRunStatus): string {
|
||||
@@ -109,6 +118,8 @@ export function TaskRunStatusIcon({
|
||||
className: string;
|
||||
}) {
|
||||
switch (status) {
|
||||
case "DELAYED":
|
||||
return <ClockIcon className={cn(runStatusClassNameColor(status), className)} />;
|
||||
case "PENDING":
|
||||
return <RectangleStackIcon className={cn(runStatusClassNameColor(status), className)} />;
|
||||
case "WAITING_FOR_DEPLOY":
|
||||
@@ -133,6 +144,8 @@ export function TaskRunStatusIcon({
|
||||
return <BugAntIcon className={cn(runStatusClassNameColor(status), className)} />;
|
||||
case "CRASHED":
|
||||
return <FireIcon className={cn(runStatusClassNameColor(status), className)} />;
|
||||
case "EXPIRED":
|
||||
return <TrashIcon className={cn(runStatusClassNameColor(status), className)} />;
|
||||
|
||||
default: {
|
||||
assertNever(status);
|
||||
@@ -143,6 +156,7 @@ export function TaskRunStatusIcon({
|
||||
export function runStatusClassNameColor(status: TaskRunStatus): string {
|
||||
switch (status) {
|
||||
case "PENDING":
|
||||
case "DELAYED":
|
||||
return "text-charcoal-500";
|
||||
case "WAITING_FOR_DEPLOY":
|
||||
return "text-amber-500";
|
||||
@@ -154,6 +168,7 @@ export function runStatusClassNameColor(status: TaskRunStatus): string {
|
||||
case "PAUSED":
|
||||
return "text-amber-300";
|
||||
case "CANCELED":
|
||||
case "EXPIRED":
|
||||
return "text-charcoal-500";
|
||||
case "INTERRUPTED":
|
||||
return "text-error";
|
||||
@@ -173,6 +188,8 @@ export function runStatusClassNameColor(status: TaskRunStatus): string {
|
||||
|
||||
export function runStatusTitle(status: TaskRunStatus): string {
|
||||
switch (status) {
|
||||
case "DELAYED":
|
||||
return "Delayed";
|
||||
case "PENDING":
|
||||
return "Queued";
|
||||
case "WAITING_FOR_DEPLOY":
|
||||
@@ -197,6 +214,8 @@ export function runStatusTitle(status: TaskRunStatus): string {
|
||||
return "System failure";
|
||||
case "CRASHED":
|
||||
return "Crashed";
|
||||
case "EXPIRED":
|
||||
return "Expired";
|
||||
default: {
|
||||
assertNever(status);
|
||||
}
|
||||
|
||||
@@ -118,6 +118,8 @@ export function TaskRunsTable({
|
||||
<TableHeaderCell>Duration</TableHeaderCell>
|
||||
<TableHeaderCell>Test</TableHeaderCell>
|
||||
<TableHeaderCell>Created at</TableHeaderCell>
|
||||
<TableHeaderCell>Delayed until</TableHeaderCell>
|
||||
<TableHeaderCell>TTL</TableHeaderCell>
|
||||
<TableHeaderCell>
|
||||
<span className="sr-only">Go to page</span>
|
||||
</TableHeaderCell>
|
||||
@@ -125,7 +127,7 @@ export function TaskRunsTable({
|
||||
</TableHeader>
|
||||
<TableBody>
|
||||
{total === 0 && !hasFilters ? (
|
||||
<TableBlankRow colSpan={9}>
|
||||
<TableBlankRow colSpan={10}>
|
||||
{!isLoading && <NoRuns title="No runs found" />}
|
||||
</TableBlankRow>
|
||||
) : runs.length === 0 ? (
|
||||
@@ -187,6 +189,10 @@ export function TaskRunsTable({
|
||||
<TableCell to={path}>
|
||||
{run.createdAt ? <DateTime date={run.createdAt} /> : "–"}
|
||||
</TableCell>
|
||||
<TableCell to={path}>
|
||||
{run.delayUntil ? <DateTime date={run.delayUntil} /> : "–"}
|
||||
</TableCell>
|
||||
<TableCell to={path}>{run.ttl ?? "–"}</TableCell>
|
||||
<RunActionsCell run={run} path={path} />
|
||||
</TableRow>
|
||||
);
|
||||
|
||||
@@ -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>
|
||||
);
|
||||
}
|
||||
@@ -13,3 +13,4 @@ export const VERCEL_RESPONSE_TIMEOUT_STATUS_CODES = [408, 504];
|
||||
export const MAX_BATCH_TRIGGER_ITEMS = 100;
|
||||
export const MAX_TASK_RUN_ATTEMPTS = 250;
|
||||
export const BULK_ACTION_RUN_LIMIT = 250;
|
||||
export const MAX_JOB_RUN_EXECUTION_COUNT = 250;
|
||||
|
||||
@@ -40,6 +40,8 @@ export const TaskRunStatus = {
|
||||
COMPLETED_WITH_ERRORS: "COMPLETED_WITH_ERRORS",
|
||||
SYSTEM_FAILURE: "SYSTEM_FAILURE",
|
||||
CRASHED: "CRASHED",
|
||||
DELAYED: "DELAYED",
|
||||
EXPIRED: "EXPIRED",
|
||||
} as const satisfies Record<TaskRunStatusType, TaskRunStatusType>;
|
||||
|
||||
export const JobRunStatus = {
|
||||
|
||||
@@ -186,3 +186,9 @@ export { apiRateLimiter } from "./services/apiRateLimit.server";
|
||||
export { socketIo } from "./v3/handleSocketIo.server";
|
||||
export { wss } from "./v3/handleWebsockets.server";
|
||||
export { registryProxy } from "./v3/registryProxy.server";
|
||||
import { eventLoopMonitor } from "./eventLoopMonitor.server";
|
||||
import { env } from "./env.server";
|
||||
|
||||
if (env.EVENT_LOOP_MONITOR_ENABLED === "1") {
|
||||
eventLoopMonitor.enable();
|
||||
}
|
||||
|
||||
@@ -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(),
|
||||
@@ -100,6 +102,10 @@ const EnvironmentSchema = z.object({
|
||||
API_RATE_LIMIT_REQUEST_LOGS_ENABLED: z.string().default("0"),
|
||||
API_RATE_LIMIT_REJECTION_LOGS_ENABLED: z.string().default("1"),
|
||||
|
||||
//Ingesting event rate limit
|
||||
INGEST_EVENT_RATE_LIMIT_WINDOW: z.string().default("60s"),
|
||||
INGEST_EVENT_RATE_LIMIT_MAX: z.coerce.number().int().optional(),
|
||||
|
||||
//v3
|
||||
V3_ENABLED: z.string().default("false"),
|
||||
PROVIDER_SECRET: z.string().default("provider-secret"),
|
||||
@@ -111,6 +117,7 @@ const EnvironmentSchema = z.object({
|
||||
CONTAINER_REGISTRY_USERNAME: z.string().optional(),
|
||||
CONTAINER_REGISTRY_PASSWORD: z.string().optional(),
|
||||
DEPLOY_REGISTRY_HOST: z.string().optional(),
|
||||
DEPLOY_REGISTRY_NAMESPACE: z.string().default("trigger"),
|
||||
OBJECT_STORE_BASE_URL: z.string().optional(),
|
||||
OBJECT_STORE_ACCESS_KEY_ID: z.string().optional(),
|
||||
OBJECT_STORE_SECRET_ACCESS_KEY: z.string().optional(),
|
||||
@@ -144,7 +151,7 @@ const EnvironmentSchema = z.object({
|
||||
PROD_OTEL_LOG_EXPORT_TIMEOUT_MILLIS: z.string().default("30000"),
|
||||
PROD_OTEL_LOG_MAX_QUEUE_SIZE: z.string().default("512"),
|
||||
|
||||
RUNTIME_WAIT_THRESHOLD_IN_MS: z.coerce.number().int().default(30000),
|
||||
CHECKPOINT_THRESHOLD_IN_MS: z.coerce.number().int().default(30000),
|
||||
|
||||
// Internal OTEL environment variables
|
||||
INTERNAL_OTEL_TRACE_EXPORTER_URL: z.string().optional(),
|
||||
@@ -164,6 +171,44 @@ const EnvironmentSchema = z.object({
|
||||
ALERT_RESEND_API_KEY: z.string().optional(),
|
||||
|
||||
MAX_SEQUENTIAL_INDEX_FAILURE_COUNT: z.coerce.number().default(96),
|
||||
|
||||
LOOPS_API_KEY: z.string().optional(),
|
||||
MARQS_DISABLE_REBALANCING: z.coerce.boolean().default(false),
|
||||
|
||||
VERBOSE_GRAPHILE_LOGGING: z.string().default("false"),
|
||||
V2_MARQS_ENABLED: z.string().default("0"),
|
||||
V2_MARQS_CONSUMER_POOL_ENABLED: z.string().default("0"),
|
||||
V2_MARQS_CONSUMER_POOL_SIZE: z.coerce.number().int().default(10),
|
||||
V2_MARQS_CONSUMER_POLL_INTERVAL_MS: z.coerce.number().int().default(1000),
|
||||
V2_MARQS_QUEUE_SELECTION_COUNT: z.coerce.number().int().default(36),
|
||||
V2_MARQS_VISIBILITY_TIMEOUT_MS: z.coerce
|
||||
.number()
|
||||
.int()
|
||||
.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(),
|
||||
EVENT_LOOP_MONITOR_ENABLED: z.string().default("1"),
|
||||
MAXIMUM_LIVE_RELOADING_EVENTS: z.coerce.number().int().default(1000),
|
||||
MAXIMUM_TRACE_SUMMARY_VIEW_COUNT: z.coerce.number().int().default(25_000),
|
||||
TASK_PAYLOAD_OFFLOAD_THRESHOLD: z.coerce.number().int().default(524_288), // 512KB
|
||||
TASK_PAYLOAD_MAXIMUM_SIZE: z.coerce.number().int().default(3_145_728), // 3MB
|
||||
});
|
||||
|
||||
export type Environment = z.infer<typeof EnvironmentSchema>;
|
||||
|
||||
@@ -0,0 +1,77 @@
|
||||
import { createHook } from "node:async_hooks";
|
||||
import { singleton } from "./utils/singleton";
|
||||
import { tracer } from "./v3/tracer.server";
|
||||
|
||||
const THRESHOLD_NS = 1e8; // 100ms
|
||||
|
||||
const cache = new Map<number, { type: string; start?: [number, number] }>();
|
||||
|
||||
function init(asyncId: number, type: string, triggerAsyncId: number, resource: any) {
|
||||
cache.set(asyncId, {
|
||||
type,
|
||||
});
|
||||
}
|
||||
|
||||
function destroy(asyncId: number) {
|
||||
cache.delete(asyncId);
|
||||
}
|
||||
|
||||
function before(asyncId: number) {
|
||||
const cached = cache.get(asyncId);
|
||||
|
||||
if (!cached) {
|
||||
return;
|
||||
}
|
||||
|
||||
cache.set(asyncId, {
|
||||
...cached,
|
||||
start: process.hrtime(),
|
||||
});
|
||||
}
|
||||
|
||||
function after(asyncId: number) {
|
||||
const cached = cache.get(asyncId);
|
||||
|
||||
if (!cached) {
|
||||
return;
|
||||
}
|
||||
|
||||
cache.delete(asyncId);
|
||||
|
||||
if (!cached.start) {
|
||||
return;
|
||||
}
|
||||
|
||||
const diff = process.hrtime(cached.start);
|
||||
const diffNs = diff[0] * 1e9 + diff[1];
|
||||
if (diffNs > THRESHOLD_NS) {
|
||||
const time = diffNs / 1e6; // in ms
|
||||
|
||||
const newSpan = tracer.startSpan("event-loop-blocked", {
|
||||
startTime: new Date(new Date().getTime() - time),
|
||||
attributes: {
|
||||
asyncType: cached.type,
|
||||
label: "EventLoopMonitor",
|
||||
},
|
||||
});
|
||||
|
||||
newSpan.end();
|
||||
}
|
||||
}
|
||||
|
||||
export const eventLoopMonitor = singleton("eventLoopMonitor", () => {
|
||||
const hook = createHook({ init, before, after, destroy });
|
||||
|
||||
return {
|
||||
enable: () => {
|
||||
console.log("🥸 Initializing event loop monitor");
|
||||
|
||||
hook.enable();
|
||||
},
|
||||
disable: () => {
|
||||
console.log("🥸 Disabling event loop monitor");
|
||||
|
||||
hook.disable();
|
||||
},
|
||||
};
|
||||
});
|
||||
@@ -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);
|
||||
}
|
||||
|
||||
@@ -26,7 +26,7 @@ export function useEventSource(
|
||||
const eventSource = new EventSource(url, init);
|
||||
eventSource.addEventListener(event ?? "message", handler);
|
||||
|
||||
// rest data if dependencies change
|
||||
// reset data if dependencies change
|
||||
setData(null);
|
||||
|
||||
function handler(event: MessageEvent) {
|
||||
|
||||
@@ -84,7 +84,9 @@ export function createPkApiKeyForEnv(envType: RuntimeEnvironment["type"]) {
|
||||
return `pk_${envSlug(envType)}_${apiKeyId(20)}`;
|
||||
}
|
||||
|
||||
export function envSlug(environmentType: RuntimeEnvironment["type"]) {
|
||||
export type EnvSlug = "dev" | "stg" | "prod" | "prev";
|
||||
|
||||
export function envSlug(environmentType: RuntimeEnvironment["type"]): EnvSlug {
|
||||
switch (environmentType) {
|
||||
case "DEVELOPMENT": {
|
||||
return "dev";
|
||||
@@ -100,3 +102,7 @@ export function envSlug(environmentType: RuntimeEnvironment["type"]) {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
export function isEnvSlug(maybeSlug: string): maybeSlug is EnvSlug {
|
||||
return ["dev", "stg", "prod", "prev"].includes(maybeSlug);
|
||||
}
|
||||
|
||||
@@ -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";
|
||||
}
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -65,65 +65,62 @@ export async function findEnvironmentById(id: string) {
|
||||
}
|
||||
|
||||
export async function createNewSession(environment: RuntimeEnvironment, ipAddress: string) {
|
||||
return prisma.$transaction(async (tx) => {
|
||||
const session = await tx.runtimeEnvironmentSession.create({
|
||||
data: {
|
||||
environmentId: environment.id,
|
||||
ipAddress,
|
||||
},
|
||||
});
|
||||
|
||||
await tx.runtimeEnvironment.update({
|
||||
where: {
|
||||
id: environment.id,
|
||||
},
|
||||
data: {
|
||||
currentSessionId: session.id,
|
||||
},
|
||||
});
|
||||
|
||||
return session;
|
||||
const session = await prisma.runtimeEnvironmentSession.create({
|
||||
data: {
|
||||
environmentId: environment.id,
|
||||
ipAddress,
|
||||
},
|
||||
});
|
||||
|
||||
await prisma.runtimeEnvironment.update({
|
||||
where: {
|
||||
id: environment.id,
|
||||
},
|
||||
data: {
|
||||
currentSessionId: session.id,
|
||||
},
|
||||
});
|
||||
|
||||
return session;
|
||||
}
|
||||
|
||||
export async function disconnectSession(environmentId: string) {
|
||||
return prisma.$transaction(async (tx) => {
|
||||
const environment = await tx.runtimeEnvironment.findUnique({
|
||||
where: {
|
||||
id: environmentId,
|
||||
},
|
||||
});
|
||||
|
||||
if (!environment || !environment.currentSessionId) {
|
||||
return null;
|
||||
}
|
||||
|
||||
const session = await tx.runtimeEnvironmentSession.update({
|
||||
where: {
|
||||
id: environment.currentSessionId,
|
||||
},
|
||||
data: {
|
||||
disconnectedAt: new Date(),
|
||||
},
|
||||
});
|
||||
|
||||
await tx.runtimeEnvironment.update({
|
||||
where: {
|
||||
id: environment.id,
|
||||
},
|
||||
data: {
|
||||
currentSessionId: null,
|
||||
},
|
||||
});
|
||||
|
||||
return session;
|
||||
const environment = await prisma.runtimeEnvironment.findUnique({
|
||||
where: {
|
||||
id: environmentId,
|
||||
},
|
||||
});
|
||||
|
||||
if (!environment || !environment.currentSessionId) {
|
||||
return null;
|
||||
}
|
||||
|
||||
const session = await prisma.runtimeEnvironmentSession.update({
|
||||
where: {
|
||||
id: environment.currentSessionId,
|
||||
},
|
||||
data: {
|
||||
disconnectedAt: new Date(),
|
||||
},
|
||||
});
|
||||
|
||||
await prisma.runtimeEnvironment.update({
|
||||
where: {
|
||||
id: environment.id,
|
||||
},
|
||||
data: {
|
||||
currentSessionId: null,
|
||||
},
|
||||
});
|
||||
|
||||
return session;
|
||||
}
|
||||
|
||||
type DisplayableInputEnvironment = Prisma.RuntimeEnvironmentGetPayload<{
|
||||
select: {
|
||||
id: true;
|
||||
type: true;
|
||||
slug: true;
|
||||
orgMember: {
|
||||
select: {
|
||||
user: {
|
||||
@@ -138,17 +135,24 @@ type DisplayableInputEnvironment = Prisma.RuntimeEnvironmentGetPayload<{
|
||||
};
|
||||
}>;
|
||||
|
||||
export function displayableEnvironments(
|
||||
export function displayableEnvironment(
|
||||
environment: DisplayableInputEnvironment,
|
||||
userId: string | undefined
|
||||
) {
|
||||
let userName: string | undefined = undefined;
|
||||
|
||||
if (environment.type === "DEVELOPMENT") {
|
||||
if (!environment.orgMember) {
|
||||
userName = "Deleted";
|
||||
} else if (environment.orgMember.user.id !== userId) {
|
||||
userName = getUsername(environment.orgMember.user);
|
||||
}
|
||||
}
|
||||
|
||||
return {
|
||||
id: environment.id,
|
||||
type: environment.type,
|
||||
userName: environment.orgMember
|
||||
? environment.orgMember.user.id === userId
|
||||
? undefined
|
||||
: getUsername(environment.orgMember.user)
|
||||
: undefined,
|
||||
slug: environment.slug,
|
||||
userName,
|
||||
};
|
||||
}
|
||||
|
||||
@@ -118,6 +118,7 @@ export function batchTaskRunItemStatusForRunStatus(
|
||||
case TaskRunStatus.COMPLETED_WITH_ERRORS:
|
||||
case TaskRunStatus.SYSTEM_FAILURE:
|
||||
case TaskRunStatus.CRASHED:
|
||||
case TaskRunStatus.EXPIRED:
|
||||
return BatchTaskRunItemStatus.FAILED;
|
||||
case TaskRunStatus.PENDING:
|
||||
case TaskRunStatus.WAITING_FOR_DEPLOY:
|
||||
@@ -125,6 +126,7 @@ export function batchTaskRunItemStatusForRunStatus(
|
||||
case TaskRunStatus.RETRYING_AFTER_FAILURE:
|
||||
case TaskRunStatus.EXECUTING:
|
||||
case TaskRunStatus.PAUSED:
|
||||
case TaskRunStatus.DELAYED:
|
||||
return BatchTaskRunItemStatus.PENDING;
|
||||
default:
|
||||
assertNever(status);
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -10,7 +10,12 @@ import type {
|
||||
TaskSpec,
|
||||
WorkerUtils,
|
||||
} from "graphile-worker";
|
||||
import { run as graphileRun, makeWorkerUtils, parseCronItems } from "graphile-worker";
|
||||
import {
|
||||
run as graphileRun,
|
||||
makeWorkerUtils,
|
||||
parseCronItems,
|
||||
Logger as GraphileLogger,
|
||||
} from "graphile-worker";
|
||||
import { SpanKind, trace } from "@opentelemetry/api";
|
||||
|
||||
import omit from "lodash.omit";
|
||||
@@ -19,6 +24,7 @@ import { $replica, PrismaClient, PrismaClientOrTransaction } from "~/db.server";
|
||||
import { PgListenService } from "~/services/db/pgListen.server";
|
||||
import { workerLogger as logger } from "~/services/logger.server";
|
||||
import { flattenAttributes } from "@trigger.dev/core/v3";
|
||||
import { env } from "~/env.server";
|
||||
|
||||
const tracer = trace.getTracer("zodWorker", "3.0.0.dp.1");
|
||||
|
||||
@@ -56,13 +62,16 @@ const AddJobResultsSchema = z.array(GraphileJobSchema);
|
||||
|
||||
export type ZodTasks<TConsumerSchema extends MessageCatalogSchema> = {
|
||||
[K in keyof TConsumerSchema]: {
|
||||
queueName?: string | ((payload: z.infer<TConsumerSchema[K]>) => string);
|
||||
jobKey?: string | ((payload: z.infer<TConsumerSchema[K]>) => string | undefined);
|
||||
priority?: number;
|
||||
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>;
|
||||
};
|
||||
};
|
||||
|
||||
@@ -75,11 +84,17 @@ 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>;
|
||||
};
|
||||
};
|
||||
|
||||
export type ZodWorkerEnqueueOptions = TaskSpec & {
|
||||
type ZodTaskSpec = Omit<TaskSpec, "queueName">;
|
||||
|
||||
export type ZodWorkerEnqueueOptions = ZodTaskSpec & {
|
||||
tx?: PrismaClientOrTransaction;
|
||||
};
|
||||
|
||||
@@ -162,12 +177,25 @@ export class ZodWorker<TMessageCatalog extends MessageCatalogSchema> {
|
||||
|
||||
this.#workerUtils = await makeWorkerUtils(this.#runnerOptions);
|
||||
|
||||
const graphileLogger = new GraphileLogger((scope) => {
|
||||
return (level, message, meta) => {
|
||||
if (env.VERBOSE_GRAPHILE_LOGGING !== "true") return;
|
||||
|
||||
logger.debug(`[graphile-worker][${this.#name}][${level}] ${message}`, {
|
||||
scope,
|
||||
meta,
|
||||
workerName: this.#name,
|
||||
});
|
||||
};
|
||||
});
|
||||
|
||||
this.#runner = await graphileRun({
|
||||
...this.#runnerOptions,
|
||||
noHandleSignals: true,
|
||||
taskList: this.#createTaskListFromTasks(),
|
||||
parsedCronItems,
|
||||
forbiddenFlags: this.#rateLimiter?.forbiddenFlags.bind(this.#rateLimiter),
|
||||
logger: graphileLogger,
|
||||
});
|
||||
|
||||
if (!this.#runner) {
|
||||
@@ -237,6 +265,20 @@ export class ZodWorker<TMessageCatalog extends MessageCatalogSchema> {
|
||||
this.#logDebug("stop");
|
||||
});
|
||||
|
||||
this.#runner?.events.on("worker:getJob:error", ({ worker, error }) => {
|
||||
this.#logDebug("worker:getJob:error", { workerId: worker.workerId, error });
|
||||
});
|
||||
|
||||
this.#runner?.events.on("worker:getJob:start", ({ worker }) => {
|
||||
if (env.VERBOSE_GRAPHILE_LOGGING !== "true") return;
|
||||
this.#logDebug("worker:getJob:start", { workerId: worker.workerId });
|
||||
});
|
||||
|
||||
this.#runner?.events.on("job:start", ({ worker, job }) => {
|
||||
if (env.VERBOSE_GRAPHILE_LOGGING !== "true") return;
|
||||
this.#logDebug("job:start", { workerId: worker.workerId, job });
|
||||
});
|
||||
|
||||
process.on("SIGTERM", this._handleSignal.bind(this));
|
||||
process.on("SIGINT", this._handleSignal.bind(this));
|
||||
|
||||
@@ -250,16 +292,18 @@ export class ZodWorker<TMessageCatalog extends MessageCatalogSchema> {
|
||||
|
||||
this.#shuttingDown = true;
|
||||
|
||||
this.#logDebug(
|
||||
`Received ${signal}, shutting down zodWorker with timeout ${this.#shutdownTimeoutInMs}ms`
|
||||
);
|
||||
|
||||
if (this.#shutdownTimeoutInMs) {
|
||||
setTimeout(() => {
|
||||
this.#logDebug("Shutdown timeout reached, exiting process");
|
||||
this.#logDebug(`Shutdown timeout of ${this.#shutdownTimeoutInMs} reached, exiting process`);
|
||||
|
||||
process.exit(0);
|
||||
}, this.#shutdownTimeoutInMs);
|
||||
}
|
||||
|
||||
this.#logDebug(`Received ${signal}, shutting down zodWorker...`);
|
||||
|
||||
this.stop().finally(() => {
|
||||
this.#logDebug("zodWorker stopped");
|
||||
});
|
||||
@@ -286,10 +330,6 @@ export class ZodWorker<TMessageCatalog extends MessageCatalogSchema> {
|
||||
...optionsWithoutTx,
|
||||
};
|
||||
|
||||
if (typeof task.queueName === "function") {
|
||||
spec.queueName = task.queueName(payload);
|
||||
}
|
||||
|
||||
if (typeof task.jobKey === "function") {
|
||||
const jobKey = task.jobKey(payload);
|
||||
|
||||
@@ -298,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,
|
||||
@@ -345,17 +379,15 @@ export class ZodWorker<TMessageCatalog extends MessageCatalogSchema> {
|
||||
`SELECT * FROM ${this.graphileWorkerSchema}.add_job(
|
||||
identifier => $1::text,
|
||||
payload => $2::json,
|
||||
queue_name => $3::text,
|
||||
run_at => $4::timestamptz,
|
||||
max_attempts => $5::int,
|
||||
job_key => $6::text,
|
||||
priority => $7::int,
|
||||
flags => $8::text[],
|
||||
job_key_mode => $9::text
|
||||
run_at => $3::timestamptz,
|
||||
max_attempts => $4::int,
|
||||
job_key => $5::text,
|
||||
priority => $6::int,
|
||||
flags => $7::text[],
|
||||
job_key_mode => $8::text
|
||||
)`,
|
||||
identifier,
|
||||
JSON.stringify(payload),
|
||||
spec.queueName || null,
|
||||
spec.runAt || null,
|
||||
spec.maxAttempts || null,
|
||||
spec.jobKey || null,
|
||||
@@ -447,33 +479,15 @@ export class ZodWorker<TMessageCatalog extends MessageCatalogSchema> {
|
||||
return taskList;
|
||||
}
|
||||
|
||||
async #getQueueName(queueId: number | null) {
|
||||
if (queueId === null) {
|
||||
return;
|
||||
}
|
||||
|
||||
const schema = z.array(z.object({ queue_name: z.string() }));
|
||||
|
||||
const rawQueueNameResults = await $replica.$queryRawUnsafe(
|
||||
`SELECT queue_name FROM ${this.graphileWorkerSchema}._private_job_queues WHERE id = $1`,
|
||||
queueId
|
||||
);
|
||||
|
||||
const queueNameResults = schema.parse(rawQueueNameResults);
|
||||
|
||||
return queueNameResults[0]?.queue_name;
|
||||
}
|
||||
|
||||
async #rescheduleTask(payload: unknown, helpers: JobHelpers) {
|
||||
this.#logDebug("Rescheduling task", { payload, job: helpers.job });
|
||||
|
||||
await this.enqueue(helpers.job.task_identifier, payload, {
|
||||
runAt: new Date(Date.now() + 1000 * 10),
|
||||
queueName: await this.#getQueueName(helpers.job.job_queue_id),
|
||||
runAt: helpers.job.run_at,
|
||||
priority: helpers.job.priority,
|
||||
jobKey: helpers.job.key ?? undefined,
|
||||
flags: Object.keys(helpers.job.flags ?? []),
|
||||
maxAttempts: helpers.job.max_attempts,
|
||||
maxAttempts: helpers.job.max_attempts - (helpers.job.attempts - 1),
|
||||
});
|
||||
}
|
||||
|
||||
@@ -569,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);
|
||||
@@ -650,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`;
|
||||
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
import { PrismaClient, prisma } from "~/db.server";
|
||||
import { Project } from "~/models/project.server";
|
||||
import { displayableEnvironments } from "~/models/runtimeEnvironment.server";
|
||||
import { displayableEnvironment } from "~/models/runtimeEnvironment.server";
|
||||
import { User } from "~/models/user.server";
|
||||
import { sortEnvironments } from "~/utils/environmentSort";
|
||||
|
||||
@@ -86,7 +86,7 @@ export class ProjectPresenter {
|
||||
httpEndpointCount: project._count.httpEndpoints,
|
||||
environments: sortEnvironments(
|
||||
project.environments.map((environment) => ({
|
||||
...displayableEnvironments(environment, userId),
|
||||
...displayableEnvironment(environment, userId),
|
||||
userId: environment.orgMember?.user.id,
|
||||
}))
|
||||
),
|
||||
|
||||
Some files were not shown because too many files have changed in this diff Show More
Reference in New Issue
Block a user