androidinterview.com

Low Level Design (LLD) Interview Questions

Implement your own LiveData

Tier: EssentialDifficulty: MediumAsked of: Mid, Senior

Design a value holder that updates a screen when its value changes. An observer is the callback the screen registers to receive those updates.

The problem

A screen receives count 1 while it is active. It stops, and the count changes to 2 and then 3. When it starts again, it should receive 3 once. Destroying the screen should remove its observer.

For the interview version, start with an initial value and run all operations on the main thread. Use three lifecycle states, stopped, started and destroyed. Started represents both Android STARTED and RESUMED. Background posting, an initially unset value and combining sources are follow-ups.

How to explain the design

“I keep the latest value and a list of observers. Each observer belongs to a screen's lifecycle. When the value changes, I notify only active screens. When a stopped screen starts, I send it the latest value if it missed an update. When it is destroyed, I remove its observer.”

There are two small classes. Lifecycle tells listeners when the screen starts, stops or is destroyed. SimpleLiveData stores the value and observer entries. Each entry keeps a callback, its lifecycle and the last version it received.

A version is just a counter. Increase it on each write. Comparing that counter with an observer's last version tells us whether the observer has already received this update.

Walk through one update

  1. Register an observer. If the screen is started, send the current value.
  2. Call setValue(2). Store 2 and increase the version.
  3. Notify each started observer and record the version it received.
  4. Skip stopped observers. On restart, their older version tells us to send the latest value.
  5. On destruction, remove the entry and its lifecycle listener.

Interview implementation

Write observe, setValue and the small delivery check first. The function returned by observe lets the caller unsubscribe early. A new registration is independent of earlier registrations, even if it uses the same callback. Callbacks should return normally and should not write another value while delivery is running.

Java

SimpleLiveData.java

package interview.livedata;

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

public class SimpleLiveData<T> {
    public enum State { STOPPED, STARTED, DESTROYED }

    public static class Lifecycle {
        private State state = State.STOPPED;
        private final List<Runnable> listeners = new ArrayList<>();

        public Runnable listen(Runnable listener) {
            listeners.add(listener);
            return () -> listeners.remove(listener);
        }

        public void moveTo(State next) {
            if (state == State.DESTROYED) throw new IllegalStateException("Owner is destroyed");
            state = next;
            for (Runnable listener : new ArrayList<>(listeners)) {
                if (listeners.contains(listener)) listener.run();
            }
            if (next == State.DESTROYED) listeners.clear();
        }
    }

    private class Entry {
        final Lifecycle owner;
        final Consumer<T> callback;
        long lastVersion = -1;
        Runnable detach = () -> {};

        Entry(Lifecycle owner, Consumer<T> callback) {
            this.owner = owner;
            this.callback = callback;
        }
    }

    private T value;
    private long version;
    private final List<Entry> observers = new ArrayList<>();

    public SimpleLiveData(T initialValue) { value = initialValue; }

    // Call these methods on the main thread.
    public Runnable observe(Lifecycle owner, Consumer<T> callback) {
        if (owner.state == State.DESTROYED) return () -> {};
        Entry entry = new Entry(owner, callback);
        observers.add(entry);
        entry.detach = owner.listen(() -> {
            if (owner.state == State.DESTROYED) remove(entry);
            else notifyEntry(entry);
        });
        notifyEntry(entry);
        return () -> remove(entry);
    }

    public void setValue(T next) {
        value = next;
        version++;
        for (Entry entry : new ArrayList<>(observers)) notifyEntry(entry);
    }

    private void notifyEntry(Entry entry) {
        if (!observers.contains(entry) || entry.owner.state != State.STARTED) return;
        if (entry.lastVersion == version) return;
        entry.lastVersion = version;
        entry.callback.accept(value);
    }

    private void remove(Entry entry) {
        observers.remove(entry);
        entry.detach.run();
    }
}
package interview.livedata;

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

public class SimpleLiveData<T> {
    public enum State { STOPPED, STARTED, DESTROYED }

    public static class Lifecycle {
        private State state = State.STOPPED;
        private final List<Runnable> listeners = new ArrayList<>();

        public Runnable listen(Runnable listener) {
            listeners.add(listener);
            return () -> listeners.remove(listener);
        }

        public void moveTo(State next) {
            if (state == State.DESTROYED) throw new IllegalStateException("Owner is destroyed");
            state = next;
            for (Runnable listener : new ArrayList<>(listeners)) {
                if (listeners.contains(listener)) listener.run();
            }
            if (next == State.DESTROYED) listeners.clear();
        }
    }

    private class Entry {
        final Lifecycle owner;
        final Consumer<T> callback;
        long lastVersion = -1;
        Runnable detach = () -> {};

        Entry(Lifecycle owner, Consumer<T> callback) {
            this.owner = owner;
            this.callback = callback;
        }
    }

    private T value;
    private long version;
    private final List<Entry> observers = new ArrayList<>();

    public SimpleLiveData(T initialValue) { value = initialValue; }

    public Runnable observe(Lifecycle owner, Consumer<T> callback) {
        if (owner.state == State.DESTROYED) return () -> {};
        Entry entry = new Entry(owner, callback);
        observers.add(entry);
        entry.detach = owner.listen(() -> {
            if (owner.state == State.DESTROYED) remove(entry);
            else notifyEntry(entry);
        });
        notifyEntry(entry);
        return () -> remove(entry);
    }

    public void setValue(T next) {
        value = next;
        version++;
        for (Entry entry : new ArrayList<>(observers)) notifyEntry(entry);
    }

    private void notifyEntry(Entry entry) {
        if (!observers.contains(entry) || entry.owner.state != State.STARTED) return;
        if (entry.lastVersion == version) return;
        entry.lastVersion = version;
        entry.callback.accept(value);
    }

    private void remove(Entry entry) {
        observers.remove(entry);
        entry.detach.run();
    }
}

Kotlin

SimpleLiveData.kt

package interview.livedata

enum class State { STOPPED, STARTED, DESTROYED }

class Lifecycle {
    var state = State.STOPPED
        private set
    private val listeners = mutableListOf<() -> Unit>()

    fun listen(listener: () -> Unit): () -> Unit {
        listeners.add(listener)
        return { listeners.remove(listener) }
    }

    fun moveTo(next: State) {
        check(state != State.DESTROYED)
        state = next
        for (listener in listeners.toList()) {
            if (listener in listeners) listener()
        }
        if (next == State.DESTROYED) listeners.clear()
    }
}

// Call these methods on the main thread.
class SimpleLiveData<T>(initialValue: T) {
    private var value = initialValue
    private var version = 0L
    private val observers = mutableListOf<Entry>()

    private inner class Entry(val owner: Lifecycle, val callback: (T) -> Unit) {
        var lastVersion = -1L
        var detach: () -> Unit = {}
    }

    fun observe(owner: Lifecycle, callback: (T) -> Unit): () -> Unit {
        if (owner.state == State.DESTROYED) return {}
        val entry = Entry(owner, callback)
        observers.add(entry)
        entry.detach = owner.listen {
            if (owner.state == State.DESTROYED) remove(entry)
            else notify(entry)
        }
        notify(entry)
        return { remove(entry) }
    }

    fun setValue(next: T) {
        value = next
        version++
        for (entry in observers.toList()) notify(entry)
    }

    private fun notify(entry: Entry) {
        if (entry !in observers || entry.owner.state != State.STARTED) return
        if (entry.lastVersion == version) return
        entry.lastVersion = version
        entry.callback(value)
    }

    private fun remove(entry: Entry) {
        observers.remove(entry)
        entry.detach()
    }
}
package interview.livedata

enum class State { STOPPED, STARTED, DESTROYED }

class Lifecycle {
    var state = State.STOPPED
        private set
    private val listeners = mutableListOf<() -> Unit>()

    fun listen(listener: () -> Unit): () -> Unit {
        listeners.add(listener)
        return { listeners.remove(listener) }
    }

    fun moveTo(next: State) {
        check(state != State.DESTROYED)
        state = next
        for (listener in listeners.toList()) {
            if (listener in listeners) listener()
        }
        if (next == State.DESTROYED) listeners.clear()
    }
}

class SimpleLiveData<T>(initialValue: T) {
    private var value = initialValue
    private var version = 0L
    private val observers = mutableListOf<Entry>()

    private inner class Entry(val owner: Lifecycle, val callback: (T) -> Unit) {
        var lastVersion = -1L
        var detach: () -> Unit = {}
    }

    fun observe(owner: Lifecycle, callback: (T) -> Unit): () -> Unit {
        if (owner.state == State.DESTROYED) return {}
        val entry = Entry(owner, callback)
        observers.add(entry)
        entry.detach = owner.listen {
            if (owner.state == State.DESTROYED) remove(entry)
            else notify(entry)
        }
        notify(entry)
        return { remove(entry) }
    }

    fun setValue(next: T) {
        value = next
        version++
        for (entry in observers.toList()) notify(entry)
    }

    private fun notify(entry: Entry) {
        if (entry !in observers || entry.owner.state != State.STARTED) return
        if (entry.lastVersion == version) return
        entry.lastVersion = version
        entry.callback(value)
    }

    private fun remove(entry: Entry) {
        observers.remove(entry)
        entry.detach()
    }
}

Follow-up questions

Why not notify every observer?

“I notify only started screens because a stopped screen should not update its UI. I still keep the latest value, so the screen can catch up when it starts again.” A destroyed screen is removed completely, which also releases its callback.

