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
68 changes: 68 additions & 0 deletions Flipcash/Core/Controllers/BackfillQueue.swift
Original file line number Diff line number Diff line change
@@ -0,0 +1,68 @@
//
// BackfillQueue.swift
// Flipcash
//
// Copyright © 2026 Code Inc. All rights reserved.
//

import Foundation
import FlipcashCore

/// The transcript fetches still waiting to run, in the order workers take them. A chat the user
/// opens fetches for itself and drops out of the queue, which puts it ahead of everything waiting.
struct BackfillQueue: Sendable {

/// What a conversation needs to bring its local transcript level with the server.
enum Kind: Equatable, Sendable {
/// The local cursor lags the server's head: stream the missed window from it.
case delta
/// Nothing is cached at all: fetch the newest page, which also seats the cursor to head.
case newestPage
}

struct Item: Equatable, Sendable {
let conversationID: ConversationID
let kind: Kind
}

private var items: [Item] = []
private var workers = 0

var count: Int { items.count }

func contains(_ conversationID: ConversationID) -> Bool {
items.contains { $0.conversationID == conversationID }
}

/// Appends `newItems`, skipping any conversation already queued.
mutating func enqueue(_ newItems: [Item]) {
for item in newItems where !contains(item.conversationID) {
items.append(item)
}
}

/// Drops a queued conversation, for a caller that is fetching it itself. Returns false when it isn't queued.
@discardableResult
mutating func remove(_ conversationID: ConversationID) -> Bool {
guard let index = items.firstIndex(where: { $0.conversationID == conversationID }) else { return false }
items.remove(at: index)
return true
}

/// Removes and returns the next item to run.
mutating func popNext() -> Item? {
items.isEmpty ? nil : items.removeFirst()
}

/// Takes the worker slots a caller may start to drain the queue, up to `limit` running at once.
/// Each started worker calls ``releaseWorker()`` when it finishes.
mutating func claimWorkers(limit: Int) -> Int {
let claimed = max(0, min(limit - workers, items.count))
workers += claimed
return claimed
}

mutating func releaseWorker() {
workers -= 1
}
}
261 changes: 192 additions & 69 deletions Flipcash/Core/Controllers/ConversationController.swift

Large diffs are not rendered by default.

100 changes: 100 additions & 0 deletions Flipcash/Core/Controllers/FeedReconcile.swift
Original file line number Diff line number Diff line change
@@ -0,0 +1,100 @@
//
// FeedReconcile.swift
// Flipcash
//
// Copyright © 2026 Code Inc. All rights reserved.
//

import Foundation
import FlipcashCore

/// How long the launch reconcile waits before it applies what has landed.
struct FeedReconcileTiming: Sendable {
/// The longest the reconcile waits for the slowest of the three feeds.
var feedCap: Duration
/// The longest it then waits for unread counts to resolve.
var unreadCap: Duration

/// Matches the Android launch sync.
static let launch = FeedReconcileTiming(feedCap: .seconds(2), unreadCap: .milliseconds(1500))
}

/// Waits for a fixed number of units of work to finish, or for a deadline, whichever comes first.
/// The work itself is never cancelled by the deadline: it keeps running and reports late through
/// ``isClosed``.
@MainActor
final class BoundedWait {

private var remaining: Int
private var isDone = false
private var waiter: CheckedContinuation<Void, Never>?
private var deadline: Task<Void, Never>?

/// Whether ``wait(upTo:)`` has already returned, so a unit finishing now is late.
private(set) var isClosed = false

init(expecting count: Int) {
remaining = count
isDone = count <= 0
}

/// Records one unit finished.
func signal() {
remaining -= 1
if remaining <= 0 { finish() }
}

/// Suspends until every unit has signalled or `cap` has passed, then closes the wait.
func wait(upTo cap: Duration) async {
if !isDone {
deadline = Task { [weak self] in
try? await Task.sleep(for: cap)
guard !Task.isCancelled else { return }
self?.finish()
}
await withCheckedContinuation { waiter = $0 }
}
deadline?.cancel()
deadline = nil
isClosed = true
}

private func finish() {
isDone = true
waiter?.resume()
waiter = nil
}
}

