diff --git a/.github/workflows/maven.yml b/.github/workflows/maven.yml index d85367d..372737f 100644 --- a/.github/workflows/maven.yml +++ b/.github/workflows/maven.yml @@ -109,4 +109,4 @@ jobs: uses: MobilityDB/MEOS-API/.github/actions/check-test-outcome@master with: log: ${{ runner.temp }}/build.log - min-tests: "19" + min-tests: "28" diff --git a/GENERATION.md b/GENERATION.md index 24127ff..2722a9b 100644 --- a/GENERATION.md +++ b/GENERATION.md @@ -49,6 +49,25 @@ MobilityDB @ master CI stages the master-derived catalog to `tools/meos-idl.json` before the build; the catalog is gitignored, not committed. +## The Flink SQL surface + +`codegen_jvm.py --engine flink-sql` emits the MobilityDB SQL surface for Flink SQL from the +same catalog and jar: one Flink function per SQL name the catalog states (a signature's +`sqlName`, else the function's `@sqlfn`), each overload an `eval` method, and one RAW type per +MEOS value type the signatures use. A MEOS value crosses Flink as the serialized form its +catalog codec writes, hex WKB where the catalog states a WKB decoder and an `asHexWKB` +encoder and text otherwise, never as a native pointer, so Flink can copy, checkpoint and group +it. A SQL array argument stands for the C array the catalog pairs with a count +(`shape.inputArrays`), and the eval passes the length of the Flink array as that count. A +signature whose types the surface cannot carry (`Datum`, aggregate state, an array element with +no codec) is skipped and counted in the generator's report. + +`org.mobilitydb.flink.sql.MobilityFlinkSql.registerAll(tEnv)` registers every function as a +catalog function under its MobilityDB SQL name. A Flink built-in of the same name (`lower`, +`round`, `abs`, …) resolves before a catalog function, and a name Flink's parser reserves +(`overlaps`, `contains`, `union`, …) is written quoted, `` `overlaps`(a, b) ``. Maven +`generate-sources` runs the engine into `target/generated-sql`, beside the facades. + ## Generate-then-retire — the green-CI version is the probe Hand-written facades/glue are replaced by the generated forwarders **family by family, diff --git a/benchmark/src/main/java/org/mobilitydb/flink/meos/wirings/MeosCrossStreamJoin.java b/benchmark/src/main/java/org/mobilitydb/flink/meos/wirings/MeosCrossStreamJoin.java index edcd5f4..dc1e4cc 100644 --- a/benchmark/src/main/java/org/mobilitydb/flink/meos/wirings/MeosCrossStreamJoin.java +++ b/benchmark/src/main/java/org/mobilitydb/flink/meos/wirings/MeosCrossStreamJoin.java @@ -51,7 +51,7 @@ * *

Typical usage — per-vehicle-pair "did they come within * 100m of each other in the last 5 minutes?" via - * {@code MeosOpsTGeo.edwithin_tgeo_tgeo} (tier = {@code cross-stream}): + * {@code MeosOpsFreeGeo.edwithin_tgeo_tgeo} (tier = {@code cross-stream}): * *

