Development guide for the lemline-runner ecosystem...
Guide development across the modular lemline-runner ecosystem - a collection of feature-focused modules built on shared infrastructure. The main lemline-runner module provides messaging, CLI, and configuration, while feature modules (lemline-runner-*) implement specific workflow capabilities.
Core Documentation:
| Module | Purpose |
|---|---|
lemline-runner-common |
Shared infrastructure: outbox pattern, cleaner pattern, repository abstractions, model interfaces |
lemline-runner |
Main runtime: messaging (commands/events), CLI, configuration, activity runners |
| Module | Purpose | Tables |
|---|---|---|
lemline-runner-definitions |
Workflow definition storage and cache sync | lemline_definitions |
lemline-runner-waits |
Wait/sleep task implementation | lemline_waits |
lemline-runner-retries |
Task retry scheduling with exponential backoff | lemline_retries |
lemline-runner-parents |
Parent-child workflow relationships (run task) | lemline_parents |
lemline-runner-forks |
Parallel branch execution (fork task) | lemline_forks, lemline_fork_branches |
lemline-runner-schedules |
Scheduled workflow execution (cron/interval/after) | lemline_schedules |
lemline-runner-listeners |
CloudEvent listeners (listen task) | lemline_listeners, lemline_listener_events |
lemline-runner-failures |
Failed workflow tracking and dead letter storage | lemline_failures |
Each feature module has its own README.md with detailed architecture, file reference, and usage patterns.
A Specific Feature:
Read the module's README first:
Messaging:
First, read runner-messaging.md
Shared Infrastructure:
Read lemline-runner-common/README.md first
AbstractOutbox<T> in lemline-runner-commonAbstractCleaner<T> in lemline-runner-commonlemline-runner-common/repositories/ops/lemline-runner-common/models/Configuration:
First read runner-configuration.md
CLI:
First read runner-cli.md
suspend functions for all database operations (Kotlin coroutines, NOT Mutiny)FOR UPDATE SKIP LOCKED for outbox queries to prevent double-processingsuspend functions insteadāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāā
ā lemline-runner (core) ā
ā āāāāāāāāāāāāāāāāāāāā āāāāāāāāāāāāāāāāāāāā āāāāāāāāāāāāāāāāāāāā ā
ā ā Messaging ā ā CLI Commands ā ā Configuration ā ā
ā ā (commands/events)ā ā (Picocli) ā ā (Quarkus) ā ā
ā āāāāāāāāāāāāāāāāāāāā āāāāāāāāāāāāāāāāāāāā āāāāāāāāāāāāāāāāāāāā ā
āāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāā
ā
āāāāāāāāāāāāāāāāā“āāāāāāāāāāāāāāāāā
ā lemline-runner-common ā
ā (shared infrastructure) ā
ā ⢠Outbox pattern ā
ā ⢠Cleaner pattern ā
ā ⢠Repository abstractions ā
ā ⢠Model interfaces ā
āāāāāāāāāāāāāāāāā¬āāāāāāāāāāāāāāāāā
ā
āāāāāāāāāāāāāāāāāāāāāāāāā¼āāāāāāāāāāāāāāāāāāāāāāāāā
ā ā ā
āāāāāāāāā¼āāāāāāā āāāāāāāāāāāāā¼āāāāāāā āāāāāāāāāāāāāā¼āāāāāā
ā Feature ā ā Feature ā ā Feature ā
ā Modules ā ā Modules ā ā Modules ā
ā ā ā ā ā ā
ā ⢠waits ā ā ⢠schedules ā ā ⢠definitions ā
ā ⢠retries ā ā ⢠listeners ā ā ⢠failures ā
ā ⢠parents ā ā ā ā ā
ā ⢠forks ā ā ā ā ā
āāāāāāāāāāāāāāāā āāāāāāāāāāāāāāāāāāāāā āāāāāāāāāāāāāāāāāāāā
āāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāā
ā COMMANDS CHANNEL (high-throughput) ā
ā commands āāāŗ WorkflowCommandHandler āāāŗ commands ā
ā ā² ā ā
ā ā ā (needs persistence) ā
āāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāā
ā ā
āāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāā
ā ā ā¼ EVENTS CHANNEL ā
ā ā events āāāŗ WorkflowEventHandler āāāŗ Feature Modules ā
ā ā ā ā
ā ā āāāāāāāāāāā“āāāāāāāāāā ā
ā ā ā¼ ā¼ ā
ā ā āāāāāāāāāāāāāāāāāāā āāāāāāāāāāāāāāāāāāā ā
ā ā ā Feature Service ā ā Feature Service ā ā
ā ā ā (e.g., Waits) ā ā (e.g., Parents) ā ā
ā ā āāāāāāāāāā¬āāāāāāāāā āāāāāāāāāā¬āāāāāāāāā ā
ā ā ā ā ā
ā ā ā¼ ā¼ ā
ā ā āāāāāāāāāāāāāāāāāāā āāāāāāāāāāāāāāāāāāā ā
ā ā ā DB Table ā ā DB Table ā ā
ā ā ā (lemline_waits) ā ā(lemline_parents)ā ā
ā ā āāāāāāāāāā¬āāāāāāāāā āāāāāāāāāāāāāāāāāāā ā
ā ā ā ā
ā ā ā¼ ā
ā ā āāāāāāāāāāāāāāāāāāā ā
ā ā ā Outbox Relay ā ā
ā āāāāāāāāāāāāāāāā (scheduled poll)ā ā
ā āāāāāāāāāāāāāāāāāāā ā
āāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāāā
Key principle: State travels with messages. Database only used when necessary. Each feature is self-contained.
| Purpose | File |
|---|---|
| Step execution | StepByStepRunner.kt |
| Command handling | WorkflowCommandHandler.kt |
| Event handling | WorkflowEventHandler.kt |
| Message structure | InstanceMessage.kt |
| Outbox base | AbstractOutbox.kt |
| Cleaner base | AbstractCleaner.kt |
1. Module Structure:
lemline-runner-myfeature/
āāā src/
ā āāā main/kotlin/com/lemline/runner/myfeature/
ā ā āāā MyFeatureService.kt ā Business logic
ā ā āāā MyFeatureModel.kt ā Database entity
ā ā āāā MyFeatureRepository.kt ā Database operations
ā ā āāā MyFeatureOutbox.kt ā Outbox processor (optional)
ā ā āāā MyFeatureCleaner.kt ā Cleanup scheduler (optional)
ā ā āāā MyFeatureConfig.kt ā Configuration
ā āāā test/kotlin/ ā Tests for all databases
ā āāā testFixtures/kotlin/ ā Test utilities
āāā build.gradle.kts
āāā README.md ā Architecture and usage
2. Dependencies in build.gradle.kts:
dependencies {
implementation(project(":lemline-common"))
implementation(project(":lemline-core"))
implementation(project(":lemline-runner-common")) // Always required!
// Add other feature modules if needed
}
3. Model with Interfaces:
// Compose behavior from lemline-runner-common interfaces
data class MyFeatureModel(
override val id: IDV7,
override val instanceMessage: InstanceMessage<MyEvent>,
// Outbox fields
override val outboxScheduledFor: Instant,
override var outboxDelayedUntil: Instant? = outboxScheduledFor,
override var outboxAttemptCount: Int = 0,
override var outboxCompletedAt: Instant? = null,
override var outboxFailedAt: Instant? = null,
// ... other outbox fields
// Cleanup field
override var cleanupAfter: Instant? = null,
) : WithId, WithInstanceMessage, WithOutbox, WithCleanup
4. Repository:
@ApplicationScoped
class MyFeatureRepository : Repository<MyFeatureModel>(),
WithIdRepository<MyFeatureModel>,
OutboxRepository<MyFeatureModel>,
CleanerRepository<MyFeatureModel> {
override suspend fun findByUUID(uuid: IDV7): MyFeatureModel? {
// Use pool from base Repository
return pool.withConnection { conn ->
conn.preparedQuery("SELECT * FROM lemline_myfeature WHERE id = $1")
.execute(Tuple.of(uuid.value))
.awaitSuspending()
.firstOrNull()
?.let { MyFeatureModel.fromRow(it) }
}
}
// OutboxRepository provides findPendingWithLock() automatically
// CleanerRepository provides findOldCompleted() automatically
}
5. Service:
@ApplicationScoped
class MyFeatureService @Inject constructor(
private val repository: MyFeatureRepository,
private val commandEmitter: CommandEmitter
) {
suspend fun handleMyEventStarted(message: InstanceMessage<MyEventStarted>) {
val model = MyFeatureModel.from(message)
repository.insert(model)
}
}
6. Outbox (if needed):
@ApplicationScoped
class MyFeatureOutbox @Inject constructor(
private val repository: MyFeatureRepository,
private val emitter: WorkflowCommandEmitter,
private val config: MyFeatureConfig
) : AbstractOutbox<MyFeatureModel>(
name = "MyFeature",
config = config.outbox
) {
override suspend fun findEntitiesToProcess(limit: Int) =
repository.findPendingWithLock(limit)
override suspend fun process(entity: MyFeatureModel) {
emitter.send(createResumeCommand(entity))
}
override suspend fun markCompleted(entity: MyFeatureModel) {
repository.markCompleted(entity.id)
}
}
7. Register in WorkflowEventHandler:
// In lemline-runner/src/.../messaging/events/WorkflowEventHandler.kt
when (val state = message.state) {
is MyEventStarted -> myFeatureService.handleMyEventStarted(message)
// ...
}
// In feature module service
suspend fun handleMyEvent(message: InstanceMessage<MyEvent>) {
val event = message.state
// 1. Create model from event
val model = MyFeatureModel.from(message, event)
// 2. Persist to database
repository.insert(model)
// 3. If immediate resume needed (no outbox), emit command
if (!needsDelay) {
commandEmitter.send(createResumeCommand(message))
}
// Otherwise, outbox processor will handle it later
}
Test all 3 databases: PostgreSQL, MySQL, H2
// Base test with test logic
abstract class MyRepositoryTestBase : FunSpec({
lateinit var repository: MyRepository
test("should find by UUID") {
val model = createTestModel()
repository.insert(model)
val found = repository.findByUUID(model.id)
found shouldNotBe null
found?.id shouldBe model.id
}
})
// PostgreSQL test
@QuarkusTest
@TestProfile(PostgresProfile::class)
class MyRepositoryPostgresTest : MyRepositoryTestBase() {
@Inject
override lateinit var repository: MyRepository
}
// MySQL test
@QuarkusTest
@TestProfile(MySQLProfile::class)
class MyRepositoryMySQLTest : MyRepositoryTestBase() {
@Inject
override lateinit var repository: MyRepository
}
// H2 test
@QuarkusTest
@TestProfile(H2Profile::class)
class MyRepositoryH2Test : MyRepositoryTestBase() {
@Inject
override lateinit var repository: MyRepository
}
@QuarkusTest
class MyFeatureServiceTest : FunSpec({
@Inject
lateinit var service: MyFeatureService
@Inject
lateinit var repository: MyFeatureRepository
test("should handle event") {
val message = createTestMessage()
service.handleMyEvent(message)
// Verify database state
val stored = repository.findByUUID(message.workflowId)
stored shouldNotBe null
}
})
@QuarkusTest
class MyFeatureOutboxTest : FunSpec({
@Inject
lateinit var outbox: MyFeatureOutbox
@Inject
lateinit var repository: MyFeatureRepository
test("should process pending entities") {
// Insert pending entity
val model = createPendingModel()
repository.insert(model)
// Process
outbox.doWork()
// Verify marked completed
val processed = repository.findByUUID(model.id)
processed?.outboxCompletedAt shouldNotBe null
}
})
IMPORTANT: Place migrations in the feature module, not in lemline-runner!
lemline-runner-myfeature/
āāā src/main/resources/db/migration/
āāā postgresql/
ā āāā V{N}__Create_myfeature_table.sql
āāā mysql/
ā āāā V{N}__Create_myfeature_table.sql
āāā h2/
āāā V{N}__Create_myfeature_table.sql
Naming: V{N}__Description.sql where N is the next available version number across all modules.
Check existing versions first:
# Find highest version number
find . -name "V*.sql" | sort
-- PostgreSQL version
CREATE TABLE lemline_myfeature
(
id UUID PRIMARY KEY,
-- Instance message (serialized workflow state)
instance_message TEXT NOT NULL,
-- Feature-specific columns
my_custom_field VARCHAR(255),
-- Outbox columns (if using outbox pattern)
outbox_scheduled_for TIMESTAMP NOT NULL,
outbox_delayed_until TIMESTAMP NOT NULL,
outbox_attempt_count INT NOT NULL DEFAULT 0,
outbox_completed_at TIMESTAMP,
outbox_failed_at TIMESTAMP,
outbox_error_class VARCHAR(255),
outbox_error_message VARCHAR(500),
outbox_error_stacktrace TEXT,
-- Cleanup column (if using cleaner pattern)
cleanup_after TIMESTAMP,
-- Timestamps
created_at TIMESTAMP NOT NULL DEFAULT NOW()
);
-- Index for outbox queries (FOR UPDATE SKIP LOCKED)
CREATE INDEX idx_lemline_myfeature_pending
ON lemline_myfeature (outbox_delayed_until)
WHERE outbox_completed_at IS NULL AND outbox_failed_at IS NULL;
-- Index for cleanup queries
CREATE INDEX idx_lemline_myfeature_cleanup
ON lemline_myfeature (cleanup_after)
WHERE cleanup_after IS NOT NULL;
PostgreSQL:
UUID typeTEXT for long stringsWHERE clauseMySQL:
CHAR(36) instead of UUIDLONGTEXT for long stringsH2:
UUID typeCLOB for very long strings# Test with PostgreSQL
./gradlew :lemline-runner-myfeature:test -Dquarkus.test.profile=postgres
# Test with MySQL
./gradlew :lemline-runner-myfeature:test -Dquarkus.test.profile=mysql
# Test with H2
./gradlew :lemline-runner-myfeature:test -Dquarkus.test.profile=h2
# All tests in main runner
./gradlew :lemline-runner:test
# All tests in feature module
./gradlew :lemline-runner-myfeature:test
# Specific test class
./gradlew :lemline-runner-myfeature:test --tests "com.lemline.runner.myfeature.MyTest"
# Test specific module with specific database
./gradlew :lemline-runner-waits:test -Dquarkus.test.profile=postgres
./gradlew :lemline-runner-waits:test -Dquarkus.test.profile=mysql
./gradlew :lemline-runner-waits:test -Dquarkus.test.profile=h2
# Test all runner modules
./gradlew test -p lemline-runner -p lemline-runner-common -p lemline-runner-waits -p lemline-runner-retries ...
Example: Adding support for a new task type that requires persistence
Add task to lemline-core (see core-dev skill)
lemline-core/src/.../models/tasks/lemline-core/src/.../processors/AsyncTaskException if needs persistenceCreate feature module
mkdir -p lemline-runner-myfeature/src/{main,test}/kotlin/com/lemline/runner/myfeature
mkdir -p lemline-runner-myfeature/src/main/resources/db/migration/{postgresql,mysql,h2}
Add to settings.gradle.kts
include("lemline-runner-myfeature")
Create build.gradle.kts
dependencies {
implementation(project(":lemline-common"))
implementation(project(":lemline-core"))
implementation(project(":lemline-runner-common"))
}
Implement feature (Model, Repository, Service, Outbox, Cleaner)
Add migrations for all 3 databases
Register in WorkflowEventHandler
// In lemline-runner
when (val state = message.state) {
is MyFeatureStarted -> myFeatureService.handleStarted(message)
// ...
}
Add tests for all databases
Add to main runner dependencies
// In lemline-runner/build.gradle.kts
implementation(project(":lemline-runner-myfeature"))
Document in README.md in the feature module
lemline-common
ā
ā
lemline-core
ā
ā
lemline-runner-common āāāāāā
ā ā
ā ā
āāāāāāāāāāāāāāāāāāāāāāāā¼āāāāāāāāāāāāāāāāāā
ā ā ā
lemline-runner-* lemline-runner-* lemline-runner-*
(feature modules) (feature modules) (feature modules)
ā ā ā
ā ā ā
āāāāāāāāāāāāāāāāāāāāāāāā“āāāāāāāāāāāāāāāāāā
ā
lemline-runner
(main runtime)
Key principle: Feature modules depend only on common modules, not on each other (except rare cases like schedules/parents).