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
- Register an observer. If the screen is started, send the current value.
- Call
setValue(2). Store 2 and increase the version. - Notify each started observer and record the version it received.
- Skip stopped observers. On restart, their older version tells us to send the latest value.
- 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