Skip to content

Commit 0b5a2db

Browse files
authored
Merge pull request #989 from polywrap/nk/improve-workflow
feat(cli): add job status in job result object
2 parents f10af9f + fb1a828 commit 0b5a2db

7 files changed

Lines changed: 92 additions & 38 deletions

File tree

packages/cli/lang/en.json

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -182,6 +182,7 @@
182182
"commands_run_error_noApi": "API needs to be initialized",
183183
"commands_run_error_readFail": "Failed to read query {query}",
184184
"commands_run_error_unsupportedOutputFileExt": "Unsupported outputFile extention: ${outputFileExt}",
185+
"commands_run_error_cueDoesNotExist": "Require cue to run validator, checkout https://cuelang.org/ for more information",
185186
"commands_run_error_noWorkflowScriptFound": "Workflow script not found at path: {path}",
186187
"commands_run_error_noTestEnvFound": "polywrap test-env not found, please run 'polywrap infra up --modules=eth-ens-ipfs'",
187188
"commands_testEnv_description": "Manage a test environment for Polywrap",

packages/cli/lang/es.json

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -182,6 +182,7 @@
182182
"commands_run_error_noApi": "API needs to be initialized",
183183
"commands_run_error_readFail": "Failed to read query {query}",
184184
"commands_run_error_unsupportedOutputFileExt": "Unsupported outputFile extention: ${outputFileExt}",
185+
"commands_run_error_cueDoesNotExist": "Require cue to run validator, checkout https://cuelang.org/ for more information",
185186
"commands_run_error_noWorkflowScriptFound": "Workflow script not found at path: {path}",
186187
"commands_run_error_noTestEnvFound": "polywrap test-env not found, please run 'polywrap infra up --modules=eth-ens-ipfs'",
187188
"commands_testEnv_description": "Manage a test environment for Polywrap",

packages/cli/src/commands/run.ts

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -7,7 +7,7 @@ import {
77
validateOutput,
88
} from "../lib";
99

