Compare commits

...

2 Commits

Author SHA1 Message Date
Diego Freitas ceba226614 otimizado weed worker e trocado senha do ntrip ibge 2026-09-14 08:27:36 -03:00
Diego Freitas 68ff42e0b7 raw_processor_core radiometric_single 2026-09-14 07:25:32 -03:00
15 changed files with 6719 additions and 3471 deletions

View File

@ -16,7 +16,7 @@ from camera_worker.oak_fcc3_core.oak_fcc3_client import OakFcc3Client
from camera_worker.camera_imu import IMUCamera
CAMERA_MULTISPECTRAL_VERSION = "production_v1_2026_08_24"
CAMERA_MULTISPECTRAL_VERSION = "production_v1_2026_09_12_raw_save_light"
class CameraMultispectral:
@ -116,6 +116,7 @@ class CameraMultispectral:
self.ultimo_tensor_multispec = None
self.ultimo_frame_rgb = None
self.ultimo_raw_multi = None
self.ultimo_raw_meta = None
self.ultimo_meta = None
self.ultimo_decoded = None
@ -584,6 +585,7 @@ class CameraMultispectral:
self.ultimo_tensor_multispec = None
self.ultimo_frame_rgb = None
self.ultimo_raw_multi = None
self.ultimo_raw_meta = None
self.ultimo_meta = None
self.ultimo_decoded = None
@ -889,6 +891,9 @@ class CameraMultispectral:
# Arrays do pacote são imutáveis após a captura.
# Mantemos referências e copiamos somente no salvamento.
self.ultimo_raw_multi = dict(frame)
# Meta científico pertence EXATAMENTE ao mesmo pacote RAW.
# Mantido no mesmo lock para impedir RAW N + meta N+1.
self.ultimo_raw_meta = dict(meta)
self.timestamp_ultimo_raw_multi = ts_raw_perf
self._ultimo_resultado_raw = {
"erro": None,
@ -1008,7 +1013,7 @@ class CameraMultispectral:
}
def requisitar_frame_raw_multi(self, force: bool = False, max_age_s: float = None):
"""Retorna cópia do último RAW capturado pelo fluxo operacional."""
"""Retorna snapshot por referência do último RAW operacional imutável."""
try:
if max_age_s is None or float(max_age_s) <= 0:
max_age_s = max(1.0, self._cache_max_age_s)
@ -1023,10 +1028,10 @@ class CameraMultispectral:
if idade > float(max_age_s):
raise RuntimeError(f"cache RAW antigo: {idade:.2f}s")
raw_frame = {
cam_id: arr.copy()
for cam_id, arr in self.ultimo_raw_multi.items()
}
# Snapshot barato: arrays do pacote operacional são imutáveis
# após a captura. Copiamos somente o dicionário de referências.
# A thread de salvamento mantém os ndarrays vivos enquanto escreve.
raw_frame = dict(self.ultimo_raw_multi)
resultado = dict(self._ultimo_resultado_raw)
resultado["cache_age_s"] = float(idade)
resultado["force_ignorado"] = bool(force)
@ -1785,7 +1790,12 @@ class CameraMultispectral:
def _ts_name(self) -> str:
return datetime.now().strftime("%Y%m%d_%H%M%S_%f")[:-3]
def requisitar_bundle_raw_multispec(self, force: bool = True, max_age_s: float = None):
def requisitar_bundle_raw_multispec(
self,
force: bool = True,
max_age_s: float = None,
include_preview: bool = False,
):
"""
Monta o bundle a partir do cache do fluxo operacional.
@ -1796,16 +1806,30 @@ class CameraMultispectral:
if max_age_s is None or float(max_age_s) <= 0:
max_age_s = max(1.0, self._cache_max_age_s)
raw_frame, resultado = self.requisitar_frame_raw_multi(
force=False,
max_age_s=max_age_s,
)
if raw_frame is None:
raise RuntimeError(resultado.get("erro") or "RAW cache indisponível")
agora = time.perf_counter()
with self._lock:
raw_meta = dict(self.ultimo_meta or {})
if self.ultimo_raw_multi is None or self.timestamp_ultimo_raw_multi is None:
raise RuntimeError("cache RAW ainda não disponível")
idade = agora - self.timestamp_ultimo_raw_multi
if idade > float(max_age_s):
raise RuntimeError(f"cache RAW antigo: {idade:.2f}s")
# Snapshot científico ATÔMICO. Os ndarrays continuam por
# referência, sem copiar megabytes dentro do lock.
raw_frame = dict(self.ultimo_raw_multi)
raw_meta = dict(self.ultimo_raw_meta or {})
resultado = dict(self._ultimo_resultado_raw)
resultado["cache_age_s"] = float(idade)
resultado["force_ignorado"] = bool(force)
frame_id_raw = int(resultado.get("frame_id", 0) or 0)
frame_id_meta = int(raw_meta.get("frame_id", 0) or 0)
if frame_id_raw and frame_id_meta and frame_id_raw != frame_id_meta:
raise RuntimeError(
"snapshot RAW/meta inconsistente: "
f"raw_frame_id={frame_id_raw} meta_frame_id={frame_id_meta}"
)
# Reforça o contrato científico no próprio stream_meta salvo.
raw_meta.setdefault("frame_type", "RAW_BRUTO")
@ -1830,6 +1854,9 @@ class CameraMultispectral:
f"Presentes: {sorted(presentes)}"
)
preview_bgr = None
preview_method = "disabled_runtime_save"
if include_preview:
preview_bgr, preview_method = self._build_preview_raw_multispec(
raw_frame=raw_frame,
raw_meta=raw_meta,
@ -1920,8 +1947,7 @@ class CameraMultispectral:
"""
Salva pacote RAW_BRUTO multiespectral no mesmo espírito do capture de dataset.
Saída:
<nome>.png
Saída operacional:
<nome>.json
<nome>_CAM_A.bin
<nome>_CAM_B.bin
@ -1939,6 +1965,7 @@ class CameraMultispectral:
bundle, resultado = self.requisitar_bundle_raw_multispec(
force=True,
max_age_s=0.0,
include_preview=False,
)
if not resultado.get("frame_valido", False):
@ -1977,17 +2004,10 @@ class CameraMultispectral:
if not payload_files:
raise RuntimeError("Nenhum payload RAW foi salvo.")
# Em operação, PosProcessamento salva APENAS os RAW .bin + JSON.
# Preview bonito é reconstruído offline a partir do próprio bundle.
caminho_preview = None
if preview_bgr is not None and hasattr(preview_bgr, "size") and preview_bgr.size > 0:
caminho_preview = os.path.join(
pasta,
f"{nome_base}.png"
)
cv2.imwrite(caminho_preview, preview_bgr)
caminhos.append(caminho_preview)
meta_save = {
"ts": datetime.now().isoformat(timespec="milliseconds"),
"source": "operacao_robo",

View File

@ -45,6 +45,7 @@ class OakFcc3Client:
imu_modo="rotation_vector",
imu_freq_hz=200,
evaluate_quality=True,
require_product_contract=False,
hardware_sync_enabled=None,
frame_sync_master=None,
@ -65,6 +66,30 @@ class OakFcc3Client:
self.module_calibration_json = module_calibration_json
self.module_params = self._load_module_params(module_calibration_json)
self.fusion_config = self.module_params.get("fusion_config", {}) or {}
self.require_product_contract = bool(require_product_contract)
assembly = self.module_params.get("assembly_metadata", {}) or {}
self.product_contract = bool(
self.module_params.get("schema") == "multispec_module_params_v3"
and assembly.get("schema") == "multispec_module_params_assembly_v1"
)
if self.require_product_contract and not self.product_contract:
raise RuntimeError(
"OakFcc3Client exige module_params de produção homologado: "
f"schema={self.module_params.get('schema')!r} "
f"assembly={assembly.get('schema')!r}"
)
# Em produto, Bayer e raster RGB nativo vêm do MP, não de fallbacks do caller.
mp_bayer = str(self.module_params.get("bayer_pattern", "") or "").upper()
if mp_bayer:
self.bayer = mp_bayer
rgb_size = (self.module_params.get("sensor_size_by_role", {}) or {}).get("rgb")
if isinstance(rgb_size, (list, tuple)) and len(rgb_size) == 2:
self.width = int(rgb_size[0])
self.height = int(rgb_size[1])
self.imu_modo = str(imu_modo).strip().lower()
self.imu_freq_hz = int(imu_freq_hz)
@ -84,14 +109,15 @@ class OakFcc3Client:
self.svc = OakFcc3Service(
timeout=10,
fps=fps,
width=width,
height=height,
width=self.width,
height=self.height,
frame_type=frame_type,
output_dtype=output_dtype,
capture_mode=capture_mode,
raw_policy=raw_policy,
mx_id=self.mx_id,
module_calibration_json=module_calibration_json,
require_product_contract=self.require_product_contract,
imu_modo=self.imu_modo,
imu_freq_hz=self.imu_freq_hz,
@ -108,15 +134,15 @@ class OakFcc3Client:
self.applied_camera_controls = {}
self.radiometric_controller = None
self.core = RawProcessorCore(
sensor_width=width,
sensor_height=height,
bayer_pattern=bayer,
sensor_width=self.width,
sensor_height=self.height,
bayer_pattern=self.bayer,
calibration_json_path=module_calibration_json,
)
self.preview = RawProcessorPreview(
sensor_width=width,
sensor_height=height,
bayer_pattern=bayer,
sensor_width=self.width,
sensor_height=self.height,
bayer_pattern=self.bayer,
)
def __enter__(self):
@ -133,6 +159,46 @@ class OakFcc3Client:
with open(path, "r", encoding="utf-8") as f:
return json.load(f)
def get_contract(self):
"""
Contrato estático/runtime consumido por CameraMultispectral.
O module_params é autoridade de hardware/calibração. O target final
pode ser null no contrato novo; quando existir é apenas um default
compatível/legado. O Weed CameraManager passa o target ONNX ao caller.
"""
try:
service_contract = self.svc.get_contract() or {}
except Exception:
service_contract = {}
mp = self.module_params or {}
fusion = mp.get("fusion_config", {}) or {}
target = fusion.get("target_size")
default_target = None
if isinstance(target, (list, tuple)) and len(target) == 2:
tw, th = int(target[0]), int(target[1])
if tw > 0 and th > 0:
default_target = [tw, th]
sensor_sizes = mp.get("sensor_size_by_role", {}) or {}
camera_hw = mp.get("camera_hardware", {}) or {}
out = dict(service_contract)
out.update({
"product_contract": bool(self.product_contract),
"require_product_contract": bool(self.require_product_contract),
"module_params_schema": mp.get("schema"),
"module_calibration_json": self.module_calibration_json,
"sensor_size_by_role": sensor_sizes,
"camera_hardware": camera_hw,
"bayer_pattern": mp.get("bayer_pattern") or out.get("bayer_pattern"),
"frame_type": self.frame_type,
"raw_policy": self.raw_policy,
"default_target_size": default_target,
})
return out
def apply_module_camera_settings(self):
camera_settings = self.module_params.get("camera_settings", {}) or {}
@ -374,8 +440,8 @@ class OakFcc3Client:
evaluate_quality = self.evaluate_quality
evaluate_quality = bool(evaluate_quality)
tensor = self.core.fuse_multispec_cameras(decoded, meta, channels_expected)
tensor = self.core.resize_tensor_chw(tensor, target_size=target_size)
tensor = self.core.fuse_multispec_cameras(decoded, meta, channels_expected, target_size=target_size)
#tensor = self.core.resize_tensor_chw(tensor, target_size=target_size)
# Mantém paridade com build_infer_tensor_from_stream: se a calibração
# habilitar patch normalization, ela também vale no caminho decoded.

View File

@ -53,14 +53,15 @@ class OakFcc3Manager:
mx_id=None,
module_calibration_json=None,
module_params=None,
require_product_contract=False,
imu_modo="rotation_vector",
imu_freq_hz=200,
):
self.fps = fps
# Para compatibilidade, mantemos width/height.
# No RAW_BRUTO isso não muda o sensor, pois usamos 800p fixo.
# No MULTISPEC isso representa a saída final alinhada da OAK.
# No RAW_BRUTO a resolução física é definida pelo module_params/hardware.
# No MULTISPEC width/height representa a saída final alinhada da OAK.
self.width = int(width)
self.height = int(height)
self.size = (self.width, self.height)
@ -120,6 +121,52 @@ class OakFcc3Manager:
self.fusion_config = (self.module_params or {}).get("fusion_config", {}) or {}
self.aligned_geometry = None
# Contrato produto: MP descreve hardware/calibração. target_size pode ser null.
self.require_product_contract = bool(require_product_contract)
self.module_params_schema = (self.module_params or {}).get("schema")
assembly = (self.module_params or {}).get("assembly_metadata", {}) or {}
self.product_contract = bool(
self.module_params_schema == "multispec_module_params_v3"
and assembly.get("schema") == "multispec_module_params_assembly_v1"
)
self.sensor_size_by_role = copy.deepcopy(
(self.module_params or {}).get("sensor_size_by_role", {}) or {}
)
self.camera_hardware_expected = copy.deepcopy(
(self.module_params or {}).get("camera_hardware", {}) or {}
)
self.bayer_pattern = str(
(self.module_params or {}).get("bayer_pattern", "") or ""
).upper() or None
if self.require_product_contract and not self.product_contract:
raise RuntimeError(
"OakFcc3Manager exige module_params de produção homologado: "
f"schema={self.module_params_schema!r} assembly={assembly.get('schema')!r}"
)
if self.product_contract:
self._validate_static_product_contract()
roles_from_mp = {}
for role in ("rgb", "re", "nir"):
hw = self.camera_hardware_expected.get(role, {}) or {}
socket_name = str(hw.get("socket", "") or "").upper()
if socket_name:
roles_from_mp[socket_name] = role
if roles_from_mp:
self.roles = roles_from_mp
rgb_size = self.sensor_size_by_role.get("rgb")
if isinstance(rgb_size, (list, tuple)) and len(rgb_size) == 2:
self.sensor_width = int(rgb_size[0])
self.sensor_height = int(rgb_size[1])
self.camera_controls = {
cam_id: self._default_controls_for_role(role)
for cam_id, role in self.roles.items()
}
sync_cfg = (self.module_params.get("capture_synchronization", {}) if isinstance(self.module_params, dict) else {})
if hardware_sync_enabled is None: hardware_sync_enabled = sync_cfg.get("hardware_sync_enabled", False)
if frame_sync_master is None: frame_sync_master = sync_cfg.get("frame_sync_master", "CAM_A")
@ -135,6 +182,17 @@ class OakFcc3Manager:
self.async_capture_enabled = True
self.async_capture_mode = "latest" # latest | queue
self.async_capture_max_queue = 2
# Default de produção: materialização antecipada no producer,
# pois apresentou menor latência ponta a ponta.
# Para testar materialização diferida:
# OAK_FCC3_DEFER_RAW_MATERIALIZATION=1
_defer_env = str(
os.environ.get("OAK_FCC3_DEFER_RAW_MATERIALIZATION", "0")
).strip().lower()
self.defer_raw_materialization = _defer_env not in (
"0", "false", "no", "off", "disabled"
)
self._capture_thread = None
self._capture_stop_event = threading.Event()
self._capture_lock = threading.RLock()
@ -183,6 +241,113 @@ class OakFcc3Manager:
with open(path, "r", encoding="utf-8") as f:
return json.load(f)
def _validate_static_product_contract(self):
sizes = self.sensor_size_by_role or {}
hardware = self.camera_hardware_expected or {}
for role in ("rgb", "re", "nir"):
size = sizes.get(role)
hw = hardware.get(role)
if not (isinstance(size, (list, tuple)) and len(size) == 2):
raise RuntimeError(f"module_params sem sensor_size_by_role.{role}")
w, h = int(size[0]), int(size[1])
if w <= 0 or h <= 0:
raise RuntimeError(f"sensor_size_by_role.{role} inválido: {size}")
if not isinstance(hw, dict):
raise RuntimeError(f"module_params sem camera_hardware.{role}")
socket_name = str(hw.get("socket", "") or "").upper()
sensor_name = str(hw.get("sensor", "") or "").upper()
if not socket_name or not sensor_name:
raise RuntimeError(
f"camera_hardware.{role} precisa conter socket e sensor: {hw}"
)
hw_size = hw.get("size")
if hw_size is not None:
got = [int(hw_size[0]), int(hw_size[1])]
if got != [w, h]:
raise RuntimeError(
f"camera_hardware.{role}.size={got} diverge de sensor_size_by_role={size}"
)
if self.bayer_pattern not in ("RGGB", "BGGR", "GRBG", "GBRG"):
raise RuntimeError(
f"bayer_pattern inválido no module_params: {self.bayer_pattern!r}"
)
def _expected_native_size_for_role(self, role):
size = (self.sensor_size_by_role or {}).get(str(role).lower())
if isinstance(size, (list, tuple)) and len(size) == 2:
return [int(size[0]), int(size[1])]
return None
def _set_resolution_by_native_size(self, cam, role, is_color):
size = self._expected_native_size_for_role(role)
if size is None:
if is_color:
names = ("THE_800_P", "THE_1080_P")
enum_cls = dai.ColorCameraProperties.SensorResolution
else:
names = ("THE_800_P", "THE_720_P", "THE_400_P")
enum_cls = dai.MonoCameraProperties.SensorResolution
for name in names:
value = getattr(enum_cls, name, None)
if value is not None:
try:
cam.setResolution(value)
return
except Exception:
continue
return
key = tuple(size)
if is_color:
names = {
(1280, 800): "THE_800_P",
(1920, 1200): "THE_1200_P",
(1920, 1080): "THE_1080_P",
(1280, 720): "THE_720_P",
}
enum_cls = dai.ColorCameraProperties.SensorResolution
else:
names = {
(1280, 800): "THE_800_P",
(1280, 720): "THE_720_P",
(640, 400): "THE_400_P",
}
enum_cls = dai.MonoCameraProperties.SensorResolution
enum_name = names.get(key)
enum_value = getattr(enum_cls, enum_name, None) if enum_name else None
if enum_value is None:
raise RuntimeError(
f"Resolução nativa {size} da role={role} não possui enum DepthAI suportado "
f"nesta versão do runtime (esperado={enum_name})."
)
cam.setResolution(enum_value)
def _validate_connected_hardware(self, features):
if not self.product_contract:
return True
by_socket = {str(f.socket.name).upper(): f for f in features}
for role in ("rgb", "re", "nir"):
hw = self.camera_hardware_expected.get(role, {}) or {}
socket_name = str(hw.get("socket", "") or "").upper()
expected_sensor = str(hw.get("sensor", "") or "").upper()
feature = by_socket.get(socket_name)
if feature is None:
raise RuntimeError(
f"Hardware calibrado ausente: role={role} socket={socket_name}"
)
actual_sensor = str(getattr(feature, "sensorName", "") or "").upper()
if expected_sensor and actual_sensor and actual_sensor != expected_sensor:
raise RuntimeError(
f"Sensor divergente em {socket_name}/{role}: "
f"detectado={actual_sensor} homologado={expected_sensor}"
)
return True
def _default_controls_for_role(self, role: str):
role = str(role).lower()
@ -293,11 +458,11 @@ class OakFcc3Manager:
"""
Fluxo clássico.
RGB/OV9782:
ColorCamera raw para RAW_BRUTO.
RGB (OV9782/AR0234):
ColorCamera raw na resolução nativa homologada pelo module_params.
MONO/OV9282:
MonoCamera raw quando disponível.
MonoCamera raw na resolução nativa homologada pelo module_params.
"""
sensor_name_u = str(sensor_name or "").upper()
role_u = str(role or "").lower()
@ -311,14 +476,7 @@ class OakFcc3Manager:
if is_rgb:
cam = self.pipeline.createColorCamera()
cam.setBoardSocket(socket)
try:
cam.setResolution(dai.ColorCameraProperties.SensorResolution.THE_800_P)
except Exception:
try:
cam.setResolution(dai.ColorCameraProperties.SensorResolution.THE_1080_P)
except Exception:
pass
self._set_resolution_by_native_size(cam, role_u, is_color=True)
cam.setInterleaved(False)
cam.setColorOrder(dai.ColorCameraProperties.ColorOrder.RGB)
@ -335,17 +493,7 @@ class OakFcc3Manager:
mono = self.pipeline.create(dai.node.MonoCamera)
mono.setBoardSocket(socket)
try:
mono.setResolution(dai.MonoCameraProperties.SensorResolution.THE_800_P)
except Exception:
try:
mono.setResolution(dai.MonoCameraProperties.SensorResolution.THE_720_P)
except Exception:
try:
mono.setResolution(dai.MonoCameraProperties.SensorResolution.THE_400_P)
except Exception:
pass
self._set_resolution_by_native_size(mono, role_u, is_color=False)
mono.setFps(float(self.fps))
@ -490,7 +638,7 @@ class OakFcc3Manager:
def _create_color_camera_multispec(self, socket):
cam = self.pipeline.create(dai.node.ColorCamera)
cam.setBoardSocket(socket)
cam.setResolution(dai.ColorCameraProperties.SensorResolution.THE_800_P)
self._set_resolution_by_native_size(cam, "rgb", is_color=True)
cam.setFps(float(self.fps))
cam.setInterleaved(False)
@ -500,10 +648,10 @@ class OakFcc3Manager:
cam.setVideoSize(int(self.sensor_width), int(self.sensor_height))
return cam
def _create_mono_camera_multispec(self, socket):
def _create_mono_camera_multispec(self, socket, role):
cam = self.pipeline.create(dai.node.MonoCamera)
cam.setBoardSocket(socket)
cam.setResolution(dai.MonoCameraProperties.SensorResolution.THE_800_P)
self._set_resolution_by_native_size(cam, role, is_color=False)
cam.setFps(float(self.fps))
return cam
@ -625,7 +773,7 @@ class OakFcc3Manager:
bit_depth = 8
raw_format = "BGR888p"
elif role == "re":
cam = self._create_mono_camera_multispec(socket)
cam = self._create_mono_camera_multispec(socket, "re")
src_output = cam.out
quad = self.aligned_geometry["quad_re"]
out_type = dai.ImgFrame.Type.GRAY8
@ -634,7 +782,7 @@ class OakFcc3Manager:
bit_depth = 8
raw_format = "GRAY8"
elif role == "nir":
cam = self._create_mono_camera_multispec(socket)
cam = self._create_mono_camera_multispec(socket, "nir")
src_output = cam.out
quad = self.aligned_geometry["quad_nir"]
out_type = dai.ImgFrame.Type.GRAY8
@ -896,6 +1044,7 @@ class OakFcc3Manager:
self.pipeline = dai.Pipeline()
features = self.device.getConnectedCameraFeatures()
self._validate_connected_hardware(features)
self.queues.clear()
self.buffers.clear()
@ -943,11 +1092,14 @@ class OakFcc3Manager:
self.buffers[cam_id] = deque(maxlen=self.buffer_size)
self.control_queues[cam_id] = None
native_size = self._expected_native_size_for_role(role)
self.camera_info[cam_id] = {
"id": cam_id,
"socket": socket_name,
"sensor": f.sensorName,
"role": role,
"native_size": native_size,
"bayer_pattern": self.bayer_pattern if str(role).lower() == "rgb" else None,
}
self._validate_capture_mode()
@ -1143,6 +1295,14 @@ class OakFcc3Manager:
return {
"mx_id": self.mx_id,
"backend": "oak_fcc3",
"manager_version": "production_v2_2026_09_12",
"product_contract": bool(self.product_contract),
"require_product_contract": bool(self.require_product_contract),
"module_params_schema": self.module_params_schema,
"module_calibration_json": self.module_calibration_json,
"sensor_size_by_role": copy.deepcopy(self.sensor_size_by_role),
"camera_hardware_expected": copy.deepcopy(self.camera_hardware_expected),
"bayer_pattern": self.bayer_pattern,
"running": bool(self.running),
"fps": self.fps,
"width": self.width,
@ -1201,9 +1361,10 @@ class OakFcc3Manager:
t0 = time.perf_counter()
deadline = t0 + float(timeout)
packet = None
with self._capture_cond:
while time.perf_counter() < deadline:
packet = None
mode = str(getattr(self, "async_capture_mode", "latest")).lower()
@ -1218,14 +1379,75 @@ class OakFcc3Manager:
if packet is not None:
seq = int(packet.get("seq", 0))
self._last_consumed_packet_seq = seq
break
frames = packet["frames"]
meta = dict(packet["meta"])
remaining = deadline - time.perf_counter()
if remaining <= 0:
break
age_ms = (time.perf_counter() - float(packet.get("created_perf_counter", time.perf_counter()))) * 1000.0
self._capture_cond.wait(timeout=min(0.005, remaining))
if packet is not None:
# IMPORTANTE: materialização RAW acontece FORA do lock da captura.
# Assim a thread OAK pode continuar drenando/sincronizando enquanto
# o consumidor transforma somente a tripleta que realmente usará.
seq = int(packet.get("seq", 0))
age_ms = (
time.perf_counter()
- float(packet.get("created_perf_counter", time.perf_counter()))
) * 1000.0
get_wait_ms = (time.perf_counter() - t0) * 1000.0
if bool(packet.get("deferred_raw_packet", False)):
cp = dict(packet.get("capture_perf", {}) or {})
t_mat0 = time.perf_counter()
selected_items = packet.get("selected_items", {}) or {}
frames = {}
frame_controls = {}
for cam_id, item in selected_items.items():
frame, controls = self._materialize_selected_item(
cam_id=cam_id,
item=item,
perf=cp,
)
if frame is None:
raise RuntimeError(
f"Frame RAW assíncrono selecionado inválido: {cam_id}"
)
frames[cam_id] = frame
frame_controls[cam_id] = controls or {}
consumer_materialize_ms = (
time.perf_counter() - t_mat0
) * 1000.0
t_meta0 = time.perf_counter()
meta = self._build_meta(
frames,
packet.get("timestamps", {}) or {},
float(packet.get("sync_dt_ms", 0.0) or 0.0),
bool(packet.get("sync_ok", False)),
frame_controls=frame_controls,
frame_id_override=packet.get("frame_id"),
)
consumer_meta_ms = (
time.perf_counter() - t_meta0
) * 1000.0
cp["async_consumer_materialize_ms"] = float(
consumer_materialize_ms
)
cp["async_consumer_meta_ms"] = float(consumer_meta_ms)
cp["async_deferred_packet"] = True
else:
frames = packet["frames"]
meta = dict(packet["meta"])
cp = dict(meta.get("capture_perf", {}) or {})
cp["async_deferred_packet"] = False
cp["async_consumer"] = True
cp["async_packet_seq"] = seq
cp["async_packet_age_ms"] = float(age_ms)
@ -1235,12 +1457,6 @@ class OakFcc3Manager:
return frames, meta
remaining = deadline - time.perf_counter()
if remaining <= 0:
break
self._capture_cond.wait(timeout=min(0.005, remaining))
raise TimeoutError(
f"Timeout aguardando pacote assíncrono do OAK-FFC-3. "
f"status={self.get_async_capture_status()}"
@ -1456,21 +1672,65 @@ class OakFcc3Manager:
}
else:
if bool(getattr(self, "defer_raw_materialization", True)):
# Não toca no payload RAW aqui. O ImgFrame permanece vivo
# no deque e só será materializado se entrar na tripleta
# escolhida pelo sincronizador. Frames descartados custam
# apenas metadados/timestamp.
frame = None
frame_controls = None
deferred_raw = True
if perf is not None:
perf["deferred_raw_msgs"] += 1
else:
frame, frame_controls = self._materialize_raw_msg(
cam_id=cam_id,
msg=msg,
perf=perf,
metric_prefix="drain",
)
deferred_raw = False
if self._is_preview_mode() or self._is_multispec_mode():
t0_ctrl = self._cap_now_ms()
frame_controls = self._extract_frame_controls(msg)
if perf is not None:
perf["drain_controls_ms"] += self._cap_now_ms() - t0_ctrl
deferred_raw = False
self.buffers[cam_id].append({
"frame": frame,
"msg": msg if deferred_raw else None,
"deferred_raw": bool(deferred_raw),
"timestamp": ts_start,
"timestamp_start": ts_start,
"timestamp_end": ts_end,
"controls": frame_controls,
})
if perf is not None:
perf["drained_total"] += 1
perf["drained_by_cam"][cam_id] += 1
def _materialize_raw_msg(self, cam_id, msg, perf=None, metric_prefix="selected"):
"""Materializa um RAW10 packed somente quando ele realmente será usado."""
if msg is None:
raise RuntimeError(f"ImgFrame ausente para materialização RAW: {cam_id}")
t_all = self._cap_now_ms()
t0_data = self._cap_now_ms()
data = msg.getData()
if perf is not None:
perf["drain_get_data_ms"] += self._cap_now_ms() - t0_data
data_ms = self._cap_now_ms() - t0_data
t0_copy = self._cap_now_ms()
raw = np.frombuffer(data, dtype=np.uint8).copy()
if perf is not None:
perf["drain_frombuffer_copy_ms"] += self._cap_now_ms() - t0_copy
copy_ms = self._cap_now_ms() - t0_copy
t0_shape = self._cap_now_ms()
h = int(msg.getHeight())
w = int(msg.getWidth())
stride = self._get_imgframe_stride(msg, raw.size, h, w)
expected = h * stride
if raw.size < expected:
@ -1480,8 +1740,7 @@ class OakFcc3Manager:
)
frame = raw[:expected].reshape((h, stride))
if perf is not None:
perf["drain_reshape_ms"] += self._cap_now_ms() - t0_shape
shape_ms = self._cap_now_ms() - t0_shape
self._last_raw_dims[cam_id] = {
"sensor_width": w,
@ -1492,20 +1751,40 @@ class OakFcc3Manager:
t0_ctrl = self._cap_now_ms()
frame_controls = self._extract_frame_controls(msg)
if perf is not None:
perf["drain_controls_ms"] += self._cap_now_ms() - t0_ctrl
self.buffers[cam_id].append({
"frame": frame,
"timestamp": ts_start,
"timestamp_start": ts_start,
"timestamp_end": ts_end,
"controls": frame_controls,
})
controls_ms = self._cap_now_ms() - t0_ctrl
if perf is not None:
perf["drained_total"] += 1
perf["drained_by_cam"][cam_id] += 1
if metric_prefix == "drain":
perf["drain_get_data_ms"] += data_ms
perf["drain_frombuffer_copy_ms"] += copy_ms
perf["drain_reshape_ms"] += shape_ms
perf["drain_controls_ms"] += controls_ms
else:
perf["selected_materialize_ms"] += self._cap_now_ms() - t_all
perf["selected_get_data_ms"] += data_ms
perf["selected_frombuffer_copy_ms"] += copy_ms
perf["selected_reshape_ms"] += shape_ms
perf["selected_controls_ms"] += controls_ms
perf["selected_materialized_frames"] += 1
perf["selected_materialized_bytes"] += int(expected)
return frame, frame_controls
def _materialize_selected_item(self, cam_id, item, perf=None):
if not bool(item.get("deferred_raw", False)):
return item.get("frame"), item.get("controls", {}) or {}
frame, controls = self._materialize_raw_msg(
cam_id=cam_id,
msg=item.get("msg"),
perf=perf,
metric_prefix="selected",
)
item["frame"] = frame
item["controls"] = controls
item["deferred_raw"] = False
item["msg"] = None
return frame, controls
def _get_imgframe_stride(self, msg, raw_size: int, h: int, w: int) -> int:
try:
@ -1521,7 +1800,7 @@ class OakFcc3Manager:
return int(np.ceil(w * 5.0 / 4.0))
def _try_get_synced_packet(self, perf=None):
def _try_get_synced_packet(self, perf=None, materialize=True):
required_cam_ids = self._get_required_cam_ids()
if not required_cam_ids:
@ -1549,11 +1828,6 @@ class OakFcc3Manager:
for cam_id, item in selected.items()
}
frame_controls = {
cam_id: item.get("controls", {})
for cam_id, item in selected.items()
}
ts_values = list(timestamps.values())
sync_dt_ms = (
@ -1570,11 +1844,6 @@ class OakFcc3Manager:
for cam_id, ts in timestamps.items()
}
perf["selected_seq_by_cam"] = {
cam_id: item.get("controls", {}).get("sequence_num")
for cam_id, item in selected.items()
}
perf["sync_dt_ms"] = float(sync_dt_ms)
perf["sync_ok"] = bool(sync_ok)
@ -1601,10 +1870,41 @@ class OakFcc3Manager:
# Em best/best_effort, entrega a melhor combinação disponível,
# mesmo quando estiver fora da tolerância.
frames = {
cam_id: item["frame"]
#
# No caminho assíncrono/latest, materialize=False mantém apenas as
# referências ImgFrame escolhidas. Se esse pacote for substituído por
# outro antes do consumidor pegá-lo, nenhum payload RAW foi copiado.
# A materialização fica para get_next_frame(), uma única vez, no pacote
# que realmente alimentará o tensor.
if materialize:
payload = {}
frame_controls = {}
for cam_id, item in selected.items():
frame, controls = self._materialize_selected_item(
cam_id=cam_id,
item=item,
perf=perf,
)
if frame is None:
raise RuntimeError(f"Frame selecionado inválido: {cam_id}")
payload[cam_id] = frame
frame_controls[cam_id] = controls or {}
if perf is not None:
perf["selected_seq_by_cam"] = {
cam_id: frame_controls.get(cam_id, {}).get("sequence_num")
for cam_id in selected.keys()
}
else:
# Shallow copy do descritor. O ImgFrame permanece vivo por sua
# referência Python, sem copiar os megabytes do RAW.
payload = {
cam_id: dict(item)
for cam_id, item in selected.items()
}
frame_controls = {}
if perf is not None:
perf["selected_deferred_for_consumer"] = len(payload)
# Consome todos os frames anteriores e o próprio frame selecionado.
for cam_id, used_item in selected.items():
@ -1622,7 +1922,7 @@ class OakFcc3Manager:
)
return (
frames,
payload,
timestamps,
sync_dt_ms,
sync_ok,
@ -1669,7 +1969,15 @@ class OakFcc3Manager:
# Meta
# ============================================================
def _build_meta(self, frames, timestamps, sync_dt_ms, sync_ok, frame_controls):
def _build_meta(
self,
frames,
timestamps,
sync_dt_ms,
sync_ok,
frame_controls,
frame_id_override=None,
):
payload_sources = list(frames.keys())
shapes = {
@ -1730,7 +2038,11 @@ class OakFcc3Manager:
camera_info[cam_id] = item
meta = {
"frame_id": self.frame_id,
"frame_id": (
int(self.frame_id)
if frame_id_override is None
else int(frame_id_override)
),
"backend": "oak_fcc3",
"frame_type": self.frame_type,
"capture_mode": self.capture_mode,
@ -1974,6 +2286,17 @@ class OakFcc3Manager:
"drain_frombuffer_copy_ms": 0.0,
"drain_reshape_ms": 0.0,
"drain_controls_ms": 0.0,
"deferred_raw_msgs": 0,
"selected_materialize_ms": 0.0,
"selected_get_data_ms": 0.0,
"selected_frombuffer_copy_ms": 0.0,
"selected_reshape_ms": 0.0,
"selected_controls_ms": 0.0,
"selected_materialized_frames": 0,
"selected_materialized_bytes": 0,
"selected_deferred_for_consumer": 0,
"async_consumer_materialize_ms": 0.0,
"async_consumer_meta_ms": 0.0,
"sync_select_ms": 0.0,
"meta_ms": 0.0,
"sleep_ms": 0.0,
@ -2093,32 +2416,64 @@ class OakFcc3Manager:
# Importante: esta thread é a única que mexe nas queues/buffers.
self._drain_queues_to_buffers(perf=perf)
synced = self._try_get_synced_packet(perf=perf)
# Em RAW_BRUTO + async, o produtor publica apenas descritores
# ImgFrame da tripleta sincronizada. O consumidor materializa
# somente o pacote latest que realmente pegar. Isso elimina
# getData()/copy() de pacotes completos que seriam sobrescritos.
defer_packet_to_consumer = bool(
self._is_raw_mode()
and getattr(self, "defer_raw_materialization", True)
)
synced = self._try_get_synced_packet(
perf=perf,
materialize=not defer_packet_to_consumer,
)
if synced is None:
# Dorme curto. Pode testar 0.0005 se quiser reduzir latência.
time.sleep(0.001)
continue
frames, timestamps, sync_dt_ms, sync_ok, frame_controls = synced
payload, timestamps, sync_dt_ms, sync_ok, frame_controls = synced
self.frame_id += 1
packet_frame_id = int(self.frame_id)
meta = None
if not defer_packet_to_consumer:
meta = self._build_meta(
frames,
payload,
timestamps,
sync_dt_ms,
sync_ok,
frame_controls=frame_controls,
frame_id_override=packet_frame_id,
)
now = time.perf_counter()
packet = {
"seq": int(self._latest_packet_seq + 1),
"created_perf_counter": float(now),
"frames": frames,
"meta": meta,
"frame_id": packet_frame_id,
"deferred_raw_packet": bool(defer_packet_to_consumer),
}
if defer_packet_to_consumer:
packet.update({
"selected_items": payload,
"timestamps": timestamps,
"sync_dt_ms": float(sync_dt_ms),
"sync_ok": bool(sync_ok),
"capture_perf": perf or {},
})
else:
packet.update({
"frames": payload,
"meta": meta,
})
# Adiciona perf de captura assíncrona no meta.
if perf is not None:
perf["async_thread"] = True
@ -2126,7 +2481,10 @@ class OakFcc3Manager:
perf["sync_dt_ms"] = float(sync_dt_ms)
perf["sync_ok"] = bool(sync_ok)
perf["wait_reason"] = "async_packet_ready"
if meta is not None:
meta["capture_perf"] = perf
else:
packet["capture_perf"] = perf
with self._capture_cond:
self._latest_packet_seq += 1
@ -2195,6 +2553,7 @@ class OakFcc3Manager:
st.update({
"enabled": bool(getattr(self, "async_capture_enabled", True)),
"mode": str(getattr(self, "async_capture_mode", "latest")),
"defer_raw_materialization": bool(getattr(self, "defer_raw_materialization", True)),
"thread_alive": bool(self._capture_thread is not None and self._capture_thread.is_alive()),
"latest_seq": int(getattr(self, "_latest_packet_seq", 0)),
"last_consumed_seq": int(getattr(self, "_last_consumed_packet_seq", 0)),

View File

@ -872,7 +872,7 @@ class RawProcessorCore:
# OAK_CORE_PERF_LOG_INTERVAL_S=1.0
# OAK_CORE_SHAPES_LOG_INTERVAL_S=5.0
self.core_perf_debug = str(
os.getenv("OAK_CORE_PERF_DEBUG", "1")
os.getenv("OAK_CORE_PERF_DEBUG", "0")
).strip().lower() not in ("0", "false", "no", "off")
try:
self.core_perf_log_interval_s = max(0.2, float(

View File

@ -1,6 +1,8 @@
import base64
import datetime
import json
import os
import queue
import threading
import time
from pathlib import Path
@ -26,7 +28,7 @@ from visual_worker.utils import converter_valores_numpy
from shared.perf_monitor import VisualPerfMonitor
CAMERA_MANAGER_VERSION = "production_v1_2026_08_24"
CAMERA_MANAGER_VERSION = "production_v1_2026_09_12_runtime_priority_v3"
PRODUCT_ASSEMBLY_SCHEMA = "multispec_module_params_assembly_v1"
PRODUCT_MODULE_SCHEMA = "multispec_module_params_v3"
@ -87,6 +89,8 @@ class CameraManager:
self._loop_deteccao_iniciado = False
self._loop_analise_iniciado = False
self._loop_stream_iniciado = False
self._loop_preview_iniciado = False
self._loop_replay_writer_iniciado = False
self._loop_publicacao_iniciado = False
self._vida_lock = threading.RLock()
@ -101,6 +105,25 @@ class CameraManager:
self._pub_lock = threading.RLock()
self._preview_lock = threading.RLock()
# Preview/replay são explicitamente secundários ao runtime.
# Produção visual ocorre em loop próprio, em baixa frequência,
# e consumidores só leem o cache pronto.
self.preview_fps = 1.0
self.preview_max_width = 480
self.preview_jpeg_quality = 45
self._preview_encoded_cache = {}
# Snapshot local do estado Redis usado pelo hot path.
# Operação/controle/contexto são atualizados em até 100 ms; equipamento
# é quase estático e pode ser atualizado bem mais devagar.
self.runtime_snapshot_interval_s = 0.10
self.runtime_equipment_interval_s = 2.0
self.runtime_snapshot_max_age_s = 0.50
# Escrita de replay é assíncrona e descartável sob pressão.
# Nunca bloqueia tensor/inferência/detecção por causa de disco.
self._replay_save_queue = queue.Queue(maxsize=32)
self.perf = VisualPerfMonitor(janela=180)
self._ultimo_posproc_save_ts = 0.0
@ -164,6 +187,8 @@ class CameraManager:
self._ultimo_preview_seg_ts = 0.0
self._ultimo_preview_overlay_ts = 0.0
self._ultimo_preview_debug_ts = 0.0
with self._preview_lock:
self._preview_encoded_cache = {}
self._fps_infer_last_ts = None
self._fps_infer_ema = 0.0
@ -172,6 +197,16 @@ class CameraManager:
self._ultimo_loop_analise_fps = 0.0
self._ultimo_config_update_ts = 0.0
self._ultimo_runtime_snapshot_ts = 0.0
self._ultimo_equipment_snapshot_ts = 0.0
self._runtime_equipment_cache = {}
self._runtime_ctx_snapshot = {
"ts": 0.0,
"operacao": {},
"controle": {},
"contexto": {},
"equipamento": {},
}
# Diagnostico fino do pipeline. Estes campos sao somente observabilidade:
# nao alteram frequencias, caches, gates ou comandos dos bicos.
@ -352,10 +387,20 @@ class CameraManager:
)
fusion = mp.get("fusion_config", {}) or {}
mp_target = self._normalizar_size_wh(
fusion.get("target_size"),
mp_target_raw = fusion.get("target_size")
# Contrato atual:
# - fusion_config.target_size pode ser None no module_params;
# - ONNX estático é a autoridade do tamanho final;
# - se um MP legado trouxer target_size, ele continua validado.
mp_target = (
None
if mp_target_raw is None
else self._normalizar_size_wh(
mp_target_raw,
"module_params.fusion_config.target_size",
)
)
sizes_raw = mp.get("sensor_size_by_role", {}) or {}
hardware = mp.get("camera_hardware", {}) or {}
@ -408,11 +453,9 @@ class CameraManager:
"""
Fecha a autoridade da resolução antes de abrir a OAK:
ONNX input [W,H]
==
module_params.fusion_config.target_size
==
seg_config.ia_resolution
ONNX input [W,H] = autoridade final do runtime
module_params.fusion_config.target_size = opcional/legado
seg_config.ia_resolution = opcional, mas deve casar com ONNX
camera_width/camera_height deixam de participar da decisão física.
"""
@ -422,9 +465,16 @@ class CameraManager:
mx_id,
)
mp_target = list(mp_info["target_size"])
mp_target_raw = mp_info.get("target_size")
mp_target = (
None
if mp_target_raw is None
else list(mp_target_raw)
)
if model_target != mp_target:
# MP novo pode omitir o target final. Nesse caso o ONNX manda.
# MP legado com target explícito continua fail-closed se divergir.
if mp_target is not None and model_target != mp_target:
raise RuntimeError(
"Contrato de resolução incompatível entre ONNX e module_params: "
f"ONNX={model_target} module_params={mp_target}"
@ -439,7 +489,7 @@ class CameraManager:
if cfg_target != model_target:
raise RuntimeError(
"Contrato de resolução incompatível entre config, ONNX e MP: "
"Contrato de resolução incompatível entre config e ONNX: "
f"config={cfg_target} ONNX={model_target} MP={mp_target}"
)
else:
@ -763,7 +813,8 @@ class CameraManager:
self.mx_id = str(mx_id)
from weed_worker.config import load_seg_config
self.seg_config = load_seg_config()
runtime_snapshot = self._atualizar_snapshot_runtime(force=True)
self.seg_config = load_seg_config(runtime_snapshot=runtime_snapshot)
self.posproc_intervalo_min_s = float(
self.seg_config.get("posproc_intervalo_min_s", 5.0)
)
@ -1067,7 +1118,7 @@ class CameraManager:
if self.weed_detector is not None:
return
self.weed_detector = WeedDetector()
self.weed_detector = WeedDetector(config=self.seg_config)
def _iniciar_loops_se_necessario(self, seg_config):
freq_analise = float(seg_config.get("analise_fps", 15.0))
@ -1076,6 +1127,20 @@ class CameraManager:
freq_inferencia = float(seg_config.get("inferencia_fps", freq_tensor))
freq_deteccao = float(seg_config.get("deteccao_fps", freq_inferencia))
# Parâmetros visuais são de baixa prioridade e podem ser ajustados
# sem tocar no contrato científico/runtime.
self.preview_fps = max(0.2, float(seg_config.get("preview_fps", 1.0) or 1.0))
self.preview_max_width = max(160, int(seg_config.get("preview_max_width", 480) or 480))
self.preview_jpeg_quality = max(20, min(85, int(seg_config.get("preview_jpeg_quality", 45) or 45)))
if not self._loop_preview_iniciado:
self._iniciar_loop_preview_cache(freq=self.preview_fps)
self._loop_preview_iniciado = True
if not self._loop_replay_writer_iniciado:
self._iniciar_loop_replay_writer()
self._loop_replay_writer_iniciado = True
if not self._loop_stream_iniciado:
stream = getattr(self.camera, "stream", None)
stream_fps = float(getattr(stream, "_op_fps", 2.0) or 2.0)
@ -1261,7 +1326,7 @@ class CameraManager:
if detector is None or reset_nome is None:
# Fallback seguro: recria apenas o detector, nunca câmera/modelo.
self.weed_detector = WeedDetector()
self.weed_detector = WeedDetector(config=self.seg_config)
if self.seg_config and hasattr(self.weed_detector, "atualizar_config"):
self.weed_detector.atualizar_config(self.seg_config)
reset_nome = "recriado"
@ -1540,6 +1605,63 @@ class CameraManager:
except Exception as e:
self.mostrar_log(f"[saude] erro: {e}")
def _atualizar_snapshot_runtime(self, intervalo_s=None, force=False):
"""
Atualiza um snapshot local de Redis fora do caminho frame-a-frame.
Atribuição do dict final é atômica em CPython: leitores do detector
enxergam o snapshot antigo ou o novo, nunca um objeto parcialmente
montado. Em caso de falha, preserva o último snapshot válido; o gate
fail-closed verifica a idade antes de liberar pulverização.
"""
agora = time.time()
if intervalo_s is None:
intervalo_s = float(self.runtime_snapshot_interval_s)
atual = self._runtime_ctx_snapshot or {}
if (
not force
and float(atual.get("ts", 0.0) or 0.0) > 0.0
and (agora - float(self._ultimo_runtime_snapshot_ts or 0.0)) < float(intervalo_s)
):
return atual
try:
operacao = ContextoGlobalRedis.get_operacao() or {}
controle = ContextoGlobalRedis.get_controle() or {}
contexto = ContextoGlobalRedis.get_contexto() or {}
equipamento = self._runtime_equipment_cache
precisa_equip = (
force
or not equipamento
or (agora - float(self._ultimo_equipment_snapshot_ts or 0.0))
>= float(self.runtime_equipment_interval_s)
)
if precisa_equip:
equipamento_novo = ContextoGlobalRedis.get_equipamento() or {}
if equipamento_novo:
equipamento = equipamento_novo
self._runtime_equipment_cache = equipamento_novo
self._ultimo_equipment_snapshot_ts = agora
novo = {
"ts": float(agora),
"operacao": operacao,
"controle": controle,
"contexto": contexto,
"equipamento": equipamento or {},
}
self._runtime_ctx_snapshot = novo
self._ultimo_runtime_snapshot_ts = agora
return novo
except Exception as e:
# Não transforma uma falha transitória de Redis em exceção dentro
# do detector. Snapshot antigo será rejeitado pelo gate se envelhecer.
self.perf.inc("runtime_snapshot_erros")
return atual
def _atualizar_config_dinamica_detector(self, intervalo_s=0.20):
"""
Atualiza config do detector e a cabeça ONNX em baixa frequência.
@ -1560,7 +1682,10 @@ class CameraManager:
try:
from weed_worker.config import load_seg_config
cfg = load_seg_config()
snapshot = self._atualizar_snapshot_runtime(
intervalo_s=self.runtime_snapshot_interval_s
)
cfg = load_seg_config(runtime_snapshot=snapshot)
cfg = self._validar_config_dinamica_contrato(cfg)
self.seg_config = cfg
self.qtd_bicos = int(cfg.get("qtd_bicos", self.qtd_bicos or 7) or 7)
@ -1773,73 +1898,151 @@ class CameraManager:
return None, None, None
def get_selected_frame(self, _frame_type: TipoFrameCamera):
agora = time.time()
"""
Compatibilidade externa: retorna SOMENTE o último preview cacheado.
try:
Regra de produto:
- nunca chama preview_infer_cached();
- nunca monta overlay/debug sob demanda;
- nunca força inferência;
- nunca toca no cursor da câmera.
O produtor de preview é _iniciar_loop_preview_cache().
"""
if _frame_type not in self.FRAME_TYPES_PREVIEW:
return None
# Se ainda não tem runtime/tensor, tenta devolver o último preview válido.
if self.model_svc is None or self._ultimo_raw_input is None:
return self._get_cached_preview_frame(_frame_type)
# Janela curta: usa cache recente.
if (agora - self._ultimo_preview_ts) < 0.20:
cached = self._get_cached_preview_frame(_frame_type)
if cached is not None:
return cached
rgb_frame, seg_frame, overlay_frame, _, _ = self.model_svc.preview_infer_cached(
self._ultimo_raw_input,
self._ultimo_predictions,
alpha=0.5,
)
debug_frame = None
if _frame_type == TipoFrameCamera.Debug:
debug_frame = self.get_debug_frame(
mostrar=False,
overlay_bgr=overlay_frame,
)
# Atualiza o timestamp da tentativa de preview.
self._ultimo_preview_ts = agora
# Importante:
# só atualiza cache se o frame novo for válido.
# Nunca apaga último válido com None.
self._set_cached_preview_frame(TipoFrameCamera.Rgb, rgb_frame, agora)
self._set_cached_preview_frame(TipoFrameCamera.Segmentacao, seg_frame, agora)
self._set_cached_preview_frame(TipoFrameCamera.Overlay, overlay_frame, agora)
self._set_cached_preview_frame(TipoFrameCamera.Debug, debug_frame, agora)
# Retorna o frame pedido.
# Se o novo veio None, devolve o último válido cacheado.
return self._get_cached_preview_frame(_frame_type)
except Exception as e:
self.mostrar_log(f"[weed] erro em get_selected_frame: {e}")
# Mesmo em erro, tenta devolver último válido.
try:
return self._get_cached_preview_frame(_frame_type)
except Exception:
def _get_cached_preview_frame_ref(self, frame_type):
"""Retorna referência imutável do cache visual, sem cópia."""
with self._preview_lock:
if frame_type == TipoFrameCamera.Rgb:
return self._ultimo_preview_rgb
if frame_type == TipoFrameCamera.Segmentacao:
return self._ultimo_preview_seg
if frame_type == TipoFrameCamera.Overlay:
return self._ultimo_preview_overlay
if frame_type == TipoFrameCamera.Debug:
return self._ultimo_preview_debug
return None
def _get_cached_preview_frame(self, frame_type):
with self._preview_lock:
if frame_type == TipoFrameCamera.Rgb:
frame = self._ultimo_preview_rgb
elif frame_type == TipoFrameCamera.Segmentacao:
frame = self._ultimo_preview_seg
elif frame_type == TipoFrameCamera.Overlay:
frame = self._ultimo_preview_overlay
elif frame_type == TipoFrameCamera.Debug:
frame = self._ultimo_preview_debug
else:
frame = None
# API legada pode receber uma cópia pequena (preview já reduzido).
# O hot path interno usa _get_cached_preview_frame_ref().
return self._copiar_frame(self._get_cached_preview_frame_ref(frame_type))
return self._copiar_frame(frame)
def get_cached_preview_payload(self, frame_type):
"""
Retorna payload já comprimido/convertido para Base64.
É o caminho do GetCameraFrame via Redis e não executa OpenCV.
"""
if frame_type not in self.FRAME_TYPES_PREVIEW:
return None
with self._preview_lock:
item = self._preview_encoded_cache.get(frame_type)
if not item:
return None
return {
"base64": item.get("base64"),
"jpeg": item.get("jpeg"),
"source_ts": float(item.get("source_ts", 0.0) or 0.0),
"width": int(item.get("width", 0) or 0),
"height": int(item.get("height", 0) or 0),
"quality": int(item.get("quality", self.preview_jpeg_quality) or self.preview_jpeg_quality),
}
def get_cached_preview_jpeg(self, frame_type):
"""Retorna bytes JPEG imutáveis do cache para replay em disco."""
with self._preview_lock:
item = self._preview_encoded_cache.get(frame_type)
if not item:
return None, 0.0
return item.get("jpeg"), float(item.get("source_ts", 0.0) or 0.0)
def _preparar_preview_leve(self, frame):
if not self._frame_valido(frame):
return None
arr = frame
if arr.dtype != np.uint8:
arr = self._normalizar_frame_para_uint8(arr)
h, w = arr.shape[:2]
max_w = int(self.preview_max_width)
if w > max_w:
escala = max_w / float(w)
novo_h = max(1, int(round(h * escala)))
arr = cv2.resize(arr, (max_w, novo_h), interpolation=cv2.INTER_AREA)
else:
# Isola o cache de buffers reaproveitados pelo gerador de debug.
arr = arr.copy()
return np.ascontiguousarray(arr)
def _cache_preview_bundle(self, frames_por_tipo: dict, source_ts: float):
"""
Reduz, comprime e publica atomicamente o bundle visual.
Este método é chamado SOMENTE pelo preview producer.
"""
novos = {}
quality = int(self.preview_jpeg_quality)
for frame_type, frame in frames_por_tipo.items():
leve = self._preparar_preview_leve(frame)
if not self._frame_valido(leve):
continue
ok, enc = cv2.imencode(
".jpg",
leve,
[int(cv2.IMWRITE_JPEG_QUALITY), quality],
)
if not ok:
continue
jpeg_bytes = enc.tobytes()
b64 = base64.b64encode(jpeg_bytes).decode("ascii")
h, w = leve.shape[:2]
novos[frame_type] = {
"frame": leve,
"jpeg": jpeg_bytes,
"base64": b64,
"source_ts": float(source_ts),
"width": int(w),
"height": int(h),
"quality": quality,
}
if not novos:
return False
with self._preview_lock:
self._preview_encoded_cache.update(novos)
item = novos.get(TipoFrameCamera.Rgb)
if item:
self._ultimo_preview_rgb = item["frame"]
self._ultimo_preview_rgb_ts = source_ts
item = novos.get(TipoFrameCamera.Segmentacao)
if item:
self._ultimo_preview_seg = item["frame"]
self._ultimo_preview_seg_ts = source_ts
item = novos.get(TipoFrameCamera.Overlay)
if item:
self._ultimo_preview_overlay = item["frame"]
self._ultimo_preview_overlay_ts = source_ts
item = novos.get(TipoFrameCamera.Debug)
if item:
self._ultimo_preview_debug = item["frame"]
self._ultimo_preview_debug_ts = source_ts
self._ultimo_preview_ts = float(source_ts)
return True
def get_debug_frame(self, mostrar=False, overlay_bgr=None):
overlay = overlay_bgr if overlay_bgr is not None else self._ultimo_preview_overlay
@ -1854,7 +2057,9 @@ class CameraManager:
"fps_loop": self._ultimo_loop_analise_fps,
}
self._ultimo_preview_debug = self._montar_debug_overlay(
# Não publica diretamente no cache: o preview producer faz a
# publicação atômica depois de reduzir/comprimir o bundle inteiro.
return self._montar_debug_overlay(
overlay_bgr=overlay,
atuacao_bicos=self._ultimo_controle or {},
config=self.seg_config or {},
@ -1862,8 +2067,6 @@ class CameraManager:
mostrar=mostrar,
)
return self._ultimo_preview_debug
def _montar_debug_overlay(
self,
overlay_bgr,
@ -2012,42 +2215,151 @@ class CameraManager:
return frame
def _set_cached_preview_frame(self, frame_type, frame, ts=None):
if not self._frame_valido(frame):
return False
ts = time.time() if ts is None else ts
frame_copy = self._copiar_frame(frame)
if not self._frame_valido(frame_copy):
return False
with self._preview_lock:
if frame_type == TipoFrameCamera.Rgb:
self._ultimo_preview_rgb = frame_copy
self._ultimo_preview_rgb_ts = ts
return True
if frame_type == TipoFrameCamera.Segmentacao:
self._ultimo_preview_seg = frame_copy
self._ultimo_preview_seg_ts = ts
return True
if frame_type == TipoFrameCamera.Overlay:
self._ultimo_preview_overlay = frame_copy
self._ultimo_preview_overlay_ts = ts
return True
if frame_type == TipoFrameCamera.Debug:
self._ultimo_preview_debug = frame_copy
self._ultimo_preview_debug_ts = ts
return True
return False
"""Compatibilidade interna; publica um único frame no cache leve."""
ts = time.time() if ts is None else float(ts)
return self._cache_preview_bundle({frame_type: frame}, ts)
# ============================================================
# Loops
# ============================================================
def _iniciar_loop_preview_cache(self, freq=1.0):
"""
Produtor visual de baixa prioridade.
Ele reaproveita o ÚLTIMO tensor + ÚLTIMA predição já calculados pelo
runtime e monta os previews em frequência baixa. Consumidores nunca
provocam esta construção.
"""
def loop():
ultimo_pred_marker = None
while True:
t0 = time.time()
t_perf0 = time.perf_counter()
did_work = False
try:
if self._em_warmup or not self.operante:
time.sleep(0.20)
continue
model_svc = self.model_svc
with self._pred_lock:
pred_cache = self._pred_cache
predictions = pred_cache.get("predictions")
raw_input = pred_cache.get("raw_input")
pred_ts = float(pred_cache.get("ts", 0.0) or 0.0)
tensor_ts = float(pred_cache.get("tensor_ts", 0.0) or 0.0)
generation = int(pred_cache.get("generation", -1))
if raw_input is None:
raw_input = self._ultimo_raw_input
if model_svc is None or raw_input is None or predictions is None:
time.sleep(0.10)
continue
# Marcador lógico do frame, independente do endereço Python
# do ndarray. Evita falso "mesmo frame" por reutilização de id().
pred_marker = (generation, pred_ts, tensor_ts)
if pred_marker == ultimo_pred_marker:
time.sleep(min(0.10, max(0.02, 1.0 / max(freq, 0.2))))
continue
did_work = True
rgb_frame, seg_frame, overlay_frame, _, _ = model_svc.preview_infer_cached(
raw_input,
predictions,
alpha=0.5,
)
debug_frame = self.get_debug_frame(
mostrar=False,
overlay_bgr=overlay_frame,
)
source_ts = pred_ts if pred_ts > 0.0 else time.time()
ok = self._cache_preview_bundle(
{
TipoFrameCamera.Rgb: rgb_frame,
TipoFrameCamera.Segmentacao: seg_frame,
TipoFrameCamera.Overlay: overlay_frame,
TipoFrameCamera.Debug: debug_frame,
},
source_ts=source_ts,
)
if ok:
ultimo_pred_marker = pred_marker
except Exception as e:
self.perf.inc("preview_cache_erros")
self.mostrar_log(f"[weed] erro no preview producer: {e}")
finally:
t_perf1 = time.perf_counter()
if did_work:
self.perf.tick(
"preview_cache",
latencia_ms=(t_perf1 - t_perf0) * 1000.0,
)
dt = time.time() - t0
time.sleep(max(0.0, (1.0 / max(freq, 0.2)) - dt))
threading.Thread(
target=loop,
daemon=True,
name="weed-preview-cache",
).start()
def _iniciar_loop_replay_writer(self):
"""Único escritor assíncrono dos JPEGs de replay."""
def loop():
while True:
try:
caminho, jpeg_bytes = self._replay_save_queue.get()
try:
os.makedirs(os.path.dirname(caminho), exist_ok=True)
with open(caminho, "wb") as f:
f.write(jpeg_bytes)
finally:
self._replay_save_queue.task_done()
except Exception as e:
self.mostrar_log(f"[weed][REPLAY_IO] erro: {e}")
time.sleep(0.05)
threading.Thread(
target=loop,
daemon=True,
name="weed-replay-writer",
).start()
def _enfileirar_replay_jpeg(self, caminho: str, jpeg_bytes: bytes) -> bool:
if not jpeg_bytes:
return False
try:
self._replay_save_queue.put_nowait((caminho, jpeg_bytes))
return True
except queue.Full:
# Replay é best-effort. Sob pressão, descartamos o mais antigo
# para nunca sacrificar o runtime científico.
try:
self._replay_save_queue.get_nowait()
self._replay_save_queue.task_done()
except Exception:
pass
try:
self._replay_save_queue.put_nowait((caminho, jpeg_bytes))
self.perf.inc("replay_queue_drop_oldest")
return True
except queue.Full:
self.perf.inc("replay_queue_drop")
return False
def _iniciar_loop_captura_tensor(self, freq=25.0):
def loop():
periodo = 1.0 / max(float(freq), 0.1)
@ -2273,6 +2585,7 @@ class CameraManager:
"ts": pred_ts,
"tensor_ts": tensor_ts,
"predictions": predictions,
"raw_input": tensor5,
"res": res,
"infer_ms": float(infer_ms),
"infer_gpu_ms": float(infer_forward_ms or 0.0),
@ -2368,6 +2681,9 @@ class CameraManager:
t_cfg0 = time.perf_counter()
cfg_cpu0 = time.thread_time()
self._atualizar_snapshot_runtime(
intervalo_s=self.runtime_snapshot_interval_s
)
self._atualizar_config_dinamica_detector(intervalo_s=0.20)
t_cfg1 = time.perf_counter()
config_ms = (t_cfg1 - t_cfg0) * 1000.0
@ -2419,10 +2735,6 @@ class CameraManager:
analise = analise_completa.get("dados_visuais", {})
detector_ms = (t_det1 - t_det0) * 1000.0
t_conv0 = time.perf_counter()
analise_convertida = converter_valores_numpy(analise)
t_conv1 = time.perf_counter()
# ====================================================
# Gate de pulverização
# ====================================================
@ -2435,6 +2747,12 @@ class CameraManager:
analise["pulverizacao"] = debug_pulverizacao
t_ctrl1 = time.perf_counter()
# Converte SOMENTE depois do gate, para que telemetria e
# comando publiquem exatamente o mesmo estado final.
t_conv0 = time.perf_counter()
analise_convertida = converter_valores_numpy(analise)
t_conv1 = time.perf_counter()
# Publicação atômica em relação a fechar/substituir câmera.
# Se a geração mudou, nenhum dado antigo chega aos caches/bicos.
with self._vida_lock:
@ -2668,11 +2986,15 @@ class CameraManager:
frame_type = TipoFrameCamera(
camera_ctx.get("frame_type", TipoFrameCamera.Rgb.value)
)
frame = self.get_selected_frame(frame_type)
# Stream também é consumidor de cache. Não monta preview.
frame = self._get_cached_preview_frame_ref(frame_type)
self.camera.enviar_frame_tcp(frame)
if self.debug_visual:
self.get_debug_frame(mostrar=True)
dbg = self._get_cached_preview_frame_ref(TipoFrameCamera.Debug)
if dbg is not None:
cv2.imshow("Debug Weed Worker", dbg)
cv2.waitKey(1)
except Exception as e:
self.mostrar_log(f"Erro no loop de stream: {e}")
@ -2923,9 +3245,23 @@ class CameraManager:
}
try:
operacao = ContextoGlobalRedis.get_operacao()
controle = ContextoGlobalRedis.get_controle()
contexto = ContextoGlobalRedis.get_contexto()
snapshot = self._runtime_ctx_snapshot or {}
snapshot_ts = float(snapshot.get("ts", 0.0) or 0.0)
snapshot_age_s = (time.time() - snapshot_ts) if snapshot_ts > 0.0 else float("inf")
debug["snapshot_age_ms"] = (
snapshot_age_s * 1000.0
if np.isfinite(snapshot_age_s)
else None
)
# Segurança: nunca libera pulverização a partir de estado velho.
if snapshot_age_s > float(self.runtime_snapshot_max_age_s):
debug["motivo"] = f"runtime snapshot antigo: {snapshot_age_s:.3f}s"
return False, debug
operacao = snapshot.get("operacao") or {}
controle = snapshot.get("controle") or {}
contexto = snapshot.get("contexto") or {}
traj = (contexto.get("Trajetoria", {}) or {})
gerais = (contexto.get("Gerais", {}) or {})
@ -3232,6 +3568,12 @@ class CameraManager:
}
resumo["pub_debug"] = getattr(self, "_ultimo_pub_debug", {})
resumo["preview_cache"] = {
"fps_config": float(self.preview_fps),
"max_width": int(self.preview_max_width),
"jpeg_quality": int(self.preview_jpeg_quality),
"replay_queue_size": int(self._replay_save_queue.qsize()),
}
self._set_pub_cache("performance_weed", resumo)
@ -3248,6 +3590,7 @@ class CameraManager:
det = loops.get("deteccao", {})
pub = loops.get("publicacao", {})
stream = loops.get("stream", {})
preview = loops.get("preview_cache", {})
pub_dbg = getattr(self, "_ultimo_pub_debug", {})
@ -3257,7 +3600,8 @@ class CameraManager:
f"inf={self._fps(inf):.1f} "
f"det={self._fps(det):.1f} "
f"pub={self._fps(pub):.1f} "
f"stream={self._fps(stream):.1f} | "
f"stream={self._fps(stream):.1f} "
f"preview={self._fps(preview):.1f} | "
f"period inf={self._fmt(self._per(inf))}ms "
f"det={self._fmt(self._per(det))}ms "
f"tensor={self._fmt(self._per(tensor))}ms"
@ -3285,7 +3629,9 @@ class CameraManager:
f"redis={self._fmt(self._m(pub, 'redis_ms'))} "
f"debug_pub={pub_dbg.get('publicou', 0)} "
f"campos={pub_dbg.get('campos', 0)} "
f"redis_dbg={pub_dbg.get('redis_ms', 0):.1f}ms"
f"redis_dbg={pub_dbg.get('redis_ms', 0):.1f}ms | "
f"PREVIEW total={self._fmt(self._lat(preview))} "
f"q={self._replay_save_queue.qsize()}"
)
self.mostrar_log(
@ -3468,38 +3814,26 @@ class CameraManager:
continue
# ====================================================
# 2) FLUXO ATUAL: imagens para replay
# 2) REPLAY VISUAL: grava exatamente o JPEG já cacheado.
# Nenhum resize/overlay/imencode é executado sob demanda.
# ====================================================
frame = self.get_selected_frame(tipo_enum)
jpeg_bytes, source_ts = self.get_cached_preview_jpeg(tipo_enum)
if not self._frame_valido(frame):
if not jpeg_bytes:
self.mostrar_log(
f"⚠️ Frame não salvo | "
f"tipo={tipo_nome} "
f"nome={nome_frame} "
f"motivo=sem_frame_valido_em_cache"
f"motivo=sem_jpeg_cacheado"
)
continue
frame = self._copiar_frame(frame)
caminho = os.path.join(pasta, f"{nome_frame}.jpg")
# Replay leve: JPEG comprimido.
# Mantém baixo uso de disco durante operação longa.
if frame.dtype != "uint8":
frame_salvar = self._normalizar_frame_para_uint8(frame)
else:
frame_salvar = frame
ok = cv2.imwrite(caminho, frame_salvar, [int(cv2.IMWRITE_JPEG_QUALITY), 70])
if not ok or not os.path.exists(caminho):
ok = self._enfileirar_replay_jpeg(caminho, jpeg_bytes)
if not ok:
self.mostrar_log(
f"❌ Falha ao salvar frame | "
f"tipo={tipo_nome} "
f"nome={nome_frame} "
f"caminho={caminho}"
f"⚠️ Replay descartado para proteger runtime | "
f"tipo={tipo_nome} nome={nome_frame}"
)
continue

View File

@ -28,9 +28,9 @@ Hardware/calibração:
module_params.json homologado
Resolução final do tensor:
ONNX input H/W
==
module_params.fusion_config.target_size
ONNX input H/W é a autoridade operacional.
module_params.fusion_config.target_size pode ser null no contrato novo;
quando preenchido, é apenas compatibilidade/default e deve casar com o ONNX.
Ordem nominal dos canais:
input_channels deste config
@ -171,7 +171,7 @@ WEED_DEFAULT_CONFIG = {
# ========================================================
"debug_visual": False,
"debug_perf": False,
"detector_debug_perf": True,
"detector_debug_perf": False,
# RAW científico para pós-processamento.
"posproc_intervalo_min_s": 5.0,
@ -577,10 +577,12 @@ def _positive(
def aplicar_overrides_redis(
cfg: dict,
runtime_snapshot: dict | None = None,
) -> dict:
snapshot = runtime_snapshot if isinstance(runtime_snapshot, dict) else None
operacao = (
ContextoGlobalRedis
.get_operacao()
(snapshot.get("operacao") if snapshot is not None else ContextoGlobalRedis.get_operacao())
or {}
)
@ -593,14 +595,12 @@ def aplicar_overrides_redis(
)
contexto = (
ContextoGlobalRedis
.get_contexto()
(snapshot.get("contexto") if snapshot is not None else ContextoGlobalRedis.get_contexto())
or {}
)
equipamento = (
ContextoGlobalRedis
.get_equipamento()
(snapshot.get("equipamento") if snapshot is not None else ContextoGlobalRedis.get_equipamento())
or {}
)
@ -1011,18 +1011,16 @@ def normalizar_config_runtime(
# API
# ============================================================
def load_seg_config():
def load_seg_config(runtime_snapshot: dict | None = None):
"""
Retorna snapshot NOVO a cada chamada.
Retorna um snapshot NOVO de configuração.
Não existe cache de valores dinâmicos:
- velocidade;
- cabeça;
- classe semântica;
- qtd_bicos;
- zona de atuação.
Se ``runtime_snapshot`` for informado, usa os valores já lidos pelo
CameraManager (operacao/contexto/equipamento) e NÃO toca no Redis.
Isso permite atualizar configuração em baixa frequência sem colocar
leituras externas no hot path de detecção.
O lock só impede leituras concorrentes inconsistentes durante a montagem.
Sem snapshot, preserva o comportamento legado e lê diretamente do Redis.
"""
with _CONFIG_LOCK:
cfg = deepcopy(
@ -1030,7 +1028,8 @@ def load_seg_config():
)
cfg = aplicar_overrides_redis(
cfg
cfg,
runtime_snapshot=runtime_snapshot,
)
cfg = normalizar_config_runtime(

View File

@ -11,7 +11,7 @@ os.environ.setdefault("OMP_WAIT_POLICY", "PASSIVE")
os.environ.setdefault("KMP_BLOCKTIME", "0")
# Limita kernels Numba, caso o Visual Worker utilize Numba
# direta ou indiretamente.
os.environ.setdefault("NUMBA_NUM_THREADS", "4")
os.environ.setdefault("NUMBA_NUM_THREADS", "16")
def main():
@ -21,12 +21,16 @@ def main():
# Importar e configurar OpenCV antes dos módulos do Visual Worker,
# pois eles podem carregar OpenCV internamente.
import cv2
cv2.setNumThreads(2)
opencv_threads = max(1, int(os.environ.get("WEED_OPENCV_THREADS", "8")))
cv2.setNumThreads(opencv_threads)
cv2.ocl.setUseOpenCL(False)
print(
f"[weed][PERF] OpenCV threads requested={opencv_threads} "
f"effective={cv2.getNumThreads()} preview_cache=True"
)
from weed_worker.config import mostrar_log, get_camera_manager, iniciar_camera_manager
from shared.enums import WeedWorkerCommandType, TipoFrameCamera
from shared.utils import encode_image_base64
from shared.contexto_global_redis import ContextoGlobalRedis, CmdKey, CtxKey
def loop_ativo():
@ -65,16 +69,28 @@ def main():
get_camera_manager().atualizar_saude_camera()
elif acao == WeedWorkerCommandType.GetCameraFrame:
tipo = TipoFrameCamera(dados.get("params", TipoFrameCamera.Rgb.value))
frame = get_camera_manager().get_selected_frame(tipo)
if frame is not None:
base64_img = encode_image_base64(frame)
if base64_img is not None:
# Leitura barata: JPEG/Base64 já foi produzido pelo loop de preview.
# Este comando nunca monta imagem, nunca roda OpenCV e nunca infere.
payload = get_camera_manager().get_cached_preview_payload(tipo)
if payload is not None and payload.get("base64"):
resposta = {
"frame": base64_img,
"frame": payload["base64"],
# Mantém compatibilidade com o C#: timestamp da RESPOSTA
# precisa ser posterior ao instante em que ele enviou o request.
"timestamp": time.time(),
"tipo": tipo.value
"frame_timestamp": payload.get("source_ts", 0.0),
"tipo": tipo.value,
"width": payload.get("width", 0),
"height": payload.get("height", 0),
}
ContextoGlobalRedis.publicar_comando(CmdKey.WeedWorkerTx, { "cmd": WeedWorkerCommandType.GetCameraFrame.value, "params": resposta })
ContextoGlobalRedis.publicar_comando(
CmdKey.WeedWorkerTx,
{
"cmd": WeedWorkerCommandType.GetCameraFrame.value,
"params": resposta,
},
)
elif acao == WeedWorkerCommandType.SaveCameraFrames:
nome = dados.get("params", {}).get("nome", "")
pasta = dados.get("params", {}).get("caminho", "frames_salvos")

View File

@ -42,10 +42,15 @@ class WeedDetector:
CONTRATO_OFICIAL = "target_binary"
TARGET_ID = 1
def __init__(self):
def __init__(self, config: Optional[dict] = None):
# O CameraManager já mantém um snapshot de configuração. Quando ele é
# fornecido, não existe motivo para o detector reler Redis durante a
# construção. Mantemos o fallback legado para usos isolados/testes.
if config is None:
from weed_worker.config import load_seg_config
config = load_seg_config()
self.config = load_seg_config()
self.config = dict(config or {})
self.prediction_contract = str(
self.config.get("prediction_contract", self.CONTRATO_OFICIAL)
@ -120,6 +125,12 @@ class WeedDetector:
self._morf_cache_k: Optional[int] = None
self._morf_kernel = None
# Geometria discreta da grade depende somente de H/W, quantidade de
# bicos e quantidade de células. Evita reconstruir linspace/áreas em
# todo frame. O cache é invalidado automaticamente quando a geometria
# muda, pois _inicializar_estado() é chamado nesses casos.
self._grid_geometry_cache = {}
# Tempo/movimento.
self._last_update_ts: Optional[float] = None
self._distancia_deslocada_total_cm = 0.0
@ -305,11 +316,19 @@ class WeedDetector:
# mas aqui tratamos como velocidade em m/s.
vel_norm = vel_mps
# O contrato operacional é SEMPRE binário, inclusive quando a
# cabeça ONNX selecionada é semantic: o SegFormerService converte
# a classe semantic_target_class escolhida para 0/1 antes daqui.
# Construímos a máscara uma única vez e a reutilizamos no radar e
# na memória espacial.
mask_target = (predictions == self.TARGET_ID)
t_rad0 = time.perf_counter()
target_no_radar_frame, estat_target = self._decidir_target_no_radar(
predictions,
cfg,
vel_norm=vel_norm,
mask_target=mask_target,
)
t_rad1 = time.perf_counter()
@ -319,6 +338,7 @@ class WeedDetector:
cfg=cfg,
vel_mps=vel_mps,
agora=agora,
mask_target=mask_target,
)
t_bic1 = time.perf_counter()
@ -377,6 +397,7 @@ class WeedDetector:
predictions: np.ndarray,
cfg: dict,
vel_norm: float = 0.0,
mask_target: Optional[np.ndarray] = None,
):
on_global = float(cfg.get("min_frac_erva_global_on", 0.0020))
off_global = float(cfg.get("min_frac_erva_global_off", 0.0015))
@ -393,7 +414,12 @@ class WeedDetector:
h, w = predictions.shape[:2]
inv_total = 1.0 / max(1, h * w)
target_sum = int(self._is_target[predictions].sum())
if mask_target is None:
mask_target = (predictions == self.TARGET_ID)
# count_nonzero evita indexação por LUT + array temporário adicional
# quando a máscara já foi montada pelo ciclo principal.
target_sum = int(np.count_nonzero(mask_target))
frac_global = target_sum * inv_total
alpha = float(cfg.get("erva_frac_ema", 0.30))
@ -430,6 +456,7 @@ class WeedDetector:
cfg: dict,
vel_mps: float,
agora: float,
mask_target: Optional[np.ndarray] = None,
):
qtd_bicos = self.qtd_bicos
@ -460,6 +487,7 @@ class WeedDetector:
frac_grid, grid_info = self._calcular_frac_grid_por_bico_cell(
predictions=predictions,
cfg=cfg,
mask_target=mask_target,
)
self._frac_grid_last = frac_grid
@ -657,21 +685,92 @@ class WeedDetector:
return shifted
def _get_grid_geometry(self, h: int, w: int):
key = (int(h), int(w), int(self.qtd_bicos), int(self.num_cells))
cached = self._grid_geometry_cache.get(key)
if cached is not None:
return cached
x_edges = np.linspace(0, w, self.qtd_bicos + 1, dtype=np.int32)
y_edges = np.linspace(0, h, self.num_cells + 1, dtype=np.int32)
x_widths = np.diff(x_edges).astype(np.int32, copy=False)
y_heights = np.diff(y_edges).astype(np.int32, copy=False)
areas = (
y_heights[:, None].astype(np.int64)
* x_widths[None, :].astype(np.int64)
)
cached = {
"x_edges": x_edges,
"y_edges": y_edges,
"x_widths": x_widths,
"y_heights": y_heights,
"areas": areas,
"valid_cell_count": int(np.count_nonzero(y_heights > 0)),
"reduceat_safe": bool(
np.all(x_widths > 0)
and np.all(y_heights > 0)
),
}
# Uma única geometria é usada no runtime normal. Limita crescimento
# acidental caso algum caller varie resolução em teste.
if len(self._grid_geometry_cache) >= 4:
self._grid_geometry_cache.clear()
self._grid_geometry_cache[key] = cached
return cached
def _calcular_frac_grid_por_bico_cell(
self,
predictions: np.ndarray,
cfg: dict,
mask_target: Optional[np.ndarray] = None,
):
h, w = predictions.shape[:2]
mask_target = self._is_target[predictions]
if mask_target is None:
mask_target = (predictions == self.TARGET_ID)
mask_target = self._aplicar_filtros_opcionais(
mask_target=mask_target,
cfg=cfg,
)
# Integral image em inteiro com sinal.
# Evita overflow nos cálculos A - B - C + D.
geom = self._get_grid_geometry(h, w)
x_edges = geom["x_edges"]
y_edges = geom["y_edges"]
areas = geom["areas"]
if geom["reduceat_safe"]:
# Soma retangular em duas reduções segmentadas. Mantém exatamente
# os mesmos limites discretos de np.linspace usados pela versão
# anterior, mas evita integral int64 HxW + 700 loops Python.
# bool -> reduceat produz contagem inteira sem alterar a máscara.
sums_y = np.add.reduceat(
mask_target,
y_edges[:-1],
axis=0,
)
counts_yx = np.add.reduceat(
sums_y,
x_edges[:-1],
axis=1,
)
frac_cells_bicos = np.divide(
counts_yx,
areas,
out=np.zeros_like(areas, dtype=np.float32),
where=areas > 0,
).astype(np.float32, copy=False)
frac_grid = np.ascontiguousarray(
frac_cells_bicos.T,
dtype=np.float32,
)
else:
# Fallback para geometrias patológicas em que há bins vazios
# (ex.: mais células que pixels). Preserva a semântica antiga.
ii = np.pad(
mask_target.astype(np.int64, copy=False)
.cumsum(axis=0, dtype=np.int64)
@ -680,46 +779,28 @@ class WeedDetector:
mode="constant",
constant_values=0,
)
x_edges = np.linspace(0, w, self.qtd_bicos + 1, dtype=np.int32)
y_edges = np.linspace(0, h, self.num_cells + 1, dtype=np.int32)
frac_grid = np.zeros(
(self.qtd_bicos, self.num_cells),
dtype=np.float32,
)
valid_cell_count = 0
for c in range(self.num_cells):
y_top = int(y_edges[c])
y_bot = int(y_edges[c + 1])
if y_bot <= y_top:
continue
valid_cell_count += 1
cell_h = y_bot - y_top
for b in range(self.qtd_bicos):
x0 = int(x_edges[b])
x1 = int(x_edges[b + 1])
if x1 <= x0:
continue
area = float(cell_h * (x1 - x0))
area = int((y_bot - y_top) * (x1 - x0))
if area <= 0:
continue
total = (
ii[y_bot, x1]
- ii[y_top, x1]
- ii[y_bot, x0]
+ ii[y_top, x0]
)
frac_grid[b, c] = float(total) / area
frac_grid[b, c] = np.float32(float(total) / float(area))
# A segmentação é naturalmente da esquerda para a direita.
# Quando habilitado, inverte associação região imagem -> bico físico.
@ -727,12 +808,17 @@ class WeedDetector:
frac_grid = frac_grid[::-1, :].copy()
return frac_grid, {
"valid_cell_count": int(valid_cell_count),
"valid_cell_count": int(geom["valid_cell_count"]),
"height": int(h),
"width": int(w),
"x_edges": x_edges,
"y_edges": y_edges,
"zona_coord": "frac_top_to_bottom",
"grid_backend": (
"reduceat_segmented"
if geom["reduceat_safe"]
else "integral_fallback"
),
}
def _atualizar_score_memoria(

View File

@ -722,10 +722,11 @@
},
"flatfield_config": {
"enabled": true,
"npz_file": "calibration/flatfield_maps_v1.npz",
"npz_file": "flatfield_maps_v1.npz",
"apply_before_fusion": true,
"apply_after_decode": true,
"apply_space": "native_camera_space",
"saturation_guard_enabled": false,
"apply_space": "final_tensor_space",
"map_type": "gain",
"channels": [
"R",
@ -758,7 +759,7 @@
},
"subtract_dark": false,
"clip_output": true,
"json_file": "calibration/flatfield_maps_v1.json",
"json_file": "flatfield_maps_v1.json",
"schema": "multispec_flatfield_production_v2",
"created_at": "2026-09-08 16:31:07"
},

View File

@ -848,8 +848,8 @@
"created_at": "2026-05-08 15:26:18",
"json_file": "flatfield_maps_v1.json",
"npz_file": "flatfield_maps_v1.npz",
"apply_before_fusion": false,
"apply_after_decode": false,
"apply_before_fusion": true,
"apply_after_decode": true,
"apply_space": "final_tensor_space",
"map_type": "gain",
"formula": "channel_corrected = max(channel_linear - dark, 0) * gain_map",

View File

@ -45,6 +45,7 @@ class OakFcc3Client:
imu_modo="rotation_vector",
imu_freq_hz=200,
evaluate_quality=True,
require_product_contract=False,
hardware_sync_enabled=None,
frame_sync_master=None,
@ -65,6 +66,30 @@ class OakFcc3Client:
self.module_calibration_json = module_calibration_json
self.module_params = self._load_module_params(module_calibration_json)
self.fusion_config = self.module_params.get("fusion_config", {}) or {}
self.require_product_contract = bool(require_product_contract)
assembly = self.module_params.get("assembly_metadata", {}) or {}
self.product_contract = bool(
self.module_params.get("schema") == "multispec_module_params_v3"
and assembly.get("schema") == "multispec_module_params_assembly_v1"
)
if self.require_product_contract and not self.product_contract:
raise RuntimeError(
"OakFcc3Client exige module_params de produção homologado: "
f"schema={self.module_params.get('schema')!r} "
f"assembly={assembly.get('schema')!r}"
)
# Em produto, Bayer e raster RGB nativo vêm do MP, não de fallbacks do caller.
mp_bayer = str(self.module_params.get("bayer_pattern", "") or "").upper()
if mp_bayer:
self.bayer = mp_bayer
rgb_size = (self.module_params.get("sensor_size_by_role", {}) or {}).get("rgb")
if isinstance(rgb_size, (list, tuple)) and len(rgb_size) == 2:
self.width = int(rgb_size[0])
self.height = int(rgb_size[1])
self.imu_modo = str(imu_modo).strip().lower()
self.imu_freq_hz = int(imu_freq_hz)
@ -84,14 +109,15 @@ class OakFcc3Client:
self.svc = OakFcc3Service(
timeout=10,
fps=fps,
width=width,
height=height,
width=self.width,
height=self.height,
frame_type=frame_type,
output_dtype=output_dtype,
capture_mode=capture_mode,
raw_policy=raw_policy,
mx_id=self.mx_id,
module_calibration_json=module_calibration_json,
require_product_contract=self.require_product_contract,
imu_modo=self.imu_modo,
imu_freq_hz=self.imu_freq_hz,
@ -108,15 +134,15 @@ class OakFcc3Client:
self.applied_camera_controls = {}
self.radiometric_controller = None
self.core = RawProcessorCore(
sensor_width=width,
sensor_height=height,
bayer_pattern=bayer,
sensor_width=self.width,
sensor_height=self.height,
bayer_pattern=self.bayer,
calibration_json_path=module_calibration_json,
)
self.preview = RawProcessorPreview(
sensor_width=width,
sensor_height=height,
bayer_pattern=bayer,
sensor_width=self.width,
sensor_height=self.height,
bayer_pattern=self.bayer,
)
def __enter__(self):
@ -133,6 +159,46 @@ class OakFcc3Client:
with open(path, "r", encoding="utf-8") as f:
return json.load(f)
def get_contract(self):
"""
Contrato estático/runtime consumido por CameraMultispectral.
O module_params é autoridade de hardware/calibração. O target final
pode ser null no contrato novo; quando existir é apenas um default
compatível/legado. O Weed CameraManager passa o target ONNX ao caller.
"""
try:
service_contract = self.svc.get_contract() or {}
except Exception:
service_contract = {}
mp = self.module_params or {}
fusion = mp.get("fusion_config", {}) or {}
target = fusion.get("target_size")
default_target = None
if isinstance(target, (list, tuple)) and len(target) == 2:
tw, th = int(target[0]), int(target[1])
if tw > 0 and th > 0:
default_target = [tw, th]
sensor_sizes = mp.get("sensor_size_by_role", {}) or {}
camera_hw = mp.get("camera_hardware", {}) or {}
out = dict(service_contract)
out.update({
"product_contract": bool(self.product_contract),
"require_product_contract": bool(self.require_product_contract),
"module_params_schema": mp.get("schema"),
"module_calibration_json": self.module_calibration_json,
"sensor_size_by_role": sensor_sizes,
"camera_hardware": camera_hw,
"bayer_pattern": mp.get("bayer_pattern") or out.get("bayer_pattern"),
"frame_type": self.frame_type,
"raw_policy": self.raw_policy,
"default_target_size": default_target,
})
return out
def apply_module_camera_settings(self):
camera_settings = self.module_params.get("camera_settings", {}) or {}
@ -374,8 +440,8 @@ class OakFcc3Client:
evaluate_quality = self.evaluate_quality
evaluate_quality = bool(evaluate_quality)
tensor = self.core.fuse_multispec_cameras(decoded, meta, channels_expected)
tensor = self.core.resize_tensor_chw(tensor, target_size=target_size)
tensor = self.core.fuse_multispec_cameras(decoded, meta, channels_expected, target_size=target_size)
#tensor = self.core.resize_tensor_chw(tensor, target_size=target_size)
# Mantém paridade com build_infer_tensor_from_stream: se a calibração
# habilitar patch normalization, ela também vale no caminho decoded.

File diff suppressed because it is too large Load Diff

File diff suppressed because it is too large Load Diff

File diff suppressed because it is too large Load Diff

File diff suppressed because it is too large Load Diff