androidinterview.com

Low Level Design (LLD) Interview Questions

Implement an Event Bus

Tier: CommonDifficulty: EasyAsked of: Junior, Mid

An event bus lets one part of an app announce an event to interested parts of the app. A publisher sends the event. Subscribers register callbacks to receive it.

The problem

A cart publishes ItemAdded. A badge updates its count and an analytics callback records the action. Unsubscribing the badge should stop its updates while analytics continues receiving events.

Start with one thread, immediate delivery and exact event type matching. A subscriber to a parent class does not receive subclass events in this version. Old events are not replayed. If a callback throws, the error reaches the publisher and the remaining callbacks are not called.

How to explain the design

“I keep a map from event type to a list of subscribers. Publishing an event finds that list and calls each subscriber. Subscribing returns a function that removes that subscription. I copy the list before delivery so callbacks can unsubscribe without breaking the loop.”

Each subscription has an active flag. If one callback removes another subscription during delivery, the removed subscription is skipped even though it was in the copied list.

Walk through a publish

  1. Register two callbacks for ItemAdded.
  2. Publish an ItemAdded event. Both callbacks receive it.
  3. Call the badge's unsubscribe function.
  4. Publish again. Only the analytics callback receives it.

Interview implementation

EventBus owns the subscriber map. Subscription stores the callback and whether it is still active. The caller is responsible for unsubscribing when its screen finishes.

Java

EventBus.java

package interview.eventbus;

import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.function.Consumer;

public class EventBus {
    private static class Subscription {
        final Consumer<Object> receive;
        boolean active = true;
        Subscription(Consumer<Object> receive) { this.receive = receive; }
    }
    private final Map<Class<?>, List<Subscription>> subscribers = new HashMap<>();

    public <T> Runnable subscribe(Class<T> type, Consumer<T> handler) {
        Subscription subscription = new Subscription(event -> handler.accept(type.cast(event)));
        subscribers.computeIfAbsent(type, key -> new ArrayList<>()).add(subscription);
        return () -> {
            subscription.active = false;
            List<Subscription> list = subscribers.get(type);
            if (list != null) {
                list.remove(subscription);
                if (list.isEmpty()) subscribers.remove(type);
            }
        };
    }

    // Deliver on the caller's thread to this exact event type.
    public void publish(Object event) {
        var snapshot = new ArrayList<>(subscribers.getOrDefault(event.getClass(), List.of()));
        for (Subscription subscription : snapshot) {
            if (subscription.active) subscription.receive.accept(event);
        }
    }
}
package interview.eventbus;

import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.function.Consumer;

public class EventBus {
    private static class Subscription {
        final Consumer<Object> receive;
        boolean active = true;
        Subscription(Consumer<Object> receive) { this.receive = receive; }
    }
    private final Map<Class<?>, List<Subscription>> subscribers = new HashMap<>();

    public <T> Runnable subscribe(Class<T> type, Consumer<T> handler) {
        Subscription subscription = new Subscription(event -> handler.accept(type.cast(event)));
        subscribers.computeIfAbsent(type, key -> new ArrayList<>()).add(subscription);
        return () -> {
            subscription.active = false;
            List<Subscription> list = subscribers.get(type);
            if (list != null) {
                list.remove(subscription);
                if (list.isEmpty()) subscribers.remove(type);
            }
        };
    }

    public void publish(Object event) {
        var snapshot = new ArrayList<>(subscribers.getOrDefault(event.getClass(), List.of()));
        for (Subscription subscription : snapshot) {
            if (subscription.active) subscription.receive.accept(event);
        }
    }
}

Kotlin

EventBus.kt

package interview.eventbus

class EventBus {
    private class Subscription(val receive: (Any) -> Unit) {
        var active = true
    }
    private val subscribers = mutableMapOf<Class<*>, MutableList<Subscription>>()

    fun <T : Any> subscribe(type: Class<T>, handler: (T) -> Unit): () -> Unit {
        val subscription = Subscription { event -> handler(type.cast(event)) }
        subscribers.getOrPut(type) { mutableListOf() }.add(subscription)
        return {
            subscription.active = false
            subscribers[type]?.remove(subscription)
            if (subscribers[type]?.isEmpty() == true) subscribers.remove(type)
        }
    }

    // Deliver on the caller's thread to this exact event type.
    fun publish(event: Any) {
        val snapshot = subscribers[event.javaClass]?.toList() ?: return
        for (subscription in snapshot) {
            if (subscription.active) subscription.receive(event)
        }
    }
}
package interview.eventbus

class EventBus {
    private class Subscription(val receive: (Any) -> Unit) {
        var active = true
    }
    private val subscribers = mutableMapOf<Class<*>, MutableList<Subscription>>()

    fun <T : Any> subscribe(type: Class<T>, handler: (T) -> Unit): () -> Unit {
        val subscription = Subscription { event -> handler(type.cast(event)) }
        subscribers.getOrPut(type) { mutableListOf() }.add(subscription)
        return {
            subscription.active = false
            subscribers[type]?.remove(subscription)
            if (subscribers[type]?.isEmpty() == true) subscribers.remove(type)
        }
    }

    fun publish(event: Any) {
        val snapshot = subscribers[event.javaClass]?.toList() ?: return
        for (subscription in snapshot) {
            if (subscription.active) subscription.receive(event)
        }
    }
}

Follow-up questions

Delivery on another thread?

