diff --git a/agents/common/src/main/java/com/dbx/agent/AgentProtocol.java b/agents/common/src/main/java/com/dbx/agent/AgentProtocol.java index 6af23152a..b29ad89d2 100644 --- a/agents/common/src/main/java/com/dbx/agent/AgentProtocol.java +++ b/agents/common/src/main/java/com/dbx/agent/AgentProtocol.java @@ -48,6 +48,7 @@ public final class AgentProtocol { public static final String MONGO_METHOD_LIST_DATABASES = "list_databases"; public static final String MONGO_METHOD_LIST_COLLECTIONS = "list_collections"; public static final String MONGO_METHOD_FIND_DOCUMENTS = "find_documents"; + public static final String MONGO_METHOD_FIND_ONE = "find_one"; /** * MongoDB read path that returns documents as relaxed Extended JSON for transfer. */ @@ -233,6 +234,7 @@ public final class AgentProtocol { MONGO_METHOD_LIST_DATABASES, MONGO_METHOD_LIST_COLLECTIONS, MONGO_METHOD_FIND_DOCUMENTS, + MONGO_METHOD_FIND_ONE, MONGO_METHOD_FIND_DOCUMENTS_EXTENDED_JSON, MONGO_METHOD_COUNT_DOCUMENTS, MONGO_METHOD_SERVER_VERSION, diff --git a/agents/common/src/main/resources/agent-protocol-v1.json b/agents/common/src/main/resources/agent-protocol-v1.json index ca4dc704b..eadc94527 100644 --- a/agents/common/src/main/resources/agent-protocol-v1.json +++ b/agents/common/src/main/resources/agent-protocol-v1.json @@ -39,6 +39,6 @@ "disconnect", "shutdown" ], - "mongoLegacyMethods": ["list_databases", "list_collections", "find_documents", "find_documents_extended_json", "count_documents", "server_version", "create_index", "drop_indexes", "drop_collection", "drop_database", "insert_document", "update_document", "update_documents", "delete_document", "delete_documents"], + "mongoLegacyMethods": ["list_databases", "list_collections", "find_documents", "find_one", "find_documents_extended_json", "count_documents", "server_version", "create_index", "drop_indexes", "drop_collection", "drop_database", "insert_document", "update_document", "update_documents", "delete_document", "delete_documents"], "kvMethods": ["kv_list_prefix", "kv_get", "kv_put", "kv_delete", "kv_rename", "kv_history", "kv_status", "etcd_compact", "etcd_defrag", "etcd_watch_start", "etcd_watch_poll", "etcd_watch_stop", "etcd_lease_list", "etcd_lease_get", "etcd_lease_grant", "etcd_lease_keepalive_once", "etcd_lease_revoke", "etcd_auth_user_list", "etcd_auth_user_get", "etcd_auth_user_add", "etcd_auth_user_delete", "etcd_auth_user_change_password", "etcd_auth_user_grant_role", "etcd_auth_user_revoke_role", "etcd_auth_role_list", "etcd_auth_role_get", "etcd_auth_role_add", "etcd_auth_role_delete", "etcd_auth_role_grant_permission", "etcd_auth_role_revoke_permission"] } diff --git a/agents/common/src/main/resources/agent-protocol-v2.json b/agents/common/src/main/resources/agent-protocol-v2.json index 079cebc27..0aaa374a7 100644 --- a/agents/common/src/main/resources/agent-protocol-v2.json +++ b/agents/common/src/main/resources/agent-protocol-v2.json @@ -43,7 +43,7 @@ "disconnect", "shutdown" ], - "mongoLegacyMethods": ["list_databases", "list_collections", "find_documents", "find_documents_extended_json", "count_documents", "server_version", "create_index", "drop_indexes", "drop_collection", "drop_database", "insert_document", "update_document", "update_documents", "delete_document", "delete_documents"], + "mongoLegacyMethods": ["list_databases", "list_collections", "find_documents", "find_one", "find_documents_extended_json", "count_documents", "server_version", "create_index", "drop_indexes", "drop_collection", "drop_database", "insert_document", "update_document", "update_documents", "delete_document", "delete_documents"], "kvMethods": ["kv_list_prefix", "kv_get", "kv_put", "kv_delete", "kv_rename", "kv_history", "kv_status", "etcd_compact", "etcd_defrag", "etcd_watch_start", "etcd_watch_poll", "etcd_watch_stop", "etcd_lease_list", "etcd_lease_get", "etcd_lease_grant", "etcd_lease_keepalive_once", "etcd_lease_revoke", "etcd_auth_user_list", "etcd_auth_user_get", "etcd_auth_user_add", "etcd_auth_user_delete", "etcd_auth_user_change_password", "etcd_auth_user_grant_role", "etcd_auth_user_revoke_role", "etcd_auth_role_list", "etcd_auth_role_get", "etcd_auth_role_add", "etcd_auth_role_delete", "etcd_auth_role_grant_permission", "etcd_auth_role_revoke_permission"], "sessionField": "agentSessionId", "cursorSessionField": "sessionId" diff --git a/agents/drivers/mongodb/src/main/java/com/dbx/agent/mongodb/MongoAgent.java b/agents/drivers/mongodb/src/main/java/com/dbx/agent/mongodb/MongoAgent.java index e537c1d06..154c50986 100644 --- a/agents/drivers/mongodb/src/main/java/com/dbx/agent/mongodb/MongoAgent.java +++ b/agents/drivers/mongodb/src/main/java/com/dbx/agent/mongodb/MongoAgent.java @@ -455,6 +455,50 @@ public final class MongoAgent { return documentQueryResult(documents, total); } + private static Object findOne(JsonObject params) { + MongoClient c = requireClient(); + String database = params.get("database").getAsString(); + String collection = params.get("collection").getAsString(); + Document filter = documentOrNull(params, "filter"); + Document projection = documentOrNull(params, "projection"); + Document options = documentOrNull(params, "options"); + Document sort = null; + + if (options != null) { + for (String key : options.keySet()) { + if (!"sort".equals(key)) { + throw new IllegalArgumentException("Unsupported findOne option: " + key); + } + } + Object rawSort = options.get("sort"); + if (rawSort != null) { + if (!(rawSort instanceof Document sortDocument)) { + throw new IllegalArgumentException("Invalid findOne option sort: expected an object"); + } + sort = sortDocument; + } + } + + var iterable = c.getDatabase(database).getCollection(collection).find(filter == null ? new Document() : filter); + if (projection != null) { + iterable = iterable.projection(projection); + } + if (sort != null) { + iterable = iterable.sort(sort); + } + Document document = iterable.limit(1).first(); + + List> documents = new ArrayList<>(); + List extendedDocuments = new ArrayList<>(); + if (document != null) { + documents.add(bsonToJson(document)); + extendedDocuments.add(bsonToExtendedJson(document)); + } + Map result = documentQueryResult(documents, new CollectionTotal(documents.size(), true)); + result.put("extended_documents", extendedDocuments); + return result; + } + /** * MongoDB Extended JSON read path for transfer; output follows the driver's * relaxed Extended JSON representation rather than the UI display format. @@ -1249,6 +1293,7 @@ public final class MongoAgent { case AgentProtocol.MONGO_METHOD_LIST_COLLECTIONS -> listCollections(params); case AgentProtocol.METHOD_LIST_INDEXES -> listIndexes(params); case AgentProtocol.MONGO_METHOD_FIND_DOCUMENTS -> findDocuments(params); + case AgentProtocol.MONGO_METHOD_FIND_ONE -> findOne(params); case AgentProtocol.MONGO_METHOD_FIND_DOCUMENTS_EXTENDED_JSON -> findDocumentsExtendedJson(params); case AgentProtocol.MONGO_METHOD_COUNT_DOCUMENTS -> countDocuments(params); case AgentProtocol.MONGO_METHOD_SERVER_VERSION -> serverVersion(params); diff --git a/agents/drivers/mongodb/src/test/java/com/dbx/agent/mongodb/MongoAgentTest.java b/agents/drivers/mongodb/src/test/java/com/dbx/agent/mongodb/MongoAgentTest.java index 9602bd12a..1f66fa1a1 100644 --- a/agents/drivers/mongodb/src/test/java/com/dbx/agent/mongodb/MongoAgentTest.java +++ b/agents/drivers/mongodb/src/test/java/com/dbx/agent/mongodb/MongoAgentTest.java @@ -13,6 +13,7 @@ import com.google.gson.JsonArray; import com.google.gson.JsonObject; import com.google.gson.JsonParser; import com.mongodb.MongoClientSettings; +import com.mongodb.client.FindIterable; import com.mongodb.client.MongoClient; import com.mongodb.client.MongoCollection; import com.mongodb.client.MongoDatabase; @@ -171,6 +172,74 @@ class MongoAgentTest { assertFalse(json.getAsJsonObject("error").get("message").getAsString().contains("Unknown method")); } + @Test + void findOneUsesOneBoundedReadWithoutCounting() { + List calls = new ArrayList<>(); + MongoClient client = recordingFindOneMongoClient( + calls, + new Document("name", "latest").append("createdAt", 2) + ); + String response = MongoAgent.handleRequest( + "{\"jsonrpc\":\"2.0\",\"id\":16,\"method\":\"find_one\"," + + "\"params\":{\"database\":\"app\",\"collection\":\"orders\"," + + "\"filter\":\"{\\\"status\\\":\\\"open\\\"}\"," + + "\"projection\":\"{\\\"secret\\\":0}\"," + + "\"options\":\"{\\\"sort\\\":{\\\"createdAt\\\":-1}}\"}}", + client + ); + + JsonObject json = JsonParser.parseString(response).getAsJsonObject(); + assertFalse(json.has("error"), json.toString()); + JsonObject result = json.getAsJsonObject("result"); + assertEquals(1, result.get("total").getAsInt()); + assertFalse(result.has("total_is_exact")); + assertEquals("latest", result.getAsJsonArray("documents").get(0).getAsJsonObject().get("name").getAsString()); + assertEquals(1, result.getAsJsonArray("extended_documents").size()); + assertEquals( + List.of( + "find:{\"status\": \"open\"}", + "projection:{\"secret\": 0}", + "sort:{\"createdAt\": -1}", + "limit:1", + "first" + ), + calls + ); + } + + @Test + void findOneReturnsEmptyResultWhenNoDocumentMatches() { + List calls = new ArrayList<>(); + MongoClient client = recordingFindOneMongoClient(calls, null); + String response = MongoAgent.handleRequest( + "{\"jsonrpc\":\"2.0\",\"id\":17,\"method\":\"find_one\"," + + "\"params\":{\"database\":\"app\",\"collection\":\"orders\"}}", + client + ); + + JsonObject result = JsonParser.parseString(response).getAsJsonObject().getAsJsonObject("result"); + assertEquals(0, result.get("total").getAsInt()); + assertEquals(0, result.getAsJsonArray("documents").size()); + assertEquals(0, result.getAsJsonArray("extended_documents").size()); + assertEquals(List.of("find:{}", "limit:1", "first"), calls); + } + + @Test + void findOneRejectsUnsupportedOptionsBeforeReading() { + List calls = new ArrayList<>(); + MongoClient client = recordingFindOneMongoClient(calls, null); + String response = MongoAgent.handleRequest( + "{\"jsonrpc\":\"2.0\",\"id\":18,\"method\":\"find_one\"," + + "\"params\":{\"database\":\"app\",\"collection\":\"orders\"," + + "\"options\":\"{\\\"hint\\\":{\\\"createdAt\\\":1}}\"}}", + client + ); + + JsonObject error = JsonParser.parseString(response).getAsJsonObject().getAsJsonObject("error"); + assertEquals("Unsupported findOne option: hint", error.get("message").getAsString()); + assertTrue(calls.isEmpty()); + } + @Test void collectionTotalUsesEstimatedCountForEmptyFilter() { List calls = new ArrayList<>(); @@ -934,6 +1003,66 @@ class MongoAgentTest { ); } + @SuppressWarnings("unchecked") + private static MongoClient recordingFindOneMongoClient(List calls, Document firstDocument) { + FindIterable[] iterableRef = new FindIterable[1]; + FindIterable iterable = (FindIterable) Proxy.newProxyInstance( + FindIterable.class.getClassLoader(), + new Class[] {FindIterable.class}, + (proxy, method, args) -> { + if ("projection".equals(method.getName()) || "sort".equals(method.getName())) { + calls.add(method.getName() + ":" + ((Document) args[0]).toJson()); + return iterableRef[0]; + } + if ("limit".equals(method.getName())) { + calls.add("limit:" + args[0]); + return iterableRef[0]; + } + if ("first".equals(method.getName())) { + calls.add("first"); + return firstDocument; + } + throw new UnsupportedOperationException(method.getName()); + } + ); + iterableRef[0] = iterable; + + MongoCollection collection = (MongoCollection) Proxy.newProxyInstance( + MongoCollection.class.getClassLoader(), + new Class[] {MongoCollection.class}, + (proxy, method, args) -> { + if ("find".equals(method.getName())) { + calls.add("find:" + ((Document) args[0]).toJson()); + return iterable; + } + throw new UnsupportedOperationException(method.getName()); + } + ); + MongoDatabase database = (MongoDatabase) Proxy.newProxyInstance( + MongoDatabase.class.getClassLoader(), + new Class[] {MongoDatabase.class}, + (proxy, method, args) -> { + if ("getCollection".equals(method.getName())) { + return collection; + } + throw new UnsupportedOperationException(method.getName()); + } + ); + return (MongoClient) Proxy.newProxyInstance( + MongoClient.class.getClassLoader(), + new Class[] {MongoClient.class}, + (proxy, method, args) -> { + if ("getDatabase".equals(method.getName())) { + return database; + } + if ("close".equals(method.getName())) { + return null; + } + throw new UnsupportedOperationException(method.getName()); + } + ); + } + private static JsonObject minimalConnection() { JsonObject conn = new JsonObject(); conn.addProperty("host", "127.0.0.1"); diff --git a/crates/dbx-core/assets/agent-protocol-v1.json b/crates/dbx-core/assets/agent-protocol-v1.json index ca4dc704b..eadc94527 100644 --- a/crates/dbx-core/assets/agent-protocol-v1.json +++ b/crates/dbx-core/assets/agent-protocol-v1.json @@ -39,6 +39,6 @@ "disconnect", "shutdown" ], - "mongoLegacyMethods": ["list_databases", "list_collections", "find_documents", "find_documents_extended_json", "count_documents", "server_version", "create_index", "drop_indexes", "drop_collection", "drop_database", "insert_document", "update_document", "update_documents", "delete_document", "delete_documents"], + "mongoLegacyMethods": ["list_databases", "list_collections", "find_documents", "find_one", "find_documents_extended_json", "count_documents", "server_version", "create_index", "drop_indexes", "drop_collection", "drop_database", "insert_document", "update_document", "update_documents", "delete_document", "delete_documents"], "kvMethods": ["kv_list_prefix", "kv_get", "kv_put", "kv_delete", "kv_rename", "kv_history", "kv_status", "etcd_compact", "etcd_defrag", "etcd_watch_start", "etcd_watch_poll", "etcd_watch_stop", "etcd_lease_list", "etcd_lease_get", "etcd_lease_grant", "etcd_lease_keepalive_once", "etcd_lease_revoke", "etcd_auth_user_list", "etcd_auth_user_get", "etcd_auth_user_add", "etcd_auth_user_delete", "etcd_auth_user_change_password", "etcd_auth_user_grant_role", "etcd_auth_user_revoke_role", "etcd_auth_role_list", "etcd_auth_role_get", "etcd_auth_role_add", "etcd_auth_role_delete", "etcd_auth_role_grant_permission", "etcd_auth_role_revoke_permission"] } diff --git a/crates/dbx-core/assets/agent-protocol-v2.json b/crates/dbx-core/assets/agent-protocol-v2.json index 079cebc27..0aaa374a7 100644 --- a/crates/dbx-core/assets/agent-protocol-v2.json +++ b/crates/dbx-core/assets/agent-protocol-v2.json @@ -43,7 +43,7 @@ "disconnect", "shutdown" ], - "mongoLegacyMethods": ["list_databases", "list_collections", "find_documents", "find_documents_extended_json", "count_documents", "server_version", "create_index", "drop_indexes", "drop_collection", "drop_database", "insert_document", "update_document", "update_documents", "delete_document", "delete_documents"], + "mongoLegacyMethods": ["list_databases", "list_collections", "find_documents", "find_one", "find_documents_extended_json", "count_documents", "server_version", "create_index", "drop_indexes", "drop_collection", "drop_database", "insert_document", "update_document", "update_documents", "delete_document", "delete_documents"], "kvMethods": ["kv_list_prefix", "kv_get", "kv_put", "kv_delete", "kv_rename", "kv_history", "kv_status", "etcd_compact", "etcd_defrag", "etcd_watch_start", "etcd_watch_poll", "etcd_watch_stop", "etcd_lease_list", "etcd_lease_get", "etcd_lease_grant", "etcd_lease_keepalive_once", "etcd_lease_revoke", "etcd_auth_user_list", "etcd_auth_user_get", "etcd_auth_user_add", "etcd_auth_user_delete", "etcd_auth_user_change_password", "etcd_auth_user_grant_role", "etcd_auth_user_revoke_role", "etcd_auth_role_list", "etcd_auth_role_get", "etcd_auth_role_add", "etcd_auth_role_delete", "etcd_auth_role_grant_permission", "etcd_auth_role_revoke_permission"], "sessionField": "agentSessionId", "cursorSessionField": "sessionId" diff --git a/crates/dbx-core/src/db/agent_driver.rs b/crates/dbx-core/src/db/agent_driver.rs index 29efbbf4d..6d06344ad 100644 --- a/crates/dbx-core/src/db/agent_driver.rs +++ b/crates/dbx-core/src/db/agent_driver.rs @@ -1272,6 +1272,7 @@ pub enum MongoAgentMethod { ListDatabases, ListCollections, FindDocuments, + FindOne, FindDocumentsExtendedJson, CountDocuments, ServerVersion, @@ -1287,10 +1288,11 @@ pub enum MongoAgentMethod { } impl MongoAgentMethod { - pub const ALL: [Self; 15] = [ + pub const ALL: [Self; 16] = [ Self::ListDatabases, Self::ListCollections, Self::FindDocuments, + Self::FindOne, Self::FindDocumentsExtendedJson, Self::CountDocuments, Self::ServerVersion, @@ -1310,6 +1312,7 @@ impl MongoAgentMethod { Self::ListDatabases => "list_databases", Self::ListCollections => "list_collections", Self::FindDocuments => "find_documents", + Self::FindOne => "find_one", Self::FindDocumentsExtendedJson => "find_documents_extended_json", Self::CountDocuments => "count_documents", Self::ServerVersion => "server_version", @@ -2464,6 +2467,10 @@ impl AgentDriverClient { self.call_mongo_method(MongoAgentMethod::FindDocuments, params).await } + pub async fn mongo_find_one(&mut self, params: Value) -> Result { + self.call_mongo_method(MongoAgentMethod::FindOne, params).await + } + /// Calls the Mongo agent read method that returns MongoDB relaxed Extended JSON. pub async fn mongo_find_documents_extended_json( &mut self, @@ -4314,6 +4321,7 @@ for line in sys.stdin: assert_eq!(MongoAgentMethod::ListDatabases.as_str(), "list_databases"); assert_eq!(MongoAgentMethod::ListCollections.as_str(), "list_collections"); assert_eq!(MongoAgentMethod::FindDocuments.as_str(), "find_documents"); + assert_eq!(MongoAgentMethod::FindOne.as_str(), "find_one"); assert_eq!(MongoAgentMethod::FindDocumentsExtendedJson.as_str(), "find_documents_extended_json"); assert_eq!(MongoAgentMethod::CountDocuments.as_str(), "count_documents"); assert_eq!(MongoAgentMethod::ServerVersion.as_str(), "server_version"); @@ -4381,6 +4389,7 @@ for line in sys.stdin: let _mongo_list_collections = AgentDriverClient::mongo_list_collections::; let _mongo_list_collection_specs = AgentDriverClient::mongo_list_collection_specs::; let _mongo_find_documents = AgentDriverClient::mongo_find_documents::; + let _mongo_find_one = AgentDriverClient::mongo_find_one::; let _mongo_find_documents_extended_json = AgentDriverClient::mongo_find_documents_extended_json::; let _mongo_server_version = AgentDriverClient::mongo_server_version::; diff --git a/crates/dbx-core/src/mongo_ops.rs b/crates/dbx-core/src/mongo_ops.rs index 9f34994ef..69f2a61f2 100644 --- a/crates/dbx-core/src/mongo_ops.rs +++ b/crates/dbx-core/src/mongo_ops.rs @@ -170,8 +170,18 @@ pub async fn mongo_find_one_core( PoolKind::MongoDb(client) => { mongo_driver::find_one(client, database, collection, filter, projection, options).await } - // The legacy agent only exposes paginated find, which also performs a count. - PoolKind::Agent(_) => Err("MongoDB legacy agent does not support the bounded findOne path".to_string()), + PoolKind::Agent(client) => { + let mut client = client.lock().await; + client + .mongo_find_one(serde_json::json!({ + "database": database, + "collection": collection, + "filter": filter, + "projection": projection, + "options": options, + })) + .await + } _ => Err("Not a MongoDB connection".to_string()), } }