This commit is contained in:
Dwi Swandhana
2026-08-06 19:32:16 +07:00
parent f0dc7987a9
commit b82a799e20
+70 -94
View File
@@ -14,7 +14,7 @@ import datetime
import traceback import traceback
import serial # type: ignore 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 DateTime as SqDateTime # type: ignore
from sqlalchemy import Date as SqDate # type: ignore from sqlalchemy import Date as SqDate # type: ignore
from sqlalchemy.orm import declarative_base, sessionmaker # type: ignore from sqlalchemy.orm import declarative_base, sessionmaker # type: ignore
@@ -55,10 +55,6 @@ ERROR_LOG_KEYWORDS = (
# Global Variables # Global Variables
active_genexpert_connections = {} active_genexpert_connections = {}
connection_lock = threading.Lock() 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 # Network Configuration
TCP_LISTENER_PORT = 6001 # PC GeneXpert set ke mode Client, konek ke IP:PORT ini 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 SERVER_HOST = '0.0.0.0' # Listen di semua interface
@@ -339,6 +335,56 @@ def get_pending_orders(ip_addr):
finally: finally:
session.close() 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): def parse_hl7_segments(hl7_message):
return [segment for segment in str(hl7_message or "").split('\r') if segment] 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 return
reply = create_genexpert_astm_order_message(selected_orders, ip_addr=ip_addr, query_tag=query_tag) 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: 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: finally:
session.close() session.close()
@@ -636,27 +685,7 @@ def process_genexpert_hl7_message(conn, ip_addr, clean_hl7, response_framing):
# 3. Update Database PaslabOrder # 3. Update Database PaslabOrder
if rnoreg: if rnoreg:
# Cari kolom flag yang cocok dengan IP yang sedang terkoneksi mark_genexpert_order_flag(rnoreg, ip_addr, reason="astm-comment-duplicate-or-rejected")
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}")
print(f"[GENEXPERT-ASTM] Transaksi penolakan order selesai diproses.") print(f"[GENEXPERT-ASTM] Transaksi penolakan order selesai diproses.")
return return
@@ -735,7 +764,6 @@ def process_genexpert_hl7_message(conn, ip_addr, clean_hl7, response_framing):
return return
if "QCN^J01" in clean_hl7: 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) 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") send_genexpert_response(conn, ip_addr, ack_msg, response_framing, label="qcn-ack")
print(f"[GENEXPERT] Menerima konfirmasi query dari {ip_addr}.") print(f"[GENEXPERT] Menerima konfirmasi query dari {ip_addr}.")
@@ -937,10 +965,6 @@ def handle_genexpert_client(conn, addr):
with connection_lock: with connection_lock:
if active_genexpert_connections.get(client_ip) is conn: if active_genexpert_connections.get(client_ip) is conn:
del active_genexpert_connections[client_ip] 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: try:
conn.close() conn.close()
except Exception: except Exception:
@@ -1052,13 +1076,8 @@ def parse_genexpert_astm_records(astm_string, device_name):
) )
session.add(new_result) session.add(new_result)
source_ip = str(device_name or "").replace("GeneXpert-", "").strip() 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() 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!") print(f"[GENEXPERT-DB-SUCCESS] Hasil lab untuk Order {seq_no_safe} berhasil disimpan ke LisPhoenix!")
except Exception as e: 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"): def send_all_orders(conn, ip_addr, hl7_msg, msg_id, response_framing="mllp"):
orders = get_genexpert_query_orders(ip_addr, hl7_msg) orders = get_genexpert_query_orders(ip_addr, hl7_msg)
scheduled_orders = []
if not orders: if not orders:
print(f"[GENEXPERT] Tidak ada order pending untuk {ip_addr}") print(f"[GENEXPERT] Tidak ada order pending untuk {ip_addr}")
rsp = create_genexpert_rsp_z02_response([], hl7_msg, ip_addr=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) rsp = create_genexpert_rsp_z02_response(orders, hl7_msg, ip_addr=ip_addr)
first_accnumber = str(orders[0].rnoreg or "").strip() if orders else "" first_accnumber = str(orders[0].rnoreg or "").strip() if orders else ""
debug_genexpert_order_message(rsp, ip_addr=ip_addr) 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: for order in orders:
print(f"[GENEXPERT] Order ditawarkan ke {ip_addr}: {order.rnoreg}") rnoreg = str(order.rnoreg or "").strip()
scheduled_orders.append({ print(f"[GENEXPERT] Order ditawarkan ke {ip_addr}: {rnoreg}, sent_ok={sent_ok}")
"accnumber": str(order.rnoreg or "").strip(), if sent_ok:
"register_no": str(order.rnoreg or "").strip(), mark_genexpert_order_flag(rnoreg, ip_addr, reason="hl7-order-sent")
"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"],
)
def log_genexpert_hl7(direction, ip_addr, hl7_message, label=""): def log_genexpert_hl7(direction, ip_addr, hl7_message, label=""):
message_type = extract_message_type(hl7_message) or "UNKNOWN" 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"[GENEXPERT-DEBUG] Send response ip={ip_addr}, framing={framing}, "
f"response_mode={response_mode}, label={label}, bytes=multi-frame" f"response_mode={response_mode}, label={label}, bytes=multi-frame"
) )
send_genexpert_astm_frame(conn, ip_addr, hl7_message, label=label) return send_genexpert_astm_frame(conn, ip_addr, hl7_message, label=label)
return
payload = frame_genexpert_response(hl7_message, framing) payload = frame_genexpert_response(hl7_message, framing)
if framing == "astm": 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"[GENEXPERT-DEBUG] Send response ip={ip_addr}, framing={framing}, "
f"response_mode={response_mode}, label={label}, bytes={len(payload)}" 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"): def send_genexpert_transport_ack(conn, ip_addr, framing, reason="frame-received"):
if framing != "astm": if framing != "astm":
@@ -1554,43 +1567,6 @@ def get_active_genexpert_ips():
with connection_lock: with connection_lock:
return list(active_genexpert_connections.keys()) 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): def create_hl7_dsr_response(order, msg_control_id, qrd_segment):
""" """
Membuat pesan balasan DSR^Q03 (Data Response) untuk GeneXpert. Membuat pesan balasan DSR^Q03 (Data Response) untuk GeneXpert.