Compare commits
274 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| e85fc501a6 | |||
| 0bfeb0816f | |||
| b207601732 | |||
| 3913e57ef4 | |||
| 26f310397a | |||
| 0a845767a0 | |||
| ed2a26c865 | |||
| 4a68e71583 | |||
| b657eb6555 | |||
| 62c9a5b712 | |||
| f339b41ef3 | |||
| ae40ce3995 | |||
| 374edef020 | |||
| b82db67b81 | |||
| 26093896d2 | |||
| e7bd1ee676 | |||
| 584c7da5df | |||
| c9e1a3e9c5 | |||
| 2f5b4a8471 | |||
| 69f6891687 | |||
| acd7681e58 | |||
| 180a5ef01d | |||
| 62544d3234 | |||
| 9b045071d8 | |||
| 618d2207c6 | |||
| 36d8bee14a | |||
| 44e1b87547 | |||
| 71b0ef8f77 | |||
| f66543c9c7 | |||
| 03b104a3d5 | |||
| fde939a30e | |||
| 4986bfda2e | |||
| a3abe4ca08 | |||
| fad32a79dd | |||
| 9d843caf83 | |||
| 5eaf7bca34 | |||
| 4a0368dec3 | |||
| f10f120e55 | |||
| d3a18fbdf6 | |||
| b82a07ad1c | |||
| fa81bf356a | |||
| fec4dc3bef | |||
| a80089c88f | |||
| 63a643b7c9 | |||
| efd970a901 | |||
| 89a5e9f7da | |||
| 404931e224 | |||
| 57bf98307a | |||
| 7551adfeb6 | |||
| 69117e8ba1 | |||
| b580d53b59 | |||
| d81e21d2ec | |||
| 826a64fe6f | |||
| d1849b0ea9 | |||
| 3aa634d2b8 | |||
| 7de3fb01c7 | |||
| 4ae9a3feac | |||
| 5b18298fef | |||
| 7a9bd18ba2 | |||
| 03af44e545 | |||
| 5776257663 | |||
| d3997c9fc6 | |||
| 0676ea9668 | |||
| 7a4721e122 | |||
| a5aed0d139 | |||
| f1480595d5 | |||
| 328947dbfd | |||
| 279717b092 | |||
| 702d198445 | |||
| 803f3c15ab | |||
| 1c24348f7d | |||
| f854cb90eb | |||
| 7268f17b00 | |||
| 624ddce32f | |||
| 9be1557bb7 | |||
| 6ce6f8e3ad | |||
| d462b7a51b | |||
| f2894c177a | |||
| e35f29764a | |||
| 1207efbbad | |||
| 7ea8532cce | |||
| 6642228f26 | |||
| d39145d810 | |||
| 8886bb76e0 | |||
| 4b72726078 | |||
| 6dcfeadaca | |||
| ae839ebe11 | |||
| 5d0d71c2ae | |||
| 1239a3ceb9 | |||
| dd31b1e668 | |||
| a707446989 | |||
| eb050f6730 | |||
| 56d9bf7c67 | |||
| 73e469daf5 | |||
| 0382cf8719 | |||
| 0fe835492d | |||
| 29b69160e1 | |||
| ae9efe3d8b | |||
| 3feb5ffb5f | |||
| 7fb482de64 | |||
| 6f11584aaa | |||
| 43e240cd50 | |||
| afe7f410c7 | |||
| 28837f39b3 | |||
| 5c64cefaf0 | |||
| dc5ed0a0cd | |||
| e7e7397ed9 | |||
| 49184c7189 | |||
| eb60126284 | |||
| d876c358d2 | |||
| 6e9c8ab555 | |||
| 9a3eb289cd | |||
| 51315fc3c8 | |||
| 4047f00562 | |||
| 545f85b44e | |||
| abe202e7b7 | |||
| ae27fd83af | |||
| c702d6a9ca | |||
| 9af2570da6 | |||
| a946797d95 | |||
| 8c4df326cc | |||
| 11b997d2bf | |||
| 8694e573f5 | |||
| b271742dca | |||
| 6f9f25481e | |||
| 67aaffb6fc | |||
| e3cf456c69 | |||
| bf7827e7b8 | |||
| bc020a3ffe | |||
| b361afbfe4 | |||
| a3d809740d | |||
| f93eae300e | |||
| a2365e406d | |||
| 42d319c2d1 | |||
| b66d5525ef | |||
| 719c0a0b94 | |||
| f1c768a255 | |||
| d9c9e80bc4 | |||
| d39932ebf7 | |||
| 9bcb8cb42a | |||
| 2374f8e8ac | |||
| a22b5869e4 | |||
| f1571cbfab | |||
| 222dfed4d2 | |||
| 89c0e68fb5 | |||
| 5b745dc1a4 | |||
| 2dea8dec35 | |||
| 252f5b3967 | |||
| 014d233339 | |||
| e398bdf859 | |||
| 426316aa3c | |||
| 4f9b8d6721 | |||
| a345ad2d57 | |||
| e2b67400ea | |||
| c4dc2849af | |||
| a8263dec86 | |||
| 53054d0146 | |||
| 74236ae7c1 | |||
| 055d468b38 | |||
| fc473105e3 | |||
| 3a35d60091 | |||
| f11c0f991b | |||
| 1cf181fd02 | |||
| 412d0149d5 | |||
| d614c8d2ca | |||
| 4329d1130a | |||
| 17f9fc8273 | |||
| db611bebba | |||
| 215b60c742 | |||
| 9c89ab40aa | |||
| b342e06583 | |||
| 4418de515b | |||
| 0d2a71c5cf | |||
| 9e02b45808 | |||
| 395abe1b92 | |||
| b35eebb666 | |||
| 9ecf07731a | |||
| b890a64046 | |||
| 4266b308e7 | |||
| 67c62f7d25 | |||
| c6d6a9191c | |||
| cd8f6b9af0 | |||
| a49b659701 | |||
| 61c9a7c7d6 | |||
| de302858f8 | |||
| cf9d466b9e | |||
| 7b3ffe040b | |||
| 737ff9928c | |||
| fc4b95d69e | |||
| 6c546936f3 | |||
| 08e4fc002f | |||
| 0e347b001b | |||
| 2e56af9742 | |||
| dc5eb68a0b | |||
| 891278ab1f | |||
| 41e1ac9e0b | |||
| c5aaabdb4e | |||
| affa3f25ee | |||
| 86c24afb44 | |||
| 6afda40c35 | |||
| b585d1f2c5 | |||
| f26ba9a5cd | |||
| ba0ffacb15 | |||
| e55ca0279b | |||
| 9304c447be | |||
| ebec14093d | |||
| 67cac97b26 | |||
| 0f32dee0eb | |||
| 479d355b7b | |||
| 49e9dd8cb2 | |||
| cae3298564 | |||
| 1f190fe680 | |||
| 85ec1af198 | |||
| 52c364d691 | |||
| a5d83971aa | |||
| 5b9c67aeef | |||
| 1c98eb2b61 | |||
| f424c81086 | |||
| 0a33bf7206 | |||
| 819b663ad7 | |||
| 478ce006cb | |||
| fe9435f63f | |||
| d7cc2c9d6c | |||
| 851b1769cb | |||
| e8381d5341 | |||
| 20aa1cbe9b | |||
| 883feab257 | |||
| 83304f86f8 | |||
| 5d78bba5e3 | |||
| e87a0cba1d | |||
| 8852d5b388 | |||
| 79ea8edc7c | |||
| 553370a9e5 | |||
| 47c18d9146 | |||
| c385db63be | |||
| 86b2674e68 | |||
| 161b2b3323 | |||
| 9b6f8f9238 | |||
| f50cee66d9 | |||
| 44d443a048 | |||
| 2a1273b683 | |||
| 64401ee8ec | |||
| b7845685eb | |||
| 52c9d485f8 | |||
| ca47e6bdcd | |||
| 933cf0f119 | |||
| b6517af522 | |||
| 479445f741 | |||
| 1fabcf68fa | |||
| d7c79ee651 | |||
| 8d6e710dc2 | |||
| 90a302f9f7 | |||
| 07efef405d | |||
| dd63fe6e9f | |||
| d8c645a3f0 | |||
| 24600d2c01 | |||
| 584936fa29 | |||
| a54bdb23bf | |||
| 07f7c65054 | |||
| d04c8e89f6 | |||
| b69902dd84 | |||
| 2aabb994cf | |||
| 310d51bac0 | |||
| 6be1099a12 | |||
| baf3a84cda | |||
| 70f9bd0d70 | |||
| a1860dbaae | |||
| d69e4e712d | |||
| 7fae67c47d | |||
| 9ae0ca64af | |||
| da69fa0613 | |||
| 336029b842 | |||
| 29edcd3df9 | |||
| a9ff32418e |
@@ -0,0 +1,9 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
"@trigger.dev/core": patch
|
||||
---
|
||||
|
||||
- Fix additionalFiles that aren't decendants
|
||||
- Stop swallowing uncaught exceptions in prod
|
||||
- Improve warnings and errors, fail early on critical warnings
|
||||
- New arg to --save-logs even for successful builds
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
---
|
||||
|
||||
v3 CLI update command and package manager detection fix
|
||||
@@ -0,0 +1,6 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
"@trigger.dev/core": patch
|
||||
---
|
||||
|
||||
better handle task metadata parse errors, and display nicely formatted errors
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
---
|
||||
|
||||
Vastly improved dev command output
|
||||
@@ -0,0 +1,8 @@
|
||||
---
|
||||
"@trigger.dev/sdk": patch
|
||||
"@trigger.dev/core-apps": patch
|
||||
"trigger.dev": patch
|
||||
"@trigger.dev/core": patch
|
||||
---
|
||||
|
||||
add machine config and secure zod connection
|
||||
@@ -4,7 +4,8 @@
|
||||
"commit": false,
|
||||
"fixed": [
|
||||
[
|
||||
"@trigger.dev/*"
|
||||
"@trigger.dev/*",
|
||||
"trigger.dev"
|
||||
]
|
||||
],
|
||||
"linked": [],
|
||||
@@ -16,7 +17,10 @@
|
||||
"emails",
|
||||
"proxy",
|
||||
"yalt",
|
||||
"@trigger.dev/database"
|
||||
"@trigger.dev/database",
|
||||
"coordinator",
|
||||
"docker-provider",
|
||||
"kubernetes-provider"
|
||||
],
|
||||
"___experimentalUnsafeOptions_WILL_CHANGE_IN_PATCH": {
|
||||
"onlyUpdatePeerDependentsWhenOutOfRange": true
|
||||
|
||||
@@ -0,0 +1,7 @@
|
||||
---
|
||||
"@trigger.dev/sdk": patch
|
||||
"trigger.dev": patch
|
||||
"@trigger.dev/core": patch
|
||||
---
|
||||
|
||||
Configurable log levels in the config file and via env var
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
---
|
||||
|
||||
Don’t swallow some error messages when deploying
|
||||
@@ -0,0 +1,6 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
"@trigger.dev/core": patch
|
||||
---
|
||||
|
||||
Add option to print console logs in the dev CLI locally (issue #1014)
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"@trigger.dev/core": patch
|
||||
---
|
||||
|
||||
Fixed batch otel flushing
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
---
|
||||
|
||||
Fix permissions inside node_modules
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"@trigger.dev/react": patch
|
||||
---
|
||||
|
||||
Fix for shared queryKey between useRunDetails and useRunStatuses
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"@trigger.dev/core": patch
|
||||
---
|
||||
|
||||
Retry 429, 500, and connection error API requests to the trigger.dev server
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"@trigger.dev/sdk": patch
|
||||
---
|
||||
|
||||
Export queue from the SDK
|
||||
@@ -0,0 +1,7 @@
|
||||
---
|
||||
"@trigger.dev/core-apps": patch
|
||||
"trigger.dev": patch
|
||||
"@trigger.dev/core": patch
|
||||
---
|
||||
|
||||
Display errors for runs and deployments
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
---
|
||||
|
||||
v3: fix digest extraction
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"@trigger.dev/nestjs": patch
|
||||
---
|
||||
|
||||
fix: [nestjs integration] fastify HTTP adapter detection now works correctly for response headers
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
---
|
||||
|
||||
Changed "Worker" to "Version" in the dev command key
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
---
|
||||
|
||||
Correctly handle self-hosted deploy command errors
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
---
|
||||
|
||||
Default to retrying enabled in dev when running init
|
||||
@@ -0,0 +1,6 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
"@trigger.dev/core": patch
|
||||
---
|
||||
|
||||
Handle string and non-stringifiable outputs like functions
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
---
|
||||
|
||||
Improve error messages during dev/deploy and handle deploy image build issues
|
||||
@@ -0,0 +1,6 @@
|
||||
---
|
||||
"@trigger.dev/sdk": patch
|
||||
"@trigger.dev/core": patch
|
||||
---
|
||||
|
||||
Added replayRun function to the SDK
|
||||
@@ -0,0 +1,6 @@
|
||||
---
|
||||
"@trigger.dev/sdk": patch
|
||||
"@trigger.dev/core": patch
|
||||
---
|
||||
|
||||
Improve the SDK function types and expose a new APIError instead of the APIResult type
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
---
|
||||
|
||||
Added a Node.js runtime check for the CLI
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
---
|
||||
|
||||
Fix post start hooks
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
---
|
||||
|
||||
Use the dashboard url instead of the API url for the View logs link
|
||||
@@ -0,0 +1,98 @@
|
||||
{
|
||||
"mode": "pre",
|
||||
"tag": "beta",
|
||||
"initialVersions": {
|
||||
"coordinator": "0.0.1",
|
||||
"docker-provider": "0.0.1",
|
||||
"kubernetes-provider": "0.0.1",
|
||||
"proxy": "0.0.11",
|
||||
"webapp": "1.0.0",
|
||||
"yalt": "0.0.1",
|
||||
"@trigger.dev/airtable": "2.3.18",
|
||||
"@trigger.dev/github": "2.3.18",
|
||||
"@trigger.dev/linear": "2.3.18",
|
||||
"@trigger.dev/openai": "2.3.18",
|
||||
"@trigger.dev/plain": "2.3.18",
|
||||
"@trigger.dev/replicate": "2.3.18",
|
||||
"@trigger.dev/resend": "2.3.18",
|
||||
"@trigger.dev/sendgrid": "2.3.18",
|
||||
"@trigger.dev/shopify": "2.3.18",
|
||||
"@trigger.dev/slack": "2.3.18",
|
||||
"@trigger.dev/stripe": "2.3.18",
|
||||
"@trigger.dev/supabase": "2.3.18",
|
||||
"@trigger.dev/typeform": "2.3.18",
|
||||
"@trigger.dev/astro": "2.3.18",
|
||||
"@trigger.dev/cli": "2.3.18",
|
||||
"trigger.dev": "2.3.18",
|
||||
"@trigger.dev/core": "2.3.18",
|
||||
"@trigger.dev/core-apps": "0.0.0",
|
||||
"@trigger.dev/core-backend": "2.3.18",
|
||||
"@trigger.dev/database": "0.0.1",
|
||||
"emails": "1.0.0",
|
||||
"@trigger.dev/eslint-plugin": "2.3.18",
|
||||
"@trigger.dev/express": "2.3.18",
|
||||
"@trigger.dev/hono": "2.3.18",
|
||||
"@trigger.dev/integration-kit": "2.3.18",
|
||||
"@trigger.dev/nestjs": "2.3.18",
|
||||
"@trigger.dev/nextjs": "2.3.18",
|
||||
"@trigger.dev/otlp-importer": "2.3.11",
|
||||
"@trigger.dev/react": "2.3.18",
|
||||
"@trigger.dev/remix": "2.3.18",
|
||||
"@trigger.dev/sveltekit": "2.3.18",
|
||||
"@trigger.dev/testing": "2.3.18",
|
||||
"@trigger.dev/sdk": "2.3.18",
|
||||
"@trigger.dev/yalt": "2.3.18"
|
||||
},
|
||||
"changesets": [
|
||||
"angry-eagles-trade",
|
||||
"beige-pens-dance",
|
||||
"breezy-gorillas-mate",
|
||||
"chilled-hornets-move",
|
||||
"clean-pianos-listen",
|
||||
"cool-glasses-bake",
|
||||
"cuddly-feet-approve",
|
||||
"dry-walls-check",
|
||||
"eight-pumas-float",
|
||||
"few-students-share",
|
||||
"green-bags-wink",
|
||||
"khaki-apricots-design",
|
||||
"khaki-poems-lay",
|
||||
"late-icons-lie",
|
||||
"late-steaks-behave",
|
||||
"lemon-jobs-repair",
|
||||
"light-bulldogs-press",
|
||||
"light-dragons-complain",
|
||||
"loud-actors-remember",
|
||||
"many-ligers-pump",
|
||||
"mighty-camels-joke",
|
||||
"new-rivers-tell",
|
||||
"ninety-pets-travel",
|
||||
"odd-poets-own",
|
||||
"polite-ducks-switch",
|
||||
"poor-flowers-cross",
|
||||
"rare-roses-float",
|
||||
"real-planets-stare",
|
||||
"rotten-dryers-exercise",
|
||||
"shaggy-spoons-taste",
|
||||
"sharp-emus-compare",
|
||||
"sharp-zebras-serve",
|
||||
"shiny-coats-cry",
|
||||
"silly-suits-switch",
|
||||
"smart-needles-move",
|
||||
"smart-olives-eat",
|
||||
"spicy-lamps-smoke",
|
||||
"strange-ghosts-matter",
|
||||
"strong-lemons-add",
|
||||
"stupid-bulldogs-applaud",
|
||||
"sweet-lizards-press",
|
||||
"swift-dragons-peel",
|
||||
"tall-bees-wave",
|
||||
"tame-guests-know",
|
||||
"tender-oranges-rhyme",
|
||||
"tidy-balloons-suffer",
|
||||
"tidy-dryers-sleep",
|
||||
"tiny-doors-type",
|
||||
"tiny-elephants-scream",
|
||||
"tricky-bulldogs-heal"
|
||||
]
|
||||
}
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
---
|
||||
|
||||
Add openssl to prod worker image and allow passing auth token via env var for deploy
|
||||
@@ -0,0 +1,7 @@
|
||||
---
|
||||
"@trigger.dev/sdk": patch
|
||||
"trigger.dev": patch
|
||||
"@trigger.dev/core": patch
|
||||
---
|
||||
|
||||
Fixed incorrect span timings around checkpoints by implementing a precise wall clock that resets after restores
|
||||
@@ -0,0 +1,7 @@
|
||||
---
|
||||
"@trigger.dev/sdk": patch
|
||||
"trigger.dev": patch
|
||||
"@trigger.dev/core": patch
|
||||
---
|
||||
|
||||
Adding task with a triggerSource of schedule
|
||||
@@ -0,0 +1,56 @@
|
||||
---
|
||||
"@trigger.dev/sdk": patch
|
||||
"@trigger.dev/core": patch
|
||||
---
|
||||
|
||||
Updates the `trigger`, `batchTrigger` and their `*AndWait` variants to use the first parameter for the payload/items, and the second parameter for options.
|
||||
|
||||
Before:
|
||||
|
||||
```ts
|
||||
await yourTask.trigger({ payload: { foo: "bar" }, options: { idempotencyKey: "key_1234" } });
|
||||
await yourTask.triggerAndWait({ payload: { foo: "bar" }, options: { idempotencyKey: "key_1234" } });
|
||||
|
||||
await yourTask.batchTrigger({ items: [{ payload: { foo: "bar" } }, { payload: { foo: "baz" } }] });
|
||||
await yourTask.batchTriggerAndWait({ items: [{ payload: { foo: "bar" } }, { payload: { foo: "baz" } }] });
|
||||
```
|
||||
|
||||
After:
|
||||
|
||||
```ts
|
||||
await yourTask.trigger({ foo: "bar" }, { idempotencyKey: "key_1234" });
|
||||
await yourTask.triggerAndWait({ foo: "bar" }, { idempotencyKey: "key_1234" });
|
||||
|
||||
await yourTask.batchTrigger([{ payload: { foo: "bar" } }, { payload: { foo: "baz" } }]);
|
||||
await yourTask.batchTriggerAndWait([{ payload: { foo: "bar" } }, { payload: { foo: "baz" } }]);
|
||||
```
|
||||
|
||||
We've also changed the API of the `triggerAndWait` result. Before, if the subtask that was triggered finished with an error, we would automatically "rethrow" the error in the parent task.
|
||||
|
||||
Now instead we're returning a `TaskRunResult` object that allows you to discriminate between successful and failed runs in the subtask:
|
||||
|
||||
Before:
|
||||
|
||||
```ts
|
||||
try {
|
||||
const result = await yourTask.triggerAndWait({ foo: "bar" });
|
||||
|
||||
// result is the output of your task
|
||||
console.log("result", result);
|
||||
|
||||
} catch (error) {
|
||||
// handle subtask errors here
|
||||
}
|
||||
```
|
||||
|
||||
After:
|
||||
|
||||
```ts
|
||||
const result = await yourTask.triggerAndWait({ foo: "bar" });
|
||||
|
||||
if (result.ok) {
|
||||
console.log(`Run ${result.id} succeeded with output`, result.output);
|
||||
} else {
|
||||
console.log(`Run ${result.id} failed with error`, result.error);
|
||||
}
|
||||
```
|
||||
@@ -0,0 +1,13 @@
|
||||
---
|
||||
"@trigger.dev/sdk": patch
|
||||
"trigger.dev": patch
|
||||
"@trigger.dev/core": patch
|
||||
---
|
||||
|
||||
When using idempotency keys, triggerAndWait and batchTriggerAndWait will still work even if the existing runs have already been completed (or even partially completed, in the case of batchTriggerAndWait)
|
||||
|
||||
- TaskRunExecutionResult.id is now the run friendlyId, not the attempt friendlyId
|
||||
- A single TaskRun can now have many batchItems, in the case of batchTriggerAndWait while using idempotency keys
|
||||
- A run’s idempotencyKey is now added to the ctx as well as the TaskEvent and displayed in the span view
|
||||
- When resolving batchTriggerAndWait, the runtimes no longer reject promises, leading to an error in the parent task
|
||||
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
---
|
||||
|
||||
Update trigger.dev CLI for new batch otel support
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
---
|
||||
|
||||
Add additional logging around cleaning up dev workers, and always kill them after 5 seconds if they haven't already exited
|
||||
@@ -0,0 +1,6 @@
|
||||
---
|
||||
"@trigger.dev/otlp-importer": patch
|
||||
"trigger.dev": patch
|
||||
---
|
||||
|
||||
Fix package builds and CLI commands on Windows
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"@trigger.dev/sdk": patch
|
||||
---
|
||||
|
||||
Remove unimplemented batchOptions
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
---
|
||||
|
||||
Fix CLI logout and add list-profiles command
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
---
|
||||
|
||||
Fixing an issue with bundling @trigger.dev/core/v3 in dev when using pnpm
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
---
|
||||
|
||||
Added DEBUG to the ignored env vars
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"@trigger.dev/core": patch
|
||||
---
|
||||
|
||||
Make optional schedule object fields nullish
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
---
|
||||
|
||||
Fixed an issue where the trigger.dev package was not being built before publishing to npm
|
||||
@@ -0,0 +1,8 @@
|
||||
---
|
||||
"trigger.dev": major
|
||||
"@trigger.dev/core": major
|
||||
"@trigger.dev/otlp-importer": major
|
||||
"@trigger.dev/sdk": major
|
||||
---
|
||||
|
||||
Updates to support Trigger.dev v3
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"@trigger.dev/core": patch
|
||||
---
|
||||
|
||||
Added JSDocs to the schedule SDK types
|
||||
@@ -0,0 +1,7 @@
|
||||
---
|
||||
"@trigger.dev/sdk": patch
|
||||
"trigger.dev": patch
|
||||
"@trigger.dev/core": patch
|
||||
---
|
||||
|
||||
Dynamically import superjson and fix some bundling issues
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
---
|
||||
|
||||
Changed the binary name from trigger.dev to triggerdev to fix a Windows issue
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
---
|
||||
|
||||
Stop swallowing deployment errors and display them better
|
||||
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
---
|
||||
|
||||
Init command was failing on Windows because of bad template paths
|
||||
@@ -0,0 +1,10 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
"@trigger.dev/core": patch
|
||||
---
|
||||
|
||||
Fixes an issue that caused failed tasks when resuming after calling `triggerAndWait` or `batchTriggerAndWait` in prod/staging (this doesn't effect dev).
|
||||
|
||||
The version of Node.js we use for deployed workers (latest 20) would crash with an out-of-memory error when the checkpoint was restored. This crash does not happen on Node 18x or Node21x, so we've decided to upgrade the worker version to Node.js21x, to mitigate this issue.
|
||||
|
||||
You'll need to re-deploy to production to fix the issue.
|
||||
@@ -0,0 +1,6 @@
|
||||
---
|
||||
"@trigger.dev/sdk": patch
|
||||
"@trigger.dev/core": patch
|
||||
---
|
||||
|
||||
Added cancelRun to the SDK
|
||||
@@ -0,0 +1,7 @@
|
||||
---
|
||||
"trigger.dev": patch
|
||||
"@trigger.dev/core": patch
|
||||
---
|
||||
|
||||
- Add graceful exit for prod workers
|
||||
- Prevent overflow in long waits
|
||||
@@ -0,0 +1,7 @@
|
||||
---
|
||||
"@trigger.dev/sdk": patch
|
||||
"trigger.dev": patch
|
||||
"@trigger.dev/core": patch
|
||||
---
|
||||
|
||||
Added a new global - Task Catalog - to better handle task metadata
|
||||
+31
-1
@@ -31,6 +31,9 @@ NODE_ENV=development
|
||||
# FROM_EMAIL=
|
||||
# REPLY_TO_EMAIL=
|
||||
|
||||
# Remove the following line to enable logging telemetry traces to the console
|
||||
LOG_TELEMETRY="false"
|
||||
|
||||
# CLOUD VARIABLES
|
||||
POSTHOG_PROJECT_KEY=
|
||||
PLAIN_API_KEY=
|
||||
@@ -42,4 +45,31 @@ CLOUD_LINEAR_CLIENT_ID=
|
||||
CLOUD_LINEAR_CLIENT_SECRET=
|
||||
CLOUD_SLACK_APP_HOST=
|
||||
CLOUD_SLACK_CLIENT_ID=
|
||||
CLOUD_SLACK_CLIENT_SECRET=
|
||||
CLOUD_SLACK_CLIENT_SECRET=
|
||||
|
||||
# v3 variables
|
||||
PROVIDER_SECRET=provider-secret # generate the actual secret with `openssl rand -hex 32`
|
||||
COORDINATOR_SECRET=coordinator-secret # generate the actual secret with `openssl rand -hex 32`
|
||||
|
||||
# Uncomment the following line to enable the registry proxy
|
||||
# ENABLE_REGISTRY_PROXY=true
|
||||
# DEPOT_TOKEN=<Depot org token>
|
||||
# DEPOT_PROJECT_ID=<Depot project id>
|
||||
# DEPLOY_REGISTRY_HOST=${APP_ORIGIN} # This is the host that the deploy CLI will use to push images to the registry
|
||||
# CONTAINER_REGISTRY_ORIGIN=<Container registry origin e.g. https://registry.digitalocean.com>
|
||||
# CONTAINER_REGISTRY_USERNAME=<Container registry username e.g. Digital ocean email address>
|
||||
# CONTAINER_REGISTRY_PASSWORD=<Container registry password e.g. Digital ocean PAT>
|
||||
# DEV_OTEL_EXPORTER_OTLP_ENDPOINT="http://0.0.0.0:4318"
|
||||
# These are needed for the object store (for handling large payloads/outputs)
|
||||
# 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
|
||||
|
||||
# These control the server-side internal telemetry
|
||||
# INTERNAL_OTEL_TRACE_EXPORTER_URL=<URL to send traces to>
|
||||
# INTERNAL_OTEL_TRACE_EXPORTER_AUTH_HEADER_NAME=<Header name for the auth token>
|
||||
# INTERNAL_OTEL_TRACE_EXPORTER_AUTH_HEADER_VALUE=<Auth token value>
|
||||
# INTERNAL_OTEL_TRACE_LOGGING_ENABLED=1
|
||||
# INTERNAL_OTEL_TRACE_SAMPING_RATE=20 # this means 1/20 traces or 5% of traces will be sampled (sampled = recorded)
|
||||
# INTERNAL_OTEL_TRACE_INSTRUMENT_PRISMA_ENABLED=0,
|
||||
@@ -0,0 +1,21 @@
|
||||
name: OpenTelemetry Auto-Instrumentation Request
|
||||
description: Suggest an SDK that you'd like to be auto-instrumented in the Run log view
|
||||
title: "auto-instrumentation: "
|
||||
labels: ["🌟 enhancement"]
|
||||
body:
|
||||
- type: textarea
|
||||
attributes:
|
||||
label: What API or SDK would you to have automatic spans for?
|
||||
description: A clear description of which API or SDK you'd like, and links to it.
|
||||
validations:
|
||||
required: true
|
||||
- type: textarea
|
||||
attributes:
|
||||
label: Is there an existing OpenTelemetry auto-instrumentation package?
|
||||
description: You can search for existing ones – https://opentelemetry.io/ecosystem/registry/?component=instrumentation&language=js
|
||||
validations:
|
||||
required: true
|
||||
- type: textarea
|
||||
attributes:
|
||||
label: Additional information
|
||||
description: Add any other information related to the feature here. If your feature request is related to any issues or discussions, link them here.
|
||||
@@ -4,31 +4,36 @@ on:
|
||||
jobs:
|
||||
e2e:
|
||||
name: "🧪 E2E Tests"
|
||||
runs-on: buildjet-4vcpu-ubuntu-2204
|
||||
runs-on: buildjet-16vcpu-ubuntu-2204
|
||||
steps:
|
||||
- name: 🐳 Login to Docker Hub
|
||||
if: github.event_name == 'push'
|
||||
uses: docker/login-action@v2
|
||||
with:
|
||||
username: ${{ secrets.DOCKERHUB_USERNAME }}
|
||||
password: ${{ secrets.DOCKERHUB_TOKEN }}
|
||||
username: ${{ secrets.DOCKERHUB_USERNAME || vars.DOCKERHUB_USERNAME }}
|
||||
password: ${{ secrets.DOCKERHUB_TOKEN || vars.DOCKERHUB_TOKEN }}
|
||||
|
||||
- name: ⬇️ Checkout repo
|
||||
uses: actions/checkout@v3
|
||||
with:
|
||||
fetch-depth: 0
|
||||
submodules: recursive
|
||||
|
||||
- name: ⎔ Setup pnpm
|
||||
uses: pnpm/action-setup@v2.2.4
|
||||
with:
|
||||
version: 7.18
|
||||
version: 8.15.5
|
||||
|
||||
- name: ⎔ Setup node
|
||||
uses: buildjet/setup-node@v3
|
||||
with:
|
||||
node-version: 18
|
||||
node-version: 20.11.1
|
||||
cache: "pnpm"
|
||||
|
||||
- name: Install Protoc
|
||||
uses: arduino/setup-protoc@v3
|
||||
with:
|
||||
repo-token: ${{ secrets.GITHUB_TOKEN }}
|
||||
|
||||
- name: 📥 Download deps
|
||||
run: pnpm install --frozen-lockfile
|
||||
|
||||
@@ -43,7 +48,10 @@ jobs:
|
||||
|
||||
# Build packages
|
||||
pnpm run build --filter @references/nextjs-test^...
|
||||
cd apps/webapp && pnpm run build:server
|
||||
cd ../..
|
||||
pnpm --filter @trigger.dev/database generate
|
||||
pnpm --filter @trigger.dev/otlp-importer generate
|
||||
|
||||
# Move trigger-cli bin to correct place
|
||||
pnpm install --frozen-lockfile
|
||||
|
||||
@@ -3,23 +3,18 @@ on:
|
||||
workflow_call:
|
||||
jobs:
|
||||
publish:
|
||||
strategy:
|
||||
fail-fast: true # when a job fails, all remaining ones will be cancelled
|
||||
matrix:
|
||||
runs-on: [buildjet-4vcpu-ubuntu-2204, buildjet-4vcpu-ubuntu-2204-arm]
|
||||
name: ${{matrix.runs-on}}
|
||||
runs-on: ${{matrix.runs-on}}
|
||||
runs-on: ubuntu-latest
|
||||
outputs:
|
||||
version: ${{ steps.get_version.outputs.version }}
|
||||
short_sha: ${{ steps.get_commit.outputs.sha_short }}
|
||||
steps:
|
||||
- name: 🐳 Login to Docker Hub
|
||||
uses: docker/login-action@v2
|
||||
with:
|
||||
username: ${{ secrets.DOCKERHUB_USERNAME }}
|
||||
password: ${{ secrets.DOCKERHUB_TOKEN }}
|
||||
- name: Setup Depot CLI
|
||||
uses: depot/setup-action@v1
|
||||
|
||||
- name: ⬇️ Checkout repo
|
||||
uses: actions/checkout@v3
|
||||
with:
|
||||
submodules: recursive
|
||||
|
||||
- name: 🆚 Get the version
|
||||
id: get_version
|
||||
@@ -41,19 +36,12 @@ jobs:
|
||||
echo "Invalid reference: ${GITHUB_REF}"
|
||||
exit 1
|
||||
fi
|
||||
if [[ ${{matrix.runs-on}} == *-arm ]]; then
|
||||
IMAGE_TAG="${IMAGE_TAG}-arm"
|
||||
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: 🐳 Build Docker Image
|
||||
run: |
|
||||
docker build -t release_build_image -f ./docker/Dockerfile .
|
||||
|
||||
- name: 🐙 Login to GitHub Container Registry
|
||||
uses: docker/login-action@v2
|
||||
with:
|
||||
@@ -61,24 +49,11 @@ jobs:
|
||||
username: ${{ github.repository_owner }}
|
||||
password: ${{ secrets.GITHUB_TOKEN }}
|
||||
|
||||
- name: 🐙 Push to GitHub Container Registry
|
||||
run: |
|
||||
docker tag release_build_image $REGISTRY/$REPOSITORY:$IMAGE_TAG
|
||||
docker push $REGISTRY/$REPOSITORY:$IMAGE_TAG
|
||||
env:
|
||||
REGISTRY: ghcr.io/triggerdotdev
|
||||
REPOSITORY: trigger.dev
|
||||
IMAGE_TAG: ${{ steps.get_version.outputs.version }}
|
||||
|
||||
- name: 🐙 Push 'latest' to GitHub Container Registry
|
||||
if: startsWith(github.ref, 'refs/tags/v.docker')
|
||||
run: |
|
||||
LATEST=latest
|
||||
if [[ ${{matrix.runs-on}} == *-arm ]]; then
|
||||
LATEST="${LATEST}-arm"
|
||||
fi
|
||||
docker tag release_build_image $REGISTRY/$REPOSITORY:$LATEST
|
||||
docker push $REGISTRY/$REPOSITORY:$LATEST
|
||||
env:
|
||||
REGISTRY: ghcr.io/triggerdotdev
|
||||
REPOSITORY: trigger.dev
|
||||
- name: 🐳 Build image and push to GitHub Container Registry
|
||||
uses: depot/build-push-action@v1
|
||||
with:
|
||||
file: ./docker/Dockerfile
|
||||
platforms: linux/amd64,linux/arm64
|
||||
tags: |
|
||||
ghcr.io/triggerdotdev/trigger.dev:${{ steps.get_version.outputs.version }}
|
||||
push: true
|
||||
|
||||
@@ -0,0 +1,94 @@
|
||||
name: "🚢 Publish Infra Images"
|
||||
|
||||
on:
|
||||
push:
|
||||
tags:
|
||||
- "infra-dev-*"
|
||||
- "infra-test-*"
|
||||
- "infra-prod-*"
|
||||
paths:
|
||||
- ".github/workflows/publish.yml"
|
||||
- "packages/**"
|
||||
- "!packages/**/*.md"
|
||||
- "!packages/**/*.eslintrc"
|
||||
- "apps/**"
|
||||
- "!apps/**/*.md"
|
||||
- "!apps/**/*.eslintrc"
|
||||
- "integrations/**"
|
||||
- "!integrations/**/*.md"
|
||||
- "!integrations/**/*.eslintrc"
|
||||
- "pnpm-lock.yaml"
|
||||
- "pnpm-workspace.yaml"
|
||||
- "turbo.json"
|
||||
- "docker/Dockerfile"
|
||||
- "docker/scripts/**"
|
||||
- "tests/**"
|
||||
|
||||
permissions:
|
||||
id-token: write
|
||||
packages: write
|
||||
contents: read
|
||||
|
||||
concurrency:
|
||||
group: ${{ github.workflow }}-${{ github.ref }}
|
||||
|
||||
env:
|
||||
AWS_REGION: us-east-1
|
||||
|
||||
jobs:
|
||||
build:
|
||||
strategy:
|
||||
matrix:
|
||||
package: [coordinator, kubernetes-provider]
|
||||
runs-on: buildjet-16vcpu-ubuntu-2204
|
||||
env:
|
||||
DOCKER_BUILDKIT: "1"
|
||||
steps:
|
||||
- uses: actions/checkout@v4
|
||||
|
||||
- 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)
|
||||
if [[ "${{ matrix.package }}" == *-provider ]]; then
|
||||
provider_type=$(echo ${{ matrix.package }} | cut -d- -f1)
|
||||
repository=provider/${provider_type}
|
||||
else
|
||||
repository=${{ matrix.package }}
|
||||
fi
|
||||
echo "IMAGE_TAG=${env}-${sha}-${ts}" >> "$GITHUB_OUTPUT"
|
||||
echo "REPOSITORY=${repository}" >> "$GITHUB_OUTPUT"
|
||||
|
||||
- name: Set up Docker Buildx
|
||||
uses: docker/setup-buildx-action@v3
|
||||
|
||||
# ..to avoid rate limits when pulling images
|
||||
- name: Login to DockerHub
|
||||
uses: docker/login-action@v3
|
||||
with:
|
||||
username: ${{ secrets.DOCKERHUB_USERNAME }}
|
||||
password: ${{ secrets.DOCKERHUB_TOKEN }}
|
||||
|
||||
- name: 🚢 Build Container Image
|
||||
run: |
|
||||
docker build -t infra_image -f ./apps/${{ matrix.package }}/Containerfile .
|
||||
|
||||
# ..to push image
|
||||
- name: 🐙 Login to GitHub Container Registry
|
||||
uses: docker/login-action@v3
|
||||
with:
|
||||
registry: ghcr.io
|
||||
username: ${{ github.repository_owner }}
|
||||
password: ${{ secrets.GITHUB_TOKEN }}
|
||||
|
||||
- name: 🐙 Push to GitHub Container Registry
|
||||
run: |
|
||||
docker tag infra_image $REGISTRY/$REPOSITORY:$IMAGE_TAG
|
||||
docker push $REGISTRY/$REPOSITORY:$IMAGE_TAG
|
||||
env:
|
||||
REGISTRY: ghcr.io/triggerdotdev
|
||||
REPOSITORY: ${{ steps.prep.outputs.REPOSITORY }}
|
||||
IMAGE_TAG: ${{ steps.prep.outputs.IMAGE_TAG }}
|
||||
@@ -9,6 +9,10 @@ on:
|
||||
- "build-*"
|
||||
paths:
|
||||
- ".github/workflows/publish.yml"
|
||||
- ".github/workflows/typecheck.yml"
|
||||
- ".github/workflows/unit-tests.yml"
|
||||
- ".github/workflows/e2e.yml"
|
||||
- ".github/workflows/publish-docker.yml"
|
||||
- "packages/**"
|
||||
- "!packages/**/*.md"
|
||||
- "!packages/**/*.eslintrc"
|
||||
|
||||
@@ -8,12 +8,11 @@ on:
|
||||
- "**.md"
|
||||
- ".github/CODEOWNERS"
|
||||
- ".github/ISSUE_TEMPLATE/**"
|
||||
|
||||
|
||||
jobs:
|
||||
release:
|
||||
name: 🦋 Changesets Release
|
||||
runs-on: buildjet-4vcpu-ubuntu-2204
|
||||
runs-on: buildjet-8vcpu-ubuntu-2204
|
||||
if: |
|
||||
github.repository == 'triggerdotdev/trigger.dev'
|
||||
outputs:
|
||||
@@ -31,14 +30,19 @@ jobs:
|
||||
- name: ⎔ Setup pnpm
|
||||
uses: pnpm/action-setup@v2.2.4
|
||||
with:
|
||||
version: 7.18
|
||||
version: 8.15.5
|
||||
|
||||
- name: ⎔ Setup node
|
||||
uses: buildjet/setup-node@v3
|
||||
with:
|
||||
node-version: 18
|
||||
node-version: 20.11.1
|
||||
cache: "pnpm"
|
||||
|
||||
- name: Install Protoc
|
||||
uses: arduino/setup-protoc@v3
|
||||
with:
|
||||
repo-token: ${{ secrets.GITHUB_TOKEN }}
|
||||
|
||||
- name: 📥 Download deps
|
||||
run: pnpm install --frozen-lockfile
|
||||
|
||||
@@ -46,7 +50,7 @@ jobs:
|
||||
run: pnpm run generate
|
||||
|
||||
- name: 🔎 Type check
|
||||
run: pnpm run typecheck
|
||||
run: pnpm run typecheck --filter "@trigger.dev/*" --filter "trigger.dev"
|
||||
|
||||
- name: 🔐 Setup npm auth
|
||||
run: |
|
||||
|
||||
@@ -3,7 +3,7 @@ on:
|
||||
workflow_call:
|
||||
jobs:
|
||||
typecheck:
|
||||
runs-on: buildjet-4vcpu-ubuntu-2204
|
||||
runs-on: buildjet-8vcpu-ubuntu-2204
|
||||
|
||||
steps:
|
||||
- name: ⬇️ Checkout repo
|
||||
@@ -14,14 +14,19 @@ jobs:
|
||||
- name: ⎔ Setup pnpm
|
||||
uses: pnpm/action-setup@v2.2.4
|
||||
with:
|
||||
version: 7.18
|
||||
version: 8.15.5
|
||||
|
||||
- name: ⎔ Setup node
|
||||
uses: buildjet/setup-node@v3
|
||||
with:
|
||||
node-version: 18
|
||||
node-version: 20.11.1
|
||||
cache: "pnpm"
|
||||
|
||||
- name: Install Protoc
|
||||
uses: arduino/setup-protoc@v3
|
||||
with:
|
||||
repo-token: ${{ secrets.GITHUB_TOKEN }}
|
||||
|
||||
- name: 📥 Download deps
|
||||
run: pnpm install --frozen-lockfile
|
||||
|
||||
|
||||
@@ -4,7 +4,7 @@ on:
|
||||
jobs:
|
||||
unitTests:
|
||||
name: "🧪 Unit Tests"
|
||||
runs-on: buildjet-4vcpu-ubuntu-2204
|
||||
runs-on: buildjet-8vcpu-ubuntu-2204
|
||||
steps:
|
||||
- name: ⬇️ Checkout repo
|
||||
uses: actions/checkout@v3
|
||||
@@ -14,14 +14,19 @@ jobs:
|
||||
- name: ⎔ Setup pnpm
|
||||
uses: pnpm/action-setup@v2.2.4
|
||||
with:
|
||||
version: 7.18
|
||||
version: 8.15.5
|
||||
|
||||
- name: ⎔ Setup node
|
||||
uses: buildjet/setup-node@v3
|
||||
with:
|
||||
node-version: 18
|
||||
node-version: 20.11.1
|
||||
cache: "pnpm"
|
||||
|
||||
- name: Install Protoc
|
||||
uses: arduino/setup-protoc@v3
|
||||
with:
|
||||
repo-token: ${{ secrets.GITHUB_TOKEN }}
|
||||
|
||||
- name: ⎔ Setup Deno
|
||||
uses: denoland/setup-deno@v1
|
||||
with:
|
||||
|
||||
+2
-1
@@ -53,4 +53,5 @@ apps/**/public/build
|
||||
/playwright-report/
|
||||
/playwright/.cache/
|
||||
|
||||
.cosine
|
||||
.cosine
|
||||
.trigger/
|
||||
@@ -0,0 +1,3 @@
|
||||
[submodule "packages/otlp-importer/protos"]
|
||||
path = packages/otlp-importer/protos
|
||||
url = https://github.com/open-telemetry/opentelemetry-proto.git
|
||||
Vendored
+26
-2
@@ -9,7 +9,7 @@
|
||||
"request": "launch",
|
||||
"name": "Debug WebApp",
|
||||
"command": "pnpm run dev --filter webapp",
|
||||
"envFile": "${workspaceFolder}/apps/webapp/.env",
|
||||
"envFile": "${workspaceFolder}/.env",
|
||||
"cwd": "${workspaceFolder}",
|
||||
"sourceMaps": true
|
||||
},
|
||||
@@ -23,11 +23,35 @@
|
||||
{
|
||||
"type": "node-terminal",
|
||||
"request": "launch",
|
||||
"name": "Debug BYO Auth",
|
||||
"name": "Debug v2 job catalog",
|
||||
"command": "pnpm run byo-auth",
|
||||
"envFile": "${workspaceFolder}/references/job-catalog/.env",
|
||||
"cwd": "${workspaceFolder}/references/job-catalog",
|
||||
"sourceMaps": true
|
||||
},
|
||||
{
|
||||
"type": "node-terminal",
|
||||
"request": "launch",
|
||||
"name": "Debug V3 Dev CLI",
|
||||
"command": "pnpm exec triggerdev dev --log-level debug",
|
||||
"cwd": "${workspaceFolder}/references/v3-catalog",
|
||||
"sourceMaps": true
|
||||
},
|
||||
{
|
||||
"type": "node-terminal",
|
||||
"request": "launch",
|
||||
"name": "Debug V3 Deploy CLI",
|
||||
"command": "pnpm exec triggerdev deploy",
|
||||
"cwd": "${workspaceFolder}/references/v3-catalog",
|
||||
"sourceMaps": true
|
||||
},
|
||||
{
|
||||
"type": "node",
|
||||
"request": "attach",
|
||||
"name": "Attach to Trigger.dev CLI (v3)",
|
||||
"port": 9229,
|
||||
"restart": true,
|
||||
"skipFiles": ["<node_internals>/**"]
|
||||
}
|
||||
]
|
||||
}
|
||||
|
||||
+1
-1
@@ -47,7 +47,7 @@ pnpm exec changeset version --snapshot prerelease
|
||||
3. Build the packages:
|
||||
|
||||
```sh
|
||||
pnpm run build --filter "@trigger.dev/*"
|
||||
pnpm run build --filter "@trigger.dev/*" --filter "trigger.dev"
|
||||
```
|
||||
|
||||
4. Publish the snapshot (replace "dev" with your tag)
|
||||
|
||||
@@ -7,7 +7,7 @@
|
||||
|
||||
### The open source background jobs framework
|
||||
|
||||
[Discord](https://discord.gg/JtBAxBr2m3) | [Website](https://trigger.dev) | [Issues](https://github.com/triggerdotdev/trigger.dev/issues) | [Docs](https://trigger.dev/docs)
|
||||
[Discord](https://trigger.dev/discord) | [Website](https://trigger.dev) | [Issues](https://github.com/triggerdotdev/trigger.dev/issues) | [Docs](https://trigger.dev/docs)
|
||||
|
||||
[](https://twitter.com/triggerdotdev)
|
||||
[](https://github.com/triggerdotdev/trigger.dev)
|
||||
|
||||
@@ -0,0 +1,4 @@
|
||||
HTTP_SERVER_PORT=8020
|
||||
PLATFORM_ENABLED=true
|
||||
PLATFORM_WS_PORT=3030
|
||||
SECURE_CONNECTION=false
|
||||
@@ -0,0 +1,3 @@
|
||||
dist/
|
||||
node_modules/
|
||||
.env
|
||||
@@ -0,0 +1,60 @@
|
||||
# syntax=docker/dockerfile:labs
|
||||
|
||||
FROM node:18-bullseye-slim@sha256:a4edd54dcfdcacc8a4100fee71498e8671d99556a1acf5614539214a70092426 AS node-18
|
||||
|
||||
WORKDIR /app
|
||||
|
||||
FROM node-18 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
|
||||
|
||||
RUN apt-get update \
|
||||
&& apt-get install -y buildah ca-certificates dumb-init \
|
||||
&& rm -rf /var/lib/apt/lists/*
|
||||
|
||||
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 coordinator build:bundle
|
||||
|
||||
FROM alpine AS cri-tools
|
||||
|
||||
WORKDIR /cri-tools
|
||||
|
||||
ARG CRICTL_VERSION=v1.29.0
|
||||
ARG CRICTL_CHECKSUM=sha256:d16a1ffb3938f5a19d5c8f45d363bd091ef89c0bc4d44ad16b933eede32fdcbb
|
||||
ADD --checksum=${CRICTL_CHECKSUM} \
|
||||
https://github.com/kubernetes-sigs/cri-tools/releases/download/${CRICTL_VERSION}/crictl-${CRICTL_VERSION}-linux-amd64.tar.gz .
|
||||
RUN tar zxvf crictl-${CRICTL_VERSION}-linux-amd64.tar.gz
|
||||
|
||||
FROM base AS runner
|
||||
|
||||
RUN corepack enable
|
||||
ENV NODE_ENV production
|
||||
|
||||
COPY --from=cri-tools --chown=node:node /cri-tools/crictl /usr/local/bin
|
||||
COPY --from=builder --chown=node:node /app/apps/coordinator/dist/index.mjs ./index.mjs
|
||||
|
||||
EXPOSE 8000
|
||||
|
||||
CMD [ "/usr/bin/dumb-init", "--", "/usr/local/bin/node", "./index.mjs" ]
|
||||
@@ -0,0 +1,3 @@
|
||||
# Coordinator
|
||||
|
||||
Sits between the platform and tasks. Facilitates communication and checkpointing, amongst other things.
|
||||
@@ -0,0 +1,34 @@
|
||||
{
|
||||
"name": "coordinator",
|
||||
"private": true,
|
||||
"version": "0.0.1",
|
||||
"description": "",
|
||||
"main": "dist/index.cjs",
|
||||
"scripts": {
|
||||
"build": "npm run build:bundle",
|
||||
"build:bundle": "esbuild src/index.ts --bundle --outfile=dist/index.mjs --platform=node --format=esm --target=esnext --banner:js=\"import { createRequire } from 'module';const require = createRequire(import.meta.url);\"",
|
||||
"build:image": "docker build -f Containerfile . -t coordinator",
|
||||
"dev": "tsx --no-warnings=ExperimentalWarning --require dotenv/config --watch src/index.ts",
|
||||
"start": "tsx src/index.ts",
|
||||
"typecheck": "tsc --noEmit"
|
||||
},
|
||||
"keywords": [],
|
||||
"author": "",
|
||||
"license": "MIT",
|
||||
"dependencies": {
|
||||
"@trigger.dev/core": "workspace:*",
|
||||
"@trigger.dev/core-apps": "workspace:*",
|
||||
"execa": "^8.0.1",
|
||||
"nanoid": "^5.0.6",
|
||||
"prom-client": "^15.1.0",
|
||||
"socket.io": "^4.7.4",
|
||||
"socket.io-client": "^4.7.4"
|
||||
},
|
||||
"devDependencies": {
|
||||
"@types/node": "^18",
|
||||
"dotenv": "^16.4.2",
|
||||
"esbuild": "^0.19.11",
|
||||
"tsx": "^4.7.0",
|
||||
"typescript": "^5.3.3"
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,983 @@
|
||||
import { createServer } from "node:http";
|
||||
import { $ } from "execa";
|
||||
import { nanoid } from "nanoid";
|
||||
import { Server } from "socket.io";
|
||||
import {
|
||||
CoordinatorToPlatformMessages,
|
||||
CoordinatorToProdWorkerMessages,
|
||||
PlatformToCoordinatorMessages,
|
||||
ProdWorkerSocketData,
|
||||
ProdWorkerToCoordinatorMessages,
|
||||
ZodNamespace,
|
||||
ZodSocketConnection,
|
||||
} from "@trigger.dev/core/v3";
|
||||
import { HttpReply, getTextBody, SimpleLogger } from "@trigger.dev/core-apps";
|
||||
|
||||
import { collectDefaultMetrics, register, Gauge } from "prom-client";
|
||||
collectDefaultMetrics();
|
||||
|
||||
const HTTP_SERVER_PORT = Number(process.env.HTTP_SERVER_PORT || 8020);
|
||||
const NODE_NAME = process.env.NODE_NAME || "coordinator";
|
||||
const DEFAULT_RETRY_DELAY_THRESHOLD_IN_MS = 30_000;
|
||||
|
||||
const REGISTRY_HOST = process.env.REGISTRY_HOST || "localhost:5000";
|
||||
const CHECKPOINT_PATH = process.env.CHECKPOINT_PATH || "/checkpoints";
|
||||
const REGISTRY_TLS_VERIFY = process.env.REGISTRY_TLS_VERIFY === "false" ? "false" : "true";
|
||||
|
||||
const PLATFORM_ENABLED = ["1", "true"].includes(process.env.PLATFORM_ENABLED ?? "true");
|
||||
const PLATFORM_HOST = process.env.PLATFORM_HOST || "127.0.0.1";
|
||||
const PLATFORM_WS_PORT = process.env.PLATFORM_WS_PORT || 3030;
|
||||
const PLATFORM_SECRET = process.env.PLATFORM_SECRET || "coordinator-secret";
|
||||
const SECURE_CONNECTION = ["1", "true"].includes(process.env.SECURE_CONNECTION ?? "false");
|
||||
|
||||
const logger = new SimpleLogger(`[${NODE_NAME}]`);
|
||||
|
||||
type CheckpointerInitializeReturn = {
|
||||
canCheckpoint: boolean;
|
||||
willSimulate: boolean;
|
||||
};
|
||||
|
||||
type CheckpointAndPushOptions = {
|
||||
runId: string;
|
||||
leaveRunning?: boolean;
|
||||
projectRef: string;
|
||||
deploymentVersion: string;
|
||||
};
|
||||
|
||||
type CheckpointData = {
|
||||
location: string;
|
||||
docker: boolean;
|
||||
};
|
||||
|
||||
class Checkpointer {
|
||||
#initialized = false;
|
||||
#canCheckpoint = false;
|
||||
#dockerMode = !process.env.KUBERNETES_PORT;
|
||||
|
||||
#logger = new SimpleLogger("[checkptr]");
|
||||
#abortControllers = new Map<string, AbortController>();
|
||||
|
||||
constructor(private opts = { forceSimulate: false }) {}
|
||||
|
||||
async initialize(): Promise<CheckpointerInitializeReturn> {
|
||||
if (this.#initialized) {
|
||||
return this.#getInitializeReturn();
|
||||
}
|
||||
|
||||
this.#logger.log(`${this.#dockerMode ? "Docker" : "Kubernetes"} mode`);
|
||||
|
||||
if (this.#dockerMode) {
|
||||
try {
|
||||
await $`criu --version`;
|
||||
} catch (error) {
|
||||
this.#logger.error("No checkpoint support: Missing CRIU binary");
|
||||
this.#logger.error("Will simulate instead");
|
||||
this.#canCheckpoint = false;
|
||||
this.#initialized = true;
|
||||
|
||||
return this.#getInitializeReturn();
|
||||
}
|
||||
|
||||
try {
|
||||
await $`docker checkpoint`;
|
||||
} catch (error) {
|
||||
this.#logger.error(
|
||||
"No checkpoint support: Docker needs to have experimental features enabled"
|
||||
);
|
||||
this.#logger.error("Will simulate instead");
|
||||
this.#canCheckpoint = false;
|
||||
this.#initialized = true;
|
||||
|
||||
return this.#getInitializeReturn();
|
||||
}
|
||||
} else {
|
||||
try {
|
||||
await $`buildah login --get-login ${REGISTRY_HOST}`;
|
||||
} catch (error) {
|
||||
this.#logger.error(`No checkpoint support: Not logged in to registry ${REGISTRY_HOST}`);
|
||||
this.#canCheckpoint = false;
|
||||
this.#initialized = true;
|
||||
|
||||
return this.#getInitializeReturn();
|
||||
}
|
||||
}
|
||||
|
||||
this.#logger.log(
|
||||
`Full checkpoint support${
|
||||
this.#dockerMode && this.opts.forceSimulate ? " with forced simulation enabled." : "!"
|
||||
}`
|
||||
);
|
||||
|
||||
this.#initialized = true;
|
||||
this.#canCheckpoint = true;
|
||||
|
||||
return this.#getInitializeReturn();
|
||||
}
|
||||
|
||||
#getInitializeReturn(): CheckpointerInitializeReturn {
|
||||
return {
|
||||
canCheckpoint: this.#canCheckpoint,
|
||||
willSimulate: this.#dockerMode && (!this.#canCheckpoint || this.opts.forceSimulate),
|
||||
};
|
||||
}
|
||||
|
||||
#getImageRef(projectRef: string, deploymentVersion: string, shortCode: string) {
|
||||
return `${REGISTRY_HOST}/trigger/${projectRef}:${deploymentVersion}.prod-${shortCode}`;
|
||||
}
|
||||
|
||||
#getExportLocation(projectRef: string, deploymentVersion: string, shortCode: string) {
|
||||
const basename = `${projectRef}-${deploymentVersion}-${shortCode}`;
|
||||
|
||||
if (this.#dockerMode) {
|
||||
return basename;
|
||||
} else {
|
||||
return `${CHECKPOINT_PATH}/${basename}.tar`;
|
||||
}
|
||||
}
|
||||
|
||||
async checkpointAndPush(opts: CheckpointAndPushOptions): Promise<CheckpointData | undefined> {
|
||||
const start = performance.now();
|
||||
logger.log(`checkpointAndPush() start`, { start, opts });
|
||||
|
||||
const result = await this.#checkpointAndPush(opts);
|
||||
|
||||
const end = performance.now();
|
||||
logger.log(`checkpointAndPush() end`, {
|
||||
start,
|
||||
end,
|
||||
diff: end - start,
|
||||
opts,
|
||||
success: !!result,
|
||||
});
|
||||
|
||||
return result;
|
||||
}
|
||||
|
||||
isCheckpointing(runId: string) {
|
||||
return this.#abortControllers.has(runId);
|
||||
}
|
||||
|
||||
cancelCheckpoint(runId: string) {
|
||||
const controller = this.#abortControllers.get(runId);
|
||||
|
||||
if (!controller) {
|
||||
logger.debug("Nothing to cancel", { runId });
|
||||
return;
|
||||
}
|
||||
|
||||
controller.abort("cancelCheckpointing()");
|
||||
this.#abortControllers.delete(runId);
|
||||
}
|
||||
|
||||
async #checkpointAndPush({
|
||||
runId,
|
||||
leaveRunning = true, // This mirrors kubernetes behaviour more accurately
|
||||
projectRef,
|
||||
deploymentVersion,
|
||||
}: CheckpointAndPushOptions): Promise<CheckpointData | undefined> {
|
||||
await this.initialize();
|
||||
|
||||
if (!this.#dockerMode && !this.#canCheckpoint) {
|
||||
this.#logger.error("No checkpoint support. Simulation requires docker.");
|
||||
return;
|
||||
}
|
||||
|
||||
if (this.#abortControllers.has(runId)) {
|
||||
logger.error("Checkpoint procedure already in progress", {
|
||||
options: {
|
||||
runId,
|
||||
leaveRunning,
|
||||
projectRef,
|
||||
deploymentVersion,
|
||||
},
|
||||
});
|
||||
return;
|
||||
}
|
||||
|
||||
const controller = new AbortController();
|
||||
this.#abortControllers.set(runId, controller);
|
||||
|
||||
const $$ = $({ signal: controller.signal });
|
||||
|
||||
try {
|
||||
const shortCode = nanoid(8);
|
||||
const imageRef = this.#getImageRef(projectRef, deploymentVersion, shortCode);
|
||||
const exportLocation = this.#getExportLocation(projectRef, deploymentVersion, shortCode);
|
||||
|
||||
this.#logger.log("Checkpointing:", {
|
||||
options: {
|
||||
runId,
|
||||
leaveRunning,
|
||||
projectRef,
|
||||
deploymentVersion,
|
||||
},
|
||||
});
|
||||
|
||||
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 (leaveRunning) {
|
||||
this.#logger.debug(
|
||||
await $$`docker checkpoint create --leave-running ${containterName} ${exportLocation}`
|
||||
);
|
||||
} else {
|
||||
this.#logger.debug(
|
||||
await $$`docker checkpoint create ${containterName} ${exportLocation}`
|
||||
);
|
||||
}
|
||||
}
|
||||
} catch (error: any) {
|
||||
this.#logger.error(error.stderr);
|
||||
return;
|
||||
}
|
||||
|
||||
this.#logger.log("checkpoint created:", {
|
||||
runId,
|
||||
location: exportLocation,
|
||||
});
|
||||
|
||||
return {
|
||||
location: exportLocation,
|
||||
docker: true,
|
||||
};
|
||||
}
|
||||
|
||||
// Create checkpoint (CRI)
|
||||
if (!this.#canCheckpoint) {
|
||||
throw new Error("No checkpoint support in kubernetes mode.");
|
||||
}
|
||||
|
||||
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) {
|
||||
throw new Error("could not find container id");
|
||||
}
|
||||
|
||||
this.#logger.debug(await $$`crictl checkpoint --export=${exportLocation} ${containerId}`);
|
||||
|
||||
// Create image from checkpoint
|
||||
const container = this.#logger.debug(await $$`buildah from scratch`);
|
||||
this.#logger.debug(await $$`buildah add ${container} ${exportLocation} /`);
|
||||
this.#logger.debug(
|
||||
await $$`buildah config --annotation=io.kubernetes.cri-o.annotations.checkpoint.name=counter ${container}`
|
||||
);
|
||||
this.#logger.debug(await $$`buildah commit ${container} ${imageRef}`);
|
||||
this.#logger.debug(await $$`buildah rm ${container}`);
|
||||
|
||||
// Push checkpoint image
|
||||
this.#logger.debug(await $$`buildah push --tls-verify=${REGISTRY_TLS_VERIFY} ${imageRef}`);
|
||||
|
||||
this.#logger.log("Checkpointed and pushed image to:", { location: imageRef });
|
||||
|
||||
try {
|
||||
await $$`rm ${exportLocation}`;
|
||||
this.#logger.log("Deleted checkpoint archive", { exportLocation });
|
||||
|
||||
// Disabled for now as this will increase restore time by having to pull the image again
|
||||
// await $`buildah rmi ${imageRef}`;
|
||||
// this.#logger.log("Deleted checkpoint image", { imageRef });
|
||||
} catch (error) {
|
||||
this.#logger.error("Failed during checkpoint cleanup", { exportLocation });
|
||||
this.#logger.debug(error);
|
||||
}
|
||||
|
||||
return {
|
||||
location: imageRef,
|
||||
docker: false,
|
||||
};
|
||||
} catch (error) {
|
||||
this.#logger.error("checkpoint failed", {
|
||||
options: {
|
||||
runId,
|
||||
leaveRunning,
|
||||
projectRef,
|
||||
deploymentVersion,
|
||||
},
|
||||
error,
|
||||
});
|
||||
return;
|
||||
} finally {
|
||||
this.#abortControllers.delete(runId);
|
||||
}
|
||||
}
|
||||
|
||||
#getRunContainerName(suffix: string) {
|
||||
return `task-run-${suffix}`;
|
||||
}
|
||||
}
|
||||
|
||||
class TaskCoordinator {
|
||||
#httpServer: ReturnType<typeof createServer>;
|
||||
#checkpointer = new Checkpointer({ forceSimulate: true });
|
||||
|
||||
#prodWorkerNamespace: ZodNamespace<
|
||||
typeof ProdWorkerToCoordinatorMessages,
|
||||
typeof CoordinatorToProdWorkerMessages,
|
||||
typeof ProdWorkerSocketData
|
||||
>;
|
||||
#platformSocket?: ZodSocketConnection<
|
||||
typeof CoordinatorToPlatformMessages,
|
||||
typeof PlatformToCoordinatorMessages
|
||||
>;
|
||||
|
||||
#checkpointableTasks = new Map<
|
||||
string,
|
||||
{ resolve: (value: void) => void; reject: (err?: any) => void }
|
||||
>();
|
||||
|
||||
#delayThresholdInMs: number;
|
||||
|
||||
constructor(
|
||||
private port: number,
|
||||
private host = "0.0.0.0"
|
||||
) {
|
||||
this.#httpServer = this.#createHttpServer();
|
||||
this.#checkpointer.initialize();
|
||||
this.#delayThresholdInMs = this.#getDelayThreshold();
|
||||
|
||||
if (process.env.DELAY_THRESHOLD_IN_MS) {
|
||||
this.#delayThresholdInMs = this.#getDelayThreshold();
|
||||
}
|
||||
|
||||
const io = new Server(this.#httpServer);
|
||||
this.#prodWorkerNamespace = this.#createProdWorkerNamespace(io);
|
||||
|
||||
this.#platformSocket = this.#createPlatformSocket();
|
||||
|
||||
const connectedTasksTotal = new Gauge({
|
||||
name: "daemon_connected_tasks_total", // don't change this without updating dashboard config
|
||||
help: "The number of tasks currently connected.",
|
||||
collect: () => {
|
||||
connectedTasksTotal.set(this.#prodWorkerNamespace.namespace.sockets.size);
|
||||
},
|
||||
});
|
||||
register.registerMetric(connectedTasksTotal);
|
||||
}
|
||||
|
||||
#getDelayThreshold() {
|
||||
if (!process.env.RETRY_DELAY_THRESHOLD_IN_MS) {
|
||||
return DEFAULT_RETRY_DELAY_THRESHOLD_IN_MS;
|
||||
}
|
||||
|
||||
const threshold = parseInt(process.env.RETRY_DELAY_THRESHOLD_IN_MS);
|
||||
|
||||
if (isNaN(threshold)) {
|
||||
logger.log(
|
||||
"RETRY_DELAY_THRESHOLD_IN_MS parses as NaN, must supply integer. Will use default instead.",
|
||||
{
|
||||
RETRY_DELAY_THRESHOLD_IN_MS: process.env.RETRY_DELAY_THRESHOLD_IN_MS,
|
||||
DEFAULT_DELAY_THRESHOLD_IN_MS: DEFAULT_RETRY_DELAY_THRESHOLD_IN_MS,
|
||||
}
|
||||
);
|
||||
return DEFAULT_RETRY_DELAY_THRESHOLD_IN_MS;
|
||||
}
|
||||
|
||||
return threshold;
|
||||
}
|
||||
|
||||
#createPlatformSocket() {
|
||||
if (!PLATFORM_ENABLED) {
|
||||
console.log("INFO: platform connection disabled");
|
||||
return;
|
||||
}
|
||||
|
||||
const platformConnection = new ZodSocketConnection({
|
||||
namespace: "coordinator",
|
||||
host: PLATFORM_HOST,
|
||||
port: Number(PLATFORM_WS_PORT),
|
||||
secure: SECURE_CONNECTION,
|
||||
clientMessages: CoordinatorToPlatformMessages,
|
||||
serverMessages: PlatformToCoordinatorMessages,
|
||||
authToken: PLATFORM_SECRET,
|
||||
handlers: {
|
||||
RESUME_AFTER_DEPENDENCY: async (message) => {
|
||||
const taskSocket = await this.#getAttemptSocket(message.attemptFriendlyId);
|
||||
|
||||
if (!taskSocket) {
|
||||
logger.log("Socket for attempt not found", {
|
||||
attemptFriendlyId: message.attemptFriendlyId,
|
||||
});
|
||||
return;
|
||||
}
|
||||
|
||||
// In case the task resumed faster than we could checkpoint
|
||||
this.#cancelCheckpoint(message.runId);
|
||||
|
||||
taskSocket.emit("RESUME_AFTER_DEPENDENCY", message);
|
||||
},
|
||||
RESUME_AFTER_DURATION: async (message) => {
|
||||
const taskSocket = await this.#getAttemptSocket(message.attemptFriendlyId);
|
||||
|
||||
if (!taskSocket) {
|
||||
logger.log("Socket for attempt not found", {
|
||||
attemptFriendlyId: message.attemptFriendlyId,
|
||||
});
|
||||
return;
|
||||
}
|
||||
|
||||
taskSocket.emit("RESUME_AFTER_DURATION", message);
|
||||
},
|
||||
REQUEST_ATTEMPT_CANCELLATION: async (message) => {
|
||||
const taskSocket = await this.#getAttemptSocket(message.attemptFriendlyId);
|
||||
|
||||
if (!taskSocket) {
|
||||
logger.log("Socket for attempt not found", {
|
||||
attemptFriendlyId: message.attemptFriendlyId,
|
||||
});
|
||||
return;
|
||||
}
|
||||
|
||||
taskSocket.emit("REQUEST_ATTEMPT_CANCELLATION", message);
|
||||
},
|
||||
READY_FOR_RETRY: async (message) => {
|
||||
const taskSocket = await this.#getRunSocket(message.runId);
|
||||
|
||||
if (!taskSocket) {
|
||||
logger.log("Socket for attempt not found", {
|
||||
runId: message.runId,
|
||||
});
|
||||
return;
|
||||
}
|
||||
|
||||
taskSocket.emit("READY_FOR_RETRY", message);
|
||||
},
|
||||
},
|
||||
});
|
||||
|
||||
return platformConnection;
|
||||
}
|
||||
|
||||
async #getRunSocket(runId: string) {
|
||||
const sockets = await this.#prodWorkerNamespace.fetchSockets();
|
||||
|
||||
for (const socket of sockets) {
|
||||
if (socket.data.runId === runId) {
|
||||
return socket;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async #getAttemptSocket(attemptFriendlyId: string) {
|
||||
const sockets = await this.#prodWorkerNamespace.fetchSockets();
|
||||
|
||||
for (const socket of sockets) {
|
||||
if (socket.data.attemptFriendlyId === attemptFriendlyId) {
|
||||
return socket;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#createProdWorkerNamespace(io: Server) {
|
||||
const provider = new ZodNamespace({
|
||||
io,
|
||||
name: "prod-worker",
|
||||
clientMessages: ProdWorkerToCoordinatorMessages,
|
||||
serverMessages: CoordinatorToProdWorkerMessages,
|
||||
socketData: ProdWorkerSocketData,
|
||||
postAuth: async (socket, next, logger) => {
|
||||
function setSocketDataFromHeader(
|
||||
dataKey: keyof typeof socket.data,
|
||||
headerName: string,
|
||||
required: boolean = true
|
||||
) {
|
||||
const value = socket.handshake.headers[headerName];
|
||||
|
||||
if (value) {
|
||||
socket.data[dataKey] = Array.isArray(value) ? value[0] : value;
|
||||
return;
|
||||
}
|
||||
|
||||
if (required) {
|
||||
logger.error("missing required header", { headerName });
|
||||
throw new Error("missing header");
|
||||
}
|
||||
}
|
||||
|
||||
try {
|
||||
setSocketDataFromHeader("podName", "x-pod-name");
|
||||
setSocketDataFromHeader("contentHash", "x-trigger-content-hash");
|
||||
setSocketDataFromHeader("projectRef", "x-trigger-project-ref");
|
||||
setSocketDataFromHeader("runId", "x-trigger-run-id");
|
||||
setSocketDataFromHeader("attemptFriendlyId", "x-trigger-attempt-friendly-id", false);
|
||||
setSocketDataFromHeader("envId", "x-trigger-env-id");
|
||||
setSocketDataFromHeader("deploymentId", "x-trigger-deployment-id");
|
||||
setSocketDataFromHeader("deploymentVersion", "x-trigger-deployment-version");
|
||||
} catch (error) {
|
||||
logger.error("setSocketDataFromHeader error", { error });
|
||||
socket.disconnect(true);
|
||||
return;
|
||||
}
|
||||
|
||||
logger.debug("success", socket.data);
|
||||
|
||||
next();
|
||||
},
|
||||
onConnection: async (socket, handler, sender) => {
|
||||
const logger = new SimpleLogger(`[prod-worker][${socket.id}]`);
|
||||
|
||||
const checkpointInProgress = () => {
|
||||
return this.#checkpointableTasks.has(socket.data.runId);
|
||||
};
|
||||
|
||||
const readyToCheckpoint = async (): Promise<
|
||||
{ success: true } | { success: false; reason?: string }
|
||||
> => {
|
||||
if (checkpointInProgress()) {
|
||||
return {
|
||||
success: false,
|
||||
reason: "checkpoint in progress",
|
||||
};
|
||||
}
|
||||
|
||||
const isCheckpointable = new Promise((resolve, reject) => {
|
||||
// We set a reasonable timeout to prevent waiting forever
|
||||
// TODO: We may also want to cancel the task as it's unlikely to recover
|
||||
setTimeout(() => reject("timeout"), 10_000);
|
||||
|
||||
this.#checkpointableTasks.set(socket.data.runId, { resolve, reject });
|
||||
});
|
||||
|
||||
try {
|
||||
await isCheckpointable;
|
||||
this.#checkpointableTasks.delete(socket.data.runId);
|
||||
|
||||
return {
|
||||
success: true,
|
||||
};
|
||||
} catch (error) {
|
||||
logger.error("Error while waiting for checkpointable state", { error });
|
||||
|
||||
return {
|
||||
success: false,
|
||||
reason: typeof error === "string" ? error : "unknown",
|
||||
};
|
||||
}
|
||||
};
|
||||
|
||||
this.#platformSocket?.send("LOG", {
|
||||
metadata: socket.data,
|
||||
text: "connected",
|
||||
});
|
||||
|
||||
socket.on("LOG", (message, callback) => {
|
||||
logger.log("[LOG]", message.text);
|
||||
|
||||
callback();
|
||||
|
||||
this.#platformSocket?.send("LOG", {
|
||||
version: "v1",
|
||||
metadata: socket.data,
|
||||
text: message.text,
|
||||
});
|
||||
});
|
||||
|
||||
socket.on("READY_FOR_EXECUTION", async (message) => {
|
||||
logger.log("[READY_FOR_EXECUTION]", message);
|
||||
|
||||
try {
|
||||
const executionAck = await this.#platformSocket?.sendWithAck(
|
||||
"READY_FOR_EXECUTION",
|
||||
message
|
||||
);
|
||||
|
||||
if (!executionAck) {
|
||||
logger.error("no execution ack", { runId: socket.data.runId });
|
||||
|
||||
socket.emit("REQUEST_EXIT", {
|
||||
version: "v1",
|
||||
});
|
||||
|
||||
return;
|
||||
}
|
||||
|
||||
if (!executionAck.success) {
|
||||
logger.error("failed to get execution payload", { runId: socket.data.runId });
|
||||
|
||||
socket.emit("REQUEST_EXIT", {
|
||||
version: "v1",
|
||||
});
|
||||
|
||||
return;
|
||||
}
|
||||
|
||||
socket.emit("EXECUTE_TASK_RUN", {
|
||||
version: "v1",
|
||||
executionPayload: executionAck.payload,
|
||||
});
|
||||
|
||||
socket.data.attemptFriendlyId = executionAck.payload.execution.attempt.id;
|
||||
} catch (error) {
|
||||
logger.error("Error", { error });
|
||||
}
|
||||
});
|
||||
|
||||
socket.on("READY_FOR_RESUME", async (message) => {
|
||||
logger.log("[READY_FOR_RESUME]", message);
|
||||
|
||||
socket.data.attemptFriendlyId = message.attemptFriendlyId;
|
||||
this.#platformSocket?.send("READY_FOR_RESUME", message);
|
||||
});
|
||||
|
||||
socket.on("TASK_RUN_COMPLETED", async ({ completion, execution }, callback) => {
|
||||
logger.log("completed task", { completionId: completion.id });
|
||||
|
||||
const completeWithoutCheckpoint = (shouldExit: boolean) => {
|
||||
this.#platformSocket?.send("TASK_RUN_COMPLETED", {
|
||||
version: "v1",
|
||||
execution,
|
||||
completion,
|
||||
});
|
||||
callback({ willCheckpointAndRestore: false, shouldExit });
|
||||
};
|
||||
|
||||
if (completion.ok) {
|
||||
completeWithoutCheckpoint(true);
|
||||
return;
|
||||
}
|
||||
|
||||
if (
|
||||
completion.error.type === "INTERNAL_ERROR" &&
|
||||
completion.error.code === "TASK_RUN_CANCELLED"
|
||||
) {
|
||||
completeWithoutCheckpoint(true);
|
||||
return;
|
||||
}
|
||||
|
||||
if (completion.retry === undefined) {
|
||||
completeWithoutCheckpoint(true);
|
||||
return;
|
||||
}
|
||||
|
||||
if (completion.retry.delay < this.#delayThresholdInMs) {
|
||||
completeWithoutCheckpoint(false);
|
||||
return;
|
||||
}
|
||||
|
||||
const { canCheckpoint, willSimulate } = await this.#checkpointer.initialize();
|
||||
|
||||
const willCheckpointAndRestore = canCheckpoint || willSimulate;
|
||||
|
||||
if (!willCheckpointAndRestore) {
|
||||
completeWithoutCheckpoint(false);
|
||||
return;
|
||||
}
|
||||
|
||||
// The worker will then put itself in a checkpointable state
|
||||
callback({ willCheckpointAndRestore: true, shouldExit: false });
|
||||
|
||||
const ready = await readyToCheckpoint();
|
||||
|
||||
if (!ready.success) {
|
||||
logger.error("Failed to become checkpointable", {
|
||||
runId: socket.data.runId,
|
||||
reason: ready.reason,
|
||||
});
|
||||
return;
|
||||
}
|
||||
|
||||
const checkpoint = await this.#checkpointer.checkpointAndPush({
|
||||
runId: socket.data.runId,
|
||||
projectRef: socket.data.projectRef,
|
||||
deploymentVersion: socket.data.deploymentVersion,
|
||||
});
|
||||
|
||||
if (!checkpoint) {
|
||||
logger.error("Failed to checkpoint", { runId: socket.data.runId });
|
||||
completeWithoutCheckpoint(false);
|
||||
return;
|
||||
}
|
||||
|
||||
this.#platformSocket?.send("TASK_RUN_COMPLETED", {
|
||||
version: "v1",
|
||||
execution,
|
||||
completion,
|
||||
checkpoint,
|
||||
});
|
||||
|
||||
if (!checkpoint.docker || !willSimulate) {
|
||||
socket.emit("REQUEST_EXIT", {
|
||||
version: "v1",
|
||||
});
|
||||
}
|
||||
});
|
||||
|
||||
socket.on("READY_FOR_CHECKPOINT", async (message) => {
|
||||
logger.log("[READY_FOR_CHECKPOINT]", message);
|
||||
|
||||
const checkpointable = this.#checkpointableTasks.get(socket.data.runId);
|
||||
|
||||
if (!checkpointable) {
|
||||
logger.error("No checkpoint scheduled", { runId: socket.data.runId });
|
||||
return;
|
||||
}
|
||||
|
||||
checkpointable.resolve();
|
||||
});
|
||||
|
||||
socket.on("CANCEL_CHECKPOINT", async (message) => {
|
||||
logger.log("[CANCEL_CHECKPOINT]", message);
|
||||
|
||||
this.#cancelCheckpoint(socket.data.runId);
|
||||
});
|
||||
|
||||
socket.on("WAIT_FOR_DURATION", async (message, callback) => {
|
||||
logger.log("[WAIT_FOR_DURATION]", message);
|
||||
|
||||
if (checkpointInProgress()) {
|
||||
logger.error("Checkpoint already in progress", { runId: socket.data.runId });
|
||||
callback({ willCheckpointAndRestore: false });
|
||||
return;
|
||||
}
|
||||
|
||||
const { canCheckpoint, willSimulate } = await this.#checkpointer.initialize();
|
||||
|
||||
const willCheckpointAndRestore = canCheckpoint || willSimulate;
|
||||
|
||||
callback({ willCheckpointAndRestore });
|
||||
|
||||
if (!willCheckpointAndRestore) {
|
||||
return;
|
||||
}
|
||||
|
||||
const ready = await readyToCheckpoint();
|
||||
|
||||
if (!ready.success) {
|
||||
logger.error("Failed to become checkpointable", {
|
||||
runId: socket.data.runId,
|
||||
reason: ready.reason,
|
||||
});
|
||||
return;
|
||||
}
|
||||
|
||||
const checkpoint = await this.#checkpointer.checkpointAndPush({
|
||||
runId: socket.data.runId,
|
||||
projectRef: socket.data.projectRef,
|
||||
deploymentVersion: socket.data.deploymentVersion,
|
||||
});
|
||||
|
||||
if (!checkpoint) {
|
||||
// The task container will keep running until the wait duration has elapsed
|
||||
logger.error("Failed to checkpoint", { runId: socket.data.runId });
|
||||
return;
|
||||
}
|
||||
|
||||
if (!checkpoint.docker || !willSimulate) {
|
||||
socket.emit("REQUEST_EXIT", {
|
||||
version: "v1",
|
||||
});
|
||||
}
|
||||
|
||||
this.#platformSocket?.send("CHECKPOINT_CREATED", {
|
||||
version: "v1",
|
||||
attemptFriendlyId: message.attemptFriendlyId,
|
||||
docker: checkpoint.docker,
|
||||
location: checkpoint.location,
|
||||
reason: {
|
||||
type: "WAIT_FOR_DURATION",
|
||||
ms: message.ms,
|
||||
now: message.now,
|
||||
},
|
||||
});
|
||||
});
|
||||
|
||||
socket.on("WAIT_FOR_TASK", async (message, callback) => {
|
||||
logger.log("[WAIT_FOR_TASK]", message);
|
||||
|
||||
const { canCheckpoint, willSimulate } = await this.#checkpointer.initialize();
|
||||
|
||||
const willCheckpointAndRestore = canCheckpoint || willSimulate;
|
||||
|
||||
callback({ willCheckpointAndRestore });
|
||||
|
||||
if (!willCheckpointAndRestore) {
|
||||
return;
|
||||
}
|
||||
|
||||
const checkpoint = await this.#checkpointer.checkpointAndPush({
|
||||
runId: socket.data.runId,
|
||||
projectRef: socket.data.projectRef,
|
||||
deploymentVersion: socket.data.deploymentVersion,
|
||||
});
|
||||
|
||||
if (!checkpoint) {
|
||||
logger.error("Failed to checkpoint", { runId: socket.data.runId });
|
||||
return;
|
||||
}
|
||||
|
||||
if (!checkpoint.docker || !willSimulate) {
|
||||
socket.emit("REQUEST_EXIT", {
|
||||
version: "v1",
|
||||
});
|
||||
}
|
||||
|
||||
this.#platformSocket?.send("CHECKPOINT_CREATED", {
|
||||
version: "v1",
|
||||
attemptFriendlyId: message.attemptFriendlyId,
|
||||
docker: checkpoint.docker,
|
||||
location: checkpoint.location,
|
||||
reason: {
|
||||
type: "WAIT_FOR_TASK",
|
||||
friendlyId: message.friendlyId,
|
||||
},
|
||||
});
|
||||
});
|
||||
|
||||
socket.on("WAIT_FOR_BATCH", async (message, callback) => {
|
||||
logger.log("[WAIT_FOR_BATCH]", message);
|
||||
|
||||
const { canCheckpoint, willSimulate } = await this.#checkpointer.initialize();
|
||||
|
||||
const willCheckpointAndRestore = canCheckpoint || willSimulate;
|
||||
|
||||
callback({ willCheckpointAndRestore });
|
||||
|
||||
if (!willCheckpointAndRestore) {
|
||||
return;
|
||||
}
|
||||
|
||||
const checkpoint = await this.#checkpointer.checkpointAndPush({
|
||||
runId: socket.data.runId,
|
||||
projectRef: socket.data.projectRef,
|
||||
deploymentVersion: socket.data.deploymentVersion,
|
||||
});
|
||||
|
||||
if (!checkpoint) {
|
||||
logger.error("Failed to checkpoint", { runId: socket.data.runId });
|
||||
return;
|
||||
}
|
||||
|
||||
if (!checkpoint.docker || !willSimulate) {
|
||||
socket.emit("REQUEST_EXIT", {
|
||||
version: "v1",
|
||||
});
|
||||
}
|
||||
|
||||
this.#platformSocket?.send("CHECKPOINT_CREATED", {
|
||||
version: "v1",
|
||||
attemptFriendlyId: message.attemptFriendlyId,
|
||||
docker: checkpoint.docker,
|
||||
location: checkpoint.location,
|
||||
reason: {
|
||||
type: "WAIT_FOR_BATCH",
|
||||
batchFriendlyId: message.batchFriendlyId,
|
||||
runFriendlyIds: message.runFriendlyIds,
|
||||
},
|
||||
});
|
||||
});
|
||||
|
||||
socket.on("INDEX_TASKS", async (message, callback) => {
|
||||
logger.log("[INDEX_TASKS]", message);
|
||||
|
||||
const workerAck = await this.#platformSocket?.sendWithAck("CREATE_WORKER", {
|
||||
version: "v1",
|
||||
projectRef: socket.data.projectRef,
|
||||
envId: socket.data.envId,
|
||||
deploymentId: message.deploymentId,
|
||||
metadata: {
|
||||
contentHash: socket.data.contentHash,
|
||||
packageVersion: message.packageVersion,
|
||||
tasks: message.tasks,
|
||||
},
|
||||
});
|
||||
|
||||
if (!workerAck) {
|
||||
logger.debug("no worker ack while indexing", message);
|
||||
}
|
||||
|
||||
callback({ success: !!workerAck?.success });
|
||||
});
|
||||
|
||||
socket.on("INDEXING_FAILED", async (message) => {
|
||||
logger.log("[INDEXING_FAILED]", message);
|
||||
|
||||
this.#platformSocket?.send("INDEXING_FAILED", {
|
||||
version: "v1",
|
||||
deploymentId: message.deploymentId,
|
||||
error: message.error,
|
||||
});
|
||||
});
|
||||
},
|
||||
onDisconnect: async (socket, handler, sender, logger) => {
|
||||
this.#platformSocket?.send("LOG", {
|
||||
metadata: socket.data,
|
||||
text: "disconnect",
|
||||
});
|
||||
},
|
||||
handlers: {
|
||||
TASK_HEARTBEAT: async (message) => {
|
||||
this.#platformSocket?.send("TASK_HEARTBEAT", message);
|
||||
},
|
||||
},
|
||||
});
|
||||
|
||||
return provider;
|
||||
}
|
||||
|
||||
#cancelCheckpoint(runId: string) {
|
||||
const checkpointWait = this.#checkpointableTasks.get(runId);
|
||||
|
||||
if (checkpointWait) {
|
||||
// Stop waiting for task to reach checkpointable state
|
||||
checkpointWait.reject("Checkpoint cancelled");
|
||||
}
|
||||
|
||||
// Cancel checkpointing procedure
|
||||
this.#checkpointer.cancelCheckpoint(runId);
|
||||
}
|
||||
|
||||
#createHttpServer() {
|
||||
const httpServer = createServer(async (req, res) => {
|
||||
logger.log(`[${req.method}]`, req.url);
|
||||
|
||||
const reply = new HttpReply(res);
|
||||
|
||||
switch (req.url) {
|
||||
case "/health": {
|
||||
return reply.text("ok");
|
||||
}
|
||||
case "/metrics": {
|
||||
return reply.text(await register.metrics(), 200, register.contentType);
|
||||
}
|
||||
case "/whoami": {
|
||||
return reply.text(NODE_NAME);
|
||||
}
|
||||
case "/checkpoint": {
|
||||
const body = await getTextBody(req);
|
||||
// await this.#checkpointer.checkpointAndPush(body);
|
||||
return reply.text(`sent restore request: ${body}`);
|
||||
}
|
||||
default: {
|
||||
return reply.empty(404);
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
httpServer.on("clientError", (err, socket) => {
|
||||
socket.end("HTTP/1.1 400 Bad Request\r\n\r\n");
|
||||
});
|
||||
|
||||
httpServer.on("listening", () => {
|
||||
logger.log("server listening on port", HTTP_SERVER_PORT);
|
||||
});
|
||||
|
||||
return httpServer;
|
||||
}
|
||||
|
||||
listen() {
|
||||
this.#httpServer.listen(this.port, this.host);
|
||||
}
|
||||
}
|
||||
|
||||
const coordinator = new TaskCoordinator(HTTP_SERVER_PORT);
|
||||
coordinator.listen();
|
||||
@@ -0,0 +1,19 @@
|
||||
{
|
||||
"include": ["./src/**/*.ts"],
|
||||
"exclude": ["node_modules"],
|
||||
"compilerOptions": {
|
||||
"target": "es2016",
|
||||
"module": "commonjs",
|
||||
"esModuleInterop": true,
|
||||
"resolveJsonModule": true,
|
||||
"forceConsistentCasingInFileNames": true,
|
||||
"strict": true,
|
||||
"skipLibCheck": true,
|
||||
"paths": {
|
||||
"@trigger.dev/core/v3": ["../../packages/core/src/v3"],
|
||||
"@trigger.dev/core/v3/*": ["../../packages/core/src/v3/*"],
|
||||
"@trigger.dev/core-apps": ["../../packages/core-apps/src"],
|
||||
"@trigger.dev/core-apps/*": ["../../packages/core-apps/src/*"]
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,9 @@
|
||||
HTTP_SERVER_PORT=8050
|
||||
|
||||
PLATFORM_WS_PORT=3030
|
||||
PLATFORM_SECRET=provider-secret
|
||||
SECURE_CONNECTION=false
|
||||
|
||||
# Use this if you are on macOS
|
||||
# COORDINATOR_HOST="host.docker.internal"
|
||||
# OTEL_EXPORTER_OTLP_ENDPOINT="http://host.docker.internal:4318"
|
||||
@@ -0,0 +1,3 @@
|
||||
dist/
|
||||
node_modules/
|
||||
.env
|
||||
@@ -0,0 +1,16 @@
|
||||
# syntax=docker/dockerfile:labs
|
||||
|
||||
FROM node:18-slim AS base
|
||||
|
||||
RUN apt-get update \
|
||||
&& apt-get install -y dumb-init
|
||||
|
||||
FROM base
|
||||
|
||||
WORKDIR /app
|
||||
|
||||
COPY --chown=node dist/index.mjs /app/
|
||||
|
||||
EXPOSE 8000
|
||||
|
||||
ENTRYPOINT [ "/usr/bin/dumb-init", "--", "/usr/local/bin/node", "/app/index.mjs" ]
|
||||
@@ -0,0 +1,3 @@
|
||||
# Docker provider
|
||||
|
||||
The `docker-provider` allows the platform to be orchestrator-agnostic. The platform can perform actions such as `INDEX_TASKS` or `INVOKE_TASK` which the provider translates into Docker actions.
|
||||
@@ -0,0 +1,31 @@
|
||||
{
|
||||
"name": "docker-provider",
|
||||
"private": true,
|
||||
"version": "0.0.1",
|
||||
"description": "",
|
||||
"main": "dist/index.cjs",
|
||||
"scripts": {
|
||||
"build": "npm run build:bundle",
|
||||
"build:bundle": "esbuild src/index.ts --bundle --outfile=dist/index.mjs --platform=node --format=esm --target=esnext --banner:js=\"import { createRequire } from 'module';const require = createRequire(import.meta.url);\"",
|
||||
"build:image": "docker build -f Containerfile . -t docker-provider",
|
||||
"dev": "tsx --no-warnings=ExperimentalWarning --require dotenv/config --watch src/index.ts",
|
||||
"start": "tsx src/index.ts",
|
||||
"typecheck": "tsc --noEmit"
|
||||
},
|
||||
"keywords": [],
|
||||
"author": "",
|
||||
"license": "MIT",
|
||||
"dependencies": {
|
||||
"@trigger.dev/core": "workspace:*",
|
||||
"@trigger.dev/core-apps": "workspace:*",
|
||||
"execa": "^8.0.1",
|
||||
"socket.io-client": "^4.7.4"
|
||||
},
|
||||
"devDependencies": {
|
||||
"@types/node": "^18.19.8",
|
||||
"dotenv": "^16.4.2",
|
||||
"esbuild": "^0.19.11",
|
||||
"tsx": "^4.7.0",
|
||||
"typescript": "^5.3.3"
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,302 @@
|
||||
import { $, type ExecaChildProcess, execa } from "execa";
|
||||
import {
|
||||
SimpleLogger,
|
||||
TaskOperations,
|
||||
ProviderShell,
|
||||
TaskOperationsRestoreOptions,
|
||||
TaskOperationsCreateOptions,
|
||||
TaskOperationsIndexOptions,
|
||||
} from "@trigger.dev/core-apps";
|
||||
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 logger = new SimpleLogger(`[${MACHINE_NAME}]`);
|
||||
|
||||
type InitializeReturn = {
|
||||
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> {
|
||||
if (this.#initialized) {
|
||||
return this.#getInitializeReturn();
|
||||
}
|
||||
|
||||
logger.log("Initializing task operations");
|
||||
|
||||
if (this.opts.forceSimulate) {
|
||||
logger.log("Forced simulation enabled. Will simulate regardless of checkpoint support.");
|
||||
}
|
||||
|
||||
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();
|
||||
}
|
||||
|
||||
#getInitializeReturn(): InitializeReturn {
|
||||
return {
|
||||
canCheckpoint: this.#canCheckpoint,
|
||||
willSimulate: !this.#canCheckpoint || this.opts.forceSimulate,
|
||||
};
|
||||
}
|
||||
|
||||
async index(opts: TaskOperationsIndexOptions) {
|
||||
await this.#initialize();
|
||||
|
||||
const containerName = this.#getIndexContainerName(opts.shortCode);
|
||||
|
||||
logger.log(`Indexing task ${opts.imageRef}`, {
|
||||
host: COORDINATOR_HOST,
|
||||
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,
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
async create(opts: TaskOperationsCreateOptions) {
|
||||
await this.#initialize();
|
||||
|
||||
const containerName = this.#getRunContainerName(opts.runId);
|
||||
|
||||
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}`,
|
||||
])
|
||||
);
|
||||
} catch (error) {
|
||||
if (!isExecaChildProcess(error)) {
|
||||
throw error;
|
||||
}
|
||||
|
||||
logger.error("Create failed:", {
|
||||
opts,
|
||||
exitCode: error.exitCode,
|
||||
escapedCommand: error.escapedCommand,
|
||||
stdout: error.stdout,
|
||||
stderr: error.stderr,
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
async restore(opts: TaskOperationsRestoreOptions) {
|
||||
await this.#initialize();
|
||||
|
||||
const containerName = this.#getRunContainerName(opts.runId);
|
||||
|
||||
if (!this.#canCheckpoint || this.opts.forceSimulate) {
|
||||
logger.log("Simulating restore");
|
||||
|
||||
const unpause = logger.debug(await $`docker unpause ${containerName}`);
|
||||
|
||||
if (unpause.exitCode !== 0) {
|
||||
throw new Error("docker unpause command failed");
|
||||
}
|
||||
|
||||
await this.#sendPostStart(containerName);
|
||||
return;
|
||||
}
|
||||
|
||||
const { exitCode } = logger.debug(
|
||||
await $`docker start --checkpoint=${opts.checkpointRef} ${containerName}`
|
||||
);
|
||||
|
||||
if (exitCode !== 0) {
|
||||
throw new Error("docker start command failed");
|
||||
}
|
||||
|
||||
await this.#sendPostStart(containerName);
|
||||
}
|
||||
|
||||
async delete(opts: { runId: string }) {
|
||||
await this.#initialize();
|
||||
|
||||
const containerName = this.#getRunContainerName(opts.runId);
|
||||
await this.#sendPreStop(containerName);
|
||||
|
||||
logger.log("noop: delete");
|
||||
}
|
||||
|
||||
async get(opts: { runId: string }) {
|
||||
await this.#initialize();
|
||||
|
||||
logger.log("noop: get");
|
||||
}
|
||||
|
||||
#getIndexContainerName(suffix: string) {
|
||||
return `task-index-${suffix}`;
|
||||
}
|
||||
|
||||
#getRunContainerName(suffix: string) {
|
||||
return `task-run-${suffix}`;
|
||||
}
|
||||
|
||||
async #sendPostStart(containerName: string): Promise<void> {
|
||||
try {
|
||||
const port = await this.#getHttpServerPort(containerName);
|
||||
logger.debug(await this.#runLifecycleCommand(containerName, port, "postStart", "restore"));
|
||||
} catch (error) {
|
||||
logger.error("postStart error", { error });
|
||||
throw new Error("postStart command failed");
|
||||
}
|
||||
}
|
||||
|
||||
async #sendPreStop(containerName: string): Promise<void> {
|
||||
try {
|
||||
const port = await this.#getHttpServerPort(containerName);
|
||||
logger.debug(await this.#runLifecycleCommand(containerName, port, "preStop", "terminate"));
|
||||
} catch (error) {
|
||||
logger.error("preStop error", { error });
|
||||
throw new Error("preStop command failed");
|
||||
}
|
||||
}
|
||||
|
||||
async #getHttpServerPort(containerName: string): Promise<number> {
|
||||
// We first get the correct port, which is random during dev as we run with host networking and need to avoid clashes
|
||||
// FIXME: Skip this in prod
|
||||
const logs = logger.debug(await $`docker logs ${containerName}`);
|
||||
const matches = logs.stdout.match(/http server listening on port (?<port>[0-9]+)/);
|
||||
|
||||
const port = Number(matches?.groups?.port);
|
||||
|
||||
if (!port) {
|
||||
throw new Error("failed to extract port from logs");
|
||||
}
|
||||
|
||||
return port;
|
||||
}
|
||||
|
||||
async #runLifecycleCommand<THookType extends "postStart" | "preStop">(
|
||||
containerName: string,
|
||||
port: number,
|
||||
type: THookType,
|
||||
cause: THookType extends "postStart" ? PostStartCauses : PreStopCauses,
|
||||
retryCount = 0
|
||||
): Promise<ExecaChildProcess> {
|
||||
try {
|
||||
return await execa("docker", [
|
||||
"exec",
|
||||
containerName,
|
||||
"busybox",
|
||||
"wget",
|
||||
"-q",
|
||||
"-O-",
|
||||
`127.0.0.1:${port}/${type}?cause=${cause}`,
|
||||
]);
|
||||
} catch (error: any) {
|
||||
if (type === "postStart" && retryCount < 6) {
|
||||
logger.debug(`retriable ${type} error`, { retryCount, message: error?.message });
|
||||
await setTimeout(exponentialBackoff(retryCount + 1, 2, 50, 1150, 50));
|
||||
|
||||
return this.#runLifecycleCommand(containerName, port, type, cause, retryCount + 1);
|
||||
}
|
||||
|
||||
logger.error(`final ${type} error`, { message: error?.message });
|
||||
throw new Error(`${type} command failed after ${retryCount - 1} retries`);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
const provider = new ProviderShell({
|
||||
tasks: new DockerTaskOperations({ forceSimulate: true }),
|
||||
type: "docker",
|
||||
});
|
||||
|
||||
provider.listen();
|
||||
|
||||
function exponentialBackoff(
|
||||
retryCount: number,
|
||||
exponential: number,
|
||||
minDelay: number,
|
||||
maxDelay: number,
|
||||
jitter: number
|
||||
): number {
|
||||
// Calculate the delay using the exponential backoff formula
|
||||
const delay = Math.min(Math.pow(exponential, retryCount) * minDelay, maxDelay);
|
||||
|
||||
// Calculate the jitter
|
||||
const jitterValue = Math.random() * jitter;
|
||||
|
||||
// Return the calculated delay with jitter
|
||||
return delay + jitterValue;
|
||||
}
|
||||
@@ -0,0 +1,17 @@
|
||||
{
|
||||
"compilerOptions": {
|
||||
"target": "es2016",
|
||||
"module": "commonjs",
|
||||
"esModuleInterop": true,
|
||||
"forceConsistentCasingInFileNames": true,
|
||||
"resolveJsonModule": true,
|
||||
"strict": true,
|
||||
"skipLibCheck": true,
|
||||
"paths": {
|
||||
"@trigger.dev/core/v3": ["../../packages/core/src/v3"],
|
||||
"@trigger.dev/core/v3/*": ["../../packages/core/src/v3/*"],
|
||||
"@trigger.dev/core-apps": ["../../packages/core-apps/src"],
|
||||
"@trigger.dev/core-apps/*": ["../../packages/core-apps/src/*"]
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,9 @@
|
||||
HTTP_SERVER_PORT=8060
|
||||
|
||||
PLATFORM_WS_PORT=3030
|
||||
PLATFORM_SECRET=provider-secret
|
||||
SECURE_CONNECTION=false
|
||||
|
||||
# Use this if you are on macOS
|
||||
# COORDINATOR_HOST="host.docker.internal"
|
||||
# OTEL_EXPORTER_OTLP_ENDPOINT="http://host.docker.internal:4318"
|
||||
@@ -0,0 +1,3 @@
|
||||
dist/
|
||||
node_modules/
|
||||
.env
|
||||
@@ -0,0 +1,47 @@
|
||||
FROM node:18-alpine@sha256:ca9f6cb0466f9638e59e0c249d335a07c867cd50c429b5c7830dda1bed584649 AS node-18-alpine
|
||||
|
||||
WORKDIR /app
|
||||
|
||||
FROM node-18-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
|
||||
|
||||
RUN apk add --no-cache dumb-init
|
||||
|
||||
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 kubernetes-provider build:bundle
|
||||
|
||||
FROM base AS runner
|
||||
|
||||
RUN corepack enable
|
||||
ENV NODE_ENV production
|
||||
|
||||
COPY --from=builder --chown=node:node /app/apps/kubernetes-provider/dist/index.mjs ./index.mjs
|
||||
|
||||
EXPOSE 8000
|
||||
|
||||
USER node
|
||||
|
||||
CMD [ "/usr/bin/dumb-init", "--", "/usr/local/bin/node", "./index.mjs" ]
|
||||
@@ -0,0 +1,3 @@
|
||||
# Kubernetes provider
|
||||
|
||||
The `kubernetes-provider` allows the platform to be orchestrator-agnostic. The platform can perform actions such as `INDEX_TASKS` or `INVOKE_TASK` which the provider translates into Kubernetes actions.
|
||||
@@ -0,0 +1,31 @@
|
||||
{
|
||||
"name": "kubernetes-provider",
|
||||
"private": true,
|
||||
"version": "0.0.1",
|
||||
"description": "",
|
||||
"main": "dist/index.cjs",
|
||||
"scripts": {
|
||||
"build": "npm run build:bundle",
|
||||
"build:bundle": "esbuild src/index.ts --bundle --outfile=dist/index.mjs --platform=node --format=esm --target=esnext --banner:js=\"import { createRequire } from 'module';const require = createRequire(import.meta.url);\"",
|
||||
"build:image": "docker build -f Containerfile . -t kubernetes-provider",
|
||||
"dev": "tsx --no-warnings=ExperimentalWarning --require dotenv/config --watch src/index.ts",
|
||||
"start": "tsx src/index.ts",
|
||||
"typecheck": "tsc --noEmit"
|
||||
},
|
||||
"keywords": [],
|
||||
"author": "",
|
||||
"license": "MIT",
|
||||
"dependencies": {
|
||||
"@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"
|
||||
},
|
||||
"devDependencies": {
|
||||
"dotenv": "^16.4.2",
|
||||
"esbuild": "^0.19.11",
|
||||
"tsx": "^4.7.0",
|
||||
"typescript": "^5.3.3"
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,544 @@
|
||||
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";
|
||||
import { randomUUID } from "crypto";
|
||||
import { TaskMonitor } from "./taskMonitor";
|
||||
|
||||
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 logger = new SimpleLogger(`[${NODE_NAME}]`);
|
||||
logger.log(`running in ${RUNTIME_ENV} mode`);
|
||||
|
||||
type Namespace = {
|
||||
metadata: {
|
||||
name: string;
|
||||
};
|
||||
};
|
||||
|
||||
type ComputeResources = {
|
||||
[K in "cpu" | "memory" | "ephemeral-storage"]?: string;
|
||||
};
|
||||
|
||||
class KubernetesTaskOperations implements TaskOperations {
|
||||
#namespace: Namespace;
|
||||
#k8sApi: {
|
||||
core: k8s.CoreV1Api;
|
||||
batch: k8s.BatchV1Api;
|
||||
};
|
||||
|
||||
constructor(namespace = "default") {
|
||||
this.#namespace = {
|
||||
metadata: {
|
||||
name: namespace,
|
||||
},
|
||||
};
|
||||
|
||||
this.#k8sApi = this.#createK8sApi();
|
||||
}
|
||||
|
||||
async index(opts: TaskOperationsIndexOptions) {
|
||||
await this.#createJob(
|
||||
{
|
||||
metadata: {
|
||||
name: this.#getIndexContainerName(opts.shortCode),
|
||||
namespace: this.#namespace.metadata.name,
|
||||
},
|
||||
spec: {
|
||||
completions: 1,
|
||||
backoffLimit: 0,
|
||||
ttlSecondsAfterFinished: 300,
|
||||
template: {
|
||||
metadata: {
|
||||
labels: {
|
||||
...this.#getSharedLabels(opts),
|
||||
app: "task-index",
|
||||
"app.kubernetes.io/part-of": "trigger-worker",
|
||||
"app.kubernetes.io/component": "index",
|
||||
deployment: opts.deploymentId,
|
||||
},
|
||||
},
|
||||
spec: {
|
||||
...this.#defaultPodSpec,
|
||||
containers: [
|
||||
{
|
||||
name: this.#getIndexContainerName(opts.shortCode),
|
||||
image: opts.imageRef,
|
||||
ports: [
|
||||
{
|
||||
containerPort: 8000,
|
||||
},
|
||||
],
|
||||
resources: {
|
||||
limits: {
|
||||
cpu: "250m",
|
||||
memory: "0.5G",
|
||||
"ephemeral-storage": "2Gi",
|
||||
},
|
||||
},
|
||||
lifecycle: {
|
||||
preStop: {
|
||||
exec: {
|
||||
command: this.#getLifecycleCommand("preStop", "terminate"),
|
||||
},
|
||||
},
|
||||
},
|
||||
env: [
|
||||
...this.#getSharedEnv(opts.envId),
|
||||
{
|
||||
name: "INDEX_TASKS",
|
||||
value: "true",
|
||||
},
|
||||
{
|
||||
name: "TRIGGER_SECRET_KEY",
|
||||
value: opts.apiKey,
|
||||
},
|
||||
{
|
||||
name: "TRIGGER_API_URL",
|
||||
value: opts.apiUrl,
|
||||
},
|
||||
],
|
||||
},
|
||||
],
|
||||
},
|
||||
},
|
||||
},
|
||||
},
|
||||
this.#namespace
|
||||
);
|
||||
}
|
||||
|
||||
async create(opts: TaskOperationsCreateOptions) {
|
||||
await this.#createPod(
|
||||
{
|
||||
metadata: {
|
||||
name: this.#getRunContainerName(opts.runId),
|
||||
namespace: this.#namespace.metadata.name,
|
||||
labels: {
|
||||
...this.#getSharedLabels(opts),
|
||||
app: "task-run",
|
||||
"app.kubernetes.io/part-of": "trigger-worker",
|
||||
"app.kubernetes.io/component": "create",
|
||||
run: opts.runId,
|
||||
},
|
||||
},
|
||||
spec: {
|
||||
...this.#defaultPodSpec,
|
||||
terminationGracePeriodSeconds: 60 * 60,
|
||||
containers: [
|
||||
{
|
||||
name: this.#getRunContainerName(opts.runId),
|
||||
image: opts.image,
|
||||
ports: [
|
||||
{
|
||||
containerPort: 8000,
|
||||
},
|
||||
],
|
||||
resources: {
|
||||
requests: {
|
||||
...this.#defaultResourceRequests,
|
||||
},
|
||||
limits: {
|
||||
...this.#defaultResourceLimits,
|
||||
...this.#getResourcesFromMachineConfig(opts.machine),
|
||||
},
|
||||
},
|
||||
lifecycle: {
|
||||
preStop: {
|
||||
exec: {
|
||||
command: this.#getLifecycleCommand("preStop", "terminate"),
|
||||
},
|
||||
},
|
||||
},
|
||||
env: [
|
||||
...this.#getSharedEnv(opts.envId),
|
||||
{
|
||||
name: "TRIGGER_RUN_ID",
|
||||
value: opts.runId,
|
||||
},
|
||||
],
|
||||
volumeMounts: [
|
||||
{
|
||||
name: "taskinfo",
|
||||
mountPath: "/etc/taskinfo",
|
||||
},
|
||||
],
|
||||
},
|
||||
],
|
||||
volumes: [
|
||||
{
|
||||
name: "taskinfo",
|
||||
emptyDir: {},
|
||||
},
|
||||
],
|
||||
},
|
||||
},
|
||||
this.#namespace
|
||||
);
|
||||
}
|
||||
|
||||
async restore(opts: TaskOperationsRestoreOptions) {
|
||||
await this.#createPod(
|
||||
{
|
||||
metadata: {
|
||||
name: `${this.#getRunContainerName(opts.runId)}-${randomUUID().slice(0, 8)}`,
|
||||
namespace: this.#namespace.metadata.name,
|
||||
labels: {
|
||||
...this.#getSharedLabels(opts),
|
||||
app: "task-run",
|
||||
"app.kubernetes.io/part-of": "trigger-worker",
|
||||
"app.kubernetes.io/component": "restore",
|
||||
run: opts.runId,
|
||||
checkpoint: opts.checkpointId,
|
||||
},
|
||||
},
|
||||
spec: {
|
||||
...this.#defaultPodSpec,
|
||||
initContainers: [
|
||||
{
|
||||
name: "pull-base-image",
|
||||
image: opts.imageRef,
|
||||
command: ["sleep", "0"],
|
||||
},
|
||||
{
|
||||
name: "populate-taskinfo",
|
||||
image: "busybox",
|
||||
command: ["/bin/sh", "-c"],
|
||||
args: ["printenv COORDINATOR_HOST | tee /etc/taskinfo/coordinator-host"],
|
||||
env: [
|
||||
{
|
||||
name: "COORDINATOR_HOST",
|
||||
valueFrom: {
|
||||
fieldRef: {
|
||||
fieldPath: "status.hostIP",
|
||||
},
|
||||
},
|
||||
},
|
||||
],
|
||||
volumeMounts: [
|
||||
{
|
||||
name: "taskinfo",
|
||||
mountPath: "/etc/taskinfo",
|
||||
},
|
||||
],
|
||||
},
|
||||
],
|
||||
containers: [
|
||||
{
|
||||
name: this.#getRunContainerName(opts.runId),
|
||||
image: opts.checkpointRef,
|
||||
ports: [
|
||||
{
|
||||
containerPort: 8000,
|
||||
},
|
||||
],
|
||||
resources: {
|
||||
requests: {
|
||||
...this.#defaultResourceRequests,
|
||||
},
|
||||
limits: {
|
||||
...this.#defaultResourceLimits,
|
||||
...this.#getResourcesFromMachineConfig(opts.machine),
|
||||
},
|
||||
},
|
||||
lifecycle: {
|
||||
postStart: {
|
||||
exec: {
|
||||
command: this.#getLifecycleCommand("postStart", "restore"),
|
||||
},
|
||||
},
|
||||
preStop: {
|
||||
exec: {
|
||||
command: this.#getLifecycleCommand("preStop", "terminate"),
|
||||
},
|
||||
},
|
||||
},
|
||||
volumeMounts: [
|
||||
{
|
||||
name: "taskinfo",
|
||||
mountPath: "/etc/taskinfo",
|
||||
},
|
||||
],
|
||||
},
|
||||
],
|
||||
volumes: [
|
||||
{
|
||||
name: "taskinfo",
|
||||
emptyDir: {},
|
||||
},
|
||||
],
|
||||
},
|
||||
},
|
||||
this.#namespace
|
||||
);
|
||||
}
|
||||
|
||||
async delete(opts: { runId: string }) {
|
||||
await this.#deletePod({
|
||||
runId: opts.runId,
|
||||
namespace: this.#namespace,
|
||||
});
|
||||
}
|
||||
|
||||
async get(opts: { runId: string }) {
|
||||
await this.#getPod(opts.runId, this.#namespace);
|
||||
}
|
||||
|
||||
#envTypeToLabelValue(type: EnvironmentType) {
|
||||
switch (type) {
|
||||
case "PRODUCTION":
|
||||
return "prod";
|
||||
case "STAGING":
|
||||
return "stg";
|
||||
case "DEVELOPMENT":
|
||||
return "dev";
|
||||
case "PREVIEW":
|
||||
return "preview";
|
||||
}
|
||||
}
|
||||
|
||||
get #defaultPodSpec(): Omit<k8s.V1PodSpec, "containers"> {
|
||||
return {
|
||||
restartPolicy: "Never",
|
||||
automountServiceAccountToken: false,
|
||||
imagePullSecrets: [
|
||||
{
|
||||
name: "registry-trigger",
|
||||
},
|
||||
],
|
||||
nodeSelector: {
|
||||
nodetype: "worker",
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
get #defaultResourceRequests(): ComputeResources {
|
||||
return {
|
||||
"ephemeral-storage": "2Gi",
|
||||
};
|
||||
}
|
||||
|
||||
get #defaultResourceLimits(): ComputeResources {
|
||||
return {
|
||||
"ephemeral-storage": "10Gi",
|
||||
};
|
||||
}
|
||||
|
||||
#getSharedEnv(envId: string): k8s.V1EnvVar[] {
|
||||
return [
|
||||
{
|
||||
name: "TRIGGER_ENV_ID",
|
||||
value: envId,
|
||||
},
|
||||
{
|
||||
name: "DEBUG",
|
||||
value: process.env.DEBUG ? "1" : "0",
|
||||
},
|
||||
{
|
||||
name: "HTTP_SERVER_PORT",
|
||||
value: "8000",
|
||||
},
|
||||
{
|
||||
name: "OTEL_EXPORTER_OTLP_ENDPOINT",
|
||||
value: OTEL_EXPORTER_OTLP_ENDPOINT,
|
||||
},
|
||||
{
|
||||
name: "POD_NAME",
|
||||
valueFrom: {
|
||||
fieldRef: {
|
||||
fieldPath: "metadata.name",
|
||||
},
|
||||
},
|
||||
},
|
||||
{
|
||||
name: "COORDINATOR_HOST",
|
||||
valueFrom: {
|
||||
fieldRef: {
|
||||
fieldPath: "status.hostIP",
|
||||
},
|
||||
},
|
||||
},
|
||||
{
|
||||
name: "MACHINE_NAME",
|
||||
valueFrom: {
|
||||
fieldRef: {
|
||||
fieldPath: "spec.nodeName",
|
||||
},
|
||||
},
|
||||
},
|
||||
];
|
||||
}
|
||||
|
||||
#getSharedLabels(
|
||||
opts: TaskOperationsIndexOptions | TaskOperationsCreateOptions | TaskOperationsRestoreOptions
|
||||
): Record<string, string> {
|
||||
return {
|
||||
env: opts.envId,
|
||||
envtype: this.#envTypeToLabelValue(opts.envType),
|
||||
org: opts.orgId,
|
||||
project: opts.projectId,
|
||||
};
|
||||
}
|
||||
|
||||
#getResourcesFromMachineConfig(config: Machine): ComputeResources {
|
||||
return {
|
||||
cpu: `${config.cpu}`,
|
||||
memory: `${config.memory}G`,
|
||||
};
|
||||
}
|
||||
|
||||
#getLifecycleCommand<THookType extends "postStart" | "preStop">(
|
||||
type: THookType,
|
||||
cause: THookType extends "postStart" ? PostStartCauses : PreStopCauses
|
||||
) {
|
||||
const retries = 5;
|
||||
|
||||
// This will retry sending the lifecycle hook up to `retries` times
|
||||
// The sleep is required as this may start running before the HTTP server is up
|
||||
const exec = [
|
||||
"/bin/sh",
|
||||
"-c",
|
||||
`for i in $(seq ${retries}); do sleep 1; busybox wget -q -O- 127.0.0.1:8000/${type}?cause=${cause} && break; done`,
|
||||
];
|
||||
|
||||
logger.debug("getLifecycleCommand()", { exec });
|
||||
|
||||
return exec;
|
||||
}
|
||||
|
||||
#getIndexContainerName(suffix: string) {
|
||||
return `task-index-${suffix}`;
|
||||
}
|
||||
|
||||
#getRunContainerName(suffix: string) {
|
||||
return `task-run-${suffix}`;
|
||||
}
|
||||
|
||||
#createK8sApi() {
|
||||
const kubeConfig = new k8s.KubeConfig();
|
||||
|
||||
if (RUNTIME_ENV === "local") {
|
||||
kubeConfig.loadFromDefault();
|
||||
} else if (RUNTIME_ENV === "kubernetes") {
|
||||
kubeConfig.loadFromCluster();
|
||||
} else {
|
||||
throw new Error(`Unsupported runtime environment: ${RUNTIME_ENV}`);
|
||||
}
|
||||
|
||||
return {
|
||||
core: kubeConfig.makeApiClient(k8s.CoreV1Api),
|
||||
batch: kubeConfig.makeApiClient(k8s.BatchV1Api),
|
||||
};
|
||||
}
|
||||
|
||||
async #createPod(pod: k8s.V1Pod, namespace: Namespace) {
|
||||
try {
|
||||
const res = await this.#k8sApi.core.createNamespacedPod(namespace.metadata.name, pod);
|
||||
logger.debug(res.body);
|
||||
} catch (err: unknown) {
|
||||
this.#handleK8sError(err);
|
||||
}
|
||||
}
|
||||
|
||||
async #deletePod(opts: { runId: string; namespace: Namespace }) {
|
||||
try {
|
||||
const res = await this.#k8sApi.core.deleteNamespacedPod(
|
||||
opts.runId,
|
||||
opts.namespace.metadata.name
|
||||
);
|
||||
logger.debug(res.body);
|
||||
} catch (err: unknown) {
|
||||
this.#handleK8sError(err);
|
||||
}
|
||||
}
|
||||
|
||||
async #getPod(runId: string, namespace: Namespace) {
|
||||
try {
|
||||
const res = await this.#k8sApi.core.readNamespacedPod(runId, namespace.metadata.name);
|
||||
logger.debug(res.body);
|
||||
return res.body;
|
||||
} catch (err: unknown) {
|
||||
this.#handleK8sError(err);
|
||||
}
|
||||
}
|
||||
|
||||
async #createJob(job: k8s.V1Job, namespace: Namespace) {
|
||||
try {
|
||||
const res = await this.#k8sApi.batch.createNamespacedJob(namespace.metadata.name, job);
|
||||
logger.debug(res.body);
|
||||
} catch (err: unknown) {
|
||||
this.#handleK8sError(err);
|
||||
}
|
||||
}
|
||||
|
||||
#throwUnlessRecord(candidate: unknown): asserts candidate is Record<string, unknown> {
|
||||
if (typeof candidate !== "object" || candidate === null) {
|
||||
throw candidate;
|
||||
}
|
||||
}
|
||||
|
||||
#handleK8sError(err: unknown) {
|
||||
this.#throwUnlessRecord(err);
|
||||
|
||||
if ("body" in err && err.body) {
|
||||
logger.error(err.body);
|
||||
this.#throwUnlessRecord(err.body);
|
||||
|
||||
if (typeof err.body.message === "string") {
|
||||
throw new Error(err.body?.message);
|
||||
} else {
|
||||
throw err.body;
|
||||
}
|
||||
} else {
|
||||
logger.error(err);
|
||||
throw err;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
const provider = new ProviderShell({
|
||||
tasks: new KubernetesTaskOperations(),
|
||||
type: "kubernetes",
|
||||
});
|
||||
|
||||
provider.listen();
|
||||
|
||||
const taskMonitor = new TaskMonitor({
|
||||
runtimeEnv: RUNTIME_ENV,
|
||||
onIndexFailure: async (deploymentId, failureInfo) => {
|
||||
logger.log("Indexing failed", { deploymentId, failureInfo });
|
||||
|
||||
try {
|
||||
provider.platformSocket.send("INDEXING_FAILED", {
|
||||
deploymentId,
|
||||
error: {
|
||||
name: `Crashed with exit code ${failureInfo.exitCode}`,
|
||||
message: failureInfo.reason,
|
||||
stack: failureInfo.logs,
|
||||
},
|
||||
});
|
||||
} catch (error) {
|
||||
logger.error(error);
|
||||
}
|
||||
},
|
||||
onRunFailure: async (runId, failureInfo) => {
|
||||
logger.log("Run failed:", { runId, failureInfo });
|
||||
|
||||
try {
|
||||
provider.platformSocket.send("WORKER_CRASHED", { runId, ...failureInfo });
|
||||
} catch (error) {
|
||||
logger.error(error);
|
||||
}
|
||||
},
|
||||
});
|
||||
|
||||
taskMonitor.start();
|
||||
@@ -0,0 +1,442 @@
|
||||
import * as k8s from "@kubernetes/client-node";
|
||||
import { SimpleLogger } from "@trigger.dev/core-apps";
|
||||
import { setTimeout } from "timers/promises";
|
||||
import PQueue from "p-queue";
|
||||
|
||||
type IndexFailureHandler = (
|
||||
deploymentId: string,
|
||||
failureInfo: {
|
||||
exitCode: number;
|
||||
reason: string;
|
||||
logs: string;
|
||||
}
|
||||
) => Promise<any>;
|
||||
|
||||
type RunFailureHandler = (
|
||||
runId: string,
|
||||
failureInfo: {
|
||||
exitCode: number;
|
||||
reason: string;
|
||||
logs: string;
|
||||
}
|
||||
) => Promise<any>;
|
||||
|
||||
type TaskMonitorOptions = {
|
||||
runtimeEnv: "local" | "kubernetes";
|
||||
onIndexFailure?: IndexFailureHandler;
|
||||
onRunFailure?: RunFailureHandler;
|
||||
namespace?: string;
|
||||
};
|
||||
|
||||
export class TaskMonitor {
|
||||
#enabled = false;
|
||||
#logger = new SimpleLogger("[TaskMonitor]");
|
||||
#taskInformer: ReturnType<typeof k8s.makeInformer<k8s.V1Pod>>;
|
||||
#processedPods = new Map<string, number>();
|
||||
#queue = new PQueue({ concurrency: 10 });
|
||||
#k8sClient: {
|
||||
core: k8s.CoreV1Api;
|
||||
kubeConfig: k8s.KubeConfig;
|
||||
};
|
||||
|
||||
private namespace = "default";
|
||||
private fieldSelector = "status.phase=Failed";
|
||||
private labelSelector = "app in (task-index, task-run)";
|
||||
|
||||
constructor(private opts: TaskMonitorOptions) {
|
||||
this.#k8sClient = this.#createK8sClient();
|
||||
|
||||
this.#taskInformer = this.#createTaskInformer();
|
||||
this.#taskInformer.on("connect", this.#onInformerConnected.bind(this));
|
||||
this.#taskInformer.on("error", this.#onInformerError.bind(this));
|
||||
this.#taskInformer.on("update", this.#enqueueOnPodUpdated.bind(this));
|
||||
}
|
||||
|
||||
#createTaskInformer() {
|
||||
const listTasks = () =>
|
||||
this.#k8sClient.core.listNamespacedPod(
|
||||
this.namespace,
|
||||
undefined,
|
||||
undefined,
|
||||
undefined,
|
||||
this.fieldSelector,
|
||||
this.labelSelector
|
||||
);
|
||||
|
||||
// Uses watch with local caching
|
||||
// https://kubernetes.io/docs/reference/using-api/api-concepts/#efficient-detection-of-changes
|
||||
const informer = k8s.makeInformer(
|
||||
this.#k8sClient.kubeConfig,
|
||||
`/api/v1/namespaces/${this.namespace}/pods`,
|
||||
listTasks,
|
||||
this.labelSelector,
|
||||
this.fieldSelector
|
||||
);
|
||||
|
||||
return informer;
|
||||
}
|
||||
|
||||
async #onInformerConnected() {
|
||||
this.#logger.log("Connected");
|
||||
}
|
||||
|
||||
async #onInformerError(error: any) {
|
||||
this.#logger.error("Error:", error);
|
||||
|
||||
// Automatic reconnect
|
||||
await setTimeout(2_000);
|
||||
this.#taskInformer.start();
|
||||
}
|
||||
|
||||
#enqueueOnPodUpdated(pod: k8s.V1Pod) {
|
||||
this.#queue.add(async () => {
|
||||
try {
|
||||
// It would be better to only pass the cache key, but the pod may already be removed from the cache by the time we process it
|
||||
await this.#onPodUpdated(pod);
|
||||
} catch (error) {
|
||||
this.#logger.error("Caught onPodUpdated() error:", error);
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
async #onPodUpdated(pod: k8s.V1Pod) {
|
||||
this.#logger.debug(`Updated: ${pod.metadata?.name}`);
|
||||
this.#logger.debug("Updated", JSON.stringify(pod, null, 2));
|
||||
|
||||
// We only care about failures
|
||||
if (pod.status?.phase !== "Failed") {
|
||||
return;
|
||||
}
|
||||
|
||||
const podName = pod.metadata?.name;
|
||||
|
||||
if (!podName) {
|
||||
this.#logger.error("Pod is nameless", { pod });
|
||||
return;
|
||||
}
|
||||
|
||||
const containerStatus = pod.status.containerStatuses?.[0];
|
||||
|
||||
if (!containerStatus?.state) {
|
||||
this.#logger.error("Pod failed, but container status doesn't have state", {
|
||||
status: pod.status,
|
||||
});
|
||||
return;
|
||||
}
|
||||
|
||||
if (this.#processedPods.has(podName)) {
|
||||
this.#logger.debug("Pod update already processed", {
|
||||
podName,
|
||||
timestamp: this.#processedPods.get(podName),
|
||||
});
|
||||
return;
|
||||
}
|
||||
|
||||
this.#processedPods.set(podName, Date.now());
|
||||
|
||||
const podStatus = this.#getPodStatusSummary(pod.status);
|
||||
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) {
|
||||
return;
|
||||
}
|
||||
|
||||
const rawLogs = await this.#getLogTail(podName);
|
||||
|
||||
this.#logger.log(`${podName} failed with:`, {
|
||||
podStatus,
|
||||
containerState,
|
||||
rawLogs,
|
||||
});
|
||||
|
||||
const rawReason = podStatus.reason ?? containerState.reason ?? "";
|
||||
const message = podStatus.message ?? containerState.message ?? "";
|
||||
|
||||
let reason = rawReason || "Unknown error";
|
||||
let logs = rawLogs || "";
|
||||
|
||||
switch (rawReason) {
|
||||
case "Error":
|
||||
reason = "Unknown error.";
|
||||
break;
|
||||
case "Evicted":
|
||||
if (message.startsWith("Pod ephemeral local storage usage")) {
|
||||
reason = "Storage limit exceeded.";
|
||||
} else if (message) {
|
||||
reason = `Evicted: ${message}`;
|
||||
} else {
|
||||
reason = "Evicted for unknown reason.";
|
||||
}
|
||||
|
||||
if (logs.startsWith("failed to try resolving symlinks")) {
|
||||
logs = "";
|
||||
}
|
||||
break;
|
||||
case "OOMKilled":
|
||||
reason = "Out of memory! Try increasing the memory on this task.";
|
||||
break;
|
||||
default:
|
||||
break;
|
||||
}
|
||||
|
||||
const failureInfo = {
|
||||
exitCode,
|
||||
reason,
|
||||
logs,
|
||||
};
|
||||
|
||||
const app = pod.metadata?.labels?.app;
|
||||
|
||||
switch (app) {
|
||||
case "task-index":
|
||||
const deploymentId = pod.metadata?.labels?.deployment;
|
||||
|
||||
if (!deploymentId) {
|
||||
this.#logger.error("Index is missing ID", { pod });
|
||||
return;
|
||||
}
|
||||
|
||||
if (this.opts.onIndexFailure) {
|
||||
await this.opts.onIndexFailure(deploymentId, failureInfo);
|
||||
}
|
||||
break;
|
||||
case "task-run":
|
||||
const runId = pod.metadata?.labels?.run;
|
||||
|
||||
if (!runId) {
|
||||
this.#logger.error("Run is missing ID", { pod });
|
||||
return;
|
||||
}
|
||||
|
||||
if (this.opts.onRunFailure) {
|
||||
await this.opts.onRunFailure(runId, failureInfo);
|
||||
}
|
||||
break;
|
||||
default:
|
||||
this.#logger.error("Pod has invalid app label", { pod });
|
||||
return;
|
||||
}
|
||||
|
||||
await this.#deletePod(podName);
|
||||
}
|
||||
|
||||
async #getLogTail(podName: string) {
|
||||
try {
|
||||
const logs = await this.#k8sClient.core.readNamespacedPodLog(
|
||||
podName,
|
||||
this.namespace,
|
||||
undefined,
|
||||
undefined,
|
||||
undefined,
|
||||
1024, // limitBytes
|
||||
undefined,
|
||||
undefined,
|
||||
undefined,
|
||||
20 // tailLines
|
||||
);
|
||||
|
||||
const responseBody = logs.body ?? "";
|
||||
|
||||
if (responseBody.startsWith("unable to retrieve container logs")) {
|
||||
return "";
|
||||
}
|
||||
|
||||
// Type is wrong, body may be undefined
|
||||
return responseBody;
|
||||
} catch (error) {
|
||||
this.#logger.error("Log tail error:", error instanceof Error ? error.message : "unknown");
|
||||
return "";
|
||||
}
|
||||
}
|
||||
|
||||
#getPodStatusSummary(status: k8s.V1PodStatus) {
|
||||
return {
|
||||
reason: status.reason,
|
||||
message: status.message,
|
||||
};
|
||||
}
|
||||
|
||||
#getContainerStateSummary(state: k8s.V1ContainerState) {
|
||||
return {
|
||||
reason: state.terminated?.reason,
|
||||
exitCode: state.terminated?.exitCode,
|
||||
message: state.terminated?.message,
|
||||
};
|
||||
}
|
||||
|
||||
#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 });
|
||||
}
|
||||
|
||||
#printStats(includeMoreDetails = false) {
|
||||
this.#logger.log("Stats:", {
|
||||
cacheSize: this.#taskInformer.list().length,
|
||||
totalProcessed: this.#processedPods.size,
|
||||
...(includeMoreDetails && {
|
||||
processedPods: this.#processedPods,
|
||||
}),
|
||||
});
|
||||
}
|
||||
|
||||
async #deletePod(name: string) {
|
||||
this.#logger.debug("Deleting pod:", name);
|
||||
|
||||
await this.#k8sClient.core
|
||||
.deleteNamespacedPod(name, this.namespace)
|
||||
.catch(this.#handleK8sError.bind(this));
|
||||
}
|
||||
|
||||
async start() {
|
||||
this.#enabled = true;
|
||||
|
||||
const interval = setInterval(() => {
|
||||
if (!this.#enabled) {
|
||||
clearInterval(interval);
|
||||
return;
|
||||
}
|
||||
|
||||
this.#printStats();
|
||||
}, 300_000);
|
||||
|
||||
await this.#taskInformer.start();
|
||||
|
||||
// this.#launchTests();
|
||||
}
|
||||
|
||||
async stop() {
|
||||
if (!this.#enabled) {
|
||||
return;
|
||||
}
|
||||
|
||||
this.#enabled = false;
|
||||
this.#logger.log("Shutting down..");
|
||||
|
||||
await this.#taskInformer.stop();
|
||||
|
||||
this.#printStats(true);
|
||||
}
|
||||
|
||||
async #launchTests() {
|
||||
const createPod = async (
|
||||
container: k8s.V1Container,
|
||||
name: string,
|
||||
labels?: Record<string, string>
|
||||
) => {
|
||||
this.#logger.log("Creating pod:", name);
|
||||
|
||||
const pod = {
|
||||
metadata: {
|
||||
name,
|
||||
labels,
|
||||
},
|
||||
spec: {
|
||||
restartPolicy: "Never",
|
||||
automountServiceAccountToken: false,
|
||||
terminationGracePeriodSeconds: 1,
|
||||
containers: [container],
|
||||
},
|
||||
} satisfies k8s.V1Pod;
|
||||
|
||||
await this.#k8sClient.core
|
||||
.createNamespacedPod(this.namespace, pod)
|
||||
.catch(this.#handleK8sError.bind(this));
|
||||
};
|
||||
|
||||
const createOomPod = async (name: string, labels?: Record<string, string>) => {
|
||||
const container = {
|
||||
name,
|
||||
image: "polinux/stress",
|
||||
resources: {
|
||||
limits: {
|
||||
memory: "100Mi",
|
||||
},
|
||||
},
|
||||
command: ["stress"],
|
||||
args: ["--vm", "1", "--vm-bytes", "150M", "--vm-hang", "1"],
|
||||
} satisfies k8s.V1Container;
|
||||
|
||||
await createPod(container, name, labels);
|
||||
};
|
||||
|
||||
const createNonZeroExitPod = async (name: string, labels?: Record<string, string>) => {
|
||||
const container = {
|
||||
name,
|
||||
image: "busybox",
|
||||
command: ["sh"],
|
||||
args: ["-c", "exit 1"],
|
||||
} satisfies k8s.V1Container;
|
||||
|
||||
await createPod(container, name, labels);
|
||||
};
|
||||
|
||||
const createOoDiskPod = async (name: string, labels?: Record<string, string>) => {
|
||||
const container = {
|
||||
name,
|
||||
image: "busybox",
|
||||
command: ["sh"],
|
||||
args: [
|
||||
"-c",
|
||||
"echo creating huge-file..; head -c 1000m /dev/zero > huge-file; ls -lh huge-file; sleep infinity",
|
||||
],
|
||||
resources: {
|
||||
limits: {
|
||||
"ephemeral-storage": "500Mi",
|
||||
},
|
||||
},
|
||||
} satisfies k8s.V1Container;
|
||||
|
||||
await createPod(container, name, labels);
|
||||
};
|
||||
|
||||
await createNonZeroExitPod("non-zero-exit-task", { app: "task-run", run: "123" });
|
||||
await createOomPod("oom-task", { app: "task-index", deployment: "456" });
|
||||
await createOoDiskPod("ood-task", { app: "task-run", run: "abc" });
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,17 @@
|
||||
{
|
||||
"compilerOptions": {
|
||||
"target": "es2016",
|
||||
"module": "commonjs",
|
||||
"esModuleInterop": true,
|
||||
"forceConsistentCasingInFileNames": true,
|
||||
"resolveJsonModule": true,
|
||||
"strict": true,
|
||||
"skipLibCheck": true,
|
||||
"paths": {
|
||||
"@trigger.dev/core/v3": ["../../packages/core/src/v3"],
|
||||
"@trigger.dev/core/v3/*": ["../../packages/core/src/v3/*"],
|
||||
"@trigger.dev/core-apps": ["../../packages/core-apps/src"],
|
||||
"@trigger.dev/core-apps/*": ["../../packages/core-apps/src/*"]
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,49 +0,0 @@
|
||||
import type { StorybookConfig } from "@storybook/react-webpack5";
|
||||
import path from "path";
|
||||
|
||||
const root = path.resolve(__dirname, "../app");
|
||||
|
||||
const config: StorybookConfig = {
|
||||
webpackFinal: async (config) => {
|
||||
return {
|
||||
...config,
|
||||
resolve: {
|
||||
...config.resolve,
|
||||
alias: {
|
||||
...(config.resolve?.alias ?? {}),
|
||||
"~": root,
|
||||
},
|
||||
extensions: [
|
||||
...(config.resolve?.extensions ?? []),
|
||||
...[".ts", ".tsx", ".js", ".jsx", ".mdx"],
|
||||
],
|
||||
},
|
||||
};
|
||||
},
|
||||
stories: ["../app/**/stories/*.mdx", "../app/**/stories/*.stories.@(js|jsx|ts|tsx)"],
|
||||
addons: [
|
||||
"@storybook/addon-links",
|
||||
"@storybook/addon-essentials",
|
||||
"@storybook/addon-interactions",
|
||||
"storybook-addon-variants",
|
||||
"storybook-addon-designs",
|
||||
"@storybook/addon-docs",
|
||||
{
|
||||
name: "@storybook/addon-styling",
|
||||
options: {
|
||||
// Check out https://github.com/storybookjs/addon-styling/blob/main/docs/api.md
|
||||
// For more details on this addon's options.
|
||||
postCss: true,
|
||||
},
|
||||
},
|
||||
],
|
||||
framework: {
|
||||
name: "@storybook/react-webpack5",
|
||||
options: {},
|
||||
},
|
||||
docs: {
|
||||
autodocs: "tag",
|
||||
},
|
||||
staticDirs: [path.resolve("public")],
|
||||
};
|
||||
export default config;
|
||||
@@ -1,50 +0,0 @@
|
||||
import type { Preview } from "@storybook/react";
|
||||
import "../app/tailwind.css";
|
||||
import { createRemixStub } from "@remix-run/testing";
|
||||
import React from "react";
|
||||
import { LocaleContextProvider } from "../app/components/primitives/LocaleProvider";
|
||||
import { OperatingSystemContextProvider } from "../app/components/primitives/OperatingSystemProvider";
|
||||
|
||||
const preview: Preview = {
|
||||
parameters: {
|
||||
actions: { argTypesRegex: "^on[A-Z].*" },
|
||||
controls: {
|
||||
matchers: {
|
||||
color: /(background|color)$/i,
|
||||
date: /Date$/,
|
||||
},
|
||||
},
|
||||
backgrounds: {
|
||||
default: "App background",
|
||||
values: [
|
||||
{
|
||||
name: "App background",
|
||||
value: "#0B1018",
|
||||
},
|
||||
],
|
||||
},
|
||||
layout: "fullscreen",
|
||||
},
|
||||
decorators: [
|
||||
(Story) => {
|
||||
const RemixStub = createRemixStub([
|
||||
{
|
||||
path: "/*",
|
||||
action: () => ({ redirect: "/" }),
|
||||
loader: () => ({ redirect: "/" }),
|
||||
Component: Story,
|
||||
},
|
||||
]);
|
||||
|
||||
return (
|
||||
<OperatingSystemContextProvider platform="mac">
|
||||
<LocaleContextProvider locales={window.navigator.languages as string[]}>
|
||||
<RemixStub initialEntries={["/"]} />
|
||||
</LocaleContextProvider>
|
||||
</OperatingSystemContextProvider>
|
||||
);
|
||||
},
|
||||
],
|
||||
};
|
||||
|
||||
export default preview;
|
||||
@@ -0,0 +1,53 @@
|
||||
export function AISparkleIcon({ className }: { className?: string }) {
|
||||
return (
|
||||
<svg className={className} viewBox="0 0 18 18" fill="none" xmlns="http://www.w3.org/2000/svg">
|
||||
<path
|
||||
d="M14.9806 0.803884C14.8871 0.33646 14.4767 0 14 0C13.5233 0 13.1129 0.33646 13.0194 0.803884L12.7809 1.99644C12.7017 2.3923 12.3923 2.70174 11.9964 2.78091L10.8039 3.01942C10.3365 3.1129 10 3.52332 10 4C10 4.47668 10.3365 4.8871 10.8039 4.98058L11.9964 5.21909C12.3923 5.29826 12.7017 5.6077 12.7809 6.00356L13.0194 7.19612C13.1129 7.66354 13.5233 8 14 8C14.4767 8 14.8871 7.66354 14.9806 7.19612L15.2191 6.00356C15.2983 5.6077 15.6077 5.29826 16.0036 5.21909L17.1961 4.98058C17.6635 4.8871 18 4.47668 18 4C18 3.52332 17.6635 3.1129 17.1961 3.01942L16.0036 2.78091C15.6077 2.70174 15.2983 2.3923 15.2191 1.99644L14.9806 0.803884Z"
|
||||
fill="url(#paint0_linear_11402_36656)"
|
||||
/>
|
||||
<path
|
||||
d="M5.94868 4.68377C5.81257 4.27543 5.43043 4 5 4C4.56957 4 4.18743 4.27543 4.05132 4.68377L3.36754 6.73509C3.26801 7.03369 3.03369 7.26801 2.73509 7.36754L0.683772 8.05132C0.27543 8.18743 0 8.56957 0 9C0 9.43043 0.27543 9.81257 0.683772 9.94868L2.73509 10.6325C3.03369 10.732 3.26801 10.9663 3.36754 11.2649L4.05132 13.3162C4.18743 13.7246 4.56957 14 5 14C5.43043 14 5.81257 13.7246 5.94868 13.3162L6.63246 11.2649C6.73199 10.9663 6.96631 10.732 7.26491 10.6325L9.31623 9.94868C9.72457 9.81257 10 9.43043 10 9C10 8.56957 9.72457 8.18743 9.31623 8.05132L7.26491 7.36754C6.96631 7.26801 6.73199 7.03369 6.63246 6.73509L5.94868 4.68377Z"
|
||||
fill="url(#paint1_linear_11402_36656)"
|
||||
/>
|
||||
<path
|
||||
d="M12.9487 12.6838C12.8126 12.2754 12.4304 12 12 12C11.5696 12 11.1874 12.2754 11.0513 12.6838L10.8675 13.2351C10.768 13.5337 10.5337 13.768 10.2351 13.8675L9.68377 14.0513C9.27543 14.1874 9 14.5696 9 15C9 15.4304 9.27543 15.8126 9.68377 15.9487L10.2351 16.1325C10.5337 16.232 10.768 16.4663 10.8675 16.7649L11.0513 17.3162C11.1874 17.7246 11.5696 18 12 18C12.4304 18 12.8126 17.7246 12.9487 17.3162L13.1325 16.7649C13.232 16.4663 13.4663 16.232 13.7649 16.1325L14.3162 15.9487C14.7246 15.8126 15 15.4304 15 15C15 14.5696 14.7246 14.1874 14.3162 14.0513L13.7649 13.8675C13.4663 13.768 13.232 13.5337 13.1325 13.2351L12.9487 12.6838Z"
|
||||
fill="url(#paint2_linear_11402_36656)"
|
||||
/>
|
||||
<defs>
|
||||
<linearGradient
|
||||
id="paint0_linear_11402_36656"
|
||||
x1="9"
|
||||
y1="0"
|
||||
x2="9"
|
||||
y2="18"
|
||||
gradientUnits="userSpaceOnUse"
|
||||
>
|
||||
<stop stopColor="#E543FF" />
|
||||
<stop offset="1" stopColor="#286399" />
|
||||
</linearGradient>
|
||||
<linearGradient
|
||||
id="paint1_linear_11402_36656"
|
||||
x1="9"
|
||||
y1="0"
|
||||
x2="9"
|
||||
y2="18"
|
||||
gradientUnits="userSpaceOnUse"
|
||||
>
|
||||
<stop stopColor="#E543FF" />
|
||||
<stop offset="1" stopColor="#286399" />
|
||||
</linearGradient>
|
||||
<linearGradient
|
||||
id="paint2_linear_11402_36656"
|
||||
x1="9"
|
||||
y1="0"
|
||||
x2="9"
|
||||
y2="18"
|
||||
gradientUnits="userSpaceOnUse"
|
||||
>
|
||||
<stop stopColor="#E543FF" />
|
||||
<stop offset="1" stopColor="#286399" />
|
||||
</linearGradient>
|
||||
</defs>
|
||||
</svg>
|
||||
);
|
||||
}
|
||||
@@ -0,0 +1,29 @@
|
||||
export function AttemptIcon({ className }: { className?: string }) {
|
||||
return (
|
||||
<svg
|
||||
className={className}
|
||||
width="16"
|
||||
height="16"
|
||||
viewBox="0 0 16 16"
|
||||
fill="none"
|
||||
xmlns="http://www.w3.org/2000/svg"
|
||||
>
|
||||
<g clipPath="url(#clip0_9964_113464)">
|
||||
<path
|
||||
fillRule="evenodd"
|
||||
clipRule="evenodd"
|
||||
d="M16 0H0V16H16V0ZM7.09906 4.4L4.53906 11.5H6.11906L6.63906 10H9.35906L9.87906 11.5H11.4591L8.89906 4.4H7.09906ZM7.99906 6L8.92906 8.73H7.06906L7.99906 6Z"
|
||||
fill="currentColor"
|
||||
/>
|
||||
</g>
|
||||
<defs>
|
||||
<clipPath id="clip0_9964_113464">
|
||||
<path
|
||||
d="M0 2C0 0.895431 0.895431 0 2 0H14C15.1046 0 16 0.895431 16 2V14C16 15.1046 15.1046 16 14 16H2C0.895431 16 0 15.1046 0 14V2Z"
|
||||
fill="white"
|
||||
/>
|
||||
</clipPath>
|
||||
</defs>
|
||||
</svg>
|
||||
);
|
||||
}
|
||||
@@ -0,0 +1,14 @@
|
||||
export function ExitIcon({ className }: { className?: string }) {
|
||||
return (
|
||||
<svg className={className} viewBox="0 0 16 16" fill="none" xmlns="http://www.w3.org/2000/svg">
|
||||
<line x1="3.5" y1="8" x2="11.5" y2="8" stroke="currentColor" strokeLinecap="round" />
|
||||
<line x1="15.5" y1="1.5" x2="15.5" y2="14.5" stroke="currentColor" strokeLinecap="round" />
|
||||
<path
|
||||
d="M8.5 4.5L12 8L8.5 11.5"
|
||||
stroke="currentColor"
|
||||
strokeLinecap="round"
|
||||
strokeLinejoin="round"
|
||||
/>
|
||||
</svg>
|
||||
);
|
||||
}
|
||||
@@ -0,0 +1,25 @@
|
||||
export function ShowParentIcon({ className }: { className?: string }) {
|
||||
return (
|
||||
<svg className={className} viewBox="0 0 16 16" fill="none" xmlns="http://www.w3.org/2000/svg">
|
||||
<path
|
||||
fill="currentColor"
|
||||
fillRule="evenodd"
|
||||
clipRule="evenodd"
|
||||
d="M3.125 16H3.625H4.125V15H3.125H2C1.89733 15 1.8 14.9848 1.70919 14.9573L1.19765 14.8024L1.04267 14.2908C1.01515 14.2 1 14.1027 1 14V12.875V11.875H0V12.375V12.875V14C0 14.202 0.0299409 14.397 0.0856265 14.5808C0.135047 14.7439 0.204744 14.8982 0.291786 15.0407C0.363386 15.158 0.446722 15.2673 0.540161 15.3671C0.570065 15.399 0.601004 15.4299 0.632924 15.4598C0.732663 15.5533 0.841982 15.6366 0.95925 15.7082C1.10181 15.7953 1.25612 15.865 1.41923 15.9144C1.60303 15.9701 1.79802 16 2 16H3.125ZM10.125 15V16H9.625H9.125H6.875H6.375H5.875V15H6.875H9.125H10.125ZM15 5.875H16V6.375V6.875V9.125V9.625V10.125H15V9.125V6.875V5.875ZM5.875 1V0H6.375H6.875H9.125H9.625H10.125V1H9.125H6.875H5.875ZM1 10.125H0V9.625V9.125V6.875V6.375V5.875H1V6.875V9.125V10.125ZM0 3.625V4.125H1V3.125V2C1 1.89733 1.01515 1.8 1.04267 1.70919L1.19765 1.19765L1.70919 1.04267C1.8 1.01515 1.89733 1 2 1H3.125H4.125V0H3.625H3.125H2C1.79802 0 1.60303 0.0299409 1.41923 0.0856264C1.25612 0.135047 1.10181 0.204744 0.95925 0.291786C0.841983 0.363386 0.732663 0.446722 0.632924 0.540161C0.601004 0.570065 0.570065 0.601004 0.540161 0.632924C0.446722 0.732663 0.363386 0.841982 0.291786 0.95925C0.204744 1.10181 0.135047 1.25612 0.0856265 1.41923C0.0299409 1.60303 0 1.79802 0 2V3.125V3.625ZM12.375 0H11.875V1H12.875H14C14.1027 1 14.2 1.01515 14.2908 1.04267L14.8024 1.19765L14.9573 1.70919C14.9848 1.8 15 1.89733 15 2V3.125V4.125H16V3.625V3.125V2C16 1.79802 15.9701 1.60303 15.9144 1.41923C15.865 1.25612 15.7953 1.10181 15.7082 0.95925C15.6366 0.841982 15.5533 0.732663 15.4598 0.632924C15.4299 0.601004 15.399 0.570065 15.3671 0.540161C15.2673 0.446722 15.158 0.363386 15.0407 0.291786C14.8982 0.204744 14.7439 0.135047 14.5808 0.0856265C14.397 0.0299409 14.202 0 14 0H12.875H12.375ZM16 12.375V11.875H15V12.875V14C15 14.1027 14.9848 14.2 14.9573 14.2908L14.8024 14.8024L14.2908 14.9573C14.2 14.9848 14.1027 15 14 15H12.875H11.875V16H12.375H12.875H14C14.202 16 14.397 15.9701 14.5808 15.9144C14.7439 15.865 14.8982 15.7953 15.0407 15.7082C15.158 15.6366 15.2673 15.5533 15.3671 15.4598C15.399 15.4299 15.4299 15.399 15.4598 15.3671C15.5533 15.2673 15.6366 15.158 15.7082 15.0407C15.7953 14.8982 15.865 14.7439 15.9144 14.5808C15.9701 14.397 16 14.202 16 14V12.875V12.375ZM8.75 6.26758L10.4523 7.96991C10.7452 8.26281 11.2201 8.26281 11.513 7.96991C11.8059 7.67702 11.8059 7.20215 11.513 6.90925L8.59099 3.98725C8.43431 3.83057 8.22556 3.7577 8.02045 3.76865C7.81604 3.75834 7.60821 3.83124 7.45209 3.98736L4.53009 6.90937C4.23719 7.20226 4.23719 7.67713 4.53009 7.97003C4.82298 8.26292 5.29785 8.26292 5.59075 7.97003L7.25 6.31077L7.25 11.25C7.25 11.6642 7.58579 12 8 12C8.41421 12 8.75 11.6642 8.75 11.25L8.75 6.26758Z"
|
||||
/>
|
||||
</svg>
|
||||
);
|
||||
}
|
||||
|
||||
export function ShowParentIconSelected({ className }: { className?: string }) {
|
||||
return (
|
||||
<svg className={className} viewBox="0 0 16 16" fill="none" xmlns="http://www.w3.org/2000/svg">
|
||||
<path
|
||||
fill="currentColor"
|
||||
fillRule="evenodd"
|
||||
clipRule="evenodd"
|
||||
d="M2 0C0.895431 0 0 0.895431 0 2V14C0 15.1046 0.895431 16 2 16H14C15.1046 16 16 15.1046 16 14V2C16 0.895431 15.1046 0 14 0H2ZM8.75 6.26758L10.4523 7.96991C10.7452 8.26281 11.2201 8.26281 11.513 7.96991C11.8059 7.67702 11.8059 7.20215 11.513 6.90925L8.59099 3.98725C8.43344 3.8297 8.22324 3.7569 8.01703 3.76884C7.8138 3.75958 7.60752 3.83254 7.45234 3.98773L4.53033 6.90973C4.23744 7.20263 4.23744 7.6775 4.53033 7.97039C4.82322 8.26329 5.2981 8.26329 5.59099 7.97039L7.25 6.31138L7.25 11.25C7.25 11.6642 7.58579 12 8 12C8.41421 12 8.75 11.6642 8.75 11.25L8.75 6.26758Z"
|
||||
/>
|
||||
</svg>
|
||||
);
|
||||
}
|
||||
@@ -0,0 +1,29 @@
|
||||
export function TaskIcon({ className }: { className?: string }) {
|
||||
return (
|
||||
<svg
|
||||
className={className}
|
||||
width="16"
|
||||
height="16"
|
||||
viewBox="0 0 16 16"
|
||||
fill="none"
|
||||
xmlns="http://www.w3.org/2000/svg"
|
||||
>
|
||||
<g clipPath="url(#clip0_9221_99732)">
|
||||
<path
|
||||
fillRule="evenodd"
|
||||
clipRule="evenodd"
|
||||
d="M16 0H0V16H16V0ZM10.8901 5.73995V4.44995H5.11011V5.73995H7.23011V11.55H8.77011V5.73995H10.8901Z"
|
||||
fill="currentColor"
|
||||
/>
|
||||
</g>
|
||||
<defs>
|
||||
<clipPath id="clip0_9221_99732">
|
||||
<path
|
||||
d="M0 2C0 0.895431 0.895431 0 2 0H14C15.1046 0 16 0.895431 16 2V14C16 15.1046 15.1046 16 14 16H2C0.895431 16 0 15.1046 0 14V2Z"
|
||||
fill="white"
|
||||
/>
|
||||
</clipPath>
|
||||
</defs>
|
||||
</svg>
|
||||
);
|
||||
}
|
||||
Binary file not shown.
|
After Width: | Height: | Size: 166 B |
Binary file not shown.
|
After Width: | Height: | Size: 170 B |
@@ -0,0 +1,21 @@
|
||||
export function ATAndTLogo({ className }: { className?: string }) {
|
||||
return (
|
||||
<svg
|
||||
className={className}
|
||||
width="66"
|
||||
height="33"
|
||||
viewBox="0 0 66 33"
|
||||
fill="none"
|
||||
xmlns="http://www.w3.org/2000/svg"
|
||||
>
|
||||
<path
|
||||
d="M57.0405 21.156C57.0064 21.1581 56.9723 21.1529 56.9403 21.1408C56.9084 21.1287 56.8794 21.1099 56.8552 21.0858C56.8311 21.0616 56.8123 21.0326 56.8002 21.0007C56.7881 20.9687 56.7829 20.9346 56.785 20.9005V13.4102H54.2407C54.2068 13.4114 54.173 13.4056 54.1415 13.3932C54.1099 13.3808 54.0813 13.362 54.0573 13.338C54.0333 13.3141 54.0145 13.2854 54.0021 13.2538C53.9897 13.2223 53.984 13.1885 53.9852 13.1546V12.1226C53.9839 12.0887 53.9896 12.0549 54.002 12.0233C54.0143 11.9916 54.0331 11.9629 54.0571 11.9389C54.0811 11.9148 54.1098 11.896 54.1414 11.8836C54.173 11.8711 54.2068 11.8653 54.2407 11.8666H61.0127C61.0468 11.8646 61.0809 11.8698 61.1129 11.882C61.1448 11.8942 61.1738 11.913 61.1979 11.9372C61.222 11.9614 61.2407 11.9904 61.2527 12.0224C61.2647 12.0544 61.2698 12.0886 61.2677 12.1226V13.1557C61.2689 13.1896 61.2632 13.2233 61.2508 13.2548C61.2384 13.2863 61.2197 13.315 61.1958 13.3389C61.1718 13.3629 61.1432 13.3817 61.1118 13.3941C61.0803 13.4066 61.0465 13.4124 61.0127 13.4112H58.4684V20.9005C58.4696 20.9344 58.4639 20.9681 58.4515 20.9997C58.4391 21.0313 58.4203 21.0599 58.3963 21.0839C58.3723 21.1079 58.3437 21.1266 58.3121 21.1391C58.2805 21.1515 58.2468 21.1572 58.2129 21.156H57.0394M37.5459 17.4163L36.2102 13.5845L34.8629 17.4163H37.5459ZM40.5081 20.8528C40.5662 21.0041 40.4735 21.1555 40.3107 21.1555H39.1087C39.0327 21.1609 38.9572 21.1394 38.8955 21.0947C38.8338 21.0499 38.7899 20.9848 38.7715 20.9109L38.0861 18.9369H34.3343L33.6489 20.9109C33.6305 20.9847 33.5867 21.0498 33.5251 21.0945C33.4635 21.1392 33.3881 21.1608 33.3122 21.1555H32.1744C32.0231 21.1555 31.9183 21.0041 31.977 20.8528L35.1245 12.0985C35.1826 11.9357 35.2873 11.866 35.4612 11.866H37.0184C37.1928 11.866 37.3091 11.9351 37.3672 12.0985L40.5147 20.8528M49.5381 19.9014C50.2811 19.9014 50.7812 19.5422 51.1875 18.9265L49.3067 16.9058C48.5862 17.3127 48.1212 17.7185 48.1212 18.5316C48.1212 19.3322 48.7716 19.9025 49.5387 19.9025M50.0613 13.1031C49.4581 13.1031 49.1088 13.4869 49.1088 13.9969C49.1088 14.3917 49.3172 14.7399 49.7942 15.2509C50.6189 14.7739 50.9677 14.4843 50.9677 13.9733C50.9677 13.4962 50.6661 13.096 50.0618 13.096M55.2025 20.8303C55.3533 20.9931 55.2606 21.1555 55.0742 21.1555H53.5937C53.511 21.1627 53.4278 21.1483 53.3524 21.1136C53.277 21.0788 53.212 21.025 53.1638 20.9575L52.2864 19.9826C51.6942 20.7722 50.869 21.3644 49.4981 21.3644C47.8021 21.3644 46.4658 20.3428 46.4658 18.5898C46.4658 17.2425 47.1863 16.5225 48.2786 15.9194C47.744 15.3041 47.5 14.6538 47.5 14.0852C47.5 12.6447 48.5106 11.6582 50.0316 11.6582C51.5889 11.6582 52.5408 12.5761 52.5408 13.9338C52.5408 15.0952 51.7046 15.7444 50.8218 16.2325L52.123 17.6379L52.8545 16.3602C52.9477 16.2094 53.0519 16.1519 53.2383 16.1519H54.3646C54.5511 16.1519 54.6553 16.2802 54.5401 16.477L53.2449 18.706L55.208 20.8314M43.6161 21.1566C43.65 21.1579 43.6838 21.1521 43.7155 21.1398C43.7471 21.1274 43.7758 21.1086 43.7998 21.0846C43.8239 21.0607 43.8427 21.032 43.8552 21.0004C43.8676 20.9688 43.8734 20.935 43.8721 20.901V13.4107H46.4164C46.4503 13.4119 46.4841 13.4062 46.5157 13.3938C46.5472 13.3813 46.5759 13.3626 46.5999 13.3386C46.6238 13.3146 46.6426 13.2859 46.655 13.2544C46.6674 13.2228 46.6732 13.1891 46.672 13.1552V12.1226C46.6733 12.0887 46.6676 12.0549 46.6552 12.0233C46.6428 11.9916 46.624 11.9629 46.6001 11.9389C46.5761 11.9148 46.5474 11.896 46.5158 11.8836C46.4842 11.8711 46.4504 11.8653 46.4164 11.8666H39.6406C39.6067 11.8653 39.5728 11.8711 39.5412 11.8836C39.5097 11.896 39.481 11.9148 39.457 11.9389C39.433 11.9629 39.4142 11.9916 39.4019 12.0233C39.3895 12.0549 39.3838 12.0887 39.3851 12.1226V13.1557C39.3838 13.1896 39.3896 13.2234 39.402 13.2549C39.4144 13.2865 39.4332 13.3152 39.4572 13.3391C39.4812 13.3631 39.5098 13.3819 39.5414 13.3943C39.5729 13.4067 39.6067 13.4125 39.6406 13.4112H42.1849V20.9005C42.1837 20.9344 42.1894 20.9681 42.2018 20.9997C42.2143 21.0313 42.233 21.0599 42.257 21.0839C42.281 21.1079 42.3096 21.1266 42.3412 21.1391C42.3728 21.1515 42.4065 21.1572 42.4404 21.156L43.6161 21.1566Z"
|
||||
fill="currentColor"
|
||||
/>
|
||||
<path
|
||||
d="M9.2351 25.6737C11.2728 27.2567 13.7799 28.1157 16.3602 28.1149C19.296 28.1149 21.9725 27.0248 24.0151 25.2361C24.0397 25.2142 24.0277 25.1999 24.003 25.2142C23.0862 25.8261 20.4739 27.1624 16.3602 27.1624C12.7851 27.1624 10.5259 26.3646 9.2499 25.6529C9.22523 25.6408 9.217 25.6583 9.23455 25.6748M17.1487 26.2686C20.0083 26.2686 23.1503 25.4889 25.0295 23.9464C25.5438 23.5258 26.033 22.9665 26.4716 22.2148C26.74 21.7483 26.9742 21.263 27.1724 20.7628C27.1812 20.7381 27.1669 20.726 27.1477 20.754C25.4002 23.3306 20.339 24.9296 15.1188 24.9296C11.4252 24.9296 7.45135 23.7485 5.89571 21.4931C5.88035 21.4723 5.865 21.4811 5.87377 21.5052C7.3181 24.5863 11.7152 26.2686 17.1476 26.2686M14.0232 21.1581C8.07591 21.1581 5.2717 18.389 4.76338 16.4983C4.7579 16.4709 4.73926 16.4764 4.73926 16.5016C4.73926 17.1377 4.80287 17.9591 4.91253 18.5041C4.96463 18.7695 5.1867 19.1857 5.49706 19.5186C6.9392 21.0166 10.5286 23.1212 16.7468 23.1212C25.2187 23.1212 27.1554 20.2994 27.5508 19.3705C27.8337 18.7125 27.9801 17.5073 27.9801 16.5C27.9803 16.29 27.9751 16.0801 27.9648 15.8705C27.9648 15.8392 27.9467 15.8376 27.9406 15.8672C27.5173 18.1373 20.2792 21.157 14.0249 21.157M5.85458 11.5172C5.49236 12.2876 5.21385 13.0946 5.02385 13.9244C4.96901 14.1766 4.99533 14.2989 5.07868 14.4876C5.79152 15.9999 9.39686 18.4192 17.8068 18.4192C22.9376 18.4192 26.9235 17.158 27.5689 14.8582C27.6878 14.4349 27.6939 13.9875 27.5414 13.3854C27.3769 12.712 27.0512 11.9268 26.7803 11.3757C26.7716 11.3576 26.7557 11.3604 26.7584 11.3812C26.8588 14.3982 18.4456 16.3426 14.2003 16.3426C9.60194 16.3426 5.76685 14.5111 5.76685 12.1971C5.76685 11.9751 5.81291 11.7585 5.87652 11.521C5.882 11.4991 5.86445 11.4964 5.85458 11.5156M24.0348 7.81261C24.0859 7.89265 24.1114 7.9864 24.1077 8.0813C24.1077 9.37209 20.158 11.6548 13.8702 11.6548C9.25045 11.6548 8.38572 9.94127 8.38572 8.85117C8.38572 8.46733 8.53487 8.0632 8.86442 7.65798C8.88252 7.63385 8.86716 7.62508 8.84633 7.64262C8.24498 8.15186 7.6971 8.72105 7.21118 9.34138C6.97978 9.63365 6.83666 9.89246 6.83666 10.0476C6.83666 12.3068 12.4999 13.9441 17.7958 13.9441C23.4437 13.9441 25.9562 12.1017 25.9562 10.4896C25.9562 9.91111 25.7369 9.57388 25.1556 8.91861C24.7816 8.49255 24.428 8.15094 24.0589 7.80438C24.0408 7.78958 24.0282 7.80164 24.0408 7.81974M22.3108 6.52949C20.5693 5.48545 18.547 4.8916 16.3668 4.8916C14.1713 4.8916 12.0887 5.50574 10.3351 6.57775C9.81086 6.90017 9.51585 7.15899 9.51585 7.49128C9.51585 8.47117 11.8057 9.52453 15.8673 9.52453C19.8866 9.52453 23.005 8.37082 23.005 7.25988C23.005 6.99503 22.7731 6.80915 22.3042 6.52949"
|
||||
fill="currentColor"
|
||||
/>
|
||||
</svg>
|
||||
);
|
||||
}
|
||||
File diff suppressed because one or more lines are too long
Some files were not shown because too many files have changed in this diff Show More
Reference in New Issue
Block a user