“I would accept an executor and submit each callback to it.” An executor decides where a task runs, such as a worker thread or the UI thread. Define cancellation first. A useful rule is to skip a queued callback if it has not started, while allowing a callback already running to finish. Protect the subscriber map and use a thread safe active flag when publishing and unsubscribing can happen on different threads.

One callback fails?

“The current version stops and lets the caller see the exception. If subscribers should be independent, I would report that failure and continue to the next subscriber.” Replace the delivery loop with this fragment. The injected report callback records the exception and must return normally.

Kotlin

for (subscription in snapshot) {
    if (!subscription.active) continue
    try {
        subscription.receive(event)
    } catch (error: Exception) {
        report(error)
    }
}

Java

for (Subscription subscription : snapshot) {
    if (!subscription.active) continue;
    try {
        subscription.receive.accept(event);
    } catch (Exception error) {
        report.accept(error);
    }
}

A callback publishes another event?

“This version delivers the nested event immediately, before finishing the first event's remaining callbacks.” If that ordering is confusing, put new events in a queue and let one loop drain it. When delivery is already running, publishing only adds to the queue. This avoids growing the call stack, though a subscriber that keeps publishing forever still needs to be fixed.

What should I test?

“Two subscribers for the same exact event type should each receive it once, and subscribers for another type should receive nothing.” Publishing with no subscribers should do nothing. Unsubscribing twice should be safe. If an earlier callback removes a later subscriber, that later callback should be skipped. For the continue-on-error extension, a throwing callback should be reported and the next callback should still run.

Extended implementation and optional features

This reference explores a larger scope. Use it after you can explain and write the interview version. Its extra types and features are not required for the scope above.

Java

com.androidinterview.eventbus.EventBus.java

package com.androidinterview.eventbus;

import java.util.ArrayDeque;
import java.util.ArrayList;
import java.util.Collections;
import java.util.Deque;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.CopyOnWriteArrayList;
import java.util.concurrent.Executor;

import com.androidinterview.eventbus.lifecycle.LifecycleOwner;
import com.androidinterview.eventbus.thread.DirectExecutor;

// The bus. Three pieces of state and nothing else.
//
// subscribers is the registry, one list per event type. It is a
// ConcurrentHashMap of copy on write lists because the read to write ratio is
// extreme. Posts happen constantly, subscriptions happen when a screen opens.
// A copy on write list gives every post a stable snapshot to walk with no
// lock at all, which is also what makes a handler safe to subscribe or
// unsubscribe from inside a delivery.
//
// typeCache holds the flattened supertypes of each concrete event class, so
// the class hierarchy is walked once per event class rather than once per
// post.
//
// sticky holds the last event posted for each concrete type. It is a plain
// LinkedHashMap guarded by its own monitor, because the store and the replay
// have to be one decision and iteration order has to be stable.
public final class EventBus {

    private final Map<Class<?>, CopyOnWriteArrayList<Subscriber<?>>> subscribers = new ConcurrentHashMap<>();
    private final Map<Class<?>, List<Class<?>>> typeCache = new ConcurrentHashMap<>();
    private final Map<Class<?>, Object> sticky = new LinkedHashMap<>();

    private final Executor defaultExecutor;

    public EventBus() {
        this(DirectExecutor.INSTANCE);
    }

    public EventBus(Executor defaultExecutor) {
        this.defaultExecutor = defaultExecutor;
    }

    public <T> Subscription subscribe(Class<T> type, EventHandler<? super T> handler) {
        return subscribe(type, defaultExecutor, handler);
    }

    // The registration and the sticky replay happen under one lock, and that
    // pairing is the whole trick. Either this subscriber is already in the
    // registry when a sticky event is stored, in which case the post finds it
    // and the replay below sees the older value, or it registers afterwards
    // and the replay hands it the new one. It cannot see both and it cannot
    // miss both.
    //
    // Delivery itself happens outside the lock, because a handler is caller
    // code and holding a lock across it invites a deadlock.
    public <T> Subscription subscribe(Class<T> type, Executor on, EventHandler<? super T> handler) {
        Subscriber<T> subscriber = new Subscriber<>(this, type, on, handler);
        List<Object> replay = new ArrayList<>();
        synchronized (sticky) {
            subscribers.computeIfAbsent(type, key -> new CopyOnWriteArrayList<>()).add(subscriber);
            for (Object event : sticky.values()) {
                if (type.isInstance(event)) {
                    replay.add(event);
                }
            }
        }
        for (Object event : replay) {
            subscriber.deliver(event);
        }
        return subscriber;
    }

    public <T> Subscription subscribe(LifecycleOwner owner, Class<T> type, EventHandler<? super T> handler) {
        return subscribe(owner, type, defaultExecutor, handler);
    }

    // The overload every screen should be using. Forgetting to unsubscribe is
    // the classic event bus leak, because the bus is a long lived object and
    // a handler written inside an Activity holds that Activity. Tying the
    // subscription to the owner makes the leak impossible to write rather
    // than merely documented.
    public <T> Subscription subscribe(
            LifecycleOwner owner, Class<T> type, Executor on, EventHandler<? super T> handler) {
        Subscription subscription = subscribe(type, on, handler);
        owner.lifecycle().addOnDestroy(subscription::unsubscribe);
        return subscription;
    }

    // Collect the targets first, then deliver. Collecting produces a snapshot,
    // so a handler that subscribes or unsubscribes during delivery changes
    // what the next post sees and never what this one is halfway through.
    public void post(Object event) {
        for (Subscriber<?> subscriber : targetsOf(event)) {
            subscriber.deliver(event);
        }
    }

