Development guide for the lemline-core module...
Guide development of the lemline-core module - the pure, stateless workflow execution engine implementing the Serverless Workflow DSL v1.0 specification.
Documentation:
Orchestration:
First, read core-orchestrators.md
Nodes:
First, read core-nodes.md
Processors:
First, read core-processors.md
processors/DoProcessor.kt, ForProcessor.kt, SwitchProcessor.ktprocessors/WaitProcessor.kt, CallProcessor.kt, RunProcessor.ktStates:
First, read core-states.md
states/ directory, extend TaskStatescope property in the relevant state classTaskStates.ktError Handling:
First, read core-errors.md
Expressions:
First, read core-expressions.md
scope propertyFork/Parallel:
First, read core-fork.md
forkBranchCompleted/forkBranchFailed in StepByStepOrchestrator.ktDSL Parsing:
First, read core-overview.md
NextStepInfoAsyncTaskException for wait/fork/runWorkflowTaskState subclasses must be @Serializablenull in stateUpdates when leaving a nodescope property when state provides expression variablesby lazy for children propertyFlowDirective - Continue, End, or Then(target) for navigationsuspend functionsAsyncTaskException subtypesDirection parameter - behavior differs based on entry directionWorkflowCommand
│
▼
Orchestrator.runByTask()
│
├── Check if condition (skip if false)
├── Transform input (inputFrom)
├── Get processor for node
└── Call processor.getNextStepInfo()
│
├── AsyncTaskException ──► WaitStarted/ForkStarted/RunWorkflowStarted
│
└── NextStepInfo ──► completeTask()
├── Transform output (outputAs)
├── Export to context (exportAs)
└── Navigate to next
│
└── TaskScheduled/WorkflowCompleted/WorkflowFailed
| Purpose | File |
|---|---|
| Step orchestration | orchestrator/StepByStepOrchestrator.kt |
| Full execution | orchestrator/FullOrchestrator.kt |
| Node structure | nodes/Node.kt |
| Position addressing | nodes/NodePosition.kt |
| Processor interface | processors/NodeProcessor.kt |
| Base state | states/TaskState.kt |
| Commands/Events | orchestrator/WorkflowState.kt |
| JQ evaluation | expressions/JQExpression.kt |
| DSL parsing | definitions/DefinitionCache.kt |
// 1. State class
@Serializable
data class CustomState(
override val startedAt: Instant = Clock.System.now(),
val customField: String = ""
) : TaskState() {
// Optional: provide scope variables
override val scope: Scope get() = buildJsonObject {
put("custom", JsonPrimitive(customField))
}
}
// 2. Processor
class CustomProcessor(override val node: Node<CustomTask>) : NodeProcessor<CustomTask, CustomState> {
override fun createInitialState() = CustomState()
override fun getNextStepInfo(state: CustomState, dataset: JsonElement, scope: Scope, direction: Direction): NextStepInfo<CustomState> {
return when (direction) {
FROM_PARENT -> {
// Process task
NextStepInfo(
state = state,
rawOutput = result,
stateUpdates = mapOf(node.position to null), // Clean up
flowDirective = FlowDirective.Continue
)
}
FROM_CHILD -> { /* Handle child completion */ }
else -> { /* Handle other directions */ }
}
}
}
// 3. Register in factory
fun createProcessor(node: Node<*>) = when (node.task) {
is CustomTask -> CustomProcessor(node as Node<CustomTask>)
// ...
}
// For activities that need orchestrator coordination
override fun getNextStepInfo(...): NextStepInfo<WaitState> {
throw AsyncTaskException.WaitStartedException(
state = state,
transformedInput = dataset,
config = WaitStartedException.Config(waitUntil = calculateWaitUntil())
)
}
override fun getNextStepInfo(state, dataset, scope, direction) = when (direction) {
FROM_PARENT -> {
// First entry - initialize and go to first child
NextStepInfo(state = initialState, flowDirective = Continue)
}
FROM_CHILD -> {
if (hasMoreChildren) {
NextStepInfo(state = nextState, flowDirective = Continue)
} else {
// Done - clean up state and return to parent
NextStepInfo(stateUpdates = mapOf(node.position to null), flowDirective = Continue)
}
}
}
@Test
fun `should execute workflow`() = runTest {
val yaml = """
document:
name: test
version: "1.0"
do:
- myTask:
set:
result: "success"
""".trimIndent()
val workflow = DefinitionCache.parse(yaml)
val orchestrator = FullOrchestrator(activityRunner, definitionLoader)
val result = orchestrator.start(workflow, JsonObject(mapOf()))
assertEquals("success", result.jsonObject["result"]?.jsonPrimitive?.content)
}
@Test
fun `DoProcessor should iterate children`() {
val node = createDoNode(childCount = 3)
val processor = DoProcessor(node)
val result = processor.getNextStepInfo(
state = processor.createInitialState(),
dataset = JsonObject(mapOf()),
scope = JsonObject(mapOf()),
direction = Direction.FROM_PARENT
)
assertEquals(0, (result.state as DoState).index)
}
# All tests
./gradlew :lemline-core:test
# Specific test class
./gradlew :lemline-core:test --tests "com.lemline.core.tests.MyTest"
# Specific test method
./gradlew :lemline-core:test --tests "com.lemline.core.tests.MyTest.should do something"