File size: 29,144 Bytes
199bfa3
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
"""
Operation Registry and Dispatcher for Advanced Data Explorer.

All post-processing operations the execution engine can dispatch.
The LLM QueryPlanner references this registry to build execution plans.
Operations are whitelisted — no arbitrary code execution.
"""
import pandas as pd
import numpy as np
from typing import Any, Dict, List, Optional
from dataclasses import dataclass, field

from analysis.advanced_analysis import AdvancedAnalyzer


# ── Result Dataclass ──────────────────────────────────────────────

@dataclass
class OperationOutput:
    """Result of a single operation step."""
    op_name: str
    label: str
    data: Any                                   # Primary result (dict, list, DataFrame)
    metadata: Dict = field(default_factory=dict) # Stats about the operation
    exportable_df: Optional[pd.DataFrame] = None # If this step produces downloadable data


# ── Operation Registry ────────────────────────────────────────────
# Each entry describes one whitelisted operation for the LLM prompt.

OPERATION_REGISTRY = {
    "find_peak_pressure_blocks": {
        "description": "Find contiguous time blocks where a pressure sensor exceeds a threshold",
        "params": {
            "pressure_tag": {"type": "str", "required": True, "description": "Sensor tag (e.g. PT130)"},
            "threshold_bar": {"type": "float", "required": True, "description": "Pressure threshold in bar"},
        },
        "returns": "Dict with blocks list: start_time, end_time, duration, peak, avg for each block",
    },
    "detect_constant_periods": {
        "description": "Find periods where a single sensor value is constant within tolerance for a minimum duration",
        "params": {
            "tag": {"type": "str", "required": True, "description": "Sensor tag name"},
            "tolerance": {"type": "float", "required": False, "default": 0.01},
            "min_duration_minutes": {"type": "float", "required": False, "default": 2.0},
        },
        "returns": "List of constant period dicts: start, end, duration_minutes, value",
    },
    "detect_constant_periods_multi": {
        "description": "Find periods where multiple sensors are simultaneously constant (overlapping constant windows)",
        "params": {
            "tags": {"type": "List[str]", "required": True, "description": "Sensor tags that must all be constant"},
            "tolerance": {"type": "float", "required": False, "default": 0.01},
            "min_duration_minutes": {"type": "float", "required": False, "default": 2.0},
        },
        "returns": "List of overlapping constant period dicts",
    },
    "downsample_frequency": {
        "description": "Reduce data frequency (e.g. 10Hz to 2Hz). Requires raw resolution data.",
        "params": {
            "target_hz": {"type": "float", "required": True, "description": "Target frequency in Hz"},
        },
        "returns": "Downsampled DataFrame",
    },
    "detect_pressure_ramps": {
        "description": "Detect windows where pressure ramps from one level to another",
        "params": {
            "pressure_tag": {"type": "str", "required": True, "description": "Pressure sensor tag"},
            "start_bar": {"type": "float", "required": True, "description": "Ramp start pressure in bar"},
            "end_bar": {"type": "float", "required": True, "description": "Ramp end pressure in bar"},
            "max_ramps": {"type": "int", "required": False, "default": None, "description": "Max ramps to return"},
        },
        "returns": "List of ramp dicts: start_time, end_time, start_pressure, peak_pressure, duration_minutes",
    },
    "extract_windows": {
        "description": "Extract sub-DataFrames for each time window from a previous operation's result",
        "params": {
            "source": {"type": "str", "required": True, "description": "Must be 'previous_result'"},
        },
        "returns": "List of DataFrames, one per detected window",
    },
    "filter_by_condition": {
        "description": "Filter data rows where a sensor meets a condition",
        "params": {
            "tag": {"type": "str", "required": True, "description": "Sensor tag"},
            "operator": {"type": "str", "required": True, "description": "One of: gt, gte, lt, lte, eq"},
            "value": {"type": "float", "required": True, "description": "Threshold value"},
        },
        "returns": "Filtered DataFrame",
    },
    "compute_statistics": {
        "description": "Compute summary statistics (min, max, mean, std, count) for sensors",
        "params": {
            "tags": {"type": "List[str]", "required": False, "default": None, "description": "Tags to summarize (default: all)"},
        },
        "returns": "Statistics dict per sensor",
    },
    "detect_and_analyze_fills": {
        "description": "Detect fill/test cycles using motor speed, then compute per-fill metrics: kWh/kg, total mass, pump strokes, peak pressure. Can filter by kWh/kg threshold or specific fill ID.",
        "params": {
            "kwh_per_kg_threshold": {"type": "float", "required": False, "default": None, "description": "Only return fills where kWh/kg exceeds this"},
            "fill_id": {"type": "int", "required": False, "default": None, "description": "Analyze only this fill number"},
            "include_strokes": {"type": "bool", "required": False, "default": True, "description": "Whether to compute stroke count"},
        },
        "returns": "List of fill dicts: fill_id, start_time, end_time, duration_min, peak_pressure, avg_flow, total_mass_kg, kwh_per_kg, total_strokes",
    },
    "count_pump_strokes": {
        "description": "Count total pump strokes (motor revolutions) in the time window by integrating RPM over time",
        "params": {
            "speed_tag": {"type": "str", "required": False, "default": None, "description": "Motor speed tag (defaults to first available)"},
        },
        "returns": "Dict: total_revolutions, avg_rpm, active_minutes, peak_rpm, note",
    },
}


