Skip to content

Repository files navigation

WARNING: This is an experimental project. Do not use in a production environment

Durable workflows using swift-jobs

// Create queue to run activities on
let activityJobQueue = JobQueue(
    .valkey(valkeyClient, configuration: .init(queueName: "Activities"), logger: logger),
    logger: logger
)
// run job queue processor in background
async let _ = activityQueue.processor().run()

// create event source and dispatch manager
let eventSource = ValkeyEventSource(client: valkeyClient, namespace: "Activities")
let dispatchManager = WorkflowDispatchManager(activityQueueDriver: activityJobQueue, eventSource: eventSource)
// register activities with dispatch manager
let activity = dispatchManager.registerActivity(name: "plus") { (input: Int, _) in
    input + 1
}
let activity = dispatchManager.registerActivity(name: "multiple") { (input: Int, _) in
    input * 2
}
// create workflow
let basicWorkflow = WorkflowDefinition(name: "basic") { (input: Int, context) in
    async let result = context.runActivity(activity1, input: input)
    async let result2 = context.runActivity(activity2, input: input)
    return try await Output(first: result, second: result2)
}
// run workflow
let result = try await dispatchManager.runWorkflow(
    workflow: basicWorkflow,
    workflowID: WorkflowID(),
    input: 25,
    logger: logger
)

You can also setup a WorkflowQueue which uses a JobQueue to queue workflows. NB The workflow job queue creates two separate job queues one for workflows and one for activities. If workflows and activities were to run on the same queue there is a good chance the job queue would get stuck running only workflows.

let workflowQueue = try await WorkflowJobQueue(
    valkeyClient: valkeyClient,
    configuration: .init(name: "TestWorkflow"),
    logger: logger
)
// run job queue processors and valkey client
async let _ = valkeyClient.run()
async let _ = workflowQueue.processor().run()
let workflow = workflowQueue.registerWorkflow(name: "basic", input: Int.self) { input, context in
    async let result = context.runActivity(activity1, input: input)
    async let result2 = context.runActivity(activity2, input: input)
    return try await Output(first: result, second: result2)
}
let output = try await workflowQueue.runWorkflow(workflow, parameters: 25)

About

Durable workflows using Swift Jobs

Topics

Resources

Stars

4 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages