From 0c4ae7caf68cab8bd0d6da1b6af4cdc1e3a44741 Mon Sep 17 00:00:00 2001 From: t8y2 <1156263951@qq.com> Date: Wed, 1 Jul 2026 20:23:51 +0800 Subject: [PATCH] fix(mongodb): support shell delete commands --- .../main/java/com/dbx/agent/AgentProtocol.java | 5 ++++- .../src/main/resources/agent-protocol-v1.json | 4 +++- .../java/com/dbx/agent/mongodb/MongoAgent.java | 16 ++++++++++++++++ .../com/dbx/agent/mongodb/MongoAgentTest.java | 14 ++++++++++++++ crates/dbx-core/assets/agent-protocol-v1.json | 4 +++- crates/dbx-core/src/db/agent_driver.rs | 17 ++++++++++++++++- crates/dbx-core/src/mongo_ops.rs | 13 ++++++++++++- 7 files changed, 68 insertions(+), 5 deletions(-) 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 c894003b2..9afc9de99 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_UPDATE_DOCUMENT = "update_document"; public static final String MONGO_METHOD_UPDATE_DOCUMENTS = "update_documents"; public static final String MONGO_METHOD_DELETE_DOCUMENT = "delete_document"; + public static final String MONGO_METHOD_DELETE_DOCUMENTS = "delete_documents"; public static final String KV_METHOD_LIST_PREFIX = "kv_list_prefix"; public static final String KV_METHOD_GET = "kv_get"; @@ -123,7 +124,9 @@ public final class AgentProtocol { MONGO_METHOD_SERVER_VERSION, MONGO_METHOD_INSERT_DOCUMENT, MONGO_METHOD_UPDATE_DOCUMENT, - MONGO_METHOD_DELETE_DOCUMENT + MONGO_METHOD_UPDATE_DOCUMENTS, + MONGO_METHOD_DELETE_DOCUMENT, + MONGO_METHOD_DELETE_DOCUMENTS )); public static final List KV_METHODS = Collections.unmodifiableList(Arrays.asList( diff --git a/agents/common/src/main/resources/agent-protocol-v1.json b/agents/common/src/main/resources/agent-protocol-v1.json index 6fd3e71fe..d9f88a126 100644 --- a/agents/common/src/main/resources/agent-protocol-v1.json +++ b/agents/common/src/main/resources/agent-protocol-v1.json @@ -72,7 +72,9 @@ "server_version", "insert_document", "update_document", - "delete_document" + "update_documents", + "delete_document", + "delete_documents" ], "kvMethods": [ "kv_list_prefix", 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 839b82eda..73843d864 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 @@ -536,6 +536,21 @@ public final class MongoAgent { return Collections.singletonMap("deleted_count", result.getDeletedCount()); } + private static Object deleteDocuments(JsonObject params) { + MongoClient c = requireClient(); + String database = params.get("database").getAsString(); + String collection = params.get("collection").getAsString(); + String filterJson = params.get("filter_json").getAsString(); + boolean many = params.get("many").getAsBoolean(); + + var col = c.getDatabase(database).getCollection(collection); + Document filter = documentForWrite(filterJson); + // Shell deleteOne/deleteMany use a filter document, unlike the row-view + // delete path which always targets a single _id. + var result = many ? col.deleteMany(filter) : col.deleteOne(filter); + return Collections.singletonMap("deleted_count", result.getDeletedCount()); + } + private static Map bsonToJson(Document doc) { Map result = new LinkedHashMap<>(); for (Map.Entry entry : doc.entrySet()) { @@ -662,6 +677,7 @@ public final class MongoAgent { case AgentProtocol.MONGO_METHOD_UPDATE_DOCUMENT -> updateDocument(params); case AgentProtocol.MONGO_METHOD_UPDATE_DOCUMENTS -> updateDocuments(params); case AgentProtocol.MONGO_METHOD_DELETE_DOCUMENT -> deleteDocument(params); + case AgentProtocol.MONGO_METHOD_DELETE_DOCUMENTS -> deleteDocuments(params); case AgentProtocol.METHOD_DISCONNECT, AgentProtocol.METHOD_SHUTDOWN -> { if (client != null) { client.close(); 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 47934942b..7c5aff74c 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 @@ -149,6 +149,20 @@ class MongoAgentTest { assertFalse(json.getAsJsonObject("error").get("message").getAsString().contains("Unknown method")); } + @Test + void deleteDocumentsMethodIsRecognizedOverJsonRpc() { + String response = MongoAgent.handleRequest( + "{\"jsonrpc\":\"2.0\",\"id\":11,\"method\":\"delete_documents\"," + + "\"params\":{\"database\":\"app\",\"collection\":\"orders\"," + + "\"filter_json\":\"{\\\"status\\\":\\\"draft\\\"}\",\"many\":true}}"); + + JsonObject json = JsonParser.parseString(response).getAsJsonObject(); + assertEquals(11, json.get("id").getAsInt()); + assertEquals("Not connected", json.getAsJsonObject("error").get("message").getAsString()); + assertFalse(json.getAsJsonObject("error").get("message").getAsString().contains("Unknown method")); + assertTrue(AgentProtocol.MONGO_LEGACY_METHODS.contains(AgentProtocol.MONGO_METHOD_DELETE_DOCUMENTS)); + } + @Test void extractsServerVersionFromBuildInfo() { assertEquals("4.4.29", MongoAgent.serverVersionFromBuildInfo(new Document("version", "4.4.29"))); diff --git a/crates/dbx-core/assets/agent-protocol-v1.json b/crates/dbx-core/assets/agent-protocol-v1.json index 6fd3e71fe..d9f88a126 100644 --- a/crates/dbx-core/assets/agent-protocol-v1.json +++ b/crates/dbx-core/assets/agent-protocol-v1.json @@ -72,7 +72,9 @@ "server_version", "insert_document", "update_document", - "delete_document" + "update_documents", + "delete_document", + "delete_documents" ], "kvMethods": [ "kv_list_prefix", diff --git a/crates/dbx-core/src/db/agent_driver.rs b/crates/dbx-core/src/db/agent_driver.rs index 5eef772e5..9cc2d6af0 100644 --- a/crates/dbx-core/src/db/agent_driver.rs +++ b/crates/dbx-core/src/db/agent_driver.rs @@ -249,10 +249,11 @@ pub enum MongoAgentMethod { UpdateDocument, UpdateDocuments, DeleteDocument, + DeleteDocuments, } impl MongoAgentMethod { - pub const ALL: [Self; 8] = [ + pub const ALL: [Self; 10] = [ Self::ListDatabases, Self::ListCollections, Self::FindDocuments, @@ -260,7 +261,9 @@ impl MongoAgentMethod { Self::ServerVersion, Self::InsertDocument, Self::UpdateDocument, + Self::UpdateDocuments, Self::DeleteDocument, + Self::DeleteDocuments, ]; pub fn as_str(self) -> &'static str { @@ -274,6 +277,7 @@ impl MongoAgentMethod { Self::UpdateDocument => "update_document", Self::UpdateDocuments => "update_documents", Self::DeleteDocument => "delete_document", + Self::DeleteDocuments => "delete_documents", } } } @@ -972,6 +976,13 @@ impl AgentDriverClient { self.call_mongo_method(MongoAgentMethod::DeleteDocument, params).await } + pub async fn mongo_delete_documents( + &mut self, + params: Value, + ) -> Result { + self.call_mongo_method(MongoAgentMethod::DeleteDocuments, params).await + } + pub async fn try_optional_handshake(&mut self, app_version: &str) -> Option { match self.call_method::(AgentMethod::Handshake, agent_handshake_params(app_version)).await { Ok(handshake) => { @@ -1495,7 +1506,9 @@ mod tests { assert_eq!(MongoAgentMethod::ServerVersion.as_str(), "server_version"); assert_eq!(MongoAgentMethod::InsertDocument.as_str(), "insert_document"); assert_eq!(MongoAgentMethod::UpdateDocument.as_str(), "update_document"); + assert_eq!(MongoAgentMethod::UpdateDocuments.as_str(), "update_documents"); assert_eq!(MongoAgentMethod::DeleteDocument.as_str(), "delete_document"); + assert_eq!(MongoAgentMethod::DeleteDocuments.as_str(), "delete_documents"); } #[test] @@ -1537,7 +1550,9 @@ mod tests { let _mongo_server_version = AgentDriverClient::mongo_server_version::; let _mongo_insert_document = AgentDriverClient::mongo_insert_document::; let _mongo_update_document = AgentDriverClient::mongo_update_document::; + let _mongo_update_documents = AgentDriverClient::mongo_update_documents::; let _mongo_delete_document = AgentDriverClient::mongo_delete_document::; + let _mongo_delete_documents = AgentDriverClient::mongo_delete_documents::; } #[test] diff --git a/crates/dbx-core/src/mongo_ops.rs b/crates/dbx-core/src/mongo_ops.rs index 310370f1c..881539908 100644 --- a/crates/dbx-core/src/mongo_ops.rs +++ b/crates/dbx-core/src/mongo_ops.rs @@ -318,7 +318,18 @@ pub async fn mongo_delete_documents_core( PoolKind::MongoDb(client) => { mongo_driver::delete_documents(client, database, collection, filter_json, many).await } - PoolKind::Agent(_) => Err("MongoDB legacy agent does not support bulk deleteOne/deleteMany writes".to_string()), + PoolKind::Agent(client) => { + let mut client = client.lock().await; + let result: serde_json::Value = client + .mongo_delete_documents(serde_json::json!({ + "database": database, + "collection": collection, + "filter_json": filter_json, + "many": many, + })) + .await?; + Ok(result.get("deleted_count").and_then(|v| v.as_u64()).unwrap_or(0)) + } _ => Err("Not a MongoDB connection".to_string()), } }