Compare commits

...

24 Commits

Author SHA1 Message Date
nicktrn ed03f4bc15 fix lockfile 2024-05-01 11:33:10 +01:00
github-actions[bot] 9feb0f70b0 chore: Update version for release (beta) (#1079)
Co-authored-by: github-actions[bot] <github-actions[bot]@users.noreply.github.com>
2024-05-01 11:32:15 +01:00
nicktrn 25de9e59af run workflows 2024-05-01 10:57:09 +01:00
nicktrn 83dc871550 add changeset 2024-05-01 10:48:27 +01:00
nicktrn 62700245a3 v3: fix consecutive waits (#1073)
* fix spacing for delete hints

* don't try to resume deps on wait resume

* sending duration wait resumes is not an error anymore

* set correct status with new wait resume flow

* cancel checkpoint schema v2

* don't mix messages and schemas

* prevent unintended case fallthrough in tree view

* completely switch to platform-led duration wait resumes

* prevent infinite restores

* some entries for the catalog

* add pg to additional packages

* add checkpoint safe timeout

* prevent duplicate spans after restore

* wait for post start

* add sdk version to deploy tab

* fail on impossible checkpoint scenarios

* remove debug logs
2024-05-01 10:33:30 +01:00
Matt Aitken 6ce820cb45 Test tasks that return different types 2024-04-30 19:07:44 +01:00
Matt Aitken 0f0a6884e8 Environment variables pasting uses dotenv (#1075)
* Add dotenv package to the webapp frontend, required some polyfills

* Use dotenv to parse the pasted env vars. Make the panel wider on larger screens
2024-04-30 14:23:49 +01:00
Matt Aitken 7ff8f0ebab Task and run page improvements (#1076)
* TaskListPresenter: if there are no tasks then don’t do stats queries

* RunListPresenter, use BasePresenter and the read replica

* Added populate script

* Simplified the Runs list query, added live timer

* Added TaskRun indexes for the RunList

* Status can’t be null now we’re using the TaskRun status

* Use defer so the page loads and shows a spinner

* Improved the loading style

* Get rid of latest run info from the tasks table super slow

* Fix for the activity graph tooltip getting clipped

* Added a code comment crediting the GitHub issue with the portal fix

* Add search to the tasks list

* Padding

* Fix for the schedules columns not being UTC

* Remove unused function
2024-04-30 14:19:50 +01:00
Matt Aitken 68455c796a Fixes/run filtered keyboard nav (#1074)
* Ensure each switch condition returns a state

* Fix for scrolling when filtered

* Fix for up/down navigation when filtered
2024-04-29 18:52:46 +01:00
Eric Allam b0a2c42e0e Add some additional logging around nacking messages 2024-04-29 17:53:38 +01:00
Andreas Thomas cac3c32f6a docs: add 'retry' import in code snippet (#1071)
* docs: add 'retry' import in code snippet

* docs: import all primitives
2024-04-29 17:42:36 +01:00
Matt Aitken ed8d24fd3d Run page performance improvements (#1072)
* lotsOfLogs task now outputs much larger logs

* Removed tree view collapse/expand animation

* JSDocs for useDebounce

* WIP moving filtering into the state

* Reworked the reducer to do the filtering
2024-04-29 17:05:26 +01:00
nicktrn d0ef36260a stop trying to pull init image when already present 2024-04-29 15:33:46 +01:00
Eric Allam 8eb68dd852 Fix pnpm lock file 2024-04-29 14:50:28 +01:00
github-actions[bot] a42037da03 chore: Update version for release (beta) (#1070)
Co-authored-by: github-actions[bot] <github-actions[bot]@users.noreply.github.com>
2024-04-29 14:48:01 +01:00
Eric Allam 43bc7ed94e Hoist uncaughtException handler to the top of workers to better report error messages 2024-04-29 14:22:55 +01:00
Matt Aitken 4fdb7f8288 Bulk insert environment variables (#1069)
* Fix for #1066. Correct environment username if dev

* WIP changing the form

* WIP on the form

* WIP on repository

* Adding environment variables en masse is working

* WIP on pasting

* useList hook with reducer

* Got bulk insert working with pasting… it was a pain

* Allow overwriting of values

* Set a max height on the new env var form
2024-04-28 18:31:31 +01:00
Matt Aitken 37b9b056c4 Upload payload packets outside of the db transaction 2024-04-28 18:18:33 +01:00
Matt Aitken 801c86bf73 Added search to the test tasks list 2024-04-28 16:10:46 +01:00
Matt Aitken a1de11a001 Default the test page to the “DEV” tab 2024-04-28 15:59:45 +01:00
Praveen Pendyala 29e9e372ee Added example of github workflow for deploy to staging (#1068) 2024-04-26 18:39:34 +01:00
Eric Allam affc128161 Try to fix prisma generate errors 2024-04-26 18:28:17 +01:00
Eric Allam 40ba8ad0ee Fix the test page when a scheduled task has recent run payloads that aren’t scheduled payloads 2024-04-26 18:20:04 +01:00
Matt Aitken 96168eb383 Fixes for test page problems 2024-04-26 17:45:06 +01:00
136 changed files with 3174 additions and 1510 deletions
+5
View File
@@ -0,0 +1,5 @@
---
"trigger.dev": patch
---
Hoist uncaughtException handler to the top of workers to better report error messages
+6
View File
@@ -0,0 +1,6 @@
---
"trigger.dev": patch
"@trigger.dev/core": patch
---
Fix issues with consecutive waits
+2
View File
@@ -46,6 +46,7 @@
"changesets": [
"angry-eagles-trade",
"beige-pens-dance",
"big-tomatoes-deliver",
"breezy-gorillas-mate",
"chilled-hornets-move",
"clean-pianos-listen",
@@ -66,6 +67,7 @@
"lemon-jobs-repair",
"light-bulldogs-press",
"light-dragons-complain",
"little-crabs-cross",
"loud-actors-remember",
"many-ligers-pump",
"mighty-camels-joke",
+3 -1
View File
@@ -4,13 +4,15 @@ on:
jobs:
publish:
runs-on: ubuntu-latest
env:
PRISMA_ENGINES_CHECKSUM_IGNORE_MISSING: 1
outputs:
version: ${{ steps.get_version.outputs.version }}
short_sha: ${{ steps.get_commit.outputs.sha_short }}
steps:
- name: Setup Depot CLI
uses: depot/setup-action@v1
- name: ⬇️ Checkout repo
uses: actions/checkout@v3
with:
+17 -5
View File
@@ -157,16 +157,18 @@ class Checkpointer {
return this.#abortControllers.has(runId);
}
cancelCheckpoint(runId: string) {
cancelCheckpoint(runId: string): boolean {
const controller = this.#abortControllers.get(runId);
if (!controller) {
logger.debug("Nothing to cancel", { runId });
return;
return false;
}
controller.abort("cancelCheckpointing()");
this.#abortControllers.delete(runId);
return true;
}
async #checkpointAndPush({
@@ -725,10 +727,18 @@ class TaskCoordinator {
checkpointable.resolve();
});
socket.on("CANCEL_CHECKPOINT", async (message) => {
socket.on("CANCEL_CHECKPOINT", async (message, callback) => {
logger.log("[CANCEL_CHECKPOINT]", message);
this.#cancelCheckpoint(socket.data.runId);
if (message.version === "v1") {
this.#cancelCheckpoint(socket.data.runId);
// v1 has no callback
return;
}
const checkpointCanceled = this.#cancelCheckpoint(socket.data.runId);
callback({ version: "v2", checkpointCanceled });
});
socket.on("WAIT_FOR_DURATION", async (message, callback) => {
@@ -933,7 +943,9 @@ class TaskCoordinator {
}
// Cancel checkpointing procedure
this.#checkpointer.cancelCheckpoint(runId);
const checkpointCanceled = this.#checkpointer.cancelCheckpoint(runId);
return checkpointCanceled;
}
#createHttpServer() {
+1
View File
@@ -212,6 +212,7 @@ class KubernetesTaskOperations implements TaskOperations {
{
name: "populate-taskinfo",
image: "docker.io/library/busybox",
imagePullPolicy: "IfNotPresent",
command: ["/bin/sh", "-c"],
args: ["printenv COORDINATOR_HOST | tee /etc/taskinfo/coordinator-host"],
env: [
+2 -1
View File
@@ -17,4 +17,5 @@ build-storybook.log
.storybook-out
storybook-static
/prisma/seed.js
/prisma/seed.js
/prisma/populate.js
@@ -1,5 +1,9 @@
import { Paragraph } from "./Paragraph";
export function Hint({ children }: { children: React.ReactNode }) {
return <Paragraph variant="extra-small">{children}</Paragraph>;
export function Hint({ children, className }: { children: React.ReactNode; className?: string }) {
return (
<Paragraph variant="extra-small" className={className}>
{children}
</Paragraph>
);
}
@@ -0,0 +1,103 @@
import type { VirtualElement as IVirtualElement } from "@popperjs/core";
import { ReactNode, useEffect, useState } from "react";
import { createPortal } from "react-dom";
import { usePopper } from "react-popper";
import { useEvent } from "react-use";
import useLazyRef from "~/hooks/useLazyRef";
// Recharts 3.x will have portal support, but until then we're using this:
//https://github.com/recharts/recharts/issues/2458#issuecomment-1063463873
export interface PopperPortalProps {
active?: boolean;
children: ReactNode;
}
export default function TooltipPortal({ active = true, children }: PopperPortalProps) {
const [portalElement, setPortalElement] = useState<HTMLDivElement>();
const [popperElement, setPopperElement] = useState<HTMLDivElement | null>();
const virtualElementRef = useLazyRef(() => new VirtualElement());
const { styles, attributes, update } = usePopper(
virtualElementRef.current,
popperElement,
POPPER_OPTIONS
);
useEffect(() => {
const el = document.createElement("div");
document.body.appendChild(el);
setPortalElement(el);
return () => el.remove();
}, []);
useEvent("mousemove", ({ clientX: x, clientY: y }) => {
virtualElementRef.current?.update(x, y);
if (!active) return;
update?.();
});
useEffect(() => {
if (!active) return;
update?.();
}, [active, update]);
if (!portalElement) return null;
return createPortal(
<div
ref={setPopperElement}
{...attributes.popper}
style={{
...styles.popper,
zIndex: 1000,
display: active ? "block" : "none",
}}
>
{children}
</div>,
portalElement
);
}
class VirtualElement implements IVirtualElement {
private rect = {
width: 0,
height: 0,
top: 0,
right: 0,
bottom: 0,
left: 0,
x: 0,
y: 0,
toJSON() {
return this;
},
};
update(x: number, y: number) {
this.rect.y = y;
this.rect.top = y;
this.rect.bottom = y;
this.rect.x = x;
this.rect.left = x;
this.rect.right = x;
}
getBoundingClientRect(): DOMRect {
return this.rect;
}
}
const POPPER_OPTIONS: Parameters<typeof usePopper>[2] = {
placement: "right-start",
modifiers: [
{
name: "offset",
options: {
offset: [8, 8],
},
},
],
};
@@ -1,10 +1,9 @@
import { VirtualItem, Virtualizer, useVirtualizer } from "@tanstack/react-virtual";
import { motion } from "framer-motion";
import { MutableRefObject, RefObject, useCallback, useEffect, useReducer, useRef } from "react";
import { UnmountClosed } from "react-collapse";
import { cn } from "~/utils/cn";
import { NodeState, NodesState, reducer } from "./reducer";
import { applyFilterToState, concreteStateFromInput, selectedIdFromState } from "./utils";
import { concreteStateFromInput, selectedIdFromState } from "./utils";
export type TreeViewProps<TData> = {
tree: FlatTree<TData>;
@@ -104,23 +103,22 @@ export function TreeView<TData>({
if (!node) return null;
const state = nodes[node.id];
if (!state) return null;
if (!state.visible) return null;
return (
<div
key={node.id}
data-index={virtualItem.index}
ref={virtualizer.measureElement}
className="overflow-clip [&_.ReactCollapse--collapse]:transition-all"
className="overflow-clip"
{...getNodeProps(node.id)}
>
<UnmountClosed key={node.id} isOpened={state.visible}>
{renderNode({
node,
state,
index: virtualItem.index,
virtualizer: virtualizer,
virtualItem,
})}
</UnmountClosed>
{renderNode({
node,
state,
index: virtualItem.index,
virtualizer: virtualizer,
virtualItem,
})}
</div>
);
})}
@@ -130,19 +128,23 @@ export function TreeView<TData>({
);
}
type TreeStateHookProps<TData> = {
export type Filter<TData, TFilterValue> = {
value?: TFilterValue;
fn: (value: TFilterValue, node: FlatTreeItem<TData>) => boolean;
};
type TreeStateHookProps<TData, TFilterValue> = {
tree: FlatTree<TData>;
selectedId?: string;
collapsedIds?: string[];
onSelectedIdChanged?: (selectedId: string | undefined) => void;
onCollapsedIdsChanged?: (collapsedIds: string[]) => void;
estimatedRowHeight: (params: {
node: FlatTreeItem<TData>;
state: NodeState;
index: number;
}) => number;
parentRef: RefObject<any>;
filter?: (node: FlatTreeItem<TData>) => boolean;
filter?: Filter<TData, TFilterValue>;
};
//this is so Framer Motion can be used to render the components
@@ -178,24 +180,24 @@ export type UseTreeStateOutput = {
scrollToNode: (id: string) => void;
};
export function useTree<TData>({
export function useTree<TData, TFilterValue>({
tree,
selectedId,
collapsedIds,
onSelectedIdChanged,
onCollapsedIdsChanged,
parentRef,
estimatedRowHeight,
filter,
}: TreeStateHookProps<TData>): UseTreeStateOutput {
}: TreeStateHookProps<TData, TFilterValue>): UseTreeStateOutput {
const previousNodeCount = useRef(tree.length);
const previousSelectedId = useRef<string | undefined>(selectedId);
const [state, dispatch] = useReducer(
reducer,
concreteStateFromInput({ tree, selectedId, collapsedIds })
concreteStateFromInput({ tree, selectedId, collapsedIds, filter })
);
//fire onSelectedIdChanged()
useEffect(() => {
const selectedId = selectedIdFromState(state.nodes);
if (selectedId !== previousSelectedId.current) {
@@ -204,12 +206,7 @@ export function useTree<TData>({
}
}, [state.changes.selectedId]);
useEffect(() => {
if (state.changes.collapsedIds) {
onCollapsedIdsChanged?.(state.changes.collapsedIds);
}
}, [state.changes.collapsedIds]);
//update tree when the number of nodes changes
useEffect(() => {
if (tree.length !== previousNodeCount.current) {
previousNodeCount.current = tree.length;
@@ -217,9 +214,25 @@ export function useTree<TData>({
}
}, [previousNodeCount.current, tree.length]);
//update the filter, if it's changed
const previousFilter = useRef(filter);
useEffect(() => {
//check if the value (not reference) of the filter is the same
const previousValue = previousFilter.current
? JSON.stringify(previousFilter.current.value)
: undefined;
const newValue = filter ? JSON.stringify(filter.value) : undefined;
previousFilter.current = filter;
if (previousValue !== newValue) {
dispatch({ type: "UPDATE_FILTER", payload: { filter } });
}
}, [filter?.value]);
const virtualizer = useVirtualizer({
count: tree.length,
getItemKey: (index) => tree[index].id,
count: state.visibleNodeIds.length,
getItemKey: (index) => state.visibleNodeIds[index],
getScrollElement: () => parentRef.current,
estimateSize: (index: number) => {
return estimatedRowHeight({
@@ -233,7 +246,7 @@ export function useTree<TData>({
const scrollToNodeFn = useCallback(
(id: string) => {
const itemIndex = tree.findIndex((node) => node.id === id);
const itemIndex = state.visibleNodeIds.findIndex((n) => n === id);
if (itemIndex !== -1) {
virtualizer.scrollToIndex(itemIndex, { align: "auto" });
@@ -269,21 +282,21 @@ export function useTree<TData>({
const expandNode = useCallback(
(id: string, scrollToNode = true) => {
dispatch({ type: "EXPAND_NODE", payload: { id, tree, scrollToNode, scrollToNodeFn } });
dispatch({ type: "EXPAND_NODE", payload: { id, scrollToNode, scrollToNodeFn } });
},
[state]
);
const collapseNode = useCallback(
(id: string) => {
dispatch({ type: "COLLAPSE_NODE", payload: { id, tree } });
dispatch({ type: "COLLAPSE_NODE", payload: { id } });
},
[state]
);
const toggleExpandNode = useCallback(
(id: string, scrollToNode = true) => {
dispatch({ type: "TOGGLE_EXPAND_NODE", payload: { id, tree, scrollToNode, scrollToNodeFn } });
dispatch({ type: "TOGGLE_EXPAND_NODE", payload: { id, scrollToNode, scrollToNodeFn } });
},
[state]
);
@@ -292,7 +305,7 @@ export function useTree<TData>({
(scrollToNode = true) => {
dispatch({
type: "SELECT_FIRST_VISIBLE_NODE",
payload: { tree, scrollToNode, scrollToNodeFn },
payload: { scrollToNode, scrollToNodeFn },
});
},
[tree, state]
@@ -302,7 +315,7 @@ export function useTree<TData>({
(scrollToNode = true) => {
dispatch({
type: "SELECT_LAST_VISIBLE_NODE",
payload: { tree, scrollToNode, scrollToNodeFn },
payload: { scrollToNode, scrollToNodeFn },
});
},
[tree, state]
@@ -312,7 +325,7 @@ export function useTree<TData>({
(scrollToNode = true) => {
dispatch({
type: "SELECT_NEXT_VISIBLE_NODE",
payload: { tree, scrollToNode, scrollToNodeFn },
payload: { scrollToNode, scrollToNodeFn },
});
},
[state]
@@ -322,7 +335,7 @@ export function useTree<TData>({
(scrollToNode = true) => {
dispatch({
type: "SELECT_PREVIOUS_VISIBLE_NODE",
payload: { tree, scrollToNode, scrollToNodeFn },
payload: { scrollToNode, scrollToNodeFn },
});
},
[state]
@@ -332,7 +345,7 @@ export function useTree<TData>({
(scrollToNode = true) => {
dispatch({
type: "SELECT_PARENT_NODE",
payload: { tree, scrollToNode, scrollToNodeFn },
payload: { scrollToNode, scrollToNodeFn },
});
},
[state]
@@ -340,35 +353,35 @@ export function useTree<TData>({
const expandAllBelowDepth = useCallback(
(depth: number) => {
dispatch({ type: "EXPAND_ALL_BELOW_DEPTH", payload: { tree, depth } });
dispatch({ type: "EXPAND_ALL_BELOW_DEPTH", payload: { depth } });
},
[state]
);
const collapseAllBelowDepth = useCallback(
(depth: number) => {
dispatch({ type: "COLLAPSE_ALL_BELOW_DEPTH", payload: { tree, depth } });
dispatch({ type: "COLLAPSE_ALL_BELOW_DEPTH", payload: { depth } });
},
[state]
);
const expandLevel = useCallback(
(level: number) => {
dispatch({ type: "EXPAND_LEVEL", payload: { tree, level } });
dispatch({ type: "EXPAND_LEVEL", payload: { level } });
},
[state]
);
const collapseLevel = useCallback(
(level: number) => {
dispatch({ type: "COLLAPSE_LEVEL", payload: { tree, level } });
dispatch({ type: "COLLAPSE_LEVEL", payload: { level } });
},
[state]
);
const toggleExpandLevel = useCallback(
(level: number) => {
dispatch({ type: "TOGGLE_EXPAND_LEVEL", payload: { tree, level } });
dispatch({ type: "TOGGLE_EXPAND_LEVEL", payload: { level } });
},
[state]
);
@@ -480,7 +493,7 @@ export function useTree<TData>({
return {
selected: selectedIdFromState(state.nodes),
nodes: filter ? applyFilterToState(tree, state.nodes, filter) : state.nodes,
nodes: state.nodes,
getTreeProps,
getNodeProps,
selectNode,
@@ -1,5 +1,7 @@
import { FlatTree } from "./TreeView";
import assertNever from "assert-never";
import { Filter, FlatTree } from "./TreeView";
import {
applyFilterToState,
applyVisibility,
collapsedIdsFromState,
concreteStateFromInput,
@@ -18,12 +20,15 @@ export type NodeState = {
export type Changes = {
selectedId: string | undefined;
collapsedIds: string[] | undefined;
};
export type TreeState = {
tree: FlatTree<any>;
nodes: NodesState;
filteredNodes: NodesState;
changes: Changes;
filter: Filter<any, any> | undefined;
visibleNodeIds: string[];
};
export type NodesState = Record<string, NodeState>;
@@ -71,7 +76,6 @@ type ExpandNodeAction = {
type: "EXPAND_NODE";
payload: {
id: string;
tree: FlatTree<any>;
} & WithScrollToNode;
};
@@ -79,7 +83,6 @@ type CollapseNodeAction = {
type: "COLLAPSE_NODE";
payload: {
id: string;
tree: FlatTree<any>;
};
};
@@ -87,7 +90,6 @@ type ToggleExpandNodeAction = {
type: "TOGGLE_EXPAND_NODE";
payload: {
id: string;
tree: FlatTree<any>;
} & WithScrollToNode;
};
@@ -95,7 +97,6 @@ type ExpandAllBelowDepthAction = {
type: "EXPAND_ALL_BELOW_DEPTH";
payload: {
depth: number;
tree: FlatTree<any>;
};
};
@@ -103,7 +104,6 @@ type CollapseAllBelowDepthAction = {
type: "COLLAPSE_ALL_BELOW_DEPTH";
payload: {
depth: number;
tree: FlatTree<any>;
};
};
@@ -111,7 +111,6 @@ type ExpandLevelAction = {
type: "EXPAND_LEVEL";
payload: {
level: number;
tree: FlatTree<any>;
};
};
@@ -119,7 +118,6 @@ type CollapseLevelAction = {
type: "COLLAPSE_LEVEL";
payload: {
level: number;
tree: FlatTree<any>;
};
};
@@ -127,43 +125,39 @@ type ToggleExpandLevelAction = {
type: "TOGGLE_EXPAND_LEVEL";
payload: {
level: number;
tree: FlatTree<any>;
};
};
type SelectFirstVisibleNodeAction = {
type: "SELECT_FIRST_VISIBLE_NODE";
payload: {
tree: FlatTree<any>;
} & WithScrollToNode;
payload: {} & WithScrollToNode;
};
type SelectLastVisibleNodeAction = {
type: "SELECT_LAST_VISIBLE_NODE";
payload: {
tree: FlatTree<any>;
} & WithScrollToNode;
payload: {} & WithScrollToNode;
};
type SelectNextVisibleNodeAction = {
type: "SELECT_NEXT_VISIBLE_NODE";
payload: {
tree: FlatTree<any>;
} & WithScrollToNode;
payload: {} & WithScrollToNode;
};
type SelectPreviousVisibleNodeAction = {
type: "SELECT_PREVIOUS_VISIBLE_NODE";
payload: {
tree: FlatTree<any>;
} & WithScrollToNode;
payload: {} & WithScrollToNode;
};
type SelectParentNodeAction = {
type: "SELECT_PARENT_NODE";
payload: {} & WithScrollToNode;
};
type UpdateFilterAction = {
type: "UPDATE_FILTER";
payload: {
tree: FlatTree<any>;
} & WithScrollToNode;
filter: Filter<any, any> | undefined;
};
};
export type Action =
@@ -184,7 +178,8 @@ export type Action =
| SelectLastVisibleNodeAction
| SelectNextVisibleNodeAction
| SelectPreviousVisibleNodeAction
| SelectParentNodeAction;
| SelectParentNodeAction
| UpdateFilterAction;
export function reducer(state: TreeState, action: Action): TreeState {
switch (action.type) {
@@ -204,7 +199,12 @@ export function reducer(state: TreeState, action: Action): TreeState {
action.payload.scrollToNodeFn(action.payload.id);
}
return { nodes: newNodes, changes: generateChanges(state.nodes, newNodes) };
return applyFilterToState({
...state,
tree: state.tree,
nodes: newNodes,
changes: generateChanges(state.nodes, newNodes),
});
}
case "DESELECT_NODE": {
const nodes = {
@@ -212,28 +212,36 @@ export function reducer(state: TreeState, action: Action): TreeState {
[action.payload.id]: { ...state.nodes[action.payload.id], selected: false },
};
return { nodes, changes: generateChanges(state.nodes, nodes) };
return applyFilterToState({
...state,
nodes,
changes: generateChanges(state.nodes, nodes),
});
}
case "DESELECT_ALL_NODES": {
const nodes = Object.fromEntries(
Object.entries(state.nodes).map(([key, value]) => [key, { ...value, selected: false }])
);
return { nodes, changes: generateChanges(state.nodes, nodes) };
return applyFilterToState({
...state,
nodes,
changes: generateChanges(state.nodes, nodes),
});
}
case "TOGGLE_NODE_SELECTION": {
const currentlySelected = state.nodes[action.payload.id]?.selected ?? false;
if (currentlySelected) {
return reducer(state, { type: "DESELECT_NODE", payload: { id: action.payload.id } });
} else {
return reducer(state, {
type: "SELECT_NODE",
payload: {
id: action.payload.id,
scrollToNode: action.payload.scrollToNode,
scrollToNodeFn: action.payload.scrollToNodeFn,
},
});
}
return reducer(state, {
type: "SELECT_NODE",
payload: {
id: action.payload.id,
scrollToNode: action.payload.scrollToNode,
scrollToNodeFn: action.payload.scrollToNodeFn,
},
});
}
case "EXPAND_NODE": {
const newNodes = {
@@ -245,37 +253,44 @@ export function reducer(state: TreeState, action: Action): TreeState {
action.payload.scrollToNodeFn(action.payload.id);
}
const visibleNodes = applyVisibility(action.payload.tree, newNodes);
return { nodes: visibleNodes, changes: generateChanges(state.nodes, visibleNodes) };
const visibleNodes = applyVisibility(state.tree, newNodes);
return applyFilterToState({
...state,
nodes: visibleNodes,
changes: generateChanges(state.nodes, visibleNodes),
});
}
case "COLLAPSE_NODE": {
const visibleNodes = applyVisibility(action.payload.tree, {
const visibleNodes = applyVisibility(state.tree, {
...state.nodes,
[action.payload.id]: { ...state.nodes[action.payload.id], expanded: false },
});
return { nodes: visibleNodes, changes: generateChanges(state.nodes, visibleNodes) };
return applyFilterToState({
...state,
nodes: visibleNodes,
changes: generateChanges(state.nodes, visibleNodes),
});
}
case "TOGGLE_EXPAND_NODE": {
const currentlyExpanded = state.nodes[action.payload.id]?.expanded ?? true;
if (currentlyExpanded) {
return reducer(state, {
type: "COLLAPSE_NODE",
payload: { id: action.payload.id, tree: action.payload.tree },
});
} else {
return reducer(state, {
type: "EXPAND_NODE",
payload: {
id: action.payload.id,
tree: action.payload.tree,
scrollToNode: action.payload.scrollToNode,
scrollToNodeFn: action.payload.scrollToNodeFn,
},
payload: { id: action.payload.id },
});
}
return reducer(state, {
type: "EXPAND_NODE",
payload: {
id: action.payload.id,
scrollToNode: action.payload.scrollToNode,
scrollToNodeFn: action.payload.scrollToNodeFn,
},
});
}
case "EXPAND_ALL_BELOW_DEPTH": {
const nodesToExpand = action.payload.tree.filter(
const nodesToExpand = state.tree.filter(
(n) => n.level >= action.payload.depth && n.hasChildren
);
@@ -289,11 +304,15 @@ export function reducer(state: TreeState, action: Action): TreeState {
])
);
const visibleNodes = applyVisibility(action.payload.tree, newNodes);
return { nodes: visibleNodes, changes: generateChanges(state.nodes, visibleNodes) };
const visibleNodes = applyVisibility(state.tree, newNodes);
return applyFilterToState({
...state,
nodes: visibleNodes,
changes: generateChanges(state.nodes, visibleNodes),
});
}
case "COLLAPSE_ALL_BELOW_DEPTH": {
const nodesToCollapse = action.payload.tree.filter(
const nodesToCollapse = state.tree.filter(
(n) => n.level >= action.payload.depth && n.hasChildren
);
@@ -307,11 +326,15 @@ export function reducer(state: TreeState, action: Action): TreeState {
])
);
const visibleNodes = applyVisibility(action.payload.tree, newNodes);
return { nodes: visibleNodes, changes: generateChanges(state.nodes, visibleNodes) };
const visibleNodes = applyVisibility(state.tree, newNodes);
return applyFilterToState({
...state,
nodes: visibleNodes,
changes: generateChanges(state.nodes, visibleNodes),
});
}
case "EXPAND_LEVEL": {
const nodesToExpand = action.payload.tree.filter(
const nodesToExpand = state.tree.filter(
(n) => n.level <= action.payload.level && n.hasChildren
);
@@ -325,11 +348,15 @@ export function reducer(state: TreeState, action: Action): TreeState {
])
);
const visibleNodes = applyVisibility(action.payload.tree, newNodes);
return { nodes: visibleNodes, changes: generateChanges(state.nodes, visibleNodes) };
const visibleNodes = applyVisibility(state.tree, newNodes);
return applyFilterToState({
...state,
nodes: visibleNodes,
changes: generateChanges(state.nodes, visibleNodes),
});
}
case "COLLAPSE_LEVEL": {
const nodesToCollapse = action.payload.tree.filter(
const nodesToCollapse = state.tree.filter(
(n) => n.level === action.payload.level && n.hasChildren
);
@@ -343,13 +370,17 @@ export function reducer(state: TreeState, action: Action): TreeState {
])
);
const visibleNodes = applyVisibility(action.payload.tree, newNodes);
return { nodes: visibleNodes, changes: generateChanges(state.nodes, visibleNodes) };
const visibleNodes = applyVisibility(state.tree, newNodes);
return applyFilterToState({
...state,
nodes: visibleNodes,
changes: generateChanges(state.nodes, visibleNodes),
});
}
case "TOGGLE_EXPAND_LEVEL": {
//first get the first item at that level in the tree. If it is expanded, collapse all nodes at that level
//if it is collapsed, expand all nodes at that level
const nodesAtLevel = action.payload.tree.filter(
const nodesAtLevel = state.tree.filter(
(n) => n.level === action.payload.level && n.hasChildren
);
const firstNode = nodesAtLevel[0];
@@ -364,21 +395,19 @@ export function reducer(state: TreeState, action: Action): TreeState {
type: "COLLAPSE_LEVEL",
payload: {
level: action.payload.level,
tree: action.payload.tree,
},
});
} else {
return reducer(state, {
type: "EXPAND_LEVEL",
payload: {
level: action.payload.level,
tree: action.payload.tree,
},
});
}
return reducer(state, {
type: "EXPAND_LEVEL",
payload: {
level: action.payload.level,
},
});
}
case "SELECT_FIRST_VISIBLE_NODE": {
const node = firstVisibleNode(action.payload.tree, state.nodes);
const node = firstVisibleNode(state.tree, state.filteredNodes);
if (node) {
return reducer(state, {
type: "SELECT_NODE",
@@ -389,9 +418,11 @@ export function reducer(state: TreeState, action: Action): TreeState {
},
});
}
return state;
}
case "SELECT_LAST_VISIBLE_NODE": {
const node = lastVisibleNode(action.payload.tree, state.nodes);
const node = lastVisibleNode(state.tree, state.filteredNodes);
if (node) {
return reducer(state, {
type: "SELECT_NODE",
@@ -402,6 +433,8 @@ export function reducer(state: TreeState, action: Action): TreeState {
},
});
}
return state;
}
case "SELECT_NEXT_VISIBLE_NODE": {
const selected = selectedIdFromState(state.nodes);
@@ -409,14 +442,13 @@ export function reducer(state: TreeState, action: Action): TreeState {
return reducer(state, {
type: "SELECT_FIRST_VISIBLE_NODE",
payload: {
tree: action.payload.tree,
scrollToNode: action.payload.scrollToNode,
scrollToNodeFn: action.payload.scrollToNodeFn,
},
});
}
const visible = visibleNodes(action.payload.tree, state.nodes);
const visible = visibleNodes(state.tree, state.filteredNodes);
const selectedIndex = visible.findIndex((node) => node.id === selected);
const nextNode = visible[selectedIndex + 1];
if (nextNode) {
@@ -429,6 +461,8 @@ export function reducer(state: TreeState, action: Action): TreeState {
},
});
}
return state;
}
case "SELECT_PREVIOUS_VISIBLE_NODE": {
const selected = selectedIdFromState(state.nodes);
@@ -437,16 +471,15 @@ export function reducer(state: TreeState, action: Action): TreeState {
return reducer(state, {
type: "SELECT_FIRST_VISIBLE_NODE",
payload: {
tree: action.payload.tree,
scrollToNode: action.payload.scrollToNode,
scrollToNodeFn: action.payload.scrollToNodeFn,
},
});
}
const visible = visibleNodes(action.payload.tree, state.nodes);
const visible = visibleNodes(state.tree, state.filteredNodes);
const selectedIndex = visible.findIndex((node) => node.id === selected);
const previousNode = visible[selectedIndex - 1];
const previousNode = visible[Math.max(0, selectedIndex - 1)];
if (previousNode) {
return reducer(state, {
type: "SELECT_NODE",
@@ -467,19 +500,18 @@ export function reducer(state: TreeState, action: Action): TreeState {
return reducer(state, {
type: "SELECT_FIRST_VISIBLE_NODE",
payload: {
tree: action.payload.tree,
scrollToNode: action.payload.scrollToNode,
scrollToNodeFn: action.payload.scrollToNodeFn,
},
});
}
const selectedNode = action.payload.tree.find((node) => node.id === selected);
const selectedNode = state.tree.find((node) => node.id === selected);
if (!selectedNode) {
return state;
}
const parentNode = action.payload.tree.find((node) => node.id === selectedNode.parentId);
const parentNode = state.tree.find((node) => node.id === selectedNode.parentId);
if (parentNode) {
return reducer(state, {
type: "SELECT_NODE",
@@ -498,12 +530,23 @@ export function reducer(state: TreeState, action: Action): TreeState {
const selectedId = selectedIdFromState(state.nodes);
const collapsedIds = collapsedIdsFromState(state.nodes);
const newState = concreteStateFromInput({
...state,
tree: action.payload.tree,
selectedId,
collapsedIds,
});
return newState;
}
case "UPDATE_FILTER": {
const newState = applyFilterToState({
...state,
filter: action.payload.filter,
});
return newState;
}
default: {
assertNever(action);
}
}
throw new Error(`Unhandled action type: ${(action as any).type}`);
@@ -1,4 +1,4 @@
import { FlatTree, FlatTreeItem } from "./TreeView";
import { Filter, FlatTree, FlatTreeItem } from "./TreeView";
import { Changes, NodeState, NodesState, TreeState } from "./reducer";
type PartialNodeState = Record<string, Partial<NodeState>>;
@@ -8,10 +8,12 @@ const defaultExpanded = true;
export function concreteStateFromInput({
tree,
filter,
selectedId,
collapsedIds,
}: {
tree: FlatTree<any>;
filter: Filter<any, any> | undefined;
selectedId: string | undefined;
collapsedIds: string[] | undefined;
}): TreeState {
@@ -35,10 +37,15 @@ export function concreteStateFromInput({
}
}
}
const nodes = concreteStateFromPartialState(tree, state);
return {
nodes: concreteStateFromPartialState(tree, state),
changes: { selectedId, collapsedIds: [] },
tree,
nodes,
changes: { selectedId },
filter,
filteredNodes: nodes,
visibleNodeIds: visibleNodes(tree, nodes).map((node) => node.id),
};
}
@@ -82,28 +89,48 @@ export function selectedIdFromState(state: NodesState): string | undefined {
return selected?.[0];
}
export function applyFilterToState<TData>(
tree: FlatTree<TData>,
inputNodes: NodesState,
filter: (node: FlatTreeItem<TData>) => boolean
): NodesState {
export function applyFilterToState<TData>({
tree,
nodes,
filter,
visibleNodeIds,
changes,
}: TreeState): TreeState {
if (!filter || !filter.value) {
return {
tree,
nodes,
filteredNodes: nodes,
changes,
filter,
visibleNodeIds: visibleNodes(tree, nodes).map((node) => node.id),
};
}
//we need to do two passes, first collect all the nodes that are results
const newFilteredOut = new Set<string>();
for (const node of tree) {
if (!filter(node)) {
if (!filter.fn(filter.value, node)) {
newFilteredOut.add(node.id);
}
}
//nothing is filtered out
if (newFilteredOut.size === 0) {
return inputNodes;
return {
tree,
nodes,
filteredNodes: nodes,
changes,
filter,
visibleNodeIds: visibleNodes(tree, nodes).map((node) => node.id),
};
}
//copy of nodes
const nodes = { ...inputNodes };
const filteredNodes = { ...nodes };
const selected = selectedIdFromState(nodes);
const selected = selectedIdFromState(filteredNodes);
const visible = new Set<string>();
const expanded = new Set<string>();
@@ -148,28 +175,35 @@ export function applyFilterToState<TData>(
//now set the visibility and expanded state
for (const id of hidden) {
nodes[id] = { ...nodes[id], visible: false };
filteredNodes[id] = { ...filteredNodes[id], visible: false };
}
for (const id of visible) {
nodes[id] = { ...nodes[id], visible: true };
filteredNodes[id] = { ...filteredNodes[id], visible: true };
}
for (const id of collapsed) {
nodes[id] = { ...nodes[id], expanded: false };
filteredNodes[id] = { ...filteredNodes[id], expanded: false };
}
for (const id of expanded) {
nodes[id] = { ...nodes[id], expanded: true };
filteredNodes[id] = { ...filteredNodes[id], expanded: true };
}
if (selected) {
if (visible.has(selected)) {
nodes[selected] = { ...nodes[selected], selected: true };
filteredNodes[selected] = { ...filteredNodes[selected], selected: true };
} else {
nodes[selected] = { ...nodes[selected], selected: false };
filteredNodes[selected] = { ...filteredNodes[selected], selected: false };
}
}
return nodes;
return {
tree,
nodes,
filteredNodes,
changes,
filter,
visibleNodeIds: visibleNodes(tree, filteredNodes).map((node) => node.id),
};
}
export function visibleNodes(tree: FlatTree<any>, nodes: NodesState) {
@@ -215,6 +249,5 @@ export function generateChanges(a: NodesState, b: NodesState): Changes {
return {
selectedId: selectedIdA !== selectedIdB ? selectedIdB : undefined,
collapsedIds: collapsedChanges.length > 0 ? collapsedChanges : undefined,
};
}
@@ -1,18 +1,14 @@
import { formatDuration } from "@trigger.dev/core/v3";
import { useState, useEffect } from "react";
import { Paragraph } from "~/components/primitives/Paragraph";
import { cn } from "~/utils/cn";
import { useEffect, useState } from "react";
export function LiveTimer({
startTime,
endTime,
updateInterval = 250,
className,
}: {
startTime: Date;
endTime?: Date;
updateInterval?: number;
className?: string;
}) {
const [now, setNow] = useState<Date>();
@@ -30,13 +26,13 @@ export function LiveTimer({
}, [startTime]);
return (
<Paragraph variant="extra-small" className={cn("whitespace-nowrap tabular-nums", className)}>
<>
{formatDuration(startTime, now, {
style: "short",
maxDecimalPoints: 0,
units: ["d", "h", "m", "s"],
})}
</Paragraph>
</>
);
}
@@ -38,6 +38,15 @@ export const RUNNING_STATUSES: TaskRunStatus[] = [
"WAITING_TO_RESUME",
];
export const FINISHED_STATUSES: TaskRunStatus[] = [
"COMPLETED_SUCCESSFULLY",
"CANCELED",
"COMPLETED_WITH_ERRORS",
"INTERRUPTED",
"SYSTEM_FAILURE",
"CRASHED",
];
export function descriptionForTaskRunStatus(status: TaskRunStatus): string {
return taskRunStatusDescriptions[status];
}
@@ -3,7 +3,6 @@ import { StopIcon } from "@heroicons/react/24/outline";
import { BeakerIcon, BookOpenIcon, CheckIcon } from "@heroicons/react/24/solid";
import { useLocation } from "@remix-run/react";
import { formatDuration } from "@trigger.dev/core/v3";
import { User } from "@trigger.dev/database";
import { Button, LinkButton } from "~/components/primitives/Buttons";
import { Dialog, DialogTrigger } from "~/components/primitives/Dialog";
import { useEnvironments } from "~/hooks/useEnvironments";
@@ -28,6 +27,7 @@ import {
import { CancelRunDialog } from "./CancelRunDialog";
import { ReplayRunDialog } from "./ReplayRunDialog";
import { TaskRunStatusCombo } from "./TaskRunStatus";
import { LiveTimer } from "./LiveTimer";
type RunsTableProps = {
total: number;
@@ -94,9 +94,15 @@ export function TaskRunsTable({
{run.startedAt ? <DateTime date={run.startedAt} /> : ""}
</TableCell>
<TableCell to={path}>
{formatDuration(run.startedAt, run.completedAt, {
style: "short",
})}
{run.startedAt && run.finishedAt ? (
formatDuration(new Date(run.startedAt), new Date(run.finishedAt), {
style: "short",
})
) : run.startedAt ? (
<LiveTimer startTime={new Date(run.startedAt)} />
) : (
""
)}
</TableCell>
<TableCell to={path}>
{run.isTest ? (
@@ -195,7 +201,12 @@ function BlankState({ isLoading, filters }: Pick<RunsTableProps, "isLoading" | "
{environment ? (
<>
{" "}
in <EnvironmentLabel environment={environment} size="large" />
in{" "}
<EnvironmentLabel
environment={environment}
userName={environment.userName}
size="large"
/>
</>
) : null}
</Paragraph>
+6 -1
View File
@@ -1,6 +1,11 @@
import { useRef } from "react";
//a function that you call with a debounce delay, the function will only be called after the delay has passed
/**
* A function that you call with a debounce delay, the function will only be called after the delay has passed
*
* @param fn The function to debounce
* @param delay In ms
*/
export function useDebounce<T extends (...args: any[]) => any>(fn: T, delay: number) {
const timeout = useRef<ReturnType<typeof setTimeout>>();
+11
View File
@@ -0,0 +1,11 @@
import { useRef, MutableRefObject } from "react";
const useLazyRef = <T>(initialValFunc: () => T) => {
const ref: MutableRefObject<T | null> = useRef(null);
if (ref.current === null) {
ref.current = initialValFunc();
}
return ref;
};
export default useLazyRef;
+72
View File
@@ -0,0 +1,72 @@
import { Reducer, useReducer } from "react";
export type ListState<T> = {
items: T[];
};
type AppendAction<T> = {
type: "append";
items: T[];
};
type UpdateAction<T> = {
type: "update";
index: number;
item: T;
};
type DeleteAction<T> = {
type: "delete";
index: number;
};
type InsertAfter<T> = {
type: "insertAfter";
index: number;
items: T[];
};
type Action<T> = AppendAction<T> | UpdateAction<T> | DeleteAction<T> | InsertAfter<T>;
function reducer<T>(state: ListState<T>, action: Action<T>): ListState<T> {
switch (action.type) {
case "append":
return { items: [...state.items, ...action.items] };
case "update":
return {
items: state.items.map((v, i) => (i === action.index ? action.item : v)),
};
case "delete":
return { items: state.items.filter((_, i) => i !== action.index) };
case "insertAfter":
return {
items: [
...state.items.slice(0, action.index + 1),
...action.items,
...state.items.slice(action.index + 1),
],
};
}
}
type HookReturn<T> = {
items: T[];
append: (items: T[]) => void;
update: (index: number, item: T) => void;
delete: (index: number) => void;
insertAfter: (index: number, items: T[]) => void;
};
export function useList<T>(initialItems: T[]): HookReturn<T> {
const [state, dispatch] = useReducer<Reducer<ListState<T>, Action<T>>>(reducer, {
items: initialItems,
});
return {
items: state.items,
append: (items: T[]) => dispatch({ type: "append", items }),
update: (index: number, item: T) => dispatch({ type: "update", index, item }),
delete: (index: number) => dispatch({ type: "delete", index }),
insertAfter: (index: number, items: T[]) => dispatch({ type: "insertAfter", index, items }),
};
}
@@ -103,6 +103,7 @@ export class DeploymentPresenter {
exportName: "asc",
},
},
sdkVersion: true,
},
},
triggeredBy: {
@@ -135,6 +136,7 @@ export class DeploymentPresenter {
},
deployedBy: deployment.triggeredBy,
errorData: this.#prepareErrorData(deployment.errorData),
sdkVersion: deployment.worker?.sdkVersion,
},
};
}
@@ -1,9 +1,10 @@
import { Prisma, TaskRunStatus } from "@trigger.dev/database";
import { Direction } from "~/components/runs/RunStatuses";
import { sqlDatabaseSchema, PrismaClient, prisma } from "~/db.server";
import { FINISHED_STATUSES } from "~/components/runs/v3/TaskRunStatus";
import { sqlDatabaseSchema } from "~/db.server";
import { displayableEnvironments } from "~/models/runtimeEnvironment.server";
import { getUsername } from "~/utils/username";
import { CANCELLABLE_STATUSES } from "~/v3/services/cancelTaskRun.server";
import { BasePresenter } from "./basePresenter.server";
type RunListOptions = {
userId?: string;
@@ -28,13 +29,7 @@ export type RunList = Awaited<ReturnType<RunListPresenter["call"]>>;
export type RunListItem = RunList["runs"][0];
export type RunListAppliedFilters = RunList["filters"];
export class RunListPresenter {
#prismaClient: PrismaClient;
constructor(prismaClient: PrismaClient = prisma) {
this.#prismaClient = prismaClient;
}
export class RunListPresenter extends BasePresenter {
public async call({
userId,
projectSlug,
@@ -60,7 +55,7 @@ export class RunListPresenter {
to !== undefined;
// Find the project scoped to the organization
const project = await this.#prismaClient.project.findFirstOrThrow({
const project = await this._replica.project.findFirstOrThrow({
select: {
id: true,
environments: {
@@ -88,7 +83,7 @@ export class RunListPresenter {
});
//get all possible tasks
const possibleTasks = await this.#prismaClient.backgroundWorkerTask.findMany({
const possibleTasks = await this._replica.backgroundWorkerTask.findMany({
distinct: ["slug"],
where: {
projectId: project.id,
@@ -96,7 +91,7 @@ export class RunListPresenter {
});
//get the runs
let runs = await this.#prismaClient.$queryRaw<
let runs = await this._replica.$queryRaw<
{
id: string;
number: BigInt;
@@ -107,10 +102,9 @@ export class RunListPresenter {
status: TaskRunStatus;
createdAt: Date;
lockedAt: Date | null;
completedAt: Date | null;
updatedAt: Date;
isTest: boolean;
spanId: string;
attempts: BigInt;
}[]
>`
SELECT
@@ -123,20 +117,13 @@ export class RunListPresenter {
tr.status AS status,
tr."createdAt" AS "createdAt",
tr."lockedAt" AS "lockedAt",
tra."completedAt" AS "completedAt",
tr."updatedAt" AS "updatedAt",
tr."isTest" AS "isTest",
tr."spanId" AS "spanId",
COUNT(tra.id) AS attempts
tr."spanId" AS "spanId"
FROM
${sqlDatabaseSchema}."TaskRun" tr
LEFT JOIN
(
SELECT *,
ROW_NUMBER() OVER (PARTITION BY "taskRunId" ORDER BY "createdAt" DESC) rn
FROM ${sqlDatabaseSchema}."TaskRunAttempt"
) tra ON tr.id = tra."taskRunId" AND tra.rn = 1
LEFT JOIN
${sqlDatabaseSchema}."BackgroundWorker" bw ON tra."backgroundWorkerId" = bw.id
${sqlDatabaseSchema}."BackgroundWorker" bw ON tr."lockedToVersionId" = bw.id
WHERE
-- project
tr."projectId" = ${project.id}
@@ -154,15 +141,11 @@ export class RunListPresenter {
? Prisma.sql`AND tr."taskIdentifier" IN (${Prisma.join(tasks)})`
: Prisma.empty
}
${hasStatusFilters ? Prisma.sql`AND (` : Prisma.empty}
${
statuses && statuses.length > 0
? Prisma.sql`tr.status = ANY(ARRAY[${Prisma.join(statuses)}]::"TaskRunStatus"[])`
? Prisma.sql`AND tr.status = ANY(ARRAY[${Prisma.join(statuses)}]::"TaskRunStatus"[])`
: Prisma.empty
}
${statuses && statuses.length > 0 && hasStatusFilters ? Prisma.sql` OR ` : Prisma.empty}
${hasStatusFilters ? Prisma.sql`tr.status IS NULL` : Prisma.empty}
${hasStatusFilters ? Prisma.sql`) ` : Prisma.empty}
${
environments && environments.length > 0
? Prisma.sql`AND tr."runtimeEnvironmentId" IN (${Prisma.join(environments)})`
@@ -179,8 +162,6 @@ export class RunListPresenter {
? Prisma.sql`AND tr."createdAt" <= ${new Date(to).toISOString()}::timestamp`
: Prisma.empty
}
GROUP BY
tr."friendlyId", tr."taskIdentifier", tr."runtimeEnvironmentId", tr.id, bw.version, tra.status, tr."createdAt", tra."startedAt", tra."completedAt"
ORDER BY
${direction === "forward" ? Prisma.sql`tr.id DESC` : Prisma.sql`tr.id ASC`}
LIMIT ${pageSize + 1}`;
@@ -219,19 +200,21 @@ export class RunListPresenter {
throw new Error(`Environment not found for TaskRun ${run.id}`);
}
const hasFinished = FINISHED_STATUSES.includes(run.status);
return {
id: run.id,
friendlyId: run.runFriendlyId,
number: Number(run.number),
createdAt: run.createdAt,
startedAt: run.lockedAt,
completedAt: run.completedAt,
createdAt: run.createdAt.toISOString(),
startedAt: run.lockedAt ? run.lockedAt.toISOString() : undefined,
hasFinished,
finishedAt: hasFinished ? run.updatedAt.toISOString() : undefined,
isTest: run.isTest,
status: run.status,
version: run.version,
taskIdentifier: run.taskIdentifier,
spanId: run.spanId,
attempts: Number(run.attempts),
isReplayable: true,
isCancellable: CANCELLABLE_STATUSES.includes(run.status),
environment: displayableEnvironments(environment, userId),
@@ -4,16 +4,15 @@ import {
TaskRunStatus,
TaskTriggerSource,
} from "@trigger.dev/database";
import { PrismaClient, prisma, sqlDatabaseSchema } from "~/db.server";
import { QUEUED_STATUSES, RUNNING_STATUSES } from "~/components/runs/v3/TaskRunStatus";
import { sqlDatabaseSchema } from "~/db.server";
import { Organization } from "~/models/organization.server";
import { Project } from "~/models/project.server";
import { displayableEnvironments } from "~/models/runtimeEnvironment.server";
import { User } from "~/models/user.server";
import { sortEnvironments } from "~/services/environmentSort.server";
import { logger } from "~/services/logger.server";
import { getUsername } from "~/utils/username";
import { BasePresenter } from "./basePresenter.server";
import { QUEUED_STATUSES, RUNNING_STATUSES } from "~/components/runs/v3/TaskRunStatus";
import { displayableEnvironments } from "~/models/runtimeEnvironment.server";
export type Task = {
slug: string;
@@ -26,10 +25,6 @@ export type Task = {
type: RuntimeEnvironmentType;
userName?: string;
}[];
latestRun?: {
createdAt: Date;
status: TaskRunStatus;
};
};
type Return = Awaited<ReturnType<TaskListPresenter["call"]>>;
@@ -98,39 +93,8 @@ export class TaskListPresenter extends BasePresenter {
JOIN ${sqlDatabaseSchema}."BackgroundWorkerTask" tasks ON tasks."workerId" = workers.id
ORDER BY slug ASC;`;
let latestRuns = [] as {
createdAt: Date;
status: TaskRunStatus;
taskIdentifier: string;
}[];
if (tasks.length > 0) {
const uniqueTaskSlugs = new Set(tasks.map((t) => t.slug));
latestRuns = await this._replica.$queryRaw<
{
createdAt: Date;
status: TaskRunStatus;
taskIdentifier: string;
}[]
>`
SELECT * FROM (
SELECT
"createdAt",
"status",
"taskIdentifier",
ROW_NUMBER() OVER (PARTITION BY "taskIdentifier" ORDER BY "updatedAt" DESC) AS rn
FROM
${sqlDatabaseSchema}."TaskRun"
WHERE
"taskIdentifier" IN(${Prisma.join(Array.from(uniqueTaskSlugs))})
AND "projectId" = ${project.id}
) t
WHERE rn = 1;`;
}
//group by the task identifier (task.slug). Add the latestRun and add all the environments.
const outputTasks = tasks.reduce((acc, task) => {
const latestRun = latestRuns.find((r) => r.taskIdentifier === task.slug);
const environment = project.environments.find((env) => env.id === task.runtimeEnvironmentId);
if (!environment) {
throw new Error(`Environment not found for TaskRun ${task.id}`);
@@ -151,13 +115,6 @@ export class TaskListPresenter extends BasePresenter {
//order the environments
existingTask.environments = sortEnvironments(existingTask.environments);
existingTask.latestRun = latestRun
? {
createdAt: latestRun.createdAt,
status: latestRun.status,
}
: undefined;
return acc;
}, [] as Task[]);
@@ -186,6 +143,10 @@ export class TaskListPresenter extends BasePresenter {
}
async #getActivity(tasks: string[], projectId: string) {
if (tasks.length === 0) {
return {};
}
const activity = await this._replica.$queryRaw<
{
taskIdentifier: string;
@@ -257,6 +218,10 @@ export class TaskListPresenter extends BasePresenter {
}
async #getRunningStats(tasks: string[], projectId: string) {
if (tasks.length === 0) {
return {};
}
const statuses = await this._replica.$queryRaw<
{
taskIdentifier: string;
@@ -305,6 +270,10 @@ export class TaskListPresenter extends BasePresenter {
}
async #getAverageDurations(tasks: string[], projectId: string) {
if (tasks.length === 0) {
return {};
}
const durations = await this._replica.$queryRaw<
{
taskIdentifier: string;
@@ -63,17 +63,18 @@ export class TestPresenter {
const searchParams = createSearchParams(url, TestSearchParams);
//no environmentId
if (!searchParams.success || !searchParams.params.get("environment")) {
if (!searchParams.success) {
return {
hasSelectedEnvironment: false as const,
environments,
};
}
//default to dev environment
const environment = searchParams.params.get("environment") ?? "dev";
//is the environmentId valid?
const matchingEnvironment = project.environments.find(
(env) => env.slug === searchParams.params.get("environment")
);
const matchingEnvironment = project.environments.find((env) => env.slug === environment);
if (!matchingEnvironment) {
return {
hasSelectedEnvironment: false as const,
@@ -101,7 +102,7 @@ export class TestPresenter {
WHERE "runtimeEnvironmentId" = ${matchingEnvironment.id}
),
latest_workers AS (SELECT * FROM workers WHERE rn = 1)
SELECT bwt.id, version, slug as "taskIdentifier", "filePath", "exportName", bwt."friendlyId"
SELECT bwt.id, version, slug as "taskIdentifier", "filePath", "exportName", bwt."friendlyId", bwt."triggerSource"
FROM latest_workers
JOIN ${sqlDatabaseSchema}."BackgroundWorkerTask" bwt ON bwt."workerId" = latest_workers.id
ORDER BY bwt."exportName" ASC;
@@ -171,17 +171,23 @@ export class TestTaskPresenter {
return {
triggerSource: "SCHEDULED",
task: taskWithEnvironment,
runs: await Promise.all(
latestRuns.map(async (r) => {
const number = Number(r.number);
runs: (
await Promise.all(
latestRuns.map(async (r) => {
const number = Number(r.number);
return {
...r,
number,
payload: await getScheduleTaskRunPayload(r),
};
})
),
const payload = await getScheduleTaskRunPayload(r);
if (payload.success) {
return {
...r,
number,
payload: payload.data,
};
}
})
)
).filter(Boolean),
};
}
}
@@ -189,6 +195,6 @@ export class TestTaskPresenter {
async function getScheduleTaskRunPayload(run: RawRun) {
const payload = await parsePacket({ data: run.payload, dataType: run.payloadType });
const parsed = ScheduledTaskPayload.parse(payload);
const parsed = ScheduledTaskPayload.safeParse(payload);
return parsed;
}
@@ -232,7 +232,7 @@ export default function Page() {
<FormError id={projectSlug.errorId}>{projectSlug.error}</FormError>
<FormError>{deleteForm.error}</FormError>
<Hint>
This change is irreversible, so please be certain. Type in the Project slug
This change is irreversible, so please be certain. Type in the Project slug{" "}
<InlineCode variant="extra-small">{project.slug}</InlineCode> and then press
Delete.
</Hint>
@@ -1,10 +1,10 @@
import { ChatBubbleLeftRightIcon, ChevronDownIcon, ChevronUpIcon } from "@heroicons/react/20/solid";
import { useRevalidator } from "@remix-run/react";
import { LoaderFunctionArgs } from "@remix-run/server-runtime";
import { formatDuration, formatDurationMilliseconds } from "@trigger.dev/core/v3";
import { formatDurationMilliseconds } from "@trigger.dev/core/v3";
import { TaskRunStatus } from "@trigger.dev/database";
import { Fragment, Suspense, useEffect, useState } from "react";
import { Bar, BarChart, ResponsiveContainer, Tooltip, TooltipProps, XAxis, YAxis } from "recharts";
import { Bar, BarChart, ResponsiveContainer, Tooltip, TooltipProps } from "recharts";
import { TypedAwait, typeddefer, useTypedLoaderData } from "remix-typedjson";
import { Feedback } from "~/components/Feedback";
import { InitCommandV3, TriggerDevStepV3, TriggerLoginStepV3 } from "~/components/SetupCommands";
@@ -14,8 +14,9 @@ import { EnvironmentLabel } from "~/components/environments/EnvironmentLabel";
import { MainCenteredContainer, PageBody, PageContainer } from "~/components/layout/AppLayout";
import { Button } from "~/components/primitives/Buttons";
import { Callout } from "~/components/primitives/Callout";
import { DateTime, formatDateTime } from "~/components/primitives/DateTime";
import { formatDateTime } from "~/components/primitives/DateTime";
import { Header1, Header2, Header3 } from "~/components/primitives/Headers";
import { Input } from "~/components/primitives/Input";
import { NavBar, PageTitle } from "~/components/primitives/PageHeader";
import { Paragraph } from "~/components/primitives/Paragraph";
import { Spinner } from "~/components/primitives/Spinner";
@@ -31,13 +32,9 @@ import {
TableRow,
} from "~/components/primitives/Table";
import { SimpleTooltip } from "~/components/primitives/Tooltip";
import TooltipPortal from "~/components/primitives/TooltipPortal";
import { TaskFunctionName } from "~/components/runs/v3/TaskPath";
import {
TaskRunStatusCombo,
TaskRunStatusIcon,
runStatusClassNameColor,
runStatusTitle,
} from "~/components/runs/v3/TaskRunStatus";
import { TaskRunStatusCombo } from "~/components/runs/v3/TaskRunStatus";
import {
TaskTriggerSourceIcon,
taskTriggerSourceDescription,
@@ -45,8 +42,8 @@ import {
import { useEventSource } from "~/hooks/useEventSource";
import { useOrganization } from "~/hooks/useOrganizations";
import { useProject } from "~/hooks/useProject";
import { useUser } from "~/hooks/useUser";
import { TaskActivity, TaskListPresenter } from "~/presenters/v3/TaskListPresenter.server";
import { useTextFilter } from "~/hooks/useTextFilter";
import { Task, TaskActivity, TaskListPresenter } from "~/presenters/v3/TaskListPresenter.server";
import { requireUserId } from "~/services/session.server";
import { cn } from "~/utils/cn";
import { ProjectParamSchema, v3RunsPath, v3TasksStreamingPath } from "~/utils/pathBuilder";
@@ -84,6 +81,31 @@ export default function Page() {
const project = useProject();
const { tasks, userHasTasks, activity, runningStats, durations } =
useTypedLoaderData<typeof loader>();
const { filterText, setFilterText, filteredItems } = useTextFilter<Task>({
items: tasks,
filter: (task, text) => {
if (task.slug.toLowerCase().includes(text.toLowerCase())) {
return true;
}
if (
task.exportName.toLowerCase().includes(text.toLowerCase().replace("(", "").replace(")", ""))
) {
return true;
}
if (task.filePath.toLowerCase().includes(text.toLowerCase())) {
return true;
}
if (task.triggerSource === "SCHEDULED" && "scheduled".includes(text.toLowerCase())) {
return true;
}
return false;
},
});
const hasTasks = tasks.length > 0;
//live reload the page when the tasks change
@@ -105,11 +127,22 @@ export default function Page() {
<PageTitle title="Tasks" />
</NavBar>
<PageBody>
<div className={cn("grid h-full grid-cols-1 gap-4")}>
<div className="h-full">
{hasTasks ? (
<div className="flex flex-col gap-4 pb-4">
{!userHasTasks && <UserHasNoTasks />}
<div className={cn("grid h-full grid-rows-1")}>
{hasTasks ? (
<div className="flex flex-col gap-4 pb-4">
{!userHasTasks && <UserHasNoTasks />}
<div className="pb-4">
<div className="h-8">
<Input
placeholder="Search tasks"
variant="tertiary"
icon="search"
fullWidth={true}
value={filterText}
onChange={(e) => setFilterText(e.target.value)}
autoFocus
/>
</div>
<Table>
<TableHeader>
<TableRow>
@@ -120,13 +153,12 @@ export default function Page() {
<TableHeaderCell>Activity (7d)</TableHeaderCell>
<TableHeaderCell>Avg. duration</TableHeaderCell>
<TableHeaderCell>Environments</TableHeaderCell>
<TableHeaderCell>Last run</TableHeaderCell>
<TableHeaderCell hiddenLabel>Go to page</TableHeaderCell>
</TableRow>
</TableHeader>
<TableBody>
{tasks.length > 0 ? (
tasks.map((task) => {
{filteredItems.length > 0 ? (
filteredItems.map((task) => {
const path = v3RunsPath(organization, project, {
tasks: [task.slug],
});
@@ -218,30 +250,12 @@ export default function Page() {
))}
</div>
</TableCell>
<TableCell to={path}>
{task.latestRun ? (
<div
className={cn(
"flex items-center gap-1",
runStatusClassNameColor(task.latestRun.status)
)}
>
<TaskRunStatusIcon
status={task.latestRun.status}
className="h-4 w-4"
/>
<DateTime date={task.latestRun.createdAt} />
</div>
) : (
"Never run"
)}
</TableCell>
<TableCellChevron to={path} />
</TableRow>
);
})
) : (
<TableBlankRow colSpan={6}>
<TableBlankRow colSpan={8}>
<Paragraph variant="small" className="flex items-center justify-center">
No tasks match your filters
</Paragraph>
@@ -250,12 +264,12 @@ export default function Page() {
</TableBody>
</Table>
</div>
) : (
<MainCenteredContainer className="max-w-prose">
<CreateTaskInstructions />
</MainCenteredContainer>
)}
</div>
</div>
) : (
<MainCenteredContainer className="max-w-prose">
<CreateTaskInstructions />
</MainCenteredContainer>
)}
</div>
</PageBody>
</PageContainer>
@@ -362,7 +376,9 @@ function TaskActivityGraph({ activity }: { activity: TaskActivity }) {
content={<CustomTooltip />}
allowEscapeViewBox={{ x: true, y: true }}
wrapperStyle={{ zIndex: 1000 }}
animationDuration={0}
/>
{/* The background */}
<Bar
dataKey="bg"
@@ -425,18 +441,21 @@ const CustomTooltip = ({ active, payload, label }: TooltipProps<number, string>)
}));
const title = payload[0].payload.day as string;
const formattedDate = formatDateTime(new Date(title), "UTC", [], false, false);
return (
<div className="rounded-sm border border-grid-bright bg-background-dimmed px-3 py-2">
<Header3 className="border-b-charcoal-650 border-b pb-2">{formattedDate}</Header3>
<div className="mt-2 grid grid-cols-[1fr_auto] gap-2 text-xs text-text-bright">
{items.map((item) => (
<Fragment key={item.status}>
<TaskRunStatusCombo status={item.status} />
<p>{item.value}</p>
</Fragment>
))}
<TooltipPortal active={active}>
<div className="rounded-sm border border-grid-bright bg-background-dimmed px-3 py-2">
<Header3 className="border-b-charcoal-650 border-b pb-2">{formattedDate}</Header3>
<div className="mt-2 grid grid-cols-[1fr_auto] gap-2 text-xs text-text-bright">
{items.map((item) => (
<Fragment key={item.status}>
<TaskRunStatusCombo status={item.status} />
<p>{item.value}</p>
</Fragment>
))}
</div>
</div>
</div>
</TooltipPortal>
);
}
@@ -92,6 +92,9 @@ export default function Page() {
<DeploymentStatus status={deployment.status} className="text-sm" />
</Property>
<Property label="Tasks">{deployment.tasks ? deployment.tasks.length : ""}</Property>
<Property label="SDK Version">
{deployment.sdkVersion ? deployment.sdkVersion : ""}
</Property>
<Property label="Started at">
<Paragraph variant="small/bright">
<DateTimeAccurate date={deployment.createdAt} /> UTC
@@ -1,35 +1,44 @@
import { Submission, conform, useForm } from "@conform-to/react";
import {
FieldConfig,
list,
requestIntent,
useFieldList,
useFieldset,
useForm,
} from "@conform-to/react";
import { parse } from "@conform-to/zod";
import { Form, useActionData, useLocation, useNavigate, useNavigation } from "@remix-run/react";
import { PlusIcon, XMarkIcon } from "@heroicons/react/20/solid";
import { Form, useActionData, useNavigate, useNavigation } from "@remix-run/react";
import { ActionFunctionArgs, LoaderFunctionArgs, json } from "@remix-run/server-runtime";
import { Fragment, useEffect, useRef, useState } from "react";
import { RefObject, useCallback, useEffect, useRef, useState } from "react";
import { redirect, typedjson, useTypedLoaderData } from "remix-typedjson";
import { z } from "zod";
import { InlineCode } from "~/components/code/InlineCode";
import { EnvironmentLabel } from "~/components/environments/EnvironmentLabel";
import {
environmentTextClassName,
environmentTitle,
} from "~/components/environments/EnvironmentLabel";
import { Button, LinkButton } from "~/components/primitives/Buttons";
import { Callout } from "~/components/primitives/Callout";
import { Checkbox } from "~/components/primitives/Checkbox";
import { Dialog, DialogContent, DialogHeader } from "~/components/primitives/Dialog";
import { Fieldset } from "~/components/primitives/Fieldset";
import { FormButtons } from "~/components/primitives/FormButtons";
import { FormError } from "~/components/primitives/FormError";
import { Hint } from "~/components/primitives/Hint";
import { Input } from "~/components/primitives/Input";
import { InputGroup } from "~/components/primitives/InputGroup";
import { Label } from "~/components/primitives/Label";
import { Switch } from "~/components/primitives/Switch";
import { prisma } from "~/db.server";
import { useList } from "~/hooks/useList";
import { useOrganization } from "~/hooks/useOrganizations";
import { useProject } from "~/hooks/useProject";
import { redirectWithSuccessMessage } from "~/models/message.server";
import { EnvironmentVariablesPresenter } from "~/presenters/v3/EnvironmentVariablesPresenter.server";
import { requireUserId } from "~/services/session.server";
import {
ProjectParamSchema,
v3EnvironmentVariablesPath,
v3NewEnvironmentVariablesPath,
} from "~/utils/pathBuilder";
import { cn } from "~/utils/cn";
import { ProjectParamSchema, v3EnvironmentVariablesPath } from "~/utils/pathBuilder";
import { EnvironmentVariablesRepository } from "~/v3/environmentVariables/environmentVariablesRepository.server";
import { CreateEnvironmentVariable } from "~/v3/environmentVariables/repository";
import { EnvironmentVariableKey } from "~/v3/environmentVariables/repository";
import dotenv from "dotenv";
export const loader = async ({ request, params }: LoaderFunctionArgs) => {
const userId = await requireUserId(request);
@@ -47,7 +56,6 @@ export const loader = async ({ request, params }: LoaderFunctionArgs) => {
environments,
});
} catch (error) {
console.error(error);
throw new Response(undefined, {
status: 400,
statusText: "Something went wrong, if this problem persists please contact support.",
@@ -55,9 +63,39 @@ export const loader = async ({ request, params }: LoaderFunctionArgs) => {
}
};
const Variable = z.object({
key: EnvironmentVariableKey,
value: z.string().nonempty("Value is required"),
});
type Variable = z.infer<typeof Variable>;
const schema = z.object({
action: z.enum(["create", "create-more"]),
...CreateEnvironmentVariable.shape,
overwrite: z.preprocess((i) => {
if (i === "true") return true;
if (i === "false") return false;
return;
}, z.boolean()),
environmentIds: z.preprocess((i) => {
if (typeof i === "string") return [i];
if (Array.isArray(i)) {
const ids = i.filter((v) => typeof v === "string" && v !== "");
if (ids.length === 0) {
return;
}
return ids;
}
return;
}, z.array(z.string(), { required_error: "At least one environment is required" })),
variables: z.preprocess((i) => {
if (!Array.isArray(i)) {
return [];
}
return i;
}, Variable.array().nonempty("At least one variable is required")),
});
export const action = async ({ request, params }: ActionFunctionArgs) => {
@@ -92,22 +130,22 @@ export const action = async ({ request, params }: ActionFunctionArgs) => {
const result = await repository.create(project.id, userId, submission.value);
if (!result.success) {
submission.error.key = result.error;
if (result.variableErrors) {
for (const { key, error } of result.variableErrors) {
const index = submission.value.variables.findIndex((v) => v.key === key);
if (index !== -1) {
submission.error[`variables[${index}].key`] = error;
}
}
} else {
submission.error.variables = result.error;
}
return json(submission);
}
switch (submission.value.action) {
case "create":
return redirect(
v3EnvironmentVariablesPath({ slug: organizationSlug }, { slug: projectParam })
);
case "create-more":
return redirectWithSuccessMessage(
v3NewEnvironmentVariablesPath({ slug: organizationSlug }, { slug: projectParam }),
request,
`Created ${submission.value.key} environment variable`
);
}
return redirect(v3EnvironmentVariablesPath({ slug: organizationSlug }, { slug: projectParam }));
};
export default function Page() {
@@ -118,15 +156,11 @@ export default function Page() {
const navigate = useNavigate();
const organization = useOrganization();
const project = useProject();
const keyFieldRef = useRef<HTMLInputElement>(null);
const isLoading =
navigation.state !== "idle" &&
navigation.formMethod === "post" &&
navigation.formData?.get("action") === "create";
const isLoading = navigation.state !== "idle" && navigation.formMethod === "post";
const [form, { key }] = useForm({
id: "create-environment-variable",
const [form, { environmentIds, variables }] = useForm({
id: "create-environment-variables",
// TODO: type this
lastSubmission: lastSubmission as any,
onValidate({ formData }) {
@@ -141,14 +175,6 @@ export default function Page() {
setIsOpen(true);
}, []);
useEffect(() => {
if (navigation.state !== "idle") return;
if (lastSubmission !== undefined) return;
form.ref.current?.reset();
keyFieldRef.current?.focus();
}, [navigation.state, lastSubmission]);
return (
<Dialog
open={isOpen}
@@ -158,61 +184,64 @@ export default function Page() {
}
}}
>
<DialogContent>
<DialogHeader>New environment variable</DialogHeader>
<Form method="post" {...form.props}>
<DialogContent className="md:max-w-2xl lg:max-w-3xl">
<DialogHeader>New environment variables</DialogHeader>
<Form
method="post"
{...form.props}
className="max-h-[70vh] overflow-y-auto scrollbar-thin scrollbar-track-transparent scrollbar-thumb-charcoal-600"
>
<Fieldset className="mt-2">
<InputGroup fullWidth>
<Label>Key</Label>
<Input
{...conform.input(key)}
placeholder="e.g. CLIENT_KEY"
autoFocus
ref={keyFieldRef}
/>
</InputGroup>
<InputGroup fullWidth>
<div className="flex items-center justify-between">
<Label>Values</Label>
<Switch
variant="small"
label="Reveal values"
checked={revealAll}
onCheckedChange={(e) => setRevealAll(e.valueOf())}
/>
</div>
<div className="grid grid-cols-[auto_1fr] gap-x-2 gap-y-2">
{environments.map((environment, index) => {
return (
<Fragment key={environment.id}>
<input
type="hidden"
name={`values[${index}].environmentId`}
value={environment.id}
/>
<label
className="flex items-center justify-end"
htmlFor={`values[${index}].value`}
<Label>Environments</Label>
<div className="flex flex-wrap items-center gap-2">
{environments.map((environment) => (
<Checkbox
key={environment.id}
id={environment.id}
value={environment.id}
name="environmentIds"
type="radio"
label={
<span
className={cn("text-xs uppercase", environmentTextClassName(environment))}
>
<EnvironmentLabel environment={environment} className="h-5 px-2" />
</label>
<Input
type={revealAll ? "text" : "password"}
name={`values[${index}].value`}
placeholder="Not set"
/>
</Fragment>
);
})}
{environmentTitle(environment)}
</span>
}
variant="button"
/>
))}
</div>
<FormError id={environmentIds.errorId}>{environmentIds.error}</FormError>
<Hint>
Dev environment variables specified here will be overridden by ones in your .env
file when running locally.
</Hint>
</InputGroup>
<Hint>Tip: Paste your .env into this form to populate it:</Hint>
<InputGroup fullWidth>
<FieldLayout>
<Label>Keys</Label>
<div className="flex justify-between gap-1">
<Label>Values</Label>
<Switch
variant="small"
label="Reveal"
checked={revealAll}
onCheckedChange={(e) => setRevealAll(e.valueOf())}
/>
</div>
</FieldLayout>
<VariableFields
revealValues={revealAll}
formId={form.id}
formRef={form.ref}
variablesFields={variables}
/>
<FormError id={variables.errorId}>{variables.error}</FormError>
</InputGroup>
<Callout variant="info" className="inline-flex">
Dev environment variables specified here will be overridden by ones in your{" "}
<InlineCode variant="extra-small">.env</InlineCode> file when running locally.
</Callout>
<FormError id={key.errorId}>{key.error}</FormError>
<FormError>{form.error}</FormError>
<FormButtons
confirmButton={
@@ -221,18 +250,18 @@ export default function Page() {
type="submit"
variant="primary/small"
disabled={isLoading}
name="action"
value="create-more"
name="overwrite"
value="false"
>
{isLoading ? "Saving" : "Save and add another"}
{isLoading ? "Saving" : "Save"}
</Button>
<Button
variant="secondary/small"
disabled={isLoading}
name="action"
value="create"
name="overwrite"
value="true"
>
{isLoading ? "Saving" : "Save"}
{isLoading ? "Overwriting" : "Overwrite"}
</Button>
</div>
}
@@ -251,3 +280,156 @@ export default function Page() {
</Dialog>
);
}
function FieldLayout({ children }: { children: React.ReactNode }) {
return <div className="grid w-full grid-cols-[1fr_1fr_2rem] gap-2">{children}</div>;
}
function VariableFields({
revealValues,
formId,
variablesFields,
formRef,
}: {
revealValues: boolean;
formId?: string;
variablesFields: FieldConfig<any>;
formRef: RefObject<HTMLFormElement>;
}) {
const {
items,
append,
update,
delete: remove,
insertAfter,
} = useList<Variable>([{ key: "", value: "" }]);
const handlePaste = useCallback((index: number, e: React.ClipboardEvent<HTMLInputElement>) => {
const clipboardData = e.clipboardData;
if (!clipboardData) return;
let text = clipboardData.getData("text");
if (!text) return;
const variables = dotenv.parse(text);
const keyValuePairs = Object.entries(variables).map(([key, value]) => ({ key, value }));
//do the default paste
if (keyValuePairs.length === 0) return;
//prevent default pasting
e.preventDefault();
const [firstPair, ...rest] = keyValuePairs;
update(index, firstPair);
for (const pair of rest) {
requestIntent(formRef.current ?? undefined, list.append(variablesFields.name));
}
insertAfter(index, rest);
}, []);
const fields = useFieldList(formRef, variablesFields);
return (
<>
{fields.map((field, index) => {
const item = items[index];
return (
<VariableField
formId={formId}
key={index}
index={index}
value={item}
onChange={(value) => update(index, value)}
onPaste={(e) => handlePaste(index, e)}
onDelete={() => {
requestIntent(
formRef.current ?? undefined,
list.remove(variablesFields.name, { index })
);
remove(index);
}}
showDeleteButton={items.length > 1}
showValue={revealValues}
config={field}
/>
);
})}
<Button
variant="tertiary/medium"
type="button"
onClick={() => {
requestIntent(formRef.current ?? undefined, list.append(variablesFields.name));
append([{ key: "", value: "" }]);
}}
LeadingIcon={PlusIcon}
>
Add another
</Button>
</>
);
}
function VariableField({
formId,
index,
value,
onChange,
onPaste,
onDelete,
showDeleteButton,
showValue,
config,
}: {
formId?: string;
index: number;
value: Variable;
onChange: (value: Variable) => void;
onPaste: (e: React.ClipboardEvent<HTMLInputElement>) => void;
onDelete: () => void;
showDeleteButton: boolean;
showValue: boolean;
config: FieldConfig<Variable>;
}) {
const ref = useRef<HTMLFieldSetElement>(null);
const fields = useFieldset(ref, config);
const baseFieldName = `variables[${index}]`;
return (
<fieldset ref={ref}>
<FieldLayout>
<Input
id={`${formId}-${baseFieldName}.key`}
name={`${baseFieldName}.key`}
placeholder="e.g. CLIENT_KEY"
value={value.key}
onChange={(e) => onChange({ ...value, key: e.currentTarget.value })}
autoFocus={index === 0}
onPaste={onPaste}
/>
<Input
id={`${formId}-${baseFieldName}.value`}
name={`${baseFieldName}.value`}
type={showValue ? "text" : "password"}
placeholder="Not set"
value={value.value}
onChange={(e) => onChange({ ...value, value: e.currentTarget.value })}
/>
{showDeleteButton && (
<Button
variant="minimal/medium"
type="button"
onClick={() => onDelete()}
LeadingIcon={XMarkIcon}
/>
)}
</FieldLayout>
<div className="space-y-2">
<FormError id={fields.key.errorId}>{fields.key.error}</FormError>
<FormError id={fields.value.errorId}>{fields.value.error}</FormError>
</div>
</fieldset>
);
}
@@ -187,8 +187,9 @@ export default function Page() {
to={v3NewEnvironmentVariablesPath(organization, project)}
variant="primary/small"
LeadingIcon={PlusIcon}
shortcut={{ key: "n" }}
>
New environment variable
Add new
</LinkButton>
</div>
<Table>
@@ -1,6 +1,4 @@
import {
ArrowsPointingInIcon,
ArrowsPointingOutIcon,
ChevronDownIcon,
ChevronRightIcon,
MagnifyingGlassMinusIcon,
@@ -17,7 +15,8 @@ import {
} from "@trigger.dev/core/v3";
import { RuntimeEnvironmentType } from "@trigger.dev/database";
import { motion } from "framer-motion";
import { useEffect, useRef, useState } from "react";
import { useCallback, useEffect, useRef, useState } from "react";
import { useHotkeys } from "react-hotkeys-hook";
import { typedjson, useTypedLoaderData } from "remix-typedjson";
import { ShowParentIcon, ShowParentIconSelected } from "~/assets/icons/ShowParentIcon";
import tileBgPath from "~/assets/images/error-banner-tile@2x.png";
@@ -26,11 +25,13 @@ import { InlineCode } from "~/components/code/InlineCode";
import { EnvironmentLabel } from "~/components/environments/EnvironmentLabel";
import { MainCenteredContainer, PageBody } from "~/components/layout/AppLayout";
import { Badge } from "~/components/primitives/Badge";
import { Button, LinkButton } from "~/components/primitives/Buttons";
import { LinkButton } from "~/components/primitives/Buttons";
import { Callout } from "~/components/primitives/Callout";
import { Header3 } from "~/components/primitives/Headers";
import { Input } from "~/components/primitives/Input";
import { NavBar, PageAccessories, PageTitle } from "~/components/primitives/PageHeader";
import { Paragraph } from "~/components/primitives/Paragraph";
import { Popover, PopoverArrowTrigger, PopoverContent } from "~/components/primitives/Popover";
import {
ResizableHandle,
ResizablePanel,
@@ -66,10 +67,6 @@ import {
v3RunsPath,
} from "~/utils/pathBuilder";
import { SpanView } from "../resources.orgs.$organizationSlug.projects.v3.$projectParam.runs.$runParam.spans.$spanParam/route";
import { number } from "zod";
import { useHotkeys } from "react-hotkeys-hook";
import { Popover, PopoverArrowTrigger, PopoverContent } from "~/components/primitives/Popover";
import { Header3 } from "~/components/primitives/Headers";
export const loader = async ({ request, params }: LoaderFunctionArgs) => {
const userId = await requireUserId(request);
@@ -268,29 +265,25 @@ function TasksTreeView({
onSelectedIdChanged,
estimatedRowHeight: () => 32,
parentRef,
filter: (node) => {
const nodePassesErrorTest = (errorsOnly && node.data.isError) || !errorsOnly;
if (!nodePassesErrorTest) return false;
filter: {
value: { text: filterText, errorsOnly },
fn: (value, node) => {
const nodePassesErrorTest = (value.errorsOnly && node.data.isError) || !value.errorsOnly;
if (!nodePassesErrorTest) return false;
if (filterText === "") return true;
if (node.data.message.toLowerCase().includes(filterText.toLowerCase())) {
return true;
}
return false;
if (value.text === "") return true;
if (node.data.message.toLowerCase().includes(value.text.toLowerCase())) {
return true;
}
return false;
},
},
});
return (
<div className="grid h-full grid-rows-[2.5rem_1fr_3.25rem] overflow-hidden">
<div className="mx-3 flex items-center justify-between gap-2 border-b border-grid-dimmed">
<Input
placeholder="Search log"
variant="tertiary"
icon="search"
fullWidth={true}
value={filterText}
onChange={(e) => setFilterText(e.target.value)}
/>
<SearchField onChange={setFilterText} />
<div className="flex items-center gap-2">
<Switch
variant="small"
@@ -1004,3 +997,27 @@ function NumberShortcuts({ toggleLevel }: { toggleLevel: (depth: number) => void
</div>
);
}
function SearchField({ onChange }: { onChange: (value: string) => void }) {
const [value, setValue] = useState("");
const updateFilterText = useDebounce((text: string) => {
onChange(text);
}, 250);
const updateValue = useCallback((value: string) => {
setValue(value);
updateFilterText(value);
}, []);
return (
<Input
placeholder="Search log"
variant="tertiary"
icon="search"
fullWidth={true}
value={value}
onChange={(e) => updateValue(e.target.value)}
/>
);
}
@@ -1,7 +1,7 @@
import { BeakerIcon, BookOpenIcon } from "@heroicons/react/24/solid";
import { useNavigation } from "@remix-run/react";
import { LoaderFunctionArgs } from "@remix-run/server-runtime";
import { typedjson, useTypedLoaderData } from "remix-typedjson";
import { TypedAwait, typeddefer, typedjson, useTypedLoaderData } from "remix-typedjson";
import { TaskIcon } from "~/assets/icons/TaskIcon";
import { BlankstateInstructions } from "~/components/BlankstateInstructions";
import { StepContentContainer } from "~/components/StepContentContainer";
@@ -22,6 +22,8 @@ import { cn } from "~/utils/cn";
import { ProjectParamSchema, v3ProjectPath, v3TestPath } from "~/utils/pathBuilder";
import { ListPagination } from "../../components/ListPagination";
import { TextLink } from "~/components/primitives/TextLink";
import { Spinner } from "~/components/primitives/Spinner";
import { Suspense } from "react";
export const loader = async ({ request, params }: LoaderFunctionArgs) => {
const userId = await requireUserId(request);
@@ -33,7 +35,7 @@ export const loader = async ({ request, params }: LoaderFunctionArgs) => {
TaskRunListSearchFilters.parse(s);
const presenter = new RunListPresenter();
const list = await presenter.call({
const list = presenter.call({
userId,
projectSlug: projectParam,
tasks,
@@ -46,13 +48,13 @@ export const loader = async ({ request, params }: LoaderFunctionArgs) => {
cursor: cursor,
});
return typedjson({
list,
return typeddefer({
data: list,
});
};
export default function Page() {
const { list } = useTypedLoaderData<typeof loader>();
const { data } = useTypedLoaderData<typeof loader>();
const navigation = useNavigation();
const isLoading = navigation.state !== "idle";
const project = useProject();
@@ -64,36 +66,53 @@ export default function Page() {
<PageTitle title="Runs" />
</NavBar>
<PageBody>
{list.runs.length === 0 && !list.hasFilters ? (
list.possibleTasks.length === 0 ? (
<CreateFirstTaskInstructions />
) : (
<RunTaskInstructions />
)
) : (
<div className={cn("grid h-fit grid-cols-1 gap-4")}>
<div>
<div className="mb-2 flex items-center justify-between gap-x-2">
<RunsFilters
possibleEnvironments={project.environments}
possibleTasks={list.possibleTasks}
/>
<div className="flex items-center justify-end gap-x-2">
<ListPagination list={list} />
</div>
<Suspense
fallback={
<div className="flex items-center justify-center py-2">
<div className="mx-auto flex items-center gap-2">
<Spinner />
<Paragraph variant="small">Loading runs</Paragraph>
</div>
<TaskRunsTable
total={list.runs.length}
hasFilters={list.hasFilters}
filters={list.filters}
runs={list.runs}
isLoading={isLoading}
/>
<ListPagination list={list} className="mt-2 justify-end" />
</div>
</div>
)}
}
>
<TypedAwait resolve={data}>
{(list) => (
<>
{list.runs.length === 0 && !list.hasFilters ? (
list.possibleTasks.length === 0 ? (
<CreateFirstTaskInstructions />
) : (
<RunTaskInstructions />
)
) : (
<div className={cn("grid h-fit grid-cols-1 gap-4")}>
<div>
<div className="mb-2 flex items-center justify-between gap-x-2">
<RunsFilters
possibleEnvironments={project.environments}
possibleTasks={list.possibleTasks}
/>
<div className="flex items-center justify-end gap-x-2">
<ListPagination list={list} />
</div>
</div>
<TaskRunsTable
total={list.runs.length}
hasFilters={list.hasFilters}
filters={list.filters}
runs={list.runs}
isLoading={isLoading}
/>
<ListPagination list={list} className="mt-2 justify-end" />
</div>
</div>
)}
</>
)}
</TypedAwait>
</Suspense>
</PageBody>
</>
);
@@ -261,10 +261,10 @@ function SchedulesTable({
{schedule.userProvidedDeduplicationKey ? schedule.deduplicationKey : ""}
</TableCell>
<TableCell to={path} className={cellClass}>
<DateTime date={schedule.nextRun} />
<DateTime date={schedule.nextRun} timeZone="utc" />
</TableCell>
<TableCell to={path} className={cellClass}>
{schedule.lastRun ? <DateTime date={schedule.lastRun} /> : ""}
{schedule.lastRun ? <DateTime date={schedule.lastRun} timeZone="utc" /> : ""}
</TableCell>
<TableCell to={path} className={cellClass}>
<div className="flex gap-1">
@@ -9,6 +9,7 @@ import {
} from "~/components/environments/EnvironmentLabel";
import { PageBody, PageContainer } from "~/components/layout/AppLayout";
import { Header2 } from "~/components/primitives/Headers";
import { Input } from "~/components/primitives/Input";
import { NavBar, PageTitle } from "~/components/primitives/PageHeader";
import { Paragraph } from "~/components/primitives/Paragraph";
import { RadioButtonCircle } from "~/components/primitives/RadioButton";
@@ -20,6 +21,7 @@ import {
import { Spinner } from "~/components/primitives/Spinner";
import {
Table,
TableBlankRow,
TableBody,
TableCell,
TableHeader,
@@ -32,6 +34,7 @@ import { useLinkStatus } from "~/hooks/useLinkStatus";
import { useOptimisticLocation } from "~/hooks/useOptimisticLocation";
import { useOrganization } from "~/hooks/useOrganizations";
import { useProject } from "~/hooks/useProject";
import { useTextFilter } from "~/hooks/useTextFilter";
import {
SelectedEnvironment,
TaskListItem,
@@ -67,7 +70,7 @@ export default function Page() {
//get optimistic location for the segment control
const optimisticLocation = useOptimisticLocation();
const environment = new URLSearchParams(optimisticLocation.search).get("environment");
const environment = new URLSearchParams(optimisticLocation.search).get("environment") ?? "dev";
const navigation = useNavigation();
@@ -150,11 +153,50 @@ function TaskSelector({
tasks: TaskListItem[];
environmentSlug: string;
}) {
const organization = useOrganization();
const project = useProject();
const { filterText, setFilterText, filteredItems } = useTextFilter<TaskListItem>({
items: tasks,
filter: (task, text) => {
if (task.taskIdentifier.toLowerCase().includes(text.toLowerCase())) {
return true;
}
if (task.exportName.toLowerCase().includes(text.toLowerCase())) {
return true;
}
if (task.filePath.toLowerCase().includes(text.toLowerCase())) {
return true;
}
if (task.id.toLowerCase().includes(text.toLowerCase())) {
return true;
}
if (task.friendlyId.toLowerCase().includes(text.toLowerCase())) {
return true;
}
if (task.triggerSource === "SCHEDULED" && "scheduled".includes(text.toLowerCase())) {
return true;
}
return false;
},
});
return (
<div className="divide-y divide-charcoal-800 overflow-y-auto scrollbar-thin scrollbar-track-transparent scrollbar-thumb-charcoal-600">
<div className="px-2 pb-2">
<Input
placeholder="Search tasks"
variant="medium"
icon="search"
fullWidth={true}
value={filterText}
autoFocus
onChange={(e) => setFilterText(e.target.value)}
/>
</div>
<Table>
<TableHeader>
<TableRow>
@@ -166,42 +208,17 @@ function TaskSelector({
</TableRow>
</TableHeader>
<TableBody>
{tasks.map((t) => {
const path = v3TestTaskPath(organization, project, t, environmentSlug);
const { isActive, isPending } = useLinkStatus(path);
return (
<TableRow
key={t.taskIdentifier}
className={cn(
(isActive || isPending) &&
"z-20 rounded-sm outline outline-1 outline-offset-[-1px] outline-secondary"
)}
>
<TableCell to={path} actionClassName="pl-2.5 pr-1 py-1">
<RadioButtonCircle checked={isActive || isPending} />
</TableCell>
<TableCell to={path} actionClassName="pl-1 pr-2 py-1">
<div className="flex flex-col gap-0.5">
<TaskFunctionName
variant="extra-small"
functionName={t.exportName}
className="-ml-1 inline-flex"
/>
<div className="flex items-start gap-1">
<TaskTriggerSourceIcon source={t.triggerSource} className="size-3.5" />
<Paragraph variant="extra-small" className="text-text-dimmed">
{t.taskIdentifier}
</Paragraph>
</div>
</div>
</TableCell>
<TableCell to={path} actionClassName="px-2 py-1">
{t.filePath}
</TableCell>
</TableRow>
);
})}
{filteredItems.length > 0 ? (
filteredItems.map((t) => (
<TaskRow key={t.friendlyId} task={t} environmentSlug={environmentSlug} />
))
) : (
<TableBlankRow colSpan={3}>
<Paragraph spacing variant="small">
No tasks match "{filterText}"
</Paragraph>
</TableBlankRow>
)}
</TableBody>
</Table>
</div>
@@ -217,3 +234,43 @@ function NoTaskInstructions({ environment }: { environment?: SelectedEnvironment
</div>
);
}
function TaskRow({ task, environmentSlug }: { task: TaskListItem; environmentSlug: string }) {
const organization = useOrganization();
const project = useProject();
const path = v3TestTaskPath(organization, project, task, environmentSlug);
const { isActive, isPending } = useLinkStatus(path);
return (
<TableRow
key={task.taskIdentifier}
className={cn(
(isActive || isPending) &&
"z-20 rounded-sm outline outline-1 outline-offset-[-1px] outline-secondary"
)}
>
<TableCell to={path} actionClassName="pl-2.5 pr-1 py-1">
<RadioButtonCircle checked={isActive || isPending} />
</TableCell>
<TableCell to={path} actionClassName="pl-1 pr-2 py-1">
<div className="flex flex-col gap-0.5">
<TaskFunctionName
variant="extra-small"
functionName={task.exportName}
className="-ml-1 inline-flex"
/>
<div className="flex items-start gap-1">
<TaskTriggerSourceIcon source={task.triggerSource} className="size-3.5" />
<Paragraph variant="extra-small" className="text-text-dimmed">
{task.taskIdentifier}
</Paragraph>
</div>
</div>
</TableCell>
<TableCell to={path} actionClassName="px-2 py-1">
{task.filePath}
</TableCell>
</TableRow>
);
}
@@ -239,7 +239,7 @@ export default function Page() {
<FormError id={organizationSlug.errorId}>{organizationSlug.error}</FormError>
<FormError>{deleteForm.error}</FormError>
<Hint>
This change is irreversible, so please be certain. Type in the Organization slug
This change is irreversible, so please be certain. Type in the Organization slug{" "}
<InlineCode variant="extra-small">{organization.slug}</InlineCode> and then
press Delete.
</Hint>
@@ -342,7 +342,9 @@ function Timeline({ startTime, duration, inProgress, isError }: TimelineProps) {
<DateTimeAccurate date={startTime} />
</Paragraph>
{state === "pending" ? (
<LiveTimer startTime={startTime} className="" />
<Paragraph variant="extra-small" className={cn("whitespace-nowrap tabular-nums")}>
<LiveTimer startTime={startTime} />
</Paragraph>
) : (
<Paragraph variant="small">
<DateTimeAccurate
@@ -157,17 +157,17 @@ function TreeViewParent({
onSelectedIdChanged: (id) => {
console.log("onSelectedIdChanged", id);
},
onCollapsedIdsChanged: (ids) => {
console.log("onCollapsedIdsChanged", ids);
},
estimatedRowHeight: () => 32,
parentRef,
filter: (node) => {
if (filterText === "") return true;
if (node.data.title.toLowerCase().includes(filterText.toLowerCase())) {
return true;
}
return false;
filter: {
value: filterText,
fn: (text, node) => {
if (text === "") return true;
if (node.data.title.toLowerCase().includes(text.toLowerCase())) {
return true;
}
return false;
},
},
});
@@ -1,10 +1,17 @@
import { Prisma, PrismaClient } from "@trigger.dev/database";
import { Prisma, PrismaClient, RuntimeEnvironmentType } from "@trigger.dev/database";
import { z } from "zod";
import { environmentTitle } from "~/components/environments/EnvironmentLabel";
import { $transaction, prisma } from "~/db.server";
import { env } from "~/env.server";
import { getSecretStore } from "~/services/secrets/secretStore.server";
import { generateFriendlyId } from "../friendlyIdentifiers";
import { EnvironmentVariable, ProjectEnvironmentVariable, Repository, Result } from "./repository";
import { env } from "~/env.server";
import {
CreateResult,
EnvironmentVariable,
ProjectEnvironmentVariable,
Repository,
Result,
} from "./repository";
function secretKeyProjectPrefix(projectId: string) {
return `environmentvariable:${projectId}:`;
@@ -35,8 +42,15 @@ export class EnvironmentVariablesRepository implements Repository {
async create(
projectId: string,
userId: string,
options: { key: string; values: { value: string; environmentId: string }[] }
): Promise<Result> {
options: {
overwrite: boolean;
environmentIds: string[];
variables: {
key: string;
value: string;
}[];
}
): Promise<CreateResult> {
const project = await this.prismaClient.project.findUnique({
where: {
id: projectId,
@@ -55,6 +69,18 @@ export class EnvironmentVariablesRepository implements Repository {
id: true,
},
},
environmentVariables: {
select: {
key: true,
values: {
select: {
environment: {
select: { id: true, type: true },
},
},
},
},
},
},
});
@@ -62,58 +88,109 @@ export class EnvironmentVariablesRepository implements Repository {
return { success: false as const, error: "Project not found" };
}
if (options.values.every((v) => !project.environments.some((e) => e.id === v.environmentId))) {
if (options.environmentIds.every((v) => !project.environments.some((e) => e.id === v))) {
return { success: false as const, error: `Environment not found` };
}
//get rid of empty strings
const values = options.values.filter((v) => v.value.trim() !== "");
//get rid of empty variables
const values = options.variables.filter((v) => v.key.trim() !== "" && v.value.trim() !== "");
if (values.length === 0) {
return { success: false as const, error: `You must set at least one value` };
}
//check if any of them exist in an environment we're setting
if (!options.overwrite) {
const existingVariableKeys: { key: string; environments: RuntimeEnvironmentType[] }[] = [];
for (const variable of values) {
const existingVariable = project.environmentVariables.find((v) => v.key === variable.key);
if (
existingVariable &&
existingVariable.values.some((v) => options.environmentIds.includes(v.environment.id))
) {
existingVariableKeys.push({
key: variable.key,
environments: existingVariable.values
.filter((v) => options.environmentIds.includes(v.environment.id))
.map((v) => v.environment.type),
});
}
}
if (existingVariableKeys.length > 0) {
return {
success: false as const,
error: `Some of the variables are already set for these environments`,
variableErrors: existingVariableKeys.map((val) => ({
key: val.key,
error: `Variable already set in ${val.environments
.map((e) => environmentTitle({ type: e }))
.join(", ")}.`,
})),
};
}
}
try {
const result = await $transaction(this.prismaClient, async (tx) => {
const environmentVariable = await tx.environmentVariable.create({
data: {
key: options.key,
friendlyId: generateFriendlyId("envvar"),
project: {
connect: {
id: projectId,
for (const variable of values) {
const environmentVariable = await tx.environmentVariable.upsert({
where: {
projectId_key: {
key: variable.key,
projectId,
},
},
},
});
const secretStore = getSecretStore("DATABASE", {
prismaClient: tx,
});
//create the secret values and references
for (const value of values) {
const key = secretKey(projectId, value.environmentId, options.key);
//create the secret reference
const secretReference = await tx.secretReference.create({
data: {
key,
provider: "DATABASE",
create: {
key: variable.key,
friendlyId: generateFriendlyId("envvar"),
project: {
connect: {
id: projectId,
},
},
},
update: {},
});
const variableValue = await tx.environmentVariableValue.create({
data: {
variableId: environmentVariable.id,
environmentId: value.environmentId,
valueReferenceId: secretReference.id,
},
const secretStore = getSecretStore("DATABASE", {
prismaClient: tx,
});
await secretStore.setSecret<{ secret: string }>(key, {
secret: value.value,
});
//set the secret values and references
for (const environmentId of options.environmentIds) {
const key = secretKey(projectId, environmentId, variable.key);
//create the secret reference
const secretReference = await tx.secretReference.upsert({
where: {
key,
},
create: {
key,
provider: "DATABASE",
},
update: {},
});
const variableValue = await tx.environmentVariableValue.upsert({
where: {
variableId_environmentId: {
variableId: environmentVariable.id,
environmentId,
},
},
create: {
variableId: environmentVariable.id,
environmentId: environmentId,
valueReferenceId: secretReference.id,
},
update: {},
});
await secretStore.setSecret<{ secret: string }>(key, {
secret: variable.value,
});
}
}
});
@@ -126,7 +203,7 @@ export class EnvironmentVariablesRepository implements Repository {
if (error.code === "P2002") {
return {
success: false as const,
error: `There's already an environment variable called ${options.key}.`,
error: `There was already an existing field`,
};
}
}
@@ -1,22 +1,27 @@
import { RuntimeEnvironmentType } from "@trigger.dev/database";
import { z } from "zod";
const EnvironmentVariable = z
export const EnvironmentVariableKey = z
.string()
.nonempty("Environment variable key is required")
.regex(/^\w+$/, "Environment variables can only contain alphanumeric characters and underscores");
.nonempty("Key is required")
.regex(/^\w+$/, "Keys can only use alphanumeric characters and underscores");
export const CreateEnvironmentVariable = z.object({
key: EnvironmentVariable,
values: z.array(
z.object({
environmentId: z.string(),
value: z.string(),
})
),
export const CreateEnvironmentVariables = z.object({
environmentIds: z.array(z.string()),
variables: z.array(z.object({ key: EnvironmentVariableKey, value: z.string() })),
});
export type CreateEnvironmentVariable = z.infer<typeof CreateEnvironmentVariable>;
export type CreateEnvironmentVariables = z.infer<typeof CreateEnvironmentVariables>;
export type CreateResult =
| {
success: true;
}
| {
success: false;
error: string;
variableErrors?: { key: string; error: string }[];
};
export const EditEnvironmentVariable = z.object({
id: z.string(),
@@ -60,7 +65,11 @@ export type EnvironmentVariable = {
};
export interface Repository {
create(projectId: string, userId: string, options: CreateEnvironmentVariable): Promise<Result>;
create(
projectId: string,
userId: string,
options: CreateEnvironmentVariables
): Promise<CreateResult>;
edit(projectId: string, userId: string, options: EditEnvironmentVariable): Promise<Result>;
getProject(projectId: string, userId: string): Promise<ProjectEnvironmentVariable[]>;
getEnvironment(
+65 -17
View File
@@ -231,7 +231,7 @@ export class MarQS {
return;
}
const message = await this.#readMessage(messageData.messageId);
const message = await this.readMessage(messageData.messageId);
if (message) {
span.setAttributes({
@@ -308,7 +308,7 @@ export class MarQS {
return;
}
const message = await this.#readMessage(messageData.messageId);
const message = await this.readMessage(messageData.messageId);
if (message) {
span.setAttributes({
@@ -336,7 +336,7 @@ export class MarQS {
return this.#trace(
"acknowledgeMessage",
async (span) => {
const message = await this.#readMessage(messageId);
const message = await this.readMessage(messageId);
if (!message) {
return;
@@ -374,12 +374,13 @@ export class MarQS {
public async replaceMessage(
messageId: string,
messageData: Record<string, unknown>,
timestamp?: number
timestamp?: number,
inplace?: boolean
) {
return this.#trace(
"replaceMessage",
async (span) => {
const oldMessage = await this.#readMessage(messageId);
const oldMessage = await this.readMessage(messageId);
if (!oldMessage) {
return;
@@ -392,6 +393,27 @@ export class MarQS {
[SemanticAttributes.PARENT_QUEUE]: oldMessage.parentQueue,
});
const traceContext = {
traceparent: oldMessage.data.traceparent,
tracestate: oldMessage.data.tracestate,
};
const newMessage: MessagePayload = {
version: "1",
// preserve original trace context
data: { ...messageData, ...traceContext },
queue: oldMessage.queue,
concurrencyKey: oldMessage.concurrencyKey,
timestamp: timestamp ?? Date.now(),
messageId,
parentQueue: oldMessage.parentQueue,
};
if (inplace) {
await this.#callReplaceMessage(newMessage);
return;
}
await this.#callAcknowledgeMessage({
parentQueue: oldMessage.parentQueue,
messageKey: this.keys.messageKey(messageId),
@@ -403,16 +425,6 @@ export class MarQS {
messageId,
});
const newMessage: MessagePayload = {
version: "1",
data: messageData,
queue: oldMessage.queue,
concurrencyKey: oldMessage.concurrencyKey,
timestamp: timestamp ?? Date.now(),
messageId,
parentQueue: oldMessage.parentQueue,
};
await this.#callEnqueueMessage(newMessage);
},
{
@@ -455,7 +467,7 @@ export class MarQS {
return this.#trace(
"nackMessage",
async (span) => {
const message = await this.#readMessage(messageId);
const message = await this.readMessage(messageId);
if (!message) {
return;
@@ -505,7 +517,7 @@ export class MarQS {
return this.options.visibilityTimeoutInMs ?? 300000;
}
async #readMessage(messageId: string) {
async readMessage(messageId: string) {
return this.#trace(
"readMessage",
async (span) => {
@@ -881,6 +893,17 @@ export class MarQS {
};
}
async #callReplaceMessage(message: MessagePayload) {
logger.debug("Calling replaceMessage", {
messagePayload: message,
});
return this.redis.replaceMessage(
this.keys.messageKey(message.messageId),
JSON.stringify(message)
);
}
async #callAcknowledgeMessage({
parentQueue,
messageKey,
@@ -1185,6 +1208,25 @@ return {messageId, messageScore} -- Return message details
`,
});
this.redis.defineCommand("replaceMessage", {
numberOfKeys: 1,
lua: `
local messageKey = KEYS[1]
local messageData = ARGV[1]
-- Check if message exists
local existingMessage = redis.call('GET', messageKey)
-- Do nothing if it doesn't
if #existingMessage == nil then
return nil
end
-- Replace the message
redis.call('SET', messageKey, messageData, 'GET')
`,
});
this.redis.defineCommand("acknowledgeMessage", {
numberOfKeys: 7,
lua: `
@@ -1406,6 +1448,12 @@ declare module "ioredis" {
callback?: Callback<[string, string]>
): Result<[string, string] | null, Context>;
replaceMessage(
messageKey: string,
messageData: string,
callback?: Callback<void>
): Result<void, Context>;
acknowledgeMessage(
parentQueue: string,
messageKey: string,
@@ -28,13 +28,14 @@ import { socketIo } from "../handleSocketIo.server";
import { findCurrentWorkerDeployment } from "../models/workerDeployment.server";
import { RestoreCheckpointService } from "../services/restoreCheckpoint.server";
import { tracer } from "../tracer.server";
import { CrashTaskRunService } from "../services/crashTaskRun.server";
const WithTraceContext = z.object({
traceparent: z.string().optional(),
tracestate: z.string().optional(),
});
const MessageBody = z.discriminatedUnion("type", [
export const SharedQueueMessageBody = z.discriminatedUnion("type", [
WithTraceContext.extend({
type: z.literal("EXECUTE"),
taskIdentifier: z.string(),
@@ -51,8 +52,14 @@ const MessageBody = z.discriminatedUnion("type", [
resumableAttemptId: z.string(),
checkpointEventId: z.string(),
}),
WithTraceContext.extend({
type: z.literal("FAIL"),
reason: z.string(),
}),
]);
export type SharedQueueMessageBody = z.infer<typeof SharedQueueMessageBody>;
type BackgroundWorkerWithTasks = BackgroundWorker & { tasks: BackgroundWorkerTask[] };
export type SharedQueueConsumerOptions = {
@@ -233,7 +240,7 @@ export class SharedQueueConsumer {
logger.log("dequeueMessageInSharedQueue()", { queueMessage: message });
const messageBody = MessageBody.safeParse(message.data);
const messageBody = SharedQueueMessageBody.safeParse(message.data);
if (!messageBody.success) {
logger.error("Failed to parse message", {
@@ -411,11 +418,21 @@ export class SharedQueueConsumer {
});
if (!queue) {
logger.debug("SharedQueueConsumer queue not found, so nacking message", {
queueMessage: message,
taskRunQueue: lockedTaskRun.queue,
runtimeEnvironmentId: lockedTaskRun.runtimeEnvironmentId,
});
await this.#nackAndDoMoreWork(message.messageId, this._options.nextTickInterval);
return;
}
if (!this._enabled) {
logger.debug("SharedQueueConsumer not enabled, so nacking message", {
queueMessage: message,
});
await marqs?.nackMessage(message.messageId);
return;
}
@@ -521,6 +538,11 @@ export class SharedQueueConsumer {
}),
]);
logger.error("SharedQueueConsumer errored, so nacking message", {
queueMessage: message,
error: e instanceof Error ? { name: e.name, message: e.message, stack: e.stack } : e,
});
await this.#nackAndDoMoreWork(message.messageId);
return;
}
@@ -739,6 +761,34 @@ export class SharedQueueConsumer {
break;
}
// Fail for whatever reason, usually runs that have been resumed but stopped heartbeating
case "FAIL": {
const existingTaskRun = await prisma.taskRun.findUnique({
where: {
id: message.messageId,
},
});
if (!existingTaskRun) {
logger.error("No existing task run to fail", {
queueMessage: messageBody,
messageId: message.messageId,
});
await this.#ackAndDoMoreWork(message.messageId);
return;
}
// TODO: Consider failing the attempt and retrying instead. This may not be a good idea, as dequeued FAIL messages tend to point towards critical, persistent errors.
const service = new CrashTaskRunService();
await service.call(existingTaskRun.id, {
crashAttempts: true,
reason: messageBody.data.reason,
});
await this.#ackAndDoMoreWork(message.messageId);
return;
}
}
this.#doMoreWork();
@@ -10,6 +10,7 @@ import { generateFriendlyId } from "../friendlyIdentifiers";
import { marqs } from "~/v3/marqs/index.server";
import { CreateCheckpointRestoreEventService } from "./createCheckpointRestoreEvent.server";
import { BaseService } from "./baseService.server";
import { CrashTaskRunService } from "./crashTaskRun.server";
const FREEZABLE_RUN_STATUSES: TaskRunStatus[] = ["EXECUTING", "RETRYING_AFTER_FAILURE"];
const FREEZABLE_ATTEMPT_STATUSES: TaskRunAttemptStatus[] = ["EXECUTING", "FAILED"];
@@ -61,6 +62,14 @@ export class CreateCheckpointService extends BaseService {
status: attempt.taskRun.status,
},
});
// This should only affect CLIs < beta.24, in very limited scenarios
const service = new CrashTaskRunService(this._prisma);
await service.call(attempt.taskRunId, {
crashAttempts: true,
reason: "Unfreezable state: Please upgrade your CLI",
});
return;
}
@@ -2,13 +2,14 @@ import {
CoordinatorToPlatformMessages,
TaskRunExecution,
TaskRunExecutionResult,
WaitReason,
} from "@trigger.dev/core/v3";
import type { InferSocketMessageSchema } from "@trigger.dev/core/v3/zodSocket";
import { $transaction, PrismaClientOrTransaction } from "~/db.server";
import { logger } from "~/services/logger.server";
import { marqs } from "~/v3/marqs/index.server";
import { socketIo } from "../handleSocketIo.server";
import { sharedQueueTasks } from "../marqs/sharedQueueConsumer.server";
import { SharedQueueMessageBody, sharedQueueTasks } from "../marqs/sharedQueueConsumer.server";
import { BaseService } from "./baseService.server";
import { TaskRunAttempt } from "@trigger.dev/database";
@@ -91,12 +92,13 @@ export class ResumeAttemptService extends BaseService {
switch (params.type) {
case "WAIT_FOR_DURATION": {
logger.error(
"Attempt requested resume after duration wait, this is unexpected and likely a bug",
{ attemptId: attempt.id }
);
logger.debug("Sending duration wait resume message", {
attemptId: attempt.id,
attemptFriendlyId: params.attemptFriendlyId,
});
await this.#setPostResumeStatuses(attempt, tx);
// Attempts should not request resume for duration waits, this is just here as a backup
socketIo.coordinatorNamespace.emit("RESUME_AFTER_DURATION", {
version: "v1",
attemptId: attempt.id,
@@ -119,6 +121,9 @@ export class ResumeAttemptService extends BaseService {
logger.error("No task dependency", { attemptId: attempt.id });
return;
}
await this.#handleDependencyResume(attempt, completedAttemptIds, tx);
break;
}
case "WAIT_FOR_BATCH": {
@@ -136,6 +141,9 @@ export class ResumeAttemptService extends BaseService {
logger.error("No batch dependency", { attemptId: attempt.id });
return;
}
await this.#handleDependencyResume(attempt, completedAttemptIds, tx);
break;
}
default: {
@@ -143,7 +151,8 @@ export class ResumeAttemptService extends BaseService {
}
}
await this.#handleDependencyResume(attempt, completedAttemptIds, tx);
// Prevent infinite restores by failing runs that don't heartbeat after post-restore resume requests
await this.#replaceResumeWithFailMessage(attempt.taskRunId, params.type);
});
}
@@ -215,7 +224,20 @@ export class ResumeAttemptService extends BaseService {
executions.push(executionPayload.execution);
}
const updated = await tx.taskRunAttempt.update({
await this.#setPostResumeStatuses(attempt, tx);
socketIo.coordinatorNamespace.emit("RESUME_AFTER_DEPENDENCY", {
version: "v1",
runId: attempt.taskRunId,
attemptId: attempt.id,
attemptFriendlyId: attempt.friendlyId,
completions,
executions,
});
}
async #setPostResumeStatuses(attempt: TaskRunAttempt, tx: PrismaClientOrTransaction) {
return await tx.taskRunAttempt.update({
where: {
id: attempt.id,
},
@@ -230,14 +252,51 @@ export class ResumeAttemptService extends BaseService {
},
},
});
}
socketIo.coordinatorNamespace.emit("RESUME_AFTER_DEPENDENCY", {
version: "v1",
runId: attempt.taskRunId,
attemptId: attempt.id,
attemptFriendlyId: attempt.friendlyId,
completions,
executions,
});
async #replaceResumeWithFailMessage(messageId: string, waitReason: WaitReason) {
const currentMessage = await marqs?.readMessage(messageId);
if (!currentMessage) {
logger.debug("No message to replace", { messageId, waitReason });
return;
}
const currentBody = SharedQueueMessageBody.safeParse(currentMessage.data);
if (!currentBody.success) {
logger.debug("Invalid message body", { messageId, waitReason, currentBody });
return;
}
const currentType = currentBody.data.type;
if (currentType !== "RESUME" && currentType !== "RESUME_AFTER_DURATION") {
logger.debug("Not a resume message", { messageId, waitReason, currentBody });
return;
}
let reason = "Worker unresponsive after restore";
switch (waitReason) {
case "WAIT_FOR_DURATION":
reason = "Worker unresponsive after waiting for duration";
break;
case "WAIT_FOR_TASK":
reason = "Worker unresponsive after waiting for task";
break;
case "WAIT_FOR_BATCH":
reason = "Worker unresponsive after waiting for batch task";
break;
default:
break;
}
const failMessage: SharedQueueMessageBody = {
type: "FAIL",
reason,
};
return await marqs?.replaceMessage(messageId, failMessage, undefined, true);
}
}
@@ -76,6 +76,15 @@ export class TriggerTaskService extends BaseService {
immediate: true,
},
async (event, traceContext) => {
const runFriendlyId = generateFriendlyId("run");
const payloadPacket = await this.#handlePayloadPacket(
body.payload,
body.options?.payloadType ?? "application/json",
runFriendlyId,
environment
);
const lockId = taskIdentifierToLockId(taskId);
const run = await $transaction(this._prisma, async (tx) => {
@@ -105,15 +114,6 @@ export class TriggerTaskService extends BaseService {
event.setAttribute("queueName", queueName);
span.setAttribute("queueName", queueName);
const runFriendlyId = generateFriendlyId("run");
const payloadPacket = await this.#handlePayloadPacket(
body.payload,
body.options?.payloadType ?? "application/json",
runFriendlyId,
environment
);
const taskRun = await tx.taskRun.create({
data: {
status: "PENDING",
+6 -1
View File
@@ -16,6 +16,8 @@
"typecheck": "tsc -p ./tsconfig.check.json",
"db:seed": "node prisma/seed.js",
"db:seed:local": "ts-node prisma/seed.ts",
"build:db:populate": "esbuild --platform=node --bundle --minify --format=cjs ./prisma/populate.ts --outdir=prisma",
"db:populate": "node prisma/populate.js --",
"generate:sourcemaps": "remix build --sourcemap",
"clean:sourcemaps": "run-s clean:sourcemaps:*",
"clean:sourcemaps:public": "rimraf ./build/**/*.map",
@@ -58,6 +60,7 @@
"@opentelemetry/sdk-trace-base": "^1.22.0",
"@opentelemetry/sdk-trace-node": "^1.22.0",
"@opentelemetry/semantic-conventions": "^1.22.0",
"@popperjs/core": "^2.11.8",
"@prisma/instrumentation": "^5.11.0",
"@radix-ui/react-alert-dialog": "^1.0.4",
"@radix-ui/react-dialog": "^1.0.3",
@@ -105,6 +108,7 @@
"cronstrue": "^2.21.0",
"cross-env": "^7.0.3",
"cuid": "^2.1.8",
"dotenv": "^16.4.5",
"emails": "workspace:*",
"evt": "^2.4.13",
"express": "^4.18.1",
@@ -136,6 +140,7 @@
"react-collapse": "^5.1.1",
"react-dom": "^18.2.0",
"react-hotkeys-hook": "^4.4.1",
"react-popper": "^2.3.0",
"react-resizable-panels": "^2.0.9",
"react-stately": "^3.29.1",
"react-use": "^17.4.0",
@@ -227,4 +232,4 @@
"engines": {
"node": ">=16.0.0"
}
}
}
+101
View File
@@ -0,0 +1,101 @@
// Bulk adds data to the database for testing
// Call it like this
// 1. pnpm run build:db:populate
// 2. pnpm run db:populate -- --projectRef=proj_liazlkfgmfcusswwgohl --taskIdentifier=child-task --runCount=100000
import { generateFriendlyId } from "~/v3/friendlyIdentifiers";
import { prisma } from "../app/db.server";
async function populate() {
if (process.env.NODE_ENV !== "development") {
return;
}
const projectRef = getArg("projectRef");
if (!projectRef) {
throw new Error("projectRef is required");
}
const project = await prisma.project.findUnique({
include: {
environments: true,
},
where: {
externalRef: projectRef,
},
});
if (!project) {
throw new Error("Project not found");
}
const taskIdentifier = getArg("taskIdentifier");
if (!taskIdentifier) {
throw new Error("taskIdentifier is required");
}
const runCount = parseInt(getArg("runCount") || "100");
const task = await prisma.backgroundWorkerTask.findFirst({
where: {
projectId: project.id,
slug: taskIdentifier,
},
orderBy: {
createdAt: "desc",
},
});
if (!task) {
throw new Error("Task not found");
}
const runs = await prisma.taskRun.createMany({
data: Array(runCount)
.fill(0)
.map((_, index) => {
const friendlyId = generateFriendlyId("run");
return {
status: "CANCELED",
number: index + 1,
friendlyId,
runtimeEnvironmentId: project.environments[randomIndex(project.environments)].id,
projectId: project.id,
taskIdentifier,
payload: JSON.stringify({ foo: "bar" }),
traceId: "traceId",
spanId: "spanId",
queue: "task/${taskIdentifier}",
};
}),
skipDuplicates: true,
});
console.log(`Added ${runs.count} runs`);
}
function getArg(name: string) {
const args = process.argv.slice(2);
let value = "";
args.forEach((val) => {
if (val.startsWith(`--${name}=`)) {
value = val.split("=")[1];
}
});
return !value ? undefined : value;
}
function randomIndex<T>(array: T[]) {
return Math.floor(Math.random() * array.length);
}
populate()
.catch((e) => {
console.error(e);
process.exit(1);
})
.finally(async () => {
await prisma.$disconnect();
});
+1
View File
@@ -21,6 +21,7 @@ module.exports = {
"random-words",
"superjson",
],
browserNodeBuiltinsPolyfill: { modules: { path: true, os: true, crypto: true } },
watchPaths: async () => {
return [
"../../packages/core/src/**/*",
+2
View File
@@ -66,6 +66,8 @@ We provide some useful functions that you can use to retry smaller parts of a ta
You can retry a block of code that can throw an error, with the same retry settings as a task.
```ts /trigger/retry-on-throw.ts
import { task, logger, retry } from "@trigger.dev/sdk/v3"
export const retryOnThrow = task({
id: "retry-on-throw",
run: async (payload: any) => {
+35 -2
View File
@@ -7,8 +7,10 @@ This simple GitHub action file will deploy you Trigger.dev tasks when new code i
<Warning>The deploy step will fail if any version mismatches are detected. Please see the [version pinning](/v3/github-actions#version-pinning) section for more details.</Warning>
```yaml .github/workflows/release-trigger.yml
name: Deploy to Trigger.dev
<CodeGroup>
```yaml .github/workflows/release-trigger-prod.yml
name: Deploy to Trigger.dev (prod)
on:
push:
@@ -39,6 +41,37 @@ jobs:
npx trigger.dev@beta deploy
```
```yaml .github/workflows/release-trigger-staging.yml
name: Deploy to Trigger.dev (staging)
# Requires manually calling the workflow from a branch / commit to deploy to staging
on:
workflow_dispatch:
jobs:
deploy:
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v2
- name: Use Node.js 20.x
uses: actions/setup-node@v4
with:
node-version: "20.x"
- name: Install dependencies
run: npm install
- name: 🚀 Deploy Trigger.dev
env:
TRIGGER_ACCESS_TOKEN: ${{ secrets.TRIGGER_ACCESS_TOKEN }}
run: |
npx trigger.dev@beta deploy --env staging
```
</CodeGroup>
If you already have a GitHub action file, you can just add the final step "🚀 Deploy Trigger.dev" to your existing file.
You need to add the `TRIGGER_ACCESS_TOKEN` secret to your repository. You can create a new access token by going to your profile page and then clicking on the "Personal Access Tokens" tab.
+14
View File
@@ -1,5 +1,19 @@
# @trigger.dev/airtable
## 3.0.0-beta.24
### Patch Changes
- @trigger.dev/integration-kit@3.0.0-beta.24
- @trigger.dev/sdk@3.0.0-beta.24
## 3.0.0-beta.23
### Patch Changes
- @trigger.dev/integration-kit@3.0.0-beta.23
- @trigger.dev/sdk@3.0.0-beta.23
## 3.0.0-beta.22
### Patch Changes
+3 -3
View File
@@ -1,6 +1,6 @@
{
"name": "@trigger.dev/airtable",
"version": "3.0.0-beta.22",
"version": "3.0.0-beta.24",
"description": "Trigger.dev integration for airtable",
"main": "./dist/index.js",
"types": "./dist/index.d.ts",
@@ -25,8 +25,8 @@
"typecheck": "tsc --noEmit"
},
"dependencies": {
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.22",
"@trigger.dev/sdk": "workspace:^3.0.0-beta.22",
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.24",
"@trigger.dev/sdk": "workspace:^3.0.0-beta.24",
"airtable": "^0.12.1",
"zod": "3.22.3"
},
+14
View File
@@ -1,5 +1,19 @@
# @trigger.dev/github
## 3.0.0-beta.24
### Patch Changes
- @trigger.dev/integration-kit@3.0.0-beta.24
- @trigger.dev/sdk@3.0.0-beta.24
## 3.0.0-beta.23
### Patch Changes
- @trigger.dev/integration-kit@3.0.0-beta.23
- @trigger.dev/sdk@3.0.0-beta.23
## 3.0.0-beta.22
### Patch Changes
+3 -3
View File
@@ -1,6 +1,6 @@
{
"name": "@trigger.dev/github",
"version": "3.0.0-beta.22",
"version": "3.0.0-beta.24",
"description": "The official GitHub integration for Trigger.dev",
"main": "./dist/index.js",
"types": "./dist/index.d.ts",
@@ -30,8 +30,8 @@
"@octokit/request-error": "^5.0.1",
"@octokit/webhooks": "^12.0.10",
"octokit": "^3.1.2",
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.22",
"@trigger.dev/sdk": "workspace:^3.0.0-beta.22",
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.24",
"@trigger.dev/sdk": "workspace:^3.0.0-beta.24",
"zod": "3.22.3"
},
"engines": {
+14
View File
@@ -1,5 +1,19 @@
# @trigger.dev/linear
## 3.0.0-beta.24
### Patch Changes
- @trigger.dev/integration-kit@3.0.0-beta.24
- @trigger.dev/sdk@3.0.0-beta.24
## 3.0.0-beta.23
### Patch Changes
- @trigger.dev/integration-kit@3.0.0-beta.23
- @trigger.dev/sdk@3.0.0-beta.23
## 3.0.0-beta.22
### Patch Changes
+3 -3
View File
@@ -1,6 +1,6 @@
{
"name": "@trigger.dev/linear",
"version": "3.0.0-beta.22",
"version": "3.0.0-beta.24",
"description": "Trigger.dev integration for @linear/sdk",
"main": "./dist/index.js",
"types": "./dist/index.d.ts",
@@ -26,8 +26,8 @@
},
"dependencies": {
"@linear/sdk": "^8.0.0",
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.22",
"@trigger.dev/sdk": "workspace:^3.0.0-beta.22",
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.24",
"@trigger.dev/sdk": "workspace:^3.0.0-beta.24",
"zod": "3.22.3"
},
"engines": {
+14
View File
@@ -1,5 +1,19 @@
# @trigger.dev/slack
## 3.0.0-beta.24
### Patch Changes
- @trigger.dev/integration-kit@3.0.0-beta.24
- @trigger.dev/sdk@3.0.0-beta.24
## 3.0.0-beta.23
### Patch Changes
- @trigger.dev/integration-kit@3.0.0-beta.23
- @trigger.dev/sdk@3.0.0-beta.23
## 3.0.0-beta.22
### Patch Changes
+3 -3
View File
@@ -1,6 +1,6 @@
{
"name": "@trigger.dev/openai",
"version": "3.0.0-beta.22",
"version": "3.0.0-beta.24",
"description": "The official OpenAI integration for Trigger.dev",
"main": "./dist/index.js",
"types": "./dist/index.d.ts",
@@ -42,8 +42,8 @@
},
"dependencies": {
"openai": "^4.16.1",
"@trigger.dev/sdk": "workspace:^3.0.0-beta.22",
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.22"
"@trigger.dev/sdk": "workspace:^3.0.0-beta.24",
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.24"
},
"engines": {
"node": ">=18.0.0"
+14
View File
@@ -1,5 +1,19 @@
# @trigger.dev/plain
## 3.0.0-beta.24
### Patch Changes
- @trigger.dev/integration-kit@3.0.0-beta.24
- @trigger.dev/sdk@3.0.0-beta.24
## 3.0.0-beta.23
### Patch Changes
- @trigger.dev/integration-kit@3.0.0-beta.23
- @trigger.dev/sdk@3.0.0-beta.23
## 3.0.0-beta.22
### Patch Changes
+3 -3
View File
@@ -1,6 +1,6 @@
{
"name": "@trigger.dev/plain",
"version": "3.0.0-beta.22",
"version": "3.0.0-beta.24",
"description": "The official Plain.com integration for Trigger.dev",
"main": "./dist/index.js",
"types": "./dist/index.d.ts",
@@ -24,8 +24,8 @@
"build:tsup": "tsup"
},
"dependencies": {
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.22",
"@trigger.dev/sdk": "workspace:^3.0.0-beta.22",
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.24",
"@trigger.dev/sdk": "workspace:^3.0.0-beta.24",
"@team-plain/typescript-sdk": "^2.7.0"
},
"engines": {
+14
View File
@@ -1,5 +1,19 @@
# @trigger.dev/replicate
## 3.0.0-beta.24
### Patch Changes
- @trigger.dev/integration-kit@3.0.0-beta.24
- @trigger.dev/sdk@3.0.0-beta.24
## 3.0.0-beta.23
### Patch Changes
- @trigger.dev/integration-kit@3.0.0-beta.23
- @trigger.dev/sdk@3.0.0-beta.23
## 3.0.0-beta.22
### Patch Changes
+3 -3
View File
@@ -1,6 +1,6 @@
{
"name": "@trigger.dev/replicate",
"version": "3.0.0-beta.22",
"version": "3.0.0-beta.24",
"description": "Trigger.dev integration for replicate",
"main": "./dist/index.js",
"types": "./dist/index.d.ts",
@@ -25,8 +25,8 @@
"typecheck": "tsc --noEmit"
},
"dependencies": {
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.22",
"@trigger.dev/sdk": "workspace:^3.0.0-beta.22",
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.24",
"@trigger.dev/sdk": "workspace:^3.0.0-beta.24",
"replicate": "^0.18.1",
"zod": "3.22.3"
},
+14
View File
@@ -1,5 +1,19 @@
# @trigger.dev/resend
## 3.0.0-beta.24
### Patch Changes
- @trigger.dev/integration-kit@3.0.0-beta.24
- @trigger.dev/sdk@3.0.0-beta.24
## 3.0.0-beta.23
### Patch Changes
- @trigger.dev/integration-kit@3.0.0-beta.23
- @trigger.dev/sdk@3.0.0-beta.23
## 3.0.0-beta.22
### Patch Changes
+3 -3
View File
@@ -1,6 +1,6 @@
{
"name": "@trigger.dev/resend",
"version": "3.0.0-beta.22",
"version": "3.0.0-beta.24",
"description": "The official Resend.com integration for Trigger.dev",
"main": "./dist/index.js",
"types": "./dist/index.d.ts",
@@ -24,8 +24,8 @@
"build:tsup": "tsup"
},
"dependencies": {
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.22",
"@trigger.dev/sdk": "workspace:^3.0.0-beta.22",
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.24",
"@trigger.dev/sdk": "workspace:^3.0.0-beta.24",
"resend": "^2.1.0"
},
"engines": {
+14
View File
@@ -1,5 +1,19 @@
# @trigger.dev/sendgrid
## 3.0.0-beta.24
### Patch Changes
- @trigger.dev/integration-kit@3.0.0-beta.24
- @trigger.dev/sdk@3.0.0-beta.24
## 3.0.0-beta.23
### Patch Changes
- @trigger.dev/integration-kit@3.0.0-beta.23
- @trigger.dev/sdk@3.0.0-beta.23
## 3.0.0-beta.22
### Patch Changes
+3 -3
View File
@@ -1,6 +1,6 @@
{
"name": "@trigger.dev/sendgrid",
"version": "3.0.0-beta.22",
"version": "3.0.0-beta.24",
"description": "Trigger.dev integration for @sendgrid/mail",
"main": "./dist/index.js",
"types": "./dist/index.d.ts",
@@ -26,8 +26,8 @@
},
"dependencies": {
"@sendgrid/mail": "^7.7.0",
"@trigger.dev/sdk": "workspace:^3.0.0-beta.22",
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.22"
"@trigger.dev/sdk": "workspace:^3.0.0-beta.24",
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.24"
},
"engines": {
"node": ">=16.8.0"
+14
View File
@@ -1,5 +1,19 @@
# @trigger.dev/shopify
## 3.0.0-beta.24
### Patch Changes
- @trigger.dev/integration-kit@3.0.0-beta.24
- @trigger.dev/sdk@3.0.0-beta.24
## 3.0.0-beta.23
### Patch Changes
- @trigger.dev/integration-kit@3.0.0-beta.23
- @trigger.dev/sdk@3.0.0-beta.23
## 3.0.0-beta.22
### Patch Changes
+3 -3
View File
@@ -1,6 +1,6 @@
{
"name": "@trigger.dev/shopify",
"version": "3.0.0-beta.22",
"version": "3.0.0-beta.24",
"description": "Trigger.dev integration for @shopify/shopify-api",
"main": "./dist/index.js",
"types": "./dist/index.d.ts",
@@ -26,8 +26,8 @@
},
"dependencies": {
"@shopify/shopify-api": "^8.0.2",
"@trigger.dev/sdk": "workspace:^3.0.0-beta.22",
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.22",
"@trigger.dev/sdk": "workspace:^3.0.0-beta.24",
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.24",
"zod": "3.22.3"
},
"engines": {
+12
View File
@@ -1,5 +1,17 @@
# @trigger.dev/slack
## 3.0.0-beta.24
### Patch Changes
- @trigger.dev/sdk@3.0.0-beta.24
## 3.0.0-beta.23
### Patch Changes
- @trigger.dev/sdk@3.0.0-beta.23
## 3.0.0-beta.22
### Patch Changes
+2 -2
View File
@@ -1,6 +1,6 @@
{
"name": "@trigger.dev/slack",
"version": "3.0.0-beta.22",
"version": "3.0.0-beta.24",
"description": "The official Slack integration for Trigger.dev",
"main": "./dist/index.js",
"types": "./dist/index.d.ts",
@@ -25,7 +25,7 @@
},
"dependencies": {
"@slack/web-api": "^6.8.1",
"@trigger.dev/sdk": "workspace:^3.0.0-beta.22",
"@trigger.dev/sdk": "workspace:^3.0.0-beta.24",
"zod": "3.22.3"
},
"engines": {
+14
View File
@@ -1,5 +1,19 @@
# @trigger.dev/stripe
## 3.0.0-beta.24
### Patch Changes
- @trigger.dev/integration-kit@3.0.0-beta.24
- @trigger.dev/sdk@3.0.0-beta.24
## 3.0.0-beta.23
### Patch Changes
- @trigger.dev/integration-kit@3.0.0-beta.23
- @trigger.dev/sdk@3.0.0-beta.23
## 3.0.0-beta.22
### Patch Changes
+3 -3
View File
@@ -1,6 +1,6 @@
{
"name": "@trigger.dev/stripe",
"version": "3.0.0-beta.22",
"version": "3.0.0-beta.24",
"description": "Trigger.dev integration for stripe",
"main": "./dist/index.js",
"types": "./dist/index.d.ts",
@@ -25,8 +25,8 @@
"typecheck": "tsc --noEmit"
},
"dependencies": {
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.22",
"@trigger.dev/sdk": "workspace:^3.0.0-beta.22",
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.24",
"@trigger.dev/sdk": "workspace:^3.0.0-beta.24",
"stripe": "^12.14.0",
"zod": "3.22.3"
},
+14
View File
@@ -1,5 +1,19 @@
# @trigger.dev/supabase
## 3.0.0-beta.24
### Patch Changes
- @trigger.dev/integration-kit@3.0.0-beta.24
- @trigger.dev/sdk@3.0.0-beta.24
## 3.0.0-beta.23
### Patch Changes
- @trigger.dev/integration-kit@3.0.0-beta.23
- @trigger.dev/sdk@3.0.0-beta.23
## 3.0.0-beta.22
### Patch Changes
+3 -3
View File
@@ -1,6 +1,6 @@
{
"name": "@trigger.dev/supabase",
"version": "3.0.0-beta.22",
"version": "3.0.0-beta.24",
"description": "Trigger.dev integration for @supabase/supabase-js",
"main": "./dist/index.js",
"types": "./dist/index.d.ts",
@@ -26,8 +26,8 @@
},
"dependencies": {
"@supabase/supabase-js": "^2.26.0",
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.22",
"@trigger.dev/sdk": "workspace:^3.0.0-beta.22",
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.24",
"@trigger.dev/sdk": "workspace:^3.0.0-beta.24",
"supabase-management-js": "^1.0.0",
"zod": "3.22.3"
},
+14
View File
@@ -1,5 +1,19 @@
# @trigger.dev/typeform
## 3.0.0-beta.24
### Patch Changes
- @trigger.dev/integration-kit@3.0.0-beta.24
- @trigger.dev/sdk@3.0.0-beta.24
## 3.0.0-beta.23
### Patch Changes
- @trigger.dev/integration-kit@3.0.0-beta.23
- @trigger.dev/sdk@3.0.0-beta.23
## 3.0.0-beta.22
### Patch Changes
+3 -3
View File
@@ -1,6 +1,6 @@
{
"name": "@trigger.dev/typeform",
"version": "3.0.0-beta.22",
"version": "3.0.0-beta.24",
"description": "The official Typeform integration for Trigger.dev",
"main": "./dist/index.js",
"types": "./dist/index.d.ts",
@@ -24,8 +24,8 @@
"typecheck": "tsc --noEmit"
},
"dependencies": {
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.22",
"@trigger.dev/sdk": "workspace:^3.0.0-beta.22",
"@trigger.dev/integration-kit": "workspace:^3.0.0-beta.24",
"@trigger.dev/sdk": "workspace:^3.0.0-beta.24",
"@typeform/api-client": "^1.8.0",
"zod": "3.22.3"
},
+2 -1
View File
@@ -14,6 +14,7 @@
"db:migrate": "turbo run db:migrate:deploy generate",
"db:seed": "turbo run db:seed",
"db:studio": "turbo run db:studio",
"db:populate": "turbo run db:populate",
"dev": "turbo run dev --parallel",
"i:dev": "infisical run -- turbo run dev --parallel",
"format": "prettier . --write --config prettier.config.js",
@@ -72,4 +73,4 @@
"engine.io-parser@5.2.2": "patches/engine.io-parser@5.2.2.patch"
}
}
}
}
+12
View File
@@ -1,5 +1,17 @@
# @trigger.dev/astro
## 3.0.0-beta.24
### Patch Changes
- @trigger.dev/sdk@3.0.0-beta.24
## 3.0.0-beta.23
### Patch Changes
- @trigger.dev/sdk@3.0.0-beta.23
## 3.0.0-beta.22
### Patch Changes
+2 -2
View File
@@ -1,7 +1,7 @@
{
"name": "@trigger.dev/astro",
"description": "An Astro-native integration for Trigger.dev background jobs platform",
"version": "3.0.0-beta.22",
"version": "3.0.0-beta.24",
"main": "./dist/index.js",
"types": "./dist/index.d.ts",
"files": [
@@ -20,7 +20,7 @@
"build:tsup": "tsup"
},
"peerDependencies": {
"@trigger.dev/sdk": "workspace:^3.0.0-beta.22"
"@trigger.dev/sdk": "workspace:^3.0.0-beta.24"
},
"devDependencies": {
"astro": "^3.0.12",
+15
View File
@@ -1,5 +1,20 @@
# trigger.dev
## 3.0.0-beta.24
### Patch Changes
- 83dc87155: Fix issues with consecutive waits
- Updated dependencies [83dc87155]
- @trigger.dev/core@3.0.0-beta.24
## 3.0.0-beta.23
### Patch Changes
- 43bc7ed94: Hoist uncaughtException handler to the top of workers to better report error messages
- @trigger.dev/core@3.0.0-beta.23
## 3.0.0-beta.22
### Patch Changes
+2 -2
View File
@@ -1,6 +1,6 @@
{
"name": "trigger.dev",
"version": "3.0.0-beta.22",
"version": "3.0.0-beta.24",
"description": "A Command-Line Interface for Trigger.dev (v3) projects",
"main": "./dist/index.js",
"types": "./dist/index.d.ts",
@@ -86,7 +86,7 @@
"@opentelemetry/sdk-trace-base": "^1.22.0",
"@opentelemetry/sdk-trace-node": "^1.22.0",
"@opentelemetry/semantic-conventions": "^1.22.0",
"@trigger.dev/core": "workspace:3.0.0-beta.22",
"@trigger.dev/core": "workspace:3.0.0-beta.24",
"@types/degit": "^2.8.3",
"chalk": "^5.2.0",
"chokidar": "^3.5.3",
+10
View File
@@ -202,6 +202,13 @@ async function _deployCommand(dir: string, options: DeployCommandOptions) {
projectRef: options.projectRef,
});
if (resolvedConfig.status === "error") {
logger.error("Failed to read config:", resolvedConfig.error);
span && recordSpanException(span, resolvedConfig.error);
throw new SkipLoggingError("Failed to read config");
}
logger.debug("Resolved config", { resolvedConfig });
span?.setAttributes({
@@ -1126,6 +1133,9 @@ async function compileProject(
format: "cjs", // This is needed to support opentelemetry instrumentation that uses module patching
target: ["node18", "es2020"],
outdir: "out",
banner: {
js: `process.on("uncaughtException", function(error, origin) { if (error instanceof Error) { process.send && process.send({ type: "EVENT", message: { type: "UNCAUGHT_EXCEPTION", payload: { error: { name: error.name, message: error.message, stack: error.stack }, origin }, version: "v1" } }); } else { process.send && process.send({ type: "EVENT", message: { type: "UNCAUGHT_EXCEPTION", payload: { error: { name: "Error", message: typeof error === "string" ? error : JSON.stringify(error) }, origin }, version: "v1" } }); } });`,
},
define: {
TRIGGER_API_URL: `"${config.triggerUrl}"`,
__PROJECT_CONFIG__: JSON.stringify(config),
+13 -5
View File
@@ -153,6 +153,11 @@ async function startDev(
logger.debug("Initial config", { config });
if (config.status === "error") {
logger.error("Failed to read config", config.error);
process.exit(1);
}
async function getDevReactElement(
configParam: ResolvedConfig,
authorization: { apiUrl: string; accessToken: string },
@@ -164,18 +169,18 @@ async function startDev(
apiClient = new CliApiClient(apiUrl, accessToken);
const devEnv = await apiClient.getProjectEnv({
projectRef: config.config.project,
projectRef: configParam.project,
env: "dev",
});
if (!devEnv.success) {
if (devEnv.error === "Project not found") {
logger.error(
`Project not found: ${config.config.project}. Ensure you are using the correct project ref and CLI profile (use --profile). Currently using the "${options.profile}" profile, which points to ${authorization.apiUrl}`
`Project not found: ${configParam.project}. Ensure you are using the correct project ref and CLI profile (use --profile). Currently using the "${options.profile}" profile, which points to ${authorization.apiUrl}`
);
} else {
logger.error(
`Failed to initialize dev environment: ${devEnv.error}. Using project ref ${config.config.project}`
`Failed to initialize dev environment: ${devEnv.error}. Using project ref ${configParam.project}`
);
}
@@ -388,6 +393,9 @@ function useDev({
resolveDir: process.cwd(),
sourcefile: "__entryPoint.ts",
},
banner: {
js: `process.on("uncaughtException", function(error, origin) { if (error instanceof Error) { process.send && process.send({ type: "UNCAUGHT_EXCEPTION", payload: { error: { name: error.name, message: error.message, stack: error.stack }, origin }, version: "v1" }); } else { process.send && process.send({ type: "UNCAUGHT_EXCEPTION", payload: { error: { name: "Error", message: typeof error === "string" ? error : JSON.stringify(error) }, origin }, version: "v1" }); } });`,
},
bundle: true,
metafile: true,
write: false,
@@ -602,10 +610,10 @@ function useDev({
} else {
}
if (e.originalError.stack) {
if (e.originalError.message || e.originalError.stack) {
logger.log(
`${chalkError("X Error:")} Worker failed to start`,
e.originalError.stack
e.originalError.stack ?? e.originalError.message
);
}
+24 -13
View File
@@ -126,6 +126,10 @@ export type ReadConfigResult =
| {
status: "in-memory";
config: ResolvedConfig;
}
| {
status: "error";
error: unknown;
};
export async function readConfig(
@@ -182,22 +186,29 @@ export async function readConfig(
],
});
// import the config file
const userConfigModule = await import(builtConfigFileHref);
try {
// import the config file
const userConfigModule = await import(builtConfigFileHref);
// The --project-ref CLI arg will always override the project specified in the config file
const rawConfig = await normalizeConfig(
userConfigModule?.config,
options?.projectRef ? { project: options?.projectRef } : undefined
);
// The --project-ref CLI arg will always override the project specified in the config file
const rawConfig = await normalizeConfig(
userConfigModule?.config,
options?.projectRef ? { project: options?.projectRef } : undefined
);
const config = Config.parse(rawConfig);
const config = Config.parse(rawConfig);
return {
status: "file",
config: await resolveConfig(absoluteDir, config),
path: configPath,
};
return {
status: "file",
config: await resolveConfig(absoluteDir, config),
path: configPath,
};
} catch (error) {
return {
status: "error",
error,
};
}
}
export async function resolveConfig(path: string, config: Config): Promise<ResolvedConfig> {
@@ -33,19 +33,4 @@ export const sender = new ZodMessageSender({
},
});
process.on("uncaughtException", (error, origin) => {
sender
.send("UNCAUGHT_EXCEPTION", {
error: {
name: error.name,
message: error.message,
stack: error.stack,
},
origin,
})
.catch((err) => {
console.error("Failed to send UNCAUGHT_EXCEPTION message", err);
});
});
taskCatalog.setGlobalTaskCatalog(new StandardTaskCatalog());
@@ -13,6 +13,7 @@ import {
TaskRunExecution,
TaskRunExecutionPayload,
TaskRunExecutionResult,
WaitReason,
correctErrorStackTrace,
} from "@trigger.dev/core/v3";
import { ZodIpcConnection } from "@trigger.dev/core/v3/zodIpc";
@@ -68,8 +69,10 @@ export class ProdBackgroundWorker {
> = new Evt();
public preCheckpointNotification = Evt.create<{ willCheckpointAndRestore: boolean }>();
public checkpointCanceledNotification = Evt.create<{ checkpointCanceled: boolean }>();
public onReadyForCheckpoint = Evt.create<{ version?: "v1" }>();
public onCancelCheckpoint = Evt.create<{ version?: "v1" }>();
public onCancelCheckpoint = Evt.create<{ version?: "v1" | "v2"; reason?: WaitReason }>();
private _onClose: Evt<void> = new Evt();
@@ -251,6 +254,9 @@ export class ProdBackgroundWorker {
this.preCheckpointNotification.attach((message) => {
taskRunProcess.preCheckpointNotification.post(message);
});
this.checkpointCanceledNotification.attach((message) => {
taskRunProcess.checkpointCanceledNotification.post(message);
});
await taskRunProcess.initialize();
@@ -377,8 +383,10 @@ class TaskRunProcess {
> = new Evt();
public preCheckpointNotification = Evt.create<{ willCheckpointAndRestore: boolean }>();
public checkpointCanceledNotification = Evt.create<{ checkpointCanceled: boolean }>();
public onReadyForCheckpoint = Evt.create<{ version?: "v1" }>();
public onCancelCheckpoint = Evt.create<{ version?: "v1" }>();
public onCancelCheckpoint = Evt.create<{ version?: "v1" | "v2"; reason?: WaitReason }>();
constructor(
private execution: ProdTaskRunExecution,
@@ -434,26 +442,62 @@ class TaskRunProcess {
this.onTaskHeartbeat.post(message.id);
},
TASKS_READY: async (message) => {},
WAIT_FOR_TASK: async (message) => {
this.onWaitForTask.post(message);
},
WAIT_FOR_BATCH: async (message) => {
this.onWaitForBatch.post(message);
},
WAIT_FOR_DURATION: async (message) => {
// Post to coordinator
this.onWaitForDuration.post(message);
// The coordinator will let us know if a checkpoint is about to happen
// We then pass this back down to the runtime in the child process
const { willCheckpointAndRestore } = await this.preCheckpointNotification.waitFor();
try {
// ..and wait for response
const { willCheckpointAndRestore } = await this.preCheckpointNotification.waitFor(
30_000
);
return { willCheckpointAndRestore };
},
WAIT_FOR_TASK: async (message) => {
this.onWaitForTask.post(message);
return {
willCheckpointAndRestore,
};
} catch (error) {
console.error("Error while waiting for pre-checkpoint notification", error);
// Assume we won't get checkpointed
return {
willCheckpointAndRestore: false,
};
}
},
READY_FOR_CHECKPOINT: async (message) => {
this.onReadyForCheckpoint.post(message);
},
CANCEL_CHECKPOINT: async (message) => {
const version = "v2";
// Post to coordinator
this.onCancelCheckpoint.post(message);
try {
// ..and wait for response
const { checkpointCanceled } = await this.checkpointCanceledNotification.waitFor(
30_000
);
return {
version,
checkpointCanceled,
};
} catch (error) {
console.error("Error while waiting for checkpoint cancellation", error);
// Assume it's been canceled
return {
version,
checkpointCanceled: true,
};
}
},
},
});
+52 -35
View File
@@ -79,21 +79,33 @@ class ProdWorker {
});
this.#backgroundWorker.onReadyForCheckpoint.attach(async (message) => {
// Flush before checkpointing so we don't flush the same spans again after restore
await this.#backgroundWorker.flushTelemetry();
this.#coordinatorSocket.socket.emit("READY_FOR_CHECKPOINT", { version: "v1" });
});
// Currently, this is only used for duration waits. Might need adjusting for other use cases.
this.#backgroundWorker.onCancelCheckpoint.attach(async (message) => {
logger.log("onCancelCheckpoint() clearing paused state, don't wait for post start hook", {
paused: this.paused,
nextResumeAfter: this.nextResumeAfter,
waitForPostStart: this.waitForPostStart,
});
logger.log("onCancelCheckpoint", { message });
this.paused = false;
this.nextResumeAfter = undefined;
this.waitForPostStart = false;
const { checkpointCanceled } = await this.#coordinatorSocket.socket.emitWithAck(
"CANCEL_CHECKPOINT",
{
version: "v2",
reason: message.reason,
}
);
this.#coordinatorSocket.socket.emit("CANCEL_CHECKPOINT", { version: "v1" });
if (checkpointCanceled) {
if (message.reason === "WAIT_FOR_DURATION") {
// Worker will resume immediately
this.paused = false;
this.nextResumeAfter = undefined;
this.waitForPostStart = false;
}
}
this.#backgroundWorker.checkpointCanceledNotification.post({ checkpointCanceled });
});
this.#backgroundWorker.onWaitForDuration.attach(async (message) => {
@@ -216,7 +228,7 @@ class ProdWorker {
}
}
#prepareForWait(reason: WaitReason, willCheckpointAndRestore: boolean) {
async #prepareForWait(reason: WaitReason, willCheckpointAndRestore: boolean) {
logger.log(`prepare for ${reason}`, { willCheckpointAndRestore });
this.#backgroundWorker.preCheckpointNotification.post({ willCheckpointAndRestore });
@@ -225,6 +237,12 @@ class ProdWorker {
this.paused = true;
this.nextResumeAfter = reason;
this.waitForPostStart = true;
if (reason === "WAIT_FOR_TASK" || reason === "WAIT_FOR_BATCH") {
// Flush before checkpointing so we don't flush the same spans again after restore
// Duration waits do this via the "ready for checkpoint" event instead
await this.#backgroundWorker.flushTelemetry();
}
}
}
@@ -253,6 +271,7 @@ class ProdWorker {
#resumeAfterDuration() {
this.paused = false;
this.nextResumeAfter = undefined;
this.waitForPostStart = false;
this.#backgroundWorker.waitCompletedNotification();
}
@@ -267,6 +286,7 @@ class ProdWorker {
return headers;
}
// FIXME: If the the worker can't connect for a while, this runs MANY times - it should only run once
#createCoordinatorSocket(host: string) {
const extraHeaders = this.#returnValidatedExtraHeaders({
"x-machine-name": MACHINE_NAME,
@@ -342,6 +362,7 @@ class ProdWorker {
this.paused = false;
this.nextResumeAfter = undefined;
this.waitForPostStart = false;
for (let i = 0; i < message.completions.length; i++) {
const completion = message.completions[i];
@@ -428,6 +449,25 @@ class ProdWorker {
return;
}
if (this.paused) {
if (!this.nextResumeAfter) {
return;
}
if (!this.attemptFriendlyId) {
logger.error("Missing friendly ID");
return;
}
socket.emit("READY_FOR_RESUME", {
version: "v1",
attemptFriendlyId: this.attemptFriendlyId,
type: this.nextResumeAfter,
});
return;
}
if (process.env.INDEX_TASKS === "true") {
try {
const taskResources = await this.#initializeWorker();
@@ -519,30 +559,6 @@ class ProdWorker {
}
}
if (this.paused) {
if (!this.nextResumeAfter) {
return;
}
if (!this.attemptFriendlyId) {
logger.error("Missing friendly ID");
return;
}
if (this.nextResumeAfter === "WAIT_FOR_DURATION") {
this.#resumeAfterDuration();
return;
}
socket.emit("READY_FOR_RESUME", {
version: "v1",
attemptFriendlyId: this.attemptFriendlyId,
type: this.nextResumeAfter,
});
return;
}
if (this.executing) {
return;
}
@@ -587,7 +603,8 @@ class ProdWorker {
case "/status": {
return reply.json({
executing: this.executing,
pause: this.paused,
paused: this.paused,
completed: this.completed.size,
nextResumeAfter: this.nextResumeAfter,
});
}
@@ -174,7 +174,7 @@ const zodIpc = new ZodIpcConnection({
prodRuntimeManager.resumeTask(completion, execution);
},
WAIT_COMPLETED_NOTIFICATION: async () => {
prodRuntimeManager.resumeAfterRestore();
prodRuntimeManager.resumeAfterDuration();
},
CLEANUP: async ({ flush, kill }, sender) => {
if (kill) {
@@ -19,22 +19,4 @@ export const tracingSDK = new TracingSDK({
diagLogLevel: (process.env.OTEL_LOG_LEVEL as TracingDiagnosticLogLevel) ?? "none",
});
process.on("uncaughtException", (error, origin) => {
process.send?.({
type: "EVENT",
message: {
type: "UNCAUGHT_EXCEPTION",
payload: {
error: {
name: error.name,
message: error.message,
stack: error.stack,
},
origin,
},
version: "v1",
},
});
});
taskCatalog.setGlobalTaskCatalog(new StandardTaskCatalog());
+15
View File
@@ -1,5 +1,20 @@
# create-trigger
## 3.0.0-beta.24
### Patch Changes
- Updated dependencies [83dc87155]
- @trigger.dev/core@3.0.0-beta.24
- @trigger.dev/yalt@3.0.0-beta.24
## 3.0.0-beta.23
### Patch Changes
- @trigger.dev/core@3.0.0-beta.23
- @trigger.dev/yalt@3.0.0-beta.23
## 3.0.0-beta.22
### Patch Changes
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@trigger.dev/cli",
"version": "3.0.0-beta.22",
"version": "3.0.0-beta.24",
"description": "The Trigger.dev CLI",
"main": "./dist/index.js",
"types": "./dist/index.d.ts",
+4
View File
@@ -1,5 +1,9 @@
# @trigger.dev/core-apps
## 3.0.0-beta.24
## 3.0.0-beta.23
## 3.0.0-beta.22
## 3.0.0-beta.21
+1 -1
View File
@@ -1,7 +1,7 @@
{
"name": "@trigger.dev/core-apps",
"description": "Backend core code used across apps",
"version": "3.0.0-beta.22",
"version": "3.0.0-beta.24",
"private": true,
"license": "MIT",
"main": "./dist/index.js",
+4
View File
@@ -1,5 +1,9 @@
# @trigger.dev/core-backend
## 3.0.0-beta.24
## 3.0.0-beta.23
## 3.0.0-beta.22
## 3.0.0-beta.21
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@trigger.dev/core-backend",
"version": "3.0.0-beta.22",
"version": "3.0.0-beta.24",
"description": "Core code used across `@trigger.dev/sdk` and Trigger.dev server",
"license": "MIT",
"main": "./dist/index.js",
+8
View File
@@ -1,5 +1,13 @@
# internal-platform
## 3.0.0-beta.24
### Patch Changes
- 83dc87155: Fix issues with consecutive waits
## 3.0.0-beta.23
## 3.0.0-beta.22
## 3.0.0-beta.21
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@trigger.dev/core",
"version": "3.0.0-beta.22",
"version": "3.0.0-beta.24",
"description": "Core code used across the Trigger.dev SDK and platform",
"license": "MIT",
"main": "./dist/index.js",
@@ -1,4 +1,5 @@
import { clock } from "../clock-api";
import { logger } from "../logger-api";
import {
BatchTaskRunExecutionResult,
ProdChildToWorkerMessages,
@@ -7,7 +8,7 @@ import {
TaskRunExecution,
TaskRunExecutionResult,
} from "../schemas";
import { unboundedTimeout } from "../utils/timers";
import { checkpointSafeTimeout, unboundedTimeout } from "../utils/timers";
import { ZodIpcConnection } from "../zodIpc";
import { RuntimeManager } from "./manager";
@@ -23,7 +24,9 @@ export class ProdRuntimeManager implements RuntimeManager {
{ resolve: (value: BatchTaskRunExecutionResult) => void; reject: (err?: any) => void }
> = new Map();
_waitForRestore: { resolve: (value: "restore") => void; reject: (err?: any) => void } | undefined;
_waitForDuration:
| { resolve: (value: "external") => void; reject: (err?: any) => void }
| undefined;
constructor(
private ipc: ZodIpcConnection<
@@ -40,15 +43,16 @@ export class ProdRuntimeManager implements RuntimeManager {
async waitForDuration(ms: number): Promise<void> {
const now = Date.now();
const resolveAfterDuration = unboundedTimeout(ms, "duration" as const);
const internalTimeout = unboundedTimeout(ms, "internal" as const);
const checkpointSafeInternalTimeout = checkpointSafeTimeout(ms);
if (ms <= this.waitThresholdInMs) {
await resolveAfterDuration;
await internalTimeout;
return;
}
const waitForRestore = new Promise<"restore">((resolve, reject) => {
this._waitForRestore = { resolve, reject };
const externalResume = new Promise<"external">((resolve, reject) => {
this._waitForDuration = { resolve, reject };
});
const { willCheckpointAndRestore } = await this.ipc.sendWithAck("WAIT_FOR_DURATION", {
@@ -57,29 +61,54 @@ export class ProdRuntimeManager implements RuntimeManager {
});
if (!willCheckpointAndRestore) {
await resolveAfterDuration;
await internalTimeout;
return;
}
this.ipc.send("READY_FOR_CHECKPOINT", {});
// Don't wait for checkpoint beyond the requested wait duration
await Promise.race([waitForRestore, resolveAfterDuration]);
// The coordinator can then cancel any in-progress checkpoints
this.ipc.send("CANCEL_CHECKPOINT", {});
}
resumeAfterRestore(): void {
if (!this._waitForRestore) {
return;
}
// internalTimeout acts as a backup and will be accurate if the checkpoint never happens
// checkpointSafeInternalTimeout is accurate even after non-simulated restores
await Promise.race([internalTimeout, checkpointSafeInternalTimeout]);
// Resets the clock to the current time
clock.reset();
this._waitForRestore.resolve("restore");
this._waitForRestore = undefined;
// The coordinator should cancel any in-progress checkpoints
const { checkpointCanceled, version } = await this.ipc.sendWithAck("CANCEL_CHECKPOINT", {
version: "v2",
reason: "WAIT_FOR_DURATION",
});
if (checkpointCanceled) {
// There won't be a checkpoint or external resume and we've already completed our internal timeout
return;
}
// No checkpoint was canceled, so we were checkpointed. We need to wait for the external resume message.
await externalResume;
}
resumeAfterDuration(): void {
if (!this._waitForDuration) {
return;
}
process.stdout.write("pre");
process.stdout.write(JSON.stringify(clock.preciseNow()));
console.log("pre", clock.preciseNow());
// Resets the clock to the current time
clock.reset();
console.log("post", clock.preciseNow());
process.stdout.write("post");
process.stdout.write(JSON.stringify(clock.preciseNow()));
this._waitForDuration.resolve("external");
this._waitForDuration = undefined;
}
async waitUntil(date: Date): Promise<void> {
+1 -1
View File
@@ -1,6 +1,6 @@
import { z } from "zod";
import { BackgroundWorkerMetadata, ImageDetailsMetadata } from "./resources";
import { QueueOptions } from "./messages";
import { QueueOptions } from "./schemas";
export const WhoAmIResponseSchema = z.object({
userId: z.string(),
+1 -1
View File
@@ -1,5 +1,5 @@
import { z } from "zod";
import { RetryOptions } from "./messages";
import { RetryOptions } from "./schemas";
import { EventFilter } from "./eventFilter";
import { Prettify } from "../types";
+497 -176
View File
@@ -1,54 +1,15 @@
import { z } from "zod";
import { TaskRunExecution, TaskRunExecutionResult } from "./common";
export const EnvironmentType = z.enum(["PRODUCTION", "STAGING", "DEVELOPMENT", "PREVIEW"]);
export type EnvironmentType = z.infer<typeof EnvironmentType>;
export const MachineCpu = z
.union([z.literal(0.25), z.literal(0.5), z.literal(1), z.literal(2), z.literal(4)])
.default(0.5);
export type MachineCpu = z.infer<typeof MachineCpu>;
export const MachineMemory = z
.union([z.literal(0.25), z.literal(0.5), z.literal(1), z.literal(2), z.literal(4), z.literal(8)])
.default(1);
export type MachineMemory = z.infer<typeof MachineMemory>;
export const Machine = z.object({
version: z.literal("v1").default("v1"),
cpu: MachineCpu,
memory: MachineMemory,
});
export type Machine = z.infer<typeof Machine>;
export const TaskRunExecutionPayload = z.object({
execution: TaskRunExecution,
traceContext: z.record(z.unknown()),
environment: z.record(z.string()).optional(),
});
export type TaskRunExecutionPayload = z.infer<typeof TaskRunExecutionPayload>;
export const ProdTaskRunExecution = TaskRunExecution.extend({
worker: z.object({
id: z.string(),
contentHash: z.string(),
version: z.string(),
}),
});
export type ProdTaskRunExecution = z.infer<typeof ProdTaskRunExecution>;
export const ProdTaskRunExecutionPayload = z.object({
execution: ProdTaskRunExecution,
traceContext: z.record(z.unknown()),
environment: z.record(z.string()).optional(),
});
export type ProdTaskRunExecutionPayload = z.infer<typeof ProdTaskRunExecutionPayload>;
import {
EnvironmentType,
Machine,
ProdTaskRunExecution,
ProdTaskRunExecutionPayload,
TaskMetadataWithFilePath,
TaskRunExecutionPayload,
WaitReason,
} from "./schemas";
import { TaskResource } from "./resources";
export const BackgroundWorkerServerMessages = z.discriminatedUnion("type", [
z.object({
@@ -148,131 +109,6 @@ export const workerToChildMessages = {
}),
};
export const FixedWindowRateLimit = z.object({
type: z.literal("fixed-window"),
limit: z.number(),
window: z.union([
z.object({
seconds: z.number(),
}),
z.object({
minutes: z.number(),
}),
z.object({
hours: z.number(),
}),
]),
});
export const SlidingWindowRateLimit = z.object({
type: z.literal("sliding-window"),
limit: z.number(),
window: z.union([
z.object({
seconds: z.number(),
}),
z.object({
minutes: z.number(),
}),
z.object({
hours: z.number(),
}),
]),
});
export const RateLimitOptions = z.discriminatedUnion("type", [
FixedWindowRateLimit,
SlidingWindowRateLimit,
]);
export const RetryOptions = z.object({
/** The number of attempts before giving up */
maxAttempts: z.number().int().optional(),
/** The exponential factor to use when calculating the next retry time.
*
* Each subsequent retry will be calculated as `previousTimeout * factor`
*/
factor: z.number().optional(),
/** The minimum time to wait before retrying */
minTimeoutInMs: z.number().int().optional(),
/** The maximum time to wait before retrying */
maxTimeoutInMs: z.number().int().optional(),
/** Randomize the timeout between retries.
*
* This can be useful to prevent the thundering herd problem where all retries happen at the same time.
*/
randomize: z.boolean().optional(),
});
export type RetryOptions = z.infer<typeof RetryOptions>;
export type RateLimitOptions = z.infer<typeof RateLimitOptions>;
export const QueueOptions = z.object({
/** You can define a shared queue and then pass the name in to your task.
*
* @example
*
* ```ts
* const myQueue = queue({
name: "my-queue",
concurrencyLimit: 1,
});
export const task1 = task({
id: "task-1",
queue: {
name: "my-queue",
},
run: async (payload: { message: string }) => {
// ...
},
});
export const task2 = task({
id: "task-2",
queue: {
name: "my-queue",
},
run: async (payload: { message: string }) => {
// ...
},
});
* ```
*/
name: z.string().optional(),
/** An optional property that specifies the maximum number of concurrent run executions.
*
* If this property is omitted, the task can potentially use up the full concurrency of an environment. */
concurrencyLimit: z.number().int().min(0).max(1000).optional(),
/** @deprecated This feature is coming soon */
rateLimit: RateLimitOptions.optional(),
});
export type QueueOptions = z.infer<typeof QueueOptions>;
export const TaskMetadata = z.object({
id: z.string(),
packageVersion: z.string(),
queue: QueueOptions.optional(),
retry: RetryOptions.optional(),
machine: Machine.partial().optional(),
triggerSource: z.string().optional(),
});
export type TaskMetadata = z.infer<typeof TaskMetadata>;
export const TaskFileMetadata = z.object({
filePath: z.string(),
exportName: z.string(),
});
export type TaskFileMetadata = z.infer<typeof TaskFileMetadata>;
export const TaskMetadataWithFilePath = TaskMetadata.merge(TaskFileMetadata);
export type TaskMetadataWithFilePath = z.infer<typeof TaskMetadataWithFilePath>;
export const UncaughtExceptionMessage = z.object({
version: z.literal("v1").default("v1"),
error: z.object({
@@ -355,8 +191,22 @@ export const ProdChildToWorkerMessages = {
}),
},
CANCEL_CHECKPOINT: {
message: z.object({
version: z.literal("v1").default("v1"),
message: z
.discriminatedUnion("version", [
z.object({
version: z.literal("v1"),
}),
z.object({
version: z.literal("v2"),
reason: WaitReason.optional(),
}),
])
.default({ version: "v1" }),
callback: z.object({
// TODO: Figure out how best to handle callback schema parsing in zod IPC
version: z.literal("v2") /* .default("v2") */,
checkpointCanceled: z.boolean(),
reason: WaitReason.optional(),
}),
},
WAIT_FOR_DURATION: {
@@ -417,3 +267,474 @@ export const ProdWorkerToChildMessages = {
}),
},
};
export const ProviderToPlatformMessages = {
LOG: {
message: z.object({
version: z.literal("v1").default("v1"),
data: z.string(),
}),
},
LOG_WITH_ACK: {
message: z.object({
version: z.literal("v1").default("v1"),
data: z.string(),
}),
callback: z.object({
status: z.literal("ok"),
}),
},
WORKER_CRASHED: {
message: z.object({
version: z.literal("v1").default("v1"),
runId: z.string(),
reason: z.string().optional(),
exitCode: z.number().optional(),
message: z.string().optional(),
logs: z.string().optional(),
}),
},
INDEXING_FAILED: {
message: z.object({
version: z.literal("v1").default("v1"),
deploymentId: z.string(),
error: z.object({
name: z.string(),
message: z.string(),
stack: z.string().optional(),
}),
}),
},
};
export const PlatformToProviderMessages = {
HEALTH: {
message: z.object({
version: z.literal("v1").default("v1"),
}),
callback: z.object({
status: z.literal("ok"),
}),
},
INDEX: {
message: z.object({
version: z.literal("v1").default("v1"),
imageTag: z.string(),
shortCode: z.string(),
apiKey: z.string(),
apiUrl: z.string(),
// identifiers
envId: z.string(),
envType: EnvironmentType,
orgId: z.string(),
projectId: z.string(),
deploymentId: z.string(),
}),
callback: z.discriminatedUnion("success", [
z.object({
success: z.literal(false),
error: z.object({
name: z.string(),
message: z.string(),
stack: z.string().optional(),
}),
}),
z.object({
success: z.literal(true),
}),
]),
},
// TODO: this should be a shared queue message instead
RESTORE: {
message: z.object({
version: z.literal("v1").default("v1"),
type: z.enum(["DOCKER", "KUBERNETES"]),
location: z.string(),
reason: z.string().optional(),
imageRef: z.string(),
machine: Machine,
// identifiers
checkpointId: z.string(),
envId: z.string(),
envType: EnvironmentType,
orgId: z.string(),
projectId: z.string(),
runId: z.string(),
}),
},
DELETE: {
message: z.object({
version: z.literal("v1").default("v1"),
name: z.string(),
}),
callback: z.object({
message: z.string(),
}),
},
GET: {
message: z.object({
version: z.literal("v1").default("v1"),
name: z.string(),
}),
},
};
export const CoordinatorToPlatformMessages = {
LOG: {
message: z.object({
version: z.literal("v1").default("v1"),
metadata: z.any(),
text: z.string(),
}),
},
CREATE_WORKER: {
message: z.object({
version: z.literal("v1").default("v1"),
projectRef: z.string(),
envId: z.string(),
deploymentId: z.string(),
metadata: z.object({
cliPackageVersion: z.string().optional(),
contentHash: z.string(),
packageVersion: z.string(),
tasks: TaskResource.array(),
}),
}),
callback: z.discriminatedUnion("success", [
z.object({
success: z.literal(false),
}),
z.object({
success: z.literal(true),
}),
]),
},
READY_FOR_EXECUTION: {
message: z.object({
version: z.literal("v1").default("v1"),
runId: z.string(),
totalCompletions: z.number(),
}),
callback: z.discriminatedUnion("success", [
z.object({
success: z.literal(false),
}),
z.object({
success: z.literal(true),
payload: ProdTaskRunExecutionPayload,
}),
]),
},
READY_FOR_RESUME: {
message: z.object({
version: z.literal("v1").default("v1"),
attemptFriendlyId: z.string(),
type: WaitReason,
}),
},
TASK_RUN_COMPLETED: {
message: z.object({
version: z.literal("v1").default("v1"),
execution: ProdTaskRunExecution,
completion: TaskRunExecutionResult,
checkpoint: z
.object({
docker: z.boolean(),
location: z.string(),
})
.optional(),
}),
},
TASK_HEARTBEAT: {
message: z.object({
version: z.literal("v1").default("v1"),
attemptFriendlyId: z.string(),
}),
},
CHECKPOINT_CREATED: {
message: z.object({
version: z.literal("v1").default("v1"),
attemptFriendlyId: z.string(),
docker: z.boolean(),
location: z.string(),
reason: z.discriminatedUnion("type", [
z.object({
type: z.literal("WAIT_FOR_DURATION"),
ms: z.number(),
now: z.number(),
}),
z.object({
type: z.literal("WAIT_FOR_BATCH"),
batchFriendlyId: z.string(),
runFriendlyIds: z.string().array(),
}),
z.object({
type: z.literal("WAIT_FOR_TASK"),
friendlyId: z.string(),
}),
z.object({
type: z.literal("RETRYING_AFTER_FAILURE"),
attemptNumber: z.number(),
}),
]),
}),
},
INDEXING_FAILED: {
message: z.object({
version: z.literal("v1").default("v1"),
deploymentId: z.string(),
error: z.object({
name: z.string(),
message: z.string(),
stack: z.string().optional(),
}),
}),
},
};
export const PlatformToCoordinatorMessages = {
RESUME_AFTER_DEPENDENCY: {
message: z.object({
version: z.literal("v1").default("v1"),
runId: z.string(),
attemptId: z.string(),
attemptFriendlyId: z.string(),
completions: TaskRunExecutionResult.array(),
executions: TaskRunExecution.array(),
}),
},
RESUME_AFTER_DURATION: {
message: z.object({
version: z.literal("v1").default("v1"),
attemptId: z.string(),
attemptFriendlyId: z.string(),
}),
},
REQUEST_ATTEMPT_CANCELLATION: {
message: z.object({
version: z.literal("v1").default("v1"),
attemptId: z.string(),
attemptFriendlyId: z.string(),
}),
},
READY_FOR_RETRY: {
message: z.object({
version: z.literal("v1").default("v1"),
runId: z.string(),
}),
},
};
export const ClientToSharedQueueMessages = {
READY_FOR_TASKS: {
message: z.object({
version: z.literal("v1").default("v1"),
backgroundWorkerId: z.string(),
}),
},
BACKGROUND_WORKER_DEPRECATED: {
message: z.object({
version: z.literal("v1").default("v1"),
backgroundWorkerId: z.string(),
}),
},
BACKGROUND_WORKER_MESSAGE: {
message: z.object({
version: z.literal("v1").default("v1"),
backgroundWorkerId: z.string(),
data: BackgroundWorkerClientMessages,
}),
},
};
export const SharedQueueToClientMessages = {
SERVER_READY: {
message: z.object({
version: z.literal("v1").default("v1"),
id: z.string(),
}),
},
BACKGROUND_WORKER_MESSAGE: {
message: z.object({
version: z.literal("v1").default("v1"),
backgroundWorkerId: z.string(),
data: BackgroundWorkerServerMessages,
}),
},
};
export const ProdWorkerToCoordinatorMessages = {
LOG: {
message: z.object({
version: z.literal("v1").default("v1"),
text: z.string(),
}),
callback: z.void(),
},
INDEX_TASKS: {
message: z.object({
version: z.literal("v1").default("v1"),
deploymentId: z.string(),
tasks: TaskResource.array(),
packageVersion: z.string(),
}),
callback: z.discriminatedUnion("success", [
z.object({
success: z.literal(false),
}),
z.object({
success: z.literal(true),
}),
]),
},
READY_FOR_EXECUTION: {
message: z.object({
version: z.literal("v1").default("v1"),
runId: z.string(),
totalCompletions: z.number(),
}),
},
READY_FOR_RESUME: {
message: z.object({
version: z.literal("v1").default("v1"),
attemptFriendlyId: z.string(),
type: WaitReason,
}),
},
READY_FOR_CHECKPOINT: {
message: z.object({
version: z.literal("v1").default("v1"),
}),
},
CANCEL_CHECKPOINT: {
message: z
.discriminatedUnion("version", [
z.object({
version: z.literal("v1"),
}),
z.object({
version: z.literal("v2"),
reason: WaitReason.optional(),
}),
])
.default({ version: "v1" }),
callback: z.object({
version: z.literal("v2").default("v2"),
checkpointCanceled: z.boolean(),
reason: WaitReason.optional(),
}),
},
TASK_HEARTBEAT: {
message: z.object({
version: z.literal("v1").default("v1"),
attemptFriendlyId: z.string(),
}),
},
TASK_RUN_COMPLETED: {
message: z.object({
version: z.literal("v1").default("v1"),
execution: ProdTaskRunExecution,
completion: TaskRunExecutionResult,
}),
callback: z.object({
willCheckpointAndRestore: z.boolean(),
shouldExit: z.boolean(),
}),
},
WAIT_FOR_DURATION: {
message: z.object({
version: z.literal("v1").default("v1"),
ms: z.number(),
now: z.number(),
attemptFriendlyId: z.string(),
}),
callback: z.object({
willCheckpointAndRestore: z.boolean(),
}),
},
WAIT_FOR_TASK: {
message: z.object({
version: z.literal("v1").default("v1"),
friendlyId: z.string(),
// This is the attempt that is waiting
attemptFriendlyId: z.string(),
}),
callback: z.object({
willCheckpointAndRestore: z.boolean(),
}),
},
WAIT_FOR_BATCH: {
message: z.object({
version: z.literal("v1").default("v1"),
batchFriendlyId: z.string(),
runFriendlyIds: z.string().array(),
// This is the attempt that is waiting
attemptFriendlyId: z.string(),
}),
callback: z.object({
willCheckpointAndRestore: z.boolean(),
}),
},
INDEXING_FAILED: {
message: z.object({
version: z.literal("v1").default("v1"),
deploymentId: z.string(),
error: z.object({
name: z.string(),
message: z.string(),
stack: z.string().optional(),
}),
}),
},
};
export const CoordinatorToProdWorkerMessages = {
RESUME_AFTER_DEPENDENCY: {
message: z.object({
version: z.literal("v1").default("v1"),
attemptId: z.string(),
completions: TaskRunExecutionResult.array(),
executions: TaskRunExecution.array(),
}),
},
RESUME_AFTER_DURATION: {
message: z.object({
version: z.literal("v1").default("v1"),
attemptId: z.string(),
}),
},
EXECUTE_TASK_RUN: {
message: z.object({
version: z.literal("v1").default("v1"),
executionPayload: ProdTaskRunExecutionPayload,
}),
},
REQUEST_ATTEMPT_CANCELLATION: {
message: z.object({
version: z.literal("v1").default("v1"),
attemptId: z.string(),
}),
},
REQUEST_EXIT: {
message: z.object({
version: z.literal("v1").default("v1"),
}),
},
READY_FOR_RETRY: {
message: z.object({
version: z.literal("v1").default("v1"),
runId: z.string(),
}),
},
};
export const ProdWorkerSocketData = z.object({
contentHash: z.string(),
projectRef: z.string(),
envId: z.string(),
runId: z.string(),
attemptFriendlyId: z.string().optional(),
podName: z.string(),
deploymentId: z.string(),
deploymentVersion: z.string(),
});
+1 -1
View File
@@ -1,5 +1,5 @@
import { z } from "zod";
import { Machine, QueueOptions, RetryOptions } from "./messages";
import { QueueOptions, RetryOptions, Machine } from "./schemas";
export const TaskResource = z.object({
id: z.string(),
+179 -469
View File
@@ -1,16 +1,184 @@
import { z } from "zod";
import { RequireKeys } from "../types";
import { TaskRunExecution, TaskRunExecutionResult } from "./common";
import {
BackgroundWorkerClientMessages,
BackgroundWorkerServerMessages,
ProdTaskRunExecution,
ProdTaskRunExecutionPayload,
RetryOptions,
Machine,
EnvironmentType,
} from "./messages";
import { TaskResource } from "./resources";
import { TaskRunExecution } from "./common";
/*
WARNING: Never import anything from ./messages here. If it's needed in both, put it here instead.
*/
export const EnvironmentType = z.enum(["PRODUCTION", "STAGING", "DEVELOPMENT", "PREVIEW"]);
export type EnvironmentType = z.infer<typeof EnvironmentType>;
export const MachineCpu = z
.union([z.literal(0.25), z.literal(0.5), z.literal(1), z.literal(2), z.literal(4)])
.default(0.5);
export type MachineCpu = z.infer<typeof MachineCpu>;
export const MachineMemory = z
.union([z.literal(0.25), z.literal(0.5), z.literal(1), z.literal(2), z.literal(4), z.literal(8)])
.default(1);
export type MachineMemory = z.infer<typeof MachineMemory>;
export const Machine = z.object({
version: z.literal("v1").default("v1"),
cpu: MachineCpu,
memory: MachineMemory,
});
export type Machine = z.infer<typeof Machine>;
export const TaskRunExecutionPayload = z.object({
execution: TaskRunExecution,
traceContext: z.record(z.unknown()),
environment: z.record(z.string()).optional(),
});
export type TaskRunExecutionPayload = z.infer<typeof TaskRunExecutionPayload>;
export const ProdTaskRunExecution = TaskRunExecution.extend({
worker: z.object({
id: z.string(),
contentHash: z.string(),
version: z.string(),
}),
});
export type ProdTaskRunExecution = z.infer<typeof ProdTaskRunExecution>;
export const ProdTaskRunExecutionPayload = z.object({
execution: ProdTaskRunExecution,
traceContext: z.record(z.unknown()),
environment: z.record(z.string()).optional(),
});
export type ProdTaskRunExecutionPayload = z.infer<typeof ProdTaskRunExecutionPayload>;
export const FixedWindowRateLimit = z.object({
type: z.literal("fixed-window"),
limit: z.number(),
window: z.union([
z.object({
seconds: z.number(),
}),
z.object({
minutes: z.number(),
}),
z.object({
hours: z.number(),
}),
]),
});
export const SlidingWindowRateLimit = z.object({
type: z.literal("sliding-window"),
limit: z.number(),
window: z.union([
z.object({
seconds: z.number(),
}),
z.object({
minutes: z.number(),
}),
z.object({
hours: z.number(),
}),
]),
});
export const RateLimitOptions = z.discriminatedUnion("type", [
FixedWindowRateLimit,
SlidingWindowRateLimit,
]);
export type RateLimitOptions = z.infer<typeof RateLimitOptions>;
export const RetryOptions = z.object({
/** The number of attempts before giving up */
maxAttempts: z.number().int().optional(),
/** The exponential factor to use when calculating the next retry time.
*
* Each subsequent retry will be calculated as `previousTimeout * factor`
*/
factor: z.number().optional(),
/** The minimum time to wait before retrying */
minTimeoutInMs: z.number().int().optional(),
/** The maximum time to wait before retrying */
maxTimeoutInMs: z.number().int().optional(),
/** Randomize the timeout between retries.
*
* This can be useful to prevent the thundering herd problem where all retries happen at the same time.
*/
randomize: z.boolean().optional(),
});
export type RetryOptions = z.infer<typeof RetryOptions>;
export const QueueOptions = z.object({
/** You can define a shared queue and then pass the name in to your task.
*
* @example
*
* ```ts
* const myQueue = queue({
name: "my-queue",
concurrencyLimit: 1,
});
export const task1 = task({
id: "task-1",
queue: {
name: "my-queue",
},
run: async (payload: { message: string }) => {
// ...
},
});
export const task2 = task({
id: "task-2",
queue: {
name: "my-queue",
},
run: async (payload: { message: string }) => {
// ...
},
});
* ```
*/
name: z.string().optional(),
/** An optional property that specifies the maximum number of concurrent run executions.
*
* If this property is omitted, the task can potentially use up the full concurrency of an environment. */
concurrencyLimit: z.number().int().min(0).max(1000).optional(),
/** @deprecated This feature is coming soon */
rateLimit: RateLimitOptions.optional(),
});
export type QueueOptions = z.infer<typeof QueueOptions>;
export const TaskMetadata = z.object({
id: z.string(),
packageVersion: z.string(),
queue: QueueOptions.optional(),
retry: RetryOptions.optional(),
machine: Machine.partial().optional(),
triggerSource: z.string().optional(),
});
export type TaskMetadata = z.infer<typeof TaskMetadata>;
export const TaskFileMetadata = z.object({
filePath: z.string(),
exportName: z.string(),
});
export type TaskFileMetadata = z.infer<typeof TaskFileMetadata>;
export const TaskMetadataWithFilePath = TaskMetadata.merge(TaskFileMetadata);
export type TaskMetadataWithFilePath = z.infer<typeof TaskMetadataWithFilePath>;
export const PostStartCauses = z.enum(["index", "create", "restore"]);
export type PostStartCauses = z.infer<typeof PostStartCauses>;
@@ -55,461 +223,3 @@ export type ResolvedConfig = RequireKeys<
export const WaitReason = z.enum(["WAIT_FOR_DURATION", "WAIT_FOR_TASK", "WAIT_FOR_BATCH"]);
export type WaitReason = z.infer<typeof WaitReason>;
export const ProviderToPlatformMessages = {
LOG: {
message: z.object({
version: z.literal("v1").default("v1"),
data: z.string(),
}),
},
LOG_WITH_ACK: {
message: z.object({
version: z.literal("v1").default("v1"),
data: z.string(),
}),
callback: z.object({
status: z.literal("ok"),
}),
},
WORKER_CRASHED: {
message: z.object({
version: z.literal("v1").default("v1"),
runId: z.string(),
reason: z.string().optional(),
exitCode: z.number().optional(),
message: z.string().optional(),
logs: z.string().optional(),
}),
},
INDEXING_FAILED: {
message: z.object({
version: z.literal("v1").default("v1"),
deploymentId: z.string(),
error: z.object({
name: z.string(),
message: z.string(),
stack: z.string().optional(),
}),
}),
},
};
export const PlatformToProviderMessages = {
HEALTH: {
message: z.object({
version: z.literal("v1").default("v1"),
}),
callback: z.object({
status: z.literal("ok"),
}),
},
INDEX: {
message: z.object({
version: z.literal("v1").default("v1"),
imageTag: z.string(),
shortCode: z.string(),
apiKey: z.string(),
apiUrl: z.string(),
// identifiers
envId: z.string(),
envType: EnvironmentType,
orgId: z.string(),
projectId: z.string(),
deploymentId: z.string(),
}),
callback: z.discriminatedUnion("success", [
z.object({
success: z.literal(false),
error: z.object({
name: z.string(),
message: z.string(),
stack: z.string().optional(),
}),
}),
z.object({
success: z.literal(true),
}),
]),
},
// TODO: this should be a shared queue message instead
RESTORE: {
message: z.object({
version: z.literal("v1").default("v1"),
type: z.enum(["DOCKER", "KUBERNETES"]),
location: z.string(),
reason: z.string().optional(),
imageRef: z.string(),
machine: Machine,
// identifiers
checkpointId: z.string(),
envId: z.string(),
envType: EnvironmentType,
orgId: z.string(),
projectId: z.string(),
runId: z.string(),
}),
},
DELETE: {
message: z.object({
version: z.literal("v1").default("v1"),
name: z.string(),
}),
callback: z.object({
message: z.string(),
}),
},
GET: {
message: z.object({
version: z.literal("v1").default("v1"),
name: z.string(),
}),
},
};
export const CoordinatorToPlatformMessages = {
LOG: {
message: z.object({
version: z.literal("v1").default("v1"),
metadata: z.any(),
text: z.string(),
}),
},
CREATE_WORKER: {
message: z.object({
version: z.literal("v1").default("v1"),
projectRef: z.string(),
envId: z.string(),
deploymentId: z.string(),
metadata: z.object({
cliPackageVersion: z.string().optional(),
contentHash: z.string(),
packageVersion: z.string(),
tasks: TaskResource.array(),
}),
}),
callback: z.discriminatedUnion("success", [
z.object({
success: z.literal(false),
}),
z.object({
success: z.literal(true),
}),
]),
},
READY_FOR_EXECUTION: {
message: z.object({
version: z.literal("v1").default("v1"),
runId: z.string(),
totalCompletions: z.number(),
}),
callback: z.discriminatedUnion("success", [
z.object({
success: z.literal(false),
}),
z.object({
success: z.literal(true),
payload: ProdTaskRunExecutionPayload,
}),
]),
},
READY_FOR_RESUME: {
message: z.object({
version: z.literal("v1").default("v1"),
attemptFriendlyId: z.string(),
type: WaitReason,
}),
},
TASK_RUN_COMPLETED: {
message: z.object({
version: z.literal("v1").default("v1"),
execution: ProdTaskRunExecution,
completion: TaskRunExecutionResult,
checkpoint: z
.object({
docker: z.boolean(),
location: z.string(),
})
.optional(),
}),
},
TASK_HEARTBEAT: {
message: z.object({
version: z.literal("v1").default("v1"),
attemptFriendlyId: z.string(),
}),
},
CHECKPOINT_CREATED: {
message: z.object({
version: z.literal("v1").default("v1"),
attemptFriendlyId: z.string(),
docker: z.boolean(),
location: z.string(),
reason: z.discriminatedUnion("type", [
z.object({
type: z.literal("WAIT_FOR_DURATION"),
ms: z.number(),
now: z.number(),
}),
z.object({
type: z.literal("WAIT_FOR_BATCH"),
batchFriendlyId: z.string(),
runFriendlyIds: z.string().array(),
}),
z.object({
type: z.literal("WAIT_FOR_TASK"),
friendlyId: z.string(),
}),
z.object({
type: z.literal("RETRYING_AFTER_FAILURE"),
attemptNumber: z.number(),
}),
]),
}),
},
INDEXING_FAILED: {
message: z.object({
version: z.literal("v1").default("v1"),
deploymentId: z.string(),
error: z.object({
name: z.string(),
message: z.string(),
stack: z.string().optional(),
}),
}),
},
};
export const PlatformToCoordinatorMessages = {
RESUME_AFTER_DEPENDENCY: {
message: z.object({
version: z.literal("v1").default("v1"),
runId: z.string(),
attemptId: z.string(),
attemptFriendlyId: z.string(),
completions: TaskRunExecutionResult.array(),
executions: TaskRunExecution.array(),
}),
},
RESUME_AFTER_DURATION: {
message: z.object({
version: z.literal("v1").default("v1"),
attemptId: z.string(),
attemptFriendlyId: z.string(),
}),
},
REQUEST_ATTEMPT_CANCELLATION: {
message: z.object({
version: z.literal("v1").default("v1"),
attemptId: z.string(),
attemptFriendlyId: z.string(),
}),
},
READY_FOR_RETRY: {
message: z.object({
version: z.literal("v1").default("v1"),
runId: z.string(),
}),
},
};
export const ClientToSharedQueueMessages = {
READY_FOR_TASKS: {
message: z.object({
version: z.literal("v1").default("v1"),
backgroundWorkerId: z.string(),
}),
},
BACKGROUND_WORKER_DEPRECATED: {
message: z.object({
version: z.literal("v1").default("v1"),
backgroundWorkerId: z.string(),
}),
},
BACKGROUND_WORKER_MESSAGE: {
message: z.object({
version: z.literal("v1").default("v1"),
backgroundWorkerId: z.string(),
data: BackgroundWorkerClientMessages,
}),
},
};
export const SharedQueueToClientMessages = {
SERVER_READY: {
message: z.object({
version: z.literal("v1").default("v1"),
id: z.string(),
}),
},
BACKGROUND_WORKER_MESSAGE: {
message: z.object({
version: z.literal("v1").default("v1"),
backgroundWorkerId: z.string(),
data: BackgroundWorkerServerMessages,
}),
},
};
export const ProdWorkerToCoordinatorMessages = {
LOG: {
message: z.object({
version: z.literal("v1").default("v1"),
text: z.string(),
}),
callback: z.void(),
},
INDEX_TASKS: {
message: z.object({
version: z.literal("v1").default("v1"),
deploymentId: z.string(),
tasks: TaskResource.array(),
packageVersion: z.string(),
}),
callback: z.discriminatedUnion("success", [
z.object({
success: z.literal(false),
}),
z.object({
success: z.literal(true),
}),
]),
},
READY_FOR_EXECUTION: {
message: z.object({
version: z.literal("v1").default("v1"),
runId: z.string(),
totalCompletions: z.number(),
}),
},
READY_FOR_RESUME: {
message: z.object({
version: z.literal("v1").default("v1"),
attemptFriendlyId: z.string(),
type: WaitReason,
}),
},
READY_FOR_CHECKPOINT: {
message: z.object({
version: z.literal("v1").default("v1"),
}),
},
CANCEL_CHECKPOINT: {
message: z.object({
version: z.literal("v1").default("v1"),
}),
},
TASK_HEARTBEAT: {
message: z.object({
version: z.literal("v1").default("v1"),
attemptFriendlyId: z.string(),
}),
},
TASK_RUN_COMPLETED: {
message: z.object({
version: z.literal("v1").default("v1"),
execution: ProdTaskRunExecution,
completion: TaskRunExecutionResult,
}),
callback: z.object({
willCheckpointAndRestore: z.boolean(),
shouldExit: z.boolean(),
}),
},
WAIT_FOR_DURATION: {
message: z.object({
version: z.literal("v1").default("v1"),
ms: z.number(),
now: z.number(),
attemptFriendlyId: z.string(),
}),
callback: z.object({
willCheckpointAndRestore: z.boolean(),
}),
},
WAIT_FOR_TASK: {
message: z.object({
version: z.literal("v1").default("v1"),
friendlyId: z.string(),
// This is the attempt that is waiting
attemptFriendlyId: z.string(),
}),
callback: z.object({
willCheckpointAndRestore: z.boolean(),
}),
},
WAIT_FOR_BATCH: {
message: z.object({
version: z.literal("v1").default("v1"),
batchFriendlyId: z.string(),
runFriendlyIds: z.string().array(),
// This is the attempt that is waiting
attemptFriendlyId: z.string(),
}),
callback: z.object({
willCheckpointAndRestore: z.boolean(),
}),
},
INDEXING_FAILED: {
message: z.object({
version: z.literal("v1").default("v1"),
deploymentId: z.string(),
error: z.object({
name: z.string(),
message: z.string(),
stack: z.string().optional(),
}),
}),
},
};
export const CoordinatorToProdWorkerMessages = {
RESUME_AFTER_DEPENDENCY: {
message: z.object({
version: z.literal("v1").default("v1"),
attemptId: z.string(),
completions: TaskRunExecutionResult.array(),
executions: TaskRunExecution.array(),
}),
},
RESUME_AFTER_DURATION: {
message: z.object({
version: z.literal("v1").default("v1"),
attemptId: z.string(),
}),
},
EXECUTE_TASK_RUN: {
message: z.object({
version: z.literal("v1").default("v1"),
executionPayload: ProdTaskRunExecutionPayload,
}),
},
REQUEST_ATTEMPT_CANCELLATION: {
message: z.object({
version: z.literal("v1").default("v1"),
attemptId: z.string(),
}),
},
REQUEST_EXIT: {
message: z.object({
version: z.literal("v1").default("v1"),
}),
},
READY_FOR_RETRY: {
message: z.object({
version: z.literal("v1").default("v1"),
runId: z.string(),
}),
},
};
export const ProdWorkerSocketData = z.object({
contentHash: z.string(),
projectRef: z.string(),
envId: z.string(),
runId: z.string(),
attemptFriendlyId: z.string().optional(),
podName: z.string(),
deploymentId: z.string(),
deploymentVersion: z.string(),
});

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