Why keep a version?

“The version tells me whether this observer has already received the latest write.” If a screen received version 3, stops and starts without another write, there is nothing new to send. If versions 4 and 5 arrive while it is stopped, it receives only version 5 on restart.

This is the delivery check inside the interview implementation. Record the version before calling the observer.

Kotlin

if (entry !in observers || entry.owner.state != State.STARTED) return
if (entry.lastVersion == version) return
entry.lastVersion = version
entry.callback(value)

Java

if (!observers.contains(entry) || entry.owner.state != State.STARTED) return;
if (entry.lastVersion == version) return;
entry.lastVersion = version;
entry.callback.accept(value);

What about null?

“Null can be a real value, such as no item being selected. It should still be delivered.” To also support an initially unset value, add a separate hasValue flag. Leave it false until the first write, then set it true even if that write is null.

What about background writes?

“I would queue the write on the main thread, so observers still run on the main thread.” Inject a java.util.concurrent.Executor called mainExecutor that queues work there, then add this method inside SimpleLiveData.

Kotlin

fun postValue(next: T) {
    mainExecutor.execute { setValue(next) }
}

Java

public void postValue(T next) {
    mainExecutor.execute(() -> setValue(next));
}

This small extension queues every write. To match Android's pending write behavior, keep one pending value and schedule only one delivery task. Writes arriving before that task runs replace the pending value, so posting 2 then 3 delivers only 3. Protect the pending value and task flag with the same lock.

What should I test?

“I would test what each observer actually receives.” A started observer gets the initial value once. A stopped observer gets nothing during writes of 2 and 3, then gets only 3 when started. Restarting again without a write adds no callback. After unsubscribe or destruction, later writes produce no callbacks. A write of null should still be delivered to an active observer.

The lifecycle behavior follows Android's LiveData overview. This small exercise models that behavior without recreating the full Android library.

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.livedata.core.LiveData.java

package com.androidinterview.livedata.core;

import java.util.ArrayList;
import java.util.LinkedHashMap;
import java.util.Map;

import com.androidinterview.livedata.lifecycle.LifecycleObserver;
import com.androidinterview.livedata.lifecycle.LifecycleOwner;
import com.androidinterview.livedata.lifecycle.LifecycleState;
import com.androidinterview.livedata.thread.MainThread;

// The observable holder. Read only on purpose, because the UI is allowed to
// watch a value and not to set one.
//
// Two counters carry the design. The holder keeps a version that goes up on
// every write, and every observer remembers the version it last saw. A
// delivery happens only when the observer is behind, so a late observer gets
// the current value exactly once, a stopped observer catches up when it starts
// again, and nobody ever sees the same value twice.
public class LiveData<T> {

    static final int START_VERSION = -1;

    // A sentinel, because null is a legal value, so "nothing yet" has to be a
    // different thing from "the value is null".
    private static final Object NOT_SET = new Object();

    private final MainThread mainThread;
    private final Map<Observer<? super T>, ObserverWrapper> observers = new LinkedHashMap<>();
    private final Object postLock = new Object();

    private Object data;
    private Object pending = NOT_SET;
    private int version;
    private int activeCount;
    private boolean dispatching;
    private boolean dispatchInvalidated;

    protected LiveData(MainThread mainThread) {
        this.mainThread = mainThread;
        this.data = NOT_SET;
        this.version = START_VERSION;
    }

    protected LiveData(MainThread mainThread, T initial) {
        this.mainThread = mainThread;
        this.data = initial;
        this.version = START_VERSION + 1;
    }

    @SuppressWarnings("unchecked")
    public T getValue() {
        return data == NOT_SET ? null : (T) data;
    }

    public boolean hasObservers() {
        return !observers.isEmpty();
    }

    public boolean hasActiveObservers() {
        return activeCount > 0;
    }

    // Bind an observer to a screen. The wrapper it registers hears DESTROYED
    // and removes itself, so nothing survives the screen. That is the whole
    // reason this exists instead of a plain listener list.
    public void observe(LifecycleOwner owner, Observer<? super T> observer) {
        assertMainThread("observe");
        if (owner.getLifecycle().currentState() == LifecycleState.DESTROYED) {
            return;
        }
        LifecycleBoundObserver wrapper = new LifecycleBoundObserver(owner, observer);
        ObserverWrapper existing = observers.putIfAbsent(observer, wrapper);
        if (existing != null) {
            if (existing.owner() != owner) {
                throw new IllegalArgumentException("that observer is already bound to another owner");
            }
            return;
        }
        owner.getLifecycle().addObserver(wrapper);
    }

    // No owner, so always active, which is what one holder feeding another
    // needs. The caveat is the important part. Nothing will ever remove this
    // observer for you, so the caller owns a matching removeObserver, and
    // forgetting it leaks the observer and everything it captured.
    public void observeForever(Observer<? super T> observer) {
        assertMainThread("observeForever");
        AlwaysActiveObserver wrapper = new AlwaysActiveObserver(observer);
        ObserverWrapper existing = observers.putIfAbsent(observer, wrapper);
        if (existing != null) {
            if (existing.owner() != null) {
                throw new IllegalArgumentException("that observer is already bound to a lifecycle");
            }
            return;
        }
        wrapper.activeStateChanged(true);
    }

    public void removeObserver(Observer<? super T> observer) {
        assertMainThread("removeObserver");
        ObserverWrapper wrapper = observers.remove(observer);
        if (wrapper == null) {
            return;
        }
        wrapper.detach();
        wrapper.activeStateChanged(false);
    }

    // Main thread only, and it fails loudly. Delivery is synchronous, so by the
    // time this returns every active observer has run.
    protected void setValue(T value) {
        assertMainThread("setValue");
        version++;
        data = value;
        dispatchValue(null);
    }

    // Callable from anywhere, and it coalesces. Post three values before the
    // main thread gets a turn and the queue carries one task that delivers the
    // third. A feature for progress, a trap for anything where every value
    // matters, and the first thing to say about it out loud.
    protected void postValue(T value) {
        boolean shouldPost;
        synchronized (postLock) {
            shouldPost = pending == NOT_SET;
            pending = value;
        }
        if (!shouldPost) {
            return;
        }
        mainThread.post(this::drainPending);
    }

    @SuppressWarnings("unchecked")
    private void drainPending() {
        Object next;
        synchronized (postLock) {
            next = pending;
            pending = NOT_SET;
        }
        setValue((T) next);
    }

    // Called on the first active observer and on losing the last one. A holder
    // that owns something expensive, a location listener or a socket, starts it
    // here and stops it there, so the cost exists only while somebody looks.
    protected void onActive() {
    }

    protected void onInactive() {
    }

    MainThread mainThread() {
        return mainThread;
    }

    int version() {
        return version;
    }

    private void assertMainThread(String operation) {
        if (!mainThread.isMainThread()) {
            throw new IllegalStateException(operation + " must be called from the main thread");
        }
    }

    private void changeActiveCount(int delta) {
        int previous = activeCount;
        activeCount += delta;
        if (previous == 0 && activeCount > 0) {
            onActive();
        }
        if (previous > 0 && activeCount == 0) {
            onInactive();
        }
    }

    // Delivery, with a re-entrancy guard. An observer may set a new value from
    // inside onChanged, and without the guard the two dispatches interleave and
    // observers see values out of order. The flag says a newer value arrived,
    // so abandon this pass and start again.
    private void dispatchValue(ObserverWrapper initiator) {
        if (dispatching) {
            dispatchInvalidated = true;
            return;
        }
        dispatching = true;
        try {
            do {
                dispatchInvalidated = false;
                if (initiator != null) {
                    considerNotify(initiator);
                    initiator = null;
                } else {
                    for (ObserverWrapper wrapper : new ArrayList<>(observers.values())) {
                        considerNotify(wrapper);
                        if (dispatchInvalidated) {
                            break;
                        }
                    }
                }
            } while (dispatchInvalidated);
        } finally {
            dispatching = false;
        }
    }

    // The three questions that decide whether a value reaches an observer, and
    // the order matters. Is it active, is it still allowed to be active, and is
    // it behind. Only the third one moves the version forward.
    @SuppressWarnings("unchecked")
    private void considerNotify(ObserverWrapper wrapper) {
        if (!wrapper.active) {
            return;
        }
        if (!wrapper.shouldBeActive()) {
            wrapper.activeStateChanged(false);
            return;
        }
        if (wrapper.lastVersion >= version) {
            return;
        }
        wrapper.lastVersion = version;
        wrapper.observer.onChanged((T) data);
    }

    // What the holder stores per observer. The observer itself is a plain
    // callback and knows none of this.
    private abstract class ObserverWrapper {

        final Observer<? super T> observer;
        int lastVersion = START_VERSION;
        boolean active;

        ObserverWrapper(Observer<? super T> observer) {
            this.observer = observer;
        }

        abstract boolean shouldBeActive();

        LifecycleOwner owner() {
            return null;
        }

        void detach() {
        }

        // The one place active flips, so the counter and the onActive hook
        // cannot drift apart. Going active tries a delivery straight away,
        // which is how a stopped screen catches up the moment it starts.
        void activeStateChanged(boolean newActive) {
            if (newActive == active) {
                return;
            }
            active = newActive;
            changeActiveCount(active ? 1 : -1);
            if (active) {
                dispatchValue(this);
            }
        }
    }

