Skip to content
Merged
44 changes: 37 additions & 7 deletions docs/impulse/docs/config/configuration.md
Original file line number Diff line number Diff line change
Expand Up @@ -161,7 +161,8 @@ in `solver_config` so the solver renames each table's columns at read time.
the internal names.
- `filters` (`dict[str, str]`): equality filters applied **after** renaming. Keys are internal column
names; values are literals to match. Useful for project/toolbox scoping where a single value should
always be enforced.
always be enforced. Values are strings but are coerced to the target column's type, so a boolean
column takes `"true"` / `"false"` (an unparseable value errors at read time under ANSI mode).

Top-level fields on `SolverConfig`:

Expand All @@ -179,7 +180,7 @@ Per-table sections (each a `TableConfig`):
| `channel_tags` | when `channel_tags_table` is configured | Tag key/value column renames |
| `channel_metrics` | always | Custom channel_id column, custom value/timestamp columns |
| `channel_mapping` | when `channel_mapping_table` is configured | Alias-table column renames; `priority` column; optional `join_keys` for non-default alias-resolution composite keys |
| `channels` | always | RLE column renames (`tstart`/`tend`/`value`) |
| `channels` | always | RLE column renames (`tstart`/`tend`/`value`); `filters` supported (applied at read time, before non-solve columns are dropped) |
| `unit_conversion` | when `unit_conversion_table` is configured | Unit-conversion table column renames (`unit`, `group_id`, `conversion_factor`) |

Internal column names that mappings can target:
Expand All @@ -189,6 +190,8 @@ Internal column names that mappings can target:
| `container_id` | Container identifier |
| `channel_id` | Channel identifier |
| `tstart`, `tend`| Sample interval start/end on the `channels` table (RLE) |
| `timestamp` | Raw sample timestamp on the `channels` table (RAW mode; encoded into `tstart`/`tend`) |
| `is_plausible` | Boolean plausibility flag on the `channels` table (RAW mode); consumed by `drop_implausible_data` |
| `start_ts`, `stop_ts` | Measurement start/stop epoch timestamps on the `container_metrics` table — referenced by `ContainerEvent` to derive event-fact start/end |
| `value` | Sample value (or attribute value on the EAV tag table) |
| `key` | Attribute key on the EAV `container_tags` table |
Expand All @@ -206,10 +209,36 @@ Internal column names that mappings can target:

:::note Feature support

`DefaultSolver` consumes every section of `solver_config`: per-table
`column_name_mapping`, per-table `filters`, top-level `project_id`, and the
`channel_mapping` / `unit_conversion` sections. Sections for tables you do
not configure (e.g. `channel_tags`, `channel_mapping`) are simply unused.
`DefaultSolver` consumes every section's `column_name_mapping`, plus the
top-level `project_id` and the `channel_mapping` / `unit_conversion` sections.
Per-table `filters` are applied for `container_tags`, `container_metrics`,
`channel_mapping`, and `channels`. Filters on `channel_tags`, `channel_metrics`,
and `poi_channels` are accepted for forward compatibility but **not yet applied**.
Sections for tables you do not configure (e.g. `channel_tags`, `channel_mapping`)
are simply unused.

:::

:::caution Channels filters in RAW mode

`channels.filters` are applied **before** raw encoding. A filter that removes
samples from the middle of a channel therefore *bridges* the surrounding interval
(the last good value is held across the gap) rather than splitting it.

**Intended scope:** use `channels.filters` in RAW mode **only for whole-channel
scoping** (a value constant across all of a channel's samples, e.g. a single
`source`/stream or project scoping), never for per-sample cleaning (value ranges,
quality flags, NaN drops). Per-sample cleaning leaves interior gaps that get bridged,
distorting the signal. To drop implausible samples with correct boundaries, use
`drop_implausible_data`.

When `data_type = RAW`, the validator **rejects** a `channels.filters` entry on
`is_plausible` (unambiguously per-sample) and **warns** on any other entry, since it
cannot tell scoping from cleaning statically.

In RLE mode there is no such restriction: the `channels` table is already encoded, so
a filter simply drops the matching `[tstart, tend)` interval rows without bridging, and
may target any column.

:::