    // Same as post, and it also remembers the event so a subscriber that
    // arrives later is told the current answer instead of waiting for the
    // next change. Useful for state, wrong for anything that reads as a
    // one off instruction, because that instruction will fire again.
    public void postSticky(Object event) {
        List<Subscriber<?>> targets;
        synchronized (sticky) {
            sticky.put(event.getClass(), event);
            targets = targetsOf(event);
        }
        for (Subscriber<?> subscriber : targets) {
            subscriber.deliver(event);
        }
    }

    // Whoever posts a sticky event owns clearing it. A dialog request that is
    // never removed is replayed to the next screen that subscribes, and that
    // is the bug the sticky section of the answer is about.
    public void removeSticky(Class<?> type) {
        synchronized (sticky) {
            sticky.keySet().removeIf(type::isAssignableFrom);
        }
    }

    void remove(Subscriber<?> subscriber) {
        CopyOnWriteArrayList<Subscriber<?>> list = subscribers.get(subscriber.type());
        if (list != null) {
            list.remove(subscriber);
        }
    }

    // Most specific first. Subscribers to the concrete class hear about the
    // event before subscribers to its supertypes, which is the order a reader
    // expects when a specific handler and a generic logger both exist.
    private List<Subscriber<?>> targetsOf(Object event) {
        List<Subscriber<?>> targets = new ArrayList<>();
        for (Class<?> type : typeCache.computeIfAbsent(event.getClass(), EventBus::flatten)) {
            CopyOnWriteArrayList<Subscriber<?>> list = subscribers.get(type);
            if (list != null) {
                targets.addAll(list);
            }
        }
        return targets;
    }

    // A breadth first walk of the class and its interfaces, so a subscription
    // to a sealed interface receives every implementation. Object is skipped,
    // because a subscription to Object is a subscription to everything and it
    // is never what anybody meant.
    private static List<Class<?>> flatten(Class<?> eventType) {
        List<Class<?>> flat = new ArrayList<>();
        Deque<Class<?>> queue = new ArrayDeque<>();
        queue.add(eventType);
        while (!queue.isEmpty()) {
            Class<?> next = queue.poll();
            if (next == Object.class || flat.contains(next)) {
                continue;
            }
            flat.add(next);
            if (next.getSuperclass() != null) {
                queue.add(next.getSuperclass());
            }
            Collections.addAll(queue, next.getInterfaces());
        }
        return List.copyOf(flat);
    }
}
package com.androidinterview.eventbus;

import java.util.ArrayDeque;
import java.util.ArrayList;
import java.util.Collections;
import java.util.Deque;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.CopyOnWriteArrayList;
import java.util.concurrent.Executor;

import com.androidinterview.eventbus.lifecycle.LifecycleOwner;
import com.androidinterview.eventbus.thread.DirectExecutor;

public final class EventBus {

    private final Map<Class<?>, CopyOnWriteArrayList<Subscriber<?>>> subscribers = new ConcurrentHashMap<>();
    private final Map<Class<?>, List<Class<?>>> typeCache = new ConcurrentHashMap<>();
    private final Map<Class<?>, Object> sticky = new LinkedHashMap<>();

    private final Executor defaultExecutor;

    public EventBus() {
        this(DirectExecutor.INSTANCE);
    }

    public EventBus(Executor defaultExecutor) {
        this.defaultExecutor = defaultExecutor;
    }

    public <T> Subscription subscribe(Class<T> type, EventHandler<? super T> handler) {
        return subscribe(type, defaultExecutor, handler);
    }

    public <T> Subscription subscribe(Class<T> type, Executor on, EventHandler<? super T> handler) {
        Subscriber<T> subscriber = new Subscriber<>(this, type, on, handler);
        List<Object> replay = new ArrayList<>();
        synchronized (sticky) {
            subscribers.computeIfAbsent(type, key -> new CopyOnWriteArrayList<>()).add(subscriber);
            for (Object event : sticky.values()) {
                if (type.isInstance(event)) {
                    replay.add(event);
                }
            }
        }
        for (Object event : replay) {
            subscriber.deliver(event);
        }
        return subscriber;
    }

    public <T> Subscription subscribe(LifecycleOwner owner, Class<T> type, EventHandler<? super T> handler) {
        return subscribe(owner, type, defaultExecutor, handler);
    }

    public <T> Subscription subscribe(
            LifecycleOwner owner, Class<T> type, Executor on, EventHandler<? super T> handler) {
        Subscription subscription = subscribe(type, on, handler);
        owner.lifecycle().addOnDestroy(subscription::unsubscribe);
        return subscription;
    }

    public void post(Object event) {
        for (Subscriber<?> subscriber : targetsOf(event)) {
            subscriber.deliver(event);
        }
    }

    public void postSticky(Object event) {
        List<Subscriber<?>> targets;
        synchronized (sticky) {
            sticky.put(event.getClass(), event);
            targets = targetsOf(event);
        }
        for (Subscriber<?> subscriber : targets) {
            subscriber.deliver(event);
        }
    }

    public void removeSticky(Class<?> type) {
        synchronized (sticky) {
            sticky.keySet().removeIf(type::isAssignableFrom);
        }
    }

