Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
46 changes: 39 additions & 7 deletions apps/swift-ios/App/NativeFeatureClient.swift
Original file line number Diff line number Diff line change
Expand Up @@ -3898,7 +3898,7 @@ final class NativeFeatureClient: FeatureClient, FeatureDeviceManaging,
guard !Task.isCancelled,
self?.isCurrentDetail(route, generation: streamGeneration) == true,
self?.environmentGeneration == sessionGeneration else { return }
let subscription = try await route.client.threadEvents(
let subscription = try await route.client.threadEventBatches(
threadID: route.wireID,
after: sequence,
turnLimit: supportsPagination ? Self.initialThreadUserTurnLimit : nil
Expand All @@ -3914,8 +3914,8 @@ final class NativeFeatureClient: FeatureClient, FeatureDeviceManaging,
self?.continuation.yield(.threadSync(id: route.uiID, state: .catchingUp))
}
self?.activeDetailConnectionID = subscriptionConnectionID
for try await item in subscription.events {
if case .synchronized = item {
for try await items in subscription.events {
if items.contains(where: { if case .synchronized = $0 { true } else { false } }) {
let connectionID = await route.client.currentConnectionID()
guard !Task.isCancelled, let self,
self.isCurrentDetail(route, generation: streamGeneration),
Expand All @@ -3933,8 +3933,8 @@ final class NativeFeatureClient: FeatureClient, FeatureDeviceManaging,
self.continuation.yield(.threadSync(id: route.uiID, state: .catchingUp))
self.ensureDetailCatchUpFallback(route, generation: streamGeneration)
}
self.consumeDetailStreamItem(
item, route: route, subscriptionEpoch: subscriptionEpoch
self.consumeDetailStreamBatch(
items, route: route, subscriptionEpoch: subscriptionEpoch
)
}
// A thread subscription stays open until its owner leaves.
Expand Down Expand Up @@ -4070,10 +4070,41 @@ final class NativeFeatureClient: FeatureClient, FeatureDeviceManaging,
ensureDetailCatchUpFallback(route, generation: detailStreamGeneration)
}

private func consumeDetailStreamBatch(
_ items: [ThreadStreamItem],
route: NativeThreadRoute,
subscriptionEpoch: Int
) {
// Keep snapshot replacement and page-watermark merges at their event
// positions. Ordinary replay batches share one legacy sync publication.
let previousSequence = activeThreadSequence
let needsIndividualUpdates = items.count == 1 || activeRawThread == nil
|| pendingOlderThreadPage != nil || items.contains { item in
switch item {
case .snapshot: true
case let .event(event):
event["type"]?.stringValue == "thread.reverted"
|| event["type"]?.stringValue == "thread.deleted"
case .synchronized: false
}
}
for item in items {
consumeDetailStreamItem(
item, route: route, subscriptionEpoch: subscriptionEpoch,
synchronizeLegacy: needsIndividualUpdates
)
}
if !needsIndividualUpdates, activeThreadSequence != previousSequence,
serverConfigsByEnvironmentID[route.environmentID]?.threadResumeCompletionMarker != true {
markDetailSynchronized(route)
}
}

private func consumeDetailStreamItem(
_ item: ThreadStreamItem,
route: NativeThreadRoute,
subscriptionEpoch: Int
subscriptionEpoch: Int,
synchronizeLegacy: Bool
) {
switch item {
case .synchronized:
Expand Down Expand Up @@ -4144,7 +4175,8 @@ final class NativeFeatureClient: FeatureClient, FeatureDeviceManaging,
scheduleDetailRefresh(threadID: route.uiID, client: route.client, force: true)
}
}
if serverConfigsByEnvironmentID[route.environmentID]?.threadResumeCompletionMarker != true {
if synchronizeLegacy,
serverConfigsByEnvironmentID[route.environmentID]?.threadResumeCompletionMarker != true {
markDetailSynchronized(route)
}
}
Expand Down
6 changes: 3 additions & 3 deletions apps/swift-ios/Core/T3Client.swift
Original file line number Diff line number Diff line change
Expand Up @@ -625,18 +625,18 @@ public actor T3Client {
)
}

public func threadEvents(
public func threadEventBatches(
threadID: String,
after sequence: Int? = nil,
turnLimit: Int? = nil
) async throws -> (events: AsyncThrowingStream<ThreadStreamItem, Error>, connectionID: UUID) {
) async throws -> (events: AsyncThrowingStream<[ThreadStreamItem], Error>, connectionID: UUID) {
var payload: [String: JSONValue] = [
"threadId": .string(threadID),
"requestCompletionMarker": .bool(true),
]
if let sequence { payload["afterSequence"] = .number(Double(sequence)) }
if let turnLimit { payload["turnLimit"] = .number(Double(turnLimit)) }
return try await rpc.subscribeOnCurrentConnection(
return try await rpc.subscribeBatchesOnCurrentConnection(
RPCMethod.subscribeThread.rawValue,
payload: .object(payload),
as: ThreadStreamItem.self
Expand Down
23 changes: 19 additions & 4 deletions apps/swift-ios/Core/WebSocketRPC.swift
Original file line number Diff line number Diff line change
Expand Up @@ -494,16 +494,31 @@ public actor WebSocketRPCClient {
_ tag: String,
payload: JSONValue = .object([:]),
as type: Value.Type
) async throws -> (events: AsyncThrowingStream<Value, Error>, connectionID: UUID) {
try await subscribeOnCurrentConnection {
subscribe(tag, payload: payload, reconnect: false, as: type)
}
}

public func subscribeBatchesOnCurrentConnection<Value: Decodable & Sendable>(
_ tag: String,
payload: JSONValue = .object([:]),
as type: Value.Type
) async throws -> (events: AsyncThrowingStream<[Value], Error>, connectionID: UUID) {
try await subscribeOnCurrentConnection {
subscribeBatches(tag, payload: payload, reconnect: false, as: type)
}
}

private func subscribeOnCurrentConnection<Value: Sendable>(
makeStream: () -> AsyncThrowingStream<Value, Error>
) async throws -> (events: AsyncThrowingStream<Value, Error>, connectionID: UUID) {
try Task.checkCancellation()
start()
while true {
try Task.checkCancellation()
if let id = connectionID {
return (
subscribe(tag, payload: payload, reconnect: false, as: type),
id
)
return (makeStream(), id)
}
_ = try await waitForConnection(after: nil)
}
Expand Down
14 changes: 8 additions & 6 deletions apps/swift-ios/Tests/CoreTests/WebSocketRPCRaceTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -31,11 +31,13 @@ final class WebSocketRPCRaceTests: XCTestCase {
connector: SequencedConnector(connections: [connection]),
endpointProvider: { URL(string: "wss://studio.example/ws")! }
)
let stream = await client.subscribeBatches("thread.events", reconnect: false, as: Int.self)
let subscription = try await client.subscribeBatchesOnCurrentConnection("thread.events", as: Int.self)
let connectionID = await client.currentConnectionID()
XCTAssertEqual(subscription.connectionID, connectionID)
await connection.waitUntilSubscriptionStarted()
try await connection.sendSubscriptionValues((0..<600).map { .number(Double($0)) })

var iterator = stream.makeAsyncIterator()
var iterator = subscription.events.makeAsyncIterator()
var received: [Int] = []
var batchSizes: [Int] = []
while received.count < 600, let batch = try await iterator.next() {
Expand All @@ -60,10 +62,10 @@ final class WebSocketRPCRaceTests: XCTestCase {
subscriptionBufferLimit: 2,
endpointProvider: { URL(string: "wss://studio.example/ws")! }
)
let stream = await client.subscribeBatches("thread.events", reconnect: false, as: String.self)
let subscription = try await client.subscribeBatchesOnCurrentConnection("thread.events", as: String.self)
// A unary request is not needed: the connection closes only after overflow.
await connection.waitUntilClosed()
var iterator = stream.makeAsyncIterator()
var iterator = subscription.events.makeAsyncIterator()
let buffered = try await iterator.next()
XCTAssertEqual(buffered, ["first", "second"])
do {
Expand Down Expand Up @@ -157,15 +159,15 @@ final class WebSocketRPCRaceTests: XCTestCase {
}


func testColdSubscriptionRetainsFailedSocketIdentity() async throws {
func testColdBatchSubscriptionRetainsFailedSocketIdentity() async throws {
let connection = SubscriptionTrafficConnection(sendsInvalidSubscriptionValue: true)
let connector = GatedConnector(connection: connection)
let client = WebSocketRPCClient(
connector: connector,
endpointProvider: { URL(string: "wss://studio.example/ws")! }
)
let pending = Task {
try await client.subscribeOnCurrentConnection("thread.events", as: Int.self)
try await client.subscribeBatchesOnCurrentConnection("thread.events", as: Int.self)
}
await connector.waitUntilConnectStarted()
let beforeConnection = await client.currentConnectionID()
Expand Down
Loading
Loading