diff --git a/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/AfterAllStateMachine.java b/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/AfterAllStateMachine.java index e50f95a6b4a4..afedf3d0635d 100644 --- a/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/AfterAllStateMachine.java +++ b/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/AfterAllStateMachine.java @@ -28,9 +28,6 @@ * A {@link TriggerStateMachine} that fires and finishes once after all of its sub-triggers have * fired. */ -@SuppressWarnings({ - "nullness" // TODO(https://github.com/apache/beam/issues/20497) -}) public class AfterAllStateMachine extends TriggerStateMachine { private AfterAllStateMachine(List subTriggers) { @@ -117,7 +114,7 @@ public void onFire(TriggerContext context) throws Exception { @Override public String toString() { StringBuilder builder = new StringBuilder("AfterAll.of("); - Joiner.on(", ").appendTo(builder, subTriggers); + Joiner.on(", ").appendTo(builder, subTriggers()); builder.append(")"); return builder.toString(); } diff --git a/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/AfterDelayFromFirstElementStateMachine.java b/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/AfterDelayFromFirstElementStateMachine.java index e8ca0990559b..8c4f0d1c7bea 100644 --- a/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/AfterDelayFromFirstElementStateMachine.java +++ b/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/AfterDelayFromFirstElementStateMachine.java @@ -32,6 +32,7 @@ import org.apache.beam.sdk.transforms.SerializableFunction; import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableList; import org.checkerframework.checker.nullness.qual.Nullable; +import org.checkerframework.dataflow.qual.Pure; import org.joda.time.Duration; import org.joda.time.Instant; import org.joda.time.format.PeriodFormat; @@ -45,9 +46,6 @@ */ // This class should be inlined to subclasses and deleted, simplifying them too // https://github.com/apache/beam/issues/18117 -@SuppressWarnings({ - "nullness" // TODO(https://github.com/apache/beam/issues/20497) -}) public abstract class AfterDelayFromFirstElementStateMachine extends TriggerStateMachine { protected static final List> IDENTITY = ImmutableList.of(); @@ -61,6 +59,7 @@ public abstract class AfterDelayFromFirstElementStateMachine extends TriggerStat private static final PeriodFormatter PERIOD_FORMATTER = PeriodFormat.wordBased(Locale.ENGLISH); /** To complete an implementation, return the desired time from the TriggerContext. */ + @Pure public abstract @Nullable Instant getCurrentTime(TriggerStateMachine.TriggerContext context); /** diff --git a/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/AfterEachStateMachine.java b/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/AfterEachStateMachine.java index e16e859ba66f..edbdccefbe96 100644 --- a/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/AfterEachStateMachine.java +++ b/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/AfterEachStateMachine.java @@ -17,6 +17,7 @@ */ package org.apache.beam.runners.core.triggers; +import static org.apache.beam.sdk.util.Preconditions.checkStateNotNull; import static org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkArgument; import java.util.Arrays; @@ -41,9 +42,6 @@ * Repeatedly.forever(a)}, since the repeated trigger never finishes. * */ -@SuppressWarnings({ - "nullness" // TODO(https://github.com/apache/beam/issues/20497) -}) public class AfterEachStateMachine extends TriggerStateMachine { private AfterEachStateMachine(List subTriggers) { @@ -72,7 +70,7 @@ public static TriggerStateMachine inOrder(IterableA composite trigger is never invoked once all of its subtriggers are finished, so a null + * here means a caller has violated that invariant. + */ + private ExecutableTriggerStateMachine firstUnfinishedSubTrigger(TriggerContext context) { + return checkStateNotNull( + context.trigger().firstUnfinishedSubTrigger(), + "%s invoked after all of its subtriggers finished", + this); + } + private void updateFinishedState(TriggerContext context) { context.trigger().setFinished(context.trigger().firstUnfinishedSubTrigger() == null); } diff --git a/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/AfterFirstStateMachine.java b/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/AfterFirstStateMachine.java index 8d2056540d8e..45888a5c77aa 100644 --- a/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/AfterFirstStateMachine.java +++ b/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/AfterFirstStateMachine.java @@ -28,9 +28,6 @@ * Create a composite {@link TriggerStateMachine} that fires once after at least one of its * sub-triggers have fired. */ -@SuppressWarnings({ - "nullness" // TODO(https://github.com/apache/beam/issues/20497) -}) public class AfterFirstStateMachine extends TriggerStateMachine { AfterFirstStateMachine(List subTriggers) { @@ -114,7 +111,7 @@ public void onFire(TriggerContext context) throws Exception { @Override public String toString() { StringBuilder builder = new StringBuilder("AfterFirst.of("); - Joiner.on(", ").appendTo(builder, subTriggers); + Joiner.on(", ").appendTo(builder, subTriggers()); builder.append(")"); return builder.toString(); diff --git a/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/AfterWatermarkStateMachine.java b/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/AfterWatermarkStateMachine.java index 760b88a58975..5368b62434e1 100644 --- a/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/AfterWatermarkStateMachine.java +++ b/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/AfterWatermarkStateMachine.java @@ -17,6 +17,7 @@ */ package org.apache.beam.runners.core.triggers; +import static org.apache.beam.sdk.util.Preconditions.checkStateNotNull; import static org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkNotNull; import java.util.Objects; @@ -51,9 +52,6 @@ * AfterWatermark.pastEndOfWindow.withEarlyFirings(OnceTrigger)} or {@code * AfterWatermark.pastEndOfWindow.withEarlyFirings(OnceTrigger)}. */ -@SuppressWarnings({ - "nullness" // TODO(https://github.com/apache/beam/issues/20497) -}) public class AfterWatermarkStateMachine { private static final String TO_STRING = "AfterWatermark.pastEndOfWindow()"; @@ -77,7 +75,7 @@ public static class AfterWatermarkEarlyAndLate extends TriggerStateMachine { @SuppressWarnings("unchecked") private AfterWatermarkEarlyAndLate( - TriggerStateMachine earlyTrigger, TriggerStateMachine lateTrigger) { + TriggerStateMachine earlyTrigger, @Nullable TriggerStateMachine lateTrigger) { super( lateTrigger == null ? ImmutableList.of(earlyTrigger) @@ -109,7 +107,11 @@ public void onElement(OnElementContext c) throws Exception { if (!c.trigger().isMerging()) { // If merges can never happen, we just run the unfinished subtrigger - c.trigger().firstUnfinishedSubTrigger().invokeOnElement(c); + checkStateNotNull( + c.trigger().firstUnfinishedSubTrigger(), + "%s invoked after all of its subtriggers finished", + this) + .invokeOnElement(c); } else { // If merges can happen, we run for all subtriggers because they might be // de-activated or re-activated diff --git a/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/DefaultTriggerStateMachine.java b/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/DefaultTriggerStateMachine.java index 116da24ba0a9..fb14b8d726c8 100644 --- a/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/DefaultTriggerStateMachine.java +++ b/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/DefaultTriggerStateMachine.java @@ -24,9 +24,6 @@ * {@link RepeatedlyStateMachine#forever} and {@link AfterWatermarkStateMachine#pastEndOfWindow} for * more details. */ -@SuppressWarnings({ - "nullness" // TODO(https://github.com/apache/beam/issues/20497) -}) public class DefaultTriggerStateMachine extends TriggerStateMachine { private DefaultTriggerStateMachine() { diff --git a/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/ExecutableTriggerStateMachine.java b/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/ExecutableTriggerStateMachine.java index 2d64986f51a0..0b3fe8e30e57 100644 --- a/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/ExecutableTriggerStateMachine.java +++ b/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/ExecutableTriggerStateMachine.java @@ -17,6 +17,7 @@ */ package org.apache.beam.runners.core.triggers; +import static org.apache.beam.sdk.util.Preconditions.checkStateNotNull; import static org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkNotNull; import static org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkState; @@ -29,9 +30,6 @@ * times (both in the same trigger expression and in other trigger expressions), the {@code * ExecutableTrigger} wrapped around them forms a tree (only one occurrence). */ -@SuppressWarnings({ - "nullness" // TODO(https://github.com/apache/beam/issues/20497) -}) public class ExecutableTriggerStateMachine implements Serializable { /** Store the index assigned to this trigger. */ @@ -60,12 +58,10 @@ private ExecutableTriggerStateMachine(TriggerStateMachine trigger, int nextUnuse this.trigger = checkNotNull(trigger, "trigger must not be null"); this.triggerIndex = nextUnusedIndex++; - if (trigger.subTriggers() != null) { - for (TriggerStateMachine subTrigger : trigger.subTriggers()) { - ExecutableTriggerStateMachine subExecutable = create(subTrigger, nextUnusedIndex); - subTriggers.add(subExecutable); - nextUnusedIndex = subExecutable.firstIndexAfterSubtree; - } + for (TriggerStateMachine subTrigger : trigger.subTriggers()) { + ExecutableTriggerStateMachine subExecutable = create(subTrigger, nextUnusedIndex); + subTriggers.add(subExecutable); + nextUnusedIndex = subExecutable.firstIndexAfterSubtree; } firstIndexAfterSubtree = nextUnusedIndex; } @@ -99,18 +95,19 @@ public boolean isCompatible(ExecutableTriggerStateMachine other) { } public ExecutableTriggerStateMachine getSubTriggerContaining(int index) { - checkNotNull(subTriggers); checkState( index > triggerIndex && index < firstIndexAfterSubtree, "Cannot find sub-trigger containing index not in this tree."); ExecutableTriggerStateMachine previous = null; for (ExecutableTriggerStateMachine subTrigger : subTriggers) { if (index < subTrigger.triggerIndex) { - return previous; + break; } previous = subTrigger; } - return previous; + // The bounds check above places the index within this subtree but past this trigger itself, + // so it always falls inside one of the sub-triggers. + return checkStateNotNull(previous, "No sub-trigger of %s contains index %s", trigger, index); } public void invokePrefetchOnElement(TriggerStateMachine.PrefetchContext c) { diff --git a/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/OrFinallyStateMachine.java b/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/OrFinallyStateMachine.java index e7f5f8754174..b1c7d7c74ada 100644 --- a/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/OrFinallyStateMachine.java +++ b/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/OrFinallyStateMachine.java @@ -23,9 +23,6 @@ /** * Executes the {@code actual} trigger until it finishes or until the {@code until} trigger fires. */ -@SuppressWarnings({ - "nullness" // TODO(https://github.com/apache/beam/issues/20497) -}) class OrFinallyStateMachine extends TriggerStateMachine { private static final int ACTUAL = 0; @@ -96,7 +93,7 @@ public void onFire(TriggerStateMachine.TriggerContext context) throws Exception @Override public String toString() { - return String.format("%s.orFinally(%s)", subTriggers.get(ACTUAL), subTriggers.get(UNTIL)); + return String.format("%s.orFinally(%s)", subTriggers().get(ACTUAL), subTriggers().get(UNTIL)); } private void updateFinishedState(TriggerContext c) throws Exception { diff --git a/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/RepeatedlyStateMachine.java b/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/RepeatedlyStateMachine.java index 0dfd2bb3fd6e..0a652f96e8e8 100644 --- a/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/RepeatedlyStateMachine.java +++ b/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/RepeatedlyStateMachine.java @@ -31,9 +31,6 @@ *

{@code Repeatedly.forever(someTrigger)} behaves like an infinite {@code * AfterEach.inOrder(someTrigger, someTrigger, someTrigger, ...)}. */ -@SuppressWarnings({ - "nullness" // TODO(https://github.com/apache/beam/issues/20497) -}) public class RepeatedlyStateMachine extends TriggerStateMachine { private static final int REPEATED = 0; @@ -97,7 +94,7 @@ public void onFire(TriggerContext context) throws Exception { @Override public String toString() { - return String.format("Repeatedly.forever(%s)", subTriggers.get(REPEATED)); + return String.format("Repeatedly.forever(%s)", subTriggers().get(REPEATED)); } private ExecutableTriggerStateMachine getRepeated(TriggerContext context) { diff --git a/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/TriggerStateMachine.java b/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/TriggerStateMachine.java index 8216c674ba80..a92a87bb2f4a 100644 --- a/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/TriggerStateMachine.java +++ b/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/TriggerStateMachine.java @@ -18,6 +18,7 @@ package org.apache.beam.runners.core.triggers; import java.io.Serializable; +import java.util.Collections; import java.util.List; import java.util.Objects; import org.apache.beam.runners.core.MergingStateAccessor; @@ -27,7 +28,9 @@ import org.apache.beam.sdk.transforms.windowing.Window; import org.apache.beam.sdk.transforms.windowing.WindowFn; import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Joiner; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.MoreObjects; import org.checkerframework.checker.nullness.qual.Nullable; +import org.checkerframework.dataflow.qual.Pure; import org.joda.time.Instant; /** @@ -93,9 +96,6 @@ * invocations of the callbacks. All important values should be persisted using state before the * callback returns. */ -@SuppressWarnings({ - "nullness" // TODO(https://github.com/apache/beam/issues/20497) -}) public abstract class TriggerStateMachine implements Serializable { /** @@ -130,7 +130,11 @@ public interface TriggerInfo { /** Returns an iterable over the unfinished sub-triggers of the current trigger. */ Iterable unfinishedSubTriggers(); - /** Returns the first unfinished sub-trigger. */ + /** + * Returns the first unfinished sub-trigger, or {@code null} if all sub-triggers of the current + * trigger are finished. + */ + @Nullable ExecutableTriggerStateMachine firstUnfinishedSubTrigger(); /** @@ -192,6 +196,7 @@ public abstract static class TriggerContext { public abstract @Nullable Instant currentSynchronizedProcessingTime(); /** The current event time for the input or {@code null} if unknown. */ + @Pure public abstract @Nullable Instant currentEventTime(); } @@ -346,8 +351,8 @@ public void clear(TriggerContext c) throws Exception { } } - public Iterable subTriggers() { - return subTriggers; + public List subTriggers() { + return MoreObjects.firstNonNull(subTriggers, Collections.emptyList()); } /** Returns whether this performs the same triggering as the given {@code Trigger}. */ diff --git a/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/TriggerStateMachineContextFactory.java b/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/TriggerStateMachineContextFactory.java index 9e491c17b73e..60fb8474ab0e 100644 --- a/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/TriggerStateMachineContextFactory.java +++ b/runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/TriggerStateMachineContextFactory.java @@ -46,7 +46,6 @@ *

These contexts are highly interdependent and share many fields; it is inadvisable to create * them via any means other than this factory class. */ -@SuppressWarnings({"nullness", "keyfor"}) // TODO(https://github.com/apache/beam/issues/20497) public class TriggerStateMachineContextFactory { private final WindowFn windowFn; @@ -256,17 +255,17 @@ public boolean finishedInAllMergingWindows() { } } + private StateNamespace namespaceFor(W window, int triggerIndex) { + return StateNamespaces.windowAndTrigger(windowCoder, window, triggerIndex); + } + private class StateAccessorImpl implements StateAccessor { protected final int triggerIndex; protected final StateNamespace windowNamespace; public StateAccessorImpl(W window, ExecutableTriggerStateMachine trigger) { this.triggerIndex = trigger.getTriggerIndex(); - this.windowNamespace = namespaceFor(window); - } - - protected StateNamespace namespaceFor(W window) { - return StateNamespaces.windowAndTrigger(windowCoder, window, triggerIndex); + this.windowNamespace = namespaceFor(window, triggerIndex); } @Override @@ -295,7 +294,8 @@ public Map accessInEachMergingWindow( StateTag address) { ImmutableMap.Builder builder = ImmutableMap.builder(); for (W mergingWindow : activeToBeMerged) { - StateT stateForWindow = stateInternals.state(namespaceFor(mergingWindow), address); + StateT stateForWindow = + stateInternals.state(namespaceFor(mergingWindow, triggerIndex), address); builder.put(mergingWindow, stateForWindow); } return builder.build(); @@ -309,6 +309,9 @@ private class TriggerContextImpl extends TriggerStateMachine.TriggerContext { private final Timers timers; private final TriggerInfoImpl triggerInfo; + // A context and its TriggerInfo refer to each other, so the TriggerInfo necessarily sees this + // context before its fields are all assigned. It only stores the reference for later use. + @SuppressWarnings({"nullness:assignment", "nullness:argument"}) private TriggerContextImpl( W window, Timers timers, @@ -372,6 +375,9 @@ private class OnElementContextImpl extends TriggerStateMachine.OnElementContext private final TriggerInfoImpl triggerInfo; private final Instant eventTimestamp; + // A context and its TriggerInfo refer to each other, so the TriggerInfo necessarily sees this + // context before its fields are all assigned. It only stores the reference for later use. + @SuppressWarnings({"nullness:assignment", "nullness:argument"}) private OnElementContextImpl( W window, Timers timers, @@ -447,6 +453,9 @@ private class OnMergeContextImpl extends TriggerStateMachine.OnMergeContext { private final Timers timers; private final MergingTriggerInfoImpl triggerInfo; + // A context and its TriggerInfo refer to each other, so the TriggerInfo necessarily sees this + // context before its fields are all assigned. It only stores the reference for later use. + @SuppressWarnings({"nullness:assignment", "nullness:argument"}) private OnMergeContextImpl( W window, Timers timers, diff --git a/runners/core-java/src/test/java/org/apache/beam/runners/core/triggers/AfterEachStateMachineTest.java b/runners/core-java/src/test/java/org/apache/beam/runners/core/triggers/AfterEachStateMachineTest.java index 301cf7d2c204..e48c9ac05ece 100644 --- a/runners/core-java/src/test/java/org/apache/beam/runners/core/triggers/AfterEachStateMachineTest.java +++ b/runners/core-java/src/test/java/org/apache/beam/runners/core/triggers/AfterEachStateMachineTest.java @@ -17,8 +17,11 @@ */ package org.apache.beam.runners.core.triggers; +import static org.hamcrest.MatcherAssert.assertThat; +import static org.hamcrest.Matchers.containsString; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertThrows; import static org.junit.Assert.assertTrue; import org.apache.beam.runners.core.triggers.TriggerStateMachineTester.SimpleTriggerStateMachineTester; @@ -95,6 +98,57 @@ public void testAfterEachInSequence() throws Exception { assertTrue(tester.isMarkedFinished(window)); } + /** + * Once the last subtrigger finishes, the {@link AfterEachStateMachine} itself is marked finished + * and the trigger machinery stops delivering elements to it. + */ + @Test + public void testAfterEachFinishesWithItsLastSubtrigger() throws Exception { + tester = + TriggerStateMachineTester.forTrigger( + AfterEachStateMachine.inOrder( + AfterPaneStateMachine.elementCountAtLeast(1), + AfterPaneStateMachine.elementCountAtLeast(1)), + FixedWindows.of(Duration.millis(10))); + + IntervalWindow window = new IntervalWindow(new Instant(0), new Instant(10)); + + tester.injectElements(1); + tester.fireIfShouldFire(window); + assertFalse(tester.isMarkedFinished(window)); + + tester.injectElements(2); + tester.fireIfShouldFire(window); + assertTrue(tester.isMarkedFinished(window)); + + // The window is closed, so this element must not re-enter the trigger. + tester.injectElements(3); + assertTrue(tester.isMarkedFinished(window)); + } + + /** + * {@link AfterEachStateMachine} has no subtrigger to delegate to once they are all finished, and + * nothing in the trigger machinery enforces that it is not invoked in that state. If the + * invariant is ever broken, it must fail with a message identifying the trigger. + */ + @Test + public void testInvokedWithAllSubtriggersFinished() throws Exception { + tester = + TriggerStateMachineTester.forTrigger( + AfterEachStateMachine.inOrder( + AfterPaneStateMachine.elementCountAtLeast(1), + AfterPaneStateMachine.elementCountAtLeast(1)), + FixedWindows.of(Duration.millis(10))); + + IntervalWindow window = new IntervalWindow(new Instant(0), new Instant(10)); + tester.setSubTriggerFinishedForWindow(0, window, true); + tester.setSubTriggerFinishedForWindow(1, window, true); + + IllegalStateException thrown = + assertThrows(IllegalStateException.class, () -> tester.shouldFire(window)); + assertThat(thrown.getMessage(), containsString("AfterEach.inOrder")); + } + @Test public void testToString() { TriggerStateMachine trigger =