added a watchdog to force the SSE stream open
This commit is contained in:
1 parent
abbed438ec
commit
ca622d3a39
2 files changed
+44
-3
No files matched your search
@@ -8,6 +8,7 @@ import iris.protocol.ServerCaps
|
||||
import iris.protocol.messageSendFrame
|
||||
import iris.protocol.syncFrame
|
||||
import iris.util.IrisLog
|
||||
import iris.util.nowMillis
|
||||
import kotlinx.coroutines.CancellationException
|
||||
import kotlinx.coroutines.CompletableDeferred
|
||||
import kotlinx.coroutines.CoroutineScope
|
||||
@@ -110,6 +111,16 @@ class GatewayClient(
|
||||
// reconnect state race in the connect loop.
|
||||
private var lastAck: HelloAckPayload? = null
|
||||
|
||||
// Epoch ms of the last frame delivered by the receive stream (SSE or
|
||||
// long-poll). The watchdog uses this to detect a STALE stream: a live
|
||||
// connection (heartbeats / poll answers keep it open, so the read timeout
|
||||
// never fires) that has stopped delivering frames — the state a gateway
|
||||
// restart can leave the app in, where sent messages sit at "sending…"
|
||||
// because their echo is never delivered (issue #6). Set when the receive
|
||||
// loop (re)starts so the first window isn't treated as stale.
|
||||
@Volatile
|
||||
private var lastFrameMs: Long = 0L
|
||||
|
||||
/**
|
||||
* Fired promptly the moment the SSE hello (hello.ack) is received — on
|
||||
* every (re)connect. Used for time-critical work that must not wait for
|
||||
@@ -186,6 +197,7 @@ class GatewayClient(
|
||||
sseFailures = 0
|
||||
usingLongPoll = false
|
||||
httpCursor = store.syncCursor
|
||||
lastFrameMs = nowMs()
|
||||
// Provisional Connected state (previous caps/channels) until the
|
||||
// SSE hello arrives with the real ones.
|
||||
val prev = _state.value
|
||||
@@ -217,6 +229,16 @@ class GatewayClient(
|
||||
receiveJob.cancel()
|
||||
break
|
||||
}
|
||||
// Stale-stream watchdog: the gateway is up (health ok) but
|
||||
// the receive stream has delivered no frames for a while —
|
||||
// a live-but-dead connection the read timeout can't see.
|
||||
// Force a fresh (re)connect so parked frames (e.g. our own
|
||||
// echo) are re-delivered from the outbox.
|
||||
if (lastFrameMs > 0L && nowMs() - lastFrameMs > STALE_STREAM_TIMEOUT_MS) {
|
||||
IrisLog.w("receive stream stale (no frames for ${STALE_STREAM_TIMEOUT_MS}ms); forcing reconnect")
|
||||
receiveJob.cancel()
|
||||
break
|
||||
}
|
||||
}
|
||||
receiveJob.join()
|
||||
}
|
||||
@@ -321,16 +343,24 @@ class GatewayClient(
|
||||
|
||||
/** The receive stream is open again (long-poll answered): the link is
|
||||
* back. Long-poll has no hello, so restore the Connected state from the
|
||||
* last hello.ack (the SSE path gets a fresh one). */
|
||||
* last hello.ack (the SSE path gets a fresh one). Fire [onHelloAck] too:
|
||||
* the long-poll path has no `event: hello`, so without this the app's
|
||||
* "on connected" side effects (resend queued offline sends, load the
|
||||
* active lane's history) never run when a gateway restart lands while
|
||||
* the app is on the long-poll fallback — queued messages would sit at
|
||||
* "sending…" forever until an app restart (issue #6). */
|
||||
private fun restoreConnected() {
|
||||
if (_state.value !is State.Reconnecting) return
|
||||
_state.value =
|
||||
val connected =
|
||||
lastAck?.let { State.Connected(it.serverCaps, it.channels, it.lastPushedCursor) }
|
||||
?: State.Connected(ServerCaps(), emptyList())
|
||||
_state.value = connected
|
||||
onHelloAck?.invoke(connected)
|
||||
}
|
||||
|
||||
/** Deliver an HTTP-leg frame to the same sinks as any other frame. */
|
||||
private fun emitHttpFrame(frame: Frame) {
|
||||
lastFrameMs = nowMs()
|
||||
_events.tryEmit(frame)
|
||||
frame.id?.let { id ->
|
||||
pending[id]?.complete(frame)
|
||||
@@ -410,6 +440,12 @@ class GatewayClient(
|
||||
companion object {
|
||||
const val UPLOAD_TIMEOUT_MS = 120_000L
|
||||
const val PULL_TIMEOUT_MS = 300_000L
|
||||
|
||||
/** No frames from the receive stream for this long (gateway still
|
||||
* healthy) = stale stream: force a fresh (re)connect. 3 min keeps an
|
||||
* idle app from churning while still recovering a stuck stream well
|
||||
* before a user notices (issue #6). */
|
||||
const val STALE_STREAM_TIMEOUT_MS = 180_000L
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -518,4 +554,6 @@ class GatewayClient(
|
||||
val capped = minOf(base, 30_000L)
|
||||
return capped + Random.nextLong(0, 500)
|
||||
}
|
||||
|
||||
private fun nowMs(): Long = nowMillis()
|
||||
}
|
||||
@@ -892,7 +892,10 @@ class IrisController(
|
||||
* refreshes (skipped on a plain reconnect via historyLoaded).
|
||||
*/
|
||||
private fun onConnectedLane(connected: GatewayClient.State.Connected) {
|
||||
channels.setAll(connected.channels)
|
||||
// Never wipe the directory with an empty list: the long-poll restore
|
||||
// path carries no channels when there was no prior SSE hello (lastAck
|
||||
// null), and the cached directory is still valid then.
|
||||
if (connected.channels.isNotEmpty()) channels.setAll(connected.channels)
|
||||
val home = connected.channels.firstOrNull { it.isDefault }?.chatId
|
||||
_homeChannel.value = home ?: ChatStore.DEFAULT_LANE
|
||||
if (home == null) return
|
||||
|
||||
Reference in new issue
Block a user