    private final class AlwaysActiveObserver extends ObserverWrapper {

        AlwaysActiveObserver(Observer<? super T> observer) {
            super(observer);
        }

        @Override
        boolean shouldBeActive() {
            return true;
        }
    }

    // The lifecycle aware half, and it is this small. Active means at least
    // STARTED, DESTROYED means remove yourself, and every state change asks
    // the question again.
    private final class LifecycleBoundObserver extends ObserverWrapper implements LifecycleObserver {

        private final LifecycleOwner owner;

        LifecycleBoundObserver(LifecycleOwner owner, Observer<? super T> observer) {
            super(observer);
            this.owner = owner;
        }

        @Override
        boolean shouldBeActive() {
            return owner.getLifecycle().currentState().isAtLeast(LifecycleState.STARTED);
        }

        @Override
        LifecycleOwner owner() {
            return owner;
        }

        @Override
        void detach() {
            owner.getLifecycle().removeObserver(this);
        }

        @Override
        public void onStateChanged(LifecycleState state) {
            if (state == LifecycleState.DESTROYED) {
                removeObserver(observer);
                return;
            }
            activeStateChanged(shouldBeActive());
        }
    }
}
package com.androidinterview.livedata.core;

import java.util.ArrayList;
import java.util.LinkedHashMap;
import java.util.Map;

import com.androidinterview.livedata.lifecycle.LifecycleObserver;
import com.androidinterview.livedata.lifecycle.LifecycleOwner;
import com.androidinterview.livedata.lifecycle.LifecycleState;
import com.androidinterview.livedata.thread.MainThread;

public class LiveData<T> {

    static final int START_VERSION = -1;

    private static final Object NOT_SET = new Object();

    private final MainThread mainThread;
    private final Map<Observer<? super T>, ObserverWrapper> observers = new LinkedHashMap<>();
    private final Object postLock = new Object();

    private Object data;
    private Object pending = NOT_SET;
    private int version;
    private int activeCount;
    private boolean dispatching;
    private boolean dispatchInvalidated;

    protected LiveData(MainThread mainThread) {
        this.mainThread = mainThread;
        this.data = NOT_SET;
        this.version = START_VERSION;
    }

    protected LiveData(MainThread mainThread, T initial) {
        this.mainThread = mainThread;
        this.data = initial;
        this.version = START_VERSION + 1;
    }

    @SuppressWarnings("unchecked")
    public T getValue() {
        return data == NOT_SET ? null : (T) data;
    }

    public boolean hasObservers() {
        return !observers.isEmpty();
    }

    public boolean hasActiveObservers() {
        return activeCount > 0;
    }

    public void observe(LifecycleOwner owner, Observer<? super T> observer) {
        assertMainThread("observe");
        if (owner.getLifecycle().currentState() == LifecycleState.DESTROYED) {
            return;
        }
        LifecycleBoundObserver wrapper = new LifecycleBoundObserver(owner, observer);
        ObserverWrapper existing = observers.putIfAbsent(observer, wrapper);
        if (existing != null) {
            if (existing.owner() != owner) {
                throw new IllegalArgumentException("that observer is already bound to another owner");
            }
            return;
        }
        owner.getLifecycle().addObserver(wrapper);
    }

    public void observeForever(Observer<? super T> observer) {
        assertMainThread("observeForever");
        AlwaysActiveObserver wrapper = new AlwaysActiveObserver(observer);
        ObserverWrapper existing = observers.putIfAbsent(observer, wrapper);
        if (existing != null) {
            if (existing.owner() != null) {
                throw new IllegalArgumentException("that observer is already bound to a lifecycle");
            }
            return;
        }
        wrapper.activeStateChanged(true);
    }

    public void removeObserver(Observer<? super T> observer) {
        assertMainThread("removeObserver");
        ObserverWrapper wrapper = observers.remove(observer);
        if (wrapper == null) {
            return;
        }
        wrapper.detach();
        wrapper.activeStateChanged(false);
    }

    protected void setValue(T value) {
        assertMainThread("setValue");
        version++;
        data = value;
        dispatchValue(null);
    }

    protected void postValue(T value) {
        boolean shouldPost;
        synchronized (postLock) {
            shouldPost = pending == NOT_SET;
            pending = value;
        }
        if (!shouldPost) {
            return;
        }
        mainThread.post(this::drainPending);
    }

    @SuppressWarnings("unchecked")
    private void drainPending() {
        Object next;
        synchronized (postLock) {
            next = pending;
            pending = NOT_SET;
        }
        setValue((T) next);
    }

    protected void onActive() {
    }

    protected void onInactive() {
    }

    MainThread mainThread() {
        return mainThread;
    }

    int version() {
        return version;
    }

    private void assertMainThread(String operation) {
        if (!mainThread.isMainThread()) {
            throw new IllegalStateException(operation + " must be called from the main thread");
        }
    }

    private void changeActiveCount(int delta) {
        int previous = activeCount;
        activeCount += delta;
        if (previous == 0 && activeCount > 0) {
            onActive();
        }
        if (previous > 0 && activeCount == 0) {
            onInactive();
        }
    }

    private void dispatchValue(ObserverWrapper initiator) {
        if (dispatching) {
            dispatchInvalidated = true;
            return;
        }
        dispatching = true;
        try {
            do {
                dispatchInvalidated = false;
                if (initiator != null) {
                    considerNotify(initiator);
                    initiator = null;
                } else {
                    for (ObserverWrapper wrapper : new ArrayList<>(observers.values())) {
                        considerNotify(wrapper);
                        if (dispatchInvalidated) {
                            break;
                        }
                    }
                }
            } while (dispatchInvalidated);
        } finally {
            dispatching = false;
        }
    }

    @SuppressWarnings("unchecked")
    private void considerNotify(ObserverWrapper wrapper) {
        if (!wrapper.active) {
            return;
        }
        if (!wrapper.shouldBeActive()) {
            wrapper.activeStateChanged(false);
            return;
        }
        if (wrapper.lastVersion >= version) {
            return;
        }
        wrapper.lastVersion = version;
        wrapper.observer.onChanged((T) data);
    }

    private abstract class ObserverWrapper {

        final Observer<? super T> observer;
        int lastVersion = START_VERSION;
        boolean active;

        ObserverWrapper(Observer<? super T> observer) {
            this.observer = observer;
        }

        abstract boolean shouldBeActive();

        LifecycleOwner owner() {
            return null;
        }

        void detach() {
        }

        void activeStateChanged(boolean newActive) {
            if (newActive == active) {
                return;
            }
            active = newActive;
            changeActiveCount(active ? 1 : -1);
            if (active) {
                dispatchValue(this);
            }
        }
    }

    private final class AlwaysActiveObserver extends ObserverWrapper {

        AlwaysActiveObserver(Observer<? super T> observer) {
            super(observer);
        }

        @Override
        boolean shouldBeActive() {
            return true;
        }
    }

    private final class LifecycleBoundObserver extends ObserverWrapper implements LifecycleObserver {

        private final LifecycleOwner owner;

        LifecycleBoundObserver(LifecycleOwner owner, Observer<? super T> observer) {
            super(observer);
            this.owner = owner;
        }

        @Override
        boolean shouldBeActive() {
            return owner.getLifecycle().currentState().isAtLeast(LifecycleState.STARTED);
        }

        @Override
        LifecycleOwner owner() {
            return owner;
        }

        @Override
        void detach() {
            owner.getLifecycle().removeObserver(this);
        }

        @Override
        public void onStateChanged(LifecycleState state) {
            if (state == LifecycleState.DESTROYED) {
                removeObserver(observer);
                return;
            }
            activeStateChanged(shouldBeActive());
        }
    }
}

com.androidinterview.livedata.core.MediatorLiveData.java

package com.androidinterview.livedata.core;

import java.util.ArrayList;
import java.util.LinkedHashMap;
import java.util.Map;

import com.androidinterview.livedata.thread.MainThread;

// A holder that watches other holders. It is the answer to combining two
// sources, and it is also what map and switchMap are built out of.
//
// The trick is in onActive and onInactive. A mediator subscribes to its sources
// only while somebody is watching it, so a chain of five transformations over a
// database query costs nothing when the screen is in the background. Subscribe
// eagerly in addSource instead and you have built a component that keeps every
// upstream source hot forever.
public class MediatorLiveData<T> extends MutableLiveData<T> {

    private final Map<LiveData<?>, Source<?>> sources = new LinkedHashMap<>();

    public MediatorLiveData(MainThread mainThread) {
        super(mainThread);
    }

    public <S> void addSource(LiveData<S> source, Observer<? super S> onChanged) {
        Source<S> plug = new Source<>(source, onChanged);
        Source<?> existing = sources.putIfAbsent(source, plug);
        if (existing != null) {
            if (existing.onChanged != onChanged) {
                throw new IllegalArgumentException("that source is already added with a different observer");
            }
            return;
        }
        if (hasActiveObservers()) {
            plug.plug();
        }
    }

    public void removeSource(LiveData<?> source) {
        Source<?> plug = sources.remove(source);
        if (plug != null) {
            plug.unplug();
        }
    }

    @Override
    protected void onActive() {
        for (Source<?> source : new ArrayList<>(sources.values())) {
            source.plug();
        }
    }

    @Override
    protected void onInactive() {
        for (Source<?> source : new ArrayList<>(sources.values())) {
            source.unplug();
        }
    }

