Skip to content
Merged
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 @@ -14,6 +14,7 @@ import org.opensearch.alerting.action.GetEmailGroupAction
import org.opensearch.alerting.action.GetRemoteIndexesAction
import org.opensearch.alerting.action.SearchEmailAccountAction
import org.opensearch.alerting.action.SearchEmailGroupAction
import org.opensearch.alerting.actionv2.DeleteMonitorV2Action
import org.opensearch.alerting.actionv2.GetMonitorV2Action
import org.opensearch.alerting.actionv2.IndexMonitorV2Action
import org.opensearch.alerting.actionv2.SearchMonitorV2Action
Expand Down Expand Up @@ -55,6 +56,7 @@ import org.opensearch.alerting.resthandler.RestSearchAlertingCommentAction
import org.opensearch.alerting.resthandler.RestSearchEmailAccountAction
import org.opensearch.alerting.resthandler.RestSearchEmailGroupAction
import org.opensearch.alerting.resthandler.RestSearchMonitorAction
import org.opensearch.alerting.resthandlerv2.RestDeleteMonitorV2Action
import org.opensearch.alerting.resthandlerv2.RestGetMonitorV2Action
import org.opensearch.alerting.resthandlerv2.RestIndexMonitorV2Action
import org.opensearch.alerting.resthandlerv2.RestSearchMonitorV2Action
Expand Down Expand Up @@ -90,6 +92,7 @@ import org.opensearch.alerting.transport.TransportSearchAlertingCommentAction
import org.opensearch.alerting.transport.TransportSearchEmailAccountAction
import org.opensearch.alerting.transport.TransportSearchEmailGroupAction
import org.opensearch.alerting.transport.TransportSearchMonitorAction
import org.opensearch.alerting.transportv2.TransportDeleteMonitorV2Action
import org.opensearch.alerting.transportv2.TransportGetMonitorV2Action
import org.opensearch.alerting.transportv2.TransportIndexMonitorV2Action
import org.opensearch.alerting.transportv2.TransportSearchMonitorV2Action
Expand Down Expand Up @@ -233,6 +236,7 @@ internal class AlertingPlugin : PainlessExtension, ActionPlugin, ScriptPlugin, R

