diff --git a/dsl-template/src/commonMain/kotlin/command/BulkWrite.kt b/dsl-template/src/commonMain/kotlin/command/BulkWrite.kt index c33480e1..c5449e80 100644 --- a/dsl-template/src/commonMain/kotlin/command/BulkWrite.kt +++ b/dsl-template/src/commonMain/kotlin/command/BulkWrite.kt @@ -484,7 +484,7 @@ class BulkWrite private constructor( when (operation) { is InsertOne<*> -> { writeInt32("insert", 0) - writeSafe("document", operation.document) + writeSafe("document", operation.document, operation.documentType) operation.options.writeTo(this) } @@ -505,7 +505,7 @@ class BulkWrite private constructor( writeDocument("filter") { operation.filter.writeTo(this) } - writeSafe("updateMods", operation.document) + writeSafe("updateMods", operation.document, operation.documentType) writeBoolean("multi", false) operation.options.writeTo(this) } @@ -515,7 +515,7 @@ class BulkWrite private constructor( writeDocument("filter") { operation.filter.writeTo(this) } - writeSafe("updateMods", operation.document) + writeSafe("updateMods", operation.document, operation.documentType) writeBoolean("multi", false) writeBoolean("upsert", true) operation.options.writeTo(this) diff --git a/dsl-template/src/commonMain/kotlin/command/CreateCollection.kt b/dsl-template/src/commonMain/kotlin/command/CreateCollection.kt index eccd20b9..3dbe7855 100644 --- a/dsl-template/src/commonMain/kotlin/command/CreateCollection.kt +++ b/dsl-template/src/commonMain/kotlin/command/CreateCollection.kt @@ -59,6 +59,7 @@ class CreateCollection private constructor( class CreateCollectionOptions(context: BsonContext) : Options by OptionsHolder(context), HasCapped, + HasTimeSeries, HasValidation, HasWriteConcern, HasComment diff --git a/dsl-template/src/commonMain/kotlin/options/ExpiresAfterOption.kt b/dsl-template/src/commonMain/kotlin/options/ExpiresAfterOption.kt new file mode 100644 index 00000000..c799e9ab --- /dev/null +++ b/dsl-template/src/commonMain/kotlin/options/ExpiresAfterOption.kt @@ -0,0 +1,61 @@ +/* + * Copyright (c) 2026, OpenSavvy and contributors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package opensavvy.ktmongo.dsl.options + +import opensavvy.ktmongo.bson.BsonValueWriter +import opensavvy.ktmongo.dsl.BsonContext +import opensavvy.ktmongo.dsl.LowLevelApi +import kotlin.time.Duration +import kotlin.time.Duration.Companion.seconds + +/** + * The duration after which the document expires. + * + * See [TimeSeriesDsl.expiresAfter]. + */ +class ExpiresAfterOption private constructor( + val duration: Duration, + context: BsonContext, + @Suppress("unused") overloadMarker: Unit, +) : AbstractOption("expireAfterSeconds", context) { + + constructor( + duration: Duration, + context: BsonContext, + ) : this( + if (duration == Duration.ZERO) duration else duration.coerceAtLeast(1.seconds), + context, + Unit, + ) + + @LowLevelApi + override fun write(writer: BsonValueWriter) = with(writer) { + writeInt64(duration.inWholeSeconds) + } + + @LowLevelApi + override fun merge(other: Option): Option { + require(other is ExpiresAfterOption) { "Cannot merge sort options of different types: ${this::class} and ${other::class}" } + + // If this option is specified for both time-series and clustered collections, + // the shortest specified time should apply + return ExpiresAfterOption( + duration = minOf(this.duration, other.duration), + context, + ) + } +} diff --git a/dsl-template/src/commonMain/kotlin/options/TimeSeries.kt b/dsl-template/src/commonMain/kotlin/options/TimeSeries.kt new file mode 100644 index 00000000..e50c8e8e --- /dev/null +++ b/dsl-template/src/commonMain/kotlin/options/TimeSeries.kt @@ -0,0 +1,512 @@ +/* + * Copyright (c) 2026, OpenSavvy and contributors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package opensavvy.ktmongo.dsl.options + +import opensavvy.ktmongo.bson.BsonFieldWriter +import opensavvy.ktmongo.bson.BsonValueWriter +import opensavvy.ktmongo.dsl.BsonContext +import opensavvy.ktmongo.dsl.DangerousMongoApi +import opensavvy.ktmongo.dsl.KtMongoDsl +import opensavvy.ktmongo.dsl.LowLevelApi +import opensavvy.ktmongo.dsl.path.Field +import opensavvy.ktmongo.dsl.path.FieldDsl +import opensavvy.ktmongo.dsl.path.Path +import opensavvy.ktmongo.dsl.tree.AbstractBsonNode +import opensavvy.ktmongo.dsl.tree.AbstractCompoundBsonNode +import opensavvy.ktmongo.dsl.tree.BsonNode +import opensavvy.ktmongo.dsl.tree.CompoundBsonNode +import kotlin.reflect.KProperty1 +import kotlin.time.Duration +import kotlin.time.Duration.Companion.seconds +import kotlin.time.Instant + +/** + * The options to define a time-series collection. + * + * For more information, see [HasTimeSeries.timeSeries]. + */ +class TimeSeriesOption( + val node: BsonNode, + context: BsonContext, +) : AbstractOption("timeseries", context) { + + init { + @OptIn(LowLevelApi::class) + node.freeze() + } + + @LowLevelApi + override fun write(writer: BsonValueWriter) = with(writer) { + writeDocument { + node.writeTo(this) + } + } +} + +/** + * Declares time-series collections. + * + * See [timeSeries]. + */ +@KtMongoDsl +interface HasTimeSeries : Options { + + /** + * Time-series are special collections optimized for cases where data is created over time. + * + * Examples: + * - A temperature sensor that makes readings every minute. + * - An audit log where a new document is inserted each time a user makes an action. + * - A log-in log where a new document is inserted on every log-in attempt. + * + * Instead of storing documents independently like in normal collections, documents are grouped together + * by their source (see [TimeSeriesDsl.metaField]) in buckets of some time period (see [TimeSeriesDsl.granularity] + * and [TimeSeriesDsl.bucketMaxSpan]), which allows documents from the source and a similar time to be queried together + * much more efficiently. Additionally, MongoDB can use [compression](https://www.mongodb.com/docs/manual/core/timeseries/timeseries-compression/) + * to greatly reduce the disk space used by such a collection. + * + * Time-series documents do not require an `_id` field. + * + * ### Example + * + * ```kotlin + * class Notification( + * val userId: ObjectId, + * val createdAt: Instant, + * val text: String, + * ) + * + * notifications.create { + * timeSeries { + * timeField(Notification::createdAt) + * metaField(Notification::userId) + * granularity(Minutes) + * expiresAfter(7.days) + * } + * } + * ``` + * + * ### External resources + * + * - [Official documentation](https://www.mongodb.com/docs/manual/core/timeseries-collections/) + * + * @param block The specification of this time-series collection. + * - [TimeSeriesDsl.timeField] The field which contains the timestamp used to organize this time series. **Required**. + * - [TimeSeriesDsl.metaField] A field used to group documents together. + * - [TimeSeriesDsl.expiresAfter] Documents will be automatically removed from the collection after this duration. + * - [TimeSeriesDsl.granularity] How often documents are created in a given bucket. + * - [TimeSeriesDsl.bucketMaxSpan] A new bucket will be created after this duration. + */ + @OptIn(DangerousMongoApi::class, LowLevelApi::class) + fun timeSeries( + block: TimeSeriesDsl.() -> Unit, + ) { + val options = TimeSeriesDslImpl(this, context).apply(block) + accept(TimeSeriesOption(options, context)) + } + + private class TimeSeriesDslImpl( + private val optionsBlock: Options, + context: BsonContext, + ) : AbstractCompoundBsonNode(context), TimeSeriesDsl { + + @OptIn(DangerousMongoApi::class, LowLevelApi::class) + override fun timeField(field: Field) { + accept(TimeField(field.path, context)) + } + + @OptIn(DangerousMongoApi::class, LowLevelApi::class) + override fun expiresAfter(duration: Duration) { + optionsBlock.accept(ExpiresAfterOption(duration, context)) + } + + @OptIn(DangerousMongoApi::class, LowLevelApi::class) + override fun metaField(field: Field) { + accept(MetaField(field.path, context)) + } + + @OptIn(DangerousMongoApi::class, LowLevelApi::class) + override fun granularity(granularity: TimeSeriesDsl.Granularity) { + accept(GranularityField(granularity, context)) + } + + @OptIn(DangerousMongoApi::class, LowLevelApi::class) + override fun bucketMaxSpan(maxSpan: Duration) { + val span = + if (maxSpan == Duration.ZERO) maxSpan + else maxSpan.coerceAtLeast(1.seconds) + + accept(BucketMaxSpan(span, context)) + accept(BucketRounding(span, context)) + } + } + + @LowLevelApi + private class TimeField( + private val path: Path, + context: BsonContext, + ) : AbstractBsonNode(context) { + + @OptIn(LowLevelApi::class) + override fun write(writer: BsonFieldWriter) = with(writer) { + writeString("timeField", path.toString()) + } + } + + @LowLevelApi + private class MetaField( + private val path: Path, + context: BsonContext, + ) : AbstractBsonNode(context) { + + @OptIn(LowLevelApi::class) + override fun write(writer: BsonFieldWriter) = with(writer) { + writeString("metaField", path.toString()) + } + } + + @LowLevelApi + private class GranularityField( + private val granularity: TimeSeriesDsl.Granularity, + context: BsonContext, + ) : AbstractBsonNode(context) { + + @OptIn(LowLevelApi::class) + override fun write(writer: BsonFieldWriter) = with(writer) { + writeString("granularity", granularity.bsonName) + } + } + + @LowLevelApi + private class BucketMaxSpan( + private val duration: Duration, + context: BsonContext, + ) : AbstractBsonNode(context) { + + @OptIn(LowLevelApi::class) + override fun write(writer: BsonFieldWriter) = with(writer) { + writeInt64("bucketMaxSpanSeconds", duration.inWholeSeconds) + } + } + + @LowLevelApi + private class BucketRounding( + private val duration: Duration, + context: BsonContext, + ) : AbstractBsonNode(context) { + + @OptIn(LowLevelApi::class) + override fun write(writer: BsonFieldWriter) = with(writer) { + writeInt64("bucketRoundingSeconds", duration.inWholeSeconds) + } + } +} + +/** + * Options for [HasTimeSeries.timeSeries]. + */ +@KtMongoDsl +interface TimeSeriesDsl : CompoundBsonNode, FieldDsl { + + /** + * Specifies the field which contains the date used to organize this time-series collection. + * + * ### Example + * + * ```kotlin + * class Notification( + * val userId: ObjectId, + * val createdAt: Instant, + * val text: String, + * ) + * + * notifications.create { + * timeSeries { + * timeField(Notification::createdAt) + * metaField(Notification::userId) + * granularity(Minutes) + * expiresAfter(7.days) + * } + * } + * ``` + * + * ### External resources + * + * - [Official documentation](https://www.mongodb.com/docs/manual/core/timeseries/timeseries-bucketing/) + */ + fun timeField(field: Field) + + /** + * Specifies the field which contains the date used to organize this time-series collection. + * + * ### Example + * + * ```kotlin + * class Notification( + * val userId: ObjectId, + * val createdAt: Instant, + * val text: String, + * ) + * + * notifications.create { + * timeSeries { + * timeField(Notification::createdAt) + * metaField(Notification::userId) + * granularity(Minutes) + * expiresAfter(7.days) + * } + * } + * ``` + * + * ### External resources + * + * - [Official documentation](https://www.mongodb.com/docs/manual/core/timeseries/timeseries-bucketing/) + */ + fun timeField(field: KProperty1) { + timeField(field.field) + } + + /** + * Specifies how long documents are retained in the collection. + * + * Once this time has passed after the value of the [timeField], MongoDB automatically removes the document. + * + * The minimum resolution is one second, so values less than one second are rounded up to one second. + * + * **Note.** This method sets the [ExpiresAfterOption], which does not appear in the `toString()` representation + * of this DSL, because MongoDB places this option at the root of the `create {}` command instead of within + * the `timeseries` block. + * + * ### Example + * + * ```kotlin + * class Notification( + * val userId: ObjectId, + * val createdAt: Instant, + * val text: String, + * ) + * + * notifications.create { + * timeSeries { + * timeField(Notification::createdAt) + * metaField(Notification::userId) + * granularity(Minutes) + * expiresAfter(7.days) + * } + * } + * ``` + * + * ### External resources + * + * - [Official documentation](https://www.mongodb.com/docs/manual/core/timeseries/timeseries-automatic-removal) + */ + fun expiresAfter(duration: Duration) + + /** + * Controls how MongoDB buckets documents. + * + * To improve performance, MongoDB stores time-series documents in buckets, with documents closely related to each other + * in the same bucket for improved cache performance. + * + * Specifying a meta-field is optional but recommended when multiple documents that originated at around the same time + * are likely to be queried together. + * + * If two documents have the same meta-field value and a similar time-field value, MongoDB stores them together. + * + * A good meta-field is a field that appears in almost all queries and values that are commonly used, for + * example, the origin of the event or its destination. + * + * ### Example + * + * ```kotlin + * class Notification( + * val userId: ObjectId, + * val createdAt: Instant, + * val text: String, + * ) + * + * notifications.create { + * timeSeries { + * timeField(Notification::createdAt) + * metaField(Notification::userId) + * granularity(Minutes) + * } + * } + * ``` + * + * In this example, we store notifications destined for some users. + * We bucket the data by the user each notification is destined for, to ensure we can query all notifications for a given user quickly. + * + * ### External resources + * + * - [How bucketing works](https://www.mongodb.com/docs/manual/core/timeseries/timeseries-bucketing/#how-bucketing-works) + * - [Best practices](https://www.mongodb.com/docs/manual/core/timeseries/timeseries-considerations/) + */ + fun metaField(field: Field) + + /** + * Controls how MongoDB buckets documents. + * + * To improve performance, MongoDB stores time-series documents in buckets, with documents closely related to each other + * in the same bucket for improved cache performance. + * + * Specifying a meta-field is optional but recommended when multiple documents that originated at around the same time + * are likely to be queried together. + * + * If two documents have the same meta-field value and a similar time-field value, MongoDB stores them together. + * + * A good meta-field is a field that appears in almost all queries and values that are commonly used, for + * example, the origin of the event or its destination. + * + * ### Example + * + * ```kotlin + * class Notification( + * val userId: ObjectId, + * val createdAt: Instant, + * val text: String, + * ) + * + * notifications.create { + * timeSeries { + * timeField(Notification::createdAt) + * metaField(Notification::userId) + * granularity(Minutes) + * } + * } + * ``` + * + * In this example, we store notifications destined for some users. + * We bucket the data by the user each notification is destined for, to ensure we can query all notifications for a given user quickly. + * + * ### External resources + * + * - [How bucketing works](https://www.mongodb.com/docs/manual/core/timeseries/timeseries-bucketing/#how-bucketing-works) + * - [Best practices](https://www.mongodb.com/docs/manual/core/timeseries/timeseries-considerations/) + */ + fun metaField(field: KProperty1) { + metaField(field.field) + } + + /** + * Specifies the maximum size of a bucket. + * + * To learn more about buckets, see [metaField]. + * + * This option is not compatible with [bucketMaxSpan]. + * + * This option communicates to MongoDB how many documents will be created in each bucket. + * For example, if [metaField] is a company ID, and a new document is created by each company every few minutes, choose [Granularity.Minutes]. + * + * If this option is not specified, [Granularity.Seconds] is used. + * + * ### Example + * + * ```kotlin + * class Notification( + * val userId: ObjectId, + * val createdAt: Instant, + * val text: String, + * ) + * + * notifications.create { + * timeSeries { + * timeField(Notification::createdAt) + * metaField(Notification::userId) + * granularity(Minutes) + * } + * } + * ``` + * + * ### External resources + * + * - [Official documentation](https://www.mongodb.com/docs/manual/core/timeseries/timeseries-considerations/#granularity) + */ + fun granularity(granularity: Granularity) + + /** + * Specifies the maximum time between two timestamps in the same bucket. + * + * To learn more about buckets, see [metaField]. + * + * This option is not compatible with [granularity]. + * + * This option communicates to MongoDB how long to keep a bucket open. + * Once this time has passed since the creation of the bucket, MongoDB will close the bucket and create a new one. + * + * The minimum resolution is one second, so values less than one second are rounded up to one second. + * + * This function sets both the `bucketMaxSpanSeconds` and `bucketRoundingSeconds` options. + * + * ### Example + * + * ```kotlin + * class Notification( + * val userId: ObjectId, + * val createdAt: Instant, + * val text: String, + * ) + * + * notifications.create { + * timeSeries { + * timeField(Notification::createdAt) + * metaField(Notification::userId) + * bucketMaxSpan(12.hours) + * } + * } + * ``` + * + * ### External resources + * + * - [Official documentation](https://www.mongodb.com/docs/manual/core/timeseries/timeseries-granularity/#using-custom-bucketing-parameters) + */ + fun bucketMaxSpan( + maxSpan: Duration, + ) + + /** + * The granularity of the time-series buckets. + * + * For more information, see [metaField]. + * + * ### External resources + * + * - [Official documentation](https://www.mongodb.com/docs/manual/core/timeseries/timeseries-considerations/#granularity) + */ + enum class Granularity(val bsonName: String) { + /** + * Select this granularity if a new document is created every few seconds for a given [metaField] value. + * + * A bucket will store up to an hour of data. + */ + Seconds("seconds"), + + /** + * Select this granularity if a new document is created every few minutes for a given [metaField] value. + * + * A bucket will store up to 24 hours of data. + */ + Minutes("minutes"), + + /** + * Select this granularity if a new document is created every few hours for a given [metaField] value. + * + * A bucket will store up to a month of data. + */ + Hours("hours"), + } +} diff --git a/dsl/src/commonMain/kotlin/command/BulkWrite.kt b/dsl/src/commonMain/kotlin/command/BulkWrite.kt index bb076a1d..abe76793 100644 --- a/dsl/src/commonMain/kotlin/command/BulkWrite.kt +++ b/dsl/src/commonMain/kotlin/command/BulkWrite.kt @@ -487,7 +487,7 @@ class BulkWrite private constructor( when (operation) { is InsertOne<*> -> { writeInt32("insert", 0) - writeSafe("document", operation.document) + writeSafe("document", operation.document, operation.documentType) operation.options.writeTo(this) } @@ -508,7 +508,7 @@ class BulkWrite private constructor( writeDocument("filter") { operation.filter.writeTo(this) } - writeSafe("updateMods", operation.document) + writeSafe("updateMods", operation.document, operation.documentType) writeBoolean("multi", false) operation.options.writeTo(this) } @@ -518,7 +518,7 @@ class BulkWrite private constructor( writeDocument("filter") { operation.filter.writeTo(this) } - writeSafe("updateMods", operation.document) + writeSafe("updateMods", operation.document, operation.documentType) writeBoolean("multi", false) writeBoolean("upsert", true) operation.options.writeTo(this) diff --git a/dsl/src/commonMain/kotlin/command/CreateCollection.kt b/dsl/src/commonMain/kotlin/command/CreateCollection.kt index ce2fba2f..44b16730 100644 --- a/dsl/src/commonMain/kotlin/command/CreateCollection.kt +++ b/dsl/src/commonMain/kotlin/command/CreateCollection.kt @@ -62,6 +62,7 @@ class CreateCollection private constructor( class CreateCollectionOptions(context: BsonContext) : Options by OptionsHolder(context), HasCapped, + HasTimeSeries, HasValidation, HasWriteConcern, HasComment diff --git a/dsl/src/commonMain/kotlin/options/ExpiresAfterOption.kt b/dsl/src/commonMain/kotlin/options/ExpiresAfterOption.kt new file mode 100644 index 00000000..4667b133 --- /dev/null +++ b/dsl/src/commonMain/kotlin/options/ExpiresAfterOption.kt @@ -0,0 +1,64 @@ +/* + * Copyright (c) 2026, OpenSavvy and contributors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +// This file is generated from dsl-template/src/commonMain/kotlin/options/ExpiresAfterOption.kt +// DO NOT EDIT THIS FILE DIRECTLY. To learn more, read dsl-template/README.md. + +package opensavvy.ktmongo.dsl.options + +import opensavvy.ktmongo.bson.BsonValueWriter +import opensavvy.ktmongo.dsl.BsonContext +import opensavvy.ktmongo.dsl.LowLevelApi +import kotlin.time.Duration +import kotlin.time.Duration.Companion.seconds + +/** + * The duration after which the document expires. + * + * See [TimeSeriesDsl.expiresAfter]. + */ +class ExpiresAfterOption private constructor( + val duration: Duration, + context: BsonContext, + @Suppress("unused") overloadMarker: Unit, +) : AbstractOption("expireAfterSeconds", context) { + + constructor( + duration: Duration, + context: BsonContext, + ) : this( + if (duration == Duration.ZERO) duration else duration.coerceAtLeast(1.seconds), + context, + Unit, + ) + + @LowLevelApi + override fun write(writer: BsonValueWriter) = with(writer) { + writeInt64(duration.inWholeSeconds) + } + + @LowLevelApi + override fun merge(other: Option): Option { + require(other is ExpiresAfterOption) { "Cannot merge sort options of different types: ${this::class} and ${other::class}" } + + // If this option is specified for both time-series and clustered collections, + // the shortest specified time should apply + return ExpiresAfterOption( + duration = minOf(this.duration, other.duration), + context, + ) + } +} diff --git a/dsl/src/commonMain/kotlin/options/TimeSeries.kt b/dsl/src/commonMain/kotlin/options/TimeSeries.kt new file mode 100644 index 00000000..baea1ac4 --- /dev/null +++ b/dsl/src/commonMain/kotlin/options/TimeSeries.kt @@ -0,0 +1,515 @@ +/* + * Copyright (c) 2026, OpenSavvy and contributors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +// This file is generated from dsl-template/src/commonMain/kotlin/options/TimeSeries.kt +// DO NOT EDIT THIS FILE DIRECTLY. To learn more, read dsl-template/README.md. + +package opensavvy.ktmongo.dsl.options + +import opensavvy.ktmongo.bson.BsonFieldWriter +import opensavvy.ktmongo.bson.BsonValueWriter +import opensavvy.ktmongo.dsl.BsonContext +import opensavvy.ktmongo.dsl.DangerousMongoApi +import opensavvy.ktmongo.dsl.KtMongoDsl +import opensavvy.ktmongo.dsl.LowLevelApi +import opensavvy.ktmongo.dsl.path.Field +import opensavvy.ktmongo.dsl.path.FieldDsl +import opensavvy.ktmongo.dsl.path.Path +import opensavvy.ktmongo.dsl.tree.AbstractBsonNode +import opensavvy.ktmongo.dsl.tree.AbstractCompoundBsonNode +import opensavvy.ktmongo.dsl.tree.BsonNode +import opensavvy.ktmongo.dsl.tree.CompoundBsonNode +import kotlin.reflect.KProperty1 +import kotlin.time.Duration +import kotlin.time.Duration.Companion.seconds +import kotlin.time.Instant + +/** + * The options to define a time-series collection. + * + * For more information, see [HasTimeSeries.timeSeries]. + */ +class TimeSeriesOption( + val node: BsonNode, + context: BsonContext, +) : AbstractOption("timeseries", context) { + + init { + @OptIn(LowLevelApi::class) + node.freeze() + } + + @LowLevelApi + override fun write(writer: BsonValueWriter) = with(writer) { + writeDocument { + node.writeTo(this) + } + } +} + +/** + * Declares time-series collections. + * + * See [timeSeries]. + */ +@KtMongoDsl +interface HasTimeSeries : Options { + + /** + * Time-series are special collections optimized for cases where data is created over time. + * + * Examples: + * - A temperature sensor that makes readings every minute. + * - An audit log where a new document is inserted each time a user makes an action. + * - A log-in log where a new document is inserted on every log-in attempt. + * + * Instead of storing documents independently like in normal collections, documents are grouped together + * by their source (see [TimeSeriesDsl.metaField]) in buckets of some time period (see [TimeSeriesDsl.granularity] + * and [TimeSeriesDsl.bucketMaxSpan]), which allows documents from the source and a similar time to be queried together + * much more efficiently. Additionally, MongoDB can use [compression](https://www.mongodb.com/docs/manual/core/timeseries/timeseries-compression/) + * to greatly reduce the disk space used by such a collection. + * + * Time-series documents do not require an `_id` field. + * + * ### Example + * + * ```kotlin + * class Notification( + * val userId: ObjectId, + * val createdAt: Instant, + * val text: String, + * ) + * + * notifications.create { + * timeSeries { + * timeField(Notification::createdAt) + * metaField(Notification::userId) + * granularity(Minutes) + * expiresAfter(7.days) + * } + * } + * ``` + * + * ### External resources + * + * - [Official documentation](https://www.mongodb.com/docs/manual/core/timeseries-collections/) + * + * @param block The specification of this time-series collection. + * - [TimeSeriesDsl.timeField] The field which contains the timestamp used to organize this time series. **Required**. + * - [TimeSeriesDsl.metaField] A field used to group documents together. + * - [TimeSeriesDsl.expiresAfter] Documents will be automatically removed from the collection after this duration. + * - [TimeSeriesDsl.granularity] How often documents are created in a given bucket. + * - [TimeSeriesDsl.bucketMaxSpan] A new bucket will be created after this duration. + */ + @OptIn(DangerousMongoApi::class, LowLevelApi::class) + fun timeSeries( + block: TimeSeriesDsl.() -> Unit, + ) { + val options = TimeSeriesDslImpl(this, context).apply(block) + accept(TimeSeriesOption(options, context)) + } + + private class TimeSeriesDslImpl( + private val optionsBlock: Options, + context: BsonContext, + ) : AbstractCompoundBsonNode(context), TimeSeriesDsl { + + @OptIn(DangerousMongoApi::class, LowLevelApi::class) + override fun timeField(field: Field) { + accept(TimeField(field.path, context)) + } + + @OptIn(DangerousMongoApi::class, LowLevelApi::class) + override fun expiresAfter(duration: Duration) { + optionsBlock.accept(ExpiresAfterOption(duration, context)) + } + + @OptIn(DangerousMongoApi::class, LowLevelApi::class) + override fun metaField(field: Field) { + accept(MetaField(field.path, context)) + } + + @OptIn(DangerousMongoApi::class, LowLevelApi::class) + override fun granularity(granularity: TimeSeriesDsl.Granularity) { + accept(GranularityField(granularity, context)) + } + + @OptIn(DangerousMongoApi::class, LowLevelApi::class) + override fun bucketMaxSpan(maxSpan: Duration) { + val span = + if (maxSpan == Duration.ZERO) maxSpan + else maxSpan.coerceAtLeast(1.seconds) + + accept(BucketMaxSpan(span, context)) + accept(BucketRounding(span, context)) + } + } + + @LowLevelApi + private class TimeField( + private val path: Path, + context: BsonContext, + ) : AbstractBsonNode(context) { + + @OptIn(LowLevelApi::class) + override fun write(writer: BsonFieldWriter) = with(writer) { + writeString("timeField", path.toString()) + } + } + + @LowLevelApi + private class MetaField( + private val path: Path, + context: BsonContext, + ) : AbstractBsonNode(context) { + + @OptIn(LowLevelApi::class) + override fun write(writer: BsonFieldWriter) = with(writer) { + writeString("metaField", path.toString()) + } + } + + @LowLevelApi + private class GranularityField( + private val granularity: TimeSeriesDsl.Granularity, + context: BsonContext, + ) : AbstractBsonNode(context) { + + @OptIn(LowLevelApi::class) + override fun write(writer: BsonFieldWriter) = with(writer) { + writeString("granularity", granularity.bsonName) + } + } + + @LowLevelApi + private class BucketMaxSpan( + private val duration: Duration, + context: BsonContext, + ) : AbstractBsonNode(context) { + + @OptIn(LowLevelApi::class) + override fun write(writer: BsonFieldWriter) = with(writer) { + writeInt64("bucketMaxSpanSeconds", duration.inWholeSeconds) + } + } + + @LowLevelApi + private class BucketRounding( + private val duration: Duration, + context: BsonContext, + ) : AbstractBsonNode(context) { + + @OptIn(LowLevelApi::class) + override fun write(writer: BsonFieldWriter) = with(writer) { + writeInt64("bucketRoundingSeconds", duration.inWholeSeconds) + } + } +} + +/** + * Options for [HasTimeSeries.timeSeries]. + */ +@KtMongoDsl +interface TimeSeriesDsl : CompoundBsonNode, FieldDsl { + + /** + * Specifies the field which contains the date used to organize this time-series collection. + * + * ### Example + * + * ```kotlin + * class Notification( + * val userId: ObjectId, + * val createdAt: Instant, + * val text: String, + * ) + * + * notifications.create { + * timeSeries { + * timeField(Notification::createdAt) + * metaField(Notification::userId) + * granularity(Minutes) + * expiresAfter(7.days) + * } + * } + * ``` + * + * ### External resources + * + * - [Official documentation](https://www.mongodb.com/docs/manual/core/timeseries/timeseries-bucketing/) + */ + fun timeField(field: Field) + + /** + * Specifies the field which contains the date used to organize this time-series collection. + * + * ### Example + * + * ```kotlin + * class Notification( + * val userId: ObjectId, + * val createdAt: Instant, + * val text: String, + * ) + * + * notifications.create { + * timeSeries { + * timeField(Notification::createdAt) + * metaField(Notification::userId) + * granularity(Minutes) + * expiresAfter(7.days) + * } + * } + * ``` + * + * ### External resources + * + * - [Official documentation](https://www.mongodb.com/docs/manual/core/timeseries/timeseries-bucketing/) + */ + fun timeField(field: KProperty1) { + timeField(field.field) + } + + /** + * Specifies how long documents are retained in the collection. + * + * Once this time has passed after the value of the [timeField], MongoDB automatically removes the document. + * + * The minimum resolution is one second, so values less than one second are rounded up to one second. + * + * **Note.** This method sets the [ExpiresAfterOption], which does not appear in the `toString()` representation + * of this DSL, because MongoDB places this option at the root of the `create {}` command instead of within + * the `timeseries` block. + * + * ### Example + * + * ```kotlin + * class Notification( + * val userId: ObjectId, + * val createdAt: Instant, + * val text: String, + * ) + * + * notifications.create { + * timeSeries { + * timeField(Notification::createdAt) + * metaField(Notification::userId) + * granularity(Minutes) + * expiresAfter(7.days) + * } + * } + * ``` + * + * ### External resources + * + * - [Official documentation](https://www.mongodb.com/docs/manual/core/timeseries/timeseries-automatic-removal) + */ + fun expiresAfter(duration: Duration) + + /** + * Controls how MongoDB buckets documents. + * + * To improve performance, MongoDB stores time-series documents in buckets, with documents closely related to each other + * in the same bucket for improved cache performance. + * + * Specifying a meta-field is optional but recommended when multiple documents that originated at around the same time + * are likely to be queried together. + * + * If two documents have the same meta-field value and a similar time-field value, MongoDB stores them together. + * + * A good meta-field is a field that appears in almost all queries and values that are commonly used, for + * example, the origin of the event or its destination. + * + * ### Example + * + * ```kotlin + * class Notification( + * val userId: ObjectId, + * val createdAt: Instant, + * val text: String, + * ) + * + * notifications.create { + * timeSeries { + * timeField(Notification::createdAt) + * metaField(Notification::userId) + * granularity(Minutes) + * } + * } + * ``` + * + * In this example, we store notifications destined for some users. + * We bucket the data by the user each notification is destined for, to ensure we can query all notifications for a given user quickly. + * + * ### External resources + * + * - [How bucketing works](https://www.mongodb.com/docs/manual/core/timeseries/timeseries-bucketing/#how-bucketing-works) + * - [Best practices](https://www.mongodb.com/docs/manual/core/timeseries/timeseries-considerations/) + */ + fun metaField(field: Field) + + /** + * Controls how MongoDB buckets documents. + * + * To improve performance, MongoDB stores time-series documents in buckets, with documents closely related to each other + * in the same bucket for improved cache performance. + * + * Specifying a meta-field is optional but recommended when multiple documents that originated at around the same time + * are likely to be queried together. + * + * If two documents have the same meta-field value and a similar time-field value, MongoDB stores them together. + * + * A good meta-field is a field that appears in almost all queries and values that are commonly used, for + * example, the origin of the event or its destination. + * + * ### Example + * + * ```kotlin + * class Notification( + * val userId: ObjectId, + * val createdAt: Instant, + * val text: String, + * ) + * + * notifications.create { + * timeSeries { + * timeField(Notification::createdAt) + * metaField(Notification::userId) + * granularity(Minutes) + * } + * } + * ``` + * + * In this example, we store notifications destined for some users. + * We bucket the data by the user each notification is destined for, to ensure we can query all notifications for a given user quickly. + * + * ### External resources + * + * - [How bucketing works](https://www.mongodb.com/docs/manual/core/timeseries/timeseries-bucketing/#how-bucketing-works) + * - [Best practices](https://www.mongodb.com/docs/manual/core/timeseries/timeseries-considerations/) + */ + fun metaField(field: KProperty1) { + metaField(field.field) + } + + /** + * Specifies the maximum size of a bucket. + * + * To learn more about buckets, see [metaField]. + * + * This option is not compatible with [bucketMaxSpan]. + * + * This option communicates to MongoDB how many documents will be created in each bucket. + * For example, if [metaField] is a company ID, and a new document is created by each company every few minutes, choose [Granularity.Minutes]. + * + * If this option is not specified, [Granularity.Seconds] is used. + * + * ### Example + * + * ```kotlin + * class Notification( + * val userId: ObjectId, + * val createdAt: Instant, + * val text: String, + * ) + * + * notifications.create { + * timeSeries { + * timeField(Notification::createdAt) + * metaField(Notification::userId) + * granularity(Minutes) + * } + * } + * ``` + * + * ### External resources + * + * - [Official documentation](https://www.mongodb.com/docs/manual/core/timeseries/timeseries-considerations/#granularity) + */ + fun granularity(granularity: Granularity) + + /** + * Specifies the maximum time between two timestamps in the same bucket. + * + * To learn more about buckets, see [metaField]. + * + * This option is not compatible with [granularity]. + * + * This option communicates to MongoDB how long to keep a bucket open. + * Once this time has passed since the creation of the bucket, MongoDB will close the bucket and create a new one. + * + * The minimum resolution is one second, so values less than one second are rounded up to one second. + * + * This function sets both the `bucketMaxSpanSeconds` and `bucketRoundingSeconds` options. + * + * ### Example + * + * ```kotlin + * class Notification( + * val userId: ObjectId, + * val createdAt: Instant, + * val text: String, + * ) + * + * notifications.create { + * timeSeries { + * timeField(Notification::createdAt) + * metaField(Notification::userId) + * bucketMaxSpan(12.hours) + * } + * } + * ``` + * + * ### External resources + * + * - [Official documentation](https://www.mongodb.com/docs/manual/core/timeseries/timeseries-granularity/#using-custom-bucketing-parameters) + */ + fun bucketMaxSpan( + maxSpan: Duration, + ) + + /** + * The granularity of the time-series buckets. + * + * For more information, see [metaField]. + * + * ### External resources + * + * - [Official documentation](https://www.mongodb.com/docs/manual/core/timeseries/timeseries-considerations/#granularity) + */ + enum class Granularity(val bsonName: String) { + /** + * Select this granularity if a new document is created every few seconds for a given [metaField] value. + * + * A bucket will store up to an hour of data. + */ + Seconds("seconds"), + + /** + * Select this granularity if a new document is created every few minutes for a given [metaField] value. + * + * A bucket will store up to 24 hours of data. + */ + Minutes("minutes"), + + /** + * Select this granularity if a new document is created every few hours for a given [metaField] value. + * + * A bucket will store up to a month of data. + */ + Hours("hours"), + } +} diff --git a/dsl/src/commonTest/kotlin/options/TimeSeriesTest.kt b/dsl/src/commonTest/kotlin/options/TimeSeriesTest.kt new file mode 100644 index 00000000..61d8fe49 --- /dev/null +++ b/dsl/src/commonTest/kotlin/options/TimeSeriesTest.kt @@ -0,0 +1,97 @@ +/* + * Copyright (c) 2026, OpenSavvy and contributors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package opensavvy.ktmongo.dsl.options + +import opensavvy.ktmongo.bson.types.ObjectId +import opensavvy.ktmongo.dsl.command.CreateCollectionOptions +import opensavvy.ktmongo.dsl.multiContextSuite +import opensavvy.ktmongo.dsl.query.shouldBeBson +import opensavvy.ktmongo.dsl.testContext +import kotlin.time.Duration.Companion.hours +import kotlin.time.Duration.Companion.minutes +import kotlin.time.Duration.Companion.seconds +import kotlin.time.Instant + +val TimeSeriesTest by multiContextSuite { + + class Target( + val date: Instant, + val sensorId: ObjectId, + ) + + test("Simple time series") { + val options = CreateCollectionOptions(testContext()) + + options.timeSeries { + timeField(Target::date) + metaField(Target::sensorId) + } + + options.shouldBeBson(""" + { + "timeseries": { + "timeField": "date", + "metaField": "sensorId" + } + } + """.trimIndent()) + } + + test("With expiration") { + val options = CreateCollectionOptions(testContext()) + + options.timeSeries { + timeField(Target::date) + expiresAfter(1.minutes + 37.seconds) + } + + options.shouldBeBson(""" + { + "expireAfterSeconds": 97, + "timeseries": { + "timeField": "date" + } + } + """.trimIndent()) + } + + test("Bucket control") { + val options = CreateCollectionOptions(testContext()) + + options.timeSeries { + timeField(Target::date) + metaField(Target::sensorId) + + // These fields can't be used together, but it doesn't matter for this test + granularity(TimeSeriesDsl.Granularity.Hours) + bucketMaxSpan(2.hours + 30.minutes) + } + + options.shouldBeBson(""" + { + "timeseries": { + "timeField": "date", + "metaField": "sensorId", + "granularity": "hours", + "bucketMaxSpanSeconds": 9000, + "bucketRoundingSeconds": 9000 + } + } + """.trimIndent()) + } + +} diff --git a/test/src/commonMain/kotlin/operations/CollectionOperations.test.kt b/test/src/commonMain/kotlin/operations/CollectionOperations.test.kt index ff09abd0..43be7792 100644 --- a/test/src/commonMain/kotlin/operations/CollectionOperations.test.kt +++ b/test/src/commonMain/kotlin/operations/CollectionOperations.test.kt @@ -14,7 +14,7 @@ * limitations under the License. */ -@file:OptIn(LowLevelApi::class, ExperimentalCoroutinesApi::class) +@file:OptIn(LowLevelApi::class, ExperimentalCoroutinesApi::class, ExperimentalTime::class) package opensavvy.ktmongo.tests.api.operations @@ -31,6 +31,8 @@ import opensavvy.prepared.suite.SuiteDsl import opensavvy.prepared.suite.assertions.checkThrows import opensavvy.prepared.suite.now import opensavvy.prepared.suite.time +import kotlin.time.Clock +import kotlin.time.Duration.Companion.minutes import kotlin.time.ExperimentalTime import kotlin.time.Instant @@ -41,6 +43,13 @@ data class CollectionOperationsUser @OptIn(ExperimentalTime::class) constructor( val birthdate: @Serializable(with = InstantAsBsonDatetimeSerializer::class) Instant, ) +@Serializable +data class CollectionOperationsAuditLog( + val timestamp: @Serializable(with = InstantAsBsonDatetimeSerializer::class) Instant, + val user: ObjectId, + val action: String, +) + fun SuiteDsl.verifyCollectionOperations( client: Prepared, ) = suite("Collection operations") { @@ -140,4 +149,31 @@ fun SuiteDsl.verifyCollectionOperations( } } } + + val auditLog by client.collection("operation-collection-audit-log") + + test("Create a time-series collection") { + auditLog().create { + timeSeries { + timeField(CollectionOperationsAuditLog::timestamp) + metaField(CollectionOperationsAuditLog::user) + expiresAfter(30.minutes) + } + } + + val users = List(3) { auditLog().newId() } + val actions = listOf("login", "logout", "failed-login") + + auditLog().insertMany( + List(100) { + CollectionOperationsAuditLog( + timestamp = Clock.System.now(), // purposefully use real time to avoid duplicate documents + user = users.random(), + action = actions.random(), + ) + } + ) + + check(auditLog().count() == 100L) + } }