    // One upstream subscription. It keeps a version of its own, because
    // unplugging and plugging back in creates a fresh wrapper on the source
    // whose last seen version starts at nothing. Without this the screen would
    // be handed the same value again every time it came back to the foreground.
    private static final class Source<S> implements Observer<S> {

        private final LiveData<S> liveData;
        final Observer<? super S> onChanged;
        private int lastVersion = LiveData.START_VERSION;

        Source(LiveData<S> liveData, Observer<? super S> onChanged) {
            this.liveData = liveData;
            this.onChanged = onChanged;
        }

        void plug() {
            liveData.observeForever(this);
        }

        void unplug() {
            liveData.removeObserver(this);
        }

        @Override
        public void onChanged(S value) {
            if (lastVersion == liveData.version()) {
                return;
            }
            lastVersion = liveData.version();
            onChanged.onChanged(value);
        }
    }
}
package com.androidinterview.livedata.core;

import java.util.ArrayList;
import java.util.LinkedHashMap;
import java.util.Map;

import com.androidinterview.livedata.thread.MainThread;

public class MediatorLiveData<T> extends MutableLiveData<T> {

    private final Map<LiveData<?>, Source<?>> sources = new LinkedHashMap<>();

    public MediatorLiveData(MainThread mainThread) {
        super(mainThread);
    }

    public <S> void addSource(LiveData<S> source, Observer<? super S> onChanged) {
        Source<S> plug = new Source<>(source, onChanged);
        Source<?> existing = sources.putIfAbsent(source, plug);
        if (existing != null) {
            if (existing.onChanged != onChanged) {
                throw new IllegalArgumentException("that source is already added with a different observer");
            }
            return;
        }
        if (hasActiveObservers()) {
            plug.plug();
        }
    }

    public void removeSource(LiveData<?> source) {
        Source<?> plug = sources.remove(source);
        if (plug != null) {
            plug.unplug();
        }
    }

    @Override
    protected void onActive() {
        for (Source<?> source : new ArrayList<>(sources.values())) {
            source.plug();
        }
    }

    @Override
    protected void onInactive() {
        for (Source<?> source : new ArrayList<>(sources.values())) {
            source.unplug();
        }
    }

    private static final class Source<S> implements Observer<S> {

        private final LiveData<S> liveData;
        final Observer<? super S> onChanged;
        private int lastVersion = LiveData.START_VERSION;

        Source(LiveData<S> liveData, Observer<? super S> onChanged) {
            this.liveData = liveData;
            this.onChanged = onChanged;
        }

        void plug() {
            liveData.observeForever(this);
        }

        void unplug() {
            liveData.removeObserver(this);
        }

        @Override
        public void onChanged(S value) {
            if (lastVersion == liveData.version()) {
                return;
            }
            lastVersion = liveData.version();
            onChanged.onChanged(value);
        }
    }
}

com.androidinterview.livedata.core.MutableLiveData.java

package com.androidinterview.livedata.core;

import com.androidinterview.livedata.thread.MainThread;

// The writable half, and it exists only to widen two methods from protected to
// public. A view model keeps one of these private and exposes it as a LiveData,
// so the UI can watch the value and has no way to set one. Half the value of
// the whole component is this one line of type discipline.
public class MutableLiveData<T> extends LiveData<T> {

    public MutableLiveData(MainThread mainThread) {
        super(mainThread);
    }

    public MutableLiveData(MainThread mainThread, T initial) {
        super(mainThread, initial);
    }

    @Override
    public void setValue(T value) {
        super.setValue(value);
    }

    @Override
    public void postValue(T value) {
        super.postValue(value);
    }
}
package com.androidinterview.livedata.core;

import com.androidinterview.livedata.thread.MainThread;

public class MutableLiveData<T> extends LiveData<T> {

    public MutableLiveData(MainThread mainThread) {
        super(mainThread);
    }

    public MutableLiveData(MainThread mainThread, T initial) {
        super(mainThread, initial);
    }

    @Override
    public void setValue(T value) {
        super.setValue(value);
    }

    @Override
    public void postValue(T value) {
        super.postValue(value);
    }
}

com.androidinterview.livedata.core.Observer.java

package com.androidinterview.livedata.core;

// One callback, one value. The observer never learns whether it was called
// because a value arrived or because its screen came back to the foreground,
// and that ignorance is the whole point of the design.
@FunctionalInterface
public interface Observer<T> {

    void onChanged(T value);
}
package com.androidinterview.livedata.core;

@FunctionalInterface
public interface Observer<T> {

    void onChanged(T value);
}

com.androidinterview.livedata.core.Transformations.java

package com.androidinterview.livedata.core;

import java.util.function.Function;

// Both operators are a mediator with one source and a rule. Neither needs a
// single new field on the holder, which is the point worth making in an
// interview. Get the mediator right and the operators are four lines each.
public final class Transformations {

    private Transformations() {
    }

    // Same stream, different shape.
    public static <X, Y> LiveData<Y> map(LiveData<X> source, Function<X, Y> mapper) {
        MediatorLiveData<Y> result = new MediatorLiveData<>(source.mainThread());
        result.addSource(source, value -> result.setValue(mapper.apply(value)));
        return result;
    }

    // A new stream per value, and the old one has to go. The removeSource is
    // the whole operator. Skip it and a user who types five characters into a
    // search box ends up with five live queries all writing into one result,
    // and the answer the screen shows is whichever one finished last.
    public static <X, Y> LiveData<Y> switchMap(LiveData<X> source, Function<X, LiveData<Y>> mapper) {
        MediatorLiveData<Y> result = new MediatorLiveData<>(source.mainThread());
        result.addSource(source, new Observer<X>() {

            private LiveData<Y> current;

            @Override
            public void onChanged(X value) {
                LiveData<Y> next = mapper.apply(value);
                if (current == next) {
                    return;
                }
                if (current != null) {
                    result.removeSource(current);
                }
                current = next;
                if (next != null) {
                    result.addSource(next, result::setValue);
                }
            }
        });
        return result;
    }
}
package com.androidinterview.livedata.core;

import java.util.function.Function;

public final class Transformations {

    private Transformations() {
    }

    public static <X, Y> LiveData<Y> map(LiveData<X> source, Function<X, Y> mapper) {
        MediatorLiveData<Y> result = new MediatorLiveData<>(source.mainThread());
        result.addSource(source, value -> result.setValue(mapper.apply(value)));
        return result;
    }

    public static <X, Y> LiveData<Y> switchMap(LiveData<X> source, Function<X, LiveData<Y>> mapper) {
        MediatorLiveData<Y> result = new MediatorLiveData<>(source.mainThread());
        result.addSource(source, new Observer<X>() {

            private LiveData<Y> current;

            @Override
            public void onChanged(X value) {
                LiveData<Y> next = mapper.apply(value);
                if (current == next) {
                    return;
                }
                if (current != null) {
                    result.removeSource(current);
                }
                current = next;
                if (next != null) {
                    result.addSource(next, result::setValue);
                }
            }
        });
        return result;
    }
}

com.androidinterview.livedata.lifecycle.Lifecycle.java

package com.androidinterview.livedata.lifecycle;

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

// The minimum lifecycle that makes the design work. A current state, a list of
// listeners, and a way to move.
//
// Two behaviours here are load bearing and both are copied from the real one.
// A new observer is told the current state immediately, which is how a late
// observer gets the current value without any special case in the holder. And
// moving to DESTROYED drops every observer, so a screen that is gone cannot
// keep a holder alive through a callback.
public final class Lifecycle {

    private final List<LifecycleObserver> observers = new ArrayList<>();
    private LifecycleState state = LifecycleState.CREATED;

    public LifecycleState currentState() {
        return state;
    }

    public void addObserver(LifecycleObserver observer) {
        if (state == LifecycleState.DESTROYED) {
            return;
        }
        observers.add(observer);
        // The immediate callback. Without it, an observer that arrives after
        // the screen is already started would sit inactive until the next
        // state change, and the holder would need a catch up path of its own.
        observer.onStateChanged(state);
    }

    public void removeObserver(LifecycleObserver observer) {
        observers.remove(observer);
    }

    // Iterating a copy, because an observer that hears DESTROYED removes
    // itself while we are still walking the list.
    public void moveTo(LifecycleState next) {
        state = next;
        for (LifecycleObserver observer : new ArrayList<>(observers)) {
            observer.onStateChanged(next);
        }
        if (next == LifecycleState.DESTROYED) {
            observers.clear();
        }
    }
}
package com.androidinterview.livedata.lifecycle;

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

public final class Lifecycle {

    private final List<LifecycleObserver> observers = new ArrayList<>();
    private LifecycleState state = LifecycleState.CREATED;

    public LifecycleState currentState() {
        return state;
    }

    public void addObserver(LifecycleObserver observer) {
        if (state == LifecycleState.DESTROYED) {
            return;
        }
        observers.add(observer);
        observer.onStateChanged(state);
    }

    public void removeObserver(LifecycleObserver observer) {
        observers.remove(observer);
    }

    public void moveTo(LifecycleState next) {
        state = next;
        for (LifecycleObserver observer : new ArrayList<>(observers)) {
            observer.onStateChanged(next);
        }
        if (next == LifecycleState.DESTROYED) {
            observers.clear();
        }
    }
}

com.androidinterview.livedata.lifecycle.LifecycleObserver.java

