Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -282,6 +282,8 @@ public static AggregateClause of(String sql, AggregationType type, String innerE

static final String ANALYTICS_EVENT = "analytics_event_";

public static final String LATEST_EVENTS_CTE_PREFIX = "latest_events_";

static final String COLUMN_ENROLLMENT_GEOMETRY_GEOJSON =
String.format(
"ST_AsGeoJSON(%s)", EnrollmentAnalyticsColumnName.ENROLLMENT_GEOMETRY_COLUMN_NAME);
Expand Down Expand Up @@ -2560,17 +2562,17 @@ protected String createDefaultAlias(String formula) {
void generateFilterCTEs(
EventQueryParams params, CteContext cteContext, boolean isAggregateQuery) {

if (isAggregateQuery) {
generateAggregateFilterCTEs(params, cteContext);
return;
}

// Combine items and item filters and filter only those with an actual filter
List<QueryItem> queryItems =
Stream.concat(params.getItems().stream(), params.getItemFilters().stream())
.filter(QueryItem::hasFilter)
.toList();

if (isAggregateQuery) {
generateAggregateFilterCTEs(queryItems, params, cteContext);
return;
}

// Group query items by repeatable and non-repeatable stages
Map<RepeatableStateStatus, List<QueryItem>> itemsByRepeatableFlag =
queryItems.stream()
Expand Down Expand Up @@ -2607,27 +2609,31 @@ void generateFilterCTEs(
}

/**
* Generates filter CTEs for aggregate enrollment queries. Items are grouped by program stage UID,
* producing one CTE per stage with all dimension columns and filter conditions combined.
* Generates filter CTEs for aggregate enrollment queries. Stage items (both dimensions and item
* filters) are grouped by program stage UID, producing one CTE per stage that projects eligible
* stage dimensions and applies the combined filter conditions.
*
* <p>A {@code latest_events_<stage>} CTE is emitted only when at least one item for the stage
* carries a filter. Unfiltered stage items are still projected by the CTE when they can be read
* from the same "latest event" row without losing repeatable-stage offset semantics.
*
* @param queryItems filtered query items that have at least one filter
* @param params the {@link EventQueryParams} object
* @param cteContext the {@link CteContext} to register CTEs into
*/
private void generateAggregateFilterCTEs(
List<QueryItem> queryItems, EventQueryParams params, CteContext cteContext) {

// Collect all items that have a program stage and group by stage UID
private void generateAggregateFilterCTEs(EventQueryParams params, CteContext cteContext) {
Map<String, List<QueryItem>> itemsByStage =
queryItems.stream()
Stream.concat(params.getItems().stream(), params.getItemFilters().stream())
.filter(qi -> qi.hasProgram() && qi.hasProgramStage())
.filter(qi -> qi.hasFilter() || canProjectUnfilteredItemInAggregateFilterCte(qi))
.collect(groupingBy(qi -> qi.getProgramStage().getUid(), LinkedHashMap::new, toList()));

// For each stage, build a single CTE with all dimension columns and filters
itemsByStage.forEach(
(stageUid, stageItems) -> {
if (stageItems.stream().noneMatch(QueryItem::hasFilter)) {
return;
}
String cteSql = buildAggregateFilterCteSql(stageItems, params);
String cteKey = "latest_events_" + stageUid;
String cteKey = LATEST_EVENTS_CTE_PREFIX + stageUid;
cteContext.addCteFilter(cteKey, stageItems.get(0), cteSql);
});
}
Expand Down Expand Up @@ -2791,6 +2797,14 @@ private void buildProgramStageCte(
if (item.hasFilter()) {
return;
}
// The per-stage filter CTE projects non-offset dimensions for that stage,
// so a redundant per-item CTE would only inflate the query without adding data.
if (canProjectUnfilteredItemInAggregateFilterCte(item)
&& cteContext.getDefinitionByItemUid(
LATEST_EVENTS_CTE_PREFIX + item.getProgramStage().getUid())
!= null) {
return;
}
handleAggregatedEnrollments(cteContext, item, eventTableName, colName, params);
return;
}
Expand Down Expand Up @@ -2840,6 +2854,10 @@ protected boolean columnIsInFormula(String col) {
return col.contains("(") && col.contains(")");
}

private boolean canProjectUnfilteredItemInAggregateFilterCte(QueryItem item) {
return !item.hasRepeatableStageParams() || item.getRepeatableStageParams().isDefaultObject();
}

/**
* Computes a zero-based offset for use with the SQL <em>row_number()</em> function in CTEs that
* partition and order events by date (e.g., most recent first).
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,9 @@
*/
package org.hisp.dhis.analytics.event.data.aggregate;

import static org.hisp.dhis.analytics.event.data.AbstractJdbcEventAnalyticsManager.LATEST_EVENTS_CTE_PREFIX;
import static org.hisp.dhis.analytics.util.RepeatableStageParamsHelper.removeRepeatableStageParams;

import java.util.Map;
import java.util.Set;
import java.util.function.UnaryOperator;
Expand All @@ -44,7 +47,6 @@
* stage-specific dimensions via per-stage filter CTEs.
*/
public final class AggregatedEnrollmentHeaderColumnResolver {
private static final String LATEST_EVENTS_CTE_PREFIX = "latest_events_";

private final StageHeaderClassifier stageHeaderClassifier;

Expand Down Expand Up @@ -91,8 +93,8 @@ public void addHeaderColumns(

Map.Entry<String, CteDefinition> matchingEntry = findMatchingCte(cteDefinitionMap, colName);
if (matchingEntry != null) {
sb.addColumn(matchingEntry.getValue().getAlias() + ".value", "", matchingEntry.getKey());
sb.groupBy(matchingEntry.getKey());
sb.addColumn(matchingEntry.getValue().getAlias() + ".value", "", quotedCol);
sb.groupBy(quotedCol);
} else {
sb.addColumn(quotedCol);
sb.groupBy(quotedCol);
Expand Down Expand Up @@ -155,8 +157,12 @@ private String extractStageUid(String header) {

private Map.Entry<String, CteDefinition> findMatchingCte(
Map<String, CteDefinition> cteDefinitionMap, String colName) {
String normalizedColName =
removeRepeatableStageParams(colName.replace("\"", "").replace("`", ""));
for (Map.Entry<String, CteDefinition> entry : cteDefinitionMap.entrySet()) {
if (entry.getKey().contains(colName)) {
String normalizedKey =
removeRepeatableStageParams(entry.getKey().replace("\"", "").replace("`", ""));
if (normalizedKey.contains(normalizedColName)) {
return entry;
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -30,10 +30,12 @@
package org.hisp.dhis.analytics.event.data;

import static org.hamcrest.CoreMatchers.containsString;
import static org.hamcrest.CoreMatchers.is;
import static org.hamcrest.CoreMatchers.not;
import static org.hamcrest.MatcherAssert.assertThat;
import static org.hisp.dhis.analytics.DataType.NUMERIC;
import static org.hisp.dhis.analytics.QueryKey.NV;
import static org.hisp.dhis.analytics.table.EventAnalyticsColumnName.EVENT_STATUS_COLUMN_NAME;
import static org.hisp.dhis.analytics.table.EventAnalyticsColumnName.OCCURRED_DATE_COLUMN_NAME;
import static org.hisp.dhis.analytics.table.EventAnalyticsColumnName.OU_COLUMN_NAME;
import static org.hisp.dhis.common.DimensionConstants.OPTION_SEP;
Expand Down Expand Up @@ -78,6 +80,7 @@
import org.hisp.dhis.common.QueryFilter;
import org.hisp.dhis.common.QueryItem;
import org.hisp.dhis.common.QueryOperator;
import org.hisp.dhis.common.RepeatableStageParams;
import org.hisp.dhis.common.ValueType;
import org.hisp.dhis.dataelement.DataElement;
import org.hisp.dhis.dataelement.DataElementService;
Expand All @@ -90,6 +93,7 @@
import org.hisp.dhis.program.Program;
import org.hisp.dhis.program.ProgramIndicator;
import org.hisp.dhis.program.ProgramIndicatorService;
import org.hisp.dhis.program.ProgramStage;
import org.hisp.dhis.relationship.RelationshipConstraint;
import org.hisp.dhis.relationship.RelationshipEntity;
import org.hisp.dhis.relationship.RelationshipType;
Expand Down Expand Up @@ -411,6 +415,134 @@ void verifyAggregateEnrollmentStageOrgUnitFilterStaysInFilterCte() {
assertThat(baseCteSql, not(containsString("and \"n1rtSHYf6O6\" in ('ImspTQPwCqd')")));
}

@Test
void verifyAggregateEnrollmentUnfilteredEventStatusUsesLatestEventsCte() {
// When a stage-specific dimension without a filter (e.g. stage.EVENT_STATUS) is requested
// alongside a filtered stage dimension (stage.EVENT_DATE), the unfiltered dimension must
// also be projected by the latest_events_<stage> CTE so the outer SELECT can read it from
// the same "latest event" row. Without this, the resolver emits a reference to
// uiogg.ev_eventstatus while the CTE doesn't project the column, producing a SQL error.
String stageUid = programStage.getUid();

QueryItem dateItem =
new QueryItem(
new BaseDimensionalItemObject(OCCURRED_DATE_COLUMN_NAME),
programA,
null,
ValueType.DATE,
null,
null);
dateItem.setProgram(programA);
dateItem.setProgramStage(programStage);
dateItem.addFilter(new QueryFilter(QueryOperator.GE, "2026-01-01"));
dateItem.addFilter(new QueryFilter(QueryOperator.LE, "2026-12-31"));

QueryItem statusItem =
new QueryItem(
new BaseDimensionalItemObject(EVENT_STATUS_COLUMN_NAME),
programA,
null,
ValueType.TEXT,
null,
null);
statusItem.setProgram(programA);
statusItem.setProgramStage(programStage);

EventQueryParams.Builder params = createRequestParamsBuilder();
params.addItem(dateItem);
params.addItem(statusItem);
params.withEndpointAction(AGGREGATE);

ListGrid grid = new ListGrid();
grid.addHeader(new GridHeader("value", "Value", ValueType.NUMBER, false, false));
grid.addHeader(
new GridHeader(stageUid + ".eventdate", "Event date", ValueType.DATE, false, false));
grid.addHeader(
new GridHeader(stageUid + ".eventstatus", "Event status", ValueType.TEXT, false, false));

subject.getEnrollments(params.build(), grid, 10000);
verify(jdbcTemplate).queryForRowSet(sql.capture());

String generatedSql = noEof(sql.getValue());
String latestEventsCteSql =
generatedSql.substring(
generatedSql.indexOf("latest_events_" + stageUid + " as ("),
generatedSql.indexOf("enrollment_aggr_base as ("));

// The latest_events CTE must project both the filtered (eventdate) and unfiltered
// (eventstatus) stage dimensions so the outer SELECT can read them from the same row.
assertThat(latestEventsCteSql, containsString("ev_occurreddate"));
assertThat(latestEventsCteSql, containsString("ev_eventstatus"));

// The redundant per-item CTE must not be created when a latest_events_<stage> already exists.
assertThat(generatedSql, not(containsString(stageUid + "_" + EVENT_STATUS_COLUMN_NAME + "_0")));
}

@Test
void verifyAggregateEnrollmentRepeatableOffsetItemKeepsPerItemCteWhenLatestEventsCteExists() {
String stageUid = repeatableProgramStage.getUid();
String deUid = dataElementA.getUid();

QueryItem dateItem = createFilteredStageDateItem(repeatableProgramStage);
QueryItem offsetItem = createRepeatableOffsetDataElementItem(stageUid, -1);

EventQueryParams.Builder params = createRequestParamsBuilder();
params.addItem(dateItem);
params.addItem(offsetItem);
params.withEndpointAction(AGGREGATE);

ListGrid grid = new ListGrid();
grid.addHeader(new GridHeader("value", "Value", ValueType.NUMBER, false, false));
grid.addHeader(
new GridHeader(stageUid + ".eventdate", "Event date", ValueType.DATE, false, false));
grid.addHeader(
new GridHeader(stageUid + "[-1]." + deUid, "Offset value", ValueType.NUMBER, false, false));

subject.getEnrollments(params.build(), grid, 10000);
verify(jdbcTemplate).queryForRowSet(sql.capture());

String generatedSql = noEof(sql.getValue());
String latestEventsCteSql =
generatedSql.substring(
generatedSql.indexOf("latest_events_" + stageUid + " as ("),
generatedSql.indexOf("enrollment_aggr_base as ("));

assertThat(latestEventsCteSql, containsString("ev_occurreddate"));
assertThat(latestEventsCteSql, not(containsString("ev_" + deUid)));
assertThat(generatedSql, containsString(stageUid + "_" + deUid + "_0 as ("));
assertThat(generatedSql, containsString("as \"" + stageUid + "[-1]." + deUid + "\""));
}

@Test
void verifyAggregateEnrollmentRepeatableOffsetsAreNotDuplicatedInLatestEventsCte() {
String stageUid = repeatableProgramStage.getUid();
String deUid = dataElementA.getUid();

EventQueryParams.Builder params = createRequestParamsBuilder();
params.addItem(createFilteredStageDateItem(repeatableProgramStage));
params.addItem(createRepeatableOffsetDataElementItem(stageUid, -1));
params.addItem(createRepeatableOffsetDataElementItem(stageUid, 1));
params.withEndpointAction(AGGREGATE);

ListGrid grid = new ListGrid();
grid.addHeader(new GridHeader("value", "Value", ValueType.NUMBER, false, false));
grid.addHeader(
new GridHeader(stageUid + ".eventdate", "Event date", ValueType.DATE, false, false));

subject.getEnrollments(params.build(), grid, 10000);
verify(jdbcTemplate).queryForRowSet(sql.capture());

String generatedSql = noEof(sql.getValue());
String latestEventsCteSql =
generatedSql.substring(
generatedSql.indexOf("latest_events_" + stageUid + " as ("),
generatedSql.indexOf("enrollment_aggr_base as ("));

assertThat(countOccurrences(latestEventsCteSql, "ev_" + deUid), is(0));
assertThat(generatedSql, containsString(stageUid + "_" + deUid + "_0 as ("));
assertThat(generatedSql, containsString(stageUid + "_" + deUid + "_1 as ("));
}

@Test
void verifyAggregateEnrollmentIncludesProgramStatusFilterInSql() {
EventQueryParams.Builder params = createRequestParamsBuilder();
Expand Down Expand Up @@ -1109,6 +1241,32 @@ private EventQueryParams createAggregateEnrollmentWithStageDateParams() {
return params.build();
}

private QueryItem createFilteredStageDateItem(ProgramStage stage) {
QueryItem dateItem =
new QueryItem(
new BaseDimensionalItemObject(OCCURRED_DATE_COLUMN_NAME),
programA,
null,
ValueType.DATE,
null,
null);
dateItem.setProgram(programA);
dateItem.setProgramStage(stage);
dateItem.addFilter(new QueryFilter(QueryOperator.GE, "2026-01-01"));
dateItem.addFilter(new QueryFilter(QueryOperator.LE, "2026-12-31"));
return dateItem;
}

private QueryItem createRepeatableOffsetDataElementItem(String stageUid, int offset) {
QueryItem item = new QueryItem(new BaseDimensionalItemObject(dataElementA.getUid()));
item.setProgram(programA);
item.setProgramStage(repeatableProgramStage);
item.setValueType(ValueType.NUMBER);
item.setRepeatableStageParams(
RepeatableStageParams.of(offset, stageUid + "[" + offset + "]." + dataElementA.getUid()));
return item;
}

private EventQueryParams createAggregateEnrollmentWithEventDateParams() {
BaseDimensionalItemObject dateItem = new BaseDimensionalItemObject(OCCURRED_DATE_COLUMN_NAME);
QueryItem queryItem = new QueryItem(dateItem, programA, null, ValueType.DATE, null, null);
Expand Down Expand Up @@ -1203,4 +1361,14 @@ private void testIt(
private String noEof(String sql) {
return sql.replaceAll("\\s+", " ").trim();
}

private int countOccurrences(String value, String search) {
int count = 0;
int index = 0;
while ((index = value.indexOf(search, index)) >= 0) {
count++;
index += search.length();
}
return count;
}
}
Loading
Loading