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 @@ -35,6 +35,7 @@ import app.softnetwork.elastic.sql.query.{Asc, Criteria, Desc, SQLAggregation, S
import app.softnetwork.elastic.sql.schema.TableAlias
import app.softnetwork.elastic.sql.transform.{
AvgTransformAggregation,
BucketScriptTransformAggregation,
CardinalityTransformAggregation,
CountTransformAggregation,
Delay,
Expand Down Expand Up @@ -177,7 +178,10 @@ import org.elasticsearch.search.aggregations.metrics.{
TopHitsAggregationBuilder,
ValueCountAggregationBuilder
}
import org.elasticsearch.search.aggregations.pipeline.BucketSelectorPipelineAggregationBuilder
import org.elasticsearch.search.aggregations.pipeline.{
BucketScriptPipelineAggregationBuilder,
BucketSelectorPipelineAggregationBuilder
}
import org.elasticsearch.search.builder.{PointInTimeBuilder, SearchSourceBuilder}
import org.elasticsearch.search.slice.SliceBuilder
import org.elasticsearch.search.sort.{FieldSortBuilder, SortOrder}
Expand Down Expand Up @@ -2865,7 +2869,7 @@ trait RestHighLevelClientTransformApi extends TransformApi with RestHighLevelCli
)
}

private def convertToElasticTransformConfig(
private[client] def convertToElasticTransformConfig(
config: TransformConfig
)(implicit criteriaToNode: Criteria => JsonNode): ElasticTransformConfig = {
val builder = ElasticTransformConfig.builder()
Expand Down Expand Up @@ -2965,6 +2969,21 @@ trait RestHighLevelClientTransformApi extends TransformApi with RestHighLevelCli
topHitsBuilder.sort(sortField.name, sortOrder)
}
aggBuilder.addAggregator(topHitsBuilder)
case BucketScriptTransformAggregation(expression, bucketsPath, script, params) =>
// the model names these parameters but carries no value for them: refuse rather
// than send a script that reads them unbound
if (params.nonEmpty)
throw new UnsupportedOperationException(
s"Unsupported aggregation: the calculation $name ($expression) reads script " +
s"parameters the transform does not bind (${params.mkString(", ")})"
)
aggBuilder.addPipelineAggregator(
new BucketScriptPipelineAggregationBuilder(
name,
bucketsPath.asJava,
new org.elasticsearch.script.Script(script)
)
)
case _ =>
throw new UnsupportedOperationException(s"Unsupported aggregation: $agg")
}
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,133 @@
/*
* Copyright 2025 SOFTNETWORK
*
* 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 app.softnetwork.elastic.client

import app.softnetwork.elastic.client.rest.RestHighLevelClientApi
import app.softnetwork.elastic.sql.PainlessContextType
import app.softnetwork.elastic.sql.bridge._
import app.softnetwork.elastic.sql.transform.{
BucketScriptTransformAggregation,
Delay,
Frequency,
MaxTransformAggregation,
MinTransformAggregation,
TermsGroupBy,
TransformBucketSelectorConfig,
TransformConfig,
TransformDest,
TransformPivot,
TransformSource
}
import com.fasterxml.jackson.databind.{JsonNode, ObjectMapper}
import com.typesafe.config.{Config, ConfigFactory}
import org.elasticsearch.client.transform.PutTransformRequest
import org.elasticsearch.common.Strings
import org.scalatest.matchers.should.Matchers
import org.scalatest.wordspec.AnyWordSpec

import scala.collection.immutable.ListMap

/** The ES 7 REST client builds a transform pivot with the client library's builders, one match arm
* per `TransformAggregation` (ES 8 / ES 9 send the pivot as JSON). A materialized view's per-group
* calculation (`MAX(a) - MIN(b) AS d`, a `BucketScriptTransformAggregation`) had no arm: the
* request failed client-side, so a view with a SELECT calculation never deployed on ES 7.
*
* No Docker: `ElasticClientCompanion` builds the underlying client lazily and `apply()` is never
* called, so nothing here touches the network.
*/
class RestHighLevelClientTransformPivotSpec extends AnyWordSpec with Matchers {

private val client: RestHighLevelClientApi = new RestHighLevelClientApi {
override def config: Config = ConfigFactory.load()
}

implicit private val timestamp: Long = 1767139200000L // 2025-12-31T00:00:00Z

implicit private val context: PainlessContextType = PainlessContextType.Transform

private val mapper = new ObjectMapper()

private val calculation = BucketScriptTransformAggregation(
expression = "MAX(a) - MIN(b)",
bucketsPath = ListMap("max_a" -> "max_a", "min_b" -> "min_b"),
script = "params.max_a - params.min_b",
params = Nil
)

private def view(d: BucketScriptTransformAggregation): TransformConfig =
TransformConfig(
id = "v_d",
viewName = "v",
source = TransformSource(Seq("src"), None),
dest = TransformDest("v_d"),
pivot = Some(
TransformPivot(
groupBy = Map("g" -> TermsGroupBy("g")),
aggregations = ListMap(
"max_a" -> MaxTransformAggregation("a"),
"min_b" -> MinTransformAggregation("b"),
"d" -> d
),
bucketSelector = Some(
TransformBucketSelectorConfig(bucketsPath = Map("d" -> "d"), script = "params.d > -5")
)
)
),
delay = Delay.Default,
frequency = Frequency.Default
)

/** The body the client sends: `PutTransformRequest` renders the converted config. */
private def requestBody(config: TransformConfig): JsonNode =
mapper.readTree(
Strings.toString(new PutTransformRequest(client.convertToElasticTransformConfig(config)))
)

"the ES 7 transform request" should {

"carry the calculation as a bucket_script under its name" in {
val aggregations = requestBody(view(calculation)).path("pivot").path("aggregations")
val bucketScript = aggregations.path("d").path("bucket_script")

bucketScript.isObject shouldBe true
bucketScript.path("buckets_path").size() shouldBe 2
bucketScript.path("buckets_path").path("max_a").asText() shouldBe "max_a"
bucketScript.path("buckets_path").path("min_b").asText() shouldBe "min_b"
bucketScript.path("script").path("source").asText() shouldBe "params.max_a - params.min_b"
bucketScript.path("script").has("params") shouldBe false

// the operands it reads, and the HAVING that reads it by its name
aggregations.path("max_a").path("max").path("field").asText() shouldBe "a"
aggregations.path("min_b").path("min").path("field").asText() shouldBe "b"
aggregations
.path("having_filter")
.path("bucket_selector")
.path("buckets_path")
.path("d")
.asText() shouldBe "d"
}

"refuse a calculation whose script reads a parameter the transform does not bind" in {
val unbound = view(calculation.copy(params = Seq("__now__")))
val ex = intercept[UnsupportedOperationException](
client.convertToElasticTransformConfig(unbound)
)
ex.getMessage should include("calculation d (MAX(a) - MIN(b))")
ex.getMessage should include("__now__")
}
}
}
Loading