"""Streaming normalization of configurable digital CSV input.""" import csv from decimal import Decimal, InvalidOperation import itertools import re from .files import ImportCancelled def normalize_csv(source, destination, options, sample_rate, progress, cancel): """Write canonical CSV without expanding transition-only input into samples.""" # Нормализация отделена от упаковки SAL/DSL: каждый выходной формат # получает один и тот же CSV с временем в секундах и выбранными каналами. pattern = re.compile(r'(?:time|timestamp)\s*(?:[\[(](s|ms|us|µs|ns)[\])])?', re.I) # Decimal сохраняет десятичную сетку входа при смене единиц времени. # str() не переносит в Decimal двоичную погрешность float из Qt. rate = Decimal(str(sample_rate)) if not rate.is_finite() or rate < 0: raise ValueError('Частота должна быть конечной и неотрицательной.') mode = options.get('mode', 'auto') duration = options.get('duration') if duration is not None and (mode != 'events' or not rate): raise ValueError('Длительность требует режима переходов и заданной частоты.') # Оба потока закрываются и при отмене. Временным файлом владеет вызывающий # convert_capture: он удалит его вместе с временным каталогом. with source.open(encoding='utf-8-sig', newline='') as stream, destination.open('w', encoding='utf-8', newline='') as output: size = max(1, source.stat().st_size) consumed = 0 # Читаем последовательно, не разворачивая событийную запись в отсчёты. # Проверка отмены ограничивает задержку реакции на кнопку в интерфейсе. def lines(): nonlocal consumed for index, line in enumerate(stream): consumed += len(line) if index % 4096 == 0: if cancel(): raise ImportCancelled() progress(min(99, consumed * 100 // size)) if line.strip() and not line.lstrip().startswith(('#', ';')): yield line source_lines = lines() first = next(source_lines, '') if not first: raise ValueError('CSV не содержит данных.') # Явный разделитель имеет приоритет. Авто оценивает первую строку; # для неоднозначного CSV оператор может выбрать разделитель вручную. delimiter = options.get('delimiter') or max((',', ';', '\t'), key=first.count) first_row = next(csv.reader([first], delimiter=delimiter)) no_header = options.get('no_header', False) # При отсутствии заголовка первая строка остаётся данными. D0/D1 — # имена исходных колонок, поэтому время тоже может называться D0. names = ['D%d' % i for i in range(len(first_row))] if no_header else [x.strip() for x in first_row] if len(set(names)) != len(names) or any(not n for n in names): raise ValueError('Имена колонок должны быть непустыми и уникальными.') time_column = options.get('time_column', 'auto') # Автоматически выбираем только узнаваемое имя времени. При нескольких # кандидатах нельзя молча взять первый: это изменило бы шкалу записи. if time_column == 'auto': candidates = [n for n in names if pattern.fullmatch(n)] if len(candidates) > 1: raise ValueError('Найдено несколько колонок времени; выберите одну.') time_column = candidates[0] if candidates else None elif time_column == 'none': time_column = None if time_column is not None and time_column not in names: raise ValueError('Колонка времени не найдена: ' + time_column) if time_column is None and (not rate or mode == 'events'): raise ValueError('Без колонки времени задайте частоту; режим переходов требует времени.') # Порядок списка задаёт выходные номера каналов. Служебные колонки # не обязаны быть цифровыми, если пользователь исключил их из списка. channels = options.get('channels') or [n for n in names if n != time_column] if not 1 <= len(channels) <= 64 or len(set(channels)) != len(channels): raise ValueError('Выберите от 1 до 64 различных цифровых каналов.') if any(n not in names or n == time_column for n in channels): raise ValueError('Канал отсутствует или совпадает с колонкой времени.') # Индексы вычисляются один раз, а не поиском имён для каждой строки. indices = [names.index(n) for n in channels] time_index = names.index(time_column) if time_column is not None else None unit = options.get('time_unit', 'auto') # Явная единица позволяет читать нестандартную колонку t. Без суффикса # auto означает секунды, как и обычный импортёр цифровых записей. if unit == 'auto': match = re.search(r'[\[(](s|ms|us|µs|ns)[\])]$', time_column or '', re.I) unit = match[1].lower() if match else 's' scale = Decimal({'s': '1', 'ms': '.001', 'us': '.000001', 'µs': '.000001', 'ns': '.000000001'}[unit]) rows = csv.reader(source_lines, delimiter=delimiter) # Возвращаем уже прочитанную первую строку в ленивый итератор. if no_header: rows = itertools.chain([first_row], rows) writer = csv.writer(output) writer.writerow(['Time [s]'] + channels) previous = origin = None for index, row in enumerate(rows): try: if len(row) != len(names): raise ValueError('число колонок отличается от заголовка') # Без времени строки — последовательные отсчёты с нуля. # При наличии времени сохраняем исходное начало, даже отрицательное. timestamp = Decimal(row[time_index].strip().replace(',', '.')) * scale if time_index is not None else Decimal(index) / rate if not timestamp.is_finite() or (previous is not None and timestamp <= previous): raise ValueError('время должно быть конечным и строго возрастать') if origin is None: origin = timestamp # Проверка равномерности относится только к отсчётам с заданной # частотой; событийная запись вправе иметь длинные паузы. if mode == 'samples' and rate and abs((timestamp-origin)*rate-index) > Decimal('.00001'): raise ValueError('временные метки не соответствуют частоте отсчётов') levels = [row[i].strip() for i in indices] if any(v not in ('0', '1') for v in levels): raise ValueError('выбранные каналы должны содержать 0 или 1') writer.writerow([str(timestamp)] + levels) previous = timestamp except (ValueError, InvalidOperation) as error: raise ValueError('CSV, строка данных %d: %s' % (index+1, error)) from error if previous is None: raise ValueError('CSV не содержит данных.') # Длительность задаёт исключительную границу: N отсчётов занимают # N/rate секунд, но последний находится в (N-1)/rate. Дописываем # удержание последнего уровня, не создавая нового фронта. if duration is not None: count = Decimal(str(duration)) * rate if not count.is_finite() or count <= 0 or abs(count-count.to_integral_value()) > Decimal('.00001'): raise ValueError('Длительность должна быть положительной и попадать на сетку частоты.') last_sample = origin + (count.to_integral_value()-1)/rate if previous > last_sample: raise ValueError('Длительность заканчивается до последнего отсчёта записи.') if previous < last_sample: writer.writerow([str(last_sample)] + levels) if cancel(): raise ImportCancelled()