netflow
0.6.0indexedLightweight, flexible network library offering a clean, intuitive API for handling network requests with support for LiveData, Flow, object deserialization, customizable headers, and local data integration.
Lightweight, flexible network library offering a clean, intuitive API for handling network requests with support for LiveData, Flow, object deserialization, customizable headers, and local data integration.
A networking layer for Kotlin Multiplatform: one client, one state model, and one testing story across plain API calls, local-cache / offline-first flows, and Jetpack Paging 3.
Start with annotated interfaces for your straightforward endpoints, then drop to the call {} DSL on the same client, for the same API, when an endpoint needs caching, offline reads, or paging. The annotations are the familiar front door; the DSL is the engine behind it.
Both entry points end up on the same pipeline. Annotations are KSP-generated
code that calls the exact same RequestBuilder the DSL builds by hand:
MockNetFlowClient is a second implementation, not a wrapper
around the real one. It runs requests through the exact same
and , so interceptors, retry, and auth
are exercised in tests exactly as in production. Only the bottom step changes:
instead of hitting OkHttp/NSURLSession, your
function receives a and returns a
directly, no network involved. Every call is also recorded, so
/ / can
check what was sent without touching a socket.
dependencies {
implementation("io.github.kmpbits:netflow-core:<latest_version>")
}
Adds responsePaginated with Jetpack Paging 3 support.
dependencies {
implementation("io.github.kmpbits:netflow-core:<latest_version>")
implementation("io.github.kmpbits:netflow-paging:<latest_version>")
}
Check the latest versions on Maven Central.
Adds SettingsTokenStorage for persisting auth tokens (see Authentication).
dependencies {
implementation("io.github.kmpbits:netflow-core:<latest_version>")
implementation("io.github.kmpbits:netflow-token-storage:<latest_version>")
}
Declare your API as an annotated interface and let a KSP processor generate the implementation. Works on all Kotlin Multiplatform targets, no runtime reflection.
This is the on-ramp, not a separate library. Use annotations for the plain
request/response endpoints, the 80% case. The generated interface returns the
same ResultState / AsyncState / PagingData types the DSL uses, runs on the
same NetFlowClient, and is tested with the same . When an
endpoint needs local caching, offline reads, side effects, or
remote+local paging, write that one method with instead;
nothing else changes. You never juggle two HTTP stacks or two state models.
plugins {
id("com.google.devtools.ksp")
}
dependencies {
implementation("io.github.kmpbits:netflow-core:<latest_version>")
implementation("io.github.kmpbits:netflow-annotations:<latest_version>")
add("kspCommonMainMetadata", "io.github.kmpbits:netflow-ksp:<latest_version>")
}
kotlin.sourceSets.commonMain {
kotlin.srcDir("build/generated/ksp/metadata/commonMain/kotlin")
}
tasks.withType<org.jetbrains.kotlin.gradle.tasks.KotlinCompilationTask<*>>().configureEach {
if (name != "kspCommonMainKotlinMetadata") dependsOn("kspCommonMainKotlinMetadata")
}
Supported annotations
Supported return types
Not yet supported: the onNetworkSuccess / local {} cache hooks (including
remote+local paging). Use the call {} DSL directly for those.
Annotated functions return the API DTO. Map to your domain model in the
repository with the map helpers from netflow-core: AsyncState.map,
ResultState.map, and Flow.map:
class TodoRepository(client: NetFlowClient) {
private val api = client.createTodoApi()
suspend fun getTodos(): AsyncState<List<Todo>> =
api.getTodos(completed = null).map { dtos -> dtos.map { it.toModel() } }
fun observeTodo(id: Int): Flow<ResultState<Todo>> =
api.observeTodo(id).map { state -> state.map { it.toModel() } }
}
This is deliberate: the two-type transform stays one layer out of the
annotations, so the mapping is always explicit and compiler-checked, the same
guarantee the call {} DSL gives with its required transform parameter.
A function returning Flow<PagingData<T>> (where T extends PagingModel)
generates a network-only paged call. @Paginated(pageQueryName, pageSize)
overrides the defaults ("page", 20).
@GET("posts")
fun pagedPosts(@Query tag: String?): Flow<PagingData<PostDto>>
Map to domain in the repository with PagingData.map:
fun pagedPosts(tag: String?) = api.pagedPosts(tag).map { it.map { dto -> dto.toModel() } }
Remote + local paging (RemoteMediator, local PagingSource, insert/delete
callbacks) stays on the call {} DSL: responsePaginated { localSource(...) }.
NetFlowCall)Return NetFlowCall and the generated method stops at the request. Finish it in
the repository with any responseX function: this is how you add a local cache
or onNetworkSuccess side effects to an annotated endpoint:
@GET("todos")
fun todos(): NetFlowCall
// repository
fun getTodos(): Flow<ResultState<Todo>> =
api.todos().responseFlow<TodoDto, Todo>(transform = { it.toModel() }) {
onNetworkSuccess { db.insertTodos(it) }
local({ observe { db.todos() } }, transform = { it.toModel() })
}
@Wrapped / @Paginated are not allowed on a NetFlowCall method; those are
choices you make on the responseX call. The method can be suspend or not; the
responseX you compose on the result is already suspending either way.
client.prepareCall { … } builds a from the hand-written DSL too.
A suspend function returning a plain type maps to responseToModel<T>(): you
get the deserialized value, or an HttpException on a non-2xx response. This is
the default style, for code that prefers try/catch (or a global
handler) over a sealed state.
@GET("todos/{id}")
suspend fun getTodo(@Path id: Int): TodoDto
@GET("todos")
suspend fun getTodos(): List<TodoDto>
It returns the DTO, same as every other shape, so map to your domain type in the
repository. @Wrapped isn't supported here; use AsyncState<T> with @Wrapped
for envelope APIs, or return NetFlowCall. Unit isn't allowed either; use
AsyncState<Unit>.
val client = netflowClient {
baseUrl = "https://api.example.com"
header(Header(HttpHeader.custom("custom-header"), "value"))
header(Header(HttpHeader.CONTENT_TYPE), "application/json")
}
val response = client.call {
path = "/users"
method = HttpMethod.Get
}.response()
val user: User = client.call {
path = "/users/1"
}.responseToModel<User>()
client.call {
method = HttpMethod.Post
path = "/todos/1/attachments"
multipart {
part("caption", "before")
filePart("file", filename = "shot.png", bytes = imageBytes, contentType = "image/png")
onProgress { sent, total -> println("$sent / $total") }
}
}.response()
multipart { } holds each part's bytes in memory. It requires POST, PUT or PATCH
and cannot be combined with body(...). onProgress { sent, total -> } is optional
and reports upload byte progress; it fires on whatever thread the platform delivers
it on, so hop to your UI thread yourself if you're updating UI from it.
@NetFlowApi
interface FilesApi {
@Streaming
@GET("files/{id}")
fun download(@Path id: String): Flow<ByteArray>
}
api.download("report.zip")
.catch { e -> if (e is HttpException) showError(e.code) e }
.collect { chunk -> output.write(chunk) }
Or from the DSL: client.call { path = "files/report.zip" }.responseStream().
The body is delivered in chunks as it arrives instead of being buffered in memory
(OkHttp's body.source() on Android, an NSURLSessionDataDelegate on iOS, which
pauses the task when the collector falls behind). The flow is cold: nothing is sent
until you collect it, and cancelling the collector cancels the request.
HttpException(code, errorBody); a failure
mid-stream is propagated as-is.Download progress. onDownloadProgress { received, total -> } (or, on an annotated
interface, a @Progress onProgress: ((Long, Long) -> Unit)? parameter on the @Streaming
function) reports the cumulative bytes as they arrive from the network:
client.call {
path = "files/report.zip"
onDownloadProgress { received, total ->
if (total > 0) println("${received * 100 / total}%") else println("$received bytes")
}
}.responseStream().collect { chunk -> output.write(chunk) }
total is the Content-Length, or -1 when the server sent none (chunked or
transparently compressed responses). The callback fires on whatever thread the platform
delivers data on, restarts from 0 if the request is retried, and isn't called for a
non-2xx response. The multipart block's onProgress remains upload-only.
When your DTO and domain model are the same type, pass only one type parameter:
val flow = client.call {
path = "/users/1"
}.responseFlow<UserDto>()
When ApiType and DisplayType differ, pass transform as the first argument. The compiler enforces this: forgetting it is a build error, not a runtime crash.
val flow = client.call {
path = "/users/1"
}.responseFlow<UserDto, User>(transform = { it.toModel() })
val usersFlow = client.call {
path = "/users"
method = HttpMethod.Get
}.responseFlow<UserDto, User>(transform = { it.toModel() }) {
onNetworkSuccess { dto ->
queries.insertUser(dto.toEntity())
}
local({ observe { queries.getUser() } }, transform = { it.toModel() })
}
The transform inside local() maps from the database entity type directly to DisplayType, driving what gets shown while the network call is in flight. The transform on the function maps the network ApiType to DisplayType once the response arrives.
local({
onlyLocalCall = true
call { queries.getAllUsers() }
}, transform = { it.toModel() })
For APIs that return { "data": { ... } } instead of a plain object:
// Same type
responseWrappedFlow<UserDto>()
// Different types
responseWrappedFlow<UserDto, User>(transform = { it.toModel() })
Or set wrappedResponse = true inside the builder when using responseFlow:
responseFlow<UserDto, User>(transform = { it.toModel() }) {
wrappedResponse = true
}
// Same type
responseListFlow<UserDto>()
responseWrappedListFlow<UserDto>()
// Different types
responseListFlow<UserDto, User>(transform = { it.toModel() })
responseWrappedListFlow<UserDto, User>(transform = { it.toModel() })
lifecycleScope.launch {
usersFlow.collectLatest { state ->
when (state) {
is ResultState.Loading -> showLoading()
is ResultState.Success -> showUsers(state.data)
is ResultState.Error -> showError(state.error.message)
}
}
}
For one-shot suspending calls (no observation needed).
suspend fun deleteUser(id: Int): AsyncState<Unit> {
return client.call {
path = "users/$id"
method = HttpMethod.Delete
}.responseAsync<Unit> {
onNetworkSuccess { queries.deleteUser(id) }
}
}
suspend fun getUser(id: Int): AsyncState<User> {
return client.call {
path = "users/$id"
}.responseAsync<UserDto, User>(transform = { it.toModel() })
}
// Same type
responseListAsync<UserDto>()
responseWrappedListAsync<UserDto>()
// Different types
responseListAsync<UserDto, User>(transform = { it.toModel() })
responseWrappedListAsync<UserDto, User>(transform = { it.toModel() })
// Same type
responseWrappedAsync<UserDto>()
// Different types
responseWrappedAsync<UserDto, User>(transform = { it.toModel() })
responsePaginated integrates Jetpack Paging 3, supporting both network-only and remote+local strategies.
Your API response model must implement PagingModel:
@Serializable
data class PostDto(
val id: Int,
val title: String,
override var page: Int = 0,
override var lastUpdatedTimestamp: Long = 0L
) : PagingModel()
fun getPosts(): Flow<PagingData<Post>> = client.call {
path = "/posts"
}.responsePaginated<PostDto, Post> {
onlyApiCall = true
networkTransform { it.toModel() }
}
There are two ways to configure the local data source.
localQuery (recommended, no custom PagingSource needed)Pass countQuery, itemsQuery, and an invalidation flow. The library creates and manages the PagingSource internally. The invalidation flow triggers a reload whenever the underlying data changes: SQLDelight users pass query.asFlow(), Room users pass their Flow<List<T>>.
fun getPosts(): Flow<PagingData<Post>> = client.call {
path = "/posts"
}.responsePaginated<PostDto, Post> {
localQuery(
countQuery = { database.postQueries.countPosts().executeAsOne() },
itemsQuery = { limit, offset -> database.postQueries.selectPosts(limit, offset).executeAsList() },
invalidation = database.postQueries.selectAllPosts().asFlow(),
transform = { it.toModel() }
)
deleteOnRefresh = false
insertAll(transform = { it.toEntity() }) { posts ->
database.postQueries.transaction {
database.postQueries.deleteAll()
posts.forEach { database.postQueries.insertPost(it) }
}
}
firstItemDatabase(
itemDatabase = { database.postQueries.getFirstPost().executeAsOneOrNull() },
timestamp = { it.lastUpdatedTimestamp }
)
}
localSource / localSourceLong (custom PagingSource)Use this when you need full control over how data is loaded locally. You provide your own PagingSource<Int, E> (or PagingSource<Long, E> via localSourceLong).
Then wire it up:
fun getPosts(): Flow<PagingData<Post>> = client.call {
path = "/posts"
}.responsePaginated<PostDto, Post> {
localSource(
pagingSource = { PostPagingSource(database) },
transform = { it.toModel() }
)
// ...
}
For SQLDelight sources that use Long keys (e.g. QueryPagingSource), use localSourceLong instead; keys are bridged to Int internally.
When using a custom PagingSource (Option B), it is critical to register a listener on your database query to trigger invalidation. Without this, the UI will not update when data changes (e.g., after a network refresh or a local deletion).
If you are using SQLDelight, follow the pattern in the example above (and in the sample app's TodoPagingSource):
Query.Listener that calls invalidate() and removes itself.init.This ensures that whenever the underlying data changes, the PagingSource is marked as invalid, and the Pager will create a new one and reload the data.
val posts = repository.getPosts().cachedIn(viewModelScope)
val posts = viewModel.posts.collectAsLazyPagingItems()
LazyColumn {
items(count = posts.itemCount, key = posts.itemKey { it.id }) { index ->
posts[index]?.let { PostItem(it) }
}
}
netflow-paging ships PagingCollectionViewController, a KMP class that bridges paging data to Swift. It is designed to be used with SKIE for async sequence support.
ViewModel (Swift)
View (SwiftUI)
MockNetFlowClient implements NetFlowClient and intercepts all requests without making any real network calls. It supports response delays, request recording, and assertion helpers.
val mockClient = MockNetFlowClient { request ->
when {
request.path == "posts" && request.method == HttpMethod.Get ->
NetFlowMockResponse.success("""[{"id":1,"title":"Hello","completed":false}]""")
request.path.startsWith("posts/") && request.method == HttpMethod.Delete ->
NetFlowMockResponse.success()
request.path == "posts" && request.method == HttpMethod.Post ->
NetFlowMockResponse.success("""{"id":101,"title":"New Post","completed":false}""")
else -> NetFlowMockResponse.notFound()
}
}
NetFlowMockResponse.success(
body = """[...]""",
delay = 2.seconds // simulates slow network
)
NetFlowMockResponse.error(code = 401, errorBody = "Unauthorized")
NetFlowMockResponse.serverError("Something went wrong")
NetFlowMockResponse.notFound()
// Called at least once
mockClient.assertCalled("posts", HttpMethod.Get)
// Called exactly N times
mockClient.assertCalledTimes("posts/1", HttpMethod.Delete, times = 1)
// Never called
mockClient.assertNotCalled("posts", HttpMethod.Post)
val request = mockClient.recordedRequests.first()
assertEquals(HttpMethod.Post, request.method)
assertEquals(mapOf("title" to "New Post"), request.body)
mockClient.clearRecordedRequests()
@Test
fun `delete removes item from local database`() = runTest {
val mockClient = MockNetFlowClient { _ -> NetFlowMockResponse.success() }
val repository = PostRepositoryImpl(mockClient, database)
repository.deletePost(id = 1)
mockClient.assertCalled("posts/1", HttpMethod.Delete)
}
All helpers accept an optional delay: Duration parameter.
Configure auth { } once on the client. NetFlow adds the bearer token to every
request. On a 401 it calls your refresh block once, retries the failed request,
and holds a lock while it does, so ten requests that 401 at the same time trigger
one refresh instead of ten. You write the two lambdas. NetFlow decides when to
call them.
val client = netflowClient {
baseUrl =
auth {
loadTokens { tokenStore.read() }
refreshTokens { raw ->
res = raw.call {
path =
method = HttpMethod.Post
body(RefreshRequest(refreshToken))
}.responseAsync<TokenResponse>()
(res) {
AsyncState.Success -> BearerTokens(res..access, res..refresh)
->
}
}
}
}
client.authState
client.setTokens(tokens)
client.clearTokens()
client.call { path = ; skipAuth() }
raw is an auth-free client, so a refresh call can't recurse into another
refresh. It's a NetFlowClient too, so a generated @NetFlowApi interface works
there:
refreshTokens { raw ->
when (val res = raw.createAuthApi().refresh(RefreshRequest(refreshToken!!))) {
is AsyncState.Success -> BearerTokens(res.data.accessToken, res.data.refreshToken)
else -> null
}
}
Generated @NetFlowApi implementations run on the same client, so the token
attach and refresh-and-retry apply to them too. Put @SkipAuth on the endpoints
that don't need a token (login, sign-up, refresh) so the generated code doesn't
attach one or try to refresh on their 401:
@NetFlowApi
interface AuthApi {
@SkipAuth @POST("auth/login") suspend fun login( body: ): AsyncState<TokenResponse>
: AsyncState<TokenResponse>
}
By default the refresh is reactive: send the request, refresh on 401. Set
refreshLeeway and NetFlow also refreshes before sending, once the JWT is close
to expiring, so you skip the failed round-trip:
auth {
refreshLeeway = 30.seconds // refresh if the token expires within 30s
refreshTokens { /* ... */ }
}
It reads the token's exp claim (no signature check, that's the server's job). If
the token isn't a JWT, has no exp, or the device clock is off, it falls back to
the reactive 401 path, which is always on. null (the default) turns it off.
By default tokens live in memory only. Give auth { } a TokenStorage and
NetFlow seeds from load() on the first request and writes back on every change
(refresh, setTokens) and when the session ends (clearTokens, a failed
refresh):
interface TokenStorage {
suspend fun load(): BearerTokens?
suspend
}
client = netflowClient {
baseUrl =
auth {
storage(myTokenStorage)
refreshTokens { BearerTokens() }
}
}
loadTokens { } still works and wins over storage(...) for the initial seed.
If a storage call throws, NetFlow swallows it and keeps using the in-memory token.
netflow-core ships InMemoryTokenStorage() for tests, or for an app that wants
a fresh login every launch. It pulls in no platform storage libraries.
Persistent storage: netflow-token-storage
implementation("io.github.kmpbits:netflow-token-storage:<latest_version>")
SettingsTokenStorage(settings) stores the token pair over a
multiplatform-settings Settings. No encryption of its own, so back it with a
secure store:
// androidMain
val settings = SharedPreferencesSettings(
EncryptedSharedPreferences.create(
context, "netflow_auth",
MasterKey.Builder(context).setKeyScheme(MasterKey.KeyScheme.AES256_GCM).build(),
EncryptedSharedPreferences.PrefKeyEncryptionScheme.AES256_SIV,
EncryptedSharedPreferences.PrefValueEncryptionScheme.AES256_GCM,
)
)
// iosMain
val settings = KeychainSettings(service = "netflow_auth")
// shared
auth {
storage(SettingsTokenStorage(settings))
refreshTokens { /* ... */ }
}
androidx.security:security-crypto is deprecated and has no direct replacement.
If that bothers you, use a plain SharedPreferencesSettings (app-private and
sandboxed, but not encrypted at rest) or your own Keystore-backed store. The iOS
Keychain doesn't have this problem.
For DataStore, SQLDelight, or anything else, implement TokenStorage yourself.
It's three suspend functions.
The sample module has the full wiring: createTokenStorage() as an
expect/actual, Keychain on iOS and EncryptedSharedPreferences on Android,
under sample/.../core/di/.
Everything here also works through MockNetFlowClient(auth = { ... }), so the
whole refresh flow is testable without a network.
client.call {
path = "/secure-endpoint"
header(Header(HttpHeader.custom("Authorization"), "Bearer $token"))
}.responseFlow<SecureDataDto, SecureData>(transform = { it.toModel() })
client.call {
path = "/users"
parameter("role" to "admin")
parameter("active" to true)
}.responseFlow<UserDto, User>(transform = { it.toModel() })
client.call {
path = "/unstable-endpoint"
retry {
times = RetryTimes.THREE
delay = 1.seconds
retryOn = { it is IOException }
}
}.responseFlow<DataDto, Data>(transform = { it.toModel() })
Write it once in commonMain, it runs on both engines.
val trace = NetFlowInterceptor { chain ->
val request = chain.request.newBuilder()
.header(Header(HttpHeader.custom("X-Trace-Id"), randomTraceId()))
.build()
chain.proceed(request)
}
val timing = NetFlowInterceptor { chain ->
val start = TimeSource.Monotonic.markNow()
val response = chain.proceed(chain.request)
log("${chain.request.url} -> ${response.code} in ${start.elapsedNow()}")
response
}
val client = netflowClient {
baseUrl = "https://api.example.com"
addInterceptor(trace)
addInterceptor(timing)
}
Registration order is execution order: trace is the outermost.
Not calling chain.proceed(...) short-circuits the chain: the network is never touched:
val offline = NetFlowInterceptor { chain ->
cached(chain.request.url) ?: chain.proceed(chain.request)
}
The interceptor runs inside the retry loop and after auth: it sees the final
Authorization header and is called once per attempt.
Since the chain wraps the engine, MockNetFlowClient runs it too, so your interceptors
are testable without a network call:
val client = MockNetFlowClient(interceptors = listOf(trace)) {
NetFlowMockResponse.success(body = "[]")
}
client.call { path = "todos" }.response()
assertEquals("abc", client.recordedRequests.single()["X-Trace-Id"])
netflowClient {
baseUrl = "https://api.example.com"
pinning {
pin("api.example.com", "sha256/AAAA…=", "sha256/BBBB…=") // active + backup
pin("*.example.com", "sha256/CCCC…=")
}
}
Getting a host's SPKI hash:
openssl s_client -connect api.example.com:443 -servername api.example.com < /dev/null 2>/dev/null \
| openssl x509 -pubkey -noout \
| openssl pkey -pubin -outform der \
| openssl dgst -sha256 -binary \
| openssl enc -base64
Always declare a backup pin for the next key. Without one, certificate rotation leaves installed apps unable to connect, with no way to recover short of a new release in the store.
On iOS, the supported key types are RSA-2048, RSA-4096, EC P-256 and EC P-384; any
other type is treated as a pin failure. Multiple pins per host are supported, and
*.host matches exactly one subdomain level. Hosts with no pin declared are
unaffected, and a pin failure fails the request; there is no report-only mode.
responseToModel is the only extension that throws; all other extensions return a sealed state.
try {
val response = client.call {
path = "/might-fail"
}.responseToModel<Data>()
} catch (e: NetFlowException) {
when (e) {
is NetworkException -> { /* handle network issues */ }
is SerializationException -> { /* handle parsing errors */ }
is HttpException -> {
val code = e.code
val errorBody = e.errorBody
}
}
}
single {
netflowClient {
baseUrl = "https://api.example.com"
}
}
This project is licensed under the Apache License, Version 2.0.
Annotated interface call {} DSL
@NetFlowApi / @GET(...) client.call { ... }
│ │
│ KSP generates │
▼ ▼
┌──────────────────────────────────────────┐
│ RequestBuilder │
│ method, path, headers, body, params... │
└────────────────────┬─────────────────────┘
▼
┌──────────────────────────────────────────┐
│ NetFlowRequest │
│ .responseFlow / .responseAsync / ... │
└────────────────────┬─────────────────────┘
▼
┌──────────────────────────────────────────┐
│ InterceptingEngineAdapter │
│ your NetFlowInterceptor chain, runs │
│ again on every retried attempt │
└────────────────────┬─────────────────────┘
▼
┌──────────────────────────────────────────┐
│ TokenHolder │
│ auth { }: attach bearer, single-flight │
│ 401 → refresh → retry (no-op without │
│ auth { } configured) │
└────────────────────┬─────────────────────┘
▼
┌──────────────────────────────────────────┐
│ HttpEngineAdapter │
│ the only piece that ever changes │
└──────┬─────────────────────────┬─────────┘
▼ ▼
NetFlowClientImpl MockNetFlowClient
(real network) (test double)
│ │
▼ ▼
InternalHttpClient handler(request) →
OkHttp (Android) NetFlowMockResponse
NSURLSession (iOS)
NetFlowClientInterceptingEngineAdapterTokenHolderInternalHttpClienthandlerNetFlowMockRequestNetFlowMockResponseassertCalled(...)assertCalledTimes(...)assertNotCalled(...)call {} DSL for everything annotations can't expressnetflow-paging)ApiType) from display type (DisplayType), no trailing .map neededwrappedResponse flag for APIs that return { "data": ... } envelopes401 → refresh → retry, optional proactive refresh from the JWT exp, authState flow, pluggable TokenStoragemultipart/form-data uploads: multipart { } on the DSL, @Multipart / @Part on annotated interfacesFlow<ByteArray> of the body as it arrives, unbuffered (@Streaming on annotated interfaces, responseStream() on the DSL)MockNetFlowClient for testing: no real network calls, with response delays and request historyPagingCollectionViewControllerMockNetFlowClientonNetworkSuccessclient.call { … }
interface TodoApi {
suspend fun getTodos( completed: Boolean?): AsyncState<List<TodoDto>>
fun observeTodo( id: Int): Flow<ResultState<TodoDto>>
suspend fun create( request: CreateTodoRequest): AsyncState<TodoDto>
suspend fun get( id: Int): AsyncState<TodoDto>
fun pagedTodos(): Flow<PagingData<TodoDto>>
fun todosCall(): NetFlowCall // request only, compose the response yourself
suspend fun upload(
id: Int,
caption: String,
file: FilePart,
): AsyncState<Unit>
suspend fun download( url: String): AsyncState<ByteArray>
suspend fun search(
filters: Map<String, Any?>?,
extraHeaders: Map<String, Any?>?,
): AsyncState<List<TodoDto>>
}
val api = client.createTodoApi() // generated extension on NetFlowClient
| Annotation | Target | Notes |
|---|
@GET / @POST / @PUT / @DELETE / @PATCH | function | path defaults to "", for use with @Url |
@Path | parameter | binds to a {name} placeholder in the path |
@Query | parameter | name defaults to the parameter name, override with @Query("user_id"); a null value omits it |
@Header | parameter | same naming and null-omits-it rules as @Query |
@Authorization | parameter (String) | binds to Authorization: Bearer <value>; a null value omits the header |
@AcceptLanguage | parameter (String) | binds to Accept-Language; a null value omits the header |
@Headers("Name: Value", ...) | function | static headers |
@Body | parameter | any @Serializable type or Map<String, Any> |
@Wrapped | function | routes the response through the { "data": ... } envelope; @Wrapped(false) overrides an interface-level @NetFlowApi(wrapped = true) |
@SkipAuth | function | opts out of the client's auth { } (login, sign-up, refresh endpoints) |
@Multipart + @Part | function + parameter | sends a multipart/form-data body; each @Part is a FilePart (file), a primitive (text field), or a @Serializable value (JSON field); a null @Part is omitted; mutually exclusive with @Body |
@Progress | parameter ((Long, Long) -> Unit) | reports upload progress on @Multipart, or download progress on @Streaming |
@Streaming | function, non-suspend, returns Flow<ByteArray> | streams the response body (see Streaming); cannot combine with @Multipart, @Wrapped, or @Paginated |
@Url | parameter (String) | replaces the full request URL; needs an empty method path and no @Path on the same function |
@QueryMap / @HeaderMap | parameter (Map<String, Any?>) | dynamic query parameters / headers on top of any individual @Query/@Header; a null map omits everything, a null value omits that entry |
| Return type | Suspend | Behavior |
|---|
Flow<ResultState<T>> / Flow<ResultState<List<T>>> | no | reactive, non-suspending |
AsyncState<T> / AsyncState<List<T>> | yes | one-shot |
bare T / List<T> | yes | returns the value or throws HttpException |
Flow<PagingData<T>> | no | network-only paging |
NetFlowCall | either | the request without a response strategy; compose it yourself in the repository |
NetFlowCall401 → refresh → retry) and retry { } apply only
before the first byte. Once a chunk has been delivered it is never replayed.chain.proceed(...)'s response unchanged; a short-circuited
2xx emits its body as a single UTF-8 chunk and any other status throws
HttpException.class PostPagingSource(private val database: AppDatabase) : PagingSource<Int, PostEntity>() {
private val query = database.postQueries.selectPosts()
private val listener = object : Query.Listener {
override fun queryResultsChanged() {
invalidate()
query.removeListener(this)
}
}
init {
query.addListener(listener)
}
override fun getRefreshKey(state: PagingState<Int, PostEntity>): Int? {
return state.anchorPosition?.let { anchor ->
state.closestPageToPosition(anchor)?.prevKey?.plus(1)
?: state.closestPageToPosition(anchor)?.nextKey?.minus(1)
}
}
override suspend fun load(params: LoadParams<Int>): LoadResult<Int, PostEntity> {
// your load implementation
}
}
| Property | Default | Description |
|---|
defaultPageSize | 20 | Items loaded per page |
pageQueryName | "page" | URL query parameter name for the page number |
onlyApiCall | false | true for network-only (no local DB) |
wrappedResponse | false | true when API returns { "data": [...] } |
deleteOnRefresh | true | Clear local DB before inserting on REFRESH. Set to false when handling delete inside insertAll |
refresh | false | Force refresh on start, ignoring cache timeout |
cacheTimeout | 1 hour | How long before re-fetching from the network |
import netflowCore // or your KMP framework name
final class PostListViewModel: ObservableObject {
private let viewModel = // your KMP ViewModel from DI
private(set) var posts: [Post] = []
private(set) var isLoading: Bool = false
private let delegate = PagingCollectionViewController<Post>()
init() {
observeData()
observeLoadStates()
observePagingData()
}
func loadNextPage() { delegate.loadNextPage() }
private func observePagingData() {
Task {
for await pagingData in viewModel.posts {
delegate.submitData(pagingData: pagingData)
}
}
}
private func observeData() {
Task {
for await _ in delegate.onPagesUpdatedFlow {
self.posts = delegate.getItems()
self.isLoading = false
}
}
}
private func observeLoadStates() {
Task {
for await loadState in delegate.loadStateFlow {
guard let loadState else { continue }
switch loadState.refresh {
case _ as Paging_commonLoadStateLoading:
self.isLoading = true
default:
self.isLoading = false
}
}
}
}
deinit { delegate.clearScope() }
}
struct PostListView: View {
private var viewModel = PostListViewModel()
var body: some View {
NavigationStack {
Group {
if viewModel.isLoading && viewModel.posts.isEmpty {
ProgressView()
.frame(maxWidth: .infinity, maxHeight: .infinity)
} else {
List {
ForEach(viewModel.posts, id: \.id) { post in
PostItemView(post: post)
.onAppear {
if post.id == viewModel.posts.last?.id {
viewModel.loadNextPage()
}
}
}
}
.listStyle(.plain)
}
}
.navigationTitle("Posts")
}
}
}
| Helper | Code | Description |
|---|
NetFlowMockResponse.success(body) | 200 | Successful response with optional body |
NetFlowMockResponse.error(code, errorBody) | custom | Client error |
NetFlowMockResponse.notFound() | 404 | Not found |
NetFlowMockResponse.serverError(errorBody) | 500 | Server error |
NetFlowMockResponse.stream(chunks) | 200 | Chunks emitted by a streaming call (a plain success(body) streams as one chunk; a Content-Length header feeds onDownloadProgress's total) |
Surfaced from shared tags and platforms — no rankings paid for.