Skip to content

Commit 303154d

Browse files
authored
Merge pull request #11 from 0xprathamesh/feat/dag_engine
update:schema for Workflow DAG
2 parents 025749e + def03c2 commit 303154d

7 files changed

Lines changed: 333 additions & 24 deletions

File tree

Lines changed: 174 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,174 @@
1+
/*
2+
Warnings:
3+
4+
- You are about to drop the column `workflowId` on the `Job` table. All the data in the column will be lost.
5+
- The `status` column on the `Workflow` table would be dropped and recreated. This will lead to data loss if there is data in the column.
6+
- A unique constraint covering the columns `[key,version]` on the table `Workflow` will be added. If there are existing duplicate values, this will fail.
7+
- Added the required column `key` to the `Workflow` table without a default value. This is not possible if the table is not empty.
8+
- Added the required column `name` to the `Workflow` table without a default value. This is not possible if the table is not empty.
9+
- Added the required column `updatedAt` to the `Workflow` table without a default value. This is not possible if the table is not empty.
10+
- Added the required column `version` to the `Workflow` table without a default value. This is not possible if the table is not empty.
11+
12+
*/
13+
-- CreateEnum
14+
CREATE TYPE "WorkflowLifecycleStatus" AS ENUM ('DRAFT', 'ACTIVE', 'ARCHIVED');
15+
16+
-- CreateEnum
17+
CREATE TYPE "WorkflowRunStatus" AS ENUM ('PENDING', 'RUNNING', 'SUCCESS', 'FAILED', 'CANCELLED');
18+
19+
-- CreateEnum
20+
CREATE TYPE "NodeRunStatus" AS ENUM ('PENDING', 'READY', 'QUEUED', 'RUNNING', 'SUCCESS', 'FAILED', 'SKIPPED', 'CANCELLED');
21+
22+
-- CreateEnum
23+
CREATE TYPE "EdgeConditionType" AS ENUM ('ON_SUCCESS', 'ON_FAILURE', 'ALWAYS');
24+
25+
-- AlterEnum
26+
ALTER TYPE "JobStatus" ADD VALUE 'CANCELLED';
27+
28+
-- DropForeignKey
29+
ALTER TABLE "Job" DROP CONSTRAINT "Job_workflowId_fkey";
30+
31+
-- DropIndex
32+
DROP INDEX "Job_workflowId_idx";
33+
34+
-- AlterTable
35+
ALTER TABLE "Job" DROP COLUMN "workflowId",
36+
ADD COLUMN "workflowRunId" TEXT;
37+
38+
-- AlterTable
39+
ALTER TABLE "Workflow" ADD COLUMN "description" TEXT,
40+
ADD COLUMN "key" TEXT NOT NULL,
41+
ADD COLUMN "name" TEXT NOT NULL,
42+
ADD COLUMN "updatedAt" TIMESTAMP(3) NOT NULL,
43+
ADD COLUMN "version" INTEGER NOT NULL,
44+
DROP COLUMN "status",
45+
ADD COLUMN "status" "WorkflowLifecycleStatus" NOT NULL DEFAULT 'ACTIVE';
46+
47+
-- CreateTable
48+
CREATE TABLE "WorkflowNode" (
49+
"id" TEXT NOT NULL,
50+
"workflowId" TEXT NOT NULL,
51+
"nodeKey" TEXT NOT NULL,
52+
"type" "JobType" NOT NULL,
53+
"payloadTemplate" JSONB NOT NULL,
54+
"maxRetries" INTEGER NOT NULL DEFAULT 3,
55+
"timeoutMs" INTEGER,
56+
"continueOnFailure" BOOLEAN NOT NULL,
57+
"createdAt" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP,
58+
"updatedAt" TIMESTAMP(3) NOT NULL,
59+
60+
CONSTRAINT "WorkflowNode_pkey" PRIMARY KEY ("id")
61+
);
62+
63+
-- CreateTable
64+
CREATE TABLE "WorkflowEdge" (
65+
"id" TEXT NOT NULL,
66+
"workflowId" TEXT NOT NULL,
67+
"fromNodeId" TEXT NOT NULL,
68+
"toNodeId" TEXT NOT NULL,
69+
"condition" "EdgeConditionType" NOT NULL DEFAULT 'ON_SUCCESS',
70+
71+
CONSTRAINT "WorkflowEdge_pkey" PRIMARY KEY ("id")
72+
);
73+
74+
-- CreateTable
75+
CREATE TABLE "WorkflowRun" (
76+
"id" TEXT NOT NULL,
77+
"workflowId" TEXT NOT NULL,
78+
"status" "WorkflowRunStatus" NOT NULL DEFAULT 'PENDING',
79+
"input" JSONB,
80+
"output" JSONB,
81+
"errorMessage" TEXT,
82+
"startedAt" TIMESTAMP(3),
83+
"completedAt" TIMESTAMP(3),
84+
"createdAt" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP,
85+
"updatedAt" TIMESTAMP(3) NOT NULL,
86+
87+
CONSTRAINT "WorkflowRun_pkey" PRIMARY KEY ("id")
88+
);
89+
90+
-- CreateTable
91+
CREATE TABLE "WorkflowNodeRun" (
92+
"id" TEXT NOT NULL,
93+
"workflowRunId" TEXT NOT NULL,
94+
"nodeId" TEXT NOT NULL,
95+
"status" "NodeRunStatus" NOT NULL DEFAULT 'PENDING',
96+
"attempt" INTEGER NOT NULL DEFAULT 0,
97+
"payload" JSONB NOT NULL,
98+
"result" JSONB,
99+
"errorMessage" TEXT,
100+
"queuedAt" TIMESTAMP(3),
101+
"startedAt" TIMESTAMP(3),
102+
"completedAt" TIMESTAMP(3),
103+
"createdAt" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP,
104+
"updatedAt" TIMESTAMP(3) NOT NULL,
105+
"jobId" TEXT,
106+
107+
CONSTRAINT "WorkflowNodeRun_pkey" PRIMARY KEY ("id")
108+
);
109+
110+
-- CreateIndex
111+
CREATE INDEX "WorkflowNode_workflowId_idx" ON "WorkflowNode"("workflowId");
112+
113+
-- CreateIndex
114+
CREATE UNIQUE INDEX "WorkflowNode_workflowId_nodeKey_key" ON "WorkflowNode"("workflowId", "nodeKey");
115+
116+
-- CreateIndex
117+
CREATE INDEX "WorkflowEdge_workflowId_fromNodeId_idx" ON "WorkflowEdge"("workflowId", "fromNodeId");
118+
119+
-- CreateIndex
120+
CREATE INDEX "WorkflowEdge_workflowId_toNodeId_idx" ON "WorkflowEdge"("workflowId", "toNodeId");
121+
122+
-- CreateIndex
123+
CREATE UNIQUE INDEX "WorkflowEdge_workflowId_fromNodeId_toNodeId_condition_key" ON "WorkflowEdge"("workflowId", "fromNodeId", "toNodeId", "condition");
124+
125+
-- CreateIndex
126+
CREATE INDEX "WorkflowRun_workflowId_status_createdAt_idx" ON "WorkflowRun"("workflowId", "status", "createdAt");
127+
128+
-- CreateIndex
129+
CREATE UNIQUE INDEX "WorkflowNodeRun_jobId_key" ON "WorkflowNodeRun"("jobId");
130+
131+
-- CreateIndex
132+
CREATE INDEX "WorkflowNodeRun_workflowRunId_status_idx" ON "WorkflowNodeRun"("workflowRunId", "status");
133+
134+
-- CreateIndex
135+
CREATE INDEX "WorkflowNodeRun_nodeId_status_idx" ON "WorkflowNodeRun"("nodeId", "status");
136+
137+
-- CreateIndex
138+
CREATE UNIQUE INDEX "WorkflowNodeRun_workflowRunId_nodeId_key" ON "WorkflowNodeRun"("workflowRunId", "nodeId");
139+
140+
-- CreateIndex
141+
CREATE INDEX "Job_workflowRunId_idx" ON "Job"("workflowRunId");
142+
143+
-- CreateIndex
144+
CREATE INDEX "Workflow_key_status_idx" ON "Workflow"("key", "status");
145+
146+
-- CreateIndex
147+
CREATE UNIQUE INDEX "Workflow_key_version_key" ON "Workflow"("key", "version");
148+
149+
-- AddForeignKey
150+
ALTER TABLE "Job" ADD CONSTRAINT "Job_workflowRunId_fkey" FOREIGN KEY ("workflowRunId") REFERENCES "WorkflowRun"("id") ON DELETE SET NULL ON UPDATE CASCADE;
151+
152+
-- AddForeignKey
153+
ALTER TABLE "WorkflowNode" ADD CONSTRAINT "WorkflowNode_workflowId_fkey" FOREIGN KEY ("workflowId") REFERENCES "Workflow"("id") ON DELETE RESTRICT ON UPDATE CASCADE;
154+
155+
-- AddForeignKey
156+
ALTER TABLE "WorkflowEdge" ADD CONSTRAINT "WorkflowEdge_workflowId_fkey" FOREIGN KEY ("workflowId") REFERENCES "Workflow"("id") ON DELETE CASCADE ON UPDATE CASCADE;
157+
158+
-- AddForeignKey
159+
ALTER TABLE "WorkflowEdge" ADD CONSTRAINT "WorkflowEdge_fromNodeId_fkey" FOREIGN KEY ("fromNodeId") REFERENCES "WorkflowNode"("id") ON DELETE CASCADE ON UPDATE CASCADE;
160+
161+
-- AddForeignKey
162+
ALTER TABLE "WorkflowEdge" ADD CONSTRAINT "WorkflowEdge_toNodeId_fkey" FOREIGN KEY ("toNodeId") REFERENCES "WorkflowNode"("id") ON DELETE CASCADE ON UPDATE CASCADE;
163+
164+
-- AddForeignKey
165+
ALTER TABLE "WorkflowRun" ADD CONSTRAINT "WorkflowRun_workflowId_fkey" FOREIGN KEY ("workflowId") REFERENCES "Workflow"("id") ON DELETE RESTRICT ON UPDATE CASCADE;
166+
167+
-- AddForeignKey
168+
ALTER TABLE "WorkflowNodeRun" ADD CONSTRAINT "WorkflowNodeRun_workflowRunId_fkey" FOREIGN KEY ("workflowRunId") REFERENCES "WorkflowRun"("id") ON DELETE CASCADE ON UPDATE CASCADE;
169+
170+
-- AddForeignKey
171+
ALTER TABLE "WorkflowNodeRun" ADD CONSTRAINT "WorkflowNodeRun_nodeId_fkey" FOREIGN KEY ("nodeId") REFERENCES "WorkflowNode"("id") ON DELETE RESTRICT ON UPDATE CASCADE;
172+
173+
-- AddForeignKey
174+
ALTER TABLE "WorkflowNodeRun" ADD CONSTRAINT "WorkflowNodeRun_jobId_fkey" FOREIGN KEY ("jobId") REFERENCES "Job"("id") ON DELETE SET NULL ON UPDATE CASCADE;

