Skip to content

Commit 663f046

Browse files
Add application id extension to lifecycle events (#1750)
Signed-off-by: Matheus André <matheusandr2@gmail.com>
1 parent 2c9e734 commit 663f046

3 files changed

Lines changed: 44 additions & 16 deletions

File tree

‎impl/core/src/main/java/io/serverlessworkflow/impl/LifecycleEvents.java‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -66,4 +66,10 @@ private LifecycleEvents() {}
6666
/** Notifies about the change of a workflow's status phase. */
6767
public static final String WORKFLOW_STATUS_CHANGED =
6868
"io.serverlessworkflow.workflow.status-changed.v1";
69+
70+
/**
71+
* CloudEvent extension attribute carried by every lifecycle event, holding the {@link
72+
* WorkflowApplication#id()} of the application that produced it.
73+
*/
74+
public static final String APPLICATION_ID_EXTENSION = "applicationid";
6975
}

‎impl/core/src/main/java/io/serverlessworkflow/impl/lifecycle/ce/AbstractLifeCyclePublisher.java‎

Lines changed: 18 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,7 @@
1515
*/
1616
package io.serverlessworkflow.impl.lifecycle.ce;
1717

18+
import static io.serverlessworkflow.impl.LifecycleEvents.APPLICATION_ID_EXTENSION;
1819
import static io.serverlessworkflow.impl.LifecycleEvents.TASK_CANCELLED;
1920
import static io.serverlessworkflow.impl.LifecycleEvents.TASK_COMPLETED;
2021
import static io.serverlessworkflow.impl.LifecycleEvents.TASK_FAULTED;
@@ -89,7 +90,7 @@ public void onTaskStarted(TaskStartedEvent event) {
8990
event,
9091
ev ->
9192
factory.build(
92-
builder()
93+
builder(event)
9394
.withData(cloudEventData(factory.build(event), this::convert))
9495
.withType(TASK_STARTED)));
9596
}
@@ -101,7 +102,7 @@ public void onTaskRetried(TaskRetriedEvent event) {
101102
event,
102103
ev ->
103104
factory.build(
104-
builder()
105+
builder(event)
105106
.withData(cloudEventData(factory.build(event), this::convert))
106107
.withType(TASK_RETRIED)));
107108
}
@@ -113,7 +114,7 @@ public void onTaskCompleted(TaskCompletedEvent event) {
113114
event,
114115
ev ->
115116
factory.build(
116-
builder()
117+
builder(event)
117118
.withData(cloudEventData(factory.build(event), this::convert))
118119
.withType(TASK_COMPLETED)));
119120
}
@@ -125,7 +126,7 @@ public void onTaskSuspended(TaskSuspendedEvent event) {
125126
event,
126127
ev ->
127128
factory.build(
128-
builder()
129+
builder(event)
129130
.withData(cloudEventData(factory.build(event), this::convert))
130131
.withType(TASK_SUSPENDED)));
131132
}
@@ -137,7 +138,7 @@ public void onTaskResumed(TaskResumedEvent event) {
137138
event,
138139
ev ->
139140
factory.build(
140-
builder()
141+
builder(event)
141142
.withData(cloudEventData(factory.build(event), this::convert))
142143
.withType(TASK_RESUMED)));
143144
}
@@ -149,7 +150,7 @@ public void onTaskCancelled(TaskCancelledEvent event) {
149150
event,
150151
ev ->
151152
factory.build(
152-
builder()
153+
builder(event)
153154
.withData(cloudEventData(factory.build(event), this::convert))
154155
.withType(TASK_CANCELLED)));
155156
}
@@ -161,7 +162,7 @@ public void onTaskFailed(TaskFailedEvent event) {
161162
event,
162163
ev ->
163164
factory.build(
164-
builder()
165+
builder(event)
165166
.withData(cloudEventData(factory.build(event), this::convert))
166167
.withType(TASK_FAULTED)));
167168
}
@@ -173,7 +174,7 @@ public void onWorkflowStarted(WorkflowStartedEvent event) {
173174
event,
174175
ev ->
175176
factory.build(
176-
builder()
177+
builder(event)
177178
.withData(cloudEventData(factory.build(event), this::convert))
178179
.withType(WORKFLOW_STARTED)));
179180
}
@@ -185,7 +186,7 @@ public void onWorkflowSuspended(WorkflowSuspendedEvent event) {
185186
event,
186187
ev ->
187188
factory.build(
188-
builder()
189+
builder(event)
189190
.withData(cloudEventData(factory.build(event), this::convert))
190191
.withType(WORKFLOW_SUSPENDED)));
191192
}
@@ -197,7 +198,7 @@ public void onWorkflowCancelled(WorkflowCancelledEvent event) {
197198
event,
198199
ev ->
199200
factory.build(
200-
builder()
201+
builder(event)
201202
.withData(cloudEventData(factory.build(event), this::convert))
202203
.withType(WORKFLOW_CANCELLED)));
203204
}
@@ -209,7 +210,7 @@ public void onWorkflowResumed(WorkflowResumedEvent event) {
209210
event,
210211
ev ->
211212
factory.build(
212-
builder()
213+
builder(event)
213214
.withData(cloudEventData(factory.build(event), this::convert))
214215
.withType(WORKFLOW_RESUMED)));
215216
}
@@ -221,7 +222,7 @@ public void onWorkflowCompleted(WorkflowCompletedEvent event) {
221222
event,
222223
ev ->
223224
factory.build(
224-
builder()
225+
builder(event)
225226
.withData(cloudEventData(factory.build(event), this::convert))
226227
.withType(WORKFLOW_COMPLETED)));
227228
}
@@ -233,7 +234,7 @@ public void onWorkflowFailed(WorkflowFailedEvent event) {
233234
event,
234235
ev ->
235236
factory.build(
236-
builder()
237+
builder(event)
237238
.withData(cloudEventData(factory.build(event), this::convert))
238239
.withType(WORKFLOW_FAULTED)));
239240
}
@@ -246,7 +247,7 @@ public void onWorkflowStatusChanged(WorkflowStatusEvent event) {
246247
event,
247248
ev ->
248249
factory.build(
249-
builder()
250+
builder(event)
250251
.withData(cloudEventData(factory.build(event), this::convert))
251252
.withType(WORKFLOW_STATUS_CHANGED)));
252253
}
@@ -328,11 +329,12 @@ private static <T> CloudEventData cloudEventData(T data, ToBytes<T> toBytes) {
328329
return PojoCloudEventData.wrap(data, toBytes);
329330
}
330331

