# Assignment 3: A plant stream, collected and made trustworthy
**Module:** Week 03 · **Released:** L6 · **Due:** ~1 week later · **6 points on Canvas.**

## Goal
Collect ten minutes of a live chemical plant stream, then build a batch pipeline that turns it into a table you would compute on. The stream arrives late, out of order, duplicated, with gaps and suspect readings. You will land it, validate it, deduplicate it, fill a regular grid, and find the disturbance in it.

You need nothing beyond your project and `uv`. There is no account, and the connection runs over port 443, so it works from any network.

**You should be able to:**
- land an unbounded stream durably before processing it;
- window over event time, and choose a watermark and state its cost;
- make a pipeline idempotent under at-least-once delivery;
- validate with pandera and quarantine bad records instead of deleting them;
- reshape long to wide, fill a ragged grid, and record which cells you invented;
- write a stage against a Polars lazy query plan and read `.explain()`.

## The stream
A [Tennessee Eastman](https://doi.org/10.1016/0098-1354(93)80018-I) plant simulation, the same plant L5 used as files, published live over MQTT.

| | |
|---|---|
| Host | `kitchin-services.cheme.cmu.edu`, port 443, WebSocket path `/mqtt` |
| Topics | `plant/tep/birth` (one retained message), `plant/tep/telemetry` (samples) |
| Speed | plant runs 144× wall clock; one sample per 180 plant seconds, one message per 1.25 s |
| Ten minutes | one plant day: 480 messages, ~25,000 readings, ~3 MB |
| Tags | 53: 41 measurements `XMEAS_*` and 12 manipulated variables `XMV_*` |
| Disturbances | one episode every 12 plant hours, lasting 1.5 to 4 h, then clears; your ten minutes holds two |

L5's files had 52 columns because they drop `XMV_12`, the agitator speed, which sits constant. Expect it here.

### The birth message is your data contract
It lists every tag with its description and unit, plus the clock conversion. Build your schema from it rather than typing the tag list.

```json
{"type": "birth", "plant": "tep", "accelerationFactor": 144,
 "plantEpoch": "2026-09-09T14:55:46.480Z", "sampleIntervalSeconds": 180,
 "tags": [{"name": "XMEAS_1_A_feed", "description": "A Feed (stream 1)", "unit": "kscmh"},
          {"name": "XMEAS_2_D_feed", "description": "D Feed (stream 2)", "unit": "kg/hr"}]}
```

Two of the 53 tags are shown; the real message carries all of them.

Use the descriptions in your report: "the separator pressure moved" is a finding, "column 13 moved" is not.

### A telemetry message, as the collector saves it
```json
{"receivedAt": "2026-09-09T14:56:29.712Z", "topic": "plant/tep/telemetry",
 "payload": {"seq": 197, "timestamp": "2026-09-09T16:37:45.479Z", "isHistorical": false,
   "metrics": [
     {"name": "XMEAS_1_A_feed",  "value": 0.253256,  "statusCode": "Good",
      "sourceTimestamp": "2026-09-09T16:37:45.479Z"},
     {"name": "XMEAS_23_feed_A", "value": 31.551008, "statusCode": "Good",
      "sourceTimestamp": "2026-09-09T16:34:45.479Z"}]}}
```

### Three clocks

| Field | Clock | What it is |
|---|---|---|
| `sourceTimestamp` (per metric) | plant | **event time**, when the instrument measured it. Window on this. |
| `timestamp` (per message) | plant | when the batch was assembled |
| `receivedAt` (envelope) | your laptop | **processing time**, when it landed |

The first two run on the plant clock, 144× fast from `plantEpoch`. Subtracting `receivedAt` from either gives a meaningless number until you put `receivedAt` on the plant clock too. Both clocks agree at `plantEpoch`, and every wall second after it is 144 plant seconds:

```text
receivedAt on the plant clock = plantEpoch + (receivedAt − plantEpoch) × accelerationFactor
```

For the message below, `receivedAt` is 43.232 wall seconds after `plantEpoch`, which is 6,225.4 plant seconds, so the message landed at 16:39:31.9 on the plant clock, 106 plant seconds after its `timestamp`. In Polars a datetime minus a datetime is a duration, and a duration times an integer is a duration, so this is one expression.

**Nineteen tags report the past.** `XMEAS_23` to `XMEAS_41` are composition analysers. Each reports its previous sample, so its `sourceTimestamp` lags the message by 0 to 720 s and repeats across several messages. Assume one timestamp per message and those nineteen channels are wrong by up to twelve minutes.

### What goes wrong on the link
The link drops out, the publisher buffers, and the backlog drains when it reconnects, so late data comes **in bursts**. Over 6,000 messages at these settings:

| | |
|---|---|
| records arriving after a newer one | 13.6% |
| lateness, plant minutes | median 51, p95 117, max 198 |
| samples never delivered | 20 in 5,961 |
| duplicate deliveries | 59 |
| `Uncertain` / `Bad` status codes | 648 / 146 in 318,000 |

Your numbers will differ. Every late record carries `isHistorical: true`.

## Tasks

### 1. Collect
Run the collector for **at least ten minutes**:

```bash
uv run --no-project https://kitchingroup.cheme.cmu.edu/f26-06763/a03-collect.py --minutes 10
```

It appends to `raw/stream-<date>.ndjson` and never truncates, so rerunning adds data. **Treat that file as read-only.** Every later stage replays from it, and it is the one thing you cannot regenerate.

### 2. Decode and reshape (Polars, lazy)
Explode `metrics[]` into one row per reading and write `processed/tidy.parquet` with **exactly these columns**:

| Column | Type | From |
|---|---|---|
| `tag` | string | `metrics[].name` |
| `value` | double | `metrics[].value` |
| `status` | string | `metrics[].statusCode` |
| `event_time` | timestamp, UTC | `metrics[].sourceTimestamp` |
| `server_time` | timestamp, UTC | the message `timestamp` |
| `received_at` | timestamp, UTC | the envelope `receivedAt` |
| `seq` | int | the message `seq` |
| `is_historical` | bool | the message `isHistorical` |

- Write this stage in **Polars on a `LazyFrame`**: scan, express the explode and casts as expressions, and `.collect()` once. Print `.explain()` and keep the output, for example in `results/explain.txt`. Later stages may use pandas or Polars.
- **Getting at the nesting.** Keep only `plant/tep/telemetry` lines. `payload` is a struct, so `pl.col("payload").struct.field("seq")` pulls out one field (`.unnest("payload")` pulls out all of them). `metrics` is a list of structs: `.explode("metrics")` gives one row per reading, and `pl.col("metrics").struct.field("name")` then reads each field. Parse the three timestamps with `.str.to_datetime(...)` to UTC. The L6 notebook does the same reshape in pandas with `pd.json_normalize`, if you want to see the shape of the result first.
- **Watch out:** `pl.scan_ndjson` infers the schema from the first 100 lines. Your file mixes birth and telemetry messages, and a birth message past line 100 loses its fields silently. Pass `infer_schema_length=None` or declare the schema.

### 3. Validate and quarantine
- Generate a pandera schema **from the birth message's `tags` list**: per-tag dtype and plausible range, `status` in the set the stream uses, `event_time` inside your collection window.
- A reading's range depends on its unit, and the birth message gives each tag's unit, so build a lookup from tag to bounds and check `value` against it with a check on the whole frame (`pa.DataFrameSchema(columns, checks=pa.Check(fn))`), as the L6 notebook does with its unit ranges.
- Take the collection window from when you collected, `receivedAt` on the plant clock, with some slack for the replayed backlog. A window computed from the event times it checks can never fail.
- Validate with `lazy=True`, so one run reports every failure. `SchemaErrors.failure_cases` gives you the offending rows: its `index` column points back at them.
- Write every `Bad_*` and `Uncertain_*` reading, and everything the schema rejects, to `processed/quarantine.parquet` with a `reason` column. Remove them from `tidy.parquet` and report the counts.
- Show the gate works: copy a few hundred clean rows, break a handful on purpose (a tag the birth message does not declare, a value outside its unit's range, a missing value, an unknown status, a timestamp from another year), validate the copy with `lazy=True`, and show which checks caught which rows.

Install `pandera[pandas]` and `import pandera.pandas as pa`. The old `import pandera as pa` still works but warns that it will stop.

```python
import pandera.pandas as pa

schema = pa.DataFrameSchema({
    "tag":    pa.Column(str, pa.Check.isin(tag_names)),
    "value":  pa.Column(float, nullable=False),
    "status": pa.Column(str),
})
try:
    schema.validate(df, lazy=True)
except pa.errors.SchemaErrors as e:
    bad = e.failure_cases        # column, check, failure_case, index
```

### 4. Deduplicate and choose a watermark
- **Deduplicate on `(tag, event_time)`** and report how many rows that removed. Do not key on `seq`: it is one byte and wraps at 256 within your collection. Expect fewer than 53 rows per sample afterwards, because a latched analyser's repeats collapse to one measurement.
- **Show idempotence:** run the whole pipeline twice over the same file and show nothing changes.
- **Measure lateness.** A reading's lateness is how far its `event_time` sits behind the *frontier*, the newest `event_time` in any message that arrived before its own, or zero if it is not behind. The L6 notebook computes exactly this in pandas.
- **Sweep a watermark** over allowances of 0, 15, 30, 60 and 120 plant minutes (more if you like). Write `results/watermark.csv` with one row per allowance and at least these columns:
  - `allowance_plant_min`, the allowance;
  - `completeness_pct`, the percentage of readings whose lateness is within the allowance, so they arrive in time to be counted;
  - `still_waiting`, the readings not yet final when you stopped collecting: those whose `event_time` is within the allowance of the newest `event_time` you received.
- **Compare clocks:** compute one aggregate in 30-minute windows on `event_time`, and again on `received_at` converted to the plant clock (the formula is under Three clocks, above), and report where and by how much they disagree. Polars `group_by_dynamic` or pandas `resample` works on either column.

### 5. Fill the grid, window, find the disturbance
**Grid.** Pivot the deduplicated, quarantined table onto a regular grid, one row per 180-plant-second instant from first to last sample and one column per tag. Save it as `processed/grid.parquet`.
- A pivot only has rows for instants that delivered something. Build the spine of instants first (`pl.datetime_range` or `pd.date_range`) and join the pivot onto it, so a lost sample is an empty row rather than a missing one, and add a column for any tag that never arrived.
- Before filling, count empty cells separately for the 34 continuous channels and the 19 analysers. The birth message marks analysers with `analyserIntervalHours`. Expect roughly 1% and 58% empty. Most continuous holes sit in the grid's first few rows, which exist only because lagged analyser samples carry earlier event times than the first continuous sample; lost messages and quarantined readings make up the rest. Say where yours came from.
- Fill every cell, and record which cells you filled: a boolean mask, a `_filled` column per tag, or a long list of filled cells. Take the mask before you fill. Polars has `forward_fill`, `backward_fill` and `interpolate`; pandas has `ffill`, `bfill` and `interpolate`. The first rows have nothing before them, so decide what they get.
- The two kinds of hole differ. An empty analyser cell means no new analysis yet, so carrying forward repeats what the instrument says. An empty continuous cell means a lost message, so carrying forward invents a value. Justify your fill for each.
- Quarantine before you fill, or a `Bad` null becomes an invented number.

**Windows.** Aggregate the filled grid into tumbling 30-plant-minute windows over event time: pandas `pd.Grouper(key="event_time", freq="30min")` or `.resample()`, Polars `group_by_dynamic`. Write `results/windows.csv` in **wide** form, one row per window and one column per tag. Your last window is incomplete; count the samples in each window and say how you handled it.

**Disturbance.** Find when the plant was disturbed, when it recovered, and which tags moved. A per-window score against a quiet stretch is enough: z-score each continuous tag's window mean against the mean and standard deviation of windows you believe are quiet, average the absolute z-scores across tags, then look at the windows that stand out and the tags with the largest scores in them. A baseline that includes the disturbance inflates its own standard deviation and scores the disturbance as nearly normal. You are graded on method and honesty, **not on naming the fault**. Some disturbances are slow and genuinely hard to separate from noise, and saying so is a good answer.

### 6. Write `REPORT.md`
**Two pages maximum**, six sections in this order:

1. **The three clocks.** Which one you windowed on and why. The largest `sourceTimestamp` to message `timestamp` gap you measured, and its cause.
2. **What you collected.** Messages, records, duplicates removed, `Uncertain` and `Bad` counts, records quarantined and why, fraction late.
3. **What you filled.** Empty cells before filling, continuous against analyser, and the cause of each. Your fill for each, with the argument. How a reader tells measured from invented.
4. **The watermark.** Your sweep table, the allowance you would pick for a real plant, and why. What being wrong costs in each direction. A choice with no reason gets about half.
5. **Event time against processing time.** The same aggregate both ways, the disagreement, where it was worst, and why.
6. **The disturbance.** What, when, which tags, how confident, and what you could not tell from noise.

Also include two or three sentences on one thing `.explain()` showed the planner doing that your code did not ask for, and one line disclosing generative-AI use.

## Names to use

| Thing | Name |
|---|---|
| Raw capture | `raw/stream-*.ndjson`, append-only |
| Pipeline code | `src/…/*.py`, or flat `.py` files |
| pandera schema | any `.py` (a notebook is allowed) |
| Tidy table | `processed/tidy.parquet` |
| Filled sample grid | `processed/grid.parquet` |
| Quarantined records | `processed/quarantine.parquet` |
| Watermark sweep | `results/watermark.csv` |
| Windowed aggregate | `results/windows.csv` |
| Report | `REPORT.md` |

If you used other names, the script searches for your files and prints what it found on page one. Correct it with a flag if needed:

```bash
uv run --no-project a03-evidence.py --andrew-id yourid --tidy out/readings.parquet
```

## Submit
From your project root, after running the pipeline:

```bash
uv run --no-project https://kitchingroup.cheme.cmu.edu/f26-06763/a03-evidence.py \
    --andrew-id yourid --name "Your Name"
```

Upload the resulting `evidence.pdf` to Canvas. There is nothing to push.

- **Keep `--no-project`**, so a broken `uv.lock` in your project cannot stop the script.
- The script opens your files read-only. It re-parses your raw file independently, checks your tables against it, and embeds your code, schema and `REPORT.md`.
- **Read the PDF before uploading.** It prints your score by group. If a check fails, fix it and rerun, and submit anyway if you run out of time.
- The PDF prints the script's sha256, which matches <https://kitchingroup.cheme.cmu.edu/f26-06763/a03-evidence.py.sha256> when run from the URL.

## Grading
Five groups are scored by the script and the sixth by your TA. Within a group, checks are equally weighted.

| # | Group | Pts | Checks |
|---|---|---|---|
| 1 | **Collection** | 0.75 | raw NDJSON with the birth message; at least 20 plant hours of event time; `receivedAt` non-decreasing; raw file not rewritten |
| 2 | **Decode and reshape** | 1.0 | `tidy.parquet` with the eight columns and types; Polars `LazyFrame` with one `.collect()`; three clocks distinct and populated; more than one `event_time` per `server_time`; record count reconciles with the raw file |
| 3 | **Validate** | 1.0 | pandera schema generated from the birth tags; `lazy=True`; `Bad`/`Uncertain` in `quarantine.parquet` and not in `tidy.parquet`; report counts reconcile |
| 4 | **Dedup and watermark** | 1.25 | `(tag, event_time)` unique; duplicate count matches the script's; `watermark.csv` with several allowances and completeness and still-waiting columns; clock comparison in the report |
| 5 | **Grid and windows** | 1.0 | regular 180 s grid, one column per tag, no empty cells, fill record; pre-fill hole counts reconcile; wide `windows.csv` on 30-minute event-time buckets that the script can recompute; incomplete last window handled; disturbance interval named with tags |
| 6 | **REPORT.md** | 1.0 | the six sections and the AI-use line, read by your TA |

A check the script cannot decide is held for your TA rather than lost. Your TA can adjust in either direction when the work is clearly better or worse than the checks show. Report numbers that do not reconcile with your raw file will be asked about.

## AI use
Generative AI is allowed with disclosure in `REPORT.md`. Editing the generated PDF by hand is falsifying a submission. You must be able to explain your dedup key, your watermark choice, and why windowing on `received_at` gives a different answer.

## Stretch (not graded)
- Collect for an hour and check whether your watermark still holds over six plant days.
- Collect at the same time as a classmate and compare files: same plant, different losses.
- Reconstruct each analyser's sampling schedule from when its `sourceTimestamp` advances.
- Replay the raw file one message at a time with a bounded buffer, and report what a watermark would have emitted and when.