{@code
  * KeyedStream a = streamA.keyBy(VehiclePosition::regionId);
@@ -64,7 +64,7 @@
  *             (left, right, ctx) -> {
  *                 Pointer leftT  = left.toTGeoPointer();
  *                 Pointer rightT = right.toTGeoPointer();
- *                 if (MeosOpsTGeo.edwithin_tgeo_tgeo(leftT, rightT, 100.0) != 0) {
+ *                 if (MeosOpsFreeGeo.edwithin_tgeo_tgeo(leftT, rightT, 100.0) != 0) {
  *                     return new MeetingEvent(left.id(), right.id(), ctx.getLeftTimestamp());
  *                 }
  *                 return null;  // no output for non-matches
diff --git a/benchmark/src/main/java/org/mobilitydb/flink/meos/wirings/MeosStatelessFilter.java b/benchmark/src/main/java/org/mobilitydb/flink/meos/wirings/MeosStatelessFilter.java
index 8e3329c..0685443 100644
--- a/benchmark/src/main/java/org/mobilitydb/flink/meos/wirings/MeosStatelessFilter.java
+++ b/benchmark/src/main/java/org/mobilitydb/flink/meos/wirings/MeosStatelessFilter.java
@@ -44,14 +44,14 @@
  * independently.
  *
  * 

Typical usage: scalar-predicate filter against the - * generated {@code MeosOpsTBox.overlaps_tbox_tbox} (tier = + * generated {@code MeosOpsFreeCore.overlaps_tbox_tbox} (tier = * {@code stateless}): * *

{@code
  * DataStream in = ...;
  * DataStream overlapping = in.filter(
  *     new MeosStatelessFilter<>(
- *         pair -> MeosOpsTBox.overlaps_tbox_tbox(pair.a, pair.b)));
+ *         pair -> MeosOpsFreeCore.overlaps_tbox_tbox(pair.a, pair.b)));
  * }
* *

For int-coded predicates (JMEOS returns {@code int} for some MEOS diff --git a/benchmark/src/main/java/org/mobilitydb/flink/meos/wirings/MeosStatelessMap.java b/benchmark/src/main/java/org/mobilitydb/flink/meos/wirings/MeosStatelessMap.java index 8ef1709..dafd21c 100644 --- a/benchmark/src/main/java/org/mobilitydb/flink/meos/wirings/MeosStatelessMap.java +++ b/benchmark/src/main/java/org/mobilitydb/flink/meos/wirings/MeosStatelessMap.java @@ -52,14 +52,14 @@ * *

Typical usage: register a stateless MEOS predicate / arithmetic * call as a per-event map step in a DataStream pipeline. Example with - * the generated {@code MeosOpsTBox.overlaps_tbox_tbox} (tier = + * the generated {@code MeosOpsFreeCore.overlaps_tbox_tbox} (tier = * {@code stateless}, per the codegen manifest): * *

{@code
  * DataStream in = ...;            // (tboxA, tboxB)
  * DataStream overlap = in.map(
  *     new MeosStatelessMap<>(
- *         pair -> MeosOpsTBox.overlaps_tbox_tbox(pair.a, pair.b)));
+ *         pair -> MeosOpsFreeCore.overlaps_tbox_tbox(pair.a, pair.b)));
  * }
* *

Tier coverage: as of the codegen state on the parent PR, diff --git a/benchmark/src/main/java/org/mobilitydb/flink/meos/wirings/MeosWindowedAggregate.java b/benchmark/src/main/java/org/mobilitydb/flink/meos/wirings/MeosWindowedAggregate.java index b4d26f5..6921c9c 100644 --- a/benchmark/src/main/java/org/mobilitydb/flink/meos/wirings/MeosWindowedAggregate.java +++ b/benchmark/src/main/java/org/mobilitydb/flink/meos/wirings/MeosWindowedAggregate.java @@ -38,8 +38,8 @@ * *

The {@code windowed} tier is "output cardinality changes; needs a * window". The canonical examples are - * {@code temporal_length(tgeo)} (one length per trajectory window), - * {@code temporal_twavg(tnumber)} (one time-weighted average per + * {@code tpoint_length(tpoint)} (one length per trajectory window), + * {@code tnumber_twavg(tnumber)} (one time-weighted average per * window), and the per-class {@code _trajectory} / {@code _time} / * {@code _timespan} accessors that reduce a full sequence to a single * derived value. @@ -61,7 +61,7 @@ * per-window. * *

Typical usage — per-vehicle per-tumbling-window - * trajectory length via {@code MeosOpsTemporal.temporal_length} (tier + * trajectory length via {@code MeosOpsTPoint.tpoint_length} (tier * = {@code windowed}): * *

{@code
@@ -73,7 +73,7 @@
  *     .process(new MeosWindowedAggregate(
  *         (window, events, ctx) -> {
  *             Pointer trajectory = buildTrajectoryFromPoints(events);  // adopter helper
- *             double length = MeosOpsTemporal.temporal_length(trajectory);
+ *             double length = MeosOpsTPoint.tpoint_length(trajectory);
  *             return new VehicleLength(ctx.getCurrentKey(), window.getStart(), length);
  *         }));
  * }
diff --git a/benchmark/src/main/java/org/mobilitydb/flink/meos/wirings/README.md b/benchmark/src/main/java/org/mobilitydb/flink/meos/wirings/README.md index b79ec19..d3c0794 100644 --- a/benchmark/src/main/java/org/mobilitydb/flink/meos/wirings/README.md +++ b/benchmark/src/main/java/org/mobilitydb/flink/meos/wirings/README.md @@ -38,11 +38,11 @@ The pattern is the same across all four tiers: ```java // 1. Pick the generated MeosOps method // (Javadoc tier marker tells you which wiring to use) -boolean overlap = MeosOpsTBox.overlaps_tbox_tbox(boxA, boxB); // tier = stateless +boolean overlap = MeosOpsFreeCore.overlaps_tbox_tbox(boxA, boxB); // tier = stateless // 2. Wrap with the matching wiring MeosStatelessFilter filter = MeosStatelessFilter.fromIntPredicate( - pair -> MeosOpsTBox.overlaps_tbox_tbox(pair.a, pair.b)); + pair -> MeosOpsFreeCore.overlaps_tbox_tbox(pair.a, pair.b)); // 3. Apply to the DataStream DataStream overlapping = stream.filter(filter); @@ -61,9 +61,9 @@ through a 3-stage DataStream pipeline using two of the generated facades wired through `MeosStatelessMap` + `MeosStatelessFilter`: 1. Parse a stream of TBox WKT strings via - `MeosOpsFreeCore.tbox_in` (io-meta, no state). + `MeosOpsTBox.tbox_in` (io-meta, no state). 2. Filter to those overlapping a fixed query box via - `MeosOpsTBox.overlaps_tbox_tbox` (stateless predicate). + `MeosOpsFreeCore.overlaps_tbox_tbox` (stateless predicate). 3. Serialize each survivor to hex-WKB via `MeosOpsTBox.tbox_as_hexwkb` (io-meta, no state). diff --git a/benchmark/src/main/java/org/mobilitydb/flink/meos/wirings/demo/MeosWiringsDemoJob.java b/benchmark/src/main/java/org/mobilitydb/flink/meos/wirings/demo/MeosWiringsDemoJob.java index bb70f79..9ee612f 100644 --- a/benchmark/src/main/java/org/mobilitydb/flink/meos/wirings/demo/MeosWiringsDemoJob.java +++ b/benchmark/src/main/java/org/mobilitydb/flink/meos/wirings/demo/MeosWiringsDemoJob.java @@ -49,7 +49,7 @@ *
  • Parses each into a JMEOS {@code Pointer} via * {@code MeosOpsTBox.tbox_in} (tier = {@code io-meta}).
  • *
  • Filters to those that overlap with a fixed query TBox via - * {@code MeosOpsTBox.overlaps_tbox_tbox} wrapped as a + * {@code MeosOpsFreeCore.overlaps_tbox_tbox} wrapped as a * {@link MeosStatelessFilter} (tier = {@code stateless}).
  • *
  • Maps each surviving TBox to its serialized WKB hex via * {@code MeosOpsTBox.tbox_as_hexwkb} wrapped as a diff --git a/benchmark/src/test/java/org/mobilitydb/flink/meos/MeosCbufferSmokeTest.java b/benchmark/src/test/java/org/mobilitydb/flink/meos/MeosCbufferSmokeTest.java index 4bad5d4..83f71c1 100644 --- a/benchmark/src/test/java/org/mobilitydb/flink/meos/MeosCbufferSmokeTest.java +++ b/benchmark/src/test/java/org/mobilitydb/flink/meos/MeosCbufferSmokeTest.java @@ -59,9 +59,9 @@ static void finalizeMeos() { @Test void cbuffer() { - Pointer cb = MeosOpsFreeCbuffer.cbuffer_make(MeosOpsFreeGeo.geom_in("POINT(1 1)", 0), 0.5); + Pointer cb = MeosOpsCbuffer.cbuffer_make(MeosOpsGeometry.geom_in("POINT(1 1)", 0), 0.5); assertNotNull(cb); - assertEquals(0.5, MeosOpsFreeCbuffer.cbuffer_radius(cb), 1e-9); - assertNotNull(MeosOpsFreeCbuffer.cbuffer_out(cb, 6)); + assertEquals(0.5, MeosOpsCbuffer.cbuffer_radius(cb), 1e-9); + assertNotNull(MeosOpsCbuffer.cbuffer_out(cb, 6)); } } diff --git a/benchmark/src/test/java/org/mobilitydb/flink/meos/MeosFacadeSmokeTest.java b/benchmark/src/test/java/org/mobilitydb/flink/meos/MeosFacadeSmokeTest.java index 539dc1c..6720695 100644 --- a/benchmark/src/test/java/org/mobilitydb/flink/meos/MeosFacadeSmokeTest.java +++ b/benchmark/src/test/java/org/mobilitydb/flink/meos/MeosFacadeSmokeTest.java @@ -86,8 +86,8 @@ void geoStbox() { @Test void geoGeometry() { - Pointer geom = MeosOpsFreeGeo.geom_in("POINT(1 1)", 0); + Pointer geom = MeosOpsGeometry.geom_in("POINT(1 1)", 0); assertNotNull(geom); - assertTrue(MeosOpsFreeGeo.geo_as_text(geom, 6).toUpperCase().contains("POINT")); + assertTrue(MeosOpsGeo.geo_as_text(geom, 6).toUpperCase().contains("POINT")); } } diff --git a/benchmark/src/test/java/org/mobilitydb/flink/meos/MeosNpointSmokeTest.java b/benchmark/src/test/java/org/mobilitydb/flink/meos/MeosNpointSmokeTest.java index 0cfaf41..1be2934 100644 --- a/benchmark/src/test/java/org/mobilitydb/flink/meos/MeosNpointSmokeTest.java +++ b/benchmark/src/test/java/org/mobilitydb/flink/meos/MeosNpointSmokeTest.java @@ -58,9 +58,9 @@ static void finalizeMeos() { @Test void npoint() { - Pointer np = MeosOpsFreeNpoint.npoint_make(1, 0.5); + Pointer np = MeosOpsNpoint.npoint_make(1, 0.5); assertNotNull(np); - assertEquals(1, MeosOpsFreeNpoint.npoint_route(np)); - assertEquals(0.5, MeosOpsFreeNpoint.npoint_position(np), 1e-9); + assertEquals(1, MeosOpsNpoint.npoint_route(np)); + assertEquals(0.5, MeosOpsNpoint.npoint_position(np), 1e-9); } } diff --git a/benchmark/src/test/java/org/mobilitydb/flink/meos/MeosPoseSmokeTest.java b/benchmark/src/test/java/org/mobilitydb/flink/meos/MeosPoseSmokeTest.java index d395a55..219a36f 100644 --- a/benchmark/src/test/java/org/mobilitydb/flink/meos/MeosPoseSmokeTest.java +++ b/benchmark/src/test/java/org/mobilitydb/flink/meos/MeosPoseSmokeTest.java @@ -59,9 +59,9 @@ static void finalizeMeos() { @Test void pose() { - Pointer pose = MeosOpsFreePose.pose_in("Pose(Point(1 1), 0.5)"); + Pointer pose = MeosOpsPose.pose_in("Pose(Point(1 1), 0.5)"); assertNotNull(pose); - assertNotNull(MeosOpsFreePose.pose_out(pose, 6)); - assertEquals(0.5, MeosOpsFreePose.pose_yaw(pose), 1e-9); + assertNotNull(MeosOpsPose.pose_out(pose, 6)); + assertEquals(0.5, MeosOpsPose.pose_yaw(pose), 1e-9); } } diff --git a/binding/pom.xml b/binding/pom.xml index b15db0f..dd0bade 100644 --- a/binding/pom.xml +++ b/binding/pom.xml @@ -12,8 +12,9 @@ mobility-flink-binding - @@ -32,6 +33,17 @@ flink-clients ${flink.version} + + org.apache.flink + flink-table-api-java + ${flink.version} + + + org.apache.flink + flink-table-planner_2.12 + ${flink.version} + test + org.jmeos @@ -101,6 +113,26 @@ + + + generate-sql-surface + generate-sources + exec + + python3 + + ${project.basedir}/../tools/codegen_jvm.py + --engine + flink-sql + --catalog + ${project.basedir}/../tools/meos-idl.json + --jar + ${settings.localRepository}/org/jmeos/meos/${jmeos.version}/meos-${jmeos.version}.jar + --out + ${project.build.directory}/generated-sql + + + @@ -115,6 +147,7 @@ ${project.build.directory}/generated-facades/src/main/java + ${project.build.directory}/generated-sql/src/main/java diff --git a/binding/src/main/java/org/mobilitydb/flink/meos/wirings/MeosCrossStreamJoin.java b/binding/src/main/java/org/mobilitydb/flink/meos/wirings/MeosCrossStreamJoin.java index edcd5f4..dc1e4cc 100644 --- a/binding/src/main/java/org/mobilitydb/flink/meos/wirings/MeosCrossStreamJoin.java +++ b/binding/src/main/java/org/mobilitydb/flink/meos/wirings/MeosCrossStreamJoin.java @@ -51,7 +51,7 @@ * *

    Typical usage — per-vehicle-pair "did they come within * 100m of each other in the last 5 minutes?" via - * {@code MeosOpsTGeo.edwithin_tgeo_tgeo} (tier = {@code cross-stream}): + * {@code MeosOpsFreeGeo.edwithin_tgeo_tgeo} (tier = {@code cross-stream}): * *

    {@code
      * KeyedStream a = streamA.keyBy(VehiclePosition::regionId);
    @@ -64,7 +64,7 @@
      *             (left, right, ctx) -> {
      *                 Pointer leftT  = left.toTGeoPointer();
      *                 Pointer rightT = right.toTGeoPointer();
    - *                 if (MeosOpsTGeo.edwithin_tgeo_tgeo(leftT, rightT, 100.0) != 0) {
    + *                 if (MeosOpsFreeGeo.edwithin_tgeo_tgeo(leftT, rightT, 100.0) != 0) {
      *                     return new MeetingEvent(left.id(), right.id(), ctx.getLeftTimestamp());
      *                 }
      *                 return null;  // no output for non-matches
    diff --git a/binding/src/main/java/org/mobilitydb/flink/meos/wirings/MeosStatelessFilter.java b/binding/src/main/java/org/mobilitydb/flink/meos/wirings/MeosStatelessFilter.java
    index 8e3329c..0685443 100644
    --- a/binding/src/main/java/org/mobilitydb/flink/meos/wirings/MeosStatelessFilter.java
    +++ b/binding/src/main/java/org/mobilitydb/flink/meos/wirings/MeosStatelessFilter.java
    @@ -44,14 +44,14 @@
      * independently.
      *
      * 

    Typical usage: scalar-predicate filter against the - * generated {@code MeosOpsTBox.overlaps_tbox_tbox} (tier = + * generated {@code MeosOpsFreeCore.overlaps_tbox_tbox} (tier = * {@code stateless}): * *

    {@code
      * DataStream in = ...;
      * DataStream overlapping = in.filter(
      *     new MeosStatelessFilter<>(
    - *         pair -> MeosOpsTBox.overlaps_tbox_tbox(pair.a, pair.b)));
    + *         pair -> MeosOpsFreeCore.overlaps_tbox_tbox(pair.a, pair.b)));
      * }
    * *

    For int-coded predicates (JMEOS returns {@code int} for some MEOS diff --git a/binding/src/main/java/org/mobilitydb/flink/meos/wirings/MeosStatelessMap.java b/binding/src/main/java/org/mobilitydb/flink/meos/wirings/MeosStatelessMap.java index 8ef1709..dafd21c 100644 --- a/binding/src/main/java/org/mobilitydb/flink/meos/wirings/MeosStatelessMap.java +++ b/binding/src/main/java/org/mobilitydb/flink/meos/wirings/MeosStatelessMap.java @@ -52,14 +52,14 @@ * *

    Typical usage: register a stateless MEOS predicate / arithmetic * call as a per-event map step in a DataStream pipeline. Example with - * the generated {@code MeosOpsTBox.overlaps_tbox_tbox} (tier = + * the generated {@code MeosOpsFreeCore.overlaps_tbox_tbox} (tier = * {@code stateless}, per the codegen manifest): * *

    {@code
      * DataStream in = ...;            // (tboxA, tboxB)
      * DataStream overlap = in.map(
      *     new MeosStatelessMap<>(
    - *         pair -> MeosOpsTBox.overlaps_tbox_tbox(pair.a, pair.b)));
    + *         pair -> MeosOpsFreeCore.overlaps_tbox_tbox(pair.a, pair.b)));
      * }
    * *

    Tier coverage: as of the codegen state on the parent PR, diff --git a/binding/src/main/java/org/mobilitydb/flink/meos/wirings/MeosWindowedAggregate.java b/binding/src/main/java/org/mobilitydb/flink/meos/wirings/MeosWindowedAggregate.java index b4d26f5..6921c9c 100644 --- a/binding/src/main/java/org/mobilitydb/flink/meos/wirings/MeosWindowedAggregate.java +++ b/binding/src/main/java/org/mobilitydb/flink/meos/wirings/MeosWindowedAggregate.java @@ -38,8 +38,8 @@ * *

    The {@code windowed} tier is "output cardinality changes; needs a * window". The canonical examples are - * {@code temporal_length(tgeo)} (one length per trajectory window), - * {@code temporal_twavg(tnumber)} (one time-weighted average per + * {@code tpoint_length(tpoint)} (one length per trajectory window), + * {@code tnumber_twavg(tnumber)} (one time-weighted average per * window), and the per-class {@code _trajectory} / {@code _time} / * {@code _timespan} accessors that reduce a full sequence to a single * derived value. @@ -61,7 +61,7 @@ * per-window. * *

    Typical usage — per-vehicle per-tumbling-window - * trajectory length via {@code MeosOpsTemporal.temporal_length} (tier + * trajectory length via {@code MeosOpsTPoint.tpoint_length} (tier * = {@code windowed}): * *

    {@code
    @@ -73,7 +73,7 @@
      *     .process(new MeosWindowedAggregate(
      *         (window, events, ctx) -> {
      *             Pointer trajectory = buildTrajectoryFromPoints(events);  // adopter helper
    - *             double length = MeosOpsTemporal.temporal_length(trajectory);
    + *             double length = MeosOpsTPoint.tpoint_length(trajectory);
      *             return new VehicleLength(ctx.getCurrentKey(), window.getStart(), length);
      *         }));
      * }
    diff --git a/binding/src/main/java/org/mobilitydb/flink/meos/wirings/README.md b/binding/src/main/java/org/mobilitydb/flink/meos/wirings/README.md index b79ec19..d3c0794 100644 --- a/binding/src/main/java/org/mobilitydb/flink/meos/wirings/README.md +++ b/binding/src/main/java/org/mobilitydb/flink/meos/wirings/README.md @@ -38,11 +38,11 @@ The pattern is the same across all four tiers: ```java // 1. Pick the generated MeosOps method // (Javadoc tier marker tells you which wiring to use) -boolean overlap = MeosOpsTBox.overlaps_tbox_tbox(boxA, boxB); // tier = stateless +boolean overlap = MeosOpsFreeCore.overlaps_tbox_tbox(boxA, boxB); // tier = stateless // 2. Wrap with the matching wiring MeosStatelessFilter filter = MeosStatelessFilter.fromIntPredicate( - pair -> MeosOpsTBox.overlaps_tbox_tbox(pair.a, pair.b)); + pair -> MeosOpsFreeCore.overlaps_tbox_tbox(pair.a, pair.b)); // 3. Apply to the DataStream DataStream overlapping = stream.filter(filter); @@ -61,9 +61,9 @@ through a 3-stage DataStream pipeline using two of the generated facades wired through `MeosStatelessMap` + `MeosStatelessFilter`: 1. Parse a stream of TBox WKT strings via - `MeosOpsFreeCore.tbox_in` (io-meta, no state). + `MeosOpsTBox.tbox_in` (io-meta, no state). 2. Filter to those overlapping a fixed query box via - `MeosOpsTBox.overlaps_tbox_tbox` (stateless predicate). + `MeosOpsFreeCore.overlaps_tbox_tbox` (stateless predicate). 3. Serialize each survivor to hex-WKB via `MeosOpsTBox.tbox_as_hexwkb` (io-meta, no state). diff --git a/binding/src/main/java/org/mobilitydb/flink/meos/wirings/demo/MeosWiringsDemoJob.java b/binding/src/main/java/org/mobilitydb/flink/meos/wirings/demo/MeosWiringsDemoJob.java index bb70f79..9ee612f 100644 --- a/binding/src/main/java/org/mobilitydb/flink/meos/wirings/demo/MeosWiringsDemoJob.java +++ b/binding/src/main/java/org/mobilitydb/flink/meos/wirings/demo/MeosWiringsDemoJob.java @@ -49,7 +49,7 @@ *
  • Parses each into a JMEOS {@code Pointer} via * {@code MeosOpsTBox.tbox_in} (tier = {@code io-meta}).
  • *
  • Filters to those that overlap with a fixed query TBox via - * {@code MeosOpsTBox.overlaps_tbox_tbox} wrapped as a + * {@code MeosOpsFreeCore.overlaps_tbox_tbox} wrapped as a * {@link MeosStatelessFilter} (tier = {@code stateless}).
  • *
  • Maps each surviving TBox to its serialized WKB hex via * {@code MeosOpsTBox.tbox_as_hexwkb} wrapped as a diff --git a/binding/src/test/java/org/mobilitydb/flink/meos/MeosCbufferSmokeTest.java b/binding/src/test/java/org/mobilitydb/flink/meos/MeosCbufferSmokeTest.java index 4bad5d4..83f71c1 100644 --- a/binding/src/test/java/org/mobilitydb/flink/meos/MeosCbufferSmokeTest.java +++ b/binding/src/test/java/org/mobilitydb/flink/meos/MeosCbufferSmokeTest.java @@ -59,9 +59,9 @@ static void finalizeMeos() { @Test void cbuffer() { - Pointer cb = MeosOpsFreeCbuffer.cbuffer_make(MeosOpsFreeGeo.geom_in("POINT(1 1)", 0), 0.5); + Pointer cb = MeosOpsCbuffer.cbuffer_make(MeosOpsGeometry.geom_in("POINT(1 1)", 0), 0.5); assertNotNull(cb); - assertEquals(0.5, MeosOpsFreeCbuffer.cbuffer_radius(cb), 1e-9); - assertNotNull(MeosOpsFreeCbuffer.cbuffer_out(cb, 6)); + assertEquals(0.5, MeosOpsCbuffer.cbuffer_radius(cb), 1e-9); + assertNotNull(MeosOpsCbuffer.cbuffer_out(cb, 6)); } } diff --git a/binding/src/test/java/org/mobilitydb/flink/meos/MeosFacadeSmokeTest.java b/binding/src/test/java/org/mobilitydb/flink/meos/MeosFacadeSmokeTest.java index 539dc1c..6720695 100644 --- a/binding/src/test/java/org/mobilitydb/flink/meos/MeosFacadeSmokeTest.java +++ b/binding/src/test/java/org/mobilitydb/flink/meos/MeosFacadeSmokeTest.java @@ -86,8 +86,8 @@ void geoStbox() { @Test void geoGeometry() { - Pointer geom = MeosOpsFreeGeo.geom_in("POINT(1 1)", 0); + Pointer geom = MeosOpsGeometry.geom_in("POINT(1 1)", 0); assertNotNull(geom); - assertTrue(MeosOpsFreeGeo.geo_as_text(geom, 6).toUpperCase().contains("POINT")); + assertTrue(MeosOpsGeo.geo_as_text(geom, 6).toUpperCase().contains("POINT")); } } diff --git a/binding/src/test/java/org/mobilitydb/flink/meos/MeosNpointSmokeTest.java b/binding/src/test/java/org/mobilitydb/flink/meos/MeosNpointSmokeTest.java index 0cfaf41..1be2934 100644 --- a/binding/src/test/java/org/mobilitydb/flink/meos/MeosNpointSmokeTest.java +++ b/binding/src/test/java/org/mobilitydb/flink/meos/MeosNpointSmokeTest.java @@ -58,9 +58,9 @@ static void finalizeMeos() { @Test void npoint() { - Pointer np = MeosOpsFreeNpoint.npoint_make(1, 0.5); + Pointer np = MeosOpsNpoint.npoint_make(1, 0.5); assertNotNull(np); - assertEquals(1, MeosOpsFreeNpoint.npoint_route(np)); - assertEquals(0.5, MeosOpsFreeNpoint.npoint_position(np), 1e-9); + assertEquals(1, MeosOpsNpoint.npoint_route(np)); + assertEquals(0.5, MeosOpsNpoint.npoint_position(np), 1e-9); } } diff --git a/binding/src/test/java/org/mobilitydb/flink/meos/MeosPoseSmokeTest.java b/binding/src/test/java/org/mobilitydb/flink/meos/MeosPoseSmokeTest.java index d395a55..219a36f 100644 --- a/binding/src/test/java/org/mobilitydb/flink/meos/MeosPoseSmokeTest.java +++ b/binding/src/test/java/org/mobilitydb/flink/meos/MeosPoseSmokeTest.java @@ -59,9 +59,9 @@ static void finalizeMeos() { @Test void pose() { - Pointer pose = MeosOpsFreePose.pose_in("Pose(Point(1 1), 0.5)"); + Pointer pose = MeosOpsPose.pose_in("Pose(Point(1 1), 0.5)"); assertNotNull(pose); - assertNotNull(MeosOpsFreePose.pose_out(pose, 6)); - assertEquals(0.5, MeosOpsFreePose.pose_yaw(pose), 1e-9); + assertNotNull(MeosOpsPose.pose_out(pose, 6)); + assertEquals(0.5, MeosOpsPose.pose_yaw(pose), 1e-9); } } diff --git a/binding/src/test/java/org/mobilitydb/flink/sql/GeneratedSqlSurfaceTest.java b/binding/src/test/java/org/mobilitydb/flink/sql/GeneratedSqlSurfaceTest.java new file mode 100644 index 0000000..87092ed --- /dev/null +++ b/binding/src/test/java/org/mobilitydb/flink/sql/GeneratedSqlSurfaceTest.java @@ -0,0 +1,168 @@ +/***************************************************************************** + * + * This MobilityDB code is provided under The PostgreSQL License. + * Copyright (c) 2020-2026, Université libre de Bruxelles and MobilityDB + * contributors + * + * Permission to use, copy, modify, and distribute this software and its + * documentation for any purpose, without fee, and without a written + * agreement is hereby granted, provided that the above copyright notice and + * this paragraph and the following two paragraphs appear in all copies. + * + * IN NO EVENT SHALL UNIVERSITE LIBRE DE BRUXELLES BE LIABLE TO ANY PARTY FOR + * DIRECT, INDIRECT, SPECIAL, INCIDENTAL, OR CONSEQUENTIAL DAMAGES, INCLUDING + * LOST PROFITS, ARISING OUT OF THE USE OF THIS SOFTWARE AND ITS DOCUMENTATION, + * EVEN IF UNIVERSITE LIBRE DE BRUXELLES HAS BEEN ADVISED OF THE POSSIBILITY + * OF SUCH DAMAGE. + * + * UNIVERSITE LIBRE DE BRUXELLES SPECIFICALLY DISCLAIMS ANY WARRANTIES, + * INCLUDING, BUT NOT LIMITED TO, THE IMPLIED WARRANTIES OF MERCHANTABILITY + * AND FITNESS FOR A PARTICULAR PURPOSE. THE SOFTWARE PROVIDED HEREUNDER IS ON + * AN "AS IS" BASIS, AND UNIVERSITE LIBRE DE BRUXELLES HAS NO OBLIGATIONS TO + * PROVIDE MAINTENANCE, SUPPORT, UPDATES, ENHANCEMENTS, OR MODIFICATIONS. + * + *****************************************************************************/ + +package org.mobilitydb.flink.sql; + +import functions.GeneratedFunctions; +import java.time.Duration; +import java.time.Instant; +import java.util.ArrayList; +import java.util.List; +import org.apache.flink.table.api.EnvironmentSettings; +import org.apache.flink.table.api.TableEnvironment; +import org.apache.flink.types.Row; +import org.apache.flink.util.CloseableIterator; +import org.junit.jupiter.api.AfterAll; +import org.junit.jupiter.api.BeforeAll; +import org.junit.jupiter.api.Test; +import org.mobilitydb.flink.sql.types.TFloat; +import org.mobilitydb.flink.sql.types.TGeomPoint; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * Runs the generated Flink SQL surface end to end against libmeos: each query plans against + * the generated functions, executes on a Flink mini-cluster and reads back its answer. The + * temporal value is a two-instant float sequence spanning two days in UTC, carried as the + * hex WKB its generated type writes, against a libmeos on the load path. + */ +class GeneratedSqlSurfaceTest { + + private static TableEnvironment tEnv; + private static String hexLiteral; + private static String tfloat; + + @BeforeAll + static void init() { + // No-op error handler so a parse error returns rather than terminating the JVM. + GeneratedFunctions.meos_initialize_error_handler((level, code, message) -> { }); + GeneratedFunctions.meos_initialize(); + String hex = TFloat.encode(GeneratedFunctions.tfloat_in( + "[1@2020-01-01 00:00:00+00, 3@2020-01-03 00:00:00+00]")).form; + hexLiteral = "'" + hex + "'"; + tfloat = "tfloatFromHexWKB(" + hexLiteral + ")"; + tEnv = TableEnvironment.create(EnvironmentSettings.inStreamingMode()); + MobilityFlinkSql.registerAll(tEnv); + } + + @AfterAll + static void finalizeMeos() { + GeneratedFunctions.meos_finalize(); + } + + /** Every row the query emits; a grouped query's last row is its final answer. */ + private static List rows(String sql) throws Exception { + List out = new ArrayList<>(); + try (CloseableIterator it = tEnv.executeSql(sql).collect()) { + it.forEachRemaining(out::add); + } + return out; + } + + private static Object scalar(String sql) throws Exception { + return rows(sql).get(0).getField(0); + } + + @Test + void accessorsReadTheValue() throws Exception { + assertEquals(2, scalar("SELECT numInstants(" + tfloat + ")")); + assertEquals(1.0, scalar("SELECT startValue(" + tfloat + ")")); + } + + @Test + void overloadsResolveByArgumentType() throws Exception { + assertEquals(3.0, scalar("SELECT startValue(tAdd(" + tfloat + ", 2.0))")); + assertEquals(2, scalar("SELECT numInstants(tAdd(" + tfloat + ", " + tfloat + "))")); + } + + @Test + void timeValuesCrossAsFlinkTypes() throws Exception { + assertEquals(Duration.ofHours(48), scalar("SELECT duration(" + tfloat + ")")); + assertEquals(Instant.parse("2020-01-01T00:00:00Z"), + scalar("SELECT startTimestamp(" + tfloat + ")")); + assertEquals(Instant.parse("2020-01-02T00:00:00Z"), + scalar("SELECT startTimestamp(shiftTime(" + tfloat + ", INTERVAL '1' DAY))")); + assertEquals(1, scalar("SELECT numInstants(atTime(" + tfloat + + ", TO_TIMESTAMP_LTZ(1577836800000, 3)))")); + } + + @Test + void constructorsTakeFlinkScalars() throws Exception { + String inst = "tfloat(CAST(1.5 AS DOUBLE), TO_TIMESTAMP_LTZ(1577836800000, 3))"; + assertEquals(1, scalar("SELECT numInstants(" + inst + ")")); + assertEquals(1.5, scalar("SELECT startValue(" + inst + ")")); + assertEquals(Instant.parse("2020-01-01T00:00:00Z"), scalar("SELECT startTimestamp(" + inst + ")")); + } + + @Test + void textConstructorsAndOutput() throws Exception { + String box = (String) scalar( + "SELECT asText(tbox_in('TBOXFLOAT XT([1, 2],[2020-01-01, 2020-01-02])'))"); + assertTrue(box.startsWith("TBOXFLOAT XT([1, 2]"), box); + } + + @Test + void aReservedNameIsCalledQuoted() throws Exception { + assertEquals(true, + scalar("SELECT `overlaps`(floatspan_in('[1, 3]'), floatspan_in('[2, 4]'))")); + } + + @Test + void aBuiltinNameResolvesToTheBuiltinUnlessQualified() throws Exception { + assertThrows(Exception.class, () -> scalar("SELECT lower(floatspan_in('[1, 3]'))")); + assertEquals(1.0, scalar("SELECT default_database.`lower`(floatspan_in('[1, 3]'))")); + } + + private static String tgeompoint(String text) { + return "tgeompointFromHexEWKB('" + + TGeomPoint.encode(GeneratedFunctions.tgeompoint_in(text)).form + "')"; + } + + @Test + void arraysCrossAsFlinkArrays() throws Exception { + assertEquals("{1, 2, 3}", scalar("SELECT intset_out(`set`(ARRAY[3, 1, 2]))")); + assertEquals(Instant.parse("2020-01-01T00:00:00Z"), scalar("SELECT startValue(`set`(ARRAY[" + + "TO_TIMESTAMP_LTZ(1577923200000, 3), TO_TIMESTAMP_LTZ(1577836800000, 3)]))")); + String origin = tgeompoint("[Point(0 0)@2020-01-01, Point(0 0)@2020-01-02]"); + String far = tgeompoint("[Point(3 4)@2020-01-01, Point(3 4)@2020-01-02]"); + String near = tgeompoint("[Point(0 1)@2020-01-01, Point(0 1)@2020-01-02]"); + assertEquals(5.0, scalar("SELECT minDistance(ARRAY[" + origin + "], ARRAY[" + far + "])")); + assertEquals(1.0, scalar("SELECT minDistance(ARRAY[" + origin + "], ARRAY[" + far + ", " + + near + "])")); + assertThrows(Exception.class, + () -> scalar("SELECT intset_out(`set`(ARRAY[1, CAST(NULL AS INT)]))")); + } + + @Test + void valuesGroupAcrossAShuffle() throws Exception { + List r = rows("SELECT startValue(v), COUNT(*) FROM (SELECT tfloatFromHexWKB(h) AS v" + + " FROM (VALUES (" + hexLiteral + "), (" + hexLiteral + ")) AS s(h)) GROUP BY v"); + Row last = r.get(r.size() - 1); + assertEquals(1.0, last.getField(0)); + assertEquals(2L, last.getField(1)); + } +}