Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
41 changes: 41 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -213,6 +213,47 @@ client.flagTracker.subscribe("game-runner") { e -> /* ... */ }
// unsubscribe with the same listener reference
```

### Observing SDK health

`FBClient.events()` returns a hot `Flow<FBEvent>` carrying lifecycle + health signals so an
application can surface transport / decode failures to its own observability stack (Sentry,
Timber, structured logs) — the SDK itself never captures third-party observability.

```kotlin
// e.g. in Application.onCreate() alongside the FBLifecycleConnector wiring
appScope.launch {
client.events().collect { event ->
when (event) {
FBEvent.Ready -> Timber.d("FeatBit ready.")
FBEvent.Reconnecting -> Timber.d("FeatBit reconnecting.")
is FBEvent.SyncError -> if (event.recoverable) {
Timber.w(event.cause, "FeatBit sync error (recoverable).")
} else {
Timber.e(event.cause, "FeatBit sync failed.")
Sentry.captureException(event.cause)
}
is FBEvent.TransportError -> Timber.w(event.cause, "FeatBit transport error.")
}
}
}
```

`SyncError` carries a `recoverable: Boolean`. It's `true` when the sync layer skipped a
single malformed flag entry within an otherwise-valid batch or when the client is already
initialised and the failure only affects the current batch; `false` when the very first
initialisation failed and evaluations are still returning caller defaults. Subscribe
**before** `start()` to observe the initial `Ready` emission.

### Resilience to malformed server payloads

`StreamingJson` uses `coerceInputValues = true`, and both the streaming and polling paths
apply per-flag `try/catch` isolation around `store.upsert`. If the server ever ships a
malformed variation for one flag entry (e.g. a `null` on a non-null string field, which we
have observed for json-typed flags with an unset variation slot), the SDK logs the failure,
emits `FBEvent.SyncError(recoverable = true)`, drops that entry, and continues processing
the rest of the batch — instead of aborting the whole payload and leaving the client stuck
in a not-ready state.

### Offline mode & bootstrapping

```kotlin
Expand Down
9 changes: 9 additions & 0 deletions featbit-client/src/main/kotlin/co/featbit/client/FBClient.kt
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ import co.featbit.client.changetracker.FlagTracker
import co.featbit.client.evaluation.EvalDetail
import co.featbit.client.model.FBUser
import co.featbit.client.model.FeatureFlag
import kotlinx.coroutines.flow.Flow
import java.io.Closeable
import kotlin.time.Duration
import kotlin.time.Duration.Companion.seconds
Expand All @@ -22,6 +23,14 @@ public interface FBClient : Closeable {
/** The tracker used to subscribe to feature flag changes. */
public val flagTracker: FlagTracker

/**
* Hot flow of SDK lifecycle + health events. Applications subscribe to surface transport /
* decode failures to their own observability stack (Sentry breadcrumb, Timber log). See
* [FBEvent] for the event catalogue. Subscribe before [start] to observe the initial
* [FBEvent.Ready] emission; late subscribers observe subsequent events only.
*/
public fun events(): Flow<FBEvent>

/**
* Starts the client and suspends until it is ready or [timeout] elapses.
*
Expand Down
20 changes: 18 additions & 2 deletions featbit-client/src/main/kotlin/co/featbit/client/FBClientImpl.kt
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,10 @@ import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.SupervisorJob
import kotlinx.coroutines.cancel
import kotlinx.coroutines.channels.BufferOverflow
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.MutableSharedFlow
import kotlinx.coroutines.flow.asSharedFlow
import kotlinx.coroutines.launch
import kotlinx.coroutines.withTimeoutOrNull
import kotlin.time.Duration
Expand Down Expand Up @@ -54,6 +58,16 @@ public class FBClientImpl(
},
)

// Hot event bus for SDK lifecycle + health signals. extraBufferCapacity + DROP_OLDEST so a
// slow subscriber cannot back-pressure the sync layer; a burst of errors during reconnect
// storms is expected and coalescing to the latest is the desired behaviour.
private val eventBus = MutableSharedFlow<FBEvent>(
replay = 0,
extraBufferCapacity = 64,
onBufferOverflow = BufferOverflow.DROP_OLDEST,
)
private val emitEvent: (FBEvent) -> Unit = { eventBus.tryEmit(it) }

@Volatile
private var user: FBUser = initialUser

Expand All @@ -67,12 +81,14 @@ public class FBClientImpl(

override val flagTracker: FlagTracker get() = flagTrackerImpl

override fun events(): Flow<FBEvent> = eventBus.asSharedFlow()

private fun newDataSynchronizer(forUser: FBUser): DataSynchronizer = when {
options.offline -> NullDataSynchronizer()
options.dataSyncMode == DataSyncMode.Streaming ->
StreamingDataSynchronizer(options, forUser, store)
StreamingDataSynchronizer(options, forUser, store, emitEvent)
options.dataSyncMode == DataSyncMode.Polling ->
PollingDataSynchronizer(options, forUser, store)
PollingDataSynchronizer(options, forUser, store, emitEvent)
else -> NullDataSynchronizer()
}

Expand Down
37 changes: 37 additions & 0 deletions featbit-client/src/main/kotlin/co/featbit/client/FBEvent.kt
Original file line number Diff line number Diff line change
@@ -0,0 +1,37 @@
package co.featbit.client

/**
* Lifecycle + health events emitted by [FBClient.events]. Applications wire a collector to
* surface transport / decode failures to their own observability stack (Sentry breadcrumb,
* Timber log, etc.) — the SDK itself never captures third-party observability.
*
* Events are emitted on a hot channel: subscribe before calling [FBClient.start] to avoid
* missing the first [Ready]. Late subscribers observe subsequent events only.
*/
public sealed interface FBEvent {

/** Fired the first time the sync layer completes an initial payload for the current user. */
public data object Ready : FBEvent

/**
* Fired when the streaming websocket is (re)connecting after a drop or the initial handshake.
* Not fired for the very first connection attempt.
*/
public data object Reconnecting : FBEvent

/**
* Fired when a streaming payload fails to decode or a single flag entry within a batch is
* malformed. [recoverable] is `true` when the sync loop continues after the failure — either
* because a single flag was skipped or because a reconnect is expected shortly.
*
* @param cause the throwable that surfaced from the sync layer.
* @param recoverable whether flag evaluations continue against the last-known snapshot.
*/
public data class SyncError(val cause: Throwable, val recoverable: Boolean) : FBEvent

/**
* Fired for transport-level failures (websocket close with a non-normal code, network
* unavailable, HTTP polling failure). The SDK will attempt to reconnect autonomously.
*/
public data class TransportError(val cause: Throwable) : FBEvent
}
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
package co.featbit.client.datasynchronizer

import co.featbit.client.FBEvent
import co.featbit.client.internal.GetUserFlags
import co.featbit.client.model.FBUser
import co.featbit.client.options.FBOptions
Expand All @@ -25,6 +26,7 @@ internal class PollingDataSynchronizer(
options: FBOptions,
user: FBUser,
private val store: MemoryStore,
private val emitEvent: (FBEvent) -> Unit = {},
private val scope: CoroutineScope = CoroutineScope(SupervisorJob() + Dispatchers.IO),
private val getUserFlags: GetUserFlags = GetUserFlags(options, user),
) : DataSynchronizer {
Expand Down Expand Up @@ -68,29 +70,41 @@ internal class PollingDataSynchronizer(
logger.error(
"Polling data synchronizer encountered fatal HTTP error ${response.statusCode}. Stop polling...",
)
emitEvent(FBEvent.TransportError(RuntimeException("polling fatal HTTP ${response.statusCode}")))
startTask.complete(false)
close()
return
}

if (response.isError) {
logger.warn("Polling data synchronizer encountered transient HTTP error ${response.statusCode}.")
emitEvent(FBEvent.TransportError(RuntimeException("polling transient HTTP ${response.statusCode}")))
return
}

timestamp = System.currentTimeMillis()
logger.debug { "Polling received ${response.flags.size} flags." }

response.flags.forEach { store.upsert(it) }
// Per-flag isolation: skip malformed entries rather than aborting the whole batch.
response.flags.forEach { flag ->
try {
store.upsert(flag)
} catch (ex: Exception) {
logger.warn("Skipping malformed flag entry from polling response: ${ex.javaClass.simpleName}: ${ex.message}")
emitEvent(FBEvent.SyncError(ex, recoverable = true))
}
}

if (initializedFlag.compareAndSet(false, true)) {
startTask.complete(true)
logger.info("Polling data synchronizer initialized for user $userKey.")
emitEvent(FBEvent.Ready)
}
} catch (ex: CancellationException) {
throw ex
} catch (ex: Exception) {
logger.error("Exception occurred while polling data.", ex)
emitEvent(FBEvent.SyncError(ex, recoverable = initializedFlag.get()))
}
}

