diff --git a/dsl-template/src/commonMain/kotlin/options/WriteConcern.kt b/dsl-template/src/commonMain/kotlin/options/WriteConcern.kt index 96c698e9..03429247 100644 --- a/dsl-template/src/commonMain/kotlin/options/WriteConcern.kt +++ b/dsl-template/src/commonMain/kotlin/options/WriteConcern.kt @@ -379,4 +379,53 @@ interface HasWriteConcern : Options { accept(WriteConcernOption(concern, context)) } + /** + * Specifies the [WriteConcern] for this operation. + * + * The write concern specifies which nodes must acknowledge having applied this write operation. + * The stronger the write concern, the less chance of data loss, but the higher the latency. + * + * To learn more about the different options, see: + * - [acknowledgement]: how many nodes should acknowledge this request? + * - [writeToJournal]: should the nodes also acknowledge synchronizing their journal? + * - [writeTimeout]: after how long should we give up on the acknowledgment, if it doesn't succeed? + * + * ### Example + * + * ```kotlin + * collections.updateMany( + * options = { + * writeConcern(Majority, writeTimeout = 2.seconds) + * } + * ) { + * User::age inc 1 + * } + * ``` + * + * ### Transactions + * + * In multi-document transactions, only specify a write concern at the transaction level, and not at the level + * of individual operation. + * + * ### External resources + * + * - [Official documentation](https://www.mongodb.com/docs/manual/reference/write-concern) + */ + fun writeConcern( + /** + * Describes how many nodes must acknowledge the write operation. See [WriteConcern.acknowledgment]. + */ + acknowledgement: WriteAcknowledgment? = null, + /** + * Specifies whether the write operation must acknowledge being written to the journal. See [WriteConcern.writeToJournal]. + */ + writeToJournal: Boolean? = null, + /** + * Specifies a time limit for a write operation to propagate to enough members to achieve [acknowledgement] and [writeToJournal]. See [WriteConcern.writeTimeout]. + */ + writeTimeout: Duration? = null, + ) { + writeConcern(WriteConcern(acknowledgement, writeToJournal, writeTimeout)) + } + } diff --git a/dsl/src/commonMain/kotlin/options/WriteConcern.kt b/dsl/src/commonMain/kotlin/options/WriteConcern.kt index 8f699199..29c993ae 100644 --- a/dsl/src/commonMain/kotlin/options/WriteConcern.kt +++ b/dsl/src/commonMain/kotlin/options/WriteConcern.kt @@ -382,4 +382,53 @@ interface HasWriteConcern : Options { accept(WriteConcernOption(concern, context)) } + /** + * Specifies the [WriteConcern] for this operation. + * + * The write concern specifies which nodes must acknowledge having applied this write operation. + * The stronger the write concern, the less chance of data loss, but the higher the latency. + * + * To learn more about the different options, see: + * - [acknowledgement]: how many nodes should acknowledge this request? + * - [writeToJournal]: should the nodes also acknowledge synchronizing their journal? + * - [writeTimeout]: after how long should we give up on the acknowledgment, if it doesn't succeed? + * + * ### Example + * + * ```kotlin + * collections.updateMany( + * options = { + * writeConcern(Majority, writeTimeout = 2.seconds) + * } + * ) { + * User::age inc 1 + * } + * ``` + * + * ### Transactions + * + * In multi-document transactions, only specify a write concern at the transaction level, and not at the level + * of individual operation. + * + * ### External resources + * + * - [Official documentation](https://www.mongodb.com/docs/manual/reference/write-concern) + */ + fun writeConcern( + /** + * Describes how many nodes must acknowledge the write operation. See [WriteConcern.acknowledgment]. + */ + acknowledgement: WriteAcknowledgment? = null, + /** + * Specifies whether the write operation must acknowledge being written to the journal. See [WriteConcern.writeToJournal]. + */ + writeToJournal: Boolean? = null, + /** + * Specifies a time limit for a write operation to propagate to enough members to achieve [acknowledgement] and [writeToJournal]. See [WriteConcern.writeTimeout]. + */ + writeTimeout: Duration? = null, + ) { + writeConcern(WriteConcern(acknowledgement, writeToJournal, writeTimeout)) + } + } -- 2.51.2 From 467f50d4c3d88376ea1b8dcd65650167431af107 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Ivan=20=E2=80=9CCLOVIS=E2=80=9D=20Canet?= Date: Sun, 13 Sep 2026 10:56:37 +0200 Subject: [PATCH 2/4] feat(driver-api): Improved MongoSyntaxException for insert operations --- .../kotlin/CoroutineMongoCollectionImpl.kt | 14 +++++-- .../src/commonMain/kotlin/ErrorHandling.kt | 22 +++++++++++ .../MultiplatformMongoCollectionImpl.kt | 22 ++++++----- .../src/jvmMain/kotlin/Exceptions.jvm.kt | 21 ++++++++++ .../jvmMain/kotlin/SyncMongoCollectionImpl.kt | 14 +++++-- .../kotlin/command/errors/Exceptions.kt | 24 +++++++++++- .../kotlin/command/errors/Exceptions.kt | 24 +++++++++++- .../operations/InsertOperations.test.kt | 39 ++++++++++++++++++- 8 files changed, 159 insertions(+), 21 deletions(-) diff --git a/driver-coroutines/src/jvmMain/kotlin/CoroutineMongoCollectionImpl.kt b/driver-coroutines/src/jvmMain/kotlin/CoroutineMongoCollectionImpl.kt index 88472526..f3489db8 100644 --- a/driver-coroutines/src/jvmMain/kotlin/CoroutineMongoCollectionImpl.kt +++ b/driver-coroutines/src/jvmMain/kotlin/CoroutineMongoCollectionImpl.kt @@ -154,6 +154,8 @@ private class CoroutineMongoCollectionImpl( ) } catch (e: com.mongodb.MongoWriteException) { throw e.toKtMongo(model, fullyQualifiedName) + } catch (e: com.mongodb.MongoCommandException) { + throw e.toKtMongo(model, fullyQualifiedName, factory) } } @@ -163,10 +165,14 @@ private class CoroutineMongoCollectionImpl( model.options.options() - inner.withWriteConcern(model.options).insertMany( - model.documents, - model.options.toJava(), - ) + try { + inner.withWriteConcern(model.options).insertMany( + model.documents, + model.options.toJava(), + ) + } catch (e: com.mongodb.MongoCommandException) { + throw e.toKtMongo(model, fullyQualifiedName, factory) + } } // endregion diff --git a/driver-multiplatform/src/commonMain/kotlin/ErrorHandling.kt b/driver-multiplatform/src/commonMain/kotlin/ErrorHandling.kt index 1a1a67b2..5afd3261 100644 --- a/driver-multiplatform/src/commonMain/kotlin/ErrorHandling.kt +++ b/driver-multiplatform/src/commonMain/kotlin/ErrorHandling.kt @@ -20,6 +20,7 @@ import opensavvy.ktmongo.bson.BsonType import opensavvy.ktmongo.bson.multiplatform.BsonDocument import opensavvy.ktmongo.dsl.command.Command import opensavvy.ktmongo.dsl.command.errors.MongoException +import opensavvy.ktmongo.dsl.command.errors.MongoSyntaxException import opensavvy.ktmongo.dsl.command.errors.MongoWriteException private class WriteErrorDataImpl( @@ -57,3 +58,24 @@ internal fun checkNoWriteErrors( errors = errors.map { WriteErrorDataImpl(it) }, ) } + +internal fun checkNoSyntaxErrors( + doc: BsonDocument, + command: Command, + collection: MultiplatformMongoCollection<*>, + server: MongoException.ServerAddress = collection.database.client.serverAddress, +) { + if (doc["ok"]?.decodeDouble() == 1.0) { + return // No errors, nothing to do + } + + throw MongoSyntaxException( + errorMessage = doc["errmsg"]?.decodeString() ?: "No error message were provided", + code = doc["code"]?.decodeInt32() ?: -1, + codeName = doc["codeName"]?.decodeString() ?: "", + fullResponse = doc, + command = command, + server = server, + namespace = collection.fullyQualifiedName, + ) +} diff --git a/driver-multiplatform/src/commonMain/kotlin/MultiplatformMongoCollectionImpl.kt b/driver-multiplatform/src/commonMain/kotlin/MultiplatformMongoCollectionImpl.kt index 0625aa29..dd8c4785 100644 --- a/driver-multiplatform/src/commonMain/kotlin/MultiplatformMongoCollectionImpl.kt +++ b/driver-multiplatform/src/commonMain/kotlin/MultiplatformMongoCollectionImpl.kt @@ -86,31 +86,33 @@ internal class MultiplatformMongoCollectionImpl( ) message as Message.OpMsg - check(message.body.document["ok"]?.decodeDouble() == 1.0) + checkNoSyntaxErrors(message.body.document, model, this) checkNoWriteErrors(message.body.document, model, this) } override suspend fun insertMany(documents: Iterable, options: InsertManyOptions.() -> Unit) { + val model = InsertMany( + context = database.client.context, + documents = documents.toList(), + documentType = type, + ).apply { + this.options.options() + } + val message = database.client.sendSingle( database.client.createDriverMessage { document { writeString("insert", name) writeString($$"$db", database.name) - InsertMany( - context = database.client.context, - documents = documents.toList(), - documentType = type, - ).apply { - this.options.options() - }.writeTo(this) + model.writeTo(this) } } ) message as Message.OpMsg - check(message.body.document["ok"]?.decodeDouble() == 1.0) { "Message is not OK: $message" } - check(message.body.document["writeErrors"] == null) { "Write errors occurred: $message" } + checkNoSyntaxErrors(message.body.document, model, this) + checkNoWriteErrors(message.body.document, model, this) } override fun filter(filter: FilterQuery.() -> Unit): MultiplatformMongoCollection = diff --git a/driver-shared-official/src/jvmMain/kotlin/Exceptions.jvm.kt b/driver-shared-official/src/jvmMain/kotlin/Exceptions.jvm.kt index dd9b4b90..e6937d07 100644 --- a/driver-shared-official/src/jvmMain/kotlin/Exceptions.jvm.kt +++ b/driver-shared-official/src/jvmMain/kotlin/Exceptions.jvm.kt @@ -16,10 +16,13 @@ package opensavvy.ktmongo.official +import opensavvy.ktmongo.bson.official.BsonFactory import opensavvy.ktmongo.dsl.LowLevelApi import opensavvy.ktmongo.dsl.command.Command import opensavvy.ktmongo.dsl.command.errors.MongoException +import opensavvy.ktmongo.dsl.command.errors.MongoSyntaxException import opensavvy.ktmongo.dsl.command.errors.MongoWriteException +import com.mongodb.MongoCommandException as OfficialMongoCommandException import com.mongodb.MongoWriteException as OfficialMongoWriteException import com.mongodb.ServerAddress as OfficialServerAddress @@ -60,3 +63,21 @@ fun OfficialMongoWriteException.toKtMongo( cause = this, ) } + +@LowLevelApi +fun OfficialMongoCommandException.toKtMongo( + command: Command, + namespace: String, + factory: BsonFactory, +): MongoSyntaxException { + return MongoSyntaxException( + errorMessage = errorMessage, + code = errorCode, + codeName = errorCodeName, + fullResponse = factory.readDocument(response), + server = serverAddress.toKtMongo(), + command = command, + namespace = namespace, + cause = this, + ) +} diff --git a/driver-sync/src/jvmMain/kotlin/SyncMongoCollectionImpl.kt b/driver-sync/src/jvmMain/kotlin/SyncMongoCollectionImpl.kt index 595ec8a0..905440ff 100644 --- a/driver-sync/src/jvmMain/kotlin/SyncMongoCollectionImpl.kt +++ b/driver-sync/src/jvmMain/kotlin/SyncMongoCollectionImpl.kt @@ -154,6 +154,8 @@ private class SyncMongoCollectionImpl( ) } catch (e: com.mongodb.MongoWriteException) { throw e.toKtMongo(model, fullyQualifiedName) + } catch (e: com.mongodb.MongoCommandException) { + throw e.toKtMongo(model, fullyQualifiedName, factory) } } @@ -163,10 +165,14 @@ private class SyncMongoCollectionImpl( model.options.options() - inner.withWriteConcern(model.options).insertMany( - model.documents, - model.options.toJava(), - ) + try { + inner.withWriteConcern(model.options).insertMany( + model.documents, + model.options.toJava(), + ) + } catch (e: com.mongodb.MongoCommandException) { + throw e.toKtMongo(model, fullyQualifiedName, factory) + } } // endregion diff --git a/dsl-template/src/commonMain/kotlin/command/errors/Exceptions.kt b/dsl-template/src/commonMain/kotlin/command/errors/Exceptions.kt index d4236475..f846c28c 100644 --- a/dsl-template/src/commonMain/kotlin/command/errors/Exceptions.kt +++ b/dsl-template/src/commonMain/kotlin/command/errors/Exceptions.kt @@ -16,6 +16,7 @@ package opensavvy.ktmongo.dsl.command.errors +import opensavvy.ktmongo.bson.BsonDocument import opensavvy.ktmongo.dsl.command.Command /** @@ -37,6 +38,27 @@ sealed class MongoException( } } +/** + * MongoDB refused to execute a command because it is malformed. + */ +class MongoSyntaxException( + val errorMessage: String, + val code: Int, + val codeName: String, + val fullResponse: BsonDocument, + val server: ServerAddress, + val command: Command, + val namespace: String, + cause: Throwable? = null, +) : MongoException( + message = buildString { + appendLine("$code $codeName • $errorMessage") + appendLine("\tat ${command::class.simpleName} $command") + append("\tat $server $namespace") + }, + cause = cause, +) + /** * A write operation failed. */ @@ -56,7 +78,7 @@ class MongoWriteException( } appendLine("\tat ${command::class.simpleName} $command") - appendLine("\tat $server $namespace") + append("\tat $server $namespace") }, cause = cause, ) { diff --git a/dsl/src/commonMain/kotlin/command/errors/Exceptions.kt b/dsl/src/commonMain/kotlin/command/errors/Exceptions.kt index 3fe8d000..c60275b9 100644 --- a/dsl/src/commonMain/kotlin/command/errors/Exceptions.kt +++ b/dsl/src/commonMain/kotlin/command/errors/Exceptions.kt @@ -19,6 +19,7 @@ package opensavvy.ktmongo.dsl.command.errors +import opensavvy.ktmongo.bson.BsonDocument import opensavvy.ktmongo.dsl.command.Command /** @@ -40,6 +41,27 @@ sealed class MongoException( } } +/** + * MongoDB refused to execute a command because it is malformed. + */ +class MongoSyntaxException( + val errorMessage: String, + val code: Int, + val codeName: String, + val fullResponse: BsonDocument, + val server: ServerAddress, + val command: Command, + val namespace: String, + cause: Throwable? = null, +) : MongoException( + message = buildString { + appendLine("$code $codeName • $errorMessage") + appendLine("\tat ${command::class.simpleName} $command") + append("\tat $server $namespace") + }, + cause = cause, +) + /** * A write operation failed. */ @@ -59,7 +81,7 @@ class MongoWriteException( } appendLine("\tat ${command::class.simpleName} $command") - appendLine("\tat $server $namespace") + append("\tat $server $namespace") }, cause = cause, ) { diff --git a/test/src/commonMain/kotlin/operations/InsertOperations.test.kt b/test/src/commonMain/kotlin/operations/InsertOperations.test.kt index 931494f2..559d8124 100644 --- a/test/src/commonMain/kotlin/operations/InsertOperations.test.kt +++ b/test/src/commonMain/kotlin/operations/InsertOperations.test.kt @@ -25,8 +25,11 @@ import opensavvy.ktmongo.bson.BsonDocument import opensavvy.ktmongo.bson.decode import opensavvy.ktmongo.bson.types.ObjectId import opensavvy.ktmongo.dsl.LowLevelApi +import opensavvy.ktmongo.dsl.command.InsertMany import opensavvy.ktmongo.dsl.command.InsertOne +import opensavvy.ktmongo.dsl.command.errors.MongoSyntaxException import opensavvy.ktmongo.dsl.command.errors.MongoWriteException +import opensavvy.ktmongo.dsl.options.WriteAcknowledgment import opensavvy.ktmongo.dsl.options.WriteConcern import opensavvy.ktmongo.dsl.path.Field import opensavvy.ktmongo.tests.api.collection @@ -78,7 +81,7 @@ fun SuiteDsl.verifyInsertOperations( ) } - test("Cannot insert two documents with the same ID") { + test("insertOne • Cannot insert two documents with the same ID") { val id = collection().newId() val alice = InsertOperationsUser( @@ -101,6 +104,40 @@ fun SuiteDsl.verifyInsertOperations( check(e.errors[0].code == 11000) check(e.namespace == collection().fullyQualifiedName) } + + test("insertOne • Cannot set invalid options") { + val e = checkThrows { + collection().insertOne( + InsertOperationsUser( + _id = collection().newId(), + name = "Alice", + ), + options = { + // MongoDB doesn't accept a query with this many nodes + writeConcern(WriteAcknowledgment.Nodes(100)) + } + ) + } + check(e.command is InsertOne<*>) + check(e.codeName == "FailedToParse") + } + + test("insertMany • Cannot set invalid options") { + val e = checkThrows { + collection().insertMany( + InsertOperationsUser( + _id = collection().newId(), + name = "Alice", + ), + options = { + // MongoDB doesn't accept a query with this many nodes + writeConcern(WriteAcknowledgment.Nodes(100)) + } + ) + } + check(e.command is InsertMany<*>) + check(e.codeName == "FailedToParse") + } } test("Decreased type safety") { -- 2.51.2 From 7a23626a363c21a47f071aad068400545268b77c Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Ivan=20=E2=80=9CCLOVIS=E2=80=9D=20Canet?= Date: Sun, 13 Sep 2026 11:01:57 +0200 Subject: [PATCH 3/4] feat(driver-multiplatform): Improved MongoSyntaxException for all operations --- .../MultiplatformMongoCollectionImpl.kt | 140 ++++++++++-------- 1 file changed, 75 insertions(+), 65 deletions(-) diff --git a/driver-multiplatform/src/commonMain/kotlin/MultiplatformMongoCollectionImpl.kt b/driver-multiplatform/src/commonMain/kotlin/MultiplatformMongoCollectionImpl.kt index dd8c4785..68a57ca5 100644 --- a/driver-multiplatform/src/commonMain/kotlin/MultiplatformMongoCollectionImpl.kt +++ b/driver-multiplatform/src/commonMain/kotlin/MultiplatformMongoCollectionImpl.kt @@ -66,7 +66,7 @@ internal class MultiplatformMongoCollectionImpl( document: Document, options: InsertOneOptions.() -> Unit, ) { - val model = InsertOne( + val command = InsertOne( context = database.client.context, document = document, documentType = type, @@ -80,18 +80,18 @@ internal class MultiplatformMongoCollectionImpl( writeString("insert", name) writeString($$"$db", database.name) - model.writeTo(this) + command.writeTo(this) } } ) message as Message.OpMsg - checkNoSyntaxErrors(message.body.document, model, this) - checkNoWriteErrors(message.body.document, model, this) + checkNoSyntaxErrors(message.body.document, command, this) + checkNoWriteErrors(message.body.document, command, this) } override suspend fun insertMany(documents: Iterable, options: InsertManyOptions.() -> Unit) { - val model = InsertMany( + val command = InsertMany( context = database.client.context, documents = documents.toList(), documentType = type, @@ -105,14 +105,14 @@ internal class MultiplatformMongoCollectionImpl( writeString("insert", name) writeString($$"$db", database.name) - model.writeTo(this) + command.writeTo(this) } } ) message as Message.OpMsg - checkNoSyntaxErrors(message.body.document, model, this) - checkNoWriteErrors(message.body.document, model, this) + checkNoSyntaxErrors(message.body.document, command, this) + checkNoWriteErrors(message.body.document, command, this) } override fun filter(filter: FilterQuery.() -> Unit): MultiplatformMongoCollection = @@ -122,43 +122,47 @@ internal class MultiplatformMongoCollectionImpl( MultiplatformMongoAggregationPipelineImpl(this, PipelineChainLink(context)) override suspend fun create(options: CreateCollectionOptions.() -> Unit) { + val command = CreateCollection( + context = database.client.context, + ).apply { + this.options.options() + } + val message = database.client.sendSingle( database.client.createDriverMessage { document { writeString("create", name) writeString($$"$db", database.name) - CreateCollection( - context = database.client.context, - ).apply { - this.options.options() - }.writeTo(this) + command.writeTo(this) } } ) message as Message.OpMsg - check(message.body.document["ok"]?.decodeDouble() == 1.0) + checkNoSyntaxErrors(message.body.document, command, this) } override suspend fun drop(options: DropOptions.() -> Unit) { + val command = Drop( + context = database.client.context, + ).apply { + this.options.options() + } + val message = database.client.sendSingle( database.client.createDriverMessage { document { writeString("drop", name) writeString($$"$db", database.name) - Drop( - context = database.client.context, - ).apply { - this.options.options() - }.writeTo(this) + command.writeTo(this) } } ) message as Message.OpMsg - check(message.body.document["ok"]?.decodeDouble() == 1.0) + checkNoSyntaxErrors(message.body.document, command, this) } @OptIn(DangerousMongoApi::class) @@ -206,21 +210,23 @@ internal class MultiplatformMongoCollectionImpl( } override suspend fun countEstimated(): Long { + val command = Count( + context = database.client.context, + ) + val message = database.client.sendSingle( database.client.createDriverMessage { document { writeString("count", name) writeString($$"$db", database.name) - Count( - context = database.client.context, - ).writeTo(this) + command.writeTo(this) } } ) message as Message.OpMsg - check(message.body.document["ok"]?.decodeDouble() == 1.0) + checkNoSyntaxErrors(message.body.document, command, this) return message.body.document["n"] ?.decodeLong(message) @@ -228,45 +234,49 @@ internal class MultiplatformMongoCollectionImpl( } override suspend fun deleteOne(options: DeleteOneOptions.() -> Unit, filter: FilterQuery.() -> Unit) { + val command = DeleteOne( + context = database.client.context, + ).apply { + this.options.options() + this.filter.filter() + } + val message = database.client.sendSingle( database.client.createDriverMessage { document { writeString("delete", name) writeString($$"$db", database.name) - DeleteOne( - context = database.client.context, - ).apply { - this.options.options() - this.filter.filter() - }.writeTo(this) + command.writeTo(this) } } ) message as Message.OpMsg - check(message.body.document["ok"]?.decodeDouble() == 1.0) + checkNoSyntaxErrors(message.body.document, command, this) } override suspend fun deleteMany(options: DeleteManyOptions.() -> Unit, filter: FilterQuery.() -> Unit) { + val command = DeleteMany( + context = database.client.context, + ).apply { + this.options.options() + this.filter.filter() + } + val message = database.client.sendSingle( database.client.createDriverMessage { document { writeString("delete", name) writeString($$"$db", database.name) - DeleteMany( - context = database.client.context, - ).apply { - this.options.options() - this.filter.filter() - }.writeTo(this) + command.writeTo(this) } } ) message as Message.OpMsg - check(message.body.document["ok"]?.decodeDouble() == 1.0) + checkNoSyntaxErrors(message.body.document, command, this) } override fun find(): MultiplatformMongoIterable = @@ -288,7 +298,7 @@ internal class MultiplatformMongoCollectionImpl( ) override suspend fun updateMany(options: UpdateOptions.() -> Unit, filter: FilterQuery.() -> Unit, update: UpdateQuery.() -> Unit): UpdateOperations.UpdateResult { - val model = UpdateMany( + val command = UpdateMany( context = database.client.context, ).apply { this.options.options() @@ -301,13 +311,13 @@ internal class MultiplatformMongoCollectionImpl( document { writeString("update", name) writeString($$"$db", database.name) - model.writeTo(this) + command.writeTo(this) } } ) message as Message.OpMsg - check(message.body.document["ok"]?.decodeDouble() == 1.0) + checkNoSyntaxErrors(message.body.document, command, this) // TODO handle unacknowledged updates @@ -315,7 +325,7 @@ internal class MultiplatformMongoCollectionImpl( } override suspend fun updateOne(options: UpdateOptions.() -> Unit, filter: FilterQuery.() -> Unit, update: UpdateQuery.() -> Unit): UpdateOperations.UpdateResult { - val model = UpdateOne( + val command = UpdateOne( context = database.client.context, ).apply { this.options.options() @@ -328,13 +338,13 @@ internal class MultiplatformMongoCollectionImpl( document { writeString("update", name) writeString($$"$db", database.name) - model.writeTo(this) + command.writeTo(this) } } ) message as Message.OpMsg - check(message.body.document["ok"]?.decodeDouble() == 1.0) + checkNoSyntaxErrors(message.body.document, command, this) // TODO handle unacknowledged updates @@ -342,7 +352,7 @@ internal class MultiplatformMongoCollectionImpl( } override suspend fun upsertOne(options: UpdateOptions.() -> Unit, filter: FilterQuery.() -> Unit, update: UpsertQuery.() -> Unit): MultiplatformMongoCollection.UpsertResult { - val model = UpsertOne( + val command = UpsertOne( context = database.client.context, ).apply { this.options.options() @@ -355,13 +365,13 @@ internal class MultiplatformMongoCollectionImpl( document { writeString("update", name) writeString($$"$db", database.name) - model.writeTo(this) + command.writeTo(this) } } ) message as Message.OpMsg - check(message.body.document["ok"]?.decodeDouble() == 1.0) + checkNoSyntaxErrors(message.body.document, command, this) // TODO handle unacknowledged updates @@ -369,7 +379,7 @@ internal class MultiplatformMongoCollectionImpl( } override suspend fun replaceOne(options: ReplaceOptions.() -> Unit, filter: FilterQuery.() -> Unit, document: Document) { - val model = ReplaceOne( + val command = ReplaceOne( context = database.client.context, document = document, documentType = type, @@ -383,17 +393,17 @@ internal class MultiplatformMongoCollectionImpl( document { writeString("update", name) writeString($$"$db", database.name) - model.writeTo(this) + command.writeTo(this) } } ) message as Message.OpMsg - check(message.body.document["ok"]?.decodeDouble() == 1.0) + checkNoSyntaxErrors(message.body.document, command, this) } override suspend fun repsertOne(options: ReplaceOptions.() -> Unit, filter: FilterQuery.() -> Unit, document: Document) { - val model = RepsertOne( + val command = RepsertOne( context = database.client.context, document = document, documentType = type, @@ -407,13 +417,13 @@ internal class MultiplatformMongoCollectionImpl( document { writeString("update", name) writeString($$"$db", database.name) - model.writeTo(this) + command.writeTo(this) } } ) message as Message.OpMsg - check(message.body.document["ok"]?.decodeDouble() == 1.0) + checkNoSyntaxErrors(message.body.document, command, this) } override suspend fun findOneAndUpdate(options: UpdateOptions.() -> Unit, filter: FilterQuery.() -> Unit, update: UpdateQuery.() -> Unit): Document? { @@ -421,7 +431,7 @@ internal class MultiplatformMongoCollectionImpl( } override suspend fun bulkWrite(options: BulkWriteOptions.() -> Unit, filter: FilterQuery.() -> Unit, operations: BulkWrite.() -> Unit) { - val model = BulkWrite( + val command = BulkWrite( context = database.client.context, documentType = type, globalFilter = filter, @@ -442,17 +452,17 @@ internal class MultiplatformMongoCollectionImpl( } } - model.writeTo(this) + command.writeTo(this) } } ) message as Message.OpMsg - check(message.body.document["ok"]?.decodeDouble() == 1.0) + checkNoSyntaxErrors(message.body.document, command, this) } override suspend fun updateManyWithPipeline(options: UpdateOptions.() -> Unit, filter: FilterQuery.() -> Unit, update: UpdateWithPipelineQuery.() -> Unit): UpdateOperations.UpdateResult { - val model = UpdateManyWithPipeline( + val command = UpdateManyWithPipeline( context = database.client.context, ).apply { this.options.options() @@ -465,13 +475,13 @@ internal class MultiplatformMongoCollectionImpl( document { writeString("update", name) writeString($$"$db", database.name) - model.writeTo(this) + command.writeTo(this) } } ) message as Message.OpMsg - check(message.body.document["ok"]?.decodeDouble() == 1.0) + checkNoSyntaxErrors(message.body.document, command, this) // TODO handle unacknowledged updates @@ -479,7 +489,7 @@ internal class MultiplatformMongoCollectionImpl( } override suspend fun updateOneWithPipeline(options: UpdateOptions.() -> Unit, filter: FilterQuery.() -> Unit, update: UpdateWithPipelineQuery.() -> Unit): UpdateOperations.UpdateResult { - val model = UpdateOneWithPipeline( + val command = UpdateOneWithPipeline( context = database.client.context, ).apply { this.options.options() @@ -492,13 +502,13 @@ internal class MultiplatformMongoCollectionImpl( document { writeString("update", name) writeString($$"$db", database.name) - model.writeTo(this) + command.writeTo(this) } } ) message as Message.OpMsg - check(message.body.document["ok"]?.decodeDouble() == 1.0) + checkNoSyntaxErrors(message.body.document, command, this) // TODO handle unacknowledged updates @@ -506,7 +516,7 @@ internal class MultiplatformMongoCollectionImpl( } override suspend fun upsertOneWithPipeline(options: UpdateOptions.() -> Unit, filter: FilterQuery.() -> Unit, update: UpdateWithPipelineQuery.() -> Unit): MultiplatformMongoCollection.UpsertResult { - val model = UpsertOneWithPipeline( + val command = UpsertOneWithPipeline( context = database.client.context, ).apply { this.options.options() @@ -519,13 +529,13 @@ internal class MultiplatformMongoCollectionImpl( document { writeString("update", name) writeString($$"$db", database.name) - model.writeTo(this) + command.writeTo(this) } } ) message as Message.OpMsg - check(message.body.document["ok"]?.decodeDouble() == 1.0) + checkNoSyntaxErrors(message.body.document, command, this) check(message.body.document["writeErrors"] == null) { "There were write errors: ${message.body.document["writeErrors"]}" } // TODO handle unacknowledged updates -- 2.51.2 From e31120c89b67365decb828cbca312c4fed39b902 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Ivan=20=E2=80=9CCLOVIS=E2=80=9D=20Canet?= Date: Sun, 13 Sep 2026 11:07:32 +0200 Subject: [PATCH 4/4] feat(driver-api): Improved MongoWriteException for insertMany --- .../kotlin/CoroutineMongoCollectionImpl.kt | 2 ++ .../src/jvmMain/kotlin/Exceptions.jvm.kt | 20 +++++++++++++++++ .../jvmMain/kotlin/SyncMongoCollectionImpl.kt | 2 ++ .../operations/InsertOperations.test.kt | 22 +++++++++++++++++++ 4 files changed, 46 insertions(+) diff --git a/driver-coroutines/src/jvmMain/kotlin/CoroutineMongoCollectionImpl.kt b/driver-coroutines/src/jvmMain/kotlin/CoroutineMongoCollectionImpl.kt index f3489db8..13b75ea0 100644 --- a/driver-coroutines/src/jvmMain/kotlin/CoroutineMongoCollectionImpl.kt +++ b/driver-coroutines/src/jvmMain/kotlin/CoroutineMongoCollectionImpl.kt @@ -170,6 +170,8 @@ private class CoroutineMongoCollectionImpl( model.documents, model.options.toJava(), ) + } catch (e: com.mongodb.MongoBulkWriteException) { + throw e.toKtMongo(model, fullyQualifiedName) } catch (e: com.mongodb.MongoCommandException) { throw e.toKtMongo(model, fullyQualifiedName, factory) } diff --git a/driver-shared-official/src/jvmMain/kotlin/Exceptions.jvm.kt b/driver-shared-official/src/jvmMain/kotlin/Exceptions.jvm.kt index e6937d07..26da51b7 100644 --- a/driver-shared-official/src/jvmMain/kotlin/Exceptions.jvm.kt +++ b/driver-shared-official/src/jvmMain/kotlin/Exceptions.jvm.kt @@ -22,6 +22,7 @@ import opensavvy.ktmongo.dsl.command.Command import opensavvy.ktmongo.dsl.command.errors.MongoException import opensavvy.ktmongo.dsl.command.errors.MongoSyntaxException import opensavvy.ktmongo.dsl.command.errors.MongoWriteException +import com.mongodb.MongoBulkWriteException as OfficialMongoBulkWriteException import com.mongodb.MongoCommandException as OfficialMongoCommandException import com.mongodb.MongoWriteException as OfficialMongoWriteException import com.mongodb.ServerAddress as OfficialServerAddress @@ -64,6 +65,25 @@ fun OfficialMongoWriteException.toKtMongo( ) } +@LowLevelApi +fun OfficialMongoBulkWriteException.toKtMongo( + command: Command, + namespace: String, +): MongoWriteException { + return MongoWriteException( + server = serverAddress.toKtMongo(), + command = command, + errors = writeErrors.map { + WriteErrorDataImpl( + code = it.code, + message = it.message, + ) + }, + namespace = namespace, + cause = this, + ) +} + @LowLevelApi fun OfficialMongoCommandException.toKtMongo( command: Command, diff --git a/driver-sync/src/jvmMain/kotlin/SyncMongoCollectionImpl.kt b/driver-sync/src/jvmMain/kotlin/SyncMongoCollectionImpl.kt index 905440ff..5258a365 100644 --- a/driver-sync/src/jvmMain/kotlin/SyncMongoCollectionImpl.kt +++ b/driver-sync/src/jvmMain/kotlin/SyncMongoCollectionImpl.kt @@ -170,6 +170,8 @@ private class SyncMongoCollectionImpl( model.documents, model.options.toJava(), ) + } catch (e: com.mongodb.MongoBulkWriteException) { + throw e.toKtMongo(model, fullyQualifiedName) } catch (e: com.mongodb.MongoCommandException) { throw e.toKtMongo(model, fullyQualifiedName, factory) } diff --git a/test/src/commonMain/kotlin/operations/InsertOperations.test.kt b/test/src/commonMain/kotlin/operations/InsertOperations.test.kt index 559d8124..9404d0d5 100644 --- a/test/src/commonMain/kotlin/operations/InsertOperations.test.kt +++ b/test/src/commonMain/kotlin/operations/InsertOperations.test.kt @@ -105,6 +105,28 @@ fun SuiteDsl.verifyInsertOperations( check(e.namespace == collection().fullyQualifiedName) } + test("insertMany • Cannot insert two documents with the same ID") { + val id = collection().newId() + + val alice = InsertOperationsUser( + _id = id, + name = "Alice", + ) + + val bob = InsertOperationsUser( + _id = id, + name = "Bob", + ) + + val e = checkThrows { + collection().insertMany(alice, bob) + } + check((e.command as? InsertMany<*>)?.documents == listOf(alice, bob)) + check(e.errors.size == 1) + check(e.errors[0].code == 11000) + check(e.namespace == collection().fullyQualifiedName) + } + test("insertOne • Cannot set invalid options") { val e = checkThrows { collection().insertOne(