From b61a3c95bcc3d326fe09b165e21ca40ca741991b Mon Sep 17 00:00:00 2001 From: Vincent Koc Date: Sat, 1 Aug 2026 08:06:51 +0800 Subject: [PATCH] feat(rpc): add rowid-forward message pagination --- CHANGELOG.md | 3 + README.md | 2 +- Sources/IMsgCore/MessageStore+Messages.swift | 78 ++++++ .../IMsgCore/MessageStore+URLPreviews.swift | 122 +++++++++ Sources/IMsgCore/MessagesAfterPage.swift | 11 + .../RPCServer+MessagesAfterHandlers.swift | 107 ++++++++ Sources/imsg/RPCServer.swift | 6 + .../MessageStoreMessagesAfterPageTests.swift | 142 +++++++++++ Tests/imsgTests/RPCMessagesAfterTests.swift | 235 ++++++++++++++++++ docs/rpc.md | 37 +++ 10 files changed, 742 insertions(+), 1 deletion(-) create mode 100644 Sources/IMsgCore/MessagesAfterPage.swift create mode 100644 Sources/imsg/RPCServer+MessagesAfterHandlers.swift create mode 100644 Tests/IMsgCoreTests/MessageStoreMessagesAfterPageTests.swift create mode 100644 Tests/imsgTests/RPCMessagesAfterTests.swift diff --git a/CHANGELOG.md b/CHANGELOG.md index 67ac9245..c63c9aa7 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,9 @@ ## 0.13.5 - Unreleased +### JSON-RPC +- feat: add bounded `messages.after` pagination with a stable ROWID cursor and authoritative continuation signal (#200). + ## 0.13.4 - 2026-07-27 ### Highlights diff --git a/README.md b/README.md index 4eca570a..e50b290c 100644 --- a/README.md +++ b/README.md @@ -205,7 +205,7 @@ specific outgoing Apple ID phone number or inline reply target. It is intended for agents and long-running integrations that want a single process for chats, history, send, and watch. -Read methods: `chats.list`, `messages.history`, `messages.stats`, `messages.scheduled`, `watch.subscribe`, +Read methods: `chats.list`, `messages.history`, `messages.after`, `messages.stats`, `messages.scheduled`, `watch.subscribe`, `watch.unsubscribe`, `message.send_status`. Mutating: `send`, `poll.send`. Bridge introspection: `handles.check`. See [docs/rpc.md](docs/rpc.md) for request and response shapes. diff --git a/Sources/IMsgCore/MessageStore+Messages.swift b/Sources/IMsgCore/MessageStore+Messages.swift index 7bcc17ef..5bc0da0c 100644 --- a/Sources/IMsgCore/MessageStore+Messages.swift +++ b/Sources/IMsgCore/MessageStore+Messages.swift @@ -121,12 +121,90 @@ extension MessageStore { } } + public func messagesAfterPage( + afterRowID: Int64, + chatID: Int64?, + limit: Int, + includeReactions: Bool = false + ) throws -> MessagesAfterPage { + guard limit > 0 else { + return MessagesAfterPage(messages: [], nextRowID: afterRowID, hasMore: false) + } + + return try withConnection { db in + var physicalLimit = limit == Int.max ? limit : limit + 1 + + while true { + let query = MessagesAfterQuery( + store: self, + afterRowID: MessageID(rawValue: afterRowID), + chatID: chatID.map { ChatID(rawValue: $0) }, + limit: physicalLimit, + includeReactions: includeReactions + ) + var physicalMessages: [Message] = [] + var parentCache: ReplyParentCache = [:] + var pollOptionCache = PollOptionTextCache() + let rows = try db.prepareRowIterator(query.sql, bindings: query.bindings) + while let row = try rows.failableNext() { + let decoded = try decodeMessageRow( + row, + columns: query.selection.columns, + fallbackChatID: query.fallbackChatID + ) + physicalMessages.append( + try message( + from: decoded, + db, + parentCache: &parentCache, + pollOptionCache: &pollOptionCache + )) + } + + let visibleMessages = try pageVisibleMessages(physicalMessages, db: db) + if visibleMessages.count > limit { + let overflowRowID = visibleMessages[limit].rowID + let consumed = physicalMessages.prefix { $0.rowID < overflowRowID } + let pageMessages = try pageVisibleMessages(Array(consumed), db: db) + let nextRowID = consumed.last?.rowID ?? afterRowID + return MessagesAfterPage( + messages: try enrichMessagesWithTrailingURLPreviews( + pageMessages, + afterRowID: nextRowID, + db: db + ), + nextRowID: nextRowID, + hasMore: true + ) + } + if physicalMessages.count < physicalLimit || physicalLimit == Int.max { + return MessagesAfterPage( + messages: visibleMessages, + nextRowID: physicalMessages.last?.rowID ?? afterRowID, + hasMore: false + ) + } + guard let nextLimit = nextHistoryPhysicalLimit(after: physicalLimit) else { + return MessagesAfterPage( + messages: visibleMessages, + nextRowID: physicalMessages.last?.rowID ?? afterRowID, + hasMore: false + ) + } + physicalLimit = nextLimit + } + } + } + func messagesAfterBatch( afterRowID: Int64, chatID: Int64?, limit: Int, includeReactions: Bool ) throws -> MessagesAfterBatch { + guard limit > 0 else { + return MessagesAfterBatch(messages: [], maxScannedRowID: afterRowID) + } let query = MessagesAfterQuery( store: self, afterRowID: MessageID(rawValue: afterRowID), diff --git a/Sources/IMsgCore/MessageStore+URLPreviews.swift b/Sources/IMsgCore/MessageStore+URLPreviews.swift index c51df019..c58640b9 100644 --- a/Sources/IMsgCore/MessageStore+URLPreviews.swift +++ b/Sources/IMsgCore/MessageStore+URLPreviews.swift @@ -1,4 +1,5 @@ import Foundation +import SQLite enum URLPreviewCoalescingFallback { case suppress @@ -96,6 +97,127 @@ extension MessageStore { message.balloonBundleID == MessageStore.urlPreviewBalloonBundleID } + func pageVisibleMessages(_ messages: [Message], db: Connection) throws -> [Message] { + try coalesceURLPreviewMessages( + messages, + validateExistingCoalescence: { text, preview in + try self.precedingTextMessageForURLPreview(preview, db: db)?.rowID == text.rowID + }, + fallbackForUnmatchedPreview: { preview in + guard try self.precedingTextMessageForURLPreview(preview, db: db) != nil else { + return nil + } + return .suppress + } + ) + } + + func enrichMessagesWithTrailingURLPreviews( + _ messages: [Message], + afterRowID: Int64, + db: Connection + ) throws -> [Message] { + guard schema.hasBalloonBundleIDColumn, !messages.isEmpty else { return messages } + + var enriched = messages + let indexByRowID = Dictionary( + uniqueKeysWithValues: messages.enumerated().map { ($0.element.rowID, $0.offset) } + ) + var lastBaseByChatID: [Int64: Message] = [:] + for message in messages where !isURLPreviewBalloon(message) && !message.isReaction { + lastBaseByChatID[message.chatID] = message + } + let pageBases = lastBaseByChatID.values.sorted { $0.rowID < $1.rowID } + guard !pageBases.isEmpty else { return messages } + + let reactionFilter = + schema.hasReactionColumns + ? """ + AND ( + next.associated_message_type IS NULL + OR next.associated_message_type < 2000 + OR next.associated_message_type > 3006 + ) + """ + : "" + let selection = MessageRowSelection(store: self, includeChatID: true) + + // Keep each VALUES block below SQLite's historical 999-variable limit. + for start in stride(from: 0, to: pageBases.count, by: 400) { + let end = min(start + 400, pageBases.count) + let batch = pageBases[start.. ? + AND next_cmj.chat_id = page_base.chat_id + AND COALESCE(next.balloon_bundle_id, '') <> ? + \(reactionFilter) + ORDER BY next_cmj.message_id ASC + LIMIT 1 + ) AS boundary_rowid + FROM page_base + ) + SELECT \(selection.selectList), + preview_window.parent_rowid AS preview_parent_rowid + FROM preview_window + JOIN chat_message_join cmj ON cmj.chat_id = preview_window.chat_id + JOIN message m ON m.ROWID = cmj.message_id + LEFT JOIN handle h ON m.handle_id = h.ROWID + WHERE m.ROWID > ? + AND (preview_window.boundary_rowid IS NULL OR m.ROWID < preview_window.boundary_rowid) + AND m.balloon_bundle_id = ? + ORDER BY m.ROWID ASC + """ + var bindings: [Binding?] = [] + for message in batch { + bindings.append(message.rowID) + bindings.append(message.chatID) + } + bindings.append(afterRowID) + bindings.append(MessageStore.urlPreviewBalloonBundleID) + bindings.append(afterRowID) + bindings.append(MessageStore.urlPreviewBalloonBundleID) + + var parentCache: ReplyParentCache = [:] + var pollOptionCache = PollOptionTextCache() + let rows = try db.prepareRowIterator(sql, bindings: bindings) + while let row = try rows.failableNext() { + let parentRowID = try int64Value(row, "preview_parent_rowid") + guard + let parentRowID, + let index = indexByRowID[parentRowID] + else { + continue + } + let decoded = try decodeMessageRow( + row, + columns: selection.columns, + fallbackChatID: enriched[index].chatID + ) + let preview = try message( + from: decoded, + db, + parentCache: &parentCache, + pollOptionCache: &pollOptionCache + ) + guard + try precedingTextMessageForURLPreview(preview, db: db)?.rowID == parentRowID + else { + continue + } + enriched[index] = enriched[index].withURLPreview(urlPreviewMetadata(from: preview)) + } + } + return enriched + } + private func previousMessageInSameChat( _ chronological: [(offset: Int, element: Message)], before position: Int, diff --git a/Sources/IMsgCore/MessagesAfterPage.swift b/Sources/IMsgCore/MessagesAfterPage.swift new file mode 100644 index 00000000..603a05bb --- /dev/null +++ b/Sources/IMsgCore/MessagesAfterPage.swift @@ -0,0 +1,11 @@ +public struct MessagesAfterPage: Sendable, Equatable { + public let messages: [Message] + public let nextRowID: Int64 + public let hasMore: Bool + + public init(messages: [Message], nextRowID: Int64, hasMore: Bool) { + self.messages = messages + self.nextRowID = nextRowID + self.hasMore = hasMore + } +} diff --git a/Sources/imsg/RPCServer+MessagesAfterHandlers.swift b/Sources/imsg/RPCServer+MessagesAfterHandlers.swift new file mode 100644 index 00000000..da9e8594 --- /dev/null +++ b/Sources/imsg/RPCServer+MessagesAfterHandlers.swift @@ -0,0 +1,107 @@ +import CoreFoundation +import Foundation +import IMsgCore + +extension RPCServer { + func handleMessagesAfter(id: Any?, params: [String: Any]) async throws { + let supportedParams: Set = [ + "since_rowid", + "chat_id", + "limit", + "attachments", + "convert_attachments", + "include_reactions", + ] + if let unknown = params.keys.filter({ !supportedParams.contains($0) }).sorted().first { + throw RPCError.invalidParams("unknown messages.after param: \(unknown)") + } + + guard let sinceRowID = strictMessagesAfterInt64(params["since_rowid"]), sinceRowID >= 0 else { + throw RPCError.invalidParams("since_rowid must be a non-negative integer") + } + + let chatID: Int64? + if let rawChatID = params["chat_id"] { + guard let parsed = strictMessagesAfterInt64(rawChatID), parsed > 0 else { + throw RPCError.invalidParams("chat_id must be a positive integer") + } + chatID = parsed + } else { + chatID = nil + } + + let limit: Int + if let rawLimit = params["limit"] { + guard let parsed = strictMessagesAfterInt(rawLimit), (1...500).contains(parsed) else { + throw RPCError.invalidParams("limit must be an integer between 1 and 500") + } + limit = parsed + } else { + limit = 100 + } + + let includeAttachments = try strictMessagesAfterBool( + params["attachments"], + name: "attachments" + ) + let attachmentOptions = AttachmentQueryOptions( + convertUnsupported: try strictMessagesAfterBool( + params["convert_attachments"], + name: "convert_attachments" + )) + let page = try store.messagesAfterPage( + afterRowID: sinceRowID, + chatID: chatID, + limit: limit, + includeReactions: try strictMessagesAfterBool( + params["include_reactions"], + name: "include_reactions" + ) + ) + let reactionsByMessageID = try store.reactions(for: page.messages) + var payloads: [[String: Any]] = [] + payloads.reserveCapacity(page.messages.count) + for message in page.messages { + payloads.append( + try await buildMessagePayload( + store: store, + cache: cache, + message: message, + includeAttachments: includeAttachments, + includeReactions: true, + prefetchedReactions: reactionsByMessageID[message.rowID] ?? [], + attachmentOptions: attachmentOptions, + contactResolver: contactResolver + )) + } + + respond( + id: id, + result: [ + "messages": payloads, + "next_rowid": page.nextRowID, + "has_more": page.hasMore, + ] + ) + } +} + +private func strictMessagesAfterInt64(_ value: Any?) -> Int64? { + guard let number = value as? NSNumber else { return nil } + guard CFGetTypeID(number) != CFBooleanGetTypeID() else { return nil } + return Int64(number.stringValue) +} + +private func strictMessagesAfterInt(_ value: Any?) -> Int? { + guard let number = value as? NSNumber else { return nil } + guard CFGetTypeID(number) != CFBooleanGetTypeID() else { return nil } + return Int(number.stringValue) +} + +private func strictMessagesAfterBool(_ value: Any?, name: String) throws -> Bool { + guard let value else { return false } + guard let number = value as? NSNumber, CFGetTypeID(number) == CFBooleanGetTypeID() else { + throw RPCError.invalidParams("\(name) must be a boolean") + } + return number.boolValue +} diff --git a/Sources/imsg/RPCServer.swift b/Sources/imsg/RPCServer.swift index cb3ee6e6..ca860267 100644 --- a/Sources/imsg/RPCServer.swift +++ b/Sources/imsg/RPCServer.swift @@ -34,6 +34,7 @@ let kSupportedRPCMethods: [String] = [ "chats.markUnread", "messages.stats", "messages.history", + "messages.after", "watch.subscribe", "watch.unsubscribe", "send", @@ -168,6 +169,11 @@ final class RPCServer { try await handleMessagesStats(id: id, params: params) case "messages.history": try await handleMessagesHistory(id: id, params: params) + case "messages.after": + guard request.paramsAreNamed else { + throw RPCError.invalidParams("messages.after params must be an object") + } + try await handleMessagesAfter(id: id, params: params) case "watch.subscribe": try await handleWatchSubscribe(id: id, params: params) case "watch.unsubscribe": diff --git a/Tests/IMsgCoreTests/MessageStoreMessagesAfterPageTests.swift b/Tests/IMsgCoreTests/MessageStoreMessagesAfterPageTests.swift new file mode 100644 index 00000000..d650caec --- /dev/null +++ b/Tests/IMsgCoreTests/MessageStoreMessagesAfterPageTests.swift @@ -0,0 +1,142 @@ +import Foundation +import SQLite +import Testing + +@testable import IMsgCore + +@Test +func messagesAfterPageAdvancesPastCoalescedPhysicalRows() throws { + let db = try makeURLPreviewTestDB() + let now = Date() + try db.run("INSERT INTO handle(ROWID, id) VALUES (1, '+123')") + try insertURLPreviewTestMessage( + db, + rowID: 1, + text: "Dump https://example.com", + guid: "text-guid", + date: now + ) + try insertURLPreviewTestMessage( + db, + rowID: 2, + text: "https://example.com", + guid: "preview-guid", + balloonBundleID: MessageStore.urlPreviewBalloonBundleID, + date: now.addingTimeInterval(1) + ) + + let store = try MessageStore(connection: db, path: ":memory:") + let page = try store.messagesAfterPage(afterRowID: 0, chatID: 1, limit: 1) + + #expect(page.messages.map(\.rowID) == [1]) + #expect(page.messages.first?.urlPreview?.rowID == 2) + #expect(page.nextRowID == 2) + #expect(page.hasMore == false) +} + +@Test +func messagesAfterPageCanAdvanceAnEmptySuppressedPage() throws { + let db = try makeURLPreviewTestDB() + let now = Date() + try db.run("INSERT INTO handle(ROWID, id) VALUES (1, '+123')") + try insertURLPreviewTestMessage( + db, + rowID: 1, + text: "Dump https://example.com", + guid: "text-guid", + date: now + ) + try insertURLPreviewTestMessage( + db, + rowID: 2, + text: "https://example.com", + guid: "preview-guid", + balloonBundleID: MessageStore.urlPreviewBalloonBundleID, + date: now.addingTimeInterval(1) + ) + let store = try MessageStore(connection: db, path: ":memory:") + let previewPage = try store.messagesAfterPage(afterRowID: 1, chatID: 1, limit: 1) + #expect(previewPage.messages.isEmpty) + #expect(previewPage.nextRowID == 2) + #expect(previewPage.hasMore == false) +} + +@Test +func messagesAfterPageFillsPastSuppressedPreviewRows() throws { + let db = try makeURLPreviewTestDB() + let now = Date() + try db.run("INSERT INTO handle(ROWID, id) VALUES (1, '+123')") + try insertURLPreviewTestMessage( + db, + rowID: 1, + text: "Dump https://example.com", + guid: "text-guid", + date: now + ) + try insertURLPreviewTestMessage( + db, + rowID: 2, + text: "https://example.com", + guid: "preview-guid", + balloonBundleID: MessageStore.urlPreviewBalloonBundleID, + date: now.addingTimeInterval(1) + ) + try insertURLPreviewTestMessage( + db, + rowID: 3, + text: "after", + guid: "after-guid", + date: now.addingTimeInterval(2) + ) + + let store = try MessageStore(connection: db, path: ":memory:") + let page = try store.messagesAfterPage(afterRowID: 1, chatID: 1, limit: 1) + + #expect(page.messages.map(\.rowID) == [3]) + #expect(page.nextRowID == 3) + #expect(page.hasMore == false) +} + +@Test +func messagesAfterPageResolvesPreviewPastAnInterleavedChatRow() throws { + let db = try makeURLPreviewTestDB() + let now = Date() + try db.run("INSERT INTO handle(ROWID, id) VALUES (1, '+123')") + try insertURLPreviewTestMessage( + db, + rowID: 1, + chatID: 1, + text: "Dump https://example.com", + guid: "text-guid", + date: now + ) + try insertURLPreviewTestMessage( + db, + rowID: 2, + chatID: 2, + text: "other chat", + guid: "other-guid", + date: now.addingTimeInterval(1) + ) + try insertURLPreviewTestMessage( + db, + rowID: 3, + chatID: 1, + text: "https://example.com", + guid: "preview-guid", + balloonBundleID: MessageStore.urlPreviewBalloonBundleID, + date: now.addingTimeInterval(2) + ) + + let store = try MessageStore(connection: db, path: ":memory:") + let first = try store.messagesAfterPage(afterRowID: 0, chatID: nil, limit: 1) + #expect(first.messages.map(\.rowID) == [1]) + #expect(first.messages.first?.urlPreview?.rowID == 3) + #expect(first.nextRowID == 1) + #expect(first.hasMore == true) + + let second = try store.messagesAfterPage(afterRowID: first.nextRowID, chatID: nil, limit: 1) + #expect(second.messages.map(\.rowID) == [2]) + #expect(second.nextRowID == 3) + #expect(second.hasMore == false) +} diff --git a/Tests/imsgTests/RPCMessagesAfterTests.swift b/Tests/imsgTests/RPCMessagesAfterTests.swift new file mode 100644 index 00000000..7b0eb49b --- /dev/null +++ b/Tests/imsgTests/RPCMessagesAfterTests.swift @@ -0,0 +1,235 @@ +import Foundation +import SQLite +import Testing + +@testable import IMsgCore +@testable import imsg + +@Test +func rpcMessagesAfterPagesMoreThanFiveHundredEqualTimestampRows() async throws { + let timestamp = Date(timeIntervalSince1970: 1_700_000_000) + let store = try makeMessagesAfterStore( + rows: (1...502).map { (Int64($0), Int64(1), timestamp) } + ) + let output = TestRPCOutput() + let server = RPCServer(store: store, verbose: false, output: output) + + await server.handleLineForTesting( + #"{"jsonrpc":"2.0","id":"first","method":"messages.after","params":{"since_rowid":0,"limit":500}}"# + ) + + let first = try #require(output.responses.first?["result"] as? [String: Any]) + let firstMessages = try #require(first["messages"] as? [[String: Any]]) + #expect(firstMessages.compactMap { testInt64($0["id"]) } == (1...500).map(Int64.init)) + #expect(testInt64(first["next_rowid"]) == 500) + #expect(first["has_more"] as? Bool == true) + + await server.handleLineForTesting( + #"{"jsonrpc":"2.0","id":"second","method":"messages.after","params":{"since_rowid":500,"limit":500}}"# + ) + + let second = try #require(output.responses.last?["result"] as? [String: Any]) + let secondMessages = try #require(second["messages"] as? [[String: Any]]) + #expect(secondMessages.compactMap { testInt64($0["id"]) } == [501, 502]) + #expect(testInt64(second["next_rowid"]) == 502) + #expect(second["has_more"] as? Bool == false) +} + +@Test +func rpcMessagesAfterFiltersInterleavedChatsWithoutGuessingFromGlobalRowIDs() async throws { + let timestamp = Date(timeIntervalSince1970: 1_700_000_000) + let store = try makeMessagesAfterStore( + rows: (1...6).map { rowID in + (Int64(rowID), rowID.isMultiple(of: 2) ? Int64(2) : Int64(1), timestamp) + } + ) + let output = TestRPCOutput() + let server = RPCServer(store: store, verbose: false, output: output) + + await server.handleLineForTesting( + #"{"jsonrpc":"2.0","id":"first","method":"messages.after","params":{"since_rowid":0,"chat_id":2,"limit":2}}"# + ) + + let first = try #require(output.responses.first?["result"] as? [String: Any]) + let firstMessages = try #require(first["messages"] as? [[String: Any]]) + #expect(firstMessages.compactMap { testInt64($0["id"]) } == [2, 4]) + #expect(testInt64(first["next_rowid"]) == 4) + #expect(first["has_more"] as? Bool == true) + + await server.handleLineForTesting( + #"{"jsonrpc":"2.0","id":"second","method":"messages.after","params":{"since_rowid":4,"chat_id":2,"limit":2}}"# + ) + + let second = try #require(output.responses.last?["result"] as? [String: Any]) + let secondMessages = try #require(second["messages"] as? [[String: Any]]) + #expect(secondMessages.compactMap { testInt64($0["id"]) } == [6]) + #expect(testInt64(second["next_rowid"]) == 6) + #expect(second["has_more"] as? Bool == false) +} + +@Test +func rpcMessagesAfterIncludesAttachmentMetadataWhenRequested() async throws { + let store = try CommandTestDatabase.makeStoreForRPCWithAttachment( + filename: "/tmp/example.jpg", + transferName: "example.jpg", + uti: "public.jpeg", + mimeType: "image/jpeg" + ) + let output = TestRPCOutput() + let server = RPCServer(store: store, verbose: false, output: output) + + await server.handleLineForTesting( + #"{"jsonrpc":"2.0","id":"attachments","method":"messages.after","params":{"since_rowid":0,"attachments":true}}"# + ) + + let result = try #require(output.responses.first?["result"] as? [String: Any]) + let messages = try #require(result["messages"] as? [[String: Any]]) + let attachments = try #require(messages.first?["attachments"] as? [[String: Any]]) + #expect(attachments.first?["transfer_name"] as? String == "example.jpg") +} + +@Test +func rpcMessagesAfterCanPageReactionEventsWithoutCursorLoss() async throws { + let store = try makeMessagesAfterReactionStore() + let output = TestRPCOutput() + let server = RPCServer(store: store, verbose: false, output: output) + + await server.handleLineForTesting( + #"{"jsonrpc":"2.0","id":"reaction","method":"messages.after","params":{"since_rowid":5,"limit":1,"include_reactions":true}}"# + ) + + let first = try #require(output.responses.first?["result"] as? [String: Any]) + let firstMessages = try #require(first["messages"] as? [[String: Any]]) + #expect(firstMessages.compactMap { testInt64($0["id"]) } == [6]) + #expect(firstMessages.first?["is_reaction"] as? Bool == true) + #expect(testInt64(first["next_rowid"]) == 6) + #expect(first["has_more"] as? Bool == true) + + await server.handleLineForTesting( + #"{"jsonrpc":"2.0","id":"message","method":"messages.after","params":{"since_rowid":6,"limit":1,"include_reactions":true}}"# + ) + + let second = try #require(output.responses.last?["result"] as? [String: Any]) + let secondMessages = try #require(second["messages"] as? [[String: Any]]) + #expect(secondMessages.compactMap { testInt64($0["id"]) } == [7]) + #expect(secondMessages.first?["is_reaction"] == nil) + #expect(testInt64(second["next_rowid"]) == 7) + #expect(second["has_more"] as? Bool == false) +} + +@Test +func rpcMessagesAfterReturnsStableEmptyPageAndAdvertisesCapability() async throws { + let store = try CommandTestDatabase.makeStoreForRPC() + let output = TestRPCOutput() + let server = RPCServer(store: store, verbose: false, output: output) + + await server.handleLineForTesting( + #"{"jsonrpc":"2.0","id":"empty","method":"messages.after","params":{"since_rowid":999}}"# + ) + + let result = try #require(output.responses.first?["result"] as? [String: Any]) + #expect((result["messages"] as? [[String: Any]])?.isEmpty == true) + #expect(testInt64(result["next_rowid"]) == 999) + #expect(result["has_more"] as? Bool == false) + #expect(kSupportedRPCMethods.contains("messages.after")) +} + +@Test( + arguments: [ + #"{"jsonrpc":"2.0","id":1,"method":"messages.after","params":{}}"#, + #"{"jsonrpc":"2.0","id":1,"method":"messages.after","params":{"since_rowid":-1}}"#, + #"{"jsonrpc":"2.0","id":1,"method":"messages.after","params":{"since_rowid":true}}"#, + #"{"jsonrpc":"2.0","id":1,"method":"messages.after","params":{"since_rowid":1.5}}"#, + #"{"jsonrpc":"2.0","id":1,"method":"messages.after","params":{"since_rowid":0,"chat_id":0}}"#, + #"{"jsonrpc":"2.0","id":1,"method":"messages.after","params":{"since_rowid":0,"limit":0}}"#, + #"{"jsonrpc":"2.0","id":1,"method":"messages.after","params":{"since_rowid":0,"limit":501}}"#, + #"{"jsonrpc":"2.0","id":1,"method":"messages.after","params":{"since_rowid":0,"attachments":"yes"}}"#, + #"{"jsonrpc":"2.0","id":1,"method":"messages.after","params":{"since_rowid":0,"convert_attachments":1}}"#, + #"{"jsonrpc":"2.0","id":1,"method":"messages.after","params":{"since_rowid":0,"include_reactions":"yes"}}"#, + #"{"jsonrpc":"2.0","id":1,"method":"messages.after","params":{"since_rowid":0,"extra":true}}"#, + #"{"jsonrpc":"2.0","id":1,"method":"messages.after","params":[]}"#, + ] +) +func rpcMessagesAfterRejectsInvalidParams(_ request: String) async throws { + let store = try CommandTestDatabase.makeStoreForRPC() + let output = TestRPCOutput() + let server = RPCServer(store: store, verbose: false, output: output) + + await server.handleLineForTesting(request) + + let error = try #require(output.errors.first?["error"] as? [String: Any]) + #expect(testInt64(error["code"]) == -32602) +} + +private func makeMessagesAfterStore(rows: [(Int64, Int64, Date)]) throws -> MessageStore { + let db = try Connection(.inMemory) + try CommandTestDatabase.createSchema(db, includeChatHandleJoin: true) + try db.run( + """ + INSERT INTO chat( + ROWID, chat_identifier, guid, display_name, service_name, + account_id, account_login, last_addressed_handle + ) + VALUES + (1, 'chat-one', 'iMessage;+;chat-one', 'Chat One', 'iMessage', '', '', ''), + (2, 'chat-two', 'iMessage;+;chat-two', 'Chat Two', 'iMessage', '', '', '') + """ + ) + try db.run("INSERT INTO handle(ROWID, id) VALUES (1, '+15550000001')") + try db.run("INSERT INTO chat_handle_join(chat_id, handle_id) VALUES (1, 1), (2, 1)") + for (rowID, chatID, timestamp) in rows { + try db.run( + """ + INSERT INTO message(ROWID, handle_id, text, date, is_from_me, service) + VALUES (?, 1, ?, ?, 0, 'iMessage') + """, + rowID, + "message-\(rowID)", + CommandTestDatabase.appleEpoch(timestamp) + ) + try db.run( + "INSERT INTO chat_message_join(chat_id, message_id) VALUES (?, ?)", + chatID, + rowID + ) + } + return try MessageStore( + connection: db, + path: ":memory:", + hasAttributedBody: false, + hasReactionColumns: false + ) +} + +private func makeMessagesAfterReactionStore() throws -> MessageStore { + let db = try Connection(.inMemory) + try CommandTestDatabase.createSchema( + db, + includeChatHandleJoin: true, + includeReactionColumns: true + ) + try CommandTestDatabase.seedRPCChat(db) + try db.run("UPDATE message SET guid = 'parent-guid' WHERE ROWID = 5") + try db.run( + """ + INSERT INTO message( + ROWID, handle_id, text, guid, associated_message_guid, associated_message_type, + date, is_from_me, service + ) + VALUES + (6, 1, '', 'reaction-guid', 'p:0/parent-guid', 2001, ?, 0, 'iMessage'), + (7, 1, 'after reaction', 'message-guid', NULL, NULL, ?, 0, 'iMessage') + """, + CommandTestDatabase.appleEpoch(Date(timeIntervalSince1970: 1_700_000_001)), + CommandTestDatabase.appleEpoch(Date(timeIntervalSince1970: 1_700_000_002)) + ) + try db.run("INSERT INTO chat_message_join(chat_id, message_id) VALUES (1, 6), (1, 7)") + return try MessageStore(connection: db, path: ":memory:") +} + +private func testInt64(_ value: Any?) -> Int64? { + if let value = value as? Int64 { return value } + if let value = value as? Int { return Int64(value) } + if let value = value as? NSNumber { return value.int64Value } + return nil +} diff --git a/docs/rpc.md b/docs/rpc.md index 71ae48a7..3235c1f4 100644 --- a/docs/rpc.md +++ b/docs/rpc.md @@ -80,6 +80,43 @@ Result: { "messages": [Message] } ``` +### `messages.after` + +Reads a bounded page in stable message ROWID order. This is the resumable +history surface for message catchup; unlike `messages.history`, it does not +order by timestamp or return the newest rows first. + +Params: + +- `since_rowid` (int, required) — exclusive, non-negative cursor. +- `chat_id` (int, optional) — omit to page across all chats. +- `limit` (int, default 100, maximum 500) +- `attachments` (bool, default `false`) +- `convert_attachments` (bool, default `false`) +- `include_reactions` (bool, default `false`) — include standalone reaction + events in the ordered scan. + +Result: + +```json +{ + "messages": [Message], + "next_rowid": 500, + "has_more": true +} +``` + +Messages are ordered by `message.ROWID ASC`, including when timestamps are +equal. `limit` bounds the returned user-visible messages; the scan can consume +additional URL-preview rows while coalescing or suppressing them. `next_rowid` +is the authoritative physical scan cursor and may therefore advance past the +final returned message. A page can be empty when only suppressed preview rows +remain. Persist `next_rowid` after every response, then request another page +while `has_more` is true. Do not infer pagination state from the message count +or final message id. Set `include_reactions` to `true` when the cursor must also +cover reaction events; with the default, the cursor tracks user-visible message +catchup only. + ### `messages.scheduled` Reads future outbound Send Later rows from `chat.db`. This method is read-only and does not require the IMCore bridge.