From 7ddb5715c894904a97071d9ae00c83fe9fe4621c Mon Sep 17 00:00:00 2001 From: Pratyush Sharma <56130065+pratyush618@users.noreply.github.com> Date: Sun, 28 Jun 2026 19:32:25 +0530 Subject: [PATCH 1/2] feat(java): canvas workflow shortcuts (chain/group/chord) Canvas builds a Workflow from links: chain (sequential), group (parallel), chord (callback after a parallel group). --- .../byteveda/taskito/workflows/Canvas.java | 70 +++++++++++++++++++ 1 file changed, 70 insertions(+) create mode 100644 sdks/java/src/main/java/org/byteveda/taskito/workflows/Canvas.java diff --git a/sdks/java/src/main/java/org/byteveda/taskito/workflows/Canvas.java b/sdks/java/src/main/java/org/byteveda/taskito/workflows/Canvas.java new file mode 100644 index 00000000..04190257 --- /dev/null +++ b/sdks/java/src/main/java/org/byteveda/taskito/workflows/Canvas.java @@ -0,0 +1,70 @@ +package org.byteveda.taskito.workflows; + +import org.byteveda.taskito.task.Task; + +/** + * Celery-style shortcuts that build a {@link Workflow} from a few links: + * {@link #chain} (sequential), {@link #group} (parallel), and {@link #chord} + * (a callback after a parallel group). Submit the result like any workflow. + * + *

A {@code chord}'s callback runs once every group step completes; like the + * rest of the static-DAG model it receives its own payload, not the group's + * aggregated results (use a fan-out + fan-in for aggregation). + */ +public final class Canvas { + private Canvas() {} + + /** A named step bound to a task and payload, for use in canvas shortcuts. */ + public static Link link(String name, Task task, T payload) { + return new Link(name, task.name(), payload); + } + + /** Run {@code links} one after another (each depends on the previous). */ + public static Workflow chain(String name, Link... links) { + Workflow workflow = Workflow.named(name); + String previous = null; + for (Link link : links) { + workflow.step(link.toStep(previous == null ? new String[0] : new String[] {previous})); + previous = link.name; + } + return workflow; + } + + /** Run all {@code links} in parallel (no dependencies between them). */ + public static Workflow group(String name, Link... links) { + Workflow workflow = Workflow.named(name); + for (Link link : links) { + workflow.step(link.toStep(new String[0])); + } + return workflow; + } + + /** Run {@code group} in parallel, then {@code callback} once all of them complete. */ + public static Workflow chord(String name, Link callback, Link... group) { + Workflow workflow = Workflow.named(name); + String[] groupNames = new String[group.length]; + for (int i = 0; i < group.length; i++) { + workflow.step(group[i].toStep(new String[0])); + groupNames[i] = group[i].name; + } + workflow.step(callback.toStep(groupNames)); + return workflow; + } + + /** A canvas building block: a step name, its task, and its payload. */ + public static final class Link { + private final String name; + private final String taskName; + private final Object payload; + + Link(String name, String taskName, Object payload) { + this.name = name; + this.taskName = taskName; + this.payload = payload; + } + + Step toStep(String[] after) { + return Step.of(name, taskName, payload).after(after).build(); + } + } +} From 1ac316f33d5cd43a37191d6c2bdaa9ab2bfb217d Mon Sep 17 00:00:00 2001 From: Pratyush Sharma <56130065+pratyush618@users.noreply.github.com> Date: Sun, 28 Jun 2026 19:32:25 +0530 Subject: [PATCH 2/2] test(java): cover canvas shortcuts --- .../java/org/byteveda/taskito/CanvasTest.java | 40 +++++++++++++++++++ 1 file changed, 40 insertions(+) create mode 100644 sdks/java/src/test/java/org/byteveda/taskito/CanvasTest.java diff --git a/sdks/java/src/test/java/org/byteveda/taskito/CanvasTest.java b/sdks/java/src/test/java/org/byteveda/taskito/CanvasTest.java new file mode 100644 index 00000000..e0943932 --- /dev/null +++ b/sdks/java/src/test/java/org/byteveda/taskito/CanvasTest.java @@ -0,0 +1,40 @@ +package org.byteveda.taskito; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.util.List; +import java.util.Set; +import org.byteveda.taskito.task.Task; +import org.byteveda.taskito.workflows.Canvas; +import org.byteveda.taskito.workflows.Workflow; +import org.byteveda.taskito.workflows.WorkflowAnalysis; +import org.junit.jupiter.api.Test; + +class CanvasTest { + + private static final Task T = Task.of("cv.t", Integer.class); + + @Test + void chainIsLinear() { + Workflow wf = Canvas.chain("c", Canvas.link("a", T, 1), Canvas.link("b", T, 2), Canvas.link("c", T, 3)); + assertEquals(List.of("a", "b", "c"), WorkflowAnalysis.topologicalOrder(wf)); + assertEquals(List.of("a"), WorkflowAnalysis.roots(wf)); + assertEquals(List.of("c"), WorkflowAnalysis.leaves(wf)); + } + + @Test + void groupIsParallel() { + Workflow wf = Canvas.group("g", Canvas.link("a", T, 1), Canvas.link("b", T, 2), Canvas.link("c", T, 3)); + assertEquals(Set.of("a", "b", "c"), Set.copyOf(WorkflowAnalysis.roots(wf))); + assertEquals(Set.of("a", "b", "c"), Set.copyOf(WorkflowAnalysis.leaves(wf))); + } + + @Test + void chordCallbackRunsAfterGroup() { + Workflow wf = Canvas.chord("ch", Canvas.link("cb", T, 0), Canvas.link("a", T, 1), Canvas.link("b", T, 2)); + assertEquals(Set.of("a", "b"), WorkflowAnalysis.ancestors(wf, "cb")); + assertEquals(List.of("cb"), WorkflowAnalysis.leaves(wf)); + assertTrue(WorkflowAnalysis.roots(wf).containsAll(Set.of("a", "b"))); + } +}