Effect-agnostic MongoDB client and repository layer for Scala 3 — any runtime, any BSON codec, AST-free by default.
libraryDependencies += "org.mongo4s" %% "mongo4s-cats" % "4.0.0"
On another runtime, swap the module: mongo4s-zio,
mongo4s-kyo or mongo4s-rapid. Codecs and repositories are separate
modules — see Quick start.
No hardcoded cats-effect or fs2. The runtime (cats-effect / ZIO / Kyo / rapid) and the
BSON codec (your own derivation, or medeia / zio-bson / calypso) are independent modules, each wired in through
a
given import. core depends on neither.
mongo4s wraps the official mongodb-driver-reactivestreams directly — no mongo4cats
underneath. A type-safe Field/Filter/Update builder replaces string-keyed
queries, PrimaryKey turns an entity into single- or compound-key lookups, and
BaseMongoRepository gives you CRUD/batch operations over a collection for free. Every piece is
interpretable against a real MongoDB and an in-memory fake, so repositories are unit-testable
without a running database.
Pick a runtime. The default codec (bson-direct) derives straight from your case class — AST-free, no
third-party codec library, no extra dependency beyond mongo4s-core itself:
libraryDependencies ++= Seq(
"org.mongo4s" %% "mongo4s-cats" % "4.0.0", // mongo4s-core + cats-effect integration
"org.mongo4s" %% "mongo4s-bson-direct" % "4.0.0", // ast-free bson codecs
"org.mongo4s" %% "mongo4s-bson-cats-data" % "4.0.0", // if you need NonEmptyList etc. codec instances
"org.mongo4s" %% "mongo4s-repositories" % "4.0.0", // if you need auto-generated CRUD repository ops for your model
)
import cats.effect.{IO, IOApp}
import mongo4s.cats.{CatsStream, MongoClientResource}
import mongo4s.bson.direct.WireCodec
import mongo4s.{Field, PrimaryKey}
import mongo4s.repositories.BaseMongoRepository
import mongo4s.cats.CatsInstances.given
final case class User(id: String, name: String, age: Int) derives WireCodec
object User:
given PrimaryKey[User, String] = PrimaryKey.single("id")(_.id)
object Main extends IOApp.Simple:
def run: IO[Unit] =
MongoClientResource.fromConnectionString[IO]("mongodb://localhost:27017").use { client =>
for
db <- client.getDatabase("myapp")
collection <- db.getDirectCollection[User]("users")
users = BaseMongoRepository(collection)
_ <- users.insertOne(User("1", "Alice", 30))
alice <- users.findOne("1")
adults <- users.findByFilter(Field.of[User, Int](_.age).gte(18))
yield ()
}
Swap mongo4s-cats for mongo4s-zio / mongo4s-kyo /
mongo4s-rapid and the matching *Instances.given import to change runtime — the body
of the for is identical. What changes is how you take the client: cats gives you a
Resource to .use, ZIO a scoped ZIO and kyo a Scope effect,
and rapid takes the body as a function. BsonEncoder/BsonDecoder for built-in types
(String, Int,
Option, List, Vector, Set, Seq, …) resolve with
no import at all.
Already have a model on medeia, zio-schema, or calypso? Swap derives WireCodec +
getDirectCollection for derives MedeiaDocumentCodec/etc. + getCollection
(and BaseMongoRepository.create(db, "users") instead of constructing it from a collection
directly) — see BSON codecs below for all four backends.
Already have a MongoClientSettings built elsewhere (connection pool tuning, read/write concerns,
TLS, credentials, …)? Use MongoClientResource.fromSettings instead of
fromConnectionString. MongoClient.fromClient/fromSettings/fromConnectionString
give you the same thing unwrapped, if you'd rather own .close yourself. A driver
CodecRegistry is not how you plug a codec into mongo4s, though — see
Codecs and the driver's registry:
val settings = MongoClientSettings.builder().applyConnectionString(ConnectionString("mongodb://localhost:27017")).build()
MongoClientResource.fromSettings[IO](settings).use { client => ... }
MongoClientSettings sets read/write concerns for the whole client. To narrow them,
MongoDatabase and MongoCollection both carry withReadConcern/
withWriteConcern/withReadPreference, each returning a new handle
rather than mutating the one you have — so per-operation control is just chaining. Concerns inherit database →
collection as they do in the driver, and a derived collection keeps the codec it was opened with, including
the WireCodec that getDirectCollection registers.
val durable = collection.withWriteConcern(WriteConcern.MAJORITY)
durable.insertOne(user) // majority-acknowledged
collection.insertOne(other) // still the client's default
For more examples see:
examples/src/main/scala/mongo4s/examples
— a shared domain model (opaque types, enums, nested case classes) run through every runtime/codec combination
(cats+medeia, ZIO+zio-bson, kyo+medeia, rapid+calypso), a repository example covering all three
BaseMongoRepository construction styles against bson-direct, and a sessions/transactions + typed
aggregation pipeline example on cats+medeia. Most of what this page shows is compiled there too —
ReadmeSnippets.scala
walks the same ground section by section, so if the API moves and these docs don't, CI fails.
MongoClient[F, S] → MongoDatabase[F, S] → MongoCollection[F, S, A] mirror
the
driver's own hierarchy, wrapped in your effect F[_] and stream type S[_].
Field.of[E, A](_.someField) is a macro that reads a field selector at compile time — no strings, no
reflection — and gives you a typed path to build filters, updates, and sorts:
val adults = Field.of[User, Int](_.age).gte(18)
val named = Field.of[User, String](_.name).equalTo("Jenna") && adults
val setAge = Field.of[User, Int](_.age).set(27)
val city = Field.of[Order, String](_.address.city).equalTo("Barcelona") // dotted paths from nested selectors
Each segment is checked against the case class it is selected from, so _.name.length and
_.items.head.sku are compile errors rather than paths that render fine and match nothing.
Selector-derived names are spelled through the collection's FieldNaming. Names that are already
what the document stores — a map key, an array position, _id — go through at,
/ or Field.stored, and are used verbatim:
val totals = Field.of[Order, Map[String, Int]](_.totals)
val eur = totals.at("EUR") // "totals" is renamed, "EUR" is not
val first: Field[Order, Item] = itemsField / "0" // any other stored segment
val id = Field.stored[Order, ObjectId]("_id")
ageField.gte(13) && ageField.lte(19) // $and
ageField.notIn(List(40, 41, 42)) // $nin
tagsField.contains("urgent") // array membership
tagsField.containsAll(List("a", "b")) // $all
tagsField.hasSize(3) // $size
scoreField.hasType(BsonTypeName.Int) // $type — an enum, not a string
Filter.text[User]("scala") // $text
Filter.expr[User](someBsonDocument) // $expr
itemsField.elemMatch(skuField.equalTo("abc") && qtyField.gt(2)) // one element satisfies both
Filter.and/or fold Filter.all and Filter.none away rather
than emitting a one-element $and. That is what makes an empty list safe:
field.in(Nil) is Filter.none, not "match everything", and a chain
a && b && c renders as one flat $and. The comparisons have symbolic
aliases where they read better — ===, =!=, >, >=,
<, <=. An operator that only means something for one kind of field compiles only
on it: regex on text, mod on a number, hasSize on a collection.
ageField.set(31) // $set
ageField.inc(1) // $inc — also $mul, $min, $max
tagsField.push("vip") // $push — also $pull, $addToSet, and the $each variants
tagsField.pushAll(List("vip"), PushOptions.default[String].sortedAscending.withSlice(10))
Update.setOnInsert(nameField, "") // $setOnInsert, for merge-style upserts
Update.rename(oldField, newField) // $rename — also currentDate, popFirst/popLast
Numeric operators only apply to numeric fields, so Update.inc(nameField, 1) does not compile. An
Option-typed field takes the unwrapped value — scoreField.inc(5L) on a
Field[User, Option[Long]] — because None has no numeric encoding, and the obvious
stand-in, $inc by zero, is a write that quietly does nothing. Update.combine merges
operators of the same name into one sub-document; Update.Raw carries operators the AST does not
model and merges the same way, without mutating the document you handed it. One operator document holds one
change per path, so an update that changes a path twice is refused when it renders rather than losing the
first change, and so is one whose paths overlap across operators. bulkWrite takes a
Seq[WriteCommand[E]] — InsertOne/ReplaceOne/UpdateOne/UpdateMany/DeleteOne/DeleteMany,
carrying the same Filter and Update values, run in the order given unless
ordered = false. A BulkWriteFailed names each failing command by its position, and its
result is what the others wrote.
An update that computes from the document's own values is a pipeline, not operators.
Update.pipeline takes the stages the server accepts in one — $addFields,
$project, $unset, $replaceRoot or a Stage.raw — and every
update call takes it:
val exclaimed = Update.pipeline(
Stage.addFields[User]("name" -> BsonDocument("$concat", BsonArray(java.util.List.of(BsonString("$name"), BsonString("!")))))
)
collection.updateMany(Filter.all, exclaimed)
The positional operators are path segments, so / builds them. $[] updates every
element and needs nothing else; $[identifier] updates only what an array filter selects,
and those go in arrayFilters on
updateOne/updateMany/findOneAndUpdate:
val lowQty: Field[Order, Int] = Field.of[Order, List[Item]](_.items) / "$[low]" / "qty"
val elementQty: Field[Item, Int] = Field.stored("low.qty")
collection.updateOne(orderId.equalTo("1"), Update.set(lowQty, 100), UpdateOptions.default.withArrayFilters(Seq(elementQty.lt(3))))
Array-filter paths are written against the stored document — the identifier is not a field of
your model, so build them with Field.stored and spell the field names as they appear on the wire.
Everything reached through / is a stored segment too, so neither the identifier nor
$[] is touched by the collection's FieldNaming.
updateOne/updateMany take an UpdateOptions, replaceOne a
ReplaceOptions. Both are immutable, built by chaining from default, and exist so new
options stay additive — a write method's signature never has to grow another parameter again. Each operation
family gets its own type rather than one shared bag, so upsert cannot be set on something that
has no upsert to give.
UpdateOptions.default // nothing set
UpdateOptions.upsert // shorthand for default.withUpsert
UpdateOptions.upsert.withArrayFilters(Seq(elementQty.lt(3)))
collection.deleteMany(filter, DeleteOptions.default.withCollation(caseInsensitive))
collection.count(filter, CountOptions.default.withSkip(20).withLimit(10))
What each type carries today: collation, hint and comment everywhere;
bypassDocumentValidation on the two that write whole documents; limit,
skip and maxTime on CountOptions.
The findOneAnd* trio works the same way, except its options carry the entity type because
sort and projection do. returnUpdated defaults to true, and
returningPrevious asks for the document as it was:
findOneAndUpdate(filter, update, FindOneAndUpdateOptions.default[User].returningPrevious). They
carry collation, hint, maxTime and comment too, and the
update and replace ones bypassDocumentValidation.
find returns a builder; nothing is sent until first, all,
stream or attempting:
collection
.find(adults)
.sort(Sort.asc(nameField))
.projection(Projection.empty[User].include(ageField).withoutId)
.skip(20).limit(10)
.hint(indexKeys)
.collation(Collation.builder().locale("en").build())
.maxTime(5.seconds)
.batchSize(100)
.comment("adults page 3")
.all
filter narrows what is already there; sort, projection,
skip and limit replace. The same options are on aggregate, which adds
allowDiskUse, and collation/maxTime/batchSize on
distinct.
Projection keeps inclusion and exclusion apart, since MongoDB rejects a projection mixing them.
They are separate types, so chaining exclude onto an inclusion projection does not compile.
_id is the exception, via withoutId.
The numbers are checked where you write them, because MongoDB reads the edge cases as something else: a
limit or batchSize below one — 0 is no limit at all — a negative skip,
and a maxTime under a millisecond, which truncates to no time limit. distinct over a
field typed as an array answers with the elements' type, because that is what the server answers with.
Writes report what actually happened. UpdateResult carries matchedCount,
modifiedCount and upsertedId, which is the only way to tell "matched but unchanged"
from "nothing matched", or to recover the id an upsert generated:
val result = collection.updateOne(filter, update, UpdateOptions.upsert)
result.map(r => if r.wasUpserted then r.upsertedId else None)
upsertedId is the document's _id, not your primary key, and it is reported only when
the upsert actually inserted. The two coincide only when the PrimaryKey names _id
(storedId, or WithId); on a key over any other field the server generates an
ObjectId and reports that, matching K in neither value nor type. What lands in
_id is the entity codec's decision, and PrimaryKey only decides what the repository
filters on — let them disagree and upsertedId is neither.
Under an unacknowledged write concern (w=0) the server sends nothing back, so every write result
comes back empty — zero counts, None for the ids. That is indistinguishable from a write that
matched nothing, which is the trade w=0 makes. Nothing throws, so the write path stays usable;
if you need to tell the two apart, do not use w=0.
find(...).all fails the whole query if any document does not decode. A collection is rarely
written by one version of one service, though, and when it isn't, a document that doesn't fit is a fact about
the data rather than a reason to lose the rest of the page. attempting reports each document
separately:
val readable: IO[List[User]] =
collection.find().attempting.all.map(_.collect { case Right(user) => user })
val everything: S[DecodeResult[User]] = collection.find().attempting.stream
DecodeResult[A] is Either[BsonError, A]. It's on find,
aggregate and distinct, and as watchAttempting on a collection's
watch — watchAsAttempting at the client and database level. Transport errors still
fail the effect — only decoding is made per-document. On a change stream the result sits on the document,
inside the event — ChangeEvent[DecodeResult[A]] — so an event whose document fails still arrives
with its resume token, and a consumer can checkpoint past it.
A failed write does not arrive as a driver exception. Every failure the driver reports is translated into
MongoError, so branching on one is a pattern match rather than a code comparison:
import mongo4s.MongoError
collection.insertOne(user).recoverWith {
case MongoError.DuplicateKey(_) => collection.replaceOne(byEmail, user).void
case MongoError.WriteConflict(_) => retry
}
| Case | Raised for |
|---|---|
DuplicateKey |
a unique index violation |
WriteConflict |
two transactions touching the same document |
ExecutionTimeout |
maxTime expired, or the server's own limit |
Unauthorized |
the credentials do not allow the operation |
Unavailable |
the socket failed, or no server could be selected |
BulkWriteFailed |
one bulk command, several failures — failures keeps them all,
duplicateKeys filters
|
Failed |
anything else the server reported, with its code intact |
Every case carries cause, the driver's own exception, so nothing is lost — along with
code, labels and hasLabel, which is how withTransaction
decides what to retry. Translation happens on the Publisher before any runtime sees it, so
IO, Task, KIO, rapid.Task and every stream report the same
type for the same failure. Errors mongo4s raises itself are untouched:
BsonError.DecodingFailure for a document that does not fit, RsBridgeError for the
stream bridge.
PrimaryKey[E, K] turns an entity into a key-based filter — a single field, a native
_id (ObjectId or your own encoder), or a compound key of any width:
given PrimaryKey[User, String] = PrimaryKey.single("id")(_.id)
given PrimaryKey[Note, ObjectId] = PrimaryKey.storedId(_.id) // keys on "_id" — see WithId
given PrimaryKey[Order, (userId: String, seq: Int)] =
PrimaryKey.compound(o => (userId = o.userId, seq = o.seq), FieldNaming.snakeCase)
PrimaryKey.id(_.id) is single("id") spelled short, for the common case. A compound
key is a named tuple, so its labels are the field names: no second list of strings to keep in
step with the extractors, and no arity ceiling. The optional FieldNaming spells those labels the
way the collection stores them, so userId above is written as user_id while the key
still reads as ordinary Scala at every call site —
users.findOne((userId = "u1", seq = 3)).
Field names are given separately from the extractors so the key knows them without a key value — that is what
lets repository.ensureKeyIndex build the unique index that makes the key a key. Without one, two
concurrent upserts on the same key can both miss and both insert. These are stored
names, used verbatim: under snakeCase write "user_id", not "userId".
inFilter on a compound key produces an $or of $ands; on a single field
it's
a plain $in — and an empty key list always produces Filter.none, so
findMany(Nil)/deleteMany(Nil) are safe no-ops instead of matching every document.
client.withTransaction opens a session, runs the body in a transaction on it, and closes it — on
every path. The session is given implicitly to everything inside, so a collection or repository call joins the
transaction without being told to:
client.withTransaction {
users.insertOne(User("2", "Bob", 41)) // Option[ClientSession] is already given here
}
It commits on success and rolls back on failure and on cancellation: Effect[F]
carries a guaranteeCase that sees how the action ended, so an interrupted transaction does not
linger on the server until it is reaped. A rollback that itself fails is attached as a suppressed exception
rather than replacing the error that caused it.
Reusing one already-open session across more than one transaction? The same behaviour is on the session itself — it commits and rolls back the same way, but leaves the session's own lifetime to you:
import mongo4s.withTransaction // the extension on ClientSession
for
session <- client.startSession
_ <- session.withTransaction(users.insertOne(User("2", "Bob", 41)))
_ <- session.withTransaction(users.insertOne(User("3", "Carol", 29))) // same session, second transaction
_ <- IO.delay(session.close())
yield ()
Every MongoClient/MongoDatabase/MongoCollection/Repository
method that reaches the server accepts the same (using session: Option[ClientSession] = None), so a
BaseMongoRepository call can join an existing transaction the same way a raw
collection call can.
MongoSession.startTransaction/commitTransaction/abortTransaction
are there for the fully manual path, where you pass the session explicitly with (using
Some(session))
at each call site — nothing is automatic there, including the rollback. Transactions need a replica set or
sharded cluster.
withTransaction also retries, the way the driver's own does. A
TransientTransactionError means nothing was committed, so the whole transaction — body
included — is started again; an UnknownTransactionCommitResult means the commit may already
have landed, so only the commit is asked again rather than the work redone. Both are bounded by one deadline
taken before the first attempt, 120 seconds by default:
import mongo4s.operations.TransactionOptions
client.withTransaction(users.insertOne(User("2", "Bob", 41)), TransactionOptions.default.withRetryTimeout(10.seconds))
client.withTransaction(body, TransactionOptions.withoutRetries) // report the first failure, as 2.x did
TransactionOptions also carries the transaction's own readConcern,
writeConcern, readPreference and maxCommitTime. Errors the server did
not label are never retried — a failure in your own code fails the transaction immediately. The deadline
is measured with Effect.monotonic, which the cats and zio backends
override with their runtime's own clock, so TestControl/TestClock can drive the
retry window in a test without waiting on a real one.
Two rules follow the driver specification. Between attempts of the whole transaction there is a short random
pause, up to 5 ms growing by half each attempt to a 500 ms ceiling, so clients that conflicted do not retry in
lockstep; and a commit that ran out of maxCommitTime is not asked again. A call takes its session
when it is written, not when it runs, so val write = users.insertOne(user) built outside
the block commits on its own even when it runs inside.
A session is useful without a transaction too — to read your own writes causally, or one point in time from a
snapshot. client.withSession gives it to every call in the body and closes it however the body
ends:
import mongo4s.operations.SessionOptions
client.withSession(
for
before <- users.count()
rows <- users.find().all // the same point in time as the count
yield (before, rows),
SessionOptions.snapshot,
)
aggregate takes a Seq[Stage[A]] — a typed pipeline-stage AST, mirroring
Filter/Update/Sort, instead of raw BsonDocuments — built
with
the same Field.of selectors:
import mongo4s.operations.{Sort, Stage}
val pipeline = Seq(
Stage.matching(Field.of[User, Int](_.age).gte(18)),
Stage.sortBy(Sort.asc(Field.of[User, String](_.name))),
Stage.limit(10),
)
val adults: IO[List[User]] = collection.aggregate[User](pipeline).all
Grouping uses typed accumulators:
import mongo4s.operations.Accumulator
val byAge = Seq(
Stage.groupBy(Field.of[User, Int](_.age))(
"count" -> Accumulator.count[User],
"names" -> Accumulator.push(Field.of[User, String](_.name)),
)
)
Stage covers
$match/$project/$sort/$limit/$skip/$count/$unwind/$lookup/$graphLookup/$group/$addFields/$unset/$replaceRoot/$facet/$sample/$sortByCount/$unionWith/$bucket/$bucketAuto/$densify/$setWindowFields/$geoNear/$out/$merge,
with Stage.raw(document) as the escape hatch for anything else.
Stage.geoNear(locationField, Geometry.Point(13.405, 52.52), "distance") // first stage only
Stage.unset(ageField).and(tagsField)
Stage.sortByCount(cityField)
Stage.bucketAuto(nameField, 4)("count" -> Accumulator.count[Person])
Stage.bucketAutoRounded(ageField, 5, BucketGranularity.R10)() // numeric fields only
$geoNear is $near as a stage: it sorts by distance and writes each distance, in
metres, into the field you name, and GeoNearOptions carries its query, distance bounds and
multiplier. The server takes it only as a pipeline's first stage, over a geospatial index.
bucketAutoRounded rounds the boundaries to a preferred-number series, which the server only does
for numbers, so it compiles only on a numeric field. $unset takes more fields with
and, because fields of unrelated types have no single Field type to share.
$lookup has two forms: Stage.lookup is the equality join,
Stage.lookupWith
takes a sub-pipeline over the foreign collection with let binding values from the outer document.
Stage.graphLookup walks a self-referencing hierarchy. Both take a second type parameter, because
the sub-pipeline and the connect fields belong to the foreign collection rather than to
A.
$out and $merge write the pipeline's result into a collection, and both take options
that stay out of the way until you need them — with none set they render as the bare collection name:
Stage.out[Order]("archive", OutOptions.default.inDatabase("cold"))
Stage.merge[Order]("archive", MergeOptions.default
.onFields(List("user_id", "seq"))
.whenMatched(MergeOptions.WhenMatched.KeepExisting)
.whenNotMatched(MergeOptions.WhenNotMatched.Insert))
on names fields of the target collection, so it is a list of stored names rather
than Field values — the target's shape is not A. A single field renders as a string
and several as an array, as the server expects.
Two things about the typing. A is the pipeline's starting document type and never changes
down the pipeline — a $group or $project invents a new shape, but the stages after it
are still typed against A, and the real output type is stated once, at aggregate[B].
And a field belonging to a stage's own output — _id after a $group, a facet's counter
— has no Field[A, _] to name it, so it goes through Stage.raw:
val buckets = Seq(
Stage.groupBy(ageField)("count" -> Accumulator.count[User]),
Stage.raw[User](BsonDocument("$sort", BsonDocument("_id", BsonInt32(1)))), // "_id" is the group's, not User's
)
For an output shape you'd rather not model, BsonDocumentCodec[BsonDocument] is in scope by
default, so aggregate[BsonDocument] just works. For a bson-direct entity the output
codec comes from DocumentCodecBridge.toDocumentCodec[User] — aggregate and
distinct decode through BsonDocumentCodec, not WireCodec.
aggregateDirect[B] is the AST-free counterpart: it asks for a WireCodec[B] rather
than a BsonDocumentDecoder[B] and, on a collection opened with getDirectCollection,
decodes the pipeline's output straight off the wire with no BsonDocument in between — the
same trip find already makes there:
final case class ByAge(_id: Int, total: Int) derives WireCodec
collection.aggregateDirect[ByAge](Seq(Stage.groupBy(ageField)("total" -> Accumulator.count[User]))).all
It carries the direct path's strictness with it — every field the output type declares has to be
present, so a model with a field the pipeline does not produce is a decode error rather than a default. On a
collection opened with getCollection the same call still works, with the WireCodec
bridged to a document codec: the method is about the codec you have, not about which constructor you used.
explain answers the question a typed filter otherwise leaves open — did it use an index? It
is on find and on aggregate, runs the query the builder had already produced, and
returns the server's plan as a BsonDocument:
collection.find(adults).explain(ExplainVerbosity.EXECUTION_STATS).map(ExplainSummary.of)
// summary.indexes List("age_1")
// summary.usedIndex true
// summary.scannedCollection false
// summary.sortedInMemory false — a SORT stage means the order was not index-provided
// summary.execution Some(ExecutionSummary(returned, keysExamined, docsExamined, durationMillis))
The plan stays a document because an explain result is not one shape: the same
aggregate(...).explain() returns a top-level queryPlanner when MongoDB collapses the
pipeline into a plain query and a top-level stages array when it does not, over an open set of
stage names, with a layer per shard on a sharded cluster — which is why the server stamps
explainVersion on it. ExplainSummary types the part that is stable, reading
by walking for field names rather than a fixed path, so it answers the same way for a find, a pipeline the
server optimised into a find, and a pipeline it kept as stages. When it cannot recognise something it says
nothing rather than guessing, and the full document is still there.
Index[E] is built from the same field selectors, and carries the options MongoDB attaches to an
index:
import mongo4s.operations.Index
collection.createIndex(Index.ascending(nameField).descending(ageField).named("name_age"))
collection.createIndex(Index.unique(idField))
collection.createIndex(Index.ascending(createdAtField).expiringAfter(30.days)) // TTL
collection.createIndex(Index.ascending(ageField).where(ageField.gte(18))) // partial
collection.createIndex(Index.empty[User].text(bioField).withSparse)
collection.createIndex(Index.hashed(idField)) // sharding
collection.createIndex(Index.geo2DSphere(Field.stored[User, Any]("location"))) // also geo2D
collection.createIndex(Index.ascending(ageField).withHidden) // ignored by the planner
collection.createIndex(Index.ascending(Field.stored[User, Any]("$**"))) // wildcard
collection.listIndexes // F[List[BsonDocument]], as the server reports them
collection.dropIndex("name_age")
collection.createIndexes(List(Index.ascending(nameField), Index.descending(ageField))) // one command
collection.dropIndexes // all but _id's
collection.renameCollection("people_archive", dropTarget = true)
Keys are ordered — a compound index is only usable by queries that respect that order.
createIndex returns the name the server gave it, and is idempotent, so it's safe on every start.
createIndexes returns the names in the order given, and the server builds them together, so if one
cannot be built, none is. renameCollection refuses to replace an existing collection unless
dropTarget says to, and leaves the handle it was called on pointing at the old, now empty, name.
expireAfterSeconds is a whole number on the server, so a sub-second TTL is rejected rather than
silently truncated to "expire immediately"; when immediately is what you mean — each document carries its own
expiry time — expiringAtFieldTime asks for a TTL of zero on purpose. From a repository,
ensureKeyIndex builds the unique index the PrimaryKey describes without you restating
its fields.
watch exists at all three levels, matching the driver's own scope hierarchy —
MongoClient.watch (the whole deployment, needs a replica set/sharded cluster),
MongoDatabase.watch (one database), and MongoCollection.watch (one collection).
Every event is a ChangeEvent[A], not a bare BsonDocument:
final case class ChangeEvent[A](
operationType: OperationType, // com.mongodb's own enum — INSERT/UPDATE/DELETE/...
namespace: Option[Namespace], // the database and collection it happened in
destination: Option[Namespace], // where a rename moved the collection to
documentKey: Option[BsonDocument],
fullDocument: Option[A],
fullDocumentBeforeChange: Option[A],
updateDescription: Option[UpdateDescription], // updated/removed field paths, for UPDATE events
resumeToken: BsonDocument,
clusterTime: Option[BsonTimestamp],
wallTime: Option[BsonDateTime], // when the change was applied, in server wall-clock time
splitEvent: Option[SplitEvent], // set only on a fragment of an event too large for one message
)
MongoCollection.watch decodes through the collection's own codec (ChangeEvent[A]);
MongoClient/MongoDatabase.watch span more than one document shape so they default to
ChangeEvent[BsonDocument] — use watchAs[A] when you know the events all decode the
same way. At those two levels namespace is how an event says which collection it came from.
Everything else is WatchOptions[A]:
import mongo4s.changestream.WatchOptions
collection.watch(
WatchOptions
.resumeAfter[User](token) // shorthand for default[User].resumingAfter(token)
.withFullDocument(FullDocument.DEFAULT)
.withMaxAwaitTime(2.seconds)
.withBatchSize(64)
)
WatchOptions.default[User].startingAfter(token) // the other two starting points, same builder
WatchOptions.default[User].startingAt(timestamp)
fullDocument defaults to UPDATE_LOOKUP, not the server's own default — MongoDB fills
the document in only for inserts and replaces, which leaves the most common question ("what does this document
look like now?") unanswered on updates. It is always None for deletes; there is no document left
to look up. fullDocumentBeforeChange needs pre-images enabled on the collection.
resumingAfter and startingAfter are alternatives — each clears the other, since the
server rejects a stream carrying both resumeAfter and startAfter.
resumeToken is on every event so a consumer that dies mid-stream can restart from just after the
last event it actually handled, rather than from now:
collection.watch().evalTap(handle).evalTap(e => saveToken(e.resumeToken))
// later, on restart:
collection.watch(WatchOptions.resumeAfter[User](savedToken))
WatchOptions.pipeline filters the change stream itself, and matches against the change
event's own shape ({operationType, fullDocument, ns, ...}), not the collection's document
shape. A Field.of path is therefore the wrong tool: it would render "age" where the
event needs "fullDocument.age". Use Stage.raw, or Field.stored for a
path under fullDocument:
val insertsOnly = WatchOptions
.default[User]
.withPipeline(Seq(Stage.raw(BsonDocument("$match", BsonDocument("operationType", BsonString("insert"))))))
collection.watch(insertsOnly)
A change stream never completes on its own — take from it, or interrupt it. A document that does not decode
ends it with an error; watchAttempting — watchAsAttempting at the client and
database level — reports the failure and carries on instead, which is what you want for a long-lived
subscription. Its elements are ChangeEvent[DecodeResult[A]], so the event around a bad document,
resume token included, still arrives. UPDATE_LOOKUP is a second read per update event and the
whole document on the wire; a consumer that only needs what changed can ask for FullDocument.DEFAULT
and read updateDescription. All three levels need a replica set or sharded cluster.
Every codec ultimately produces a BsonDocumentCodec[A] (entity ⇄ org.bson.BsonDocument)
or, for the AST-free path below, a WireCodec[A]. mongo4s never registers a global
CodecProvider, so backends never collide inside one process.
| Module | Backend | Notes |
|---|---|---|
mongo4s-bson-medeia |
medeia | derives BsonDocumentCodec |
mongo4s-bson-zio |
zio-bson | zio.bson.BsonCodec (add zio-schema-bson yourself to derive one from a Schema) |
mongo4s-bson-calypso |
calypso | hand-written forProductN codecs |
mongo4s-bson-direct |
mongo4s itself | WireCodec[A], AST-free — see below |
mongo4s-bson-cats-data |
cats-core | BsonEncoder/BsonDecoder/WireCodec for
NonEmptyList/Chain/NonEmptyVector/NonEmptySet/NonEmptyMap,
WireCodec for Ior
|
WireCodec
medeia/zio-bson/calypso all build an intermediate org.bson.BsonValue tree before it ever reaches
the
driver. WireCodec[A] skips that: derived via Mirror, it writes straight to the
driver's
own streaming BsonWriter/BsonReader — the same low-level SPI jsoniter-scala uses for
JSON — no BsonDocument is ever built, on either side. No third-party codec dependency needed;
bson-direct is transitively pulled in by mongo4s-core.
import mongo4s.bson.direct.WireCodec
final case class Address(city: String, zip: String) derives WireCodec
final case class Person(id: String, name: String, tags: List[String], address: Address) derives WireCodec
sealed trait Shape derives WireCodec
object Shape:
final case class Circle(radius: Double) extends Shape derives WireCodec
final case class Rectangle(width: Double, height: Double) extends Shape derives WireCodec
Products, Option, Either, nested case classes, sealed traits/enums (via a
_type discriminator field, written first), and self-/mutually-recursive types all derive directly —
recursive derivation defers behind a lazy val internally so a type's own given never
forces itself mid-construction. Anything with an existing BsonEncoder/BsonDecoder
bridges automatically, at the cost of one BsonValue per field instead of zero.
List/Vector/Seq/Set/Array write a real BSON
array, and Map[String, A] a real BSON document keyed by its own keys —
String being the key type isn't a limitation of the general mechanism, it's just what BSON's
own field names are. Every other Iterable collection with a
scala.collection.Factory (Queue, ArraySeq, ListSet,
LazyList, SortedSet/TreeSet given an Ordering, …) gets an
array-shaped WireCodec too, generically — no dedicated given needed per type.
getDirectCollection registers the derived codec with the driver via
CodecRegistries.fromCodecs(...), so insert/find/replace/update/delete/bulkWrite decode straight to
A — genuinely zero BsonDocument construction on the hot path, not just a thinner
bridge:
import mongo4s.bson.direct.WireCodec
import mongo4s.repositories.BaseMongoRepository
final case class Person(id: String, name: String, age: Int) derives WireCodec
object Person:
given PrimaryKey[Person, String] = PrimaryKey.single("id")(_.id)
for
db <- client.getDatabase("myapp")
collection <- db.getDirectCollection[Person]("people")
repo = BaseMongoRepository(collection)
_ <- repo.insertOne(Person("1", "bob", 30))
yield ()
The trade for that speed is strictness: derivation requires every modelled field to be present, unless its
decoder supplies a default — which Option does. So a projection that drops a field the
entity declares cannot be read back through a direct collection. Use getCollection with a
BsonDocumentCodec, or model the projected shape as its own type, when you need partial reads.
Strictness extends to BSON numeric types. A Long field is read with readInt64, an
Int with readInt32, a Double with readDouble — the exact
type the encoder writes. So a document that stores 42 as an Int32 (written by
mongosh, by a $inc, or by another client) decodes fine through
getCollection, whose BsonDecoder[Long] accepts any whole number, and
fails through getDirectCollection. The failure is an ordinary
BsonError, so attempting reports it per document like any other decode error. If a
collection holds mixed numeric widths for the same field, read it through getCollection.
aggregate/distinct still go through
BsonDocumentCodec/BsonDecoder
on a direct collection (their output shape isn't A, and they're not the hot path); everything
else — Filter/Update/Field construction — is identical regardless of
which codec backs the collection.
By default derives WireCodec writes field and discriminator names exactly as they appear in the
Scala source. To match an existing collection's naming convention (snake_case, say), bring a
WireCodecConfig into scope before deriving — it reuses the same FieldNaming the query
layer already uses, rather than a second, independent naming mechanism:
import mongo4s.bson.direct.WireCodecConfig
given WireCodecConfig = WireCodecConfig.SnakeCase
final case class Person(firstName: String, lastName: String) derives WireCodec // writes "first_name"/"last_name"
Field naming and discriminator naming (the _type value for sealed traits/enums) are independently
configurable: start from WireCodecConfig.Default (or SnakeCase) and chain
withFieldNaming, withDiscriminatorNaming, withEncodeEmptyCasesAsString,
withOmitNoneFields — the same shape WatchOptions uses, and the reason new derivation
options can be added in a minor release without breaking binary compatibility. getDirectCollection
takes no naming of its own: the codec
is what writes the field names, so it is what decides how a query spells them, and a derived
WireCodec — a product, an enum or a sealed trait, or one mapped from it with imap/
iemap — reports the WireCodecConfig it was built with.
getCollection still takes the parameter, because a BsonDocumentCodec from
medeia, calypso or zio-bson carries its own naming configuration mongo4s cannot read; there it still has
to match the codec. A collection has one naming, so derive the entity and what it contains under one
WireCodecConfig; a nested type derived under another is reached with Field.stored.
encodeEmptyCasesAsString = true writes a parameterless case as a bare BSON string instead of a
{"_type": "..."} document — only safe when nested inside another document, not at the root.
A field holding None is left out of the document entirely rather than stored as an explicit
null — the key costs nothing on disk, and a collection of mostly-empty optional fields gets
materially smaller. Reads are unaffected either way: a missing field decodes to None, and a
document already carrying an explicit null still decodes to None, so a collection
written before this keeps working unchanged.
final case class Contact(name: String, email: Option[String]) derives WireCodec
Contact("bob", None) // {"name": "bob"} — no "email" key at all
{field: null} as a filter still matches, since MongoDB reads that predicate as "null or missing",
but $exists: true no longer matches a None and a sparse index on that field no longer
includes those documents — WireCodecConfig.Default.withOmitNoneFields(false) restores the explicit
null where either matters. The flag governs fields of a derived product only: a None
inside an array is still written as null, because dropping an element would shift every position
after it, and Update.set(field, None) still writes null because
$set: null and $unset are different intents (field.unset is the latter).
A nested Option does not survive a round-trip on either codec path, with or without the flag:
Some(None) is written as null and reads back as None, because BSON has
one null and both layers map onto it. To tell "absent" from "present but empty", model it as an
enum or a wrapper case class rather than Option[Option[A]].
Most types don't need derives WireCodec at all — WireCodec[A] is just
WireEncoder[A] with WireDecoder[A], and both halves compose the same way
BsonEncoder/BsonDecoder already do elsewhere in mongo4s:
import mongo4s.bson.direct.WireCodec
enum Provider(val value: String):
case Stripe extends Provider("stripe")
case Adyen extends Provider("adyen")
object Provider:
def from(value: String): Option[Provider] = Provider.values.find(_.value == value)
given WireCodec[Provider] =
WireCodec[String].iemap(raw => from(raw).toRight(s"Unsupported provider: $raw"))(_.value)
imap (total in both directions) and iemap (decode can fail — reported as
BsonError.InvalidValue) replace a hand-rolled WireCodec.instance(...) for the
common case of one type wrapping another: no BsonWriter/BsonReader calls to write
by hand, and the wrap/unwrap logic lives in exactly one place instead of being duplicated across the encode
and decode sides. Both keep what the codec underneath knows: a wrapper over an Option leaves the
empty value out of the document and reads the missing field back as that value, and a derived entity wrapped
to validate what it reads keeps its naming.
Resolution order matters when a type qualifies for more than one instance. An existing
BsonEncoder/BsonDecoder pair — hand-written, or coming from
mongo4s-bson-medeia/-zio/-calypso — is bridged into a
WireCodec and takes precedence over automatic Mirror derivation, so a case class
carrying those instances resolves without needing derives WireCodec. Writing
derives WireCodec on the type still wins over both: it puts a concrete instance in the companion,
which is more specific than either generic given.
An opaque type over a primitive (opaque type UserId = String) typically needs a
WireCodec[UserId] for getDirectCollection and a
BsonEncoder[UserId] for Field.of[...].equalTo/PrimaryKey.single — two
typeclasses, normally two independently hand-written codecs. ScalarWireCodec[A] — the type
every primitive in WirePrimitiveInstances actually has — closes that gap:
import mongo4s.bson.direct.ScalarWireCodec
import mongo4s.bson.BsonEncoder
opaque type UserId = String
object UserId:
def apply(value: String): UserId = value
extension (id: UserId) def value: String = id
given ScalarWireCodec[UserId] = ScalarWireCodec[String].imap(UserId.apply)(_.value)
given BsonEncoder[UserId] = summon[ScalarWireCodec[UserId]].toBsonEncoder
toBsonEncoder builds one BsonValue per call — cheap on the query-construction path
BsonEncoder actually runs on, unlike materializing one for every field of every document, which
is exactly what WireCodec exists to avoid in the first place. ScalarWireCodec is
deliberately narrower than WireCodec: a derived case class or sum type
(derives WireCodec) is genuinely document-shaped, so it's never typed as
ScalarWireCodec — asking for one (ScalarWireCodec[SomeCaseClass]) fails to compile
instead of misbehaving at runtime.
WireCodec[Either[A, B]] derives given WireCodec[A] and WireCodec[B],
flat like a derived sealed trait ({"_type": "Circle", "radius": 2.0}) rather than wrapped in a
Left/Right envelope — the discriminator is each branch's own runtime type name,
and a document-shaped branch (a case class) has its fields inlined directly instead of nested under a
"value" key:
final case class ValidationError(message: String)
final case class Approved(reference: String)
// Right(Approved("abc")) -> {"_type": "Approved", "reference": "abc"}
// Left(ValidationError("bad input")) -> {"_type": "ValidationError", "message": "bad input"}
val result: Either[ValidationError, Approved] = ...
A scalar branch (no fields to inline, e.g. a bare String) keeps a "value" field,
since a bare BSON scalar can't also carry the discriminator in the same slot. The two branches must resolve
to distinguishable type names — Either[Foo, Foo], or two differently-named types that happen to
share a runtime class name once generics are erased, throws when the codec is summoned rather than risking a
silent wrong-branch decode later. So does a branch with a field of its own named _type, which
inlined next to the discriminator would overwrite it.
The discriminator is ClassTag[A].runtimeClass.getSimpleName — for
Int/Long/Double/Boolean/etc, that's the JVM
primitive's name ("int", not "Integer"), since ClassTag[Int]'s
runtimeClass is the primitive class, not the boxed one. Not a bug, just worth knowing if one
branch is a bare numeric/boolean type.
cats.data support
NonEmptyList/Chain/NonEmptyVector/NonEmptySet/NonEmptyMap
get BsonEncoder/BsonDecoder and WireCodec instances from
mongo4s-bson-cats-data, delegating to the already-existing
List/Vector/Set/Map
instances rather than reimplementing BSON encoding:
import mongo4s.bson.catsdata.CatsDataBsonInstances.given // or CatsDataWireInstances.given for bson-direct
import cats.data.NonEmptyList
final case class Team(members: NonEmptyList[String]) derives MedeiaDocumentCodec // or your codec backend's own derives
NonEmptySet/NonEmptyMap additionally need a cats.Order for their
element/key type — the same Order you'd already need to construct one of these types directly.
Ior[A, B] gets a WireCodec too — flat and discriminated by each branch's own type
name, the same idea as bson-direct's own Either[A, B] above, not
"Left"/"Right"/"Both". Both doesn't have a single "own"
type — it holds an A and a B at once — so its discriminator is the two names
joined ("String+Foo"), with each side nested under its own "left"/"right"
key rather than inlined, to avoid a silent field-name collision if A and B happen
to share a field.
The driver resolves codecs from a process- or client-wide CodecRegistry. mongo4s does not: a codec
is resolved per collection, from the BsonDocumentCodec[A] or WireCodec[A] in implicit
scope at the getCollection / getDirectCollection call. That is what lets medeia,
zio-bson, calypso and bson-direct coexist in one process without colliding.
So a CodecRegistry you set on MongoClientSettings is not where an
entity codec belongs — mongo4s never asks it for one, on either path:
getCollection path, mongo4s asks the driver for BsonDocument and does its
own encode/decode. So do listIndexes, listCollections, runCommand,
aggregate, distinct and watch — all of them read
BsonDocument/BsonValue. Your registry is never asked about A.
getDirectCollection path, mongo4s registers the derived WireCodec[A]
ahead of the client's registry, so a Codec[A] registered there does not
silently shadow the codec the collection was opened with.
That does not make it inert, though. It is still load-bearing in two ways:
BsonDocument codec. Since mongo4s
reads and writes BsonDocument everywhere, a registry that replaces the defaults
instead of extending them fails every single operation with
CodecConfigurationException: Can't find a codec for … BsonDocument. Always build yours as
CodecRegistries.fromRegistries(yours, MongoClientSettings.getDefaultCodecRegistry).
underlying.
MongoClient/MongoDatabase/MongoCollection each expose the driver
object they wrap, and there the driver's rules apply in full —
collection.underlying.withDocumentClass(classOf[Foo]) resolves Foo from your
registry, mongo4s not involved.
So: set a registry for driver-level defaults, or for a type you handle through the escape hatch — not to route the entities mongo4s already has codecs for.
BaseMongoRepository[F, S, E, K] implements count/find/insert/upsert/update/delete/bulkWrite,
batched
by batchSize (default 500), over any MongoCollection[F, S, E], from either codec path:
BaseMongoRepository(collection) // over a collection you already have
BaseMongoRepository.create[F, S, E, K](db, "collection") // reads return the document as stored, _id included
BaseMongoRepository.createDirect[F, S, E, K](db, "collection") // the same over getDirectCollection and a WireCodec
BaseMongoRepository.withoutId[F, S, E, K](db, "collection") // strips _id from reads — for entities that do not model it
BaseMongoRepository.objectId[F, S, E](db, "collection") // WithId[ObjectId, E], auto _id round-trip
create is the default, and it is the right one more often than it looks: none of the four
codec backends rejects a document carrying an _id the entity does not model — medeia,
zio-bson, calypso and a bridged WireCodec all decode it and ignore the field. Reach for
withoutId when a codec you wrote is strict about unknown fields, or when you would rather
not carry _id over the wire at all. Never for an entity keyed on _id: the projection
strips the key, and it comes back missing.
Reads and writes that take a condition rather than a key come in two spellings, and the name says
which: …ByField takes a field and a value, …ByFilter takes a whole
Filter. Each …ByField call is exactly its …ByFilter counterpart over
field.equalTo(value):
users.findByField(nameField, "alice") // List[User]
users.getByField(nameField, "alice") // the same, as a stream
users.updateByField(ageField, 30, birthday) // every 30-year-old
users.deleteByField(ageField, 30) // likewise
insertOne returns the entity's K and insertMany the K of
every entity, in insertion order and across batches. A repository carries a PrimaryKey[E, K], so
it names what it just inserted rather than handing back a raw BsonValue to decode.
MongoCollection has no key type and keeps returning the driver's
InsertOneResult/InsertManyResult — that is where the server's own
_id still lives.
A batched write that fails in a later batch is reported as the whole write, the way the driver reports its own
batches: BulkWriteFailed names the failing command by its position in the sequence you passed, and
its result counts what the earlier batches wrote. findMany(keys) answers the way a
query does — in no particular order, without the keys that matched nothing. A findOneAndUpdate
projection you pass wins over the repository's defaultProjection.
Paging goes through Page:
import mongo4s.repositories.Page
users.findByFilter(adults, Page.sortedBy(Sort.asc(nameField)).skipping(20).taking(10))
users.getByFilter(adults, Page.first(100)) // same, as a stream
Skip-based paging re-scans what it skips, so it degrades on deep pages — and a row deleted earlier
shifts everything after it, so a walk can miss rows it never saw. findPage avoids both by asking
for what comes after a key rather than for an offset. The repository already knows the key's fields,
so it builds both halves itself — the sort that gives the page an order, and the comparison that steps
past the cursor, lexicographic for a compound key. ensureKeyIndex builds the index that makes the
whole walk a range scan.
val firstPage = users.findPage(100) // KeysetPage(items, next)
val nextPage = users.findPage(100, after = firstPage.next)
users.findPage(100, after = Some(lastKey), filter = adults) // narrowed, still paged by key
users.findPage(20, descending = true) // newest key first
next is the key the following page starts after, or None on the last page — the
repository reads one row past the page to know which. The key has to be unique, one BSON type per field, and
indexed.
ensureKeyIndex builds the unique index the PrimaryKey describes.
WithId[Id, E] wraps an entity with a separately-typed id (type Oid[E] = WithId[ObjectId,
E]) and ships its own PrimaryKey/BsonDocumentCodec instances, for entities
that don't carry their own id field — it lives in mongo4s-core, so it's available whether or not
you use the repository layer.
Like MongoCollection, every Repository method takes (using session:
Option[ClientSession] = None) — a repository call can join a transaction the same
way a raw collection call can. BaseMongoRepository is open: adding a
domain query means declaring it against collection/Filter/Field
directly, not reimplementing what's already there.
For unit tests, FakeMongoCollection (in mongo4s-testkit, a published module — add it
as "org.mongo4s" %% "mongo4s-testkit" % "4.0.0" % Test) implements
MongoCollection in memory — the exact
same Filter/Update/Field AST the real driver interprets is interpreted
against an in-memory buffer instead, so repository logic is testable without a running MongoDB. Filters,
every update operator, sorting, paging, projections, distinct and a subset of
aggregate are simulated, and so are unique indexes, listIndexes as the server
describes them, the way a bulk write reports a duplicate, and how the server compares numbers of different
widths — each checked against a real server by a parity spec. watch, explain,
renameCollection, $text, $expr, Filter.Raw, a pipeline
update, a collation and an update carrying arrayFilters throw
UnsupportedOperationException naming what was asked for, rather than quietly answering wrong.
It streams through the runtime's RsBridge, reading the collection when the stream runs, so
FakeMongoCollection[IO, S, User](codec) and FakeRepository[IO, S, User, String]() are
all it takes.
Replace-based upserts — what upsert/upsertMany go through — insert on a miss the way
the server does; an update-based UpdateOptions.upsert that matches nothing throws
instead of
guessing what the operators would have built.
Each runtime module provides given Effect[F] (sequencing, failure, and a finalizer that sees how
the action ended) and given RsBridge[F, S] (Reactive-Streams Publisher →
F/S):
| Module | Effect | Stream | Notes |
|---|---|---|---|
mongo4s-cats |
any F with cats.effect.kernel.Async |
fs2.Stream | via fs2.interop.reactivestreams |
mongo4s-zio |
zio.Task | zio.stream.ZStream | via zio-interop-reactivestreams |
mongo4s-kyo |
A < (Async & Abort[Throwable]) |
kyo.Stream | via kyo-reactive-streams |
mongo4s-rapid |
rapid.Task | rapid.Stream |
Beyond the members a runtime has to supply, Effect carries defaults it may override:
monotonic and sleep — the pause between transaction retries — read the
runtime's own clock and scheduler on all four backends instead of System.nanoTime and a blocked
thread.
MongoClient.fromClient/fromSettings/fromConnectionString return a bare
F[MongoClient[F, S]] — you own calling .close. Every runtime module also ships
MongoClientResource with the same three constructor names, wrapped in that runtime's own
resource-safety idiom rather than one type copy-pasted across all four:
// cats — cats.effect.Resource
import mongo4s.cats.MongoClientResource
MongoClientResource.fromConnectionString[IO]("mongodb://localhost:27017").use { client =>
// logic
}
// zio — ZIO.acquireRelease, released when the enclosing ZIO.scoped block exits
import mongo4s.zio.MongoClientResource
ZIO.scoped:
for
client <- MongoClientResource.fromConnectionString("mongodb://localhost:27017")
...
yield ()
// kyo — Scope.acquireRelease, released by Scope.run
import mongo4s.kyo.MongoClientResource
Scope.run:
for
client <- MongoClientResource.fromConnectionString("mongodb://localhost:27017")
...
yield ()
// rapid has no Resource/Scope type — Task.guarantee is its only finalizer primitive, attached to an
// already-known computation — so this is bracket-shaped (a `use` callback) instead of a composable value
import mongo4s.rapid.MongoClientResource
MongoClientResource.fromConnectionString("mongodb://localhost:27017") { client =>
// logic
}
RsBridgeConfig controls how a driver Publisher becomes your F and
S. A given of your own overrides the default wherever an RsBridge is
summoned:
import mongo4s.RsBridgeConfig
given RsBridgeConfig = RsBridgeConfig(
bufferSize = 512, // outstanding demand for streaming reads
timeout = Some(5.seconds), // per non-streaming operation; unset by default
strictSingleResult = false, // fail on a second result instead of taking the first
)
bufferSize bounds memory against a fast cursor; it does not apply to all, which asks
for everything by definition. timeout is a backstop for a cursor that stops signalling entirely —
the driver has its own timeouts, and streams are deliberately excluded, since a change stream sitting idle is
working rather than stuck. Operations that expect at most one document read two, not the whole cursor — enough
to notice a second result under strictSingleResult, and no more;
AggregateQuery.first pushes a $limit into the pipeline for the same reason.
Published for Scala 3 under org.mongo4s:
"org.mongo4s" %% "mongo4s-<module>" % "4.0.0"
| Kind | Module | Notes |
|---|---|---|
| core | mongo4s-core |
client/database/collection, Field/Filter/Update, PrimaryKey, WithId |
| bson | mongo4s-bson-core |
the scalar + document codec seam |
mongo4s-bson-direct |
WireCodec — AST-free, no third-party dependency | |
mongo4s-bson-cats-data |
cats.data (NonEmptyList/Chain/NonEmptyVector/NonEmptySet/NonEmptyMap/Ior) instances | |
mongo4s-bson-medeia |
bridges medeia | |
mongo4s-bson-zio |
bridges zio-bson (add zio-schema-bson yourself for the Schema-derived route) | |
mongo4s-bson-calypso |
bridges calypso | |
| runtime | mongo4s-cats |
cats-effect 3 + fs2 |
mongo4s-zio |
ZIO 2 + zio-streams | |
mongo4s-kyo |
kyo 1.0.0-RC6 | |
mongo4s-rapid |
rapid | |
| repositories | mongo4s-repositories |
BaseMongoRepository, Repository, Page |
| testkit | mongo4s-testkit |
FakeMongoCollection, FakeRepository — in-memory doubles for unit tests |
What the artifacts promise:
versionScheme := "semver-spec" describes what the artifacts promise.
Effect and RsBridge will carry default implementations, so
implementing either typeclass yourself keeps working across minor releases.
4.0.0 is a breaking release. An Update that changes one path twice is refused
rather than keeping only the last change. findPage returns a KeysetPage with the key
the next page starts after. watchAttempting yields ChangeEvent[DecodeResult[A]], so an
event whose document fails keeps its resume token. distinct over an array answers with its
elements, and regex, mod, hasSize and popFirst/
popLast compile only on the fields they apply to. The testkit's fake streams through the
RsBridge instead of taking an emit function, and its listIndexes answers
the way the server does. MongoCollection gains members an implementation of your own has to
provide, and Stage cases an exhaustive match has to cover.
3.0.0 was the last breaking release before it: WireCodec stopped deriving itself for
any case class it met, driver exceptions arrive as MongoError, Repository's condition
methods say what they take, and its inserts return the entity's key.
The full migration guides, 3.x → 4.0 and 2.x → 3.0, including the corrected behaviours, are in COMPATIBILITY.md; what changed in each release is in CHANGELOG.md; what is not covered yet and how it will land is in ROADMAP.md.
3.0.0 moved the whole build onto Scala 3.9 LTS and dropped
3.3 LTS. TASTy is not forward compatible, and that is a break MiMa cannot see,
which is why it took a major rather than a minor. What it buys is one Scala version across every
module: mongo4s-bson-calypso, mongo4s-kyo and mongo4s-rapid,
which used to be pinned to a fast-release 3.8, are on the LTS line with the rest.
mongo4s-kyo depends on a kyo release candidate and sits outside the binary compatibility
promise the other artifacts make until kyo reaches 1.0.0 final. Compiling any module that touches
kyo requires JDK 25: kyo's Frame macro runs
inside the compiler and its class files target Java 25. Since 3.0.0 the build enforces
that for every module and refuses to load on an older JDK, naming the version it found,
rather than letting four modules compile and the fifth fail with
UnsupportedClassVersionError: class file version 69.0. One JDK across the build is
deliberate — the alternative is publishing artifacts compiled against different JDKs. It is a
build-time requirement; the published artifacts are unaffected.
Six JMH harnesses in benchmarks/, against mongo4cats as a reference point. Everything below was
measured in one sitting on an Apple M4 Max (14 cores, 36 GB), macOS 26.6.2,
OpenJDK 25.0.4.1, Scala 3.9.0, JMH 1.37, with mongo:8.2 in Docker on
localhost:27018. Directional ballparks from one laptop, not hardware-independent authorities.
org.bson.BsonDocument| Codec | Encode ops/s | Decode ops/s |
|---|---|---|
| mongo4s-bson-direct (via DocumentCodecBridge) | ~2.74M | ~2.07M |
| calypso (forProductN) | ~4.22M | ~3.67M |
| medeia (derives) | ~2.42M | ~2.72M |
| zio-bson (zio-schema derived) | ~1.31M | ~2.71M |
| mongo4cats-zio-json | ~1.25M | ~1.12M |
| mongo4cats-circe | ~1.21M | ~795k |
Every mongo4s codec bridge beats both mongo4cats codecs on decode, by 1.8× at the narrowest
(bson-direct against zio-json) and 4.6× at the widest (calypso against circe) — the cost of
case class ↔ circe/zio-json ↔ mongo4cats.Bson ↔ org.Bson instead of straight to
org.bson. bson-direct appears here through DocumentCodecBridge, forced to
materialize the BsonDocument it normally skips, so it is a middling backend in this table; its
AST-free numbers are below.
The same entity, but all the way to real BSON wire bytes rather than stopping at a
BsonDocument. That is the trip the driver actually makes, and the only axis on which an AST-free
codec can be compared to an AST-based one at all.
| Backend | Throughput | Alloc |
|---|---|---|
| hand-written, straight to the wire | ~2.48M ops/s | 1824 B/op |
| WireCodec | ~2.08M ops/s | 1880 B/op |
| calypso | ~1.42M ops/s | 3211 B/op |
| medeia | ~1.23M ops/s | 4592 B/op |
| zio-bson | ~866k ops/s | 4785 B/op |
| mongo4cats-circe | ~717k ops/s | 6411 B/op |
| mongo4cats-zio-json | ~651k ops/s | 5528 B/op |
| Backend | Throughput | Alloc |
|---|---|---|
| WireCodec | ~1.60M ops/s | 1248 B/op |
| hand-written, straight to the wire | ~1.55M ops/s | 1088 B/op |
| calypso | ~1.03M ops/s | 3183 B/op |
| medeia | ~954k ops/s | 2861 B/op |
| zio-bson | ~864k ops/s | 2840 B/op |
| mongo4cats-zio-json | ~599k ops/s | 6313 B/op |
| mongo4cats-circe | ~502k ops/s | 8552 B/op |
WireCodec wins both directions against every AST backend — 1.47× the best of
them on encode (calypso) and 1.68× on decode (medeia) — and allocates up to
6.9× less than mongo4cats-circe on decode. The hand-written row is the ceiling derivation is
measured against: a codec written by hand straight into the BsonWriter encodes
1.19× faster than the derived one, and on decode the two are level within error. In a
typical CRUD workload the round trip dwarfs all of this; it matters on bulk paths.
A type bson-direct has no ScalarWireCodec for falls back to its
BsonEncoder/BsonDecoder, materializing one org.bson.BsonValue per
field — the very allocation the path exists to avoid.
| Direction | Native ScalarWireCodec | Through the BsonValue bridge |
|---|---|---|
| Encode | ~2.58M ops/s, 1768 B/op | ~2.24M ops/s, 2064 B/op |
| Decode | ~2.05M ops/s, 1248 B/op | ~1.60M ops/s, 1736 B/op |
Decode is where it shows: 1.28× the throughput and 28% less garbage, four fields out of seven. Both columns write byte-identical BSON — the bridge is not a different format, only a slower way to the same bytes. The allocation figures are the durable half: bytes per operation are a property of the code rather than of the machine, and three of these four came back byte-for-byte identical to a run on different hardware, while every throughput ratio moved.
Since 3.0.0 every Publisher the driver hands back is wrapped so failures arrive as
MongoError rather than as driver exceptions — one extra Subscriber per subscription
and one extra virtual call per element, on the hot path of every operation.
ErrorTranslationBenchmark drains a synchronous publisher with and without that wrapper, at 1, 100
and 10,000 elements: every difference is inside the error bars, and the wrapper's own
allocation is one object per subscription rather than per element (16 B on a 24 B baseline, and only
visible at 10,000). Typed errors are not paid for on the happy path.
The three sections that follow all cross a socket to MongoDB, and on this machine that socket is the whole story. Docker Desktop on macOS is a Linux VM, and a round trip through it costs milliseconds. The floor was measured from a shell inside the container, with no port mapping in the path:
empty loop 0.000 ms/op
ping 1.35–2.08 ms/op
findOne 1.45–1.97 ms/op
An empty loop is free, so that 2.4 ms is the server round trip itself. mongo4s doing real work from the host lands around 200 ops/s — roughly half the rate a bare shell gets for a ping in the same container. So read the tables below as comparisons, not as rates: which stack is faster than which is what they measure, while the absolute operations-per-second are a property of this environment. A latency-bound measurement also compresses every ratio, because the constant round trip is a larger share of each result — where a table shows a gap, the real gap is at least that wide. The codec sections above have no server in them and are unaffected.
| Documents returned | aggregateDirect | aggregate + bridge | aggregate + medeia |
|---|---|---|---|
| 10 | 262 ops/s | 268 ops/s | 258 ops/s |
| 10000 | 69 ops/s | 50 ops/s | 58 ops/s |
Both rows are the point. At ten documents the three sit inside each other's error bars, because the round trip
is everything and the codec is noise — the honest answer for most CRUD. At ten thousand the decode starts to
tell, and aggregateDirect is 1.37× the bridged path. Pick it for cursors that
return a lot; below that it is a wash and either call is fine.
| Operation | mongo4s+medeia | mongo4s+bson-direct | mongo4cats+circe | mongo4cats+zio-json |
|---|---|---|---|---|
| insertOne | 253 | 251 | 256 | 253 |
| insertMany (10 docs) | 224 | 217 | 219 | 221 |
| findOneById | 250 | 255 | 254 | 253 |
| findAll (~100 docs) | 242 | 240 | 229 | 231 |
| findStream (~100 docs) | 232 | 238 | 10 | 10 |
| updateOne | 262 | 264 | 253 | 255 |
| count | 245 | 238 | 245 | 242 |
Every column lands in the same band for every operation — the round trip dominates and none of the six
RsBridge/collection wrappers stands out. Run-to-run error is ±4–41%, far wider than the spread
between the columns, so the highlighted cells mark the measured extreme rather than a ranking. The one real
outlier is mongo4cats-cats's find(filter).stream, a 22× gap against its own
.all that no round trip explains: it bridges through a hand-rolled
cats.effect.std.Queue-backed Subscriber instead of
fs2.interop.reactivestreams, which mongo4s-cats uses. mongo4cats-zio does not share the
problem, so this is one bridge rather than the library.
| Operation | mongo4s+medeia | mongo4s+bson-direct | mongo4cats+circe | mongo4cats+zio-json |
|---|---|---|---|---|
| insertOne | 203 | 201 | 201 | 202 |
| insertMany (10 docs) | 192 | 191 | 192 | 186 |
| findOneById | 200 | 198 | 203 | 202 |
| findAll (~100 docs) | 199 | 197 | 208 | 195 |
| findStream (~100 docs) | 193 | 191 | 14 | 14 |
| updateOne | 202 | 199 | 200 | 200 |
| count | 199 | 196 | 202 | 200 |
On throughput the two mongo4s codecs are indistinguishable here, and so are they from mongo4cats on everything but the stream: at ~200 ops/s a codec has nowhere to show. The allocation table is the one to read.
| Operation | mongo4s+medeia | mongo4s+bson-direct | mongo4cats+circe | mongo4cats+zio-json |
|---|---|---|---|---|
| insertOne | 25.7 | 22.8 | 26.4 | 25.9 |
| insertMany (10 docs) | 67.1 | 39.3 | 79.4 | 72.2 |
| findOneById | 37.2 | 35.1 | 40.6 | 37.9 |
| findAll (~100 docs) | 326.4 | 157.5 | 928.4 | 643.9 |
| findStream (~100 docs) | 379.7 | 211.2 | 2850.6 | 2573.1 |
| updateOne | 23.3 | 23.2 | 23.6 | 23.5 |
| count | 32.8 | 32.9 | 32.5 | 32.5 |
bson-direct allocates less than bson-medeia wherever a document actually passes through the
codec — the AST-free advantage measured in isolation above survives end-to-end through a real
driver
round trip. The gap tracks how many documents a call encodes or decodes: 11% for one,
41% on insertMany, 44–52% on
findAll/findStream. On updateOne and count the two are
level to three digits, the expected result rather than a surprise, since neither encodes or decodes an entity.
On the bulk reads mongo4cats allocates up to 13× more than either mongo4s config — its own
mongo4cats.bson.BsonValue wrapper adds a full extra tree per document on top of org.bson's, and
that cost multiplies with document count.
Why mongo4s is shaped the way it is — what was chosen, what was rejected, and what it cost. Several of these decisions look arbitrary until you know the failure they were a response to, so the failures are here too.
given-based throughout, derivation via Mirror and
inline, no runtime reflection. Not cross-building to 2.13 is what buys macro field selectors,
derives, opaque types and typed contextual parameters — the whole reason the API can be
type-safe where a cross-built one cannot.
underlying is public on client, database and
collection. Anything mongo4s does not model is still reachable, and reaching for it is not a defeat.
MongoClient[F, S] → MongoDatabase[F, S] → MongoCollection[F, S, A]
follow the official driver's hierarchy one-to-one. Someone who knows
mongodb-driver-reactivestreams can guess where things are; someone who does not can read the
driver's documentation and have it apply. Two type parameters, not one, because no effect system bundles
effect and stream — cats-effect pairs with fs2, ZIO ships its own ZStream, and the pairing is a
user decision.
Every method that reaches the server takes (using session: Option[ClientSession] = None).
Sessions are invisible until you opt in, and joining a transaction requires no change at the call site — the
session is given implicitly by withTransaction and picked up by every call inside the block,
including repository calls that know nothing about transactions. The alternative — a second overload of every
method taking an explicit session — doubles the API surface and makes a transactional call textually
different from a non-transactional one.
mongo4s-core depends on the Mongo driver and nothing else. A runtime module is two
given instances.
trait Effect[F[*]]:
def pure[A](a: A): F[A]
def delay[A](a: => A): F[A]
def map[A, B](fa: F[A])(f: A => B): F[B]
def flatMap[A, B](fa: F[A])(f: A => F[B]): F[B]
def raiseError[A](error: Throwable): F[A]
def handleErrorWith[A](fa: F[A])(f: Throwable => F[A]): F[A]
def guaranteeCase[A](fa: F[A])(finalizer: ExitCase => F[Unit]): F[A]
// suspend, unit, void, attempt, guarantee, onError, bracket, bracketCase — derived, overridable
Seven abstract members, and that number is a deliberate ceiling: everything derivable is derived with a default, so supporting a new effect system is an afternoon, and adding capability later does not break implementors.
guaranteeCase is the load-bearing one. It sees
ExitCase.Succeeded | Errored | Canceled, and the Canceled case is why an interrupted
transaction rolls back immediately instead of lingering on the server until it is reaped, and why a canceled
read releases its cursor. An abstraction built on Monad + MonadError alone cannot
express that — which is the reason Effect exists rather than the library taking a
cats.effect.Async constraint. A finalizer that itself fails does not replace the error that
caused it: it is attached with addSuppressed, because the rollback error is what you want to see
but the error that triggered the rollback is what you need to debug.
trait RsBridge[F[*], S[*]]:
def one[A](publisher: => Publisher[A]): F[A]
def option[A](publisher: => Publisher[A]): F[Option[A]]
def list[A](publisher: => Publisher[A]): F[List[A]]
def unit[A](publisher: => Publisher[A]): F[Unit]
def stream[A](publisher: => Publisher[A])(using Streamable[S, A]): S[A]
def liveStream[A](publisher: => Publisher[A])(using Streamable[S, A]): S[A] = stream(publisher)
Every Publisher is taken by name, and that is not a style choice: the driver
opens a cursor when a publisher is subscribed to, and building an F[A] value is not running it. A
by-value parameter would open a cursor for a computation that may never run, or run twice for one that is
retried. Streamable[S, A] is a marker typeclass with no methods; it exists so every streaming
call names its element type at the call site, which is what lets the kyo backend supply the per-element
Tag its stream type requires.
The driver's model is a process-wide CodecRegistry. mongo4s does not use it for entity codecs at
all: a codec is resolved as a given at the call site that creates the collection, and attached to
that collection alone. The cost is one more implicit parameter. What it buys is that medeia, zio-bson,
calypso and bson-direct can all be in use in the same process on different collections — and that a missing
codec is a compile error naming the type, rather than a CodecConfigurationException under load.
The driver's registry still matters twice, and both are load-bearing:
BsonDocument codec. A registry that
replaces the defaults instead of extending them fails every operation, because mongo4s reads and
writes BsonDocument underneath.
underlying. The escape hatch is the user's, and their registry is
what applies there.
On the direct path the derived WireCodec[A] is registered ahead of the client's
registry. The other order was tried and was wrong: a Codec[A] in the user's registry silently
shadowed the derived codec, so the collection wrote a shape nobody had asked for and no error was raised.
Both behaviours are pinned by integration specs.
On a direct collection the codec also decides how queries spell fields, which only holds if every codec tells
the truth about its naming. A derived sum type used to report identity whatever it wrote, and
imap/iemap forgot the naming of the codec underneath — both wrote
first_name and let queries say firstName. Both carry it now, and a law spec checks that
the naming a codec reports is the spelling it writes.
Every other Scala Mongo library encodes through an intermediate tree: case class → some JSON or BSON AST →
BsonDocument → wire bytes. bson-direct derives a WireCodec[A] that
writes into the driver's BsonWriter directly and reads from its BsonReader, with no
intermediate representation at all. That is where the performance comes from — but the reason it is the
default codec is dependency footprint: it lives in mongo4s itself, so the default path pulls in no
third-party codec library.
It is strict by construction, and that is a trade, not an oversight. Derivation requires
every modelled field to be present unless its decoder supplies a default (Option does), so a
projection that drops a modelled field cannot be read back through a direct collection — use
getCollection with a BsonDocumentCodec, or model the projected shape as its own
type. Strictness catches the far more common bug, which is a document that quietly lost a field.
aggregate and distinct still go through BsonDocumentCodec, because
their output shape is not A and they are not the hot path.
Filter and Update are real enum ADTs that mongo4s interprets itself,
not thin wrappers over the driver's opaque Bson builders. That is what makes a second
interpreter possible at all, and what lets the AST be normalized rather than passed through.
Field paths track derived-vs-stored per segment, not per path. Selector-derived names are
spelled through the collection's FieldNaming; names that are already what the document stores — a
map key, an array index, _id, a PrimaryKey's field names, a $lookup's
foreign field — are used verbatim. Per-segment is the only granularity that works:
totals.at("EUR") has to rename totals under snake_case while leaving the map key
alone.
Filter.and/or fold all and none away
instead of emitting a one-element $and. Not cosmetic: it is what makes
field.in(Nil) produce Filter.none rather than an empty predicate, so
deleteMany(Nil) is a safe no-op instead of deleting the collection. An empty list is the most
dangerous input a query builder takes, and it gets a defined answer here rather than an emergent one.
Update merges operators of the same name into one sub-document —
{$set: {a: 1}} combined with {$set: {b: 2}} is {$set: {a: 1, b: 2}},
not two $set keys the server would reject. Update.Raw merges the same way, but must
clone the caller's document first: not cloning meant a rendered update wrote into the
caller's own BsonDocument, so a shared Raw came back permanently carrying fields it
never declared, and rendering it twice under different namings emitted both spellings. An update that would
render empty throws rather than sending {}.
Merging cannot keep two changes to one path, so it refuses to try. A sub-document holds one
value per key, and the second $inc on a path used to replace the first silently — while the
fake, applying both, answered differently. An update touching a path twice, or overlapping paths across
operators, is refused at render, and the fake renders before it applies. A pipeline update is a case of
Update rather than a separate type, so everything that takes an update takes one; it renders as
stages and is refused where the two forms would mix.
Numeric operators go through NumericOf[C, A], so an Option[Long] field takes a plain
Long. There is deliberately no NumericValue[Option[A]]: None has no
numeric encoding, and the obvious stand-in — $inc by zero — is a write that silently does
nothing. The same capability pattern now bounds regex (TextOf), mod and
hasSize/$pop, which used to render filters the server accepts and that match nothing.
all fails the whole query on the first undecodable document, because a decode failure usually
means the model and the collection disagree, and finding out immediately beats processing 99 documents and
silently dropping one. But a collection written by more than one version of more than one service will
contain documents this version cannot read, and failing a whole page over one of them is useless — so
attempting reports each document as Either[BsonError, A]. Transport errors still fail
the effect; only decoding is made recoverable, because only decoding is a per-document property.
watchAttempting puts the result on the event's documents rather than around the event, because
wrapping the event lost its resume token with it — and a durable consumer that cannot checkpoint past a bad
document reads it again after every restart.
Operations expecting at most one document read two, not the whole cursor — exactly enough to
notice a second result under strictSingleResult, and no more.
AggregateQuery.first pushes a $limit into the pipeline for the same reason:
bounding work at the server, not in the client.
liveStream exists because fs2 chunks.
fromPublisher(p, bufferSize) will not emit an element until the chunk is full or the publisher
finishes — its own scaladoc says so. With the default bufferSize = 256, a change stream
delivered nothing until 256 events had piled up, and a change stream never finishes.
find(...).stream only ever worked by accident: its cursor completes, which flushes the partial
chunk. So the bridge has two stream methods — finite reads keep stream and its buffering, which
is what makes them fast, while every watch goes through liveStream, which the cats
backend overrides to fromPublisher(p, 1). The default implementation delegates to
stream, because zio (queue-based), rapid and kyo have no chunk-fill behaviour to work around.
This was found only because change streams had never actually been exercised — the integration spec was
calling cancel on itself and reporting green. Two lessons kept: never let an integration spec
cancel itself into looking green, and never time a change-stream test with sleep — take the
server's operationTime from hello and pass startingAt(now), which is
race-free no matter when the stream actually subscribes. RsBridgeConfig's per-operation timeout
deliberately does not apply to streams: a change stream sitting idle is working, not stuck.
withTransaction commits on success and rolls back on failure and on cancellation
— the cancellation half is the whole reason guaranteeCase has the shape it does. A rollback that
itself fails is attached as a suppressed exception rather than replacing the original error. The manual path
is still there, and nothing is automatic on it, including the rollback: the safe version should be the easy
one, not the only one.
The retry rules are the specification's. A commit that ran out of maxCommitTime
carries the "unknown commit result" label — the driver adds it — so it used to be asked again until
the retry window closed, which turned the limit into a delay; it is reported instead. A whole-transaction retry
pauses first, for a jittered interval growing from 5 ms to a 500 ms ceiling, which needed a sleep
Effect did not have. And the body is built inside the protection, not before it: one that threw
before returning an F used to escape the finalizer and leave the transaction open.
Every event is a ChangeEvent[A] — a real case class, not a raw BsonDocument to dig
through — and watch exists at all three driver scopes. fullDocument defaults to
UPDATE_LOOKUP, not the server's default: MongoDB fills the document in only for
inserts and replaces, which leaves the most common question ("what does this document look like now?")
unanswered on exactly the events you are usually watching for. resumeAfter and
startAfter clear each other, because the server rejects a stream carrying both.
One known footgun, accepted deliberately: WatchOptions[E].pipeline is typed
Seq[Stage[E]], but a change stream pipeline matches against the change event envelope, not the
document — so a Field.of path renders "age" where the event needs
"fullDocument.age", and Stage.raw is the right tool. The alternative was a second,
parallel Stage type for events; that was judged worse than one type with a documented sharp
edge.
BaseMongoRepository is open, not generated — extra domain queries are declared
directly against collection/Filter/Field. A generated repository would
have to either predict every query or make you drop out of the abstraction to write one. Batching is
empty-list-safe throughout, which follows from Filter's folding.
BulkWriteResult.upsertedIds is keyed by command position, so
combine deliberately does not rebase indices; a caller splitting one logical write into batches
calls shiftUpsertedIds(offset) per batch first, since only the caller knows how many commands
preceded it. Getting this wrong collapsed every batch onto index 0 — a bug that only appears once a write
exceeds one batch, which is to say in production and not in tests. The same rebasing was later found missing
for failures: a batched failure now describes the whole write, as the driver's own batching does — global
indices, and a result counting everything that landed.
Because Filter/Update are an AST rather than driver builders, they can be
interpreted twice: into real Bson for the server, and against an in-memory buffer.
FakeMongoCollection is that second interpreter — filters, every update operator, sorting, paging,
projections, distinct and a subset of aggregate are simulated, each checked against
a real server by a parity spec, while watch, explain, $text,
$expr, Filter.Raw, a pipeline update and an update carrying arrayFilters
throw UnsupportedOperationException naming what was asked for rather than quietly answering
wrong. A fake that lies is worse than no fake.
The fake streams through the RsBridge rather than an emit function the caller
supplies. emit took the results as a list, computed when the stream was built, so a stream set up
before an insert missed it — the one thing a cursor never does — and nothing but find
could stream, because nobody supplied an emitter for the other types.
It ships as its own published module, mongo4s-testkit, so it is usable from a consumer's own
tests rather than only inside this build. Alongside it, FakeRepository is a
BaseMongoRepository over a FakeMongoCollection — the real repository logic, only the
collection faked, so the fake cannot drift from what production does.
Three commitments, and the mechanics that make each keepable. New Effect/
RsBridge members carry default implementations, so implementing either typeclass
yourself keeps compiling across minor releases — liveStream was added that way.
Binary compatibility is checked, not asserted: MiMa runs in CI against the previous release,
and every deliberate break is a visible filter plus an entry in the compatibility notes.
Configuration types are builders, not case classes — learned the expensive way. Adding a
field with a default to a case class is source-compatible but never binary-compatible: default
arguments are resolved at the call site, so the JVM only ever sees the old arity, and apply,
copy and the constructor all change. Since derivation options keep arriving,
WireCodecConfig became a final class with a private constructor and
withX methods — one break taken once, every option after it additive.
Either/Ior codec. Dropping
_type and trying branch A then branch B via getMark()/reset() is
technically feasible, and was rejected because it is silently unsafe whenever the two branches are
structurally similar enough that both decode "successfully" — the wrong branch wins with no error.
CodecProvider convenience. It would remove one implicit parameter
and reintroduce process-wide coupling between unrelated collections.
Stream.take(n). Not a mongo4s decision,
but worth knowing: rapid decides Step.Stop inside transform, so it needs element
n+1 before it stops. On an infinite change stream carrying exactly n events it blocks forever.
mongo4s-bson-circe module, and there is not going
to be one. circe is a JSON codec: routing an entity through it is the same trip the benchmarks measure at
2.5–5× slower on decode, the cost mongo4s exists to remove. A model already on circe should get a
BsonDocumentCodec[A] written against BSON directly.
insertMany and bulkWrite take a
Seq. Taking an S would need RsBridge to fold a stream into
F, which the four runtimes share no shape for, and the loop it saves is three lines in each.
$bucketAuto, $geoNear and renameCollection in the
fake. The first balances its buckets by the server's own rules, the second needs spherical
distances the fake does not compute, and the third a second namespace a single-collection fake does not
have. Each is refused by name.
BsonDocument the API accepts. Raw filters, stages, hints and
resume tokens are taken as given: mongo4s never writes into them. Update.Raw is the exception,
cloned because rendering used to write into it.