|
3 | 3 | import io.kurrent.dbclient.*; |
4 | 4 | import com.fasterxml.jackson.databind.JsonNode; |
5 | 5 | import com.fasterxml.jackson.databind.node.ObjectNode; |
| 6 | +import io.kurrentdb.protocol.streams.v2.AppendStreamFailure.ErrorCase; |
| 7 | +import io.opentelemetry.api.common.AttributeKey; |
6 | 8 | import io.opentelemetry.api.trace.SpanContext; |
| 9 | +import io.opentelemetry.api.trace.SpanKind; |
| 10 | +import io.opentelemetry.api.trace.StatusCode; |
7 | 11 | import io.opentelemetry.sdk.trace.ReadableSpan; |
8 | 12 | import org.junit.jupiter.api.Assertions; |
| 13 | +import org.junit.jupiter.api.Assumptions; |
9 | 14 | import org.junit.jupiter.api.Test; |
10 | 15 | import org.junit.jupiter.api.Timeout; |
11 | 16 |
|
12 | | -import java.util.List; |
13 | | -import java.util.UUID; |
| 17 | +import java.util.*; |
14 | 18 | import java.util.concurrent.CountDownLatch; |
15 | 19 | import java.util.concurrent.ExecutionException; |
16 | 20 | import java.util.concurrent.TimeUnit; |
@@ -289,4 +293,120 @@ public void onEvent(Subscription subscription, ResolvedEvent event) { |
289 | 293 | List<ReadableSpan> subscribeSpans = getSpansForOperation(ClientTelemetryConstants.Operations.SUBSCRIBE); |
290 | 294 | Assertions.assertTrue(subscribeSpans.isEmpty(), "No spans should be recorded for deleted events"); |
291 | 295 | } |
| 296 | + |
| 297 | + @Test |
| 298 | + default void testMultiStreamAppendIsInstrumentedWithTracingAsExpected() throws Throwable { |
| 299 | + KurrentDBClient client = getDefaultClient(); |
| 300 | + |
| 301 | + Optional<ServerVersion> version = client.getServerVersion().get(); |
| 302 | + |
| 303 | + Assumptions.assumeTrue( |
| 304 | + version.isPresent() && version.get().isGreaterOrEqualThan(25, 0), |
| 305 | + "Multi-stream append is not supported server versions below 25.0.0" |
| 306 | + ); |
| 307 | + |
| 308 | + String streamName1 = generateName(); |
| 309 | + String streamName2 = generateName(); |
| 310 | + |
| 311 | + EventData event1 = EventData.builderAsJson("TestEvent", mapper.writeValueAsBytes(new Foo())) |
| 312 | + .eventId(UUID.randomUUID()) |
| 313 | + .build(); |
| 314 | + |
| 315 | + EventData event2 = EventData.builderAsJson("TestEvent", mapper.writeValueAsBytes(new Foo())) |
| 316 | + .eventId(UUID.randomUUID()) |
| 317 | + .build(); |
| 318 | + |
| 319 | + AppendStreamRequest request1 = new AppendStreamRequest( |
| 320 | + streamName1, |
| 321 | + Collections.singletonList(event1).iterator(), |
| 322 | + StreamState.noStream() |
| 323 | + ); |
| 324 | + |
| 325 | + AppendStreamRequest request2 = new AppendStreamRequest( |
| 326 | + streamName2, |
| 327 | + Collections.singletonList(event2).iterator(), |
| 328 | + StreamState.noStream() |
| 329 | + ); |
| 330 | + |
| 331 | + MultiAppendWriteResult result = client.multiStreamAppend( |
| 332 | + Arrays.asList(request1, request2).iterator() |
| 333 | + ).get(); |
| 334 | + |
| 335 | + Assertions.assertNotNull(result); |
| 336 | + Assertions.assertTrue(result.getSuccesses().isPresent()); |
| 337 | + |
| 338 | + List<ReadableSpan> spans = getSpansForOperation(ClientTelemetryConstants.Operations.MULTI_APPEND); |
| 339 | + Assertions.assertEquals(1, spans.size()); |
| 340 | + |
| 341 | + assertSpanAttributeEquals(spans.get(0), ClientTelemetryAttributes.Database.SYSTEM, ClientTelemetryConstants.INSTRUMENTATION_NAME); |
| 342 | + assertSpanAttributeEquals(spans.get(0), ClientTelemetryAttributes.Database.OPERATION, ClientTelemetryConstants.Operations.MULTI_APPEND); |
| 343 | + assertSpanAttributeEquals(spans.get(0), ClientTelemetryAttributes.Database.USER, "admin"); |
| 344 | + Assertions.assertEquals(StatusCode.OK, spans.get(0).toSpanData().getStatus().getStatusCode()); |
| 345 | + Assertions.assertEquals(SpanKind.CLIENT, spans.get(0).getKind()); |
| 346 | + } |
| 347 | + |
| 348 | + @Test |
| 349 | + default void testMultiStreamAppendIsInstrumentedWithFailures() throws Throwable { |
| 350 | + KurrentDBClient client = getDefaultClient(); |
| 351 | + |
| 352 | + Optional<ServerVersion> version = client.getServerVersion().get(); |
| 353 | + |
| 354 | + Assumptions.assumeTrue( |
| 355 | + version.isPresent() && version.get().isGreaterOrEqualThan(25, 0), |
| 356 | + "Multi-stream append is not supported server versions below 25.0.0" |
| 357 | + ); |
| 358 | + |
| 359 | + String streamName1 = generateName(); |
| 360 | + String streamName2 = generateName(); |
| 361 | + |
| 362 | + EventData event1 = EventData.builderAsJson("TestEvent", mapper.writeValueAsBytes(new Foo())) |
| 363 | + .eventId(UUID.randomUUID()) |
| 364 | + .build(); |
| 365 | + |
| 366 | + EventData event2 = EventData.builderAsJson("TestEvent", mapper.writeValueAsBytes(new Foo())) |
| 367 | + .eventId(UUID.randomUUID()) |
| 368 | + .build(); |
| 369 | + |
| 370 | + AppendStreamRequest request1 = new AppendStreamRequest( |
| 371 | + streamName1, |
| 372 | + Collections.singletonList(event1).iterator(), |
| 373 | + StreamState.noStream() |
| 374 | + ); |
| 375 | + |
| 376 | + AppendStreamRequest request2 = new AppendStreamRequest( |
| 377 | + streamName2, |
| 378 | + Collections.singletonList(event2).iterator(), |
| 379 | + StreamState.streamExists() |
| 380 | + ); |
| 381 | + |
| 382 | + MultiAppendWriteResult result = client.multiStreamAppend( |
| 383 | + Arrays.asList(request1, request2).iterator() |
| 384 | + ).get(); |
| 385 | + |
| 386 | + Assertions.assertNotNull(result); |
| 387 | + Assertions.assertFalse(result.getSuccesses().isPresent()); |
| 388 | + Assertions.assertTrue(result.getFailures().isPresent()); |
| 389 | + |
| 390 | + List<ReadableSpan> spans = getSpansForOperation(ClientTelemetryConstants.Operations.MULTI_APPEND); |
| 391 | + Assertions.assertEquals(1, spans.size()); |
| 392 | + |
| 393 | + ReadableSpan span = spans.get(0); |
| 394 | + |
| 395 | + Assertions.assertEquals(StatusCode.ERROR, span.toSpanData().getStatus().getStatusCode()); |
| 396 | + Assertions.assertEquals("", span.toSpanData().getStatus().getDescription()); |
| 397 | + |
| 398 | + List<io.opentelemetry.sdk.trace.data.EventData> events = span.toSpanData().getEvents(); |
| 399 | + |
| 400 | + Assertions.assertEquals(1, events.size()); |
| 401 | + |
| 402 | + io.opentelemetry.sdk.trace.data.EventData failureEvent = events.get(0); |
| 403 | + |
| 404 | + Assertions.assertEquals("exception", failureEvent.getName()); |
| 405 | + |
| 406 | + Assertions.assertEquals(ErrorCase.STREAM_REVISION_CONFLICT.toString(), |
| 407 | + failureEvent.getAttributes().get(AttributeKey.stringKey("exception.type"))); |
| 408 | + |
| 409 | + Assertions.assertNotNull(failureEvent.getAttributes().get(AttributeKey.longKey("exception.revision"))); |
| 410 | + Assertions.assertEquals(-1L, failureEvent.getAttributes().get(AttributeKey.longKey("exception.revision"))); |
| 411 | + } |
292 | 412 | } |
0 commit comments