Subjects¶
This page covers LazyFlowSubject building block for on-demand data loading along
with the metadata system that threads cross-cutting information through containers.
Table of Contents¶
- LazyFlowSubject
- Basic Usage
- Listening and Caching Behaviour
- Reloading Data
- Pushing Values Directly
- Replacing the Loader
- Load Triggers
- ContainerConfiguration
- listenReloadable
- whenActive
- SubjectFactory
- Testability
- Convenience Factory Functions
- LoaderDecorator
- Terminating a Load from a Decorator
- Metadata
- ContainerMetadata
- SourceType
- ReloadFunctionMetadata
- BackgroundLoadMetadata
- LoadTrigger
- Custom Metadata
- Flow Dependencies in Loader Functions
- dependsOnContainerFlow
- dependsOnFlow
- Key Stability
- Reload Configuration
LazyFlowSubject¶
LazyFlowSubject<T> converts a suspend loader function into a
Flow<Container<T>>. The loader runs lazily and its result is automatically
cached for subsequent subscribers.
Basic Usage¶
class ProductRepository(
private val localDataSource: ProductsLocalDataSource,
private val remoteDataSource: ProductsRemoteDataSource,
) {
private val productsSubject = LazyFlowSubject.create {
// Step 1: emit local (cached) data immediately if available:
val local = localDataSource.getProducts()
if (local != null) emit(local)
// Step 2: fetch from remote and update the local cache:
val remote = remoteDataSource.getProducts()
localDataSource.save(remote)
emit(remote)
}
// ListContainerFlow<T> is an alias for Flow<Container<List<T>>>
fun listenProducts(): ListContainerFlow<Product> = productsSubject.listen()
fun reload() = productsSubject.reloadAsync()
}
The lambda passed to LazyFlowSubject.create is the loader function. Inside
it, you have access to the Emitter<T> receiver, which provides:
emit(value): emit a loaded valueloadTrigger: LoadTrigger: why the loader was called (see Load Triggers)metadata: ContainerMetadata: metadata from the triggering calldependsOnFlow { }/dependsOnContainerFlow { }: subscribe to external flows (see Flow Dependencies)
Listening and Caching Behaviour¶
- The loader starts running the first time
listen()is called and a subscriber begins collecting - The latest loaded value is cached; new subscribers receive it immediately without re-running the loader
- When the last subscriber stops collecting, a timer starts (default: 1 second). If no new subscriber appears before the timer expires, the cached value is cleared and the loader will run again on the next subscription
- You can customise the timeout:
You can also configure how the initial load behaves via loadConfig and
metadata, the same options accepted by reload:
LazyFlowSubject.create(
loadConfig = LoadConfig.SilentLoading,
metadata = MyCustomMetadata,
) {
emit(loadData())
}
Reloading Data¶
Re-run the previous loader without changing it:
// fire-and-forget:
subject.reloadAsync()
// or observe the reload result:
subject.reload().collect { value -> ... }
You can reload silently (the existing cached value is kept visible while the new load runs in the background):
A config passed to reload / reloadAsync applies to that load only. The
subject keeps its own configuration - the one it was created with, or the one a
later newLoad established - so the next reloadAsync() with no config behaves
normally again. newLoad is the opposite: the config it receives becomes the
subject's new default (see Replacing the Loader).
reload / reloadAsync also accept an optional metadata: ContainerMetadata
that is merged into the container emitted by this reload - a convenient channel
for tagging why the reload happened. The same argument is available on
Container.reload(config, metadata) and LazyCache.reload(arg, config, metadata):
Attach a metadata type that implements ContainerMetadata.OneShot
when it should apply only to this single reload and not stick to later loads.
LoadConfigOneShotMetadata is a built-in one-shot type that carries a
LoadConfig for a single load. Passing a config to reload is exactly this,
spelled shorter. The metadata form exists for code paths that can pass metadata
but no explicit config - a custom metadata pipeline, or a dependency change (see
Reload Configuration, which uses this type
internally):
// this reload is silent...
subject.reloadAsync(metadata = LoadConfigOneShotMetadata(LoadConfig.SilentLoading))
// ...the next one is back to the subject's own config:
subject.reloadAsync()
Pushing Values Directly¶
Place a value into the subject without running the loader:
subject.updateWith(successContainer("immediate value"))
// update the existing container:
subject.updateWith { oldContainer ->
oldContainer.map { it + " (updated)" }
}
// shorthand that only runs if the current container is Success:
subject.updateIfSuccess { oldValue ->
oldValue.copy(isFavorite = true)
}
After updateWith, calling reload() re-runs the previous loader function.
Replacing the Loader¶
You can replace the loader function at any time:
// assign a new multi-value loader (returns a Flow<T> of the new results):
val flow: Flow<T> = subject.newLoad {
emit(step1())
emit(step2())
}
// fire-and-forget variant:
subject.newAsyncLoad { emit(loadData()) }
// single-value loader (suspends until the value is loaded):
val result: T = subject.newSimpleLoad { fetchData() }
// single-value loader (fire-and-forget):
subject.newSimpleAsyncLoad { fetchData() }
The newLoad / newSimpleLoad functions also accept a config argument and
an optional metadata argument:
config = LoadConfig.SilentLoading: the existing cached value stays visible while the new load runs in the background. Emitted containers carrybackgroundLoadState = BackgroundLoadState.LoadingwhenemitBackgroundLoadsis enabled (e.g. vialistenReloadable()).metadata: arbitrary data attached to this load call; accessible inside the loader function via themetadataproperty on theEmitterreceiver. See Metadata for available types and how to define custom ones.
subject.newAsyncLoad(
config = LoadConfig.SilentLoading,
metadata = SourceTypeMetadata(RemoteSourceType),
) {
emit(fetchRemote())
}
Load Triggers¶
Inside the loader function, loadTrigger tells you why the loader was
called. Use it to skip unnecessary steps, e.g. the local-cache check on
explicit reloads:
private val productsSubject = LazyFlowSubject.create {
if (loadTrigger != LoadTrigger.Reload) {
val local = localDataSource.getProducts()
if (local != null) emit(local)
}
val remote = remoteDataSource.getProducts()
localDataSource.save(remote)
emit(remote)
}
| Value | Meaning |
|---|---|
LoadTrigger.NewLoad |
Loader was set with newLoad() or create {} |
LoadTrigger.Reload |
Loader was re-triggered by reload() / reloadAsync() |
LoadTrigger.CacheExpired |
Cache timeout elapsed; the next subscriber triggered a fresh load |
ContainerConfiguration¶
Pass a ContainerConfiguration to listen() to control what extra metadata
is attached to emitted containers:
val flow: Flow<Container<List<Product>>> = subject.listen(
configuration = ContainerConfiguration(
emitReloadFunction = true, // attach a reload function to each container
emitBackgroundLoads = true, // set BackgroundLoadState metadata value to Loading while reloading silently
)
)
emitReloadFunction- each emittedContainer.Success/Container.Errorcarries areloadFunctionthat, when called, triggersreloadAsync()on the subject. Useful for UI components that need to offer a "retry" button without knowing about the subject directly.emitBackgroundLoads- while a silent reload is in progress, emitted containers havebackgroundLoadStatemetadata property. Useful when you want to display the current loaded data along with an additional indication that something is being loaded right now (for example, PullToRefresh behavior)
listenReloadable¶
listenReloadable() is a convenience shorthand that enables both flags:
// equivalent to listen(ContainerConfiguration(emitReloadFunction = true, emitBackgroundLoads = true))
fun listenProducts(): Flow<Container<List<Product>>> = subject.listenReloadable()
whenActive¶
whenActive() optional call registers a suspendable block of code that is
executed whenever the subject becomes active (when at least one subscriber starts
collecting a flow returned by the subject):
val productsSubject = LazyFlowSubject
.create<List<Product>> {
emit(loadProducts())
}
.whenActive {
// e.g. subscribe to other flow and update the subject when
// something is changed
}
SubjectFactory¶
SubjectFactory is an interface for creating LazyFlowSubject instances.
Use it instead of calling LazyFlowSubject.create {} directly.
Testability¶
Injecting SubjectFactory via DI (e.g. Hilt) makes it straightforward to
replace the real factory with a test double that returns pre-configured or mock
subjects:
class ProductRepository(
private val subjectFactory: SubjectFactory = SubjectFactory,
) {
private val subject = subjectFactory.createSubject {
delay(1000)
emit("my-item")
}
fun listen(): ContainerFlow<String> = subject.listenReloadable()
}
In production, bind DefaultSubjectFactory with your chosen cache timeout:
@Provides
@Singleton
fun provideSubjectFactory(): SubjectFactory =
DefaultSubjectFactory(cacheTimeoutMillis = 60_000L)
In tests, replace it with any fake implementation of SubjectFactory or a mock.
The global default can also be overridden for tests:
Convenience Factory Functions¶
SubjectFactory provides several extension functions to reduce boilerplate:
// Create a LazyFlowSubject with a simple (single-value) loader:
val subject: LazyFlowSubject<String> = subjectFactory.createSimpleSubject(
sourceType = RemoteSourceType,
) { fetchString() }
// Create a Flow directly (backed by a LazyFlowSubject internally):
val flow: Flow<Container<String>> = subjectFactory.createFlow {
emit(fetchData())
}
// Create a reloadable Flow (emitReloadFunction + emitBackgroundLoads enabled):
val flow: Flow<Container<String>> = subjectFactory.createReloadableFlow {
emit(fetchData())
}
LoaderDecorator¶
A LoaderDecorator wraps every loader function of a subject or a cache, so
cross-cutting logic (session checks, logging, error mapping) lives in
one place instead of being repeated in each loader:
public fun interface LoaderDecorator {
public suspend fun DecoratedFlowComposer.decorate(originLoader: suspend () -> Unit)
}
The receiver is a DecoratedFlowComposer, which extends FlowComposer, so a
decorator can declare its own
flow dependencies. A typical use
case is failing every load while there is no valid session, and re-running all
loaders as soon as a new token appears:
val sessionDecorator = LoaderDecorator { originLoader ->
val token: String = dependsOnFlow("session-token") { sessionManager.tokenFlow }
if (token.isBlank()) throw NoSessionException()
originLoader()
}
Two rules:
- The implementation must call
originLoader(), otherwise nothing is emitted and the load fails with anIllegalStateException. Throwing your own exception instead is fine - it fails the load like any error raised by the loader itself. The only other way out is to terminate the load explicitly. - Choose dependency keys that cannot clash with the keys used by the loaders being decorated (see Key Stability).
Terminating a Load from a Decorator¶
Throwing an exception from a decorator fails the load like any other error, so
it still obeys the current load configuration: with
LoadConfig.SilentLoadingAndError the previously cached value is kept and the
error is only reported as background state. When a decorator needs to override
that policy, DecoratedFlowComposer offers two terminating functions:
public interface DecoratedFlowComposer : FlowComposer {
public fun completeWithFailure(exception: Exception): Nothing
public fun completeWithCacheCleanUp(): Nothing
}
completeWithFailure(exception)finishes the load with an error container. The exception reaches collectors regardless of the silent error policy.completeWithCacheCleanUp()finishes the load with a pending container, so any cached value is dropped and collectors go back to the loading state.
Both functions return Nothing: they unwind the decorator body immediately, so
the origin loader is not executed if it has not been called yet. Calling them
after originLoader() discards whatever the loader emitted.
A sign-out decorator that must not leave stale data behind:
val sessionDecorator = LoaderDecorator { originLoader ->
val session = dependsOnFlow("session") { sessionManager.sessionFlow }
when (session) {
// no session at all: wipe cached data, show the loading state
is Session.SignedOut -> completeWithCacheCleanUp()
// expired token: always surface the error, even for silent loads
is Session.Expired -> completeWithFailure(SessionExpiredException())
is Session.Active -> originLoader()
}
}
For page loaders the same functions terminate the whole paging session, not only the page being loaded.
Install it wherever a loader is configured:
// a single subject:
LazyFlowSubject.create(loaderDecorator = sessionDecorator) {
emit(loadData())
}
// a cache (applies to the subject of every argument):
LazyCache.create(loaderDecorator = sessionDecorator) { arg ->
emit(loadData(arg))
}
// every subject and cache produced by a factory:
DefaultSubjectFactory(loaderDecorator = sessionDecorator)
For page loaders, the decorator wraps each page load separately, not the paging session as a whole - so a decorator that waits for a valid token does so before every page request.
Stores expose the same hook via setLoaderDecorator(...) on any store builder,
or via SimpleStoreFactory(loaderDecorator = ...); see the
Store documentation.
Metadata¶
ContainerMetadata¶
ContainerMetadata is an immutable bag attached to Container.Success and
Container.Error. Multiple metadata instances can be combined:
val meta = SourceTypeMetadata(RemoteSourceType) + ReloadFunctionMetadata { _, _ -> loadItems() }
val container = successContainer("data", meta)
When two metadata instances of the same type are combined, the second one replaces the first:
val combined = SourceTypeMetadata(LocalSourceType) + SourceTypeMetadata(RemoteSourceType)
// result: only RemoteSourceType is kept
Access a specific metadata type:
Or use the shorthand extension properties available on containers:
val source: SourceType = container.sourceType
val bgLoadState: BackgroundLoadState = container.backgroundLoadState
val reloadFn = container.metadata.reloadFunction
SourceType¶
SourceType communicates where the data came from. The built-in values are:
| Value | Meaning |
|---|---|
LocalSourceType |
Loaded from a local/on-device data source |
RemoteSourceType |
Fetched from a remote/network data source |
ImmediateSourceType |
Set directly, not via a loader |
FakeSourceType |
Provided by a test double or fake implementation |
UnknownSourceType |
Source is not known |
Emit with a source type inside the loader:
LazyFlowSubject.create {
emit(localDataSource.get(), LocalSourceType)
// isLastValue arg is optional, but it can improve performance a bit:
emit(remoteDataSource.get(), RemoteSourceType, isLastValue = true)
}
Set a source type on an existing container:
val container = successContainer("data", SourceTypeMetadata(RemoteSourceType))
// or update metadata on an existing container:
val updated = container.update { sourceType = RemoteSourceType }
Read the source type:
ReloadFunctionMetadata¶
Attaching a reload function to a container lets UI components trigger a reload without needing a direct reference to the subject / view-model, or any other components:
// The listenReloadable() shorthand does this automatically:
fun listenProducts() = productsSubject.listenReloadable()
// Or attach manually on an existing container:
val container = successContainer("data") + ReloadFunctionMetadata { _, _ -> loadMyData() }
// Or override the reload function in a flow:
val flow = source.containerUpdate {
val originalReload = reloadFunction
reloadFunction = { config, metadata ->
println("Reloading...")
originalReload(config, metadata) // call the original reload function if needed
}
}
Calling the reload function from UI code:
val container: Container<String> = ...
container.fold(
onError = { ex -> Button(onClick = { container.reload() }) { Text("Retry") } },
onSuccess = { value -> /* ... */ },
)
BackgroundLoadMetadata¶
When a silent reload is in progress and emitBackgroundLoads = true is set,
emitted containers carry backgroundLoadState = Loading metadata value. UI can use this to
show an indicator (e.g. pull-to-refresh) while still displaying the stale data:
container.fold(
onSuccess = { value ->
if (backgroundLoadState == BackgroundLoadState.Loading) ShowRefreshIndicator()
ShowContent(value)
},
)
LoadTrigger¶
Available inside the loader function via loadTrigger: LoadTrigger:
LazyFlowSubject.create {
when (loadTrigger) {
LoadTrigger.NewLoad -> { /* first-ever load or newLoad() called */ }
LoadTrigger.Reload -> { /* explicit reload() call */ }
LoadTrigger.CacheExpired -> { /* fresh load after cache timed out */ }
}
}
Custom Metadata¶
You can define your own metadata types by implementing ContainerMetadata:
data class TimestampMetadata(val timestamp: Long) : ContainerMetadata
// attach:
val container = successContainer("data", TimestampMetadata(System.currentTimeMillis()))
// read:
val ts: Long? = container.metadata.get<TimestampMetadata>()?.timestamp
Implement ContainerMetadata.Hidden to prevent the metadata from being seen
by downstream collectors (it is still passed through internally and visible from
the loader function):
Implement ContainerMetadata.OneShot for metadata that is relevant only to the
single load request it was attached to (for example, a value passed to
reloadAsync(metadata = …)). A one-shot value is emitted with that load's
container and stays attached while the value lives in the in-memory cache
(including re-emission to a new collector that subscribes before the cache
expires), but it is not carried into any new load - the next reload, query
change, dependency update, or post-cache-expiry reload produces containers
without it. This keeps transient signals (e.g. "this refresh came from a push")
from sticking to later, unrelated loads.
data object PushRefreshMetadata : ContainerMetadata, ContainerMetadata.OneShot
subject.reloadAsync(metadata = PushRefreshMetadata) // emitted container carries it...
subject.reloadAsync() // ...the next reload does not
The behaviour can be disabled per instance by overriding isOneShot to return
false, in which case the metadata behaves like an ordinary one and survives
subsequent reloads.
Flow Dependencies in Loader Functions¶
Starting from v2.0.0-beta13, loader functions can subscribe to external Kotlin flows. When the subscribed flow emits a new value, the loader function is automatically re-executed.
dependsOnContainerFlow¶
Use dependsOnContainerFlow to depend on a Flow<Container<T>>. If the
dependent flow emits Container.Error, the current load is failed with the
same exception. If it emits Container.Pending, the load waits:
interface SessionProvider {
fun getCurrentUserFlow(): Flow<Container<User>>
}
private val itemsSubject = LazyFlowSubject.create {
// The loader re-runs whenever the current user changes:
val currentUser: User = dependsOnContainerFlow("getCurrentUser") {
sessionProvider.getCurrentUserFlow()
}
val items = remoteDataSource.getItems(currentUser)
emit(items, RemoteSourceType)
}
dependsOnFlow¶
For plain Flow<T> dependencies (not wrapped in Container):
Key Stability¶
Every dependsOnFlow / dependsOnContainerFlow call must be given a stable
key (or key + arguments) that uniquely identifies the flow instance. The keys
are used to cache the subscribed flows across re-executions of the loader:
// simple key:
val user: User = dependsOnContainerFlow("getCurrentUser") {
sessionProvider.getCurrentUserFlow()
}
// key with arguments (important if the argument changes the flow):
val userId: String = sessionProvider.getCurrentUserId()
val user: User = dependsOnContainerFlow("getUserById", userId) {
userRepository.getUserById(userId)
}
If you call dependsOnFlow or dependsOnContainerFlow with the same key
twice within one loader execution, the second call is ignored and returns the
cached result from the first call:
val a: String = dependsOnContainerFlow("key") { getFlow1() }
val b: String = dependsOnContainerFlow("key") { getFlow2() } // getFlow2 is ignored
// a == b, both refer to the results from the first call
Reload Configuration¶
A reload caused by a dependency change behaves like an ordinary reload: it uses
the load config the subject is currently working with, so by default the
container goes back to Pending while the loader re-runs.
Pass a FlowComposer.Config as one of the keys to configure the reloads
triggered by that particular dependency:
private val starsSubject = LazyFlowSubject.create {
val config = FlowComposer.Config(LoadConfig.SilentLoading)
val filter: StarFilter = dependsOnFlow("filter", config) { filterFlow }
emit(starsDataSource.fetchStars(filter))
}
| Property | Default | Meaning |
|---|---|---|
loadConfig |
null |
Load config applied to the reload triggered by this dependency. null keeps the config the subject is already using; LoadConfig.SilentLoading keeps the currently loaded value visible while the loader re-runs |
reloadDependencies |
false |
When true, every flow dependency of the loader is asked to reload itself before the loader re-runs, by invoking the reloadFunction attached to its last container. Dependencies whose containers carry no reload function are unaffected |
The config is part of the dependency key, so keep it a stable value (a val or
a constant) exactly like the other keys.
Note. Explicit
reload()/reloadAsync()calls always reload the dependencies, regardless of this setting.Behaviour change in 3.5.0. Before 3.5.0, dependency-triggered reloads were always silent and ignored the subject's own load config. They now follow the subject's config unless a
FlowComposer.Configsays otherwise. To restore the previous behaviour, passFlowComposer.Config(LoadConfig.SilentLoading)as shown above.