def get_registry_for_prompt() -> str:
    """Serialize the operation registry as text for LLM prompt injection."""
    lines = ["AVAILABLE OPERATIONS:"]
    for op_name, info in OPERATION_REGISTRY.items():
        lines.append(f"\n  {op_name}: {info['description']}")
        lines.append("    Parameters:")
        for pname, pinfo in info["params"].items():
            req = "REQUIRED" if pinfo.get("required") else f"optional, default={pinfo.get('default')}"
            lines.append(f"      {pname} ({pinfo['type']}, {req}): {pinfo.get('description', '')}")
        lines.append(f"    Returns: {info['returns']}")
    return "\n".join(lines)


# ── Operation Dispatcher ──────────────────────────────────────────

OPERATOR_MAP = {"gt": ">", "gte": ">=", "lt": "<", "lte": "<=", "eq": "=="}
MOTOR_SPEED_TAGS = ["M130_Speed", "MC130_VFD_Speed", "M130_RPM"]


class OperationDispatcher:
    """Dispatch whitelisted operations against a DataFrame."""

    def __init__(self):
        self.analyzer = AdvancedAnalyzer()

    def run(
        self,
        op_name: str,
        df: pd.DataFrame,
        params: Dict,
        prev_result: Any = None,
        db_connector=None,
    ) -> OperationOutput:
        """Dispatch a single operation by name."""
        handler = getattr(self, f"_run_{op_name}", None)
        if handler is None:
            raise ValueError(f"Unknown operation: {op_name}")

        # Some operations need extra context
        if op_name in ("detect_and_analyze_fills",):
            return handler(df, params, db_connector=db_connector)
        if op_name in ("extract_windows",):
            return handler(df, params, prev_result=prev_result)

        return handler(df, params)

    # ── Handlers ──────────────────────────────────────────────────

    def _run_find_peak_pressure_blocks(self, df: pd.DataFrame, params: Dict) -> OperationOutput:
        pressure_tag = params["pressure_tag"]
        threshold = float(params["threshold_bar"])
        result = self.analyzer.find_peak_pressure_blocks(df, pressure_tag, threshold)

        # Build exportable table from blocks
        export_df = None
        blocks = result.get("blocks", [])
        if blocks:
            export_df = pd.DataFrame(blocks)

        return OperationOutput(
            op_name="find_peak_pressure_blocks",
            label=f"{pressure_tag} > {threshold} bar",
            data=result,
            metadata={"num_blocks": result.get("num_blocks", 0), "total_points": result.get("total_points_above", 0)},
            exportable_df=export_df,
        )

    def _run_detect_constant_periods(self, df: pd.DataFrame, params: Dict) -> OperationOutput:
        tag = params["tag"]
        tolerance = float(params.get("tolerance", 0.01))
        min_dur = float(params.get("min_duration_minutes", 2.0))
        periods = self.analyzer.detect_constant_periods(df, tag, tolerance, min_dur)

        export_df = pd.DataFrame(periods) if periods else None

        return OperationOutput(
            op_name="detect_constant_periods",
            label=f"{tag} constant periods (±{tolerance}, ≥{min_dur} min)",
            data=periods,
            metadata={"num_periods": len(periods)},
            exportable_df=export_df,
        )

    def _run_detect_constant_periods_multi(self, df: pd.DataFrame, params: Dict) -> OperationOutput:
        tags = params["tags"]
        tolerance = float(params.get("tolerance", 0.01))
        min_dur = float(params.get("min_duration_minutes", 2.0))

        # Detect constant periods per tag
        per_tag_periods = {}
        for tag in tags:
            if tag in df.columns:
                per_tag_periods[tag] = self.analyzer.detect_constant_periods(df, tag, tolerance, min_dur)
            else:
                per_tag_periods[tag] = []

        # Find overlapping intervals across all tags
        overlapping = self._intersect_periods(per_tag_periods, min_dur)

        export_df = pd.DataFrame(overlapping) if overlapping else None

        return OperationOutput(
            op_name="detect_constant_periods_multi",
            label=f"Simultaneous constant: {', '.join(tags)} (≥{min_dur} min)",
            data=overlapping,
            metadata={"num_overlaps": len(overlapping), "per_tag_counts": {t: len(p) for t, p in per_tag_periods.items()}},
            exportable_df=export_df,
        )

    def _intersect_periods(self, per_tag: Dict[str, List[Dict]], min_dur: float) -> List[Dict]:
        """Intersect constant periods across multiple tags — return only overlapping windows."""
        tags = list(per_tag.keys())
        if not tags or any(len(per_tag[t]) == 0 for t in tags):
            return []

        # Start with first tag's periods
        current_intervals = [(p["start"], p["end"]) for p in per_tag[tags[0]]]

        # Intersect with each subsequent tag
        for tag in tags[1:]:
            other = [(p["start"], p["end"]) for p in per_tag[tag]]
            current_intervals = self._pairwise_intersect(current_intervals, other)
            if not current_intervals:
                return []

        # Filter by minimum duration and build result dicts
        result = []
        for start, end in current_intervals:
            dur_min = (end - start).total_seconds() / 60
            if dur_min >= min_dur:
                # Gather each tag's value during this window
                tag_values = {}
                for tag in tags:
                    for p in per_tag[tag]:
                        if p["start"] <= start and p["end"] >= end:
                            tag_values[tag] = p["value"]
                            break

                result.append({
                    "start": start,
                    "end": end,
                    "duration_minutes": dur_min,
                    "tag_values": tag_values,
                })
        return result

    @staticmethod
    def _pairwise_intersect(intervals_a, intervals_b):
        """Compute pairwise intersection of two sorted interval lists."""
        result = []
        i, j = 0, 0
        while i < len(intervals_a) and j < len(intervals_b):
            a_start, a_end = intervals_a[i]
            b_start, b_end = intervals_b[j]
            # Overlap
            start = max(a_start, b_start)
            end = min(a_end, b_end)
            if start < end:
                result.append((start, end))
            # Advance the interval that ends first
            if a_end <= b_end:
                i += 1
            else:
                j += 1
        return result

    def _run_downsample_frequency(self, df: pd.DataFrame, params: Dict) -> OperationOutput:
        target_hz = float(params["target_hz"])

        if "timestamp" not in df.columns or len(df) < 2:
            return OperationOutput("downsample_frequency", f"Downsample to {target_hz} Hz", df, {"error": "Insufficient data"})

        # Estimate current frequency
        sorted_df = df.sort_values("timestamp")
        dt = sorted_df["timestamp"].diff().dt.total_seconds().dropna()
        current_hz = 1.0 / dt.median() if dt.median() > 0 else 1.0

        if target_hz >= current_hz:
            # Already at or below target frequency
            return OperationOutput(
                "downsample_frequency", f"Data already at {current_hz:.1f} Hz (target: {target_hz} Hz)",
                sorted_df, {"current_hz": current_hz, "target_hz": target_hz, "rows_before": len(sorted_df), "rows_after": len(sorted_df)},
                exportable_df=sorted_df,
            )

        step = max(1, int(round(current_hz / target_hz)))
        downsampled = sorted_df.iloc[::step].reset_index(drop=True)

        return OperationOutput(
            op_name="downsample_frequency",
            label=f"Downsampled {current_hz:.1f} Hz → {target_hz} Hz",
            data=downsampled,
            metadata={"current_hz": current_hz, "target_hz": target_hz, "step": step,
                       "rows_before": len(sorted_df), "rows_after": len(downsampled)},
            exportable_df=downsampled,
        )

    def _run_detect_pressure_ramps(self, df: pd.DataFrame, params: Dict) -> OperationOutput:
        pressure_tag = params["pressure_tag"]
        start_bar = float(params["start_bar"])
        end_bar = float(params["end_bar"])
        max_ramps = params.get("max_ramps")

        if pressure_tag not in df.columns:
            return OperationOutput("detect_pressure_ramps", f"Ramps in {pressure_tag}", [],
                                   {"error": f"Tag {pressure_tag} not found"})

        sorted_df = df[["timestamp", pressure_tag]].dropna(subset=[pressure_tag]).sort_values("timestamp")
        timestamps = sorted_df["timestamp"].values
        values = sorted_df[pressure_tag].values

        # State machine: IDLE → RAMPING → COMPLETE
        ramps = []
        state = "IDLE"
        ramp_start_idx = 0
        ramp_max_val = 0.0

        for i in range(len(values)):
            if state == "IDLE":
                # Look for crossing start_bar from below
                if values[i] >= start_bar * 0.9 and values[i] <= start_bar * 1.3:
                    state = "RAMPING"
                    ramp_start_idx = i
                    ramp_max_val = values[i]
            elif state == "RAMPING":
                ramp_max_val = max(ramp_max_val, values[i])

                # Ramp complete: crossed end_bar
                if values[i] >= end_bar:
                    ramps.append({
                        "start_time": pd.Timestamp(timestamps[ramp_start_idx]),
                        "end_time": pd.Timestamp(timestamps[i]),
                        "start_pressure": float(values[ramp_start_idx]),
                        "peak_pressure": float(ramp_max_val),
                        "duration_minutes": (pd.Timestamp(timestamps[i]) - pd.Timestamp(timestamps[ramp_start_idx])).total_seconds() / 60,
                    })
                    state = "IDLE"
                    continue

                # Ramp aborted: pressure dropped significantly below start
                if values[i] < start_bar * 0.5:
                    state = "IDLE"

        if max_ramps and len(ramps) > int(max_ramps):
            ramps = ramps[:int(max_ramps)]

        export_df = pd.DataFrame(ramps) if ramps else None

        return OperationOutput(
            op_name="detect_pressure_ramps",
            label=f"{pressure_tag} ramps {start_bar}{end_bar} bar",
            data=ramps,
            metadata={"num_ramps": len(ramps)},
            exportable_df=export_df,
        )

    def _run_extract_windows(self, df: pd.DataFrame, params: Dict, prev_result: Any = None) -> OperationOutput:
        """Slice DataFrame into sub-DataFrames per detected window from previous operation."""
        windows = []
        if isinstance(prev_result, OperationOutput):
            data = prev_result.data
            if isinstance(data, list):
                windows = data
            elif isinstance(data, dict):
                windows = data.get("blocks", data.get("ramps", []))

        if not windows:
            return OperationOutput("extract_windows", "Extract windows", [],
                                   {"error": "No windows found in previous result"})

        dfs = []
        labels = []
        for i, w in enumerate(windows):
            start = w.get("start_time", w.get("start"))
            end = w.get("end_time", w.get("end"))
            if start is None or end is None:
                continue

            mask = (df["timestamp"] >= pd.Timestamp(start)) & (df["timestamp"] <= pd.Timestamp(end))
            window_df = df[mask].copy()
            if not window_df.empty:
                dfs.append(window_df)
                labels.append(f"Window {i + 1}")

        return OperationOutput(
            op_name="extract_windows",
            label=f"Extracted {len(dfs)} windows",
            data=dfs,
            metadata={"num_windows": len(dfs), "labels": labels,
                       "rows_per_window": [len(d) for d in dfs]},
            exportable_df=None,  # Multi-window: handled by export_utils
        )

    def _run_filter_by_condition(self, df: pd.DataFrame, params: Dict) -> OperationOutput:
        tag = params["tag"]
        operator = params["operator"]
        value = float(params["value"])

        if tag not in df.columns:
            return OperationOutput("filter_by_condition", f"Filter {tag}", df,
                                   {"error": f"Tag {tag} not found"})

        op_str = OPERATOR_MAP.get(operator)
        if not op_str:
            return OperationOutput("filter_by_condition", f"Filter {tag}", df,
                                   {"error": f"Unknown operator: {operator}. Use: gt, gte, lt, lte, eq"})

        if operator == "gt":
            mask = df[tag] > value
        elif operator == "gte":
            mask = df[tag] >= value
        elif operator == "lt":
            mask = df[tag] < value
        elif operator == "lte":
            mask = df[tag] <= value
        else:
            mask = df[tag] == value

        filtered = df[mask].copy()

        return OperationOutput(
            op_name="filter_by_condition",
            label=f"{tag} {op_str} {value}",
            data=filtered,
            metadata={"rows_before": len(df), "rows_after": len(filtered)},
            exportable_df=filtered,
        )

    def _run_compute_statistics(self, df: pd.DataFrame, params: Dict) -> OperationOutput:
        tags = params.get("tags")
        numeric_cols = [c for c in df.columns if c != "timestamp" and pd.api.types.is_numeric_dtype(df[c])]
        if tags:
            numeric_cols = [c for c in numeric_cols if c in tags]

        stats = {}
        for col in numeric_cols:
            data = df[col].dropna()
            if len(data) > 0:
                stats[col] = {
                    "count": int(len(data)),
                    "min": float(data.min()),
                    "max": float(data.max()),
                    "mean": float(data.mean()),
                    "std": float(data.std()),
                }

        export_df = pd.DataFrame(stats).T.reset_index().rename(columns={"index": "Sensor"}) if stats else None

        return OperationOutput(
            op_name="compute_statistics",
            label=f"Statistics for {len(numeric_cols)} sensors",
            data=stats,
            metadata={"num_sensors": len(stats)},
            exportable_df=export_df,
        )

    def _run_detect_and_analyze_fills(self, df: pd.DataFrame, params: Dict, db_connector=None) -> OperationOutput:
        """Detect fill cycles and compute per-fill metrics."""
        from analysis.pump_cycle_detector import PumpCycleDetector

        kwh_threshold = params.get("kwh_per_kg_threshold")
        target_fill_id = params.get("fill_id")
        include_strokes = params.get("include_strokes", True)

        # Find motor speed tag in DataFrame
        speed_tag = None
        for tag in MOTOR_SPEED_TAGS:
            if tag in df.columns:
                speed_tag = tag
                break

        if not speed_tag:
            return OperationOutput("detect_and_analyze_fills", "Detect fills", [],
                                   {"error": "No motor speed tag found in data. Include M130_Speed or MC130_VFD_Speed."})

        # Prepare motor speed series for cycle detector
        motor_df = df[["timestamp", speed_tag]].dropna(subset=[speed_tag]).rename(columns={speed_tag: "value"})

        detector = PumpCycleDetector()
        raw_cycles = detector.detect_cycles(motor_df)

        if not raw_cycles:
            return OperationOutput("detect_and_analyze_fills", "Detect fills", [],
                                   {"error": "No cycles detected in the time range"})

        # Enrich with pressure/flow/temp from the fetched data (no additional DB query needed)
        enriched = self._enrich_fills_from_df(raw_cycles, df, db_connector)

        # Compute per-fill advanced metrics
        fill_results = []
        for c in enriched:
            fill_dict = {
                "fill_id": c["cycle_id"],
                "start_time": c["start_time"],
                "end_time": c["end_time"],
                "duration_min": c.get("duration_minutes", 0),
                "peak_pressure_bar": c.get("peak_pressure", None),
                "avg_flow_kg_min": c.get("avg_flow", None),
                "total_mass_kg": c.get("total_mass_kg", None),
                "kwh_per_kg": None,
                "total_revolutions": None,
            }

            # kWh/kg computation
            if "avg_power_kw" in c and c.get("avg_flow") and c["avg_flow"] > 0:
                # kWh/kg = AvgPower [kW] * 60 / FT140 [kg/min]
                kwh_kg = (c["avg_power_kw"] * 60) / c["avg_flow"] if c["avg_flow"] > 0.001 else None
                fill_dict["kwh_per_kg"] = round(kwh_kg, 4) if kwh_kg is not None else None

            # Stroke count
            if include_strokes and "avg_rpm" in c and c.get("duration_minutes"):
                revolutions = c["avg_rpm"] * c["duration_minutes"]
                fill_dict["total_revolutions"] = int(round(revolutions))

            fill_results.append(fill_dict)

        # Apply filters
        if target_fill_id is not None:
            target_id = int(target_fill_id)
            matched = [f for f in fill_results if f["fill_id"] == target_id]
            if not matched:
                available = [f["fill_id"] for f in fill_results]
                return OperationOutput("detect_and_analyze_fills", f"Fill {target_id}", [],
                                       {"error": f"Fill {target_id} not found. Available fill IDs: {available}"})
            fill_results = matched

        if kwh_threshold is not None:
            threshold = float(kwh_threshold)
            fill_results = [f for f in fill_results if f.get("kwh_per_kg") is not None and f["kwh_per_kg"] > threshold]

        export_df = pd.DataFrame(fill_results) if fill_results else None

        note = "Revolutions = RPM × minutes. For simplex: strokes = revolutions. For triplex: strokes = 3 × revolutions."

        return OperationOutput(
            op_name="detect_and_analyze_fills",
            label=f"Fill analysis ({len(fill_results)} fills)",
            data=fill_results,
            metadata={"num_fills": len(fill_results), "note": note},
            exportable_df=export_df,
        )

    def _enrich_fills_from_df(self, cycles: List[Dict], df: pd.DataFrame, db_connector) -> List[Dict]:
        """Enrich cycle dicts with metrics computed from the already-fetched wide DataFrame."""
        enriched = []
        for c in cycles:
            ec = c.copy()
            mask = (df["timestamp"] >= c["start_time"]) & (df["timestamp"] <= c["end_time"])
            cdf = df[mask]

            if cdf.empty:
                enriched.append(ec)
                continue

            # Peak pressure
            if "PT130" in cdf.columns:
                pt130 = cdf["PT130"].dropna()
                if len(pt130) > 0:
                    ec["peak_pressure"] = float(pt130.max())
                    ec["initial_pressure"] = float(pt130.iloc[0])

            # Average flow
            if "FT140" in cdf.columns:
                ft140 = cdf["FT140"].dropna()
                if len(ft140) > 0:
                    ec["avg_flow"] = float(ft140.mean())
                    # Total mass via trapezoidal integration
                    flow_sorted = cdf[["timestamp", "FT140"]].dropna().sort_values("timestamp")
                    if len(flow_sorted) > 1:
                        dt_min = flow_sorted["timestamp"].diff().dt.total_seconds().fillna(0) / 60
                        ec["total_mass_kg"] = float((flow_sorted["FT140"] * dt_min).sum())

            # Average power
            if "AvgPower" in cdf.columns:
                power = cdf["AvgPower"].dropna()
                if len(power) > 0:
                    ec["avg_power_kw"] = float(power.mean())

            # Motor speed
            for tag in MOTOR_SPEED_TAGS:
                if tag in cdf.columns:
                    speed = cdf[tag].dropna()
                    if len(speed) > 0:
                        ec["avg_rpm"] = float(speed.mean())
                        ec["peak_rpm"] = float(speed.max())
                        break

            # Peak temperatures
            for tag in ["TT110", "TT130"]:
                if tag in cdf.columns:
                    data = cdf[tag].dropna()
                    if len(data) > 0:
                        ec[f"peak_{tag}"] = float(data.max())

            enriched.append(ec)
        return enriched

    def _run_count_pump_strokes(self, df: pd.DataFrame, params: Dict) -> OperationOutput:
        """Count total pump revolutions by integrating RPM over time."""
        speed_tag = params.get("speed_tag")

        if not speed_tag:
            for tag in MOTOR_SPEED_TAGS:
                if tag in df.columns:
                    speed_tag = tag
                    break

        if not speed_tag or speed_tag not in df.columns:
            return OperationOutput("count_pump_strokes", "Pump strokes", {},
                                   {"error": "No motor speed tag found. Include M130_Speed or MC130_VFD_Speed."})

        speed_df = df[["timestamp", speed_tag]].dropna(subset=[speed_tag]).sort_values("timestamp")

        if len(speed_df) < 2:
            return OperationOutput("count_pump_strokes", "Pump strokes", {},
                                   {"error": "Insufficient motor speed data"})

        rpm_values = speed_df[speed_tag].values
        timestamps = speed_df["timestamp"].values

        # Integrate: revolutions = sum(RPM_i × Δt_i) where Δt in minutes
        dt_minutes = pd.Series(timestamps).diff().dt.total_seconds().fillna(0).values / 60
        total_revolutions = float(np.sum(rpm_values * dt_minutes))
        active_mask = rpm_values > 10  # RPM > idle threshold
        active_minutes = float(np.sum(dt_minutes[active_mask]))

        result = {
            "total_revolutions": int(round(total_revolutions)),
            "avg_rpm": float(np.mean(rpm_values[active_mask])) if active_mask.any() else 0,
            "peak_rpm": float(np.max(rpm_values)) if len(rpm_values) > 0 else 0,
            "active_minutes": round(active_minutes, 1),
            "speed_tag_used": speed_tag,
            "note": "Revolutions = RPM × time. For simplex: strokes = revolutions. For triplex: strokes = 3 × revolutions.",
        }

        return OperationOutput(
            op_name="count_pump_strokes",
            label=f"Pump strokes ({speed_tag})",
            data=result,
            metadata={"total_revolutions": result["total_revolutions"]},
        )