Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ package org.opensearch.commons.alerting.action
import org.opensearch.commons.alerting.model.DocumentLevelTriggerRunResult
import org.opensearch.commons.alerting.model.InputRunResults
import org.opensearch.commons.alerting.util.AlertingException
import org.opensearch.commons.alerting.util.readMapAsMutableMap
import org.opensearch.core.action.ActionResponse
import org.opensearch.core.common.io.stream.StreamInput
import org.opensearch.core.common.io.stream.StreamOutput
Expand All @@ -30,9 +31,9 @@ class DocLevelMonitorFanOutResponse : ActionResponse, ToXContentObject {
nodeId = sin.readString(),
executionId = sin.readString(),
monitorId = sin.readString(),
lastRunContexts = sin.readMap()!! as MutableMap<String, Any>,
lastRunContexts = sin.readMapAsMutableMap() as MutableMap<String, Any>,
inputResults = InputRunResults.readFrom(sin),
triggerResults = suppressWarning(sin.readMap(StreamInput::readString, DocumentLevelTriggerRunResult::readFrom)),
triggerResults = readTriggerResults(sin),
exception = sin.readException()
)

Expand Down Expand Up @@ -84,9 +85,11 @@ class DocLevelMonitorFanOutResponse : ActionResponse, ToXContentObject {
}