// Alerting V2
RestIndexMonitorV2Action(),
RestDeleteMonitorV2Action(),
RestGetMonitorV2Action(),
RestSearchMonitorV2Action(settings, clusterService),
)
Expand Down Expand Up @@ -273,6 +277,7 @@ internal class AlertingPlugin : PainlessExtension, ActionPlugin, ScriptPlugin, R
ActionPlugin.ActionHandler(IndexMonitorV2Action.INSTANCE, TransportIndexMonitorV2Action::class.java),
ActionPlugin.ActionHandler(GetMonitorV2Action.INSTANCE, TransportGetMonitorV2Action::class.java),
ActionPlugin.ActionHandler(SearchMonitorV2Action.INSTANCE, TransportSearchMonitorV2Action::class.java),
ActionPlugin.ActionHandler(DeleteMonitorV2Action.INSTANCE, TransportDeleteMonitorV2Action::class.java),
)
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,8 @@ import org.apache.lucene.search.TotalHits
import org.apache.lucene.search.TotalHits.Relation
import org.opensearch.action.search.SearchResponse
import org.opensearch.action.search.ShardSearchFailure
import org.opensearch.alerting.AlertingPlugin.Companion.MONITOR_BASE_URI
import org.opensearch.alerting.AlertingPlugin.Companion.MONITOR_V2_BASE_URI
import org.opensearch.alerting.modelv2.MonitorV2
import org.opensearch.commons.alerting.model.Monitor
import org.opensearch.commons.alerting.model.ScheduledJob
Expand All @@ -27,9 +29,16 @@ object AlertingV2Utils {
// returns the exception to pass into actionListener.onFailure if not.
fun validateMonitorV1(scheduledJob: ScheduledJob): Exception? {
if (scheduledJob is MonitorV2) {
return IllegalStateException("The ID given corresponds to a V2 Monitor, but a V1 Monitor was expected")
return IllegalStateException(
"The ID given corresponds to an Alerting V2 Monitor, but a V1 Monitor was expected. " +
"If you wish to operate on a V1 Monitor (e.g. Per Query, Per Document, etc), please use " +
"the Alerting V1 APIs with endpoint prefix: $MONITOR_BASE_URI."
)
} else if (scheduledJob !is Monitor && scheduledJob !is Workflow) {
return IllegalStateException("The ID given corresponds to a scheduled job of unknown type: ${scheduledJob.javaClass.name}")
return IllegalStateException(
"The ID given corresponds to a scheduled job of unknown type: ${scheduledJob.javaClass.name}. " +
"Please validate the ID and ensure it corresponds to a valid Monitor."
)
}
return null
}
Expand All @@ -38,9 +47,16 @@ object AlertingV2Utils {
// returns the exception to pass into actionListener.onFailure if not.
fun validateMonitorV2(scheduledJob: ScheduledJob): Exception? {
if (scheduledJob is Monitor || scheduledJob is Workflow) {
return IllegalStateException("The ID given corresponds to a V1 Monitor, but a V2 Monitor was expected")
return IllegalStateException(
"The ID given corresponds to an Alerting V1 Monitor, but a V2 Monitor was expected. " +
"If you wish to operate on a V2 Monitor (e.g. PPL Monitor), please use " +
"the Alerting V2 APIs with endpoint prefix: $MONITOR_V2_BASE_URI."
)
} else if (scheduledJob !is MonitorV2) {
return IllegalStateException("The ID given corresponds to a scheduled job of unknown type: ${scheduledJob.javaClass.name}")
return IllegalStateException(
"The ID given corresponds to a scheduled job of unknown type: ${scheduledJob.javaClass.name}. " +
"Please validate the ID and ensure it corresponds to a valid Monitor."
)
}
return null
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ import org.opensearch.action.admin.indices.exists.indices.IndicesExistsRequest
import org.opensearch.action.admin.indices.exists.indices.IndicesExistsResponse
import org.opensearch.action.search.SearchRequest
import org.opensearch.action.search.SearchResponse
import org.opensearch.alerting.AlertingV2Utils.validateMonitorV1
import org.opensearch.alerting.opensearchapi.suspendUntil
import org.opensearch.common.xcontent.LoggingDeprecationHandler
import org.opensearch.common.xcontent.XContentType
Expand Down Expand Up @@ -132,7 +133,12 @@ class WorkflowService(
xContentRegistry,
LoggingDeprecationHandler.INSTANCE, hit.sourceAsString
).use { hitsParser ->
val monitor = ScheduledJob.parse(hitsParser, hit.id, hit.version) as Monitor
val scheduledJob = ScheduledJob.parse(hitsParser, hit.id, hit.version)
validateMonitorV1(scheduledJob)?.let {
throw OpenSearchException(it)
}

val monitor = scheduledJob as Monitor
monitors.add(monitor)
}
}
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,15 @@
/*
* Copyright OpenSearch Contributors
* SPDX-License-Identifier: Apache-2.0
*/

package org.opensearch.alerting.actionv2

import org.opensearch.action.ActionType

class DeleteMonitorV2Action private constructor() : ActionType<DeleteMonitorV2Response>(NAME, ::DeleteMonitorV2Response) {
companion object {
val INSTANCE = DeleteMonitorV2Action()
const val NAME = "cluster:admin/opensearch/alerting/v2/monitor/delete"
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,39 @@
/*
* Copyright OpenSearch Contributors
* SPDX-License-Identifier: Apache-2.0
*/

package org.opensearch.alerting.actionv2

import org.opensearch.action.ActionRequest
import org.opensearch.action.ActionRequestValidationException
import org.opensearch.action.support.WriteRequest
import org.opensearch.core.common.io.stream.StreamInput
import org.opensearch.core.common.io.stream.StreamOutput
import java.io.IOException

class DeleteMonitorV2Request : ActionRequest {
val monitorV2Id: String
val refreshPolicy: WriteRequest.RefreshPolicy

constructor(monitorV2Id: String, refreshPolicy: WriteRequest.RefreshPolicy) : super() {
this.monitorV2Id = monitorV2Id
this.refreshPolicy = refreshPolicy
}

@Throws(IOException::class)
constructor(sin: StreamInput) : this(
monitorV2Id = sin.readString(),
refreshPolicy = WriteRequest.RefreshPolicy.readFrom(sin)
)

override fun validate(): ActionRequestValidationException? {
return null
}

@Throws(IOException::class)
override fun writeTo(out: StreamOutput) {
out.writeString(monitorV2Id)
refreshPolicy.writeTo(out)
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,43 @@
/*
* Copyright OpenSearch Contributors
* SPDX-License-Identifier: Apache-2.0
*/

package org.opensearch.alerting.actionv2

import org.opensearch.commons.alerting.util.IndexUtils
import org.opensearch.commons.notifications.action.BaseResponse
import org.opensearch.core.common.io.stream.StreamInput
import org.opensearch.core.common.io.stream.StreamOutput
import org.opensearch.core.xcontent.ToXContent
import org.opensearch.core.xcontent.XContentBuilder

class DeleteMonitorV2Response : BaseResponse {
var id: String
var version: Long

constructor(
id: String,
version: Long
) : super() {
this.id = id
this.version = version
}

constructor(sin: StreamInput) : this(
sin.readString(), // id
sin.readLong() // version
)

override fun writeTo(out: StreamOutput) {
out.writeString(id)
out.writeLong(version)
}

override fun toXContent(builder: XContentBuilder, params: ToXContent.Params): XContentBuilder {
return builder.startObject()
.field(IndexUtils._ID, id)
.field(IndexUtils._VERSION, version)
.endObject()
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,52 @@
/*
* Copyright OpenSearch Contributors
* SPDX-License-Identifier: Apache-2.0
*/

package org.opensearch.alerting.resthandlerv2

import org.apache.logging.log4j.LogManager
import org.apache.logging.log4j.Logger
import org.opensearch.action.support.WriteRequest.RefreshPolicy
import org.opensearch.alerting.AlertingPlugin
import org.opensearch.alerting.actionv2.DeleteMonitorV2Action
import org.opensearch.alerting.actionv2.DeleteMonitorV2Request
import org.opensearch.alerting.util.REFRESH
import org.opensearch.rest.BaseRestHandler
import org.opensearch.rest.RestHandler.Route
import org.opensearch.rest.RestRequest
import org.opensearch.rest.RestRequest.Method.DELETE
import org.opensearch.rest.action.RestToXContentListener
import org.opensearch.transport.client.node.NodeClient
import java.io.IOException

private val log: Logger = LogManager.getLogger(RestDeleteMonitorV2Action::class.java)

class RestDeleteMonitorV2Action : BaseRestHandler() {

override fun getName(): String {
return "delete_monitor_v2_action"
}

override fun routes(): List<Route> {
return mutableListOf(
Route(
DELETE,
"${AlertingPlugin.MONITOR_V2_BASE_URI}/{monitorV2Id}"
)
)
}

@Throws(IOException::class)
override fun prepareRequest(request: RestRequest, client: NodeClient): RestChannelConsumer {
val monitorV2Id = request.param("monitorV2Id")
log.info("${request.method()} ${AlertingPlugin.MONITOR_V2_BASE_URI}/$monitorV2Id")

val refreshPolicy = RefreshPolicy.parse(request.param(REFRESH, RefreshPolicy.IMMEDIATE.value))
val deleteMonitorV2Request = DeleteMonitorV2Request(monitorV2Id, refreshPolicy)

return RestChannelConsumer { channel ->
client.execute(DeleteMonitorV2Action.INSTANCE, deleteMonitorV2Request, RestToXContentListener(channel))
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ import org.opensearch.action.support.IndicesOptions
import org.opensearch.action.support.WriteRequest.RefreshPolicy
import org.opensearch.action.support.clustermanager.AcknowledgedResponse
import org.opensearch.alerting.MonitorMetadataService
import org.opensearch.alerting.actionv2.DeleteMonitorV2Response
import org.opensearch.alerting.core.lock.LockModel
import org.opensearch.alerting.core.lock.LockService
import org.opensearch.alerting.opensearchapi.suspendUntil
Expand Down Expand Up @@ -74,6 +75,19 @@ object DeleteMonitorService :
return DeleteMonitorResponse(deleteResponse.id, deleteResponse.version)
}

/**
* Deletes the monitorV2, which does not come with other metadata and queries
* like doc level monitors
* @param monitorV2Id monitorV2 ID to be deleted
* @param refreshPolicy
*/
suspend fun deleteMonitorV2(monitorV2Id: String, refreshPolicy: RefreshPolicy): DeleteMonitorV2Response {
val deleteResponse = deleteMonitor(monitorV2Id, refreshPolicy)
deleteLock(monitorV2Id)
return DeleteMonitorV2Response(deleteResponse.id, deleteResponse.version)
}

// both Alerting v1 and v2 workflows flow through this function
private suspend fun deleteMonitor(monitorId: String, refreshPolicy: RefreshPolicy): DeleteResponse {
val deleteMonitorRequest = DeleteRequest(ScheduledJob.SCHEDULED_JOBS_INDEX, monitorId)
.setRefreshPolicy(refreshPolicy)
Expand Down Expand Up @@ -167,7 +181,12 @@ object DeleteMonitorService :
}

private suspend fun deleteLock(monitor: Monitor) {
client.suspendUntil<Client, Boolean> { lockService.deleteLock(LockModel.generateLockId(monitor.id), it) }
deleteLock(monitor.id)
}

// both Alerting v1 and v2 workflows flow through this function
private suspend fun deleteLock(monitorId: String) {
client.suspendUntil<Client, Boolean> { lockService.deleteLock(LockModel.generateLockId(monitorId), it) }
}

/**
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ import org.opensearch.action.get.GetResponse
import org.opensearch.action.support.ActionFilters
import org.opensearch.action.support.HandledTransportAction
import org.opensearch.action.support.WriteRequest.RefreshPolicy
import org.opensearch.alerting.AlertingV2Utils.validateMonitorV1
import org.opensearch.alerting.opensearchapi.suspendUntil
import org.opensearch.alerting.service.DeleteMonitorService
import org.opensearch.alerting.settings.AlertingSettings
Expand Down Expand Up @@ -87,7 +88,7 @@ class TransportDeleteMonitorAction @Inject constructor(
) {
suspend fun resolveUserAndStart(refreshPolicy: RefreshPolicy) {
try {
val monitor = getMonitor()
val monitor = getMonitor() ?: return // null means there was an issue retrieving the Monitor

val canDelete = user == null || !doFilterForUser(user) ||
checkUserPermissionsWithResource(user, monitor.user, actionListener, "monitor", monitorId)
Expand Down Expand Up @@ -115,11 +116,11 @@ class TransportDeleteMonitorAction @Inject constructor(
}
}

private suspend fun getMonitor(): Monitor {
private suspend fun getMonitor(): Monitor? {
val getRequest = GetRequest(ScheduledJob.SCHEDULED_JOBS_INDEX, monitorId)

val getResponse: GetResponse = client.suspendUntil { get(getRequest, it) }
if (getResponse.isExists == false) {
if (!getResponse.isExists) {
actionListener.onFailure(
AlertingException.wrap(
OpenSearchStatusException("Monitor with $monitorId is not found", RestStatus.NOT_FOUND)
Expand All @@ -130,7 +131,16 @@ class TransportDeleteMonitorAction @Inject constructor(
xContentRegistry, LoggingDeprecationHandler.INSTANCE,
getResponse.sourceAsBytesRef, XContentType.JSON
)
return ScheduledJob.parse(xcp, getResponse.id, getResponse.version) as Monitor
val scheduledJob = ScheduledJob.parse(xcp, getResponse.id, getResponse.version)

validateMonitorV1(scheduledJob)?.let {
actionListener.onFailure(AlertingException.wrap(it))
return null
}

val monitor = scheduledJob as Monitor

return monitor
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@ import org.opensearch.action.search.SearchResponse
import org.opensearch.action.support.ActionFilters
import org.opensearch.action.support.HandledTransportAction
import org.opensearch.action.support.WriteRequest.RefreshPolicy
import org.opensearch.alerting.AlertingV2Utils.validateMonitorV1
import org.opensearch.alerting.core.lock.LockModel
import org.opensearch.alerting.core.lock.LockService
import org.opensearch.alerting.opensearchapi.addFilter
Expand Down Expand Up @@ -299,7 +300,13 @@ class TransportDeleteWorkflowAction @Inject constructor(
xContentRegistry,
LoggingDeprecationHandler.INSTANCE, hit.sourceAsString
).use { hitsParser ->
val monitor = ScheduledJob.parse(hitsParser, hit.id, hit.version) as Monitor
val scheduledJob = ScheduledJob.parse(hitsParser, hit.id, hit.version)

validateMonitorV1(scheduledJob)?.let {
throw OpenSearchException(it)
}

val monitor = scheduledJob as Monitor
deletableMonitors.add(monitor)
}
}
Expand Down Expand Up @@ -327,12 +334,17 @@ class TransportDeleteWorkflowAction @Inject constructor(
)
}

private fun parseWorkflow(getResponse: GetResponse): Workflow {
private fun parseWorkflow(getResponse: GetResponse): Workflow? {
val xcp = XContentHelper.createParser(
xContentRegistry, LoggingDeprecationHandler.INSTANCE,
getResponse.sourceAsBytesRef, XContentType.JSON
)
return ScheduledJob.parse(xcp, getResponse.id, getResponse.version) as Workflow
val scheduledJob = ScheduledJob.parse(xcp, getResponse.id, getResponse.version)
validateMonitorV1(scheduledJob)?.let {
actionListener.onFailure(AlertingException.wrap(it))
return null
}
return scheduledJob as Workflow
}

private suspend fun deleteWorkflow(deleteRequest: DeleteRequest): DeleteResponse {
Expand Down
Loading
Loading