diff --git a/.gitignore b/.gitignore index 00d29414..b632a0b9 100644 --- a/.gitignore +++ b/.gitignore @@ -44,6 +44,10 @@ captures/ .externalNativeBuild/ .cxx/ +# Local dev convenience script (machine-specific paths) +Makefile +CLAUDE.md + # Temp files temp_opencode_analysis/ # Images (except screenshots/ and app resources) diff --git a/app/src/main/java/dev/blazelight/p4oc/core/network/OpenCodeEventSource.kt b/app/src/main/java/dev/blazelight/p4oc/core/network/OpenCodeEventSource.kt index 771af017..4714e761 100644 --- a/app/src/main/java/dev/blazelight/p4oc/core/network/OpenCodeEventSource.kt +++ b/app/src/main/java/dev/blazelight/p4oc/core/network/OpenCodeEventSource.kt @@ -17,12 +17,14 @@ import kotlinx.coroutines.Dispatchers import kotlinx.coroutines.SupervisorJob import kotlinx.coroutines.cancel import kotlinx.coroutines.channels.Channel +import kotlinx.coroutines.delay import kotlinx.coroutines.flow.MutableSharedFlow import kotlinx.coroutines.flow.MutableStateFlow import kotlinx.coroutines.flow.SharedFlow import kotlinx.coroutines.flow.StateFlow import kotlinx.coroutines.flow.asSharedFlow import kotlinx.coroutines.flow.asStateFlow +import kotlinx.coroutines.isActive import kotlinx.coroutines.launch import kotlinx.serialization.json.Json import okhttp3.OkHttpClient @@ -44,6 +46,12 @@ class OpenCodeEventSource( companion object { private const val TAG = "OpenCodeEventSource" private const val MAX_CONSECUTIVE_ERRORS = 15 + + // Liveness watchdog: the server emits frequent heartbeat comments, so a Connected socket + // that receives no frame (message OR heartbeat) for this long is presumed dead (silent + // zombie socket — readTimeout is 0). Conservative vs. the heartbeat cadence. (Issue #14.) + private const val SSE_STALE_MS = 60_000L + private const val WATCHDOG_INTERVAL_MS = 20_000L } private val eventPumpScope = CoroutineScope(SupervisorJob() + Dispatchers.IO) @@ -73,6 +81,15 @@ class OpenCodeEventSource( private val consecutiveErrors = AtomicInteger(0) + // Timestamp of the last received SSE frame (message or heartbeat comment). 0 = none yet. + // Read by the watchdog off the event thread, so @Volatile. + @Volatile + private var lastFrameAtMs: Long = 0L + + private fun markFrameReceived() { + lastFrameAtMs = System.currentTimeMillis() + } + init { eventPumpScope.launch { for (queued in eventChannel) { @@ -83,8 +100,37 @@ class OpenCodeEventSource( } } } + startLivenessWatchdog() } + /** + * Force a reconnect when the socket claims Connected but has gone silent past [SSE_STALE_MS]. + * With readTimeout(0) the library never surfaces a silently-dead socket, so nothing else would + * detect it — the app would show "connected" while receiving zero events (issue #14). + */ + private fun startLivenessWatchdog() { + eventPumpScope.launch { + while (isActive) { + delay(WATCHDOG_INTERVAL_MS) + if (isStale(System.currentTimeMillis(), lastFrameAtMs, _connectionState.value)) { + val idle = System.currentTimeMillis() - lastFrameAtMs + AppLog.w(TAG, "SSE idle ${idle}ms > ${SSE_STALE_MS}ms while Connected – forcing reconnect") + reconnect() + } + } + } + } + + /** + * True when the socket claims [ConnectionState.Connected] but has received no frame + * (message or heartbeat) within [SSE_STALE_MS]. `lastFrameAtMs == 0L` means no frame yet + * (connect in progress) — never stale. Extracted for deterministic unit testing. + */ + internal fun isStale(nowMs: Long, lastFrameAtMs: Long, state: ConnectionState): Boolean = + lastFrameAtMs != 0L && + state is ConnectionState.Connected && + (nowMs - lastFrameAtMs) > SSE_STALE_MS + fun connect() { val besRef: BackgroundEventSource synchronized(lock) { @@ -100,6 +146,7 @@ class OpenCodeEventSource( AppLog.d(TAG, "connect() – creating BackgroundEventSource") _connectionState.value = ConnectionState.Connecting consecutiveErrors.set(0) + markFrameReceived() val gen = ++generation besRef = createBackgroundEventSource(gen) @@ -132,6 +179,7 @@ class OpenCodeEventSource( toClose = detachLocked() _connectionState.value = ConnectionState.Connecting consecutiveErrors.set(0) + markFrameReceived() val gen = ++generation besRef = createBackgroundEventSource(gen) @@ -155,6 +203,9 @@ class OpenCodeEventSource( _connectionState.value = ConnectionState.Disconnected } toClose?.closeSafely() + // Stop the event pump + liveness watchdog so they don't outlive this instance. + eventPumpScope.cancel("OpenCodeEventSource shut down") + eventChannel.close() } /** @@ -279,19 +330,23 @@ class OpenCodeEventSource( AppLog.d(TAG, "SSE connected (onOpen)") errorFiredSinceOpen = false consecutiveErrors.set(0) + markFrameReceived() _connectionState.value = ConnectionState.Connected enqueueEvent(OpenCodeEvent.Connected, gen) } override fun onMessage(event: String, messageEvent: MessageEvent) { if (!isActiveGeneration(gen)) return + markFrameReceived() val data = messageEvent.data AppLog.v(TAG, "SSE message received: event=$event, length=${data.length}") parseAndEmitEvent(data, gen) } override fun onComment(comment: String) { - // Keepalive — nothing to do + // Keepalive heartbeat — no event to emit, but it proves the socket is alive. + if (!isActiveGeneration(gen)) return + markFrameReceived() } override fun onClosed() { diff --git a/app/src/main/java/dev/blazelight/p4oc/data/session/SessionRepositoryImpl.kt b/app/src/main/java/dev/blazelight/p4oc/data/session/SessionRepositoryImpl.kt index ce072a1c..671229f9 100644 --- a/app/src/main/java/dev/blazelight/p4oc/data/session/SessionRepositoryImpl.kt +++ b/app/src/main/java/dev/blazelight/p4oc/data/session/SessionRepositoryImpl.kt @@ -191,6 +191,10 @@ class SessionRepositoryImpl( override fun acceptEvent(event: OpenCodeEvent) { if (event is OpenCodeEvent.Connected) { hydrateAfterReconnect() + // Message refresh must run on EVERY reconnect, independent of the snapshot-hydrate + // inFlight guard above (which only fires once). Messages are a separate concern from + // the session-list snapshot, so refetch them directly here (issue #14). + scope.launch { refreshOpenSessionMessages() } scope.launch { reconcileObservedPendingPermissions() } return } @@ -345,6 +349,21 @@ class SessionRepositoryImpl( } } + /** + * Re-fetch messages for every actively-observed session after an SSE reconnect. + * Called directly from the [OpenCodeEvent.Connected] path (NOT gated by the snapshot-hydrate + * inFlight guard, which only permits one hydration) so an open conversation self-heals on every + * reconnect (issue #14: had to leave and re-enter to see updates). + * Reuses the same REST refetch/merge path [loadMessages] that ChatViewModel uses on entry. + */ + private suspend fun refreshOpenSessionMessages() { + val openSessions = synchronized(messageStates) { messageStates.keys.toList() } + for (id in openSessions) { + runCatching { loadMessages(SessionId(id), limit = null) } + .onFailure { AppLog.w(TAG, "Post-reconnect message refresh failed for $id: ${it.message}") } + } + } + override fun messages(sessionId: SessionId): StateFlow> = messageState( sessionId.value ).asStateFlow() diff --git a/app/src/main/java/dev/blazelight/p4oc/data/session/SessionRepositoryProvider.kt b/app/src/main/java/dev/blazelight/p4oc/data/session/SessionRepositoryProvider.kt index 7f5065f7..e4ede6d6 100644 --- a/app/src/main/java/dev/blazelight/p4oc/data/session/SessionRepositoryProvider.kt +++ b/app/src/main/java/dev/blazelight/p4oc/data/session/SessionRepositoryProvider.kt @@ -1,12 +1,16 @@ package dev.blazelight.p4oc.data.session +import dev.blazelight.p4oc.core.log.AppLog import dev.blazelight.p4oc.core.network.ConnectionManager +import dev.blazelight.p4oc.core.network.ConnectionState import dev.blazelight.p4oc.data.remote.mapper.MessageMapper import dev.blazelight.p4oc.data.server.ActiveServerApiProvider import dev.blazelight.p4oc.data.workspace.WorkspaceClient +import dev.blazelight.p4oc.domain.model.OpenCodeEvent import dev.blazelight.p4oc.domain.server.ServerGeneration import dev.blazelight.p4oc.domain.server.WorkspaceKey import dev.blazelight.p4oc.domain.workspace.Workspace +import kotlinx.coroutines.CancellationException import kotlinx.coroutines.CoroutineDispatcher import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.Dispatchers @@ -41,6 +45,46 @@ class SessionRepositoryProvider( private val scope = CoroutineScope(SupervisorJob() + dispatcher) private val entries = mutableMapOf() + private companion object { + const val TAG = "SessionRepositoryProvider" + } + + init { + // Deliver SSE reconnects to every live workspace repository so open conversations + // self-heal without navigation (issue #14). The synthetic OpenCodeEvent.Connected never + // reaches the per-workspace fan-out (it carries no directory), so we broadcast it here on + // each non-Connected → Connected transition of the shared connection. + scope.launch { + var wasConnected = false + connectionManager.connectionState.collect { state -> + val nowConnected = state is ConnectionState.Connected + if (nowConnected && !wasConnected) { + broadcastReconnect() + } + wasConnected = nowConnected + } + } + } + + @Suppress("TooGenericExceptionCaught") // resilience guard: one bad repo must not stop the broadcast + private fun broadcastReconnect() { + // Generation is a single global counter in ConnectionManager, so one generation == one + // connection == one server: filtering by generation alone correctly targets exactly the + // repositories on the reconnected connection, without depending on server-URL normalization. + val generation = connectionManager.currentGeneration ?: return + val repositories = synchronized(this) { + entries.filterKeys { it.generation == generation.value } + .map { it.value.repository } + } + repositories.forEach { repository -> + try { + repository.acceptEvent(OpenCodeEvent.Connected) + } catch (e: Exception) { + AppLog.e(TAG, "Reconnect broadcast failed for a workspace: ${e.message}", e) + } + } + } + fun acquire(workspace: Workspace, generation: ServerGeneration): Lease = synchronized(this) { val key = workspace.toProviderKey(generation) val entry = entries.getOrPut(key) { @@ -70,6 +114,7 @@ class SessionRepositoryProvider( repositoryToClose.repository.close() } + @Suppress("TooGenericExceptionCaught") // resilience guard: one bad event must not kill the collector private fun collectWorkspaceEvents( workspace: Workspace, generation: ServerGeneration, @@ -80,7 +125,19 @@ class SessionRepositoryProvider( scopedEvent.generation == generation && scopedEvent.workspaceKey == workspace.key ) { - repository.acceptEvent(scopedEvent.event) + try { + repository.acceptEvent(scopedEvent.event) + } catch (ce: CancellationException) { + throw ce + } catch (e: Exception) { + // A single malformed/unexpected event must never permanently kill this + // workspace's live event delivery (issue #14: chat froze until re-entry). + AppLog.e( + TAG, + "Dropping event ${scopedEvent.event::class.simpleName} for ${workspace.key}: ${e.message}", + e, + ) + } } } } diff --git a/app/src/test/java/dev/blazelight/p4oc/core/network/OpenCodeEventSourceTest.kt b/app/src/test/java/dev/blazelight/p4oc/core/network/OpenCodeEventSourceTest.kt index 894884c5..91c16fa9 100644 --- a/app/src/test/java/dev/blazelight/p4oc/core/network/OpenCodeEventSourceTest.kt +++ b/app/src/test/java/dev/blazelight/p4oc/core/network/OpenCodeEventSourceTest.kt @@ -43,6 +43,35 @@ class OpenCodeEventSourceTest { unmockkObject(AppLog) } + @Test + fun `isStale flags a silent Connected socket past the threshold`() { + val source = newSource() + try { + val now = 1_000_000L + val staleFrame = now - 61_000L // > SSE_STALE_MS (60s) + val freshFrame = now - 30_000L + + // Connected + no frame for > 60s => stale (the zombie-socket case, issue #14). + assertEquals(true, source.isStale(now, staleFrame, ConnectionState.Connected)) + // Connected but a recent frame => not stale. + assertEquals(false, source.isStale(now, freshFrame, ConnectionState.Connected)) + // No frame yet (0L) => never stale, even if "Connected". + assertEquals(false, source.isStale(now, 0L, ConnectionState.Connected)) + // Not Connected => watchdog leaves reconnect to the normal error/escalation path. + assertEquals(false, source.isStale(now, staleFrame, ConnectionState.Disconnected)) + assertEquals(false, source.isStale(now, staleFrame, ConnectionState.Connecting)) + } finally { + source.shutdown() + } + } + + private fun newSource(): OpenCodeEventSource = OpenCodeEventSource( + okHttpClient = OkHttpClient(), + json = json, + baseUrl = "http://127.0.0.1:1", + eventMapper = EventMapper(json, MessageMapper(json)), + ) + @Test fun `slow collector receives more than previous delta buffer capacity without loss`() = runTest { val source = OpenCodeEventSource( diff --git a/app/src/test/java/dev/blazelight/p4oc/data/session/SessionRepositoryProviderTest.kt b/app/src/test/java/dev/blazelight/p4oc/data/session/SessionRepositoryProviderTest.kt index 20589515..9e6da564 100644 --- a/app/src/test/java/dev/blazelight/p4oc/data/session/SessionRepositoryProviderTest.kt +++ b/app/src/test/java/dev/blazelight/p4oc/data/session/SessionRepositoryProviderTest.kt @@ -1,24 +1,32 @@ package dev.blazelight.p4oc.data.session import dev.blazelight.p4oc.core.network.ConnectionManager +import dev.blazelight.p4oc.core.network.ConnectionState import dev.blazelight.p4oc.core.network.OpenCodeApi import dev.blazelight.p4oc.data.remote.mapper.MessageMapper import dev.blazelight.p4oc.data.server.ActiveServerApiProvider +import dev.blazelight.p4oc.domain.model.Message import dev.blazelight.p4oc.domain.model.OpenCodeEvent import dev.blazelight.p4oc.domain.model.Session +import dev.blazelight.p4oc.domain.model.TokenUsage import dev.blazelight.p4oc.domain.server.ScopedEvent import dev.blazelight.p4oc.domain.server.ServerGeneration import dev.blazelight.p4oc.domain.server.ServerRef +import dev.blazelight.p4oc.domain.session.SessionId import dev.blazelight.p4oc.domain.workspace.Workspace +import io.mockk.coVerify import io.mockk.every import io.mockk.mockk import kotlinx.coroutines.CoroutineDispatcher import kotlinx.coroutines.flow.Flow +import kotlinx.coroutines.flow.MutableStateFlow +import kotlinx.coroutines.flow.StateFlow import kotlinx.coroutines.flow.emptyFlow import kotlinx.coroutines.flow.flowOf import kotlinx.coroutines.test.StandardTestDispatcher import kotlinx.coroutines.test.runTest import kotlinx.serialization.json.Json +import org.junit.Assert.assertEquals import org.junit.Assert.assertNotSame import org.junit.Assert.assertSame import org.junit.Assert.assertTrue @@ -96,14 +104,91 @@ class SessionRepositoryProviderTest { assertTrue(lease.repository.state.value is RepoState.Hydrating) } + @Test + fun `throwing event does not permanently kill workspace event collection`() = runTest { + // First event throws inside acceptEvent; the collector must survive and process the next. + val bad = mockk { every { sessionID } throws RuntimeException("boom") } + val provider = provider( + scopedEvents = flowOf( + scoped(OpenCodeEvent.MessageUpdated(bad)), + scoped(OpenCodeEvent.MessageUpdated(assistantMessage("m1"))), + ), + dispatcher = StandardTestDispatcher(testScheduler), + ) + + val lease = provider.acquire(workspace, generation) + testScheduler.advanceUntilIdle() + + assertEquals( + listOf("m1"), + lease.repository.messages(SessionId("s1")).value.map { it.message.id }, + ) + } + + @Test + fun `reconnect refreshes messages for open sessions on every reconnect`() = runTest { + val api = mockk(relaxed = true) + val connectionState = MutableStateFlow(ConnectionState.Disconnected) + val provider = provider( + api = api, + connectionState = connectionState, + dispatcher = StandardTestDispatcher(testScheduler), + ) + + val lease = provider.acquire(workspace, generation) + lease.repository.messages(SessionId("s1")) // register an actively-observed conversation + testScheduler.advanceUntilIdle() + + // First reconnect: non-Connected -> Connected. This also runs the one-shot snapshot + // hydration, which sets (and never clears) the repository's inFlight guard. + connectionState.value = ConnectionState.Connected + testScheduler.advanceUntilIdle() + + // Second reconnect. Advance on the intermediate Error so the collector observes the + // non-Connected dwell (in production the server is down for seconds). The message refresh + // must NOT be suppressed by the now-set inFlight guard — otherwise an open conversation + // goes stale after the first reconnect (issue #14). + connectionState.value = ConnectionState.Error("dropped") + testScheduler.advanceUntilIdle() + connectionState.value = ConnectionState.Connected + testScheduler.advanceUntilIdle() + + // The open conversation self-heals via a REST message refetch on both reconnects. + coVerify(atLeast = 2) { api.getMessages("s1", null, "/repo") } + } + + private fun scoped(event: OpenCodeEvent): ScopedEvent = ScopedEvent( + serverRef = server, + generation = generation, + workspaceKey = workspace.key, + event = event, + ) + + private fun assistantMessage(id: String): Message.Assistant = Message.Assistant( + id = id, + sessionID = "s1", + createdAt = 1L, + parentID = "", + providerID = "provider", + modelID = "model", + mode = "chat", + agent = "assistant", + cost = 0.0, + tokens = TokenUsage(input = 0, output = 0), + ) + private fun provider( scopedEvents: Flow = emptyFlow(), + connectionState: StateFlow = MutableStateFlow(ConnectionState.Disconnected), + api: OpenCodeApi = mockk(relaxed = true), dispatcher: CoroutineDispatcher = StandardTestDispatcher(), ): SessionRepositoryProvider = SessionRepositoryProvider( - activeServerApiProvider = ActiveServerApiProvider { _, _ -> mockk(relaxed = true) }, + activeServerApiProvider = ActiveServerApiProvider { _, _ -> api }, messageMapper = MessageMapper(Json { ignoreUnknownKeys = true }), connectionManager = mockk { every { this@mockk.scopedEvents } returns scopedEvents + every { this@mockk.connectionState } returns connectionState + every { this@mockk.currentGeneration } returns generation }, dispatcher = dispatcher, ) diff --git a/docs/design-locks/B-sse-hydrate-race.md b/docs/design-locks/B-sse-hydrate-race.md index 51dbcec6..7e905320 100644 --- a/docs/design-locks/B-sse-hydrate-race.md +++ b/docs/design-locks/B-sse-hydrate-race.md @@ -31,6 +31,7 @@ 9. Cancellation during hydrate (workspace scope dies) → coroutine cancels, buffer GC'd with scope, no state commit. 10. Retry: user-triggered or auto on next reconnect. `hydrate(force=true)` rebuilds snapshot, then merges with current `Stale.snapshot` taking REST as source of truth for snapshot fields, but preserving any post-hydrate live events received during the new hydrate window (recursive — same buffering rules, capped at 3 retries). 11. App background during hydrate: `WorkspaceViewModel` survives backgrounding; hydrate continues. If SSE drops, on resume `OpenCodeEventSource` reconnects and `Stale` may trigger re-hydrate. +12. **Reconnect self-heal covers per-session message state, not just the session/project snapshot.** The synthetic `OpenCodeEvent.Connected` carries no directory, so it never reaches the per-workspace event fan-out (`ScopedEvent` would tag it `WorkspaceKey.Global`). Instead, `SessionRepositoryProvider` observes `ConnectionManager.connectionState` and, on every non-`Connected → Connected` transition, delivers `OpenCodeEvent.Connected` to every live repository for the current server generation. `hydrateAfterReconnect()` then re-fetches messages for every actively-observed session (those present in `messageStates`) via the same REST path as initial screen entry (`loadMessages`), so an open conversation catches up on anything missed during the outage without the user leaving and re-entering. A liveness watchdog in `OpenCodeEventSource` (no frame — message or heartbeat — for > 60 s while `Connected`) forces the reconnect that drives this path when a socket dies silently. (Issue #14.) ## Rejected alternatives @@ -49,6 +50,8 @@ | App backgrounds during hydrate | Hydrate continues. SSE may drop; reconnect on resume; `Stale` triggers re-hydrate. | | Buffer overflow (>512 events during slow hydrate) | Drop oldest, warn log, `overflowed=true`. On `Live` transition, schedule `hydrate(force=true)`. | | `Disconnected` during `Hydrating` | `connectionState=Disconnected` immediately, buffer untouched. On subsequent `Connected`, hydrate retries. | +| SSE reconnects while a conversation is open | Provider broadcasts `Connected` to the workspace repo; `hydrateAfterReconnect()` rebuilds the snapshot **and** re-fetches messages for every open session. Open chat catches up with no navigation (issue #14). | +| Socket dies silently (no `onError`/`onClosed`, `readTimeout=0`) | Liveness watchdog sees no frame for > 60 s while `Connected` and forces `reconnect()`, which then drives the reconnect self-heal above. | ## Files affected