    void remove(Subscriber<?> subscriber) {
        CopyOnWriteArrayList<Subscriber<?>> list = subscribers.get(subscriber.type());
        if (list != null) {
            list.remove(subscriber);
        }
    }

    private List<Subscriber<?>> targetsOf(Object event) {
        List<Subscriber<?>> targets = new ArrayList<>();
        for (Class<?> type : typeCache.computeIfAbsent(event.getClass(), EventBus::flatten)) {
            CopyOnWriteArrayList<Subscriber<?>> list = subscribers.get(type);
            if (list != null) {
                targets.addAll(list);
            }
        }
        return targets;
    }

    private static List<Class<?>> flatten(Class<?> eventType) {
        List<Class<?>> flat = new ArrayList<>();
        Deque<Class<?>> queue = new ArrayDeque<>();
        queue.add(eventType);
        while (!queue.isEmpty()) {
            Class<?> next = queue.poll();
            if (next == Object.class || flat.contains(next)) {
                continue;
            }
            flat.add(next);
            if (next.getSuperclass() != null) {
                queue.add(next.getSuperclass());
            }
            Collections.addAll(queue, next.getInterfaces());
        }
        return List.copyOf(flat);
    }
}

com.androidinterview.eventbus.EventHandler.java

package com.androidinterview.eventbus;

// What a subscriber actually is, one method taking one event. Keeping it a
// single method interface means a lambda works everywhere a handler is asked
// for, and it keeps the bus from ever needing to know about annotations or
// reflection.
//
// The Kotlin side has no equivalent file, because a function type says the
// same thing without a declaration.
@FunctionalInterface
public interface EventHandler<T> {
    void onEvent(T event);
}
package com.androidinterview.eventbus;

@FunctionalInterface
public interface EventHandler<T> {
    void onEvent(T event);
}

com.androidinterview.eventbus.Subscriber.java

package com.androidinterview.eventbus;

import java.util.concurrent.Executor;
import java.util.concurrent.atomic.AtomicBoolean;

// One registration. It knows the type it asked for, the handler to call and
// the executor to call it on. It is package private on purpose, because
// callers only ever hold it as a Subscription.
final class Subscriber<T> implements Subscription {

    private final EventBus bus;
    private final Class<T> type;
    private final Executor executor;
    private final EventHandler<? super T> handler;
    private final AtomicBoolean alive = new AtomicBoolean(true);

    Subscriber(EventBus bus, Class<T> type, Executor executor, EventHandler<? super T> handler) {
        this.bus = bus;
        this.type = type;
        this.executor = executor;
        this.handler = handler;
    }

    Class<T> type() {
        return type;
    }

    // The alive flag is read twice, and that is deliberate. Once here, so a
    // cancelled subscription is skipped at post time, and once inside the
    // task, because an executor may run the task long after unsubscribe
    // returned. Without the second check, a screen that unsubscribed in
    // onDestroy can still be handed an event on the main thread afterwards.
    void deliver(Object event) {
        if (!alive.get()) {
            return;
        }
        T typed = type.cast(event);
        executor.execute(() -> {
            if (alive.get()) {
                handler.onEvent(typed);
            }
        });
    }

    // Compare and set, so a double unsubscribe is harmless and the removal
    // from the registry happens exactly once.
    @Override
    public void unsubscribe() {
        if (alive.compareAndSet(true, false)) {
            bus.remove(this);
        }
    }

    @Override
    public boolean isActive() {
        return alive.get();
    }
}
package com.androidinterview.eventbus;

import java.util.concurrent.Executor;
import java.util.concurrent.atomic.AtomicBoolean;

final class Subscriber<T> implements Subscription {

    private final EventBus bus;
    private final Class<T> type;
    private final Executor executor;
    private final EventHandler<? super T> handler;
    private final AtomicBoolean alive = new AtomicBoolean(true);

    Subscriber(EventBus bus, Class<T> type, Executor executor, EventHandler<? super T> handler) {
        this.bus = bus;
        this.type = type;
        this.executor = executor;
        this.handler = handler;
    }

    Class<T> type() {
        return type;
    }

    void deliver(Object event) {
        if (!alive.get()) {
            return;
        }
        T typed = type.cast(event);
        executor.execute(() -> {
            if (alive.get()) {
                handler.onEvent(typed);
            }
        });
    }

    @Override
    public void unsubscribe() {
        if (alive.compareAndSet(true, false)) {
            bus.remove(this);
        }
    }

    @Override
    public boolean isActive() {
        return alive.get();
    }
}

com.androidinterview.eventbus.Subscription.java

package com.androidinterview.eventbus;

// The handle that subscribe hands back, and the reason the bus has no
// unsubscribe method that takes a type and a handler.
//
// A bus that asks you to hand back what you registered forces every caller to
// keep the type and the lambda in fields just so it can find its own
// registration again, and the moment a lambda is written inline the caller
// cannot cancel at all. One object you can cancel removes both problems.
public interface Subscription {

    void unsubscribe();

    boolean isActive();
}
package com.androidinterview.eventbus;

public interface Subscription {

    void unsubscribe();

    boolean isActive();
}

com.androidinterview.eventbus.events.AppEvent.java

package com.androidinterview.eventbus.events;

// A small event hierarchy, here only to show what supertype delivery buys.
// Subscribe to SignedIn and you get one kind of event. Subscribe to AppEvent
// and you get all three, which is how an analytics logger or a crash
// breadcrumb trail is written without touching any publisher.
//
// Events are values. They carry no callbacks and no references to a screen,
// because a bus holds the last sticky one for as long as the process lives.
public sealed interface AppEvent {

