diff --git a/CHANGELOG.md b/CHANGELOG.md index 56fdef005a..2c42b2621f 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -86,6 +86,8 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - Literal backticks in cloudflared, cloud-sql-proxy, SSH config, remote command and dump tool install messages. - Tunnel command preview showing port 0 or the wrong host when Port is blank or the connection uses a host list. - SSH tab host-list warning naming replica set failover for Redis and Kafka, and implying Sentinel works through a tunnel. +- Weaviate exports failing after the first 10,000 objects of a collection. +- Weaviate Raw Filter row sent as a filter on a property named `__RAW__`. - iPhone and iPad reading a Safe Mode level they do not recognize from iCloud as Off. ### Security diff --git a/Packages/TableProCore/Sources/TableProWeaviateCore/WeaviateFilterBuilder.swift b/Packages/TableProCore/Sources/TableProWeaviateCore/WeaviateFilterBuilder.swift index 8d5edb30d1..6a29940eba 100644 --- a/Packages/TableProCore/Sources/TableProWeaviateCore/WeaviateFilterBuilder.swift +++ b/Packages/TableProCore/Sources/TableProWeaviateCore/WeaviateFilterBuilder.swift @@ -10,6 +10,7 @@ public enum WeaviateFilterError: Error, LocalizedError, Equatable { case comparisonNeedsNumberOrDate(column: String, op: String) case textMatchNeedsText(column: String, op: String) case vectorNotFilterable(column: String) + case rawFilterUnsupported public var errorDescription: String? { switch self { @@ -46,6 +47,10 @@ public enum WeaviateFilterError: Error, LocalizedError, Equatable { ) case .vectorNotFilterable(let column): return String(format: String(localized: "Weaviate cannot filter on %@."), column) + case .rawFilterUnsupported: + return String( + localized: "Raw filters aren't available for Weaviate. Write the where filter in a GraphQL query in the editor." + ) } } } @@ -96,6 +101,8 @@ public enum WeaviateValueKind: String, Sendable, Equatable { } public enum WeaviateFilterBuilder { + private static let rawFilterColumn = "__RAW__" + public static func graphQLWhere( filters: [WeaviateFilterSpec], logicMode: String, @@ -110,6 +117,9 @@ public enum WeaviateFilterBuilder { public static func operand(for filter: WeaviateFilterSpec, types: [String: String]) throws -> String { let column = filter.column + guard column != rawFilterColumn else { + throw WeaviateFilterError.rawFilterUnsupported + } guard column != WeaviateSchema.vectorColumn else { throw WeaviateFilterError.vectorNotFilterable(column: column) } diff --git a/Packages/TableProCore/Sources/TableProWeaviateCore/Wire/WeaviateClient.swift b/Packages/TableProCore/Sources/TableProWeaviateCore/Wire/WeaviateClient.swift index cb56e8146b..307eff015e 100644 --- a/Packages/TableProCore/Sources/TableProWeaviateCore/Wire/WeaviateClient.swift +++ b/Packages/TableProCore/Sources/TableProWeaviateCore/Wire/WeaviateClient.swift @@ -57,11 +57,37 @@ public final class WeaviateClient: @unchecked Sendable { offset: Int, includeVector: Bool = true ) async throws -> [WeaviateObject] { - var query = [ - "class": collection, - "limit": String(max(limit, 0)), - "offset": String(max(offset, 0)) - ] + try await listObjects( + collection: collection, + limit: limit, + position: ["offset": String(max(offset, 0))], + includeVector: includeVector + ) + } + + public func objects( + collection: String, + limit: Int, + after: String?, + includeVector: Bool = true + ) async throws -> [WeaviateObject] { + try await listObjects( + collection: collection, + limit: limit, + position: after.map { ["after": $0] } ?? [:], + includeVector: includeVector + ) + } + + private func listObjects( + collection: String, + limit: Int, + position: [String: String], + includeVector: Bool + ) async throws -> [WeaviateObject] { + var query = position + query["class"] = collection + query["limit"] = String(max(limit, 0)) if includeVector { query["include"] = "vector" } diff --git a/Packages/TableProCore/Sources/TableProWeaviateCore/Wire/WeaviateObjectPages.swift b/Packages/TableProCore/Sources/TableProWeaviateCore/Wire/WeaviateObjectPages.swift new file mode 100644 index 0000000000..833bc9c40b --- /dev/null +++ b/Packages/TableProCore/Sources/TableProWeaviateCore/Wire/WeaviateObjectPages.swift @@ -0,0 +1,74 @@ +import Foundation + +public struct WeaviateObjectPages: AsyncSequence, Sendable { + public typealias Element = [WeaviateObject] + + private let client: WeaviateClient + private let collection: String + private let pageSize: Int + private let includeVector: Bool + + init(client: WeaviateClient, collection: String, pageSize: Int, includeVector: Bool) { + self.client = client + self.collection = collection + self.pageSize = Swift.max(pageSize, 1) + self.includeVector = includeVector + } + + public func makeAsyncIterator() -> AsyncIterator { + AsyncIterator(pages: self) + } + + public struct AsyncIterator: AsyncIteratorProtocol { + private let pages: WeaviateObjectPages + private var cursor: String? + private var isExhausted = false + + init(pages: WeaviateObjectPages) { + self.pages = pages + } + + public mutating func next() async throws -> [WeaviateObject]? { + guard !isExhausted else { return nil } + try Task.checkCancellation() + let page = try await pages.client.objects( + collection: pages.collection, + limit: pages.pageSize, + after: cursor, + includeVector: pages.includeVector + ) + isExhausted = page.count < pages.pageSize + guard !page.isEmpty else { return nil } + if !isExhausted { + cursor = try nextCursor(after: page) + } + return page + } + + private func nextCursor(after page: [WeaviateObject]) throws -> String { + guard let last = page.last?.uuid, !last.isEmpty else { + throw WeaviateError.malformedResponse(String( + format: String(localized: "Weaviate returned an object of %@ without an id, so the rest of the collection cannot be read."), + pages.collection + )) + } + guard last != cursor else { + throw WeaviateError.malformedResponse(String( + format: String(localized: "Weaviate returned the same page of %@ twice. Reading a whole collection needs Weaviate 1.18 or later."), + pages.collection + )) + } + return last + } + } +} + +public extension WeaviateClient { + func objectPages( + collection: String, + pageSize: Int, + includeVector: Bool = true + ) -> WeaviateObjectPages { + WeaviateObjectPages(client: self, collection: collection, pageSize: pageSize, includeVector: includeVector) + } +} diff --git a/Packages/TableProCore/Tests/TableProSyncTests/SyncRecordMapperTests.swift b/Packages/TableProCore/Tests/TableProSyncTests/SyncRecordMapperTests.swift index 0463cdcad4..2fa2afdfb4 100644 --- a/Packages/TableProCore/Tests/TableProSyncTests/SyncRecordMapperTests.swift +++ b/Packages/TableProCore/Tests/TableProSyncTests/SyncRecordMapperTests.swift @@ -86,4 +86,25 @@ struct SyncRecordMapperTests { let decoded = try #require(SyncRecordMapper.toConnection(record)) #expect(decoded.safeModeLevel == .confirmWrites) } + + @Test("An unrecognized wire value keeps the legacy read-only restriction") + func unknownWireValuePreservesReadOnly() throws { + let record = makeRawRecord(safeModeLevelRaw: "someFutureLevel", isReadOnly: true) + let decoded = try #require(SyncRecordMapper.toConnection(record)) + #expect(decoded.safeModeLevel == .readOnly) + } + + @Test("A rename preserves an unrecognized wire value and requires confirmation") + func renamePreservesUnknownWireValue() throws { + let record = makeRawRecord(safeModeLevelRaw: "someFutureLevel") + var connection = try #require(SyncRecordMapper.toConnection(record)) + connection.name = "Renamed" + + SyncRecordMapper.updateRecord(record, with: connection) + + #expect(record["safeModeLevel"] as? String == "someFutureLevel") + let decoded = try #require(SyncRecordMapper.toConnection(record)) + #expect(decoded.name == "Renamed") + #expect(decoded.safeModeLevel == .confirmWrites) + } } diff --git a/Packages/TableProCore/Tests/TableProWeaviateCoreTests/WeaviateFilterAndPathTests.swift b/Packages/TableProCore/Tests/TableProWeaviateCoreTests/WeaviateFilterAndPathTests.swift index 78cc4eab39..26362c7e19 100644 --- a/Packages/TableProCore/Tests/TableProWeaviateCoreTests/WeaviateFilterAndPathTests.swift +++ b/Packages/TableProCore/Tests/TableProWeaviateCoreTests/WeaviateFilterAndPathTests.swift @@ -171,6 +171,31 @@ struct WeaviateFilterOperatorTests { } } + @Test("The raw filter row is refused instead of filtering a property named __RAW__") + func rawFilterRowIsRefused() { + #expect(throws: WeaviateFilterError.rawFilterUnsupported) { + _ = try operand("__RAW__", "=", "wordCount > 10") + } + #expect(throws: WeaviateFilterError.rawFilterUnsupported) { + _ = try WeaviateFilterBuilder.graphQLWhere( + filters: [ + WeaviateFilterSpec(column: "title", op: "=", value: "a"), + WeaviateFilterSpec(column: "__RAW__", op: "=", value: "{ path: [\"title\"] }") + ], + logicMode: "AND", + types: articleTypes + ) + } + } + + @Test("The raw filter refusal points at the GraphQL editor") + func rawFilterRefusalNamesTheEditor() { + #expect( + WeaviateFilterError.rawFilterUnsupported.errorDescription + == "Raw filters aren't available for Weaviate. Write the where filter in a GraphQL query in the editor." + ) + } + @Test("A quote in a value cannot break out of the GraphQL string") func valuesAreEscaped() throws { let clause = try operand("title", "=", "a\"b\\c") diff --git a/Packages/TableProCore/Tests/TableProWeaviateCoreTests/WeaviateObjectPagesTests.swift b/Packages/TableProCore/Tests/TableProWeaviateCoreTests/WeaviateObjectPagesTests.swift new file mode 100644 index 0000000000..33a96410ef --- /dev/null +++ b/Packages/TableProCore/Tests/TableProWeaviateCoreTests/WeaviateObjectPagesTests.swift @@ -0,0 +1,163 @@ +import Foundation +@testable import TableProWeaviateCore +import Testing + +@Suite("Weaviate object pages") +struct WeaviateObjectPagesTests { + @Test("Reading a collection past the query maximum yields every object once, in order") + func readsPastQueryMaximum() async throws { + let transport = CursorPagingWeaviateTransport(objectCount: 12_000) + let client = testClient(transport: transport) + + var uuids: [String] = [] + for try await page in client.objectPages(collection: "Article", pageSize: 500) { + uuids += page.map(\.uuid) + } + + #expect(uuids.count == 12_000) + #expect(uuids == transport.uuids) + } + + @Test("The first page is a plain listing, so an object at the nil uuid is not skipped") + func firstPageKeepsTheNilUUID() async throws { + let transport = CursorPagingWeaviateTransport(objectCount: 1_200) + let client = testClient(transport: transport) + + var uuids: [String] = [] + for try await page in client.objectPages(collection: "Article", pageSize: 500) { + uuids += page.map(\.uuid) + } + + #expect(uuids.first == CursorPagingWeaviateTransport.nilUUID) + #expect(uuids == transport.uuids) + } + + @Test("Each later page starts after the last uuid of the page before it") + func pagesByCursor() async throws { + let transport = CursorPagingWeaviateTransport(objectCount: 1_200) + let client = testClient(transport: transport) + + var pages: [[WeaviateObject]] = [] + for try await page in client.objectPages(collection: "Article", pageSize: 500) { + pages.append(page) + } + + #expect(pages.map(\.count) == [500, 500, 200]) + #expect(transport.queries.map { $0["after"] } == [nil, pages[0].last?.uuid, pages[1].last?.uuid]) + #expect(transport.queries.allSatisfy { $0["offset"] == nil }) + #expect(transport.queries.allSatisfy { $0["class"] == "Article" && $0["include"] == "vector" }) + } + + @Test("A collection that fills its last page exactly ends on the empty page after it") + func endsOnEmptyPage() async throws { + let transport = CursorPagingWeaviateTransport(objectCount: 1_000) + let client = testClient(transport: transport) + + var counts: [Int] = [] + for try await page in client.objectPages(collection: "Article", pageSize: 500) { + counts.append(page.count) + } + + #expect(counts == [500, 500]) + #expect(transport.queries.count == 3) + } + + @Test("A server that ignores the cursor fails the read instead of repeating the first page") + func refusesACursorThatDoesNotMove() async throws { + let transport = CursorPagingWeaviateTransport(objectCount: 1_200, honorsCursor: false) + let client = testClient(transport: transport) + + let expected = WeaviateError.malformedResponse( + "Weaviate returned the same page of Article twice. Reading a whole collection needs Weaviate 1.18 or later." + ) + await #expect(throws: expected) { + for try await _ in client.objectPages(collection: "Article", pageSize: 500) {} + } + #expect(transport.queries.count == 2) + } + + @Test("A full page whose last object has no id cannot be paged past, so the read fails") + func refusesAPageWithoutACursor() async throws { + let transport = FakeWeaviateTransport() + transport.respond( + method: "GET", + path: "/v1/objects", + status: 200, + json: ["objects": [["class": "Article", "properties": ["title": "Hello"]]]] + ) + let client = testClient(transport: transport) + + let expected = WeaviateError.malformedResponse( + "Weaviate returned an object of Article without an id, so the rest of the collection cannot be read." + ) + await #expect(throws: expected) { + for try await _ in client.objectPages(collection: "Article", pageSize: 1) {} + } + #expect(transport.requests.count == 1) + } + + @Test("Leaving out the vector leaves out the include parameter") + func omitsVectorWhenAsked() async throws { + let transport = CursorPagingWeaviateTransport(objectCount: 10) + let client = testClient(transport: transport) + + for try await _ in client.objectPages(collection: "Article", pageSize: 500, includeVector: false) {} + + #expect(transport.queries.map { $0["include"] } == [nil]) + } +} + +final class CursorPagingWeaviateTransport: WeaviateTransport, @unchecked Sendable { + static let nilUUID = "00000000-0000-0000-0000-000000000000" + + let uuids: [String] + private let queryMaximumResults: Int + private let honorsCursor: Bool + private(set) var queries: [[String: String]] = [] + + init(objectCount: Int, queryMaximumResults: Int = 10_000, honorsCursor: Bool = true) { + uuids = (0.. WeaviateHTTPResponse { + let items = URLComponents(url: request.url, resolvingAgainstBaseURL: false)?.queryItems ?? [] + let query = Dictionary(items.map { ($0.name, $0.value ?? "") }, uniquingKeysWith: { _, last in last }) + queries.append(query) + + let limit = Int(query["limit"] ?? "") ?? 25 + let offset = Int(query["offset"] ?? "") ?? 0 + if honorsCursor, let after = query["after"] { + guard offset == 0 else { + return refusal("offset cannot be set with after and limit parameters") + } + let start = uuids.firstIndex { $0 > after } ?? uuids.count + return page(start: start, limit: limit) + } + guard offset + limit <= queryMaximumResults else { + return refusal("query maximum results exceeded") + } + return page(start: offset, limit: limit) + } + + func cancelAll() {} + + private func page(start: Int, limit: Int) -> WeaviateHTTPResponse { + let lower = min(max(start, 0), uuids.count) + let upper = min(lower + max(limit, 0), uuids.count) + let objects: [[String: Any]] = uuids[lower.. WeaviateHTTPResponse { + let json: [String: Any] = ["error": [["message": "msg:offset or limit code:400 err:\(message)"]]] + let body = (try? JSONSerialization.data(withJSONObject: json)) ?? Data() + return WeaviateHTTPResponse(statusCode: 422, body: body) + } +} diff --git a/Packages/TableProCore/Tests/TableProWeaviateCoreTests/WeaviateTestSupport.swift b/Packages/TableProCore/Tests/TableProWeaviateCoreTests/WeaviateTestSupport.swift index cd6a8ba115..d1ccaf1dcc 100644 --- a/Packages/TableProCore/Tests/TableProWeaviateCoreTests/WeaviateTestSupport.swift +++ b/Packages/TableProCore/Tests/TableProWeaviateCoreTests/WeaviateTestSupport.swift @@ -57,7 +57,7 @@ func testSettings(auth: WeaviateAuth = WeaviateAuth(method: .none)) -> WeaviateC } func testClient( - transport: FakeWeaviateTransport, + transport: any WeaviateTransport, auth: WeaviateAuth = WeaviateAuth(method: .none) ) -> WeaviateClient { WeaviateClient(settings: testSettings(auth: auth), transport: transport, timeout: { 30 }) diff --git a/Plugins/WeaviateDriverPlugin/WeaviatePluginDriver+Execution.swift b/Plugins/WeaviateDriverPlugin/WeaviatePluginDriver+Execution.swift index c6e983719e..0142ec23db 100644 --- a/Plugins/WeaviateDriverPlugin/WeaviatePluginDriver+Execution.swift +++ b/Plugins/WeaviateDriverPlugin/WeaviatePluginDriver+Execution.swift @@ -74,20 +74,12 @@ extension WeaviatePluginDriver { estimatedRowCount: nil ))) - var offset = 0 - while true { - try Task.checkCancellation() - let objects = try await client.objects( - collection: name, limit: Self.exportPageSize, offset: offset, includeVector: true - ) - guard !objects.isEmpty else { break } + for try await objects in client.objectPages(collection: name, pageSize: Self.exportPageSize) { continuation.yield(.rows(objects.map { object in WeaviateObjectCodec.row(for: object, columns: columns).map { value in value.map(PluginCellValue.text) ?? .null } })) - if objects.count < Self.exportPageSize { break } - offset += objects.count } } diff --git a/TablePro/Resources/Localizable.xcstrings b/TablePro/Resources/Localizable.xcstrings index 0535a15bcc..059422bdf8 100644 --- a/TablePro/Resources/Localizable.xcstrings +++ b/TablePro/Resources/Localizable.xcstrings @@ -183049,6 +183049,15 @@ }, "Weaviate does not support views." : { + }, + "Raw filters aren't available for Weaviate. Write the where filter in a GraphQL query in the editor." : { + + }, + "Weaviate returned an object of %@ without an id, so the rest of the collection cannot be read." : { + + }, + "Weaviate returned the same page of %@ twice. Reading a whole collection needs Weaviate 1.18 or later." : { + }, "The first sheet has no header row to map columns from." : {