diff --git a/implementation/revapi.json b/implementation/revapi.json index d37de1ac5..517f49f09 100644 --- a/implementation/revapi.json +++ b/implementation/revapi.json @@ -51,7 +51,63 @@ "criticality" : "highlight", "minSeverity" : "POTENTIALLY_BREAKING", "minCriticality" : "documented", - "differences" : [ ] + "differences" : [ + { + "ignore": true, + "code": "java.method.parameterTypeChanged", + "old": "parameter boolean io.smallrye.mutiny.Context::contains(===java.lang.String===)", + "new": "parameter boolean io.smallrye.mutiny.Context::contains(===java.lang.Object===)", + "parameterIndex": "0", + "justification": "ADD YOUR EXPLANATION FOR THE NECESSITY OF THIS CHANGE" + }, + { + "ignore": true, + "code": "java.method.parameterTypeChanged", + "old": "parameter io.smallrye.mutiny.Context io.smallrye.mutiny.Context::delete(===java.lang.String===)", + "new": "parameter io.smallrye.mutiny.Context io.smallrye.mutiny.Context::delete(===java.lang.Object===)", + "parameterIndex": "0", + "justification": "ADD YOUR EXPLANATION FOR THE NECESSITY OF THIS CHANGE" + }, + { + "ignore": true, + "code": "java.method.parameterTypeParameterChanged", + "old": "parameter io.smallrye.mutiny.Context io.smallrye.mutiny.Context::from(===java.util.Map===)", + "new": "parameter io.smallrye.mutiny.Context io.smallrye.mutiny.Context::from(===java.util.Map===)", + "parameterIndex": "0", + "justification": "ADD YOUR EXPLANATION FOR THE NECESSITY OF THIS CHANGE" + }, + { + "ignore": true, + "code": "java.method.parameterTypeChanged", + "old": "parameter T io.smallrye.mutiny.Context::get(===java.lang.String===) throws java.util.NoSuchElementException", + "new": "parameter T io.smallrye.mutiny.Context::get(===java.lang.Object===) throws java.util.NoSuchElementException", + "parameterIndex": "0", + "justification": "ADD YOUR EXPLANATION FOR THE NECESSITY OF THIS CHANGE" + }, + { + "ignore": true, + "code": "java.method.parameterTypeChanged", + "old": "parameter T io.smallrye.mutiny.Context::getOrElse(===java.lang.String===, java.util.function.Supplier)", + "new": "parameter T io.smallrye.mutiny.Context::getOrElse(===java.lang.Object===, java.util.function.Supplier)", + "parameterIndex": "0", + "justification": "ADD YOUR EXPLANATION FOR THE NECESSITY OF THIS CHANGE" + }, + { + "ignore": true, + "code": "java.method.returnTypeTypeParametersChanged", + "old": "method java.util.Set io.smallrye.mutiny.Context::keys()", + "new": "method java.util.Set io.smallrye.mutiny.Context::keys()", + "justification": "ADD YOUR EXPLANATION FOR THE NECESSITY OF THIS CHANGE" + }, + { + "ignore": true, + "code": "java.method.parameterTypeChanged", + "old": "parameter io.smallrye.mutiny.Context io.smallrye.mutiny.Context::put(===java.lang.String===, java.lang.Object)", + "new": "parameter io.smallrye.mutiny.Context io.smallrye.mutiny.Context::put(===java.lang.Object===, java.lang.Object)", + "parameterIndex": "0", + "justification": "ADD YOUR EXPLANATION FOR THE NECESSITY OF THIS CHANGE" + } + ] } }, { "extension" : "revapi.reporter.json", diff --git a/implementation/src/main/java/io/smallrye/mutiny/Context.java b/implementation/src/main/java/io/smallrye/mutiny/Context.java index 4ede85360..f6387078d 100644 --- a/implementation/src/main/java/io/smallrye/mutiny/Context.java +++ b/implementation/src/main/java/io/smallrye/mutiny/Context.java @@ -16,8 +16,6 @@ * A context is provided by a {@link io.smallrye.mutiny.subscription.UniSubscriber} or {@link Flow.Subscriber} * that implements {@link io.smallrye.mutiny.subscription.ContextSupport}. *

- * Context keys are represented as {@link String} while values can be from heterogeneous types. - *

* {@link Context} instances are thread-safe. * Internal storage is not allocated until the first entry is being added. *

@@ -57,9 +55,9 @@ public static Context of(Object... entries) { if (entries.length % 2 != 0) { throw new IllegalArgumentException("Arguments must be balanced to form (key, value) pairs"); } - HashMap map = new HashMap<>(); + HashMap map = new HashMap<>(); for (int i = 0; i < entries.length; i = i + 2) { - String key = nonNull(entries[i], "key").toString(); + Object key = nonNull(entries[i], "key"); Object value = nonNull(entries[i + 1], "value"); map.put(key, value); } @@ -73,17 +71,17 @@ public static Context of(Object... entries) { * @return the new context * @throws NullPointerException when {@code entries} is null */ - public static Context from(Map entries) { + public static Context from(Map entries) { return new Context(requireNonNull(entries, "The entries map cannot be null")); } - private volatile ConcurrentHashMap entries; + private volatile ConcurrentHashMap entries; private Context() { this.entries = null; } - private Context(Map initialEntries) { + private Context(Map initialEntries) { this.entries = new ConcurrentHashMap<>(initialEntries); } @@ -93,7 +91,7 @@ private Context(Map initialEntries) { * @param key the key * @return {@code true} when there is an entry for {@code key}, {@code false} otherwise */ - public boolean contains(String key) { + public boolean contains(Object key) { if (entries == null) { return false; } else { @@ -110,7 +108,7 @@ public boolean contains(String key) { * @throws NoSuchElementException when there is no entry for {@code key} */ @SuppressWarnings("unchecked") - public T get(String key) throws NoSuchElementException { + public T get(Object key) throws NoSuchElementException { if (entries == null) { throw new NoSuchElementException("The context is empty"); } @@ -130,7 +128,7 @@ public T get(String key) throws NoSuchElementException { * @return the value */ @SuppressWarnings("unchecked") - public T getOrElse(String key, Supplier alternativeSupplier) { + public T getOrElse(Object key, Supplier alternativeSupplier) { if (entries != null) { T value = (T) entries.get(key); if (value != null) { @@ -147,7 +145,7 @@ public T getOrElse(String key, Supplier alternativeSupplier) { * @param value the value, cannot be {@code null} * @return this context */ - public Context put(String key, Object value) { + public Context put(Object key, Object value) { if (entries == null) { synchronized (this) { if (entries == null) { @@ -165,7 +163,7 @@ public Context put(String key, Object value) { * @param key the key * @return this context */ - public Context delete(String key) { + public Context delete(Object key) { if (entries != null) { entries.remove(key); } @@ -188,18 +186,33 @@ public boolean isEmpty() { * * @return the set of keys */ - public Set keys() { + public Set keys() { if (this.entries == null) { return Collections.emptySet(); } - HashSet set = new HashSet<>(); - Enumeration enumeration = entries.keys(); + HashSet set = new HashSet<>(); + Enumeration enumeration = entries.keys(); while (enumeration.hasMoreElements()) { set.add(enumeration.nextElement()); } return set; } + /** + * Fork a context from the current context. + * + * @return a new context with the current context data as initial values + */ + public Context fork() { + synchronized (this) { + if (this.isEmpty()) { + return Context.empty(); + } else { + return Context.from(this.entries); + } + } + } + @Override public boolean equals(Object other) { if (this == other) { diff --git a/implementation/src/main/java/io/smallrye/mutiny/Multi.java b/implementation/src/main/java/io/smallrye/mutiny/Multi.java index 1e52f570c..9c974b96d 100644 --- a/implementation/src/main/java/io/smallrye/mutiny/Multi.java +++ b/implementation/src/main/java/io/smallrye/mutiny/Multi.java @@ -739,4 +739,14 @@ default Multi capDemandsTo(long max) { default > MultiSplitter split(Class keyType, Function splitter) { return new MultiSplitter<>(this, keyType, splitter); } + + @CheckReturnValue + default Multi forkContext() { + return forkContext(newContext -> {}); + } + + @CheckReturnValue + default Multi forkContext(Consumer additionalSteps) { + throw new UnsupportedOperationException("Default method added to limit binary incompatibility"); + } } diff --git a/implementation/src/main/java/io/smallrye/mutiny/Uni.java b/implementation/src/main/java/io/smallrye/mutiny/Uni.java index 6d561bf3d..84a9fd5bb 100644 --- a/implementation/src/main/java/io/smallrye/mutiny/Uni.java +++ b/implementation/src/main/java/io/smallrye/mutiny/Uni.java @@ -855,4 +855,14 @@ default Uni withContext(BiFunction, Context, Uni> builder) { default Uni> attachContext() { return this.withContext((uni, ctx) -> uni.onItem().transform(item -> new ItemWithContext<>(ctx, item))); } + + @CheckReturnValue + default Uni forkContext() { + return forkContext(newContext -> {}); + } + + @CheckReturnValue + default Uni forkContext(Consumer additionalSteps) { + throw new UnsupportedOperationException("Default method added to limit binary incompatibility"); + } } diff --git a/implementation/src/main/java/io/smallrye/mutiny/operators/AbstractMulti.java b/implementation/src/main/java/io/smallrye/mutiny/operators/AbstractMulti.java index 90ae45f19..fc749d215 100644 --- a/implementation/src/main/java/io/smallrye/mutiny/operators/AbstractMulti.java +++ b/implementation/src/main/java/io/smallrye/mutiny/operators/AbstractMulti.java @@ -6,6 +6,7 @@ import java.util.concurrent.Executor; import java.util.concurrent.Flow; import java.util.function.BiFunction; +import java.util.function.Consumer; import java.util.function.LongFunction; import java.util.function.Predicate; @@ -32,12 +33,7 @@ import io.smallrye.mutiny.groups.MultiSubscribe; import io.smallrye.mutiny.helpers.ParameterValidation; import io.smallrye.mutiny.infrastructure.Infrastructure; -import io.smallrye.mutiny.operators.multi.MultiCacheOp; -import io.smallrye.mutiny.operators.multi.MultiDemandCapping; -import io.smallrye.mutiny.operators.multi.MultiEmitOnOp; -import io.smallrye.mutiny.operators.multi.MultiLogger; -import io.smallrye.mutiny.operators.multi.MultiSubscribeOnOp; -import io.smallrye.mutiny.operators.multi.MultiWithContext; +import io.smallrye.mutiny.operators.multi.*; import io.smallrye.mutiny.operators.multi.processors.BroadcastProcessor; import io.smallrye.mutiny.subscription.MultiSubscriber; import io.smallrye.mutiny.subscription.MultiSubscriberAdapter; @@ -216,4 +212,9 @@ public MultiDemandPausing pauseDemand() { public Multi capDemandsUsing(LongFunction function) { return Infrastructure.onMultiCreation(new MultiDemandCapping<>(this, nonNull(function, "function"))); } + + @Override + public Multi forkContext(Consumer additionalSteps) { + return Infrastructure.onMultiCreation(new MultiContextForkingOperator<>(this, nonNull(additionalSteps, "additionalSteps"))); + } } diff --git a/implementation/src/main/java/io/smallrye/mutiny/operators/AbstractUni.java b/implementation/src/main/java/io/smallrye/mutiny/operators/AbstractUni.java index a616f000a..be37d5d24 100644 --- a/implementation/src/main/java/io/smallrye/mutiny/operators/AbstractUni.java +++ b/implementation/src/main/java/io/smallrye/mutiny/operators/AbstractUni.java @@ -4,6 +4,7 @@ import java.util.concurrent.Executor; import java.util.function.BiFunction; +import java.util.function.Consumer; import java.util.function.Predicate; import io.smallrye.mutiny.Context; @@ -148,4 +149,9 @@ public Uni log() { public Uni withContext(BiFunction, Context, Uni> builder) { return Infrastructure.onUniCreation(new UniWithContext<>(this, nonNull(builder, "builder"))); } + + @Override + public Uni forkContext(Consumer additionalSteps) { + return Infrastructure.onUniCreation(new UniContextForkingOperator<>(this, nonNull(additionalSteps, "additionalSteps"))); + } } diff --git a/implementation/src/main/java/io/smallrye/mutiny/operators/multi/MultiContextForkingOperator.java b/implementation/src/main/java/io/smallrye/mutiny/operators/multi/MultiContextForkingOperator.java new file mode 100644 index 000000000..d585a441e --- /dev/null +++ b/implementation/src/main/java/io/smallrye/mutiny/operators/multi/MultiContextForkingOperator.java @@ -0,0 +1,43 @@ +package io.smallrye.mutiny.operators.multi; + +import io.smallrye.mutiny.Context; +import io.smallrye.mutiny.Multi; +import io.smallrye.mutiny.subscription.ContextSupport; +import io.smallrye.mutiny.subscription.MultiSubscriber; + +import java.util.function.Consumer; + +public final class MultiContextForkingOperator extends AbstractMultiOperator { + + private final Consumer additionalSteps; + + public MultiContextForkingOperator(Multi upstream, Consumer additionalSteps) { + super(upstream); + this.additionalSteps = additionalSteps; + } + + @Override + public void subscribe(MultiSubscriber subscriber) { + upstream.subscribe().withSubscriber(new MultiContextForkingOperatorProcessor(subscriber, additionalSteps)); + } + + private static class MultiContextForkingOperatorProcessor extends MultiOperatorProcessor { + + private final Context forkedContext; + + public MultiContextForkingOperatorProcessor(MultiSubscriber downstream, Consumer additionalSteps) { + super(downstream); + if (downstream instanceof ContextSupport provider) { + this.forkedContext = provider.context().fork(); + } else { + this.forkedContext = Context.empty(); + } + additionalSteps.accept(this.forkedContext); + } + + @Override + public Context context() { + return forkedContext; + } + } +} diff --git a/implementation/src/main/java/io/smallrye/mutiny/operators/multi/MultiRetryWhenOp.java b/implementation/src/main/java/io/smallrye/mutiny/operators/multi/MultiRetryWhenOp.java index d94f51ccb..21c4e5164 100644 --- a/implementation/src/main/java/io/smallrye/mutiny/operators/multi/MultiRetryWhenOp.java +++ b/implementation/src/main/java/io/smallrye/mutiny/operators/multi/MultiRetryWhenOp.java @@ -41,7 +41,14 @@ public MultiRetryWhenOp(Multi upstream, Predicate void subscribe(MultiSubscriber downstream, Predicate onFailurePredicate, Function, ? extends Publisher> triggerStreamFactory, Multi upstream) { - TriggerSubscriber other = new TriggerSubscriber(); + Context context; + if (downstream instanceof ContextSupport provider) { + context = provider.context(); + } else { + context = Context.empty(); + } + + TriggerSubscriber other = new TriggerSubscriber(context); Subscriber signaller = new SerializedSubscriber<>(other.processor); signaller.onSubscribe(Subscriptions.empty()); MultiSubscriber serialized = new SerializedSubscriber<>(downstream); @@ -176,7 +183,11 @@ static final class TriggerSubscriber extends AbstractMulti implements Multi, Subscriber, ContextSupport { RetryWhenOperator operator; private final Flow.Processor processor = UnicastProcessor. create().serialized(); - private Context context; + private final Context context; + + TriggerSubscriber(Context context) { + this.context = context; + } @Override public void onSubscribe(Flow.Subscription s) { @@ -200,11 +211,6 @@ public void onComplete() { @Override public void subscribe(Subscriber actual) { - if (actual instanceof ContextSupport) { - this.context = ((ContextSupport) actual).context(); - } else { - this.context = Context.empty(); - } processor.subscribe(actual); } diff --git a/implementation/src/main/java/io/smallrye/mutiny/operators/uni/UniAndCombination.java b/implementation/src/main/java/io/smallrye/mutiny/operators/uni/UniAndCombination.java index 8bf85ac0d..bfefed62c 100644 --- a/implementation/src/main/java/io/smallrye/mutiny/operators/uni/UniAndCombination.java +++ b/implementation/src/main/java/io/smallrye/mutiny/operators/uni/UniAndCombination.java @@ -67,7 +67,7 @@ private class AndSupervisor implements UniSubscription { Context context = subscriber.context(); for (Uni uni : unis) { - UniHandler result = new UniHandler(this, uni, context); + UniHandler result = new UniHandler(this, uni, context.fork()); handlers.add(result); } diff --git a/implementation/src/main/java/io/smallrye/mutiny/operators/uni/UniContextForkingOperator.java b/implementation/src/main/java/io/smallrye/mutiny/operators/uni/UniContextForkingOperator.java new file mode 100644 index 000000000..b083e47aa --- /dev/null +++ b/implementation/src/main/java/io/smallrye/mutiny/operators/uni/UniContextForkingOperator.java @@ -0,0 +1,40 @@ +package io.smallrye.mutiny.operators.uni; + +import io.smallrye.mutiny.Context; +import io.smallrye.mutiny.Uni; +import io.smallrye.mutiny.operators.AbstractUni; +import io.smallrye.mutiny.operators.UniOperator; +import io.smallrye.mutiny.subscription.UniSubscriber; + +import java.util.function.Consumer; + +public final class UniContextForkingOperator extends UniOperator { + + private final Consumer additionalSteps; + + public UniContextForkingOperator(Uni upstream, Consumer additionalSteps) { + super(upstream); + this.additionalSteps = additionalSteps; + } + + @Override + public void subscribe(UniSubscriber subscriber) { + AbstractUni.subscribe(upstream(), new UniContextForkingOperatorProcessor<>(additionalSteps, subscriber)); + } + + private static class UniContextForkingOperatorProcessor extends UniOperatorProcessor { + + private final Context forkedContext; + + public UniContextForkingOperatorProcessor(Consumer additionalSteps, UniSubscriber downstream) { + super(downstream); + this.forkedContext = downstream.context().fork(); + additionalSteps.accept(this.forkedContext); + } + + @Override + public Context context() { + return forkedContext; + } + } +} diff --git a/implementation/src/main/java/io/smallrye/mutiny/operators/uni/UniOrCombination.java b/implementation/src/main/java/io/smallrye/mutiny/operators/uni/UniOrCombination.java index 15ebd83fb..8a987627c 100644 --- a/implementation/src/main/java/io/smallrye/mutiny/operators/uni/UniOrCombination.java +++ b/implementation/src/main/java/io/smallrye/mutiny/operators/uni/UniOrCombination.java @@ -45,7 +45,7 @@ public void subscribe(UniSubscriber subscriber) { List> futures = new ArrayList<>(); challengers.forEach(uni -> { - CompletableFuture future = uni.subscribe().asCompletionStage(subscriber.context()); + CompletableFuture future = uni.subscribe().asCompletionStage(subscriber.context().fork()); futures.add(future); }); diff --git a/implementation/src/main/java/io/smallrye/mutiny/operators/uni/builders/UniJoinAll.java b/implementation/src/main/java/io/smallrye/mutiny/operators/uni/builders/UniJoinAll.java index 4bf579d61..206b9e671 100644 --- a/implementation/src/main/java/io/smallrye/mutiny/operators/uni/builders/UniJoinAll.java +++ b/implementation/src/main/java/io/smallrye/mutiny/operators/uni/builders/UniJoinAll.java @@ -74,7 +74,7 @@ private boolean trySubscribe(int index, Uni uni) { if (proceed) { uni.onSubscription() .invoke(subscription -> this.onSubscribe(index, subscription)) - .subscribe().with(subscriber.context(), item -> this.onItem(index, item), this::onFailure); + .subscribe().with(subscriber.context().fork(), item -> this.onItem(index, item), this::onFailure); } return proceed; } diff --git a/implementation/src/main/java/io/smallrye/mutiny/operators/uni/builders/UniJoinFirst.java b/implementation/src/main/java/io/smallrye/mutiny/operators/uni/builders/UniJoinFirst.java index 8348743ce..2d636699d 100644 --- a/implementation/src/main/java/io/smallrye/mutiny/operators/uni/builders/UniJoinFirst.java +++ b/implementation/src/main/java/io/smallrye/mutiny/operators/uni/builders/UniJoinFirst.java @@ -72,7 +72,7 @@ private boolean trySubscribe(int index) { Uni uni = unis.get(index); uni.onSubscription() .invoke(subscription -> this.onSubscribe(index, subscription)) - .subscribe().with(subscriber.context(), this::onItem, this::onFailure); + .subscribe().with(subscriber.context().fork(), this::onItem, this::onFailure); } return proceed; } diff --git a/implementation/src/test/java/io/smallrye/mutiny/ContextTest.java b/implementation/src/test/java/io/smallrye/mutiny/ContextTest.java index 7517c50c2..ae26614e1 100644 --- a/implementation/src/test/java/io/smallrye/mutiny/ContextTest.java +++ b/implementation/src/test/java/io/smallrye/mutiny/ContextTest.java @@ -66,7 +66,7 @@ void unbalancedOf() { @Test void from() { - HashMap map = new HashMap() { + HashMap map = new HashMap<>() { { put("foo", "bar"); put("abc", "def"); @@ -162,11 +162,27 @@ void getOnMissingKey() { @Test void keysetIsACopy() { Context context = Context.of("foo", "bar", "123", 456); - Set k1 = context.keys(); + Set k1 = context.keys(); context.put("bar", "baz"); - Set k2 = context.keys(); + Set k2 = context.keys(); assertThat(k1).isNotSameAs(k2); } + + @Test + void emptyFork() { + Context empty = Context.empty(); + Context fork = empty.fork(); + assertThat(empty.isEmpty()); + assertThat(fork.isEmpty()); + } + + @Test + void fork() { + Context root = Context.of("foo", "bar"); + Context fork = root.fork().put("foo", "yolo"); + assertThat(root. get("foo")).isEqualTo("bar"); + assertThat(fork. get("foo")).isEqualTo("yolo"); + } } @Nested @@ -309,100 +325,99 @@ void uniOperatorSubclassesPropagateContext() { @Test void joinAllAndAttachContext() { - Context context = Context.of("foo", "bar", "baz", "baz"); + Context context = Context.of(58, "58", 63, "63", 69, "69"); Uni a = Uni.createFrom().item(58) - .withContext((uni, ctx) -> uni.onItem().invoke(n -> ctx.put(n.toString(), n))); + .withContext((uni, ctx) -> uni.onItem().invoke(n -> ctx.put("!!!", "!!!"))); Uni b = Uni.createFrom().item(63) - .withContext((uni, ctx) -> uni.onItem().invoke(n -> ctx.put(n.toString(), n))); + .withContext((uni, ctx) -> uni.onItem().invoke(n -> ctx.put("!!!", "!!!"))); Uni c = Uni.createFrom().item(69) - .withContext((uni, ctx) -> uni.onItem().invoke(n -> ctx.put(n.toString(), n))); + .withContext((uni, ctx) -> uni.onItem().invoke(n -> ctx.put("!!!", "!!!"))); UniAssertSubscriber sub = Uni.join().all(a, b, c).andFailFast() .attachContext() .onItem().transform(itemsWithContext -> { + Context ctx = itemsWithContext.context(); - return itemsWithContext.get().toString() + "::" + ctx.get("58") + "::" + ctx.get("63") + "::" - + ctx.get("69"); + return itemsWithContext.get().toString() + "::" + ctx.getOrElse("!!!", () -> "~"); }) .subscribe().withSubscriber(UniAssertSubscriber.create(context)); - sub.assertCompleted().assertItem("[58, 63, 69]::58::63::69"); + sub.assertCompleted().assertItem("[58, 63, 69]::~"); } @Test void joinFirstAndAttachContext() { - Context context = Context.of("foo", "bar", "baz", "baz"); + Context context = Context.of(58, "58", 63, "63", 69, "69"); Uni a = Uni.createFrom().item(58) - .withContext((uni, ctx) -> uni.onItem().invoke(n -> ctx.put(n.toString(), n))); + .withContext((uni, ctx) -> uni.onItem().invoke(n -> ctx.put("!!!", "!!!"))); Uni b = Uni.createFrom().item(63) - .withContext((uni, ctx) -> uni.onItem().invoke(n -> ctx.put(n.toString(), n))); + .withContext((uni, ctx) -> uni.onItem().invoke(n -> ctx.put("!!!", "!!!"))); Uni c = Uni.createFrom().item(69) - .withContext((uni, ctx) -> uni.onItem().invoke(n -> ctx.put(n.toString(), n))); + .withContext((uni, ctx) -> uni.onItem().invoke(n -> ctx.put("!!!", "!!!"))); UniAssertSubscriber sub = Uni.join().first(a, b, c).withItem() .attachContext() .onItem().transform(item -> { Context ctx = item.context(); - return item.get().toString() + "::" + ctx.get("58"); + return item.get().toString() + "::" + ctx.getOrElse("!!!", () -> "~"); }) .subscribe().withSubscriber(UniAssertSubscriber.create(context)); - sub.assertCompleted().assertItem("58::58"); + sub.assertCompleted().assertItem("58::~"); } @Test void combineAllAndAttachContext() { - Context context = Context.of("foo", "bar", "baz", "baz"); + Context context = Context.of(58, "58", 63, "63", 69, "69"); Uni a = Uni.createFrom().item(58) - .withContext((uni, ctx) -> uni.onItem().invoke(n -> ctx.put(n.toString(), n))); + .withContext((uni, ctx) -> uni.onItem().invoke(n -> ctx.put("!!!", "!!!"))); Uni b = Uni.createFrom().item(63) - .withContext((uni, ctx) -> uni.onItem().invoke(n -> ctx.put(n.toString(), n))); + .withContext((uni, ctx) -> uni.onItem().invoke(n -> ctx.put("!!!", "!!!"))); Uni c = Uni.createFrom().item(69) - .withContext((uni, ctx) -> uni.onItem().invoke(n -> ctx.put(n.toString(), n))); + .withContext((uni, ctx) -> uni.onItem().invoke(n -> ctx.put("!!!", "!!!"))); UniAssertSubscriber sub = Uni.combine().all().unis(a, b, c).asTuple() .attachContext() .onItem().transform(itemsWithContext -> { Context ctx = itemsWithContext.context(); - return itemsWithContext.get().toString() + "::" + ctx.get("58") + "::" + ctx.get("63") + "::" - + ctx.get("69"); + return itemsWithContext.get().toString() + "::" + ctx.getOrElse("!!!", () -> "~"); }) .subscribe().withSubscriber(UniAssertSubscriber.create(context)); - sub.assertCompleted().assertItem("Tuple{item1=58,item2=63,item3=69}::58::63::69"); + sub.assertCompleted().assertItem("Tuple{item1=58,item2=63,item3=69}::~"); } @Test void combineAnyAndAttachContext() { - Context context = Context.of("foo", "bar", "baz", "baz"); + Context context = Context.of(58, "58", 63, "63", 69, "69"); Uni a = Uni.createFrom().item(58) - .withContext((uni, ctx) -> uni.onItem().invoke(n -> ctx.put(n.toString(), n))); + .withContext((uni, ctx) -> uni.onItem().invoke(n -> ctx.put("!!!", "!!!"))); Uni b = Uni.createFrom().item(63) - .withContext((uni, ctx) -> uni.onItem().invoke(n -> ctx.put(n.toString(), n))); + .withContext((uni, ctx) -> uni.onItem().invoke(n -> ctx.put("!!!", "!!!"))); Uni c = Uni.createFrom().item(69) - .withContext((uni, ctx) -> uni.onItem().invoke(n -> ctx.put(n.toString(), n))); + .withContext((uni, ctx) -> uni.onItem().invoke(n -> ctx.put("!!!", "!!!"))); UniAssertSubscriber sub = Uni.combine().any().of(a, b, c) .attachContext() .onItem().transform(item -> { Context ctx = item.context(); - return item.get().toString() + "::" + ctx.get("58"); + return item.get().toString() + "::" + ctx.getOrElse("!!!", () -> "~"); }) .subscribe().withSubscriber(UniAssertSubscriber.create(context)); - sub.assertCompleted().assertItem("58::58"); + sub.assertCompleted().assertItem("58::~"); } @Test diff --git a/implementation/src/test/java/io/smallrye/mutiny/operators/MultiContextForkTest.java b/implementation/src/test/java/io/smallrye/mutiny/operators/MultiContextForkTest.java new file mode 100644 index 000000000..92dfefa70 --- /dev/null +++ b/implementation/src/test/java/io/smallrye/mutiny/operators/MultiContextForkTest.java @@ -0,0 +1,56 @@ +package io.smallrye.mutiny.operators; + +import io.smallrye.mutiny.Context; +import io.smallrye.mutiny.Uni; +import org.junit.jupiter.api.Test; + +import java.time.Duration; +import java.util.ArrayList; +import java.util.List; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +public class MultiContextForkTest { + + @Test + void checkBoundaries() { + ArrayList trace = new ArrayList<>(); + Context rootContext = Context.of(123, "yolo"); + + List result = Uni.createFrom().item(69).toMulti() + .withContext((multi, ctx) -> { + trace.add(ctx.getOrElse(123, () -> "!!!")); + trace.add(ctx.getOrElse(456, () -> "!!!")); + return multi.onItem().invoke(() -> { + trace.add("-" + ctx.getOrElse(123, () -> "!!!")); + trace.add("-" + ctx.getOrElse(456, () -> "!!!")); + }); + }) + .forkContext() + .forkContext(ctx -> ctx.put(456, "bar")) + .withContext((multi, ctx) -> { + trace.add(ctx.getOrElse(123, () -> "!!!")); + return multi.onItem().invoke(() -> { + trace.add("-" + ctx.getOrElse(123, () -> "!!!")); + }); + }) + .forkContext(ctx -> ctx.put(123, "foo")) + .collect().asList().awaitUsing(rootContext).atMost(Duration.ofSeconds(1)); + + assertThat(result).containsExactly(69); + assertThat(trace).containsExactly("foo", "foo", "bar", "-foo", "-bar", "-foo"); + assertThat(rootContext.get(123)).isEqualTo("yolo"); + } + + @Test + void throwingConsumer() { + assertThatThrownBy(() -> + Uni.createFrom().item(69).toMulti() + .forkContext(ctx -> { + throw new RuntimeException("boom"); + }).collect().asList().await().atMost(Duration.ofSeconds(1))) + .isInstanceOf(RuntimeException.class) + .hasMessage("boom"); + } +} diff --git a/implementation/src/test/java/io/smallrye/mutiny/operators/UniContextForkTest.java b/implementation/src/test/java/io/smallrye/mutiny/operators/UniContextForkTest.java new file mode 100644 index 000000000..a284dd10f --- /dev/null +++ b/implementation/src/test/java/io/smallrye/mutiny/operators/UniContextForkTest.java @@ -0,0 +1,55 @@ +package io.smallrye.mutiny.operators; + +import io.smallrye.mutiny.Context; +import io.smallrye.mutiny.Uni; +import org.junit.jupiter.api.Test; + +import java.time.Duration; +import java.util.ArrayList; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +class UniContextForkTest { + + @Test + void checkBoundaries() { + ArrayList trace = new ArrayList<>(); + Context rootContext = Context.of(123, "yolo"); + + Integer result = Uni.createFrom().item(69) + .withContext((uni, ctx) -> { + trace.add(ctx.getOrElse(123, () -> "!!!")); + trace.add(ctx.getOrElse(456, () -> "!!!")); + return uni.onItem().invoke(() -> { + trace.add("-" + ctx.getOrElse(123, () -> "!!!")); + trace.add("-" + ctx.getOrElse(456, () -> "!!!")); + }); + }) + .forkContext() + .forkContext(ctx -> ctx.put(456, "bar")) + .withContext((uni, ctx) -> { + trace.add(ctx.getOrElse(123, () -> "!!!")); + return uni.onItem().invoke(() -> { + trace.add("-" + ctx.getOrElse(123, () -> "!!!")); + }); + }) + .forkContext(ctx -> ctx.put(123, "foo")) + .awaitUsing(rootContext).atMost(Duration.ofSeconds(1)); + + assertThat(result).isEqualTo(69); + assertThat(trace).containsExactly("foo", "foo", "bar", "-foo", "-bar", "-foo"); + assertThat(rootContext.get(123)).isEqualTo("yolo"); + } + + @Test + void throwingConsumer() { + assertThatThrownBy(() -> + Uni.createFrom().item(69) + .forkContext(ctx -> { + throw new RuntimeException("boom"); + }).await().atMost(Duration.ofSeconds(1))) + .isInstanceOf(RuntimeException.class) + .hasMessage("boom"); + } +}