package com.androidinterview.livedata.lifecycle;

// Something that wants to hear about state changes. One method, because a
// callback per event would mean a new method every time a state is added.
@FunctionalInterface
public interface LifecycleObserver {

    void onStateChanged(LifecycleState state);
}
package com.androidinterview.livedata.lifecycle;

@FunctionalInterface
public interface LifecycleObserver {

    void onStateChanged(LifecycleState state);
}

com.androidinterview.livedata.lifecycle.LifecycleOwner.java

package com.androidinterview.livedata.lifecycle;

// An activity or a fragment, reduced to the only thing our holder needs from
// it. Taking the owner rather than the activity is what keeps this code
// testable, because a test can hand over a plain object with a lifecycle it
// drives by hand.
public interface LifecycleOwner {

    Lifecycle getLifecycle();
}
package com.androidinterview.livedata.lifecycle;

public interface LifecycleOwner {

    Lifecycle getLifecycle();
}

com.androidinterview.livedata.lifecycle.LifecycleState.java

package com.androidinterview.livedata.lifecycle;

// The four states we need, written in order from dead to fully visible. The
// order is the point, because "is this observer active" is the single question
// isAtLeast answers, and DESTROYED sitting at the bottom means a destroyed
// screen can never be at least STARTED.
//
// The real Lifecycle has INITIALIZED too. It buys nothing here, so it is left
// out rather than copied.
public enum LifecycleState {

    DESTROYED,
    CREATED,
    STARTED,
    RESUMED;

    public boolean isAtLeast(LifecycleState other) {
        return compareTo(other) >= 0;
    }
}
package com.androidinterview.livedata.lifecycle;

public enum LifecycleState {

    DESTROYED,
    CREATED,
    STARTED,
    RESUMED;

    public boolean isAtLeast(LifecycleState other) {
        return compareTo(other) >= 0;
    }
}

com.androidinterview.livedata.sample.Location.java

package com.androidinterview.livedata.sample;

public record Location(double latitude, double longitude) {}
package com.androidinterview.livedata.sample;

public record Location(double latitude, double longitude) {}

com.androidinterview.livedata.sample.LocationLiveData.java

package com.androidinterview.livedata.sample;

import java.util.function.Consumer;

import com.androidinterview.livedata.core.LiveData;
import com.androidinterview.livedata.thread.MainThread;

// Why onActive and onInactive are worth having. A location listener costs
// battery, so it runs while a screen is watching and stops when the last one
// goes away. Rotation is the case that makes this subtle. The old activity is
// destroyed and the new one observes moments later, so a naive design would
// stop the GPS and start it again for no reason. Real code softens that with a
// short delay before releasing, and it is a good thing to mention.
public final class LocationLiveData extends LiveData<Location> {

    private final LocationProvider provider;
    private final Consumer<Location> listener;

    public LocationLiveData(MainThread mainThread, LocationProvider provider) {
        super(mainThread);
        this.provider = provider;
        // postValue, not setValue, because a location callback arrives on
        // whichever thread the provider felt like using.
        this.listener = this::postValue;
    }

    @Override
    protected void onActive() {
        provider.start(listener);
    }

    @Override
    protected void onInactive() {
        provider.stop(listener);
    }
}
package com.androidinterview.livedata.sample;

import java.util.function.Consumer;

import com.androidinterview.livedata.core.LiveData;
import com.androidinterview.livedata.thread.MainThread;

public final class LocationLiveData extends LiveData<Location> {

    private final LocationProvider provider;
    private final Consumer<Location> listener;

    public LocationLiveData(MainThread mainThread, LocationProvider provider) {
        super(mainThread);
        this.provider = provider;
        this.listener = this::postValue;
    }

    @Override
    protected void onActive() {
        provider.start(listener);
    }

    @Override
    protected void onInactive() {
        provider.stop(listener);
    }
}

com.androidinterview.livedata.sample.LocationProvider.java

package com.androidinterview.livedata.sample;

import java.util.function.Consumer;

// Stands in for the platform location client. It calls back on a thread of its
// own choosing, which is exactly the case postValue exists for.
public interface LocationProvider {

    void start(Consumer<Location> listener);

    void stop(Consumer<Location> listener);
}
package com.androidinterview.livedata.sample;

import java.util.function.Consumer;

public interface LocationProvider {

    void start(Consumer<Location> listener);

    void stop(Consumer<Location> listener);
}

com.androidinterview.livedata.sample.SingleLiveEvent.java

package com.androidinterview.livedata.sample;

import java.util.concurrent.atomic.AtomicBoolean;

import com.androidinterview.livedata.core.MutableLiveData;
import com.androidinterview.livedata.core.Observer;
import com.androidinterview.livedata.lifecycle.LifecycleOwner;
import com.androidinterview.livedata.thread.MainThread;

// The famous workaround, written out so you can see why it is a workaround and
// not a fix. The holder replays its current value to whoever starts observing,
// which is right for a name on a screen and wrong for "show a toast", because
// rotation brings a new observer that is behind. The flag breaks the replay.
//
// Two honest flaws. Only one observer can win the flag, so a second observer on
// the same event silently gets nothing, and the wrapping means the caller can
// no longer remove the observer it registered. The bug is that a state holder
// was asked to carry events, which is why one shot signals moved to a channel.
public final class SingleLiveEvent<T> extends MutableLiveData<T> {

    private final AtomicBoolean pending = new AtomicBoolean(false);

    public SingleLiveEvent(MainThread mainThread) {
        super(mainThread);
    }

    @Override
    public void observe(LifecycleOwner owner, Observer<? super T> observer) {
        super.observe(owner, value -> {
            if (pending.compareAndSet(true, false)) {
                observer.onChanged(value);
            }
        });
    }

    @Override
    public void setValue(T value) {
        pending.set(true);
        super.setValue(value);
    }
}
package com.androidinterview.livedata.sample;

import java.util.concurrent.atomic.AtomicBoolean;

import com.androidinterview.livedata.core.MutableLiveData;
import com.androidinterview.livedata.core.Observer;
import com.androidinterview.livedata.lifecycle.LifecycleOwner;
import com.androidinterview.livedata.thread.MainThread;

public final class SingleLiveEvent<T> extends MutableLiveData<T> {

    private final AtomicBoolean pending = new AtomicBoolean(false);

    public SingleLiveEvent(MainThread mainThread) {
        super(mainThread);
    }

    @Override
    public void observe(LifecycleOwner owner, Observer<? super T> observer) {
        super.observe(owner, value -> {
            if (pending.compareAndSet(true, false)) {
                observer.onChanged(value);
            }
        });
    }

    @Override
    public void setValue(T value) {
        pending.set(true);
        super.setValue(value);
    }
}

com.androidinterview.livedata.thread.MainThread.java

package com.androidinterview.livedata.thread;

// The main thread, as an interface, so the holder never touches Looper or
// Handler and a test never needs a device.
//
// Two questions is all the holder asks. Am I on it, which guards setValue, and
// please run this on it, which is what postValue uses.
public interface MainThread {

    boolean isMainThread();

    void post(Runnable task);
}
package com.androidinterview.livedata.thread;

public interface MainThread {

    boolean isMainThread();

    void post(Runnable task);
}

com.androidinterview.livedata.thread.QueuedMainThread.java

package com.androidinterview.livedata.thread;

import java.util.ArrayDeque;
import java.util.Deque;

// A main thread you drain by hand, which is exactly what Looper does in a real
// app and exactly what a test wants.
//
// Work posted from any thread lands on the queue. Whichever thread calls drain
// becomes the main thread for as long as the drain runs, so isMainThread is a
// real answer rather than a hardcoded true.
public final class QueuedMainThread implements MainThread {

    private final Deque<Runnable> queue = new ArrayDeque<>();
    private volatile Thread draining;

    @Override
    public boolean isMainThread() {
        return Thread.currentThread() == draining;
    }

    @Override
    public void post(Runnable task) {
        synchronized (queue) {
            queue.addLast(task);
        }
    }

    // Run everything queued, on this thread, in the order it was posted. A
    // task that posts another task is picked up by the same loop, which is the
    // behaviour a real message queue has.
    public void drain() {
        Thread previous = draining;
        draining = Thread.currentThread();
        try {
            while (true) {
                Runnable task;
                synchronized (queue) {
                    task = queue.pollFirst();
                }
                if (task == null) {
                    return;
                }
                task.run();
            }
        } finally {
            draining = previous;
        }
    }

    // The convenience a caller wants. Run this block as if it were on the main
    // thread, then let everything it posted run too.
    public void runOnMain(Runnable block) {
        post(block);
        drain();
    }
}
package com.androidinterview.livedata.thread;

import java.util.ArrayDeque;
import java.util.Deque;

public final class QueuedMainThread implements MainThread {

    private final Deque<Runnable> queue = new ArrayDeque<>();
    private volatile Thread draining;

    @Override
    public boolean isMainThread() {
        return Thread.currentThread() == draining;
    }

    @Override
    public void post(Runnable task) {
        synchronized (queue) {
            queue.addLast(task);
        }
    }

    public void drain() {
        Thread previous = draining;
        draining = Thread.currentThread();
        try {
            while (true) {
                Runnable task;
                synchronized (queue) {
                    task = queue.pollFirst();
                }
                if (task == null) {
                    return;
                }
                task.run();
            }
        } finally {
            draining = previous;
        }
    }