Expand All @@ -235,7 +264,8 @@ not configure (e.g. `channel_tags`, `channel_mapping`) are simply unused.
"filters": {"toolbox_id": "my_toolbox"}
},
"channels": {
"column_name_mapping": {}
"column_name_mapping": {},
"filters": {"source": "live"}
}
}
}
Expand Down
2 changes: 2 additions & 0 deletions skills/impulse-analyze/SKILL.md
Original file line number Diff line number Diff line change
Expand Up @@ -131,6 +131,8 @@ DefaultSolver(spark, config=None, is_raw_data=False, drop_implausible_data=False
`is_raw_data=True`; the ad-hoc equivalent of `query_engine.raw_encoder`. Import from
`impulse_query_engine.analyze.query.solvers.solver_config`.
- `drop_implausible_data=True` drops rows where `is_plausible = false` (requires `is_raw_data=True`).
A `channels.filters` entry on `is_plausible` with `is_raw_data=True` is rejected at construction; use
`drop_implausible_data` instead.
- `config` takes a `SolverConfig` for column-name remapping / project scoping — the same object
described under `solver_config` in `impulse-config`.

Expand Down
4 changes: 4 additions & 0 deletions skills/impulse-config/SKILL.md
Original file line number Diff line number Diff line change
Expand Up @@ -117,6 +117,10 @@ With `data_type="RAW"`, the `channels` table additionally uses the internal name
per-sample timestamp) and — only when `drop_implausible_data` is on — `is_plausible`. Remap them the
same way, e.g. `"channels": {"column_name_mapping": {"ts_raw": "timestamp"}}`.

**`channels.filters` in RAW mode** run before raw encoding, so they bridge intervals across dropped
samples. Use them only for whole-channel scoping, not per-sample cleaning: an `is_plausible` filter is
rejected (use `drop_implausible_data`), any other channels filter warns.