‎prisma/schema.prisma‎

Lines changed: 140 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -16,19 +16,58 @@ enum JobStatus {
1616
RUNNING
1717
SUCCESS
1818
FAILED
19+
CANCELLED
1920
}
2021

2122
enum JobType {
2223
AGENT_TASK
2324
WEBHOOK
2425
GENERIC
2526
}
27+
enum WorkflowLifecycleStatus {
28+
DRAFT
29+
ACTIVE
30+
ARCHIVED
31+
}
32+
enum WorkflowRunStatus {
33+
PENDING
34+
RUNNING
35+
SUCCESS
36+
FAILED
37+
CANCELLED
38+
}
39+
enum NodeRunStatus {
40+
PENDING
41+
READY
42+
QUEUED
43+
RUNNING
44+
SUCCESS
45+
FAILED
46+
SKIPPED
47+
CANCELLED
48+
}
2649

50+
enum EdgeConditionType {
51+
ON_SUCCESS
52+
ON_FAILURE
53+
ALWAYS
54+
}
2755
model Workflow {
28-
id String @id @default(uuid())
29-
status String
56+
id String @id @default(uuid())
57+
key String
58+
version Int
59+
name String
60+
description String?
61+
status WorkflowLifecycleStatus @default(ACTIVE)
62+
3063
createdAt DateTime @default(now())
31-
jobs Job[]
64+
updatedAt DateTime @updatedAt
65+
nodes WorkflowNode[]
66+
edges WorkflowEdge[]
67+
runs WorkflowRun[]
68+
69+
@@unique([key, version])
70+
@@index([key,status])
3271
}
3372

