otimizado weed worker e trocado senha do ntrip ibge

This commit is contained in:
Diego Freitas 2026-09-14 08:27:36 -03:00
commit ceba226614
17 changed files with 7455 additions and 3626 deletions

View File

@ -1857,7 +1857,7 @@ namespace AgroBase.Services
// "AGRO_NTRIP_USERNAME" // "AGRO_NTRIP_USERNAME"
//) ?? string.Empty; //) ?? string.Empty;
string password = "c3pc7*9N"; string password = "kY3zd$*5";
//Environment.GetEnvironmentVariable( //Environment.GetEnvironmentVariable(
// "AGRO_NTRIP_PASSWORD" // "AGRO_NTRIP_PASSWORD"
//) ?? string.Empty; //) ?? string.Empty;

View File

@ -16,7 +16,7 @@ from camera_worker.oak_fcc3_core.oak_fcc3_client import OakFcc3Client
from camera_worker.camera_imu import IMUCamera 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: class CameraMultispectral:
@ -116,6 +116,7 @@ class CameraMultispectral:
self.ultimo_tensor_multispec = None self.ultimo_tensor_multispec = None
self.ultimo_frame_rgb = None self.ultimo_frame_rgb = None
self.ultimo_raw_multi = None self.ultimo_raw_multi = None
self.ultimo_raw_meta = None
self.ultimo_meta = None self.ultimo_meta = None
self.ultimo_decoded = None self.ultimo_decoded = None
@ -584,6 +585,7 @@ class CameraMultispectral:
self.ultimo_tensor_multispec = None self.ultimo_tensor_multispec = None
self.ultimo_frame_rgb = None self.ultimo_frame_rgb = None
self.ultimo_raw_multi = None self.ultimo_raw_multi = None
self.ultimo_raw_meta = None
self.ultimo_meta = None self.ultimo_meta = None
self.ultimo_decoded = None self.ultimo_decoded = None
@ -889,6 +891,9 @@ class CameraMultispectral:
# Arrays do pacote são imutáveis após a captura. # Arrays do pacote são imutáveis após a captura.
# Mantemos referências e copiamos somente no salvamento. # Mantemos referências e copiamos somente no salvamento.
self.ultimo_raw_multi = dict(frame) 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.timestamp_ultimo_raw_multi = ts_raw_perf
self._ultimo_resultado_raw = { self._ultimo_resultado_raw = {
"erro": None, "erro": None,
@ -1008,7 +1013,7 @@ class CameraMultispectral:
} }
def requisitar_frame_raw_multi(self, force: bool = False, max_age_s: float = None): 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: try:
if max_age_s is None or float(max_age_s) <= 0: if max_age_s is None or float(max_age_s) <= 0:
max_age_s = max(1.0, self._cache_max_age_s) max_age_s = max(1.0, self._cache_max_age_s)
@ -1023,10 +1028,10 @@ class CameraMultispectral:
if idade > float(max_age_s): if idade > float(max_age_s):
raise RuntimeError(f"cache RAW antigo: {idade:.2f}s") raise RuntimeError(f"cache RAW antigo: {idade:.2f}s")
raw_frame = { # Snapshot barato: arrays do pacote operacional são imutáveis
cam_id: arr.copy() # após a captura. Copiamos somente o dicionário de referências.
for cam_id, arr in self.ultimo_raw_multi.items() # 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 = dict(self._ultimo_resultado_raw)
resultado["cache_age_s"] = float(idade) resultado["cache_age_s"] = float(idade)
resultado["force_ignorado"] = bool(force) resultado["force_ignorado"] = bool(force)
@ -1785,7 +1790,12 @@ class CameraMultispectral:
def _ts_name(self) -> str: def _ts_name(self) -> str:
return datetime.now().strftime("%Y%m%d_%H%M%S_%f")[:-3] 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. 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: if max_age_s is None or float(max_age_s) <= 0:
max_age_s = max(1.0, self._cache_max_age_s) max_age_s = max(1.0, self._cache_max_age_s)
raw_frame, resultado = self.requisitar_frame_raw_multi( agora = time.perf_counter()
force=False,
max_age_s=max_age_s,
)
if raw_frame is None:
raise RuntimeError(resultado.get("erro") or "RAW cache indisponível")
with self._lock: 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. # Reforça o contrato científico no próprio stream_meta salvo.
raw_meta.setdefault("frame_type", "RAW_BRUTO") raw_meta.setdefault("frame_type", "RAW_BRUTO")
@ -1830,10 +1854,13 @@ class CameraMultispectral:
f"Presentes: {sorted(presentes)}" f"Presentes: {sorted(presentes)}"
) )
preview_bgr, preview_method = self._build_preview_raw_multispec( preview_bgr = None
raw_frame=raw_frame, preview_method = "disabled_runtime_save"
raw_meta=raw_meta, if include_preview:
) preview_bgr, preview_method = self._build_preview_raw_multispec(
raw_frame=raw_frame,
raw_meta=raw_meta,
)
resultado = dict(resultado) resultado = dict(resultado)
resultado.update({ resultado.update({
@ -1920,8 +1947,7 @@ class CameraMultispectral:
""" """
Salva pacote RAW_BRUTO multiespectral no mesmo espírito do capture de dataset. Salva pacote RAW_BRUTO multiespectral no mesmo espírito do capture de dataset.
Saída: Saída operacional:
<nome>.png
<nome>.json <nome>.json
<nome>_CAM_A.bin <nome>_CAM_A.bin
<nome>_CAM_B.bin <nome>_CAM_B.bin
@ -1939,6 +1965,7 @@ class CameraMultispectral:
bundle, resultado = self.requisitar_bundle_raw_multispec( bundle, resultado = self.requisitar_bundle_raw_multispec(
force=True, force=True,
max_age_s=0.0, max_age_s=0.0,
include_preview=False,
) )
if not resultado.get("frame_valido", False): if not resultado.get("frame_valido", False):
@ -1977,17 +2004,10 @@ class CameraMultispectral:
if not payload_files: if not payload_files:
raise RuntimeError("Nenhum payload RAW foi salvo.") 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 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 = { meta_save = {
"ts": datetime.now().isoformat(timespec="milliseconds"), "ts": datetime.now().isoformat(timespec="milliseconds"),
"source": "operacao_robo", "source": "operacao_robo",

View File

@ -45,6 +45,7 @@ class OakFcc3Client:
imu_modo="rotation_vector", imu_modo="rotation_vector",
imu_freq_hz=200, imu_freq_hz=200,
evaluate_quality=True, evaluate_quality=True,
require_product_contract=False,
hardware_sync_enabled=None, hardware_sync_enabled=None,
frame_sync_master=None, frame_sync_master=None,
@ -65,6 +66,30 @@ class OakFcc3Client:
self.module_calibration_json = module_calibration_json self.module_calibration_json = module_calibration_json
self.module_params = self._load_module_params(module_calibration_json) self.module_params = self._load_module_params(module_calibration_json)
self.fusion_config = self.module_params.get("fusion_config", {}) or {} 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_modo = str(imu_modo).strip().lower()
self.imu_freq_hz = int(imu_freq_hz) self.imu_freq_hz = int(imu_freq_hz)
@ -84,14 +109,15 @@ class OakFcc3Client:
self.svc = OakFcc3Service( self.svc = OakFcc3Service(
timeout=10, timeout=10,
fps=fps, fps=fps,
width=width, width=self.width,
height=height, height=self.height,
frame_type=frame_type, frame_type=frame_type,
output_dtype=output_dtype, output_dtype=output_dtype,
capture_mode=capture_mode, capture_mode=capture_mode,
raw_policy=raw_policy, raw_policy=raw_policy,
mx_id=self.mx_id, mx_id=self.mx_id,
module_calibration_json=module_calibration_json, module_calibration_json=module_calibration_json,
require_product_contract=self.require_product_contract,
imu_modo=self.imu_modo, imu_modo=self.imu_modo,
imu_freq_hz=self.imu_freq_hz, imu_freq_hz=self.imu_freq_hz,
@ -108,15 +134,15 @@ class OakFcc3Client:
self.applied_camera_controls = {} self.applied_camera_controls = {}
self.radiometric_controller = None self.radiometric_controller = None
self.core = RawProcessorCore( self.core = RawProcessorCore(
sensor_width=width, sensor_width=self.width,
sensor_height=height, sensor_height=self.height,
bayer_pattern=bayer, bayer_pattern=self.bayer,
calibration_json_path=module_calibration_json, calibration_json_path=module_calibration_json,
) )
self.preview = RawProcessorPreview( self.preview = RawProcessorPreview(
sensor_width=width, sensor_width=self.width,
sensor_height=height, sensor_height=self.height,
bayer_pattern=bayer, bayer_pattern=self.bayer,
) )
def __enter__(self): def __enter__(self):
@ -133,6 +159,46 @@ class OakFcc3Client:
with open(path, "r", encoding="utf-8") as f: with open(path, "r", encoding="utf-8") as f:
return json.load(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): def apply_module_camera_settings(self):
camera_settings = self.module_params.get("camera_settings", {}) or {} camera_settings = self.module_params.get("camera_settings", {}) or {}
@ -374,8 +440,8 @@ class OakFcc3Client:
evaluate_quality = self.evaluate_quality evaluate_quality = self.evaluate_quality
evaluate_quality = bool(evaluate_quality) evaluate_quality = bool(evaluate_quality)
tensor = self.core.fuse_multispec_cameras(decoded, meta, channels_expected) 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) #tensor = self.core.resize_tensor_chw(tensor, target_size=target_size)
# Mantém paridade com build_infer_tensor_from_stream: se a calibração # Mantém paridade com build_infer_tensor_from_stream: se a calibração
# habilitar patch normalization, ela também vale no caminho decoded. # habilitar patch normalization, ela também vale no caminho decoded.

View File

@ -53,14 +53,15 @@ class OakFcc3Manager:
mx_id=None, mx_id=None,
module_calibration_json=None, module_calibration_json=None,
module_params=None, module_params=None,
require_product_contract=False,
imu_modo="rotation_vector", imu_modo="rotation_vector",
imu_freq_hz=200, imu_freq_hz=200,
): ):
self.fps = fps self.fps = fps
# Para compatibilidade, mantemos width/height. # Para compatibilidade, mantemos width/height.
# No RAW_BRUTO isso não muda o sensor, pois usamos 800p fixo. # No RAW_BRUTO a resolução física é definida pelo module_params/hardware.
# No MULTISPEC isso representa a saída final alinhada da OAK. # No MULTISPEC width/height representa a saída final alinhada da OAK.
self.width = int(width) self.width = int(width)
self.height = int(height) self.height = int(height)
self.size = (self.width, self.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.fusion_config = (self.module_params or {}).get("fusion_config", {}) or {}
self.aligned_geometry = None 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 {}) 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 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") 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_enabled = True
self.async_capture_mode = "latest" # latest | queue self.async_capture_mode = "latest" # latest | queue
self.async_capture_max_queue = 2 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_thread = None
self._capture_stop_event = threading.Event() self._capture_stop_event = threading.Event()
self._capture_lock = threading.RLock() self._capture_lock = threading.RLock()
@ -183,6 +241,113 @@ class OakFcc3Manager:
with open(path, "r", encoding="utf-8") as f: with open(path, "r", encoding="utf-8") as f:
return json.load(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): def _default_controls_for_role(self, role: str):
role = str(role).lower() role = str(role).lower()
@ -293,11 +458,11 @@ class OakFcc3Manager:
""" """
Fluxo clássico. Fluxo clássico.
RGB/OV9782: RGB (OV9782/AR0234):
ColorCamera raw para RAW_BRUTO. ColorCamera raw na resolução nativa homologada pelo module_params.
MONO/OV9282: 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() sensor_name_u = str(sensor_name or "").upper()
role_u = str(role or "").lower() role_u = str(role or "").lower()
@ -311,14 +476,7 @@ class OakFcc3Manager:
if is_rgb: if is_rgb:
cam = self.pipeline.createColorCamera() cam = self.pipeline.createColorCamera()
cam.setBoardSocket(socket) cam.setBoardSocket(socket)
self._set_resolution_by_native_size(cam, role_u, is_color=True)
try:
cam.setResolution(dai.ColorCameraProperties.SensorResolution.THE_800_P)
except Exception:
try:
cam.setResolution(dai.ColorCameraProperties.SensorResolution.THE_1080_P)
except Exception:
pass
cam.setInterleaved(False) cam.setInterleaved(False)
cam.setColorOrder(dai.ColorCameraProperties.ColorOrder.RGB) cam.setColorOrder(dai.ColorCameraProperties.ColorOrder.RGB)
@ -335,17 +493,7 @@ class OakFcc3Manager:
mono = self.pipeline.create(dai.node.MonoCamera) mono = self.pipeline.create(dai.node.MonoCamera)
mono.setBoardSocket(socket) mono.setBoardSocket(socket)
self._set_resolution_by_native_size(mono, role_u, is_color=False)
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
mono.setFps(float(self.fps)) mono.setFps(float(self.fps))
@ -490,7 +638,7 @@ class OakFcc3Manager:
def _create_color_camera_multispec(self, socket): def _create_color_camera_multispec(self, socket):
cam = self.pipeline.create(dai.node.ColorCamera) cam = self.pipeline.create(dai.node.ColorCamera)
cam.setBoardSocket(socket) 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.setFps(float(self.fps))
cam.setInterleaved(False) cam.setInterleaved(False)
@ -500,10 +648,10 @@ class OakFcc3Manager:
cam.setVideoSize(int(self.sensor_width), int(self.sensor_height)) cam.setVideoSize(int(self.sensor_width), int(self.sensor_height))
return cam 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 = self.pipeline.create(dai.node.MonoCamera)
cam.setBoardSocket(socket) 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)) cam.setFps(float(self.fps))
return cam return cam
@ -625,7 +773,7 @@ class OakFcc3Manager:
bit_depth = 8 bit_depth = 8
raw_format = "BGR888p" raw_format = "BGR888p"
elif role == "re": elif role == "re":
cam = self._create_mono_camera_multispec(socket) cam = self._create_mono_camera_multispec(socket, "re")
src_output = cam.out src_output = cam.out
quad = self.aligned_geometry["quad_re"] quad = self.aligned_geometry["quad_re"]
out_type = dai.ImgFrame.Type.GRAY8 out_type = dai.ImgFrame.Type.GRAY8
@ -634,7 +782,7 @@ class OakFcc3Manager:
bit_depth = 8 bit_depth = 8
raw_format = "GRAY8" raw_format = "GRAY8"
elif role == "nir": elif role == "nir":
cam = self._create_mono_camera_multispec(socket) cam = self._create_mono_camera_multispec(socket, "nir")
src_output = cam.out src_output = cam.out
quad = self.aligned_geometry["quad_nir"] quad = self.aligned_geometry["quad_nir"]
out_type = dai.ImgFrame.Type.GRAY8 out_type = dai.ImgFrame.Type.GRAY8
@ -896,6 +1044,7 @@ class OakFcc3Manager:
self.pipeline = dai.Pipeline() self.pipeline = dai.Pipeline()
features = self.device.getConnectedCameraFeatures() features = self.device.getConnectedCameraFeatures()
self._validate_connected_hardware(features)
self.queues.clear() self.queues.clear()
self.buffers.clear() self.buffers.clear()
@ -943,11 +1092,14 @@ class OakFcc3Manager:
self.buffers[cam_id] = deque(maxlen=self.buffer_size) self.buffers[cam_id] = deque(maxlen=self.buffer_size)
self.control_queues[cam_id] = None self.control_queues[cam_id] = None
native_size = self._expected_native_size_for_role(role)
self.camera_info[cam_id] = { self.camera_info[cam_id] = {
"id": cam_id, "id": cam_id,
"socket": socket_name, "socket": socket_name,
"sensor": f.sensorName, "sensor": f.sensorName,
"role": role, "role": role,
"native_size": native_size,
"bayer_pattern": self.bayer_pattern if str(role).lower() == "rgb" else None,
} }
self._validate_capture_mode() self._validate_capture_mode()
@ -1143,6 +1295,14 @@ class OakFcc3Manager:
return { return {
"mx_id": self.mx_id, "mx_id": self.mx_id,
"backend": "oak_fcc3", "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), "running": bool(self.running),
"fps": self.fps, "fps": self.fps,
"width": self.width, "width": self.width,
@ -1201,9 +1361,10 @@ class OakFcc3Manager:
t0 = time.perf_counter() t0 = time.perf_counter()
deadline = t0 + float(timeout) deadline = t0 + float(timeout)
packet = None
with self._capture_cond: with self._capture_cond:
while time.perf_counter() < deadline: while time.perf_counter() < deadline:
packet = None
mode = str(getattr(self, "async_capture_mode", "latest")).lower() mode = str(getattr(self, "async_capture_mode", "latest")).lower()
@ -1218,22 +1379,7 @@ class OakFcc3Manager:
if packet is not None: if packet is not None:
seq = int(packet.get("seq", 0)) seq = int(packet.get("seq", 0))
self._last_consumed_packet_seq = seq self._last_consumed_packet_seq = seq
break
frames = packet["frames"]
meta = dict(packet["meta"])
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
cp = dict(meta.get("capture_perf", {}) or {})
cp["async_consumer"] = True
cp["async_packet_seq"] = seq
cp["async_packet_age_ms"] = float(age_ms)
cp["async_get_wait_ms"] = float(get_wait_ms)
cp["async_status"] = self.get_async_capture_status()
meta["capture_perf"] = cp
return frames, meta
remaining = deadline - time.perf_counter() remaining = deadline - time.perf_counter()
if remaining <= 0: if remaining <= 0:
@ -1241,6 +1387,76 @@ class OakFcc3Manager:
self._capture_cond.wait(timeout=min(0.005, remaining)) 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)
cp["async_get_wait_ms"] = float(get_wait_ms)
cp["async_status"] = self.get_async_capture_status()
meta["capture_perf"] = cp
return frames, meta
raise TimeoutError( raise TimeoutError(
f"Timeout aguardando pacote assíncrono do OAK-FFC-3. " f"Timeout aguardando pacote assíncrono do OAK-FFC-3. "
f"status={self.get_async_capture_status()}" f"status={self.get_async_capture_status()}"
@ -1456,47 +1672,36 @@ class OakFcc3Manager:
} }
else: else:
t0_data = self._cap_now_ms() if bool(getattr(self, "defer_raw_materialization", True)):
data = msg.getData() # Não toca no payload RAW aqui. O ImgFrame permanece vivo
if perf is not None: # no deque e só será materializado se entrar na tripleta
perf["drain_get_data_ms"] += self._cap_now_ms() - t0_data # escolhida pelo sincronizador. Frames descartados custam
# apenas metadados/timestamp.
t0_copy = self._cap_now_ms() frame = None
raw = np.frombuffer(data, dtype=np.uint8).copy() frame_controls = None
if perf is not None: deferred_raw = True
perf["drain_frombuffer_copy_ms"] += self._cap_now_ms() - t0_copy if perf is not None:
perf["deferred_raw_msgs"] += 1
t0_shape = self._cap_now_ms() else:
h = int(msg.getHeight()) frame, frame_controls = self._materialize_raw_msg(
w = int(msg.getWidth()) cam_id=cam_id,
stride = self._get_imgframe_stride(msg, raw.size, h, w) msg=msg,
perf=perf,
expected = h * stride metric_prefix="drain",
if raw.size < expected:
raise RuntimeError(
f"RAW menor que esperado: raw.size={raw.size}, esperado={expected}, "
f"w={w}, h={h}, stride={stride}"
) )
deferred_raw = False
frame = raw[:expected].reshape((h, stride)) 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: if perf is not None:
perf["drain_reshape_ms"] += self._cap_now_ms() - t0_shape perf["drain_controls_ms"] += self._cap_now_ms() - t0_ctrl
deferred_raw = False
self._last_raw_dims[cam_id] = {
"sensor_width": w,
"sensor_height": h,
"stride": stride,
"packed_width": stride,
}
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({ self.buffers[cam_id].append({
"frame": frame, "frame": frame,
"msg": msg if deferred_raw else None,
"deferred_raw": bool(deferred_raw),
"timestamp": ts_start, "timestamp": ts_start,
"timestamp_start": ts_start, "timestamp_start": ts_start,
"timestamp_end": ts_end, "timestamp_end": ts_end,
@ -1507,6 +1712,80 @@ class OakFcc3Manager:
perf["drained_total"] += 1 perf["drained_total"] += 1
perf["drained_by_cam"][cam_id] += 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()
data_ms = self._cap_now_ms() - t0_data
t0_copy = self._cap_now_ms()
raw = np.frombuffer(data, dtype=np.uint8).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:
raise RuntimeError(
f"RAW menor que esperado: raw.size={raw.size}, esperado={expected}, "
f"w={w}, h={h}, stride={stride}"
)
frame = raw[:expected].reshape((h, stride))
shape_ms = self._cap_now_ms() - t0_shape
self._last_raw_dims[cam_id] = {
"sensor_width": w,
"sensor_height": h,
"stride": stride,
"packed_width": stride,
}
t0_ctrl = self._cap_now_ms()
frame_controls = self._extract_frame_controls(msg)
controls_ms = self._cap_now_ms() - t0_ctrl
if perf is not None:
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: def _get_imgframe_stride(self, msg, raw_size: int, h: int, w: int) -> int:
try: try:
return int(msg.getStride()) return int(msg.getStride())
@ -1521,7 +1800,7 @@ class OakFcc3Manager:
return int(np.ceil(w * 5.0 / 4.0)) 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() required_cam_ids = self._get_required_cam_ids()
if not required_cam_ids: if not required_cam_ids:
@ -1549,11 +1828,6 @@ class OakFcc3Manager:
for cam_id, item in selected.items() 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()) ts_values = list(timestamps.values())
sync_dt_ms = ( sync_dt_ms = (
@ -1570,11 +1844,6 @@ class OakFcc3Manager:
for cam_id, ts in timestamps.items() 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_dt_ms"] = float(sync_dt_ms)
perf["sync_ok"] = bool(sync_ok) perf["sync_ok"] = bool(sync_ok)
@ -1601,10 +1870,41 @@ class OakFcc3Manager:
# Em best/best_effort, entrega a melhor combinação disponível, # Em best/best_effort, entrega a melhor combinação disponível,
# mesmo quando estiver fora da tolerância. # mesmo quando estiver fora da tolerância.
frames = { #
cam_id: item["frame"] # No caminho assíncrono/latest, materialize=False mantém apenas as
for cam_id, item in selected.items() # 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. # Consome todos os frames anteriores e o próprio frame selecionado.
for cam_id, used_item in selected.items(): for cam_id, used_item in selected.items():
@ -1622,7 +1922,7 @@ class OakFcc3Manager:
) )
return ( return (
frames, payload,
timestamps, timestamps,
sync_dt_ms, sync_dt_ms,
sync_ok, sync_ok,
@ -1669,7 +1969,15 @@ class OakFcc3Manager:
# Meta # 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()) payload_sources = list(frames.keys())
shapes = { shapes = {
@ -1730,7 +2038,11 @@ class OakFcc3Manager:
camera_info[cam_id] = item camera_info[cam_id] = item
meta = { 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", "backend": "oak_fcc3",
"frame_type": self.frame_type, "frame_type": self.frame_type,
"capture_mode": self.capture_mode, "capture_mode": self.capture_mode,
@ -1974,6 +2286,17 @@ class OakFcc3Manager:
"drain_frombuffer_copy_ms": 0.0, "drain_frombuffer_copy_ms": 0.0,
"drain_reshape_ms": 0.0, "drain_reshape_ms": 0.0,
"drain_controls_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, "sync_select_ms": 0.0,
"meta_ms": 0.0, "meta_ms": 0.0,
"sleep_ms": 0.0, "sleep_ms": 0.0,
@ -2093,32 +2416,64 @@ class OakFcc3Manager:
# Importante: esta thread é a única que mexe nas queues/buffers. # Importante: esta thread é a única que mexe nas queues/buffers.
self._drain_queues_to_buffers(perf=perf) 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: if synced is None:
# Dorme curto. Pode testar 0.0005 se quiser reduzir latência. # Dorme curto. Pode testar 0.0005 se quiser reduzir latência.
time.sleep(0.001) time.sleep(0.001)
continue 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 self.frame_id += 1
meta = self._build_meta( packet_frame_id = int(self.frame_id)
frames,
timestamps, meta = None
sync_dt_ms, if not defer_packet_to_consumer:
sync_ok, meta = self._build_meta(
frame_controls=frame_controls, payload,
) timestamps,
sync_dt_ms,
sync_ok,
frame_controls=frame_controls,
frame_id_override=packet_frame_id,
)
now = time.perf_counter() now = time.perf_counter()
packet = { packet = {
"seq": int(self._latest_packet_seq + 1), "seq": int(self._latest_packet_seq + 1),
"created_perf_counter": float(now), "created_perf_counter": float(now),
"frames": frames, "frame_id": packet_frame_id,
"meta": meta, "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. # Adiciona perf de captura assíncrona no meta.
if perf is not None: if perf is not None:
perf["async_thread"] = True perf["async_thread"] = True
@ -2126,7 +2481,10 @@ class OakFcc3Manager:
perf["sync_dt_ms"] = float(sync_dt_ms) perf["sync_dt_ms"] = float(sync_dt_ms)
perf["sync_ok"] = bool(sync_ok) perf["sync_ok"] = bool(sync_ok)
perf["wait_reason"] = "async_packet_ready" perf["wait_reason"] = "async_packet_ready"
meta["capture_perf"] = perf if meta is not None:
meta["capture_perf"] = perf
else:
packet["capture_perf"] = perf
with self._capture_cond: with self._capture_cond:
self._latest_packet_seq += 1 self._latest_packet_seq += 1
@ -2195,6 +2553,7 @@ class OakFcc3Manager:
st.update({ st.update({
"enabled": bool(getattr(self, "async_capture_enabled", True)), "enabled": bool(getattr(self, "async_capture_enabled", True)),
"mode": str(getattr(self, "async_capture_mode", "latest")), "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()), "thread_alive": bool(self._capture_thread is not None and self._capture_thread.is_alive()),
"latest_seq": int(getattr(self, "_latest_packet_seq", 0)), "latest_seq": int(getattr(self, "_latest_packet_seq", 0)),
"last_consumed_seq": int(getattr(self, "_last_consumed_packet_seq", 0)), "last_consumed_seq": int(getattr(self, "_last_consumed_packet_seq", 0)),

View File

@ -1,6 +1,8 @@
import base64
import datetime import datetime
import json import json
import os import os
import queue
import threading import threading
import time import time
from pathlib import Path from pathlib import Path
@ -26,7 +28,7 @@ from visual_worker.utils import converter_valores_numpy
from shared.perf_monitor import VisualPerfMonitor 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_ASSEMBLY_SCHEMA = "multispec_module_params_assembly_v1"
PRODUCT_MODULE_SCHEMA = "multispec_module_params_v3" PRODUCT_MODULE_SCHEMA = "multispec_module_params_v3"
@ -87,6 +89,8 @@ class CameraManager:
self._loop_deteccao_iniciado = False self._loop_deteccao_iniciado = False
self._loop_analise_iniciado = False self._loop_analise_iniciado = False
self._loop_stream_iniciado = False self._loop_stream_iniciado = False
self._loop_preview_iniciado = False
self._loop_replay_writer_iniciado = False
self._loop_publicacao_iniciado = False self._loop_publicacao_iniciado = False
self._vida_lock = threading.RLock() self._vida_lock = threading.RLock()
@ -101,6 +105,25 @@ class CameraManager:
self._pub_lock = threading.RLock() self._pub_lock = threading.RLock()
self._preview_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.perf = VisualPerfMonitor(janela=180)
self._ultimo_posproc_save_ts = 0.0 self._ultimo_posproc_save_ts = 0.0
@ -164,6 +187,8 @@ class CameraManager:
self._ultimo_preview_seg_ts = 0.0 self._ultimo_preview_seg_ts = 0.0
self._ultimo_preview_overlay_ts = 0.0 self._ultimo_preview_overlay_ts = 0.0
self._ultimo_preview_debug_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_last_ts = None
self._fps_infer_ema = 0.0 self._fps_infer_ema = 0.0
@ -172,6 +197,16 @@ class CameraManager:
self._ultimo_loop_analise_fps = 0.0 self._ultimo_loop_analise_fps = 0.0
self._ultimo_config_update_ts = 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: # Diagnostico fino do pipeline. Estes campos sao somente observabilidade:
# nao alteram frequencias, caches, gates ou comandos dos bicos. # nao alteram frequencias, caches, gates ou comandos dos bicos.
@ -352,9 +387,19 @@ class CameraManager:
) )
fusion = mp.get("fusion_config", {}) or {} fusion = mp.get("fusion_config", {}) or {}
mp_target = self._normalizar_size_wh( mp_target_raw = fusion.get("target_size")
fusion.get("target_size"),
"module_params.fusion_config.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 {} sizes_raw = mp.get("sensor_size_by_role", {}) or {}
@ -408,11 +453,9 @@ class CameraManager:
""" """
Fecha a autoridade da resolução antes de abrir a OAK: Fecha a autoridade da resolução antes de abrir a OAK:
ONNX input [W,H] ONNX input [W,H] = autoridade final do runtime
== module_params.fusion_config.target_size = opcional/legado
module_params.fusion_config.target_size seg_config.ia_resolution = opcional, mas deve casar com ONNX
==
seg_config.ia_resolution
camera_width/camera_height deixam de participar da decisão física. camera_width/camera_height deixam de participar da decisão física.
""" """
@ -422,9 +465,16 @@ class CameraManager:
mx_id, 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( raise RuntimeError(
"Contrato de resolução incompatível entre ONNX e module_params: " "Contrato de resolução incompatível entre ONNX e module_params: "
f"ONNX={model_target} module_params={mp_target}" f"ONNX={model_target} module_params={mp_target}"
@ -439,7 +489,7 @@ class CameraManager:
if cfg_target != model_target: if cfg_target != model_target:
raise RuntimeError( 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}" f"config={cfg_target} ONNX={model_target} MP={mp_target}"
) )
else: else:
@ -763,7 +813,8 @@ class CameraManager:
self.mx_id = str(mx_id) self.mx_id = str(mx_id)
from weed_worker.config import load_seg_config 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.posproc_intervalo_min_s = float(
self.seg_config.get("posproc_intervalo_min_s", 5.0) self.seg_config.get("posproc_intervalo_min_s", 5.0)
) )
@ -1067,7 +1118,7 @@ class CameraManager:
if self.weed_detector is not None: if self.weed_detector is not None:
return return
self.weed_detector = WeedDetector() self.weed_detector = WeedDetector(config=self.seg_config)
def _iniciar_loops_se_necessario(self, seg_config): def _iniciar_loops_se_necessario(self, seg_config):
freq_analise = float(seg_config.get("analise_fps", 15.0)) 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_inferencia = float(seg_config.get("inferencia_fps", freq_tensor))
freq_deteccao = float(seg_config.get("deteccao_fps", freq_inferencia)) 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: if not self._loop_stream_iniciado:
stream = getattr(self.camera, "stream", None) stream = getattr(self.camera, "stream", None)
stream_fps = float(getattr(stream, "_op_fps", 2.0) or 2.0) 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: if detector is None or reset_nome is None:
# Fallback seguro: recria apenas o detector, nunca câmera/modelo. # 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"): if self.seg_config and hasattr(self.weed_detector, "atualizar_config"):
self.weed_detector.atualizar_config(self.seg_config) self.weed_detector.atualizar_config(self.seg_config)
reset_nome = "recriado" reset_nome = "recriado"
@ -1540,6 +1605,63 @@ class CameraManager:
except Exception as e: except Exception as e:
self.mostrar_log(f"[saude] erro: {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): def _atualizar_config_dinamica_detector(self, intervalo_s=0.20):
""" """
Atualiza config do detector e a cabeça ONNX em baixa frequência. Atualiza config do detector e a cabeça ONNX em baixa frequência.
@ -1560,7 +1682,10 @@ class CameraManager:
try: try:
from weed_worker.config import load_seg_config 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) cfg = self._validar_config_dinamica_contrato(cfg)
self.seg_config = cfg self.seg_config = cfg
self.qtd_bicos = int(cfg.get("qtd_bicos", self.qtd_bicos or 7) or 7) 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 return None, None, None
def get_selected_frame(self, _frame_type: TipoFrameCamera): def get_selected_frame(self, _frame_type: TipoFrameCamera):
agora = time.time() """
Compatibilidade externa: retorna SOMENTE o último preview cacheado.
try: Regra de produto:
if _frame_type not in self.FRAME_TYPES_PREVIEW: - nunca chama preview_infer_cached();
return None - nunca monta overlay/debug sob demanda;
- nunca força inferência;
- nunca toca no cursor da câmera.
# Se ainda não tem runtime/tensor, tenta devolver o último preview válido. O produtor de preview é _iniciar_loop_preview_cache().
if self.model_svc is None or self._ultimo_raw_input is None: """
return self._get_cached_preview_frame(_frame_type) if _frame_type not in self.FRAME_TYPES_PREVIEW:
return None
return self._get_cached_preview_frame(_frame_type)
# Janela curta: usa cache recente. def _get_cached_preview_frame_ref(self, frame_type):
if (agora - self._ultimo_preview_ts) < 0.20: """Retorna referência imutável do cache visual, sem cópia."""
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:
return None
def _get_cached_preview_frame(self, frame_type):
with self._preview_lock: with self._preview_lock:
if frame_type == TipoFrameCamera.Rgb: if frame_type == TipoFrameCamera.Rgb:
frame = self._ultimo_preview_rgb return self._ultimo_preview_rgb
elif frame_type == TipoFrameCamera.Segmentacao: if frame_type == TipoFrameCamera.Segmentacao:
frame = self._ultimo_preview_seg return self._ultimo_preview_seg
elif frame_type == TipoFrameCamera.Overlay: if frame_type == TipoFrameCamera.Overlay:
frame = self._ultimo_preview_overlay return self._ultimo_preview_overlay
elif frame_type == TipoFrameCamera.Debug: if frame_type == TipoFrameCamera.Debug:
frame = self._ultimo_preview_debug return self._ultimo_preview_debug
else: return None
frame = None
return self._copiar_frame(frame) def _get_cached_preview_frame(self, frame_type):
# 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))
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): def get_debug_frame(self, mostrar=False, overlay_bgr=None):
overlay = overlay_bgr if overlay_bgr is not None else self._ultimo_preview_overlay 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, "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, overlay_bgr=overlay,
atuacao_bicos=self._ultimo_controle or {}, atuacao_bicos=self._ultimo_controle or {},
config=self.seg_config or {}, config=self.seg_config or {},
@ -1862,8 +2067,6 @@ class CameraManager:
mostrar=mostrar, mostrar=mostrar,
) )
return self._ultimo_preview_debug
def _montar_debug_overlay( def _montar_debug_overlay(
self, self,
overlay_bgr, overlay_bgr,
@ -2012,42 +2215,151 @@ class CameraManager:
return frame return frame
def _set_cached_preview_frame(self, frame_type, frame, ts=None): def _set_cached_preview_frame(self, frame_type, frame, ts=None):
if not self._frame_valido(frame): """Compatibilidade interna; publica um único frame no cache leve."""
return False ts = time.time() if ts is None else float(ts)
return self._cache_preview_bundle({frame_type: frame}, ts)
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
# ============================================================ # ============================================================
# Loops # 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 _iniciar_loop_captura_tensor(self, freq=25.0):
def loop(): def loop():
periodo = 1.0 / max(float(freq), 0.1) periodo = 1.0 / max(float(freq), 0.1)
@ -2273,6 +2585,7 @@ class CameraManager:
"ts": pred_ts, "ts": pred_ts,
"tensor_ts": tensor_ts, "tensor_ts": tensor_ts,
"predictions": predictions, "predictions": predictions,
"raw_input": tensor5,
"res": res, "res": res,
"infer_ms": float(infer_ms), "infer_ms": float(infer_ms),
"infer_gpu_ms": float(infer_forward_ms or 0.0), "infer_gpu_ms": float(infer_forward_ms or 0.0),
@ -2368,6 +2681,9 @@ class CameraManager:
t_cfg0 = time.perf_counter() t_cfg0 = time.perf_counter()
cfg_cpu0 = time.thread_time() 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) self._atualizar_config_dinamica_detector(intervalo_s=0.20)
t_cfg1 = time.perf_counter() t_cfg1 = time.perf_counter()
config_ms = (t_cfg1 - t_cfg0) * 1000.0 config_ms = (t_cfg1 - t_cfg0) * 1000.0
@ -2419,10 +2735,6 @@ class CameraManager:
analise = analise_completa.get("dados_visuais", {}) analise = analise_completa.get("dados_visuais", {})
detector_ms = (t_det1 - t_det0) * 1000.0 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 # Gate de pulverização
# ==================================================== # ====================================================
@ -2435,6 +2747,12 @@ class CameraManager:
analise["pulverizacao"] = debug_pulverizacao analise["pulverizacao"] = debug_pulverizacao
t_ctrl1 = time.perf_counter() 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. # Publicação atômica em relação a fechar/substituir câmera.
# Se a geração mudou, nenhum dado antigo chega aos caches/bicos. # Se a geração mudou, nenhum dado antigo chega aos caches/bicos.
with self._vida_lock: with self._vida_lock:
@ -2668,11 +2986,15 @@ class CameraManager:
frame_type = TipoFrameCamera( frame_type = TipoFrameCamera(
camera_ctx.get("frame_type", TipoFrameCamera.Rgb.value) 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) self.camera.enviar_frame_tcp(frame)
if self.debug_visual: 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: except Exception as e:
self.mostrar_log(f"Erro no loop de stream: {e}") self.mostrar_log(f"Erro no loop de stream: {e}")
@ -2923,9 +3245,23 @@ class CameraManager:
} }
try: try:
operacao = ContextoGlobalRedis.get_operacao() snapshot = self._runtime_ctx_snapshot or {}
controle = ContextoGlobalRedis.get_controle() snapshot_ts = float(snapshot.get("ts", 0.0) or 0.0)
contexto = ContextoGlobalRedis.get_contexto() 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 {}) traj = (contexto.get("Trajetoria", {}) or {})
gerais = (contexto.get("Gerais", {}) or {}) gerais = (contexto.get("Gerais", {}) or {})
@ -3232,6 +3568,12 @@ class CameraManager:
} }
resumo["pub_debug"] = getattr(self, "_ultimo_pub_debug", {}) 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) self._set_pub_cache("performance_weed", resumo)
@ -3248,6 +3590,7 @@ class CameraManager:
det = loops.get("deteccao", {}) det = loops.get("deteccao", {})
pub = loops.get("publicacao", {}) pub = loops.get("publicacao", {})
stream = loops.get("stream", {}) stream = loops.get("stream", {})
preview = loops.get("preview_cache", {})
pub_dbg = getattr(self, "_ultimo_pub_debug", {}) pub_dbg = getattr(self, "_ultimo_pub_debug", {})
@ -3257,7 +3600,8 @@ class CameraManager:
f"inf={self._fps(inf):.1f} " f"inf={self._fps(inf):.1f} "
f"det={self._fps(det):.1f} " f"det={self._fps(det):.1f} "
f"pub={self._fps(pub):.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"period inf={self._fmt(self._per(inf))}ms "
f"det={self._fmt(self._per(det))}ms " f"det={self._fmt(self._per(det))}ms "
f"tensor={self._fmt(self._per(tensor))}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"redis={self._fmt(self._m(pub, 'redis_ms'))} "
f"debug_pub={pub_dbg.get('publicou', 0)} " f"debug_pub={pub_dbg.get('publicou', 0)} "
f"campos={pub_dbg.get('campos', 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( self.mostrar_log(
@ -3468,38 +3814,26 @@ class CameraManager:
continue 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( self.mostrar_log(
f"⚠️ Frame não salvo | " f"⚠️ Frame não salvo | "
f"tipo={tipo_nome} " f"tipo={tipo_nome} "
f"nome={nome_frame} " f"nome={nome_frame} "
f"motivo=sem_frame_valido_em_cache" f"motivo=sem_jpeg_cacheado"
) )
continue continue
frame = self._copiar_frame(frame)
caminho = os.path.join(pasta, f"{nome_frame}.jpg") caminho = os.path.join(pasta, f"{nome_frame}.jpg")
ok = self._enfileirar_replay_jpeg(caminho, jpeg_bytes)
# Replay leve: JPEG comprimido. if not ok:
# 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):
self.mostrar_log( self.mostrar_log(
f"❌ Falha ao salvar frame | " f"⚠️ Replay descartado para proteger runtime | "
f"tipo={tipo_nome} " f"tipo={tipo_nome} nome={nome_frame}"
f"nome={nome_frame} "
f"caminho={caminho}"
) )
continue continue

View File

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

View File

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

View File

@ -42,10 +42,15 @@ class WeedDetector:
CONTRATO_OFICIAL = "target_binary" CONTRATO_OFICIAL = "target_binary"
TARGET_ID = 1 TARGET_ID = 1
def __init__(self): def __init__(self, config: Optional[dict] = None):
from weed_worker.config import load_seg_config # 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.prediction_contract = str(
self.config.get("prediction_contract", self.CONTRATO_OFICIAL) self.config.get("prediction_contract", self.CONTRATO_OFICIAL)
@ -120,6 +125,12 @@ class WeedDetector:
self._morf_cache_k: Optional[int] = None self._morf_cache_k: Optional[int] = None
self._morf_kernel = 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. # Tempo/movimento.
self._last_update_ts: Optional[float] = None self._last_update_ts: Optional[float] = None
self._distancia_deslocada_total_cm = 0.0 self._distancia_deslocada_total_cm = 0.0
@ -305,11 +316,19 @@ class WeedDetector:
# mas aqui tratamos como velocidade em m/s. # mas aqui tratamos como velocidade em m/s.
vel_norm = vel_mps 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() t_rad0 = time.perf_counter()
target_no_radar_frame, estat_target = self._decidir_target_no_radar( target_no_radar_frame, estat_target = self._decidir_target_no_radar(
predictions, predictions,
cfg, cfg,
vel_norm=vel_norm, vel_norm=vel_norm,
mask_target=mask_target,
) )
t_rad1 = time.perf_counter() t_rad1 = time.perf_counter()
@ -319,6 +338,7 @@ class WeedDetector:
cfg=cfg, cfg=cfg,
vel_mps=vel_mps, vel_mps=vel_mps,
agora=agora, agora=agora,
mask_target=mask_target,
) )
t_bic1 = time.perf_counter() t_bic1 = time.perf_counter()
@ -377,6 +397,7 @@ class WeedDetector:
predictions: np.ndarray, predictions: np.ndarray,
cfg: dict, cfg: dict,
vel_norm: float = 0.0, vel_norm: float = 0.0,
mask_target: Optional[np.ndarray] = None,
): ):
on_global = float(cfg.get("min_frac_erva_global_on", 0.0020)) on_global = float(cfg.get("min_frac_erva_global_on", 0.0020))
off_global = float(cfg.get("min_frac_erva_global_off", 0.0015)) off_global = float(cfg.get("min_frac_erva_global_off", 0.0015))
@ -393,7 +414,12 @@ class WeedDetector:
h, w = predictions.shape[:2] h, w = predictions.shape[:2]
inv_total = 1.0 / max(1, h * w) 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 frac_global = target_sum * inv_total
alpha = float(cfg.get("erva_frac_ema", 0.30)) alpha = float(cfg.get("erva_frac_ema", 0.30))
@ -430,6 +456,7 @@ class WeedDetector:
cfg: dict, cfg: dict,
vel_mps: float, vel_mps: float,
agora: float, agora: float,
mask_target: Optional[np.ndarray] = None,
): ):
qtd_bicos = self.qtd_bicos qtd_bicos = self.qtd_bicos
@ -460,6 +487,7 @@ class WeedDetector:
frac_grid, grid_info = self._calcular_frac_grid_por_bico_cell( frac_grid, grid_info = self._calcular_frac_grid_por_bico_cell(
predictions=predictions, predictions=predictions,
cfg=cfg, cfg=cfg,
mask_target=mask_target,
) )
self._frac_grid_last = frac_grid self._frac_grid_last = frac_grid
@ -657,69 +685,122 @@ class WeedDetector:
return shifted 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( def _calcular_frac_grid_por_bico_cell(
self, self,
predictions: np.ndarray, predictions: np.ndarray,
cfg: dict, cfg: dict,
mask_target: Optional[np.ndarray] = None,
): ):
h, w = predictions.shape[:2] 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 = self._aplicar_filtros_opcionais(
mask_target=mask_target, mask_target=mask_target,
cfg=cfg, cfg=cfg,
) )
# Integral image em inteiro com sinal. geom = self._get_grid_geometry(h, w)
# Evita overflow nos cálculos A - B - C + D. x_edges = geom["x_edges"]
ii = np.pad( y_edges = geom["y_edges"]
mask_target.astype(np.int64, copy=False) areas = geom["areas"]
.cumsum(axis=0, dtype=np.int64)
.cumsum(axis=1, dtype=np.int64),
((1, 0), (1, 0)),
mode="constant",
constant_values=0,
)
x_edges = np.linspace(0, w, self.qtd_bicos + 1, dtype=np.int32) if geom["reduceat_safe"]:
y_edges = np.linspace(0, h, self.num_cells + 1, dtype=np.int32) # 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_grid = np.zeros( frac_cells_bicos = np.divide(
(self.qtd_bicos, self.num_cells), counts_yx,
dtype=np.float32, areas,
) out=np.zeros_like(areas, dtype=np.float32),
where=areas > 0,
).astype(np.float32, copy=False)
valid_cell_count = 0 frac_grid = np.ascontiguousarray(
frac_cells_bicos.T,
for c in range(self.num_cells): dtype=np.float32,
y_top = int(y_edges[c]) )
y_bot = int(y_edges[c + 1]) else:
# Fallback para geometrias patológicas em que há bins vazios
if y_bot <= y_top: # (ex.: mais células que pixels). Preserva a semântica antiga.
continue ii = np.pad(
mask_target.astype(np.int64, copy=False)
valid_cell_count += 1 .cumsum(axis=0, dtype=np.int64)
cell_h = y_bot - y_top .cumsum(axis=1, dtype=np.int64),
((1, 0), (1, 0)),
for b in range(self.qtd_bicos): mode="constant",
x0 = int(x_edges[b]) constant_values=0,
x1 = int(x_edges[b + 1]) )
frac_grid = np.zeros(
if x1 <= x0: (self.qtd_bicos, self.num_cells),
dtype=np.float32,
)
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 continue
for b in range(self.qtd_bicos):
area = float(cell_h * (x1 - x0)) x0 = int(x_edges[b])
if area <= 0: x1 = int(x_edges[b + 1])
continue area = int((y_bot - y_top) * (x1 - x0))
if area <= 0:
total = ( continue
ii[y_bot, x1] total = (
- ii[y_top, x1] ii[y_bot, x1]
- ii[y_bot, x0] - ii[y_top, x1]
+ ii[y_top, x0] - 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. # A segmentação é naturalmente da esquerda para a direita.
# Quando habilitado, inverte associação região imagem -> bico físico. # Quando habilitado, inverte associação região imagem -> bico físico.
@ -727,12 +808,17 @@ class WeedDetector:
frac_grid = frac_grid[::-1, :].copy() frac_grid = frac_grid[::-1, :].copy()
return frac_grid, { return frac_grid, {
"valid_cell_count": int(valid_cell_count), "valid_cell_count": int(geom["valid_cell_count"]),
"height": int(h), "height": int(h),
"width": int(w), "width": int(w),
"x_edges": x_edges, "x_edges": x_edges,
"y_edges": y_edges, "y_edges": y_edges,
"zona_coord": "frac_top_to_bottom", "zona_coord": "frac_top_to_bottom",
"grid_backend": (
"reduceat_segmented"
if geom["reduceat_safe"]
else "integral_fallback"
),
} }
def _atualizar_score_memoria( def _atualizar_score_memoria(

View File

@ -129,7 +129,7 @@ namespace OperationControl.Services
public int NtripPort { get; set; } = 2101; public int NtripPort { get; set; } = 2101;
public string NtripMountpoint { get; set; } = "EESC0"; public string NtripMountpoint { get; set; } = "EESC0";
public string NtripUsername { get; set; } = "Zendion"; // Environment.GetEnvironmentVariable("AGRO_NTRIP_USERNAME") ?? string.Empty; public string NtripUsername { get; set; } = "Zendion"; // Environment.GetEnvironmentVariable("AGRO_NTRIP_USERNAME") ?? string.Empty;
public string NtripPassword { get; set; } = "c3pc7*9N"; //Environment.GetEnvironmentVariable("AGRO_NTRIP_PASSWORD") ?? string.Empty; public string NtripPassword { get; set; } = "kY3zd$*5"; //Environment.GetEnvironmentVariable("AGRO_NTRIP_PASSWORD") ?? string.Empty;
private GnssExpectedRole _expectedRole = GnssExpectedRole.PreserveCurrentConfiguration; private GnssExpectedRole _expectedRole = GnssExpectedRole.PreserveCurrentConfiguration;
private BaseFixedConfiguration _lastBaseConfiguration; private BaseFixedConfiguration _lastBaseConfiguration;

View File

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

View File

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

View File

@ -45,6 +45,7 @@ class OakFcc3Client:
imu_modo="rotation_vector", imu_modo="rotation_vector",
imu_freq_hz=200, imu_freq_hz=200,
evaluate_quality=True, evaluate_quality=True,
require_product_contract=False,
hardware_sync_enabled=None, hardware_sync_enabled=None,
frame_sync_master=None, frame_sync_master=None,
@ -65,6 +66,30 @@ class OakFcc3Client:
self.module_calibration_json = module_calibration_json self.module_calibration_json = module_calibration_json
self.module_params = self._load_module_params(module_calibration_json) self.module_params = self._load_module_params(module_calibration_json)
self.fusion_config = self.module_params.get("fusion_config", {}) or {} 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_modo = str(imu_modo).strip().lower()
self.imu_freq_hz = int(imu_freq_hz) self.imu_freq_hz = int(imu_freq_hz)
@ -84,14 +109,15 @@ class OakFcc3Client:
self.svc = OakFcc3Service( self.svc = OakFcc3Service(
timeout=10, timeout=10,
fps=fps, fps=fps,
width=width, width=self.width,
height=height, height=self.height,
frame_type=frame_type, frame_type=frame_type,
output_dtype=output_dtype, output_dtype=output_dtype,
capture_mode=capture_mode, capture_mode=capture_mode,
raw_policy=raw_policy, raw_policy=raw_policy,
mx_id=self.mx_id, mx_id=self.mx_id,
module_calibration_json=module_calibration_json, module_calibration_json=module_calibration_json,
require_product_contract=self.require_product_contract,
imu_modo=self.imu_modo, imu_modo=self.imu_modo,
imu_freq_hz=self.imu_freq_hz, imu_freq_hz=self.imu_freq_hz,
@ -108,15 +134,15 @@ class OakFcc3Client:
self.applied_camera_controls = {} self.applied_camera_controls = {}
self.radiometric_controller = None self.radiometric_controller = None
self.core = RawProcessorCore( self.core = RawProcessorCore(
sensor_width=width, sensor_width=self.width,
sensor_height=height, sensor_height=self.height,
bayer_pattern=bayer, bayer_pattern=self.bayer,
calibration_json_path=module_calibration_json, calibration_json_path=module_calibration_json,
) )
self.preview = RawProcessorPreview( self.preview = RawProcessorPreview(
sensor_width=width, sensor_width=self.width,
sensor_height=height, sensor_height=self.height,
bayer_pattern=bayer, bayer_pattern=self.bayer,
) )
def __enter__(self): def __enter__(self):
@ -133,6 +159,46 @@ class OakFcc3Client:
with open(path, "r", encoding="utf-8") as f: with open(path, "r", encoding="utf-8") as f:
return json.load(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): def apply_module_camera_settings(self):
camera_settings = self.module_params.get("camera_settings", {}) or {} camera_settings = self.module_params.get("camera_settings", {}) or {}
@ -374,8 +440,8 @@ class OakFcc3Client:
evaluate_quality = self.evaluate_quality evaluate_quality = self.evaluate_quality
evaluate_quality = bool(evaluate_quality) evaluate_quality = bool(evaluate_quality)
tensor = self.core.fuse_multispec_cameras(decoded, meta, channels_expected) 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) #tensor = self.core.resize_tensor_chw(tensor, target_size=target_size)
# Mantém paridade com build_infer_tensor_from_stream: se a calibração # Mantém paridade com build_infer_tensor_from_stream: se a calibração
# habilitar patch normalization, ela também vale no caminho decoded. # 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