Top-level `project_id` (str, optional) applies an equality filter on the `project_id` column of every
table that has one (`container_tags`, `container_metrics`, `channel_mapping`). Omit if not needed.

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -270,6 +270,9 @@ def __init__(
self.is_raw_data = is_raw_data
self.drop_implausible_data: bool = drop_implausible_data
self.raw_encoder: RawEncoder = raw_encoder
# In RAW mode an is_plausible channels filter drops samples before raw
# encoding, bridging intervals instead of splitting them. Reject it.
self.config.reject_implausible_channels_filter_in_raw(is_raw=self.is_raw_data)
self.channel_encoder: RleEncoder | IntervalEncoder = self._build_channel_encoder()

def _build_channel_encoder(self) -> RleEncoder | IntervalEncoder:
Expand Down Expand Up @@ -1098,9 +1101,11 @@ def _prepare_channels_join(self, query, channels_df) -> tuple[DataFrame, DataFra
"""Shared prelude for :meth:`solve` and :meth:`solve_calculated_channels`.

Applies optional per-channel unit conversion, reads and column-maps the
channel-data table (raw-encoding it when in raw mode), broadcast-joins it
to the channel-match frame on ``[container_id, channel_id]``, and counts
the distinct containers. Any container-metadata columns already on
channel-data table (raw-encoding it when in raw mode), applies the
per-table ``channels.filters`` as equality predicates (after
``column_name_mapping``), broadcast-joins it to the channel-match frame
on ``[container_id, channel_id]``, and counts the distinct containers.
Any container-metadata columns already on
*channels_df* (attached by :meth:`attach_container_metadata`) ride
through the broadcast join into ``joined_df``.

Expand Down Expand Up @@ -1137,6 +1142,11 @@ def _prepare_channels_join(self, query, channels_df) -> tuple[DataFrame, DataFra
q = query.db.channels(self.spark)
q = self._apply_column_mapping(q, self.config.channels.column_name_mapping)

# Equality filters on internal column names, applied before the select
# below so a filter may reference any channels column, dropped or not.
for col_name, value in self.config.channels.filters.items():
q = q.where(F.col(col_name) == value)

if self.is_raw_data:
# Encode the raw samples into intervals (RLE or interval) for the solving step.
q = self.channel_encoder.prepare_channels_df(q)
Expand Down
17 changes: 17 additions & 0 deletions src/impulse_query_engine/analyze/query/solvers/solver_config.py
Original file line number Diff line number Diff line change
Expand Up @@ -400,3 +400,20 @@ def col_map(self) -> dict[str, str]:
# per-row series_type / dtype marker column is needed in the frame.
"value_string": self.poi_value_string_col,
}

def reject_implausible_channels_filter_in_raw(self, is_raw: bool) -> None:
"""Raise if an is_plausible channels filter is set in RAW mode.

Such a filter runs before raw encoding and bridges intervals across dropped
samples instead of splitting them; use drop_implausible_data instead. No-op
when not raw.
"""
if not is_raw:
return
if self.is_plausible_col in self.channels.filters:
raise ValueError(
f"A channels.filters entry on '{self.is_plausible_col}' in RAW mode is "
"applied before raw encoding, which bridges intervals across dropped "
"samples. Use drop_implausible_data=True instead -- it drops "
"implausible points inside the encoder with correct interval boundaries."
)
41 changes: 40 additions & 1 deletion src/impulse_reporting/config/config_parser.py
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
import re
import warnings
from datetime import datetime
from enum import Enum, StrEnum
from typing import Annotated
Expand Down Expand Up @@ -422,11 +423,49 @@ def validate_drop_implausible_data_requires_raw(self):
if self.drop_implausible_data and self.data_type is not DataType.RAW:
raise ValueError(
"drop_implausible_data=True requires data_type=RAW. "
"The implausible-data filter is only applied during the RAW -> RLE "
"The implausible-data filter is only applied during the RAW -> interval "
"conversion path; RLE input is passed through unchanged."
)
return self

@model_validator(mode="after")
def reject_channels_filter_on_is_plausible_in_raw(self):
"""Reject an is_plausible channels filter in RAW mode (delegates to SolverConfig).

The shared invariant lives on ``SolverConfig`` so direct query-engine use
(``DefaultSolver``) enforces the same rule; here we surface it at config-parse time.
"""
if self.solver_config is not None:
self.solver_config.reject_implausible_channels_filter_in_raw(
is_raw=self.data_type is DataType.RAW
)
return self

@model_validator(mode="after")
def warn_channels_filters_bridge_in_raw(self):
"""Warn (not reject) on any non-plausibility channels filter in RAW mode.

Such filters run before raw encoding and bridge intervals across dropped
samples. That is fine for whole-channel scoping but corrupts per-sample cleaning,
and the two cannot be told apart statically, so we warn. ``is_plausible`` is
hard-rejected above as the one unambiguously per-sample case.
"""
if self.data_type is not DataType.RAW or self.solver_config is None:
return self
plausibility_col = self.solver_config.is_plausible_col
other = [c for c in self.solver_config.channels.filters if c != plausibility_col]
if other:
warnings.warn(
f"channels.filters {other} in RAW mode run before raw encoding and bridge "
"intervals across dropped samples. Use them only for whole-channel scoping, "
"not per-sample cleaning. To drop implausible points with correct interval "
"boundaries, use drop_implausible_data=True instead.",
# stacklevel=1 (the warn call itself): inside a pydantic model_validator the
# frames above are pydantic internals, so a higher level would mislead.
stacklevel=1,
)
return self

@model_validator(mode="after")
def default_raw_encoder_for_raw_data(self):
"""When ``data_type=RAW`` and ``raw_encoder`` is unset, default to RLE."""
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,152 @@
# pylint: disable=missing-function-docstring
"""End-to-end tests for ``SolverConfig.channels.filters``.

The channels table has no standalone ``filter_*`` stage — its per-table
equality filters are applied inside ``DefaultSolver._prepare_channels_join``
when the channel-data table is read (after ``column_name_mapping``, before the
UDF-column projection). These tests exercise the filter through a real
``solve_calculated_channels`` run against the wide-only ``basic_narrow_db``
fixture and assert on the resulting computed values, not just row counts.

Covers:
- A string filter that removes rows: only the retained sample survives, and it
carries the correct scaled value.
- No filter configured: behaviour is unchanged (full data comes through).
- A non-matching filter: zero result rows.
- A boolean channels column: string ``"true"`` matches the ``true`` rows via
Spark's literal coercion (the shared ``F.col(col) == value`` mechanism).
"""

import pyspark.sql.functions as F
import pytest

from impulse_query_engine.analyze.query.channels.calculated_channel import (
CalculatedChannel,
)
from impulse_query_engine.analyze.query.solvers.default_solver import DefaultSolver
from impulse_query_engine.analyze.query.solvers.solver_config import (
SolverConfig,
TableConfig,
)
from impulse_query_engine.measurement_db import MeasurementDB, MeasurementDBConfig
from tests.conftest import basic_narrow_db, spark # noqa: F401 (pytest fixtures)

# Known datum from tests/unit/data/basic_narrow_csv/channel_data.csv:
# container 1, channel 5 (Engine RPM), second RLE row.
_C1_RPM_TSTART = 1499929245761999
_C1_RPM_VALUE = 1081.0


def _clone_with_channel_col(db: MeasurementDB, col_name: str, expr) -> MeasurementDB:
"""Clone a ``for_debug`` db, adding *col_name* = *expr* to the channels table only."""
tables = dict(db.config.debug_tables)
tables["channels"] = tables["channels"].withColumn(col_name, expr)
return MeasurementDB(MeasurementDBConfig.for_debug(tables), ws=db.ws)


def _only_row_marker(marked_value, other_value):
"""Expression tagging just the single (container 1, _C1_RPM_TSTART) channels row.

That row gets *marked_value*; every other channels row gets *other_value*.
"""
is_target = (F.col("container_id") == 1) & (F.col("tstart") == _C1_RPM_TSTART)
return F.when(is_target, F.lit(marked_value)).otherwise(F.lit(other_value))


def _rpm_x2(query):
"""A calculated channel: Engine RPM * 2 (real, predictable per-interval values)."""
return CalculatedChannel(
query.channel(channel_name="Engine RPM") * 2,
{"channel_name": "rpm_x2", "data_key": "CALC"},
)


class TestChannelsFilter:
"""DefaultSolver applies ``config.channels.filters`` when reading channels."""

def test_string_filter_removes_rows_and_keeps_real_value(self, spark, basic_narrow_db):
"""Only the ``source == 'live'`` sample survives, with its correct scaled value."""
db = _clone_with_channel_col(
basic_narrow_db, "source", _only_row_marker("live", "archive")
)
query = db.query
cfg = SolverConfig(channels=TableConfig(filters={"source": "live"}))
result = query.select(_rpm_x2(query)).solve_calculated_channels(
spark, solver=DefaultSolver(spark, config=cfg)
)

rows = result.select("container_id", "tstart", "value").collect()
# Every channels row except the one tagged "live" was filtered out, so the
# whole solve collapses to that single retained sample.
assert len(rows) == 1
assert rows[0]["container_id"] == 1
assert rows[0]["tstart"] == _C1_RPM_TSTART
assert rows[0]["value"] == pytest.approx(_C1_RPM_VALUE * 2)

def test_no_filter_leaves_behavior_unchanged(self, spark, basic_narrow_db):
"""Default SolverConfig (no channels filters) returns the full, unfiltered data."""
query = basic_narrow_db.query
result = query.select(_rpm_x2(query)).solve_calculated_channels(
spark, solver=DefaultSolver(spark)
)

c1 = result.filter(F.col("container_id") == 1)
# The known datum is still present and correct...
target = c1.filter(F.col("tstart") == _C1_RPM_TSTART).select("value").collect()
assert len(target) == 1
assert target[0]["value"] == pytest.approx(_C1_RPM_VALUE * 2)
# ...alongside the many other samples the string-filter test removed.
assert c1.count() > 1

def test_non_matching_filter_returns_empty(self, spark, basic_narrow_db):
"""A filter value matching no channels rows yields zero results."""
db = _clone_with_channel_col(
basic_narrow_db, "source", _only_row_marker("live", "archive")
)
query = db.query
cfg = SolverConfig(channels=TableConfig(filters={"source": "does_not_exist"}))
result = query.select(_rpm_x2(query)).solve_calculated_channels(
spark, solver=DefaultSolver(spark, config=cfg)
)
assert result.count() == 0

def test_boolean_column_filter_matches_true_rows(self, spark, basic_narrow_db):
"""A boolean channels column filters on ``"true"``: Spark casts the string
literal to the column's type, so only the ``true`` rows survive."""
db = _clone_with_channel_col(basic_narrow_db, "is_valid", _only_row_marker(True, False))
query = db.query
cfg = SolverConfig(channels=TableConfig(filters={"is_valid": "true"}))
result = query.select(_rpm_x2(query)).solve_calculated_channels(
spark, solver=DefaultSolver(spark, config=cfg)
)

rows = result.select("container_id", "tstart", "value").collect()
assert len(rows) == 1
assert rows[0]["container_id"] == 1
assert rows[0]["tstart"] == _C1_RPM_TSTART
assert rows[0]["value"] == pytest.approx(_C1_RPM_VALUE * 2)


class TestChannelsFilterRawGuard:
"""The shared RAW-mode hard-reject is enforced at DefaultSolver construction.

This guarantees direct engine use cannot bypass the invariant that the reporting
config parser enforces (an is_plausible channels filter would bridge intervals
across dropped samples before raw encoding).
"""

def test_is_plausible_filter_rejected_at_construction_in_raw(self, spark):
"""Building a RAW solver with an is_plausible channels filter raises immediately."""
cfg = SolverConfig(channels=TableConfig(filters={"is_plausible": "true"}))
with pytest.raises(ValueError, match="before raw encoding"):
DefaultSolver(spark, config=cfg, is_raw_data=True)

def test_is_plausible_filter_allowed_when_not_raw(self, spark):
"""The same filter is legitimate scoping outside RAW mode, so construction succeeds."""
cfg = SolverConfig(channels=TableConfig(filters={"is_plausible": "true"}))
DefaultSolver(spark, config=cfg, is_raw_data=False)

def test_non_plausibility_filter_allowed_in_raw(self, spark):
"""The hard-reject is is_plausible-only: a scoping filter builds fine in RAW mode."""
cfg = SolverConfig(channels=TableConfig(filters={"source": "live"}))
DefaultSolver(spark, config=cfg, is_raw_data=True)
Loading
Loading