Skip to content

Commit 2fb5683

Browse files
authored
Merge pull request #1630 from tunjid/tj/network-service-monitored-items
Extract utility method for pulling items from the network and observi…
2 parents 4310486 + 0f97363 commit 2fb5683

3 files changed

Lines changed: 105 additions & 62 deletions

File tree

data/core/src/commonMain/kotlin/com/tunjid/heron/data/network/NetworkService.kt

Lines changed: 51 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,9 @@
1616

1717
package com.tunjid.heron.data.network
1818

19+
import com.tunjid.heron.data.core.models.Cursor
20+
import com.tunjid.heron.data.core.models.CursorList
21+
import com.tunjid.heron.data.core.models.canRequestData
1922
import com.tunjid.heron.data.core.types.AtProtoException
2023
import com.tunjid.heron.data.lexicons.BlueskyApi
2124
import com.tunjid.heron.data.lexicons.XrpcBlueskyApi
@@ -26,6 +29,10 @@ import io.ktor.client.HttpClient
2629
import kotlin.time.Duration
2730
import kotlin.time.Duration.Companion.milliseconds
2831
import kotlin.time.Duration.Companion.seconds
32+
import kotlinx.coroutines.flow.Flow
33+
import kotlinx.coroutines.flow.emitAll
34+
import kotlinx.coroutines.flow.emptyFlow
35+
import kotlinx.coroutines.flow.flow
2936
import sh.christian.ozone.api.response.AtpResponse
3037

3138
internal interface NetworkService {
@@ -44,6 +51,50 @@ internal interface NetworkService {
4451
): Result<T>
4552
}
4653

