#!/usr/bin/env catnip
# Journalisation d'un pipeline d'import avec logging.
# Python émet les événements ; Catnip orchestre les étapes, reçoit chaque
# enregistrement via un filtre appelable et produit le bilan.
#
# DEPS: aucune (stdlib logging)
# OFFLINE: aucun accès réseau
logging = import('logging')
# Un seul handler sur le logger racine 'batch' : les loggers enfants propagent.
# Le filtre est une lambda Catnip appelée par Python pour chaque enregistrement.
# Les logs partent sur stderr (défaut StreamHandler) : les print sont flushés
# pour garder l'ordre d'affichage quand les deux flux sont fusionnés.
handler = logging.StreamHandler()
handler.setFormatter(logging.Formatter('%(levelname)-7s %(name)-16s %(message)s'))
level_counts = dict()
count_filter = (record) => {
level_counts[record.levelname] = level_counts.get(record.levelname, 0) + 1
True
}
handler.addFilter(count_filter)
root = logging.getLogger('batch')
root.setLevel(logging.DEBUG)
root.addHandler(handler)
log_extract = logging.getLogger('batch.extract')
log_validate = logging.getLogger('batch.validate')
log_aggregate = logging.getLogger('batch.aggregate')
# Les acceptations individuelles ne sont pas journalisées, seuls les rejets.
log_validate.setLevel(logging.INFO)
struct Reading { id: str; sensor: str; value: float }
lines = list(
'R01;température;21.4',
'R02;humidité;55.2',
'R03;température;-273.0',
'R04;humidité',
'R05;pression;1013.25',
'R06;température;abc',
'R07;pression;998.0',
'R08;humidité;210.5',
'R09;luminosité;850.0',
)
# Étape 1 : extraction. Les lignes malformées sont ignorées avec un warning.
extract = (line: str): Reading | None => {
parts = line.split(';')
log_extract.debug(f"{line} → {len(parts)} champs")
result = None
if len(parts) != 3 {
log_extract.warning(f"ligne malformée ignorée : {line}")
} else {
try {
result = Reading(parts[0], parts[1], float(parts[2]))
} except {
_e: ValueError => { log_extract.warning(f"valeur non numérique ignorée : {line}") }
}
}
result
}
# Étape 2 : validation. Plage physique admise, propre à chaque sonde. Une sonde
# hors catalogue est rejetée plutôt que validée sur la plage d'une autre.
validate = (r: Reading): Reading | None => {
bounds = match r.sensor {
'température' => { tuple(-50.0, 60.0) }
'humidité' => { tuple(0.0, 100.0) }
'pression' => { tuple(950.0, 1050.0) }
_ => { None }
}
result = r
if bounds == None {
log_validate.error(f"{r.id} sonde inconnue ({r.sensor}) : exclu")
result = None
} else {
if r.value < bounds[0] or r.value > bounds[1] {
log_validate.error(f"{r.id} hors plage ({r.value}) : exclu")
result = None
}
}
result
}
# Étape 3 : agrégation. Un info par sonde résume le lot.
print("⇒ Import du lot", flush=True)
readings = lines.[extract]
valid = list()
for r in readings {
if r != None {
checked = validate(r)
if checked != None { valid.append(checked) }
}
}
sums = dict()
counts = dict()
sensors = list()
for r in valid {
if r.sensor not in sensors { sensors.append(r.sensor) }
sums[r.sensor] = sums.get(r.sensor, 0.0) + r.value
counts[r.sensor] = counts.get(r.sensor, 0) + 1
}
for sensor in sensors {
mean = round(sums[sensor] / counts[sensor], 2)
log_aggregate.info(f"{sensor}: {counts[sensor]} lecture(s), moyenne {mean}")
}
# Bilan alimenté par le filtre : Python a rappelé Catnip pour chaque ligne.
print()
print("⇒ Bilan des niveaux enregistrés", flush=True)
for level in list('DEBUG', 'INFO', 'WARNING', 'ERROR') {
print(f" {level:<7}: {level_counts.get(level, 0)}")
}
oracle_ok = level_counts.get('DEBUG', 0) == 9
oracle_ok = oracle_ok and level_counts.get('INFO', 0) == 3
oracle_ok = oracle_ok and level_counts.get('WARNING', 0) == 2
oracle_ok = oracle_ok and level_counts.get('ERROR', 0) == 3
oracle_ok = oracle_ok and len(valid) == 4
print()
print(f"⇒ Oracle (niveaux comptés + 4 lectures valides) : {oracle_ok}")