    record SignedIn(String userId) implements AppEvent {
    }

    record SignedOut(String reason) implements AppEvent {
    }

    record NetworkChanged(boolean online) implements AppEvent {
    }
}
package com.androidinterview.eventbus.events;

public sealed interface AppEvent {

    record SignedIn(String userId) implements AppEvent {
    }

    record SignedOut(String reason) implements AppEvent {
    }

    record NetworkChanged(boolean online) implements AppEvent {
    }
}

com.androidinterview.eventbus.lifecycle.Lifecycle.java

package com.androidinterview.eventbus.lifecycle;

import java.util.ArrayList;
import java.util.List;

// The smallest useful model of the Android lifecycle. The real one has a
// state machine and seven events. The bus needs two things from it, a way to
// be told when the owner is destroyed, and an honest answer to whether that
// has already happened.
public final class Lifecycle {

    private final List<Runnable> onDestroy = new ArrayList<>();
    private boolean destroyed;

    // Registering after destruction runs the action straight away. That case
    // looks unlikely and is not. A subscription made on a background thread
    // while the screen is finishing would otherwise never be cancelled, which
    // is the exact leak this class exists to stop.
    public void addOnDestroy(Runnable action) {
        synchronized (this) {
            if (!destroyed) {
                onDestroy.add(action);
                return;
            }
        }
        action.run();
    }

    // The actions run outside the lock. Unsubscribing takes the bus registry
    // locks, and holding two locks in two different orders is how a deadlock
    // gets written by accident.
    public void destroy() {
        List<Runnable> actions;
        synchronized (this) {
            if (destroyed) {
                return;
            }
            destroyed = true;
            actions = new ArrayList<>(onDestroy);
            onDestroy.clear();
        }
        for (Runnable action : actions) {
            action.run();
        }
    }

    public synchronized boolean isDestroyed() {
        return destroyed;
    }
}
package com.androidinterview.eventbus.lifecycle;

import java.util.ArrayList;
import java.util.List;

public final class Lifecycle {

    private final List<Runnable> onDestroy = new ArrayList<>();
    private boolean destroyed;

    public void addOnDestroy(Runnable action) {
        synchronized (this) {
            if (!destroyed) {
                onDestroy.add(action);
                return;
            }
        }
        action.run();
    }

    public void destroy() {
        List<Runnable> actions;
        synchronized (this) {
            if (destroyed) {
                return;
            }
            destroyed = true;
            actions = new ArrayList<>(onDestroy);
            onDestroy.clear();
        }
        for (Runnable action : actions) {
            action.run();
        }
    }

    public synchronized boolean isDestroyed() {
        return destroyed;
    }
}

com.androidinterview.eventbus.lifecycle.LifecycleOwner.java

package com.androidinterview.eventbus.lifecycle;

// Anything that can be destroyed, which on Android is an Activity, a Fragment
// or a view. The bus only ever sees this interface, so nothing in the design
// depends on the Android SDK and every test can destroy an owner on demand.
public interface LifecycleOwner {

    Lifecycle lifecycle();
}
package com.androidinterview.eventbus.lifecycle;

public interface LifecycleOwner {

    Lifecycle lifecycle();
}

com.androidinterview.eventbus.thread.DirectExecutor.java

package com.androidinterview.eventbus.thread;

import java.util.concurrent.Executor;

// Delivers on whichever thread called post. It is the cheapest option and the
// right one for a handler that only touches its own data.
//
// It is also where re-entrancy shows up. A handler that posts during delivery
// runs the second event all the way to completion inside the first, so the
// call stack grows and the observed order is depth first rather than the
// order the two events were posted in.
public final class DirectExecutor implements Executor {

    public static final DirectExecutor INSTANCE = new DirectExecutor();

    private DirectExecutor() {
    }

    @Override
    public void execute(Runnable task) {
        task.run();
    }
}
package com.androidinterview.eventbus.thread;

import java.util.concurrent.Executor;

public final class DirectExecutor implements Executor {

    public static final DirectExecutor INSTANCE = new DirectExecutor();

    private DirectExecutor() {
    }

    @Override
    public void execute(Runnable task) {
        task.run();
    }
}

com.androidinterview.eventbus.thread.MainThreadExecutor.java

package com.androidinterview.eventbus.thread;

import java.util.concurrent.Executor;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;

// Stands in for the Android main looper. A looper is one thread with a queue
// in front of it, which is exactly a single thread executor, so modelling it
// this way costs one class and keeps the bus free of the SDK.
public final class MainThreadExecutor implements Executor {

    private volatile Thread thread;

    private final ExecutorService delegate = Executors.newSingleThreadExecutor(runnable -> {
        Thread created = new Thread(runnable, "main");
        created.setDaemon(true);
        thread = created;
        return created;
    });

    // Everything goes through the queue, including work posted from the main
    // thread itself. Running those inline would let a handler that posts
    // during delivery jump ahead of events that were already waiting, and a
    // reordering like that is very hard to find later.
    @Override
    public void execute(Runnable task) {
        delegate.execute(task);
    }

    public boolean isMainThread() {
        return Thread.currentThread() == thread;
    }

    public void shutdown() {
        delegate.shutdown();
    }
}
package com.androidinterview.eventbus.thread;

import java.util.concurrent.Executor;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;

public final class MainThreadExecutor implements Executor {