companion object {
@Suppress("UNCHECKED_CAST")
fun suppressWarning(map: MutableMap<String?, Any?>?): Map<String, DocumentLevelTriggerRunResult> {
return map as Map<String, DocumentLevelTriggerRunResult>
private fun readTriggerResults(sin: StreamInput): Map<String, DocumentLevelTriggerRunResult> {
val raw = sin.readMap(StreamInput::readString, DocumentLevelTriggerRunResult::readFrom)
if (raw.isEmpty()) return mutableMapOf()
@Suppress("UNCHECKED_CAST")
return HashMap(raw as Map<String, DocumentLevelTriggerRunResult>)
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -42,7 +42,8 @@ class BucketSelectorExtAggregationBuilder :
@Throws(IOException::class)
@Suppress("UNCHECKED_CAST")
constructor(sin: StreamInput) : super(sin, NAME.preferredName) {
bucketsPathsMap = sin.readMap() as MutableMap<String, String>
@Suppress("UNCHECKED_CAST")
bucketsPathsMap = (sin.readMap()?.toMutableMap() ?: mutableMapOf()) as Map<String, String>
script = Script(sin)
gapPolicy = BucketHelpers.GapPolicy.readFrom(sin)
parentBucketPath = sin.readString()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -3,11 +3,11 @@ package org.opensearch.commons.alerting.model
import org.opensearch.Version
import org.opensearch.common.lucene.uid.Versions
import org.opensearch.commons.alerting.alerts.AlertError
import org.opensearch.commons.alerting.model.Monitor.Companion.suppressWarning
import org.opensearch.commons.alerting.util.IndexUtils.Companion.NO_SCHEMA_VERSION
import org.opensearch.commons.alerting.util.instant
import org.opensearch.commons.alerting.util.optionalTimeField
import org.opensearch.commons.alerting.util.optionalUserField
import org.opensearch.commons.alerting.util.readMapAsMutableMap
import org.opensearch.commons.authuser.User
import org.opensearch.core.common.io.stream.StreamInput
import org.opensearch.core.common.io.stream.StreamOutput
Expand Down Expand Up @@ -423,7 +423,7 @@ data class Alert(
null
},
queryResults = if (sin.version.onOrAfter(Version.V_3_7_0)) {
sin.readList { input -> suppressWarning(input.readMap()) }
sin.readList { input -> input.readMapAsMutableMap() }
} else {
listOf()
},
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@

package org.opensearch.commons.alerting.model

import org.opensearch.commons.alerting.util.readMapAsMutableMap
import org.opensearch.core.common.io.stream.StreamInput
import org.opensearch.core.common.io.stream.StreamOutput
import org.opensearch.core.xcontent.ToXContent
Expand All @@ -24,7 +25,7 @@ data class BucketLevelTriggerRunResult(
sin.readString(),
sin.readException() as Exception?, // error
sin.readMap(StreamInput::readString, ::AggregationResultBucket),
sin.readMap() as MutableMap<String, MutableMap<String, ActionRunResult>>
sin.readMapAsMutableMap() as MutableMap<String, MutableMap<String, ActionRunResult>>
)

override fun internalXContent(builder: XContentBuilder, params: ToXContent.Params): XContentBuilder {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@
package org.opensearch.commons.alerting.model

import org.opensearch.commons.alerting.alerts.AlertError
import org.opensearch.commons.alerting.util.readMapAsMutableMap
import org.opensearch.core.common.io.stream.StreamInput
import org.opensearch.core.common.io.stream.StreamOutput
import org.opensearch.core.xcontent.ToXContent
Expand All @@ -28,7 +29,7 @@ data class ChainedAlertTriggerRunResult(
triggerName = sin.readString(),
error = sin.readException(),
triggered = sin.readBoolean(),
actionResults = sin.readMap() as MutableMap<String, ActionRunResult>,
actionResults = sin.readMapAsMutableMap() as MutableMap<String, ActionRunResult>,
associatedAlertIds = sin.readStringList().toSet()
)

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@
package org.opensearch.commons.alerting.model

import org.opensearch.commons.alerting.alerts.AlertError
import org.opensearch.commons.alerting.util.readMapAsMutableMap
import org.opensearch.core.common.io.stream.StreamInput
import org.opensearch.core.common.io.stream.StreamOutput
import org.opensearch.core.common.io.stream.Writeable
Expand Down Expand Up @@ -35,7 +36,7 @@ data class ClusterMetricsTriggerRunResult(
triggerName = sin.readString(),
error = sin.readException(),
triggered = sin.readBoolean(),
actionResults = sin.readMap() as MutableMap<String, ActionRunResult>,
actionResults = sin.readMapAsMutableMap() as MutableMap<String, ActionRunResult>,
clusterTriggerResults = sin.readList((ClusterTriggerResult)::readFrom)
)

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@
package org.opensearch.commons.alerting.model

import org.opensearch.Version
import org.opensearch.commons.alerting.util.readMapAsMutableMap
import org.opensearch.core.common.io.stream.StreamInput
import org.opensearch.core.common.io.stream.StreamOutput
import org.opensearch.core.common.io.stream.Writeable
Expand All @@ -29,8 +30,8 @@ data class IndexExecutionContext(
@Throws(IOException::class)
constructor(sin: StreamInput) : this(
queries = sin.readList { DocLevelQuery(sin) },
lastRunContext = sin.readMap() as MutableMap<String, Any>,
updatedLastRunContext = sin.readMap() as MutableMap<String, Any>,
lastRunContext = sin.readMapAsMutableMap() as MutableMap<String, Any>,
updatedLastRunContext = sin.readMapAsMutableMap() as MutableMap<String, Any>,
indexName = sin.readString(),
concreteIndexName = sin.readString(),
updatedIndexNames = sin.readStringList(),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ import org.opensearch.commons.alerting.util.isBucketLevelMonitor
import org.opensearch.commons.alerting.util.isPPLMonitor
import org.opensearch.commons.alerting.util.optionalTimeField
import org.opensearch.commons.alerting.util.optionalUserField
import org.opensearch.commons.alerting.util.readMapAsMutableMap
import org.opensearch.commons.authuser.User
import org.opensearch.core.ParseField
import org.opensearch.core.common.io.stream.StreamInput
Expand Down Expand Up @@ -131,7 +132,7 @@ data class Monitor(
schemaVersion = sin.readInt(),
inputs = sin.readList((Input)::readFrom),
triggers = sin.readList((Trigger)::readFrom),
uiMetadata = suppressWarning(sin.readMap()),
uiMetadata = sin.readMapAsMutableMap(),
dataSources = if (sin.readBoolean()) {
DataSources(sin)
} else {
Expand Down Expand Up @@ -350,7 +351,7 @@ data class Monitor(
var schedule: Schedule? = null
var lastUpdateTime: Instant? = null
var enabledTime: Instant? = null
var uiMetadata: Map<String, Any> = mapOf()
var uiMetadata: Map<String, Any> = mutableMapOf()
var enabled = true
var schemaVersion = NO_SCHEMA_VERSION
val triggers: MutableList<Trigger> = mutableListOf()
Expand Down Expand Up @@ -407,7 +408,7 @@ data class Monitor(
}
ENABLED_TIME_FIELD -> enabledTime = xcp.instant()
LAST_UPDATE_TIME_FIELD -> lastUpdateTime = xcp.instant()
UI_METADATA_FIELD -> uiMetadata = xcp.map()
UI_METADATA_FIELD -> uiMetadata = xcp.map().toMutableMap()
DATA_SOURCES_FIELD -> dataSources = if (xcp.currentToken() == XContentParser.Token.VALUE_NULL) {
DataSources()
} else {
Expand Down Expand Up @@ -480,10 +481,5 @@ data class Monitor(
fun readFrom(sin: StreamInput): Monitor? {
return Monitor(sin)
}

@Suppress("UNCHECKED_CAST")
fun suppressWarning(map: MutableMap<String?, Any?>?): MutableMap<String, Any> {
return map as MutableMap<String, Any>
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ package org.opensearch.commons.alerting.model

import org.opensearch.commons.alerting.model.Monitor.Companion.NO_ID
import org.opensearch.commons.alerting.util.instant
import org.opensearch.commons.alerting.util.readMapAsMutableMap
import org.opensearch.core.common.io.stream.StreamInput
import org.opensearch.core.common.io.stream.StreamOutput
import org.opensearch.core.common.io.stream.Writeable
Expand Down Expand Up @@ -36,8 +37,8 @@ data class MonitorMetadata(
primaryTerm = sin.readLong(),
monitorId = sin.readString(),
lastActionExecutionTimes = sin.readList(ActionExecutionTime.Companion::readFrom),
lastRunContext = Monitor.suppressWarning(sin.readMap()),
sourceToQueryIndexMapping = sin.readMap() as MutableMap<String, String>
lastRunContext = sin.readMapAsMutableMap(),
sourceToQueryIndexMapping = sin.readMapAsMutableMap() as MutableMap<String, String>
)

override fun writeTo(out: StreamOutput) {
Expand All @@ -47,7 +48,7 @@ data class MonitorMetadata(
out.writeString(monitorId)
out.writeCollection(lastActionExecutionTimes)
out.writeMap(lastRunContext)
out.writeMap(sourceToQueryIndexMapping as MutableMap<String, Any>)
out.writeMap(sourceToQueryIndexMapping as Map<String, Any>)
}

override fun toXContent(builder: XContentBuilder, params: ToXContent.Params): XContentBuilder {
Expand All @@ -57,7 +58,7 @@ data class MonitorMetadata(
.field(LAST_ACTION_EXECUTION_FIELD, lastActionExecutionTimes.toTypedArray())
if (lastRunContext.isNotEmpty()) builder.field(LAST_RUN_CONTEXT_FIELD, lastRunContext)
if (sourceToQueryIndexMapping.isNotEmpty()) {
builder.field(SOURCE_TO_QUERY_INDEX_MAP_FIELD, sourceToQueryIndexMapping as MutableMap<String, Any>)
builder.field(SOURCE_TO_QUERY_INDEX_MAP_FIELD, sourceToQueryIndexMapping as Map<String, Any>)
}
if (params.paramAsBoolean("with_type", false)) builder.endObject()
return builder.endObject()
Expand All @@ -81,7 +82,7 @@ data class MonitorMetadata(
): MonitorMetadata {
lateinit var monitorId: String
val lastActionExecutionTimes = mutableListOf<ActionExecutionTime>()
var lastRunContext: Map<String, Any> = mapOf()
var lastRunContext: Map<String, Any> = mutableMapOf()
var sourceToQueryIndexMapping: MutableMap<String, String> = mutableMapOf()

XContentParserUtils.ensureExpectedToken(XContentParser.Token.START_OBJECT, xcp.currentToken(), xcp)
Expand All @@ -97,8 +98,8 @@ data class MonitorMetadata(
lastActionExecutionTimes.add(ActionExecutionTime.parse(xcp))
}
}
LAST_RUN_CONTEXT_FIELD -> lastRunContext = xcp.map()
SOURCE_TO_QUERY_INDEX_MAP_FIELD -> sourceToQueryIndexMapping = xcp.map() as MutableMap<String, String>
LAST_RUN_CONTEXT_FIELD -> lastRunContext = xcp.map().toMutableMap()
SOURCE_TO_QUERY_INDEX_MAP_FIELD -> sourceToQueryIndexMapping = (xcp.map()?.toMutableMap() ?: mutableMapOf()) as MutableMap<String, String>
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ import org.opensearch.OpenSearchException
import org.opensearch.Version
import org.opensearch.commons.alerting.alerts.AlertError
import org.opensearch.commons.alerting.util.optionalTimeField
import org.opensearch.commons.alerting.util.readMapAsMutableMap
import org.opensearch.core.common.io.stream.StreamInput
import org.opensearch.core.common.io.stream.StreamOutput
import org.opensearch.core.common.io.stream.Writeable
Expand All @@ -36,7 +37,7 @@ data class MonitorRunResult<TriggerResult : TriggerRunResult>(
sin.readInstant(), // periodEnd
sin.readException(), // error
InputRunResults.readFrom(sin), // inputResults
suppressWarning(sin.readMap()) as Map<String, TriggerResult> // triggerResults
sin.readMapAsMutableMap() as Map<String, TriggerResult> // triggerResults
)

override fun toXContent(builder: XContentBuilder, params: ToXContent.Params): XContentBuilder {
Expand Down Expand Up @@ -72,11 +73,6 @@ data class MonitorRunResult<TriggerResult : TriggerRunResult>(
fun readFrom(sin: StreamInput): MonitorRunResult<TriggerRunResult> {
return MonitorRunResult(sin)
}

@Suppress("UNCHECKED_CAST")
fun suppressWarning(map: MutableMap<String?, Any?>?): Map<String, TriggerRunResult> {
return map as Map<String, TriggerRunResult>
}
}

@Throws(IOException::class)
Expand Down Expand Up @@ -147,7 +143,7 @@ data class InputRunResults(
val count = sin.readVInt() // count
val list = mutableListOf<Map<String, Any>>()
for (i in 0 until count) {
list.add(suppressWarning(sin.readMap())) // result(map)
list.add(sin.readMapAsMutableMap()) // result(map)
}
val pplCount = if (sin.version.onOrAfter(Version.V_3_7_0)) {
sin.readVInt()
Expand All @@ -157,7 +153,7 @@ data class InputRunResults(
val pplList = mutableListOf<Map<String, Any?>>()
if (sin.version.onOrAfter(Version.V_3_7_0)) {
for (i in 0 until pplCount) {
pplList.add(suppressWarning(sin.readMap())) // pplResults
pplList.add(sin.readMapAsMutableMap()) // pplResults
}
}
val pplNumResults = if (sin.version.onOrAfter(Version.V_3_7_0)) {
Expand All @@ -168,11 +164,6 @@ data class InputRunResults(
val error = sin.readException<Exception>() // error
return InputRunResults(list, error, null, pplList, pplNumResults)
}

@Suppress("UNCHECKED_CAST")
fun suppressWarning(map: MutableMap<String?, Any?>?): Map<String, Any> {
return map as Map<String, Any>
}
}

fun afterKeysPresent(): Boolean {
Expand All @@ -196,11 +187,12 @@ data class ActionRunResult(
val error: Exception? = null
) : Writeable, ToXContent {

@Suppress("UNCHECKED_CAST")
@Throws(IOException::class)
constructor(sin: StreamInput) : this(
sin.readString(), // actionId
sin.readString(), // actionName
suppressWarning(sin.readMap()), // output
sin.readMapAsMutableMap() as Map<String, String>, // output
sin.readBoolean(), // throttled
sin.readOptionalInstant(), // executionTime
sin.readException() // error
Expand Down Expand Up @@ -233,11 +225,6 @@ data class ActionRunResult(
fun readFrom(sin: StreamInput): ActionRunResult {
return ActionRunResult(sin)
}

@Suppress("UNCHECKED_CAST")
fun suppressWarning(map: MutableMap<String?, Any?>?): MutableMap<String, String> {
return map as MutableMap<String, String>
}
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ package org.opensearch.commons.alerting.model

import org.opensearch.Version
import org.opensearch.commons.alerting.alerts.AlertError
import org.opensearch.commons.alerting.util.readMapAsMutableMap
import org.opensearch.core.common.io.stream.StreamInput
import org.opensearch.core.common.io.stream.StreamOutput
import org.opensearch.core.xcontent.ToXContent
Expand All @@ -29,7 +30,7 @@ open class QueryLevelTriggerRunResult(
triggerName = sin.readString(),
error = sin.readException(),
triggered = sin.readBoolean(),
actionResults = sin.readMap() as MutableMap<String, ActionRunResult>,
actionResults = sin.readMapAsMutableMap() as MutableMap<String, ActionRunResult>,
pplCustomQueryResults = if (sin.version.onOrAfter(Version.V_3_7_0)) {
sin.readList { it.readMap() }
} else {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -45,11 +45,4 @@ abstract class TriggerRunResult(
out.writeString(triggerName)
out.writeException(error)
}

companion object {
@Suppress("UNCHECKED_CAST")
fun suppressWarning(map: MutableMap<String?, Any?>?): MutableMap<String, ActionRunResult> {
return map as MutableMap<String, ActionRunResult>
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -294,11 +294,6 @@ data class Workflow(
return Workflow(sin)
}

@Suppress("UNCHECKED_CAST")
fun suppressWarning(map: MutableMap<String?, Any?>?): MutableMap<String, Any> {
return map as MutableMap<String, Any>
}

private const val DEFAULT_OWNER = "alerting"
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@

package org.opensearch.commons.alerting.model

import org.opensearch.commons.alerting.util.readMapAsMutableMap
import org.opensearch.core.common.io.stream.StreamInput
import org.opensearch.core.common.io.stream.StreamOutput
import org.opensearch.core.common.io.stream.Writeable
Expand Down Expand Up @@ -35,7 +36,7 @@ data class WorkflowRunResult(
executionEndTime = sin.readOptionalInstant(),
executionId = sin.readString(),
error = sin.readException(),
triggerResults = suppressWarning(sin.readMap()) as Map<String, ChainedAlertTriggerRunResult>
triggerResults = sin.readMapAsMutableMap() as Map<String, ChainedAlertTriggerRunResult>
)

override fun writeTo(out: StreamOutput) {
Expand Down Expand Up @@ -73,10 +74,5 @@ data class WorkflowRunResult(
fun readFrom(sin: StreamInput): WorkflowRunResult {
return WorkflowRunResult(sin)
}

@Suppress("UNCHECKED_CAST")
fun suppressWarning(map: MutableMap<String?, Any?>?): Map<String, TriggerRunResult> {
return map as Map<String, TriggerRunResult>
}
}
}
Loading