update
This commit is contained in:
+237
-320
@@ -19,22 +19,64 @@ 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
|
||||||
# Logging Setup
|
# Logging Setup
|
||||||
# Konfigurasi logging per hari
|
|
||||||
|
def get_app_dir():
|
||||||
|
"""
|
||||||
|
Mengembalikan direktori aplikasi.
|
||||||
|
|
||||||
|
Python:
|
||||||
|
folder tempat app.py berada
|
||||||
|
|
||||||
|
PyInstaller EXE:
|
||||||
|
folder tempat app.exe berada
|
||||||
|
"""
|
||||||
|
if getattr(sys, "frozen", False):
|
||||||
|
return os.path.dirname(os.path.abspath(sys.executable))
|
||||||
|
|
||||||
|
return os.path.dirname(os.path.abspath(__file__))
|
||||||
|
|
||||||
|
|
||||||
|
APP_DIR = get_app_dir()
|
||||||
|
|
||||||
|
THREAD_LOG_DIR = os.path.join(APP_DIR, "thread_logs")
|
||||||
|
APP_LOG_FILE = os.path.join(APP_DIR, "app.log")
|
||||||
|
|
||||||
|
thread_log_lock = threading.Lock()
|
||||||
|
|
||||||
|
|
||||||
|
# Buat folder sejak startup
|
||||||
|
try:
|
||||||
|
os.makedirs(THREAD_LOG_DIR, exist_ok=True)
|
||||||
|
except Exception as exc:
|
||||||
|
builtins.print(
|
||||||
|
f"[LOG-INIT-ERROR] Tidak bisa membuat "
|
||||||
|
f"{THREAD_LOG_DIR}: {exc}"
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
# Logging Setup
|
||||||
log_handler = TimedRotatingFileHandler(
|
log_handler = TimedRotatingFileHandler(
|
||||||
filename="app.log",
|
filename=APP_LOG_FILE,
|
||||||
when="midnight",
|
when="midnight",
|
||||||
interval=1,
|
interval=1,
|
||||||
backupCount=7,
|
backupCount=7,
|
||||||
encoding="utf-8"
|
encoding="utf-8"
|
||||||
)
|
)
|
||||||
formatter = logging.Formatter('%(asctime)s - %(levelname)s - %(threadName)s - %(message)s')
|
|
||||||
|
formatter = logging.Formatter(
|
||||||
|
'%(asctime)s - %(levelname)s - %(threadName)s - %(message)s'
|
||||||
|
)
|
||||||
|
|
||||||
log_handler.setFormatter(formatter)
|
log_handler.setFormatter(formatter)
|
||||||
log_handler.setLevel(logging.ERROR)
|
log_handler.setLevel(logging.ERROR)
|
||||||
logging.basicConfig(level=logging.ERROR, handlers=[log_handler])
|
|
||||||
|
logging.basicConfig(
|
||||||
|
level=logging.ERROR,
|
||||||
|
handlers=[log_handler]
|
||||||
|
)
|
||||||
|
|
||||||
logging.getLogger().setLevel(logging.ERROR)
|
logging.getLogger().setLevel(logging.ERROR)
|
||||||
logging.getLogger("werkzeug").setLevel(logging.ERROR)
|
logging.getLogger("werkzeug").setLevel(logging.ERROR)
|
||||||
|
|
||||||
THREAD_LOG_DIR = os.path.join(os.path.dirname(os.path.abspath(__file__)), "thread_logs")
|
|
||||||
thread_log_lock = threading.Lock()
|
thread_log_lock = threading.Lock()
|
||||||
ERROR_LOG_KEYWORDS = (
|
ERROR_LOG_KEYWORDS = (
|
||||||
"critical",
|
"critical",
|
||||||
@@ -66,7 +108,7 @@ GENEXPERT_HOST_APPLICATION_DEFAULT = "DE002"
|
|||||||
GENEXPERT_HOST_APPLICATION_BY_IP = {
|
GENEXPERT_HOST_APPLICATION_BY_IP = {
|
||||||
"10.10.120.75": "GE01", # GenExpert Kecil
|
"10.10.120.75": "GE01", # GenExpert Kecil
|
||||||
"10.10.120.108": "DE002", # GenExpert Tempat Lama
|
"10.10.120.108": "DE002", # GenExpert Tempat Lama
|
||||||
"10.10.120.73": "GE01", # GenExpert Besa
|
"10.10.120.73": "GE01", # GenExpert Besar
|
||||||
}
|
}
|
||||||
# Mapping Flag ke IP Address GeneXpert
|
# Mapping Flag ke IP Address GeneXpert
|
||||||
# Pastikan IP ini SESUAI dengan settingan "Server IP" di masing-masing alat (Client Mode)
|
# Pastikan IP ini SESUAI dengan settingan "Server IP" di masing-masing alat (Client Mode)
|
||||||
@@ -75,14 +117,14 @@ TARGET_MAPPING = {
|
|||||||
'flg_gxp2': '10.10.120.108',
|
'flg_gxp2': '10.10.120.108',
|
||||||
'flg_gxp3': '10.10.120.75'
|
'flg_gxp3': '10.10.120.75'
|
||||||
}
|
}
|
||||||
|
GENEXPERT_ORDER_LOCKS = {
|
||||||
|
"10.10.120.73": threading.Lock(),
|
||||||
|
"10.10.120.108": threading.Lock(),
|
||||||
|
"10.10.120.75": threading.Lock(),
|
||||||
|
}
|
||||||
# GeneXpert Configuration
|
# GeneXpert Configuration
|
||||||
# ==========================================
|
# ==========================================
|
||||||
# KONFIGURASI MAPPING TES (DATABASE -> GENEXPERT)
|
|
||||||
# ==========================================
|
|
||||||
# Kiri: Nama di kolom 'tes' database Anda
|
|
||||||
# Kanan: 'Host Test Code' dari Dokumen Word Anda
|
|
||||||
GENEXPERT_TEST_MAPPING = {
|
GENEXPERT_TEST_MAPPING = {
|
||||||
# Mapping untuk IP 10.10.120.75 (Multi-Assay)
|
|
||||||
"HIV": "HIV1-VL",
|
"HIV": "HIV1-VL",
|
||||||
"HBV": "HBVVL",
|
"HBV": "HBVVL",
|
||||||
"TCM TB": "MTBRIF",
|
"TCM TB": "MTBRIF",
|
||||||
@@ -130,7 +172,6 @@ GENEXPERT_TEST_MAPPING = {
|
|||||||
"8.3.8 TCM TBC (GENE EXPERT)": "MTBRIF",
|
"8.3.8 TCM TBC (GENE EXPERT)": "MTBRIF",
|
||||||
"10.3.8 TCM TBC (GENE EXPERT)": "MTBRIF",
|
"10.3.8 TCM TBC (GENE EXPERT)": "MTBRIF",
|
||||||
"11.3.9 TCM TB (GENE EXPERT)": "MTBRIF",
|
"11.3.9 TCM TB (GENE EXPERT)": "MTBRIF",
|
||||||
|
|
||||||
}
|
}
|
||||||
GENEXPERT_IP_CAPABILITIES = {
|
GENEXPERT_IP_CAPABILITIES = {
|
||||||
"10.10.120.75": ["MTBRIF", "HBVVL", "HIV1-VL", "MTB-XDR 2", "HCV", "SARSCOV2FLURSV"],
|
"10.10.120.75": ["MTBRIF", "HBVVL", "HIV1-VL", "MTB-XDR 2", "HCV", "SARSCOV2FLURSV"],
|
||||||
@@ -177,13 +218,16 @@ active_serial_ports = {}
|
|||||||
order_queues = {config['port']: Queue() for config in DEVICE_CONFIGS if config['protocol'] == 'serial'}
|
order_queues = {config['port']: Queue() for config in DEVICE_CONFIGS if config['protocol'] == 'serial'}
|
||||||
|
|
||||||
# ==========================================
|
# ==========================================
|
||||||
# 2. DATABASE MODEL
|
# 2. DATABASE CONNECTION
|
||||||
# ==========================================
|
# ==========================================
|
||||||
DATABASE_URL = "postgresql://lismikro:[email protected]:5002/lismikro"
|
DATABASE_URL = "postgresql://lismikro:[email protected]:5002/lismikro"
|
||||||
engine = create_engine(DATABASE_URL, pool_recycle=3600)
|
engine = create_engine(DATABASE_URL, pool_recycle=3600)
|
||||||
SessionLocal = sessionmaker(bind=engine)
|
SessionLocal = sessionmaker(bind=engine)
|
||||||
Base = declarative_base()
|
Base = declarative_base()
|
||||||
|
|
||||||
|
# ==========================================
|
||||||
|
# 2. DATABASE MODEL
|
||||||
|
# ==========================================
|
||||||
class Sample(Base):
|
class Sample(Base):
|
||||||
__tablename__ = 'samples'
|
__tablename__ = 'samples'
|
||||||
id = Column(Integer, primary_key=True)
|
id = Column(Integer, primary_key=True)
|
||||||
@@ -278,8 +322,12 @@ def _write_thread_log(message):
|
|||||||
with thread_log_lock:
|
with thread_log_lock:
|
||||||
with open(log_path, "a", encoding="utf-8") as fh:
|
with open(log_path, "a", encoding="utf-8") as fh:
|
||||||
fh.write(f"{timestamp} | {message}\n")
|
fh.write(f"{timestamp} | {message}\n")
|
||||||
except Exception:
|
except Exception as exc:
|
||||||
pass
|
builtins.print(
|
||||||
|
f"[THREAD-LOG-ERROR] "
|
||||||
|
f"Gagal menulis thread log ke "
|
||||||
|
f"{THREAD_LOG_DIR}: {exc}"
|
||||||
|
)
|
||||||
|
|
||||||
def print(*args, **kwargs):
|
def print(*args, **kwargs):
|
||||||
sep = kwargs.get("sep", " ")
|
sep = kwargs.get("sep", " ")
|
||||||
@@ -323,18 +371,6 @@ def get_flag_by_device(ip_addr):
|
|||||||
return flag
|
return flag
|
||||||
return None
|
return None
|
||||||
|
|
||||||
def get_pending_orders(ip_addr):
|
|
||||||
flag = get_flag_by_device(ip_addr)
|
|
||||||
if not flag:
|
|
||||||
return []
|
|
||||||
|
|
||||||
session = SessionLocal()
|
|
||||||
try:
|
|
||||||
q = session.query(PaslabOrder).filter(getattr(PaslabOrder, flag) == False)
|
|
||||||
return q.all()
|
|
||||||
finally:
|
|
||||||
session.close()
|
|
||||||
|
|
||||||
def mark_genexpert_order_flag(rnoreg, ip_addr, reason="processed"):
|
def mark_genexpert_order_flag(rnoreg, ip_addr, reason="processed"):
|
||||||
rnoreg = str(rnoreg or "").strip()
|
rnoreg = str(rnoreg or "").strip()
|
||||||
flag_name = get_flag_by_device(str(ip_addr or "").strip())
|
flag_name = get_flag_by_device(str(ip_addr or "").strip())
|
||||||
@@ -424,16 +460,6 @@ def build_hl7_segment(segment_name, fields):
|
|||||||
values = [str(fields.get(index, "")) for index in range(1, max_field + 1)]
|
values = [str(fields.get(index, "")) for index in range(1, max_field + 1)]
|
||||||
return f"{segment_name}|" + "|".join(values)
|
return f"{segment_name}|" + "|".join(values)
|
||||||
|
|
||||||
def map_hl7_sex(value):
|
|
||||||
text = str(value or "").strip().upper()
|
|
||||||
if not text:
|
|
||||||
return ""
|
|
||||||
if text.startswith("L") or "LAKI" in text or text == "M":
|
|
||||||
return "M"
|
|
||||||
if text.startswith("P") or "PEREM" in text or text == "F":
|
|
||||||
return "F"
|
|
||||||
return ""
|
|
||||||
|
|
||||||
def extract_msg_control_id(hl7_message):
|
def extract_msg_control_id(hl7_message):
|
||||||
try:
|
try:
|
||||||
segments = hl7_message.split('\r')
|
segments = hl7_message.split('\r')
|
||||||
@@ -588,60 +614,164 @@ def create_genexpert_astm_order_message(orders, ip_addr=None, query_tag=""):
|
|||||||
return message
|
return message
|
||||||
|
|
||||||
def send_all_orders_astm(conn, ip_addr, astm_msg, response_framing="astm"):
|
def send_all_orders_astm(conn, ip_addr, astm_msg, response_framing="astm"):
|
||||||
query = parse_genexpert_astm_query(astm_msg)
|
lock = GENEXPERT_ORDER_LOCKS.setdefault(ip_addr, threading.Lock())
|
||||||
requested_sample_id = str(query.get("query_sample_id") or "").strip()
|
|
||||||
query_tag = str(query.get("query_tag") or "").strip()
|
|
||||||
flag = get_flag_by_device(ip_addr)
|
|
||||||
if not flag:
|
|
||||||
print(f"[GENEXPERT] ASTM query diabaikan, flag untuk {ip_addr} tidak ditemukan.")
|
|
||||||
return
|
|
||||||
|
|
||||||
session = SessionLocal()
|
print(f"[GENEXPERT-LOCK] ip={ip_addr}, waiting")
|
||||||
try:
|
|
||||||
flag_attr = getattr(PaslabOrder, flag, None)
|
with lock:
|
||||||
if flag_attr is None:
|
print(f"[GENEXPERT-LOCK] ip={ip_addr}, acquired")
|
||||||
print(f"[GENEXPERT] ASTM query diabaikan, atribut flag {flag} tidak ada.")
|
|
||||||
|
query = parse_genexpert_astm_query(astm_msg)
|
||||||
|
requested_sample_id = str(query.get("query_sample_id") or "").strip()
|
||||||
|
query_tag = str(query.get("query_tag") or "").strip()
|
||||||
|
|
||||||
|
flag = get_flag_by_device(ip_addr)
|
||||||
|
|
||||||
|
if not flag:
|
||||||
|
print(
|
||||||
|
f"[GENEXPERT] ASTM query diabaikan, "
|
||||||
|
f"flag untuk {ip_addr} tidak ditemukan."
|
||||||
|
)
|
||||||
return
|
return
|
||||||
|
|
||||||
base_orders = session.query(PaslabOrder).filter(
|
|
||||||
(flag_attr == False) | (flag_attr == None)
|
|
||||||
).order_by(PaslabOrder.rtglast.desc().nullslast(), PaslabOrder.urut.desc()).all()
|
|
||||||
|
|
||||||
selected_orders = []
|
|
||||||
if requested_sample_id and requested_sample_id.upper() != "ALL":
|
|
||||||
for order in base_orders:
|
|
||||||
if str(order.rnoreg or "").strip() != requested_sample_id:
|
|
||||||
continue
|
|
||||||
assay_code, _, _ = resolve_genexpert_assay(order, ip_addr)
|
|
||||||
if assay_code:
|
|
||||||
selected_orders = [order]
|
|
||||||
break
|
|
||||||
else:
|
|
||||||
for order in base_orders:
|
|
||||||
assay_code, _, _ = resolve_genexpert_assay(order, ip_addr)
|
|
||||||
if assay_code:
|
|
||||||
selected_orders = [order]
|
|
||||||
break
|
|
||||||
|
|
||||||
print(
|
print(
|
||||||
f"[GENEXPERT-ASTM-QUERY] ip={ip_addr}, query_tag={query_tag}, requested_sample_id='{requested_sample_id}', "
|
f"[GENEXPERT-ORDER-SELECT] "
|
||||||
f"selected_rnoreg={[str(order.rnoreg or '').strip() for order in selected_orders]}"
|
f"ip={ip_addr}, flag={flag}, "
|
||||||
|
f"requested_sample_id='{requested_sample_id}'"
|
||||||
)
|
)
|
||||||
|
|
||||||
if not selected_orders:
|
session = SessionLocal()
|
||||||
reply = f"H|\\^&|||{sanitize_astm_field(get_genexpert_host_application(ip_addr), max_len=20)}|||||GeneXpert Host||P|1394-97|{datetime.datetime.now().strftime('%Y%m%d%H%M%S')}\rL|1|N\r"
|
|
||||||
send_genexpert_response(conn, ip_addr, reply, response_framing, label="astm-q-empty")
|
|
||||||
return
|
|
||||||
|
|
||||||
reply = create_genexpert_astm_order_message(selected_orders, ip_addr=ip_addr, query_tag=query_tag)
|
try:
|
||||||
sent_ok = send_genexpert_response(conn, ip_addr, reply, response_framing, label=f"astm-q-order:{selected_orders[0].rnoreg}")
|
flag_attr = getattr(PaslabOrder, flag, None)
|
||||||
for order in selected_orders:
|
|
||||||
rnoreg = str(order.rnoreg or "").strip()
|
if flag_attr is None:
|
||||||
print(f"[GENEXPERT] Order ASTM ditawarkan ke {ip_addr}: {rnoreg}, sent_ok={sent_ok}")
|
print(
|
||||||
if sent_ok:
|
f"[GENEXPERT] ASTM query diabaikan, "
|
||||||
mark_genexpert_order_flag(rnoreg, ip_addr, reason="astm-order-sent")
|
f"atribut flag {flag} tidak ada."
|
||||||
finally:
|
)
|
||||||
session.close()
|
return
|
||||||
|
|
||||||
|
base_orders = (
|
||||||
|
session.query(PaslabOrder)
|
||||||
|
.filter(
|
||||||
|
(flag_attr == False) | (flag_attr == None)
|
||||||
|
)
|
||||||
|
.order_by(
|
||||||
|
PaslabOrder.rtglast.desc().nullslast(),
|
||||||
|
PaslabOrder.urut.desc()
|
||||||
|
)
|
||||||
|
.all()
|
||||||
|
)
|
||||||
|
|
||||||
|
selected_orders = []
|
||||||
|
|
||||||
|
if requested_sample_id and requested_sample_id.upper() != "ALL":
|
||||||
|
|
||||||
|
for order in base_orders:
|
||||||
|
|
||||||
|
if str(order.rnoreg or "").strip() != requested_sample_id:
|
||||||
|
continue
|
||||||
|
|
||||||
|
assay_code, _, _ = resolve_genexpert_assay(
|
||||||
|
order,
|
||||||
|
ip_addr
|
||||||
|
)
|
||||||
|
|
||||||
|
if assay_code:
|
||||||
|
selected_orders = [order]
|
||||||
|
break
|
||||||
|
|
||||||
|
else:
|
||||||
|
|
||||||
|
for order in base_orders:
|
||||||
|
|
||||||
|
assay_code, _, _ = resolve_genexpert_assay(
|
||||||
|
order,
|
||||||
|
ip_addr
|
||||||
|
)
|
||||||
|
|
||||||
|
if assay_code:
|
||||||
|
selected_orders = [order]
|
||||||
|
break
|
||||||
|
|
||||||
|
print(
|
||||||
|
f"[GENEXPERT-ASTM-QUERY] "
|
||||||
|
f"ip={ip_addr}, "
|
||||||
|
f"flag={flag}, "
|
||||||
|
f"query_tag={query_tag}, "
|
||||||
|
f"requested_sample_id='{requested_sample_id}', "
|
||||||
|
f"selected_rnoreg="
|
||||||
|
f"{[str(order.rnoreg or '').strip() for order in selected_orders]}"
|
||||||
|
)
|
||||||
|
|
||||||
|
if not selected_orders:
|
||||||
|
|
||||||
|
reply = (
|
||||||
|
f"H|\\^&|||"
|
||||||
|
f"{sanitize_astm_field(get_genexpert_host_application(ip_addr), max_len=20)}"
|
||||||
|
f"|||||GeneXpert Host||P|1394-97|"
|
||||||
|
f"{datetime.datetime.now().strftime('%Y%m%d%H%M%S')}"
|
||||||
|
f"\rL|1|N\r"
|
||||||
|
)
|
||||||
|
|
||||||
|
send_genexpert_response(
|
||||||
|
conn,
|
||||||
|
ip_addr,
|
||||||
|
reply,
|
||||||
|
response_framing,
|
||||||
|
label="astm-q-empty"
|
||||||
|
)
|
||||||
|
|
||||||
|
return
|
||||||
|
|
||||||
|
reply = create_genexpert_astm_order_message(
|
||||||
|
selected_orders,
|
||||||
|
ip_addr=ip_addr,
|
||||||
|
query_tag=query_tag
|
||||||
|
)
|
||||||
|
|
||||||
|
rnoreg_first = str(
|
||||||
|
selected_orders[0].rnoreg or ""
|
||||||
|
).strip()
|
||||||
|
|
||||||
|
sent_ok = send_genexpert_response(
|
||||||
|
conn,
|
||||||
|
ip_addr,
|
||||||
|
reply,
|
||||||
|
response_framing,
|
||||||
|
label=f"astm-q-order:{rnoreg_first}"
|
||||||
|
)
|
||||||
|
|
||||||
|
for order in selected_orders:
|
||||||
|
|
||||||
|
rnoreg = str(order.rnoreg or "").strip()
|
||||||
|
|
||||||
|
print(
|
||||||
|
f"[GENEXPERT] Order ASTM ditawarkan ke "
|
||||||
|
f"{ip_addr}: {rnoreg}, sent_ok={sent_ok}"
|
||||||
|
)
|
||||||
|
|
||||||
|
if sent_ok:
|
||||||
|
|
||||||
|
flag_ok = mark_genexpert_order_flag(
|
||||||
|
rnoreg,
|
||||||
|
ip_addr,
|
||||||
|
reason="astm-order-sent"
|
||||||
|
)
|
||||||
|
|
||||||
|
print(
|
||||||
|
f"[GENEXPERT-ORDER-COMMIT] "
|
||||||
|
f"ip={ip_addr}, "
|
||||||
|
f"rnoreg={rnoreg}, "
|
||||||
|
f"flag={flag}, "
|
||||||
|
f"success={flag_ok}"
|
||||||
|
)
|
||||||
|
|
||||||
|
finally:
|
||||||
|
session.close()
|
||||||
|
|
||||||
|
print(f"[GENEXPERT-LOCK] ip={ip_addr}, released")
|
||||||
|
|
||||||
def process_genexpert_hl7_message(conn, ip_addr, clean_hl7, response_framing):
|
def process_genexpert_hl7_message(conn, ip_addr, clean_hl7, response_framing):
|
||||||
# ==========================================================
|
# ==========================================================
|
||||||
@@ -971,40 +1101,6 @@ def handle_genexpert_client(conn, addr):
|
|||||||
pass
|
pass
|
||||||
logging.info(f"[GenExpert_TCP] Koneksi {addr} ditutup.")
|
logging.info(f"[GenExpert_TCP] Koneksi {addr} ditutup.")
|
||||||
|
|
||||||
def send_order_via_active_connection(target_ip, hl7_message):
|
|
||||||
conn = None
|
|
||||||
with connection_lock:
|
|
||||||
conn = active_genexpert_connections.get(target_ip)
|
|
||||||
|
|
||||||
if not conn:
|
|
||||||
logging.warning(f"Gagal kirim Order: GeneXpert dengan IP {target_ip} BELUM TERKONEKSI ke Listener.")
|
|
||||||
print(f"Gagal kirim Order: GeneXpert dengan IP {target_ip} BELUM TERKONEKSI ke Listener.")
|
|
||||||
return False
|
|
||||||
|
|
||||||
try:
|
|
||||||
# Bungkus pesan dengan MLLP (Minimal Lower Layer Protocol) standard HL7
|
|
||||||
# Format: <VT> message <FS><CR>
|
|
||||||
mllp_msg = f"\x0b{hl7_message}\x1c\r"
|
|
||||||
|
|
||||||
logging.info(f"Mengirim Order ke {target_ip}...")
|
|
||||||
print(f"Mengirim Order ke {target_ip}...")
|
|
||||||
conn.sendall(mllp_msg.encode('utf-8'))
|
|
||||||
|
|
||||||
# Opsi: Jika ingin menunggu ACK balasan untuk Order
|
|
||||||
# Namun hati-hati ini bisa blocking jika alat lambat
|
|
||||||
ack = conn.recv(1024)
|
|
||||||
logging.info(f"Dapat ACK Order dari {target_ip}: {ack}")
|
|
||||||
print(f"Dapat ACK Order dari {target_ip}: {ack}")
|
|
||||||
return True
|
|
||||||
except Exception as e:
|
|
||||||
logging.error(f"Error mengirim ke {target_ip}: {e}")
|
|
||||||
print(f"Error mengirim ke {target_ip}: {e}")
|
|
||||||
# Jika error saat kirim, anggap koneksi rusak
|
|
||||||
with connection_lock:
|
|
||||||
if target_ip in active_genexpert_connections:
|
|
||||||
del active_genexpert_connections[target_ip]
|
|
||||||
return False
|
|
||||||
|
|
||||||
def parse_genexpert_astm_records(astm_string, device_name):
|
def parse_genexpert_astm_records(astm_string, device_name):
|
||||||
"""
|
"""
|
||||||
Parser khusus untuk membaca hasil ASTM dari instrumen GeneXpert
|
Parser khusus untuk membaca hasil ASTM dari instrumen GeneXpert
|
||||||
@@ -1100,28 +1196,35 @@ def parse_genexpert_qpd(qpd_segment):
|
|||||||
def resolve_genexpert_assay(order, ip_addr=None):
|
def resolve_genexpert_assay(order, ip_addr=None):
|
||||||
assay_name = str(getattr(order, "tes", "") or "").strip()
|
assay_name = str(getattr(order, "tes", "") or "").strip()
|
||||||
specimen_code = str(getattr(order, "kd_spesimen", "") or "").strip()
|
specimen_code = str(getattr(order, "kd_spesimen", "") or "").strip()
|
||||||
|
|
||||||
assay_code = GENEXPERT_TEST_MAPPING.get(assay_name)
|
assay_code = GENEXPERT_TEST_MAPPING.get(assay_name)
|
||||||
assay_source = "mapping:tes"
|
assay_source = "mapping:tes"
|
||||||
|
|
||||||
if not assay_code and specimen_code:
|
|
||||||
assay_code = specimen_code
|
|
||||||
assay_source = "fallback:kd_spesimen"
|
|
||||||
|
|
||||||
supported_codes = GENEXPERT_IP_CAPABILITIES.get(str(ip_addr or "").strip(), []) if ip_addr else []
|
|
||||||
capability_match = True if not supported_codes else assay_code in supported_codes
|
|
||||||
|
|
||||||
if not assay_code:
|
if not assay_code:
|
||||||
print(
|
assay_code = DEFAULT_GXP_CODE
|
||||||
f"[GENEXPERT-DEBUG] rnoreg={getattr(order, 'rnoreg', '')}, ip={ip_addr}, "
|
assay_source = "default"
|
||||||
f"tes='{assay_name}', kd_spesimen='{specimen_code}', assay_code=EMPTY"
|
|
||||||
)
|
supported_codes = (
|
||||||
return None, "none", capability_match
|
GENEXPERT_IP_CAPABILITIES.get(str(ip_addr or "").strip(), [])
|
||||||
|
if ip_addr else []
|
||||||
|
)
|
||||||
|
|
||||||
|
capability_match = (
|
||||||
|
True if not supported_codes
|
||||||
|
else assay_code in supported_codes
|
||||||
|
)
|
||||||
|
|
||||||
print(
|
print(
|
||||||
f"[GENEXPERT-DEBUG] rnoreg={getattr(order, 'rnoreg', '')}, ip={ip_addr}, "
|
f"[GENEXPERT-DEBUG] "
|
||||||
f"tes='{assay_name}', kd_spesimen='{specimen_code}', assay_code='{assay_code}', "
|
f"rnoreg={getattr(order, 'rnoreg', '')}, "
|
||||||
f"assay_source={assay_source}, capability_match={capability_match}"
|
f"ip={ip_addr}, "
|
||||||
|
f"tes='{assay_name}', "
|
||||||
|
f"kd_spesimen='{specimen_code}', "
|
||||||
|
f"assay_code='{assay_code}', "
|
||||||
|
f"assay_source={assay_source}, "
|
||||||
|
f"capability_match={capability_match}"
|
||||||
)
|
)
|
||||||
|
|
||||||
return assay_code, assay_source, capability_match
|
return assay_code, assay_source, capability_match
|
||||||
|
|
||||||
def get_genexpert_query_orders(ip_addr, hl7_msg):
|
def get_genexpert_query_orders(ip_addr, hl7_msg):
|
||||||
@@ -1563,73 +1666,6 @@ def debug_genexpert_order_message(hl7_message, ip_addr=None):
|
|||||||
f"patient_name='{current_patient_name}', raw='{segment}'"
|
f"patient_name='{current_patient_name}', raw='{segment}'"
|
||||||
)
|
)
|
||||||
|
|
||||||
def get_active_genexpert_ips():
|
|
||||||
with connection_lock:
|
|
||||||
return list(active_genexpert_connections.keys())
|
|
||||||
|
|
||||||
def create_hl7_dsr_response(order, msg_control_id, qrd_segment):
|
|
||||||
"""
|
|
||||||
Membuat pesan balasan DSR^Q03 (Data Response) untuk GeneXpert.
|
|
||||||
"""
|
|
||||||
timestamp = datetime.datetime.now().strftime('%Y%m%d%H%M%S')
|
|
||||||
|
|
||||||
# --- 1. HEADER (MSH) ---
|
|
||||||
# Field 9 (Message Type) adalah DSR^Q03
|
|
||||||
# Field 10 (Control ID) kita generate baru
|
|
||||||
# Field 6 (Receiving Fac) harusnya GeneXpert
|
|
||||||
resp_control_id = f"RESP{timestamp}"
|
|
||||||
msh = f"MSH|^~\\&|LIS|LAB|GeneXpert|Cepheid|{timestamp}||DSR^Q03|{resp_control_id}|P|2.5"
|
|
||||||
|
|
||||||
# --- 2. ACKNOWLEDGEMENT (MSA) ---
|
|
||||||
# AA = Application Accept (Kita mengerti pertanyaannya)
|
|
||||||
# msg_control_id = ID dari pesan QRY yang dikirim GeneXpert (Supaya dia tahu ini jawaban untuk pertanyaan yg mana)
|
|
||||||
msa = f"MSA|AA|{msg_control_id}"
|
|
||||||
|
|
||||||
# --- 3. QUERY DEFINITION (QRD) ---
|
|
||||||
# Kita kembalikan segmen QRD yang dikirim alat (Echo back)
|
|
||||||
# qrd_segment harus string raw dari pesan masuk
|
|
||||||
qrd = qrd_segment
|
|
||||||
|
|
||||||
# --- 4. QUERY RESPONSE STATUS (QRF) ---
|
|
||||||
# Opsional, tapi baik untuk konfirmasi
|
|
||||||
qrf = f"QRF|LIS|{timestamp}||||"
|
|
||||||
|
|
||||||
# --- 5. DATA PASIEN & ORDER (Jika Order Ditemukan) ---
|
|
||||||
if order:
|
|
||||||
# Mapping Data
|
|
||||||
sample_id = str(order.rnoreg)
|
|
||||||
pid_norm = order.norm if order.norm else ""
|
|
||||||
first_name, last_name = split_patient_name(order.nama)
|
|
||||||
pid_nama = f"{last_name}^{first_name}"
|
|
||||||
room = str(getattr(order, "ruangan", "") or "RSSA MALANG").strip().upper()
|
|
||||||
room = room.replace("|", " ").replace("^", " ")[:30]
|
|
||||||
|
|
||||||
# Gender
|
|
||||||
raw_gender = str(order.rjenis).upper()
|
|
||||||
pid_gender = 'M' if 'LAKI' in raw_gender or raw_gender == 'L' else 'F'
|
|
||||||
|
|
||||||
# Test Code (Pakai Mapping yang tadi kita buat)
|
|
||||||
nama_tes_db = str(order.tes).strip() if order.tes else ""
|
|
||||||
test_code = GENEXPERT_TEST_MAPPING.get(nama_tes_db, DEFAULT_GXP_CODE) # Default: MTBRIF
|
|
||||||
|
|
||||||
# Segmen Data
|
|
||||||
# DSP/PID/ORC/OBR tergantung setting alat.
|
|
||||||
# GeneXpert standar biasanya terima format ORM di dalam DSR atau sequence PID-ORC-OBR
|
|
||||||
|
|
||||||
pid = f"PID|1||{pid_norm}||{pid_nama}|||{pid_gender}"
|
|
||||||
pv1 = f"PV1|1|I|{room}"
|
|
||||||
orc = f"ORC|NW|{sample_id}"
|
|
||||||
obr = f"OBR|1|{sample_id}||{test_code}^{nama_tes_db}^L|||{timestamp}"
|
|
||||||
|
|
||||||
# Gabungkan
|
|
||||||
return f"{msh}\r{msa}\r{qrd}\r{qrf}\r{pid}\r{pv1}\r{orc}\r{obr}\r"
|
|
||||||
|
|
||||||
else:
|
|
||||||
# Jika TIDAK ADA ORDER (Not Found)
|
|
||||||
# Kita kirim QAK (Query Acknowledge) dengan status NF (Not Found) di QRF/QAK
|
|
||||||
# Atau cukup kirim MSA|AA tapi tanpa segmen Order
|
|
||||||
return f"{msh}\r{msa}\r{qrd}\r{qrf}\r"
|
|
||||||
|
|
||||||
# ==========================================
|
# ==========================================
|
||||||
# bioMérieux MYLA TCP Server Handler
|
# bioMérieux MYLA TCP Server Handler
|
||||||
# ==========================================
|
# ==========================================
|
||||||
@@ -1753,65 +1789,6 @@ def mark_myla_order_sent(order_id):
|
|||||||
finally:
|
finally:
|
||||||
session.close()
|
session.close()
|
||||||
|
|
||||||
def create_myla_astm_order_message(order):
|
|
||||||
"""
|
|
||||||
ASTM E1394/LIS2-A2 style:
|
|
||||||
H, P, C, O, R, L
|
|
||||||
"""
|
|
||||||
pid = sanitize_astm_field(order.norm, max_len=32)
|
|
||||||
sid = sanitize_astm_field(order.rnoreg, max_len=32)
|
|
||||||
|
|
||||||
first_name, last_name = split_patient_name(sanitize_astm_field(order.nama, max_len=80))
|
|
||||||
first_name = sanitize_astm_field(first_name, uppercase=True, max_len=20)
|
|
||||||
last_name = sanitize_astm_field(last_name, uppercase=True, max_len=20)
|
|
||||||
p_name = f"{last_name}^{first_name}"
|
|
||||||
|
|
||||||
sex_raw = sanitize_astm_field(order.rjenis, uppercase=True, max_len=10)
|
|
||||||
sex = "M" if sex_raw.startswith("L") else "F"
|
|
||||||
|
|
||||||
raw_location = sanitize_astm_field(getattr(order, 'ruangan', "UT"), uppercase=True, max_len=40)
|
|
||||||
location = re.sub(r"[^A-Z0-9]", "", raw_location)[:10] or "UT"
|
|
||||||
|
|
||||||
diagnosis = sanitize_astm_field(getattr(order, 'diagnosa', "Unspecified"), max_len=60) or "Unspecified"
|
|
||||||
specimen_type = sanitize_astm_field(order.kd_spesimen, uppercase=True, max_len=20) if order.kd_spesimen else "BLOOD"
|
|
||||||
specimen_field = f"{specimen_type}^VENA^^BAIK"
|
|
||||||
test_code = sanitize_astm_field(order.tes, uppercase=True, max_len=40) or "GENERAL"
|
|
||||||
|
|
||||||
head = r"H|\^&|||LIS|||||||||P|1"
|
|
||||||
|
|
||||||
p_rec = [""] * 36
|
|
||||||
p_rec[0] = "P"
|
|
||||||
p_rec[1] = "1"
|
|
||||||
p_rec[2] = pid
|
|
||||||
p_rec[3] = pid
|
|
||||||
p_rec[5] = p_name
|
|
||||||
p_rec[8] = sex
|
|
||||||
p_rec[25] = location
|
|
||||||
pat_str = "|".join(p_rec[:36])
|
|
||||||
|
|
||||||
com_str = f"C|1|L|{diagnosis}|G"
|
|
||||||
|
|
||||||
o_rec = [""] * 32
|
|
||||||
o_rec[0] = "O"
|
|
||||||
o_rec[1] = "1"
|
|
||||||
o_rec[2] = sid
|
|
||||||
o_rec[4] = f"^^^{test_code}"
|
|
||||||
o_rec[5] = "R"
|
|
||||||
o_rec[11] = "A"
|
|
||||||
o_rec[15] = specimen_field
|
|
||||||
ord_str = "|".join(o_rec[:32])
|
|
||||||
|
|
||||||
# Placeholder R record agar struktur record lengkap H,P,C,O,R,L.
|
|
||||||
rcd_str = "R|1|^^^ORDER_STATUS|PENDING|||N|||||F"
|
|
||||||
term = "L|1|N"
|
|
||||||
|
|
||||||
message_content = f"{head}\r{pat_str}\r{com_str}\r{ord_str}\r{rcd_str}\r{term}\r"
|
|
||||||
seq = "1"
|
|
||||||
frame_body = f"{seq}{message_content}\x03"
|
|
||||||
chk = calculate_astm_checksum(frame_body)
|
|
||||||
full_frame = f"\x02{frame_body}{chk}\r\n"
|
|
||||||
return [full_frame.encode("latin-1")]
|
|
||||||
|
|
||||||
def parse_myla_astm_records(raw_message, device_name="MYLA"):
|
def parse_myla_astm_records(raw_message, device_name="MYLA"):
|
||||||
session = SessionLocal()
|
session = SessionLocal()
|
||||||
try:
|
try:
|
||||||
@@ -1939,59 +1916,6 @@ def receive_myla_astm_transmission(conn, peer_ip, first_control=ENQ):
|
|||||||
raw_message = "".join(assembled)
|
raw_message = "".join(assembled)
|
||||||
parse_myla_astm_records(raw_message, device_name=f"MYLA-{peer_ip}")
|
parse_myla_astm_records(raw_message, device_name=f"MYLA-{peer_ip}")
|
||||||
|
|
||||||
def wait_for_astm_control(conn, expected_controls, peer_ip, timeout_seconds=MYLA_CONTROL_TIMEOUT_SECONDS):
|
|
||||||
deadline = time.time() + timeout_seconds
|
|
||||||
while time.time() < deadline:
|
|
||||||
try:
|
|
||||||
conn.settimeout(max(0.1, deadline - time.time()))
|
|
||||||
b = conn.recv(1)
|
|
||||||
if not b:
|
|
||||||
return None
|
|
||||||
if b in expected_controls:
|
|
||||||
return b
|
|
||||||
if b == ENQ:
|
|
||||||
receive_myla_astm_transmission(conn, peer_ip, first_control=ENQ)
|
|
||||||
elif b == STX:
|
|
||||||
receive_myla_astm_transmission(conn, peer_ip, first_control=STX)
|
|
||||||
except socket.timeout:
|
|
||||||
continue
|
|
||||||
except Exception as e:
|
|
||||||
logging.error(f"[MYLA-ASTM] Error wait control: {e}")
|
|
||||||
print(f"[MYLA-ASTM] Error wait control: {e}")
|
|
||||||
return None
|
|
||||||
return None
|
|
||||||
|
|
||||||
def send_order_to_myla_astm(conn, order, peer_ip):
|
|
||||||
frames = create_myla_astm_order_message(order)
|
|
||||||
conn.sendall(ENQ)
|
|
||||||
hs = wait_for_astm_control(conn, {ACK}, peer_ip, timeout_seconds=MYLA_CONTROL_TIMEOUT_SECONDS)
|
|
||||||
if hs != ACK:
|
|
||||||
logging.warning(f"[MYLA-ASTM] Handshake gagal untuk rnoreg={order.rnoreg}, respon={hs}")
|
|
||||||
print(f"[MYLA-ASTM] Handshake gagal untuk rnoreg={order.rnoreg}, respon={hs}")
|
|
||||||
return False
|
|
||||||
|
|
||||||
for i, frame in enumerate(frames, start=1):
|
|
||||||
frame_ok = False
|
|
||||||
for attempt in range(1, 4):
|
|
||||||
conn.sendall(frame)
|
|
||||||
resp = wait_for_astm_control(conn, {ACK, NAK}, peer_ip, timeout_seconds=MYLA_CONTROL_TIMEOUT_SECONDS)
|
|
||||||
if resp == ACK:
|
|
||||||
frame_ok = True
|
|
||||||
break
|
|
||||||
if resp == NAK:
|
|
||||||
logging.warning(f"[MYLA-ASTM] Frame {i} NAK rnoreg={order.rnoreg}, retry={attempt}")
|
|
||||||
print(f"[MYLA-ASTM] Frame {i} NAK rnoreg={order.rnoreg}, retry={attempt}")
|
|
||||||
time.sleep(0.5)
|
|
||||||
continue
|
|
||||||
logging.warning(f"[MYLA-ASTM] Frame {i} timeout rnoreg={order.rnoreg}, retry={attempt}")
|
|
||||||
print(f"[MYLA-ASTM] Frame {i} timeout rnoreg={order.rnoreg}, retry={attempt}")
|
|
||||||
if not frame_ok:
|
|
||||||
conn.sendall(EOT)
|
|
||||||
return False
|
|
||||||
|
|
||||||
conn.sendall(EOT)
|
|
||||||
return True
|
|
||||||
|
|
||||||
def create_myla_hl7_order_message(order, msg_control_id):
|
def create_myla_hl7_order_message(order, msg_control_id):
|
||||||
timestamp = datetime.datetime.now().strftime('%Y%m%d%H%M%S')
|
timestamp = datetime.datetime.now().strftime('%Y%m%d%H%M%S')
|
||||||
order_timestamp = format_hl7_datetime(getattr(order, "rtglast", None)) or timestamp
|
order_timestamp = format_hl7_datetime(getattr(order, "rtglast", None)) or timestamp
|
||||||
@@ -2389,13 +2313,6 @@ def start_myla_inbound_server(host, port):
|
|||||||
# VITEK PARSER
|
# VITEK PARSER
|
||||||
# ==========================================
|
# ==========================================
|
||||||
|
|
||||||
def calculate_vitek_checksum(data_str):
|
|
||||||
"""
|
|
||||||
Menghitung Checksum Vitek (Sum of bytes % 256).
|
|
||||||
"""
|
|
||||||
total = sum(ord(c) for c in data_str)
|
|
||||||
return f"{total % 256:02X}"
|
|
||||||
|
|
||||||
def parse_and_save_vitek_result(raw_data, port_name="VITEK"):
|
def parse_and_save_vitek_result(raw_data, port_name="VITEK"):
|
||||||
session = SessionLocal()
|
session = SessionLocal()
|
||||||
try:
|
try:
|
||||||
|
|||||||
Reference in New Issue
Block a user