diff --git a/listener/app.py b/listener/app.py index 863b5402..46f6b8fa 100644 --- a/listener/app.py +++ b/listener/app.py @@ -14,7 +14,7 @@ import datetime import traceback import serial # type: ignore -from sqlalchemy import create_engine, Column, Integer, String, Boolean, Text # type: ignore +from sqlalchemy import create_engine, Column, Integer, String, Boolean, Text, func # type: ignore from sqlalchemy import DateTime as SqDateTime # type: ignore from sqlalchemy import Date as SqDate # type: ignore from sqlalchemy.orm import declarative_base, sessionmaker # type: ignore @@ -55,10 +55,6 @@ ERROR_LOG_KEYWORDS = ( # Global Variables active_genexpert_connections = {} connection_lock = threading.Lock() -scheduled_result_queries = {} -scheduled_result_query_lock = threading.Lock() -genexpert_query_inflight_by_ip = {} -genexpert_query_inflight_lock = threading.Lock() # Network Configuration TCP_LISTENER_PORT = 6001 # PC GeneXpert set ke mode Client, konek ke IP:PORT ini SERVER_HOST = '0.0.0.0' # Listen di semua interface @@ -339,6 +335,56 @@ def get_pending_orders(ip_addr): finally: session.close() +def mark_genexpert_order_flag(rnoreg, ip_addr, reason="processed"): + rnoreg = str(rnoreg or "").strip() + flag_name = get_flag_by_device(str(ip_addr or "").strip()) + if not rnoreg: + print(f"[GENEXPERT-DB] Skip update flag, rnoreg kosong. reason={reason}") + return False + if not flag_name: + print(f"[GENEXPERT-DB] Skip update flag rnoreg={rnoreg}, IP {ip_addr} tidak terdaftar.") + return False + + flag_attr = getattr(PaslabOrder, flag_name, None) + if flag_attr is None: + print(f"[GENEXPERT-DB] Skip update flag rnoreg={rnoreg}, kolom {flag_name} tidak ada.") + return False + + with SessionLocal() as session: + try: + match_filter = func.trim(PaslabOrder.rnoreg) == rnoreg + matched_count = session.query(PaslabOrder).filter(match_filter).count() + if matched_count == 0: + print(f"[GENEXPERT-DB] Order {rnoreg} tidak ditemukan saat update {flag_name}. reason={reason}") + return False + + updated_count = session.query(PaslabOrder).filter(match_filter).update( + { + flag_attr: True, + PaslabOrder.updated_at: datetime.datetime.now(), + }, + synchronize_session=False, + ) + session.commit() + + true_count = session.query(PaslabOrder).filter(match_filter, flag_attr == True).count() + if true_count == matched_count: + print( + f"[GENEXPERT-DB] Order {rnoreg} set {flag_name}=TRUE. " + f"rows={updated_count}/{matched_count}, reason={reason}" + ) + return True + + print( + f"[GENEXPERT-DB-WARN] Update {flag_name} belum terverifikasi rnoreg={rnoreg}. " + f"true_rows={true_count}/{matched_count}, updated_rows={updated_count}, reason={reason}" + ) + return False + except Exception as exc: + session.rollback() + print(f"[GENEXPERT-DB-ERROR] Gagal update {flag_name} rnoreg={rnoreg}, reason={reason}, error={exc}") + return False + def parse_hl7_segments(hl7_message): return [segment for segment in str(hl7_message or "").split('\r') if segment] @@ -588,9 +634,12 @@ def send_all_orders_astm(conn, ip_addr, astm_msg, response_framing="astm"): return reply = create_genexpert_astm_order_message(selected_orders, ip_addr=ip_addr, query_tag=query_tag) - send_genexpert_response(conn, ip_addr, reply, response_framing, label=f"astm-q-order:{selected_orders[0].rnoreg}") + sent_ok = send_genexpert_response(conn, ip_addr, reply, response_framing, label=f"astm-q-order:{selected_orders[0].rnoreg}") for order in selected_orders: - print(f"[GENEXPERT] Order ASTM ditawarkan ke {ip_addr}: {order.rnoreg}") + rnoreg = str(order.rnoreg or "").strip() + print(f"[GENEXPERT] Order ASTM ditawarkan ke {ip_addr}: {rnoreg}, sent_ok={sent_ok}") + if sent_ok: + mark_genexpert_order_flag(rnoreg, ip_addr, reason="astm-order-sent") finally: session.close() @@ -636,27 +685,7 @@ def process_genexpert_hl7_message(conn, ip_addr, clean_hl7, response_framing): # 3. Update Database PaslabOrder if rnoreg: - # Cari kolom flag yang cocok dengan IP yang sedang terkoneksi - target_flag_col = None - for flag_col, mapped_ip in TARGET_MAPPING.items(): - if mapped_ip == ip_addr: - target_flag_col = flag_col - break - - if target_flag_col: - try: - # Buka sesi database dan update - with SessionLocal() as session: - order = session.query(PaslabOrder).filter(PaslabOrder.rnoreg == rnoreg).first() - if order: - # Set flag mesin tersebut menjadi True agar tidak dikirim ulang - setattr(order, target_flag_col, True) - session.commit() - print(f"[GENEXPERT-DB] Order {rnoreg} ditolak alat. Flag {target_flag_col} di-set True (Selesai).") - else: - print(f"[GENEXPERT-DB] Order {rnoreg} tidak ditemukan di database saat memproses penolakan.") - except Exception as e: - print(f"[GENEXPERT-DB-ERROR] Gagal update order duplikat {rnoreg}: {e}") + mark_genexpert_order_flag(rnoreg, ip_addr, reason="astm-comment-duplicate-or-rejected") print(f"[GENEXPERT-ASTM] Transaksi penolakan order selesai diproses.") return @@ -735,7 +764,6 @@ def process_genexpert_hl7_message(conn, ip_addr, clean_hl7, response_framing): return if "QCN^J01" in clean_hl7: - clear_genexpert_inflight_for_ip(ip_addr, reason="query-confirmation") ack_msg = create_genexpert_ack_j01_response(clean_hl7, ip_addr=ip_addr) send_genexpert_response(conn, ip_addr, ack_msg, response_framing, label="qcn-ack") print(f"[GENEXPERT] Menerima konfirmasi query dari {ip_addr}.") @@ -937,10 +965,6 @@ def handle_genexpert_client(conn, addr): with connection_lock: if active_genexpert_connections.get(client_ip) is conn: del active_genexpert_connections[client_ip] - remaining_connections = len(active_genexpert_connections) - clear_genexpert_inflight_for_ip(client_ip, reason="connection-closed") - if remaining_connections == 0: - stop_all_scheduled_result_queries(reason="no-active-genexpert") try: conn.close() except Exception: @@ -1052,13 +1076,8 @@ def parse_genexpert_astm_records(astm_string, device_name): ) session.add(new_result) source_ip = str(device_name or "").replace("GeneXpert-", "").strip() - flag_name = get_flag_by_device(source_ip) - if flag_name: - order = session.query(PaslabOrder).filter(PaslabOrder.rnoreg == seq_no_safe).first() - if order: - setattr(order, flag_name, True) - print(f"[GENEXPERT] Hasil ASTM diterima, set {flag_name}=TRUE untuk {seq_no_safe}") session.commit() + mark_genexpert_order_flag(seq_no_safe, source_ip, reason="astm-result-received") print(f"[GENEXPERT-DB-SUCCESS] Hasil lab untuk Order {seq_no_safe} berhasil disimpan ke LisPhoenix!") except Exception as e: @@ -1223,7 +1242,6 @@ def create_genexpert_rsp_z02_response(orders, incoming_hl7, ip_addr=None): def send_all_orders(conn, ip_addr, hl7_msg, msg_id, response_framing="mllp"): orders = get_genexpert_query_orders(ip_addr, hl7_msg) - scheduled_orders = [] if not orders: print(f"[GENEXPERT] Tidak ada order pending untuk {ip_addr}") rsp = create_genexpert_rsp_z02_response([], hl7_msg, ip_addr=ip_addr) @@ -1239,22 +1257,13 @@ def send_all_orders(conn, ip_addr, hl7_msg, msg_id, response_framing="mllp"): rsp = create_genexpert_rsp_z02_response(orders, hl7_msg, ip_addr=ip_addr) first_accnumber = str(orders[0].rnoreg or "").strip() if orders else "" debug_genexpert_order_message(rsp, ip_addr=ip_addr) - send_genexpert_response(conn, ip_addr, rsp, response_framing, label=f"qbp-order:{first_accnumber}") + sent_ok = send_genexpert_response(conn, ip_addr, rsp, response_framing, label=f"qbp-order:{first_accnumber}") for order in orders: - print(f"[GENEXPERT] Order ditawarkan ke {ip_addr}: {order.rnoreg}") - scheduled_orders.append({ - "accnumber": str(order.rnoreg or "").strip(), - "register_no": str(order.rnoreg or "").strip(), - "target_ip": ip_addr, - }) - - for scheduled_order in scheduled_orders: - schedule_result_query_for_order( - accnumber=scheduled_order["accnumber"], - register_no=scheduled_order["register_no"], - target_ip=scheduled_order["target_ip"], - ) + rnoreg = str(order.rnoreg or "").strip() + print(f"[GENEXPERT] Order ditawarkan ke {ip_addr}: {rnoreg}, sent_ok={sent_ok}") + if sent_ok: + mark_genexpert_order_flag(rnoreg, ip_addr, reason="hl7-order-sent") def log_genexpert_hl7(direction, ip_addr, hl7_message, label=""): message_type = extract_message_type(hl7_message) or "UNKNOWN" @@ -1424,8 +1433,7 @@ def send_genexpert_response(conn, ip_addr, hl7_message, framing, label=""): f"[GENEXPERT-DEBUG] Send response ip={ip_addr}, framing={framing}, " f"response_mode={response_mode}, label={label}, bytes=multi-frame" ) - send_genexpert_astm_frame(conn, ip_addr, hl7_message, label=label) - return + return send_genexpert_astm_frame(conn, ip_addr, hl7_message, label=label) payload = frame_genexpert_response(hl7_message, framing) if framing == "astm": @@ -1439,7 +1447,12 @@ def send_genexpert_response(conn, ip_addr, hl7_message, framing, label=""): f"[GENEXPERT-DEBUG] Send response ip={ip_addr}, framing={framing}, " f"response_mode={response_mode}, label={label}, bytes={len(payload)}" ) - conn.sendall(payload) + try: + conn.sendall(payload) + return True + except Exception as exc: + print(f"[GENEXPERT-DEBUG] Send response gagal ip={ip_addr}, label={label}, error={exc}") + return False def send_genexpert_transport_ack(conn, ip_addr, framing, reason="frame-received"): if framing != "astm": @@ -1554,43 +1567,6 @@ def get_active_genexpert_ips(): with connection_lock: return list(active_genexpert_connections.keys()) -def clear_genexpert_inflight_for_ip(ip_addr, reason="cleared"): - ip_addr = str(ip_addr or "").strip() - if not ip_addr: - return False - - with genexpert_query_inflight_lock: - state = genexpert_query_inflight_by_ip.pop(ip_addr, None) - - if state: - print( - f"[GENEXPERT-QUERY] Clear inflight ip={ip_addr}, reason={reason}, accnumber={state.get('accnumber')}" - ) - return True - return False - -def stop_all_scheduled_result_queries(reason="no-active-genexpert"): - with scheduled_result_query_lock: - scheduled_result_queries.clear() - with genexpert_query_inflight_lock: - genexpert_query_inflight_by_ip.clear() - print(f"[GENEXPERT-SCHEDULER] Nonaktif. Clear jadwal, reason={reason}") - return 0 - -def stop_scheduled_result_query(accnumber, reason="completed"): - accnumber = str(accnumber or "").strip() - with scheduled_result_query_lock: - scheduled_result_queries.pop(accnumber, None) - print(f"[GENEXPERT-SCHEDULER] Nonaktif. Stop accnumber={accnumber}, reason={reason}") - return False - -def schedule_result_query_for_order(accnumber, register_no, target_ip=None, **kwargs): - print( - f"[GENEXPERT-SCHEDULER] Nonaktif. Jadwal query hasil dilewati " - f"untuk accnumber={accnumber}, target_ip={target_ip}" - ) - return False - def create_hl7_dsr_response(order, msg_control_id, qrd_segment): """ Membuat pesan balasan DSR^Q03 (Data Response) untuk GeneXpert.