    private volatile Thread thread;

    private final ExecutorService delegate = Executors.newSingleThreadExecutor(runnable -> {
        Thread created = new Thread(runnable, "main");
        created.setDaemon(true);
        thread = created;
        return created;
    });

    @Override
    public void execute(Runnable task) {
        delegate.execute(task);
    }

    public boolean isMainThread() {
        return Thread.currentThread() == thread;
    }

    public void shutdown() {
        delegate.shutdown();
    }
}

Kotlin

com.androidinterview.eventbus.EventBus.kt

package com.androidinterview.eventbus

import com.androidinterview.eventbus.lifecycle.LifecycleOwner
import com.androidinterview.eventbus.thread.DirectExecutor
import java.util.concurrent.ConcurrentHashMap
import java.util.concurrent.CopyOnWriteArrayList
import java.util.concurrent.Executor
import java.util.concurrent.atomic.AtomicBoolean

// The handle subscribe hands back, and the reason there is no unsubscribe
// that takes a type and a handler. A bus that asks you to hand back what you
// registered forces callers to keep the lambda in a field, and a lambda
// written inline could never be cancelled at all.
interface Subscription {
    val isActive: Boolean
    fun unsubscribe()
}

// One registration. Internal because callers only ever hold it as a
// Subscription. There is no EventHandler interface anywhere in the Kotlin,
// because a function type says the same thing without a declaration.
internal class Subscriber<T : Any>(
    private val bus: EventBus,
    val type: Class<T>,
    private val executor: Executor,
    private val handler: (T) -> Unit,
) : Subscription {

    private val alive = AtomicBoolean(true)

    override val isActive: Boolean get() = alive.get()

    // The alive flag is read twice on purpose. Once here, so a cancelled
    // subscription is skipped at post time, and once inside the task, because
    // an executor may run it long after unsubscribe returned. Without the
    // second check a screen that unsubscribed in onDestroy can still be
    // handed an event on the main thread afterwards.
    fun deliver(event: Any) {
        if (!alive.get()) return
        val typed = type.cast(event)
        executor.execute { if (alive.get()) handler(typed) }
    }

    // Compare and set, so a double unsubscribe is harmless and the registry
    // is touched exactly once.
    override fun unsubscribe() {
        if (alive.compareAndSet(true, false)) bus.remove(this)
    }
}

// The bus. Three pieces of state and nothing else.
//
// subscribers is the registry, one list per event type. A ConcurrentHashMap
// of copy on write lists, because posts happen constantly and subscriptions
// happen when a screen opens. Copy on write gives every post a stable
// snapshot with no lock, which is also what makes a handler safe to subscribe
// or unsubscribe from inside a delivery.
//
// typeCache holds the flattened supertypes of each concrete event class, so
// the hierarchy is walked once per class rather than once per post.
//
// sticky holds the last event of each concrete type, in a plain LinkedHashMap
// guarded by its own monitor, because the store and the replay have to be one
// decision and the iteration order has to be stable.
class EventBus(val defaultExecutor: Executor = DirectExecutor) {

    private val subscribers = ConcurrentHashMap<Class<*>, CopyOnWriteArrayList<Subscriber<*>>>()
    private val typeCache = ConcurrentHashMap<Class<*>, List<Class<*>>>()
    private val sticky = LinkedHashMap<Class<*>, Any>()

    // Reified, so the call site reads subscribe<SignedIn> { render(it) } with
    // no class literal in it. This is the one place the Kotlin is genuinely
    // nicer than the Java rather than just shorter.
    inline fun <reified T : Any> subscribe(
        on: Executor = defaultExecutor,
        noinline handler: (T) -> Unit,
    ): Subscription = subscribe(T::class.java, on, handler)

    // The overload every screen should use. Forgetting to unsubscribe is the
    // classic event bus leak, because the bus outlives every screen and a
    // lambda written inside an Activity holds that Activity. Binding the
    // subscription to the owner makes the leak impossible to write rather
    // than merely documented.
    inline fun <reified T : Any> subscribe(
        owner: LifecycleOwner,
        on: Executor = defaultExecutor,
        noinline handler: (T) -> Unit,
    ): Subscription = subscribe(owner, T::class.java, on, handler)

    // Registration and sticky replay under one lock, and that pairing is the
    // whole trick. Either this subscriber is already in the registry when a
    // sticky event is stored, so the post finds it and the replay below sees
    // the older value, or it registers afterwards and the replay hands it the
    // new one. It cannot see both and it cannot miss both.
    //
    // Delivery happens outside the lock, because a handler is caller code and
    // holding a lock across caller code invites a deadlock.
    fun <T : Any> subscribe(type: Class<T>, on: Executor, handler: (T) -> Unit): Subscription {
        val subscriber = Subscriber(this, type, on, handler)
        val replay = synchronized(sticky) {
            subscribers.getOrPut(type) { CopyOnWriteArrayList() }.add(subscriber)
            sticky.values.filter(type::isInstance)
        }
        replay.forEach(subscriber::deliver)
        return subscriber
    }

    fun <T : Any> subscribe(
        owner: LifecycleOwner,
        type: Class<T>,
        on: Executor,
        handler: (T) -> Unit,
    ): Subscription = subscribe(type, on, handler).also { owner.lifecycle.addOnDestroy(it::unsubscribe) }

