diff --git a/CHANGELOG.md b/CHANGELOG.md index 7282b3f..06267d2 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,10 @@ # Changelog +## 0.14.0 - Unreleased + +### JSON-RPC +- feat: add bounded `messages.after` pagination with authoritative database-instance-scoped ROWID cursors, cross-chat catchup, and optional standalone reaction events (#200, #201, thanks @vincentkoc). + ## 0.13.5 - Unreleased ### JSON-RPC diff --git a/README.md b/README.md index 4eca570..e50b290 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 7bcc17e..5bc0da0 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 c51df01..c58640b 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 0000000..855cb4d --- /dev/null +++ b/Sources/IMsgCore/MessagesAfterPage.swift @@ -0,0 +1,13 @@ +public struct MessagesAfterPage: Sendable, Equatable { + public let messages: [Message] + /// Physical scan cursor scoped to the same Messages database instance. + /// Discard it after that database is replaced, restored, or recreated. + 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 0000000..da9e859 --- /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 cb3ee6e..ca86026 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 0000000..d650cae --- /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 0000000..7b0eb49 --- /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/json.md b/docs/json.md index 7663776..dcf7ad5 100644 --- a/docs/json.md +++ b/docs/json.md @@ -40,7 +40,7 @@ Returned by `imsg history`, `imsg search`, `imsg watch`, and the JSON-RPC `messa | Field | Type | Notes | |-------|------|-------| -| `id` | int | rowid. Use as the `--since-rowid` cursor in watch. | +| `id` | int | Database-instance-scoped rowid. Use as the `--since-rowid` cursor in watch, but discard saved cursors after replacing or restoring the Messages database. | | `chat_id` | int | Always present. Preferred routing handle. | | `chat_identifier` | string | Portable handle. | | `chat_guid` | string | Portable GUID. | @@ -121,7 +121,9 @@ Live watch calls do not delay the text message waiting for a preview. If the pre ### Reaction extensions -Present on `imsg watch --reactions` events: +Present on standalone reaction rows emitted by `imsg watch --reactions`, +`watch.subscribe` with `include_reactions: true`, and `messages.after` with +`include_reactions: true`: | Field | Type | Notes | |-------|------|-------| @@ -131,7 +133,10 @@ Present on `imsg watch --reactions` events: | `is_reaction_add` | bool | `true` for add, `false` for remove. | | `reacted_to_guid` | string | The message guid this tapback targets. | -`history` deliberately hides reaction rows so they don't duplicate the reacted message. Reaction events only surface in the live watch stream. +`history` deliberately hides standalone reaction rows so they don't duplicate +the reacted message. Live watch surfaces emit them only when reactions are +enabled; `messages.after` includes them in ROWID order only when +`include_reactions` is `true`. ### Native poll extension diff --git a/docs/rpc.md b/docs/rpc.md index 71ae48a..06aa41c 100644 --- a/docs/rpc.md +++ b/docs/rpc.md @@ -80,6 +80,49 @@ 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. + +ROWID cursors are scoped to the exact Messages database instance that produced +them. They are not portable between machines, accounts, or database files, and +they are not durable across replacement, restoration, or recreation of +`chat.db`. After any database replacement, discard the saved cursor and start a +new scan from a cursor appropriate for that database instance. + ### `messages.scheduled` Reads future outbound Send Later rows from `chat.db`. This method is read-only and does not require the IMCore bridge.