3473
model Job {
@@ -41,14 +80,109 @@ model Job {
4180
retries Int @default(0)
4281
maxRetries Int @default(3)
4382
queueJobId String? // BullMQ job id for correlation
44-
workflowId String?
45-
workflow Workflow? @relation(fields: [workflowId], references: [id], onDelete: SetNull)
83+
84+
//Dag run links
85+
workflowRunId String?
86+
workflowRun WorkflowRun? @relation(fields: [workflowRunId], references: [id], onDelete: SetNull)
87+
nodeRun WorkflowNodeRun? @relation("NodeRunJob")
4688
startedAt DateTime?
4789
completedAt DateTime?
4890
createdAt DateTime @default(now())
4991
updatedAt DateTime @updatedAt
5092
5193
@@index([status, type])
52-
@@index([workflowId])
94+
@@index([workflowRunId])
5395
@@index([queueJobId])
5496
}
97+
98+
model WorkflowNode {
99+
id String @id @default(uuid())
100+
workflowId String
101+
nodeKey String
102+
type JobType
103+
payloadTemplate Json
104+
maxRetries Int @default(3)
105+
timeoutMs Int?
106+
continueOnFailure Boolean
107+
108+
workflow Workflow @relation(fields: [workflowId], references: [id])
109+
outgoingEdges WorkflowEdge[] @relation("EdgeFromNode")
110+
incomingEdges WorkflowEdge[] @relation("EdgeToNode")
111+
nodeRuns WorkflowNodeRun[]
112+
113+
createdAt DateTime @default(now())
114+
updatedAt DateTime @updatedAt
115+
116+
@@unique([workflowId, nodeKey])
117+
@@index([workflowId])
118+
}
119+
120+
121+
model WorkflowEdge {
122+
id String @id @default(uuid())
123+
workflowId String
124+
fromNodeId String
125+
toNodeId String
126+
condition EdgeConditionType @default(ON_SUCCESS)
127+
128+
workflow Workflow @relation(fields: [workflowId], references: [id], onDelete: Cascade)
129+
fromNode WorkflowNode @relation("EdgeFromNode", fields: [fromNodeId], references: [id], onDelete: Cascade)
130+
toNode WorkflowNode @relation("EdgeToNode", fields: [toNodeId], references: [id], onDelete: Cascade)
131+
132+
@@unique([workflowId, fromNodeId, toNodeId, condition])
133+
@@index([workflowId, fromNodeId])
134+
@@index([workflowId, toNodeId])
135+
}
136+
137+
/// One execution instance of a workflow definition
138+
model WorkflowRun {
139+
id String @id @default(uuid())
140+
workflowId String
141+
status WorkflowRunStatus @default(PENDING)
142+
143+
input Json?
144+
output Json?
145+
errorMessage String?
146+
147+
startedAt DateTime?
148+
completedAt DateTime?
149+
createdAt DateTime @default(now())
150+
updatedAt DateTime @updatedAt
151+
152+
workflow Workflow @relation(fields: [workflowId], references: [id], onDelete: Restrict)
153+
nodeRuns WorkflowNodeRun[]
154+
jobs Job[]
155+
156+
@@index([workflowId, status, createdAt])
157+
}
158+
159+
/// Runtime state for each node within a workflow run
160+
model WorkflowNodeRun {
161+
id String @id @default(uuid())
162+
workflowRunId String
163+
nodeId String
164+
status NodeRunStatus @default(PENDING)
165+
166+
attempt Int @default(0)
167+
payload Json // resolved payload for this run
168+
result Json?
169+
errorMessage String?
170+
171+
queuedAt DateTime?
172+
startedAt DateTime?
173+
completedAt DateTime?
174+
createdAt DateTime @default(now())
175+
updatedAt DateTime @updatedAt
176+
177+
workflowRun WorkflowRun @relation(fields: [workflowRunId], references: [id], onDelete: Cascade)
178+
node WorkflowNode @relation(fields: [nodeId], references: [id], onDelete: Restrict)
179+
180+
// optional back-reference to actual queue job row
181+
jobId String? @unique
182+
job Job? @relation("NodeRunJob", fields: [jobId], references: [id], onDelete: SetNull)
183+
184+
@@unique([workflowRunId, nodeId])
185+
@@index([workflowRunId, status])
186+
@@index([nodeId, status])
187+
}
188+

‎src/modules/job/job.controller.ts‎

Lines changed: 7 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -24,8 +24,8 @@ function parseCreateBody(body: unknown): CreateJobInput | { error: string } {
2424
if (typeof payload !== "object" || payload === null) {
2525
return { error: "payload must be an object" };
2626
}
27-
const workflowId =
28-
typeof o.workflowId === "string" ? o.workflowId : undefined;
27+
const workflowRunId =
28+
typeof o.workflowRunId === "string" ? o.workflowRunId : undefined;
2929
const maxRetries =
3030
typeof o.maxRetries === "number" &&
3131
Number.isInteger(o.maxRetries) &&
@@ -35,7 +35,7 @@ function parseCreateBody(body: unknown): CreateJobInput | { error: string } {
3535
return {
3636
type: type as CreateJobInput["type"],
3737
payload: payload as CreateJobInput["payload"],
38-
workflowId,
38+
workflowRunId,
3939
maxRetries,
4040
};
4141
}
@@ -98,7 +98,7 @@ const jobController = {
9898
const {
9999
status,
100100
type,
101-
workflowId,
101+
workflowRunId,
102102
queueJobId,
103103
startedAt,
104104
completedAt,
@@ -109,7 +109,7 @@ const jobController = {
109109
const filters: {
110110
status?: JobStatus;
111111
type?: JobType;
112-
workflowId?: string;
112+
workflowRunId?: string;
113113
queueJobId?: string;
114114
startedAt?: Date;
115115
completedAt?: Date;
@@ -138,8 +138,8 @@ const jobController = {
138138
if (type) {
139139
filters.type = type as JobType;
140140
}
141-
if (workflowId) {
142-
filters.workflowId = workflowId as string;
141+
if (workflowRunId) {
142+
filters.workflowRunId = workflowRunId as string;
143143
}
144144
if (queueJobId) {
145145
filters.queueJobId = queueJobId as string;

0 commit comments

Comments
 (0)