Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
14 commits
Select commit Hold shift + click to select a range
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1,572 changes: 1,215 additions & 357 deletions apps/swift-ios/App/NativeFeatureClient.swift

Large diffs are not rendered by default.

23 changes: 23 additions & 0 deletions apps/swift-ios/Core/T3Client.swift
Original file line number Diff line number Diff line change
Expand Up @@ -638,6 +638,29 @@ public actor T3Client {
)
}

/// Passive shells start with a full snapshot on each socket subscription.
/// Registering its identity atomically prevents late socket results from
/// acquiring authority over a replacement connection.
public func shellEventsOnCurrentConnection() async throws -> (
events: AsyncThrowingStream<ShellStreamItem, Error>, connectionID: UUID
) {
try await rpc.subscribeOnCurrentConnection(
RPCMethod.subscribeShell.rawValue,
payload: .object(["requestCompletionMarker": .bool(true)]),
as: ShellStreamItem.self
)
}

public func shellEventBatchesOnCurrentConnection() async throws -> (
events: AsyncThrowingStream<[ShellStreamItem], Error>, connectionID: UUID
) {
try await rpc.subscribeBatchesOnCurrentConnection(
RPCMethod.subscribeShell.rawValue,
payload: .object(["requestCompletionMarker": .bool(true)]),
as: ShellStreamItem.self
)
}

