diff --git a/events-domain/src/main/java/io/split/android/client/events/SplitEventsManager.java b/events-domain/src/main/java/io/split/android/client/events/SplitEventsManager.java index d49b404f2..b5a9683b7 100644 --- a/events-domain/src/main/java/io/split/android/client/events/SplitEventsManager.java +++ b/events-domain/src/main/java/io/split/android/client/events/SplitEventsManager.java @@ -2,7 +2,6 @@ import static java.util.Objects.requireNonNull; -import androidx.annotation.NonNull; import androidx.annotation.VisibleForTesting; import java.util.concurrent.Executor; @@ -10,7 +9,6 @@ import io.harness.events.EventHandler; import io.harness.events.EventsManager; import io.harness.events.EventsManagers; -import io.split.android.client.SplitClient; import io.split.android.client.api.EventMetadata; import io.split.android.client.events.executors.SplitEventExecutorResources; import io.split.android.client.events.executors.SplitEventExecutorResourcesImpl; @@ -28,10 +26,6 @@ public class SplitEventsManager implements ISplitEventsManager, ListenableEvents private final DualExecutorRegistration mDualExecutorRegistration; private SplitEventExecutorResources mResources; - // Track sync completion for SDK_READY_FROM_CACHE triggering. TODO: This is a temporary adaptation before extending EventsManager requireAny. - private volatile boolean mSplitsSyncComplete = false; - private volatile boolean mSegmentsSyncComplete = false; - /** * Creates a new SplitEventsManager. * @@ -152,27 +146,20 @@ public void destroy() { /** * Notifies the synthetic composite events based on the actual internal event. * These synthetic events simplify the SDK_READY condition evaluation. + * The prerequisite configuration ensures SDK_READY_FROM_CACHE always fires before SDK_READY. *

