mirror of
https://github.com/wu736139669/hapi.git
synced 2026-10-08 19:19:42 +00:00
* 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
671 lines
27 KiB
Swift
671 lines
27 KiB
Swift
import Foundation
|
|
import HapiProtocol
|
|
|
|
/// What the chat screen renders — mirror of the web `useMessages` return
|
|
/// shape (and the Android port's `MessageWindowUiState`); cursor/generation
|
|
/// internals stay off the render path. `Equatable` compares rows by identity
|
|
/// (see ``WindowMessage``), which is exactly the change signal the web uses.
|
|
public struct MessageWindowUIState: Equatable, Sendable {
|
|
public let messages: [WindowMessage]
|
|
public let hasMore: Bool
|
|
public let isSyncingTail: Bool
|
|
public let isLoadingMore: Bool
|
|
public let warning: String?
|
|
public let viewMode: MessageViewMode
|
|
public let messagesVersion: Int
|
|
public let historyVersion: Int
|
|
public let tailRevision: Int
|
|
}
|
|
|
|
extension MessageWindowState {
|
|
public var uiState: MessageWindowUIState {
|
|
MessageWindowUIState(
|
|
messages: messages,
|
|
hasMore: hasMore,
|
|
isSyncingTail: isSyncingTail,
|
|
isLoadingMore: isLoadingMore,
|
|
warning: warning,
|
|
viewMode: viewMode,
|
|
messagesVersion: messagesVersion,
|
|
historyVersion: historyVersion,
|
|
tailRevision: tailRevision
|
|
)
|
|
}
|
|
}
|
|
|
|
/// Non-transport failure of the tail-sync loop.
|
|
struct MessageWindowSyncError: Error, LocalizedError, Equatable {
|
|
let message: String
|
|
var errorDescription: String? { message }
|
|
|
|
/// Protocol violation guard: `nextAfter` did not advance past the cursor.
|
|
static let cursorDidNotAdvance = MessageWindowSyncError(
|
|
message: "Message tail cursor did not advance"
|
|
)
|
|
}
|
|
|
|
/// Per-session message window orchestration — the async half of the web
|
|
/// reference `web/src/lib/message-window-store.ts`, driving the pure
|
|
/// transitions in `MessageWindowLogic` over actor-isolated state. Mirrors
|
|
/// the Android port's `MessageWindowStore` (`core/data/.../store/`).
|
|
///
|
|
/// Concurrency model: the web store is single-threaded JS whose interleaving
|
|
/// points are its `await`s; here every synchronous read-modify-write segment
|
|
/// between transport calls runs on the actor, and the generation counters +
|
|
/// request-baseline identity checks (ported as-is) handle whatever
|
|
/// interleaves across the provider calls — which suspend the actor. Tail
|
|
/// syncs are single-flight per session with an optional trailing re-run,
|
|
/// exactly like the web `TailSyncController`.
|
|
public actor MessageWindowController {
|
|
/// Immutable and Sendable — same-module callers may read it without
|
|
/// `await` (SE-0306); cross-module callers go through the actor.
|
|
public let sessionId: String
|
|
private let provider: any MessagesProviding
|
|
private let snapshots: WindowSnapshotStore?
|
|
|
|
private var stateValue: MessageWindowState
|
|
private var observers: [UUID: AsyncStream<MessageWindowState>.Continuation] = [:]
|
|
|
|
// MARK: - Lifecycle
|
|
|
|
public init(
|
|
sessionId: String,
|
|
provider: any MessagesProviding,
|
|
snapshots: WindowSnapshotStore? = nil,
|
|
initialState: MessageWindowState? = nil
|
|
) {
|
|
self.sessionId = sessionId
|
|
self.provider = provider
|
|
self.snapshots = snapshots
|
|
self.stateValue = initialState ?? MessageWindowLogic.createState(sessionId: sessionId)
|
|
}
|
|
|
|
/// The full window state (UI projects what it needs; see
|
|
/// ``MessageWindowState/uiState``).
|
|
public var state: MessageWindowState { stateValue }
|
|
|
|
/// Stream of state changes, starting with the current state. Change
|
|
/// detection is instance-based (identity-`Equatable` rows), like the
|
|
/// web's `next !== previous` publication gate.
|
|
public func states() -> AsyncStream<MessageWindowState> {
|
|
let (stream, continuation) = AsyncStream.makeStream(of: MessageWindowState.self)
|
|
let id = UUID()
|
|
observers[id] = continuation
|
|
continuation.yield(stateValue)
|
|
continuation.onTermination = { [weak self] _ in
|
|
Task { await self?.removeObserver(id) }
|
|
}
|
|
return stream
|
|
}
|
|
|
|
private func removeObserver(_ id: UUID) {
|
|
observers.removeValue(forKey: id)
|
|
}
|
|
|
|
@discardableResult
|
|
private func update(
|
|
_ transform: (MessageWindowState) -> MessageWindowState
|
|
) -> MessageWindowState {
|
|
let previous = stateValue
|
|
let next = transform(previous)
|
|
stateValue = next
|
|
if next != previous {
|
|
for continuation in observers.values {
|
|
continuation.yield(next)
|
|
}
|
|
}
|
|
return next
|
|
}
|
|
|
|
// MARK: - Tail controller
|
|
|
|
private var runSerial = 0
|
|
private var runningTask: Task<Void, Never>?
|
|
private var runningRunId: Int?
|
|
private var runningPrefersLatest = false
|
|
private var trailingRequested = false
|
|
/// Bumped by ``clear()``/``seedFrom(_:)`` so a finished run stops
|
|
/// chaining trailing runs (web: controller replacement).
|
|
private var controllerEpoch = 0
|
|
|
|
/// Run (or join) a tail sync (web `syncTailMessages`). With
|
|
/// `ensureAfterCurrent` the call drains: if a run is already in flight, a
|
|
/// trailing run is requested and awaited, so the caller returns only
|
|
/// after a sync that STARTED at or after this call.
|
|
public func syncTail(ensureAfterCurrent: Bool = false) async {
|
|
guard let running = runningTask, let runningId = runningRunId else {
|
|
await startTailSync().value
|
|
return
|
|
}
|
|
if stateValue.preferLatestOnActivation {
|
|
if runningPrefersLatest {
|
|
await running.value
|
|
return
|
|
}
|
|
trailingRequested = false
|
|
await startTailSync().value
|
|
return
|
|
}
|
|
guard ensureAfterCurrent else {
|
|
await running.value
|
|
return
|
|
}
|
|
trailingRequested = true
|
|
await drainTailSync(observedRunId: runningId, observed: running)
|
|
}
|
|
|
|
private func startTailSync() -> Task<Void, Never> {
|
|
runSerial += 1
|
|
let runId = runSerial
|
|
let epoch = controllerEpoch
|
|
let prefersLatest = stateValue.preferLatestOnActivation
|
|
// Synchronous generation bump: the web executes `runTailSync` to its
|
|
// first true suspension (the api call) before `startTailSync`
|
|
// returns, and the Android port replicates that with an UNDISPATCHED
|
|
// coroutine start. Here `beginTailSync` lands before the Task is even
|
|
// created, so an in-flight run this one replaces can never commit
|
|
// another page in between.
|
|
update { MessageWindowLogic.beginTailSync($0) }
|
|
let generation = stateValue.syncGeneration
|
|
let task = Task { [weak self] in
|
|
guard let self else { return }
|
|
await self.runTailSync(generation: generation)
|
|
// Completion bookkeeping is the run's LAST actor-isolated act, so
|
|
// anyone resuming from `await task.value` observes the trailing
|
|
// handoff already done (web: `finish` runs via .then before
|
|
// external awaiters).
|
|
await self.tailSyncRunCompleted(runId: runId, epoch: epoch)
|
|
}
|
|
runningTask = task
|
|
runningRunId = runId
|
|
runningPrefersLatest = prefersLatest
|
|
return task
|
|
}
|
|
|
|
private func tailSyncRunCompleted(runId: Int, epoch: Int) {
|
|
guard controllerEpoch == epoch, runningRunId == runId else { return }
|
|
runningTask = nil
|
|
runningRunId = nil
|
|
runningPrefersLatest = false
|
|
guard trailingRequested else { return }
|
|
trailingRequested = false
|
|
_ = startTailSync()
|
|
}
|
|
|
|
/// Web `waitForTailSyncDrain`: follow the chain until no newer run exists.
|
|
private func drainTailSync(observedRunId: Int, observed: Task<Void, Never>) async {
|
|
let epoch = controllerEpoch
|
|
var currentId = observedRunId
|
|
var current = observed
|
|
while true {
|
|
await current.value
|
|
guard controllerEpoch == epoch else { return }
|
|
guard let nextTask = runningTask, let nextId = runningRunId, nextId != currentId else {
|
|
return
|
|
}
|
|
currentId = nextId
|
|
current = nextTask
|
|
}
|
|
}
|
|
|
|
private func isCurrentTailSync(_ generation: Int) -> Bool {
|
|
stateValue.syncGeneration == generation
|
|
}
|
|
|
|
/// Body of one tail sync (web `runTailSync`, minus `beginTailSync`,
|
|
/// which ``startTailSync()`` already executed synchronously).
|
|
private func runTailSync(generation: Int) async {
|
|
do {
|
|
let initial = stateValue
|
|
let initialCursor = initial.newestPosition
|
|
let preferLatestOnActivation = initial.preferLatestOnActivation
|
|
let canIncrement = initialCursor != nil
|
|
&& initial.epoch != nil
|
|
&& !initial.requiresLatestReset
|
|
&& !preferLatestOnActivation
|
|
|
|
if !canIncrement {
|
|
let requestBaseline = baseline()
|
|
let response = try await provider.messages(
|
|
sessionId: sessionId,
|
|
query: .latest(limit: MessageWindowConstants.pageSize)
|
|
)
|
|
guard isCurrentTailSync(generation) else { return }
|
|
update { previous in
|
|
guard previous.syncGeneration == generation else { return previous }
|
|
var next = MessageWindowLogic.applyLatestResponse(
|
|
previous,
|
|
responseMessages: windowMessages(response),
|
|
page: response.page,
|
|
replaceServerRows: initial.requiresLatestReset
|
|
|| preferLatestOnActivation
|
|
|| response.page.reset,
|
|
requestBaseline: requestBaseline
|
|
)
|
|
next.preferLatestOnActivation = false
|
|
return next
|
|
}
|
|
finishTailSync(generation: generation, warning: nil)
|
|
return
|
|
}
|
|
|
|
var after = initialCursor!
|
|
var until: MessagePosition?
|
|
while true {
|
|
let requestBaseline = baseline()
|
|
let response = try await provider.messages(
|
|
sessionId: sessionId,
|
|
query: .after(
|
|
afterAt: after.at,
|
|
afterSeq: after.seq,
|
|
untilAt: until?.at,
|
|
untilSeq: until?.seq,
|
|
epoch: initial.epoch!,
|
|
limit: MessageWindowConstants.pageSize
|
|
)
|
|
)
|
|
guard isCurrentTailSync(generation) else { return }
|
|
|
|
if response.page.reset || response.page.direction == .latest {
|
|
update { previous in
|
|
guard previous.syncGeneration == generation else { return previous }
|
|
return MessageWindowLogic.applyLatestResponse(
|
|
previous,
|
|
responseMessages: windowMessages(response),
|
|
page: response.page,
|
|
replaceServerRows: true,
|
|
requestBaseline: requestBaseline
|
|
)
|
|
}
|
|
break
|
|
}
|
|
|
|
let nextAfter = MessageWindowLogic.pagePosition(
|
|
at: response.page.nextAfterAt,
|
|
seq: response.page.nextAfterSeq
|
|
)
|
|
let snapshotHead = MessageWindowLogic.pagePosition(
|
|
at: response.page.snapshotHeadAt,
|
|
seq: response.page.snapshotHeadSeq
|
|
)
|
|
if until == nil {
|
|
until = snapshotHead
|
|
}
|
|
|
|
update { previous in
|
|
guard previous.syncGeneration == generation else { return previous }
|
|
return MessageWindowLogic.applyAfterPage(
|
|
previous,
|
|
responseMessages: windowMessages(response),
|
|
page: response.page,
|
|
nextAfter: nextAfter
|
|
)
|
|
}
|
|
|
|
let current = stateValue
|
|
if current.requiresLatestReset
|
|
|| current.preferLatestOnActivation
|
|
|| !response.page.hasMore
|
|
|| nextAfter == nil {
|
|
break
|
|
}
|
|
if nextAfter! <= after {
|
|
throw MessageWindowSyncError.cursorDidNotAdvance
|
|
}
|
|
after = nextAfter!
|
|
}
|
|
|
|
finishTailSync(generation: generation, warning: nil)
|
|
} catch {
|
|
guard isCurrentTailSync(generation) else { return }
|
|
finishTailSync(
|
|
generation: generation,
|
|
warning: Self.warningMessage(error, fallback: "Failed to synchronize messages")
|
|
)
|
|
}
|
|
}
|
|
|
|
private func finishTailSync(generation: Int, warning: String?) {
|
|
update { MessageWindowLogic.finishTailSync($0, generation: generation, warning: warning) }
|
|
persist()
|
|
}
|
|
|
|
// MARK: - Older pages
|
|
|
|
/// Load one older page (web `fetchOlderMessages`). On an epoch mismatch
|
|
/// the window is invalidated and a full tail sync (`ensureAfterCurrent`)
|
|
/// runs before the `stopped/epoch-reset` outcome is returned.
|
|
/// `onBeforeApply` runs synchronously inside the state transition (mirror
|
|
/// of the web's synchronous updater) — it must not call back into this
|
|
/// controller.
|
|
public func fetchOlder(
|
|
onBeforeApply: (@Sendable (Int) -> Bool)? = nil
|
|
) async -> OlderLoadOutcome {
|
|
let initial = stateValue
|
|
switch MessageWindowLogic.olderLoadPrecheck(initial) {
|
|
case .stop(let reason):
|
|
return .stopped(reason)
|
|
case .proceed(let before):
|
|
let generation = initial.olderGeneration + 1
|
|
update { MessageWindowLogic.beginOlderLoad($0, generation: generation) }
|
|
do {
|
|
let response = try await provider.messages(
|
|
sessionId: sessionId,
|
|
query: .before(
|
|
beforeAt: before.at,
|
|
beforeSeq: before.seq,
|
|
limit: MessageWindowConstants.pageSize
|
|
)
|
|
)
|
|
guard stateValue.olderGeneration == generation else {
|
|
return .stopped(.invalidated)
|
|
}
|
|
|
|
if let epoch = initial.epoch, response.page.epoch != epoch {
|
|
update { MessageWindowLogic.applyOlderEpochMismatch($0, generation: generation) }
|
|
await syncTail(ensureAfterCurrent: true)
|
|
return .stopped(.epochReset)
|
|
}
|
|
|
|
var historyVersion = 0
|
|
var addedRenderableCount = 0
|
|
var applyRejected = false
|
|
update { previous in
|
|
guard previous.olderGeneration == generation else { return previous }
|
|
let incoming = windowMessages(response)
|
|
addedRenderableCount = MessageWindowLogic.countNewRenderableMessages(
|
|
previous,
|
|
incoming: incoming
|
|
)
|
|
let nextHistoryVersion = previous.historyVersion + 1
|
|
if let onBeforeApply, !onBeforeApply(nextHistoryVersion) {
|
|
applyRejected = true
|
|
return MessageWindowLogic.rejectOlderApply(previous)
|
|
}
|
|
historyVersion = nextHistoryVersion
|
|
return MessageWindowLogic.applyOlderResponse(
|
|
previous,
|
|
responseMessages: incoming,
|
|
page: response.page,
|
|
historyVersion: nextHistoryVersion
|
|
)
|
|
}
|
|
if applyRejected || historyVersion == 0 {
|
|
return .stopped(.invalidated)
|
|
}
|
|
persist()
|
|
return .applied(
|
|
historyVersion: historyVersion,
|
|
hasMore: response.page.hasMore,
|
|
addedRenderableCount: addedRenderableCount
|
|
)
|
|
} catch {
|
|
guard stateValue.olderGeneration == generation else {
|
|
return .stopped(.invalidated)
|
|
}
|
|
update {
|
|
MessageWindowLogic.failOlderLoad(
|
|
$0,
|
|
generation: generation,
|
|
warning: Self.warningMessage(error, fallback: "Failed to load older messages")
|
|
)
|
|
}
|
|
return .failed(error)
|
|
}
|
|
}
|
|
}
|
|
|
|
public func cancelOlderLoad() {
|
|
update { MessageWindowLogic.cancelOlderLoad($0) }
|
|
}
|
|
|
|
// MARK: - SSE ingest
|
|
|
|
/// Route one message-stream SSE event into this window. The caller
|
|
/// routes per session already; the guard is defensive — a foreign
|
|
/// session's rows must never corrupt this window.
|
|
/// `messages-invalidated` clears the window and starts a fresh tail
|
|
/// sync; `scheduled-matured` re-syncs so the released row (and its
|
|
/// consumption) lands even if the `message-received` frame was missed.
|
|
public func onMessageEvent(_ event: SyncEvent) {
|
|
switch event {
|
|
case .messageReceived(_, let eventSessionId, let message) where eventSessionId == sessionId:
|
|
ingestSSEMessages([WindowMessage(wire: message)])
|
|
case .messagesConsumed(_, let eventSessionId, let localIds, let invokedAt)
|
|
where eventSessionId == sessionId:
|
|
markConsumed(localIds: localIds, invokedAt: invokedAt)
|
|
case .messagesIndeterminate(_, let eventSessionId, let localIds)
|
|
where eventSessionId == sessionId:
|
|
markIndeterminate(localIds: localIds)
|
|
case .messagesRequeued(_, let eventSessionId, let localIds)
|
|
where eventSessionId == sessionId:
|
|
markRequeued(localIds: localIds)
|
|
case .messageCancelled(_, let eventSessionId, let messageId, _) where eventSessionId == sessionId:
|
|
removeMessage(localIdOrId: messageId)
|
|
case .messagesInvalidated(_, let eventSessionId) where eventSessionId == sessionId:
|
|
clear()
|
|
Task { await self.syncTail() }
|
|
case .scheduledMatured(_, let eventSessionId) where eventSessionId == sessionId:
|
|
Task { await self.syncTail(ensureAfterCurrent: true) }
|
|
default:
|
|
break
|
|
}
|
|
}
|
|
|
|
/// SSE `message-received` ingest (web `ingestIncomingMessages`).
|
|
public func ingestSSEMessages(_ messages: [WindowMessage]) {
|
|
guard !messages.isEmpty else { return }
|
|
update { MessageWindowLogic.ingestIncoming($0, incoming: messages) }
|
|
persist()
|
|
}
|
|
|
|
/// SSE `messages-consumed` (web `markMessagesConsumed`).
|
|
public func markConsumed(localIds: [String], invokedAt: Int) {
|
|
guard !localIds.isEmpty else { return }
|
|
update { MessageWindowLogic.markConsumed($0, localIds: localIds, invokedAt: invokedAt) }
|
|
persist()
|
|
}
|
|
|
|
public func markIndeterminate(localIds: [String]) {
|
|
update { MessageWindowLogic.markIndeterminate($0, localIds: localIds) }
|
|
persist()
|
|
}
|
|
|
|
public func markRequeued(localIds: [String]) {
|
|
update { MessageWindowLogic.markRequeued($0, localIds: localIds) }
|
|
persist()
|
|
}
|
|
|
|
/// SSE `message-cancelled` / optimistic DELETE removal
|
|
/// (web `removeOptimisticMessage`).
|
|
public func removeMessage(localIdOrId: String) {
|
|
update { MessageWindowLogic.removeByLocalIdOrId($0, localId: localIdOrId) }
|
|
persist()
|
|
}
|
|
|
|
// MARK: - Optimistic sends
|
|
|
|
/// Append a pre-built optimistic row (web `appendOptimisticMessage`).
|
|
public func appendOptimistic(_ message: WindowMessage) {
|
|
update { MessageWindowLogic.appendOptimistic($0, message: message) }
|
|
persist()
|
|
}
|
|
|
|
/// Append the standard optimistic row for a send
|
|
/// (`useSendMessage.createOptimisticMessage`): status `sending` until the
|
|
/// POST settles, then ``updateStatus(localId:status:)`` to
|
|
/// `queued`/`sent`/`failed`.
|
|
public func appendOptimistic(
|
|
localId: String,
|
|
text: String,
|
|
attachments: [AttachmentMetadata]? = nil,
|
|
scheduledAt: Int? = nil,
|
|
deliveryMode: String = "queue",
|
|
createdAt: Int = Int(Date().timeIntervalSince1970 * 1000)
|
|
) {
|
|
appendOptimistic(
|
|
buildOptimisticMessage(
|
|
localId: localId,
|
|
text: text,
|
|
createdAt: createdAt,
|
|
attachments: attachments,
|
|
scheduledAt: scheduledAt,
|
|
deliveryMode: deliveryMode,
|
|
status: .sending
|
|
)
|
|
)
|
|
}
|
|
|
|
public func updateStatus(localId: String, status: MessageStatus) {
|
|
update { MessageWindowLogic.updateStatus($0, localId: localId, status: status) }
|
|
persist()
|
|
}
|
|
|
|
/// A cancel DELETE answered `{"status":"invoked"}` — too late, the agent
|
|
/// consumed the row. Remove the queued snapshot and ingest the returned
|
|
/// authoritative row as `sent` (web `useCancelQueuedMessage`).
|
|
public func applyCancelInvoked(localId: String, message: WindowMessage) {
|
|
removeMessage(localIdOrId: localId)
|
|
appendOptimistic(message.withStatus(.sent))
|
|
}
|
|
|
|
// MARK: - Queued reconciliation
|
|
|
|
public func queuedReconcileCandidateLocalIds() -> [String] {
|
|
MessageWindowLogic.queuedReconcileCandidateLocalIds(stateValue)
|
|
}
|
|
|
|
public func reconcileQueuedLocalIds(candidateLocalIds: [String], queuedLocalIds: [String]) {
|
|
update {
|
|
MessageWindowLogic.reconcileQueuedLocalIds(
|
|
$0,
|
|
candidateLocalIds: candidateLocalIds,
|
|
queuedLocalIds: queuedLocalIds
|
|
)
|
|
}
|
|
persist()
|
|
}
|
|
|
|
/// Queued-state recovery after a `resume: 'gap'` reconnect (port of
|
|
/// `web/src/lib/queued-state-reconciliation.ts`): tail-sync to the drain,
|
|
/// collect candidates, ask the hub for the verdict in ≤1000-id batches,
|
|
/// stamp invoked rows like `messages-consumed`, drop deleted candidates.
|
|
public func reconcileQueuedState() async throws {
|
|
await syncTail(ensureAfterCurrent: true)
|
|
let candidateLocalIds = queuedReconcileCandidateLocalIds()
|
|
guard !candidateLocalIds.isEmpty else { return }
|
|
var queuedLocalIds: [String] = []
|
|
var indeterminateLocalIds: [String] = []
|
|
var invokedLocalMessages: [(localId: String, invokedAt: Int)] = []
|
|
var start = 0
|
|
while start < candidateLocalIds.count {
|
|
let end = min(start + Self.queuedStateBatchSize, candidateLocalIds.count)
|
|
let batch = Array(candidateLocalIds[start..<end])
|
|
let response = try await provider.queuedState(sessionId: sessionId, localIds: batch)
|
|
queuedLocalIds += response.queuedLocalIds
|
|
indeterminateLocalIds += response.indeterminateLocalIds ?? []
|
|
invokedLocalMessages += response.invokedLocalMessages.map { ($0.localId, $0.invokedAt) }
|
|
start = end
|
|
}
|
|
var timestamps: [Int] = []
|
|
var localIdsByTimestamp: [Int: [String]] = [:]
|
|
for (localId, invokedAt) in invokedLocalMessages {
|
|
if localIdsByTimestamp[invokedAt] == nil { timestamps.append(invokedAt) }
|
|
localIdsByTimestamp[invokedAt, default: []].append(localId)
|
|
}
|
|
for invokedAt in timestamps {
|
|
markConsumed(localIds: localIdsByTimestamp[invokedAt]!, invokedAt: invokedAt)
|
|
}
|
|
markIndeterminate(localIds: indeterminateLocalIds)
|
|
reconcileQueuedLocalIds(candidateLocalIds: candidateLocalIds, queuedLocalIds: queuedLocalIds + indeterminateLocalIds)
|
|
}
|
|
|
|
// MARK: - Lifecycle transitions
|
|
|
|
/// Switch view mode (web `setMessageViewMode`); tail re-entry trims.
|
|
public func setViewMode(_ mode: MessageViewMode) {
|
|
update { MessageWindowLogic.setViewMode($0, mode: mode) }
|
|
persist()
|
|
}
|
|
|
|
/// Session (re-)activation (web `activateMessageWindow`): trims back to
|
|
/// the visible window and, with usable persisted state, requests a
|
|
/// fresh-latest trailing run from an in-flight sync.
|
|
public func activate() {
|
|
var requestedLatest = false
|
|
update { previous in
|
|
let activation = MessageWindowLogic.activate(previous)
|
|
requestedLatest = activation.requestedLatest
|
|
return activation.state
|
|
}
|
|
if requestedLatest, runningTask != nil {
|
|
trailingRequested = true
|
|
}
|
|
}
|
|
|
|
/// Web `clearMessageWindow`: forget everything but keep generations
|
|
/// poisoned so in-flight work cannot commit.
|
|
public func clear() {
|
|
controllerEpoch += 1
|
|
runningTask = nil
|
|
runningRunId = nil
|
|
runningPrefersLatest = false
|
|
trailingRequested = false
|
|
snapshots?.delete(sessionId: sessionId)
|
|
update { previous in
|
|
var next = MessageWindowLogic.createState(sessionId: sessionId)
|
|
next.syncGeneration = previous.syncGeneration + 1
|
|
next.olderGeneration = previous.olderGeneration + 1
|
|
return next
|
|
}
|
|
}
|
|
|
|
/// Seed this window from another session's (web
|
|
/// `seedMessageWindowFromSession`) — resume/reopen can hand back a new
|
|
/// session id; the old rows render instantly and `requiresLatestReset`
|
|
/// forces a fresh latest page underneath.
|
|
public func seedFrom(_ source: MessageWindowController) async {
|
|
guard source.sessionId != sessionId else { return }
|
|
controllerEpoch += 1
|
|
runningTask = nil
|
|
runningRunId = nil
|
|
runningPrefersLatest = false
|
|
trailingRequested = false
|
|
let sourceState = await source.state
|
|
update { MessageWindowLogic.seededState(source: sourceState, target: $0) }
|
|
persist()
|
|
}
|
|
|
|
// MARK: - Internals
|
|
|
|
private static let queuedStateBatchSize = 1000
|
|
|
|
/// Identity baseline of the current rows, keyed by id
|
|
/// (web `requestBaseline`; duplicate ids cannot occur post-merge, but a
|
|
/// JS `Map` keeps the last insertion — mirrored here).
|
|
private func baseline() -> [String: WindowMessage] {
|
|
Dictionary(
|
|
stateValue.messages.map { ($0.id, $0) },
|
|
uniquingKeysWith: { _, last in last }
|
|
)
|
|
}
|
|
|
|
private nonisolated func windowMessages(_ response: MessagesResponse) -> [WindowMessage] {
|
|
response.messages.map { WindowMessage(wire: $0) }
|
|
}
|
|
|
|
private func persist() {
|
|
guard let snapshots else { return }
|
|
let state = stateValue
|
|
if MessageWindowLogic.shouldPersist(state) {
|
|
snapshots.save(sessionId: sessionId, snapshot: MessageWindowLogic.toPersisted(state))
|
|
} else {
|
|
snapshots.delete(sessionId: sessionId)
|
|
}
|
|
}
|
|
|
|
private static func warningMessage(_ error: any Error, fallback: String) -> String {
|
|
(error as? LocalizedError)?.errorDescription ?? fallback
|
|
}
|
|
}
|