    public void runOnMain(Runnable block) {
        post(block);
        drain();
    }
}

Kotlin

com.androidinterview.livedata.core.LiveData.kt

package com.androidinterview.livedata.core

import com.androidinterview.livedata.lifecycle.LifecycleObserver
import com.androidinterview.livedata.lifecycle.LifecycleOwner
import com.androidinterview.livedata.lifecycle.LifecycleState
import com.androidinterview.livedata.thread.MainThread

// An observer is a function, not an interface. There is nothing an interface
// would add here, and a function type means a lambda, a method reference and
// another holder's callback all fit without a wrapper.
typealias Observer<T> = (T) -> Unit

// A sentinel, because null is a perfectly good value to hold, so "no value
// yet" has to be a different thing from "the value is null".
private val NOT_SET = Any()

internal const val START_VERSION = -1

// The observable holder. Read only on purpose, because the UI is allowed to
// watch a value and is not allowed to set one.
//
// Two counters carry the whole design. The holder keeps a version that goes up
// on every write, and every observer remembers the version it last saw. A
// delivery happens only when an observer is behind, which is why a late
// observer gets the current value exactly once, an observer whose screen was
// stopped catches up the moment it starts again, and neither of them ever sees
// the same value twice.
open class LiveData<T> protected constructor(internal val mainThread: MainThread) {

    private val observers = LinkedHashMap<Observer<T>, ObserverWrapper>()
    private val postLock = Any()

    private var data: Any? = NOT_SET
    private var pending: Any? = NOT_SET
    private var activeCount = 0
    private var dispatching = false
    private var invalidated = false

    internal var version = START_VERSION
        private set

    protected constructor(mainThread: MainThread, initial: T) : this(mainThread) {
        data = initial
        version = START_VERSION + 1
    }

    @Suppress("UNCHECKED_CAST")
    val value: T?
        get() = if (data === NOT_SET) null else data as T?

    val hasObservers: Boolean get() = observers.isNotEmpty()

    val hasActiveObservers: Boolean get() = activeCount > 0

    // Bind an observer to a screen. The holder never keeps that observer alive
    // past the screen, because the wrapper it registers hears DESTROYED and
    // removes itself. That one line is the reason this exists instead of a
    // plain listener list.
    open fun observe(owner: LifecycleOwner, observer: Observer<T>) {
        assertMainThread("observe")
        if (owner.lifecycle.currentState == LifecycleState.DESTROYED) return
        val wrapper = LifecycleBoundObserver(owner, observer)
        val existing = observers.putIfAbsent(observer, wrapper)
        if (existing != null) {
            require(existing.owner === owner) { "that observer is already bound to another owner" }
            return
        }
        owner.lifecycle.addObserver(wrapper)
    }

    // No owner, so always active. Useful when the watcher is not a screen, for
    // example one holder feeding another. The caveat is the important part.
    // Nothing will ever remove this observer for you, so the caller owns a
    // matching removeObserver, and forgetting it leaks the observer and
    // everything the lambda captured.
    fun observeForever(observer: Observer<T>) {
        assertMainThread("observeForever")
        val wrapper = AlwaysActiveObserver(observer)
        val existing = observers.putIfAbsent(observer, wrapper)
        if (existing != null) {
            require(existing.owner == null) { "that observer is already bound to a lifecycle" }
            return
        }
        wrapper.activeStateChanged(true)
    }

    fun removeObserver(observer: Observer<T>) {
        assertMainThread("removeObserver")
        val wrapper = observers.remove(observer) ?: return
        wrapper.detach()
        wrapper.activeStateChanged(false)
    }

    // Main thread only, and it fails loudly rather than quietly. Delivery is
    // synchronous, so by the time this returns every active observer has run.
    protected open fun setValue(value: T) {
        assertMainThread("setValue")
        version++
        data = value
        dispatchValue(null)
    }

    // Callable from anywhere, and it coalesces. Post three values before the
    // main thread gets a turn and the queue carries one task that delivers the
    // third. A feature for progress, a trap for anything where every value
    // matters, and the first thing to say about it out loud.
    protected open fun postValue(value: T) {
        val shouldPost = synchronized(postLock) {
            val first = pending === NOT_SET
            pending = value
            first
        }
        if (!shouldPost) return
        mainThread.post {
            @Suppress("UNCHECKED_CAST")
            val next = synchronized(postLock) { pending.also { pending = NOT_SET } } as T
            setValue(next)
        }
    }

    // Called when the holder gains its first active observer and when it loses
    // its last. A holder that owns something expensive, a location listener or
    // a socket, starts it here and stops it there, so the cost exists only
    // while somebody is looking.
    protected open fun onActive() = Unit

    protected open fun onInactive() = Unit

    private fun assertMainThread(operation: String) =
        check(mainThread.isMainThread) { "$operation must be called from the main thread" }

    private fun changeActiveCount(delta: Int) {
        val previous = activeCount
        activeCount += delta
        if (previous == 0 && activeCount > 0) onActive()
        if (previous > 0 && activeCount == 0) onInactive()
    }

    // Delivery, with a re-entrancy guard. An observer may set a new value from
    // inside its own callback, and without the guard the two dispatches
    // interleave and observers see values out of order. The flag says a newer
    // value arrived, so abandon this pass and start again.
    private fun dispatchValue(initiator: ObserverWrapper?) {
        if (dispatching) {
            invalidated = true
            return
        }
        dispatching = true
        var target = initiator
        try {
            do {
                invalidated = false
                if (target != null) {
                    considerNotify(target)
                    target = null
                } else {
                    for (wrapper in observers.values.toList()) {
                        considerNotify(wrapper)
                        if (invalidated) break
                    }
                }
            } while (invalidated)
        } finally {
            dispatching = false
        }
    }

    // The three questions that decide whether a value reaches an observer, and
    // the order matters. Is it active, is it still allowed to be active, and is
    // it behind. Only the third moves its version forward.
    @Suppress("UNCHECKED_CAST")
    private fun considerNotify(wrapper: ObserverWrapper) {
        if (!wrapper.active) return
        if (!wrapper.shouldBeActive()) {
            wrapper.activeStateChanged(false)
            return
        }
        if (wrapper.lastVersion >= version) return
        wrapper.lastVersion = version
        wrapper.observer(data as T)
    }

    // What the holder stores per observer. The observer itself is a plain
    // function and knows none of this.
    private abstract inner class ObserverWrapper(val observer: Observer<T>) {

        var lastVersion = START_VERSION
        var active = false

        open val owner: LifecycleOwner? get() = null

        abstract fun shouldBeActive(): Boolean

        open fun detach() = Unit

        // The one place active flips, so the counter and the onActive hook can
        // never drift apart. Going active tries a delivery straight away, which
        // is how a stopped screen catches up the moment it starts.
        fun activeStateChanged(newActive: Boolean) {
            if (newActive == active) return
            active = newActive
            changeActiveCount(if (active) 1 else -1)
            if (active) dispatchValue(this)
        }
    }

    private inner class AlwaysActiveObserver(observer: Observer<T>) : ObserverWrapper(observer) {
        override fun shouldBeActive() = true
    }

    // The lifecycle aware half of the component, and it is this small. Active
    // means at least STARTED, DESTROYED means remove yourself, and every state
    // change asks the question again.
    private inner class LifecycleBoundObserver(
        override val owner: LifecycleOwner,
        observer: Observer<T>,
    ) : ObserverWrapper(observer), LifecycleObserver {

        override fun shouldBeActive() =
            owner.lifecycle.currentState.isAtLeast(LifecycleState.STARTED)

        override fun detach() = owner.lifecycle.removeObserver(this)

        override fun onStateChanged(state: LifecycleState) {
            if (state == LifecycleState.DESTROYED) removeObserver(observer)
            else activeStateChanged(shouldBeActive())
        }
    }
}

// The writable half, and it exists only to widen two methods. A view model
// keeps one of these private and exposes it as a LiveData, so the UI can watch
// the value and has no way to set one. Half the value of the component is that
// single line of type discipline.
open class MutableLiveData<T> : LiveData<T> {

    constructor(mainThread: MainThread) : super(mainThread)

    constructor(mainThread: MainThread, initial: T) : super(mainThread, initial)

    public override fun setValue(value: T) = super.setValue(value)

    public override fun postValue(value: T) = super.postValue(value)
}
package com.androidinterview.livedata.core

import com.androidinterview.livedata.lifecycle.LifecycleObserver
import com.androidinterview.livedata.lifecycle.LifecycleOwner
import com.androidinterview.livedata.lifecycle.LifecycleState
import com.androidinterview.livedata.thread.MainThread

typealias Observer<T> = (T) -> Unit

private val NOT_SET = Any()

internal const val START_VERSION = -1

open class LiveData<T> protected constructor(internal val mainThread: MainThread) {

    private val observers = LinkedHashMap<Observer<T>, ObserverWrapper>()
    private val postLock = Any()

    private var data: Any? = NOT_SET
    private var pending: Any? = NOT_SET
    private var activeCount = 0
    private var dispatching = false
    private var invalidated = false

    internal var version = START_VERSION
        private set

    protected constructor(mainThread: MainThread, initial: T) : this(mainThread) {
        data = initial
        version = START_VERSION + 1
    }

    @Suppress("UNCHECKED_CAST")
    val value: T?
        get() = if (data === NOT_SET) null else data as T?

