Files
hapi/ios/Packages/HapiKit/Sources/HapiProtocol/Window/MessageWindowLogic.swift
T
SSU-WEI HUANGandGitHub f0e5ba9c0f feat(codex): mid-turn Steer via app-server turn/steer (#888) (#1606)
* feat(shared): steer capability gates and live steered signal schemas

- STEERING_SUPPORTED_FLAVORS / isSteeringSupportedForSession gate which
  agents can deliver queued messages into the active turn (pi, codex,
  cursor ACP; legacy stream-json cursor excluded)
- AgentState.steeringActive, DecryptedMessage.steered and
  messages-consumed  live signal (never persisted by the hub)

* feat(cli): queue reservations and steered messages-consumed option

- MessageQueue2 gains takeByLocalId/restoreReservation/
  beginReservationDispatch/commitReservation so an async steer can reserve
  a queued row without racing the main loop's turn/start drain
- emitMessagesConsumed accepts steered: true to mark mid-turn delivery

* feat(codex): mid-turn steer via app-server turn/steer (#888)

- CodexAppServerClient.steerTurn + TurnSteerParams/Response types
- CodexRemoteLauncher registers the steer-queued-message RPC handler:
  reserves the queued row, validates it against the active turn (no
  control commands, matching mode hash), injects via turn/steer with an
  epoch guard that invalidates in-flight steers on abort/cleanup
- steeringActive agent state tracks the active-turn window
- hub syncEngine gate opens to codex; messages-consumed relays steered

* feat(web): Steered badge and steer gating for codex sessions

- HappyUserMessage shows a ↳ Steered badge fed by the live
  messages-consumed steered signal, preserved across server echoes and
  refetches (mergeMessages carries the optimistic marker)
- SessionChat gates canSteer via isSteeringSupportedForSession instead of
  the pi-only check
- clearStaleQueuedStatus normalizes a queued status on an invoked message
- fix(web): drop duplicate showSessionSummaryInChat in markdown test
  (upstream typecheck breakage)

* fix(codex,shared): address bot findings on steer gate and ambiguous turn/steer

- STEERING_SUPPORTED_FLAVORS / isSteeringSupportedForSession advertise
  codex and pi only; cursor joins when its soft-steer handler lands (#1609)
- turn/steer now splits dispatch (stdin accepted) from completion (turn
  finished): the hub RPC acks once dispatch succeeds — never on the
  concurrent turn's completion, which can exceed the 30s RPC window
- queue row commits only after the turn settles; a rejected/aborted steer
  restores the row so the message still delivers via turn/start, and a
  dispatched steer is never restored (no duplicate delivery)
- steer carries clientUserMessageId (echoed as userMessage.clientId) so
  ambiguous transport failures can reconcile the thread later
- client tests cover dispatch/complete split and stdin-write failure

* fix(codex): reconcile dispatched steers before restoring; align error copy

- A dispatched turn/steer whose completion fails (disconnect / protocol
  error) is now reconciled via thread/read by clientUserMessageId before
  the queued row is restored — the instruction is only re-delivered by
  turn/start when the thread never received it
- Reconcile targets the pinned steer thread, not whichever turn is
  current when completion fails
- syncEngine unsupported-flavor error now matches the capability gate
  (Pi and Codex only until the cursor handler lands)
- launcher tests cover steer success (ack on dispatch), reconcile-accepted
  and reconcile-rejected outcomes

* fix(codex): consume the row at dispatch; drop background reconcile

- The hub RPC acks and the queue row is consumed as soon as stdin accepts
  turn/steer; completion is background-only logging. A dispatched steer is
  never restored, so the same localId cannot be re-delivered via turn/start
  after the caller was told the steer succeeded
- Dispatch failure (stdin write error) still restores the row and reports
  failure
- steer.completed rejection is always handled (no unhandled rejection on
  the dispatch-failure path)
- tests updated: completion failure after dispatch keeps the row consumed;
  dispatch failure restores it

* fix(codex): distinguish definite rejection from indeterminate completion

- Transport-level failures (timeout, abort, disconnect, spawn, protocol)
  carry an indeterminate marker; explicit JSON-RPC error responses do not
- After a dispatched steer, turn completion resolves → commit + consumed;
  a definite app-server rejection restores the row (instruction was never
  accepted, so turn/start cannot duplicate it); an indeterminate outcome
  leaves the row reserved so it can never be delivered twice
- Completion handling registers before awaiting dispatch so the
  dispatch-failure path cannot leak an unhandled rejection
- client/launcher tests cover explicit rejection (restore), indeterminate
  outcome (row stays reserved) and dispatch failure

* fix(codex): reconcile indeterminate steers instead of a permanent reservation

- After an indeterminate completion (disconnect/protocol), reconcile the
  thread by clientUserMessageId immediately: accepted → commit + consumed,
  provably rejected → restore, still unreadable → keep the reservation and
  retry from the main-loop top on later passes (post-reconnect)
- A row never sits in dispatching forever: the hub cannot stamp it invoked
  while the instruction may never have been accepted
- tests: indeterminate keeps reserved while thread unreadable; accepted
  reconciliation consumes; rejected path restores

* fix(codex): accept all thread item shapes; retry reconcile; ack through abort

- Reconcile matcher accepts userMessage/user_message with clientId/
  client_id, matching the shapes the thread parser supports — an accepted
  steer can no longer be misclassified as rejected
- A pending reconciliation schedules a wakeLoop retry, so a temporary
  app-server outage cannot strand the reservation behind waitForTurnOrRecovery
- The success-path ACK no longer checks the steer epoch: the hub already
  reported steered on dispatch, so commit + messages-consumed must reach
  it even when an abort resets the queue in between

* fix(codex): reinit reconnected app-server; keep reconcile retries alive

- thread/read after a disconnect auto-connects a fresh app-server, which
  must be initialized before any request — reconcile now ensures
  connect + initialize (isConnected getter added to the client)
- every still-unknown loop-top reconciliation schedules the next retry,
  so recovery without external traffic is eventually observed
- launcher mock gains isConnected

* fix(codex): timer-driven reconciliation; init tracking; abort-safe ACK

- Reconciliation runs on a self-rescheduling 1s timer independent of the
  main loop (wakes it too), so idle loops and waitForTurnOrRecovery still
  observe app-server recovery; abort clears nothing implicitly — the ACK
  path commits and consumes even when the reservation was cancelled
- Absence of a durable client id is ambiguous: unmatched reads stay
  'unknown' and keep retrying instead of restoring the row
- CodexAppServerClient tracks initialized state (reset on disconnect/exit)
  so ensureAppServerInitialized re-initializes a fresh process before
  thread/read; initialize failures leave the flag false for the next retry
- tests: accepted reconciliation via scheduled timer, indeterminate
  keeps reserved, explicit rejection restores

* fix(codex): bind reconciliation to the launcher lifecycle

- runSteerReconciliation clears any armed retry timer on entry and never
  installs a second one, so loop-top and timer-driven passes cannot
  multiply
- shuttingDown is set when the main loop ends: timers are cleared and the
  pending map is dropped, so an unresolved steer can never respawn an
  app-server after cleanup (remote-to-local switch included)

* fix(codex): report steered only after app-server acceptance

- The handler now awaits steer.completed (the inject-acceptance response):
  an explicit JSON-RPC rejection surfaces as failed and restores the row
  for the normal turn/start path instead of a false steered
- Transport failure after dispatch reports 'Steer outcome is being
  reconciled' and keeps the row reserved while the timer-driven thread
  reconciliation runs
- dispatch-failure path also swallows the paired completion rejection

* fix(steer): tri-state cancel, clear-safe reservations, bounded acceptance wait

- MessageQueue2.cancelByLocalId returns 'in-flight' for a dispatching
  steer reservation: the hub neither deletes the row nor stamps invoked_at
  (new CancelMessageResponse 'busy' status; web restores the optimistic
  row); pushIsolateAndClear and reset/close share cancelReservations so
  /clear-style commands cannot have a rejected steer resurrect a discarded
  prompt
- turn/steer acceptance wait bounded at 25s (< hub 30s RPC timeout): a
  lost response is indeterminate and funnels into thread reconciliation
  instead of stranding the reservation
- tests updated for the tri-state cancel contract

* fix(codex,web): busy-aware edit flow; bound reconciliation reads

- QueuedMessagesBar edit flow treats a 'busy' cancel as unsuccessful: it
  never prefills the composer when the row is inside an async steer, so a
  second client cannot send a duplicate
- reconcileSteerByClientId bounds thread/read with a 5s timeout so a
  connected-but-silent app-server cannot hold the reservation in-flight
  indefinitely

* fix(steer): inFlight-dominated cancel acks; bounded reconciliation

- hub cancel-queued-message acks check inFlight before removed: a stale
  duplicate socket reporting removed can no longer delete the durable row
  while another socket is dispatching the steer
- reconciliation entries expire after 60s and mark delivered: after the
  rejection window, a dispatched steer that the app-server never proved
  (client ids dropped on restart) is committed instead of polling
  thread/read forever
- pre-dispatch failures (abort before write included) never enter
  reconciliation — they restore the row and report failure

* fix(steer): persist indeterminate outcomes without replay

* fix(steer): make ambiguous delivery restart-safe

* fix(steer): recover crash-held rows and preserve retry dedup

* fix(steer): ack retries and bound stdin dispatch

* fix(steer): reconcile indeterminate dispatches and serialize retries

* fix(codex): classify stdin callback failures as indeterminate

* fix(steer): recheck indeterminate cancels after ACK

* fix(steer): close retry and abort races

* fix(steer): serialize live retries and abort admission

* fix(steer): distinguish live dispatching from unknown

* fix(steer): keep ACK failures held and reconcile busy cancel

* fix(steer): distinguish held cancel from removal

* fix(store): combine schema v24 migrations

* fix(store): reserve schema v25 for steer delivery state

* fix(steer): keep held cancel state and notify requeue

* fix(steer): release explicitly cancelled unknown reservations

* fix(codex): reject cancelled reservations before native steer

* fix(codex): make reservation restore atomic with state

* fix(codex): terminate abandoned transport writes

* fix(steer): own abandoned app-server lifecycle and consume races

* fix(codex): confirm dispatch and recover abandoned turns

* test(codex): mock abandoned transport callback

* fix(codex): clear visible turn state on transport loss

* fix(steer): claim retries and cover native delivery state

* fix(native): preserve indeterminate state on Android hydration

* fix(steer): make retry claims single-winner

* fix(steer): serialize concurrent retry claims

* fix(socket): tolerate missing steer-state ACK callbacks

* fix(native): serialize retry operations

* docs(web): document unknown steer delivery and retry controls

* fix(steer): handle retry failures and abort-before-connect

* fix(steer): reinitialize after transport loss and finish iOS retry errors

* fix(steer): preserve indeterminate rows across reconnect gaps

* test(web): mock indeterminate queued recovery state

* fix(steer): recover consumed ACK tombstones

* fix(steer): expose consumed cancel tombstones
2026-08-19 20:07:39 +08:00

764 lines
30 KiB
Swift

import Foundation
/// Pure state transitions of the message window — a function-for-function
/// port of `web/src/lib/message-window-store.ts` (with the async
/// orchestration stripped out; `HapiClient`'s `MessageWindowController`
/// re-adds it), matching the Android reference port
/// (`window/MessageWindowLogic.kt`) one-to-one. Every function takes the
/// previous ``MessageWindowState`` and returns the next one; nothing here
/// touches clocks, I/O, or tasks, so the whole surface is fixture-comparable.
public enum MessageWindowLogic {
public enum TrimMode: Sendable {
case append
case prepend
}
public static func createState(sessionId: String) -> MessageWindowState {
MessageWindowState(sessionId: sessionId)
}
// MARK: - Helpers
/// Web `buildState`: swap the messages list, re-derive seq bounds, bump
/// the version. Every call site passes a freshly-built array (the web
/// counterpart always builds a new array instance too), so the version
/// bumps unconditionally.
private static func settingMessages(
_ previous: MessageWindowState,
_ messages: [WindowMessage]
) -> MessageWindowState {
var oldestSeq: Int?
var newestSeq: Int?
for message in messages {
guard let seq = message.seq else { continue }
oldestSeq = oldestSeq.map { min($0, seq) } ?? seq
newestSeq = newestSeq.map { max($0, seq) } ?? seq
}
var next = previous
next.messages = messages
next.oldestSeq = oldestSeq
next.newestSeq = newestSeq
next.messagesVersion = previous.messagesVersion + 1
return next
}
/// A row's compound position; rows without a server `seq` have none.
public static func messagePosition(_ message: WindowMessage) -> MessagePosition? {
guard let seq = message.seq else { return nil }
return MessagePosition(at: message.positionAt, seq: seq)
}
public enum PositionEnd: Sendable {
case oldest
case newest
}
public static func derivePosition(
_ messages: [WindowMessage],
_ end: PositionEnd
) -> MessagePosition? {
var selected: MessagePosition?
for message in messages {
guard let candidate = messagePosition(message) else { continue }
guard let current = selected else {
selected = candidate
continue
}
if (end == .oldest && candidate < current) || (end == .newest && candidate > current) {
selected = candidate
}
}
return selected
}
/// Pairwise page cursor (`pagePosition` in the web): both halves or nothing.
public static func pagePosition(at: Int?, seq: Int?) -> MessagePosition? {
guard let at, let seq else { return nil }
return MessagePosition(at: at, seq: seq)
}
private static func maxPosition(_ a: MessagePosition?, _ b: MessagePosition?) -> MessagePosition? {
switch (a, b) {
case (nil, let b): return b
case (let a, nil): return a
case (let a?, let b?): return a >= b ? a : b
}
}
// MARK: - Trims
private struct Trim {
let kept: [WindowMessage]
let dropped: [WindowMessage]
}
private static func sliceForTrim(
_ items: [WindowMessage],
limit: Int,
mode: TrimMode
) -> Trim {
if items.count <= limit { return Trim(kept: items, dropped: []) }
if limit <= 0 { return Trim(kept: [], dropped: items) }
switch mode {
case .prepend:
return Trim(
kept: Array(items[..<limit]),
dropped: Array(items[limit...])
)
case .append:
return Trim(
kept: Array(items[(items.count - limit)...]),
dropped: Array(items[..<(items.count - limit)])
)
}
}
/// Codex background-agent trace rows (`agent-run-*`) get their own trim
/// bucket so long traces do not evict chat (web `isCodexAgentRunMessage`).
public static func isCodexAgentRunMessage(_ message: WindowMessage) -> Bool {
guard let outer = message.content.objectValue,
outer["role"]?.stringValue == "agent",
let payload = outer["content"]?.objectValue,
payload["type"]?.stringValue == "codex",
let data = payload["data"]?.objectValue
else { return false }
let type = data["type"]?.stringValue
return type == "agent-run-start" || type == "agent-run-update" || type == "agent-run-trace"
}
/// Trim to `regularLimit` while never dropping queued rows: the regular
/// budget shrinks by the queued count, `agent-run-*` rows trim against
/// their own bucket (`agentRunWindowSize`), and queued rows are re-merged
/// afterwards (web `trimPreservingQueued`).
private static func trimPreservingQueued(
_ messages: [WindowMessage],
regularLimit: Int,
mode: TrimMode
) -> Trim {
let queued = messages.filter(\.isQueuedForInvocation)
let queuedIds = Set(queued.map(\.id))
let nonQueued = messages.filter { !queuedIds.contains($0.id) }
let agentRuns = nonQueued.filter { isCodexAgentRunMessage($0) }
let regular = nonQueued.filter { !isCodexAgentRunMessage($0) }
let regularTrim = sliceForTrim(regular, limit: max(0, regularLimit - queued.count), mode: mode)
let agentRunTrim = sliceForTrim(
agentRuns,
limit: MessageWindowConstants.agentRunWindowSize,
mode: mode
)
return Trim(
kept: MessageMerge.mergeMessages(regularTrim.kept + agentRunTrim.kept, queued),
dropped: regularTrim.dropped + agentRunTrim.dropped
)
}
// MARK: - Retention
/// Web `shouldRetainWindowMessage`: queued, or renderable by the chat
/// pipeline.
public static func shouldRetainWindowMessage(_ message: WindowMessage) -> Bool {
message.isQueuedForInvocation || MessageRetention.isRenderable(message.content)
}
/// How many of `incoming` would add a NEW renderable row (not already
/// represented by id or localId) — feeds the `applied` older-load outcome.
public static func countNewRenderableMessages(
_ previous: MessageWindowState,
incoming: [WindowMessage]
) -> Int {
var representedIds = Set(previous.messages.map(\.id))
var representedLocalIds = Set(previous.messages.compactMap(\.localId))
var count = 0
for message in incoming {
guard shouldRetainWindowMessage(message) else { continue }
if representedIds.contains(message.id) { continue }
if let localId = message.localId, representedLocalIds.contains(localId) { continue }
count += 1
representedIds.insert(message.id)
if let localId = message.localId { representedLocalIds.insert(localId) }
}
return count
}
// MARK: - Merge
/// Merge `incoming` into the window, trim per mode/limit, and repair
/// cursors/flags on overflow (web `mergeIntoWindow`): an append-side trim
/// flips `hasMore` and recomputes the older cursor from the oldest kept
/// row; a prepend-side trim drops tail rows, so the window no longer
/// reaches the live bottom — flag `requiresLatestReset` and pull the
/// newest cursor back to the newest kept row.
public static func mergeIntoWindow(
_ previous: MessageWindowState,
incoming: [WindowMessage],
mode: TrimMode? = nil,
regularLimit: Int? = nil,
advanceTailRevision: Bool = false
) -> MessageWindowState {
let retainedIncoming = incoming.filter { shouldRetainWindowMessage($0) }
if retainedIncoming.isEmpty {
return previous
}
let effectiveMode = mode
?? (previous.viewMode == .history ? .prepend : .append)
let effectiveLimit = regularLimit
?? (previous.viewMode == .history
? MessageWindowConstants.historyWindowSize
: MessageWindowConstants.visibleWindowSize)
let merged = MessageMerge.mergeMessages(previous.messages, retainedIncoming)
let trim = trimPreservingQueued(merged, regularLimit: effectiveLimit, mode: effectiveMode)
var next = settingMessages(previous, trim.kept)
if advanceTailRevision {
next.tailRevision = previous.tailRevision + 1
}
if trim.dropped.isEmpty {
return next
}
if effectiveMode == .append {
let oldest = derivePosition(trim.kept, .oldest)
next.hasMore = true
next.oldestPosition = oldest ?? next.oldestPosition
return next
}
next.requiresLatestReset = true
next.newestPosition = derivePosition(trim.kept, .newest)
return next
}
// MARK: - Latest replace
/// Apply a `latest` page (cold start, activation refresh, or a reset
/// response). With `replaceServerRows` the page is authoritative: every
/// server row captured in `requestBaseline` (by **instance identity**) is
/// discarded, while optimistic rows and rows that changed since the
/// request left (concurrent SSE) survive the swap (web
/// `applyLatestResponse`).
public static func applyLatestResponse(
_ previous: MessageWindowState,
responseMessages: [WindowMessage],
page: MessagesPage,
replaceServerRows: Bool,
requestBaseline: [String: WindowMessage]
) -> MessageWindowState {
let retainedResponseMessages = responseMessages.filter { shouldRetainWindowMessage($0) }
let concurrentServerRows = previous.messages.filter { message in
!message.isOptimistic && requestBaseline[message.id] !== message
}
let preserved = replaceServerRows
? previous.messages.filter { message in
message.isOptimistic || requestBaseline[message.id] !== message
}
: previous.messages
let authoritative = MessageMerge.mergeMessages(preserved, retainedResponseMessages)
let incoming = MessageMerge.mergeMessages(authoritative, concurrentServerRows)
let trim = trimPreservingQueued(
incoming,
regularLimit: MessageWindowConstants.visibleWindowSize,
mode: .append
)
let snapshotHead = pagePosition(at: page.snapshotHeadAt, seq: page.snapshotHeadSeq)
?? derivePosition(responseMessages, .newest)
let newestKept = derivePosition(trim.kept, .newest)
let newest = maxPosition(snapshotHead, newestKept)
let responseOldest = pagePosition(at: page.nextBeforeAt, seq: page.nextBeforeSeq)
let oldest: MessagePosition?
if !trim.dropped.isEmpty {
oldest = derivePosition(trim.kept, .oldest)
} else if replaceServerRows {
oldest = responseOldest
} else {
oldest = responseOldest ?? previous.oldestPosition
}
var next = settingMessages(previous, trim.kept)
next.hasMore = page.hasMore || (!replaceServerRows && previous.hasMore) || !trim.dropped.isEmpty
next.epoch = page.epoch
next.oldestPosition = oldest
next.newestPosition = newest
next.tailRevision = previous.tailRevision + 1
next.requiresLatestReset = false
next.isLoadingMore = replaceServerRows ? false : previous.isLoadingMore
next.olderGeneration = replaceServerRows ? previous.olderGeneration + 1 : previous.olderGeneration
next.warning = nil
return next
}
// MARK: - Tail sync
/// Web `beginTailSync`: claim a new sync generation. The older generation
/// bumps too — an older-page response captured before this point must not
/// commit while the tail request is in flight, or a reset can mistake it
/// for concurrent SSE.
public static func beginTailSync(_ previous: MessageWindowState) -> MessageWindowState {
var next = previous
next.syncGeneration = previous.syncGeneration + 1
next.olderGeneration = previous.olderGeneration + 1
next.isSyncingTail = true
next.isLoadingMore = false
next.warning = nil
return next
}
public static func finishTailSync(
_ previous: MessageWindowState,
generation: Int,
warning: String?
) -> MessageWindowState {
guard previous.syncGeneration == generation else { return previous }
var next = previous
next.isSyncingTail = false
next.warning = warning
return next
}
/// Loop body of the after-cursor catch-up (web `runTailSync` inner
/// updater): merge the page, adopt its epoch, and advance the newest
/// cursor to `max(current, nextAfter)` — unless a prepend-side trim just
/// flagged a latest reset, in which case the cursor stays pulled back.
public static func applyAfterPage(
_ previous: MessageWindowState,
responseMessages: [WindowMessage],
page: MessagesPage,
nextAfter: MessagePosition?
) -> MessageWindowState {
var merged = mergeIntoWindow(previous, incoming: responseMessages, advanceTailRevision: true)
if merged.requiresLatestReset {
merged.epoch = page.epoch
merged.warning = nil
return merged
}
merged.epoch = page.epoch
merged.newestPosition = maxPosition(nextAfter, merged.newestPosition)
merged.warning = nil
return merged
}
// MARK: - View mode
/// Re-enter tail mode (web `enterTailMode`): trim to the visible window;
/// when a history overflow flagged `requiresLatestReset`, drop epoch and
/// newest cursor so the next tail sync fetches a fresh latest page.
/// The flag itself is deliberately NOT cleared here (web quirk, pinned) —
/// `applyLatestResponse` clears it when the fresh page lands.
public static func enterTailMode(_ previous: MessageWindowState) -> MessageWindowState {
let trim = trimPreservingQueued(
previous.messages,
regularLimit: MessageWindowConstants.visibleWindowSize,
mode: .append
)
let forceLatest = previous.requiresLatestReset
let oldest = !trim.dropped.isEmpty
? derivePosition(trim.kept, .oldest)
: previous.oldestPosition
var next = settingMessages(previous, trim.kept)
next.hasMore = previous.hasMore || !trim.dropped.isEmpty
next.viewMode = .tail
next.epoch = forceLatest ? nil : previous.epoch
next.oldestPosition = oldest
next.newestPosition = forceLatest ? nil : previous.newestPosition
return next
}
public static func setViewMode(
_ previous: MessageWindowState,
mode: MessageViewMode
) -> MessageWindowState {
if previous.viewMode == mode { return previous }
if mode == .history {
var next = previous
next.viewMode = .history
return next
}
return enterTailMode(previous)
}
public struct Activation {
public let state: MessageWindowState
public let requestedLatest: Bool
}
/// Session (re-)activation (web `activateMessageWindow`): trim back to
/// the visible window, and when a usable persisted cursor exists, flag
/// `preferLatestOnActivation` — another client may have advanced the
/// session by many pages, so the next tail sync fetches the current tail
/// first and reconciles it through the reset-preservation path.
public static func activate(_ previous: MessageWindowState) -> Activation {
let trim = trimPreservingQueued(
previous.messages,
regularLimit: MessageWindowConstants.visibleWindowSize,
mode: .append
)
let forceLatest = previous.requiresLatestReset
let hasUsableCursor = previous.newestPosition != nil && previous.epoch != nil && !forceLatest
let preferLatestOnActivation = hasUsableCursor && !trim.kept.isEmpty
let requestedLatest = preferLatestOnActivation && !previous.preferLatestOnActivation
let invalidateRunningSync = requestedLatest && previous.isSyncingTail
func withActivationUpdates(_ state: MessageWindowState) -> MessageWindowState {
guard preferLatestOnActivation else { return state }
var flagged = state
flagged.preferLatestOnActivation = true
if invalidateRunningSync {
flagged.syncGeneration += 1
flagged.olderGeneration += 1
}
return flagged
}
if previous.viewMode == .tail
&& trim.kept.count == previous.messages.count
&& !forceLatest {
let state = preferLatestOnActivation ? withActivationUpdates(previous) : previous
return Activation(state: state, requestedLatest: requestedLatest)
}
let next = enterTailMode(previous)
let state = preferLatestOnActivation ? withActivationUpdates(next) : next
return Activation(state: state, requestedLatest: requestedLatest)
}
// MARK: - Older pages
public enum OlderPrecheck {
case proceed(before: MessagePosition)
case stop(OlderLoadOutcome.StopReason)
}
public static func olderLoadPrecheck(_ state: MessageWindowState) -> OlderPrecheck {
let before = state.oldestPosition
if state.isSyncingTail || state.isLoadingMore {
return .stop(.busy)
}
if !state.hasMore {
return .stop(.exhausted)
}
guard let before else {
return .stop(.unavailable)
}
return .proceed(before: before)
}
public static func beginOlderLoad(
_ previous: MessageWindowState,
generation: Int
) -> MessageWindowState {
var next = previous
next.olderGeneration = generation
next.isLoadingMore = true
next.warning = nil
return next
}
/// An older page answered with a different epoch: every cursor is
/// meaningless. Drop epoch + newest cursor and flag the latest reset;
/// the caller then runs a tail sync and reports `stopped/epoch-reset`.
public static func applyOlderEpochMismatch(
_ previous: MessageWindowState,
generation: Int
) -> MessageWindowState {
guard previous.olderGeneration == generation else { return previous }
var next = previous
next.isLoadingMore = false
next.epoch = nil
next.newestPosition = nil
next.requiresLatestReset = true
return next
}
/// The `onBeforeApply` veto path: invalidate this load, keep the window.
public static func rejectOlderApply(_ previous: MessageWindowState) -> MessageWindowState {
var next = previous
next.olderGeneration = previous.olderGeneration + 1
next.isLoadingMore = false
next.warning = nil
return next
}
public static func applyOlderResponse(
_ previous: MessageWindowState,
responseMessages: [WindowMessage],
page: MessagesPage,
historyVersion: Int
) -> MessageWindowState {
var merged = mergeIntoWindow(
previous,
incoming: responseMessages,
mode: .prepend,
regularLimit: MessageWindowConstants.olderLoadWindowSize
)
merged.hasMore = page.hasMore
merged.epoch = page.epoch
merged.oldestPosition = pagePosition(at: page.nextBeforeAt, seq: page.nextBeforeSeq)
merged.isLoadingMore = false
merged.historyVersion = historyVersion
merged.warning = nil
return merged
}
public static func failOlderLoad(
_ previous: MessageWindowState,
generation: Int,
warning: String
) -> MessageWindowState {
guard previous.olderGeneration == generation else { return previous }
var next = previous
next.isLoadingMore = false
next.warning = warning
return next
}
public static func cancelOlderLoad(_ previous: MessageWindowState) -> MessageWindowState {
guard previous.isLoadingMore else { return previous }
var next = previous
next.olderGeneration = previous.olderGeneration + 1
next.isLoadingMore = false
next.warning = nil
return next
}
// MARK: - Live rows
/// SSE `message-received` ingest (web `ingestIncomingMessages`): merge,
/// then — when an epoch is cached and no reset is pending — advance the
/// newest cursor past the incoming rows' positions, **including rows the
/// pipeline hides**, so a later tail sync does not refetch them.
public static func ingestIncoming(
_ previous: MessageWindowState,
incoming: [WindowMessage]
) -> MessageWindowState {
if incoming.isEmpty { return previous }
var merged = mergeIntoWindow(previous, incoming: incoming, advanceTailRevision: true)
if merged.epoch == nil || merged.requiresLatestReset {
return merged
}
let incomingNewest = derivePosition(incoming, .newest)
let currentNewest = merged.newestPosition
if let incomingNewest, currentNewest == nil || incomingNewest > currentNewest! {
merged.newestPosition = incomingNewest
}
return merged
}
/// Optimistic append (web `appendOptimisticMessage`): never advances the
/// newest cursor.
public static func appendOptimistic(
_ previous: MessageWindowState,
message: WindowMessage
) -> MessageWindowState {
mergeIntoWindow(
previous,
incoming: [message],
mode: previous.viewMode == .history ? .prepend : .append,
advanceTailRevision: true
)
}
public static func updateStatus(
_ previous: MessageWindowState,
localId: String,
status: MessageStatus
) -> MessageWindowState {
if localId.isEmpty { return previous }
var changed = false
let messages = previous.messages.map { message -> WindowMessage in
if message.localId != localId || message.status == status {
return message
}
changed = true
return message.withStatus(status)
}
return changed ? settingMessages(previous, messages) : previous
}
public static func markIndeterminate(
_ previous: MessageWindowState,
localIds: [String]
) -> MessageWindowState {
if localIds.isEmpty { return previous }
let idSet = Set(localIds)
var changed = false
let messages = previous.messages.map { message -> WindowMessage in
guard let localId = message.localId, idSet.contains(localId), message.status != .indeterminate else {
return message
}
changed = true
return message.withStatus(.indeterminate)
}
return changed ? settingMessages(previous, messages) : previous
}
public static func markRequeued(
_ previous: MessageWindowState,
localIds: [String]
) -> MessageWindowState {
if localIds.isEmpty { return previous }
let idSet = Set(localIds)
var changed = false
let messages = previous.messages.map { message -> WindowMessage in
guard let localId = message.localId, idSet.contains(localId), message.status == .indeterminate else {
return message
}
changed = true
return message.withStatus(.queued)
}
return changed ? settingMessages(previous, messages) : previous
}
/// `message-cancelled` / optimistic DELETE removal: matches localId OR
/// id; idempotent.
public static func removeByLocalIdOrId(
_ previous: MessageWindowState,
localId: String
) -> MessageWindowState {
if localId.isEmpty { return previous }
let messages = previous.messages.filter { $0.localId != localId && $0.id != localId }
return messages.count == previous.messages.count
? previous
: settingMessages(previous, messages)
}
/// SSE `messages-consumed` (web `markMessagesConsumed`): stamp
/// `invokedAt` and flip status to `sent` on matching rows — including
/// server rows without a client status (web quirk, pinned) — skip
/// `failed`, then re-sort: the row moves from its enqueue position to its
/// invocation position. Never advances the newest cursor.
public static func markConsumed(
_ previous: MessageWindowState,
localIds: [String],
invokedAt: Int
) -> MessageWindowState {
if localIds.isEmpty { return previous }
let idSet = Set(localIds)
var changed = false
let updated = previous.messages.map { message -> WindowMessage in
guard let localId = message.localId, idSet.contains(localId), message.status != .failed else {
return message
}
let needsStatus = message.status != .sent
let needsInvokedAt = message.hasExplicitNullInvokedAt
if !needsStatus && !needsInvokedAt { return message }
changed = true
var next = message
if needsStatus { next = next.withStatus(.sent) }
if needsInvokedAt { next = next.withInvokedAt(invokedAt) }
return next
}
if !changed { return previous }
var next = settingMessages(previous, MessageMerge.mergeMessages([], updated))
next.tailRevision = previous.tailRevision + 1
return next
}
// MARK: - Queued reconciliation
/// Web `isQueuedReconcileCandidate`.
private static func isQueuedReconcileCandidate(_ message: WindowMessage) -> Bool {
guard message.localId != nil, message.isQueuedForInvocation else { return false }
if !message.isOptimistic { return true }
return message.status == .queued || message.status == .sent
}
/// Candidate localIds for the queued-state round trip, in window order.
public static func queuedReconcileCandidateLocalIds(_ state: MessageWindowState) -> [String] {
var seen = Set<String>()
var localIds: [String] = []
for message in state.messages where isQueuedReconcileCandidate(message) {
let localId = message.localId!
if seen.insert(localId).inserted {
localIds.append(localId)
}
}
return localIds
}
/// Apply the queued-state verdict (web `reconcileQueuedLocalIds`): drop
/// candidates that are in neither the still-queued list nor (post
/// `markConsumed`) invoked — they were deleted server-side.
public static func reconcileQueuedLocalIds(
_ previous: MessageWindowState,
candidateLocalIds: [String],
queuedLocalIds: [String]
) -> MessageWindowState {
if candidateLocalIds.isEmpty { return previous }
let candidates = Set(candidateLocalIds)
let queued = Set(queuedLocalIds)
let messages = previous.messages.filter { message in
guard let localId = message.localId, candidates.contains(localId) else { return true }
return queued.contains(localId) || !isQueuedReconcileCandidate(message)
}
return messages.count == previous.messages.count
? previous
: settingMessages(previous, messages)
}
// MARK: - Persistence
/// Web `shouldPersistState`.
public static func shouldPersist(_ state: MessageWindowState) -> Bool {
!state.messages.isEmpty
|| state.hasMore
|| state.epoch != nil
|| state.oldestPosition != nil
|| state.newestPosition != nil
}
public static func toPersisted(_ state: MessageWindowState) -> PersistedMessageWindow {
PersistedMessageWindow(
messages: state.messages,
hasMore: state.hasMore,
oldestPositionAt: state.oldestPosition?.at,
oldestPositionSeq: state.oldestPosition?.seq,
newestPositionAt: state.newestPosition?.at,
newestPositionSeq: state.newestPosition?.seq,
epoch: state.epoch
)
}
/// Rebuild a state from a persisted snapshot (web `hydrateState`):
/// interrupted `sending` rows restore to `queued`/`sent`, and a snapshot
/// with rows but no usable cursor/epoch starts flagged for a latest
/// reset.
public static func hydrate(
sessionId: String,
persisted: PersistedMessageWindow
) -> MessageWindowState {
let restored = persisted.messages.map { message -> WindowMessage in
guard message.status == .sending else { return message }
return message.withStatus(message.hasExplicitNullInvokedAt ? .queued : .sent)
}
let oldest = pagePosition(at: persisted.oldestPositionAt, seq: persisted.oldestPositionSeq)
let newest = pagePosition(at: persisted.newestPositionAt, seq: persisted.newestPositionSeq)
let epoch = persisted.epoch.flatMap { $0 >= 0 ? $0 : nil }
var next = settingMessages(
createState(sessionId: sessionId),
MessageMerge.mergeMessages([], restored)
)
next.hasMore = persisted.hasMore
next.oldestPosition = oldest
next.newestPosition = newest
next.epoch = epoch
next.requiresLatestReset = !persisted.messages.isEmpty && (newest == nil || epoch == nil)
return next
}
// MARK: - Seed
/// Seed a fresh window from another session's (web
/// `seedMessageWindowFromSession`) — resume/reopen may hand back a
/// different session id; the old rows render instantly while
/// `requiresLatestReset` forces a fresh latest page underneath. Cursor
/// epoch and newest position are NOT carried over.
public static func seededState(
source: MessageWindowState,
target: MessageWindowState
) -> MessageWindowState {
var next = settingMessages(createState(sessionId: target.sessionId), source.messages)
next.hasMore = source.hasMore
next.tailRevision = source.tailRevision
next.oldestPosition = source.oldestPosition
next.requiresLatestReset = true
next.syncGeneration = target.syncGeneration + 1
next.olderGeneration = target.olderGeneration + 1
return next
}
}