54+
internal inline fun <NetworkResponse : Any, Item> NetworkService.observedItems(
55+
cursor: Cursor,
56+
crossinline responseFetcher: suspend BlueskyApi.() -> AtpResponse<NetworkResponse>,
57+
crossinline responseSaver: suspend (NetworkResponse) -> Unit,
58+
crossinline responseCursor: (NetworkResponse) -> Cursor.Next?,
59+
crossinline networkItems: (NetworkResponse, Cursor) -> List<Item>?,
60+
crossinline observedItems: (NetworkResponse, Cursor) -> Flow<CursorList<Item>>,
61+
): Flow<CursorList<Item>> =
62+
// Final or pending cursor, nothing to fetch
63+
if (!cursor.canRequestData) emptyFlow()
64+
else flow<CursorList<Item>> {
65+
val response = runCatchingWithMonitoredNetworkRetry(
66+
block = {
67+
responseFetcher()
68+
},
69+
).getOrNull()
70+
?: return@flow
71+
72+
val nextCursor = responseCursor(response)
73+
?: Cursor.Final
74+
75+
// Emit network results immediately for minimal latency
76+
networkItems(
77+
response,
78+
nextCursor,
79+
)?.let {
80+
emit(
81+
CursorList(
82+
items = it,
83+
nextCursor = nextCursor,
84+
),
85+
)
86+
}
87+
88+
responseSaver(response)
89+
90+
emitAll(
91+
observedItems(
92+
response,
93+
nextCursor,
94+
),
95+
)
96+
}
97+
4798
@Inject
4899
internal class KtorNetworkService(
49100
httpClient: HttpClient,

data/core/src/commonMain/kotlin/com/tunjid/heron/data/repository/records/BlueskyRecordOperations.kt

Lines changed: 25 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -50,7 +50,6 @@ import com.tunjid.heron.data.core.models.ListMember
5050
import com.tunjid.heron.data.core.models.ProfileWithViewerState
5151
import com.tunjid.heron.data.core.models.StarterPack
5252
import com.tunjid.heron.data.core.models.Trend
53-
import com.tunjid.heron.data.core.models.canRequestData
5453
import com.tunjid.heron.data.core.models.offset
5554
import com.tunjid.heron.data.core.models.value
5655
import com.tunjid.heron.data.core.types.FeedGeneratorUri
@@ -81,6 +80,7 @@ import com.tunjid.heron.data.network.FeedCreationService
8180
import com.tunjid.heron.data.network.GrazeResponse
8281
import com.tunjid.heron.data.network.NetworkService
8382
import com.tunjid.heron.data.network.models.profile
83+
import com.tunjid.heron.data.network.observedItems
8484
import com.tunjid.heron.data.repository.ListMemberQuery
8585
import com.tunjid.heron.data.repository.ProfilesQuery
8686
import com.tunjid.heron.data.repository.SavedStateDataSource
@@ -113,7 +113,6 @@ import kotlinx.coroutines.flow.Flow
113113
import kotlinx.coroutines.flow.SharingStarted
114114
import kotlinx.coroutines.flow.combine
115115
import kotlinx.coroutines.flow.distinctUntilChanged
116-
import kotlinx.coroutines.flow.emitAll
117116
import kotlinx.coroutines.flow.emptyFlow
118117
import kotlinx.coroutines.flow.filterNotNull
119118
import kotlinx.coroutines.flow.flow
@@ -469,31 +468,35 @@ internal class OfflineFirstBlueskyRecordOperations(
469468
cursor: Cursor,
470469
): Flow<CursorList<FeedGenerator>> =
471470
if (query.query.isBlank()) emptyFlow()
472-
else if (!cursor.canRequestData) emptyFlow()
473-
else flow {
474-
val response = networkService.runCatchingWithMonitoredNetworkRetry {
471+
else networkService.observedItems(
472+
cursor = cursor,
473+
responseFetcher = {
475474
getPopularFeedGeneratorsUnspecced(
476475
params = GetPopularFeedGeneratorsQueryParams(
477476
query = query.query,
478477
limit = query.data.limit,
479478
cursor = cursor.value,
480479
),
481480
)
482-
}
483-
.getOrNull()
484-
?: return@flow
485-
486-
multipleEntitySaverProvider.saveInTransaction {
487-
response.feeds
488-
.forEach { generatorView ->
489-
add(feedGeneratorView = generatorView)
490-
}
491-
}
492-
493-
val nextCursor = response.cursor?.let(Cursor::Next) ?: Cursor.Final
494-
val feedUris = response.feeds.map { it.uri.atUri.let(::FeedGeneratorUri) }
481+
},
482+
responseSaver = { response ->
483+
multipleEntitySaverProvider.saveInTransaction {
484+
response.feeds
485+
.forEach { generatorView ->
486+
add(feedGeneratorView = generatorView)
487+
}
488+
}
489+
},
490+
responseCursor = { response ->
491+
response.cursor?.let(Cursor::Next)
492+
},
493+
networkItems = { _, _ ->
494+
null
495+
},
496+
observedItems = { response, nextCursor ->
497+
val feedUris = response.feeds
498+
.map { it.uri.atUri.let(::FeedGeneratorUri) }
495499

496-
emitAll(
497500
feedGeneratorDao.feedGenerators(
498501
feedUris = feedUris,
499502
)
@@ -508,9 +511,9 @@ internal class OfflineFirstBlueskyRecordOperations(
508511
),
509512
nextCursor = nextCursor,
510513
)
511-
},
512-
)
513-
}
514+
}
515+
},
516+
)
514517
.flowOn(ioDispatcher)
515518

516519
override fun suggestedFeeds(): Flow<List<FeedGenerator>> =

data/core/src/commonMain/kotlin/com/tunjid/heron/data/utilities/profileLookup/ProfileLookup.kt