    // Collect the targets, then deliver. Collecting produces a snapshot, so a
    // handler that subscribes or unsubscribes during delivery changes what the
    // next post sees and never what this one is halfway through.
    fun post(event: Any) = targetsOf(event).forEach { it.deliver(event) }

    // Post, and also remember the event so a subscriber arriving later is told
    // the current answer instead of waiting for the next change. Right for
    // state, wrong for anything that reads as a one off instruction, because
    // that instruction will fire again.
    fun postSticky(event: Any) {
        val targets = synchronized(sticky) {
            sticky[event.javaClass] = event
            targetsOf(event)
        }
        targets.forEach { it.deliver(event) }
    }

    // Whoever posts a sticky event owns clearing it. A dialog request left in
    // the map is replayed to the next screen that subscribes, which is the
    // sticky bug worth naming out loud.
    fun removeSticky(type: Class<*>) {
        synchronized(sticky) { sticky.keys.removeIf(type::isAssignableFrom) }
    }

    internal fun remove(subscriber: Subscriber<*>) {
        subscribers[subscriber.type]?.remove(subscriber)
    }

    // Most specific first, so a concrete handler hears about the event before
    // a generic logger subscribed to the sealed parent.
    private fun targetsOf(event: Any): List<Subscriber<*>> =
        typesOf(event.javaClass).flatMap { subscribers[it].orEmpty() }

    // A breadth first walk of the class and its interfaces, so a subscription
    // to a sealed interface receives every implementation. Any is skipped,
    // because subscribing to it is subscribing to everything and nobody means
    // that.
    private fun typesOf(eventType: Class<*>): List<Class<*>> = typeCache.getOrPut(eventType) {
        val flat = mutableListOf<Class<*>>()
        val queue = ArrayDeque<Class<*>>().apply { add(eventType) }
        while (queue.isNotEmpty()) {
            val next = queue.removeFirst()
            if (next == Any::class.java || next in flat) continue
            flat += next
            next.superclass?.let(queue::add)
            queue.addAll(next.interfaces)
        }
        flat
    }
}
package com.androidinterview.eventbus

import com.androidinterview.eventbus.lifecycle.LifecycleOwner
import com.androidinterview.eventbus.thread.DirectExecutor
import java.util.concurrent.ConcurrentHashMap
import java.util.concurrent.CopyOnWriteArrayList
import java.util.concurrent.Executor
import java.util.concurrent.atomic.AtomicBoolean

interface Subscription {
    val isActive: Boolean
    fun unsubscribe()
}

internal class Subscriber<T : Any>(
    private val bus: EventBus,
    val type: Class<T>,
    private val executor: Executor,
    private val handler: (T) -> Unit,
) : Subscription {

    private val alive = AtomicBoolean(true)

    override val isActive: Boolean get() = alive.get()

    fun deliver(event: Any) {
        if (!alive.get()) return
        val typed = type.cast(event)
        executor.execute { if (alive.get()) handler(typed) }
    }

    override fun unsubscribe() {
        if (alive.compareAndSet(true, false)) bus.remove(this)
    }
}

class EventBus(val defaultExecutor: Executor = DirectExecutor) {

    private val subscribers = ConcurrentHashMap<Class<*>, CopyOnWriteArrayList<Subscriber<*>>>()
    private val typeCache = ConcurrentHashMap<Class<*>, List<Class<*>>>()
    private val sticky = LinkedHashMap<Class<*>, Any>()

    inline fun <reified T : Any> subscribe(
        on: Executor = defaultExecutor,
        noinline handler: (T) -> Unit,
    ): Subscription = subscribe(T::class.java, on, handler)

    inline fun <reified T : Any> subscribe(
        owner: LifecycleOwner,
        on: Executor = defaultExecutor,
        noinline handler: (T) -> Unit,
    ): Subscription = subscribe(owner, T::class.java, on, handler)

    fun <T : Any> subscribe(type: Class<T>, on: Executor, handler: (T) -> Unit): Subscription {
        val subscriber = Subscriber(this, type, on, handler)
        val replay = synchronized(sticky) {
            subscribers.getOrPut(type) { CopyOnWriteArrayList() }.add(subscriber)
            sticky.values.filter(type::isInstance)
        }
        replay.forEach(subscriber::deliver)
        return subscriber
    }

    fun <T : Any> subscribe(
        owner: LifecycleOwner,
        type: Class<T>,
        on: Executor,
        handler: (T) -> Unit,
    ): Subscription = subscribe(type, on, handler).also { owner.lifecycle.addOnDestroy(it::unsubscribe) }

    fun post(event: Any) = targetsOf(event).forEach { it.deliver(event) }

    fun postSticky(event: Any) {
        val targets = synchronized(sticky) {
            sticky[event.javaClass] = event
            targetsOf(event)
        }
        targets.forEach { it.deliver(event) }
    }

    fun removeSticky(type: Class<*>) {
        synchronized(sticky) { sticky.keys.removeIf(type::isAssignableFrom) }
    }

    internal fun remove(subscriber: Subscriber<*>) {
        subscribers[subscriber.type]?.remove(subscriber)
    }

    private fun targetsOf(event: Any): List<Subscriber<*>> =
        typesOf(event.javaClass).flatMap { subscribers[it].orEmpty() }

    private fun typesOf(eventType: Class<*>): List<Class<*>> = typeCache.getOrPut(eventType) {
        val flat = mutableListOf<Class<*>>()
        val queue = ArrayDeque<Class<*>>().apply { add(eventType) }
        while (queue.isNotEmpty()) {
            val next = queue.removeFirst()
            if (next == Any::class.java || next in flat) continue
            flat += next
            next.superclass?.let(queue::add)
            queue.addAll(next.interfaces)
        }
        flat
    }
}

