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
2 changes: 1 addition & 1 deletion .github/workflows/docker.yml
Original file line number Diff line number Diff line change
Expand Up @@ -81,7 +81,7 @@ jobs:
-DgroupId=org.jmeos -DartifactId=meos -Dversion=1.0 -Dpackaging=jar

- name: Build the application jar the image copies
run: mvn -B -Dmeos.lib.dir=/usr/local/lib -Dmeos.enabled=true -pl benchmark -am package
run: mvn -B -Dmeos.lib.dir=/usr/local/lib -pl benchmark -am package

- name: Build the image
run: docker build --progress=plain -f benchmark/Dockerfile -t flink-benchmark-ci:${{ github.sha }} benchmark/
Expand Down
2 changes: 1 addition & 1 deletion .github/workflows/maven.yml
Original file line number Diff line number Diff line change
Expand Up @@ -97,7 +97,7 @@ jobs:
# the smoke tests exercise the facades against the freshly built libmeos from /usr/local/lib.
run: |
set -o pipefail
mvn -B -Dmeos.lib.dir=/usr/local/lib -Dmeos.enabled=true clean test | tee "$RUNNER_TEMP/build.log"
mvn -B -Dmeos.lib.dir=/usr/local/lib clean test | tee "$RUNNER_TEMP/build.log"

# The build's own summary states two things this job's conclusion does
# not: a test reported SKIPPED asserts nothing while the job still reads
Expand Down
3 changes: 1 addition & 2 deletions GENERATION.md
Original file line number Diff line number Diff line change
Expand Up @@ -89,8 +89,7 @@ tools/refresh-from-master.sh --mdb ~/src/MobilityDB # refresh against a local
This repo's last leg is `tools/refresh.conf`: the `binding` Maven module, the `org.jmeos:meos:1.0`
jar coordinates, and the `mvn … clean test` command. `generate-sources` runs `tools/codegen_jvm.py
--engine flink` over the staged `tools/meos-idl.json` and the installed jar, so the facades are
regenerated by the build itself; `meos.lib.dir` is where the smoke tests find `libmeos.so` and
`meos.enabled` turns them on.
regenerated by the build itself; `meos.lib.dir` is where the smoke tests find `libmeos.so`.

The chain's per-leg commands, which `refresh-jvm-chain.sh` composes, are in `MEOS-API/GENERATION.md`
(`provision-meos.sh`) and `JMEOS/GENERATION.md` (`regen-from-catalog.sh`).
19 changes: 2 additions & 17 deletions benchmark/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -251,9 +251,6 @@
<excludes>
<exclude>${sedona.source.excludes}</exclude>
</excludes>
<testExcludes>
<testExclude>${sedona.source.excludes}</testExclude>
</testExcludes>
</configuration>
</plugin>
<plugin>
Expand Down Expand Up @@ -299,18 +296,14 @@
</execution>
</executions>
</plugin>
<!-- Run the MEOS smoke tests against the pinned libmeos: propagate
meos.enabled to the forked JVM and resolve libmeos from the repo
lib dir (LD_LIBRARY_PATH takes precedence over a stale system lib). -->
<!-- Run the MEOS smoke tests against the pinned libmeos, resolved from the
repo lib dir (LD_LIBRARY_PATH takes precedence over a stale system lib). -->
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-surefire-plugin</artifactId>
<version>3.2.5</version>
<configuration>
<reuseForks>false</reuseForks>
<systemPropertyVariables>
<meos.enabled>${meos.enabled}</meos.enabled>
</systemPropertyVariables>
<environmentVariables>
<LD_LIBRARY_PATH>${meos.lib.dir}</LD_LIBRARY_PATH>
</environmentVariables>
Expand Down Expand Up @@ -345,12 +338,4 @@
</profile>
</profiles>