/// The three feeds' results as they land, with the wait for them. A feed that failed settles
/// without a result.
@MainActor
final class FeedArrivals {

enum Source: Hashable, Sendable {
case dm(ConversationType)
case groups
}

private var feeds: [Source: [Conversation]] = [:]
private let wait: BoundedWait

init(expecting count: Int) {
wait = BoundedWait(expecting: count)
}

/// Whether ``wait(upTo:)`` has returned, so a feed landing now is late.
var isClosed: Bool { wait.isClosed }

/// Records a feed's result, `nil` when it failed.
func arrive(_ source: Source, _ feed: [Conversation]?) {
if let feed { feeds[source] = feed }
wait.signal()
}

/// The feeds that landed once all have settled or `cap` has passed.
func wait(upTo cap: Duration) async -> [Source: [Conversation]] {
await wait.wait(upTo: cap)
return feeds
}
}
25 changes: 25 additions & 0 deletions Flipcash/Core/Controllers/KnownAuthorDirectory.swift
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@
//

import Foundation
import os
import FlipcashCore
import FlipcashStore

Expand Down Expand Up @@ -48,6 +49,9 @@ final class KnownAuthorDirectory {
/// would stop a later launch asking again once they have picked a name.
@ObservationIgnored private var nameless: [UserID: ConversationMember] = [:]

/// The launch-time read's result, held for ``hydrateIfReady()``; set off the main actor.
@ObservationIgnored private let preloaded = OSAllocatedUnfairLock<[UserID: ConversationMember]?>(initialState: nil)

/// The reload in flight, so concurrent requests collapse onto one read of the same store.
@ObservationIgnored private var reloadTask: Task<Void, Never>?

Expand All @@ -69,6 +73,27 @@ final class KnownAuthorDirectory {
)
}

/// Starts the first read off the main actor and lands it in ``snapshot`` as soon as it finishes,
/// so the Chats list takes the table in the same frame it paints the cached feed. Without it the
/// snapshot is empty until a `.task` runs after that paint, and every group preview redraws with
/// its sender prefix a moment later.
///
/// The landing is pushed from here rather than left to the screen's `onAppear`, which can fire
/// before the read has finished and then never asks again.
func preload() {
Task.detached { [weak self, read, preloaded] in
guard let members = try? read() else { return }
preloaded.withLock { $0 = members }
await self?.hydrateIfReady()
}
}

/// Lands the ``preload()`` result if it has finished and nothing newer has landed. Never waits.
func hydrateIfReady() {
guard snapshot === Snapshot.empty, let members = preloaded.withLock({ $0 }), !members.isEmpty else { return }
snapshot = Snapshot(membersByUserID: members.merging(nameless) { cached, _ in cached })
}

