Skip to content

Commit 33cc1d7

Browse files
fix: aggregate enrollment stage CTE projection (#23868)
* fix: aggregate enrollment stage CTE projection
1 parent fefefbb commit 33cc1d7

5 files changed

Lines changed: 449 additions & 19 deletions

File tree

dhis-2/dhis-services/dhis-service-analytics/src/main/java/org/hisp/dhis/analytics/event/data/AbstractJdbcEventAnalyticsManager.java

Lines changed: 33 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -282,6 +282,8 @@ public static AggregateClause of(String sql, AggregationType type, String innerE
282282

283283
static final String ANALYTICS_EVENT = "analytics_event_";
284284

285+
public static final String LATEST_EVENTS_CTE_PREFIX = "latest_events_";
286+
285287
static final String COLUMN_ENROLLMENT_GEOMETRY_GEOJSON =
286288
String.format(
287289
"ST_AsGeoJSON(%s)", EnrollmentAnalyticsColumnName.ENROLLMENT_GEOMETRY_COLUMN_NAME);
@@ -2560,17 +2562,17 @@ protected String createDefaultAlias(String formula) {
25602562
void generateFilterCTEs(
25612563
EventQueryParams params, CteContext cteContext, boolean isAggregateQuery) {
25622564

2565+
if (isAggregateQuery) {
2566+
generateAggregateFilterCTEs(params, cteContext);
2567+
return;
2568+
}
2569+
25632570
// Combine items and item filters and filter only those with an actual filter
25642571
List<QueryItem> queryItems =
25652572
Stream.concat(params.getItems().stream(), params.getItemFilters().stream())
25662573
.filter(QueryItem::hasFilter)
25672574
.toList();
25682575

2569-
if (isAggregateQuery) {
2570-
generateAggregateFilterCTEs(queryItems, params, cteContext);
2571-
return;
2572-
}
2573-
25742576
// Group query items by repeatable and non-repeatable stages
25752577
Map<RepeatableStateStatus, List<QueryItem>> itemsByRepeatableFlag =
25762578
queryItems.stream()
@@ -2607,27 +2609,31 @@ void generateFilterCTEs(
26072609
}
26082610

26092611
/**
2610-
* Generates filter CTEs for aggregate enrollment queries. Items are grouped by program stage UID,
2611-
* producing one CTE per stage with all dimension columns and filter conditions combined.
2612+
* Generates filter CTEs for aggregate enrollment queries. Stage items (both dimensions and item
2613+
* filters) are grouped by program stage UID, producing one CTE per stage that projects eligible
2614+
* stage dimensions and applies the combined filter conditions.
2615+
*
2616+
* <p>A {@code latest_events_<stage>} CTE is emitted only when at least one item for the stage
2617+
* carries a filter. Unfiltered stage items are still projected by the CTE when they can be read
2618+
* from the same "latest event" row without losing repeatable-stage offset semantics.
26122619
*
2613-
* @param queryItems filtered query items that have at least one filter
26142620
* @param params the {@link EventQueryParams} object
26152621
* @param cteContext the {@link CteContext} to register CTEs into
26162622
*/
2617-
private void generateAggregateFilterCTEs(
2618-
List<QueryItem> queryItems, EventQueryParams params, CteContext cteContext) {
2619-
2620-
// Collect all items that have a program stage and group by stage UID
2623+
private void generateAggregateFilterCTEs(EventQueryParams params, CteContext cteContext) {
26212624
Map<String, List<QueryItem>> itemsByStage =
2622-
queryItems.stream()
2625+
Stream.concat(params.getItems().stream(), params.getItemFilters().stream())
26232626
.filter(qi -> qi.hasProgram() && qi.hasProgramStage())
2627+
.filter(qi -> qi.hasFilter() || canProjectUnfilteredItemInAggregateFilterCte(qi))
26242628
.collect(groupingBy(qi -> qi.getProgramStage().getUid(), LinkedHashMap::new, toList()));
26252629

2626-
// For each stage, build a single CTE with all dimension columns and filters
26272630
itemsByStage.forEach(
26282631
(stageUid, stageItems) -> {
2632+
if (stageItems.stream().noneMatch(QueryItem::hasFilter)) {
2633+
return;
2634+
}
26292635
String cteSql = buildAggregateFilterCteSql(stageItems, params);
2630-
String cteKey = "latest_events_" + stageUid;
2636+
String cteKey = LATEST_EVENTS_CTE_PREFIX + stageUid;
26312637
cteContext.addCteFilter(cteKey, stageItems.get(0), cteSql);
26322638
});
26332639
}
@@ -2791,6 +2797,14 @@ private void buildProgramStageCte(
27912797
if (item.hasFilter()) {
27922798
return;
27932799
}
2800+
// The per-stage filter CTE projects non-offset dimensions for that stage,
2801+
// so a redundant per-item CTE would only inflate the query without adding data.
2802+
if (canProjectUnfilteredItemInAggregateFilterCte(item)
2803+
&& cteContext.getDefinitionByItemUid(
2804+
LATEST_EVENTS_CTE_PREFIX + item.getProgramStage().getUid())
2805+
!= null) {
2806+
return;
2807+
}
27942808
handleAggregatedEnrollments(cteContext, item, eventTableName, colName, params);
27952809
return;
27962810
}
@@ -2840,6 +2854,10 @@ protected boolean columnIsInFormula(String col) {
28402854
return col.contains("(") && col.contains(")");
28412855
}
28422856

2857+
private boolean canProjectUnfilteredItemInAggregateFilterCte(QueryItem item) {
2858+
return !item.hasRepeatableStageParams() || item.getRepeatableStageParams().isDefaultObject();
2859+
}
2860+
28432861
/**
28442862
* Computes a zero-based offset for use with the SQL <em>row_number()</em> function in CTEs that
28452863
* partition and order events by date (e.g., most recent first).

dhis-2/dhis-services/dhis-service-analytics/src/main/java/org/hisp/dhis/analytics/event/data/aggregate/AggregatedEnrollmentHeaderColumnResolver.java

Lines changed: 10 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,9 @@
2929
*/
3030
package org.hisp.dhis.analytics.event.data.aggregate;
3131

32+
import static org.hisp.dhis.analytics.event.data.AbstractJdbcEventAnalyticsManager.LATEST_EVENTS_CTE_PREFIX;
33+
import static org.hisp.dhis.analytics.util.RepeatableStageParamsHelper.removeRepeatableStageParams;
34+
3235
import java.util.Map;
3336
import java.util.Set;
3437
import java.util.function.UnaryOperator;
@@ -44,7 +47,6 @@
4447
* stage-specific dimensions via per-stage filter CTEs.
4548
*/
4649
public final class AggregatedEnrollmentHeaderColumnResolver {
47-
private static final String LATEST_EVENTS_CTE_PREFIX = "latest_events_";
4850

4951
private final StageHeaderClassifier stageHeaderClassifier;
5052

@@ -91,8 +93,8 @@ public void addHeaderColumns(
9193

9294
Map.Entry<String, CteDefinition> matchingEntry = findMatchingCte(cteDefinitionMap, colName);
9395
if (matchingEntry != null) {
94-
sb.addColumn(matchingEntry.getValue().getAlias() + ".value", "", matchingEntry.getKey());
95-
sb.groupBy(matchingEntry.getKey());
96+
sb.addColumn(matchingEntry.getValue().getAlias() + ".value", "", quotedCol);
97+
sb.groupBy(quotedCol);
9698
} else {
9799
sb.addColumn(quotedCol);
98100
sb.groupBy(quotedCol);
@@ -155,8 +157,12 @@ private String extractStageUid(String header) {
155157

156158
private Map.Entry<String, CteDefinition> findMatchingCte(
157159
Map<String, CteDefinition> cteDefinitionMap, String colName) {
160+
String normalizedColName =
161+
removeRepeatableStageParams(colName.replace("\"", "").replace("`", ""));
158162
for (Map.Entry<String, CteDefinition> entry : cteDefinitionMap.entrySet()) {
159-
if (entry.getKey().contains(colName)) {
163+
String normalizedKey =
164+
removeRepeatableStageParams(entry.getKey().replace("\"", "").replace("`", ""));
165+
if (normalizedKey.contains(normalizedColName)) {
160166
return entry;
161167
}
162168
}

dhis-2/dhis-services/dhis-service-analytics/src/test/java/org/hisp/dhis/analytics/event/data/EnrollmentAnalyticsManagerCteTest.java

Lines changed: 168 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -30,10 +30,12 @@
3030
package org.hisp.dhis.analytics.event.data;
3131

3232
import static org.hamcrest.CoreMatchers.containsString;
33+
import static org.hamcrest.CoreMatchers.is;
3334
import static org.hamcrest.CoreMatchers.not;
3435
import static org.hamcrest.MatcherAssert.assertThat;
3536
import static org.hisp.dhis.analytics.DataType.NUMERIC;
3637
import static org.hisp.dhis.analytics.QueryKey.NV;
38+
import static org.hisp.dhis.analytics.table.EventAnalyticsColumnName.EVENT_STATUS_COLUMN_NAME;
3739
import static org.hisp.dhis.analytics.table.EventAnalyticsColumnName.OCCURRED_DATE_COLUMN_NAME;
3840
import static org.hisp.dhis.analytics.table.EventAnalyticsColumnName.OU_COLUMN_NAME;
3941
import static org.hisp.dhis.common.DimensionConstants.OPTION_SEP;
@@ -78,6 +80,7 @@
7880
import org.hisp.dhis.common.QueryFilter;
7981
import org.hisp.dhis.common.QueryItem;
8082
import org.hisp.dhis.common.QueryOperator;
83+
import org.hisp.dhis.common.RepeatableStageParams;
8184
import org.hisp.dhis.common.ValueType;
8285
import org.hisp.dhis.dataelement.DataElement;
8386
import org.hisp.dhis.dataelement.DataElementService;
@@ -90,6 +93,7 @@
9093
import org.hisp.dhis.program.Program;
9194
import org.hisp.dhis.program.ProgramIndicator;
9295
import org.hisp.dhis.program.ProgramIndicatorService;
96+
import org.hisp.dhis.program.ProgramStage;
9397
import org.hisp.dhis.relationship.RelationshipConstraint;
9498
import org.hisp.dhis.relationship.RelationshipEntity;
9599
import org.hisp.dhis.relationship.RelationshipType;
@@ -411,6 +415,134 @@ void verifyAggregateEnrollmentStageOrgUnitFilterStaysInFilterCte() {
411415
assertThat(baseCteSql, not(containsString("and \"n1rtSHYf6O6\" in ('ImspTQPwCqd')")));
412416
}
413417

418+
@Test
419+
void verifyAggregateEnrollmentUnfilteredEventStatusUsesLatestEventsCte() {
420+
// When a stage-specific dimension without a filter (e.g. stage.EVENT_STATUS) is requested
421+
// alongside a filtered stage dimension (stage.EVENT_DATE), the unfiltered dimension must
422+
// also be projected by the latest_events_<stage> CTE so the outer SELECT can read it from
423+
// the same "latest event" row. Without this, the resolver emits a reference to
424+
// uiogg.ev_eventstatus while the CTE doesn't project the column, producing a SQL error.
425+
String stageUid = programStage.getUid();
426+
427+
QueryItem dateItem =
428+
new QueryItem(
429+
new BaseDimensionalItemObject(OCCURRED_DATE_COLUMN_NAME),
430+
programA,
431+
null,
432+
ValueType.DATE,
433+
null,
434+
null);
435+
dateItem.setProgram(programA);
436+
dateItem.setProgramStage(programStage);
437+
dateItem.addFilter(new QueryFilter(QueryOperator.GE, "2026-01-01"));
438+
dateItem.addFilter(new QueryFilter(QueryOperator.LE, "2026-12-31"));
439+
440+
QueryItem statusItem =
441+
new QueryItem(
442+
new BaseDimensionalItemObject(EVENT_STATUS_COLUMN_NAME),
443+
programA,
444+
null,
445+
ValueType.TEXT,
446+
null,
447+
null);
448+
statusItem.setProgram(programA);
449+
statusItem.setProgramStage(programStage);
450+
451+
EventQueryParams.Builder params = createRequestParamsBuilder();
452+
params.addItem(dateItem);
453+
params.addItem(statusItem);
454+
params.withEndpointAction(AGGREGATE);
455+
456+
ListGrid grid = new ListGrid();
457+
grid.addHeader(new GridHeader("value", "Value", ValueType.NUMBER, false, false));
458+
grid.addHeader(
459+
new GridHeader(stageUid + ".eventdate", "Event date", ValueType.DATE, false, false));
460+
grid.addHeader(
461+
new GridHeader(stageUid + ".eventstatus", "Event status", ValueType.TEXT, false, false));
462+
463+
subject.getEnrollments(params.build(), grid, 10000);
464+
verify(jdbcTemplate).queryForRowSet(sql.capture());
465+
466+
String generatedSql = noEof(sql.getValue());
467+
String latestEventsCteSql =
468+
generatedSql.substring(
469+
generatedSql.indexOf("latest_events_" + stageUid + " as ("),
470+
generatedSql.indexOf("enrollment_aggr_base as ("));
471+
472+
// The latest_events CTE must project both the filtered (eventdate) and unfiltered
473+
// (eventstatus) stage dimensions so the outer SELECT can read them from the same row.
474+
assertThat(latestEventsCteSql, containsString("ev_occurreddate"));
475+
assertThat(latestEventsCteSql, containsString("ev_eventstatus"));
476+
477+
// The redundant per-item CTE must not be created when a latest_events_<stage> already exists.
478+
assertThat(generatedSql, not(containsString(stageUid + "_" + EVENT_STATUS_COLUMN_NAME + "_0")));
479+
}
480+
481+
@Test
482+
void verifyAggregateEnrollmentRepeatableOffsetItemKeepsPerItemCteWhenLatestEventsCteExists() {
483+
String stageUid = repeatableProgramStage.getUid();
484+
String deUid = dataElementA.getUid();
485+
486+
QueryItem dateItem = createFilteredStageDateItem(repeatableProgramStage);
487+
QueryItem offsetItem = createRepeatableOffsetDataElementItem(stageUid, -1);
488+
489+
EventQueryParams.Builder params = createRequestParamsBuilder();
490+
params.addItem(dateItem);
491+
params.addItem(offsetItem);
492+
params.withEndpointAction(AGGREGATE);
493+
494+
ListGrid grid = new ListGrid();
495+
grid.addHeader(new GridHeader("value", "Value", ValueType.NUMBER, false, false));
496+
grid.addHeader(
497+
new GridHeader(stageUid + ".eventdate", "Event date", ValueType.DATE, false, false));
498+
grid.addHeader(
499+
new GridHeader(stageUid + "[-1]." + deUid, "Offset value", ValueType.NUMBER, false, false));
500+
501+
subject.getEnrollments(params.build(), grid, 10000);
502+
verify(jdbcTemplate).queryForRowSet(sql.capture());
503+
504+
String generatedSql = noEof(sql.getValue());
505+
String latestEventsCteSql =
506+
generatedSql.substring(
507+
generatedSql.indexOf("latest_events_" + stageUid + " as ("),
508+
generatedSql.indexOf("enrollment_aggr_base as ("));
509+
510+
assertThat(latestEventsCteSql, containsString("ev_occurreddate"));
511+
assertThat(latestEventsCteSql, not(containsString("ev_" + deUid)));
512+
assertThat(generatedSql, containsString(stageUid + "_" + deUid + "_0 as ("));
513+
assertThat(generatedSql, containsString("as \"" + stageUid + "[-1]." + deUid + "\""));
514+
}
515+
516+
@Test
517+
void verifyAggregateEnrollmentRepeatableOffsetsAreNotDuplicatedInLatestEventsCte() {
518+
String stageUid = repeatableProgramStage.getUid();
519+
String deUid = dataElementA.getUid();
520+
521+
EventQueryParams.Builder params = createRequestParamsBuilder();
522+
params.addItem(createFilteredStageDateItem(repeatableProgramStage));
523+
params.addItem(createRepeatableOffsetDataElementItem(stageUid, -1));
524+
params.addItem(createRepeatableOffsetDataElementItem(stageUid, 1));
525+
params.withEndpointAction(AGGREGATE);
526+
527+
ListGrid grid = new ListGrid();
528+
grid.addHeader(new GridHeader("value", "Value", ValueType.NUMBER, false, false));
529+
grid.addHeader(
530+
new GridHeader(stageUid + ".eventdate", "Event date", ValueType.DATE, false, false));
531+
532+
subject.getEnrollments(params.build(), grid, 10000);
533+
verify(jdbcTemplate).queryForRowSet(sql.capture());
534+
535+
String generatedSql = noEof(sql.getValue());
536+
String latestEventsCteSql =
537+
generatedSql.substring(
538+
generatedSql.indexOf("latest_events_" + stageUid + " as ("),
539+
generatedSql.indexOf("enrollment_aggr_base as ("));
540+
541+
assertThat(countOccurrences(latestEventsCteSql, "ev_" + deUid), is(0));
542+
assertThat(generatedSql, containsString(stageUid + "_" + deUid + "_0 as ("));
543+
assertThat(generatedSql, containsString(stageUid + "_" + deUid + "_1 as ("));
544+
}
545+
414546
@Test
415547
void verifyAggregateEnrollmentIncludesProgramStatusFilterInSql() {
416548
EventQueryParams.Builder params = createRequestParamsBuilder();
@@ -1109,6 +1241,32 @@ private EventQueryParams createAggregateEnrollmentWithStageDateParams() {
11091241
return params.build();
11101242
}
11111243

1244+
private QueryItem createFilteredStageDateItem(ProgramStage stage) {
1245+
QueryItem dateItem =
1246+
new QueryItem(
1247+
new BaseDimensionalItemObject(OCCURRED_DATE_COLUMN_NAME),
1248+
programA,
1249+
null,
1250+
ValueType.DATE,
1251+
null,
1252+
null);
1253+
dateItem.setProgram(programA);
1254+
dateItem.setProgramStage(stage);
1255+
dateItem.addFilter(new QueryFilter(QueryOperator.GE, "2026-01-01"));
1256+
dateItem.addFilter(new QueryFilter(QueryOperator.LE, "2026-12-31"));
1257+
return dateItem;
1258+
}
1259+
1260+
private QueryItem createRepeatableOffsetDataElementItem(String stageUid, int offset) {
1261+
QueryItem item = new QueryItem(new BaseDimensionalItemObject(dataElementA.getUid()));
1262+
item.setProgram(programA);
1263+
item.setProgramStage(repeatableProgramStage);
1264+
item.setValueType(ValueType.NUMBER);
1265+
item.setRepeatableStageParams(
1266+
RepeatableStageParams.of(offset, stageUid + "[" + offset + "]." + dataElementA.getUid()));
1267+
return item;
1268+
}
1269+
11121270
private EventQueryParams createAggregateEnrollmentWithEventDateParams() {
11131271
BaseDimensionalItemObject dateItem = new BaseDimensionalItemObject(OCCURRED_DATE_COLUMN_NAME);
11141272
QueryItem queryItem = new QueryItem(dateItem, programA, null, ValueType.DATE, null, null);
@@ -1203,4 +1361,14 @@ private void testIt(
12031361
private String noEof(String sql) {
12041362
return sql.replaceAll("\\s+", " ").trim();
12051363
}
1364+
1365+
private int countOccurrences(String value, String search) {
1366+
int count = 0;
1367+
int index = 0;
1368+
while ((index = value.indexOf(search, index)) >= 0) {
1369+
count++;
1370+
index += search.length();
1371+
}
1372+
return count;
1373+
}
12061374
}

0 commit comments

Comments
 (0)