Compare commits

...

74 Commits

Author SHA1 Message Date
Eric Allam 90f52de147 Fix pnpm lock file 2023-10-04 14:01:48 +01:00
github-actions[bot] 28b05a82d8 chore: Update version for release (#538)
Co-authored-by: github-actions[bot] <github-actions[bot]@users.noreply.github.com>
2023-10-04 14:00:20 +01:00
Eric Allam b9ed7e2ced Allow blank issues 2023-10-04 13:21:38 +01:00
Eric Allam 5e651d84af Update pnpm lock file
🚀 Publish Trigger.dev Docker / ʦ TypeScript (push) Has been cancelled
🚀 Publish Trigger.dev Docker / Unit Tests (push) Has been cancelled
🚀 Publish Trigger.dev Docker / e2e Tests (push) Has been cancelled
🚀 Publish Trigger.dev Docker / publish (push) Has been cancelled
2023-10-04 11:00:55 +01:00
nicktrn 6a992a1995 Replicate integration and remote callbacks (#507)
* Support tasks with remote callbacks

* Add common integration tsconfig

* Add Replicate integration

* Basic job catalog example

* Integration catalog entry

* Check for callbackUrl during executeTask

* Fix getAll

* Improve JSDoc

* Bump version

* Remove named queue

* Simplify runTask types

* Trust the types

* Fail tasks on timeout

* Callback timeout as param

* Mess with types

* performRunExecutionV1

* Update runTask docs

* Shorten callback task methods

* Fix run method return type

* Image processing jobs

* Replicate docs

* Text output example

* Changeset

* Version bump

* Roll back ugly types

* Remove missing types

* Quicker return when waiting on remote callback

* Remote callback example

* Bump version

* Remove schema parsing

* Only schedule positive callback timeout

* Decrease callback secret length

* Explicit default timeouts

* Import deployments tasks

* JSDoc

* Deployments docs

* Fix runTask examples, mention wrappers

---------

Co-authored-by: Eric Allam <eric@trigger.dev>
2023-10-04 10:46:17 +01:00
nicktrn 81e886a1ba Add typed filters to Linear getAll helper (#517)
* Fix getAll params type

* Changeset

* Search param types

* Update changeset
2023-10-03 16:48:21 +01:00
Eric Allam ab9e4a989c Improves the performance of run resuming (#522)
* Improves the perform of run resuming

When runs resume, we try and make sure that tasks that have already been completed are cached and reused. Worst case scenario the client needs to hit the API server once for a non-cached task that is indeed completed on the server, but this can get pretty expensive when there are a larger number of tasks.

This commit does 2 different things to help:

- noop tasks are no longer “cached” using the cachedTasks strategy, instead their idempotency keys are shoved into a bloom filter and the client tests for their inclusion in the bloom filter before running them (since they don’t have any concept of output, this works)
- Additional cached tasks are lazy loaded when a task is run. This allows us to progressively fetch additional tasks to be cached on the client, which will cut down on cache misses by a decent amount

* Create warm-carrots-float.md

* Make io.yield backwards compat with older platform versions

* Better support old clients connecting to server versions that support lazy loading cached tasks

* Fixed type errors when settings headers with unknown value

* Better yield not support error message

* Rename _version to _serverVersion to be more clear
2023-10-03 16:44:18 +01:00
Matt Aitken e350659e24 Updated docs README.md 2023-10-02 14:26:54 +01:00
Eric Allam f888a49555 Updated outdated lockfile 2023-09-30 22:20:34 +01:00
Eric Allam 421c249e50 Add some documentation around canceling scheduled events 2023-09-30 22:18:52 +01:00
James Ritchie a8a6f51387 autofocus the search field on the Job page 2023-09-29 18:23:22 +01:00
James Ritchie 12e73eef22 Display framework logos on the onboarding setup pages (#519) 2023-09-29 16:44:35 +01:00
github-actions[bot] 2e33fcb16b chore: Update version for release (#521)
Co-authored-by: github-actions[bot] <github-actions[bot]@users.noreply.github.com>
2023-09-29 16:42:53 +01:00
Matt Aitken 5912cdd11c Updated the docs for the React hooks 2023-09-29 15:37:53 +01:00
Matt Aitken cc016b3ae3 CLI init: adds public key as “TRIGGER_PUBLIC_API_KEY” except for Next which overrides this 2023-09-29 15:14:52 +01:00
Matt Aitken 7760e09462 Improved CLI init Next.js middleware detection 2023-09-29 11:26:29 +01:00
James Ritchie 3ca4456c88 Youtube embedded video fits its aspect ratio instead of going full width
🚀 Publish Trigger.dev Docker / ʦ TypeScript (push) Has been cancelled
🚀 Publish Trigger.dev Docker / Unit Tests (push) Has been cancelled
🚀 Publish Trigger.dev Docker / e2e Tests (push) Has been cancelled
🚀 Publish Trigger.dev Docker / publish (push) Has been cancelled
2023-09-28 17:37:07 +01:00
James Ritchie 618b7f22da Swapped out the Homepage link in the side menu for a link to the Changelog 2023-09-28 17:27:21 +01:00
Eric Allam a12c7c3b0a Fixed sparodically failed run creations
- Use a better way of getting the latest job run number to increment
- Make the CreateRunService transaction more reliable
- Invoke dispatchers in parallel
- No longer swallow prisma errors in $transaction
2023-09-28 17:08:26 +01:00
Eric Allam a42e94c75f Add the STAGING environment by default
🚀 Publish Trigger.dev Docker / ʦ TypeScript (push) Has been cancelled
🚀 Publish Trigger.dev Docker / Unit Tests (push) Has been cancelled
🚀 Publish Trigger.dev Docker / e2e Tests (push) Has been cancelled
🚀 Publish Trigger.dev Docker / publish (push) Has been cancelled
2023-09-27 18:17:19 +01:00
Matt Aitken bc757c8ddb Latest lockfile 2023-09-27 18:12:54 +01:00
github-actions[bot] 6e11ab9183 chore: Update version for release (#513)
Co-authored-by: github-actions[bot] <github-actions[bot]@users.noreply.github.com>
2023-09-27 18:09:16 +01:00
Matt Aitken 2397fcb640 Express frameworks docs + CLI (#512)
* Manual setup docs

* Added the onboarding

* The emails package now works with Node > 18

* Added CLi support for Express, including custom init command finished messages

* Need to actually log out the installation complete message…

* Fix for an old Remix reference

* Renamed the page export

* Use resolvedOptions.triggerUrl

* Typo in manual instructions
2023-09-27 17:44:50 +01:00
Eric Allam 35d0c2a06f Implement the task output redacting to prevent redacted values from showing in the logs 2023-09-27 16:38:30 +01:00
Matt Aitken 813ec74672 Increased the intervalTrigger max from 1 day to 30 days
🚀 Publish Trigger.dev Docker / ʦ TypeScript (push) Has been cancelled
🚀 Publish Trigger.dev Docker / Unit Tests (push) Has been cancelled
🚀 Publish Trigger.dev Docker / e2e Tests (push) Has been cancelled
🚀 Publish Trigger.dev Docker / publish (push) Has been cancelled
2023-09-26 17:23:36 +01:00
Matt Aitken 03db13171a Latest lockfile 2023-09-26 16:46:41 +01:00
Matt Aitken 0ecb5129e8 Fix for incorrectly named Next.js package in manual setup 2023-09-26 12:47:49 +01:00
github-actions[bot] 44cb28c1c4 chore: Update version for release (#508)
Co-authored-by: github-actions[bot] <github-actions[bot]@users.noreply.github.com>
2023-09-26 11:49:03 +01:00
Matt Aitken 4578f6bd64 Astro CLI support (#506)
* Astro framework CLI support

* Changeset: Added Astro automatic installation

* Fixed the package name – was remix, now astro

* Fixed the export of the example

* Added IPv6 localhost to Astro hostnames

* Updated Astro onboarding to show the CLI init command, instead of manual instructions

* Astro quickstart

* Fix for type in Remix quickstart

* Next.js framework detection allows different config file extensions and “next” devDependency

* Made the dev command port more general so it works with various frameworks
2023-09-26 11:44:35 +01:00
Eric Allam 8b25e57613 hotfix 2 2023-09-24 14:43:02 -07:00
Eric Allam 8fb9ea19a3 hotfix 2023-09-24 14:04:46 -07:00
Eric Allam eb4ca0ce2d Fixed duplicate end month
🚀 Publish Trigger.dev Docker / ʦ TypeScript (push) Has been cancelled
🚀 Publish Trigger.dev Docker / Unit Tests (push) Has been cancelled
🚀 Publish Trigger.dev Docker / e2e Tests (push) Has been cancelled
🚀 Publish Trigger.dev Docker / publish (push) Has been cancelled
2023-09-22 17:24:55 -07:00
Eric Allam df24cd5b71 feat: Basic usage dashboard to show run volume (#501)
* New usage dashboard with static data

* Implement org usage dash

* Grab chart data for the last 12 months

* If no org is found just return undefined so a 404 will be shown

* Remove mock data

* Fill in missing months with 0s

---------

Co-authored-by: James Ritchie <james@jamesritchie.co.uk>
2023-09-22 17:02:53 -07:00
Matt Aitken 284324031c Latest lockfile 2023-09-22 16:23:21 -07:00
github-actions[bot] 2fdf42d444 chore: Update version for release (#481)
Co-authored-by: github-actions[bot] <github-actions[bot]@users.noreply.github.com>
2023-09-22 16:20:21 -07:00
Eric Allam 50137a6f24 Decouple zod (#500)
🚀 Publish Trigger.dev Docker / ʦ TypeScript (push) Has been cancelled
🚀 Publish Trigger.dev Docker / Unit Tests (push) Has been cancelled
🚀 Publish Trigger.dev Docker / e2e Tests (push) Has been cancelled
🚀 Publish Trigger.dev Docker / publish (push) Has been cancelled
Zod Schemas is no longer required for validating/inferring event triggers. We’ve taken inspiration from how domain-functions did it: https://github.com/seasonedcc/domain-functions/pull/114
2023-09-22 14:44:03 -07:00
Matt Aitken 3703f989fa Remix onboarding now uses the CLI init command 2023-09-22 14:39:21 -07:00
Matt Aitken 2a42942b01 The CLI now checks for a dev server API key in init and dev commands 2023-09-22 13:37:59 -07:00
Matt Aitken 7b29a92f6a Removed dynamicTrigger @internal from toJSON 2023-09-22 12:24:07 -07:00
Eric Allam 24465c9a7b Add back in Job.toJSON to fix the testing package 2023-09-22 11:21:37 -07:00
Eric Allam 42022b6ba4 Add BYO auth for oauth options 2023-09-22 10:31:40 -07:00
Matt Aitken dcc807d716 Linear getAll type error (weirdly not in VSCode…) and removed the pagination example that uses the SDK as won’t work with timeouts 2023-09-22 10:23:24 -07:00
Matt Aitken b5e37bfa61 Fix for getPathAlias typecheck failure 2023-09-22 09:59:59 -07:00
Matt Aitken 5681ebc756 Latest lockfile 2023-09-22 09:53:30 -07:00
Matt Aitken 886d6fda36 Tweaked the Linear scopes 2023-09-22 09:51:49 -07:00
Matt Aitken 1a4952720f CLI now supports multiple frameworks (with tests) (#480)
* Early work defining CLI framework support

* WIP moving CLI init logic to the Framework class

* Installing files should now work for Nextjs

* Some fixes

* WIP creating unit tests for Next.js project detection

* Delete old jest config

* Latest lockfile

* Detect use of src directory test

* Tests for detection pages/app directory

* Correct detection of Next.js project

* Renamed test file

* Create install files from template files with replacements. With tests

* Added multiple uses of the same replacement

* Created a test for the install step (it fails right now with JS)

* Removed unused import

* Another test that should pass but currently fails…

* Path alias fixed and now has tests

* New pathAlias function used

* Nextjs page install tests

* Fixed app directory install (with tests)

* Removed e2e CLI test, switched to unit testing strategy instead

* Latest lockfile

* The install files are now actual files that are copied and transformed

* Got the template files working correctly after building

* Next steps are now framework specific

* createFileFromTemplate now works with a path again. Uses mock if specified.

* Renamed apiRoute.js to pagesApiRoute.js

* Simplified pages file generation

* Next.js app API route template

* Next.js App routing support, with common files logic shared

* Dev command now uses framework default values if they exist and aren’t overridden

* Unused import

* pathAlias now works for all frameworks

* Added a test to detect Next from the next.config.js

* WIP on Remix framework support

* Tests for Remix install

* Replaced references to Next.js

* Use a green ✔️ instead of  in the CLI

* Support for multiple hostnames

* Tunneling can now use the hostname and port

* Work on multiple ports

* Improved the error messages. Added some extra pots to Next.js

* Update the Remix templates to have .server in the imports

* Remix updated to use server-runtime instead of node. Node v18+

* Frameworks can specify the watch paths and ignore paths

* Define the watch variables above, so we can easily log them for debugging

* Don’t wait for outdated package checking when running the dev command

* Improved the Remix manual setup guide

* Rewriting docs for quickstart

* Updated the Next.js quickstart

* Remix quick start

* Added a changeset

* Improved the Next.js manual setup
2023-09-22 08:55:04 -07:00
Eric Allam c0dfa8048a feat: BYO Auth (#491)
* feat: BYO Auth

Define client-side auth resolvers to be able to supply custom authentication credentials for integrations before a run is performed

- Added new defineAuthResolver
- Update all integrations to support the new auth resolvers
- Strip internal symbols from .d.ts in integrations and trigger-sdk
- Added BYO Auth docs
- Update Dynamic Schedule to support associated account IDs
- Create external accounts just-in-time
- Added Account ID field to test job when there are external auth integrations
- Show Account ID on run dashboard
- Added new Run error state called “Unresolved auth”

* Added changeset

* Remove @internal from TriggerIntegration public methods

* Add void to the result union

* DynamicTriggers now work with the new BYO auth system, and added a bunch of docs and docs changes

* Add additional key material for registering dynamic trigger task

* Add new define* instance methods to the overview
2023-09-22 08:54:31 -07:00
nicktrn 3e63a7e7a0 fix: Fail client-side on invalid Stripe event names (#492)
* Parse event names

* Add changeset
2023-09-22 08:53:04 -07:00
Eric Allam 4cc690a6ce Update sendevent.mdx 2023-09-22 08:44:14 -07:00
nicktrn 537447e318 Use absolute image paths (#490) 2023-09-21 15:45:40 -07:00
nicktrn 15f17d27e0 Going exponential with Linear (#478)
* Unleash GPT magic

* Clean up after GPT

* All the hooks

* Provisional integration catalog entry

* Sample webhook jobs

* Attachments with alpha warnings

* Remove some verbose logs

* Fix IP restrictions

* Remove tunnel

* Revert "Remove tunnel"

This reverts commit c5b69ce6524e3b40c26b66cdc56e087b8576b8c6.

* Resolve event name clashes

* Remove circular dependency

* Use correct payload uuid

* Schema fixes

* Fix webhook event name

* Start to Linearify catalog entry

* Remove todo

* More catalog updates

* Make OAuth work

* Rename webhook helper

* Schema juggling

* More discrimination

* Add Issue SLA event

* Simplify triggers

* Handle rate limits

* Fix Project schema

* Payload examples

* Improve event props

* Remove redundant source metadata

* One type to rule them all

* Recursive WithoutFunctions type

* Linear output serializer

* Some tasks

* Update catalog entry

* Dynamic usage sample

* Bump version

* Remove tunnel

* More tasks

* Add optional skipRetrying on runTask errors

* Fail fast on user errors

* Entity getter tasks

* Another couple of tasks

* Token to apiKey

* Sort tasks

* Add filtered issue SLA triggers

* Add docs

* Type fixes

* Job catalog examples

* Serialization helper docs

* Add changeset

* Refactor webhooks

* Enhance properties

* Pagination helper and docs

* Clean up imports

* Change misc catalog job

---------

Co-authored-by: Matt Aitken <matt@mattaitken.com>
2023-09-21 15:14:54 -07:00
Matt Aitken 91fc1e80f3 Use the bell icon for the new status Tasks
🚀 Publish Trigger.dev Docker / ʦ TypeScript (push) Has been cancelled
🚀 Publish Trigger.dev Docker / Unit Tests (push) Has been cancelled
🚀 Publish Trigger.dev Docker / e2e Tests (push) Has been cancelled
🚀 Publish Trigger.dev Docker / publish (push) Has been cancelled
2023-09-21 15:01:35 -07:00
Gregory dfe680a390 Update introduction.mdx (#498)
🚀 Publish Trigger.dev Docker / ʦ TypeScript (push) Has been cancelled
🚀 Publish Trigger.dev Docker / Unit Tests (push) Has been cancelled
🚀 Publish Trigger.dev Docker / e2e Tests (push) Has been cancelled
🚀 Publish Trigger.dev Docker / publish (push) Has been cancelled
Changed https://github.com/triggerdotdev/examples/tree/main/resend (which was a 404) to https://github.com/triggerdotdev/examples/tree/main/resend-email-form
2023-09-21 14:06:14 -07:00
Matt Aitken f73a2a424b Revert "Upgrade to the latest remix (pre v2)"
This reverts commit 4edc7112ab.
2023-09-21 13:19:03 -07:00
Matt Aitken 6ce87f169a Fixed dependency 2023-09-21 11:35:12 -07:00
Matt Aitken ad14983e21 React status hooks (#493)
* Added stripInternal to SDK tsconfig

* Statuses can now be set from a run, and are stored in the database

* Added the key to the returned status

* Made the test job have an extra step and only pass in some of the options

* client.getRunStatuses() and the corresponding endpoint

* client.getRun() now includes status info

* Fixed circular dependency schema

* Translate null to undefined

* Added the react package to the nextjs-reference tsconfig

* Removed unused OpenAI integration from nextjs-reference project

* New hooks for getting the statuses

* Disabled most of the nextjs-reference jobs

* Updated the hooks UI

* Updated the endpoints to deal with null statuses values

* The hook is working, with an example

* Changeset: “You can create statuses in your Jobs that can then be read using React hooks”

* Changeset config is back to the old changelog style

* WIP on new React hooks guide

* Guide docs for the new hooks

* Added the status hooks to the React hooks guide

* Removed the links to the status hooks reference for now

* Re-ordered the hooks

* Fix for an error in the docs

* Set a default of a blank array for the GetRunSchema
2023-09-21 10:24:37 -07:00
Eric Allam 4edc7112ab Upgrade to the latest remix (pre v2) 2023-09-18 22:20:41 +01:00
D-K-P c2ce707f3d Fixed cal.com link 2023-09-18 13:54:53 +01:00
Eric Allam cd70a29970 Improve the Astro manual setup guide 2023-09-15 17:40:36 +01:00
Eric Allam cd94d8fee2 Fixes broken pnpm lock file 2023-09-15 16:17:53 +01:00
Eric Allam 241e40e3e2 Add a references README 2023-09-15 16:12:50 +01:00
Eric Allam 486ed20ae7 Renamed the examples dir to references (true examples are in another repo and this was confusing) 2023-09-15 16:09:04 +01:00
James Ritchie e5ffc3a3f3 Added a link to the homepage from the side menu (#479)
* Added 2 new named icons

* New side menu link to the homepage

* Updated lock file

* Removed un-used import
2023-09-15 15:11:40 +01:00
Aniket Bindhani 363c74c6da Documentation Update: Added <github_username> instead of triggerdotdev to avoid confusion while cloning the repository (#477) 2023-09-15 12:55:56 +01:00
Eric Allam f98c425186 Redirect people to discord to ask a question 2023-09-15 12:54:54 +01:00
Matt Aitken da10ba907f Added instructions for how to do Changeset snapshots 2023-09-15 12:54:38 +01:00
Eric Allam 5aeef1233b A few teaks to the templates 2023-09-15 12:52:29 +01:00
Vishesh Rawal fb5f4e308f Created Temp for pull req, bugreport & feature req 2023-09-15 12:52:29 +01:00
Wesley 6ccb4fbc63 test/368/use vitest instead of jest (#470)
* chore: add vitest dependencies

* refactor: replace jest by vitest

* test: disable broken test

* Update pnpm-lock.yaml

* Update pnpm-lock.yaml

---------

Co-authored-by: Matt Aitken <matt@mattaitken.com>
2023-09-15 12:51:25 +01:00
Matt Aitken 88ebf4a1a0 Updated the Astro docs with SSR notes 2023-09-14 15:25:43 +01:00
Matt Aitken 9a0e6412e9 Bumped package versions to 2.1.3 2023-09-14 13:37:02 +01:00
Matt Aitken a68a912d56 Docs: improved the limitations 2023-09-14 11:55:51 +01:00
Matt Aitken dde51d6d25 Astro docs improvements 2023-09-14 11:55:38 +01:00
Matt Aitken dbd094216d Updated Astro setup docs: env import 2023-09-14 11:14:35 +01:00
558 changed files with 18881 additions and 6480 deletions
+1 -6
View File
@@ -1,11 +1,6 @@
{
"$schema": "https://unpkg.com/@changesets/config@2.2.0/schema.json",
"changelog": [
"@remix-run/changelog-github",
{
"repo": "triggerdotdev/trigger.dev"
}
],
"changelog": "@changesets/cli/changelog",
"commit": false,
"fixed": [
[
+13 -4
View File
@@ -1,35 +1,43 @@
*.log
\*.log
.git
.github
# editor
.idea
.vscode
# dependencies
node_modules
.pnp
.pnp.js
# testing
coverage
# next.js
.next/
build
# packages
build
dist
packages/**/dist
packages/\*\*/dist
# misc
.DS_Store
*.pem
\*.pem
.turbo
.vercel
.cache
.output
apps/**/public/build
apps/\*\*/public/build
cypress/screenshots
cypress/videos
@@ -38,6 +46,7 @@ apps/**/styles/tailwind.css
packages/**/styles/tailwind.css
.changeset
references
examples
CHANGESETS.md
CONTRIBUTING.md
+2
View File
@@ -31,6 +31,8 @@ CLOUD_AIRTABLE_CLIENT_ID=
CLOUD_AIRTABLE_CLIENT_SECRET=
CLOUD_GITHUB_CLIENT_ID=
CLOUD_GITHUB_CLIENT_SECRET=
CLOUD_LINEAR_CLIENT_ID=
CLOUD_LINEAR_CLIENT_SECRET=
CLOUD_SLACK_APP_HOST=
CLOUD_SLACK_CLIENT_ID=
CLOUD_SLACK_CLIENT_SECRET=
+38
View File
@@ -0,0 +1,38 @@
name: 🐞 Bug Report
description: Create a bug report to help us improve
title: "bug: "
labels: ["🐞 unconfirmed bug"]
body:
- type: textarea
attributes:
label: Provide environment information
description: |
Run this command in your project root and paste the results:
```bash
npx envinfo --system --binaries
```
validations:
required: true
- type: textarea
attributes:
label: Describe the bug
description: A clear and concise description of the bug, as well as what you expected to happen when encountering it.
validations:
required: true
- type: input
attributes:
label: Reproduction repo
description: If applicable, please provide a link to a reproduction repo or a Stackblitz / CodeSandbox project. Your issue may be closed if this is not provided and we are unable to reproduce the issue. If your bug is a docs issue, link the appropriate page.
validations:
required: true
- type: textarea
attributes:
label: To reproduce
description: Describe how to reproduce your bug. Steps, code snippets, reproduction repos etc.
validations:
required: true
- type: textarea
attributes:
label: Additional information
description: Add any other information related to the bug here, screenshots if applicable.
+5
View File
@@ -0,0 +1,5 @@
blank_issues_enabled: true
contact_links:
- name: Ask a Question
url: https://trigger.dev/discord
about: Ask questions and discuss with other community members
@@ -0,0 +1,27 @@
name: Feature Request
description: Suggest an idea for this project
title: "feat: "
labels: ["🌟 enhancement"]
body:
- type: textarea
attributes:
label: Is your feature request related to a problem? Please describe.
description: A clear and concise description of what the problem is. Ex. I'm always frustrated when [...]
validations:
required: true
- type: textarea
attributes:
label: Describe the solution you'd like to see
description: A clear and concise description of what you want to happen.
validations:
required: true
- type: textarea
attributes:
label: Describe alternate solutions
description: A clear and concise description of any alternative solutions or features you've considered.
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.
+12
View File
@@ -0,0 +1,12 @@
"📌 area: cli":
- any: ["cli/**/*"]
"📌 area: t3-app":
- any: ["cli/template/**/*"]
"📚 documentation":
- any: ["www/**/*"]
- any: ["**/*.md"]
"📌 area: ci":
- any: [".github/**/*"]
+27
View File
@@ -0,0 +1,27 @@
Closes #<issue>
## ✅ Checklist
- [ ] I have followed every step in the [contributing guide](https://github.com/triggerdotdev/trigger.dev/blob/main/CONTRIBUTING.md)
- [ ] The PR title follows the convention.
- [ ] I ran and tested the code works
---
## Testing
_[Describe the steps you took to test this change]_
---
## Changelog
_[Short description of what has changed]_
---
## Screenshots
_[Screenshots]_
💯
+2 -2
View File
@@ -129,10 +129,10 @@ jobs:
run: |
# Setup environment variables
cp ./.env.example ./.env
cp ./examples/nextjs-test/.env.example ./examples/nextjs-test/.env.local
cp ./references/nextjs-test/.env.example ./references/nextjs-test/.env.local
# Build packages
pnpm run build --filter @examples/nextjs-test^...
pnpm run build --filter @references/nextjs-test^...
pnpm --filter @trigger.dev/database generate
# Move trigger-cli bin to correct place
+9
View File
@@ -19,6 +19,15 @@
"name": "Chrome webapp",
"url": "http://localhost:3030",
"webRoot": "${workspaceFolder}/apps/webapp/app"
},
{
"type": "node-terminal",
"request": "launch",
"name": "Debug BYO Auth",
"command": "pnpm run byo-auth",
"envFile": "${workspaceFolder}/references/job-catalog/.env",
"cwd": "${workspaceFolder}/references/job-catalog",
"sourceMaps": true
}
]
}
+9
View File
@@ -27,3 +27,12 @@ Please follow the best-practice of adding changesets in the same commit as the c
3. Create version `pnpm run changeset:version`
4. Release `pnpm run changeset:release`
5. Switch back to normal mode by running `pnpm run changeset:normal`
## Snapshot instructions
!MAKE SURE TO UPDATE THE TAG IN THE INSTRUCTIONS BELOW!
1. Add changesets as usual `pnpm run changeset:add`
2. Create a snapshot version (replace "dev" with your tag) `pnpm exec changeset version --snapshot dev`
3. Build the packages: `pnpm run build --filter "@trigger.dev/*"`
4. Publish the snapshot (replace "dev" with your tag) `pnpm exec changeset publish --no-git-tag --snapshot --tag dev`
+21 -13
View File
@@ -23,11 +23,11 @@ branch are tagged into a release monthly.
1. Clone the repo into a public GitHub repository or [fork the repo](https://github.com/triggerdotdev/trigger.dev/fork). If you plan to distribute the code, keep the source code public to comply with the [Apache Licence 2.0](https://github.com/triggerdotdev/trigger.dev/blob/main/LICENSE).
```
git clone https://github.com/triggerdotdev/trigger.dev.git
git clone https://github.com/<github_username>/trigger.dev.git
```
> If you are on windows, run the following command on gitbash with admin privileges:
> `git clone -c core.symlinks=true https://github.com/triggerdotdev/trigger.dev.git`
> `git clone -c core.symlinks=true https://github.com/<github_username>/trigger.dev.git`
2. Navigate to the project folder
```
@@ -133,10 +133,10 @@ pnpm run dev
2. Open a new Terminal window and run the webapp locally and then create a new project in the dashboard. Copy out the dev API key.
3. Create a new temporary Next.js app in examples directory
3. Create a new temporary Next.js app in references directory
```sh
cd ./examples
cd ./references
pnpm create next-app@latest test-cli --ts --no-eslint --tailwind --app --src-dir --import-alias "@/*"
```
@@ -149,7 +149,7 @@ pnpm create next-app@latest test-cli --ts --no-eslint --tailwind --app --src-dir
}
```
5. Back in the terminal, navigate into the example, and initialize the CLI. When prompted, select `self-hosted` and enter `localhost:3030` if you are testing against the local instance of Trigger.dev, or you can just use the Trigger.dev cloud. When asked for an API key, use the key you copied earlier.
5. Back in the terminal, navigate into the reference, and initialize the CLI. When prompted, select `self-hosted` and enter `localhost:3030` if you are testing against the local instance of Trigger.dev, or you can just use the Trigger.dev cloud. When asked for an API key, use the key you copied earlier.
```sh
cd ./test-cli
@@ -179,14 +179,14 @@ To run the end-to-end tests, follow the steps below:
```sh
cp ./.env.example ./.env
cp ./examples/nextjs-test/.env.example ./examples/nextjs-test/.env.local
cp ./references/nextjs-test/.env.example ./references/nextjs-test/.env.local
```
2. Set up dependencies
```sh
# Build packages
pnpm run build --filter @examples/nextjs-test^...
pnpm run build --filter @references/nextjs-test^...
pnpm --filter @trigger.dev/database generate
# Move trigger-cli bin to correct place
@@ -221,11 +221,11 @@ pnpm run db:studio
## Add sample jobs
The [examples/jobs-starter](./examples/jobs-starter/) project defines simple jobs you can get started with.
The [references/job-catalog](./references/job-catalog/) project defines simple jobs you can get started with.
1. `cd` into `examples/jobs-starter`
2. Create a `.env.local` file with the following content,
replacing `[TRIGGER_DEV_API_KEY]` with an actual key:
1. `cd` into `references/job-catalog`
2. Create a `.env` file with the following content,
replacing `<TRIGGER_DEV_API_KEY>` with an actual key:
```env
TRIGGER_API_KEY=[TRIGGER_DEV_API_KEY]
@@ -235,12 +235,20 @@ TRIGGER_API_URL=http://localhost:3030
`TRIGGER_API_URL` is used to configure the URL for your Trigger.dev instance,
where the jobs will be registered.
3. Run the `jobs-starter` app:
3. Run one of the the `job-catalog` files:
```sh
pnpm dev
pnpm run events
```
This will open up a local server using `express` on port 8080. Then in a new terminal window you can run the trigger-cli dev command:
```sh
pnpm run dev:trigger
```
See the [Job Catalog](./references/job-catalog/README.md) file for more.
4. Navigate to your trigger.dev instance ([http://localhost:3030](http://localhost:3030/)), to see the jobs.
You can use the test feature to trigger them.
+1 -1
View File
@@ -106,7 +106,7 @@ export function TriggerDevStep() {
</Paragraph>
<TriggerDevCommand />
<Paragraph spacing variant="small">
If youre not running on port 3000 you can specify the port by adding{" "}
If youre not running on the default you can specify the port by adding{" "}
<InlineCode variant="extra-small">--port 3001</InlineCode> to the end.
</Paragraph>
<Paragraph spacing variant="small">
@@ -51,18 +51,18 @@ export function FrameworkSelector() {
<FrameworkLink to={projectSetupNextjsPath(organization, project)} supported>
<NextjsLogo className="w-32" />
</FrameworkLink>
<FrameworkLink to={projectSetupExpressPath(organization, project)}>
<FrameworkLink to={projectSetupExpressPath(organization, project)} supported>
<ExpressLogo className="w-36" />
</FrameworkLink>
<FrameworkLink to={projectSetupRemixPath(organization, project)} supported>
<RemixLogo className="w-32" />
</FrameworkLink>
<FrameworkLink to={projectSetupRedwoodPath(organization, project)}>
<RedwoodLogo className="w-44" />
</FrameworkLink>
<FrameworkLink to={projectSetupAstroPath(organization, project)} supported>
<AstroLogo className="w-32" />
</FrameworkLink>
<FrameworkLink to={projectSetupRedwoodPath(organization, project)}>
<RedwoodLogo className="w-44" />
</FrameworkLink>
<FrameworkLink to={projectSetupNuxtPath(organization, project)}>
<NuxtLogo className="w-32" />
</FrameworkLink>
@@ -272,6 +272,21 @@ export function HowToUseApiKeysAndEndpoints() {
you should use the Test feature to trigger any scheduled Jobs.
</Callout>
</StepContentContainer>
<StepNumber
stepNumber="→"
title={
<span className="flex items-center gap-x-2">
<span>Staging</span>
<EnvironmentLabel environment={{ type: "STAGING" }} />
</span>
}
/>
<StepContentContainer>
<Paragraph spacing>
The <InlineCode>STAGING</InlineCode> environment is where your Jobs will run in a staging
environment, meant to mirror your production environment.
</Paragraph>
</StepContentContainer>
<StepNumber
stepNumber="→"
title={
@@ -7,6 +7,7 @@ import { ConnectToOAuthForm } from "./ConnectToOAuthForm";
import { Paragraph } from "../primitives/Paragraph";
import { Client } from "~/presenters/IntegrationsPresenter.server";
import { UpdateOAuthForm } from "./UpdateOAuthForm";
import { LinkButton } from "../primitives/Buttons";
export function SelectOAuthMethod({
integration,
@@ -76,7 +77,7 @@ export function SelectOAuthMethod({
id="EXTERNAL"
value="EXTERNAL"
label="Your users"
description="We will give you OAuth React components so you can connect as your users."
description="Use an external authentication provider or your own user database to provide auth credentails of your users."
variant="description"
/>
</RadioGroup>
@@ -108,14 +109,19 @@ export function SelectOAuthMethod({
)
) : (
<>
<Header2 className="mb-1 mt-4">User OAuth coming soon</Header2>
<Header2 className="mb-1 mt-4">BYO Auth</Header2>
<Paragraph spacing>
End-user OAuth is going to be released soon. If you are interested in being an early
beta tester then please{" "}
<a href="mailto:founders@trigger.dev" className="text-indigo-500 underline">
message us
</a>
.
We support external authentication providers through Auth Resolvers. Read the docs to
learn more:{" "}
<LinkButton
variant="secondary/small"
LeadingIcon={"docs"}
TrailingIcon={"external-link"}
to="https://trigger.dev/docs/documentation/guides/using-integrations-byo-auth"
target="_blank"
>
Bring your own Auth
</LinkButton>
</Paragraph>
</>
))}
@@ -197,6 +197,8 @@ function classForJobStatus(status: JobRunStatus) {
case "TIMED_OUT":
case "WAITING_ON_CONNECTIONS":
case "PENDING":
case "UNRESOLVED_AUTH":
case "INVALID_PAYLOAD":
return "text-rose-500";
default:
return "";
@@ -11,8 +11,8 @@ import {
organizationTeamPath,
projectEnvironmentsPath,
projectIntegrationsPath,
projectSetupPath,
projectPath,
projectSetupPath,
projectTriggersPath,
} from "~/utils/pathBuilder";
import { UserProfilePhoto } from "../UserProfilePhoto";
@@ -96,7 +96,7 @@ export function ProjectSideMenu() {
data-action="environments & api keys"
/>
</div>
<div className="flex flex-col">
<div className="flex flex-col gap-1">
<SideMenuItem
name="Team"
icon="team"
@@ -118,6 +118,14 @@ export function ProjectSideMenu() {
isCollapsed={isCollapsed}
data-action="onboarding"
/>
<SideMenuItem
name="Changelog"
icon="list"
to="https://trigger.dev/changelog"
isCollapsed={isCollapsed}
data-action="changelog"
target="_blank"
/>
<SideMenuItem
name="Account"
icon={UserProfilePhoto}
@@ -146,6 +154,7 @@ function SideMenuItem({
isCollapsed,
forceActive,
hasWarning = false,
target,
}: {
icon: IconNames | React.ComponentType<any>;
name: string;
@@ -153,6 +162,7 @@ function SideMenuItem({
isCollapsed: boolean;
hasWarning?: boolean;
forceActive?: boolean;
target?: string;
}) {
return (
<SimpleTooltip
@@ -164,13 +174,14 @@ function SideMenuItem({
LeadingIcon={icon}
leadingIconClassName="text-dimmed"
to={to}
target={target}
className={({ isActive, isPending }) => {
if (forceActive !== undefined) {
isActive = forceActive;
}
return cn(
"relative",
isActive
isActive || isPending
? "bg-slate-800 text-bright group-hover:bg-slate-800"
: "text-dimmed group-hover:bg-slate-850 group-hover:text-bright"
);
@@ -286,9 +286,13 @@ type NavLinkPropsType = Pick<NavLinkProps, "to" | "target"> &
Omit<React.ComponentProps<typeof ButtonContent>, "className"> & {
className?: (props: { isActive: boolean; isPending: boolean }) => string | undefined;
};
export const NavLinkButton = ({ to, className, ...props }: NavLinkPropsType) => {
export const NavLinkButton = ({ to, className, target, ...props }: NavLinkPropsType) => {
return (
<NavLink to={to} className={cn("group outline-none", props.fullWidth ? "w-full" : "")}>
<NavLink
to={to}
className={cn("group outline-none", props.fullWidth ? "w-full" : "")}
target={target}
>
{({ isActive, isPending }) => (
<ButtonContent className={className && className({ isActive, isPending })} {...props} />
)}
@@ -35,6 +35,7 @@ import {
GlobeAltIcon,
HandRaisedIcon,
HeartIcon,
HomeIcon,
KeyIcon,
LightBulbIcon,
ListBulletIcon,
@@ -50,6 +51,7 @@ import {
UserGroupIcon,
UserIcon,
UserPlusIcon,
WindowIcon,
WrenchScrewdriverIcon,
XCircleIcon,
XMarkIcon,
@@ -74,7 +76,9 @@ const icons = {
"arrow-left": (className: string) => <ArrowLeftIcon className={cn("text-white", className)} />,
background: (className: string) => <CloudIcon className={cn("text-sky-400", className)} />,
beaker: (className: string) => <BeakerIcon className={cn("text-purple-500", className)} />,
bell: (className: string) => <BellAlertIcon className={cn("text-amber-500", className)} />,
billing: (className: string) => <CreditCardIcon className={cn("text-teal-500", className)} />,
browser: (className: string) => <WindowIcon className={cn("text-dimmed", className)} />,
calendar: (className: string) => (
<CalendarDaysIcon className={cn("text-purple-500", className)} />
),
@@ -111,6 +115,7 @@ const icons = {
<HandRaisedIcon className={cn("text-amber-400", className)} />
),
heart: (className: string) => <HeartIcon className={cn("text-rose-500", className)} />,
house: (className: string) => <HomeIcon className={cn("text-dimmed", className)} />,
id: (className: string) => <FingerPrintIcon className={cn("text-rose-200", className)} />,
inactive: (className: string) => <XCircleIcon className={cn("text-rose-500", className)} />,
info: (className: string) => <InformationCircleIcon className={cn("text-blue-500", className)} />,
@@ -126,6 +131,7 @@ const icons = {
"clipboard-checked": (className: string) => (
<ClipboardDocumentCheckIcon className={cn("text-dimmed", className)} />
),
list: (className: string) => <ListBulletIcon className={cn("text-slate-400", className)} />,
log: (className: string) => (
<ChatBubbleLeftEllipsisIcon className={cn("text-slate-400", className)} />
),
@@ -2,7 +2,7 @@ import { CodeBlock } from "~/components/code/CodeBlock";
import { DateTime } from "~/components/primitives/DateTime";
import { Paragraph } from "~/components/primitives/Paragraph";
import { RunStatusIcon, RunStatusLabel } from "~/components/runs/RunStatuses";
import { MatchedRun, useRun } from "~/hooks/useRun";
import { MatchedRun } from "~/hooks/useRun";
import { formatDuration } from "~/utils";
import {
RunPanel,
@@ -167,7 +167,13 @@ export function RunOverview({ run, trigger, showRerun, paths }: RunOverviewProps
<RunPanelHeader icon={trigger.icon} title={trigger.title} />
<RunPanelBody>
<RunPanelProperties
properties={[{ label: "Event name", text: run.event.name }, ...run.properties]}
properties={[{ label: "Event name", text: run.event.name }]
.concat(
run.event.externalAccount
? [{ label: "Account ID", text: run.event.externalAccount.identifier }]
: []
)
.concat(run.properties)}
/>
</RunPanelBody>
</RunPanel>
@@ -45,6 +45,13 @@ export function TriggerDetail({
/>
)}
<RunPanelIconProperty icon="id" label="Event name" value={name} />
{trigger.externalAccount && (
<RunPanelIconProperty
icon="account"
label="Account ID"
value={trigger.externalAccount.identifier}
/>
)}
</RunPanelIconSection>
<RunPanelDivider />
<div className="mt-4 flex flex-col gap-2">
+30 -16
View File
@@ -1,15 +1,14 @@
import type { JobRunExecution, JobRunStatus } from "@trigger.dev/database";
import { NoSymbolIcon } from "@heroicons/react/20/solid";
import {
CheckCircleIcon,
ClockIcon,
ExclamationTriangleIcon,
StopIcon,
WrenchIcon,
XCircleIcon,
} from "@heroicons/react/24/solid";
import type { JobRunStatus } from "@trigger.dev/database";
import { cn } from "~/utils/cn";
import { Spinner } from "../primitives/Spinner";
import { HandRaisedIcon, NoSymbolIcon } from "@heroicons/react/20/solid";
export function hasFinished(status: JobRunStatus): boolean {
return (
@@ -17,7 +16,9 @@ export function hasFinished(status: JobRunStatus): boolean {
status === "FAILURE" ||
status === "ABORTED" ||
status === "TIMED_OUT" ||
status === "CANCELED"
status === "CANCELED" ||
status === "UNRESOLVED_AUTH" ||
status === "INVALID_PAYLOAD"
);
}
@@ -48,6 +49,9 @@ export function RunStatusIcon({ status, className }: { status: JobRunStatus; cla
return <XCircleIcon className={cn(runStatusClassNameColor(status), className)} />;
case "TIMED_OUT":
return <ExclamationTriangleIcon className={cn(runStatusClassNameColor(status), className)} />;
case "UNRESOLVED_AUTH":
case "INVALID_PAYLOAD":
return <XCircleIcon className={cn(runStatusClassNameColor(status), className)} />;
case "WAITING_ON_CONNECTIONS":
return <WrenchIcon className={cn(runStatusClassNameColor(status), className)} />;
case "ABORTED":
@@ -63,26 +67,26 @@ export type RunBasicStatus = "WAITING" | "PENDING" | "RUNNING" | "COMPLETED" | "
export function runBasicStatus(status: JobRunStatus): RunBasicStatus {
switch (status) {
case "SUCCESS":
return "COMPLETED";
case "WAITING_ON_CONNECTIONS":
case "QUEUED":
case "PREPROCESSING":
case "PENDING":
return "PENDING";
case "STARTED":
return "RUNNING";
case "QUEUED":
return "PENDING";
case "FAILURE":
return "FAILED";
case "TIMED_OUT":
return "FAILED";
case "WAITING_ON_CONNECTIONS":
return "PENDING";
case "ABORTED":
return "FAILED";
case "PREPROCESSING":
return "PENDING";
case "UNRESOLVED_AUTH":
case "CANCELED":
case "ABORTED":
case "INVALID_PAYLOAD":
return "FAILED";
case "SUCCESS":
return "COMPLETED";
default: {
const _exhaustiveCheck: never = status;
throw new Error(`Non-exhaustive match for value: ${status}`);
}
}
}
@@ -108,6 +112,14 @@ export function runStatusTitle(status: JobRunStatus): string {
return "Preprocessing";
case "CANCELED":
return "Canceled";
case "UNRESOLVED_AUTH":
return "Unresolved auth";
case "INVALID_PAYLOAD":
return "Invalid payload";
default: {
const _exhaustiveCheck: never = status;
throw new Error(`Non-exhaustive match for value: ${status}`);
}
}
}
@@ -122,6 +134,8 @@ export function runStatusClassNameColor(status: JobRunStatus): string {
case "QUEUED":
return "text-amber-300";
case "FAILURE":
case "UNRESOLVED_AUTH":
case "INVALID_PAYLOAD":
return "text-rose-500";
case "TIMED_OUT":
return "text-amber-300";
+1
View File
@@ -5,3 +5,4 @@ export const DEFAULT_MAX_CONCURRENT_RUNS = 10;
export const MAX_CONCURRENT_RUNS_LIMIT = 20;
export const PREPROCESS_RETRY_LIMIT = 2;
export const EXECUTE_JOB_RETRY_LIMIT = 10;
export const MAX_RUN_YIELDED_EXECUTIONS = 100;
+7 -5
View File
@@ -31,7 +31,7 @@ export type PrismaTransactionOptions = {
/** Sets the transaction isolation level. By default this is set to the value currently configured in your database. */
isolationLevel?: Prisma.TransactionIsolationLevel;
rethrowPrismaErrors?: boolean;
swallowPrismaErrors?: boolean;
};
export async function $transaction<R>(
@@ -55,11 +55,9 @@ export async function $transaction<R>(
name: error.name,
});
if (options?.rethrowPrismaErrors) {
throw error;
if (options?.swallowPrismaErrors) {
return;
}
return;
}
throw error;
@@ -124,6 +122,10 @@ function getClient() {
emit: "stdout",
level: "warn",
},
// {
// emit: "stdout",
// level: "query",
// },
],
});
@@ -8,7 +8,6 @@ import type {
import { customAlphabet } from "nanoid";
import slug from "slug";
import { prisma, PrismaClientOrTransaction } from "~/db.server";
import { workerQueue } from "~/services/worker.server";
import { createProject } from "./project.server";
export type { Organization };
@@ -76,6 +75,10 @@ export async function createOrganization(
},
attemptCount = 0
): Promise<Organization & { projects: Project[] }> {
if (typeof process.env.BLOCKED_USERS === "string" && process.env.BLOCKED_USERS.includes(userId)) {
throw new Error("Organization could not be created.");
}
const uniqueOrgSlug = `${slug(title)}-${nanoid(4)}`;
const orgWithSameSlug = await prisma.organization.findFirst({
@@ -172,10 +175,10 @@ function envSlug(environmentType: RuntimeEnvironment["type"]) {
return "prod";
}
case "STAGING": {
return "staging";
return "stg";
}
case "PREVIEW": {
return "preview";
return "prev";
}
}
}
+1
View File
@@ -65,6 +65,7 @@ export async function createProject(
// Create the dev and prod environments
await createEnvironment(organization, project, "PRODUCTION");
await createEnvironment(organization, project, "STAGING");
for (const member of project.organization.members) {
await createEnvironment(organization, project, "DEVELOPMENT", member);
@@ -16,7 +16,7 @@ export async function resolveRunConnections(
const result: Record<string, ConnectionAuth> = {};
for (const connection of connections) {
if (connection.integration.authSource === "LOCAL") {
if (connection.integration.authSource !== "HOSTED") {
continue;
}
+86 -1
View File
@@ -1,5 +1,5 @@
import type { Task, TaskAttempt } from "@trigger.dev/database";
import { ServerTask } from "@trigger.dev/core";
import { CachedTask, ServerTask } from "@trigger.dev/core";
export type TaskWithAttempts = Task & { attempts: TaskAttempt[] };
@@ -23,5 +23,90 @@ export function taskWithAttemptsToServerTask(task: TaskWithAttempts): ServerTask
attempts: task.attempts.length,
idempotencyKey: task.idempotencyKey,
operation: task.operation,
callbackUrl: task.callbackUrl,
};
}
export type TaskForCaching = Pick<
Task,
"id" | "status" | "idempotencyKey" | "noop" | "output" | "parentId"
>;
export function prepareTasksForCaching(
possibleTasks: TaskForCaching[],
maxSize: number
): {
tasks: CachedTask[];
cursor: string | undefined;
} {
const tasks = possibleTasks.filter((task) => task.status === "COMPLETED" && !task.noop);
// Select tasks using greedy approach
const tasksToRun: CachedTask[] = [];
let remainingSize = maxSize;
for (const task of tasks) {
const cachedTask = prepareTaskForCaching(task);
const size = calculateCachedTaskSize(cachedTask);
if (size <= remainingSize) {
tasksToRun.push(cachedTask);
remainingSize -= size;
}
}
return {
tasks: tasksToRun,
cursor: tasks.length > tasksToRun.length ? tasks[tasksToRun.length].id : undefined,
};
}
export function prepareTasksForCachingLegacy(
possibleTasks: TaskForCaching[],
maxSize: number
): {
tasks: CachedTask[];
cursor: string | undefined;
} {
const tasks = possibleTasks.filter((task) => task.status === "COMPLETED");
// Prepare tasks and calculate their sizes
const availableTasks = tasks.map((task) => {
const cachedTask = prepareTaskForCaching(task);
return { task: cachedTask, size: calculateCachedTaskSize(cachedTask) };
});
// Sort tasks in ascending order by size
availableTasks.sort((a, b) => a.size - b.size);
// Select tasks using greedy approach
const tasksToRun: CachedTask[] = [];
let remainingSize = maxSize;
for (const { task, size } of availableTasks) {
if (size <= remainingSize) {
tasksToRun.push(task);
remainingSize -= size;
}
}
return {
tasks: tasksToRun,
cursor: undefined,
};
}
function prepareTaskForCaching(task: TaskForCaching): CachedTask {
return {
id: task.idempotencyKey, // We should eventually move this back to task.id
status: task.status,
idempotencyKey: task.idempotencyKey,
noop: task.noop,
output: task.output as any,
parentId: task.parentId,
};
}
function calculateCachedTaskSize(task: CachedTask): number {
return JSON.stringify(task).length;
}
@@ -2,19 +2,19 @@ import { PrismaClient, prisma } from "~/db.server";
import { IndexEndpointStats, parseEndpointIndexStats } from "~/models/indexEndpoint.server";
import { Project } from "~/models/project.server";
import { User } from "~/models/user.server";
import {
import type {
Endpoint,
EndpointIndex,
RuntimeEnvironment,
RuntimeEnvironmentType,
} from "../../../../packages/database/src";
import { env } from "~/env.server";
} from "@trigger.dev/database";
export type Client = {
slug: string;
endpoints: {
DEVELOPMENT: ClientEndpoint;
PRODUCTION: ClientEndpoint;
STAGING?: ClientEndpoint;
};
};
@@ -133,6 +133,8 @@ export class EnvironmentsPresenter {
throw new Error("Development environment not found, this should not happen");
}
const stagingEnvironment = filtered.find((environment) => environment.type === "STAGING");
const productionEnvironment = filtered.find(
(environment) => environment.type === "PRODUCTION"
);
@@ -151,6 +153,9 @@ export class EnvironmentsPresenter {
state: "unconfigured",
environment: productionEnvironment,
},
STAGING: stagingEnvironment
? { state: "unconfigured", environment: stagingEnvironment }
: undefined,
},
};
@@ -161,6 +166,16 @@ export class EnvironmentsPresenter {
client.endpoints.DEVELOPMENT = endpointClient(devEndpoint, developmentEnvironment, baseUrl);
}
if (stagingEnvironment) {
const stagingEndpoint = stagingEnvironment.endpoints.find(
(endpoint) => endpoint.slug === slug
);
if (stagingEndpoint) {
client.endpoints.STAGING = endpointClient(stagingEndpoint, stagingEnvironment, baseUrl);
}
}
const prodEndpoint = productionEnvironment.endpoints.find(
(endpoint) => endpoint.slug === slug
);
@@ -120,8 +120,12 @@ export class IntegrationClientPresenter {
icon: integration.definition.icon,
},
authMethod: {
type: integration.authMethod?.type ?? "local",
name: integration.authMethod?.name ?? "Local Auth",
type:
integration.authMethod?.type ?? integration.authSource === "RESOLVER" ? "local" : "local",
name:
integration.authMethod?.name ?? integration.authSource === "RESOLVER"
? "Auth Resolver"
: "Local Auth",
},
help,
};
@@ -125,8 +125,9 @@ export class IntegrationsPresenter {
name: c.definition.name,
},
authMethod: {
type: c.authMethod?.type ?? "local",
name: c.authMethod?.name ?? "Local Only",
type: c.authMethod?.type ?? c.authSource === "RESOLVER" ? "resolver" : "local",
name:
c.authMethod?.name ?? c.authSource === "RESOLVER" ? "Auth Resolver" : "Local Only",
},
authSource: c.authSource,
setupStatus: c.setupStatus,
@@ -0,0 +1,198 @@
import { PrismaClient, prisma } from "~/db.server";
export class OrgUsagePresenter {
#prismaClient: PrismaClient;
constructor(prismaClient: PrismaClient = prisma) {
this.#prismaClient = prismaClient;
}
public async call({ userId, slug }: { userId: string; slug: string }) {
const organization = await this.#prismaClient.organization.findFirst({
where: {
slug,
members: {
some: {
userId,
},
},
},
});
if (!organization) {
return;
}
const startOfMonth = new Date(new Date().getFullYear(), new Date().getMonth(), 1);
const startOfLastMonth = new Date(new Date().getFullYear(), new Date().getMonth() - 1, 1); // this works for January as well
// Get count of runs since the start of the current month
const runsCount = await this.#prismaClient.jobRun.count({
where: {
organizationId: organization.id,
createdAt: {
gte: new Date(new Date().getFullYear(), new Date().getMonth(), 1),
},
},
});
// Get the count of runs for last month
const runsCountLastMonth = await this.#prismaClient.jobRun.count({
where: {
organizationId: organization.id,
createdAt: {
gte: startOfLastMonth,
lt: startOfMonth,
},
},
});
// Get the count of the runs for the last 6 months, by month. So for example we want the data shape to be:
// [
// { month: "2021-01", count: 10 },
// { month: "2021-02", count: 20 },
// { month: "2021-03", count: 30 },
// { month: "2021-04", count: 40 },
// { month: "2021-05", count: 50 },
// { month: "2021-06", count: 60 },
// ]
// This will be used to generate the chart on the usage page
// Use prisma queryRaw for this since prisma doesn't support grouping by month
const chartDataRaw = await this.#prismaClient.$queryRaw<
{
month: string;
count: number;
}[]
>`SELECT TO_CHAR("createdAt", 'YYYY-MM') as month, COUNT(*) as count FROM "JobRun" WHERE "organizationId" = ${organization.id} AND "createdAt" >= NOW() - INTERVAL '6 months' GROUP BY month ORDER BY month ASC`;
const chartData = chartDataRaw.map((obj) => ({
name: obj.month,
total: Number(obj.count), // Convert BigInt to Number
}));
const totalJobs = await this.#prismaClient.job.count({
where: {
organizationId: organization.id,
internal: false,
},
});
const totalJobsLastMonth = await this.#prismaClient.job.count({
where: {
organizationId: organization.id,
createdAt: {
lt: startOfMonth,
},
deletedAt: null,
internal: false,
},
});
const totalIntegrations = await this.#prismaClient.integration.count({
where: {
organizationId: organization.id,
},
});
const totalIntegrationsLastMonth = await this.#prismaClient.integration.count({
where: {
organizationId: organization.id,
createdAt: {
lt: startOfMonth,
},
},
});
const totalMembers = await this.#prismaClient.orgMember.count({
where: {
organizationId: organization.id,
},
});
const jobs = await this.#prismaClient.job.findMany({
where: {
organizationId: organization.id,
deletedAt: null,
internal: false,
},
select: {
id: true,
slug: true,
_count: {
select: {
runs: {
where: {
createdAt: {
gte: startOfMonth,
},
},
},
},
},
project: {
select: {
id: true,
name: true,
slug: true,
},
},
},
});
return {
id: organization.id,
runsCount,
runsCountLastMonth,
chartData: fillInMissingMonthlyData(chartData, 6),
totalJobs,
totalJobsLastMonth,
totalIntegrations,
totalIntegrationsLastMonth,
totalMembers,
jobs,
};
}
}
// This will fill in missing chart data with zeros
// So for example, if data is [{ name: "2021-01", total: 10 }, { name: "2021-03", total: 30 }] and the totalNumberOfMonths is 6
// And the current month is "2021-04", then this function will return:
// [{ name: "2020-11", total: 0 }, { name: "2020-12", total: 0 }, { name: "2021-01", total: 10 }, { name: "2021-02", total: 0 }, { name: "2021-03", total: 30 }, { name: "2021-04", total: 0 }]
function fillInMissingMonthlyData(
data: Array<{ name: string; total: number }>,
totalNumberOfMonths: number
): Array<{ name: string; total: number }> {
const currentMonth = new Date().toISOString().slice(0, 7);
const startMonth = new Date(
new Date(currentMonth).getFullYear(),
new Date(currentMonth).getMonth() - totalNumberOfMonths,
1
)
.toISOString()
.slice(0, 7);
const months = getMonthsBetween(startMonth, currentMonth);
let completeData = months.map((month) => {
let foundData = data.find((d) => d.name === month);
return foundData ? { ...foundData } : { name: month, total: 0 };
});
return completeData;
}
function getMonthsBetween(startMonth: string, endMonth: string): string[] {
const startDate = new Date(startMonth);
const endDate = new Date(endMonth);
const months = [];
let currentDate = startDate;
while (currentDate <= endDate) {
months.push(currentDate.toISOString().slice(0, 7));
currentDate = new Date(currentDate.setMonth(currentDate.getMonth() + 1));
}
return months;
}
@@ -115,6 +115,11 @@ export class RunPresenter {
payload: true,
timestamp: true,
deliveredAt: true,
externalAccount: {
select: {
identifier: true,
},
},
},
},
tasks: {
@@ -1,6 +1,8 @@
import { RedactSchema } from "@trigger.dev/core";
import { StyleSchema } from "@trigger.dev/core";
import { PrismaClient, prisma } from "~/db.server";
import { mergeProperties } from "~/utils/mergeProperties.server";
import { Redactor } from "~/utils/redactor";
type DetailsProps = {
id: string;
@@ -61,6 +63,7 @@ export class TaskDetailsPresenter {
completedAt: true,
style: true,
parentId: true,
redact: true,
attempts: {
select: {
number: true,
@@ -85,11 +88,32 @@ export class TaskDetailsPresenter {
return {
...task,
output: task.output ? JSON.stringify(task.output, null, 2) : undefined,
redact: undefined,
output: task.output
? JSON.stringify(this.#stringifyOutputWithRedactions(task.output, task.redact), null, 2)
: undefined,
connection: task.runConnection,
params: task.params as Record<string, any>,
properties: mergeProperties(task.properties, task.outputProperties),
style: task.style ? StyleSchema.parse(task.style) : undefined,
};
}
#stringifyOutputWithRedactions(output: any, redact: unknown): any {
if (!output) {
return;
}
const parsedRedact = RedactSchema.safeParse(redact);
if (!parsedRedact.success) {
return output;
}
const paths = parsedRedact.data.paths;
const redactor = new Redactor(paths);
return redactor.redact(output);
}
}
@@ -39,6 +39,15 @@ export class TestJobPresenter {
payload: true,
},
},
integrations: {
select: {
integration: {
select: {
authSource: true,
},
},
},
},
},
},
environment: {
@@ -99,6 +108,9 @@ export class TestJobPresenter {
...example,
payload: JSON.stringify(example.payload, exampleReplacer, 2),
})),
hasAuthResolver: alias.version.integrations.some(
(i) => i.integration.authSource === "RESOLVER"
),
})),
hasTestRuns: job._count.runs > 0,
};
@@ -22,6 +22,11 @@ export class TriggerDetailsPresenter {
payload: true,
timestamp: true,
deliveredAt: true,
externalAccount: {
select: {
identifier: true,
},
},
},
},
},
@@ -1,17 +1,178 @@
import { ComingSoon } from "~/components/ComingSoon";
import { PageContainer, PageBody } from "~/components/layout/AppLayout";
import { ArrowRightIcon } from "@heroicons/react/20/solid";
import {
ForwardIcon,
SquaresPlusIcon,
UsersIcon,
WrenchScrewdriverIcon,
} from "@heroicons/react/24/solid";
import { Bar, BarChart, ResponsiveContainer, Tooltip, TooltipProps, XAxis, YAxis } from "recharts";
import { PageBody, PageContainer } from "~/components/layout/AppLayout";
import { Header2 } from "~/components/primitives/Headers";
import { Paragraph } from "~/components/primitives/Paragraph";
import { TextLink } from "~/components/primitives/TextLink";
import { useOrganization } from "~/hooks/useOrganizations";
import { OrganizationParamsSchema, jobPath, organizationTeamPath } from "~/utils/pathBuilder";
import { OrgAdminHeader } from "../_app.orgs.$organizationSlug._index/OrgAdminHeader";
import { Link } from "@remix-run/react/dist/components";
import { LoaderArgs } from "@remix-run/server-runtime";
import { typedjson, useTypedLoaderData } from "remix-typedjson";
import { OrgUsagePresenter } from "~/presenters/OrgUsagePresenter.server";
import { requireUserId } from "~/services/session.server";
export async function loader({ params, request }: LoaderArgs) {
const userId = await requireUserId(request);
const { organizationSlug } = OrganizationParamsSchema.parse(params);
const presenter = new OrgUsagePresenter();
const data = await presenter.call({ userId, slug: organizationSlug });
if (!data) {
throw new Response(null, { status: 404 });
}
return typedjson(data);
}
const CustomTooltip = ({ active, payload, label }: TooltipProps<number, string>) => {
if (active && payload) {
return (
<div className="flex items-center gap-2 rounded border border-border bg-slate-900 px-4 py-2 text-sm text-dimmed">
<p className="text-white">{label}:</p>
<p className="text-white">{payload[0].value}</p>
</div>
);
}
return null;
};
export default function Page() {
const organization = useOrganization();
const loaderData = useTypedLoaderData<typeof loader>();
return (
<PageContainer>
<OrgAdminHeader />
<PageBody>
<ComingSoon
title="Usage & billing"
description="View your usage, tier and billing information. During the beta we will display usage and start billing if you exceed your limits. But don't worry, we'll give you plenty of warning."
icon="billing"
/>
<div className="mb-4 grid gap-4 md:grid-cols-2 lg:grid-cols-4">
<div className="rounded border border-border p-6">
<div className="flex flex-row items-center justify-between space-y-0 pb-2">
<Header2>Total Runs this month</Header2>
<ForwardIcon className="h-6 w-6 text-dimmed" />
</div>
<div>
<p className="text-3xl font-bold">{loaderData.runsCount.toLocaleString()}</p>
<Paragraph variant="small" className="text-dimmed">
{loaderData.runsCountLastMonth} runs last month
</Paragraph>
</div>
</div>
<div className="rounded border border-border p-6">
<div className="flex flex-row items-center justify-between space-y-0 pb-2">
<Header2>Total Jobs</Header2>
<WrenchScrewdriverIcon className="h-6 w-6 text-dimmed" />
</div>
<div>
<p className="text-3xl font-bold">{loaderData.totalJobs.toLocaleString()}</p>
<Paragraph variant="small" className="text-dimmed">
{loaderData.totalJobs === loaderData.totalJobsLastMonth ? (
<>No change since last month</>
) : loaderData.totalJobs > loaderData.totalJobsLastMonth ? (
<>+{loaderData.totalJobs - loaderData.totalJobsLastMonth} since last month</>
) : (
<>-{loaderData.totalJobsLastMonth - loaderData.totalJobs} since last month</>
)}
</Paragraph>
</div>
</div>
<div className="rounded border border-border p-6">
<div className="flex flex-row items-center justify-between space-y-0 pb-2">
<Header2>Total Integrations</Header2>
<SquaresPlusIcon className="h-6 w-6 text-dimmed" />
</div>
<div>
<p className="text-3xl font-bold">{loaderData.totalIntegrations.toLocaleString()}</p>
<Paragraph variant="small" className="text-dimmed">
{loaderData.totalIntegrations === loaderData.totalIntegrationsLastMonth ? (
<>No change since last month</>
) : loaderData.totalIntegrations > loaderData.totalIntegrationsLastMonth ? (
<>
+{loaderData.totalIntegrations - loaderData.totalIntegrationsLastMonth} since
last month
</>
) : (
<>
-{loaderData.totalIntegrationsLastMonth - loaderData.totalIntegrations} since
last month
</>
)}
</Paragraph>
</div>
</div>
<div className="rounded border border-border p-6">
<div className="flex flex-row items-center justify-between space-y-0 pb-2">
<Header2>Team members</Header2>
<UsersIcon className="h-6 w-6 text-dimmed" />
</div>
<div>
<p className="text-3xl font-bold">{loaderData.totalMembers.toLocaleString()}</p>
<TextLink
to={organizationTeamPath(organization)}
className="group text-sm text-dimmed hover:text-bright"
>
Manage
<ArrowRightIcon className="-mb-0.5 ml-0.5 h-4 w-4 text-dimmed transition group-hover:translate-x-1 group-hover:text-bright" />
</TextLink>
</div>
</div>
</div>
<div className="flex max-h-[500px] gap-x-4">
<div className="w-1/2 rounded border border-border py-6 pr-2">
<Header2 className="mb-8 pl-6">Job Runs per month</Header2>
<ResponsiveContainer width="100%" height={400}>
<BarChart data={loaderData.chartData}>
<XAxis
dataKey="name"
stroke="#888888"
fontSize={12}
tickLine={false}
axisLine={false}
/>
<YAxis
stroke="#888888"
fontSize={12}
tickLine={false}
axisLine={false}
tickFormatter={(value) => `${value}`}
/>
<Tooltip cursor={{ fill: "rgba(255,255,255,0.05)" }} content={<CustomTooltip />} />
<Bar dataKey="total" fill="#DB2777" radius={[4, 4, 0, 0]} />
</BarChart>
</ResponsiveContainer>
</div>
<div className="w-1/2 overflow-y-auto rounded border border-border px-3 py-6">
<div className="mb-2 flex items-baseline justify-between border-b border-border px-3 pb-4">
<Header2 className="">Jobs</Header2>
<Header2 className="">Runs</Header2>
</div>
<div className="space-y-2">
{loaderData.jobs.map((job) => (
<Link
to={jobPath(organization, job.project, job)}
className="flex items-center rounded px-4 py-3 transition hover:bg-slate-850"
key={job.id}
>
<div className="space-y-1">
<p className="text-sm font-medium leading-none">{job.slug}</p>
<p className="text-sm text-muted-foreground">Project: {job.project.name}</p>
</div>
<div className="ml-auto font-medium">{job._count.runs.toLocaleString()}</div>
</Link>
))}
</div>
</div>
</div>
</PageBody>
</PageContainer>
);
@@ -104,6 +104,7 @@ export default function Page() {
fullWidth={true}
value={filterText}
onChange={(e) => setFilterText(e.target.value)}
autoFocus
/>
<HelpTrigger title="Example Jobs and inspiration" />
</div>
@@ -160,7 +161,7 @@ function ExampleJobs() {
height="250"
allow="accelerometer; autoplay; clipboard-write; encrypted-media; gyroscope; picture-in-picture; web-share"
allowFullScreen
className="mb-4 w-full border-b border-slate-800"
className="mb-4 border-b border-slate-800"
/>
<Header2 spacing>How to create a Job</Header2>
<Paragraph variant="small" spacing>
@@ -85,8 +85,8 @@ export default function Page() {
const client = clients.find((c) => c.slug === selected.client);
if (!client) return undefined;
if (selected.type === "PREVIEW" || selected.type === "STAGING") {
throw new Error("PREVIEW/STAGING is not yet supported");
if (selected.type === "PREVIEW") {
throw new Error("PREVIEW is not yet supported");
}
return {
@@ -195,6 +195,18 @@ export default function Page() {
})
}
/>
{client.endpoints.STAGING && (
<EndpointRow
endpoint={client.endpoints.STAGING}
type="STAGING"
onClick={() =>
setSelected({
client: client.slug,
type: "STAGING",
})
}
/>
)}
<EndpointRow
endpoint={client.endpoints.PRODUCTION}
type="PRODUCTION"
@@ -218,7 +230,7 @@ export default function Page() {
</>
)}
</div>
{selectedEndpoint && (
{selectedEndpoint && selectedEndpoint.endpoint && (
<ConfigureEndpointSheet
slug={selectedEndpoint.clientSlug}
endpoint={selectedEndpoint.endpoint}
@@ -1,4 +1,4 @@
import { useForm } from "@conform-to/react";
import { conform, useForm } from "@conform-to/react";
import { parse } from "@conform-to/zod";
import { PopoverTrigger } from "@radix-ui/react-popover";
import { Form, useActionData, useSubmit } from "@remix-run/react";
@@ -14,6 +14,9 @@ import { Button, ButtonContent } from "~/components/primitives/Buttons";
import { Callout } from "~/components/primitives/Callout";
import { FormError } from "~/components/primitives/FormError";
import { Help, HelpContent, HelpTrigger } from "~/components/primitives/Help";
import { Input } from "~/components/primitives/Input";
import { InputGroup } from "~/components/primitives/InputGroup";
import { Label } from "~/components/primitives/Label";
import { Popover, PopoverContent } from "~/components/primitives/Popover";
import {
Select,
@@ -69,6 +72,7 @@ const schema = z.object({
}),
environmentId: z.string(),
versionId: z.string(),
accountId: z.string().optional(),
});
//todo save the chosen environment to a cookie (for that user), use it to default the env dropdown
@@ -84,11 +88,7 @@ export const action: ActionFunction = async ({ request, params }) => {
}
const testService = new TestJobService();
const run = await testService.call({
environmentId: submission.value.environmentId,
payload: submission.value.payload,
versionId: submission.value.versionId,
});
const run = await testService.call(submission.value);
if (!run) {
return redirectBackWithErrorMessage(
@@ -124,6 +124,7 @@ export default function Page() {
const [defaultJson, setDefaultJson] = useState<string>(startingJson);
const currentJson = useRef<string>(defaultJson);
const [selectedEnvironmentId, setSelectedEnvironmentId] = useState<string>(environments[0].id);
const [currentAccountId, setCurrentAccountId] = useState<string | undefined>(undefined);
const selectedEnvironment = environments.find((e) => e.id === selectedEnvironmentId);
@@ -139,6 +140,7 @@ export default function Page() {
payload: currentJson.current,
environmentId: selectedEnvironmentId,
versionId: selectedEnvironment?.versionId ?? "",
...(currentAccountId ? { accountId: currentAccountId } : {}),
},
{
action: "",
@@ -147,10 +149,10 @@ export default function Page() {
);
e.preventDefault();
},
[currentJson, selectedEnvironmentId]
[currentJson, selectedEnvironmentId, currentAccountId]
);
const [form, { environmentId, payload }] = useForm({
const [form, { environmentId, payload, accountId }] = useForm({
id: "test-job",
lastSubmission,
onValidate({ formData }) {
@@ -234,15 +236,32 @@ export default function Page() {
</div>
<HelpTrigger title="How do I run a test?" />
</div>
<div className="flex-1 overflow-auto rounded border border-slate-850 scrollbar-thin scrollbar-track-transparent scrollbar-thumb-slate-700">
<JSONEditor
defaultValue={defaultJson}
readOnly={false}
basicSetup
onChange={(v) => (currentJson.current = v)}
minHeight="150px"
/>
</div>
<InputGroup fullWidth>
<Label variant="small">Payload</Label>
<div className="flex-1 overflow-auto rounded border border-slate-850 scrollbar-thin scrollbar-track-transparent scrollbar-thumb-slate-700">
<JSONEditor
defaultValue={defaultJson}
readOnly={false}
basicSetup
onChange={(v) => (currentJson.current = v)}
minHeight="150px"
/>
</div>
</InputGroup>
{selectedEnvironment?.hasAuthResolver && (
<InputGroup fullWidth className="mb-4 mt-4">
<Label variant="small">Account ID</Label>
<Input
type="text"
fullWidth
value={currentAccountId}
placeholder={`e.g. abc_1234`}
onChange={(e) => setCurrentAccountId(e.target.value)}
/>
<FormError>{accountId.error}</FormError>
</InputGroup>
)}
<div className="flex flex-none items-center justify-between">
{payload.error ? (
<FormError id={payload.errorId}>{payload.error}</FormError>
@@ -39,9 +39,14 @@ export default function SetUpAstro() {
useProjectSetupComplete();
const devEnvironment = useDevEnvironment();
invariant(devEnvironment, "Dev environment must be defined");
const appOrigin = useAppOrigin();
return (
<PageGradient>
<div className="mx-auto max-w-3xl">
<div className="mb-12 grid place-items-center">
<AstroLogo className="w-64" />
</div>
<div className="flex items-center justify-between">
<Header1 spacing className="text-bright">
Get setup in 5 minutes
@@ -76,28 +81,16 @@ export default function SetUpAstro() {
<div>
<StepNumber
stepNumber="1"
title="Follow the steps from the Astro manual installation guide"
title="Run the CLI 'init' command in an existing Astro project"
/>
<StepContentContainer className="flex flex-col gap-2">
<Paragraph className="mt-2">Copy your server API Key to your clipboard:</Paragraph>
<div className="mb-2 flex w-full items-center justify-between">
<ClipboardField
secure
className="w-fit"
value={devEnvironment.apiKey}
variant={"secondary/medium"}
icon={<Badge variant="outline">Server</Badge>}
/>
</div>
<Paragraph>Now follow this guide:</Paragraph>
<LinkButton
to="https://trigger.dev/docs/documentation/guides/manual/astro"
variant="primary/medium"
TrailingIcon="external-link"
>
Manual installation guide
</LinkButton>
<div className="flex items-start justify-start gap-2"></div>
<StepContentContainer>
<InitCommand appOrigin={appOrigin} apiKey={devEnvironment.apiKey} />
<Paragraph spacing variant="small">
Youll notice a new folder in your project called 'jobs'. Weve added a very simple
example Job in <InlineCode variant="extra-small">example.ts</InlineCode> to help you
get started.
</Paragraph>
</StepContentContainer>
<StepNumber stepNumber="2" title="Run your Astro app" />
<StepContentContainer>
@@ -1,21 +1,120 @@
import { ChatBubbleLeftRightIcon, Squares2X2Icon } from "@heroicons/react/20/solid";
import invariant from "tiny-invariant";
import { ExpressLogo } from "~/assets/logos/ExpressLogo";
import { FrameworkComingSoon } from "~/components/frameworks/FrameworkComingSoon";
import { Feedback } from "~/components/Feedback";
import { PageGradient } from "~/components/PageGradient";
import { InitCommand, RunDevCommand, TriggerDevStep } from "~/components/SetupCommands";
import { StepContentContainer } from "~/components/StepContentContainer";
import { InlineCode } from "~/components/code/InlineCode";
import { BreadcrumbLink } from "~/components/navigation/NavBar";
import { Badge } from "~/components/primitives/Badge";
import { Button, LinkButton } from "~/components/primitives/Buttons";
import { Callout } from "~/components/primitives/Callout";
import { ClipboardField } from "~/components/primitives/ClipboardField";
import { Header1 } from "~/components/primitives/Headers";
import { Paragraph } from "~/components/primitives/Paragraph";
import { StepNumber } from "~/components/primitives/StepNumber";
import { useAppOrigin } from "~/hooks/useAppOrigin";
import { useDevEnvironment } from "~/hooks/useEnvironments";
import { useOrganization } from "~/hooks/useOrganizations";
import { useProject } from "~/hooks/useProject";
import { useProjectSetupComplete } from "~/hooks/useProjectSetupComplete";
import { Handle } from "~/utils/handle";
import { trimTrailingSlash } from "~/utils/pathBuilder";
import { projectSetupPath, trimTrailingSlash } from "~/utils/pathBuilder";
export const handle: Handle = {
breadcrumb: (match) => <BreadcrumbLink to={trimTrailingSlash(match.pathname)} title="Express" />,
};
export default function Page() {
const organization = useOrganization();
const project = useProject();
useProjectSetupComplete();
const devEnvironment = useDevEnvironment();
invariant(devEnvironment, "Dev environment must be defined");
const appOrigin = useAppOrigin();
return (
<FrameworkComingSoon
frameworkName="Express"
githubIssueUrl="https://github.com/triggerdotdev/trigger.dev/issues/451"
githubIssueNumber={451}
>
<ExpressLogo className="w-56" />
</FrameworkComingSoon>
<PageGradient>
<div className="mx-auto max-w-3xl">
<div className="mb-12 grid place-items-center">
<ExpressLogo className="w-64" />
</div>
<div className="flex items-center justify-between">
<Header1 spacing className="text-bright">
Get setup in 5 minutes
</Header1>
<div className="flex items-center gap-2">
<LinkButton
to={projectSetupPath(organization, project)}
variant="tertiary/small"
LeadingIcon={Squares2X2Icon}
>
Choose a different framework
</LinkButton>
<Feedback
button={
<Button variant="tertiary/small" LeadingIcon={ChatBubbleLeftRightIcon}>
I'm stuck!
</Button>
}
defaultValue="help"
/>
</div>
</div>
<div>
<Callout
variant={"info"}
to="https://github.com/triggerdotdev/trigger.dev/discussions/430"
className="mb-8"
>
Trigger.dev has full support for serverless. We will be adding support for long-running
servers soon.
</Callout>
<div>
<StepNumber
stepNumber="1"
title="Manually set up Trigger.dev in your existing Express project"
/>
<StepContentContainer className="flex flex-col gap-2">
<Paragraph className="mt-2">Copy your server API Key to your clipboard:</Paragraph>
<div className="mb-2 flex w-full items-center justify-between">
<ClipboardField
secure
className="w-fit"
value={devEnvironment.apiKey}
variant={"secondary/medium"}
icon={<Badge variant="outline">Server</Badge>}
/>
</div>
<Paragraph>Now follow this guide:</Paragraph>
<LinkButton
to="https://trigger.dev/docs/documentation/guides/manual/express"
variant="primary/medium"
TrailingIcon="external-link"
>
Manual installation guide
</LinkButton>
</StepContentContainer>
<StepNumber stepNumber="2" title="Run your Express app" />
<StepContentContainer>
<RunDevCommand />
<Callout variant="info">
You may be using the `start` script instead, in which case substitute `dev` in the
above commands.
</Callout>
</StepContentContainer>
<StepNumber stepNumber="3" title="Run the CLI 'dev' command" />
<StepContentContainer>
<TriggerDevStep />
</StepContentContainer>
<StepNumber stepNumber="6" title="Wait for Jobs" displaySpinner />
<StepContentContainer>
<Paragraph>This page will automatically refresh.</Paragraph>
</StepContentContainer>
</div>
</div>
</div>
</PageGradient>
);
}
@@ -28,6 +28,7 @@ import { useProject } from "~/hooks/useProject";
import { Handle } from "~/utils/handle";
import { projectSetupPath, trimTrailingSlash } from "~/utils/pathBuilder";
import { Callout } from "~/components/primitives/Callout";
import { NextjsLogo } from "~/assets/logos/NextjsLogo";
type SelectionChoices = "use-existing-project" | "create-new-next-app";
@@ -48,6 +49,9 @@ export default function SetupNextjs() {
return (
<PageGradient>
<div className="mx-auto max-w-3xl">
<div className="mb-12 grid place-items-center">
<NextjsLogo className="w-56" />
</div>
<div className="flex items-center justify-between">
<Header1 spacing className="text-bright">
Get setup in {selectedValue === "create-new-next-app" ? "5" : "2"} minutes
@@ -25,8 +25,9 @@ import { useProject } from "~/hooks/useProject";
import { Handle } from "~/utils/handle";
import { projectSetupPath, trimTrailingSlash } from "~/utils/pathBuilder";
import { Callout } from "~/components/primitives/Callout";
import { RunDevCommand, TriggerDevStep } from "~/components/SetupCommands";
import { InitCommand, RunDevCommand, TriggerDevStep } from "~/components/SetupCommands";
import { Badge } from "~/components/primitives/Badge";
import { RemixLogo } from "~/assets/logos/RemixLogo";
export const handle: Handle = {
breadcrumb: (match) => <BreadcrumbLink to={trimTrailingSlash(match.pathname)} title="Remix" />,
@@ -38,9 +39,14 @@ export default function SetUpRemix() {
useProjectSetupComplete();
const devEnvironment = useDevEnvironment();
invariant(devEnvironment, "Dev environment must be defined");
const appOrigin = useAppOrigin();
return (
<PageGradient>
<div className="mx-auto max-w-3xl">
<div className="mb-12 grid place-items-center">
<RemixLogo className="w-64" />
</div>
<div className="flex items-center justify-between">
<Header1 spacing className="text-bright">
Get setup in 5 minutes
@@ -75,28 +81,16 @@ export default function SetUpRemix() {
<div>
<StepNumber
stepNumber="1"
title="Follow the steps from the Remix manual installation guide"
title="Run the CLI 'init' command in an existing Remix project"
/>
<StepContentContainer className="flex flex-col gap-2">
<Paragraph className="mt-2">Copy your server API Key to your clipboard:</Paragraph>
<div className="mb-2 flex w-full items-center justify-between">
<ClipboardField
secure
className="w-fit"
value={devEnvironment.apiKey}
variant={"secondary/medium"}
icon={<Badge variant="outline">Server</Badge>}
/>
</div>
<Paragraph>Now follow this guide:</Paragraph>
<LinkButton
to="https://trigger.dev/docs/documentation/guides/manual/remix"
variant="primary/medium"
TrailingIcon="external-link"
>
Manual installation guide
</LinkButton>
<div className="flex items-start justify-start gap-2"></div>
<StepContentContainer>
<InitCommand appOrigin={appOrigin} apiKey={devEnvironment.apiKey} />
<Paragraph spacing variant="small">
Youll notice a new folder in your project called 'jobs'. Weve added a very simple
example Job in <InlineCode variant="extra-small">example.server.ts</InlineCode> to
help you get started.
</Paragraph>
</StepContentContainer>
<StepNumber stepNumber="2" title="Run your Remix app" />
<StepContentContainer>
@@ -81,13 +81,19 @@ class CreateExternalConnectionService {
environment: AuthenticatedEnvironment,
payload: CreateExternalConnectionBody
) {
const externalAccount = await this.#prismaClient.externalAccount.findUniqueOrThrow({
const externalAccount = await this.#prismaClient.externalAccount.upsert({
where: {
environmentId_identifier: {
environmentId: environment.id,
identifier: accountIdentifier,
},
},
create: {
environmentId: environment.id,
organizationId: environment.organizationId,
identifier: accountIdentifier,
},
update: {},
});
const integration = await this.#prismaClient.integration.findUniqueOrThrow({
@@ -0,0 +1,151 @@
import type { ActionArgs } from "@remix-run/server-runtime";
import { json } from "@remix-run/server-runtime";
import { TaskStatus } from "@trigger.dev/database";
import {
RunTaskBodyOutput,
RunTaskBodyOutputSchema,
ServerTask,
StatusHistory,
StatusHistorySchema,
StatusUpdate,
StatusUpdateData,
StatusUpdateSchema,
StatusUpdateState,
} from "@trigger.dev/core";
import { z } from "zod";
import { $transaction, PrismaClient, prisma } from "~/db.server";
import { taskWithAttemptsToServerTask } from "~/models/task.server";
import { authenticateApiRequest } from "~/services/apiAuth.server";
import { logger } from "~/services/logger.server";
import { ulid } from "~/services/ulid.server";
import { workerQueue } from "~/services/worker.server";
import { JobRunStatusRecordSchema } from "@trigger.dev/core";
const ParamsSchema = z.object({
runId: z.string(),
id: z.string(),
});
export async function action({ request, params }: ActionArgs) {
// Ensure this is a POST request
if (request.method.toUpperCase() !== "PUT") {
return { status: 405, body: "Method Not Allowed" };
}
// Next authenticate the request
const authenticationResult = await authenticateApiRequest(request);
if (!authenticationResult) {
return json({ error: "Invalid or Missing API key" }, { status: 401 });
}
const { runId, id } = ParamsSchema.parse(params);
// Now parse the request body
const anyBody = await request.json();
logger.debug("SetStatusService.call() request body", {
body: anyBody,
runId,
id,
});
const body = StatusUpdateSchema.safeParse(anyBody);
if (!body.success) {
return json({ error: "Invalid request body" }, { status: 400 });
}
const service = new SetStatusService();
try {
const statusRecord = await service.call(runId, id, body.data);
logger.debug("SetStatusService.call() response body", {
runId,
id,
statusRecord,
});
if (!statusRecord) {
return json({ error: "Something went wrong" }, { status: 500 });
}
const status = JobRunStatusRecordSchema.parse({
...statusRecord,
state: statusRecord.state ?? undefined,
history: statusRecord.history ?? undefined,
data: statusRecord.data ?? undefined,
});
return json(status);
} catch (error) {
if (error instanceof Error) {
return json({ error: error.message }, { status: 400 });
}
return json({ error: "Something went wrong" }, { status: 500 });
}
}
export class SetStatusService {
#prismaClient: PrismaClient;
constructor(prismaClient: PrismaClient = prisma) {
this.#prismaClient = prismaClient;
}
public async call(runId: string, id: string, status: StatusUpdate) {
const statusRecord = await $transaction(this.#prismaClient, async (tx) => {
const existingStatus = await tx.jobRunStatusRecord.findUnique({
where: {
runId_key: {
runId,
key: id,
},
},
});
const history: StatusHistory = [];
const historyResult = StatusHistorySchema.safeParse(existingStatus?.history);
if (historyResult.success) {
history.push(...historyResult.data);
}
if (existingStatus) {
history.push({
label: existingStatus.label,
state: (existingStatus.state ?? undefined) as StatusUpdateState,
data: (existingStatus.data ?? undefined) as StatusUpdateData,
});
}
const updatedStatus = await tx.jobRunStatusRecord.upsert({
where: {
runId_key: {
runId,
key: id,
},
},
create: {
key: id,
runId,
//this shouldn't ever use the id in reality, as the SDK makess it compulsory on the first call
label: status.label ?? id,
state: status.state,
data: status.data as any,
history: [],
},
update: {
label: status.label,
state: status.state,
data: status.data as any,
history: history as any[],
},
});
return updatedStatus;
});
return statusRecord;
}
}
@@ -0,0 +1,82 @@
import type { LoaderArgs } from "@remix-run/server-runtime";
import { json } from "@remix-run/server-runtime";
import { JobRunStatusRecordSchema } from "@trigger.dev/core";
import { z } from "zod";
import { prisma } from "~/db.server";
import { authenticateApiRequest } from "~/services/apiAuth.server";
import { logger } from "~/services/logger.server";
import { apiCors } from "~/utils/apiCors";
const ParamsSchema = z.object({
runId: z.string(),
});
const RecordsSchema = z.array(JobRunStatusRecordSchema);
export async function loader({ request, params }: LoaderArgs) {
if (request.method.toUpperCase() === "OPTIONS") {
return apiCors(request, json({}));
}
// Next authenticate the request
const authenticationResult = await authenticateApiRequest(request, { allowPublicKey: true });
if (!authenticationResult) {
return apiCors(request, json({ error: "Invalid or Missing API key" }, { status: 401 }));
}
const { runId } = ParamsSchema.parse(params);
logger.debug("Get run statuses", {
runId,
});
try {
const run = await prisma.jobRun.findUnique({
where: {
id: runId,
},
select: {
id: true,
status: true,
output: true,
statuses: {
orderBy: {
createdAt: "asc",
},
},
},
});
if (!run) {
return apiCors(request, json({ error: `No run found for id ${runId}` }, { status: 404 }));
}
const parsedStatuses = RecordsSchema.parse(
run.statuses.map((s) => ({
...s,
state: s.state ?? undefined,
data: s.data ?? undefined,
history: s.history ?? undefined,
}))
);
return apiCors(
request,
json({
run: {
id: run.id,
status: run.status,
output: run.output,
},
statuses: parsedStatuses,
})
);
} catch (error) {
if (error instanceof Error) {
return apiCors(request, json({ error: error.message }, { status: 400 }));
}
return apiCors(request, json({ error: "Something went wrong" }, { status: 500 }));
}
}
@@ -0,0 +1,124 @@
import type { ActionArgs } from "@remix-run/server-runtime";
import { json } from "@remix-run/server-runtime";
import { RuntimeEnvironmentType } from "@trigger.dev/database";
import { z } from "zod";
import { $transaction, PrismaClient, PrismaClientOrTransaction, prisma } from "~/db.server";
import { enqueueRunExecutionV2 } from "~/models/jobRunExecution.server";
import { logger } from "~/services/logger.server";
const ParamsSchema = z.object({
runId: z.string(),
id: z.string(),
secret: z.string(),
});
export async function action({ request, params }: ActionArgs) {
// Ensure this is a POST request
if (request.method.toUpperCase() !== "POST") {
return { status: 405, body: "Method Not Allowed" };
}
const { runId, id } = ParamsSchema.parse(params);
// Parse body as JSON (no schema parsing)
const body = await request.json();
const service = new CallbackRunTaskService();
try {
// Complete task with request body as output
await service.call(runId, id, body, request.url);
return json({ success: true });
} catch (error) {
if (error instanceof Error) {
logger.error("Error while processing task callback:", { error });
}
return json({ error: "Something went wrong" }, { status: 500 });
}
}
export class CallbackRunTaskService {
#prismaClient: PrismaClient;
constructor(prismaClient: PrismaClient = prisma) {
this.#prismaClient = prismaClient;
}
public async call(runId: string, id: string, taskBody: any, callbackUrl: string): Promise<void> {
const task = await findTask(prisma, id);
if (!task) {
return;
}
if (task.runId !== runId) {
return;
}
if (task.status !== "WAITING") {
return;
}
if (!task.callbackUrl) {
return;
}
if (new URL(task.callbackUrl).pathname !== new URL(callbackUrl).pathname) {
logger.error("Callback URLs don't match", { runId, taskId: id, callbackUrl });
return;
}
logger.debug("CallbackRunTaskService.call()", { task });
await this.#resumeTask(task, taskBody);
}
async #resumeTask(task: NonNullable<FoundTask>, output: any) {
await $transaction(this.#prismaClient, async (tx) => {
await tx.taskAttempt.updateMany({
where: {
taskId: task.id,
status: "PENDING",
},
data: {
status: "COMPLETED",
},
});
await tx.task.update({
where: { id: task.id },
data: {
status: "COMPLETED",
completedAt: new Date(),
output: output ? output : undefined,
},
});
await this.#resumeRunExecution(task, tx);
});
}
async #resumeRunExecution(task: NonNullable<FoundTask>, prisma: PrismaClientOrTransaction) {
await enqueueRunExecutionV2(task.run, prisma, {
skipRetrying: task.run.environment.type === RuntimeEnvironmentType.DEVELOPMENT,
});
}
}
type FoundTask = Awaited<ReturnType<typeof findTask>>;
async function findTask(prisma: PrismaClientOrTransaction, id: string) {
return prisma.task.findUnique({
where: { id },
include: {
run: {
include: {
environment: true,
queue: true,
},
},
},
});
}
@@ -3,7 +3,7 @@ import { json } from "@remix-run/server-runtime";
import type { CompleteTaskBodyOutput, ServerTask } from "@trigger.dev/core";
import { CompleteTaskBodyInputSchema } from "@trigger.dev/core";
import { z } from "zod";
import { $transaction, PrismaClient, prisma } from "~/db.server";
import { PrismaClient, prisma } from "~/db.server";
import { taskWithAttemptsToServerTask } from "~/models/task.server";
import type { AuthenticatedEnvironment } from "~/services/apiAuth.server";
import { authenticateApiRequest } from "~/services/apiAuth.server";
@@ -86,8 +86,8 @@ export class CompleteRunTaskService {
): Promise<ServerTask | undefined> {
// Using a transaction, we'll first check to see if the task already exists and return if if it does
// If it doesn't exist, we'll create it and return it
const task = await this.#prismaClient.$transaction(async (prisma) => {
const existingTask = await prisma.task.findUnique({
const task = await this.#prismaClient.$transaction(async (tx) => {
const existingTask = await tx.task.findUnique({
where: {
id,
},
@@ -129,35 +129,31 @@ export class CompleteRunTaskService {
return existingTask;
}
const task = await $transaction(prisma, async (tx) => {
if (existingTask.attempts.length === 1) {
await tx.taskAttempt.update({
where: {
id: existingTask.attempts[0].id,
},
data: {
status: "COMPLETED",
},
});
}
return await tx.task.update({
if (existingTask.attempts.length === 1) {
await tx.taskAttempt.update({
where: {
id,
id: existingTask.attempts[0].id,
},
data: {
status: "COMPLETED",
output: taskBody.output ?? undefined,
completedAt: new Date(),
outputProperties: taskBody.properties,
},
include: {
attempts: true,
},
});
});
}
return task;
return await tx.task.update({
where: {
id,
},
data: {
status: "COMPLETED",
output: taskBody.output ?? undefined,
completedAt: new Date(),
outputProperties: taskBody.properties,
},
include: {
attempts: true,
},
});
});
return task ? taskWithAttemptsToServerTask(task) : undefined;
@@ -2,7 +2,7 @@ import type { ActionArgs } from "@remix-run/server-runtime";
import { json } from "@remix-run/server-runtime";
import { FailTaskBodyInput, FailTaskBodyInputSchema, ServerTask } from "@trigger.dev/core";
import { z } from "zod";
import { $transaction, PrismaClient, prisma } from "~/db.server";
import { PrismaClient, prisma } from "~/db.server";
import { taskWithAttemptsToServerTask } from "~/models/task.server";
import type { AuthenticatedEnvironment } from "~/services/apiAuth.server";
import { authenticateApiRequest } from "~/services/apiAuth.server";
@@ -86,8 +86,8 @@ export class FailRunTaskService {
): Promise<ServerTask | undefined> {
// Using a transaction, we'll first check to see if the task already exists and return if if it does
// If it doesn't exist, we'll create it and return it
const task = await this.#prismaClient.$transaction(async (prisma) => {
const existingTask = await prisma.task.findUnique({
const task = await this.#prismaClient.$transaction(async (tx) => {
const existingTask = await tx.task.findUnique({
where: {
id,
},
@@ -129,35 +129,31 @@ export class FailRunTaskService {
return existingTask;
}
const task = await $transaction(prisma, async (tx) => {
if (existingTask.attempts.length === 1) {
await tx.taskAttempt.update({
where: {
id: existingTask.attempts[0].id,
},
data: {
status: "ERRORED",
error: formatError(taskBody.error),
},
});
}
return await prisma.task.update({
if (existingTask.attempts.length === 1) {
await tx.taskAttempt.update({
where: {
id,
id: existingTask.attempts[0].id,
},
data: {
status: "ERRORED",
output: taskBody.error ?? undefined,
completedAt: new Date(),
},
include: {
attempts: true,
error: formatError(taskBody.error),
},
});
});
}
return task;
return await tx.task.update({
where: {
id,
},
data: {
status: "ERRORED",
output: taskBody.error ?? undefined,
completedAt: new Date(),
},
include: {
attempts: true,
},
});
});
return task ? taskWithAttemptsToServerTask(task) : undefined;
@@ -1,14 +1,22 @@
import type { ActionArgs } from "@remix-run/server-runtime";
import { json } from "@remix-run/server-runtime";
import { TaskStatus } from "@trigger.dev/database";
import { RunTaskBodyOutput, RunTaskBodyOutputSchema, ServerTask } from "@trigger.dev/core";
import {
API_VERSIONS,
RunTaskBodyOutput,
RunTaskBodyOutputSchema,
RunTaskResponseWithCachedTasksBody,
ServerTask,
} from "@trigger.dev/core";
import { z } from "zod";
import { $transaction, PrismaClient, prisma } from "~/db.server";
import { taskWithAttemptsToServerTask } from "~/models/task.server";
import { prepareTasksForCaching, taskWithAttemptsToServerTask } from "~/models/task.server";
import { authenticateApiRequest } from "~/services/apiAuth.server";
import { logger } from "~/services/logger.server";
import { ulid } from "~/services/ulid.server";
import { workerQueue } from "~/services/worker.server";
import { generateSecret } from "~/services/sources/utils.server";
import { env } from "~/env.server";
const ParamsSchema = z.object({
runId: z.string(),
@@ -16,6 +24,8 @@ const ParamsSchema = z.object({
const HeadersSchema = z.object({
"idempotency-key": z.string(),
"trigger-version": z.string().optional().nullable(),
"x-cached-tasks-cursor": z.string().optional().nullable(),
});
export async function action({ request, params }: ActionArgs) {
@@ -37,7 +47,11 @@ export async function action({ request, params }: ActionArgs) {
return json({ error: "Invalid or Missing idempotency key" }, { status: 400 });
}
const { "idempotency-key": idempotencyKey } = headers.data;
const {
"idempotency-key": idempotencyKey,
"trigger-version": triggerVersion,
"x-cached-tasks-cursor": cachedTasksCursor,
} = headers.data;
const { runId } = ParamsSchema.parse(params);
@@ -48,6 +62,8 @@ export async function action({ request, params }: ActionArgs) {
body: anyBody,
runId,
idempotencyKey,
triggerVersion,
cachedTasksCursor,
});
const body = RunTaskBodyOutputSchema.safeParse(anyBody);
@@ -71,6 +87,26 @@ export async function action({ request, params }: ActionArgs) {
return json({ error: "Something went wrong" }, { status: 500 });
}
if (triggerVersion === API_VERSIONS.LAZY_LOADED_CACHED_TASKS) {
const requestMigration = new ChangeRequestLazyLoadedCachedTasks();
const responseBody = await requestMigration.call(runId, task, cachedTasksCursor);
logger.debug(
"RunTaskService.call() response migrating with ChangeRequestLazyLoadedCachedTasks",
{
responseBody,
cachedTasksCursor,
}
);
return json(responseBody, {
headers: {
"trigger-version": API_VERSIONS.LAZY_LOADED_CACHED_TASKS,
},
});
}
return json(task);
} catch (error) {
if (error instanceof Error) {
@@ -81,6 +117,51 @@ export async function action({ request, params }: ActionArgs) {
}
}
class ChangeRequestLazyLoadedCachedTasks {
#prismaClient: PrismaClient;
constructor(prismaClient: PrismaClient = prisma) {
this.#prismaClient = prismaClient;
}
public async call(
runId: string,
task: ServerTask,
cursor?: string | null
): Promise<RunTaskResponseWithCachedTasksBody> {
if (!cursor) {
return {
task,
};
}
// We need to limit the cached tasks to not be too large >2MB when serialized
const TOTAL_CACHED_TASK_BYTE_LIMIT = 2000000;
const nextTasks = await this.#prismaClient.task.findMany({
where: {
runId,
status: "COMPLETED",
noop: false,
},
take: 250,
cursor: {
id: cursor,
},
orderBy: {
id: "asc",
},
});
const preparedTasks = prepareTasksForCaching(nextTasks, TOTAL_CACHED_TASK_BYTE_LIMIT);
return {
task,
cachedTasks: preparedTasks,
};
}
}
export class RunTaskService {
#prismaClient: PrismaClient;
@@ -106,10 +187,13 @@ export class RunTaskService {
},
});
const delayUntilInFuture = taskBody.delayUntil && taskBody.delayUntil.getTime() > Date.now();
const callbackEnabled = taskBody.callback?.enabled;
if (existingTask) {
if (existingTask.status === "CANCELED") {
const existingTaskStatus =
(taskBody.delayUntil && taskBody.delayUntil.getTime() > Date.now()) || taskBody.trigger
delayUntilInFuture || callbackEnabled || taskBody.trigger
? "WAITING"
: taskBody.noop
? "COMPLETED"
@@ -154,16 +238,21 @@ export class RunTaskService {
status = "CANCELED";
} else {
status =
(taskBody.delayUntil && taskBody.delayUntil.getTime() > Date.now()) || taskBody.trigger
delayUntilInFuture || callbackEnabled || taskBody.trigger
? "WAITING"
: taskBody.noop
? "COMPLETED"
: "RUNNING";
}
const taskId = ulid();
const callbackUrl = callbackEnabled
? `${env.APP_ORIGIN}/api/v1/runs/${runId}/tasks/${taskId}/callback/${generateSecret(12)}`
: undefined;
const task = await tx.task.create({
data: {
id: ulid(),
id: taskId,
idempotencyKey,
displayKey: taskBody.displayKey,
runConnection: taskBody.connectionKey
@@ -194,6 +283,7 @@ export class RunTaskService {
properties: taskBody.properties ?? undefined,
redact: taskBody.redact ?? undefined,
operation: taskBody.operation,
callbackUrl,
style: taskBody.style ?? { style: "normal" },
attempts: {
create: {
@@ -217,6 +307,17 @@ export class RunTaskService {
},
{ tx, runAt: task.delayUntil ?? undefined }
);
} else if (task.status === "WAITING" && callbackUrl && taskBody.callback) {
if (taskBody.callback.timeoutInSeconds > 0) {
// We need to schedule the callback timeout
await workerQueue.enqueue(
"processCallbackTimeout",
{
id: task.id,
},
{ tx, runAt: new Date(Date.now() + taskBody.callback.timeoutInSeconds * 1000) }
);
}
}
return task;
+10 -2
View File
@@ -1,9 +1,8 @@
import type { LoaderArgs } from "@remix-run/server-runtime";
import { json } from "@remix-run/server-runtime";
import { cors } from "remix-utils";
import { z } from "zod";
import { prisma } from "~/db.server";
import { authenticateApiRequest, getApiKeyFromRequest } from "~/services/apiAuth.server";
import { authenticateApiRequest } from "~/services/apiAuth.server";
import { apiCors } from "~/utils/apiCors";
import { taskListToTree } from "~/utils/taskListToTree";
@@ -93,6 +92,9 @@ export async function loader({ request, params }: LoaderArgs) {
}
: undefined,
},
statuses: {
select: { key: true, label: true, state: true, data: true, history: true },
},
},
});
@@ -122,6 +124,12 @@ export async function loader({ request, params }: LoaderArgs) {
const { parentId, ...rest } = task;
return { ...rest };
}),
statuses: jobRun.statuses.map((s) => ({
...s,
state: s.state ?? undefined,
data: s.data ?? undefined,
history: s.history ?? undefined,
})),
nextCursor: nextTask ? nextTask.id : undefined,
})
);
@@ -2,28 +2,15 @@ import { parse } from "@conform-to/zod";
import { ActionArgs, json } from "@remix-run/server-runtime";
import { z } from "zod";
import { prisma } from "~/db.server";
import {
CreateEndpointError,
CreateEndpointService,
} from "~/services/endpoints/createEndpoint.server";
import { requireUserId } from "~/services/session.server";
import { RuntimeEnvironmentTypeSchema } from "@trigger.dev/core";
import { env } from "process";
import { CreateEndpointError } from "~/services/endpoints/createEndpoint.server";
import { ValidateCreateEndpointService } from "~/services/endpoints/validateCreateEndpoint.server";
const ParamsSchema = z.object({
projectId: z.string(),
});
export const bodySchema = z.object({
environmentId: z.string(),
url: z.string().url("Must be a valid URL"),
});
export async function action({ request, params }: ActionArgs) {
const userId = await requireUserId(request);
const { projectId } = ParamsSchema.parse(params);
export async function action({ request }: ActionArgs) {
const formData = await request.formData();
const submission = parse(formData, { schema: bodySchema });
@@ -48,7 +35,7 @@ export async function action({ request, params }: ActionArgs) {
}
const service = new ValidateCreateEndpointService();
const result = await service.call({
await service.call({
url: submission.value.url,
environment,
});
+23 -28
View File
@@ -1,7 +1,9 @@
import {
API_VERSIONS,
ApiEventLog,
DeliverEventResponseSchema,
DeserializedJson,
EndpointHeadersSchema,
ErrorWithStackSchema,
HttpSourceRequest,
HttpSourceResponseSchema,
@@ -89,6 +91,15 @@ export class EndpointApi {
};
}
const headers = EndpointHeadersSchema.safeParse(Object.fromEntries(response.headers.entries()));
if (headers.success && headers.data["trigger-version"]) {
return {
...pongResponse.data,
triggerVersion: headers.data["trigger-version"],
};
}
return pongResponse.data;
}
@@ -129,41 +140,15 @@ export class EndpointApi {
const anyBody = await response.json();
const data = IndexEndpointResponseSchema.parse(anyBody);
const headers = EndpointHeadersSchema.parse(Object.fromEntries(response.headers.entries()));
return {
ok: true,
data,
headers,
} as const;
}
async deliverEvent(event: ApiEventLog) {
const response = await safeFetch(this.url, {
method: "POST",
headers: {
"Content-Type": "application/json",
"x-trigger-api-key": this.apiKey,
"x-trigger-action": "DELIVER_EVENT",
},
body: JSON.stringify(event),
});
if (!response) {
throw new Error(`Could not connect to endpoint ${this.url}`);
}
if (!response.ok) {
throw new Error(`Could not connect to endpoint ${this.url}. Status code: ${response.status}`);
}
const anyBody = await response.json();
logger.debug("deliverEvent() response from endpoint", {
body: anyBody,
});
return DeliverEventResponseSchema.parse(anyBody);
}
async executeJobRequest(options: RunJobBody) {
const startTimeInMs = performance.now();
@@ -338,6 +323,15 @@ export class EndpointApi {
};
}
const headers = EndpointHeadersSchema.safeParse(Object.fromEntries(response.headers.entries()));
if (headers.success && headers.data["trigger-version"]) {
return {
...validateResponse.data,
triggerVersion: headers.data["trigger-version"],
};
}
return validateResponse.data;
}
}
@@ -359,6 +353,7 @@ function addStandardRequestOptions(options: RequestInit) {
headers: {
...options.headers,
"user-agent": "triggerdotdev-server/2.0.0",
"x-trigger-version": API_VERSIONS.LAZY_LOADED_CACHED_TASKS,
},
};
}
@@ -74,9 +74,11 @@ export class CreateEndpointService {
slug: id,
url: endpointUrl,
indexingHookIdentifier: indexingHookIdentifier(),
version: pong.triggerVersion,
},
update: {
url: endpointUrl,
version: pong.triggerVersion,
},
});
@@ -41,6 +41,7 @@ export class IndexEndpointService {
}
const { jobs, sources, dynamicTriggers, dynamicSchedules } = indexResponse.data;
const { "trigger-version": triggerVersion } = indexResponse.headers;
logger.debug("Indexing endpoint", {
endpointId: endpoint.id,
@@ -48,6 +49,7 @@ export class IndexEndpointService {
endpointSlug: endpoint.slug,
source: source,
sourceData: sourceData,
triggerVersion,
stats: {
jobs: jobs.length,
sources: sources.length,
@@ -56,6 +58,17 @@ export class IndexEndpointService {
},
});
if (triggerVersion && triggerVersion !== endpoint.version) {
await this.#prismaClient.endpoint.update({
where: {
id: endpoint.id,
},
data: {
version: triggerVersion,
},
});
}
const indexStats = {
jobs: 0,
sources: 0,
@@ -58,9 +58,11 @@ export class ValidateCreateEndpointService {
slug: validationResult.endpointId,
url: endpointUrl,
indexingHookIdentifier: indexingHookIdentifier(),
version: validationResult.triggerVersion,
},
update: {
url: endpointUrl,
version: validationResult.triggerVersion,
},
});
@@ -7,10 +7,7 @@ import { logger } from "../logger.server";
export class IngestSendEvent {
#prismaClient: PrismaClientOrTransaction;
constructor(
prismaClient: PrismaClientOrTransaction = prisma,
private deliverEvents = true
) {
constructor(prismaClient: PrismaClientOrTransaction = prisma, private deliverEvents = true) {
this.#prismaClient = prismaClient;
}
@@ -37,71 +34,55 @@ export class IngestSendEvent {
try {
const deliverAt = this.#calculateDeliverAt(options);
return await $transaction(
this.#prismaClient,
async (tx) => {
const externalAccount = options?.accountId
? await tx.externalAccount.findUniqueOrThrow({
where: {
environmentId_identifier: {
environmentId: environment.id,
identifier: options.accountId,
},
return await $transaction(this.#prismaClient, async (tx) => {
const externalAccount = options?.accountId
? await tx.externalAccount.upsert({
where: {
environmentId_identifier: {
environmentId: environment.id,
identifier: options.accountId,
},
})
: undefined;
},
create: {
environmentId: environment.id,
organizationId: environment.organizationId,
identifier: options.accountId,
},
update: {},
})
: undefined;
// Create a new event in the database
const eventLog = await tx.eventRecord.create({
data: {
organization: {
connect: {
id: environment.organizationId,
},
},
project: {
connect: {
id: environment.projectId,
},
},
environment: {
connect: {
id: environment.id,
},
},
eventId: event.id,
name: event.name,
timestamp: event.timestamp ?? new Date(),
payload: event.payload ?? {},
context: event.context ?? {},
source: event.source ?? "trigger.dev",
sourceContext,
deliverAt: deliverAt,
externalAccount: externalAccount
? {
connect: {
id: externalAccount.id,
},
}
: {},
// Create a new event in the database
const eventLog = await tx.eventRecord.create({
data: {
organizationId: environment.organizationId,
projectId: environment.projectId,
environmentId: environment.id,
eventId: event.id,
name: event.name,
timestamp: event.timestamp ?? new Date(),
payload: event.payload ?? {},
context: event.context ?? {},
source: event.source ?? "trigger.dev",
sourceContext,
deliverAt: deliverAt,
externalAccountId: externalAccount ? externalAccount.id : undefined,
},
});
if (this.deliverEvents) {
// Produce a message to the event bus
await workerQueue.enqueue(
"deliverEvent",
{
id: eventLog.id,
},
});
{ runAt: eventLog.deliverAt, tx, jobKey: `event:${eventLog.id}` }
);
}
if (this.deliverEvents) {
// Produce a message to the event bus
await workerQueue.enqueue(
"deliverEvent",
{
id: eventLog.id,
},
{ runAt: eventLog.deliverAt, tx, jobKey: `event:${eventLog.id}` }
);
}
return eventLog;
},
{ rethrowPrismaErrors: true }
);
return eventLog;
});
} catch (error) {
const prismaError = PrismaErrorSchema.safeParse(error);
@@ -1,7 +1,9 @@
import { airtable } from "./integrations/airtable";
import { github } from "./integrations/github";
import { linear } from "./integrations/linear";
import { openai } from "./integrations/openai";
import { plain } from "./integrations/plain";
import { replicate } from "./integrations/replicate";
import { resend } from "./integrations/resend";
import { sendgrid } from "./integrations/sendgrid";
import { slack } from "./integrations/slack";
@@ -33,8 +35,10 @@ export class IntegrationCatalog {
export const integrationCatalog = new IntegrationCatalog({
airtable,
github,
linear,
openai,
plain,
replicate,
resend,
slack,
stripe,
@@ -0,0 +1,109 @@
import type { HelpSample, Integration } from "../types";
function usageSample(hasApiKey: boolean): HelpSample {
return {
title: "Using the client",
code: `
import { Linear } from "@trigger.dev/linear";
const linear = new Linear({
id: "__SLUG__",${hasApiKey ? ",\n apiKey: process.env.LINEAR_API_KEY!" : ""}
});
client.defineJob({
id: "linear-react-to-new-issue",
name: "Linear - React To New Issue",
version: "0.1.0",
integrations: { linear },
trigger: linear.onIssueCreated(),
run: async (payload, io, ctx) => {
await io.linear.createComment("create-comment", {
issueId: payload.data.id,
body: "Thank's for opening this issue!"
});
await io.linear.createReaction("create-reaction", {
issueId: payload.data.id,
emoji: "+1"
});
return { payload, ctx };
},
});
`,
};
}
export const linear: Integration = {
identifier: "linear",
name: "Linear",
packageName: "@trigger.dev/linear@latest",
authenticationMethods: {
oauth2: {
name: "OAuth",
type: "oauth2",
client: {
id: {
envName: "CLOUD_LINEAR_CLIENT_ID",
},
secret: {
envName: "CLOUD_LINEAR_CLIENT_SECRET",
},
},
config: {
authorization: {
url: "https://linear.app/oauth/authorize",
scopeSeparator: ",",
},
token: {
url: "https://api.linear.app/oauth/token",
metadata: {},
},
refresh: {
url: "https://linear.app/oauth/authorize",
},
pkce: false,
},
scopes: [
{
name: "read",
description: "Read access for the user's account. This scope must always be present.",
defaultChecked: true,
},
{
name: "write",
description:
"Grants global write access to the user's account. Use a more targeted scope if you don't need full access.",
defaultChecked: true,
},
{
name: "issues:create",
description: "Grants access to create issues and attachments only.",
annotations: [{ label: "Issues" }],
},
{
name: "comments:create",
description: "Grants access to create new issue comments.",
annotations: [{ label: "Comments" }],
},
{
name: "admin",
description:
"Grants full access to admin-level endpoints. Don't use this unless you really need it.",
},
],
help: {
samples: [usageSample(false)],
},
},
apikey: {
type: "apikey",
help: {
samples: [usageSample(true)],
},
},
},
};
@@ -0,0 +1,50 @@
import type { HelpSample, Integration } from "../types";
function usageSample(hasApiKey: boolean): HelpSample {
const apiKeyPropertyName = "apiKey";
return {
title: "Using the client",
code: `
import { Replicate } from "@trigger.dev/replicate";
const replicate = new Replicate({
id: "__SLUG__",${hasApiKey ? `,\n ${apiKeyPropertyName}: process.env.REPLICATE_API_KEY!` : ""}
});
client.defineJob({
id: "replicate-create-prediction",
name: "Replicate - Create Prediction",
version: "0.1.0",
integrations: { replicate },
trigger: eventTrigger({
name: "replicate.predict",
schema: z.object({
prompt: z.string(),
version: z.string(),
}),
}),
run: async (payload, io, ctx) => {
return io.replicate.predictions.createAndAwait("await-prediction", {
version: payload.version,
input: { prompt: payload.prompt },
});
},
});
`,
};
}
export const replicate: Integration = {
identifier: "replicate",
name: "Replicate",
packageName: "@trigger.dev/replicate@latest",
authenticationMethods: {
apikey: {
type: "apikey",
help: {
samples: [usageSample(true)],
},
},
},
};
@@ -4,7 +4,14 @@ import {
SCHEDULED_EVENT,
TriggerMetadata,
} from "@trigger.dev/core";
import type { Endpoint, Integration, Job, JobIntegration, JobVersion } from "@trigger.dev/database";
import type {
Endpoint,
Integration,
Job,
JobIntegration,
JobIntegrationPayload,
JobVersion,
} from "@trigger.dev/database";
import { DEFAULT_MAX_CONCURRENT_RUNS } from "~/consts";
import type { PrismaClient } from "~/db.server";
import { prisma } from "~/db.server";
@@ -62,83 +69,7 @@ export class RegisterJobService {
});
if (!integration) {
if (jobIntegration.authSource === "LOCAL") {
integration = await this.#prismaClient.integration.upsert({
where: {
organizationId_slug: {
organizationId: environment.organizationId,
slug: jobIntegration.id,
},
},
create: {
slug: jobIntegration.id,
title: jobIntegration.metadata.name,
authSource: "LOCAL",
connectionType: "DEVELOPER",
organization: {
connect: {
id: environment.organizationId,
},
},
definition: {
connectOrCreate: {
where: {
id: jobIntegration.metadata.id,
},
create: {
id: jobIntegration.metadata.id,
name: jobIntegration.metadata.name,
instructions: jobIntegration.metadata.instructions,
},
},
},
},
update: {
title: jobIntegration.metadata.name,
authSource: "LOCAL",
connectionType: "DEVELOPER",
definition: {
connectOrCreate: {
where: {
id: jobIntegration.metadata.id,
},
create: {
id: jobIntegration.metadata.id,
name: jobIntegration.metadata.name,
instructions: jobIntegration.metadata.instructions,
},
},
},
},
});
} else {
integration = await this.#prismaClient.integration.create({
data: {
slug: jobIntegration.id,
title: jobIntegration.id,
authSource: "HOSTED",
setupStatus: "MISSING_FIELDS",
connectionType: "DEVELOPER",
organization: {
connect: {
id: environment.organizationId,
},
},
definition: {
connectOrCreate: {
where: {
id: jobIntegration.metadata.id,
},
create: {
id: jobIntegration.metadata.id,
name: jobIntegration.metadata.name,
instructions: jobIntegration.metadata.instructions,
},
},
},
},
});
}
integration = await this.#upsertIntegrationForJobIntegration(environment, jobIntegration);
}
integrations.set(jobIntegration.id, integration);
@@ -472,6 +403,7 @@ export class RegisterJobService {
key: job.id,
dispatcher: eventDispatcher,
schedule: trigger.schedule,
organizationId: job.organizationId,
});
break;
@@ -479,6 +411,145 @@ export class RegisterJobService {
}
}
async #upsertIntegrationForJobIntegration(
environment: AuthenticatedEnvironment,
jobIntegration: IntegrationConfig
): Promise<Integration> {
switch (jobIntegration.authSource) {
case "LOCAL": {
return await this.#prismaClient.integration.upsert({
where: {
organizationId_slug: {
organizationId: environment.organizationId,
slug: jobIntegration.id,
},
},
create: {
slug: jobIntegration.id,
title: jobIntegration.metadata.name,
authSource: "LOCAL",
connectionType: "DEVELOPER",
organization: {
connect: {
id: environment.organizationId,
},
},
definition: {
connectOrCreate: {
where: {
id: jobIntegration.metadata.id,
},
create: {
id: jobIntegration.metadata.id,
name: jobIntegration.metadata.name,
instructions: jobIntegration.metadata.instructions,
},
},
},
},
update: {
title: jobIntegration.metadata.name,
authSource: "LOCAL",
connectionType: "DEVELOPER",
definition: {
connectOrCreate: {
where: {
id: jobIntegration.metadata.id,
},
create: {
id: jobIntegration.metadata.id,
name: jobIntegration.metadata.name,
instructions: jobIntegration.metadata.instructions,
},
},
},
},
});
}
case "HOSTED": {
return await this.#prismaClient.integration.create({
data: {
slug: jobIntegration.id,
title: jobIntegration.id,
authSource: "HOSTED",
setupStatus: "MISSING_FIELDS",
connectionType: "DEVELOPER",
organization: {
connect: {
id: environment.organizationId,
},
},
definition: {
connectOrCreate: {
where: {
id: jobIntegration.metadata.id,
},
create: {
id: jobIntegration.metadata.id,
name: jobIntegration.metadata.name,
instructions: jobIntegration.metadata.instructions,
},
},
},
},
});
}
case "RESOLVER": {
return await this.#prismaClient.integration.upsert({
where: {
organizationId_slug: {
organizationId: environment.organizationId,
slug: jobIntegration.id,
},
},
create: {
slug: jobIntegration.id,
title: jobIntegration.metadata.name,
authSource: "RESOLVER",
connectionType: "EXTERNAL",
organization: {
connect: {
id: environment.organizationId,
},
},
definition: {
connectOrCreate: {
where: {
id: jobIntegration.metadata.id,
},
create: {
id: jobIntegration.metadata.id,
name: jobIntegration.metadata.name,
instructions: jobIntegration.metadata.instructions,
},
},
},
},
update: {
title: jobIntegration.metadata.name,
authSource: "RESOLVER",
connectionType: "EXTERNAL",
definition: {
connectOrCreate: {
where: {
id: jobIntegration.metadata.id,
},
create: {
id: jobIntegration.metadata.id,
name: jobIntegration.metadata.name,
instructions: jobIntegration.metadata.instructions,
},
},
},
},
});
}
default: {
assertExhaustive(jobIntegration.authSource);
}
}
}
async #upsertJobIntegration(
job: Job & {
integrations: Array<JobIntegration & { integration: Integration | null }>;
@@ -572,3 +643,7 @@ export class RegisterJobService {
});
}
}
function assertExhaustive(x: never): never {
throw new Error("Unexpected object: " + x);
}
@@ -13,10 +13,12 @@ export class TestJobService {
environmentId,
versionId,
payload,
accountId,
}: {
environmentId: string;
versionId: string;
payload: any;
payload?: any;
accountId?: string;
}) {
return await $transaction(
this.#prismaClient,
@@ -41,10 +43,27 @@ export class TestJobService {
},
});
const externalAccount = accountId
? await tx.externalAccount.upsert({
where: {
environmentId_identifier: {
environmentId: environment.id,
identifier: accountId,
},
},
create: {
environmentId: environment.id,
organizationId: environment.organizationId,
identifier: accountId,
},
update: {},
})
: undefined;
const event = EventSpecificationSchema.parse(version.eventSpecification);
const eventName = Array.isArray(event.name) ? event.name[0] : event.name;
const eventLog = await this.#prismaClient.eventRecord.create({
const eventLog = await tx.eventRecord.create({
data: {
organization: {
connect: {
@@ -61,6 +80,13 @@ export class TestJobService {
id: environment.id,
},
},
externalAccount: externalAccount
? {
connect: {
id: externalAccount.id,
},
}
: undefined,
eventId: `test:${eventName}:${new Date().getTime()}`,
name: eventName,
timestamp: new Date(),
@@ -2,7 +2,7 @@ import { RuntimeEnvironmentType } from "@trigger.dev/database";
import { $transaction, Prisma, PrismaClient, prisma } from "~/db.server";
import { enqueueRunExecutionV2 } from "~/models/jobRunExecution.server";
const RESUMABLE_STATUSES = ["FAILURE", "TIMED_OUT", "ABORTED", "CANCELED"];
const RESUMABLE_STATUSES = ["FAILURE", "TIMED_OUT", "UNRESOLVED_AUTH", "ABORTED", "CANCELED"];
export class ContinueRunService {
#prismaClient: PrismaClient;
@@ -42,29 +42,32 @@ export class CreateRunService {
return await $transaction(this.#prismaClient, async (tx) => {
// Get the current max number for the given jobId
const currentMaxNumber = await tx.jobRun.aggregate({
const latestJob = await tx.jobRun.findFirst({
where: { jobId: job.id },
_max: { number: true },
orderBy: { id: "desc" },
select: {
number: true,
},
});
// Increment the number for the new execution
const newNumber = (currentMaxNumber._max.number ?? 0) + 1;
const newNumber = (latestJob?.number ?? 0) + 1;
// Create the new execution with the incremented number
const run = await tx.jobRun.create({
data: {
number: newNumber,
preprocess: version.preprocessRuns,
job: { connect: { id: job.id } },
version: { connect: { id: version.id } },
event: { connect: { id: eventId } },
environment: { connect: { id: environment.id } },
organization: { connect: { id: environment.organizationId } },
project: { connect: { id: environment.projectId } },
endpoint: { connect: { id: endpoint.id } },
queue: { connect: { id: jobQueue.id } },
externalAccount: eventRecord.externalAccountId
? { connect: { id: eventRecord.externalAccountId } }
jobId: job.id,
versionId: version.id,
eventId: eventId,
environmentId: environment.id,
organizationId: environment.organizationId,
projectId: environment.projectId,
endpointId: endpoint.id,
queueId: jobQueue.id,
externalAccountId: eventRecord.externalAccountId
? eventRecord.externalAccountId
: undefined,
isTest: eventRecord.isTest,
},
@@ -1,9 +1,11 @@
import {
CachedTaskSchema,
RunJobError,
RunJobInvalidPayloadError,
RunJobResumeWithTask,
RunJobRetryWithTask,
RunJobSuccess,
RunJobUnresolvedAuthError,
RunSourceContextSchema,
} from "@trigger.dev/core";
import type { Task } from "@trigger.dev/database";
@@ -261,6 +263,7 @@ export class PerformRunExecutionV1Service {
.flat()
.filter(Boolean)
.map((t) => CachedTaskSchema.parse(t)),
yieldedExecutions: run.yieldedExecutions,
});
if (!response) {
@@ -342,6 +345,21 @@ export class PerformRunExecutionV1Service {
await this.#cancelExecution(execution);
break;
}
case "UNRESOLVED_AUTH_ERROR": {
await this.#failRunWithUnresolvedAuthError(execution, safeBody.data);
break;
}
case "INVALID_PAYLOAD": {
await this.#failRunWithInvalidPayloadError(execution, safeBody.data);
break;
}
case "YIELD_EXECUTION": {
await this.#resumeYieldedExecution(execution, safeBody.data.key);
break;
}
default: {
const _exhaustiveCheck: never = status;
throw new Error(`Non-exhaustive match for value: ${status}`);
@@ -381,6 +399,40 @@ export class PerformRunExecutionV1Service {
});
}
async #resumeYieldedExecution(execution: FoundRunExecution, key: string) {
const { run } = execution;
return await $transaction(this.#prismaClient, async (tx) => {
await tx.jobRunExecution.update({
where: {
id: execution.id,
},
data: {
status: "SUCCESS",
completedAt: new Date(),
run: {
update: {
yieldedExecutions: {
push: key,
},
},
},
},
});
const newJobExecution = await tx.jobRunExecution.create({
data: {
runId: run.id,
reason: "EXECUTE_JOB",
status: "PENDING",
retryLimit: EXECUTE_JOB_RETRY_LIMIT,
},
});
await enqueueRunExecutionV1(newJobExecution, run.queue.id, run.queue.maxJobs, tx);
});
}
async #resumeRunWithTask(execution: FoundRunExecution, data: RunJobResumeWithTask) {
const { run } = execution;
@@ -397,7 +449,9 @@ export class PerformRunExecutionV1Service {
// If the task has an operation, then the next performRunExecution will occur
// when that operation has finished
if (!data.task.operation) {
// Tasks with callbacks enabled will also get processed separately, i.e. when
// they time out, or on valid requests to their callbackUrl
if (!data.task.operation && !data.task.callbackUrl) {
const newJobExecution = await tx.jobRunExecution.create({
data: {
runId: run.id,
@@ -438,6 +492,24 @@ export class PerformRunExecutionV1Service {
});
}
async #failRunWithUnresolvedAuthError(
execution: FoundRunExecution,
data: RunJobUnresolvedAuthError
) {
return await $transaction(this.#prismaClient, async (tx) => {
await this.#failRunExecution(tx, execution, data.issues, "UNRESOLVED_AUTH");
});
}
async #failRunWithInvalidPayloadError(
execution: FoundRunExecution,
data: RunJobInvalidPayloadError
) {
return await $transaction(this.#prismaClient, async (tx) => {
await this.#failRunExecution(tx, execution, data.errors, "INVALID_PAYLOAD");
});
}
async #retryRunWithTask(execution: FoundRunExecution, data: RunJobRetryWithTask) {
const { run } = execution;
@@ -557,7 +629,7 @@ export class PerformRunExecutionV1Service {
prisma: PrismaClientOrTransaction,
execution: FoundRunExecution,
output: Record<string, any>,
status: "FAILURE" | "ABORTED" = "FAILURE"
status: "FAILURE" | "ABORTED" | "UNRESOLVED_AUTH" | "INVALID_PAYLOAD" = "FAILURE"
): Promise<void> {
const { run } = execution;
@@ -1,10 +1,17 @@
import {
CachedTask,
API_VERSIONS,
BloomFilter,
ConnectionAuth,
EndpointHeadersSchema,
RunJobError,
RunJobInvalidPayloadError,
RunJobResumeWithTask,
RunJobRetryWithTask,
RunJobSuccess,
RunJobUnresolvedAuthError,
RunSourceContext,
RunSourceContextSchema,
supportsFeature,
} from "@trigger.dev/core";
import { RuntimeEnvironmentType, type Task } from "@trigger.dev/database";
import { generateErrorMessage } from "zod-error";
@@ -16,10 +23,17 @@ import { formatError } from "~/utils/formatErrors.server";
import { safeJsonZodParse } from "~/utils/json";
import { EndpointApi } from "../endpointApi.server";
import { logger } from "../logger.server";
import { prepareTasksForCaching, prepareTasksForCachingLegacy } from "~/models/task.server";
import { MAX_RUN_YIELDED_EXECUTIONS } from "~/consts";
import { ApiEventLog } from "@trigger.dev/core";
import { RunJobBody } from "@trigger.dev/core";
type FoundRun = NonNullable<Awaited<ReturnType<typeof findRun>>>;
type FoundTask = FoundRun["tasks"][number];
// We need to limit the cached tasks to not be too large >3.5MB when serialized
const TOTAL_CACHED_TASK_BYTE_LIMIT = 3500000;
export type PerformRunExecutionV2Input = {
id: string;
reason: "PREPROCESS" | "EXECUTE_JOB";
@@ -151,6 +165,29 @@ export class PerformRunExecutionV2Service {
return;
}
try {
if (
typeof process.env.BLOCKED_ORGS === "string" &&
process.env.BLOCKED_ORGS.includes(run.organizationId)
) {
logger.debug("Skipping execution for blocked org", {
orgId: run.organizationId,
});
await this.#prismaClient.jobRun.update({
where: {
id: run.id,
},
data: {
status: "CANCELED",
completedAt: new Date(),
},
});
return;
}
} catch (e) {}
const client = new EndpointApi(run.environment.apiKey, run.endpoint.url);
const event = eventRecordToApiJson(run.event);
@@ -205,38 +242,19 @@ export class PerformRunExecutionV2Service {
const sourceContext = RunSourceContextSchema.safeParse(run.event.sourceContext);
const { response, parser, errorParser, durationInMs } = await client.executeJobRequest({
const executionBody = await this.#createExecutionBody(
run,
[run.tasks, resumedTask].flat().filter(Boolean),
startedAt,
isRetry,
connections.auth,
event,
job: {
id: run.version.job.slug,
version: run.version.version,
},
run: {
id: run.id,
isTest: run.isTest,
startedAt,
isRetry,
},
environment: {
id: run.environment.id,
slug: run.environment.slug,
type: run.environment.type,
},
organization: {
id: run.organization.id,
slug: run.organization.slug,
title: run.organization.title,
},
account: run.externalAccount
? {
id: run.externalAccount.identifier,
metadata: run.externalAccount.metadata,
}
: undefined,
connections: connections.auth,
source: sourceContext.success ? sourceContext.data : undefined,
tasks: prepareTasksForRun([run.tasks, resumedTask].flat().filter(Boolean)),
});
sourceContext.success ? sourceContext.data : undefined
);
const { response, parser, errorParser, durationInMs } = await client.executeJobRequest(
executionBody
);
if (!response) {
return await this.#failRunExecutionWithRetry({
@@ -244,6 +262,25 @@ export class PerformRunExecutionV2Service {
});
}
// Update the endpoint version if it has changed
const rawHeaders = Object.fromEntries(response.headers.entries());
const headers = EndpointHeadersSchema.safeParse(rawHeaders);
if (
headers.success &&
headers.data["trigger-version"] &&
headers.data["trigger-version"] !== run.endpoint.version
) {
await this.#prismaClient.endpoint.update({
where: {
id: run.endpoint.id,
},
data: {
version: headers.data["trigger-version"],
},
});
}
const rawBody = await response.text();
if (!response.ok) {
@@ -354,6 +391,20 @@ export class PerformRunExecutionV2Service {
await this.#cancelExecution(run);
break;
}
case "UNRESOLVED_AUTH_ERROR": {
await this.#failRunWithUnresolvedAuthError(run, safeBody.data, durationInMs);
break;
}
case "INVALID_PAYLOAD": {
await this.#failRunWithInvalidPayloadError(run, safeBody.data, durationInMs);
break;
}
case "YIELD_EXECUTION": {
await this.#resumeYieldedRun(run, safeBody.data.key, isRetry, durationInMs, executionCount);
break;
}
default: {
const _exhaustiveCheck: never = status;
throw new Error(`Non-exhaustive match for value: ${status}`);
@@ -361,6 +412,91 @@ export class PerformRunExecutionV2Service {
}
}
async #createExecutionBody(
run: FoundRun,
tasks: FoundTask[],
startedAt: Date,
isRetry: boolean,
connections: Record<string, ConnectionAuth>,
event: ApiEventLog,
source?: RunSourceContext
): Promise<RunJobBody> {
if (supportsFeature("lazyLoadedCachedTasks", run.endpoint.version)) {
const preparedTasks = prepareTasksForCaching(tasks, TOTAL_CACHED_TASK_BYTE_LIMIT);
return {
event,
job: {
id: run.version.job.slug,
version: run.version.version,
},
run: {
id: run.id,
isTest: run.isTest,
startedAt,
isRetry,
},
environment: {
id: run.environment.id,
slug: run.environment.slug,
type: run.environment.type,
},
organization: {
id: run.organization.id,
slug: run.organization.slug,
title: run.organization.title,
},
account: run.externalAccount
? {
id: run.externalAccount.identifier,
metadata: run.externalAccount.metadata,
}
: undefined,
connections,
source,
tasks: preparedTasks.tasks,
cachedTaskCursor: preparedTasks.cursor,
noopTasksSet: prepareNoOpTasksBloomFilter(tasks),
yieldedExecutions: run.yieldedExecutions,
};
}
const preparedTasks = prepareTasksForCachingLegacy(tasks, TOTAL_CACHED_TASK_BYTE_LIMIT);
return {
event,
job: {
id: run.version.job.slug,
version: run.version.version,
},
run: {
id: run.id,
isTest: run.isTest,
startedAt,
isRetry,
},
environment: {
id: run.environment.id,
slug: run.environment.slug,
type: run.environment.type,
},
organization: {
id: run.organization.id,
slug: run.organization.slug,
title: run.organization.title,
},
account: run.externalAccount
? {
id: run.externalAccount.identifier,
metadata: run.externalAccount.metadata,
}
: undefined,
connections,
source,
tasks: preparedTasks.tasks,
};
}
async #completeRunWithSuccess(run: FoundRun, data: RunJobSuccess, durationInMs: number) {
await this.#prismaClient.jobRun.update({
where: { id: run.id },
@@ -394,7 +530,9 @@ export class PerformRunExecutionV2Service {
// If the task has an operation, then the next performRunExecution will occur
// when that operation has finished
if (!data.task.operation) {
// Tasks with callbacks enabled will also get processed separately, i.e. when
// they time out, or on valid requests to their callbackUrl
if (!data.task.operation && !data.task.callbackUrl) {
await enqueueRunExecutionV2(run, tx, {
runAt: data.task.delayUntil ?? undefined,
resumeTaskId: data.task.id,
@@ -432,6 +570,90 @@ export class PerformRunExecutionV2Service {
});
}
async #failRunWithUnresolvedAuthError(
execution: FoundRun,
data: RunJobUnresolvedAuthError,
durationInMs: number
) {
return await $transaction(this.#prismaClient, async (tx) => {
await this.#failRunExecution(
tx,
"EXECUTE_JOB",
execution,
data.issues,
"UNRESOLVED_AUTH",
durationInMs
);
});
}
async #failRunWithInvalidPayloadError(
execution: FoundRun,
data: RunJobInvalidPayloadError,
durationInMs: number
) {
return await $transaction(this.#prismaClient, async (tx) => {
await this.#failRunExecution(
tx,
"EXECUTE_JOB",
execution,
data.errors,
"INVALID_PAYLOAD",
durationInMs
);
});
}
async #resumeYieldedRun(
run: FoundRun,
key: string,
isRetry: boolean,
durationInMs: number,
executionCount: number
) {
await $transaction(this.#prismaClient, async (tx) => {
if (run.yieldedExecutions.length + 1 > MAX_RUN_YIELDED_EXECUTIONS) {
return await this.#failRunExecution(
tx,
"EXECUTE_JOB",
run,
{
message: `Run has yielded too many times, the maximum is ${MAX_RUN_YIELDED_EXECUTIONS}`,
},
"FAILURE",
durationInMs
);
}
await tx.jobRun.update({
where: {
id: run.id,
},
data: {
executionDuration: {
increment: durationInMs,
},
executionCount: {
increment: 1,
},
yieldedExecutions: {
push: key,
},
},
select: {
yieldedExecutions: true,
executionCount: true,
},
});
await enqueueRunExecutionV2(run, tx, {
isRetry,
skipRetrying: run.environment.type === RuntimeEnvironmentType.DEVELOPMENT,
executionCount,
});
});
}
async #retryRunWithTask(
run: FoundRun,
data: RunJobRetryWithTask,
@@ -556,7 +778,7 @@ export class PerformRunExecutionV2Service {
reason: "EXECUTE_JOB" | "PREPROCESS",
run: FoundRun,
output: Record<string, any>,
status: "FAILURE" | "ABORTED" | "TIMED_OUT" = "FAILURE",
status: "FAILURE" | "ABORTED" | "TIMED_OUT" | "UNRESOLVED_AUTH" | "INVALID_PAYLOAD" = "FAILURE",
durationInMs: number = 0
): Promise<void> {
await $transaction(prisma, async (tx) => {
@@ -617,69 +839,16 @@ export class PerformRunExecutionV2Service {
}
}
function prepareTasksForRun(possibleTasks: FoundTask[]): CachedTask[] {
const tasks = possibleTasks.filter((task) => task.status === "COMPLETED");
function prepareNoOpTasksBloomFilter(possibleTasks: FoundTask[]): string {
const tasks = possibleTasks.filter((task) => task.status === "COMPLETED" && task.noop);
// We need to limit the cached tasks to not be too large >3.5MB when serialized
const TOTAL_CACHED_TASK_BYTE_LIMIT = 3500000;
const filter = new BloomFilter(BloomFilter.NOOP_TASK_SET_SIZE);
const cachedTasks = new Map<string, CachedTask>(); // Cache for prepared tasks
const cachedTaskSizes = new Map<string, number>(); // Cache for calculated task sizes
// Helper function to get the cached prepared task, or prepare and cache if not already cached
function getCachedTask(task: FoundTask): CachedTask {
const taskId = task.id;
if (!cachedTasks.has(taskId)) {
cachedTasks.set(taskId, prepareTaskForRun(task));
}
return cachedTasks.get(taskId)!;
for (const task of tasks) {
filter.add(task.idempotencyKey);
}
// Helper function to get the cached task size, or calculate and cache if not already cached
function getCachedTaskSize(task: CachedTask): number {
const taskId = task.id;
if (!cachedTaskSizes.has(taskId)) {
cachedTaskSizes.set(taskId, calculateCachedTaskSize(task));
}
return cachedTaskSizes.get(taskId)!;
}
// Prepare tasks and calculate their sizes
const availableTasks = tasks.map((task) => {
const cachedTask = getCachedTask(task);
return { task: cachedTask, size: getCachedTaskSize(cachedTask) };
});
// Sort tasks in ascending order by size
availableTasks.sort((a, b) => a.size - b.size);
// Select tasks using greedy approach
const tasksToRun: CachedTask[] = [];
let remainingSize = TOTAL_CACHED_TASK_BYTE_LIMIT;
for (const { task, size } of availableTasks) {
if (size <= remainingSize) {
tasksToRun.push(task);
remainingSize -= size;
}
}
return tasksToRun;
}
function prepareTaskForRun(task: FoundTask): CachedTask {
return {
id: task.idempotencyKey, // We should eventually move this back to task.id
status: task.status,
idempotencyKey: task.idempotencyKey,
noop: task.noop,
output: task.output as any,
parentId: task.parentId,
};
}
function calculateCachedTaskSize(task: CachedTask): number {
return JSON.stringify(task).length;
return filter.serialize();
}
async function findRun(prisma: PrismaClientOrTransaction, id: string) {
@@ -714,6 +883,9 @@ async function findRun(prisma: PrismaClientOrTransaction, id: string) {
output: true,
parentId: true,
},
orderBy: {
id: "asc",
},
},
event: true,
version: {
@@ -20,6 +20,7 @@ export class ReRunService {
version: true,
job: true,
event: true,
externalAccount: true,
},
where: {
id: runId,
@@ -43,6 +44,13 @@ export class ReRunService {
id: existingRun.environment.id,
},
},
externalAccount: existingRun.externalAccount
? {
connect: {
id: existingRun.externalAccount.id,
},
}
: undefined,
eventId: `${existingRun.event.eventId}:retry:${new Date().getTime()}`,
name: existingRun.event.name,
timestamp: new Date(),
@@ -50,11 +50,11 @@ export class StartRunService {
integrationId: runConnection.integration.id,
authSource: "HOSTED",
} as const)
: runConnection.result === "resolvedLocal"
: runConnection.result === "resolvedLocal" || runConnection.result === "resolvedResolver"
? ({
key,
integrationId: runConnection.integration.id,
authSource: "LOCAL",
authSource: runConnection.result === "resolvedLocal" ? "LOCAL" : "RESOLVER",
} as const)
: undefined
)
@@ -173,6 +173,7 @@ async function createRunConnections(tx: PrismaClientOrTransaction, run: FoundRun
integration: Integration;
}
| { result: "resolvedLocal"; integration: Integration }
| { result: "resolvedResolver"; integration: Integration }
| {
result: "missing";
connectionType: ConnectionType;
@@ -190,6 +191,11 @@ async function createRunConnections(tx: PrismaClientOrTransaction, run: FoundRun
result: "resolvedLocal",
integration: jobIntegration.integration,
};
} else if (jobIntegration.integration.authSource === "RESOLVER") {
acc[jobIntegration.key] = {
result: "resolvedResolver",
integration: jobIntegration.integration,
};
} else {
const connection = run.externalAccountId
? await tx.integrationConnection.findFirst({
@@ -59,6 +59,7 @@ export class RegisterScheduleService {
schedule: payload,
accountId: payload.accountId,
dynamicTrigger,
organizationId: environment.organizationId,
});
return registration;
@@ -16,24 +16,32 @@ export class RegisterScheduleSourceService {
schedule,
accountId,
dynamicTrigger,
organizationId,
}: {
key: string;
dispatcher: EventDispatcher;
schedule: ScheduleMetadata;
accountId?: string;
dynamicTrigger?: DynamicTrigger;
organizationId: string;
}) {
const validatedSchedule = validateSchedule(schedule);
return await $transaction(this.#prismaClient, async (tx) => {
const externalAccount = accountId
? await tx.externalAccount.findUniqueOrThrow({
? await tx.externalAccount.upsert({
where: {
environmentId_identifier: {
environmentId: dispatcher.environmentId,
identifier: accountId,
},
},
create: {
environmentId: dispatcher.environmentId,
organizationId: organizationId,
identifier: accountId,
},
update: {},
})
: undefined;
@@ -71,13 +71,19 @@ export class RegisterSourceServiceV1 {
}
const externalAccount = accountId
? await tx.externalAccount.findUniqueOrThrow({
? await tx.externalAccount.upsert({
where: {
environmentId_identifier: {
environmentId: environment.id,
identifier: accountId,
},
},
create: {
environmentId: environment.id,
organizationId: environment.organizationId,
identifier: accountId,
},
update: {},
})
: undefined;
@@ -71,13 +71,19 @@ export class RegisterSourceServiceV2 {
}
const externalAccount = accountId
? await tx.externalAccount.findUniqueOrThrow({
? await tx.externalAccount.upsert({
where: {
environmentId_identifier: {
environmentId: environment.id,
identifier: accountId,
},
},
create: {
environmentId: environment.id,
organizationId: environment.organizationId,
identifier: accountId,
},
update: {},
})
: undefined;
@@ -1,5 +1,5 @@
import crypto from "node:crypto";
export function generateSecret(): string {
return crypto.randomBytes(32).toString("hex");
export function generateSecret(sizeInBytes = 32): string {
return crypto.randomBytes(sizeInBytes).toString("hex");
}
@@ -0,0 +1,76 @@
import { RuntimeEnvironmentType } from "@trigger.dev/database";
import { $transaction, PrismaClient, PrismaClientOrTransaction, prisma } from "~/db.server";
import { enqueueRunExecutionV2 } from "~/models/jobRunExecution.server";
import { logger } from "../logger.server";
type FoundTask = Awaited<ReturnType<typeof findTask>>;
export class ProcessCallbackTimeoutService {
#prismaClient: PrismaClient;
constructor(prismaClient: PrismaClient = prisma) {
this.#prismaClient = prismaClient;
}
public async call(id: string) {
const task = await findTask(this.#prismaClient, id);
if (!task) {
return;
}
if (task.status !== "WAITING" || !task.callbackUrl) {
return;
}
logger.debug("ProcessCallbackTimeoutService.call", { task });
return await this.#failTask(task, "Remote callback timeout - no requests received");
}
async #failTask(task: NonNullable<FoundTask>, error: string) {
await $transaction(this.#prismaClient, async (tx) => {
await tx.taskAttempt.updateMany({
where: {
taskId: task.id,
status: "PENDING",
},
data: {
status: "ERRORED",
error
},
});
await tx.task.update({
where: { id: task.id },
data: {
status: "ERRORED",
completedAt: new Date(),
output: error,
},
});
await this.#resumeRunExecution(task, tx);
});
}
async #resumeRunExecution(task: NonNullable<FoundTask>, prisma: PrismaClientOrTransaction) {
await enqueueRunExecutionV2(task.run, prisma, {
skipRetrying: task.run.environment.type === RuntimeEnvironmentType.DEVELOPMENT,
});
}
}
async function findTask(prisma: PrismaClient, id: string) {
return prisma.task.findUnique({
where: { id },
include: {
run: {
include: {
environment: true,
queue: true,
},
},
},
});
}
@@ -24,7 +24,6 @@ export class RegisterTriggerSourceServiceV2 {
endpointSlug,
id,
key,
accountId,
registrationMetadata,
}: {
environment: AuthenticatedEnvironment;
@@ -32,7 +31,6 @@ export class RegisterTriggerSourceServiceV2 {
id: string;
endpointSlug: string;
key: string;
accountId?: string;
registrationMetadata?: any;
}): Promise<RegisterSourceEventV2 | undefined> {
const endpoint = await this.#prismaClient.endpoint.findUniqueOrThrow({
@@ -63,7 +61,7 @@ export class RegisterTriggerSourceServiceV2 {
endpoint.id,
payload.source,
dynamicTrigger.id,
accountId,
payload.accountId,
{ id: key, metadata: registrationMetadata }
);
+15 -1
View File
@@ -19,6 +19,7 @@ import { DeliverScheduledEventService } from "./schedules/deliverScheduledEvent.
import { ActivateSourceService } from "./sources/activateSource.server";
import { DeliverHttpSourceRequestService } from "./sources/deliverHttpSourceRequest.server";
import { PerformTaskOperationService } from "./tasks/performTaskOperation.server";
import { ProcessCallbackTimeoutService } from "./tasks/processCallbackTimeout";
import { addMissingVersionField } from "@trigger.dev/core";
const workerCatalog = {
@@ -30,6 +31,9 @@ const workerCatalog = {
}),
scheduleEmail: DeliverEmailSchema,
startRun: z.object({ id: z.string() }),
processCallbackTimeout: z.object({
id: z.string(),
}),
performTaskOperation: z.object({
id: z.string(),
}),
@@ -161,7 +165,8 @@ function getWorkerQueue() {
tasks: {
"events.invokeDispatcher": {
priority: 0, // smaller number = higher priority
maxAttempts: 3,
maxAttempts: 6,
queueName: (payload) => `dispatcher:${payload.id}`, // use a queue for a dispatcher so runs are created sequentially
handler: async (payload, job) => {
const service = new InvokeDispatcherService();
@@ -239,6 +244,15 @@ function getWorkerQueue() {
await service.call(payload.id);
},
},
processCallbackTimeout: {
priority: 0, // smaller number = higher priority
maxAttempts: 3,
handler: async (payload, job) => {
const service = new ProcessCallbackTimeoutService();
await service.call(payload.id);
},
},
performTaskOperation: {
priority: 0, // smaller number = higher priority
queueName: (payload) => `tasks:${payload.id}`,
+71
View File
@@ -0,0 +1,71 @@
// Redacts the given object based on the given paths
// Example:
// const redactor = new Redactor(["data.object.balance_transaction"]);
// redactor.redact({
// data: {
// object: {
// balance_transaction: "txn_1NYWgTI0XSgju2urW3aXpinM",
// },
// },
// });
// Returns:
// {
// data: {
// object: {
// balance_transaction: "[REDACTED]",
// },
// },
// }
// Does not currenly support arrays
export class Redactor {
constructor(private paths: string[]) {}
public redact(subject: unknown): unknown {
if (!Array.isArray(this.paths)) {
return subject;
}
if (this.paths.length === 0) {
return subject;
}
const clonedSubject = JSON.parse(JSON.stringify(subject));
return this.redactPathsRecursive(clonedSubject, this.paths);
}
private redactPathsRecursive(subject: any, paths: string[]): any {
for (let path of paths) {
let parts = path.split(".");
let curSubject = subject;
// Make sure curSubject is an object
if (typeof curSubject !== "object") {
break;
}
for (let i = 0; i < parts.length; i++) {
const part = parts[i];
if (Object.prototype.hasOwnProperty.call(curSubject, part) === false) {
// Path is not found in object
break;
}
if (i === parts.length - 1) {
// We're at the end of our path and have a string, redact it
curSubject[part] = "[REDACTED]";
} else if (part in curSubject && typeof curSubject[part] === "object") {
// More paths to follow, continue down the path
curSubject = curSubject[part];
} else {
// Path is not found in object or doesn't point to a string
break;
}
}
}
return subject;
}
}
+2 -2
View File
@@ -61,8 +61,8 @@
"@remix-run/server-runtime": "1.19.2-pre.0",
"@team-plain/typescript-sdk": "^2.2.0",
"@trigger.dev/companyicons": "^1.5.14",
"@trigger.dev/database": "workspace:*",
"@trigger.dev/core": "workspace:*",
"@trigger.dev/database": "workspace:*",
"@trigger.dev/sdk": "workspace:*",
"@uiw/react-codemirror": "^4.19.5",
"class-variance-authority": "^0.5.2",
@@ -73,7 +73,6 @@
"cuid": "^2.1.8",
"emails": "workspace:*",
"express": "^4.18.1",
"fast-redact": "^3.1.2",
"framer-motion": "^10.12.11",
"graphile-worker": "^0.13.0",
"highlight.run": "^7.3.4",
@@ -96,6 +95,7 @@
"react-hot-toast": "^2.4.0",
"react-hotkeys-hook": "^3.4.7",
"react-use": "^17.4.0",
"recharts": "^2.8.0",
"remix-auth": "^3.2.2",
"remix-auth-email-link": "^1.4.2",
"remix-auth-github": "^1.1.1",
+67
View File
@@ -3,6 +3,7 @@
import { integrationCatalog } from "../app/services/externalApis/integrationCatalog.server";
import { seedCloud } from "./seedCloud";
import { prisma } from "../app/db.server";
import { createEnvironment } from "~/models/organization.server";
async function seedIntegrationAuthMethods() {
for (const [_, integration] of Object.entries(integrationCatalog.getIntegrations())) {
@@ -67,12 +68,78 @@ async function seedIntegrationAuthMethods() {
}
}
async function runDataMigrations() {
await runStagingEnvironmentMigration();
}
async function runStagingEnvironmentMigration() {
try {
await prisma.$transaction(async (tx) => {
const existingDataMigration = await tx.dataMigration.findUnique({
where: {
name: "2023-09-27-AddStagingEnvironments",
},
});
if (existingDataMigration) {
return;
}
await tx.dataMigration.create({
data: {
name: "2023-09-27-AddStagingEnvironments",
},
});
console.log("Running data migration 2023-09-27-AddStagingEnvironments");
const projectsWithoutStagingEnvironments = await tx.project.findMany({
where: {
environments: {
none: {
type: "STAGING",
},
},
},
include: {
organization: true,
},
});
for (const project of projectsWithoutStagingEnvironments) {
try {
console.log(
`Creating staging environment for project ${project.slug} on org ${project.organization.slug}`
);
await createEnvironment(project.organization, project, "STAGING", undefined, tx);
} catch (error) {
console.error(error);
}
}
await tx.dataMigration.update({
where: {
name: "2023-09-27-AddStagingEnvironments",
},
data: {
completedAt: new Date(),
},
});
});
} catch (error) {
console.error(error);
}
}
async function seed() {
await seedIntegrationAuthMethods();
if (process.env.NODE_ENV === "development" && process.env.SEED_CLOUD === "enabled") {
await seedCloud(prisma);
}
await runDataMigrations();
}
seed()
+17
View File
@@ -0,0 +1,17 @@
{
"extends": "./node18.json",
"compilerOptions": {
"lib": ["DOM", "DOM.Iterable", "ES2019"],
"paths": {
"@trigger.dev/sdk/*": ["../../packages/trigger-sdk/src/*"],
"@trigger.dev/sdk": ["../../packages/trigger-sdk/src/index"],
"@trigger.dev/integration-kit/*": ["../../packages/integration-kit/src/*"],
"@trigger.dev/integration-kit": ["../../packages/integration-kit/src/index"]
},
"declaration": false,
"declarationMap": false,
"baseUrl": ".",
"stripInternal": true
},
"exclude": ["node_modules"]
}
+2 -2
View File
@@ -2,11 +2,11 @@
## Install and initial setup
`npm install`
`pnpm install`
## Running the app
`npm run dev`
`pnpm run dev --filter docs`
## View the app locally
+3
View File
@@ -0,0 +1,3 @@
<Card title="React hooks" icon="fishing-rod" href="/documentation/guides/react-hooks">
Show the live status of Job Runs in your React app
</Card>
+5
View File
@@ -0,0 +1,5 @@
<Accordion title="How do I get a Run id?">
You can call [client.getRuns()](/sdk/triggerclient/instancemethods/getruns) with a Job id to get a
list of the most recent Runs for that Job. You can then pass that run id to your frontend to use
in the hook.
</Accordion>
+15
View File
@@ -0,0 +1,15 @@
<CodeGroup>
```bash npm
npm install @trigger.dev/slack@latest
```
```bash pnpm
pnpm install @trigger.dev/slack@latest
```
```bash yarn
yarn add @trigger.dev/slack@latest
```
</CodeGroup>
+49
View File
@@ -0,0 +1,49 @@
<ParamField body="options" type="object" required>
<Expandable title="properties" defaultOpen>
<ParamField body="id" type="string" required>
The `id` property is used to uniquely identify the Job. Only change this if you want to create a new Job.
</ParamField>
<ParamField body="name" type="string" required>
The `name` of the Job that you want to appear in the dashboard and logs. You can change this without creating a new Job.
</ParamField>
<ParamField body="version" type="string" required>
The `version` property is used to version your Job. A new version will be created if you change this property. We recommend using [semantic versioning](https://www.baeldung.com/cs/semantic-versioning), e.g. `1.0.3`.
</ParamField>
<ParamField body="trigger" type="object" required>
The `trigger` property is used to define when the Job should run. There are currently the following Trigger types:
- [cronTrigger](/sdk/crontrigger)
- [intervalTrigger](/sdk/intervaltrigger)
- [eventTrigger](/sdk/eventtrigger)
- [DynamicTrigger](/sdk/dynamictrigger)
- [DynamicSchedule](/sdk/dynamicschedule)
- integration Triggers, like webhooks. See the [integrations](/integrations) page for more information.
</ParamField>
<ParamField body="run" type="function" required>
This function gets called automatically when a Run is Triggered. It has three parameters:
1. `payload` The payload that was sent to the Trigger API.
2. [io](/sdk/io) An object that contains the integrations that you specified in the `integrations` property and other useful functions like delays and running Tasks.
3. [context](/sdk/context) An object that contains information about the Organization, Job, Run and more.
This is where you put the code you want to run for a Job. You can use normal code in here and you can also use Tasks.
You can return a value from this function and it will be sent back to the Trigger API.
</ParamField>
<ParamField body="integrations" type="object">
Imports the specified integrations into the Job. The integrations will be available on the `io` object in the `run()` function with the same name as the key. For example:
<Snippet file="how-to-pass-integrations.mdx" />
</ParamField>
<ParamField body="enabled" type="boolean">
The `enabled` property is an optional property that specifies whether the Job is enabled or not. The Job will be enabled by default if you omit this property. When a job is disabled, no new runs will be triggered or resumed. In progress runs will continue to run until they are finished or delayed by using `io.wait`.
</ParamField>
<ParamField body="logLevel" type="log | error | warn | info | debug">
The `logLevel` property is an optional property that specifies the level of
logging for the Job. The level is inherited from the client if you omit this property.
- `log` - logs only essential messages
- `error` - logs error messages
- `warn` - logs errors and warning messages
- `info` - logs errors, warnings and info messages
- `debug` - logs everything with full verbosity
</ParamField>
</Expandable>
</ParamField>
+64 -92
View File
@@ -5,15 +5,15 @@ To begin, install the necessary packages in your Astro project directory. You ca
<CodeGroup>
```bash npm
npm i @trigger.dev/sdk @trigger.dev/astro
npm i @trigger.dev/sdk@latest @trigger.dev/astro@latest
```
```bash pnpm
pnpm install @trigger.dev/sdk @trigger.dev/astro
pnpm install @trigger.dev/sdk@latest @trigger.dev/astro@latest
```
```bash yarn
yarn add @trigger.dev/sdk @trigger.dev/astro
yarn add @trigger.dev/sdk@latest @trigger.dev/astro@latest
```
</CodeGroup>
@@ -32,103 +32,54 @@ Create a `.env` file at the root of your project and include your Trigger API ke
```bash
TRIGGER_API_KEY=ENTER_YOUR_DEVELOPMENT_API_KEY_HERE
TRIGGER_API_URL=https://cloud.trigger.dev
TRIGGER_API_URL=https://api.trigger.dev # this is only necessary if you are self-hosting
```
Replace `ENTER_YOUR_DEVELOPMENT_API_KEY_HERE` with the actual API key obtained from the previous step.
## Configuring the Trigger Client
To set up the Trigger Client for your project, follow these steps:
Create a file at `<root>/trigger.ts` or `<root>/src/trigger.ts`, depending on if your project uses a `src` directory, where `<root>` represents the root directory of your project.
1. **Create Configuration File:**
Next, add the following code to the file which creates and exports a new `TriggerClient`:
In your project directory, create a configuration file named `trigger.ts` or `trigger.js`, depending on whether your project uses TypeScript (`.ts`) or JavaScript (`.js`).
```typescript src/trigger.ts
import { TriggerClient } from "@trigger.dev/sdk";
2. **Choose Directory:**
export const client = new TriggerClient({
id: "my-astro-app",
apiKey: import.meta.env.TRIGGER_API_KEY,
apiUrl: import.meta.env.TRIGGER_API_URL,
});
```
Depending on your project structure, choose the appropriate directory for the configuration file. If your project uses a `src` directory, create the file within it or Otherwise, create it directly in the project root.
Replace **"my-astro-app"** with an appropriate identifier for your project.
3. **Add Configuration Code:**
## Update the astro.config file to enable SSR (Server Side Rendering)
Open the configuration file you created and add the following code:
```typescript src/trigger.(ts/js)
// trigger.ts (for TypeScript) or trigger.js (for JavaScript)
import { TriggerClient } from "@trigger.dev/sdk";
export const client = new TriggerClient({
id: "my-app",
apiKey: process.env.TRIGGER_API_KEY,
apiUrl: process.env.TRIGGER_API_URL,
});
```
Replace **"my-app"** with an appropriate identifier for your project. The **apiKey** and **apiUrl** are obtained from the environment variables you set earlier.
4. **File Location:**
- You can save the file within the **src** directory or in the project rooot.
**Example Directory Structure with src:**
```
project-root/
├── src/
├── trigger.ts
├── other files...
```
**Example Directory Structure without src:**
```
project-root/
├── trigger.ts
├── other files...
```
By following these steps, you'll configure the Trigger Client to work with your project, regardless of whether you have a separate **src** directory and whether you're using TypeScript or JavaScript files.
## update the astro.config file to enable ssr
- You need to enable ssr to use API endpoints that would be in the `pages/api` folder
- You need to enable SSR to use API endpoints (which are required by Trigger.dev).
```typescript astro.config.mjs
import { defineConfig } from "astro/config";
export default defineConfig({
//alternatively you can use "hybrid" instead of "server"
output: "server",
});
```
## Creating the API Route
To learn more about SSR, head over to the [Astro docs on SSR](https://docs.astro.build/en/guides/server-side-rendering/).
To establish an API route for interacting with Trigger.dev, follow these steps based on your project's file type and structure
## Creating an Example Job
1. Create a new file named `trigger.(ts/js)` within the `pages/api/` directory.
2. Add the following code to `trigger.(ts/js)`:
```typescript api/trigger.ts/js
import { createAstroRoute } from "@trigger.dev/astro";
import { client } from "@/trigger";
//import your jobs
import "@/jobs";
export const { POST } = createAstroRoute(client);
```
## Creating the Example Job
1. Create a folder named `Jobs` alongside your `pages` directory
2. Inside the `Jobs` folder, add two files named `example.(ts/js)` and `index.(ts/js)`.
1. Create a folder named `jobs` alongside your `pages` directory
2. Inside the `jobs` folder, add two files named `example.ts` and `index.ts`.
<CodeGroup>
```typescript example.(ts/js)
```typescript src/jobs/example.ts
import { eventTrigger } from "@trigger.dev/sdk";
import { client } from "@/trigger";
import { client } from "../trigger";
// your first job
client.defineJob({
@@ -148,25 +99,30 @@ client.defineJob({
});
```
```typescript index.ts/index.(ts/js)
// import all your job files here
export * from "./examples";
```typescript src/jobs/index.ts
// export all your job files here
export * from "./example";
```
</CodeGroup>
## Additonal Job Definitions
## Creating the API Route
You can define more job definitions by creating additional files in the `Jobs` folder and exporting them in `index` file.
To establish an API route for interacting with Trigger.dev, follow these steps based on your project's file type and structure
For example, in `index.(ts/js)`, you can export other job files like this:
1. Create a new file named `trigger.ts` within the `pages/api/` directory.
2. Add the following code to `trigger.ts`:
```typescript
// import all your job files here
```typescript src/pages/api/trigger.ts
import { createAstroRoute } from "@trigger.dev/astro";
//you may need to update this path to point at your trigger.ts file
import { client } from "../../trigger";
export * from "./examples";
export * from "./other-job-file";
//import your jobs, this could be different depending on your project structure
import "../../jobs";
export const prerender = false;
export const { POST } = createAstroRoute(client);
```
## Adding Configuration to `package.json`
@@ -175,7 +131,7 @@ Inside the `package.json` file, add the following configuration under the root o
```json
"trigger.dev": {
"endpointId": "my-app"
"endpointId": "my-astro-app"
}
```
@@ -189,12 +145,25 @@ Your `package.json` file might look something like this:
// ... other dependencies
},
"trigger.dev": {
"endpointId": "my-app"
"endpointId": "my-astro-app"
}
}
```
Replace **"my-app"** with the appropriate identifier you used during the step for creating the Trigger Client.
Replace **"my-astro-app"** with the appropriate identifier you used during the step for creating the `TriggerClient`.
## Additonal Job Definitions
You can define more job definitions by creating additional files in the `jobs` folder and exporting them in `index` file.
For example, in `index.ts`, you can export other job files like this:
```typescript
// import all your job files here
export * from "./examples";
export * from "./other-job-file";
```
## Running
@@ -225,24 +194,27 @@ In a **_separate terminal window or tab_** run:
<CodeGroup>
```bash npm
npx @trigger.dev/cli@latest dev
npx @trigger.dev/cli@latest dev --port 4321
```
```bash pnpm
pnpm dlx @trigger.dev/cli@latest dev
pnpm dlx @trigger.dev/cli@latest dev --port 4321
```
```bash yarn
yarn dlx @trigger.dev/cli@latest dev
yarn dlx @trigger.dev/cli@latest dev --port 4321
```
</CodeGroup>
<br />
<Note>
You can optionally pass the port if you're not running on 3000 by adding
`--port 4321` to the end
Astro by default runs on port 4321.
</Note>
<Note>
You can optionally pass the hostname if you're not running on localhost by adding
`--hostname <host>`. Example, in case your Astro app is running on 0.0.0.0: `--hostname 0.0.0.0`.
</Note>
### Next Steps
You should now see your example job in the Trigger.dev dashboard. You can now create additional jobs and use the Trigger.dev dashboard to test them.
+193 -1
View File
@@ -1 +1,193 @@
We're in the process of building support for the Express framework. You can follow along with progress or contribute via [this GitHub issue](https://github.com/triggerdotdev/trigger.dev/issues).
## Installing Required Packages
Start by installing the necessary packages in your Express.js project directory. You can use npm, pnpm, or yarn as your package manager.
<CodeGroup>
```bash npm
npm install @trigger.dev/sdk @trigger.dev/express
```
```bash pnpm
pnpm install @trigger.dev/sdk @trigger.dev/express
```
```bash yarn
yarn add @trigger.dev/sdk @trigger.dev/express
```
</CodeGroup>
<br />
<Note>Ensure that you execute this command within a Express project.</Note>
## Obtaining the Development Server API Key
To locate your development Server API key, login to the [Trigger.dev
dashboard](https://cloud.trigger.dev) and select the Project you want to
connect to. Then click on the Environments & API Keys tab in the left menu.
You can copy your development Server API Key from the field at the top of this page.
(Your development key will start with `tr_dev_`).
## Adding Environment Variables
Create a `.env` file at the root of your project and include your Trigger API key and URL like this:
```bash
TRIGGER_API_KEY=ENTER_YOUR_DEVELOPMENT_API_KEY_HERE
TRIGGER_API_URL=https://api.trigger.dev # this is only necessary if you are self-hosting
```
Replace `ENTER_YOUR_DEVELOPMENT_API_KEY_HERE` with the actual API key obtained from the previous step.
## Configuring the Trigger Client
Create a file for your Trigger client, in this case we create it at `<root>/trigger.(ts/js)`
```ts trigger.(ts/js)
import { TriggerClient } from "@trigger.dev/sdk";
export const client = new TriggerClient({
id: "my-app",
apiKey: process.env.TRIGGER_API_KEY!,
apiUrl: process.env.TRIGGER_API_URL!,
});
```
Replace **"my-app"** with an appropriate identifier for your project.
## Adding the API endpoint
There are a few different options depending on how your Express project is configured.
- App middleware
- Entire app for Trigger.dev (only relevant if it's the only thing your project is for)
Select the appropriate code example from below:
<CodeGroup>
```typescript app middleware
//import the client from the other file
import { client } from "./trigger";
import { createMiddleware } from "@trigger.dev/express";
//import your job files
import "./jobs/example";
//..your existing Express code
const app: Express = express();
//add the middleware
app.use(createMiddleware(client));
//..the rest of your Express code
```
```typescript entire app
//if the entire app is just for Trigger.dev
import { client } from "./trigger";
import { createExpressServer } from "@trigger.dev/express";
//import your job files
import "./jobs/example";
//this creates an app
createExpressServer(client);
```
</CodeGroup>
## Creating the Example Job
Create a Job file. In this case created `<root>/jobs/example.(ts/js)`
```typescript jobs/example.(ts/js)
import { eventTrigger } from "@trigger.dev/sdk";
import { client } from "../trigger";
// your first job
client.defineJob({
id: "example-job",
name: "Example Job",
version: "0.0.1",
trigger: eventTrigger({
name: "example.event",
}),
run: async (payload, io, ctx) => {
await io.logger.info("Hello world!", { payload });
return {
message: "Hello world!",
};
},
});
```
## Adding Configuration to `package.json`
Inside the `package.json` file, add the following configuration under the root object:
```json
"trigger.dev": {
"endpointId": "my-app"
}
```
Replace **"my-app"** with the appropriate identifier you used in the trigger.js configuration file.
## Running
### Run your Express app
Run your Express app locally, like you normally would. For example:
<CodeGroup>
```bash npm
npm run dev
```
```bash pnpm
pnpm run dev
```
```bash yarn
yarn run dev
```
</CodeGroup>
<Note>You might use `npm run start` instead of dev</Note>
### Run the CLI 'dev' command
In a **_separate terminal window or tab_** run:
<CodeGroup>
```bash npm
npx @trigger.dev/cli@latest dev
```
```bash pnpm
pnpm dlx @trigger.dev/cli@latest dev
```
```bash yarn
yarn dlx @trigger.dev/cli@latest dev
```
</CodeGroup>
<br />
<Note>
You can optionally pass the port if you're not running on 3000 by adding
`--port 3001` to the end
</Note>
<Note>
You can optionally pass the hostname if you're not running on localhost by adding
`--hostname <host>`. Example, in case your Express is running on 0.0.0.0: `--hostname 0.0.0.0`.
</Note>
+13 -33
View File
@@ -32,50 +32,30 @@ Create a `.env` file at the root of your project and include your Trigger API ke
```bash
TRIGGER_API_KEY=ENTER_YOUR_DEVELOPMENT_API_KEY_HERE
TRIGGER_API_URL=https://cloud.trigger.dev
TRIGGER_API_URL=https://api.trigger.dev # this is only necessary if you are self-hosting
```
Replace `ENTER_YOUR_DEVELOPMENT_API_KEY_HERE` with the actual API key obtained from the previous step.
## Configuring the Trigger Client
To set up the Trigger Client for your project, follow these steps:
Create a file at `<root>/app/trigger.ts`, where `<root>` represents the root directory of your project.
1. **Create Configuration File:**
Next, add the following code to the file which creates and exports a new `TriggerClient`:
In your project directory, create a configuration file named `trigger.ts` or `trigger.js`, depending on whether your project uses TypeScript (`.ts`) or JavaScript (`.js`).
```typescript app/trigger.(ts/js)
// trigger.ts (for TypeScript) or trigger.js (for JavaScript)
2. **Choose Directory:**
import { TriggerClient } from "@trigger.dev/sdk";
Create the configuration file inside the **app** directory of your project.
export const client = new TriggerClient({
id: "my-app",
apiKey: process.env.TRIGGER_API_KEY,
apiUrl: process.env.TRIGGER_API_URL,
});
```
3. **Add Configuration Code:**
Open the configuration file you created and add the following code:
```typescript app/trigger.(ts/js)
// trigger.ts (for TypeScript) or trigger.js (for JavaScript)
import { TriggerClient } from "@trigger.dev/sdk";
export const client = new TriggerClient({
id: "my-app",
apiKey: process.env.TRIGGER_API_KEY,
apiUrl: process.env.TRIGGER_API_URL,
});
```
Replace **"my-app"** with an appropriate identifier for your project. The **apiKey** and **apiUrl** are obtained from the environment variables you set earlier.
4. **Example Directory Structure :**
```
project-root/
├── app/
├── routes/
├── trigger.ts
├── other files...
```
Replace **"my-app"** with an appropriate identifier for your project.
## Creating the API Route
+26
View File
@@ -0,0 +1,26 @@
The CLI `dev` command allows the Trigger.dev service to send messages to your site. This is required for registering Jobs, triggering them and running tasks. To achieve this it creates a tunnel (using [ngrok](https://ngrok.com/)) so Trigger.dev can send messages to your machine.
You should leave the `dev` command running when you're developing.
In a **new terminal window or tab** run:
<CodeGroup>
```bash npm
npx @trigger.dev/cli@latest dev
```
```bash pnpm
pnpm dlx @trigger.dev/cli@latest dev
```
```bash yarn
yarn dlx @trigger.dev/cli@latest dev
```
</CodeGroup>
<br />
<Note>
You can optionally pass the port if you're not running on the default port by adding
`--port 3001` to the end.
</Note>
+26
View File
@@ -0,0 +1,26 @@
```typescript
// Your first job
// This Job will be triggered by an event, log a joke to the console, and then wait 5 seconds before logging the punchline
client.defineJob({
// This is the unique identifier for your Job, it must be unique across all Jobs in your project
id: "example-job",
name: "Example Job: a joke with a delay",
version: "0.0.1",
// This is triggered by an event using eventTrigger. You can also trigger Jobs with webhooks, on schedules, and more: https://trigger.dev/docs/documentation/concepts/triggers/introduction
trigger: eventTrigger({
name: "example.event",
}),
run: async (payload, io, ctx) => {
// This logs a message to the console
await io.logger.info("🧪 Example Job: a joke with a delay");
await io.logger.info("How do you comfort a JavaScript bug?");
// This waits for 5 seconds, the second parameter is the number of seconds to wait, you can add delays of up to a year
await io.wait("Wait 5 seconds for the punchline...", 5);
await io.logger.info("You console it! 🤦");
await io.logger.info(
"✨ Congratulations, You just ran your first successful Trigger.dev Job! ✨"
);
// To learn how to write much more complex (and probably funnier) Jobs, check out our docs: https://trigger.dev/docs/documentation/guides/create-a-job
},
});
```
@@ -0,0 +1,18 @@
<Step title="Triggering the Job">
There are two way to trigger this Job.
1. Use the "Test" functionality in the dashboard.
2. Use the Trigger.dev API (either via our SDK or a web request)
#### "Testing" from the dashboard
Click into the Job and then open the "Test" tab. You should see this page:
![Test Job](/images/test-job.png)
This Job doesn't have a payload schema (meaning it takes an empty object), so you can simple click the "Run test" button.
**Congratulations, you should get redirected so you can see your first Run!**
</Step>
+80
View File
@@ -0,0 +1,80 @@
<Step title="Create a Trigger.dev account">
You can either:
- Use the [Trigger.dev Cloud](https://cloud.trigger.dev).
- Or [self-host](/documentation/guides/self-hosting) the service.
</Step>
<Step title="Create your first project">
Once you've created an account, follow the steps in the app to:
1. Complete your account details.
2. Create your first Organization and Project.
</Step>
<Step title="Getting an API key">
1. Go to the "Environments & API Keys" page in your project.
![Go to the Environments & API Keys page ](/images/environments-link.png)
2. Copy the `DEV` **SERVER** API key.
![API Keys](/images/api-keys.png)
</Step>
<Step title="Run the CLI `init` command">
The easiest way to get started it to use the CLI. It will add Trigger.dev to your existing project, setup a route and give you an example file.
In a terminal window run:
<CodeGroup>
```bash npm
npx @trigger.dev/cli@latest init
```
```bash pnpm
pnpm dlx @trigger.dev/cli@latest init
```
```bash yarn
yarn dlx @trigger.dev/cli@latest init
```
</CodeGroup>
It will ask you a couple of questions
1. Are you using the [Trigger.dev Cloud](https://cloud.trigger.dev) or [self-hosting](/documentation/guides/self-hosting)?
2. Enter your development API key. Enter the key you copied earlier.
</Step>
<Step title="Run your site">
Make sure your site is running locally, we will connect to it to register your Jobs.
<Warning>You must leave this running for the rest of the steps.</Warning>
<CodeGroup>
```bash npm
npm run dev
```
```bash pnpm
pnpm run dev
```
```bash yarn
yarn run dev
```
</CodeGroup>
</Step>
+20
View File
@@ -0,0 +1,20 @@
## What's next?
<CardGroup cols={2}>
<Card title="Write your first Job" icon="hexagon-plus" href="/documentation/guides/create-a-job">
A Guide for how to create your first real Job
</Card>
<Card
title="What is Trigger.dev"
icon="wand-magic-sparkles"
href="/documentation/concepts/what-is-triggerdotdev"
>
Learn more about how Trigger.dev works and how it can help you.
</Card>
<Card title="Examples" icon="slot-machine" href="/examples">
One of the quickest ways to learn how Trigger.dev works is to view some example Jobs.
</Card>
<Card title="Get help" icon="hire-a-helper" href="/documentation/get-help">
Struggling getting setup or have a question? We're here to help.
</Card>
</CardGroup>
+24
View File
@@ -0,0 +1,24 @@
## The two types of Run progress you can use
1. Automatic updates of Run and Task progress (no extra Job code required)
2. Explicitly created and updated `statuses` (more flexible and powerful)
### Automatic updates
These require no changes inside your Job code. You can receive:
- Info about an event you sent, including the Runs it triggered.
- The overall status of the Run (in progress, success and fail statuses).
- Metadata like start and completed times.
- The Run output (what is returned or an error that failed the Job)
- Information about the Tasks that have completed/failed/are running.
### Explicit `statuses`
You can create `statuses` in your Job code. This gives you fine grained control over what you want to expose.
It allows you to:
- Show exactly what you want in your UI (with as many statuses as you want).
- Pass arbitrary data to your UI, which you can use to render elements.
- Update existing elements in your UI as the progress of the run continues.
@@ -27,6 +27,10 @@ The `DEV` environment should only be used for local development. It's where you
<Snippet file="scheduled-dev-warning.mdx" />
### Staging
The `STAGING` environment is useful for testing your Jobs against your staging server, if you have one. STAGING works identically to PROD.
### Production
The `PROD` environment is where your Jobs will run in production. It's where you can run your Jobs against real data.
+1 -1
View File
@@ -4,7 +4,7 @@ title: "Limitations"
There are a few limitations that are important to understand.
In the current beta:
In the latest version:
- Runs on localhost are limited to 5 minutes.
- On long-running servers (not serverless) Runs can be retried erroneously.
@@ -12,7 +12,7 @@ Sometimes you don't know when you write the code what the trigger or schedule wi
```typescript
//1. create a DynamicSchedule
const dynamicSchedule = new DynamicSchedule(client, {
const dynamicSchedule = client.defineDynamicSchedule({
id: "dynamicinterval",
});
@@ -53,15 +53,18 @@ client.defineJob({
}),
}),
run: async (payload, io, ctx) => {
//6. Register the DynamicSchedule
await io.registerInterval("📆", dynamicSchedule, payload.userId, {
seconds: payload.seconds,
//6. Register the DynamicSchedule (this will automatically create a task)
await dynamicSchedule.register(userId, {
type: "cron",
options: {
cron: userSchedule,
},
});
await io.wait("wait", 60);
//7. Unregister the DynamicSchedule if you want
await io.unregisterInterval("❌📆", dynamicSchedule, payload.id);
//7. Unregister the DynamicSchedule if you want (this will automatically create a task)
await dynamicSchedule.unregister(userId);
},
});
```
@@ -70,7 +73,7 @@ client.defineJob({
```typescript
//1. create a DynamicTrigger
const dynamicOnIssueOpenedTrigger = new DynamicTrigger(client, {
const dynamicOnIssueOpenedTrigger = client.defineDynamicTrigger({
id: "github-issue-opened",
event: events.onIssueOpened,
source: github.sources.repo,
@@ -96,7 +99,7 @@ client.defineJob({
//3. Register the DynamicTrigger anywhere in your app
async function registerRepo(owner: string, repo: string) {
//the first param (key) should be unique
await dynamicOnIssueOpenedTrigger.register(`${owner}/${repo}`, {
await dynamicOnIssueOpenedTrigger.register(`${owner}-${repo}`, {
owner,
repo,
});
@@ -114,15 +117,10 @@ client.defineJob({
}),
run: async (payload, io, ctx) => {
//6. Register the dynamic trigger so you get notified when an issue is opened
return await io.registerTrigger(
"register-repo",
dynamicOnIssueOpenedTrigger,
payload.repository.name,
{
owner: payload.repository.owner.login,
repo: payload.repository.name,
}
);
await dynamicOnIssueOpenedTrigger.register(`${owner}-${repo}`, {
owner,
repo,
});
},
});
```

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