public func threadEventBatches(
threadID: String,
after sequence: Int? = nil,
Expand Down
75 changes: 59 additions & 16 deletions apps/swift-ios/Features/Root/FeatureRootModel.swift
Original file line number Diff line number Diff line change
Expand Up @@ -90,6 +90,10 @@ public final class FeatureRootModel {
draftStore: draftStore
)
private var pendingSubmissionsByID: [String: FeatureQueuedSubmission] = [:]
/// Accepted messages not yet present in a server transcript, by message ID.
private var deliveredAwaitingDetail: [String: (threadID: String, message: FeatureMessage)] = [:]
/// Message IDs in each thread's latest server transcript.
private var serverMessageIDs: [String: Set<String>] = [:]
private var activeSubmissionCounts: [String: Int] = [:]
private var pendingThreadsByID: [String: FeatureThread] = [:]
private var pendingSettlementMutations: [String: PendingSettlementMutation] = [:]
Expand Down Expand Up @@ -149,6 +153,7 @@ public final class FeatureRootModel {

func applicationDidEnterBackground(at date: Date = .now) {
backgroundedAt = date
client.suspendForBackground()
}

func applicationDidBecomeActive(at date: Date = .now) async {
Expand Down Expand Up @@ -1317,6 +1322,7 @@ public final class FeatureRootModel {
}

private func removeThread(id: String) {
forgetDeliveredMessages(threadID: id)
guard let index = snapshot.threads.firstIndex(where: { $0.id == id }) else { return }
let projectID = snapshot.threads[index].projectID
snapshot.threads.remove(at: index)
Expand Down Expand Up @@ -1449,15 +1455,18 @@ public final class FeatureRootModel {
}
}

/// `isServerTranscript` is false for details built locally, such as a
/// restored outbox or a pending creation; they cannot confirm delivery.
private func store(
_ incoming: FeatureThreadDetail,
invalidatesInFlightLoad: Bool = true
invalidatesInFlightLoad: Bool = true,
isServerTranscript: Bool = true
) {
var incoming = retainingLocalAttachmentPreviews(in: incoming)
incoming.thread = retainingPendingSettlement(in: incoming.thread)
let id = incoming.thread.id
acknowledgeDeliveredMessages(incoming.messages)
let prepared = addingPendingMessages(to: incoming)
let prepared = addingPendingMessages(to: incoming, isServerTranscript: isServerTranscript)
let next = details[id].map { current in
FeatureThreadDetail(
thread: prepared.thread,
Expand All @@ -1484,7 +1493,7 @@ public final class FeatureRootModel {
incoming.thread = retainingPendingSettlement(in: incoming.thread)
let id = incoming.thread.id
acknowledgeDeliveredMessages(incoming.messages)
let next = addingPendingMessages(to: incoming)
let next = addingPendingMessages(to: incoming, isServerTranscript: true)
details[id] = next
markDetailRecentlyUsed(id)
bumpDetailLoadRevision(id: id)
Expand Down Expand Up @@ -1523,7 +1532,18 @@ public final class FeatureRootModel {
return true
}

private func removeDetail(id: String) {
private func forgetDeliveredMessages(threadID: String) {
deliveredAwaitingDetail = deliveredAwaitingDetail.filter { $0.value.threadID != threadID }
serverMessageIDs[threadID] = nil
}

private func removeDetail(id: String, forgettingDeliveredMessages: Bool = true) {
if forgettingDeliveredMessages {
forgetDeliveredMessages(threadID: id)
} else {
// The next server transcript records these again; keep the cache bounded.
serverMessageIDs[id] = nil
}
if details.removeValue(forKey: id) != nil {
detailRecency.removeAll { $0 == id }
}
Expand All @@ -1536,6 +1556,8 @@ public final class FeatureRootModel {
}

private func clearDetails() {
deliveredAwaitingDetail.removeAll()
serverMessageIDs.removeAll()
detailLoadGeneration &+= 1
detailLoadRevisions.removeAll()
storedDetailLoadRequestRevisions.removeAll()
Expand Down Expand Up @@ -1571,7 +1593,8 @@ public final class FeatureRootModel {
while details.count > Self.maximumRetainedThreadDetails,
let candidate = detailRecency.first(where: { !protected.contains($0) }) {
detailRecency.removeAll { $0 == candidate }
removeDetail(id: candidate)
// Eviction only drops the cache; accepted messages must survive a reopen.
removeDetail(id: candidate, forgettingDeliveredMessages: false)
}
}

Expand Down Expand Up @@ -1614,7 +1637,7 @@ public final class FeatureRootModel {
if snapshot.threads.contains(where: { $0.id == submission.threadID }) {
pendingSubmissionsByID[submission.id] = submission
if let detail = details[submission.threadID] {
store(addingPendingMessages(to: detail))
store(detail, isServerTranscript: false)
}
continue
}
Expand Down Expand Up @@ -1645,7 +1668,7 @@ public final class FeatureRootModel {
}
pendingSubmissionsByID[submission.id] = submission
if let detail = details[submission.threadID] {
store(addingPendingMessages(to: detail))
store(detail, isServerTranscript: false)
}
}
}
Expand Down Expand Up @@ -1708,7 +1731,7 @@ public final class FeatureRootModel {
store(FeatureThreadDetail(
thread: thread,
messages: [queuedMessage(for: submission)]
))
), isServerTranscript: false)
}

private func provider(id: String?, environmentID: String) -> FeatureProvider? {
Expand Down Expand Up @@ -1737,16 +1760,30 @@ public final class FeatureRootModel {
)
}

private func addingPendingMessages(to incoming: FeatureThreadDetail) -> FeatureThreadDetail {
private func addingPendingMessages(
to incoming: FeatureThreadDetail,
isServerTranscript: Bool
) -> FeatureThreadDetail {
let existing = Set(incoming.messages.map(\.id))
// A delivered message stays visible until a server transcript includes it.
// Its refresh no longer blocks the send, so an older read can arrive first.
if isServerTranscript {
serverMessageIDs[incoming.thread.id] = existing
for (id, delivered) in deliveredAwaitingDetail
where delivered.threadID == incoming.thread.id && existing.contains(id) {
deliveredAwaitingDetail.removeValue(forKey: id)
}
}
let delivered = deliveredAwaitingDetail.values
.filter { $0.threadID == incoming.thread.id && !existing.contains($0.message.id) }
.map(\.message)
let queued = pendingSubmissionsByID.values
.filter { $0.threadID == incoming.thread.id }
.sorted { $0.identity.createdAt < $1.identity.createdAt }
guard !queued.isEmpty else { return incoming }
.filter { $0.threadID == incoming.thread.id && !existing.contains($0.identity.messageID) }
.map(queuedMessage(for:))
guard !queued.isEmpty || !delivered.isEmpty else { return incoming }
var result = incoming
let existing = Set(result.messages.map(\.id))
result.messages.append(contentsOf: queued.lazy
.filter { !existing.contains($0.identity.messageID) }
.map(queuedMessage(for:)))
// A newer send can be delivered while an older one waits to retry; keep send order.
result.messages.append(contentsOf: (delivered + queued).sorted { $0.createdAt < $1.createdAt })
return result
}

Expand Down Expand Up @@ -1823,6 +1860,12 @@ public final class FeatureRootModel {
}
pendingCompletionSubmissionIDs.remove(submission.id)
pendingSubmissionsByID.removeValue(forKey: submission.id)
// A transcript that already includes it confirmed delivery; do not retain it.
if serverMessageIDs[submission.threadID]?.contains(submission.identity.messageID) != true {
var message = queuedMessage(for: submission)
message.state = .complete
deliveredAwaitingDetail[submission.identity.messageID] = (submission.threadID, message)
}
setAttachmentOutboxOwnership(false, for: submission)
pendingThreadsByID.removeValue(forKey: submission.threadID)
markQueuedMessageDelivered(submission)
Expand Down
3 changes: 3 additions & 0 deletions apps/swift-ios/Features/Shared/FeatureClient.swift
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ public protocol FeatureClient: AnyObject {
func backgroundSnapshot() async throws -> FeatureSnapshot
func events() -> AsyncStream<FeatureEvent>
func resumeAfterBackground(reconnect: Bool) async
func suspendForBackground()

func preuploadAttachment(
_ attachment: FeatureUploadAttachment,
Expand Down Expand Up @@ -314,6 +315,8 @@ public extension FeatureClient {

func resumeAfterBackground(reconnect: Bool) async {}

func suspendForBackground() {}

func preuploadAttachment(
_ attachment: FeatureUploadAttachment,
environmentID: String
Expand Down
Loading
Loading