id: cold-chain-temperature-excursions
namespace: company.team
description: |
Turn raw fridge, freezer and storage-room temperature logs into a list of
cold-chain excursions: when each unit went out of range, for how long, how
hot or cold it got, whether a sensor stopped reporting, and which readings
could not be read, so the affected stock can be quarantined and checked.
inputs:
- id: readings_file
type: FILE
required: false
description: CSV with the columns site, unit_id, unit_type (FRIDGE, FREEZER or
AMBIENT), reading_time (YYYY-MM-DD HH:MM) and temperature_c. Leave empty
to run against the built-in sample sensor log.
- id: min_excursion_minutes
type: INT
defaults: 30
description: Out-of-range periods shorter than this are reported as BRIEF (for
example a door left open while loading) instead of CRITICAL.
- id: max_gap_minutes
type: INT
defaults: 60
description: A unit that sends no reading for longer than this is reported as a
SENSOR_GAP, because nobody knows what the temperature was during that
time.
- id: notify_slack
type: BOOL
defaults: false
description: Post the findings to Slack when there is at least one critical
excursion, sensor gap or data error. Requires the SLACK_WEBHOOK_URL
secret.
triggers:
- id: daily_cold_chain_review
type: io.kestra.plugin.core.trigger.Schedule
description: Review the previous day's sensor log every morning at 07:00 IST,
before stock from the affected units is picked or dispatched. Shipped
disabled. Enable it once the flow reads your real sensor export instead of
the sample.
cron: "0 7 * * *"
timezone: Asia/Kolkata
disabled: true
tasks:
- id: sample_readings
type: io.kestra.plugin.core.storage.Write
description: A 12-hour sensor log for four units, one reading every 15 minutes,
built with a Pebble loop. It contains a long high excursion, a brief
door-open spike, a freezing excursion, a sensor gap, an unreadable reading
and an unknown unit type. Only used when no readings_file is provided.
extension: .csv
content: |
site,unit_id,unit_type,reading_time,temperature_c
{% set hot = [8.6, 9.1, 9.7, 10.4, 11.2, 10.6, 9.4, 8.3] -%}
{% set cold = [1.4, 0.8, 1.1] -%}
{% for i in range(0, 47) -%}
{% set t = "2026-10-01 " ~ ((i / 4) | numberFormat("00")) ~ ":" ~ ((i % 4 * 15) | numberFormat("00")) -%}
{% set fridge1 = (i >= 20 and i <= 27) ? hot[i - 20] : 4.0 + (i % 5) * 0.3 -%}
{% set fridge2 = (i == 14) ? 8.9 : ((i >= 32 and i <= 34) ? cold[i - 32] : 5.0 + (i % 3) * 0.2) -%}
Pune DC,FRG-01,FRIDGE,{{ t }},{{ fridge1 | numberFormat("0.0") }}
Pune DC,FRG-02,FRIDGE,{{ t }},{{ fridge2 | numberFormat("0.0") }}
{% if i < 16 or i > 25 -%}
Pune DC,FRZ-01,FREEZER,{{ t }},{{ (-20.0 + (i % 4) * 0.5) | numberFormat("0.0") }}
{% endif -%}
Pune DC,AMB-01,AMBIENT,{{ t }},{{ (22.0 + (i % 6) * 0.4) | numberFormat("0.0") }}
{% endfor -%}
Pune DC,FRG-01,FRIDGE,2026-10-01 06:07,ERR
Pune DC,FRG-03,CHILLER,2026-10-01 06:00,4.5
- id: find_excursions
type: io.kestra.plugin.jdbc.duckdb.Queries
description: Match every reading to the allowed range of its unit type, group
consecutive out-of-range readings into excursions, find gaps between
readings, and collect unreadable rows in DuckDB. The event report is
exported to CSV and the last query returns the summary.
inputFiles:
readings.csv: "{{ inputs.readings_file ?? outputs.sample_readings.uri }}"
outputFiles:
- excursion_report.csv
fetchType: FETCH_ONE
sql: |
CREATE TABLE ranges AS
SELECT * FROM (VALUES
('FRIDGE', 2.0, 8.0),
('FREEZER', -25.0, -15.0),
('AMBIENT', 15.0, 25.0)
) AS t(unit_type, min_c, max_c);
CREATE TABLE raw AS
SELECT
trim(site) AS site,
trim(unit_id) AS unit_id,
upper(trim(coalesce(unit_type, ''))) AS unit_type,
trim(coalesce(reading_time, '')) AS reading_time_raw,
trim(coalesce(temperature_c, '')) AS temperature_raw,
try_cast(trim(reading_time) AS TIMESTAMP) AS ts,
try_cast(trim(temperature_c) AS DOUBLE) AS temp
FROM read_csv('{{ workingDir }}/readings.csv', header = true, all_varchar = true);
CREATE TABLE readings AS
SELECT
r.*,
g.min_c,
g.max_c,
CASE WHEN r.temp > g.max_c THEN 'HIGH' WHEN r.temp < g.min_c THEN 'LOW' ELSE 'IN' END AS band
FROM raw r
JOIN ranges g ON g.unit_type = r.unit_type
WHERE r.ts IS NOT NULL AND r.temp IS NOT NULL;
CREATE TABLE sequenced AS
SELECT
*,
lead(ts) OVER (PARTITION BY site, unit_id ORDER BY ts) AS next_ts,
row_number() OVER (PARTITION BY site, unit_id ORDER BY ts)
- row_number() OVER (PARTITION BY site, unit_id, band ORDER BY ts) AS island
FROM readings;
CREATE TABLE excursions AS
SELECT
site,
unit_id,
unit_type,
CASE WHEN band = 'HIGH' THEN 'HIGH_EXCURSION' ELSE 'LOW_EXCURSION' END AS event,
min(ts) AS start_time,
coalesce(max_by(next_ts, ts), max(ts)) AS end_time,
count(*) AS readings,
CASE WHEN band = 'HIGH' THEN max(temp) ELSE min(temp) END AS peak_c,
any_value(min_c) AS min_c,
any_value(max_c) AS max_c
FROM sequenced
WHERE band <> 'IN'
GROUP BY site, unit_id, unit_type, band, island;
CREATE TABLE gaps AS
SELECT site, unit_id, unit_type, 'SENSOR_GAP' AS event, ts AS start_time, next_ts AS end_time, 0 AS readings, NULL::DOUBLE AS peak_c, min_c, max_c
FROM sequenced
WHERE date_diff('minute', ts, next_ts) > {{ inputs.max_gap_minutes }}
UNION ALL
SELECT site, unit_id, any_value(unit_type), 'SENSOR_GAP', max(ts), (SELECT max(ts) FROM readings), 0, NULL::DOUBLE, any_value(min_c), any_value(max_c)
FROM readings
GROUP BY site, unit_id
HAVING date_diff('minute', max(ts), (SELECT max(ts) FROM readings)) > {{ inputs.max_gap_minutes }};
CREATE TABLE bad_rows AS
SELECT
site,
unit_id,
unit_type,
list_filter([
CASE WHEN unit_type NOT IN (SELECT unit_type FROM ranges) THEN 'unknown unit_type "' || unit_type || '" (use FRIDGE, FREEZER or AMBIENT)' END,
CASE WHEN ts IS NULL THEN 'reading_time "' || reading_time_raw || '" is not a timestamp' END,
CASE WHEN temp IS NULL THEN 'temperature "' || temperature_raw || '" is not a number' END
], lambda x: x IS NOT NULL) AS problems
FROM raw;
CREATE TABLE report AS
SELECT * FROM (
SELECT
site,
unit_id,
unit_type,
event,
status,
start_time,
end_time,
duration_minutes,
readings,
peak_c,
allowed_range,
action
FROM (
SELECT
*,
CASE
WHEN event = 'SENSOR_GAP' THEN 'SENSOR_GAP'
WHEN duration_minutes >= {{ inputs.min_excursion_minutes }} THEN 'CRITICAL'
ELSE 'BRIEF'
END AS status,
CASE
WHEN event = 'SENSOR_GAP' THEN 'No reading for ' || duration_minutes || ' min: check the logger and treat the stock as unverified for this period'
WHEN event = 'HIGH_EXCURSION' AND duration_minutes >= {{ inputs.min_excursion_minutes }} THEN 'Above ' || max_c || ' C for ' || duration_minutes || ' min, peak ' || peak_c || ' C: quarantine the stock and check stability data before use'
WHEN event = 'LOW_EXCURSION' AND duration_minutes >= {{ inputs.min_excursion_minutes }} THEN 'Below ' || min_c || ' C for ' || duration_minutes || ' min, lowest ' || peak_c || ' C: quarantine the stock, freeze-sensitive products may be damaged'
ELSE 'Out of range for ' || duration_minutes || ' min (peak ' || peak_c || ' C): likely a door opening, check if it repeats'
END AS action
FROM (
SELECT
*,
date_diff('minute', start_time, end_time) AS duration_minutes,
min_c || ' to ' || max_c || ' C' AS allowed_range
FROM (SELECT * FROM excursions UNION ALL SELECT * FROM gaps)
)
)
UNION ALL
SELECT
site,
unit_id,
unit_type,
'UNREADABLE' AS event,
'DATA_ERROR' AS status,
NULL AS start_time,
NULL AS end_time,
NULL AS duration_minutes,
count(*) AS readings,
NULL AS peak_c,
NULL AS allowed_range,
count(*) || ' row(s) skipped: ' || array_to_string(list_distinct(flatten(list(problems))), ', ') AS action
FROM bad_rows
WHERE len(problems) > 0
GROUP BY site, unit_id, unit_type
)
ORDER BY
CASE status WHEN 'CRITICAL' THEN 0 WHEN 'SENSOR_GAP' THEN 1 WHEN 'DATA_ERROR' THEN 2 ELSE 3 END,
duration_minutes DESC NULLS LAST,
site,
unit_id;
COPY report TO '{{ outputFiles["excursion_report.csv"] }}' (HEADER, DELIMITER ',');
SELECT
(SELECT count(*) FROM raw) AS total_readings,
(SELECT count(DISTINCT site || '/' || unit_id) FROM readings) AS units_checked,
(SELECT count(DISTINCT site || '/' || unit_id) FROM readings)
- (SELECT count(DISTINCT site || '/' || unit_id) FROM report WHERE status IN ('CRITICAL', 'SENSOR_GAP', 'DATA_ERROR') AND unit_type IN (SELECT unit_type FROM ranges)) AS units_ok,
count(*) FILTER (WHERE status = 'CRITICAL') AS critical_count,
count(*) FILTER (WHERE status = 'BRIEF') AS brief_count,
count(*) FILTER (WHERE status = 'SENSOR_GAP') AS gap_count,
count(*) FILTER (WHERE status = 'DATA_ERROR') AS data_error_count,
coalesce(max(duration_minutes) FILTER (WHERE status = 'CRITICAL'), 0) AS longest_excursion_minutes,
coalesce(string_agg(site || ' / ' || unit_id || ' from ' || strftime(start_time, '%Y-%m-%d %H:%M') || ': ' || action, chr(10)) FILTER (WHERE status = 'CRITICAL'), '') AS critical_details,
coalesce(string_agg(site || ' / ' || unit_id || ' from ' || strftime(start_time, '%Y-%m-%d %H:%M') || ': ' || action, chr(10)) FILTER (WHERE status = 'SENSOR_GAP'), '') AS gap_details,
coalesce(string_agg(site || ' / ' || unit_id || ': ' || action, chr(10)) FILTER (WHERE status = 'DATA_ERROR'), '') AS data_error_details,
coalesce(string_agg(site || ' / ' || unit_id || ' from ' || strftime(start_time, '%Y-%m-%d %H:%M') || ': ' || action, chr(10)) FILTER (WHERE status = 'BRIEF'), '') AS brief_details
FROM report;
- id: route_findings
type: io.kestra.plugin.core.flow.If
description: Branch on the result. Critical excursions, sensor gaps and
unreadable data mean some stock may be compromised, so they are logged as
warnings and optionally sent to Slack. A clean log is recorded too,
together with any brief spikes.
condition: "{{ outputs.find_excursions.outputs[0].row.critical_count +
outputs.find_excursions.outputs[0].row.gap_count +
outputs.find_excursions.outputs[0].row.data_error_count > 0 }}"
then:
- id: log_findings
type: io.kestra.plugin.core.log.Log
level: WARN
description: Write every critical excursion, sensor gap, data error and brief
spike to the execution logs with the action to take.
message: |
Cold chain review: {{ outputs.find_excursions.outputs[0].row.critical_count }} critical excursion(s), {{ outputs.find_excursions.outputs[0].row.gap_count }} sensor gap(s), {{ outputs.find_excursions.outputs[0].row.data_error_count }} unit(s) with unreadable data, {{ outputs.find_excursions.outputs[0].row.brief_count }} brief spike(s). {{ outputs.find_excursions.outputs[0].row.units_ok }} of {{ outputs.find_excursions.outputs[0].row.units_checked }} units clean, {{ outputs.find_excursions.outputs[0].row.total_readings }} readings.
Critical:
{{ outputs.find_excursions.outputs[0].row.critical_details }}
Sensor gaps:
{{ outputs.find_excursions.outputs[0].row.gap_details }}
Data errors:
{{ outputs.find_excursions.outputs[0].row.data_error_details }}
Brief spikes:
{{ outputs.find_excursions.outputs[0].row.brief_details }}
- id: slack_enabled
type: io.kestra.plugin.core.flow.If
description: Only call Slack when the notify_slack input is true, so the flow
runs out of the box without any secret.
condition: "{{ inputs.notify_slack }}"
then:
- id: alert_quality_team
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
description: Tell the warehouse and quality team which units to check and which
stock to quarantine.
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
messageText: |
:thermometer: Cold chain review: {{ outputs.find_excursions.outputs[0].row.critical_count }} critical excursion(s), {{ outputs.find_excursions.outputs[0].row.gap_count }} sensor gap(s), {{ outputs.find_excursions.outputs[0].row.data_error_count }} unit(s) with unreadable data.
{{ outputs.find_excursions.outputs[0].row.critical_details }}
{{ outputs.find_excursions.outputs[0].row.gap_details }}
Full report: execution {{ execution.id }} in flow {{ flow.namespace }}.{{ flow.id }}
else:
- id: log_clean
type: io.kestra.plugin.core.log.Log
description: Record a clean review, including any brief spikes, so the execution
history shows the cold chain was checked.
message: "Cold chain review passed: {{
outputs.find_excursions.outputs[0].row.units_checked }} units, {{
outputs.find_excursions.outputs[0].row.total_readings }} readings, no
critical excursion or sensor gap, {{
outputs.find_excursions.outputs[0].row.brief_count }} brief spike(s)."
errors:
- id: log_review_failure
type: io.kestra.plugin.core.log.Log
level: ERROR
description: Make a failed review loud. If the file is missing or malformed, no
reading was checked, and a silent failure looks like a cold chain with no
excursions.
message: "Cold chain review FAILED in {{ flow.namespace }}.{{ flow.id }}
(execution {{ execution.id }}). No readings were checked. Make sure
readings_file is a CSV with the columns site, unit_id, unit_type,
reading_time and temperature_c."
outputs:
- id: excursion_report
type: FILE
description: One row per event (HIGH_EXCURSION, LOW_EXCURSION, SENSOR_GAP or
UNREADABLE) with the status (CRITICAL, BRIEF, SENSOR_GAP or DATA_ERROR),
start and end time, duration, peak temperature, allowed range and the
action to take.
value: "{{ outputs.find_excursions.outputFiles['excursion_report.csv'] }}"
- id: excursion_summary
type: JSON
description: 'Counts and details, e.g. {"units_checked": 4, "critical_count": 2,
"gap_count": 1, "data_error_count": 2, "longest_excursion_minutes": 120,
...}.'
value: "{{ outputs.find_excursions.outputs[0].row | toJson }}"