Expand Down
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
package co.featbit.client.datasynchronizer

import co.featbit.client.FBEvent
import co.featbit.client.internal.ConnectionToken
import co.featbit.client.model.EndUser
import co.featbit.client.model.FBUser
Expand Down Expand Up @@ -43,6 +44,7 @@ internal class StreamingDataSynchronizer(
options: FBOptions,
private val user: FBUser,
private val store: MemoryStore,
private val emitEvent: (FBEvent) -> Unit = {},
private val scope: CoroutineScope = CoroutineScope(SupervisorJob() + Dispatchers.IO),
private val httpClient: OkHttpClient = OkHttpClient.Builder()
.pingInterval(20, TimeUnit.SECONDS)
Expand Down Expand Up @@ -105,10 +107,14 @@ internal class StreamingDataSynchronizer(
}

override fun onClosed(webSocket: WebSocket, code: Int, reason: String) {
if (code != NORMAL_CLOSURE) scheduleReconnect("closed: $code $reason")
if (code != NORMAL_CLOSURE) {
emitEvent(FBEvent.TransportError(RuntimeException("websocket closed: $code $reason")))
scheduleReconnect("closed: $code $reason")
}
}

override fun onFailure(webSocket: WebSocket, t: Throwable, response: Response?) {
emitEvent(FBEvent.TransportError(t))
scheduleReconnect(t.message ?: "connection failure")
}
}
Expand Down Expand Up @@ -136,19 +142,44 @@ internal class StreamingDataSynchronizer(
val envelope = StreamingJson.decodeFromString(ServerEnvelope.serializer(), text)
if (envelope.messageType != "data-sync" || envelope.data == null) return

val payload = StreamingJson.decodeFromJsonElement(DataSyncPayload.serializer(), envelope.data)
payload.featureFlags.forEach(store::upsert)
val payload = decodePayloadIsolated(envelope.data) ?: return

// Per-flag isolation: one malformed entry cannot poison the whole batch. Upsert
// survives; the bad entry is dropped after being logged + surfaced via FBEvent.
payload.featureFlags.forEach { flag ->
try {
store.upsert(flag)
} catch (ex: Exception) {
logger.warn("Skipping malformed flag entry from streaming payload: ${ex.javaClass.simpleName}: ${ex.message}")
emitEvent(FBEvent.SyncError(ex, recoverable = true))
}
}
timestamp = System.currentTimeMillis()

if (initializedFlag.compareAndSet(false, true)) {
startTask.complete(true)
logger.info("Streaming data synchronizer initialized for user ${user.key}.")
emitEvent(FBEvent.Ready)
}
} catch (ex: Exception) {
logger.error("Failed to handle streaming message.", ex)
emitEvent(FBEvent.SyncError(ex, recoverable = initializedFlag.get()))
}
}

