Skip to content

Commit 994e577

Browse files
Fix link wiring
1 parent 441ee33 commit 994e577

5 files changed

Lines changed: 217 additions & 2 deletions

File tree

temporal-sdk/src/main/java/io/temporal/internal/common/LinkConverter.java

Lines changed: 75 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,8 @@ public class LinkConverter {
2222
private static final String linkPathFormat = "temporal:///namespaces/%s/workflows/%s/%s/history";
2323
private static final String nexusOperationLinkPathFormat =
2424
"temporal:///namespaces/%s/nexus-operations/%s/%s/details";
25+
private static final String activityLinkPathFormat =
26+
"temporal:///namespaces/%s/activities/%s/%s/details";
2527
private static final String linkReferenceTypeKey = "referenceType";
2628
private static final String linkEventIDKey = "eventID";
2729
private static final String linkEventTypeKey = "eventType";
@@ -35,6 +37,7 @@ public class LinkConverter {
3537
Link.WorkflowEvent.getDescriptor().getFullName();
3638
private static final String nexusOperationLinkType =
3739
Link.NexusOperation.getDescriptor().getFullName();
40+
private static final String activityLinkType = Link.Activity.getDescriptor().getFullName();
3841

3942
public static io.temporal.api.nexus.v1.Link workflowEventToNexusLink(Link.WorkflowEvent we) {
4043
try {
@@ -177,6 +180,9 @@ public static io.temporal.api.nexus.v1.Link linkToNexusLink(Link commonLink) {
177180
if (commonLink.hasNexusOperation()) {
178181
return nexusOperationToNexusLink(commonLink.getNexusOperation());
179182
}
183+
if (commonLink.hasActivity()) {
184+
return activityToNexusLink(commonLink.getActivity());
185+
}
180186
return null;
181187
}
182188

@@ -192,10 +198,79 @@ public static Link nexusLinkToLink(io.temporal.api.nexus.v1.Link nexusLink) {
192198
if (nexusOperationLinkType.equals(type)) {
193199
return nexusLinkToNexusOperation(nexusLink);
194200
}
201+
if (activityLinkType.equals(type)) {
202+
return nexusLinkToActivity(nexusLink);
203+
}
195204
log.warn("ignoring unsupported nexus link type: {}", type);
196205
return null;
197206
}
198207

208+
public static io.temporal.api.nexus.v1.Link activityToNexusLink(Link.Activity activity) {
209+
try {
210+
String url =
211+
String.format(
212+
activityLinkPathFormat,
213+
URLEncoder.encode(activity.getNamespace(), StandardCharsets.UTF_8.toString()),
214+
URLEncoder.encode(activity.getActivityId(), StandardCharsets.UTF_8.toString())
215+
.replace("+", "%20"),
216+
URLEncoder.encode(activity.getRunId(), StandardCharsets.UTF_8.toString()));
217+
return io.temporal.api.nexus.v1.Link.newBuilder()
218+
.setUrl(url)
219+
.setType(activityLinkType)
220+
.build();
221+
} catch (Exception e) {
222+
log.error("Failed to encode activity Nexus link URL", e);
223+
}
224+
return null;
225+
}
226+
227+
public static Link nexusLinkToActivity(io.temporal.api.nexus.v1.Link nexusLink) {
228+
if (!activityLinkType.equals(nexusLink.getType())) {
229+
log.error(
230+
"Failed to parse Nexus link URL: cannot parse link type {} to {}",
231+
nexusLink.getType(),
232+
activityLinkType);
233+
return null;
234+
}
235+
Link.Builder link = Link.newBuilder();
236+
try {
237+
URI uri = new URI(nexusLink.getUrl());
238+
if (!"temporal".equals(uri.getScheme())) {
239+
log.error("Failed to parse Nexus link URL: invalid scheme: {}", uri.getScheme());
240+
return null;
241+
}
242+
StringTokenizer st = new StringTokenizer(uri.getRawPath(), "/");
243+
if (!st.hasMoreTokens() || !st.nextToken().equals("namespaces")) {
244+
log.error("Failed to parse Nexus link URL: invalid path: {}", uri.getRawPath());
245+
return null;
246+
}
247+
String namespace = URLDecoder.decode(st.nextToken(), StandardCharsets.UTF_8.toString());
248+
if (!st.hasMoreTokens() || !st.nextToken().equals("activities")) {
249+
log.error("Failed to parse Nexus link URL: invalid path: {}", uri.getRawPath());
250+
return null;
251+
}
252+
String activityId = URLDecoder.decode(st.nextToken(), StandardCharsets.UTF_8.toString());
253+
if (!st.hasMoreTokens()) {
254+
log.error("Failed to parse Nexus link URL: invalid path: {}", uri.getRawPath());
255+
return null;
256+
}
257+
String runId = URLDecoder.decode(st.nextToken(), StandardCharsets.UTF_8.toString());
258+
if (!st.hasMoreTokens() || !st.nextToken().equals("details")) {
259+
log.error("Failed to parse Nexus link URL: invalid path: {}", uri.getRawPath());
260+
return null;
261+
}
262+
link.setActivity(
263+
Link.Activity.newBuilder()
264+
.setNamespace(namespace)
265+
.setActivityId(activityId)
266+
.setRunId(runId));
267+
} catch (Exception e) {
268+
log.error("Failed to parse activity Nexus link URL", e);
269+
return null;
270+
}
271+
return link.build();
272+
}
273+
199274
public static io.temporal.api.nexus.v1.Link nexusOperationToNexusLink(Link.NexusOperation no) {
200275
try {
201276
String url =

temporal-sdk/src/main/java/io/temporal/nexus/TemporalNexusClientImpl.java

Lines changed: 17 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4,13 +4,15 @@
44
import io.nexusrpc.handler.HandlerException;
55
import io.nexusrpc.handler.OperationContext;
66
import io.nexusrpc.handler.OperationStartDetails;
7+
import io.temporal.api.common.v1.Payload;
78
import io.temporal.client.ActivityClient;
89
import io.temporal.client.ActivityClientOptions;
910
import io.temporal.client.StartActivityOptions;
1011
import io.temporal.client.WorkflowClient;
1112
import io.temporal.client.WorkflowOptions;
1213
import io.temporal.client.WorkflowStub;
1314
import io.temporal.common.Experimental;
15+
import io.temporal.common.context.ContextPropagator;
1416
import io.temporal.common.interceptors.ActivityClientCallsInterceptor;
1517
import io.temporal.common.interceptors.Header;
1618
import io.temporal.internal.client.ActivityClientInternal;
@@ -28,7 +30,9 @@
2830
import java.lang.reflect.Type;
2931
import java.util.Arrays;
3032
import java.util.Collections;
33+
import java.util.HashMap;
3134
import java.util.List;
35+
import java.util.Map;
3236
import java.util.Objects;
3337
import java.util.concurrent.atomic.AtomicBoolean;
3438

@@ -491,7 +495,7 @@ private <R> TemporalOperationResult<R> startActivityImpl(
491495
activityType,
492496
args,
493497
options,
494-
Header.empty(),
498+
propagatedHeader(),
495499
request -> {
496500
ActivityClientCallsInterceptor.StartActivityInput input =
497501
new ActivityClientCallsInterceptor.StartActivityInput(
@@ -561,4 +565,16 @@ private void markAsyncOperationStarted() {
561565
+ "Use getWorkflowClient() for additional workflow interactions."));
562566
}
563567
}
568+
569+
private Header propagatedHeader() {
570+
List<ContextPropagator> propagators = client.getOptions().getContextPropagators();
571+
if (propagators.isEmpty()) {
572+
return Header.empty();
573+
}
574+
Map<String, Payload> result = new HashMap<>();
575+
for (ContextPropagator propagator : propagators) {
576+
result.putAll(propagator.serializeContext(propagator.getCurrentContext()));
577+
}
578+
return new Header(result);
579+
}
564580
}

temporal-sdk/src/test/java/io/temporal/internal/client/RootActivityClientInvokerTest.java

Lines changed: 16 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -43,7 +43,10 @@ public void setUp() {
4343
genericClient = mock(GenericWorkflowClient.class);
4444
when(genericClient.startActivity(any(StartActivityExecutionRequest.class)))
4545
.thenReturn(
46-
StartActivityExecutionResponse.newBuilder().setRunId("activity-run-id").build());
46+
StartActivityExecutionResponse.newBuilder()
47+
.setRunId("activity-run-id")
48+
.setLink(activityLink())
49+
.build());
4750
invoker =
4851
new RootActivityClientInvoker(
4952
genericClient,
@@ -97,6 +100,7 @@ public void nexusMetadataAddsCallbackLinksAndRequestId() {
97100
.getCompletionCallbacks(0)
98101
.getNexus()
99102
.getHeaderOrThrow(io.nexusrpc.Header.OPERATION_TOKEN.toLowerCase()));
103+
Assert.assertEquals(Collections.singletonList(activityLink()), nexusContext.getResponseLinks());
100104
}
101105

102106
@Test
@@ -112,6 +116,7 @@ public void nexusContextWithoutMetadataStartsOrdinaryActivity() {
112116
Assert.assertFalse(request.getRequestId().isEmpty());
113117
Assert.assertEquals(0, request.getLinksCount());
114118
Assert.assertEquals(0, request.getCompletionCallbacksCount());
119+
Assert.assertTrue(nexusContext.getResponseLinks().isEmpty());
115120
}
116121

117122
private static StartActivityInput newStartActivityInput() {
@@ -136,4 +141,14 @@ private static Link workflowEventLink() {
136141
.setEventType(EventType.EVENT_TYPE_NEXUS_OPERATION_SCHEDULED)))
137142
.build();
138143
}
144+
145+
private static Link activityLink() {
146+
return Link.newBuilder()
147+
.setActivity(
148+
Link.Activity.newBuilder()
149+
.setNamespace(NAMESPACE)
150+
.setActivityId("activity-id")
151+
.setRunId("activity-run-id"))
152+
.build();
153+
}
139154
}

temporal-sdk/src/test/java/io/temporal/internal/common/LinkConverterTest.java

Lines changed: 80 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,8 @@
11
package io.temporal.internal.common;
22

3+
import static io.temporal.internal.common.LinkConverter.activityToNexusLink;
34
import static io.temporal.internal.common.LinkConverter.linkToNexusLink;
5+
import static io.temporal.internal.common.LinkConverter.nexusLinkToActivity;
46
import static io.temporal.internal.common.LinkConverter.nexusLinkToLink;
57
import static io.temporal.internal.common.LinkConverter.nexusLinkToNexusOperation;
68
import static io.temporal.internal.common.LinkConverter.nexusLinkToWorkflowEvent;
@@ -466,6 +468,57 @@ public void testConvertNexusToNexusOperation_InvalidPathMissingDetails() {
466468
assertNull(nexusLinkToNexusOperation(input));
467469
}
468470

471+
@Test
472+
public void testConvertActivityToNexus_Valid() {
473+
Link.Activity input =
474+
Link.Activity.newBuilder()
475+
.setNamespace("ns")
476+
.setActivityId("act id/with+characters")
477+
.setRunId("run-id")
478+
.build();
479+
480+
io.temporal.api.nexus.v1.Link expected =
481+
io.temporal.api.nexus.v1.Link.newBuilder()
482+
.setUrl(
483+
"temporal:///namespaces/ns/activities/act%20id%2Fwith%2Bcharacters/run-id/details")
484+
.setType("temporal.api.common.v1.Link.Activity")
485+
.build();
486+
487+
assertEquals(expected, activityToNexusLink(input));
488+
}
489+
490+
@Test
491+
public void testConvertNexusToActivity_Valid() {
492+
io.temporal.api.nexus.v1.Link input =
493+
io.temporal.api.nexus.v1.Link.newBuilder()
494+
.setUrl(
495+
"temporal:///namespaces/ns/activities/act%20id%2Fwith%2Bcharacters/run-id/details")
496+
.setType("temporal.api.common.v1.Link.Activity")
497+
.build();
498+
499+
Link expected =
500+
Link.newBuilder()
501+
.setActivity(
502+
Link.Activity.newBuilder()
503+
.setNamespace("ns")
504+
.setActivityId("act id/with+characters")
505+
.setRunId("run-id"))
506+
.build();
507+
508+
assertEquals(expected, nexusLinkToActivity(input));
509+
}
510+
511+
@Test
512+
public void testConvertNexusToActivity_InvalidPath() {
513+
io.temporal.api.nexus.v1.Link input =
514+
io.temporal.api.nexus.v1.Link.newBuilder()
515+
.setUrl("temporal:///namespaces/ns/activities/act-id/run-id")
516+
.setType("temporal.api.common.v1.Link.Activity")
517+
.build();
518+
519+
assertNull(nexusLinkToActivity(input));
520+
}
521+
469522
@Test
470523
public void testNexusLinkToLink_WorkflowEventRoundTrip() {
471524
Link.WorkflowEvent we =
@@ -507,6 +560,19 @@ public void testNexusLinkToLink_NexusOperation() {
507560
assertEquals(expected, nexusLinkToLink(nexusLink));
508561
}
509562

563+
@Test
564+
public void testNexusLinkToLink_ActivityRoundTrip() {
565+
Link.Activity activity =
566+
Link.Activity.newBuilder()
567+
.setNamespace("ns")
568+
.setActivityId("act-id")
569+
.setRunId("run-id")
570+
.build();
571+
572+
io.temporal.api.nexus.v1.Link nexusLink = activityToNexusLink(activity);
573+
assertEquals(Link.newBuilder().setActivity(activity).build(), nexusLinkToLink(nexusLink));
574+
}
575+
510576
@Test
511577
public void testNexusLinkToLink_UnknownType() {
512578
io.temporal.api.nexus.v1.Link nexusLink =
@@ -550,6 +616,20 @@ public void testLinkToNexusLink_NexusOperation() {
550616
assertEquals(nexusOperationToNexusLink(no), actual);
551617
}
552618

619+
@Test
620+
public void testLinkToNexusLink_Activity() {
621+
Link.Activity activity =
622+
Link.Activity.newBuilder()
623+
.setNamespace("ns")
624+
.setActivityId("act-id")
625+
.setRunId("run-id")
626+
.build();
627+
628+
io.temporal.api.nexus.v1.Link actual =
629+
linkToNexusLink(Link.newBuilder().setActivity(activity).build());
630+
assertEquals(activityToNexusLink(activity), actual);
631+
}
632+
553633
@Test
554634
public void testLinkToNexusLink_Empty() {
555635
assertNull(linkToNexusLink(Link.newBuilder().build()));

temporal-sdk/src/test/java/io/temporal/nexus/TemporalNexusClientImplTest.java

Lines changed: 29 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -10,13 +10,15 @@
1010
import io.nexusrpc.handler.OperationContext;
1111
import io.nexusrpc.handler.OperationStartDetails;
1212
import io.temporal.api.common.v1.Link;
13+
import io.temporal.api.common.v1.Payload;
1314
import io.temporal.api.common.v1.WorkflowExecution;
1415
import io.temporal.client.ActivityClient;
1516
import io.temporal.client.ActivityClientOptions;
1617
import io.temporal.client.StartActivityOptions;
1718
import io.temporal.client.WorkflowClient;
1819
import io.temporal.client.WorkflowClientOptions;
1920
import io.temporal.client.WorkflowOptions;
21+
import io.temporal.common.context.ContextPropagator;
2022
import io.temporal.common.interceptors.ActivityClientCallsInterceptor;
2123
import io.temporal.internal.client.ActivityClientInternal;
2224
import io.temporal.internal.client.NexusStartWorkflowRequest;
@@ -27,6 +29,8 @@
2729
import io.temporal.serviceclient.WorkflowServiceStubs;
2830
import io.temporal.workflow.Functions;
2931
import java.time.Duration;
32+
import java.util.Collections;
33+
import java.util.concurrent.atomic.AtomicReference;
3034
import org.junit.After;
3135
import org.junit.Assert;
3236
import org.junit.Before;
@@ -50,6 +54,7 @@ public class TemporalNexusClientImplTest {
5054

5155
private TemporalNexusClientImpl client;
5256
private MockedStatic<ActivityClient> activityClientFactory;
57+
private AtomicReference<ActivityClientCallsInterceptor.StartActivityInput> activityInput;
5358

5459
@Before
5560
public void setUp() {
@@ -58,6 +63,13 @@ public void setUp() {
5863
when(workflowClient.getOptions()).thenReturn(clientOptions);
5964
when(clientOptions.getNamespace()).thenReturn(NAMESPACE);
6065
when(clientOptions.getIdentity()).thenReturn("test-identity");
66+
Payload propagatedPayload = Payload.newBuilder().build();
67+
ContextPropagator contextPropagator = mock(ContextPropagator.class);
68+
when(contextPropagator.getCurrentContext()).thenReturn("test-context");
69+
when(contextPropagator.serializeContext("test-context"))
70+
.thenReturn(Collections.singletonMap("propagated-key", propagatedPayload));
71+
when(clientOptions.getContextPropagators())
72+
.thenReturn(Collections.singletonList(contextPropagator));
6173
when(workflowClient.getWorkflowServiceStubs()).thenReturn(mock(WorkflowServiceStubs.class));
6274

6375
WorkflowClientInternal workflowClientInternal = mock(WorkflowClientInternal.class);
@@ -100,12 +112,14 @@ public void setUp() {
100112
ActivityClient activityClient =
101113
mock(ActivityClient.class, withSettings().extraInterfaces(ActivityClientInternal.class));
102114
ActivityClientCallsInterceptor activityInvoker = mock(ActivityClientCallsInterceptor.class);
115+
activityInput = new AtomicReference<>();
103116
when(((ActivityClientInternal) activityClient).getInvoker()).thenReturn(activityInvoker);
104117
when(activityInvoker.startActivity(
105118
org.mockito.ArgumentMatchers.any(
106119
ActivityClientCallsInterceptor.StartActivityInput.class)))
107120
.thenAnswer(
108121
invocation -> {
122+
activityInput.set(invocation.getArgument(0));
109123
CurrentNexusOperationContext.get().getNexusOperationMetadata().operationToken =
110124
"activity-operation-token";
111125
ActivityClientCallsInterceptor.StartActivityInput input = invocation.getArgument(0);
@@ -130,6 +144,21 @@ public void tearDown() {
130144

131145
// ---------- Activity double-start ----------
132146

147+
@Test
148+
public void startActivity_propagatesWorkflowClientContext() {
149+
StartActivityOptions options =
150+
StartActivityOptions.newBuilder()
151+
.setId("act-context")
152+
.setTaskQueue(TASK_QUEUE)
153+
.setStartToCloseTimeout(Duration.ofSeconds(10))
154+
.build();
155+
156+
client.startActivity(TestActivity.class, TestActivity::doSomething, options);
157+
158+
Assert.assertNotNull(activityInput.get());
159+
Assert.assertTrue(activityInput.get().getHeader().getValues().containsKey("propagated-key"));
160+
}
161+
133162
@Test
134163
public void doubleStartActivity_secondCallThrowsBadRequest() {
135164
StartActivityOptions options =

0 commit comments

Comments
 (0)