com.androidinterview.eventbus.events.AppEvent.kt

package com.androidinterview.eventbus.events

// A small event hierarchy, here only to show what supertype delivery buys.
// Subscribe to SignedIn and you get one kind of event. Subscribe to AppEvent
// and you get all three, which is how an analytics logger or a crash
// breadcrumb trail is written without touching any publisher.
//
// Events are values. They carry no callbacks and no reference to a screen,
// because the bus holds the last sticky one for as long as the process lives.
sealed interface AppEvent {
    data class SignedIn(val userId: String) : AppEvent
    data class SignedOut(val reason: String) : AppEvent
    data class NetworkChanged(val online: Boolean) : AppEvent
}
package com.androidinterview.eventbus.events

sealed interface AppEvent {
    data class SignedIn(val userId: String) : AppEvent
    data class SignedOut(val reason: String) : AppEvent
    data class NetworkChanged(val online: Boolean) : AppEvent
}

com.androidinterview.eventbus.lifecycle.Lifecycle.kt

package com.androidinterview.eventbus.lifecycle

// The smallest useful model of the Android lifecycle. The real one has a
// state machine and seven events. The bus needs two things, a way to be told
// when the owner is destroyed and an honest answer to whether that already
// happened.
class Lifecycle {

    private val lock = Any()
    private val onDestroy = mutableListOf<() -> Unit>()
    private var destroyed = false

    val isDestroyed: Boolean get() = synchronized(lock) { destroyed }

    // Registering after destruction runs the action straight away. That case
    // looks unlikely and is not. A subscription made on a background thread
    // while the screen is finishing would otherwise never be cancelled, which
    // is the exact leak this class exists to stop.
    fun addOnDestroy(action: () -> Unit) {
        synchronized(lock) {
            if (!destroyed) {
                onDestroy += action
                return
            }
        }
        action()
    }

    // The actions run outside the lock. Unsubscribing takes the bus registry
    // locks, and holding two locks in two different orders is how a deadlock
    // gets written by accident.
    fun destroy() {
        val actions = synchronized(lock) {
            if (destroyed) return
            destroyed = true
            onDestroy.toList().also { onDestroy.clear() }
        }
        actions.forEach { it() }
    }
}

// Anything that can be destroyed, which on Android is an Activity, a Fragment
// or a view. The bus sees only this, so nothing in the design depends on the
// SDK and a test can destroy an owner whenever it likes.
interface LifecycleOwner {
    val lifecycle: Lifecycle
}
package com.androidinterview.eventbus.lifecycle

class Lifecycle {

    private val lock = Any()
    private val onDestroy = mutableListOf<() -> Unit>()
    private var destroyed = false

    val isDestroyed: Boolean get() = synchronized(lock) { destroyed }

    fun addOnDestroy(action: () -> Unit) {
        synchronized(lock) {
            if (!destroyed) {
                onDestroy += action
                return
            }
        }
        action()
    }

    fun destroy() {
        val actions = synchronized(lock) {
            if (destroyed) return
            destroyed = true
            onDestroy.toList().also { onDestroy.clear() }
        }
        actions.forEach { it() }
    }
}

interface LifecycleOwner {
    val lifecycle: Lifecycle
}

com.androidinterview.eventbus.thread.Executors.kt

package com.androidinterview.eventbus.thread

import java.util.concurrent.Executor
import java.util.concurrent.Executors

// Stands in for the Android main looper. A looper is one thread with a queue
// in front of it, which is exactly a single thread executor, so modelling it
// this way costs a few lines and keeps the bus free of the SDK.
class MainThreadExecutor : Executor {

    @Volatile
    private var thread: Thread? = null

    private val delegate = Executors.newSingleThreadExecutor { runnable ->
        Thread(runnable, "main").also {
            it.isDaemon = true
            thread = it
        }
    }

    val isMainThread: Boolean get() = Thread.currentThread() === thread

    // Everything goes through the queue, including work posted from the main
    // thread itself. Running those inline would let a handler that posts
    // during delivery jump ahead of events already waiting, and a reordering
    // like that is very hard to find later.
    override fun execute(task: Runnable) = delegate.execute(task)

    fun shutdown() = delegate.shutdown()
}

// Delivers on whichever thread called post. Cheapest option and the right one
// for a handler that only touches its own data.
//
// It is also where re-entrancy shows up. A handler that posts during delivery
// runs the second event to completion inside the first, so the stack grows
// and the observed order is depth first rather than post order.
object DirectExecutor : Executor {
    override fun execute(task: Runnable) = task.run()
}
package com.androidinterview.eventbus.thread

import java.util.concurrent.Executor
import java.util.concurrent.Executors

class MainThreadExecutor : Executor {

    @Volatile
    private var thread: Thread? = null

    private val delegate = Executors.newSingleThreadExecutor { runnable ->
        Thread(runnable, "main").also {
            it.isDaemon = true
            thread = it
        }
    }

    val isMainThread: Boolean get() = Thread.currentThread() === thread

    override fun execute(task: Runnable) = delegate.execute(task)

    fun shutdown() = delegate.shutdown()
}

object DirectExecutor : Executor {
    override fun execute(task: Runnable) = task.run()
}

Watch