    val hasObservers: Boolean get() = observers.isNotEmpty()

    val hasActiveObservers: Boolean get() = activeCount > 0

    open fun observe(owner: LifecycleOwner, observer: Observer<T>) {
        assertMainThread("observe")
        if (owner.lifecycle.currentState == LifecycleState.DESTROYED) return
        val wrapper = LifecycleBoundObserver(owner, observer)
        val existing = observers.putIfAbsent(observer, wrapper)
        if (existing != null) {
            require(existing.owner === owner) { "that observer is already bound to another owner" }
            return
        }
        owner.lifecycle.addObserver(wrapper)
    }

    fun observeForever(observer: Observer<T>) {
        assertMainThread("observeForever")
        val wrapper = AlwaysActiveObserver(observer)
        val existing = observers.putIfAbsent(observer, wrapper)
        if (existing != null) {
            require(existing.owner == null) { "that observer is already bound to a lifecycle" }
            return
        }
        wrapper.activeStateChanged(true)
    }

    fun removeObserver(observer: Observer<T>) {
        assertMainThread("removeObserver")
        val wrapper = observers.remove(observer) ?: return
        wrapper.detach()
        wrapper.activeStateChanged(false)
    }

    protected open fun setValue(value: T) {
        assertMainThread("setValue")
        version++
        data = value
        dispatchValue(null)
    }

    protected open fun postValue(value: T) {
        val shouldPost = synchronized(postLock) {
            val first = pending === NOT_SET
            pending = value
            first
        }
        if (!shouldPost) return
        mainThread.post {
            @Suppress("UNCHECKED_CAST")
            val next = synchronized(postLock) { pending.also { pending = NOT_SET } } as T
            setValue(next)
        }
    }

    protected open fun onActive() = Unit

    protected open fun onInactive() = Unit

    private fun assertMainThread(operation: String) =
        check(mainThread.isMainThread) { "$operation must be called from the main thread" }

    private fun changeActiveCount(delta: Int) {
        val previous = activeCount
        activeCount += delta
        if (previous == 0 && activeCount > 0) onActive()
        if (previous > 0 && activeCount == 0) onInactive()
    }

    private fun dispatchValue(initiator: ObserverWrapper?) {
        if (dispatching) {
            invalidated = true
            return
        }
        dispatching = true
        var target = initiator
        try {
            do {
                invalidated = false
                if (target != null) {
                    considerNotify(target)
                    target = null
                } else {
                    for (wrapper in observers.values.toList()) {
                        considerNotify(wrapper)
                        if (invalidated) break
                    }
                }
            } while (invalidated)
        } finally {
            dispatching = false
        }
    }

    @Suppress("UNCHECKED_CAST")
    private fun considerNotify(wrapper: ObserverWrapper) {
        if (!wrapper.active) return
        if (!wrapper.shouldBeActive()) {
            wrapper.activeStateChanged(false)
            return
        }
        if (wrapper.lastVersion >= version) return
        wrapper.lastVersion = version
        wrapper.observer(data as T)
    }

    private abstract inner class ObserverWrapper(val observer: Observer<T>) {

        var lastVersion = START_VERSION
        var active = false

        open val owner: LifecycleOwner? get() = null

        abstract fun shouldBeActive(): Boolean

        open fun detach() = Unit

        fun activeStateChanged(newActive: Boolean) {
            if (newActive == active) return
            active = newActive
            changeActiveCount(if (active) 1 else -1)
            if (active) dispatchValue(this)
        }
    }

    private inner class AlwaysActiveObserver(observer: Observer<T>) : ObserverWrapper(observer) {
        override fun shouldBeActive() = true
    }

    private inner class LifecycleBoundObserver(
        override val owner: LifecycleOwner,
        observer: Observer<T>,
    ) : ObserverWrapper(observer), LifecycleObserver {

        override fun shouldBeActive() =
            owner.lifecycle.currentState.isAtLeast(LifecycleState.STARTED)

        override fun detach() = owner.lifecycle.removeObserver(this)

        override fun onStateChanged(state: LifecycleState) {
            if (state == LifecycleState.DESTROYED) removeObserver(observer)
            else activeStateChanged(shouldBeActive())
        }
    }
}

open class MutableLiveData<T> : LiveData<T> {

    constructor(mainThread: MainThread) : super(mainThread)

    constructor(mainThread: MainThread, initial: T) : super(mainThread, initial)

    public override fun setValue(value: T) = super.setValue(value)

    public override fun postValue(value: T) = super.postValue(value)
}

com.androidinterview.livedata.core.MediatorLiveData.kt

package com.androidinterview.livedata.core

import com.androidinterview.livedata.thread.MainThread

// A holder that watches other holders. It answers "combine these two sources",
// and it is also what map and switchMap are built out of.
//
// The trick is in onActive and onInactive. A mediator subscribes to its sources
// only while somebody is watching it, so a chain of five transformations over a
// database query costs nothing while the screen is in the background. Subscribe
// eagerly in addSource instead and you have built a component that keeps every
// upstream source hot forever.
open class MediatorLiveData<T>(mainThread: MainThread) : MutableLiveData<T>(mainThread) {

    private val sources = LinkedHashMap<LiveData<*>, Source<*>>()

    fun <S> addSource(source: LiveData<S>, onChanged: Observer<S>) {
        val plug = Source(source, onChanged)
        val existing = sources.putIfAbsent(source, plug)
        if (existing != null) {
            require(existing.onChanged === onChanged) {
                "that source is already added with a different observer"
            }
            return
        }
        if (hasActiveObservers) plug.plug()
    }

    fun removeSource(source: LiveData<*>) {
        sources.remove(source)?.unplug()
    }

    override fun onActive() = sources.values.toList().forEach { it.plug() }

    override fun onInactive() = sources.values.toList().forEach { it.unplug() }

    // One upstream subscription. It keeps a version of its own, because
    // unplugging and plugging back in creates a fresh wrapper on the source
    // whose last seen version starts at nothing. Without this, the screen would
    // be handed the same value again every time it returned to the foreground.
    private class Source<S>(private val liveData: LiveData<S>, val onChanged: Observer<S>) {

        private var lastVersion = START_VERSION

        private val relay: Observer<S> = { value ->
            if (lastVersion != liveData.version) {
                lastVersion = liveData.version
                onChanged(value)
            }
        }

        fun plug() = liveData.observeForever(relay)

        fun unplug() = liveData.removeObserver(relay)
    }
}

// Both operators are a mediator with one source and a rule, and neither needs
// a single new field on the holder. That is the point worth making out loud.
// Get the mediator right and the operators are three lines each.
//
// Java gathers these as static methods on a Transformations class because it
// has nowhere else to put them. Kotlin has extension functions, so map reads
// as a method on the holder and no helper class exists at all.
fun <X, Y> LiveData<X>.map(transform: (X) -> Y): LiveData<Y> =
    MediatorLiveData<Y>(mainThread).apply {
        addSource(this@map) { setValue(transform(it)) }
    }

// A new stream per value, and the old one has to go. The removeSource is the
// whole operator. Skip it and a user who types five characters into a search
// box ends up with five live queries writing into one result, and the screen
// shows whichever one finished last.
fun <X, Y> LiveData<X>.switchMap(transform: (X) -> LiveData<Y>?): LiveData<Y> =
    MediatorLiveData<Y>(mainThread).apply {
        var current: LiveData<Y>? = null
        addSource(this@switchMap) { x ->
            val next = transform(x)
            if (next !== current) {
                current?.let { removeSource(it) }
                current = next
                next?.let { source -> addSource(source) { setValue(it) } }
            }
        }
    }
package com.androidinterview.livedata.core

import com.androidinterview.livedata.thread.MainThread

open class MediatorLiveData<T>(mainThread: MainThread) : MutableLiveData<T>(mainThread) {

    private val sources = LinkedHashMap<LiveData<*>, Source<*>>()

    fun <S> addSource(source: LiveData<S>, onChanged: Observer<S>) {
        val plug = Source(source, onChanged)
        val existing = sources.putIfAbsent(source, plug)
        if (existing != null) {
            require(existing.onChanged === onChanged) {
                "that source is already added with a different observer"
            }
            return
        }
        if (hasActiveObservers) plug.plug()
    }

    fun removeSource(source: LiveData<*>) {
        sources.remove(source)?.unplug()
    }

    override fun onActive() = sources.values.toList().forEach { it.plug() }

    override fun onInactive() = sources.values.toList().forEach { it.unplug() }

    private class Source<S>(private val liveData: LiveData<S>, val onChanged: Observer<S>) {

        private var lastVersion = START_VERSION

        private val relay: Observer<S> = { value ->
            if (lastVersion != liveData.version) {
                lastVersion = liveData.version
                onChanged(value)
            }
        }

        fun plug() = liveData.observeForever(relay)

        fun unplug() = liveData.removeObserver(relay)
    }
}

fun <X, Y> LiveData<X>.map(transform: (X) -> Y): LiveData<Y> =
    MediatorLiveData<Y>(mainThread).apply {
        addSource(this@map) { setValue(transform(it)) }
    }

fun <X, Y> LiveData<X>.switchMap(transform: (X) -> LiveData<Y>?): LiveData<Y> =
    MediatorLiveData<Y>(mainThread).apply {
        var current: LiveData<Y>? = null
        addSource(this@switchMap) { x ->
            val next = transform(x)
            if (next !== current) {
                current?.let { removeSource(it) }
                current = next
                next?.let { source -> addSource(source) { setValue(it) } }
            }
        }
    }