<!-- Optional extended temporal-type families, mirroring the MobilityDB/MEOS
CMake build flags. Family inclusion is selected at build time with the
same uppercase flag names and ON|OFF (also 1|0) values as MEOS:
-DCBUFFER=ON -DNPOINT=OFF -DPOSE=ON -DRGEO=ON -DH3=ON
Defaults match MEOS: NPOINT is included by default; CBUFFER, POSE, RGEO,
and H3 are excluded unless their flag is ON|1. RGEO requires POSE (enable
both). When a family is excluded, its generated MeosOps* facade sources
and its smoke test are dropped from the build. -->
</project>
Original file line number Diff line number Diff line change
Expand Up @@ -51,7 +51,7 @@
*
* <p><b>Typical usage</b> — 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}):
*
* <pre>{@code
* KeyedStream<VehiclePosition, Integer> a = streamA.keyBy(VehiclePosition::regionId);
Expand All @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -44,14 +44,14 @@
* independently.
*
* <p><b>Typical usage</b>: scalar-predicate filter against the
* generated {@code MeosOpsTBox.overlaps_tbox_tbox} (tier =
* generated {@code MeosOpsFreeCore.overlaps_tbox_tbox} (tier =
* {@code stateless}):
*
* <pre>{@code
* DataStream<TbiePair> in = ...;
* DataStream<TbiePair> overlapping = in.filter(
* new MeosStatelessFilter<>(
* pair -> MeosOpsTBox.overlaps_tbox_tbox(pair.a, pair.b)));
* pair -> MeosOpsFreeCore.overlaps_tbox_tbox(pair.a, pair.b)));
* }</pre>
*
* <p>For int-coded predicates (JMEOS returns {@code int} for some MEOS
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -52,14 +52,14 @@
*
* <p><b>Typical usage</b>: 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):
*
* <pre>{@code
* DataStream<TbiePair> in = ...; // (tboxA, tboxB)
* DataStream<Boolean> overlap = in.map(
* new MeosStatelessMap<>(
* pair -> MeosOpsTBox.overlaps_tbox_tbox(pair.a, pair.b)));
* pair -> MeosOpsFreeCore.overlaps_tbox_tbox(pair.a, pair.b)));
* }</pre>
*
* <p><b>Tier coverage</b>: as of the codegen state on the parent PR,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -38,8 +38,8 @@
*
* <p>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.
Expand All @@ -61,7 +61,7 @@
* per-window.
*
* <p><b>Typical usage</b> — per-vehicle per-tumbling-window
* trajectory length via {@code MeosOpsTemporal.temporal_length} (tier
* trajectory length via {@code MeosOpsTPoint.tpoint_length} (tier
* = {@code windowed}):
*
* <pre>{@code
Expand All @@ -73,7 +73,7 @@
* .process(new MeosWindowedAggregate<Integer, VehiclePoint, VehicleLength, TimeWindow>(
* (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);
* }));
* }</pre>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<TboxPair> 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<TboxPair> overlapping = stream.filter(filter);
Expand All @@ -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).

Expand All @@ -72,7 +72,7 @@ Run with:
```bash
mvn -q exec:java \
-Dexec.mainClass=org.mobilitydb.flink.meos.wirings.demo.MeosWiringsDemoJob \
-Dmeos.enabled=true
-Dmobilityflink.meos.enabled=true
```

Output (expected): two `overlapping-tbox-hex` lines (the two input
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -99,7 +99,7 @@
* <pre>{@code
* mvn -q exec:java \
* -Dexec.mainClass=org.mobilitydb.flink.meos.wirings.demo.MeosAllTiersCapstoneDemo \
* -Dmeos.enabled=true
* -Dmobilityflink.meos.enabled=true
* }</pre>
*/
public final class MeosAllTiersCapstoneDemo {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -72,7 +72,7 @@
* <pre>{@code
* mvn -q exec:java \
* -Dexec.mainClass=org.mobilitydb.flink.meos.wirings.demo.MeosBoundedStateDemoJob \
* -Dmeos.enabled=true
* -Dmobilityflink.meos.enabled=true
* }</pre>
*
* <p>Expected output: 6 lines (3 per vehicle), each showing the growing
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -77,7 +77,7 @@
* <pre>{@code
* mvn -q exec:java \
* -Dexec.mainClass=org.mobilitydb.flink.meos.wirings.demo.MeosCrossStreamDemoJob \
* -Dmeos.enabled=true
* -Dmobilityflink.meos.enabled=true
* }</pre>
*/
public final class MeosCrossStreamDemoJob {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -75,7 +75,7 @@
* <pre>{@code
* mvn -q exec:java \
* -Dexec.mainClass=org.mobilitydb.flink.meos.wirings.demo.MeosWindowedDemoJob \
* -Dmeos.enabled=true
* -Dmobilityflink.meos.enabled=true
* }</pre>
*
* <p>Expected output: 4 lines (2 windows × 2 vehicles), each showing
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -49,7 +49,7 @@
* <li>Parses each into a JMEOS {@code Pointer} via
* {@code MeosOpsTBox.tbox_in} (tier = {@code io-meta}).</li>
* <li>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}).</li>
* <li>Maps each surviving TBox to its serialized WKB hex via
* {@code MeosOpsTBox.tbox_as_hexwkb} wrapped as a
Expand All @@ -62,11 +62,11 @@
* <pre>{@code
* mvn -q exec:java \
* -Dexec.mainClass=org.mobilitydb.flink.meos.wirings.demo.MeosWiringsDemoJob \
* -Dmeos.enabled=true # require libmeos loadable
* -Dmobilityflink.meos.enabled=true # require libmeos loadable
* }</pre>
*
* <p>If libmeos is not loadable on the runtime (or
* {@code -Dmeos.enabled=false}), every wrapped MeosOps
* {@code -Dmobilityflink.meos.enabled=false}), every wrapped MeosOps
* call throws {@code UnsupportedOperationException} with a clear
* message — the demo prints the throw shape and exits non-zero.
*/
Expand All @@ -89,7 +89,7 @@ public static void main(String[] args) throws Exception {
// fires the first time any MeosOps class is touched).
if (!MeosOpsTBox.MEOS_AVAILABLE) {
LOG.error("MEOS not available — the demo requires libmeos. "
+ "Set -Dmeos.enabled=true and ensure libmeos is loadable.");
+ "Set -Dmobilityflink.meos.enabled=true and ensure libmeos is loadable.");
System.exit(1);
}