/**
* Isolate payload decoding so a single bad flag entry does not abort the entire batch.
* `coerceInputValues = true` on [StreamingJson] already coerces null-on-non-null fields
* to their defaults; this catch is the belt-and-braces for anything the coercion misses.
*/
private fun decodePayloadIsolated(data: JsonElement): DataSyncPayload? = try {
StreamingJson.decodeFromJsonElement(DataSyncPayload.serializer(), data)
} catch (ex: Exception) {
logger.error("Failed to decode data-sync payload.", ex)
emitEvent(FBEvent.SyncError(ex, recoverable = initializedFlag.get()))
null
}

override fun pause() {
if (closed || paused) return
paused = true
Expand All @@ -175,6 +206,7 @@ internal class StreamingDataSynchronizer(
val attempt = ++reconnectAttempts
val backoff = min(MAX_BACKOFF_MS, BASE_BACKOFF_MS shl min(attempt, 6)) + Random.nextLong(250)
logger.warn("Streaming disconnected ($reason); reconnecting in ${backoff}ms (attempt $attempt).")
emitEvent(FBEvent.Reconnecting)
scope.launch {
delay(backoff)
reconnecting.set(false)
Expand Down Expand Up @@ -214,7 +246,15 @@ internal class StreamingDataSynchronizer(
const val MAX_BACKOFF_MS = 30_000L
const val PING_MESSAGE = """{"messageType":"ping","data":{}}"""

val StreamingJson = Json { ignoreUnknownKeys = true; encodeDefaults = true }
// coerceInputValues: server occasionally ships `"variation": null` for json-typed
// flags with an unset variation slot; without coercion the whole streaming payload
// fails to decode and FBClient never becomes ready. Coerce null → property default
// (empty string) for the defaulted non-null string fields on FeatureFlag.
val StreamingJson = Json {
ignoreUnknownKeys = true
encodeDefaults = true
coerceInputValues = true
}

/** Accepts `ws(s)://` (or `http(s)://`) and returns the `/streaming` HTTP(S) URL OkHttp uses. */
fun String.toStreamingHttpUrl() =
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -78,6 +78,7 @@ internal abstract class FbApiClient(
ignoreUnknownKeys = true
encodeDefaults = true
explicitNulls = false
coerceInputValues = true
}
}
}
2 changes: 1 addition & 1 deletion gradle.properties
Original file line number Diff line number Diff line change
Expand Up @@ -11,4 +11,4 @@ kotlin.code.style=official
# block in each library module; `artifactId` is set per-module because the group ships
# two artifacts (`featbit-client`, `featbit-client-android`).
GROUP=co.featbit
VERSION_NAME=0.1.0-SNAPSHOT
VERSION_NAME=0.1.1-SNAPSHOT
Loading