com.androidinterview.livedata.lifecycle.Lifecycle.kt

package com.androidinterview.livedata.lifecycle

// The four states we need, in order from dead to fully visible. The order is
// the point, because "is this observer active" is the one question isAtLeast
// answers, and DESTROYED at the bottom means a destroyed screen can never be
// at least STARTED.
//
// The real Lifecycle has INITIALIZED as well. It buys nothing here, so it is
// left out rather than copied.
enum class LifecycleState {
    DESTROYED,
    CREATED,
    STARTED,
    RESUMED;

    fun isAtLeast(other: LifecycleState) = ordinal >= other.ordinal
}

// One callback for every state change. A method per event would mean a new
// method every time a state is added.
fun interface LifecycleObserver {
    fun onStateChanged(state: LifecycleState)
}

// An activity or a fragment, reduced to the only thing the holder needs. A
// test hands over a plain object with a lifecycle it drives by hand.
interface LifecycleOwner {
    val lifecycle: Lifecycle
}

// The minimum lifecycle that makes the design work. A current state, a list of
// listeners, and a way to move between states.
//
// Two behaviours are load bearing and both come from the real one. A new
// observer is told the current state straight away, which is how a late
// observer gets the current value with no special case in the holder. And
// moving to DESTROYED drops every observer, so a screen that is gone cannot
// keep a holder alive through a callback.
class Lifecycle {

    private val observers = mutableListOf<LifecycleObserver>()

    var currentState: LifecycleState = LifecycleState.CREATED
        private set

    fun addObserver(observer: LifecycleObserver) {
        if (currentState == LifecycleState.DESTROYED) return
        observers += observer
        observer.onStateChanged(currentState)
    }

    fun removeObserver(observer: LifecycleObserver) {
        observers -= observer
    }

    // Walking a copy, because an observer that hears DESTROYED removes itself
    // while we are still going through the list.
    fun moveTo(next: LifecycleState) {
        currentState = next
        observers.toList().forEach { it.onStateChanged(next) }
        if (next == LifecycleState.DESTROYED) observers.clear()
    }
}
package com.androidinterview.livedata.lifecycle

enum class LifecycleState {
    DESTROYED,
    CREATED,
    STARTED,
    RESUMED;

    fun isAtLeast(other: LifecycleState) = ordinal >= other.ordinal
}

fun interface LifecycleObserver {
    fun onStateChanged(state: LifecycleState)
}

interface LifecycleOwner {
    val lifecycle: Lifecycle
}

class Lifecycle {

    private val observers = mutableListOf<LifecycleObserver>()

    var currentState: LifecycleState = LifecycleState.CREATED
        private set

    fun addObserver(observer: LifecycleObserver) {
        if (currentState == LifecycleState.DESTROYED) return
        observers += observer
        observer.onStateChanged(currentState)
    }

    fun removeObserver(observer: LifecycleObserver) {
        observers -= observer
    }

    fun moveTo(next: LifecycleState) {
        currentState = next
        observers.toList().forEach { it.onStateChanged(next) }
        if (next == LifecycleState.DESTROYED) observers.clear()
    }
}

com.androidinterview.livedata.sample.LocationLiveData.kt

package com.androidinterview.livedata.sample

import com.androidinterview.livedata.core.LiveData
import com.androidinterview.livedata.thread.MainThread

data class Location(val latitude: Double, val longitude: Double)

// Stands in for the platform location client. It calls back on a thread of its
// own choosing, which is the exact case postValue exists for.
interface LocationProvider {
    fun start(listener: (Location) -> Unit)
    fun stop(listener: (Location) -> Unit)
}

// Why onActive and onInactive are worth having. A location listener costs
// battery, so it runs while a screen is watching and stops when the last one
// goes away. Rotation is the case that makes this subtle. The old activity is
// destroyed and the new one observes moments later, so a naive design stops the
// GPS and starts it again for nothing. Real code softens that with a short
// delay before releasing, and that is a good thing to mention.
class LocationLiveData(
    mainThread: MainThread,
    private val provider: LocationProvider,
) : LiveData<Location>(mainThread) {

    // postValue, not setValue, because the callback arrives on whichever
    // thread the provider felt like using.
    private val listener: (Location) -> Unit = { postValue(it) }

    override fun onActive() = provider.start(listener)

    override fun onInactive() = provider.stop(listener)
}
package com.androidinterview.livedata.sample

import com.androidinterview.livedata.core.LiveData
import com.androidinterview.livedata.thread.MainThread

data class Location(val latitude: Double, val longitude: Double)

interface LocationProvider {
    fun start(listener: (Location) -> Unit)
    fun stop(listener: (Location) -> Unit)
}

class LocationLiveData(
    mainThread: MainThread,
    private val provider: LocationProvider,
) : LiveData<Location>(mainThread) {

    private val listener: (Location) -> Unit = { postValue(it) }

    override fun onActive() = provider.start(listener)

    override fun onInactive() = provider.stop(listener)
}

com.androidinterview.livedata.sample.SingleLiveEvent.kt

package com.androidinterview.livedata.sample

import com.androidinterview.livedata.core.MutableLiveData
import com.androidinterview.livedata.core.Observer
import com.androidinterview.livedata.lifecycle.LifecycleOwner
import com.androidinterview.livedata.thread.MainThread
import java.util.concurrent.atomic.AtomicBoolean

// The famous workaround, written out so you can see why it is a workaround and
// not a fix. The holder replays its current value to whoever starts observing,
// which is right for a name on a screen and wrong for "show a toast", because
// rotation brings a new observer that is behind. The flag breaks the replay.
//
// Two honest flaws. Only one observer can win the flag, so a second observer on
// the same event silently gets nothing, and the wrapping means the caller can
// no longer remove the function it registered. The bug is that a state holder
// was asked to carry events, which is why one shot signals moved to a channel.
class SingleLiveEvent<T>(mainThread: MainThread) : MutableLiveData<T>(mainThread) {

    private val pending = AtomicBoolean(false)

    override fun observe(owner: LifecycleOwner, observer: Observer<T>) {
        super.observe(owner) { value ->
            if (pending.compareAndSet(true, false)) observer(value)
        }
    }

    override fun setValue(value: T) {
        pending.set(true)
        super.setValue(value)
    }
}
package com.androidinterview.livedata.sample

import com.androidinterview.livedata.core.MutableLiveData
import com.androidinterview.livedata.core.Observer
import com.androidinterview.livedata.lifecycle.LifecycleOwner
import com.androidinterview.livedata.thread.MainThread
import java.util.concurrent.atomic.AtomicBoolean

class SingleLiveEvent<T>(mainThread: MainThread) : MutableLiveData<T>(mainThread) {

    private val pending = AtomicBoolean(false)

    override fun observe(owner: LifecycleOwner, observer: Observer<T>) {
        super.observe(owner) { value ->
            if (pending.compareAndSet(true, false)) observer(value)
        }
    }

    override fun setValue(value: T) {
        pending.set(true)
        super.setValue(value)
    }
}

com.androidinterview.livedata.thread.MainThread.kt

package com.androidinterview.livedata.thread

// The main thread as an interface, so the holder never touches Looper or
// Handler and a test never needs a device. Two questions is all the holder
// asks. Am I on it, which guards setValue, and please run this, which is what
// postValue uses.
interface MainThread {
    val isMainThread: Boolean
    fun post(task: () -> Unit)
}

// A main thread you drain by hand, which is what Looper does in an app and
// what a test wants. Work posted from any thread lands on the queue, and
// whichever thread calls drain is the main thread for as long as that runs, so
// isMainThread is a real answer rather than a hardcoded true.
class QueuedMainThread : MainThread {

    private val queue = ArrayDeque<() -> Unit>()

    @Volatile
    private var draining: Thread? = null

    override val isMainThread: Boolean
        get() = Thread.currentThread() === draining

    override fun post(task: () -> Unit) {
        synchronized(queue) { queue.addLast(task) }
    }

    // Run everything queued, on this thread, in the order it was posted. A task
    // that posts another task is picked up by the same loop, which is how a
    // real message queue behaves.
    fun drain() {
        val previous = draining
        draining = Thread.currentThread()
        try {
            while (true) {
                val task = synchronized(queue) { queue.removeFirstOrNull() } ?: return
                task()
            }
        } finally {
            draining = previous
        }
    }

    // Run this block as if it were on the main thread, then let everything it
    // posted run as well.
    fun runOnMain(block: () -> Unit) {
        post(block)
        drain()
    }
}
package com.androidinterview.livedata.thread

interface MainThread {
    val isMainThread: Boolean
    fun post(task: () -> Unit)
}

class QueuedMainThread : MainThread {

    private val queue = ArrayDeque<() -> Unit>()

    @Volatile
    private var draining: Thread? = null

    override val isMainThread: Boolean
        get() = Thread.currentThread() === draining

    override fun post(task: () -> Unit) {
        synchronized(queue) { queue.addLast(task) }
    }

    fun drain() {
        val previous = draining
        draining = Thread.currentThread()
        try {
            while (true) {
                val task = synchronized(queue) { queue.removeFirstOrNull() } ?: return
                task()
            }
        } finally {
            draining = previous
        }
    }

    fun runOnMain(block: () -> Unit) {
        post(block)
        drain()
    }
}

Watch