331-
private static CloudEventBuilder builder() {
332+
private static CloudEventBuilder builder(WorkflowEvent ev) {
332333
return CloudEventBuilder.v1()
333334
.withId(CloudEventUtils.id())
334335
.withSource(CloudEventUtils.source())
335-
.withTime(OffsetDateTime.now());
336+
.withTime(OffsetDateTime.now())
337+
.withExtension(APPLICATION_ID_EXTENSION, appl(ev).id());
336338
}
337339

338340
private static WorkflowApplication appl(WorkflowEvent ev) {

‎impl/test/src/test/java/io/serverlessworkflow/impl/test/LifeCycleEventsTest.java‎

Lines changed: 20 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,7 @@
1515
*/
1616
package io.serverlessworkflow.impl.test;
1717

18+
import static io.serverlessworkflow.impl.LifecycleEvents.APPLICATION_ID_EXTENSION;
1819
import static org.assertj.core.api.Assertions.assertThat;
1920
import static org.assertj.core.api.Assertions.catchThrowableOfType;
2021
import static org.awaitility.Awaitility.await;
@@ -125,6 +126,24 @@ void simpleWorkflow() throws IOException {
125126
assertThat(taskStartedEvent.startedAt()).isBefore(taskCompletedEvent.completedAt());
126127
}
127128

129+
@Test
130+
void testApplicationIdExtension() throws IOException {
131+
appl.workflowDefinition(
132+
WorkflowReader.readWorkflowFromClasspath("workflows-samples/simple-expression.yaml"))
133+
.instance(Map.of())
134+
.start()
135+
.join();
136+
assertPojoInCE("io.serverlessworkflow.workflow.completed.v1", Object.class);
137+
assertThat(appl.id()).isNotBlank();
138+
assertThat(publishedEvents)
139+
.isNotEmpty()
140+
.allSatisfy(
141+
ce ->
142+
assertThat(ce.getExtension(APPLICATION_ID_EXTENSION))
143+
.as("%s extension of %s", APPLICATION_ID_EXTENSION, ce.getType())
144+
.isEqualTo(appl.id()));
145+
}
146+
128147
@ParameterizedTest(name = "{0}")
129148
@MethodSource("waitSetWorkflowSources")
130149
void testSuspendResumeNotWait(String sourceName, Workflow workflow)
@@ -222,6 +241,7 @@ private <T> T assertPojoInCE(String type, Class<T> clazz) {
222241
() -> publishedEvents.stream().filter(ev -> ev.getType().equals(type)).findAny(),
223242
Optional::isPresent)
224243
.orElseThrow();
244+
assertThat(ce.getExtension(APPLICATION_ID_EXTENSION)).isEqualTo(appl.id());
225245
assertThat(ce.getData()).isInstanceOf(PojoCloudEventData.class);
226246
Object pojo = ((PojoCloudEventData<?>) Objects.requireNonNull(ce.getData())).getValue();
227247
assertThat(pojo).isInstanceOf(clazz);

0 commit comments

Comments
 (0)