10-
import { InvokeResult, Workflow } from "@polywrap/core-js";
10+
import { InvokeResult, Workflow, JobResult } from "@polywrap/core-js";
1111
import { PolywrapClient, PolywrapClientConfig } from "@polywrap/client-js";
1212
import path from "path";
1313
import yaml from "js-yaml";
@@ -85,7 +85,9 @@ const _run = async (workflowPath: string, options: WorkflowCommandOptions) => {
8585
workflow,
8686
config: clientConfig,
8787
ids: jobs,
88-
onExecution: async (id: string, data: unknown, error: Error) => {
88+
onExecution: async (id: string, jobResult: JobResult) => {
89+
const { data, error } = jobResult;
90+
8991
if (!quiet) {
9092
console.log("-----------------------------------");
9193
console.log(`ID: ${id}`);

packages/cli/src/lib/helpers/workflow-validator.ts

Lines changed: 13 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,5 @@
11
import { runCommand } from "../system";
2+
import { intlMsg } from "../intl";
23

34
import fs from "fs";
45
import { InvokeResult } from "@polywrap/core-js";
@@ -19,6 +20,10 @@ export async function validateOutput(
1920
result: InvokeResult,
2021
validateScriptPath: string
2122
): Promise<void> {
23+
if (!(await cueExists())) {
24+
console.warn(intlMsg.commands_run_error_cueDoesNotExist());
25+
}
26+
2227
const index = id.lastIndexOf(".");
2328
const jobId = id.substring(0, index);
2429
const stepId = id.substring(index + 1);
@@ -29,12 +34,17 @@ export async function validateOutput(
2934
await fs.promises.writeFile(jsonOutput, JSON.stringify(result, null, 2));
3035

3136
try {
32-
await runCommand(
33-
`cue vet -d ${selector} ${validateScriptPath} ${jsonOutput}`
37+
const { stderr } = await runCommand(
38+
`cue vet -d ${selector} ${validateScriptPath} ${jsonOutput}`,
39+
true
3440
);
41+
42+
if (stderr) {
43+
console.error(stderr);
44+
console.log("-----------------------------------");
45+
}
3546
} catch (e) {
3647
console.error(e);
37-
console.log(e);
3848
console.log("-----------------------------------");
3949
process.exitCode = 1;
4050
}

packages/js/client/src/__tests__/e2e/workflow.spec.ts

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,7 @@ import {
66
ensAddresses,
77
providers
88
} from "@polywrap/test-env-js";
9-
import { createPolywrapClient, PolywrapClient, PolywrapClientConfig } from "../..";
9+
import { createPolywrapClient, JobResult, PolywrapClient, PolywrapClientConfig } from "../..";
1010
import { outPropWorkflow, sanityWorkflow } from "./workflow-test-cases";
1111

1212
jest.setTimeout(200000);
@@ -89,17 +89,17 @@ describe("workflow", () => {
8989
test("sanity workflow", async () => {
9090
await client.run({
9191
workflow: sanityWorkflow,
92-
onExecution: async (id: string, data: unknown, error: unknown) => {
93-
await tests[id](data, error);
92+
onExecution: async (id: string, jobResult: JobResult) => {
93+
await tests[id](jobResult.data, jobResult.error);
9494
},
9595
});
9696
});
9797

9898
test("workflow with output propagation", async () => {
9999
await client.run({
100100
workflow: outPropWorkflow,
101-
onExecution: async (id: string, data: unknown, error: unknown) => {
102-
await tests[id](data, error);
101+
onExecution: async (id: string, jobResult: JobResult) => {
102+
await tests[id](jobResult.data, jobResult.error);
103103
},
104104
});
105105
});

packages/js/core/src/types/Workflow.ts

Lines changed: 12 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,17 @@ export type Workflow<TUri extends Uri | string = string> = {
2323
jobs: Job<TUri>;
2424
};
2525

26+
export enum JobStatus {
27+
SUCCEED,
28+
FAILED,
29+
SKIPPED,
30+
}
31+
32+
export interface JobResult<TData extends unknown = unknown>
33+
extends InvokeResult<TData> {
34+
status: JobStatus;
35+
}
36+
2637
export interface RunOptions<
2738
TData extends Record<string, unknown> = Record<string, unknown>,
2839
TUri extends Uri | string = string
@@ -32,11 +43,7 @@ export interface RunOptions<
3243
contextId?: string;
3344
ids?: string[];
3445

35-
onExecution?(
36-
id: string,
37-
data?: InvokeResult<TData>["data"],
38-
error?: InvokeResult<TData>["error"]
39-
): MaybeAsync<void>;
46+
onExecution?(id: string, jobResult: JobResult<TData>): MaybeAsync<void>;
4047
}
4148

4249
export interface WorkflowHandler {

packages/js/core/src/workflow/JobRunner.ts

Lines changed: 56 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -1,8 +1,9 @@
11
import {
22
Client,
33
executeMaybeAsyncFunction,
4-
InvokeResult,
54
Job,
5+
JobResult,
6+
JobStatus,
67
MaybeAsync,
78
Uri,
89
} from "../types";
@@ -19,14 +20,13 @@ export class JobRunner<
1920
TData extends unknown = unknown,
2021
TUri extends Uri | string = string
2122
> {
22-
private jobOutput: Map<string, InvokeResult<TData>>;
23+
private jobOutput: Map<string, JobResult<TData>>;
2324

2425
constructor(
2526
private client: Client,
2627
private onExecution?: (
2728
id: string,
28-
data?: InvokeResult<TData>["data"],
29-
error?: InvokeResult<TData>["error"]
29+
JobResult: JobResult<TData>
3030
) => MaybeAsync<void>
3131
) {
3232
this.jobOutput = new Map();
@@ -45,27 +45,47 @@ export class JobRunner<
4545
const steps = jobs[jobId].steps;
4646
if (steps) {
4747
for (let i = 0; i < steps.length; i++) {
48+
let result: JobResult<TData> | undefined;
49+
let args: Record<string, unknown> | undefined;
50+
4851
const step = steps[i];
4952
const absoluteId = parentId
5053
? `${parentId}.${jobId}.${i}`
5154
: `${jobId}.${i}`;
52-
const args = this.resolveArgs(absoluteId, step.args);
53-
const result = await this.client.invoke<TData, TUri>({
54-
uri: step.uri,
55-
method: step.method,
56-
config: step.config,
57-
args: args,
58-
});
59-
60-
this.jobOutput.set(absoluteId, result);
61-
62-
if (this.onExecution && typeof this.onExecution === "function") {
63-
await executeMaybeAsyncFunction(
64-
this.onExecution,
65-
absoluteId,
66-
result.data,
67-
result.error
68-
);
55+
try {
56+
args = this.resolveArgs(absoluteId, step.args);
57+
} catch (e) {
58+
result = {
59+
error: e,
60+
status: JobStatus.SKIPPED,
61+
};
62+
}
63+
64+
if (args) {
65+
const invokeResult = await this.client.invoke<TData, TUri>({
66+
uri: step.uri,
67+
method: step.method,
68+
config: step.config,
69+
args: args,
70+
});
71+
72+
if (invokeResult.error) {
73+
result = { ...invokeResult, status: JobStatus.FAILED };
74+
} else {
75+
result = { ...invokeResult, status: JobStatus.SUCCEED };
76+
}
77+
}
78+
79+
if (result) {
80+
this.jobOutput.set(absoluteId, result);
81+
82+
if (this.onExecution && typeof this.onExecution === "function") {
83+
await executeMaybeAsyncFunction(
84+
this.onExecution,
85+
absoluteId,
86+
result
87+
);
88+
}
6989
}
7090
}
7191
}
@@ -124,8 +144,21 @@ export class JobRunner<
124144
}
125145
}
126146
const output = outputs.get(absStepId);
127-
if (output && output[dataOrErr]) {
128-
return output[dataOrErr];
147+
if (
148+
output &&
149+
dataOrErr === "data" &&
150+
output.status === JobStatus.SUCCEED &&
151+
output.data
152+
) {
153+
return output.data;
154+
}
155+
if (
156+
output &&
157+
dataOrErr === "error" &&
158+
output.status === JobStatus.FAILED &&
159+
output.error
160+
) {
161+
return output.error;
129162
}
130163
}
131164

0 commit comments

Comments
 (0)