- * Also handles the special case where SDK_READY_FROM_CACHE should fire - * before SDK_READY when both sync events have completed. + * TODO: Remove this method once EventsManagerConfig is updated. */ private void notifySyntheticEventsIfNeeded(SplitInternalEvent internalEvent) { switch (internalEvent) { case SPLITS_UPDATED: case SPLITS_FETCHED: - mSplitsSyncComplete = true; - // Check if SDK_READY_FROM_CACHE should fire BEFORE notifying the sync complete - // to ensure correct ordering (SDK_READY_FROM_CACHE before SDK_READY) - triggerReadyFromCacheIfNeeded(); mEventsManager.notifyInternalEvent(SplitInternalEvent.SPLITS_SYNC_COMPLETE, null); break; case MY_SEGMENTS_UPDATED: case MY_SEGMENTS_FETCHED: case MY_LARGE_SEGMENTS_UPDATED: - mSegmentsSyncComplete = true; - // Check if SDK_READY_FROM_CACHE should fire BEFORE notifying the sync complete - triggerReadyFromCacheIfNeeded(); mEventsManager.notifyInternalEvent(SplitInternalEvent.SEGMENTS_SYNC_COMPLETE, null); break; @@ -182,36 +169,6 @@ private void notifySyntheticEventsIfNeeded(SplitInternalEvent internalEvent) { } } - /** - * Triggers SDK_READY_FROM_CACHE if SDK_READY is about to fire (both sync events received) - * but SDK_READY_FROM_CACHE hasn't fired yet. - *

- * TODO: This is a temporary adaptation before extending EventsManager requireAny. - */ - private void triggerReadyFromCacheIfNeeded() { - // Only trigger if both sync events have completed - if (!mSplitsSyncComplete || !mSegmentsSyncComplete) { - return; - } - - // If SDK_READY_FROM_CACHE already triggered, nothing to do - if (mEventsManager.eventAlreadyTriggered(SplitEvent.SDK_READY_FROM_CACHE)) { - return; - } - - // If SDK_READY already triggered, nothing to do (too late) - if (mEventsManager.eventAlreadyTriggered(SplitEvent.SDK_READY)) { - return; - } - - // Both sync events received and SDK_READY_FROM_CACHE not yet fired - // Trigger it by firing all required cache events - mEventsManager.notifyInternalEvent(SplitInternalEvent.SPLITS_LOADED_FROM_STORAGE, null); - mEventsManager.notifyInternalEvent(SplitInternalEvent.MY_SEGMENTS_LOADED_FROM_STORAGE, null); - mEventsManager.notifyInternalEvent(SplitInternalEvent.ATTRIBUTES_LOADED_FROM_STORAGE, null); - mEventsManager.notifyInternalEvent(SplitInternalEvent.ENCRYPTION_MIGRATION_DONE, null); - } - private void startTimeoutThread(final int blockUntilReady) { Thread timeoutThread = new Thread(new Runnable() { @Override diff --git a/events-domain/src/main/java/io/split/android/client/events/SplitEventsManagerConfigFactory.java b/events-domain/src/main/java/io/split/android/client/events/SplitEventsManagerConfigFactory.java index b6d903719..b1999a170 100644 --- a/events-domain/src/main/java/io/split/android/client/events/SplitEventsManagerConfigFactory.java +++ b/events-domain/src/main/java/io/split/android/client/events/SplitEventsManagerConfigFactory.java @@ -1,5 +1,8 @@ package io.split.android.client.events; +import java.util.HashSet; +import java.util.Set; + import io.harness.events.EventsManagerConfig; /** @@ -19,8 +22,8 @@ private SplitEventsManagerConfigFactory() { *

* Event rules: *

@@ -28,16 +31,27 @@ private SplitEventsManagerConfigFactory() { * @return the configured EventsManagerConfig */ static EventsManagerConfig create() { + // SDK_READY_FROM_CACHE fires when either: + // 1. Cache path: All cache loading events complete (AND), OR + // 2. Sync path: All sync events complete (AND) + Set cacheGroup = new HashSet<>(); + cacheGroup.add(SplitInternalEvent.SPLITS_LOADED_FROM_STORAGE); + cacheGroup.add(SplitInternalEvent.MY_SEGMENTS_LOADED_FROM_STORAGE); + cacheGroup.add(SplitInternalEvent.ATTRIBUTES_LOADED_FROM_STORAGE); + cacheGroup.add(SplitInternalEvent.ENCRYPTION_MIGRATION_DONE); + + Set syncGroup = new HashSet<>(); + syncGroup.add(SplitInternalEvent.SPLITS_SYNC_COMPLETE); + syncGroup.add(SplitInternalEvent.SEGMENTS_SYNC_COMPLETE); + return EventsManagerConfig.builder() .requireAll(SplitEvent.SDK_READY, SplitInternalEvent.SPLITS_SYNC_COMPLETE, SplitInternalEvent.SEGMENTS_SYNC_COMPLETE) - .requireAll(SplitEvent.SDK_READY_FROM_CACHE, - SplitInternalEvent.SPLITS_LOADED_FROM_STORAGE, - SplitInternalEvent.MY_SEGMENTS_LOADED_FROM_STORAGE, - SplitInternalEvent.ATTRIBUTES_LOADED_FROM_STORAGE, - SplitInternalEvent.ENCRYPTION_MIGRATION_DONE) + // SDK_READY_FROM_CACHE: OR of ANDs + // Fires when (cache group all done) OR (sync group all done) + .requireAny(SplitEvent.SDK_READY_FROM_CACHE, cacheGroup, syncGroup) .requireAny(SplitEvent.SDK_READY_TIMED_OUT, SplitInternalEvent.SDK_READY_TIMEOUT_REACHED) @@ -49,6 +63,8 @@ static EventsManagerConfig create() { SplitInternalEvent.RULE_BASED_SEGMENTS_UPDATED, SplitInternalEvent.SPLIT_KILLED_NOTIFICATION) + // SDK_READY requires SDK_READY_FROM_CACHE to fire first + .prerequisite(SplitEvent.SDK_READY, SplitEvent.SDK_READY_FROM_CACHE) .prerequisite(SplitEvent.SDK_UPDATE, SplitEvent.SDK_READY) .suppressedBy(SplitEvent.SDK_READY_TIMED_OUT, SplitEvent.SDK_READY) diff --git a/events-domain/src/main/java/io/split/android/client/events/SplitInternalEvent.java b/events-domain/src/main/java/io/split/android/client/events/SplitInternalEvent.java index b9184c3e2..48d32bec3 100644 --- a/events-domain/src/main/java/io/split/android/client/events/SplitInternalEvent.java +++ b/events-domain/src/main/java/io/split/android/client/events/SplitInternalEvent.java @@ -20,18 +20,14 @@ public enum SplitInternalEvent { /** * Synthetic event: fired when splits sync completes (either SPLITS_FETCHED or SPLITS_UPDATED). - * Used internally to simplify SDK_READY condition evaluation. - *

- * TODO: This is a temporary adaptation before extending EventsManager requireAny. + * Used internally to simplify SDK_READY and SDK_READY_FROM_CACHE condition evaluation. */ SPLITS_SYNC_COMPLETE, /** * Synthetic event: fired when segments sync completes (any of MY_SEGMENTS_FETCHED, * MY_SEGMENTS_UPDATED, or MY_LARGE_SEGMENTS_UPDATED). - * Used internally to simplify SDK_READY condition evaluation. - *

- * TODO: This is a temporary adaptation before extending EventsManager requireAny. + * Used internally to simplify SDK_READY and SDK_READY_FROM_CACHE condition evaluation. */ SEGMENTS_SYNC_COMPLETE, } diff --git a/events/src/main/java/io/harness/events/EventsManagerConfig.java b/events/src/main/java/io/harness/events/EventsManagerConfig.java index 61b20ef16..fc83b8b08 100644 --- a/events/src/main/java/io/harness/events/EventsManagerConfig.java +++ b/events/src/main/java/io/harness/events/EventsManagerConfig.java @@ -18,8 +18,8 @@ public final class EventsManagerConfig { // External events that require ALL listed internals (AND) private final Map> mRequireAll; - // External events triggered by ANY of the listed internals (OR) - private final Map> mRequireAny; + // External events triggered by ANY of the listed internal groups (OR of ANDs) + private final Map>> mRequireAny; // External-event guards: prerequisites that must have fired before External can emit private final Map> mPrerequisites; // External-event guards: if any of these have fired, suppress E @@ -31,16 +31,16 @@ public final class EventsManagerConfig { * Creates a new EventsManagerConfig. * * @param requireAll External events that require ALL listed internals (AND) - * @param requireAny External events triggered by ANY of the listed internals (OR) + * @param requireAny External events triggered by ANY of the listed internal groups (OR of ANDs) * @param prerequisites External-event guards: prerequisites that must have fired before External can emit * @param suppressedBy External-event guards: if any of these have fired, suppress E * @param executionLimits Execution policy: max executions per external event (-1 = unlimited) */ private EventsManagerConfig(Map> requireAll, - Map> requireAny, - Map> prerequisites, - Map> suppressedBy, - Map executionLimits) { + Map>> requireAny, + Map> prerequisites, + Map> suppressedBy, + Map executionLimits) { mRequireAll = requireAll == null ? Collections.emptyMap() : Collections.unmodifiableMap(new HashMap<>(requireAll)); @@ -72,7 +72,7 @@ public Map> getRequireAll() { } @NotNull - public Map> getRequireAny() { + public Map>> getRequireAny() { return mRequireAny; } @@ -110,7 +110,7 @@ public static Builder builder() { */ public static final class Builder { private final Map> mRequireAll = new HashMap<>(); - private final Map> mRequireAny = new HashMap<>(); + private final Map>> mRequireAny = new HashMap<>(); private final Map> mPrerequisites = new HashMap<>(); private final Map> mSuppressedBy = new HashMap<>(); private final Map mExecutionLimits = new HashMap<>(); @@ -121,8 +121,8 @@ private Builder() { /** * Adds a requirement that ALL specified internal events must occur for the external event to fire. * - * @param externalEvent the external event - * @param internalEvents the internal events that must ALL occur + * @param externalEvent the external event + * @param internalEvents the internal events that must ALL occur * @return this builder */ @SafeVarargs @@ -133,14 +133,45 @@ public final Builder requireAll(E externalEvent, I... internalEvents) { /** * Adds a requirement that ANY of the specified internal events will trigger the external event. + * Each internal event is treated as a group of one (singleton). * - * @param externalEvent the external event - * @param internalEvents the internal events, any of which will trigger the external event + * @param externalEvent the external event + * @param internalEvents the internal events, any of which will trigger the external event * @return this builder */ @SafeVarargs public final Builder requireAny(E externalEvent, I... internalEvents) { - mRequireAny.put(externalEvent, new HashSet<>(Arrays.asList(internalEvents))); + // Convert each individual event to a singleton Set (group of one) + Set> groups = new HashSet<>(); + for (I internalEvent : internalEvents) { + groups.add(Collections.singleton(internalEvent)); + } + mRequireAny.put(externalEvent, groups); + return this; + } + + /** + * Adds a requirement that ANY of the specified internal event groups will trigger the external event. + * Each group is an AND: all events in the group must occur. + * The external event fires when ANY group is fully satisfied (OR of ANDs). + *

+ * Example: + *

+         * .requireAny(DISH_SERVED,
+         *     Set.of(BOUGHT_INGREDIENTS, COOKED_MEAL),                    // Fresh cooking path
+         *     Set.of(ORDERED_DELIVERY, DELIVERY_ARRIVED))                 // Delivery path
+         * // Fires when: (fresh cooking done) OR (delivery arrived)
+         * 
+ * + * @param externalEvent the external event + * @param internalEventGroups the groups of internal events; all events in a group must occur (AND), + * and any group being satisfied triggers the external event (OR) + * @return this builder + */ + @SafeVarargs + public final Builder requireAny(E externalEvent, Set... internalEventGroups) { + Set> groups = new HashSet<>(Arrays.asList(internalEventGroups)); + mRequireAny.put(externalEvent, groups); return this; } diff --git a/events/src/main/java/io/harness/events/EventsManagerCore.java b/events/src/main/java/io/harness/events/EventsManagerCore.java index fbbebbdf4..bb8ae47a3 100644 --- a/events/src/main/java/io/harness/events/EventsManagerCore.java +++ b/events/src/main/java/io/harness/events/EventsManagerCore.java @@ -149,33 +149,42 @@ private void processInternal(I event, M metadata) { currentSeenInternal = new HashSet<>(mSeenInternal); } - // Evaluate AND external events - for (Map.Entry> entry : mConfig.getRequireAll().entrySet()) { - E external = entry.getKey(); - Set required = entry.getValue(); - - if (!required.isEmpty() && currentSeenInternal.containsAll(required)) { - triggerIfConditionsMet(external, metadata); - } - } + // Track events fired in this processing cycle to avoid re-firing. + // This prevents infinite loops with unlimited events. + Set firedInThisCycle = new HashSet<>(); + + // Loop until no more events fire in an iteration. + // Without this loop, events would be missed if their external prerequisites + // aren't satisfied on the first pass, but become satisfied after other events fire. + boolean anyEventFiredThisIteration; + do { + boolean requireAllEventsFired = evaluateRequireAllEvents(currentSeenInternal, firedInThisCycle, metadata); + boolean requireAnyEventsFired = evaluateRequireAnyEvents(currentSeenInternal, firedInThisCycle, metadata); + anyEventFiredThisIteration = requireAllEventsFired || requireAnyEventsFired; + + } while (anyEventFiredThisIteration); + } - // Evaluate OR external events - for (Map.Entry> entry : mConfig.getRequireAny().entrySet()) { - E external = entry.getKey(); - if (entry.getValue().contains(event)) { - triggerIfConditionsMet(external, metadata); - } + /** + * Triggers an external event if all conditions are met. + * @return true if the event was triggered, false otherwise + */ + private boolean triggerIfConditionsMet(E event, M metadata) { + if (!canEventBeTriggered(event)) { + return false; } + return trigger(event, metadata); } - private void triggerIfConditionsMet(E event, M metadata) { - if (!prerequisitesSatisfied(event) || isSuppressed(event)) { - return; - } - trigger(event, metadata); + private boolean canEventBeTriggered(E event) { + return prerequisitesSatisfied(event) && !isSuppressed(event); } - private void trigger(E event, M metadata) { + /** + * Triggers an external event. + * @return true if the event was triggered, false if it was already at max executions + */ + private boolean trigger(E event, M metadata) { Set> handlersSnapshot = Collections.emptySet(); synchronized (mLock) { @@ -184,7 +193,7 @@ private void trigger(E event, M metadata) { int triggered = count != null ? count : 0; if (max != UNLIMITED && triggered >= max) { - return; + return false; } mTriggerCount.put(event, triggered + 1); @@ -198,6 +207,7 @@ private void trigger(E event, M metadata) { for (EventHandler handler : handlersSnapshot) { mDelivery.deliver(handler, event, metadata); } + return true; } private int maxExecutions(E event) { @@ -231,4 +241,61 @@ private boolean isSuppressed(E external) { } return false; } + + /** + * Evaluates events with AND logic: fire if ALL required internal events have been seen. + */ + private boolean evaluateRequireAllEvents(Set seenInternal, Set firedInThisCycle, M metadata) { + boolean anyEventFired = false; + for (Map.Entry> entry : mConfig.getRequireAll().entrySet()) { + E externalEvent = entry.getKey(); + if (hasAlreadyFiredInCycle(externalEvent, firedInThisCycle)) { + continue; + } + Set requiredInternals = entry.getValue(); + + if (allInternalEventsSeen(requiredInternals, seenInternal) && triggerIfConditionsMet(externalEvent, metadata)) { + firedInThisCycle.add(externalEvent); + anyEventFired = true; + } + } + return anyEventFired; + } + + /** + * Evaluates events with OR-of-ANDs logic: fire if ANY group has ALL its internal events seen. + */ + private boolean evaluateRequireAnyEvents(Set seenInternal, Set firedInThisCycle, M metadata) { + boolean anyEventFired = false; + for (Map.Entry>> entry : mConfig.getRequireAny().entrySet()) { + E externalEvent = entry.getKey(); + if (hasAlreadyFiredInCycle(externalEvent, firedInThisCycle)) { + continue; + } + Set> requiredGroups = entry.getValue(); + + if (anyGroupSatisfied(requiredGroups, seenInternal) && triggerIfConditionsMet(externalEvent, metadata)) { + firedInThisCycle.add(externalEvent); + anyEventFired = true; + } + } + return anyEventFired; + } + + private boolean hasAlreadyFiredInCycle(E event, Set firedInThisCycle) { + return firedInThisCycle.contains(event); + } + + private boolean allInternalEventsSeen(Set requiredInternals, Set seenInternal) { + return !requiredInternals.isEmpty() && seenInternal.containsAll(requiredInternals); + } + + private boolean anyGroupSatisfied(Set> requiredGroups, Set seenInternal) { + for (Set group : requiredGroups) { + if (allInternalEventsSeen(group, seenInternal)) { + return true; + } + } + return false; + } } diff --git a/events/src/test/java/io/harness/events/EventsManagerConfigTest.java b/events/src/test/java/io/harness/events/EventsManagerConfigTest.java index e53ec9eba..7ab4405e1 100644 --- a/events/src/test/java/io/harness/events/EventsManagerConfigTest.java +++ b/events/src/test/java/io/harness/events/EventsManagerConfigTest.java @@ -8,6 +8,8 @@ import org.junit.Test; import java.util.Collections; +import java.util.HashSet; +import java.util.Set; public class EventsManagerConfigTest { @@ -36,8 +38,11 @@ public void builderCreatesConfigWithAllFields() { assertTrue(config.getRequireAll().get("E1").contains("I1")); assertTrue(config.getRequireAll().get("E1").contains("I2")); + // requireAny now stores Set> - single events are wrapped in singleton sets assertEquals(1, config.getRequireAny().size()); - assertTrue(config.getRequireAny().get("E2").contains("I3")); + Set> requireAnyGroups = config.getRequireAny().get("E2"); + assertEquals(1, requireAnyGroups.size()); + assertTrue(requireAnyGroups.contains(Collections.singleton("I3"))); assertEquals(1, config.getPrerequisites().size()); assertTrue(config.getPrerequisites().get("E1").contains("E0")); @@ -93,7 +98,7 @@ public void returnedMapsAreUnmodifiable() { } try { - config.getRequireAny().put("E2", Collections.singleton("I2")); + config.getRequireAny().put("E2", Collections.singleton(Collections.singleton("I2"))); Assert.fail("getRequireAny() should return an unmodifiable map"); } catch (UnsupportedOperationException expected) { // expected @@ -138,4 +143,60 @@ public void emptyMethodReturnsEmptyUnmodifiableConfig() { // expected } } + + @Test + public void requireAnyWithVarargsCreatesIndividualGroups() { + // When using requireAny(E, I...), each I should become its own singleton group + EventsManagerConfig config = EventsManagerConfig.builder() + .requireAny("E1", "I1", "I2", "I3") + .build(); + + Set> groups = config.getRequireAny().get("E1"); + assertEquals(3, groups.size()); + assertTrue(groups.contains(Collections.singleton("I1"))); + assertTrue(groups.contains(Collections.singleton("I2"))); + assertTrue(groups.contains(Collections.singleton("I3"))); + } + + @Test + public void requireAnyWithSetsCreatesAndGroups() { + // When using requireAny(E, Set...), each Set is an AND group + Set group1 = new HashSet<>(); + group1.add("I1"); + group1.add("I2"); + + Set group2 = new HashSet<>(); + group2.add("I3"); + group2.add("I4"); + + EventsManagerConfig config = EventsManagerConfig.builder() + .requireAny("E1", group1, group2) + .build(); + + Set> groups = config.getRequireAny().get("E1"); + assertEquals(2, groups.size()); + assertTrue(groups.contains(group1)); + assertTrue(groups.contains(group2)); + } + + @Test + public void requireAnyWithMixedGroupSizes() { + // Groups can have different sizes + Set singletonGroup = Collections.singleton("I1"); + + Set largeGroup = new HashSet<>(); + largeGroup.add("I2"); + largeGroup.add("I3"); + largeGroup.add("I4"); + largeGroup.add("I5"); + + EventsManagerConfig config = EventsManagerConfig.builder() + .requireAny("E1", singletonGroup, largeGroup) + .build(); + + Set> groups = config.getRequireAny().get("E1"); + assertEquals(2, groups.size()); + assertTrue(groups.contains(singletonGroup)); + assertTrue(groups.contains(largeGroup)); + } } diff --git a/events/src/test/java/io/harness/events/EventsManagerTest.java b/events/src/test/java/io/harness/events/EventsManagerTest.java index c09dbc6e7..7039a9de4 100644 --- a/events/src/test/java/io/harness/events/EventsManagerTest.java +++ b/events/src/test/java/io/harness/events/EventsManagerTest.java @@ -6,6 +6,8 @@ import org.junit.Test; +import java.util.HashSet; +import java.util.Set; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; @@ -124,10 +126,7 @@ public void seasoningAdjustedIsEmittedOnlyAfterDishServed() throws InterruptedEx eventsManager.register(CookingEvent.SEASONING_ADJUSTED, (event, metadata) -> hCount.incrementAndGet()); - // SEASONING_ADDED before DISH_SERVED - should not fire - eventsManager.notifyInternalEvent(KitchenActivity.SEASONING_ADDED, null); - - // Trigger DISH_SERVED + // Trigger DISH_SERVED first (without SEASONING_ADDED) eventsManager.notifyInternalEvent(KitchenActivity.INGREDIENTS_PREPPED, null); eventsManager.notifyInternalEvent(KitchenActivity.LEFTOVER_SAUCE_FOUND, null); eventsManager.notifyInternalEvent(KitchenActivity.OVEN_PREHEATED, null); @@ -135,7 +134,10 @@ public void seasoningAdjustedIsEmittedOnlyAfterDishServed() throws InterruptedEx // Wait for DISH_SERVED to be processed assertTrue(eventsManager.eventAlreadyTriggered(CookingEvent.DISH_SERVED)); - // Now SEASONING_ADDED should trigger SEASONING_ADJUSTED + // SEASONING_ADJUSTED should NOT have fired yet (no SEASONING_ADDED) + assertEquals(0, hCount.get()); + + // Now SEASONING_ADDED should trigger SEASONING_ADJUSTED (prerequisite is met) eventsManager.notifyInternalEvent(KitchenActivity.SEASONING_ADDED, null); assertTrue(seasoningLatch.await(TIMEOUT_MS, TimeUnit.MILLISECONDS)); @@ -443,4 +445,406 @@ public void eventAlreadyTriggeredRespectsExecutionLimits() throws InterruptedExc // Unlimited event returns false (can still fire again) assertFalse(eventsManager.eventAlreadyTriggered(CookingEvent.SEASONING_ADJUSTED)); } + + @Test + public void requireAnyWithGroupsFiresWhenFirstGroupComplete() throws InterruptedException { + // External event fires when EITHER: + // Group 1: INGREDIENTS_PREPPED AND SEASONING_ADDED + // OR + // Group 2: LEFTOVER_MEAT_FOUND AND LEFTOVER_VEGGIES_FOUND AND LEFTOVER_SAUCE_FOUND + Set group1 = new HashSet<>(); + group1.add(KitchenActivity.INGREDIENTS_PREPPED); + group1.add(KitchenActivity.SEASONING_ADDED); + + Set group2 = new HashSet<>(); + group2.add(KitchenActivity.LEFTOVER_MEAT_FOUND); + group2.add(KitchenActivity.LEFTOVER_VEGGIES_FOUND); + group2.add(KitchenActivity.LEFTOVER_SAUCE_FOUND); + + EventsManagerConfig config = EventsManagerConfig.builder() + .requireAny(CookingEvent.DISH_SERVED, group1, group2) + .executionLimit(CookingEvent.DISH_SERVED, 1) + .build(); + + CountDownLatch latch = new CountDownLatch(1); + AtomicInteger callCount = new AtomicInteger(0); + + EventDelivery delivery = (handler, event, metadata) -> { + handler.handle(event, metadata); + latch.countDown(); + }; + + EventsManager eventsManager = new EventsManagerCore<>(config, delivery); + eventsManager.register(CookingEvent.DISH_SERVED, (event, metadata) -> callCount.incrementAndGet()); + + // Complete first group + eventsManager.notifyInternalEvent(KitchenActivity.INGREDIENTS_PREPPED, null); + eventsManager.notifyInternalEvent(KitchenActivity.SEASONING_ADDED, null); + + assertTrue(latch.await(TIMEOUT_MS, TimeUnit.MILLISECONDS)); + assertEquals(1, callCount.get()); + assertTrue(eventsManager.eventAlreadyTriggered(CookingEvent.DISH_SERVED)); + } + + @Test + public void requireAnyWithGroupsFiresWhenSecondGroupComplete() throws InterruptedException { + // Same config as above, but complete the second group instead + Set group1 = new HashSet<>(); + group1.add(KitchenActivity.INGREDIENTS_PREPPED); + group1.add(KitchenActivity.SEASONING_ADDED); + + Set group2 = new HashSet<>(); + group2.add(KitchenActivity.LEFTOVER_MEAT_FOUND); + group2.add(KitchenActivity.LEFTOVER_VEGGIES_FOUND); + + EventsManagerConfig config = EventsManagerConfig.builder() + .requireAny(CookingEvent.DISH_SERVED, group1, group2) + .executionLimit(CookingEvent.DISH_SERVED, 1) + .build(); + + CountDownLatch latch = new CountDownLatch(1); + AtomicInteger callCount = new AtomicInteger(0); + + EventDelivery delivery = (handler, event, metadata) -> { + handler.handle(event, metadata); + latch.countDown(); + }; + + EventsManager eventsManager = new EventsManagerCore<>(config, delivery); + eventsManager.register(CookingEvent.DISH_SERVED, (event, metadata) -> callCount.incrementAndGet()); + + // Complete second group (not touching first group) + eventsManager.notifyInternalEvent(KitchenActivity.LEFTOVER_MEAT_FOUND, null); + eventsManager.notifyInternalEvent(KitchenActivity.LEFTOVER_VEGGIES_FOUND, null); + + assertTrue(latch.await(TIMEOUT_MS, TimeUnit.MILLISECONDS)); + assertEquals(1, callCount.get()); + assertTrue(eventsManager.eventAlreadyTriggered(CookingEvent.DISH_SERVED)); + } + + @Test + public void requireAnyWithGroupsDoesNotFireWithPartialGroup() throws InterruptedException { + Set group1 = new HashSet<>(); + group1.add(KitchenActivity.INGREDIENTS_PREPPED); + group1.add(KitchenActivity.SEASONING_ADDED); + group1.add(KitchenActivity.OVEN_PREHEATED); + + Set group2 = new HashSet<>(); + group2.add(KitchenActivity.LEFTOVER_MEAT_FOUND); + group2.add(KitchenActivity.LEFTOVER_VEGGIES_FOUND); + + EventsManagerConfig config = EventsManagerConfig.builder() + .requireAny(CookingEvent.DISH_SERVED, group1, group2) + .executionLimit(CookingEvent.DISH_SERVED, 1) + .build(); + + AtomicInteger callCount = new AtomicInteger(0); + + EventsManager eventsManager = new EventsManagerCore<>(config, SIMPLE_DELIVERY); + eventsManager.register(CookingEvent.DISH_SERVED, (event, metadata) -> callCount.incrementAndGet()); + + // Partial completion of group 1 (missing OVEN_PREHEATED) + eventsManager.notifyInternalEvent(KitchenActivity.INGREDIENTS_PREPPED, null); + eventsManager.notifyInternalEvent(KitchenActivity.SEASONING_ADDED, null); + + // Partial completion of group 2 (missing LEFTOVER_VEGGIES_FOUND) + eventsManager.notifyInternalEvent(KitchenActivity.LEFTOVER_MEAT_FOUND, null); + + // Wait for processing to complete + eventsManager.eventAlreadyTriggered(CookingEvent.DISH_SERVED); + + assertEquals(0, callCount.get()); + assertFalse(eventsManager.eventAlreadyTriggered(CookingEvent.DISH_SERVED)); + } + + @Test + public void requireAnyWithGroupsFiresOnceEvenWhenMultipleGroupsComplete() throws InterruptedException { + Set group1 = new HashSet<>(); + group1.add(KitchenActivity.INGREDIENTS_PREPPED); + + Set group2 = new HashSet<>(); + group2.add(KitchenActivity.LEFTOVER_MEAT_FOUND); + + EventsManagerConfig config = EventsManagerConfig.builder() + .requireAny(CookingEvent.DISH_SERVED, group1, group2) + .executionLimit(CookingEvent.DISH_SERVED, 1) + .build(); + + CountDownLatch latch = new CountDownLatch(1); + AtomicInteger callCount = new AtomicInteger(0); + + EventDelivery delivery = (handler, event, metadata) -> { + handler.handle(event, metadata); + latch.countDown(); + }; + + EventsManager eventsManager = new EventsManagerCore<>(config, delivery); + eventsManager.register(CookingEvent.DISH_SERVED, (event, metadata) -> callCount.incrementAndGet()); + + // Complete first group + eventsManager.notifyInternalEvent(KitchenActivity.INGREDIENTS_PREPPED, null); + + assertTrue(latch.await(TIMEOUT_MS, TimeUnit.MILLISECONDS)); + + // Now complete second group as well + eventsManager.notifyInternalEvent(KitchenActivity.LEFTOVER_MEAT_FOUND, null); + + // Wait for processing + eventsManager.eventAlreadyTriggered(CookingEvent.DISH_SERVED); + + // Should only fire once due to execution limit + assertEquals(1, callCount.get()); + } + + @Test + public void requireAnyGroupedWithPrerequisite() throws InterruptedException { + // DISH_SERVED requires simple condition + // SEASONING_ADJUSTED uses OR-of-ANDs and requires DISH_SERVED first + Set group1 = new HashSet<>(); + group1.add(KitchenActivity.SEASONING_ADDED); + + Set group2 = new HashSet<>(); + group2.add(KitchenActivity.LEFTOVER_MEAT_FOUND); + group2.add(KitchenActivity.LEFTOVER_VEGGIES_FOUND); + + EventsManagerConfig config = EventsManagerConfig.builder() + .requireAny(CookingEvent.DISH_SERVED, KitchenActivity.OVEN_PREHEATED) + .requireAny(CookingEvent.SEASONING_ADJUSTED, group1, group2) + .prerequisite(CookingEvent.SEASONING_ADJUSTED, CookingEvent.DISH_SERVED) + .executionLimit(CookingEvent.DISH_SERVED, 1) + .executionLimit(CookingEvent.SEASONING_ADJUSTED, 1) + .build(); + + CountDownLatch seasoningLatch = new CountDownLatch(1); + AtomicInteger seasoningCount = new AtomicInteger(0); + + EventDelivery delivery = (handler, event, metadata) -> { + handler.handle(event, metadata); + if (event == CookingEvent.SEASONING_ADJUSTED) { + seasoningLatch.countDown(); + } + }; + + EventsManager eventsManager = new EventsManagerCore<>(config, delivery); + eventsManager.register(CookingEvent.SEASONING_ADJUSTED, (event, metadata) -> seasoningCount.incrementAndGet()); + + // Complete group 2 for SEASONING_ADJUSTED, but DISH_SERVED not fired yet + eventsManager.notifyInternalEvent(KitchenActivity.LEFTOVER_MEAT_FOUND, null); + eventsManager.notifyInternalEvent(KitchenActivity.LEFTOVER_VEGGIES_FOUND, null); + + // Wait and verify SEASONING_ADJUSTED not fired + eventsManager.eventAlreadyTriggered(CookingEvent.SEASONING_ADJUSTED); + assertEquals(0, seasoningCount.get()); + + // Now trigger DISH_SERVED + eventsManager.notifyInternalEvent(KitchenActivity.OVEN_PREHEATED, null); + + // Wait for DISH_SERVED + assertTrue(eventsManager.eventAlreadyTriggered(CookingEvent.DISH_SERVED)); + + // Now trigger something that completes group 1 + eventsManager.notifyInternalEvent(KitchenActivity.SEASONING_ADDED, null); + + assertTrue(seasoningLatch.await(TIMEOUT_MS, TimeUnit.MILLISECONDS)); + assertEquals(1, seasoningCount.get()); + } + + @Test + public void requireAnyGroupedWithSuppressor() throws InterruptedException { + Set group1 = new HashSet<>(); + group1.add(KitchenActivity.TIMEOUT_REACHED); + + EventsManagerConfig config = EventsManagerConfig.builder() + .requireAny(CookingEvent.DISH_SERVED, KitchenActivity.OVEN_PREHEATED) + .requireAny(CookingEvent.ORDER_TIMED_OUT, group1) + .suppressedBy(CookingEvent.ORDER_TIMED_OUT, CookingEvent.DISH_SERVED) + .executionLimit(CookingEvent.DISH_SERVED, 1) + .executionLimit(CookingEvent.ORDER_TIMED_OUT, 1) + .build(); + + AtomicInteger timeoutCount = new AtomicInteger(0); + + EventsManager eventsManager = new EventsManagerCore<>(config, SIMPLE_DELIVERY); + eventsManager.register(CookingEvent.ORDER_TIMED_OUT, (event, metadata) -> timeoutCount.incrementAndGet()); + + // Trigger DISH_SERVED first + eventsManager.notifyInternalEvent(KitchenActivity.OVEN_PREHEATED, null); + assertTrue(eventsManager.eventAlreadyTriggered(CookingEvent.DISH_SERVED)); + + // Now trigger timeout - should be suppressed + eventsManager.notifyInternalEvent(KitchenActivity.TIMEOUT_REACHED, null); + + // Wait for processing + eventsManager.eventAlreadyTriggered(CookingEvent.ORDER_TIMED_OUT); + + assertEquals(0, timeoutCount.get()); + assertFalse(eventsManager.eventAlreadyTriggered(CookingEvent.ORDER_TIMED_OUT)); + } + + @Test + public void prerequisiteChainResolvedInSingleNotification() throws InterruptedException { + // DISH_SERVED fires when OVEN_PREHEATED + // SEASONING_ADJUSTED fires when OVEN_PREHEATED, but requires DISH_SERVED first + EventsManagerConfig config = EventsManagerConfig.builder() + .requireAny(CookingEvent.DISH_SERVED, KitchenActivity.OVEN_PREHEATED) + .requireAny(CookingEvent.SEASONING_ADJUSTED, KitchenActivity.OVEN_PREHEATED) + .prerequisite(CookingEvent.SEASONING_ADJUSTED, CookingEvent.DISH_SERVED) + .executionLimit(CookingEvent.DISH_SERVED, 1) + .executionLimit(CookingEvent.SEASONING_ADJUSTED, 1) + .build(); + + CountDownLatch bothFiredLatch = new CountDownLatch(2); + AtomicInteger dishServedCount = new AtomicInteger(0); + AtomicInteger seasoningCount = new AtomicInteger(0); + + EventDelivery delivery = (handler, event, metadata) -> { + handler.handle(event, metadata); + bothFiredLatch.countDown(); + }; + + EventsManager eventsManager = new EventsManagerCore<>(config, delivery); + + eventsManager.register(CookingEvent.DISH_SERVED, (event, metadata) -> dishServedCount.incrementAndGet()); + eventsManager.register(CookingEvent.SEASONING_ADJUSTED, (event, metadata) -> seasoningCount.incrementAndGet()); + + // Single notification should trigger both events (A fires, then B fires because prerequisite is now met) + eventsManager.notifyInternalEvent(KitchenActivity.OVEN_PREHEATED, null); + + assertTrue("Both events should fire from single notification", bothFiredLatch.await(TIMEOUT_MS, TimeUnit.MILLISECONDS)); + assertEquals(1, dishServedCount.get()); + assertEquals(1, seasoningCount.get()); + assertTrue(eventsManager.eventAlreadyTriggered(CookingEvent.DISH_SERVED)); + assertTrue(eventsManager.eventAlreadyTriggered(CookingEvent.SEASONING_ADJUSTED)); + } + + @Test + public void prerequisiteChainWithOrOfAndsGroups() throws InterruptedException { + // DISH_SERVED = SDK_READY_FROM_CACHE (fires when sync group completes) + // LEFTOVERS_HEATED = SDK_READY (fires when sync completes, but requires DISH_SERVED first) + + Set syncGroup = new HashSet<>(); + syncGroup.add(KitchenActivity.INGREDIENTS_PREPPED); + syncGroup.add(KitchenActivity.SEASONING_ADDED); + + Set cacheGroup = new HashSet<>(); + cacheGroup.add(KitchenActivity.LEFTOVER_MEAT_FOUND); + cacheGroup.add(KitchenActivity.LEFTOVER_VEGGIES_FOUND); + + EventsManagerConfig config = EventsManagerConfig.builder() + // DISH_SERVED fires when either sync or cache group completes + .requireAny(CookingEvent.DISH_SERVED, syncGroup, cacheGroup) + // LEFTOVERS_HEATED requires the same sync events, but also DISH_SERVED as prerequisite + .requireAll(CookingEvent.LEFTOVERS_HEATED, + KitchenActivity.INGREDIENTS_PREPPED, + KitchenActivity.SEASONING_ADDED) + .prerequisite(CookingEvent.LEFTOVERS_HEATED, CookingEvent.DISH_SERVED) + .executionLimit(CookingEvent.DISH_SERVED, 1) + .executionLimit(CookingEvent.LEFTOVERS_HEATED, 1) + .build(); + + CountDownLatch bothFiredLatch = new CountDownLatch(2); + AtomicInteger dishServedCount = new AtomicInteger(0); + AtomicInteger leftoversCount = new AtomicInteger(0); + + EventDelivery delivery = (handler, event, metadata) -> { + handler.handle(event, metadata); + bothFiredLatch.countDown(); + }; + + EventsManager eventsManager = new EventsManagerCore<>(config, delivery); + + eventsManager.register(CookingEvent.DISH_SERVED, (event, metadata) -> dishServedCount.incrementAndGet()); + eventsManager.register(CookingEvent.LEFTOVERS_HEATED, (event, metadata) -> leftoversCount.incrementAndGet()); + + // First sync event + eventsManager.notifyInternalEvent(KitchenActivity.INGREDIENTS_PREPPED, null); + + // Second sync event should trigger chain: DISH_SERVED -> LEFTOVERS_HEATED + eventsManager.notifyInternalEvent(KitchenActivity.SEASONING_ADDED, null); + + assertTrue("Both events should fire when sync completes", bothFiredLatch.await(TIMEOUT_MS, TimeUnit.MILLISECONDS)); + assertEquals(1, dishServedCount.get()); + assertEquals(1, leftoversCount.get()); + assertTrue(eventsManager.eventAlreadyTriggered(CookingEvent.DISH_SERVED)); + assertTrue(eventsManager.eventAlreadyTriggered(CookingEvent.LEFTOVERS_HEATED)); + } + + @Test + public void prerequisiteLoopTerminatesWhenNoMoreEventsCanFire() throws InterruptedException { + // Create a chain where only DISH_SERVED can fire (SEASONING_ADJUSTED requires a different trigger) + EventsManagerConfig config = EventsManagerConfig.builder() + .requireAny(CookingEvent.DISH_SERVED, KitchenActivity.OVEN_PREHEATED) + .requireAny(CookingEvent.SEASONING_ADJUSTED, KitchenActivity.SEASONING_ADDED) // Different trigger! + .prerequisite(CookingEvent.SEASONING_ADJUSTED, CookingEvent.DISH_SERVED) + .executionLimit(CookingEvent.DISH_SERVED, 1) + .executionLimit(CookingEvent.SEASONING_ADJUSTED, 1) + .build(); + + CountDownLatch dishServedLatch = new CountDownLatch(1); + AtomicInteger dishServedCount = new AtomicInteger(0); + AtomicInteger seasoningCount = new AtomicInteger(0); + + EventDelivery delivery = (handler, event, metadata) -> { + handler.handle(event, metadata); + if (event == CookingEvent.DISH_SERVED) { + dishServedLatch.countDown(); + } + }; + + EventsManager eventsManager = new EventsManagerCore<>(config, delivery); + + eventsManager.register(CookingEvent.DISH_SERVED, (event, metadata) -> dishServedCount.incrementAndGet()); + eventsManager.register(CookingEvent.SEASONING_ADJUSTED, (event, metadata) -> seasoningCount.incrementAndGet()); + + // Only DISH_SERVED should fire, loop should terminate without firing SEASONING_ADJUSTED + eventsManager.notifyInternalEvent(KitchenActivity.OVEN_PREHEATED, null); + + assertTrue(dishServedLatch.await(TIMEOUT_MS, TimeUnit.MILLISECONDS)); + assertEquals(1, dishServedCount.get()); + assertEquals(0, seasoningCount.get()); // Should NOT fire - different trigger + + // Verify processing completed (no infinite loop) + assertTrue(eventsManager.eventAlreadyTriggered(CookingEvent.DISH_SERVED)); + assertFalse(eventsManager.eventAlreadyTriggered(CookingEvent.SEASONING_ADJUSTED)); + } + + @Test + public void threeLevelPrerequisiteChain() throws InterruptedException { + // DISH_SERVED -> SEASONING_ADJUSTED -> ORDER_TIMED_OUT + // All triggered by the same internal event + EventsManagerConfig config = EventsManagerConfig.builder() + .requireAny(CookingEvent.DISH_SERVED, KitchenActivity.OVEN_PREHEATED) + .requireAny(CookingEvent.SEASONING_ADJUSTED, KitchenActivity.OVEN_PREHEATED) + .requireAny(CookingEvent.ORDER_TIMED_OUT, KitchenActivity.OVEN_PREHEATED) + .prerequisite(CookingEvent.SEASONING_ADJUSTED, CookingEvent.DISH_SERVED) + .prerequisite(CookingEvent.ORDER_TIMED_OUT, CookingEvent.SEASONING_ADJUSTED) + .executionLimit(CookingEvent.DISH_SERVED, 1) + .executionLimit(CookingEvent.SEASONING_ADJUSTED, 1) + .executionLimit(CookingEvent.ORDER_TIMED_OUT, 1) + .build(); + + CountDownLatch allFiredLatch = new CountDownLatch(3); + AtomicInteger dishServedCount = new AtomicInteger(0); + AtomicInteger seasoningCount = new AtomicInteger(0); + AtomicInteger timeoutCount = new AtomicInteger(0); + + EventDelivery delivery = (handler, event, metadata) -> { + handler.handle(event, metadata); + allFiredLatch.countDown(); + }; + + EventsManager eventsManager = new EventsManagerCore<>(config, delivery); + + eventsManager.register(CookingEvent.DISH_SERVED, (event, metadata) -> dishServedCount.incrementAndGet()); + eventsManager.register(CookingEvent.SEASONING_ADJUSTED, (event, metadata) -> seasoningCount.incrementAndGet()); + eventsManager.register(CookingEvent.ORDER_TIMED_OUT, (event, metadata) -> timeoutCount.incrementAndGet()); + + // Single notification should trigger all three events in chain + eventsManager.notifyInternalEvent(KitchenActivity.OVEN_PREHEATED, null); + + assertTrue("All three events should fire from single notification", allFiredLatch.await(TIMEOUT_MS, TimeUnit.MILLISECONDS)); + assertEquals(1, dishServedCount.get()); + assertEquals(1, seasoningCount.get()); + assertEquals(1, timeoutCount.get()); + } }