This is an automated email from the ASF dual-hosted git repository.
tkobayas pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/incubator-kie.git
The following commit(s) were added to refs/heads/main by this push:
new e3ca7ea26e6 [incubator-kie-7105] RuleUnit DSL: Follow-up CEP operators
(#7105) (#7122)
e3ca7ea26e6 is described below
commit e3ca7ea26e62d689b62acf50297b52b2e474ad7e
Author: Toshiya Kobayashi <[email protected]>
AuthorDate: Thu Sep 24 17:28:50 2026 +0900
[incubator-kie-7105] RuleUnit DSL: Follow-up CEP operators (#7105) (#7122)
* [incubator-kie-7105] RuleUnit DSL: Follow-up CEP operators (#7105)
* add more tests
---
.../drools/ruleunits/dsl/patterns/Pattern2Def.java | 186 ++++
.../ruleunits/dsl/patterns/Pattern2DefImpl.java | 155 +++
.../java/org/drools/ruleunits/dsl/CepTest.java | 1173 ++++++++++++++++++--
3 files changed, 1409 insertions(+), 105 deletions(-)
diff --git
a/drools-ruleunits/drools-ruleunits-dsl/src/main/java/org/drools/ruleunits/dsl/patterns/Pattern2Def.java
b/drools-ruleunits/drools-ruleunits-dsl/src/main/java/org/drools/ruleunits/dsl/patterns/Pattern2Def.java
index c59b14ce290..752bdc940b6 100644
---
a/drools-ruleunits/drools-ruleunits-dsl/src/main/java/org/drools/ruleunits/dsl/patterns/Pattern2Def.java
+++
b/drools-ruleunits/drools-ruleunits-dsl/src/main/java/org/drools/ruleunits/dsl/patterns/Pattern2Def.java
@@ -41,6 +41,18 @@ public interface Pattern2Def<A, B> extends PatternDef {
<V> Pattern2Def<A, B> filter(String fieldName, Function1<B, V>
leftExtractor, Index.ConstraintType constraintType, Function1<A, V>
rightExtractor);
+ /**
+ * Constrains pattern B to occur after pattern A.
+ * Both operands must be events annotated with {@code
@Role(Role.Type.EVENT)}.
+ */
+ Pattern2Def<A, B> after();
+
+ /**
+ * Constrains pattern B to occur after pattern A with at least {@code min}
interval.
+ * Both operands must be events annotated with {@code
@Role(Role.Type.EVENT)}.
+ */
+ Pattern2Def<A, B> after(long min, TimeUnit unit);
+
/**
* Constrains pattern B to occur after pattern A within inclusive bounds
{@code [min, max]}
* converted to the given time unit. The gap is measured from {@code
end(A)} to {@code start(B)}.
@@ -49,6 +61,18 @@ public interface Pattern2Def<A, B> extends PatternDef {
*/
Pattern2Def<A, B> after(long min, long max, TimeUnit unit);
+ /**
+ * Constrains pattern B to occur before pattern A.
+ * Both operands must be events annotated with {@code
@Role(Role.Type.EVENT)}.
+ */
+ Pattern2Def<A, B> before();
+
+ /**
+ * Constrains pattern B to occur before pattern A with at least {@code
min} interval.
+ * Both operands must be events annotated with {@code
@Role(Role.Type.EVENT)}.
+ */
+ Pattern2Def<A, B> before(long min, TimeUnit unit);
+
/**
* Constrains pattern B to occur before pattern A within inclusive bounds
{@code [min, max]}
* converted to the given time unit. The gap is measured from {@code
end(B)} to {@code start(A)}.
@@ -57,6 +81,168 @@ public interface Pattern2Def<A, B> extends PatternDef {
*/
Pattern2Def<A, B> before(long min, long max, TimeUnit unit);
+ /**
+ * Constrains pattern B and pattern A to start and end at the same time.
+ * Both operands must be events annotated with {@code
@Role(Role.Type.EVENT)}.
+ */
+ Pattern2Def<A, B> coincides();
+
+ /**
+ * Constrains pattern B and pattern A to start and end at the same time
within allowed deviation {@code dev}.
+ * Both operands must be events annotated with {@code
@Role(Role.Type.EVENT)}.
+ */
+ Pattern2Def<A, B> coincides(long dev, TimeUnit devUnit);
+
+ /**
+ * Constrains pattern B and pattern A to start and end at the same time
within allowed start and end deviations.
+ * Both operands must be events annotated with {@code
@Role(Role.Type.EVENT)}.
+ */
+ Pattern2Def<A, B> coincides(long startDev, TimeUnit startDevUnit, long
endDev, TimeUnit endDevUnit);
+
+ /**
+ * Constrains pattern B to occur during pattern A.
+ * Both operands must be events annotated with {@code
@Role(Role.Type.EVENT)}.
+ */
+ Pattern2Def<A, B> during();
+
+ /**
+ * Constrains pattern B to occur during pattern A with max distance.
+ * Both operands must be events annotated with {@code
@Role(Role.Type.EVENT)}.
+ */
+ Pattern2Def<A, B> during(long max, TimeUnit maxUnit);
+
+ /**
+ * Constrains pattern B to occur during pattern A within [min, max]
distance.
+ * Both operands must be events annotated with {@code
@Role(Role.Type.EVENT)}.
+ */
+ Pattern2Def<A, B> during(long min, TimeUnit minUnit, long max, TimeUnit
maxUnit);
+
+ /**
+ * Constrains pattern B to include pattern A.
+ * Both operands must be events annotated with {@code
@Role(Role.Type.EVENT)}.
+ */
+ Pattern2Def<A, B> includes();
+
+ /**
+ * Constrains pattern B to include pattern A with max distance.
+ * Both operands must be events annotated with {@code
@Role(Role.Type.EVENT)}.
+ */
+ Pattern2Def<A, B> includes(long max, TimeUnit maxUnit);
+
+ /**
+ * Constrains pattern B to include pattern A within [min, max] distance.
+ * Both operands must be events annotated with {@code
@Role(Role.Type.EVENT)}.
+ */
+ Pattern2Def<A, B> includes(long min, TimeUnit minUnit, long max, TimeUnit
maxUnit);
+
+ /**
+ * Constrains pattern B to overlap pattern A.
+ * Both operands must be events annotated with {@code
@Role(Role.Type.EVENT)}.
+ */
+ Pattern2Def<A, B> overlaps();
+
+ /**
+ * Constrains pattern B to overlap pattern A with max deviation.
+ * Both operands must be events annotated with {@code
@Role(Role.Type.EVENT)}.
+ */
+ Pattern2Def<A, B> overlaps(long maxDev, TimeUnit maxDevTimeUnit);
+
+ /**
+ * Constrains pattern B to overlap pattern A within [minDev, maxDev].
+ * Both operands must be events annotated with {@code
@Role(Role.Type.EVENT)}.
+ */
+ Pattern2Def<A, B> overlaps(long minDev, TimeUnit minDevTimeUnit, long
maxDev, TimeUnit maxDevTimeUnit);
+
+ /**
+ * Constrains pattern B to be overlapped by pattern A.
+ * Both operands must be events annotated with {@code
@Role(Role.Type.EVENT)}.
+ */
+ Pattern2Def<A, B> overlappedby();
+
+ /**
+ * Constrains pattern B to be overlapped by pattern A with max deviation.
+ * Both operands must be events annotated with {@code
@Role(Role.Type.EVENT)}.
+ */
+ Pattern2Def<A, B> overlappedby(long maxDev, TimeUnit maxDevTimeUnit);
+
+ /**
+ * Constrains pattern B to be overlapped by pattern A within [minDev,
maxDev].
+ * Both operands must be events annotated with {@code
@Role(Role.Type.EVENT)}.
+ */
+ Pattern2Def<A, B> overlappedby(long minDev, TimeUnit minDevTimeUnit, long
maxDev, TimeUnit maxDevTimeUnit);
+
+ /**
+ * Constrains pattern B to meet pattern A.
+ * Both operands must be events annotated with {@code
@Role(Role.Type.EVENT)}.
+ */
+ Pattern2Def<A, B> meets();
+
+ /**
+ * Constrains pattern B to meet pattern A with allowed deviation.
+ * Both operands must be events annotated with {@code
@Role(Role.Type.EVENT)}.
+ */
+ Pattern2Def<A, B> meets(long dev, TimeUnit devUnit);
+
+ /**
+ * Constrains pattern B to be met by pattern A.
+ * Both operands must be events annotated with {@code
@Role(Role.Type.EVENT)}.
+ */
+ Pattern2Def<A, B> metby();
+
+ /**
+ * Constrains pattern B to be met by pattern A with allowed deviation.
+ * Both operands must be events annotated with {@code
@Role(Role.Type.EVENT)}.
+ */
+ Pattern2Def<A, B> metby(long dev, TimeUnit devUnit);
+
+ /**
+ * Constrains pattern B to start pattern A.
+ * Both operands must be events annotated with {@code
@Role(Role.Type.EVENT)}.
+ */
+ Pattern2Def<A, B> starts();
+
+ /**
+ * Constrains pattern B to start pattern A with allowed deviation.
+ * Both operands must be events annotated with {@code
@Role(Role.Type.EVENT)}.
+ */
+ Pattern2Def<A, B> starts(long dev, TimeUnit devUnit);
+
+ /**
+ * Constrains pattern B to be started by pattern A.
+ * Both operands must be events annotated with {@code
@Role(Role.Type.EVENT)}.
+ */
+ Pattern2Def<A, B> startedby();
+
+ /**
+ * Constrains pattern B to be started by pattern A with allowed deviation.
+ * Both operands must be events annotated with {@code
@Role(Role.Type.EVENT)}.
+ */
+ Pattern2Def<A, B> startedby(long dev, TimeUnit devUnit);
+
+ /**
+ * Constrains pattern B to finish pattern A.
+ * Both operands must be events annotated with {@code
@Role(Role.Type.EVENT)}.
+ */
+ Pattern2Def<A, B> finishes();
+
+ /**
+ * Constrains pattern B to finish pattern A with allowed deviation.
+ * Both operands must be events annotated with {@code
@Role(Role.Type.EVENT)}.
+ */
+ Pattern2Def<A, B> finishes(long dev, TimeUnit devUnit);
+
+ /**
+ * Constrains pattern B to be finished by pattern A.
+ * Both operands must be events annotated with {@code
@Role(Role.Type.EVENT)}.
+ */
+ Pattern2Def<A, B> finishedby();
+
+ /**
+ * Constrains pattern B to be finished by pattern A with allowed deviation.
+ * Both operands must be events annotated with {@code
@Role(Role.Type.EVENT)}.
+ */
+ Pattern2Def<A, B> finishedby(long dev, TimeUnit devUnit);
+
<C> Pattern3Def<A, B, C> on(DataSource<C> dataSource);
<C> Pattern3Def<A, B, C> join(Function1<RuleFactory, Pattern1Def<C>>
patternBuilder);
diff --git
a/drools-ruleunits/drools-ruleunits-dsl/src/main/java/org/drools/ruleunits/dsl/patterns/Pattern2DefImpl.java
b/drools-ruleunits/drools-ruleunits-dsl/src/main/java/org/drools/ruleunits/dsl/patterns/Pattern2DefImpl.java
index da7108f82ed..e89685702be 100644
---
a/drools-ruleunits/drools-ruleunits-dsl/src/main/java/org/drools/ruleunits/dsl/patterns/Pattern2DefImpl.java
+++
b/drools-ruleunits/drools-ruleunits-dsl/src/main/java/org/drools/ruleunits/dsl/patterns/Pattern2DefImpl.java
@@ -74,16 +74,171 @@ public class Pattern2DefImpl<A, B> extends
SinglePatternDef<B> implements Patter
return this;
}
+ @Override
+ public Pattern2DefImpl<A, B> after() {
+ return addTemporalConstraint(DSL.after());
+ }
+
+ @Override
+ public Pattern2DefImpl<A, B> after(long min, TimeUnit unit) {
+ return addTemporalConstraint(DSL.after(min, unit));
+ }
+
@Override
public Pattern2DefImpl<A, B> after(long min, long max, TimeUnit unit) {
return addTemporalConstraint(DSL.after(min, unit, max, unit));
}
+ @Override
+ public Pattern2DefImpl<A, B> before() {
+ return addTemporalConstraint(DSL.before());
+ }
+
+ @Override
+ public Pattern2DefImpl<A, B> before(long min, TimeUnit unit) {
+ return addTemporalConstraint(DSL.before(min, unit));
+ }
+
@Override
public Pattern2DefImpl<A, B> before(long min, long max, TimeUnit unit) {
return addTemporalConstraint(DSL.before(min, unit, max, unit));
}
+ @Override
+ public Pattern2DefImpl<A, B> coincides() {
+ return addTemporalConstraint(DSL.coincides());
+ }
+
+ @Override
+ public Pattern2DefImpl<A, B> coincides(long dev, TimeUnit devUnit) {
+ return addTemporalConstraint(DSL.coincides(dev, devUnit));
+ }
+
+ @Override
+ public Pattern2DefImpl<A, B> coincides(long startDev, TimeUnit
startDevUnit, long endDev, TimeUnit endDevUnit) {
+ return addTemporalConstraint(DSL.coincides(startDev, startDevUnit,
endDev, endDevUnit));
+ }
+
+ @Override
+ public Pattern2DefImpl<A, B> during() {
+ return addTemporalConstraint(DSL.during());
+ }
+
+ @Override
+ public Pattern2DefImpl<A, B> during(long max, TimeUnit maxUnit) {
+ return addTemporalConstraint(DSL.during(max, maxUnit));
+ }
+
+ @Override
+ public Pattern2DefImpl<A, B> during(long min, TimeUnit minUnit, long max,
TimeUnit maxUnit) {
+ return addTemporalConstraint(DSL.during(min, minUnit, max, maxUnit));
+ }
+
+ @Override
+ public Pattern2DefImpl<A, B> includes() {
+ return addTemporalConstraint(DSL.includes());
+ }
+
+ @Override
+ public Pattern2DefImpl<A, B> includes(long max, TimeUnit maxUnit) {
+ return addTemporalConstraint(DSL.includes(max, maxUnit));
+ }
+
+ @Override
+ public Pattern2DefImpl<A, B> includes(long min, TimeUnit minUnit, long
max, TimeUnit maxUnit) {
+ return addTemporalConstraint(DSL.includes(min, minUnit, max, maxUnit));
+ }
+
+ @Override
+ public Pattern2DefImpl<A, B> overlaps() {
+ return addTemporalConstraint(DSL.overlaps());
+ }
+
+ @Override
+ public Pattern2DefImpl<A, B> overlaps(long maxDev, TimeUnit
maxDevTimeUnit) {
+ return addTemporalConstraint(DSL.overlaps(maxDev, maxDevTimeUnit));
+ }
+
+ @Override
+ public Pattern2DefImpl<A, B> overlaps(long minDev, TimeUnit
minDevTimeUnit, long maxDev, TimeUnit maxDevTimeUnit) {
+ return addTemporalConstraint(DSL.overlaps(minDev, minDevTimeUnit,
maxDev, maxDevTimeUnit));
+ }
+
+ @Override
+ public Pattern2DefImpl<A, B> overlappedby() {
+ return addTemporalConstraint(DSL.overlappedby());
+ }
+
+ @Override
+ public Pattern2DefImpl<A, B> overlappedby(long maxDev, TimeUnit
maxDevTimeUnit) {
+ return addTemporalConstraint(DSL.overlappedby(maxDev, maxDevTimeUnit));
+ }
+
+ @Override
+ public Pattern2DefImpl<A, B> overlappedby(long minDev, TimeUnit
minDevTimeUnit, long maxDev, TimeUnit maxDevTimeUnit) {
+ return addTemporalConstraint(DSL.overlappedby(minDev, minDevTimeUnit,
maxDev, maxDevTimeUnit));
+ }
+
+ @Override
+ public Pattern2DefImpl<A, B> meets() {
+ return addTemporalConstraint(DSL.meets());
+ }
+
+ @Override
+ public Pattern2DefImpl<A, B> meets(long dev, TimeUnit devUnit) {
+ return addTemporalConstraint(DSL.meets(dev, devUnit));
+ }
+
+ @Override
+ public Pattern2DefImpl<A, B> metby() {
+ return addTemporalConstraint(DSL.metby());
+ }
+
+ @Override
+ public Pattern2DefImpl<A, B> metby(long dev, TimeUnit devUnit) {
+ return addTemporalConstraint(DSL.metby(dev, devUnit));
+ }
+
+ @Override
+ public Pattern2DefImpl<A, B> starts() {
+ return addTemporalConstraint(DSL.starts());
+ }
+
+ @Override
+ public Pattern2DefImpl<A, B> starts(long dev, TimeUnit devUnit) {
+ return addTemporalConstraint(DSL.starts(dev, devUnit));
+ }
+
+ @Override
+ public Pattern2DefImpl<A, B> startedby() {
+ return addTemporalConstraint(DSL.startedby());
+ }
+
+ @Override
+ public Pattern2DefImpl<A, B> startedby(long dev, TimeUnit devUnit) {
+ return addTemporalConstraint(DSL.startedby(dev, devUnit));
+ }
+
+ @Override
+ public Pattern2DefImpl<A, B> finishes() {
+ return addTemporalConstraint(DSL.finishes());
+ }
+
+ @Override
+ public Pattern2DefImpl<A, B> finishes(long dev, TimeUnit devUnit) {
+ return addTemporalConstraint(DSL.finishes(dev, devUnit));
+ }
+
+ @Override
+ public Pattern2DefImpl<A, B> finishedby() {
+ return addTemporalConstraint(DSL.finishedby());
+ }
+
+ @Override
+ public Pattern2DefImpl<A, B> finishedby(long dev, TimeUnit devUnit) {
+ return addTemporalConstraint(DSL.finishedby(dev, devUnit));
+ }
+
protected Pattern2DefImpl<A, B> addTemporalConstraint(TemporalPredicate
temporalPredicate) {
patternB.constraints.add(patternDef ->
patternDef.expr(UUID.randomUUID().toString(), patternA.variable,
temporalPredicate));
return this;
diff --git
a/drools-ruleunits/drools-ruleunits-dsl/src/test/java/org/drools/ruleunits/dsl/CepTest.java
b/drools-ruleunits/drools-ruleunits-dsl/src/test/java/org/drools/ruleunits/dsl/CepTest.java
index 10b9253ea79..f3c1273a0d0 100644
---
a/drools-ruleunits/drools-ruleunits-dsl/src/test/java/org/drools/ruleunits/dsl/CepTest.java
+++
b/drools-ruleunits/drools-ruleunits-dsl/src/test/java/org/drools/ruleunits/dsl/CepTest.java
@@ -205,6 +205,40 @@ public class CepTest {
}
}
+ @Test
+ public void afterNoArg() {
+ StreamAfterNoArgUnit unit = new StreamAfterNoArgUnit();
+ RuleConfig config = RuleUnitProvider.get().newRuleConfig();
+ config.setClockType(ClockType.PSEUDO);
+ try (RuleUnitInstance<StreamAfterNoArgUnit> instance =
RuleUnitProvider.get().createRuleUnitInstance(unit, config)) {
+ SessionPseudoClock clock = instance.getClock();
+ unit.getStockTicks().append(new StockTick("DROO"));
+ clock.advanceTime(1, TimeUnit.SECONDS);
+ unit.getStockTicks().append(new StockTick("ACME"));
+ instance.fire();
+
+ assertThat(unit.getResults()).hasSize(1);
+
assertThat(unit.getResults().get(0).getCompany()).isEqualTo("ACME");
+ }
+ }
+
+ @Test
+ public void afterSingleBound() {
+ StreamAfterSingleBoundUnit unit = new StreamAfterSingleBoundUnit();
+ RuleConfig config = RuleUnitProvider.get().newRuleConfig();
+ config.setClockType(ClockType.PSEUDO);
+ try (RuleUnitInstance<StreamAfterSingleBoundUnit> instance =
RuleUnitProvider.get().createRuleUnitInstance(unit, config)) {
+ SessionPseudoClock clock = instance.getClock();
+ unit.getStockTicks().append(new StockTick("DROO"));
+ clock.advanceTime(5, TimeUnit.SECONDS);
+ unit.getStockTicks().append(new StockTick("ACME"));
+ instance.fire();
+
+ assertThat(unit.getResults()).hasSize(1);
+
assertThat(unit.getResults().get(0).getCompany()).isEqualTo("ACME");
+ }
+ }
+
// --- Temporal constraint: before ---
@Test
@@ -261,179 +295,1108 @@ public class CepTest {
}
}
- // --- Explicit expiration ---
+ @Test
+ public void beforeNoArg() {
+ StreamBeforeNoArgUnit unit = new StreamBeforeNoArgUnit();
+ RuleConfig config = RuleUnitProvider.get().newRuleConfig();
+ config.setClockType(ClockType.PSEUDO);
+ try (RuleUnitInstance<StreamBeforeNoArgUnit> instance =
RuleUnitProvider.get().createRuleUnitInstance(unit, config)) {
+ SessionPseudoClock clock = instance.getClock();
+ unit.getStockTicks().append(new StockTick("ACME"));
+ clock.advanceTime(1, TimeUnit.SECONDS);
+ unit.getStockTicks().append(new StockTick("DROO"));
+ instance.fire();
+
+ assertThat(unit.getResults()).hasSize(1);
+
assertThat(unit.getResults().get(0).getCompany()).isEqualTo("ACME");
+ }
+ }
@Test
- public void expiresRemovesEventAfterDuration() {
- StreamExpirationUnit unit = new StreamExpirationUnit();
+ public void beforeSingleBound() {
+ StreamBeforeSingleBoundUnit unit = new StreamBeforeSingleBoundUnit();
RuleConfig config = RuleUnitProvider.get().newRuleConfig();
config.setClockType(ClockType.PSEUDO);
- try (RuleUnitInstance<StreamExpirationUnit> instance =
RuleUnitProvider.get().createRuleUnitInstance(unit, config)) {
+ try (RuleUnitInstance<StreamBeforeSingleBoundUnit> instance =
RuleUnitProvider.get().createRuleUnitInstance(unit, config)) {
SessionPseudoClock clock = instance.getClock();
+ unit.getStockTicks().append(new StockTick("ACME"));
+ clock.advanceTime(5, TimeUnit.SECONDS);
unit.getStockTicks().append(new StockTick("DROO"));
instance.fire();
- EntryPoint entryPoint = findNonDefaultEntryPoint(instance);
- assertThat(entryPoint.getFactHandles()).hasSize(1);
+ assertThat(unit.getResults()).hasSize(1);
+
assertThat(unit.getResults().get(0).getCompany()).isEqualTo("ACME");
+ }
+ }
- clock.advanceTime(11, TimeUnit.SECONDS);
+ // --- Temporal constraint: coincides ---
+
+ @Test
+ public void coincidesExactMatch() {
+ StreamCoincidesUnit unit = new StreamCoincidesUnit();
+ RuleConfig config = RuleUnitProvider.get().newRuleConfig();
+ config.setClockType(ClockType.PSEUDO);
+ try (RuleUnitInstance<StreamCoincidesUnit> instance =
RuleUnitProvider.get().createRuleUnitInstance(unit, config)) {
+ unit.getStockTicks().append(new StockTick("DROO", 5000));
+ unit.getStockTicks().append(new StockTick("ACME", 5000));
instance.fire();
- assertThat(entryPoint.getFactHandles()).isEmpty();
+ assertThat(unit.getResults()).hasSize(1);
+
assertThat(unit.getResults().get(0).getCompany()).isEqualTo("ACME");
}
}
- // --- Combined CEP scenario ---
+ @Test
+ public void coincidesDoesNotMatchWhenDifferent() {
+ StreamCoincidesUnit unit = new StreamCoincidesUnit();
+ RuleConfig config = RuleUnitProvider.get().newRuleConfig();
+ config.setClockType(ClockType.PSEUDO);
+ try (RuleUnitInstance<StreamCoincidesUnit> instance =
RuleUnitProvider.get().createRuleUnitInstance(unit, config)) {
+ SessionPseudoClock clock = instance.getClock();
+ unit.getStockTicks().append(new StockTick("DROO", 5000));
+ clock.advanceTime(100, TimeUnit.MILLISECONDS);
+ unit.getStockTicks().append(new StockTick("ACME", 5000));
+ instance.fire();
+
+ assertThat(unit.getResults()).isEmpty();
+ }
+ }
@Test
- public void combinedCepScenario() {
- StreamAfterUnit unit = new StreamAfterUnit();
+ public void coincidesWithDevMatchesWithinDeviation() {
+ StreamCoincidesDevUnit unit = new StreamCoincidesDevUnit();
RuleConfig config = RuleUnitProvider.get().newRuleConfig();
config.setClockType(ClockType.PSEUDO);
- try (RuleUnitInstance<StreamAfterUnit> instance =
RuleUnitProvider.get().createRuleUnitInstance(unit, config)) {
+ try (RuleUnitInstance<StreamCoincidesDevUnit> instance =
RuleUnitProvider.get().createRuleUnitInstance(unit, config)) {
SessionPseudoClock clock = instance.getClock();
+ unit.getStockTicks().append(new StockTick("DROO", 5000));
+ // start deviation 500ms <= 1s dev, duration diff 0 <= 1s dev
+ clock.advanceTime(500, TimeUnit.MILLISECONDS);
+ unit.getStockTicks().append(new StockTick("ACME", 5000));
+ instance.fire();
- // 1. Append a zero-duration DROO event at time zero
- unit.getStockTicks().append(new StockTick("DROO"));
+ assertThat(unit.getResults()).hasSize(1);
+
assertThat(unit.getResults().get(0).getCompany()).isEqualTo("ACME");
+ }
+ }
- // 2. Advance 6 seconds and append ACME
- clock.advanceTime(6, TimeUnit.SECONDS);
- unit.getStockTicks().append(new StockTick("ACME"));
+ @Test
+ public void coincidesWithDevDoesNotMatchWhenExceedingDeviation() {
+ StreamCoincidesDevUnit unit = new StreamCoincidesDevUnit();
+ RuleConfig config = RuleUnitProvider.get().newRuleConfig();
+ config.setClockType(ClockType.PSEUDO);
+ try (RuleUnitInstance<StreamCoincidesDevUnit> instance =
RuleUnitProvider.get().createRuleUnitInstance(unit, config)) {
+ SessionPseudoClock clock = instance.getClock();
+ unit.getStockTicks().append(new StockTick("DROO", 5000));
+ // start deviation 1500ms > 1s dev
+ clock.advanceTime(1500, TimeUnit.MILLISECONDS);
+ unit.getStockTicks().append(new StockTick("ACME", 5000));
+ instance.fire();
- // 3. Fire and assert the after(5, 8, SECONDS) rule records ACME
+ assertThat(unit.getResults()).isEmpty();
+ }
+ }
+
+ @Test
+ public void coincidesWithStartEndDevMatchesWithinDeviation() {
+ StreamCoincidesStartEndDevUnit unit = new
StreamCoincidesStartEndDevUnit();
+ RuleConfig config = RuleUnitProvider.get().newRuleConfig();
+ config.setClockType(ClockType.PSEUDO);
+ try (RuleUnitInstance<StreamCoincidesStartEndDevUnit> instance =
RuleUnitProvider.get().createRuleUnitInstance(unit, config)) {
+ SessionPseudoClock clock = instance.getClock();
+ unit.getStockTicks().append(new StockTick("DROO", 5000));
+ // start dev 500ms <= 1s startDev, end dev 1500ms <= 2s endDev
(DROO ends at 5s, ACME ends at 0.5s+6s = 6.5s)
+ clock.advanceTime(500, TimeUnit.MILLISECONDS);
+ unit.getStockTicks().append(new StockTick("ACME", 6000));
instance.fire();
+
assertThat(unit.getResults()).hasSize(1);
assertThat(unit.getResults().get(0).getCompany()).isEqualTo("ACME");
+ }
+ }
- // 4. Both events should still be in the entry point
- EntryPoint entryPoint = findNonDefaultEntryPoint(instance);
- assertThat(entryPoint.getFactHandles()).hasSize(2);
-
- // 5. Advance another 5 seconds (total = 11s) and fire
- clock.advanceTime(5, TimeUnit.SECONDS);
+ @Test
+ public void coincidesWithStartEndDevDoesNotMatchWhenExceedingDeviation() {
+ StreamCoincidesStartEndDevUnit unit = new
StreamCoincidesStartEndDevUnit();
+ RuleConfig config = RuleUnitProvider.get().newRuleConfig();
+ config.setClockType(ClockType.PSEUDO);
+ try (RuleUnitInstance<StreamCoincidesStartEndDevUnit> instance =
RuleUnitProvider.get().createRuleUnitInstance(unit, config)) {
+ SessionPseudoClock clock = instance.getClock();
+ unit.getStockTicks().append(new StockTick("DROO", 5000));
+ // start dev 500ms <= 1s, but end dev = |(0.5s + 8s) - 5s| = 3.5s
> 2s endDev
+ clock.advanceTime(500, TimeUnit.MILLISECONDS);
+ unit.getStockTicks().append(new StockTick("ACME", 8000));
instance.fire();
- // 6. DROO was inserted at t=0 with @Expires("10s"), so it should
be expired
- // ACME was inserted at t=6 with @Expires("10s"), so it should
still be present
- Collection<FactHandle> remaining = entryPoint.getFactHandles();
- assertThat(remaining).hasSize(1);
- assertThat(((StockTick) ((InternalFactHandle)
remaining.iterator().next()).getObject()).getCompany())
- .isEqualTo("ACME");
+ assertThat(unit.getResults()).isEmpty();
}
}
- // --- Unsupported aggregate calls ---
+ // --- Temporal constraint: during / includes ---
@Test
- public void afterOnGroupByThrowsUnsupported() {
- assertThatThrownBy(() -> {
- new GroupByTemporalUnit();
- }).isInstanceOf(UnsupportedOperationException.class)
- .hasMessageContaining("groupBy");
+ public void duringMatches() {
+ StreamDuringUnit unit = new StreamDuringUnit();
+ RuleConfig config = RuleUnitProvider.get().newRuleConfig();
+ config.setClockType(ClockType.PSEUDO);
+ try (RuleUnitInstance<StreamDuringUnit> instance =
RuleUnitProvider.get().createRuleUnitInstance(unit, config)) {
+ SessionPseudoClock clock = instance.getClock();
+ // A starts at 0, duration 10s (ends at 10s)
+ unit.getStockTicks().append(new StockTick("DROO", 10000));
+ // B starts at 2s, duration 4s (ends at 6s) -> start dist = 2s in
[1s, 10s], end dist = 4s in [1s, 10s]
+ clock.advanceTime(2, TimeUnit.SECONDS);
+ unit.getStockTicks().append(new StockTick("ACME", 4000));
+ instance.fire();
+
+ assertThat(unit.getResults()).hasSize(1);
+
assertThat(unit.getResults().get(0).getCompany()).isEqualTo("ACME");
+ }
}
@Test
- public void afterOnAccumulateThrowsUnsupported() {
- assertThatThrownBy(() -> {
- new AccumulateTemporalUnit();
- }).isInstanceOf(UnsupportedOperationException.class)
- .hasMessageContaining("accumulate");
+ public void duringDoesNotMatchWhenNotContained() {
+ StreamDuringUnit unit = new StreamDuringUnit();
+ RuleConfig config = RuleUnitProvider.get().newRuleConfig();
+ config.setClockType(ClockType.PSEUDO);
+ try (RuleUnitInstance<StreamDuringUnit> instance =
RuleUnitProvider.get().createRuleUnitInstance(unit, config)) {
+ SessionPseudoClock clock = instance.getClock();
+ // A starts at 0, duration 5s (ends at 5s)
+ unit.getStockTicks().append(new StockTick("DROO", 5000));
+ // B starts at 2s, duration 5s (ends at 7s > end(A)) -> not during
A
+ clock.advanceTime(2, TimeUnit.SECONDS);
+ unit.getStockTicks().append(new StockTick("ACME", 5000));
+ instance.fire();
+
+ assertThat(unit.getResults()).isEmpty();
+ }
}
- // --- Helpers ---
+ @Test
+ public void includesMatches() {
+ StreamIncludesUnit unit = new StreamIncludesUnit();
+ RuleConfig config = RuleUnitProvider.get().newRuleConfig();
+ config.setClockType(ClockType.PSEUDO);
+ try (RuleUnitInstance<StreamIncludesUnit> instance =
RuleUnitProvider.get().createRuleUnitInstance(unit, config)) {
+ SessionPseudoClock clock = instance.getClock();
+ // B (ACME) starts at 0 with 10s duration (ends at 10s)
+ unit.getStockTicks().append(new StockTick("ACME", 10000));
+ // A (DROO) starts at 2s with 4s duration (ends at 6s) -> start
dist = 2s in [1s, 10s], end dist = 4s in [1s, 10s]
+ clock.advanceTime(2, TimeUnit.SECONDS);
+ unit.getStockTicks().append(new StockTick("DROO", 4000));
+ instance.fire();
- @SuppressWarnings("rawtypes")
- private ReteEvaluator getEvaluator(RuleUnitInstance<?> instance) {
- return (ReteEvaluator) ((AbstractRuleUnitInstance)
instance).getEvaluator();
+ assertThat(unit.getResults()).hasSize(1);
+
assertThat(unit.getResults().get(0).getCompany()).isEqualTo("ACME");
+ }
}
- private EntryPoint findNonDefaultEntryPoint(RuleUnitInstance<?> instance) {
- ReteEvaluator evaluator = getEvaluator(instance);
- for (EntryPoint ep : evaluator.getEntryPoints()) {
- if (!"DEFAULT".equals(ep.getEntryPointId())) {
- return ep;
- }
+ @Test
+ public void includesDoesNotMatchWhenNotIncluding() {
+ StreamIncludesUnit unit = new StreamIncludesUnit();
+ RuleConfig config = RuleUnitProvider.get().newRuleConfig();
+ config.setClockType(ClockType.PSEUDO);
+ try (RuleUnitInstance<StreamIncludesUnit> instance =
RuleUnitProvider.get().createRuleUnitInstance(unit, config)) {
+ SessionPseudoClock clock = instance.getClock();
+ // B (ACME) starts at 0 with 5s duration (ends at 5s)
+ unit.getStockTicks().append(new StockTick("ACME", 5000));
+ // A (DROO) starts at 2s with 5s duration (ends at 7s > end(B)) ->
B does not include A
+ clock.advanceTime(2, TimeUnit.SECONDS);
+ unit.getStockTicks().append(new StockTick("DROO", 5000));
+ instance.fire();
+
+ assertThat(unit.getResults()).isEmpty();
}
- throw new IllegalStateException("No non-DEFAULT entry point found");
}
- // --- Unit definitions ---
+ // --- Temporal constraint: overlaps / overlappedby ---
- @EventProcessing(EventProcessingType.STREAM)
- public static class StreamAfterUnit implements RuleUnitDefinition {
- private final DataStream<StockTick> stockTicks;
- private final List<StockTick> results = new ArrayList<>();
+ @Test
+ public void overlapsMatches() {
+ StreamOverlapsUnit unit = new StreamOverlapsUnit();
+ RuleConfig config = RuleUnitProvider.get().newRuleConfig();
+ config.setClockType(ClockType.PSEUDO);
+ try (RuleUnitInstance<StreamOverlapsUnit> instance =
RuleUnitProvider.get().createRuleUnitInstance(unit, config)) {
+ SessionPseudoClock clock = instance.getClock();
+ // B (ACME) starts at 0s, duration 5s (ends at 5s)
+ unit.getStockTicks().append(new StockTick("ACME", 5000));
+ // A (DROO) starts at 3s, duration 5s (ends at 8s)
+ // B starts before A (0 < 3), B ends after A starts but before A
ends (3 < 5 < 8)
+ // overlap dist = end(B) - start(A) = 5 - 3 = 2s in [1s, 5s]
+ clock.advanceTime(3, TimeUnit.SECONDS);
+ unit.getStockTicks().append(new StockTick("DROO", 5000));
+ instance.fire();
- public StreamAfterUnit() {
- this.stockTicks = DataSource.createStream();
+ assertThat(unit.getResults()).hasSize(1);
+
assertThat(unit.getResults().get(0).getCompany()).isEqualTo("ACME");
}
+ }
- public DataStream<StockTick> getStockTicks() { return stockTicks; }
- public List<StockTick> getResults() { return results; }
+ @Test
+ public void overlapsDoesNotMatchWhenNoOverlap() {
+ StreamOverlapsUnit unit = new StreamOverlapsUnit();
+ RuleConfig config = RuleUnitProvider.get().newRuleConfig();
+ config.setClockType(ClockType.PSEUDO);
+ try (RuleUnitInstance<StreamOverlapsUnit> instance =
RuleUnitProvider.get().createRuleUnitInstance(unit, config)) {
+ SessionPseudoClock clock = instance.getClock();
+ // B (ACME) starts at 0s, duration 2s (ends at 2s)
+ unit.getStockTicks().append(new StockTick("ACME", 2000));
+ // A (DROO) starts at 3s, duration 5s (ends at 8s) -> end(B) <
start(A), no overlap
+ clock.advanceTime(3, TimeUnit.SECONDS);
+ unit.getStockTicks().append(new StockTick("DROO", 5000));
+ instance.fire();
- @Override
- public void defineRules(RulesFactory rulesFactory) {
- // A=$a : /stockTicks [ company == "DROO" ]
- // B=$b : /stockTicks [ company == "ACME", this after[5s,8s] $a ]
- // → B after A: gap = start(B) - end(A) must be in [5s, 8s]
- rulesFactory.rule("ACME after DROO")
- .on(stockTicks) //
pattern A
- .filter(StockTick::getCompany, EQUAL, "DROO")
- .join(rule -> rule.on(stockTicks) //
pattern B
- .filter(StockTick::getCompany, EQUAL, "ACME"))
- .after(5, 8, TimeUnit.SECONDS) // B
after A
- .execute(results, (r, droo, acme) -> r.add(acme));
+ assertThat(unit.getResults()).isEmpty();
}
}
- @EventProcessing(EventProcessingType.STREAM)
- public static class StreamAfterMillisUnit implements RuleUnitDefinition {
- private final DataStream<StockTick> stockTicks;
- private final List<StockTick> results = new ArrayList<>();
+ @Test
+ public void overlappedbyMatches() {
+ StreamOverlappedbyUnit unit = new StreamOverlappedbyUnit();
+ RuleConfig config = RuleUnitProvider.get().newRuleConfig();
+ config.setClockType(ClockType.PSEUDO);
+ try (RuleUnitInstance<StreamOverlappedbyUnit> instance =
RuleUnitProvider.get().createRuleUnitInstance(unit, config)) {
+ SessionPseudoClock clock = instance.getClock();
+ // A (DROO) starts at 0s, duration 5s (ends at 5s)
+ unit.getStockTicks().append(new StockTick("DROO", 5000));
+ // B (ACME) starts at 3s, duration 5s (ends at 8s)
+ // B starts after A starts but before A ends (0 < 3 < 5), B ends
after A (8 > 5)
+ // overlap dist = end(A) - start(B) = 5 - 3 = 2s in [1s, 5s]
+ clock.advanceTime(3, TimeUnit.SECONDS);
+ unit.getStockTicks().append(new StockTick("ACME", 5000));
+ instance.fire();
- public StreamAfterMillisUnit() {
- this.stockTicks = DataSource.createStream();
+ assertThat(unit.getResults()).hasSize(1);
+
assertThat(unit.getResults().get(0).getCompany()).isEqualTo("ACME");
}
+ }
- public DataStream<StockTick> getStockTicks() { return stockTicks; }
- public List<StockTick> getResults() { return results; }
+ @Test
+ public void overlappedbyDoesNotMatchWhenNoOverlap() {
+ StreamOverlappedbyUnit unit = new StreamOverlappedbyUnit();
+ RuleConfig config = RuleUnitProvider.get().newRuleConfig();
+ config.setClockType(ClockType.PSEUDO);
+ try (RuleUnitInstance<StreamOverlappedbyUnit> instance =
RuleUnitProvider.get().createRuleUnitInstance(unit, config)) {
+ SessionPseudoClock clock = instance.getClock();
+ // A (DROO) starts at 0s, duration 2s (ends at 2s)
+ unit.getStockTicks().append(new StockTick("DROO", 2000));
+ // B (ACME) starts at 3s, duration 5s (ends at 8s) -> start(B) >
end(A), no overlap
+ clock.advanceTime(3, TimeUnit.SECONDS);
+ unit.getStockTicks().append(new StockTick("ACME", 5000));
+ instance.fire();
- @Override
- public void defineRules(RulesFactory rulesFactory) {
- // Same as StreamAfterUnit but with millisecond bounds [500ms,
1000ms]
- rulesFactory.rule("ACME after DROO millis")
- .on(stockTicks) //
pattern A
- .filter(StockTick::getCompany, EQUAL, "DROO")
- .join(rule -> rule.on(stockTicks) //
pattern B
- .filter(StockTick::getCompany, EQUAL, "ACME"))
- .after(500, 1000, TimeUnit.MILLISECONDS) // B
after A
- .execute(results, (r, droo, acme) -> r.add(acme));
+ assertThat(unit.getResults()).isEmpty();
}
}
- @EventProcessing(EventProcessingType.STREAM)
- public static class StreamBeforeUnit implements RuleUnitDefinition {
- private final DataStream<StockTick> stockTicks;
- private final List<StockTick> results = new ArrayList<>();
+ // --- Temporal constraint: meets / metby ---
- public StreamBeforeUnit() {
- this.stockTicks = DataSource.createStream();
- }
+ @Test
+ public void meetsMatchesWithinDeviation() {
+ StreamMeetsUnit unit = new StreamMeetsUnit();
+ RuleConfig config = RuleUnitProvider.get().newRuleConfig();
+ config.setClockType(ClockType.PSEUDO);
+ try (RuleUnitInstance<StreamMeetsUnit> instance =
RuleUnitProvider.get().createRuleUnitInstance(unit, config)) {
+ SessionPseudoClock clock = instance.getClock();
+ // B (ACME) starts at 0, duration 3s (ends at 3s)
+ unit.getStockTicks().append(new StockTick("ACME", 3000));
+ // A (DROO) starts at 3.5s -> |start(A) - end(B)| = 500ms <= 1s dev
+ clock.advanceTime(3500, TimeUnit.MILLISECONDS);
+ unit.getStockTicks().append(new StockTick("DROO", 3000));
+ instance.fire();
- public DataStream<StockTick> getStockTicks() { return stockTicks; }
- public List<StockTick> getResults() { return results; }
+ assertThat(unit.getResults()).hasSize(1);
+
assertThat(unit.getResults().get(0).getCompany()).isEqualTo("ACME");
+ }
+ }
- @Override
+ @Test
+ public void meetsDoesNotMatchWhenExceedingDeviation() {
+ StreamMeetsUnit unit = new StreamMeetsUnit();
+ RuleConfig config = RuleUnitProvider.get().newRuleConfig();
+ config.setClockType(ClockType.PSEUDO);
+ try (RuleUnitInstance<StreamMeetsUnit> instance =
RuleUnitProvider.get().createRuleUnitInstance(unit, config)) {
+ SessionPseudoClock clock = instance.getClock();
+ // B (ACME) starts at 0, duration 3s (ends at 3s)
+ unit.getStockTicks().append(new StockTick("ACME", 3000));
+ // A (DROO) starts at 5s -> |start(A) - end(B)| = 2s > 1s dev
+ clock.advanceTime(5, TimeUnit.SECONDS);
+ unit.getStockTicks().append(new StockTick("DROO", 3000));
+ instance.fire();
+
+ assertThat(unit.getResults()).isEmpty();
+ }
+ }
+
+ @Test
+ public void metbyMatchesWithinDeviation() {
+ StreamMetbyUnit unit = new StreamMetbyUnit();
+ RuleConfig config = RuleUnitProvider.get().newRuleConfig();
+ config.setClockType(ClockType.PSEUDO);
+ try (RuleUnitInstance<StreamMetbyUnit> instance =
RuleUnitProvider.get().createRuleUnitInstance(unit, config)) {
+ SessionPseudoClock clock = instance.getClock();
+ // A (DROO) starts at 0, duration 3s (ends at 3s)
+ unit.getStockTicks().append(new StockTick("DROO", 3000));
+ // B (ACME) starts at 3.5s -> |start(B) - end(A)| = 500ms <= 1s dev
+ clock.advanceTime(3500, TimeUnit.MILLISECONDS);
+ unit.getStockTicks().append(new StockTick("ACME", 3000));
+ instance.fire();
+
+ assertThat(unit.getResults()).hasSize(1);
+
assertThat(unit.getResults().get(0).getCompany()).isEqualTo("ACME");
+ }
+ }
+
+ @Test
+ public void metbyDoesNotMatchWhenExceedingDeviation() {
+ StreamMetbyUnit unit = new StreamMetbyUnit();
+ RuleConfig config = RuleUnitProvider.get().newRuleConfig();
+ config.setClockType(ClockType.PSEUDO);
+ try (RuleUnitInstance<StreamMetbyUnit> instance =
RuleUnitProvider.get().createRuleUnitInstance(unit, config)) {
+ SessionPseudoClock clock = instance.getClock();
+ // A (DROO) starts at 0, duration 3s (ends at 3s)
+ unit.getStockTicks().append(new StockTick("DROO", 3000));
+ // B (ACME) starts at 5s -> |start(B) - end(A)| = 2s > 1s dev
+ clock.advanceTime(5, TimeUnit.SECONDS);
+ unit.getStockTicks().append(new StockTick("ACME", 3000));
+ instance.fire();
+
+ assertThat(unit.getResults()).isEmpty();
+ }
+ }
+
+ // --- Temporal constraint: starts / startedby ---
+
+ @Test
+ public void startsMatchesWithinDeviation() {
+ StreamStartsUnit unit = new StreamStartsUnit();
+ RuleConfig config = RuleUnitProvider.get().newRuleConfig();
+ config.setClockType(ClockType.PSEUDO);
+ try (RuleUnitInstance<StreamStartsUnit> instance =
RuleUnitProvider.get().createRuleUnitInstance(unit, config)) {
+ SessionPseudoClock clock = instance.getClock();
+ // A (DROO) starts at 0, duration 5s (ends at 5s)
+ unit.getStockTicks().append(new StockTick("DROO", 5000));
+ // B (ACME) starts at 500ms (diff 500ms <= 1s dev), duration 3s
(ends at 3.5s < 5s)
+ clock.advanceTime(500, TimeUnit.MILLISECONDS);
+ unit.getStockTicks().append(new StockTick("ACME", 3000));
+ instance.fire();
+
+ assertThat(unit.getResults()).hasSize(1);
+
assertThat(unit.getResults().get(0).getCompany()).isEqualTo("ACME");
+ }
+ }
+
+ @Test
+ public void startsDoesNotMatchWhenExceedingDeviation() {
+ StreamStartsUnit unit = new StreamStartsUnit();
+ RuleConfig config = RuleUnitProvider.get().newRuleConfig();
+ config.setClockType(ClockType.PSEUDO);
+ try (RuleUnitInstance<StreamStartsUnit> instance =
RuleUnitProvider.get().createRuleUnitInstance(unit, config)) {
+ SessionPseudoClock clock = instance.getClock();
+ // A (DROO) starts at 0, duration 5s (ends at 5s)
+ unit.getStockTicks().append(new StockTick("DROO", 5000));
+ // B (ACME) starts at 1.5s (diff 1.5s > 1s dev)
+ clock.advanceTime(1500, TimeUnit.MILLISECONDS);
+ unit.getStockTicks().append(new StockTick("ACME", 2000));
+ instance.fire();
+
+ assertThat(unit.getResults()).isEmpty();
+ }
+ }
+
+ @Test
+ public void startedbyMatchesWithinDeviation() {
+ StreamStartedbyUnit unit = new StreamStartedbyUnit();
+ RuleConfig config = RuleUnitProvider.get().newRuleConfig();
+ config.setClockType(ClockType.PSEUDO);
+ try (RuleUnitInstance<StreamStartedbyUnit> instance =
RuleUnitProvider.get().createRuleUnitInstance(unit, config)) {
+ SessionPseudoClock clock = instance.getClock();
+ // A (DROO) starts at 0, duration 3s (ends at 3s)
+ unit.getStockTicks().append(new StockTick("DROO", 3000));
+ // B (ACME) starts at 500ms (diff 500ms <= 1s dev), duration 5s
(ends at 5.5s > 3s)
+ clock.advanceTime(500, TimeUnit.MILLISECONDS);
+ unit.getStockTicks().append(new StockTick("ACME", 5000));
+ instance.fire();
+
+ assertThat(unit.getResults()).hasSize(1);
+
assertThat(unit.getResults().get(0).getCompany()).isEqualTo("ACME");
+ }
+ }
+
+ @Test
+ public void startedbyDoesNotMatchWhenExceedingDeviation() {
+ StreamStartedbyUnit unit = new StreamStartedbyUnit();
+ RuleConfig config = RuleUnitProvider.get().newRuleConfig();
+ config.setClockType(ClockType.PSEUDO);
+ try (RuleUnitInstance<StreamStartedbyUnit> instance =
RuleUnitProvider.get().createRuleUnitInstance(unit, config)) {
+ SessionPseudoClock clock = instance.getClock();
+ // A (DROO) starts at 0, duration 3s (ends at 3s)
+ unit.getStockTicks().append(new StockTick("DROO", 3000));
+ // B (ACME) starts at 1.5s (diff 1.5s > 1s dev)
+ clock.advanceTime(1500, TimeUnit.MILLISECONDS);
+ unit.getStockTicks().append(new StockTick("ACME", 5000));
+ instance.fire();
+
+ assertThat(unit.getResults()).isEmpty();
+ }
+ }
+
+ // --- Temporal constraint: finishes / finishedby ---
+
+ @Test
+ public void finishesMatchesWithinDeviation() {
+ StreamFinishesUnit unit = new StreamFinishesUnit();
+ RuleConfig config = RuleUnitProvider.get().newRuleConfig();
+ config.setClockType(ClockType.PSEUDO);
+ try (RuleUnitInstance<StreamFinishesUnit> instance =
RuleUnitProvider.get().createRuleUnitInstance(unit, config)) {
+ SessionPseudoClock clock = instance.getClock();
+ // A (DROO) starts at 0, duration 5s (ends at 5s)
+ unit.getStockTicks().append(new StockTick("DROO", 5000));
+ // B (ACME) starts at 2s, duration 3.5s (ends at 5.5s, |end(B) -
end(A)| = 500ms <= 1s dev)
+ clock.advanceTime(2, TimeUnit.SECONDS);
+ unit.getStockTicks().append(new StockTick("ACME", 3500));
+ instance.fire();
+
+ assertThat(unit.getResults()).hasSize(1);
+
assertThat(unit.getResults().get(0).getCompany()).isEqualTo("ACME");
+ }
+ }
+
+ @Test
+ public void finishesDoesNotMatchWhenExceedingDeviation() {
+ StreamFinishesUnit unit = new StreamFinishesUnit();
+ RuleConfig config = RuleUnitProvider.get().newRuleConfig();
+ config.setClockType(ClockType.PSEUDO);
+ try (RuleUnitInstance<StreamFinishesUnit> instance =
RuleUnitProvider.get().createRuleUnitInstance(unit, config)) {
+ SessionPseudoClock clock = instance.getClock();
+ // A (DROO) starts at 0, duration 5s (ends at 5s)
+ unit.getStockTicks().append(new StockTick("DROO", 5000));
+ // B (ACME) starts at 2s, duration 5s (ends at 7s, |end(B) -
end(A)| = 2s > 1s dev)
+ clock.advanceTime(2, TimeUnit.SECONDS);
+ unit.getStockTicks().append(new StockTick("ACME", 5000));
+ instance.fire();
+
+ assertThat(unit.getResults()).isEmpty();
+ }
+ }
+
+ @Test
+ public void finishedbyMatchesWithinDeviation() {
+ StreamFinishedbyUnit unit = new StreamFinishedbyUnit();
+ RuleConfig config = RuleUnitProvider.get().newRuleConfig();
+ config.setClockType(ClockType.PSEUDO);
+ try (RuleUnitInstance<StreamFinishedbyUnit> instance =
RuleUnitProvider.get().createRuleUnitInstance(unit, config)) {
+ SessionPseudoClock clock = instance.getClock();
+ // B (ACME) starts at 0, duration 5s (ends at 5s)
+ unit.getStockTicks().append(new StockTick("ACME", 5000));
+ // A (DROO) starts at 2s, duration 3.5s (ends at 5.5s, |end(B) -
end(A)| = 500ms <= 1s dev)
+ clock.advanceTime(2, TimeUnit.SECONDS);
+ unit.getStockTicks().append(new StockTick("DROO", 3500));
+ instance.fire();
+
+ assertThat(unit.getResults()).hasSize(1);
+
assertThat(unit.getResults().get(0).getCompany()).isEqualTo("ACME");
+ }
+ }
+
+ @Test
+ public void finishedbyDoesNotMatchWhenExceedingDeviation() {
+ StreamFinishedbyUnit unit = new StreamFinishedbyUnit();
+ RuleConfig config = RuleUnitProvider.get().newRuleConfig();
+ config.setClockType(ClockType.PSEUDO);
+ try (RuleUnitInstance<StreamFinishedbyUnit> instance =
RuleUnitProvider.get().createRuleUnitInstance(unit, config)) {
+ SessionPseudoClock clock = instance.getClock();
+ // B (ACME) starts at 0, duration 5s (ends at 5s)
+ unit.getStockTicks().append(new StockTick("ACME", 5000));
+ // A (DROO) starts at 2s, duration 5s (ends at 7s, |end(B) -
end(A)| = 2s > 1s dev)
+ clock.advanceTime(2, TimeUnit.SECONDS);
+ unit.getStockTicks().append(new StockTick("DROO", 5000));
+ instance.fire();
+
+ assertThat(unit.getResults()).isEmpty();
+ }
+ }
+
+ // --- Explicit expiration ---
+
+ @Test
+ public void expiresRemovesEventAfterDuration() {
+ StreamExpirationUnit unit = new StreamExpirationUnit();
+ RuleConfig config = RuleUnitProvider.get().newRuleConfig();
+ config.setClockType(ClockType.PSEUDO);
+ try (RuleUnitInstance<StreamExpirationUnit> instance =
RuleUnitProvider.get().createRuleUnitInstance(unit, config)) {
+ SessionPseudoClock clock = instance.getClock();
+ unit.getStockTicks().append(new StockTick("DROO"));
+ instance.fire();
+
+ EntryPoint entryPoint = findNonDefaultEntryPoint(instance);
+ assertThat(entryPoint.getFactHandles()).hasSize(1);
+
+ clock.advanceTime(11, TimeUnit.SECONDS);
+ instance.fire();
+
+ assertThat(entryPoint.getFactHandles()).isEmpty();
+ }
+ }
+
+ // --- Combined CEP scenario ---
+
+ @Test
+ public void combinedCepScenario() {
+ StreamAfterUnit unit = new StreamAfterUnit();
+ RuleConfig config = RuleUnitProvider.get().newRuleConfig();
+ config.setClockType(ClockType.PSEUDO);
+ try (RuleUnitInstance<StreamAfterUnit> instance =
RuleUnitProvider.get().createRuleUnitInstance(unit, config)) {
+ SessionPseudoClock clock = instance.getClock();
+
+ // 1. Append a zero-duration DROO event at time zero
+ unit.getStockTicks().append(new StockTick("DROO"));
+
+ // 2. Advance 6 seconds and append ACME
+ clock.advanceTime(6, TimeUnit.SECONDS);
+ unit.getStockTicks().append(new StockTick("ACME"));
+
+ // 3. Fire and assert the after(5, 8, SECONDS) rule records ACME
+ instance.fire();
+ assertThat(unit.getResults()).hasSize(1);
+
assertThat(unit.getResults().get(0).getCompany()).isEqualTo("ACME");
+
+ // 4. Both events should still be in the entry point
+ EntryPoint entryPoint = findNonDefaultEntryPoint(instance);
+ assertThat(entryPoint.getFactHandles()).hasSize(2);
+
+ // 5. Advance another 5 seconds (total = 11s) and fire
+ clock.advanceTime(5, TimeUnit.SECONDS);
+ instance.fire();
+
+ // 6. DROO was inserted at t=0 with @Expires("10s"), so it should
be expired
+ // ACME was inserted at t=6 with @Expires("10s"), so it should
still be present
+ Collection<FactHandle> remaining = entryPoint.getFactHandles();
+ assertThat(remaining).hasSize(1);
+ assertThat(((StockTick) ((InternalFactHandle)
remaining.iterator().next()).getObject()).getCompany())
+ .isEqualTo("ACME");
+ }
+ }
+
+ // --- Unsupported aggregate calls ---
+
+ @Test
+ public void afterOnGroupByThrowsUnsupported() {
+ assertThatThrownBy(() -> {
+ new GroupByTemporalUnit();
+ }).isInstanceOf(UnsupportedOperationException.class)
+ .hasMessageContaining("groupBy");
+ }
+
+ @Test
+ public void afterOnAccumulateThrowsUnsupported() {
+ assertThatThrownBy(() -> {
+ new AccumulateTemporalUnit();
+ }).isInstanceOf(UnsupportedOperationException.class)
+ .hasMessageContaining("accumulate");
+ }
+
+ // --- Helpers ---
+
+ @SuppressWarnings("rawtypes")
+ private ReteEvaluator getEvaluator(RuleUnitInstance<?> instance) {
+ return (ReteEvaluator) ((AbstractRuleUnitInstance)
instance).getEvaluator();
+ }
+
+ private EntryPoint findNonDefaultEntryPoint(RuleUnitInstance<?> instance) {
+ ReteEvaluator evaluator = getEvaluator(instance);
+ for (EntryPoint ep : evaluator.getEntryPoints()) {
+ if (!"DEFAULT".equals(ep.getEntryPointId())) {
+ return ep;
+ }
+ }
+ throw new IllegalStateException("No non-DEFAULT entry point found");
+ }
+
+ // --- Unit definitions ---
+
+ @EventProcessing(EventProcessingType.STREAM)
+ public static class StreamAfterUnit implements RuleUnitDefinition {
+ private final DataStream<StockTick> stockTicks;
+ private final List<StockTick> results = new ArrayList<>();
+
+ public StreamAfterUnit() {
+ this.stockTicks = DataSource.createStream();
+ }
+
+ public DataStream<StockTick> getStockTicks() { return stockTicks; }
+ public List<StockTick> getResults() { return results; }
+
+ @Override
public void defineRules(RulesFactory rulesFactory) {
// A=$a : /stockTicks [ company == "DROO" ]
- // B=$b : /stockTicks [ company == "ACME", this before[5s,8s] $a ]
- // → B before A: gap = start(A) - end(B) must be in [5s, 8s]
- rulesFactory.rule("ACME before DROO")
+ // B=$b : /stockTicks [ company == "ACME", this after[5s,8s] $a ]
+ // → B after A: gap = start(B) - end(A) must be in [5s, 8s]
+ rulesFactory.rule("ACME after DROO")
.on(stockTicks) //
pattern A
.filter(StockTick::getCompany, EQUAL, "DROO")
.join(rule -> rule.on(stockTicks) //
pattern B
.filter(StockTick::getCompany, EQUAL, "ACME"))
- .before(5, 8, TimeUnit.SECONDS) // B
before A
+ .after(5, 8, TimeUnit.SECONDS) // B
after A
+ .execute(results, (r, droo, acme) -> r.add(acme));
+ }
+ }
+
+ @EventProcessing(EventProcessingType.STREAM)
+ public static class StreamAfterMillisUnit implements RuleUnitDefinition {
+ private final DataStream<StockTick> stockTicks;
+ private final List<StockTick> results = new ArrayList<>();
+
+ public StreamAfterMillisUnit() {
+ this.stockTicks = DataSource.createStream();
+ }
+
+ public DataStream<StockTick> getStockTicks() { return stockTicks; }
+ public List<StockTick> getResults() { return results; }
+
+ @Override
+ public void defineRules(RulesFactory rulesFactory) {
+ // Same as StreamAfterUnit but with millisecond bounds [500ms,
1000ms]
+ rulesFactory.rule("ACME after DROO millis")
+ .on(stockTicks) //
pattern A
+ .filter(StockTick::getCompany, EQUAL, "DROO")
+ .join(rule -> rule.on(stockTicks) //
pattern B
+ .filter(StockTick::getCompany, EQUAL, "ACME"))
+ .after(500, 1000, TimeUnit.MILLISECONDS) // B
after A
+ .execute(results, (r, droo, acme) -> r.add(acme));
+ }
+ }
+
+ @EventProcessing(EventProcessingType.STREAM)
+ public static class StreamAfterNoArgUnit implements RuleUnitDefinition {
+ private final DataStream<StockTick> stockTicks;
+ private final List<StockTick> results = new ArrayList<>();
+
+ public StreamAfterNoArgUnit() {
+ this.stockTicks = DataSource.createStream();
+ }
+
+ public DataStream<StockTick> getStockTicks() { return stockTicks; }
+ public List<StockTick> getResults() { return results; }
+
+ @Override
+ public void defineRules(RulesFactory rulesFactory) {
+ rulesFactory.rule("ACME after DROO no-arg")
+ .on(stockTicks)
+ .filter(StockTick::getCompany, EQUAL, "DROO")
+ .join(rule -> rule.on(stockTicks)
+ .filter(StockTick::getCompany, EQUAL, "ACME"))
+ .after()
+ .execute(results, (r, droo, acme) -> r.add(acme));
+ }
+ }
+
+ @EventProcessing(EventProcessingType.STREAM)
+ public static class StreamAfterSingleBoundUnit implements
RuleUnitDefinition {
+ private final DataStream<StockTick> stockTicks;
+ private final List<StockTick> results = new ArrayList<>();
+
+ public StreamAfterSingleBoundUnit() {
+ this.stockTicks = DataSource.createStream();
+ }
+
+ public DataStream<StockTick> getStockTicks() { return stockTicks; }
+ public List<StockTick> getResults() { return results; }
+
+ @Override
+ public void defineRules(RulesFactory rulesFactory) {
+ rulesFactory.rule("ACME after DROO single bound")
+ .on(stockTicks)
+ .filter(StockTick::getCompany, EQUAL, "DROO")
+ .join(rule -> rule.on(stockTicks)
+ .filter(StockTick::getCompany, EQUAL, "ACME"))
+ .after(4, TimeUnit.SECONDS)
+ .execute(results, (r, droo, acme) -> r.add(acme));
+ }
+ }
+
+ @EventProcessing(EventProcessingType.STREAM)
+ public static class StreamBeforeUnit implements RuleUnitDefinition {
+ private final DataStream<StockTick> stockTicks;
+ private final List<StockTick> results = new ArrayList<>();
+
+ public StreamBeforeUnit() {
+ this.stockTicks = DataSource.createStream();
+ }
+
+ public DataStream<StockTick> getStockTicks() { return stockTicks; }
+ public List<StockTick> getResults() { return results; }
+
+ @Override
+ public void defineRules(RulesFactory rulesFactory) {
+ // A=$a : /stockTicks [ company == "DROO" ]
+ // B=$b : /stockTicks [ company == "ACME", this before[5s,8s] $a ]
+ // → B before A: gap = start(A) - end(B) must be in [5s, 8s]
+ rulesFactory.rule("ACME before DROO")
+ .on(stockTicks) //
pattern A
+ .filter(StockTick::getCompany, EQUAL, "DROO")
+ .join(rule -> rule.on(stockTicks) //
pattern B
+ .filter(StockTick::getCompany, EQUAL, "ACME"))
+ .before(5, 8, TimeUnit.SECONDS) // B
before A
+ .execute(results, (r, droo, acme) -> r.add(acme));
+ }
+ }
+
+ @EventProcessing(EventProcessingType.STREAM)
+ public static class StreamBeforeNoArgUnit implements RuleUnitDefinition {
+ private final DataStream<StockTick> stockTicks;
+ private final List<StockTick> results = new ArrayList<>();
+
+ public StreamBeforeNoArgUnit() {
+ this.stockTicks = DataSource.createStream();
+ }
+
+ public DataStream<StockTick> getStockTicks() { return stockTicks; }
+ public List<StockTick> getResults() { return results; }
+
+ @Override
+ public void defineRules(RulesFactory rulesFactory) {
+ rulesFactory.rule("ACME before DROO no-arg")
+ .on(stockTicks)
+ .filter(StockTick::getCompany, EQUAL, "DROO")
+ .join(rule -> rule.on(stockTicks)
+ .filter(StockTick::getCompany, EQUAL, "ACME"))
+ .before()
+ .execute(results, (r, droo, acme) -> r.add(acme));
+ }
+ }
+
+ @EventProcessing(EventProcessingType.STREAM)
+ public static class StreamBeforeSingleBoundUnit implements
RuleUnitDefinition {
+ private final DataStream<StockTick> stockTicks;
+ private final List<StockTick> results = new ArrayList<>();
+
+ public StreamBeforeSingleBoundUnit() {
+ this.stockTicks = DataSource.createStream();
+ }
+
+ public DataStream<StockTick> getStockTicks() { return stockTicks; }
+ public List<StockTick> getResults() { return results; }
+
+ @Override
+ public void defineRules(RulesFactory rulesFactory) {
+ rulesFactory.rule("ACME before DROO single bound")
+ .on(stockTicks)
+ .filter(StockTick::getCompany, EQUAL, "DROO")
+ .join(rule -> rule.on(stockTicks)
+ .filter(StockTick::getCompany, EQUAL, "ACME"))
+ .before(4, TimeUnit.SECONDS)
+ .execute(results, (r, droo, acme) -> r.add(acme));
+ }
+ }
+
+ @EventProcessing(EventProcessingType.STREAM)
+ public static class StreamCoincidesUnit implements RuleUnitDefinition {
+ private final DataStream<StockTick> stockTicks;
+ private final List<StockTick> results = new ArrayList<>();
+
+ public StreamCoincidesUnit() {
+ this.stockTicks = DataSource.createStream();
+ }
+
+ public DataStream<StockTick> getStockTicks() { return stockTicks; }
+ public List<StockTick> getResults() { return results; }
+
+ @Override
+ public void defineRules(RulesFactory rulesFactory) {
+ rulesFactory.rule("ACME coincides DROO")
+ .on(stockTicks)
+ .filter(StockTick::getCompany, EQUAL, "DROO")
+ .join(rule -> rule.on(stockTicks)
+ .filter(StockTick::getCompany, EQUAL, "ACME"))
+ .coincides()
+ .execute(results, (r, droo, acme) -> r.add(acme));
+ }
+ }
+
+ @EventProcessing(EventProcessingType.STREAM)
+ public static class StreamCoincidesDevUnit implements RuleUnitDefinition {
+ private final DataStream<StockTick> stockTicks;
+ private final List<StockTick> results = new ArrayList<>();
+
+ public StreamCoincidesDevUnit() {
+ this.stockTicks = DataSource.createStream();
+ }
+
+ public DataStream<StockTick> getStockTicks() { return stockTicks; }
+ public List<StockTick> getResults() { return results; }
+
+ @Override
+ public void defineRules(RulesFactory rulesFactory) {
+ rulesFactory.rule("ACME coincides DROO dev")
+ .on(stockTicks)
+ .filter(StockTick::getCompany, EQUAL, "DROO")
+ .join(rule -> rule.on(stockTicks)
+ .filter(StockTick::getCompany, EQUAL, "ACME"))
+ .coincides(1, TimeUnit.SECONDS)
+ .execute(results, (r, droo, acme) -> r.add(acme));
+ }
+ }
+
+ @EventProcessing(EventProcessingType.STREAM)
+ public static class StreamCoincidesStartEndDevUnit implements
RuleUnitDefinition {
+ private final DataStream<StockTick> stockTicks;
+ private final List<StockTick> results = new ArrayList<>();
+
+ public StreamCoincidesStartEndDevUnit() {
+ this.stockTicks = DataSource.createStream();
+ }
+
+ public DataStream<StockTick> getStockTicks() { return stockTicks; }
+ public List<StockTick> getResults() { return results; }
+
+ @Override
+ public void defineRules(RulesFactory rulesFactory) {
+ rulesFactory.rule("ACME coincides DROO start/end dev")
+ .on(stockTicks)
+ .filter(StockTick::getCompany, EQUAL, "DROO")
+ .join(rule -> rule.on(stockTicks)
+ .filter(StockTick::getCompany, EQUAL, "ACME"))
+ .coincides(1, TimeUnit.SECONDS, 2, TimeUnit.SECONDS)
+ .execute(results, (r, droo, acme) -> r.add(acme));
+ }
+ }
+
+ @EventProcessing(EventProcessingType.STREAM)
+ public static class StreamDuringUnit implements RuleUnitDefinition {
+ private final DataStream<StockTick> stockTicks;
+ private final List<StockTick> results = new ArrayList<>();
+
+ public StreamDuringUnit() {
+ this.stockTicks = DataSource.createStream();
+ }
+
+ public DataStream<StockTick> getStockTicks() { return stockTicks; }
+ public List<StockTick> getResults() { return results; }
+
+ @Override
+ public void defineRules(RulesFactory rulesFactory) {
+ rulesFactory.rule("ACME during DROO")
+ .on(stockTicks)
+ .filter(StockTick::getCompany, EQUAL, "DROO")
+ .join(rule -> rule.on(stockTicks)
+ .filter(StockTick::getCompany, EQUAL, "ACME"))
+ .during(1, TimeUnit.SECONDS, 10, TimeUnit.SECONDS)
+ .execute(results, (r, droo, acme) -> r.add(acme));
+ }
+ }
+
+ @EventProcessing(EventProcessingType.STREAM)
+ public static class StreamIncludesUnit implements RuleUnitDefinition {
+ private final DataStream<StockTick> stockTicks;
+ private final List<StockTick> results = new ArrayList<>();
+
+ public StreamIncludesUnit() {
+ this.stockTicks = DataSource.createStream();
+ }
+
+ public DataStream<StockTick> getStockTicks() { return stockTicks; }
+ public List<StockTick> getResults() { return results; }
+
+ @Override
+ public void defineRules(RulesFactory rulesFactory) {
+ rulesFactory.rule("ACME includes DROO")
+ .on(stockTicks)
+ .filter(StockTick::getCompany, EQUAL, "DROO")
+ .join(rule -> rule.on(stockTicks)
+ .filter(StockTick::getCompany, EQUAL, "ACME"))
+ .includes(1, TimeUnit.SECONDS, 10, TimeUnit.SECONDS)
+ .execute(results, (r, droo, acme) -> r.add(acme));
+ }
+ }
+
+ @EventProcessing(EventProcessingType.STREAM)
+ public static class StreamOverlapsUnit implements RuleUnitDefinition {
+ private final DataStream<StockTick> stockTicks;
+ private final List<StockTick> results = new ArrayList<>();
+
+ public StreamOverlapsUnit() {
+ this.stockTicks = DataSource.createStream();
+ }
+
+ public DataStream<StockTick> getStockTicks() { return stockTicks; }
+ public List<StockTick> getResults() { return results; }
+
+ @Override
+ public void defineRules(RulesFactory rulesFactory) {
+ rulesFactory.rule("ACME overlaps DROO")
+ .on(stockTicks)
+ .filter(StockTick::getCompany, EQUAL, "DROO")
+ .join(rule -> rule.on(stockTicks)
+ .filter(StockTick::getCompany, EQUAL, "ACME"))
+ .overlaps(1, TimeUnit.SECONDS, 5, TimeUnit.SECONDS)
+ .execute(results, (r, droo, acme) -> r.add(acme));
+ }
+ }
+
+ @EventProcessing(EventProcessingType.STREAM)
+ public static class StreamOverlappedbyUnit implements RuleUnitDefinition {
+ private final DataStream<StockTick> stockTicks;
+ private final List<StockTick> results = new ArrayList<>();
+
+ public StreamOverlappedbyUnit() {
+ this.stockTicks = DataSource.createStream();
+ }
+
+ public DataStream<StockTick> getStockTicks() { return stockTicks; }
+ public List<StockTick> getResults() { return results; }
+
+ @Override
+ public void defineRules(RulesFactory rulesFactory) {
+ rulesFactory.rule("ACME overlappedby DROO")
+ .on(stockTicks)
+ .filter(StockTick::getCompany, EQUAL, "DROO")
+ .join(rule -> rule.on(stockTicks)
+ .filter(StockTick::getCompany, EQUAL, "ACME"))
+ .overlappedby(1, TimeUnit.SECONDS, 5, TimeUnit.SECONDS)
+ .execute(results, (r, droo, acme) -> r.add(acme));
+ }
+ }
+
+ @EventProcessing(EventProcessingType.STREAM)
+ public static class StreamMeetsUnit implements RuleUnitDefinition {
+ private final DataStream<StockTick> stockTicks;
+ private final List<StockTick> results = new ArrayList<>();
+
+ public StreamMeetsUnit() {
+ this.stockTicks = DataSource.createStream();
+ }
+
+ public DataStream<StockTick> getStockTicks() { return stockTicks; }
+ public List<StockTick> getResults() { return results; }
+
+ @Override
+ public void defineRules(RulesFactory rulesFactory) {
+ rulesFactory.rule("ACME meets DROO")
+ .on(stockTicks)
+ .filter(StockTick::getCompany, EQUAL, "DROO")
+ .join(rule -> rule.on(stockTicks)
+ .filter(StockTick::getCompany, EQUAL, "ACME"))
+ .meets(1, TimeUnit.SECONDS)
+ .execute(results, (r, droo, acme) -> r.add(acme));
+ }
+ }
+
+ @EventProcessing(EventProcessingType.STREAM)
+ public static class StreamMetbyUnit implements RuleUnitDefinition {
+ private final DataStream<StockTick> stockTicks;
+ private final List<StockTick> results = new ArrayList<>();
+
+ public StreamMetbyUnit() {
+ this.stockTicks = DataSource.createStream();
+ }
+
+ public DataStream<StockTick> getStockTicks() { return stockTicks; }
+ public List<StockTick> getResults() { return results; }
+
+ @Override
+ public void defineRules(RulesFactory rulesFactory) {
+ rulesFactory.rule("ACME metby DROO")
+ .on(stockTicks)
+ .filter(StockTick::getCompany, EQUAL, "DROO")
+ .join(rule -> rule.on(stockTicks)
+ .filter(StockTick::getCompany, EQUAL, "ACME"))
+ .metby(1, TimeUnit.SECONDS)
+ .execute(results, (r, droo, acme) -> r.add(acme));
+ }
+ }
+
+ @EventProcessing(EventProcessingType.STREAM)
+ public static class StreamStartsUnit implements RuleUnitDefinition {
+ private final DataStream<StockTick> stockTicks;
+ private final List<StockTick> results = new ArrayList<>();
+
+ public StreamStartsUnit() {
+ this.stockTicks = DataSource.createStream();
+ }
+
+ public DataStream<StockTick> getStockTicks() { return stockTicks; }
+ public List<StockTick> getResults() { return results; }
+
+ @Override
+ public void defineRules(RulesFactory rulesFactory) {
+ rulesFactory.rule("ACME starts DROO")
+ .on(stockTicks)
+ .filter(StockTick::getCompany, EQUAL, "DROO")
+ .join(rule -> rule.on(stockTicks)
+ .filter(StockTick::getCompany, EQUAL, "ACME"))
+ .starts(1, TimeUnit.SECONDS)
+ .execute(results, (r, droo, acme) -> r.add(acme));
+ }
+ }
+
+ @EventProcessing(EventProcessingType.STREAM)
+ public static class StreamStartedbyUnit implements RuleUnitDefinition {
+ private final DataStream<StockTick> stockTicks;
+ private final List<StockTick> results = new ArrayList<>();
+
+ public StreamStartedbyUnit() {
+ this.stockTicks = DataSource.createStream();
+ }
+
+ public DataStream<StockTick> getStockTicks() { return stockTicks; }
+ public List<StockTick> getResults() { return results; }
+
+ @Override
+ public void defineRules(RulesFactory rulesFactory) {
+ rulesFactory.rule("ACME startedby DROO")
+ .on(stockTicks)
+ .filter(StockTick::getCompany, EQUAL, "DROO")
+ .join(rule -> rule.on(stockTicks)
+ .filter(StockTick::getCompany, EQUAL, "ACME"))
+ .startedby(1, TimeUnit.SECONDS)
+ .execute(results, (r, droo, acme) -> r.add(acme));
+ }
+ }
+
+ @EventProcessing(EventProcessingType.STREAM)
+ public static class StreamFinishesUnit implements RuleUnitDefinition {
+ private final DataStream<StockTick> stockTicks;
+ private final List<StockTick> results = new ArrayList<>();
+
+ public StreamFinishesUnit() {
+ this.stockTicks = DataSource.createStream();
+ }
+
+ public DataStream<StockTick> getStockTicks() { return stockTicks; }
+ public List<StockTick> getResults() { return results; }
+
+ @Override
+ public void defineRules(RulesFactory rulesFactory) {
+ rulesFactory.rule("ACME finishes DROO")
+ .on(stockTicks)
+ .filter(StockTick::getCompany, EQUAL, "DROO")
+ .join(rule -> rule.on(stockTicks)
+ .filter(StockTick::getCompany, EQUAL, "ACME"))
+ .finishes(1, TimeUnit.SECONDS)
+ .execute(results, (r, droo, acme) -> r.add(acme));
+ }
+ }
+
+ @EventProcessing(EventProcessingType.STREAM)
+ public static class StreamFinishedbyUnit implements RuleUnitDefinition {
+ private final DataStream<StockTick> stockTicks;
+ private final List<StockTick> results = new ArrayList<>();
+
+ public StreamFinishedbyUnit() {
+ this.stockTicks = DataSource.createStream();
+ }
+
+ public DataStream<StockTick> getStockTicks() { return stockTicks; }
+ public List<StockTick> getResults() { return results; }
+
+ @Override
+ public void defineRules(RulesFactory rulesFactory) {
+ rulesFactory.rule("ACME finishedby DROO")
+ .on(stockTicks)
+ .filter(StockTick::getCompany, EQUAL, "DROO")
+ .join(rule -> rule.on(stockTicks)
+ .filter(StockTick::getCompany, EQUAL, "ACME"))
+ .finishedby(1, TimeUnit.SECONDS)
.execute(results, (r, droo, acme) -> r.add(acme));
}
}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]