From c50dab64cbc216e14a6231d52665314422d592bd Mon Sep 17 00:00:00 2001 From: Diego Freitas Date: Thu, 16 Apr 2026 18:38:14 -0300 Subject: [PATCH 1/2] ajustes no script de captura em raw --- Python/OAK/datasets/_0_capture_raw.py | 11 ++++++----- 1 file changed, 6 insertions(+), 5 deletions(-) diff --git a/Python/OAK/datasets/_0_capture_raw.py b/Python/OAK/datasets/_0_capture_raw.py index 26c4f02bf..48dd8f236 100644 --- a/Python/OAK/datasets/_0_capture_raw.py +++ b/Python/OAK/datasets/_0_capture_raw.py @@ -186,6 +186,11 @@ def main(): cam.configure_fps(20) cam.start_streaming() + apply_ir_comp = True + ir_k_r = 0.8 + ir_k_g = 0.4 + ir_k_b = 0.9 + while True: t0 = time.time() raw4_base, dbg = cam.grab_raw4(out_h=RAW_H, out_w=RAW_W, timeout_ms=2000, do_ae=True) @@ -196,10 +201,6 @@ def main(): gain_a = dbg.get("gain_a", None) gain_d = dbg.get("gain_d", None) - apply_ir_comp = True - ir_k_r = 0.4 - ir_k_g = 0.1 - ir_k_b = 0.5 bgr = make_bgr_preview_from_raw(raw4_base, rgirb=True, preview_fast=upscale > 0, preview_scale=upscale, apply_ir_comp=apply_ir_comp, ir_k_r=ir_k_r, ir_k_g=ir_k_g, ir_k_b=ir_k_b) # FPS @@ -305,7 +306,7 @@ def main(): }, "note": "manual", } - rgb_clean_bgr_save = make_bgr_preview_from_raw(raw4_base, rgirb=True, preview_fast=False) + rgb_clean_bgr_save = make_bgr_preview_from_raw(raw4_base, rgirb=True, preview_fast=False, apply_ir_comp=apply_ir_comp, ir_k_r=ir_k_r, ir_k_g=ir_k_g, ir_k_b=ir_k_b) raw_path, _, _ = save_sample_raw4(session_dir, raw4_base, rgb_clean_bgr_save, meta) last_msg = f"SALVO (manual): {os.path.basename(raw_path)}" last_msg_t = time.time() From 3addf89bd910f46f09f80e674ef8a26565ca19b3 Mon Sep 17 00:00:00 2001 From: Diego Freitas Date: Fri, 17 Apr 2026 17:03:16 -0300 Subject: [PATCH 2/2] =?UTF-8?q?Adicionado=20scripts=20do=20m=C3=B3dulo=20m?= =?UTF-8?q?ultiespectral=20pi?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- Python/raspi/capture_dataset.py | 645 +++++++++++++++++++++++ Python/raspi/multispectral_service.py | 121 ++++- Python/raspi/pi/camera_manager.py | 278 ++++++++++ Python/raspi/pi/frame_service.py | 239 +++++++++ Python/raspi/pi/main.py | 5 + Python/raspi/pi/protocol.py | 35 ++ Python/raspi/pi/raw_processor_core.py | 120 +++++ Python/raspi/pi/raw_processor_preview.py | 125 +++++ Python/raspi/pi/server.py | 488 +++++++++++++++++ Python/raspi/pi/state.py | 155 ++++++ Python/raspi/pi/stream_sender.py | 327 ++++++++++++ Python/raspi/pi/trigger_manager.py | 74 +++ Python/raspi/stream_receiver.py | 112 +++- Python/raspi/test_capture_raw.py | 19 - Python/raspi/test_stream_raw.py | 55 -- 15 files changed, 2681 insertions(+), 117 deletions(-) create mode 100644 Python/raspi/capture_dataset.py create mode 100644 Python/raspi/pi/camera_manager.py create mode 100644 Python/raspi/pi/frame_service.py create mode 100644 Python/raspi/pi/main.py create mode 100644 Python/raspi/pi/protocol.py create mode 100644 Python/raspi/pi/raw_processor_core.py create mode 100644 Python/raspi/pi/raw_processor_preview.py create mode 100644 Python/raspi/pi/server.py create mode 100644 Python/raspi/pi/state.py create mode 100644 Python/raspi/pi/stream_sender.py create mode 100644 Python/raspi/pi/trigger_manager.py delete mode 100644 Python/raspi/test_capture_raw.py delete mode 100644 Python/raspi/test_stream_raw.py diff --git a/Python/raspi/capture_dataset.py b/Python/raspi/capture_dataset.py new file mode 100644 index 000000000..4d96765ad --- /dev/null +++ b/Python/raspi/capture_dataset.py @@ -0,0 +1,645 @@ +import os +import time +import json +import argparse +from datetime import datetime + +import cv2 +import numpy as np + +from multispectral_service import MultiSpectralService +from stream_receiver import StreamReceiver +from pi.raw_processor_core import RawProcessorCore +from pi.raw_processor_preview import RawProcessorPreview + + +STREAM_PORT = 6001 +PI_HOST = "192.168.105.6" +PC_HOST = "192.168.105.5" + + +# ========================= +# Helpers gerais +# ========================= + +def ts_name() -> str: + return datetime.now().strftime("%Y%m%d_%H%M%S_%f")[:-3] + + +def overlay_hud( + img_bgr: np.ndarray, + lines: list[str], + base_h: int = 720, + base_font_scale: float = 0.75, + base_line_step: int = 28, +): + h, w = img_bgr.shape[:2] + + scale = h / float(base_h) + scale = max(scale, 0.4) + + font_scale = base_font_scale * scale + line_step = int(base_line_step * scale) + + thick_outline = max(1, int(3 * scale)) + thick_text = max(1, int(2 * scale)) + + y = int(24 * scale) + x = int(12 * scale) + + for s in lines: + cv2.putText(img_bgr, s, (x, y), cv2.FONT_HERSHEY_SIMPLEX, font_scale, (0, 0, 0), thick_outline, cv2.LINE_AA) + cv2.putText(img_bgr, s, (x, y), cv2.FONT_HERSHEY_SIMPLEX, font_scale, (255, 255, 255), thick_text, cv2.LINE_AA) + y += line_step + + +def save_sample( + base_dir: str, + frame_type: str, + preview_bgr: np.ndarray, + meta: dict, + raw3: np.ndarray | None = None, + packed_raw: np.ndarray | None = None, +): + """ + Salva conforme o tipo de frame: + + - RGB: + * payload em .raw float32 (3,H,W) + * preview em .png + * metadados em .json + + - RAW_BRUTO: + * payload packed RAW10 em .bin + * preview em .png + * metadados em .json + """ + os.makedirs(base_dir, exist_ok=True) + name = ts_name() + + png_path = os.path.join(base_dir, f"{name}.png") + json_path = os.path.join(base_dir, f"{name}.json") + + if frame_type == "RGB": + if raw3 is None: + raise ValueError("raw3 não pode ser None quando frame_type='RGB'") + + payload_path = os.path.join(base_dir, f"{name}.raw") + raw3.astype(np.float32).tofile(payload_path) + + meta["saved_payload_type"] = "raw3" + meta["saved_payload_path"] = os.path.basename(payload_path) + meta["saved_payload_dtype"] = "float32" + meta["saved_payload_shape"] = list(raw3.shape) + + elif frame_type == "RAW_BRUTO": + if packed_raw is None: + raise ValueError("packed_raw não pode ser None quando frame_type='RAW_BRUTO'") + + payload_path = os.path.join(base_dir, f"{name}.bin") + packed_raw.tofile(payload_path) + + meta["saved_payload_type"] = "raw10_packed" + meta["saved_payload_path"] = os.path.basename(payload_path) + meta["saved_payload_dtype"] = str(packed_raw.dtype) + meta["saved_payload_shape"] = list(packed_raw.shape) + + else: + raise ValueError(f"frame_type não suportado para save: {frame_type}") + + cv2.imwrite(png_path, preview_bgr) + + with open(json_path, "w", encoding="utf-8") as f: + json.dump(meta, f, ensure_ascii=False, indent=2) + + return payload_path, png_path, json_path + + +# ========================= +# MAIN +# ========================= + +def main(): + parser = argparse.ArgumentParser( + description="Captura de dataset RAW3 usando módulo multispectral Pi + StreamReceiver.", + formatter_class=argparse.ArgumentDefaultsHelpFormatter, + ) + + parser.add_argument("--cana", required=True, choices=["baixa", "media", "alta"], help="Estado da cana no momento da coleta.") + parser.add_argument("--horario", required=True, choices=["cedo", "meio_dia", "entardecer", "nublado"], help="Janela de iluminação / horário da coleta.") + parser.add_argument("--out_root", default="dataset", help="Pasta raiz do dataset.") + parser.add_argument("--modelo", default="imx296_pi", help="Nome do módulo/câmera para montar a pasta.") + parser.add_argument("--pi_host", default=PI_HOST, help="IP do servidor no Raspberry Pi.") + parser.add_argument("--pc_host", default=PC_HOST, help="IP local do notebook/PC que receberá o stream.") + parser.add_argument("--stream_port", type=int, default=STREAM_PORT, help="Porta TCP do receiver de stream.") + parser.add_argument("--server_port", type=int, default=5000, help="Porta TCP do servidor de comandos no Pi.") + parser.add_argument("--fps", type=int, default=20, help="FPS desejado.") + parser.add_argument("--width", type=int, default=640, help="Largura óptica da câmera.") + parser.add_argument("--height", type=int, default=480, help="Altura óptica da câmera.") + parser.add_argument("--interval", type=float, default=1.0, help="Intervalo em segundos para auto-save quando ligado.") + parser.add_argument("--preview_upscale", type=int, default=2, help="Fator de upscale visual do preview.") + parser.add_argument("--bayer", default="GBRG", choices=["GBRG", "GRBG", "RGGB", "BGGR"], help="Padrão Bayer da câmera.") + parser.add_argument("--codec_family", default="numcodecs", help="Apenas informativo no metadata local.") + parser.add_argument("--codec_name", default="blosc", help="Apenas informativo no metadata local.") + parser.add_argument("--frame_type", default="RAW_BRUTO", choices=["RAW_BRUTO", "RGB"], help="Tipo de payload pedido ao Pi.") + parser.add_argument("--output_dtype", default="uint8", choices=["uint8", "float32"], help="Dtype do payload processado no Pi.") + + args = parser.parse_args() + + raw_w = args.width + raw_h = args.height + + # dataset//brutas/cana_/// + session_dir = os.path.join( + args.modelo, + args.out_root, + "brutas", + f"cana_{args.cana}", + args.horario, + datetime.now().strftime("%Y%m%d"), + ) + os.makedirs(session_dir, exist_ok=True) + + print("============================================") + print("Coleta de dataset RAW3 - Módulo Multiespectral") + print(f"Cana : {args.cana}") + print(f"Horário : {args.horario}") + print(f"Saída : {session_dir}") + print(f"Sensor : {raw_w}x{raw_h} | Bayer={args.bayer}") + print("============================================") + + window_name = "Dataset Capture - RAW3 (C/SPACE=save | A=auto-save | M=preview scale | Q=quit)" + cv2.namedWindow(window_name, cv2.WINDOW_NORMAL) + + auto_save = False + last_auto_t = 0.0 + preview_upscale = args.preview_upscale + + t_view_fps = time.time() + view_frames = 0 + fps_view = 0.0 + + t_stream_fps = time.time() + last_stream_frame_id = None + stream_frames_accum = 0 + fps_stream = 0.0 + + last_msg = "" + last_msg_t = 0.0 + + receiver = StreamReceiver(host="0.0.0.0", port=args.stream_port) + svc = MultiSpectralService(host=args.pi_host, port=args.server_port, timeout=10) + + processor_core = RawProcessorCore( + sensor_width=raw_w, + sensor_height=raw_h, + bayer_pattern=args.bayer, + ) + processor_preview = RawProcessorPreview( + sensor_width=raw_w, + sensor_height=raw_h, + bayer_pattern=args.bayer, + ) + + last_frame_id = -1 + last_raw3 = None + last_packed_raw = None + last_preview_bgr = None + last_meta_stream = None + + try: + receiver.start() + time.sleep(0.5) + + svc.connect() + + modes_resp = svc.get_sensor_modes() + if not modes_resp.get("ok"): + raise RuntimeError("Falha ao obter sensor_modes") + + for mode in modes_resp["sensor_modes"]: + print( + f"[{mode['index']}] " + f"size={mode['size']} " + f"format={mode['format']} " + f"bit_depth={mode['bit_depth']} " + f"fps={mode['fps']}" + ) + + print("SET RES:", svc.set_resolution(raw_w, raw_h)) + print("SET FPS:", svc.set_fps(args.fps)) + print("SET BAYER:", svc.set_bayer(args.bayer)) + print("BEGIN:", svc.begin(frame_type=args.frame_type, output_dtype=args.output_dtype)) + print("START STREAM:", svc.start_stream(args.pc_host, args.stream_port, fps=args.fps)) + + camera_ctrl = svc.get_camera_controls() + + ae_enabled = bool(camera_ctrl.get("ae_enable", True)) + awb_enabled = bool(camera_ctrl.get("awb_enable", True)) + manual_exposure_us = camera_ctrl.get("exposure_time_us", None) + manual_gain = camera_ctrl.get("analogue_gain", None) + manual_colour_gains = camera_ctrl.get("colour_gains", None) + + expected_packed_w = (args.width * 10) // 8 + expected_packed_h = args.height + + while True: + t0 = time.time() + + meta = receiver.last_meta + frame = receiver.last_frame + + if meta is not None and frame is not None and meta.get("frame_id") != last_frame_id: + last_frame_id = meta["frame_id"] + #print(meta) + + if False and meta["packed_width"] != expected_packed_w or meta["packed_height"] != expected_packed_h: + def infer_sensor_size_from_packed(packed_width: int, packed_height: int): + width = int((packed_width * 8) / 10) + height = packed_height + return width, height + actual_packed_w = meta["packed_width"] + actual_packed_h = meta["packed_height"] + + padding = actual_packed_w - expected_packed_w + if actual_packed_h != args.height: + erro = True + elif actual_packed_w < expected_packed_w: + erro = True + elif padding > 64: + # margem conservadora, se quiser + erro = True + else: + erro = False + + if erro: + inferred_w = int((actual_packed_w * 8) / 10) + inferred_h = actual_packed_h + + modes_txt = [] + if modes_resp.get("ok"): + for m in modes_resp["sensor_modes"]: + size = m.get("size") + fmt = m.get("format") + fps = m.get("fps") + modes_txt.append(f"- {size[0]}x{size[1]} | {fmt} | fps={fps}") + modes_str = "\n".join(modes_txt) if modes_txt else "(não disponível)" + + RuntimeError( + f"Modo RAW inesperado.\n" + f"Solicitado: {args.width}x{args.height} (packed útil esperado {expected_packed_w})\n" + f"Recebido: packed {actual_packed_h}x{actual_packed_w}\n" + f"Obs: packed_width pode incluir padding/stride.\n" + f"Se a altura confere e o packed recebido é maior que o esperado, " + f"o modo pode estar correto com alinhamento de memória." + ) + + try: + frame_type = meta.get("frame_type", "RAW_BRUTO") + dtype_str = meta.get("dtype") or meta.get("output_dtype", "uint8") + + if frame_type == "RAW_BRUTO": + packed = frame + if packed.ndim == 3 and packed.shape[2] == 1: + packed = packed[:, :, 0] + + raw16 = processor_core.unpack_raw10_packed(packed) + preview_bgr = processor_preview.raw16_to_preview_bgr( + raw16, + bit_depth=meta.get("source_bit_depth", 10), + ) + + # continua gerando raw3 apenas para visualização/depuração local, se quiser manter + raw3 = processor_core.build_training_rgb( + raw16, + output_dtype="float32", + bit_depth=meta.get("source_bit_depth", 10), + ) + + last_packed_raw = packed.copy() + + elif frame_type == "RGB": + rgb_chw = frame + if rgb_chw.ndim != 3: + raise RuntimeError(f"Frame RGB inválido: shape={rgb_chw.shape}") + + if dtype_str == "uint8": + raw3 = rgb_chw.astype(np.float32) / 255.0 + elif dtype_str == "float32": + raw3 = rgb_chw.astype(np.float32) + else: + raise RuntimeError(f"dtype RGB não suportado: {dtype_str}") + + preview_rgb = np.transpose(raw3, (1, 2, 0)) + preview_bgr = cv2.cvtColor( + np.clip(preview_rgb * 255.0, 0, 255).astype(np.uint8), + cv2.COLOR_RGB2BGR + ) + + last_packed_raw = None + + else: + raise RuntimeError(f"frame_type não suportado neste script: {frame_type}") + + if preview_upscale and preview_upscale > 1: + preview_show = cv2.resize( + preview_bgr, + (preview_bgr.shape[1] * preview_upscale, preview_bgr.shape[0] * preview_upscale), + interpolation=cv2.INTER_NEAREST, + ) + else: + preview_show = preview_bgr.copy() + + # ========================= + # FPS do stream (frames recebidos) + # ========================= + curr_frame_id = meta.get("frame_id") + + if last_stream_frame_id is not None and curr_frame_id is not None: + delta_ids = curr_frame_id - last_stream_frame_id + if delta_ids > 0: + stream_frames_accum += delta_ids + + last_stream_frame_id = curr_frame_id + + dt_stream = time.time() - t_stream_fps + if dt_stream >= 1.0: + fps_stream = stream_frames_accum / dt_stream + stream_frames_accum = 0 + t_stream_fps = time.time() + + # ========================= + # FPS de visualização/processamento no PC + # ========================= + view_frames += 1 + dt_view = time.time() - t_view_fps + if dt_view >= 1.0: + fps_view = view_frames / dt_view + view_frames = 0 + t_view_fps = time.time() + + lines = [ + f"CANA: {args.cana} | HORA: {args.horario} | Pasta: {os.path.basename(session_dir)}", + f"AutoSave: {'ON' if auto_save else 'OFF'} | Intervalo: {args.interval:.1f}s | PreviewScale: {preview_upscale}", + f"frame_id={meta.get('frame_id')} | FPS_STREAM={fps_stream:.1f} | FPS_VIEW={fps_view:.1f}", + f"codec={meta.get('codec_name', meta.get('codec_family', '-'))} | comp={meta.get('dt_comp', 0):.4f}s | send={meta.get('dt_send_payload_prev', 0):.4f}s", + f"packed={meta.get('width')}x{meta.get('height')} | raw3_shape={list(raw3.shape)}", + f"type={meta.get('frame_type')} | layout={meta.get('output_layout')} | dtype={meta.get('dtype') or meta.get('output_dtype')}", + f"AE={'ON' if ae_enabled else 'OFF'} | AWB={'ON' if awb_enabled else 'OFF'} | EXP={manual_exposure_us} | GAIN={manual_gain}", + "Keys: C/SPACE=save | A=auto-save | E=AE | W=AWB | I/K=exp | O/L=gain | R=reset | M=preview | Q/Esc=quit", + ] + overlay_hud(preview_show, lines, base_h=raw_h) + + if last_msg and (time.time() - last_msg_t) < 2.0: + cv2.putText(preview_show, last_msg, (12, preview_show.shape[0] - 18), + cv2.FONT_HERSHEY_SIMPLEX, 0.8, (0, 255, 0), 2, cv2.LINE_AA) + + cv2.imshow(window_name, preview_show) + + last_raw3 = raw3.copy() if raw3 is not None else None + last_preview_bgr = preview_bgr.copy() if preview_bgr is not None else None + last_meta_stream = dict(meta) + + except Exception as e: + err = np.zeros((500, 1200, 3), dtype=np.uint8) + cv2.putText(err, f"Erro ao processar frame: {e}", (20, 60), + cv2.FONT_HERSHEY_SIMPLEX, 0.8, (0, 0, 255), 2, cv2.LINE_AA) + cv2.imshow(window_name, err) + print(f"[ERRO FRAME] {e}") + + now = time.time() + can_save = ( + last_meta_stream is not None + and last_preview_bgr is not None + and ( + (last_meta_stream.get("frame_type") == "RGB" and last_raw3 is not None) or + (last_meta_stream.get("frame_type") == "RAW_BRUTO" and last_packed_raw is not None) + ) + ) + + if auto_save and can_save and (now - last_auto_t) >= args.interval: + frame_type_save = last_meta_stream.get("frame_type") + + meta_save = { + "ts": datetime.now().isoformat(timespec="milliseconds"), + "cana": args.cana, + "horario": args.horario, + "sensor_width": raw_w, + "sensor_height": raw_h, + "bayer_pattern": last_meta_stream.get("source_bayer_pattern"), + "fps_target": args.fps, + "frame_type": frame_type_save, + "codec_family": last_meta_stream.get("codec_family"), + "codec_name": last_meta_stream.get("codec_name"), + "codec_params": last_meta_stream.get("codec_params"), + "stream_meta": last_meta_stream, + "note": "autosave", + } + + if frame_type_save == "RGB": + meta_save["raw3_shape"] = list(last_raw3.shape) + meta_save["raw3_dtype"] = str(last_raw3.dtype) + elif frame_type_save == "RAW_BRUTO": + meta_save["packed_shape"] = list(last_packed_raw.shape) + meta_save["packed_dtype"] = str(last_packed_raw.dtype) + meta_save["packed_height"] = int(last_packed_raw.shape[0]) + meta_save["packed_width"] = int(last_packed_raw.shape[1]) + + payload_path, _, _ = save_sample( + session_dir, + frame_type=frame_type_save, + preview_bgr=last_preview_bgr, + meta=meta_save, + raw3=last_raw3, + packed_raw=last_packed_raw, + ) + + last_msg = f"SALVO (auto): {os.path.basename(payload_path)}" + last_msg_t = now + last_auto_t = now + + k = cv2.waitKey(1) & 0xFF + if k in (ord("q"), ord("Q"), 27): + break + + elif k in (ord("a"), ord("A")): + auto_save = not auto_save + last_msg = f"AutoSave -> {'ON' if auto_save else 'OFF'}" + last_msg_t = time.time() + + elif k in (ord("m"), ord("M")): + preview_upscale = 0 if preview_upscale else args.preview_upscale + last_msg = f"Preview UPSCALE -> {preview_upscale}" + last_msg_t = time.time() + + elif k in (ord("e"), ord("E")): + ae_enabled = not ae_enabled + resp = svc.set_ae_enable(ae_enabled) + ae_enabled = bool(resp.get("ae_enable", ae_enabled)) + last_msg = f"AE -> {'ON' if ae_enabled else 'OFF'}" + last_msg_t = time.time() + + elif k in (ord("w"), ord("W")): + awb_enabled = not awb_enabled + resp = svc.set_awb_enable(awb_enabled) + awb_enabled = bool(resp.get("awb_enable", awb_enabled)) + last_msg = f"AWB -> {'ON' if awb_enabled else 'OFF'}" + last_msg_t = time.time() + + elif k in (ord("i"), ord("I")): + # aumenta exposição manual + if manual_exposure_us is None: + manual_exposure_us = 15000 + else: + manual_exposure_us = min(int(manual_exposure_us * 1.15), 200000) + + if ae_enabled: + ae_enabled = False + svc.set_ae_enable(False) + + resp = svc.set_exposure_time(int(manual_exposure_us)) + manual_exposure_us = resp.get("exposure_time_us", manual_exposure_us) + last_msg = f"ExposureTime -> {manual_exposure_us} us" + last_msg_t = time.time() + + elif k in (ord("k"), ord("K")): + # diminui exposição manual + if manual_exposure_us is None: + manual_exposure_us = 15000 + else: + manual_exposure_us = max(int(manual_exposure_us / 1.15), 100) + + if ae_enabled: + ae_enabled = False + svc.set_ae_enable(False) + + resp = svc.set_exposure_time(int(manual_exposure_us)) + manual_exposure_us = resp.get("exposure_time_us", manual_exposure_us) + last_msg = f"ExposureTime -> {manual_exposure_us} us" + last_msg_t = time.time() + + elif k in (ord("o"), ord("O")): + # aumenta ganho manual + if manual_gain is None: + manual_gain = 1.0 + else: + manual_gain = min(float(manual_gain) * 1.10, 32.0) + + if ae_enabled: + ae_enabled = False + svc.set_ae_enable(False) + + resp = svc.set_analogue_gain(float(manual_gain)) + manual_gain = resp.get("analogue_gain", manual_gain) + last_msg = f"AnalogueGain -> {manual_gain:.2f}" + last_msg_t = time.time() + + elif k in (ord("l"), ord("L")): + # diminui ganho manual + if manual_gain is None: + manual_gain = 1.0 + else: + manual_gain = max(float(manual_gain) / 1.10, 1.0) + + if ae_enabled: + ae_enabled = False + svc.set_ae_enable(False) + + resp = svc.set_analogue_gain(float(manual_gain)) + manual_gain = resp.get("analogue_gain", manual_gain) + last_msg = f"AnalogueGain -> {manual_gain:.2f}" + last_msg_t = time.time() + + elif k in (ord("r"), ord("R")): + # reset manuais + svc.clear_exposure_time() + svc.clear_analogue_gain() + svc.clear_colour_gains() + + manual_exposure_us = None + manual_gain = None + manual_colour_gains = None + + last_msg = "Manual controls resetados" + last_msg_t = time.time() + + elif k in (ord("c"), ord("C"), 32): + can_save = ( + last_meta_stream is not None + and last_preview_bgr is not None + and ( + (last_meta_stream.get("frame_type") == "RGB" and last_raw3 is not None) or + (last_meta_stream.get("frame_type") == "RAW_BRUTO" and last_packed_raw is not None) + ) + ) + + if can_save: + frame_type_save = last_meta_stream.get("frame_type") + + meta_save = { + "ts": datetime.now().isoformat(timespec="milliseconds"), + "cana": args.cana, + "horario": args.horario, + "sensor_width": raw_w, + "sensor_height": raw_h, + "bayer_pattern": last_meta_stream.get("source_bayer_pattern"), + "fps_target": args.fps, + "frame_type": frame_type_save, + "codec_family": last_meta_stream.get("codec_family"), + "codec_name": last_meta_stream.get("codec_name"), + "codec_params": last_meta_stream.get("codec_params"), + "stream_meta": last_meta_stream, + "camera_controls": { + "ae_enable": ae_enabled, + "awb_enable": awb_enabled, + "exposure_time_us": manual_exposure_us, + "analogue_gain": manual_gain, + "colour_gains": manual_colour_gains, + }, + "note": "manual", + } + + if frame_type_save == "RGB": + meta_save["raw3_shape"] = list(last_raw3.shape) + meta_save["raw3_dtype"] = str(last_raw3.dtype) + elif frame_type_save == "RAW_BRUTO": + meta_save["packed_shape"] = list(last_packed_raw.shape) + meta_save["packed_dtype"] = str(last_packed_raw.dtype) + meta_save["packed_height"] = int(last_packed_raw.shape[0]) + meta_save["packed_width"] = int(last_packed_raw.shape[1]) + + payload_path, _, _ = save_sample( + session_dir, + frame_type=frame_type_save, + preview_bgr=last_preview_bgr, + meta=meta_save, + raw3=last_raw3, + packed_raw=last_packed_raw, + ) + + last_msg = f"SALVO (manual): {os.path.basename(payload_path)}" + last_msg_t = time.time() + + dt_loop = time.time() - t0 + if dt_loop < 0.001: + time.sleep(0.001) + + finally: + try: + print("STOP STREAM:", svc.stop_stream()) + except Exception: + pass + + try: + print("STOP:", svc.stop()) + except Exception: + pass + + svc.disconnect() + receiver.stop() + cv2.destroyAllWindows() + print("Fim da captura.") + + +if __name__ == "__main__": + main() \ No newline at end of file diff --git a/Python/raspi/multispectral_service.py b/Python/raspi/multispectral_service.py index 5915c9ca6..8326cb70e 100644 --- a/Python/raspi/multispectral_service.py +++ b/Python/raspi/multispectral_service.py @@ -3,6 +3,7 @@ import json import numpy as np import base64 import time +from typing import Optional class MultiSpectralService: @@ -58,6 +59,16 @@ class MultiSpectralService: return json.loads(line.strip()) + def _numpy_dtype_from_string(self, dtype_str: str): + mapping = { + "uint8": np.uint8, + "float32": np.float32, + "uint16": np.uint16, + } + if dtype_str not in mapping: + raise RuntimeError(f"dtype não suportado recebido do Pi: {dtype_str}") + return mapping[dtype_str] + def ping(self): return self._send_command({"cmd": "ping"}) @@ -67,8 +78,12 @@ class MultiSpectralService: def get_config(self): return self._send_command({"cmd": "get_config"}) - def begin(self): - return self._send_command({"cmd": "begin"}) + def begin(self, frame_type: str = "RAW_BRUTO", output_dtype: str = "uint8"): + return self._send_command({ + "cmd": "begin", + "frame_type": frame_type, + "output_dtype": output_dtype, + }) def stop(self): return self._send_command({"cmd": "stop"}) @@ -79,6 +94,12 @@ class MultiSpectralService: def set_jpeg_quality(self, quality: int): return self._send_command({"cmd": "set_jpeg_quality", "value": quality}) + def set_bayer(self, bayer_pattern: str): + return self._send_command({ + "cmd": "set_bayer", + "pattern": bayer_pattern + }) + def set_resolution(self, width: int, height: int): return self._send_command({ "cmd": "set_resolution", @@ -86,49 +107,69 @@ class MultiSpectralService: "height": height }) - def capture_jpg_disk(self): - return self._send_command({"cmd": "capture_jpg_disk"}) - - def capture_jpg_bytes(self): - return self._send_command({"cmd": "capture_jpg_bytes"}) - - def capture_jpg_base64(self): - resp = self._send_command({"cmd": "capture_jpg_base64"}) - if not resp.get("ok"): - raise RuntimeError(resp.get("error", "Falha ao capturar JPG")) - return base64.b64decode(resp["data"]) - def capture_frame_array(self): t0 = time.perf_counter() resp = self._send_command({"cmd": "capture_frame"}) if not resp.get("ok"): raise RuntimeError(resp.get("error", "Falha ao capturar frame")) - + raw = base64.b64decode(resp["data"]) - width = resp["width"] - height = resp["height"] - channels = resp["channels"] - arr = np.frombuffer(raw, dtype=np.uint8).reshape(height, width, channels) + + frame_type = resp.get("frame_type") + output_layout = resp.get("output_layout", "HWC") + dtype_str = resp.get("dtype") or resp.get("output_dtype") or "uint8" + np_dtype = self._numpy_dtype_from_string(dtype_str) + + width = int(resp.get("output_width", resp.get("width"))) + height = int(resp.get("output_height", resp.get("height"))) + channels = int(resp.get("output_channels", resp.get("channels", 1))) + + arr = np.frombuffer(raw, dtype=np_dtype) + + if output_layout == "HW": + arr = arr.reshape(height, width) + + elif output_layout == "CHW": + arr = arr.reshape(channels, height, width) + + elif output_layout == "HWC": + arr = arr.reshape(height, width, channels) + + else: + raise RuntimeError(f"Layout não suportado recebido do Pi: {output_layout}") t1 = time.perf_counter() meta = { - "frame_type": resp.get("frame_type"), + "frame_type": frame_type, + "payload_format_version": resp.get("payload_format_version"), + "width": width, "height": height, "channels": channels, - "dtype": resp.get("dtype"), + "dtype": dtype_str, + "output_dtype": resp.get("output_dtype"), + "output_layout": output_layout, + "output_channel_names": resp.get("output_channel_names"), + + "packed_width": resp.get("packed_width"), + "packed_height": resp.get("packed_height"), + "source_width": resp.get("source_width"), + "source_height": resp.get("source_height"), + "source_bayer_pattern": resp.get("source_bayer_pattern"), + "source_bit_depth": resp.get("source_bit_depth"), + "size": resp.get("size"), - # 👇 TELEMETRIA PI "ts_pi": resp.get("ts_pi"), + "ts_pi_monotonic": resp.get("ts_pi_monotonic"), "dt_trigger": resp.get("dt_trigger"), "dt_settle": resp.get("dt_settle"), "dt_capture": resp.get("dt_capture"), + "dt_process": resp.get("dt_process"), "dt_total_pi": resp.get("dt_total_pi"), - # 👇 TELEMETRIA PC "dt_total_pc": t1 - t0, } @@ -144,3 +185,37 @@ class MultiSpectralService: def stop_stream(self): return self._send_command({"cmd": "stop_stream"}) + + def get_camera_controls(self): + return self._send_command({"cmd": "get_camera_controls"}) + + def set_ae_enable(self, value: bool): + return self._send_command({"cmd": "set_ae_enable", "value": bool(value)}) + + def set_awb_enable(self, value: bool): + return self._send_command({"cmd": "set_awb_enable", "value": bool(value)}) + + def set_exposure_time(self, exposure_time_us: Optional[int] = None): + return self._send_command({"cmd": "set_exposure_time", "value": exposure_time_us}) + + def clear_exposure_time(self): + return self._send_command({"cmd": "clear_exposure_time"}) + + def set_analogue_gain(self, gain: float | None): + return self._send_command({"cmd": "set_analogue_gain", "value": gain}) + + def clear_analogue_gain(self): + return self._send_command({"cmd": "clear_analogue_gain"}) + + def set_colour_gains(self, r_gain: float, b_gain: float): + return self._send_command({ + "cmd": "set_colour_gains", + "r_gain": r_gain, + "b_gain": b_gain + }) + + def clear_colour_gains(self): + return self._send_command({"cmd": "clear_colour_gains"}) + + def get_sensor_modes(self): + return self._send_command({"cmd": "get_sensor_modes"}) diff --git a/Python/raspi/pi/camera_manager.py b/Python/raspi/pi/camera_manager.py new file mode 100644 index 000000000..529acc715 --- /dev/null +++ b/Python/raspi/pi/camera_manager.py @@ -0,0 +1,278 @@ +from pathlib import Path +from datetime import datetime +import subprocess +from picamera2 import Picamera2 +import base64 +import io +from PIL import Image +from threading import Lock, RLock +import threading +import time +import numpy as np + + +class CameraManager: + def __init__(self, state): + self.state = state + self.picam2 = None + self.initialized = False + + self.output_dir = Path("/tmp/multispec_captures") + self.output_dir.mkdir(parents=True, exist_ok=True) + + self.last_raw_frame = None + self.last_raw_frame_id = 0 + self.last_raw_frame_ts = None + self._raw_buffer = None + self.frame_lock = Lock() + self.camera_lock = RLock() + + self._stop_event = threading.Event() + self._stream_thread = None + self._sensor_modes_cache = None + + self.current_width = None + self.current_height = None + self.current_fps = None + self._reconfigure_needed = False + + def mark_reconfigure_needed(self): + with self.camera_lock: + self._reconfigure_needed = True + + def needs_reconfigure(self): + return ( + self._reconfigure_needed or + not self.initialized or + self.current_width != self.state.width or + self.current_height != self.state.height or + self.current_fps != self.state.fps + ) + + def begin(self): + with self.camera_lock: + if self.picam2 is not None and self.needs_reconfigure(): + self.stop() + + if self.initialized and self.picam2 is not None: + return True + + self.picam2 = Picamera2() + + config = self.picam2.create_video_configuration( + main={"size": (640, 480), "format": "RGB888"}, + raw={"size": (self.state.width, self.state.height)}, + buffer_count=6 + ) + + if config is None: + self.picam2 = None + raise RuntimeError("Falha ao criar configuração da câmera. Verifique se a resolução é suportada.") + + self.picam2.configure(config) + self.picam2.start() + + self.initialized = True + self.current_width = self.state.width + self.current_height = self.state.height + self.current_fps = self.state.fps + self._reconfigure_needed = False + + self.apply_controls() + + with self.frame_lock: + self.last_raw_frame = None + self.last_raw_frame_id = 0 + self.last_raw_frame_ts = None + + self._stop_event.clear() + self._stream_thread = threading.Thread(target=self._update_loop, daemon=True) + self._stream_thread.start() + + if not self.wait_first_frame(timeout_s=2.0): + self.stop() + raise RuntimeError("Câmera inicializada, mas nenhum frame foi recebido a tempo") + + return True + + def _update_loop(self): + print("[DEBUG] Loop de atualização de frames iniciado") + + while not self._stop_event.is_set(): + request = None + try: + with self.camera_lock: + if self.picam2 is None: + time.sleep(0.05) + continue + request = self.picam2.capture_request() + + raw_array = request.make_array("raw") + self.last_raw_shape = raw_array.shape + self.last_raw_dtype = raw_array.dtype + + with self.frame_lock: + if ( + self._raw_buffer is None or + self._raw_buffer.shape != raw_array.shape or + self._raw_buffer.dtype != raw_array.dtype + ): + self._raw_buffer = raw_array.copy() + else: + np.copyto(self._raw_buffer, raw_array) + + self.last_raw_frame = self._raw_buffer + self.last_raw_frame_id += 1 + self.last_raw_frame_ts = time.perf_counter() + + except Exception as e: + print(f"[ERRO LOOP] {e}") + time.sleep(0.1) + finally: + if request is not None: + try: + request.release() + except Exception: + pass + + def stop(self): + with self.camera_lock: + self._stop_event.set() + + if self._stream_thread is not None: + self._stream_thread.join(timeout=1.5) + self._stream_thread = None + + if self.picam2 is not None: + try: + self.picam2.stop() + except Exception: + pass + try: + self.picam2.close() + except Exception: + pass + + self.picam2 = None + self.initialized = False + self._reconfigure_needed = False + + self.current_width = None + self.current_height = None + self.current_fps = None + + with self.frame_lock: + self.last_raw_frame = None + + with self.frame_lock: + self.last_raw_frame = None + self.last_raw_frame_id = 0 + self.last_raw_frame_ts = None + self._raw_buffer = None + + def capture_raw_frame(self): + with self.frame_lock: + if self.last_raw_frame is None: + return None, 0, 0, 0, 0, None + + frame = self.last_raw_frame + frame_id = self.last_raw_frame_id + frame_ts = self.last_raw_frame_ts + h, w = frame.shape[:2] + + return { + "frame": frame, + "width": w, + "height": h, + "channels": 1, + "frame_id": frame_id, + "timestamp": frame_ts, + } + + def _build_controls_from_state(self): + controls = {} + + frame_us = int(1_000_000 / max(1, self.state.fps or 10)) + controls["FrameDurationLimits"] = (frame_us, frame_us) + + controls["AeEnable"] = bool(self.state.ae_enable) + controls["AwbEnable"] = bool(self.state.awb_enable) + + if not self.state.ae_enable: + if self.state.exposure_time_us is not None: + controls["ExposureTime"] = int(self.state.exposure_time_us) + if self.state.analogue_gain is not None: + controls["AnalogueGain"] = float(self.state.analogue_gain) + + if not self.state.awb_enable and self.state.colour_gains is not None: + r_gain, b_gain = self.state.colour_gains + controls["ColourGains"] = (float(r_gain), float(b_gain)) + + return controls + + def apply_controls(self): + with self.camera_lock: + if self.picam2 is None: + return False + + controls = self._build_controls_from_state() + try: + self.picam2.set_controls(controls) + return True + except Exception as e: + print(f"[ERRO CONTROLS] {e} | controls={controls}") + return False + + def get_sensor_modes(self): + if self._sensor_modes_cache: + return self._sensor_modes_cache + + temp_picam2 = None + try: + if self.picam2 is not None: + cam = self.picam2 + else: + temp_picam2 = Picamera2() + cam = temp_picam2 + + modes = cam.sensor_modes + result = [] + + for i, m in enumerate(modes): + fmt = str(m.get("format")) if m.get("format") is not None else None + size = m.get("size") + bit_depth = m.get("bit_depth") + fps = m.get("fps") + crop_limits = m.get("crop_limits") + exposure_limits = m.get("exposure_limits") + + result.append({ + "index": i, + "format": fmt, + "size": list(size) if size is not None else None, + "bit_depth": bit_depth, + "fps": fps, + "crop_limits": list(crop_limits) if crop_limits is not None else None, + "exposure_limits": list(exposure_limits) if exposure_limits is not None else None, + }) + + self._sensor_modes_cache = result + return result + + finally: + if temp_picam2 is not None: + try: + temp_picam2.close() + except Exception: + pass + + def wait_first_frame(self, timeout_s=2.0): + t0 = time.perf_counter() + + while time.perf_counter() - t0 < timeout_s: + with self.frame_lock: + if self.last_raw_frame is not None: + return True + time.sleep(0.01) + + return False diff --git a/Python/raspi/pi/frame_service.py b/Python/raspi/pi/frame_service.py new file mode 100644 index 000000000..7708d952d --- /dev/null +++ b/Python/raspi/pi/frame_service.py @@ -0,0 +1,239 @@ +import time +import threading +import base64 + + +class FrameService: + def __init__(self, state, trigger_manager, camera_manager): + self.state = state + self.trigger = trigger_manager + self.camera = camera_manager + self._capture_lock = threading.RLock() + + self._ensure_raw_processor() + + def _ensure_raw_processor(self): + if ( + not hasattr(self, "raw_processor") or + self.raw_processor.sensor_width != self.state.width or + self.raw_processor.sensor_height != self.state.height or + self.raw_processor.bayer_pattern.upper() != self.state.source_bayer_pattern.upper() + ): + from raw_processor_core import RawProcessorCore + self.raw_processor = RawProcessorCore( + sensor_width=self.state.width, + sensor_height=self.state.height, + bayer_pattern=self.state.source_bayer_pattern, + ) + + def _capture_with_retry(self, max_attempts=10, retry_delay_s=0.02): + last_info = None + + for _ in range(max_attempts): + info = self.camera.capture_raw_frame() + + if ( + info is not None and + info.get("frame") is not None and + info.get("width", 0) > 0 and + info.get("height", 0) > 0 and + info.get("channels", 0) > 0 + ): + return info + + last_info = info + time.sleep(retry_delay_s) + + if last_info is None: + raise RuntimeError(f"Capture retornou frame vazio após {max_attempts} tentativas") + + raise RuntimeError( + f"Capture retornou frame vazio após {max_attempts} tentativas: " + f"width={last_info.get('width')}, " + f"height={last_info.get('height')}, " + f"channels={last_info.get('channels')}, " + f"camera_frame_id={last_info.get('frame_id')}, " + f"camera_frame_ts={last_info.get('timestamp')}" + ) + + def capture_frame_raw(self): + with self._capture_lock: + if not self.state.initialized: + raise RuntimeError("Módulo não inicializado") + + t0_perf = time.perf_counter() + t0_unix = time.time() + + dt_trigger = 0.0 + dt_settle = 0.0 + dt_capture = 0.0 + dt_process = 0.0 + + try: + if self.state.trigger_enabled: + trig_start = time.perf_counter() + self.trigger.pulse() + trig_end = time.perf_counter() + dt_trigger = trig_end - trig_start + + settle_ms = float(getattr(self.state, "trigger_settle_delay_ms", 0.0) or 0.0) + if settle_ms > 0: + settle_start = time.perf_counter() + time.sleep(settle_ms / 1000.0) + settle_end = time.perf_counter() + dt_settle = settle_end - settle_start + + cap_start = time.perf_counter() + raw_frame_info = self._capture_with_retry() + cap_end = time.perf_counter() + dt_capture = cap_end - cap_start + + proc_start = time.perf_counter() + if self.state.frame_type != "RAW_BRUTO": + self._ensure_raw_processor() + frame_out, meta_extra = self._build_output_frame(raw_frame_info) + proc_end = time.perf_counter() + dt_process = proc_end - proc_start + + except Exception as e: + raise RuntimeError(f"Falha durante captura/processamento de frame: {e}") from e + + if frame_out is None: + raise RuntimeError("Frame processado retornou nulo") + + total_end = time.perf_counter() + + self.state.stream_frame_id += 1 + frame_id = self.state.stream_frame_id + + meta = { + "frame_id": frame_id, + "camera_frame_id": int(raw_frame_info["frame_id"]), + "camera_frame_ts": raw_frame_info["timestamp"], + + "ts_pi": t0_unix, + "ts_pi_monotonic": t0_perf, + + "dt_trigger": dt_trigger, + "dt_settle": dt_settle, + "dt_capture": dt_capture, + "dt_process": dt_process, + "dt_total_pi": total_end - t0_perf, + + "source_bayer_pattern": self.state.source_bayer_pattern, + "source_bit_depth": self.state.source_bit_depth, + } + + meta.update(meta_extra) + + return frame_out, meta + + def capture_frame_base64(self): + frame, meta = self.capture_frame_raw() + + frame_bytes = frame.tobytes() + meta["encoding"] = "base64" + meta["size"] = len(frame_bytes) + meta["data"] = base64.b64encode(frame_bytes).decode("ascii") + return meta + + def _build_output_frame(self, raw_frame_info): + if self.state.frame_type == "RAW_BRUTO": + return self._process_raw_bruto(raw_frame_info) + + if self.state.frame_type == "RGB": + return self._process_rgb(raw_frame_info) + + if self.state.frame_type == "RGBNIR": + return self._process_rgbnir(raw_frame_info) + + raise RuntimeError(f"frame_type inválido: {self.state.frame_type}") + + def _convert_output_dtype(self, arr): + max_sensor_value = (1 << int(self.state.source_bit_depth)) - 1 + if self.state.output_dtype == "uint8": + if arr.dtype == "uint8": + return arr + if arr.dtype == "float32": + return (arr * 255.0).clip(0, 255).astype("uint8") + if arr.dtype.kind in ("u", "i"): + # assume 10-bit/16-bit vindo do pipeline + max_val = arr.max() if arr.size > 0 else 0 + if max_val <= 255: + return arr.astype("uint8") + return ((arr.astype("float32") / max_sensor_value) * 255.0).clip(0, 255).astype("uint8") + raise RuntimeError(f"dtype não suportado para uint8: {arr.dtype}") + + elif self.state.output_dtype == "float32": + if arr.dtype == "float32": + return arr + if arr.dtype == "uint8": + return arr.astype("float32") / 255.0 + if arr.dtype.kind in ("u", "i"): + max_val = arr.max() if arr.size > 0 else 0 + if max_val <= 255: + return arr.astype("float32") / 255.0 + return arr.astype("float32") / max_sensor_value + raise RuntimeError(f"dtype não suportado para float32: {arr.dtype}") + + raise RuntimeError(f"output_dtype inválido: {self.state.output_dtype}") + + def _process_raw_bruto(self, raw_frame_info): + frame = raw_frame_info["frame"] + + meta_extra = { + "frame_type": self.state.frame_type, + "payload_format_version": self.state.payload_format_version, + "output_dtype": str(frame.dtype), + "dtype": str(frame.dtype), + "output_layout": "HW", + "output_channels": 1, + "output_channel_names": ["RAW10_PACKED"], + "output_width": int(raw_frame_info["width"]), + "output_height": int(raw_frame_info["height"]), + "source_width": int(self.state.width), + "source_height": int(self.state.height), + "packed_width": int(raw_frame_info["width"]), + "packed_height": int(raw_frame_info["height"]), + "width": int(raw_frame_info["width"]), + "height": int(raw_frame_info["height"]), + "channels": 1, + } + + return frame, meta_extra + + def _process_rgb(self, raw_frame_info): + packed = raw_frame_info["frame"] + b_depth = self.state.source_bit_depth + + if packed.ndim == 3 and packed.shape[2] == 1: + packed = packed[:, :, 0] + + raw16 = self.raw_processor.unpack_raw10_packed(packed) + rgb_chw = self.raw_processor.build_training_rgb(raw16, bit_depth=b_depth) + + rgb_chw = self._convert_output_dtype(rgb_chw) + + meta_extra = { + "frame_type": self.state.frame_type, + "payload_format_version": self.state.payload_format_version, + "output_dtype": self.state.output_dtype, + "dtype": str(rgb_chw.dtype), + "output_layout": "CHW", + "output_channels": 3, + "output_channel_names": ["R", "G", "B"], + "output_width": int(rgb_chw.shape[2]), + "output_height": int(rgb_chw.shape[1]), + "source_width": int(self.state.width), + "source_height": int(self.state.height), + "width": int(rgb_chw.shape[2]), + "height": int(rgb_chw.shape[1]), + "channels": 3, + "packed_width": int(raw_frame_info["width"]), + "packed_height": int(raw_frame_info["height"]), + } + + return rgb_chw, meta_extra + + def _process_rgbnir(self, raw_frame_info): + raise NotImplementedError("RGBNIR ainda não implementado") diff --git a/Python/raspi/pi/main.py b/Python/raspi/pi/main.py new file mode 100644 index 000000000..48c85c723 --- /dev/null +++ b/Python/raspi/pi/main.py @@ -0,0 +1,5 @@ +from server import ModuleServer + +if __name__ == "__main__": + server = ModuleServer(host="0.0.0.0", port=5000) + server.start() diff --git a/Python/raspi/pi/protocol.py b/Python/raspi/pi/protocol.py new file mode 100644 index 000000000..7f8f9af8f --- /dev/null +++ b/Python/raspi/pi/protocol.py @@ -0,0 +1,35 @@ +import json + + +MAX_MESSAGE_SIZE = 1024 * 1024 # 1 MB + + +def decode_message(raw_line: str) -> dict: + if raw_line is None: + raise ValueError("Mensagem ausente") + + if len(raw_line) > MAX_MESSAGE_SIZE: + raise ValueError("Mensagem excede tamanho máximo permitido") + + raw_line = raw_line.strip() + if not raw_line: + raise ValueError("Mensagem vazia") + + try: + data = json.loads(raw_line) + except json.JSONDecodeError as e: + raise ValueError(f"JSON inválido: {e.msg}") from e + + if not isinstance(data, dict): + raise ValueError("Mensagem JSON deve ser um objeto") + + return data + + +def encode_message(data: dict) -> bytes: + if not isinstance(data, dict): + raise ValueError("Resposta deve ser um objeto dict") + + return ( + json.dumps(data, ensure_ascii=False, separators=(",", ":")) + "\n" + ).encode("utf-8") \ No newline at end of file diff --git a/Python/raspi/pi/raw_processor_core.py b/Python/raspi/pi/raw_processor_core.py new file mode 100644 index 000000000..32f264e06 --- /dev/null +++ b/Python/raspi/pi/raw_processor_core.py @@ -0,0 +1,120 @@ +import numpy as np +import math +from typing import Optional + + +class RawProcessorCore: + def __init__(self, sensor_width: int, sensor_height: int, bayer_pattern: str = "GBRG"): + self.sensor_width = sensor_width + self.sensor_height = sensor_height + self.bayer_pattern = bayer_pattern.upper() + + def unpack_raw10_packed( + self, + packed_frame: np.ndarray, + sensor_width: Optional[int] = None, + sensor_height: Optional[int] = None + ): + if packed_frame.ndim == 3 and packed_frame.shape[2] == 1: + packed_frame = packed_frame[:, :, 0] + + width = sensor_width if sensor_width is not None else self.sensor_width + height = sensor_height if sensor_height is not None else self.sensor_height + + expected_packed_width = math.ceil(width * 10 / 8) + + actual_h, actual_w = packed_frame.shape[:2] + padding = actual_w - expected_packed_width + + if actual_h != height: + raise ValueError( + f"[ERRO FRAME] Altura packed inesperada: {packed_frame.shape}, " + f"esperado altura={height}" + ) + + if actual_w < expected_packed_width: + raise ValueError( + f"[ERRO FRAME] Largura packed menor que a útil esperada: {packed_frame.shape}, " + f"esperado pelo menos ({height}, {expected_packed_width})" + ) + + if padding > 64: + raise ValueError( + f"[ERRO FRAME] Padding excessivo no packed: {packed_frame.shape}, " + f"esperado útil ({height}, {expected_packed_width}), padding={padding}" + ) + + packed_frame = packed_frame[:, :expected_packed_width] + groups = packed_frame.reshape(height, width // 4, 5).astype(np.uint16) + + b0 = groups[:, :, 0] + b1 = groups[:, :, 1] + b2 = groups[:, :, 2] + b3 = groups[:, :, 3] + b4 = groups[:, :, 4] + + p0 = (b0 << 2) | ((b4 >> 0) & 0x03) + p1 = (b1 << 2) | ((b4 >> 2) & 0x03) + p2 = (b2 << 2) | ((b4 >> 4) & 0x03) + p3 = (b3 << 2) | ((b4 >> 6) & 0x03) + + raw16 = np.empty((height, width), dtype=np.uint16) + raw16[:, 0::4] = p0 + raw16[:, 1::4] = p1 + raw16[:, 2::4] = p2 + raw16[:, 3::4] = p3 + + return raw16 + + def extract_bayer_channels(self, raw16: np.ndarray) -> dict: + p = self.bayer_pattern + + if p == "GBRG": + g1 = raw16[0::2, 0::2] + b = raw16[0::2, 1::2] + r = raw16[1::2, 0::2] + g2 = raw16[1::2, 1::2] + elif p == "GRBG": + g1 = raw16[0::2, 0::2] + r = raw16[0::2, 1::2] + b = raw16[1::2, 0::2] + g2 = raw16[1::2, 1::2] + elif p == "RGGB": + r = raw16[0::2, 0::2] + g1 = raw16[0::2, 1::2] + g2 = raw16[1::2, 0::2] + b = raw16[1::2, 1::2] + elif p == "BGGR": + b = raw16[0::2, 0::2] + g1 = raw16[0::2, 1::2] + g2 = raw16[1::2, 0::2] + r = raw16[1::2, 1::2] + else: + raise ValueError(f"Padrão Bayer não suportado: {p}") + + return {"R": r, "G1": g1, "G2": g2, "B": b} + + def build_training_rgb( + self, + raw16: np.ndarray, + output_dtype: str = "float32", + bit_depth: int = 10, + ) -> np.ndarray: + ch = self.extract_bayer_channels(raw16) + + max_val = float((1 << bit_depth) - 1) + + r = ch["R"].astype(np.float32) / max_val + g = ((ch["G1"].astype(np.float32) + ch["G2"].astype(np.float32)) * 0.5) / max_val + b = ch["B"].astype(np.float32) / max_val + + chw = np.stack([r, g, b], axis=0).astype(np.float32) + chw = np.clip(chw, 0.0, 1.0) + + if output_dtype == "float32": + return chw + + if output_dtype == "uint8": + return (chw * 255.0).clip(0, 255).astype(np.uint8) + + raise ValueError(f"output_dtype não suportado: {output_dtype}") \ No newline at end of file diff --git a/Python/raspi/pi/raw_processor_preview.py b/Python/raspi/pi/raw_processor_preview.py new file mode 100644 index 000000000..95ede77ad --- /dev/null +++ b/Python/raspi/pi/raw_processor_preview.py @@ -0,0 +1,125 @@ +import cv2 +import numpy as np +from typing import Optional + + +class RawProcessorPreview: + def __init__(self, sensor_width: int, sensor_height: int, bayer_pattern: str = "GBRG"): + self.sensor_width = sensor_width + self.sensor_height = sensor_height + self.bayer_pattern = bayer_pattern.upper() + + def raw16_to_vis8( + self, raw16: np.ndarray, + black_level: Optional[int] = None, + white_level: Optional[int] = None, + gamma: float = 2.2, + bit_depth: int = 10 + ) -> np.ndarray: + """ + Conversão para visualização: + - auto-level + - gamma + """ + max_val = float((1 << bit_depth) - 1) + + raw = raw16.astype(np.float32) + + if black_level is None: + black_level = float(raw.min()) + if white_level is None: + white_level = float(raw.max()) + + if white_level <= black_level: + norm = raw / max_val + else: + norm = (raw - black_level) / (white_level - black_level) + + norm = np.clip(norm, 0.0, 1.0) + + if gamma is not None and gamma > 0: + norm = np.power(norm, 1.0 / gamma) + + return (norm * 255.0).clip(0, 255).astype(np.uint8) + + def _debayer_code(self): + mapping = { + # Troque de BayerGB para BayerGR para inverter R e B + "GBRG": cv2.COLOR_BayerGR2BGR, + "GRBG": cv2.COLOR_BayerGB2BGR, + "RGGB": cv2.COLOR_BayerBG2BGR, + "BGGR": cv2.COLOR_BayerRG2BGR, + } + if self.bayer_pattern not in mapping: + raise ValueError(f"Padrão Bayer não suportado: {self.bayer_pattern}") + return mapping[self.bayer_pattern] + + def apply_preview_white_balance(self, bgr: np.ndarray, strength: float = 1.0) -> np.ndarray: + """ + Gray-world simples para deixar o preview mais agradável. + Não usar no raw de treino. + """ + img = bgr.astype(np.float32) + + mean_b = float(img[:, :, 0].mean()) + mean_g = float(img[:, :, 1].mean()) + mean_r = float(img[:, :, 2].mean()) + + mean_gray = (mean_b + mean_g + mean_r) / 3.0 + + eps = 1e-6 + gain_b = mean_gray / max(mean_b, eps) + gain_g = mean_gray / max(mean_g, eps) + gain_r = mean_gray / max(mean_r, eps) + + # strength=1 aplica total, strength=0 não aplica + gain_b = 1.0 + (gain_b - 1.0) * strength + gain_g = 1.0 + (gain_g - 1.0) * strength + gain_r = 1.0 + (gain_r - 1.0) * strength + + img[:, :, 0] *= gain_b + img[:, :, 1] *= gain_g + img[:, :, 2] *= gain_r + + return np.clip(img, 0, 255).astype(np.uint8) + + def apply_preview_contrast(self, bgr: np.ndarray, alpha: float = 1.08, beta: float = 0.0) -> np.ndarray: + """ + Ajuste leve de contraste/brilho para preview. + """ + out = cv2.convertScaleAbs(bgr, alpha=alpha, beta=beta) + return out + + def raw16_to_preview_bgr( + self, + raw16: np.ndarray, + gamma: float = 2.2, + wb_strength: float = 0.8, + apply_wb: bool = True, + apply_contrast: bool = True, + bit_depth: int = 10, + ) -> np.ndarray: + """ + Pipeline de preview bonito: + 1. auto-level + gamma no mosaico + 2. demosaic + 3. white balance simples + 4. leve contraste final + """ + vis8 = self.raw16_to_vis8(raw16, gamma=gamma, bit_depth=bit_depth) + bgr = cv2.cvtColor(vis8, self._debayer_code()) + + if apply_wb: + bgr = self.apply_preview_white_balance(bgr, strength=wb_strength) + + if apply_contrast: + bgr = self.apply_preview_contrast(bgr, alpha=1.08, beta=0.0) + + return bgr + + def raw16_to_preview_jpg_bytes(self, raw16: np.ndarray, jpeg_quality: int = 95) -> bytes: + bgr = self.raw16_to_preview_bgr(raw16) + ok, enc = cv2.imencode(".jpg", bgr, [int(cv2.IMWRITE_JPEG_QUALITY), int(jpeg_quality)]) + if not ok: + raise RuntimeError("Falha ao codificar preview JPG") + return enc.tobytes() diff --git a/Python/raspi/pi/server.py b/Python/raspi/pi/server.py new file mode 100644 index 000000000..1df81dcaf --- /dev/null +++ b/Python/raspi/pi/server.py @@ -0,0 +1,488 @@ +import socket +import traceback +import threading + +from protocol import decode_message, encode_message +from state import ModuleState + + +def parse_bool(value): + if isinstance(value, bool): + return value + + if isinstance(value, (int, float)): + return bool(value) + + if isinstance(value, str): + s = value.strip().lower() + if s in ("1", "true", "yes", "on"): + return True + if s in ("0", "false", "no", "off"): + return False + + raise ValueError("Valor booleano inválido") + +def parse_frame_type(value): + valid = {"RAW_BRUTO", "RGB", "RGBNIR"} + if value not in valid: + raise ValueError(f"frame_type inválido: {value}") + return value + +def parse_output_dtype(value): + valid = {"uint8", "float32"} + if value not in valid: + raise ValueError(f"output_dtype inválido: {value}") + return value + +def parse_bayer_pattern(value): + valid = {"GBRG", "GRBG", "RGGB", "BGGR"} + if value not in valid: + raise ValueError(f"bayer_pattern inválido: {value}") + return value + +class ModuleServer: + def __init__(self, host="0.0.0.0", port=5000): + self.host = host + self.port = port + self.state = ModuleState() + self.running = False + self.lock = threading.RLock() + + from camera_manager import CameraManager + from trigger_manager import TriggerManager + from frame_service import FrameService + from stream_sender import StreamSender + + self.camera = CameraManager(self.state) + self.trigger = TriggerManager( + pin=self.state.trigger_pin, + active_high=self.state.trigger_active_high, + pulse_ms=self.state.trigger_pulse_ms + ) + self.frame_service = FrameService(self.state, self.trigger, self.camera) + self.stream_sender = StreamSender(self.state, self.frame_service) + + def _mark_reconfigure_needed(self): + if hasattr(self.camera, "mark_reconfigure_needed"): + self.camera.mark_reconfigure_needed() + + def handle_command(self, msg: dict) -> dict: + with self.lock: + cmd = msg.get("cmd") + self.state.last_command = cmd + + try: + if cmd == "ping": + return {"ok": True, "reply": "pong"} + + if cmd == "get_status": + return {"ok": True, **self.state.to_dict()} + + if cmd == "get_config": + return { + "ok": True, + "fps": self.state.fps, + "jpeg_quality": self.state.jpeg_quality, + "width": self.state.width, + "height": self.state.height, + "camera_id": self.state.camera_id, + + "payload_format_version": self.state.payload_format_version, + + "source_bayer_pattern": self.state.source_bayer_pattern, + "source_bit_depth": self.state.source_bit_depth, + + "frame_type": self.state.frame_type, + "output_dtype": self.state.output_dtype, + "output_layout": self.state.output_layout, + "output_channels": self.state.output_channels, + "output_channel_names": self.state.output_channel_names, + "output_width": self.state.output_width, + "output_height": self.state.output_height, + } + + if cmd == "begin": + frame_type = parse_frame_type(msg.get("frame_type", self.state.frame_type)) + output_dtype = parse_output_dtype(msg.get("output_dtype", self.state.output_dtype)) + + # Atualiza a especificação de saída ANTES de inicializar + self.state.frame_type = frame_type + self.state.output_dtype = output_dtype + self.state.update_output_spec() + + if self.state.initialized: + return { + "ok": True, + "status": self.state.status, + "initialized": True, + + "payload_format_version": self.state.payload_format_version, + + "frame_type": self.state.frame_type, + "output_dtype": self.state.output_dtype, + "output_layout": self.state.output_layout, + "output_channels": self.state.output_channels, + "output_channel_names": self.state.output_channel_names, + "output_width": self.state.output_width, + "output_height": self.state.output_height, + + "source_bayer_pattern": self.state.source_bayer_pattern, + "source_bit_depth": self.state.source_bit_depth, + } + + self.state.status = "initializing" + self.state.status_detail = None + + self.trigger.begin() + self.camera.begin() + + self.state.initialized = True + self.state.camera_connected = True + self.state.status = "ready" + self.state.last_error = None + + return { + "ok": True, + "status": self.state.status, + "initialized": True, + + "payload_format_version": self.state.payload_format_version, + + "frame_type": self.state.frame_type, + "output_dtype": self.state.output_dtype, + "output_layout": self.state.output_layout, + "output_channels": self.state.output_channels, + "output_channel_names": self.state.output_channel_names, + "output_width": self.state.output_width, + "output_height": self.state.output_height, + + "source_bayer_pattern": self.state.source_bayer_pattern, + "source_bit_depth": self.state.source_bit_depth, + } + + if cmd == "stop": + self.state.status = "stopping" + + if getattr(self.stream_sender, "is_running", False): + self.stream_sender.stop() + + self.camera.stop() + self.trigger.stop() + + self.state.initialized = False + self.state.streaming = False + self.state.camera_connected = False + self.state.status = "idle" + self.state.status_detail = None + return {"ok": True, "status": self.state.status} + + if cmd == "capture_frame": + data = self.frame_service.capture_frame_base64() + return {"ok": True, **data} + + if cmd == "set_fps": + value = int(msg.get("value")) + if value <= 0 or value > 120: + return {"ok": False, "error": "fps inválido"} + + self.state.fps = value + self._mark_reconfigure_needed() + + return {"ok": True, "fps": self.state.fps} + + if cmd == "set_jpeg_quality": + value = int(msg.get("value")) + if value < 1 or value > 100: + return {"ok": False, "error": "jpeg_quality inválido"} + + self.state.jpeg_quality = value + return {"ok": True, "jpeg_quality": self.state.jpeg_quality} + + if cmd == "set_bayer": + pattern = parse_bayer_pattern(msg.get("pattern", self.state.source_bayer_pattern)) + self.state.source_bayer_pattern = pattern + + return { + "ok": True, + "bayer_pattern": self.state.source_bayer_pattern + } + + if cmd == "set_resolution": + width = int(msg.get("width")) + height = int(msg.get("height")) + + if width <= 0 or height <= 0: + return {"ok": False, "error": "resolução inválida"} + + self.state.width = width + self.state.height = height + self.state.update_output_spec() + self._mark_reconfigure_needed() + + return { + "ok": True, + "width": self.state.width, + "height": self.state.height, + "output_width": self.state.output_width, + "output_height": self.state.output_height, + "frame_type": self.state.frame_type, + "output_layout": self.state.output_layout, + "output_channels": self.state.output_channels, + "output_channel_names": self.state.output_channel_names, + } + + if cmd == "start_stream": + host = msg.get("host") + port = int(msg.get("port")) + fps = float(msg.get("fps", self.state.fps)) + + if not host: + return {"ok": False, "error": "host obrigatório"} + + if port <= 0 or port > 65535: + return {"ok": False, "error": "porta inválida"} + + if fps <= 0: + return {"ok": False, "error": "fps inválido"} + + if getattr(self.stream_sender, "is_running", False): + return { + "ok": True, + "streaming": True, + "host": self.state.stream_host, + "port": self.state.stream_port, + "fps": self.state.stream_fps, + + "payload_format_version": self.state.payload_format_version, + + "frame_type": self.state.frame_type, + "output_dtype": self.state.output_dtype, + "output_layout": self.state.output_layout, + "output_channels": self.state.output_channels, + "output_channel_names": self.state.output_channel_names, + "output_width": self.state.output_width, + "output_height": self.state.output_height, + + "source_bayer_pattern": self.state.source_bayer_pattern, + "source_bit_depth": self.state.source_bit_depth, + } + + self.stream_sender.start(host, port, fps) + self.state.streaming = True + self.state.stream_host = host + self.state.stream_port = port + self.state.stream_fps = fps + self.state.status = "streaming" + + return { + "ok": True, + "streaming": True, + "host": host, + "port": port, + "fps": fps, + + "payload_format_version": self.state.payload_format_version, + + "frame_type": self.state.frame_type, + "output_dtype": self.state.output_dtype, + "output_layout": self.state.output_layout, + "output_channels": self.state.output_channels, + "output_channel_names": self.state.output_channel_names, + "output_width": self.state.output_width, + "output_height": self.state.output_height, + + "source_bayer_pattern": self.state.source_bayer_pattern, + "source_bit_depth": self.state.source_bit_depth, + } + + if cmd == "stop_stream": + if getattr(self.stream_sender, "is_running", False): + self.stream_sender.stop() + + self.state.streaming = False + self.state.stream_host = None + self.state.stream_port = None + self.state.stream_fps = None + + if self.state.initialized: + self.state.status = "ready" + else: + self.state.status = "idle" + + return { + "ok": True, + "streaming": False, + "status": self.state.status + } + + if cmd == "set_ae_enable": + self.state.ae_enable = parse_bool(msg.get("value")) + + if self.camera.initialized: + self.camera.apply_controls() + + return {"ok": True, "ae_enable": self.state.ae_enable} + + if cmd == "set_awb_enable": + self.state.awb_enable = parse_bool(msg.get("value")) + + if self.camera.initialized: + self.camera.apply_controls() + + return {"ok": True, "awb_enable": self.state.awb_enable} + + if cmd == "set_exposure_time": + value = msg.get("value") + if value is not None: + value = int(value) + if value <= 0: + return {"ok": False, "error": "ExposureTime inválido"} + + self.state.exposure_time_us = value + + if self.camera.initialized: + self.camera.apply_controls() + + return {"ok": True, "exposure_time_us": self.state.exposure_time_us} + + if cmd == "set_analogue_gain": + value = msg.get("value") + if value is not None: + value = float(value) + if value <= 0: + return {"ok": False, "error": "AnalogueGain inválido"} + + self.state.analogue_gain = value + + if self.camera.initialized: + self.camera.apply_controls() + + return {"ok": True, "analogue_gain": self.state.analogue_gain} + + if cmd == "set_colour_gains": + r_gain = msg.get("r_gain") + b_gain = msg.get("b_gain") + + if r_gain is None or b_gain is None: + return {"ok": False, "error": "r_gain e b_gain são obrigatórios"} + + r_gain = float(r_gain) + b_gain = float(b_gain) + + if r_gain <= 0 or b_gain <= 0: + return {"ok": False, "error": "ColourGains inválidos"} + + self.state.colour_gains = [r_gain, b_gain] + + if self.camera.initialized: + self.camera.apply_controls() + + return {"ok": True, "colour_gains": self.state.colour_gains} + + if cmd == "clear_exposure_time": + self.state.exposure_time_us = None + if self.camera.initialized: + self.camera.apply_controls() + return {"ok": True, "exposure_time_us": None} + + if cmd == "clear_analogue_gain": + self.state.analogue_gain = None + if self.camera.initialized: + self.camera.apply_controls() + return {"ok": True, "analogue_gain": None} + + if cmd == "clear_colour_gains": + self.state.colour_gains = None + if self.camera.initialized: + self.camera.apply_controls() + return {"ok": True, "colour_gains": None} + + if cmd == "get_camera_controls": + return { + "ok": True, + "ae_enable": self.state.ae_enable, + "awb_enable": self.state.awb_enable, + "exposure_time_us": self.state.exposure_time_us, + "analogue_gain": self.state.analogue_gain, + "colour_gains": self.state.colour_gains, + "fps": self.state.fps, + } + + if cmd == "get_sensor_modes": + modes = self.camera.get_sensor_modes() + return { + "ok": True, + "sensor_modes": modes + } + + return {"ok": False, "error": f"Comando desconhecido: {cmd}"} + + except Exception as e: + self.state.set_error(str(e)) + return {"ok": False, "error": str(e)} + + def client_thread(self, conn, addr): + print(f"[INFO] Cliente conectado: {addr}") + + buffer = b"" + + try: + conn.settimeout(1.0) + + while self.running: + try: + chunk = conn.recv(4096) + except socket.timeout: + continue + + if not chunk: + break + + buffer += chunk + + while b"\n" in buffer: + line, buffer = buffer.split(b"\n", 1) + if not line.strip(): + continue + + try: + msg = decode_message(line.decode("utf-8")) + response = self.handle_command(msg) + except Exception as e: + response = {"ok": False, "error": str(e)} + + conn.sendall(encode_message(response)) + + except Exception as e: + print(f"[ERRO] Cliente {addr}: {e}") + traceback.print_exc() + finally: + try: + conn.close() + except Exception: + pass + print(f"[INFO] Cliente desconectado: {addr}") + + def start(self): + self.running = True + + with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as server_socket: + server_socket.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) + server_socket.bind((self.host, self.port)) + server_socket.listen(5) + server_socket.settimeout(1.0) + + print(f"[INFO] Servidor ouvindo em {self.host}:{self.port}") + + while self.running: + try: + conn, addr = server_socket.accept() + except socket.timeout: + continue + + t = threading.Thread( + target=self.client_thread, + args=(conn, addr), + daemon=True + ) + t.start() diff --git a/Python/raspi/pi/state.py b/Python/raspi/pi/state.py new file mode 100644 index 000000000..c08a42589 --- /dev/null +++ b/Python/raspi/pi/state.py @@ -0,0 +1,155 @@ +class ModuleState: + def __init__(self): + self.module_name = "multispectral" + self.version = "0.1.0" + self.payload_format_version = 1 + + self.status = "idle" + self.status_detail = None + self.last_error = None + self.last_command = None + + self.camera_connected = False + self.camera_id = "cam0" + + self.initialized = False + self.streaming = False + + # Configuração da captura no sensor + self.fps = 10 + self.jpeg_quality = 90 + self.width = 640 + self.height = 480 + self.source_bayer_pattern = "GBRG" + self.source_bit_depth = 10 + + # Configuração do tipo de payload de saída + self.frame_type = "RAW_BRUTO" # RAW_BRUTO | RGB | RGBNIR + self.output_dtype = "uint8" # uint8 | float32 + self.output_layout = "CHW" # sempre CHW + self.output_channels = 1 + self.output_channel_names = ["RAW10_PACKED"] + self.output_width = 640 + self.output_height = 480 + + self.trigger_enabled = True + self.trigger_pin = 18 + self.trigger_active_high = True + self.trigger_pulse_ms = 5.0 + self.trigger_settle_delay_ms = 0.0 + + self.stream_host = None + self.stream_port = None + self.stream_fps = None + self.stream_frame_id = 0 + + self.codec_family = "none" + self.codec_name = "blosc" + self.codec_params = { + "cname": "zstd", + "clevel": 3, + "shuffle": "BITSHUFFLE", + } + + self.ae_enable = True + self.awb_enable = True + self.exposure_time_us = 15000 + self.analogue_gain = 1.0 + self.colour_gains = [1.0, 1.0] + self.frame_duration_limits = None + + self.update_output_spec() + + def update_output_spec(self): + if self.frame_type == "RAW_BRUTO": + self.output_layout = "HW" + self.output_channels = 1 + self.output_channel_names = ["BAYER"] + self.output_width = self.width + self.output_height = self.height + + elif self.frame_type == "RGB": + self.output_layout = "CHW" + self.output_channels = 3 + self.output_channel_names = ["R", "G", "B"] + self.output_width = self.width // 2 + self.output_height = self.height // 2 + + elif self.frame_type == "RGBNIR": + self.output_layout = "CHW" + self.output_channels = 5 + self.output_channel_names = ["R", "G", "B", "NIR", "RE"] + self.output_width = self.width // 2 + self.output_height = self.height // 2 + + else: + raise ValueError(f"frame_type inválido: {self.frame_type}") + + def clear_error(self): + self.last_error = None + if self.status == "error": + self.status = "idle" + self.status_detail = None + + def set_error(self, error: str): + self.last_error = str(error) + self.status = "error" + self.status_detail = str(error) + + def to_dict(self): + return { + "module": self.module_name, + "version": self.version, + + "status": self.status, + "status_detail": self.status_detail, + "last_error": self.last_error, + "last_command": self.last_command, + + "camera_connected": self.camera_connected, + "camera_id": self.camera_id, + + "initialized": self.initialized, + "streaming": self.streaming, + + "fps": self.fps, + "jpeg_quality": self.jpeg_quality, + "width": self.width, + "height": self.height, + + "trigger_enabled": self.trigger_enabled, + "trigger_pin": self.trigger_pin, + "trigger_active_high": self.trigger_active_high, + "trigger_pulse_ms": self.trigger_pulse_ms, + "trigger_settle_delay_ms": self.trigger_settle_delay_ms, + + "stream_host": self.stream_host, + "stream_port": self.stream_port, + "stream_fps": self.stream_fps, + "stream_frame_id": self.stream_frame_id, + + "codec_family": self.codec_family, + "codec_name": self.codec_name, + "codec_params": self.codec_params, + + "ae_enable": self.ae_enable, + "awb_enable": self.awb_enable, + "exposure_time_us": self.exposure_time_us, + "analogue_gain": self.analogue_gain, + "colour_gains": self.colour_gains, + "frame_duration_limits": self.frame_duration_limits, + + "payload_format_version": self.payload_format_version, + + "source_bayer_pattern": self.source_bayer_pattern, + "source_bit_depth": self.source_bit_depth, + + "frame_type": self.frame_type, + "output_dtype": self.output_dtype, + "output_layout": self.output_layout, + "output_channels": self.output_channels, + "output_channel_names": self.output_channel_names, + "output_width": self.output_width, + "output_height": self.output_height, + } + \ No newline at end of file diff --git a/Python/raspi/pi/stream_sender.py b/Python/raspi/pi/stream_sender.py new file mode 100644 index 000000000..0c7d876bd --- /dev/null +++ b/Python/raspi/pi/stream_sender.py @@ -0,0 +1,327 @@ +import json +import socket +import threading +import time +import queue +from numcodecs import Blosc + + +class StreamSender: + _QUEUE_SENTINEL = object() + + def __init__(self, state, frame_service): + self.state = state + self.frame_service = frame_service + + self._lock = threading.RLock() + self._stop_event = threading.Event() + + self._sock = None + self._thread_capture = None + self._thread_send = None + + self._queue = queue.Queue(maxsize=2) + + self._codec = None + self._codec_signature = None + + self._t_json_pack = 0.0 + self._t_send_header = 0.0 + self._t_send_payload = 0.0 + + self._frames_dropped = 0 + self._capture_errors = 0 + self._send_errors = 0 + self._last_sent_camera_frame_id = 0 + + @property + def is_running(self): + return not self._stop_event.is_set() and ( + self._thread_capture is not None or self._thread_send is not None + ) + + def start(self, host: str, port: int, fps: float): + with self._lock: + if self.is_running: + raise RuntimeError("Stream já está em execução") + + if not self.state.initialized: + raise RuntimeError("Módulo não inicializado") + + self._stop_event.clear() + self._clear_queue() + + self._t_json_pack = 0.0 + self._t_send_header = 0.0 + self._t_send_payload = 0.0 + self._frames_dropped = 0 + self._capture_errors = 0 + self._send_errors = 0 + + self.state.streaming = True + self.state.stream_host = host + self.state.stream_port = port + self.state.stream_fps = fps + self._last_sent_camera_frame_id = 0 + + self._thread_capture = threading.Thread( + target=self._worker_capture, + args=(fps,), + daemon=True + ) + self._thread_send = threading.Thread( + target=self._worker_send, + args=(host, port), + daemon=True + ) + + self._thread_capture.start() + self._thread_send.start() + + def stop(self): + with self._lock: + self._stop_event.set() + + try: + self._queue.put_nowait(self._QUEUE_SENTINEL) + except queue.Full: + try: + self._queue.get_nowait() + except queue.Empty: + pass + try: + self._queue.put_nowait(self._QUEUE_SENTINEL) + except queue.Full: + pass + + sock = self._sock + self._sock = None + + if sock is not None: + try: + sock.shutdown(socket.SHUT_RDWR) + except Exception: + pass + try: + sock.close() + except Exception: + pass + + if self._thread_capture is not None: + self._thread_capture.join(timeout=2.0) + self._thread_capture = None + + if self._thread_send is not None: + self._thread_send.join(timeout=2.0) + self._thread_send = None + + self._clear_queue() + + self.state.streaming = False + self.state.stream_host = None + self.state.stream_port = None + self.state.stream_fps = None + + def _clear_queue(self): + while not self._queue.empty(): + try: + self._queue.get_nowait() + except queue.Empty: + break + + def _normalize_shuffle(self, shuffle_value): + if isinstance(shuffle_value, int): + return shuffle_value + + mapping = { + "NOSHUFFLE": Blosc.NOSHUFFLE, + "SHUFFLE": Blosc.SHUFFLE, + "BITSHUFFLE": Blosc.BITSHUFFLE, + } + + key = str(shuffle_value).upper() + if key not in mapping: + raise ValueError(f"shuffle inválido: {shuffle_value}") + + return mapping[key] + + def _build_codec_from_state(self): + family = self.state.codec_family + name = self.state.codec_name + params = dict(self.state.codec_params) + + if family == "none": + return None + + if family != "numcodecs": + raise ValueError(f"Família de codec não suportada: {family}") + + if name == "blosc": + params["shuffle"] = self._normalize_shuffle(params.get("shuffle", "SHUFFLE")) + return Blosc(**params) + + raise ValueError(f"Codec numcodecs não suportado: {name}") + + def _get_codec_signature_from_state(self): + return ( + self.state.codec_family, + self.state.codec_name, + tuple(sorted(self.state.codec_params.items())) + ) + + def _ensure_codec(self): + sig = self._get_codec_signature_from_state() + if self._codec is None or self._codec_signature != sig: + self._codec = self._build_codec_from_state() + self._codec_signature = sig + + def _compress(self, frame_bytes: bytes) -> bytes: + self._ensure_codec() + + if self.state.codec_family == "none": + return frame_bytes + + return self._codec.encode(frame_bytes) + + def _send_packet(self, sock: socket.socket, header: dict, payload: bytes): + t0 = time.perf_counter() + header_bytes = json.dumps(header, separators=(",", ":")).encode("utf-8") + t1 = time.perf_counter() + + sock.sendall(len(header_bytes).to_bytes(4, "big")) + sock.sendall(header_bytes) + sock.sendall(len(payload).to_bytes(4, "big")) + t2 = time.perf_counter() + + sock.sendall(payload) + t3 = time.perf_counter() + + self._t_json_pack = t1 - t0 + self._t_send_header = t2 - t1 + self._t_send_payload = t3 - t2 + + def _worker_capture(self, fps: float): + frame_interval = 1.0 / fps if fps > 0 else 0.0 + next_deadline = time.perf_counter() + last_frame_ts = None + + while not self._stop_event.is_set(): + t_loop0 = time.perf_counter() + + try: + frame, meta = self.frame_service.capture_frame_raw() + except Exception as e: + self._capture_errors += 1 + print(f"[WARN] Capture falhou: {e}") + time.sleep(0.05) + continue + + if frame is None: + time.sleep(0.01) + continue + + camera_frame_id = meta.get("camera_frame_id", 0) + if camera_frame_id == self._last_sent_camera_frame_id: + time.sleep(0.001) + continue + + t_frame_ready = time.perf_counter() + dt_frame_period = 0.0 if last_frame_ts is None else (t_frame_ready - last_frame_ts) + last_frame_ts = t_frame_ready + + t_bytes0 = time.perf_counter() + frame_bytes = frame.tobytes() + t_bytes1 = time.perf_counter() + + t_comp0 = time.perf_counter() + comp_bytes = self._compress(frame_bytes) + t_comp1 = time.perf_counter() + + codec_name = None if self.state.codec_family == "none" else self.state.codec_name + codec_params = {} if self.state.codec_family == "none" else dict(self.state.codec_params) + + header = { + "frame_id": meta.get("frame_id"), + "codec_family": self.state.codec_family, + "codec_name": codec_name, + "codec_params": codec_params, + "dt_frame_period": dt_frame_period, + "payload_size_raw": len(frame_bytes), + "payload_size_comp": len(comp_bytes), + "dt_bytes": t_bytes1 - t_bytes0, + "dt_comp": t_comp1 - t_comp0, + **meta + } + + queued = False + try: + self._queue.put_nowait((header, comp_bytes)) + queued = True + except queue.Full: + self._frames_dropped += 1 + try: + self._queue.get_nowait() + except queue.Empty: + pass + try: + self._queue.put_nowait((header, comp_bytes)) + queued = True + except queue.Full: + self._frames_dropped += 1 + + if queued: + self._last_sent_camera_frame_id = camera_frame_id + + if frame_interval > 0: + next_deadline += frame_interval + now = time.perf_counter() + + if next_deadline < now - frame_interval: + next_deadline = now + + sleep_time = next_deadline - now + if sleep_time > 0: + time.sleep(sleep_time) + + def _worker_send(self, host: str, port: int): + try: + with socket.create_connection((host, port), timeout=5) as sock: + sock.settimeout(10) + with self._lock: + self._sock = sock + + while not self._stop_event.is_set(): + try: + item = self._queue.get(timeout=0.2) + except queue.Empty: + continue + + if item is self._QUEUE_SENTINEL: + break + + header, payload = item + + header["dt_pack_prev"] = self._t_json_pack + header["dt_send_header_prev"] = self._t_send_header + header["dt_send_payload_prev"] = self._t_send_payload + header["stream_drops"] = self._frames_dropped + header["capture_errors"] = self._capture_errors + header["send_errors"] = self._send_errors + + self._send_packet(sock, header, payload) + self.state.stream_frame_id = header["frame_id"] + + except Exception as e: + self._send_errors += 1 + if not self._stop_event.is_set(): + print(f"[WARN] StreamSender encerrado com erro: {e}") + finally: + with self._lock: + self._sock = None + + self.state.streaming = False + self.state.stream_host = None + self.state.stream_port = None + self.state.stream_fps = None + + self._stop_event.set() \ No newline at end of file diff --git a/Python/raspi/pi/trigger_manager.py b/Python/raspi/pi/trigger_manager.py new file mode 100644 index 000000000..c98c052b6 --- /dev/null +++ b/Python/raspi/pi/trigger_manager.py @@ -0,0 +1,74 @@ +import time +import threading +import RPi.GPIO as GPIO + + +class TriggerManager: + def __init__(self, pin: int, active_high: bool = True, pulse_ms: float = 5.0): + if pin is None or int(pin) < 0: + raise ValueError("Pin inválido") + + pulse_ms = float(pulse_ms) + if pulse_ms <= 0: + raise ValueError("pulse_ms deve ser maior que zero") + + self.pin = int(pin) + self.active_high = bool(active_high) + self.pulse_ms = pulse_ms + self.initialized = False + + self._lock = threading.RLock() + self._active_level = GPIO.HIGH if self.active_high else GPIO.LOW + self._idle_level = GPIO.LOW if self.active_high else GPIO.HIGH + + def begin(self): + with self._lock: + if self.initialized: + return True + + try: + GPIO.setmode(GPIO.BCM) + GPIO.setwarnings(False) + GPIO.setup(self.pin, GPIO.OUT) + GPIO.output(self.pin, self._idle_level) + except Exception as e: + raise RuntimeError(f"Falha ao inicializar trigger GPIO {self.pin}: {e}") from e + + self.initialized = True + return True + + def stop(self): + with self._lock: + if not self.initialized: + return True + + try: + GPIO.output(self.pin, self._idle_level) + except Exception: + pass + + try: + GPIO.cleanup(self.pin) + except Exception: + pass + + self.initialized = False + return True + + def pulse(self): + with self._lock: + if not self.initialized: + raise RuntimeError("TriggerManager não inicializado") + + try: + GPIO.output(self.pin, self._active_level) + time.sleep(self.pulse_ms / 1000.0) + except Exception as e: + raise RuntimeError(f"Falha ao gerar pulso no GPIO {self.pin}: {e}") from e + finally: + try: + GPIO.output(self.pin, self._idle_level) + except Exception: + pass + + return True \ No newline at end of file diff --git a/Python/raspi/stream_receiver.py b/Python/raspi/stream_receiver.py index 88725a30c..560e672b5 100644 --- a/Python/raspi/stream_receiver.py +++ b/Python/raspi/stream_receiver.py @@ -2,8 +2,6 @@ import json import socket import threading import time -import lz4.frame -import zstandard as zstd from numcodecs import Blosc import numpy as np @@ -22,8 +20,8 @@ class StreamReceiver: self.last_meta = None self.last_receive_ts = None - self._zstd_d = zstd.ZstdDecompressor() - self._codec = Blosc(cname="lz4", clevel=1, shuffle=Blosc.SHUFFLE) + self._codec = None + self._codec_signature = None @property def is_running(self): @@ -96,16 +94,12 @@ class StreamReceiver: payload_len = int.from_bytes(self._recv_exact(client, 4), "big") payload_comp = self._recv_exact(client, payload_len) - codec = header["codec"] - - if codec == "lz4": - payload = lz4.frame.decompress(payload_comp) - elif codec == "zstd": - payload = self._zstd_d.decompress(payload_comp) - elif codec == "numcodecs": - payload = self._codec.decode(payload_comp) + + self._ensure_codec(header) + if header.get("codec_family") == "none": + payload = payload_comp else: - raise ValueError(f"Codec não suportado: {codec}") + payload = self._codec.decode(payload_comp) expected = header["payload_size_raw"] if len(payload) != expected: @@ -113,15 +107,26 @@ class StreamReceiver: f"Tamanho descomprimido inválido: {len(payload)} != {expected}" ) - height = header["height"] - width = header["width"] - channels = header["channels"] - dtype = np.uint8 if header["dtype"] == "uint8" else None + height = int(header.get("output_height", header.get("height"))) + width = int(header.get("output_width", header.get("width"))) + channels = int(header.get("output_channels", header.get("channels", 1))) + layout = header.get("output_layout", "HWC") - if dtype is None: - raise RuntimeError(f"dtype não suportado: {header['dtype']}") + dtype = self._numpy_dtype_from_header(header) - frame = np.frombuffer(payload, dtype=dtype).reshape(height, width, channels) + arr = np.frombuffer(payload, dtype=dtype) + + if layout == "HW": + frame = arr.reshape(height, width) + + elif layout == "CHW": + frame = arr.reshape(channels, height, width) + + elif layout == "HWC": + frame = arr.reshape(height, width, channels) + + else: + raise RuntimeError(f"Layout não suportado: {layout}") self.last_frame = frame self.last_meta = header @@ -137,3 +142,70 @@ class StreamReceiver: self._running = False self._client_sock = None self._server_sock = None + + + def _normalize_shuffle(self, shuffle_value): + if isinstance(shuffle_value, int): + return shuffle_value + + mapping = { + "NOSHUFFLE": Blosc.NOSHUFFLE, + "SHUFFLE": Blosc.SHUFFLE, + "BITSHUFFLE": Blosc.BITSHUFFLE, + } + + key = str(shuffle_value).upper() + if key not in mapping: + raise ValueError(f"shuffle inválido: {shuffle_value}") + + return mapping[key] + + def _build_codec_from_header(self, header: dict): + family = header.get("codec_family") + name = header.get("codec_name") + params = dict(header.get("codec_params", {})) + + if family != "numcodecs": + raise ValueError(f"Família de codec não suportada: {family}") + + if name == "blosc": + params["shuffle"] = self._normalize_shuffle(params.get("shuffle", "SHUFFLE")) + return Blosc(**params) + + raise ValueError(f"Codec numcodecs não suportado: {name}") + + def _get_codec_signature_from_header(self, header: dict): + return ( + header.get("codec_family"), + header.get("codec_name"), + tuple(sorted(dict(header.get("codec_params", {})).items())) + ) + + def _ensure_codec(self, header: dict): + family = header.get("codec_family") + + if family == "none": + self._codec = None + self._codec_signature = ("none", None, ()) + return + + sig = self._get_codec_signature_from_header(header) + if self._codec is None or self._codec_signature != sig: + self._codec = self._build_codec_from_header(header) + self._codec_signature = sig + + + def _numpy_dtype_from_header(self, header: dict): + dtype_str = header.get("dtype") or header.get("output_dtype") or "uint8" + + mapping = { + "uint8": np.uint8, + "float32": np.float32, + "uint16": np.uint16, + } + + if dtype_str not in mapping: + raise RuntimeError(f"dtype não suportado: {dtype_str}") + + return mapping[dtype_str] + diff --git a/Python/raspi/test_capture_raw.py b/Python/raspi/test_capture_raw.py deleted file mode 100644 index 67a3f5120..000000000 --- a/Python/raspi/test_capture_raw.py +++ /dev/null @@ -1,19 +0,0 @@ -import time -from multispectral_service import MultiSpectralService - -svc = MultiSpectralService(host="192.168.105.6", port=5000) -svc.connect() - -print("BEGIN:", svc.begin()) - -tempos_pi = [] -tempos_pc = [] - -for i in range(10): - frame, meta = svc.capture_frame_array() - tempos_pi.append(meta["dt_total_pi"]) - tempos_pc.append(meta["dt_total_pc"]) - print(i, meta["dt_total_pi"], meta["dt_total_pc"]) - -print("STOP:", svc.stop()) -svc.disconnect() \ No newline at end of file diff --git a/Python/raspi/test_stream_raw.py b/Python/raspi/test_stream_raw.py deleted file mode 100644 index d5fb6e62f..000000000 --- a/Python/raspi/test_stream_raw.py +++ /dev/null @@ -1,55 +0,0 @@ -import time -from multispectral_service import MultiSpectralService -from stream_receiver import StreamReceiver - - -STREAM_PORT = 6001 -PI_HOST = "192.168.105.6" -PC_HOST = "192.168.105.5" - - -def main(): - receiver = StreamReceiver(host="0.0.0.0", port=STREAM_PORT) - receiver.start() - - time.sleep(0.5) - - svc = MultiSpectralService(host=PI_HOST, port=5000, timeout=10) - svc.connect() - - try: - _fps = 15 - print("SET RES:", svc.set_resolution(640, 480)) - print("SET FPS:", svc.set_fps(_fps)) - print("BEGIN:", svc.begin()) - - print("START STREAM:", svc.start_stream(PC_HOST, STREAM_PORT, fps=_fps)) - - t0 = time.perf_counter() - last_frame_id = -1 - - while time.perf_counter() - t0 < 5: - meta = receiver.last_meta - frame = receiver.last_frame - - if meta is not None and meta["frame_id"] != last_frame_id: - last_frame_id = meta["frame_id"] - print(meta) - #print( - # f"frame_id={meta['frame_id']} " - # f"shape={frame.shape if frame is not None else None} " - # f"dt_total_pi={meta.get('dt_total_pi'):.4f}" - #) - - time.sleep(0.02) - - print("STOP STREAM:", svc.stop_stream()) - print("STOP:", svc.stop()) - - finally: - svc.disconnect() - receiver.stop() - - -if __name__ == "__main__": - main() \ No newline at end of file