From 853515e519a18c036dc28110831187fa9c30ccca Mon Sep 17 00:00:00 2001 From: Deep Shah Date: Tue, 11 Aug 2026 22:51:59 +0530 Subject: [PATCH] fix: coerce null variations + FBEvent Flow + per-flag isolation on sync Three related resilience fixes surfaced by real-world testing against the els.fluentinhealth.io FeatBit server. 1. coerceInputValues = true on both Json configs (StreamingJson + FbApiClient.json). Server occasionally ships `"variation": null` in streaming payloads for json-typed flags with an unset variation slot; without coercion the whole payload fails to decode with `Expected string value for a non-null key 'variation', got null literal`, startTask never completes, and FBClient stays not-ready for its whole lifetime. With coercion, null coerces to the property default ("") and streaming proceeds. 2. Per-flag try/catch around store.upsert in both StreamingDataSynchronizer and PollingDataSynchronizer.safePoll. If a single flag entry ever slips past coercion (or throws for any other reason), we log-and-skip that entry instead of aborting the whole batch. FBClient still becomes ready with the surviving flags. 3. New FBClient.events(): Flow API. Applications subscribe to surface transport / decode failures (Ready, Reconnecting, SyncError, TransportError) into their own observability stack (Sentry, Timber). SyncError.recoverable disambiguates one-bad-flag-in-batch from a full sync outage. README updated with the collector pattern. Version bumped to 0.1.1-SNAPSHOT so consumers requesting this branch via JitPack get a distinct artifact. Co-Authored-By: Claude Opus 4.7 (1M context) --- README.md | 41 ++++++++++++++++ .../main/kotlin/co/featbit/client/FBClient.kt | 9 ++++ .../kotlin/co/featbit/client/FBClientImpl.kt | 20 +++++++- .../main/kotlin/co/featbit/client/FBEvent.kt | 37 ++++++++++++++ .../PollingDataSynchronizer.kt | 16 ++++++- .../StreamingDataSynchronizer.kt | 48 +++++++++++++++++-- .../co/featbit/client/internal/FbApiClient.kt | 1 + gradle.properties | 2 +- 8 files changed, 166 insertions(+), 8 deletions(-) create mode 100644 featbit-client/src/main/kotlin/co/featbit/client/FBEvent.kt diff --git a/README.md b/README.md index b4a3445..f1b5698 100644 --- a/README.md +++ b/README.md @@ -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` 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 diff --git a/featbit-client/src/main/kotlin/co/featbit/client/FBClient.kt b/featbit-client/src/main/kotlin/co/featbit/client/FBClient.kt index b2b7173..0639cea 100644 --- a/featbit-client/src/main/kotlin/co/featbit/client/FBClient.kt +++ b/featbit-client/src/main/kotlin/co/featbit/client/FBClient.kt @@ -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 @@ -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 + /** * Starts the client and suspends until it is ready or [timeout] elapses. * diff --git a/featbit-client/src/main/kotlin/co/featbit/client/FBClientImpl.kt b/featbit-client/src/main/kotlin/co/featbit/client/FBClientImpl.kt index 4c49711..c3cd720 100644 --- a/featbit-client/src/main/kotlin/co/featbit/client/FBClientImpl.kt +++ b/featbit-client/src/main/kotlin/co/featbit/client/FBClientImpl.kt @@ -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 @@ -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( + replay = 0, + extraBufferCapacity = 64, + onBufferOverflow = BufferOverflow.DROP_OLDEST, + ) + private val emitEvent: (FBEvent) -> Unit = { eventBus.tryEmit(it) } + @Volatile private var user: FBUser = initialUser @@ -67,12 +81,14 @@ public class FBClientImpl( override val flagTracker: FlagTracker get() = flagTrackerImpl + override fun events(): Flow = 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() } diff --git a/featbit-client/src/main/kotlin/co/featbit/client/FBEvent.kt b/featbit-client/src/main/kotlin/co/featbit/client/FBEvent.kt new file mode 100644 index 0000000..db89bb5 --- /dev/null +++ b/featbit-client/src/main/kotlin/co/featbit/client/FBEvent.kt @@ -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 +} diff --git a/featbit-client/src/main/kotlin/co/featbit/client/datasynchronizer/PollingDataSynchronizer.kt b/featbit-client/src/main/kotlin/co/featbit/client/datasynchronizer/PollingDataSynchronizer.kt index 0367ddd..5e73bc7 100644 --- a/featbit-client/src/main/kotlin/co/featbit/client/datasynchronizer/PollingDataSynchronizer.kt +++ b/featbit-client/src/main/kotlin/co/featbit/client/datasynchronizer/PollingDataSynchronizer.kt @@ -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 @@ -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 { @@ -68,6 +70,7 @@ 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 @@ -75,22 +78,33 @@ internal class PollingDataSynchronizer( 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())) } } diff --git a/featbit-client/src/main/kotlin/co/featbit/client/datasynchronizer/StreamingDataSynchronizer.kt b/featbit-client/src/main/kotlin/co/featbit/client/datasynchronizer/StreamingDataSynchronizer.kt index f1d95d2..743ca4f 100644 --- a/featbit-client/src/main/kotlin/co/featbit/client/datasynchronizer/StreamingDataSynchronizer.kt +++ b/featbit-client/src/main/kotlin/co/featbit/client/datasynchronizer/StreamingDataSynchronizer.kt @@ -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 @@ -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) @@ -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") } } @@ -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 @@ -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) @@ -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() = diff --git a/featbit-client/src/main/kotlin/co/featbit/client/internal/FbApiClient.kt b/featbit-client/src/main/kotlin/co/featbit/client/internal/FbApiClient.kt index 457eb89..39a89b2 100644 --- a/featbit-client/src/main/kotlin/co/featbit/client/internal/FbApiClient.kt +++ b/featbit-client/src/main/kotlin/co/featbit/client/internal/FbApiClient.kt @@ -78,6 +78,7 @@ internal abstract class FbApiClient( ignoreUnknownKeys = true encodeDefaults = true explicitNulls = false + coerceInputValues = true } } } diff --git a/gradle.properties b/gradle.properties index 4e588e9..5cf2e11 100644 --- a/gradle.properties +++ b/gradle.properties @@ -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