@eventual/core
Advanced tools
Comparing version 0.7.8 to 0.7.9
@@ -46,8 +46,4 @@ "use strict"; | ||
} | ||
const workflow = (0, workflow_js_1.lookupWorkflow)(workflowName); | ||
if (workflow === undefined) { | ||
throw new Error(`no such workflow with name '${workflowName}'`); | ||
} | ||
// TODO: get workflow from execution id | ||
return orchestrateExecution(workflow, executionId, records, baseTime); | ||
return orchestrateExecution(workflowName, executionId, records, baseTime); | ||
}); | ||
@@ -66,5 +62,5 @@ logger.debug("Executions succeeded: " + | ||
}; | ||
async function orchestrateExecution(workflow, executionId, events, baseTime) { | ||
async function orchestrateExecution(workflowName, executionId, events, baseTime) { | ||
const executionLogger = logger.createChild({ | ||
persistentLogAttributes: { workflowName: workflow.name, executionId }, | ||
persistentLogAttributes: { workflowName, executionId }, | ||
}); | ||
@@ -104,2 +100,11 @@ const metrics = initializeMetrics(); | ||
}, start); | ||
const workflow = (0, workflow_js_1.lookupWorkflow)(workflowName); | ||
if (workflow === undefined) { | ||
yield (0, workflow_events_js_1.createEvent)({ | ||
type: workflow_events_js_1.WorkflowEventType.WorkflowFailed, | ||
error: "WorkflowNotFound", | ||
message: `Workflow name ${workflowName} does not exist.`, | ||
}, start); | ||
return; | ||
} | ||
const workflowContext = { | ||
@@ -143,3 +148,3 @@ name: workflow.workflowName, | ||
executionLogger.info(`Found ${newCommands.length} new commands.`); | ||
yield* await (0, utils_js_1.timed)(metrics, constants_js_1.OrchestratorMetrics.InvokeCommandsDuration, () => processCommands(newCommands)); | ||
yield* await (0, utils_js_1.timed)(metrics, constants_js_1.OrchestratorMetrics.InvokeCommandsDuration, () => processCommands(workflow, newCommands)); | ||
metrics.putMetric(constants_js_1.OrchestratorMetrics.CommandsInvoked, newCommands.length, unit_js_1.Unit.Count); | ||
@@ -269,3 +274,3 @@ // tracks the time it takes for a workflow task to be scheduled until new commands could be emitted. | ||
*/ | ||
async function processCommands(commands) { | ||
async function processCommands(workflow, commands) { | ||
console.debug("Commands to send", JSON.stringify(commands)); | ||
@@ -281,3 +286,3 @@ // register command events | ||
metrics.setDimensions({ | ||
[constants_js_1.MetricsCommon.WorkflowNameDimension]: workflow.workflowName, | ||
[constants_js_1.MetricsCommon.WorkflowNameDimension]: workflowName, | ||
}); | ||
@@ -314,2 +319,2 @@ // number of events that came from the workflow task | ||
} | ||
//# sourceMappingURL=data:application/json;base64, | ||
//# sourceMappingURL=data:application/json;base64, |
@@ -43,8 +43,4 @@ import { inspect } from "util"; | ||
} | ||
const workflow = lookupWorkflow(workflowName); | ||
if (workflow === undefined) { | ||
throw new Error(`no such workflow with name '${workflowName}'`); | ||
} | ||
// TODO: get workflow from execution id | ||
return orchestrateExecution(workflow, executionId, records, baseTime); | ||
return orchestrateExecution(workflowName, executionId, records, baseTime); | ||
}); | ||
@@ -63,5 +59,5 @@ logger.debug("Executions succeeded: " + | ||
}; | ||
async function orchestrateExecution(workflow, executionId, events, baseTime) { | ||
async function orchestrateExecution(workflowName, executionId, events, baseTime) { | ||
const executionLogger = logger.createChild({ | ||
persistentLogAttributes: { workflowName: workflow.name, executionId }, | ||
persistentLogAttributes: { workflowName, executionId }, | ||
}); | ||
@@ -101,2 +97,11 @@ const metrics = initializeMetrics(); | ||
}, start); | ||
const workflow = lookupWorkflow(workflowName); | ||
if (workflow === undefined) { | ||
yield createEvent({ | ||
type: WorkflowEventType.WorkflowFailed, | ||
error: "WorkflowNotFound", | ||
message: `Workflow name ${workflowName} does not exist.`, | ||
}, start); | ||
return; | ||
} | ||
const workflowContext = { | ||
@@ -140,3 +145,3 @@ name: workflow.workflowName, | ||
executionLogger.info(`Found ${newCommands.length} new commands.`); | ||
yield* await timed(metrics, OrchestratorMetrics.InvokeCommandsDuration, () => processCommands(newCommands)); | ||
yield* await timed(metrics, OrchestratorMetrics.InvokeCommandsDuration, () => processCommands(workflow, newCommands)); | ||
metrics.putMetric(OrchestratorMetrics.CommandsInvoked, newCommands.length, Unit.Count); | ||
@@ -266,3 +271,3 @@ // tracks the time it takes for a workflow task to be scheduled until new commands could be emitted. | ||
*/ | ||
async function processCommands(commands) { | ||
async function processCommands(workflow, commands) { | ||
console.debug("Commands to send", JSON.stringify(commands)); | ||
@@ -278,3 +283,3 @@ // register command events | ||
metrics.setDimensions({ | ||
[MetricsCommon.WorkflowNameDimension]: workflow.workflowName, | ||
[MetricsCommon.WorkflowNameDimension]: workflowName, | ||
}); | ||
@@ -310,2 +315,2 @@ // number of events that came from the workflow task | ||
} | ||
//# sourceMappingURL=data:application/json;base64, | ||
//# sourceMappingURL=data:application/json;base64, |
{ | ||
"name": "@eventual/core", | ||
"version": "0.7.8", | ||
"version": "0.7.9", | ||
"exports": { | ||
@@ -5,0 +5,0 @@ ".": { |
Sorry, the diff of this file is not supported yet
Sorry, the diff of this file is not supported yet
Sorry, the diff of this file is not supported yet
Sorry, the diff of this file is not supported yet
1481408
11725