Lines changed: 29 additions & 40 deletions
Original file line numberDiff line numberDiff line change
@@ -25,11 +25,11 @@ import com.tunjid.heron.data.core.models.Link
2525
import com.tunjid.heron.data.core.models.LinkTarget
2626
import com.tunjid.heron.data.core.models.Profile
2727
import com.tunjid.heron.data.core.models.ProfileWithViewerState
28-
import com.tunjid.heron.data.core.models.canRequestData
2928
import com.tunjid.heron.data.core.types.Id
3029
import com.tunjid.heron.data.core.types.ProfileHandleOrId
3130
import com.tunjid.heron.data.core.types.ProfileId
3231
import com.tunjid.heron.data.core.types.RecordUri
32+
import com.tunjid.heron.data.core.types.UnknownRecordUri
3333
import com.tunjid.heron.data.core.types.UnresolvableProfileException
3434
import com.tunjid.heron.data.core.types.profileId
3535
import com.tunjid.heron.data.core.types.recordKey
@@ -44,6 +44,7 @@ import com.tunjid.heron.data.lexicons.BlueskyApi
4444
import com.tunjid.heron.data.network.NetworkService
4545
import com.tunjid.heron.data.network.models.profile
4646
import com.tunjid.heron.data.network.models.profileViewerStateEntity
47+
import com.tunjid.heron.data.network.observedItems
4748
import com.tunjid.heron.data.utilities.distinctUntilChangedMap
4849
import com.tunjid.heron.data.utilities.multipleEntitysaver.MultipleEntitySaverProvider
4950
import com.tunjid.heron.data.utilities.multipleEntitysaver.add
@@ -53,9 +54,7 @@ import kotlinx.coroutines.async
5354
import kotlinx.coroutines.awaitAll
5455
import kotlinx.coroutines.coroutineScope
5556
import kotlinx.coroutines.flow.Flow
56-
import kotlinx.coroutines.flow.emitAll
5757
import kotlinx.coroutines.flow.first
58-
import kotlinx.coroutines.flow.flow
5958
import kotlinx.coroutines.flow.flowOf
6059
import kotlinx.coroutines.flow.map
6160
import sh.christian.ozone.api.Did
@@ -111,45 +110,35 @@ internal class OfflineProfileLookup(
111110
responseFetcher: suspend BlueskyApi.() -> AtpResponse<NetworkResponse>,
112111
responseProfileViews: NetworkResponse.() -> List<ProfileView>,
113112
responseCursor: NetworkResponse.() -> String?,
114-
): Flow<CursorList<ProfileWithViewerState>> = flow {
115-
// Final or pending cursor, nothing to fetch
116-
if (!cursor.canRequestData) return@flow
117-
118-
val response = networkService.runCatchingWithMonitoredNetworkRetry(
119-
block = responseFetcher,
120-
).getOrNull()
121-
?: return@flow
122-
123-
val profileViews = response.responseProfileViews()
124-
125-
val nextCursor = response.responseCursor()
126-
?.let(Cursor::Next)
127-
?: Cursor.Final
128-
129-
// Emit network results immediately for minimal latency
130-
emit(
113+
): Flow<CursorList<ProfileWithViewerState>> = networkService.observedItems(
114+
cursor = cursor,
115+
responseFetcher = responseFetcher,
116+
responseSaver = { response ->
117+
multipleEntitySaverProvider.saveInTransaction {
118+
response.responseProfileViews()
119+
.forEach { profileView ->
120+
add(
121+
viewingProfileId = signedInProfileId,
122+
profileView = profileView,
123+
)
124+
}
125+
}
126+
},
127+
responseCursor = { response ->
128+
response.responseCursor()?.let(Cursor::Next)
129+
},
130+
networkItems = { response, nextCursor ->
131131
CursorList(
132-
items = profileViews.toProfileWithViewerStates(
132+
items = response.responseProfileViews().toProfileWithViewerStates(
133133
signedInProfileId = signedInProfileId,
134134
profileMapper = ProfileView::profile,
135135
profileViewerStateMapper = ProfileView::profileViewerStateEntity,
136136
),
137137
nextCursor = nextCursor,
138-
),
139-
)
140-
141-
multipleEntitySaverProvider.saveInTransaction {
142-
profileViews
143-
.forEach { profileView ->
144-
add(
145-
viewingProfileId = signedInProfileId,
146-
profileView = profileView,
147-
)
148-
}
149-
}
150-
151-
emitAll(
152-
profileViews.observeProfileWithViewerStates(
138+
)
139+
},
140+
observedItems = { response, nextCursor ->
141+
response.responseProfileViews().observeProfileWithViewerStates(
153142
profileDao = profileDao,
154143
signedInProfileId = signedInProfileId,
155144
profileMapper = ProfileView::profile,
@@ -160,9 +149,9 @@ internal class OfflineProfileLookup(
160149
items = profileWithViewerStates,
161150
nextCursor = nextCursor,
162151
)
163-
},
164-
)
165-
}
152+
}
153+
},
154+
)
166155

167156
override suspend fun lookupProfileDid(
168157
profileId: Id.Profile,
@@ -198,7 +187,7 @@ internal class OfflineProfileLookup(
198187
override suspend fun <T : RecordUri> withDidAuthority(
199188
uri: T,
200189
): T {
201-
if (uri is com.tunjid.heron.data.core.types.UnknownRecordUri) return uri
190+
if (uri is UnknownRecordUri) return uri
202191
val profileId = uri.profileId()
203192
if (Did.Regex.matches(profileId.id)) return uri
204193
val did = lookupProfileDid(profileId) ?: return uri

0 commit comments

Comments
 (0)