/// Re-reads the local cache off the main thread, landing a new ``snapshot`` only if the table
/// changed. Safe to call whenever a roster might have moved; a repeat while one is running is
/// dropped rather than queued.
Expand Down
24 changes: 22 additions & 2 deletions Flipcash/Core/Controllers/ReadWatermarkStamps.swift
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,22 @@ final class ReadWatermarkStamps {
private var stamps: [Key: UInt64] = [:]
@ObservationIgnored private var inFlight: Set<Key> = []
@ObservationIgnored private var missing: Set<Key> = []
@ObservationIgnored private var isHolding = false
@ObservationIgnored private var held: [Key: UInt64] = [:]

/// Starts keeping fetched stamps out of view, so a batch of lookups shows up as one change at
/// ``release()`` instead of one per row.
func hold() {
isHolding = true
}

/// Publishes every stamp fetched since ``hold()`` at once; later fetches publish as they land.
func release() {
isHolding = false
guard !held.isEmpty else { return }
stamps.merge(held) { _, new in new }
held = [:]
}

/// The fetched stamp for `key`, or nil when it hasn't been fetched.
func stamp(for key: Key) -> UInt64? {
Expand All @@ -35,11 +51,15 @@ final class ReadWatermarkStamps {
/// Fetches `key`'s message unless it's already stamped, in flight, or known missing. A failed
/// fetch leaves the key unresolved, so the next call tries again.
func resolve(_ key: Key, fetch: () async throws -> ConversationMessage?) async throws {
guard stamps[key] == nil, !inFlight.contains(key), !missing.contains(key) else { return }
guard stamps[key] == nil, held[key] == nil, !inFlight.contains(key), !missing.contains(key) else { return }
inFlight.insert(key)
defer { inFlight.remove(key) }
if let message = try await fetch() {
stamps[key] = message.unreadSeq
if isHolding {
held[key] = message.unreadSeq
} else {
stamps[key] = message.unreadSeq
}
} else {
missing.insert(key)
}
Expand Down
4 changes: 3 additions & 1 deletion Flipcash/Core/Screens/Conversation/ConversationScreen.swift
Original file line number Diff line number Diff line change
Expand Up @@ -526,7 +526,9 @@ struct ConversationScreen: View {
// allowed to see —
// which is why this reads `withholdsTranscript` and not `obscuresTranscript`: a chat
// whose rules haven't landed yet is blurred without yet refusing anything.
showsGatePlaceholder: gate.withholdsTranscript && (coordinator?.items.isEmpty ?? true),
// Also stands in while an empty transcript's first load is out, so a chat with nothing
// cached opens on the placeholder and paints once with its history.
showsGatePlaceholder: (gate.withholdsTranscript || (chatExists && !didInitialRead)) && (coordinator?.items.isEmpty ?? true),
gateMintName: gateMintName,
onGateAddFunds: addFunds,
onGateJoin: joinChat,
Expand Down
3 changes: 3 additions & 0 deletions Flipcash/Core/Screens/Main/Tips/TipConversationsScreen.swift
Original file line number Diff line number Diff line change
Expand Up @@ -63,6 +63,9 @@ struct TipConversationsScreen: View {
// Takes a finished cache read before the first frame, instead of drawing empty until launch
// work lets the controller's own hydration run.
.onAppear {
// Names first: the cache hydration below invalidates the rows, and they must not
// render once without them.
sessionContainer.knownAuthors.hydrateIfReady()
conversationController.hydrateIfReady()
}
// Every counterpart, not just the rows on screen. A row's own `.task` fires when the row is
Expand Down
1 change: 1 addition & 0 deletions Flipcash/Core/Session/SessionAuthenticator.swift
Original file line number Diff line number Diff line change
Expand Up @@ -727,6 +727,7 @@ final class SessionContainer {

self.profileAvatars = ProfileAvatarStore(flipClient: flipClient, owner: session.ownerKeyPair)
self.knownAuthors = KnownAuthorDirectory(database: database, flipClient: flipClient, owner: session.ownerKeyPair)
knownAuthors.preload()
}

fileprivate func injectingEnvironment<SomeView>(into view: SomeView) -> some View where SomeView: View {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -59,7 +59,7 @@ public struct ConversationStore: Sendable {

/// Replace the feed from a paged load, sorted most-recent-activity first.
public mutating func setFeed(_ conversations: [Conversation]) {
let merged = conversations.map { keepingSelfReadPointer(seated($0)) }
let merged = conversations.map { keepingNewerActivity(keepingSelfReadPointer(seated($0))) }
self.conversations = merged.sorted { $0.lastActivity > $1.lastActivity }
}

Expand All @@ -69,7 +69,7 @@ public struct ConversationStore: Sendable {
public mutating func setFeed(_ conversations: [Conversation], type: ConversationType) {
// Only the incoming rows are merged: re-merging the other types against themselves would
// read the store's own copy as the server acknowledging an unsynced READ pointer.
let incoming = conversations.filter { $0.type == type }.map { keepingSelfReadPointer(seated($0)) }
let incoming = conversations.filter { $0.type == type }.map { keepingNewerActivity(keepingSelfReadPointer(seated($0))) }
self.conversations = (self.conversations.filter { $0.type != type } + incoming)
.sorted { $0.lastActivity > $1.lastActivity }
}
Expand Down Expand Up @@ -306,9 +306,10 @@ public struct ConversationStore: Sendable {
}

/// Bump a conversation's last activity and re-sort the feed. No-ops for a conversation not in the
/// feed.
/// feed, and for a date that is not later than the one it holds.
public mutating func advanceLastActivity(to date: Date, in conversationID: ConversationID) {
guard let index = conversations.firstIndex(where: { $0.id == conversationID }) else { return }
guard let index = conversations.firstIndex(where: { $0.id == conversationID }),
date > conversations[index].lastActivity else { return }
conversations[index].lastActivity = date
sort()
}
Expand Down Expand Up @@ -559,7 +560,7 @@ public struct ConversationStore: Sendable {
}

private mutating func upsert(_ conversation: Conversation) {
var conversation = keepingSelfReadPointer(seated(conversation))
var conversation = keepingNewerActivity(keepingSelfReadPointer(seated(conversation)))
if let index = conversations.firstIndex(where: { $0.id == conversation.id }) {
if conversation.type == .group {
conversation.members = mergedMembers(conversation.members, over: conversations[index].members)
Expand Down Expand Up @@ -606,6 +607,22 @@ public struct ConversationStore: Sendable {
return conversation
}

/// A server copy with the stored row's `lastActivity` and `lastMessage` kept when they are newer.
/// A feed or metadata response can be older than a stream event that landed while it was in
/// flight; taking it as-is would roll the row back, and a launch would show the list move twice.
private func keepingNewerActivity(_ conversation: Conversation) -> Conversation {
guard let stored = conversations.first(where: { $0.id == conversation.id }) else { return conversation }
var conversation = conversation
if stored.lastActivity > conversation.lastActivity {
conversation.lastActivity = stored.lastActivity
}
if let held = stored.lastMessage, let incoming = conversation.lastMessage,
(held.id, held.eventSequence) > (incoming.id, incoming.eventSequence) {
conversation.lastMessage = held
}
return conversation
}

/// `conversation` with the self READ pointer raised to the one the store already holds, when the
/// server's copy is behind it. The server never lowers a pointer, so a lower copy means an advance
/// it hasn't taken yet; that advance is kept and recorded as unsynced for the caller to re-send.
Expand Down
10 changes: 8 additions & 2 deletions FlipcashCore/Sources/FlipcashStore/Database+Conversations.swift
Original file line number Diff line number Diff line change
Expand Up @@ -85,11 +85,17 @@ nonisolated extension Database {
let rows = try reader.prepareRowIterator(c.table.order(c.lastActivity.desc))
return try rows.map { row in
let id = row[c.id]
let lastMessage = try latestMessage(conversationId: id)
// The row's activity only moves when the app writes the conversation. The notification
// extension writes message rows alone while the app is closed, so a chat that received a
// push would otherwise launch with its new preview at its old position. The live stream
// advances activity to a message's date the same way.
let storedActivity = Date(timeIntervalSinceReferenceDate: row[c.lastActivity])
return Conversation(
id: ConversationID(data: id),
members: membersByConversation[id] ?? [],
lastMessage: try latestMessage(conversationId: id),
lastActivity: Date(timeIntervalSinceReferenceDate: row[c.lastActivity]),
lastMessage: lastMessage,
lastActivity: max(storedActivity, lastMessage?.date ?? storedActivity),
type: ConversationType(rawValue: row[c.type]) ?? .contactDm,
isHidden: row[c.isHidden],
title: row[c.title],
Expand Down
Loading
Loading