Expand Down
6 changes: 2 additions & 4 deletions benchmark/src/test/java/berlinmod/BerlinMODBenchmarkTest.java
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,6 @@
import java.util.List;

import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.condition.EnabledIfSystemProperty;

import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
Expand All @@ -15,10 +14,9 @@
* only runnable as a {@code main}. Runs one keyed cell and one single-subtask
* ({@code keyBy(x -> 0)}) cell through {@link BerlinMODBenchmark#runCell}, which
* builds and executes a real Flink job whose spatial predicate evaluates through
* MEOS. Runs only with an extended libmeos on the loader path
* ({@code -Dmeos.enabled=true}), like the other MEOS-backed tests.
* MEOS, so it needs an extended libmeos on the loader path, like the other
* MEOS-backed tests.
*/
@EnabledIfSystemProperty(named = "meos.enabled", matches = "true")
class BerlinMODBenchmarkTest {

/** Two vehicles moving on nearby WGS84 tracks — enough for the Flink jobs to run. */
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,6 @@
import org.junit.jupiter.api.AfterAll;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.condition.EnabledIfSystemProperty;
import org.mobilitydb.meos.MeosSetSetJoin;

import java.util.HashSet;
Expand All @@ -43,10 +42,8 @@
* {@link MeosSetSetJoin} set-set family) against an independent per-pair scalar
* baseline ({@code edwithin_tgeo_tgeo} / {@code eintersects_tgeo_tgeo}). The two
* code paths must agree exactly on which trip pairs ever meet / are always
* disjoint. Runs only with {@code -Dmeos.enabled=true} and an extended libmeos
* on the library path.
* disjoint. It needs an extended libmeos on the library path.
*/
@EnabledIfSystemProperty(named = "meos.enabled", matches = "true")
class BerlinMODSetSetJoinTest {

// Four trajectory trips: T1 crosses T0's path mid-window; T3 coincides with
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -32,7 +32,6 @@
import org.junit.jupiter.api.AfterAll;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.condition.EnabledIfSystemProperty;

import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNotNull;
Expand All @@ -43,7 +42,6 @@
* family ({@code -DCBUFFER=ON}); the family requires a libmeos built with
* {@code -DCBUFFER=ON}.
*/
@EnabledIfSystemProperty(named = "meos.enabled", matches = "true")
class MeosCbufferSmokeTest {

@BeforeAll
Expand All @@ -59,9 +57,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));
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -32,22 +32,18 @@
import org.junit.jupiter.api.AfterAll;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.condition.EnabledIfSystemProperty;

import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertTrue;

/**
* Runtime check that the always-built MEOS facade families (core and geo) call
* into libmeos and return correct results. Each constructs a value through a
* {@code MeosOps*} facade method and reads it back. Runs only with
* {@code -Dmeos.enabled=true} and a libmeos on the load path. The
* optional families have their own gated smoke tests
* {@code MeosOps*} facade method and reads it back, against a libmeos on the
* load path. The optional families have their own smoke tests
* ({@link MeosCbufferSmokeTest}, {@link MeosNpointSmokeTest},
* {@link MeosPoseSmokeTest}), each compiled only when its build flag includes
* the family.
* {@link MeosPoseSmokeTest}).
*/
@EnabledIfSystemProperty(named = "meos.enabled", matches = "true")
class MeosFacadeSmokeTest {

@BeforeAll
Expand Down Expand Up @@ -86,8 +82,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"));
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -32,7 +32,6 @@
import org.junit.jupiter.api.AfterAll;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.condition.EnabledIfSystemProperty;

import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNotNull;
Expand All @@ -42,7 +41,6 @@
* correct results. Compiled and run when the build includes the npoint family
* (the default; dropped with {@code -DNPOINT=OFF}).
*/
@EnabledIfSystemProperty(named = "meos.enabled", matches = "true")
class MeosNpointSmokeTest {

@BeforeAll
Expand All @@ -58,9 +56,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);
}
}
Loading
Loading