Class MultiEventBus
-
- All Implemented Interfaces:
public final class MultiEventBusA custom event bus implementation allowing subscription to specific event IDs and publishing events to subscribers. This bus supports weak references to subscribers to avoid strong references that could prevent garbage collection.
neelabh.parui
-
-
Constructor Summary
Constructors Constructor Description MultiEventBus()
-
Method Summary
Modifier and Type Method Description final Unitsubscribe(Integer eventIds, String tag, Function1<BusEvent, Unit> subscriber)Subscribe to events based on the provided event IDs. final Unitunsubscribe(Function1<BusEvent, Unit> subscriber, String tag)Unsubscribe to all events. final Map<String, Pair<Long, Long>>snapshotAndReset(Long minIntervalMs)Snapshots the cumulative per-tag subscribe/unsubscribe counts and resets them, but only if at least minIntervalMs has elapsed since the last successful snapshot - regardless of whether the periodic call-count trigger (in recordCall) or the app-foreground trigger is asking. final UnitpublishNow(BusEvent event)Publishes the provided event to subscribers based on event filters on the current thread itself. final Unitpublish(BusEvent event)Publishes the provided event to subscribers based on event filters. final ConcurrentHashMap<Function1<BusEvent, Boolean>, WeakReference<Function1<BusEvent, Unit>>>getSubscribers()final UnitsetExecutor(ExecutorService executor)-
-
Method Detail
-
subscribe
final Unit subscribe(Integer eventIds, String tag, Function1<BusEvent, Unit> subscriber)
Subscribe to events based on the provided event IDs.
- Parameters:
eventIds- Variable number of event IDs to subscribe to.tag- Optional caller identifier used to attribute subscribe/unsubscribe counts for leak detection.subscriber- The subscriber function to be invoked when the subscribed events occur.
-
unsubscribe
final Unit unsubscribe(Function1<BusEvent, Unit> subscriber, String tag)
Unsubscribe to all events.
- Parameters:
subscriber- The subscriber function to be removedtag- Optional caller identifier used to attribute subscribe/unsubscribe counts for leak detection.
-
snapshotAndReset
final Map<String, Pair<Long, Long>> snapshotAndReset(Long minIntervalMs)
Snapshots the cumulative per-tag subscribe/unsubscribe counts and resets them, but only if at least minIntervalMs has elapsed since the last successful snapshot - regardless of whether the periodic call-count trigger (in recordCall) or the app-foreground trigger is asking. This prevents a bursty faulty session from being reported many times in quick succession: skipped windows leave counts untouched so they keep accumulating for the next allowed trigger. Never blocks the very first snapshot, since lastEmitAt starts at zero. Guarded by leakAuditLock so a concurrent subscribe/unsubscribe cannot record a call in the window between reading tags and resetting counts, which would otherwise be silently lost.
- Returns:
a map of tag to (subscribeCount, unsubscribeCount), or null if skipped due to cooldown.
-
publishNow
final Unit publishNow(BusEvent event)
Publishes the provided event to subscribers based on event filters on the current thread itself. This is a blocking call and should be called only if we want all the subscribers to execute before returning. It executes the publication on the same thread.
- Parameters:
event- The event to be published.
-
publish
final Unit publish(BusEvent event)
Publishes the provided event to subscribers based on event filters. It executes the publication on a separate thread using an Executor service.
- Parameters:
event- The event to be published.
-
getSubscribers
final ConcurrentHashMap<Function1<BusEvent, Boolean>, WeakReference<Function1<BusEvent, Unit>>> getSubscribers()
-
setExecutor
final Unit setExecutor(ExecutorService executor)
-
-
-
-