diff --git a/listener/app.py b/listener/app.py new file mode 100644 index 00000000..d570d33d --- /dev/null +++ b/listener/app.py @@ -0,0 +1,1412 @@ +from enum import Enum +from logging import config +from logging.handlers import TimedRotatingFileHandler +from queue import Queue +import re +import socket +import logging +from sqlite3 import Date +import threading +import time +import datetime +import traceback +import serial # type: ignore + +from flask import Flask, jsonify # type: ignore +from sqlalchemy import create_engine, Column, Integer, String, Boolean, Text # 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 + +# ========================================== +# 1. KONFIGURASI SISTEM +# ========================================== + +# 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 +# Mapping Flag ke IP Address GeneXpert +# Pastikan IP ini SESUAI dengan settingan "Server IP" di masing-masing alat (Client Mode) +TARGET_MAPPING = { + 'flg_gxp1': '10.10.123.73', + 'flg_gxp2': '10.10.123.74', + 'flg_gxp3': '10.10.120.75' +} +# 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 = { + # Mapping untuk IP 10.10.120.75 (Multi-Assay) + "HIV": "HIV-1_VL", # Xpert HIV-1 Viral Load XC Version 3 + "TCM TB": "MTB-RIF", # Xpert MTB-RIF Assay G4 Version 6 + "TCM TB ULTRA": "MTB-RIF_ULTRA2", # Xpert MTB-RIF Ultra Version 4 + "TCM TB XDR": "MTB-XDR", # Xpert MTB-XDR Version 1 + "COVID-19": "COV-2 2", # Xpert Xpress SARS-CoV-2 Version 2 + "HCV VL": "HCV", # Xpert HCV Viral Load Version 1 + + # Mapping untuk IP 10.10.120.74 (Khusus) + "HBV VL": "HBV", # Xpert HBV Viral Load Version 1 +} + +# Default code jika nama tes di database tidak dikenali +DEFAULT_GXP_CODE = "MTB-RIF" + +# Logging Setup +# Konfigurasi logging per hari +log_handler = TimedRotatingFileHandler( + filename="app.log", + when="midnight", + interval=1, + backupCount=7, + encoding="utf-8" +) +formatter = logging.Formatter('%(asctime)s - %(levelname)s - %(threadName)s - %(message)s') +log_handler.setFormatter(formatter) +logging.basicConfig(level=logging.INFO, handlers=[log_handler]) + +# Global Variables +app = Flask(__name__) +active_genexpert_connections = {} +connection_lock = threading.Lock() +DEVICE_CONFIGS = [ + #{ + # 'port': 'COM6', 'baud_rate': 19200, 'device_type': 'vitek', 'alat_name': 'Vitek 1', + # 'protocol': 'serial', 'flag_column': 'flg_vitek1' + #}, + { + 'port': 'COM4', 'baud_rate': 9600, 'device_type': 'vitek', 'alat_name': 'Vitek 2', + 'protocol': 'serial', 'flag_column': 'flg_vitek2' + }, + { + 'port': 'COM3', 'baud_rate': 19200, 'device_type': 'bd', 'alat_name': 'BACTEC', + 'protocol': 'serial', 'flag_column': 'flg_bd1' + }, + #BD_MGIT yang di dalam ruangan isolasi + #{ + # 'port': 'COM4', 'baud_rate': 19200, 'device_type': 'bd', 'alat_name': 'MGIT', + # 'protocol': 'serial', 'flag_column': 'flg_bd2' + #}, +] + +# Karakter kontrol standar +STX, ETX, ACK, NAK, EOT, ENQ = b'\x02', b'\x03', b'\x06', b'\x15', b'\x04', b'\x05' +RS, GS = b'\x1e', b'\x1d' # Record Separator, Group Separator +ports_lock = threading.Lock() +active_serial_ports = {} + +order_queues = {config['port']: Queue() for config in DEVICE_CONFIGS if config['protocol'] == 'serial'} + +# ========================================== +# 2. DATABASE MODEL (SESUAIKAN) +# ========================================== +DATABASE_URL = "postgresql://lismikro:lismikro@10.10.123.193:5002/lismikro" +engine = create_engine(DATABASE_URL, pool_recycle=3600) +SessionLocal = sessionmaker(bind=engine) +Base = declarative_base() + +class Sample(Base): + __tablename__ = 'samples' + id = Column(Integer, primary_key=True) + patient_name = Column(String(100)) + patient_id = Column(String(50)) + sample_id = Column(String(50), unique=True) + test_type = Column(String(100)) + result = Column(Text) + raw_message = Column(Text) + created_at = Column(SqDateTime, default=datetime.datetime.now) + +class SerialOrderQueue(Base): + __tablename__ = 'serial_order_queue' + id = Column(Integer, primary_key=True) + target_port = Column(String(50), nullable=False, index=True) + message_to_send = Column(Text, nullable=False) + status = Column(String(20), default='pending', index=True) + created_at = Column(SqDateTime, default=datetime.datetime.now) + updated_at = Column(SqDateTime, default=datetime.datetime.now, onupdate=datetime.datetime.now) + +class LisPhoenix(Base): + __tablename__ = 'lis_phoenix' + id = Column(Integer, primary_key=True) + no_id = Column(String(50)) # Patient ID + seq_no = Column(String(50), unique=True) # Isolate ID / Sample ID + rnmpas = Column(String(100)) # Patient Name + tgl_data = Column(SqDate) + rawdt = Column(Text) + organisme = Column(String(100)) + kd_orgm = Column(String(50)) + alat = Column(String(50)) + processed = Column(String(50), nullable=True) + +class LisPhoenixDtl(Base): + __tablename__ = 'lis_phoenix_dtl' + id = Column(Integer, primary_key=True) + seq_no = Column(String(50), index=True) # Isolate ID / Sample ID + kd_antibiotik = Column(String(50)) + nm_antibiotik = Column(String(100)) + keterangan = Column(String(50)) # MIC Value (e.g., <=0.5) + interpretasi = Column(String(10)) # S, I, R + no = Column(Integer) + +class PaslabOrder(Base): + __tablename__ = 'paslab' + urut = Column(Integer, primary_key=True) + rnoreg = Column(String) + nama = Column(String) + norm = Column(String) + rjenis = Column(String) + rtglast = Column(SqDateTime) + alamat = Column(String) + umur = Column(String) + namadok = Column(String) + ruangan = Column(String) + tes = Column(String) + alat = Column(String) + kd_spesimen = Column(String) + nm_spesimen = Column(String) + tgllahir = Column(SqDate) + flg_vitek1 = Column(Boolean, default=False) + flg_vitek2 = Column(Boolean, default=False) + flg_bd1 = Column(Boolean, default=False) + flg_bd2 = Column(Boolean, default=False) + flg_gxp1 = Column(Boolean, default=False) + flg_gxp2 = Column(Boolean, default=False) + flg_gxp3 = Column(Boolean, default=False) + +Base.metadata.create_all(bind=engine) + +# ========================================== +# 3. HL7 HELPER FUNCTIONS +# ========================================== + +def create_hl7_orm_message(order): + timestamp = datetime.datetime.now().strftime('%Y%m%d%H%M%S') + + # 1. Mapping Sample ID (Barcode) + # Gunakan 'rnoreg' atau 'urut' tergantung mana yang ditempel di tabung sampel + sample_id = str(order.rnoreg) + + # 2. Mapping Gender (rjenis) + # GeneXpert butuh 'M', 'F', atau 'O' + raw_gender = str(order.rjenis).upper() + if 'LAKI' in raw_gender or raw_gender == 'L': + pid_gender = 'M' + elif 'PEREMPUAN' in raw_gender or raw_gender == 'P': + pid_gender = 'F' + else: + pid_gender = 'O' + + # 3. Mapping Nama & NORM + pid_nama = order.nama if order.nama else "No Name" + pid_norm = order.norm if order.norm else "" + + # 4. Mapping Test Code + # Cek kolom 'tes'. Jika mengandung kata 'TB', set kode jadi MTB. + # Kode ini HARUS SAMA dengan "Host Test Code" di alat. + test_code = GENEXPERT_TEST_MAPPING.get(order.kd_spesimen, DEFAULT_GXP_CODE) + test_name = order.tes if order.tes else "UNKNOWN TEST" + + # Jika Anda punya tes lain (misal HIV), tambahkan if/else di sini berdasarkan order.tes + + # --- Susun Pesan HL7 --- + msh = f"MSH|^~\\&|LIS|LAB|GeneXpert|Cepheid|{timestamp}||ORM^O01|{sample_id}|P|2.3" + pid = f"PID|1||{pid_norm}||{pid_nama}|||{pid_gender}" + orc = f"ORC|NW|{sample_id}" + obr = f"OBR|1|{sample_id}||{test_code}^{test_name}^L|||{timestamp}" + + return f"{msh}\r{pid}\r{orc}\r{obr}\r" + +def parse_hl7_result(hl7_message, device_name="GeneXpert"): + session = SessionLocal() + try: + # 1. Bersihkan dan Split Pesan + # HL7 dipisahkan oleh \r (Carriage Return) + segments = hl7_message.strip().split('\r') + + # Variabel penampung + sample_id = "" # no_id + patient_id = "" # seq_no + patient_name = "" # rnmpas + result_date = None # tgl_data + results_list = [] # untuk organisme + + # Waktu default jika tidak ada di pesan + result_date = datetime.datetime.now() + + # 2. Loop setiap segmen + for segment in segments: + fields = segment.split('|') + if not fields: continue + + seg_type = fields[0] + + # --- MSH (Header) --- + if seg_type == 'MSH': + # Ambil tanggal pesan (Field 7) Format: YYYYMMDDHHMMSS + if len(fields) > 6 and fields[6]: + try: + raw_date = fields[6][:14] # Ambil 14 digit pertama + result_date = datetime.datetime.strptime(raw_date, "%Y%m%d%H%M%S") + except ValueError: + pass # Gunakan default datetime.now() jika format salah + + # --- PID (Patient ID) --- + elif seg_type == 'PID': + # PID|3 = Patient ID (seq_no) + if len(fields) > 3: + patient_id = fields[3].replace('^', '') + + # PID|5 = Patient Name (rnmpas) + if len(fields) > 5: + # Ganti caret ^ dengan spasi (Family^Name -> Family Name) + patient_name = fields[5].replace('^', ' ').strip() + + # --- OBR (Observation Request - Info Sample) --- + elif seg_type == 'OBR': + # OBR|2 atau OBR|3 biasanya berisi Sample ID / Accession No + # Kita coba ambil field 3 (Filler Order Number) dulu, kalau kosong field 2 + if len(fields) > 3 and fields[3]: + sample_id = fields[3].replace('^', '') + elif len(fields) > 2 and fields[2]: + sample_id = fields[2].replace('^', '') + + # --- OBX (Observation Result - Hasil Tes) --- + elif seg_type == 'OBX': + # OBX|3 = Test Name/Code (Misal: MTB, RIF) + # OBX|5 = Result Value (Misal: DETECTED, NOT DETECTED) + if len(fields) > 5: + test_name = fields[3].split('^')[1] if '^' in fields[3] else fields[3] + test_val = fields[5].replace('^', ' ') + + # Gabungkan nama tes dan hasil + # Contoh: "MTB: DETECTED" + results_list.append(f"{test_name}: {test_val}") + + # 3. Gabungkan semua hasil OBX menjadi satu string + if not results_list: + final_result = "No Result Found" + else: + final_result = "; ".join(results_list) + + # 4. Validasi Data Penting + if not sample_id: + print("[HL7 Parser] Sample ID tidak ditemukan. Data tidak disimpan.") + return + + print(f"[HL7 Parser] Menyimpan hasil untuk Sample: {sample_id} - {patient_name}") + + # 5. Simpan ke Database (Tabel LisPhoenix) + new_data = LisPhoenix( + no_id=sample_id, # OBR-3 + seq_no=patient_id, # PID-3 + rnmpas=patient_name, # PID-5 + tgl_data=result_date, # MSH-7 + rawdt=hl7_message, # Pesan Asli + organisme=final_result, # Gabungan OBX-5 + alat=device_name # Parameter fungsi + ) + + session.add(new_data) + session.commit() + print(f"[DB] Berhasil simpan ke LisPhoenix: {sample_id}") + + except Exception as e: + print(f"[HL7 Parser] Error menyimpan data: {e}") + session.rollback() + finally: + session.close() + +# ========================================== +# 4. NETWORK & COMMUNICATION LOGIC +# ========================================== +# ========================================== +# HL7 TCP LISTENER FOR GENEXPERT +# ========================================== +# ========================================== +# TAMBAHAN: FUNGSI TCP LISTENER (GENEXPERT) +# ========================================== +def handle_tcp_client(client_socket, addr): + """Menangani satu koneksi client GeneXpert""" + print(f"[TCP] Koneksi diterima dari {addr}") + try: + # Loop baca data dari client ini + while True: + data = client_socket.recv(4096) + if not data: + break + + # --- PROSES DATA GENEXPERT DISINI --- + # decode, parsing ASTM, save to DB + try: + msg = data.decode('latin-1', errors='ignore') + print(f"[TCP] Data dari {addr}: {msg[:50]}...") + # Panggil fungsi parser GeneXpert Anda disini + # parse_and_save_genexpert(msg) + + # Kirim ACK ASTM jika perlu (\x06) + client_socket.send(b'\x06') + except Exception as e: + print(f"[TCP] Error processing data: {e}") + + except Exception as e: + print(f"[TCP] Koneksi Error {addr}: {e}") + finally: + client_socket.close() + print(f"[TCP] Koneksi ditutup {addr}") + +def manage_tcp_server(): + """Thread Server Utama untuk GeneXpert""" + server = socket.socket(socket.AF_INET, socket.SOCK_STREAM) + # Allow reuse address agar tidak error 'Address already in use' saat restart + server.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) + + try: + server.bind((SERVER_HOST, TCP_LISTENER_PORT)) + server.listen(5) # Bisa antri 5 koneksi + print(f"[TCP-SERVER] Listening GeneXpert di port {TCP_LISTENER_PORT}...") + + while True: + # Accept koneksi baru (Blocking, tapi aman karena di thread sendiri) + client_sock, addr = server.accept() + + # Buat thread kecil untuk handle client tersebut (agar server bisa terima client lain) + client_thread = threading.Thread( + target=handle_tcp_client, + args=(client_sock, addr), + daemon=True + ) + client_thread.start() + + except Exception as e: + print(f"[TCP-SERVER] Gagal Start: {e}") + +def run_flask(): + """Wrapper untuk menjalankan Flask di Thread""" + # use_reloader=False WAJIB agar tidak membuat duplikat proses + app.run(host='0.0.0.0', port=5000, debug=False, use_reloader=False) + +def send_mllp_message(sock, hl7_msg): + """Membungkus pesan HL7 dengan MLLP (Minimal Lower Layer Protocol)""" + # Start Block: 0x0B () + # End Block: 0x1C 0x0D () + mllp_msg = b'\x0b' + hl7_msg.encode('utf-8') + b'\x1c\r' + sock.sendall(mllp_msg) + +def handle_genexpert_client(conn, addr): + ip_sender = addr[0] + logging.info(f"[CONNECT] GeneXpert terkoneksi dari IP: {ip_sender}") + + with connection_lock: + active_genexpert_connections[ip_sender] = conn + + buffer = b"" + try: + while True: + data = conn.recv(4096) + if not data: + break + + buffer += data + + # Cek apakah paket MLLP lengkap (diawali \x0b dan diakhiri \x1c\r) + # Logic sederhana: cari penutup \x1c\r + while b'\x1c\r' in buffer: + # Ekstrak pesan + start_marker = buffer.find(b'\x0b') + end_marker = buffer.find(b'\x1c\r') + + if start_marker != -1 and end_marker != -1: + raw_msg = buffer[start_marker+1 : end_marker] + hl7_str = raw_msg.decode('utf-8') + + logging.info(f"[RECV] Dari {ip_sender}: {hl7_str[:50]}...") # Log 50 karakter awal + + # Cek Tipe Pesan (MSH field 9) + if "ORU^R01" in hl7_str: + # Ini adalah Hasil + hasil = parse_hl7_result(hl7_str) + logging.info(f"[RESULT] {hasil}") + # TODO: Simpan 'hasil' ke database berdasarkan Sample ID di HL7 + + # Kirim ACK (Terima Kasih) ke Alat + # ACK dinamis + msg_control_id = hl7_str.split('|')[9] # Ambil ID pesan asli + ack_msg = f"MSH|^~\\&|LIS|LAB|GeneXpert|Cepheid|{datetime.datetime.now().strftime('%Y%m%d%H%M%S')}||ACK|{msg_control_id}|P|2.3\rMSA|AA|{msg_control_id}\r" + send_mllp_message(conn, ack_msg) + + # Hapus pesan yang sudah diproses dari buffer + buffer = buffer[end_marker+2:] + else: + break + + except Exception as e: + logging.error(f"[ERROR] Koneksi {ip_sender} terputus: {e}") + finally: + with connection_lock: + if ip_sender in active_genexpert_connections: + del active_genexpert_connections[ip_sender] + conn.close() + logging.info(f"[DISCONNECT] Koneksi dengan {ip_sender} ditutup.") + +def run_tcp_server(): + """Loop utama TCP Server""" + server = socket.socket(socket.AF_INET, socket.SOCK_STREAM) + server.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) + + try: + server.bind((SERVER_HOST, TCP_LISTENER_PORT)) + server.listen(5) # Backlog 5 koneksi + logging.info(f"TCP Listener AKTIF di port {TCP_LISTENER_PORT}. Menunggu koneksi GeneXpert...") + + while True: + conn, addr = server.accept() + # Buat thread baru untuk setiap alat yg konek + t = threading.Thread(target=handle_genexpert_client, args=(conn, addr), daemon=True) + t.start() + + except Exception as e: + logging.critical(f"Gagal menjalankan TCP Server: {e}") + +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 Python.") + return False # Indikasi gagal + + try: + # Bungkus pesan dengan MLLP (Minimal Lower Layer Protocol) standard HL7 + # Format: message + mllp_msg = f"\x0b{hl7_message}\x1c\r" + + logging.info(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}") + + return True + except Exception as e: + logging.error(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 broadcast_order_to_all_machines(hl7_message, order_id): + connected_ips = [] + with connection_lock: + connected_ips = list(active_genexpert_connections.keys()) + + if not connected_ips: + logging.warning(f"Order {order_id} GAGAL dikirim: Tidak ada GeneXpert yang terkoneksi saat ini.") + return False + + success_count = 0 + logging.info(f"Memulai broadcast Order {order_id} ke {len(connected_ips)} alat...") + + for target_ip in connected_ips: + # Panggil fungsi kirim tunggal yang sudah kita buat sebelumnya + # Note: Fungsi send_order_via_active_connection ada di kode sebelumnya + if send_order_via_active_connection(target_ip, hl7_message): + logging.info(f" -> Sukses kirim ke {target_ip}") + success_count += 1 + else: + logging.error(f" -> Gagal kirim ke {target_ip}") + + # Logika bisnis: Dianggap sukses jika minimal terkirim ke SATU alat + if success_count > 0: + return True + else: + return False +# ========================================== +# 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"): + session = SessionLocal() + try: + # --- LANGKAH 1: PEMBERSIHAN DATA --- + # Vitek sering mengirim karakter framing \x1e (RS), \x1d (GS), \x03 (ETX) + # Hapus framing characters dan newline + clean_data = raw_data.replace('\x02', '').replace('\x03', '').replace('\x1e', '').replace('\r', '').replace('\n', '') + + # Pisahkan jika ada checksum di akhir (biasanya setelah \x1d) + if '\x1d' in clean_data: + clean_data = clean_data.split('\x1d')[0] + + # --- LANGKAH 2: PARSING FIELD BERDASARKAN TAG --- + # Format Vitek: tag|value|tag|value... + fields = clean_data.split('|') + + msg_type = fields[0] # mtmpr (Order) atau mtrsl (Result) + + # Inisialisasi variabel + sample_id = None # ci (Accession Number) + patient_id = "" # pi + patient_name = "" # pn + organism_name = "" # o2 + antibiotics = [] # List untuk menyimpan hasil AB + result_date = datetime.datetime.now() + + # Penanda apakah kita sedang membaca blok antibiotik + current_ab_name = "" + current_ab_mic = "" + current_ab_int = "" + + # Loop setiap field untuk mencari Tag + for field in fields: + if not field: continue + + # --- HEADER INFO --- + if field.startswith("pi") and len(field) > 2: # Patient ID + patient_id = field[2:].strip() + elif field.startswith("pn") and len(field) > 2: # Patient Name + patient_name = field[2:].strip() + elif field.startswith("ci") and len(field) > 2: # Case ID / No Reg (KUNCI UTAMA) + # Format Vitek kadang 25-035007, kita ambil angkanya saja atau sesuai format LIS + sample_id = field[2:].strip() + + # --- ORGANISME (Bakteri) --- + elif field.startswith("o2") and len(field) > 2: + organism_name = field[2:].strip() + + # --- ANTIBIOTIK (Looping Block) --- + # Penanda blok antibiotik baru dimulai dengan tag 'ra' + elif field == "ra": + # Simpan antibiotik sebelumnya jika ada + if current_ab_name: + ab_str = f"{current_ab_name} {current_ab_mic} ({current_ab_int})" + antibiotics.append(ab_str) + # Reset temp vars + current_ab_name = "" + current_ab_mic = "" + current_ab_int = "" + + # Detail Antibiotik + elif field.startswith("a2") and len(field) > 2: # Nama Antibiotik + current_ab_name = field[2:].strip() + elif field.startswith("a3") and len(field) > 2: # Nilai MIC + current_ab_mic = field[2:].strip() + elif field.startswith("a4") and len(field) > 2: # Interpretasi (S/I/R) + current_ab_int = field[2:].strip() + elif field.startswith("an") and len(field) > 2: # Interpretasi Alternatif + if not current_ab_int: current_ab_int = field[2:].strip() + + # Jangan lupa simpan antibiotik terakhir setelah loop selesai + if current_ab_name: + ab_str = f"{current_ab_name} {current_ab_mic} ({current_ab_int})" + antibiotics.append(ab_str) + + # --- LANGKAH 3: FORMAT FINAL & SAVE --- + # Kita hanya memproses jika tipe pesan adalah 'mtrsl' (Result) dan ada Sample ID + if msg_type == 'mtrsl' and sample_id: + + # Gabungkan hasil: Organisme | AB1, AB2, AB3... + final_res_string = organism_name + if antibiotics: + final_res_string += " | " + ", ".join(antibiotics) + + # Jika tidak ada organisme tapi tipe result, mungkin Negative? + if not final_res_string and "neg" in raw_data.lower(): + final_res_string = "NEGATIVE / NO GROWTH" + + logging.info(f"[{port_name}] Save DB -> ID: {sample_id}, Pasien: {patient_name}, Bakteri: {organism_name}") + + new_entry = LisPhoenix( + no_id=sample_id, # Menggunakan tag 'ci' + seq_no=patient_id, # Menggunakan tag 'pi' + rnmpas=patient_name, # Menggunakan tag 'pn' + tgl_data=result_date, + rawdt=raw_data, + organisme=final_res_string, # Hasil gabungan + alat=port_name + ) + session.add(new_entry) + session.commit() + else: + if msg_type != 'mtrsl': + logging.info(f"[{port_name}] Mengabaikan pesan tipe {msg_type}") + else: + logging.warning(f"[{port_name}] Data tidak lengkap (No Sample ID). Raw: {raw_data[:50]}...") + + except Exception as e: + logging.error(f"Error Parsing Vitek: {e}") + session.rollback() + finally: + session.close() + +def create_vitek_order_message(order): + """ + Membuat Frame Order Vitek sesuai Manual Ref 514937. + Format: [DATA] [CS] + """ + # --- 1. ISI PESAN (CONTENT) --- + # Field Delimiter menggunakan Pipe '|' + pid = str(order.norm).strip() if order.norm else "" + sid = str(order.rnoreg).strip() if order.rnoreg else "" + p_name = str(order.nama).strip().replace('^', ' ').upper()[:20] if order.nama else "NO NAME" + specimen = str(order.kd_spesimen).upper() if order.kd_spesimen else "BLOOD" + + now = datetime.datetime.now() + date_str = now.strftime("%m/%d/%Y") + time_str = now.strftime("%H:%M") + + # Struktur mtmpr sesuai Table 1-3 & Contoh Manual + # Penting: Tidak ada Sequence Number '1' di dalam data + content_body = ( + f"mtmpr|pi{pid}|pn{p_name}" + f"|si|ss{specimen}" + f"|s1{date_str}|s2{time_str}" + f"|ci{sid}|t11|zz" + ) + + # --- 2. FRAMING & CHECKSUM (Section 2.2 & 2.5) --- + # Frame dimulai dengan STX, lalu Record dimulai dengan RS + STX = b'\x02' + RS = b'\x1e' + GS = b'\x1d' + ETX = b'\x03' + CRLF = b'\r\n' + + # Data yang dihitung Checksumnya: [RS] + [Body] + [GS] + # Manual Hal 2-10: "calculated by adding the value of all characters beginning with the first ... and ending with the " + payload_for_checksum = RS + content_body.encode('latin-1') + GS + + # Hitung Checksum + total_sum = sum(payload_for_checksum) + chk_val = total_sum % 256 + checksum_str = f"{chk_val:02X}".encode('latin-1') # Hex 2 digit uppercase + + # Rakit Frame Utuh + # [PAYLOAD+GS] [CHECKSUM] + # Perhatikan: Payload di atas sudah mengandung RS dan GS + full_frame = STX + payload_for_checksum + checksum_str + ETX + CRLF + + return [full_frame] + +def manage_vitek_port(config): + port_name = config['port'] + flag_col = config.get('flag_column') + alat_name = config.get('alat_name', 'VITEK') + logging.info(f"[{port_name}] Membuka port untuk alat {alat_name}...") + + # Buffer untuk menampung pecahan data + rx_buffer = b"" + + try: + with serial.Serial( + port=port_name, + baudrate=config['baud_rate'], + timeout=1 + ) as ser: + + while True: + has_activity = False # Penanda agar kita sleep kalau sepi + + # ========================================== + # PHASE 1: LISTENING (PRIORITAS UTAMA) + # ========================================== + try: + if ser.in_waiting > 0: + has_activity = True + data_chunk = ser.read(ser.in_waiting or 1024) + + if data_chunk: + # 1. Handle Handshake Awal (ENQ) - Alat minta izin kirim + if b'\x05' in data_chunk: + logging.info(f"[{port_name}] Terima ENQ. Kirim ACK.") + ser.write(b'\x06') # ACK + rx_buffer = b"" # Reset buffer bytes + continue + + # 2. Tampung Data + rx_buffer += data_chunk + + # 3. Cek apakah Frame Selesai? + # Vitek biasanya mengirim per baris diakhiri LF (\n) atau CR (\r) + # Atau Frame diakhiri ETX (\x03) + # Kita kirim ACK setiap kali ada tanda akhir baris/frame agar alat lanjut kirim + if b'\x03' in data_chunk or b'\n' in data_chunk or b'\r' in data_chunk: + ser.write(b'\x06') # ACK data yang baru masuk + + # 4. Handle Akhir Transmisi (EOT) - Selesai + if b'\x04' in data_chunk: + logging.info(f"[{port_name}] Terima EOT. Memproses data...") + + # Decode saat data sudah lengkap (gunakan latin-1) + try: + full_str = rx_buffer.decode('latin-1', errors='ignore') + parse_and_save_vitek_result(full_str, alat_name) + except Exception as e: + logging.error(f"[{port_name}] Error Decoding/Parsing: {e}") + + rx_buffer = b"" # Bersihkan buffer + continue + + except Exception as e: + logging.error(f"[{port_name}] Error Reading: {e}") + rx_buffer = b"" # Reset jika error parah + + + # ========================================== + # PHASE 2: SENDING ORDER (JIKA BUFFER KOSONG) + # ========================================== + # Kita hanya kirim order jika sedang tidak menerima data (buffer kosong) + if not rx_buffer and flag_col: + try: + session = SessionLocal() + # Cari order yang belum dikirim + pending_order = session.query(PaslabOrder).filter( + getattr(PaslabOrder, flag_col) == False + ).first() + + if pending_order: + has_activity = True # Jangan sleep lama-lama + logging.info(f"[{port_name}] Menemukan Order: {pending_order.rnoreg}. Memulai Handshake...") + + # --- STEP 1: HANDSHAKE (ENQ) --- + ser.reset_input_buffer() + ser.write(b'\x05') + time.sleep(0.5) + + # Baca balasan (tunggu ACK) + ack_response = ser.read(1) + + if ack_response == b'\x06': + logging.info(f"[{port_name}] Handshake dengan {alat_name}Sukses. Kirim Data...") + time.sleep(1.0) # Jeda aman + + frames = create_vitek_order_message(pending_order) + all_success = True + + # --- STEP 2: SEND FRAMES --- + for i, frame in enumerate(frames): + retry = 0 + frame_sent = False + while retry < 3: + logging.info(f"[{port_name}] Kirim Frame (Try {retry+1}): {frame}") # Log apa yang dikirim + ser.write(frame) + + # Tunggu ACK + start_wait = time.time() + response = None + + while time.time() - start_wait < 3: # Timeout 3 detik + if ser.in_waiting: + response = ser.read(1) + break + + if response == b'\x06': # ACK + logging.info(f"[{port_name}] Frame ACK (OK).") + frame_sent = True + break + elif response == b'\x15': # NAK + logging.warning(f"[{port_name}] Dibalas NAK (Checksum/Format Salah).") + time.sleep(1) + retry += 1 + elif response is None: # TIMEOUT + logging.warning(f"[{port_name}] Timeout (Alat tidak membalas). Cek Framing/Kabel.") + retry += 1 + else: # Respon Aneh + logging.warning(f"[{port_name}] Respon aneh: {response}") + retry += 1 + + if not frame_sent: + all_success = False + break + + # --- STEP 3: TERMINATE (EOT) --- + ser.write(b'\x04') + + if all_success: + logging.info(f"[{port_name}] Order {pending_order.rnoreg} SELESAI.") + setattr(pending_order, flag_col, True) + session.commit() + else: + logging.error(f"[{port_name}] Gagal kirim order {pending_order.rnoreg}.") + + else: + # Handshake gagal (Alat sibuk/Mati) + # logging.debug(f"[{port_name}] Alat sibuk/tidak balas ENQ.") + pass + + session.close() + + except Exception as e: + logging.error(f"[{port_name}] Error Sending Logic: {e}") + if 'session' in locals(): session.close() + + # ========================================== + # PHASE 3: IDLE MANAGEMENT + # ========================================== + # Jika tidak ada data masuk dan tidak ada order keluar, tidur sebentar + # Ini penting agar CPU tidak 100% dan DB tidak jebol + if not has_activity: + time.sleep(1.0) + + except Exception as e: + logging.critical(f"[{port_name}] Gagal connect Serial: {e}") + time.sleep(5) + +# ========================================== +# BECTON DICKINSON (BD) BACTEC PARSER +# ========================================== + +def parse_and_save_bd_result(raw_data, port_name="BD Bactec"): + session = SessionLocal() + try: + # --- LANGKAH 1: REASSEMBLY (JAHIT FRAME) --- + raw_frames = raw_data.split('\x02') + full_content = "" + + for frame in raw_frames: + if not frame: continue + content_chunk = frame + # Hapus Checksum & ETX/ETB + if '\x03' in frame: content_chunk = frame.split('\x03')[0] + elif '\x17' in frame: content_chunk = frame.split('\x17')[0] + + # Hapus Sequence Number di awal + if content_chunk and content_chunk[0].isdigit(): + content_chunk = content_chunk[1:] + + full_content += content_chunk + + # --- LANGKAH 2: PARSING FIELD --- + lines = full_content.split('\r') + + sample_id = None + patient_id = "" + patient_name = "" + specimen_type = "" + result_val = "" + result_date = datetime.datetime.now() + + for line in lines: + line = line.strip() + if not line: continue + fields = line.split('|') + record_type = fields[0] + + # --- PATIENT (P) --- + if record_type == 'P': + # P|Seq|?|PID||Name + if len(fields) > 3: patient_id = fields[3].strip() + if len(fields) > 5: patient_name = fields[5].replace('^', ' ').strip() + + # --- ORDER (O) --- + elif record_type == 'O': + # O|Seq|SampleID|...|...|...|...|...|...|...|...|...|...|...|Specimen + if len(fields) > 2: + sample_id = fields[2].replace('^', '').strip() + + # Ambil Spesimen (Biasanya di index 15 / Field 16) + if len(fields) > 15: + specimen_type = fields[15].strip() + + # --- RESULT (R) --- + elif record_type == 'R': + # R|Seq|Test|Result|...|...|...|...|...|...|StartDate|EndDate + if len(fields) > 3: + raw_res = fields[3].strip() + # Bersihkan hasil + if "NEGATIVE" in raw_res: result_val = "NEGATIVE" + elif "POSITIVE" in raw_res: result_val = "POSITIVE" + else: result_val = raw_res.split('^')[0] # Ambil kode depan saja + + # Ambil Tanggal Hasil Selesai (Index 12 / Field 13) + if len(fields) > 12 and len(fields[12]) >= 14: + try: + # Format BD: YYYYMMDDHHMMSS (20251006100257) + res_dt_str = fields[12][:14] + result_date = datetime.datetime.strptime(res_dt_str, "%Y%m%d%H%M%S") + except: + pass # Gunakan default jika gagal parse + + # --- LANGKAH 3: FORMAT FINAL & SAVE --- + if sample_id and result_val: + + final_res_string = result_val + if specimen_type: + final_res_string += f" ({specimen_type})" + + logging.info(f"[{port_name}] Save DB -> ID: {sample_id}, Pasien: {patient_name}, Hasil: {final_res_string}") + + new_entry = LisPhoenix( + no_id=sample_id, + seq_no=patient_id, + rnmpas=patient_name, + tgl_data=result_date, + rawdt=raw_data, + organisme=final_res_string, + alat=port_name + ) + session.add(new_entry) + session.commit() + else: + logging.warning(f"[{port_name}] Data tidak lengkap. ID: {sample_id}, Res: {result_val}") + + except Exception as e: + logging.error(f"Error Parsing BD: {e}") + session.rollback() + finally: + session.close() + +def calculate_astm_checksum(frame_content): + data_bytes = frame_content.encode('latin-1') + checksum = sum(data_bytes) % 256 + return f"{checksum:02X}" + +def create_astm_order_message(order): + """ + Membuat 1 Frame ASTM Single Block (H, P, O, L) dengan Mapping Index PRESISI. + Menghindari pergeseran kolom (shifting error). + """ + # --- 1. PERSIAPAN DATA --- + pid = str(order.norm).strip() if order.norm else "" + sid = str(order.rnoreg).strip() if order.rnoreg else "" + # Nama Pasien: Ganti karakter topi '^' dengan spasi agar tidak merusak format + p_name = str(order.nama).strip().replace('^', ' ')[:20] if order.nama else "No Name" + sex = "M" if str(order.rjenis).upper().startswith("L") else "F" + + # Lokasi / Ruangan (Field 26) + location = getattr(order, 'ruangan', "RSSA Malang") + if not location: location = "RSSA Malang" + + # Diagnosis (Clinical Info - Field 14 di ASTM standar atau 13 di beberapa varian) + # Kita pasang di Index 13 (Field 14) agar aman + diagnosis = getattr(order, 'diagnosa', "Unspecified") + if not diagnosis: diagnosis = "Unspecified" + + # Specimen Info (Field 16) + # Format: SpecimenType^BodySite^Container^Condition + specimen_type = str(order.kd_spesimen).upper() if order.kd_spesimen else "BLOOD" + body_site = "VENA" # Site + condition = "Baik" # Condition + specimen_field = f"{specimen_type}^{body_site}^^{condition}" + + # --- 2. KONSTRUKSI RECORD DENGAN INDEX PASTI --- + + # --- RECORD HEADER (H) --- + head = r"H|\^&|||LIS||||||||1" + + # --- RECORD PATIENT (P) --- + # Kita buat array kosong sebanyak 30 kolom dulu + p_rec = [""] * 30 + p_rec[0] = "P" # Field 1: Record Type + p_rec[1] = "1" # Field 2: Sequence + p_rec[2] = pid # Field 3: Patient ID (Practice) + p_rec[3] = pid # Field 4: Lab ID (Kosong) + p_rec[4] = "" # Field 5: ID 3 (Kosong) + p_rec[5] = p_name # Field 6: Patient Name (Index 5) <--- SEBELUMNYA SALAH DISINI + p_rec[7] = "" # Field 8: Birthdate + p_rec[8] = sex # Field 9: Sex (Index 8) + # ... Field 10-25 biarkan kosong ... + p_rec[25] = location # Field 26: Location (Index 25) + + # Potong array sampai index 26 saja (sisanya buang) lalu gabung + pat_str = "|".join(p_rec[:26]) + + # --- RECORD ORDER (O) --- + o_rec = [""] * 30 + o_rec[0] = "O" # Field 1 + o_rec[1] = "1" # Field 2 + o_rec[2] = sid # Field 3: Sample ID + o_rec[3] = "" # Field 4: Instrument Specimen ID + o_rec[4] = "^^^BD_BACTEC" # Field 5: Universal Test ID + o_rec[5] = "R" # Field 6: Priority + # ... Field 7-11 ... + o_rec[11] = "A" # Field 12: Action Code (A=Add, N=New) (Index 11) + o_rec[12] = diagnosis # Field 13: Clinical Info / Diagnosis (Index 12) + # ... Field 14-15 ... + o_rec[15] = specimen_field # Field 16: Specimen Source (Index 15) + + # Potong array sampai index 16 (atau lebih jika BD butuh field belakang) + # Kita ambil aman sampai field 20 + ord_str = "|".join(o_rec[:20]) + + # --- RECORD TERMINATOR (L) --- + term = "L|1|N" + + # --- 3. GABUNG FRAME --- + # Gunakan \r (Carriage Return) sebagai pemisah record + message_content = f"{head}\r{pat_str}\r{ord_str}\r{term}" + + # Sequence Frame = 1 + seq = "1" + + # Isi Frame: [Seq] [Data] [ETX] + frame_body = f"{seq}{message_content}\x03" + + # Hitung Checksum + chk = calculate_astm_checksum(frame_body) + + # Full Frame + full_frame = f"\x02{frame_body}{chk}\r\n" + + return [full_frame.encode('latin-1')] + +def manage_bd_port(config): + port_name = config['port'] + flag_col = config.get('flag_column') + alat_name = config.get('alat_name', 'BD') + logging.info(f"[{port_name}] Membuka port untuk alat {alat_name}...") + + # Buffer untuk menampung pecahan data + rx_buffer = b"" + + try: + with serial.Serial( + port=port_name, + baudrate=config['baud_rate'], + timeout=1 + ) as ser: + + while True: + has_activity = False # Penanda agar kita sleep kalau sepi + + # ========================================== + # PHASE 1: LISTENING (PRIORITAS UTAMA) + # ========================================== + try: + if ser.in_waiting > 0: + has_activity = True + data_chunk = ser.read(ser.in_waiting or 1024) + + if data_chunk: + # 1. Handle Handshake Awal (ENQ) + if b'\x05' in data_chunk: + logging.info(f"[{port_name}] Terima ENQ (Alat mau kirim data). Kirim ACK.") + ser.write(b'\x06') # ACK + rx_buffer = "" # Reset buffer untuk data baru + continue + + # 2. Handle Akhir Transmisi (EOT) + if b'\x04' in data_chunk: + logging.info(f"[{port_name}] Terima EOT (Selesai). Memproses data...") + # Proses semua data yang terkumpul di buffer + parse_and_save_bd_result(rx_buffer, alat_name) + rx_buffer = "" # Kosongkan buffer setelah save + continue + + # 3. Handle Data Frame Normal (STX ... ETX) + # ASTM butuh kita balas ACK setiap kali dikirim Frame + # Frame biasanya diawali STX (\x02) + if b'\x02' in data_chunk: + # Simpan ke buffer (decode dulu ke string) + try: + decoded_str = data_chunk.decode('utf-8', errors='ignore') + rx_buffer += decoded_str + + # WAJIB: Kirim ACK agar alat lanjut kirim baris berikutnya + ser.write(b'\x06') + # logging.debug(f"[{port_name}] Frame diterima, ACK dikirim.") + except Exception as e: + logging.error(f"Error decode frame: {e}") + + continue # Loop lagi untuk ambil sisa data + + except Exception as e: + logging.error(f"[{port_name}] Error Reading: {e}") + rx_buffer = b"" # Reset jika error parah + + + # ========================================== + # PHASE 2: SENDING ORDER (JIKA BUFFER KOSONG) + # ========================================== + # Kita hanya kirim order jika sedang tidak menerima data (buffer kosong) + if not rx_buffer and flag_col: + try: + session = SessionLocal() + # Cari order yang belum dikirim + pending_order = session.query(PaslabOrder).filter( + getattr(PaslabOrder, flag_col) == False + ).first() + + if pending_order: + has_activity = True # Jangan sleep lama-lama + logging.info(f"[{port_name}] Menemukan Order: {pending_order.rnoreg}. Memulai Handshake...") + + # --- STEP 1: HANDSHAKE (ENQ) --- + ser.reset_input_buffer() + ser.write(b'\x05') + time.sleep(0.5) + + ack_response = ser.read(1) + + if ack_response == b'\x06': + logging.info(f"[{port_name}] Handshake Sukses (Dapat ACK). Menunggu alat siap...") + + # --- PERBAIKAN 1: BERI JEDA SETELAH HANDSHAKE --- + # Mesin butuh napas sebelum terima data panjang + time.sleep(1.5) # Jeda 1.5 detik + + frames = create_astm_order_message(pending_order) + all_frames_sent = True + + # --------------------------------------------------- + # STEP 2: SEND FRAMES WITH RETRY + # --------------------------------------------------- + for i, frame in enumerate(frames): + retry_count = 0 + max_retries = 3 + frame_success = False + + while retry_count < max_retries: + ser.reset_input_buffer() + + logging.info(f"[{port_name}] Kirim Frame {i+1} (Percobaan {retry_count+1})...") + ser.write(frame) + + # Tunggu ACK + original_timeout = ser.timeout + ser.timeout = 3 # Beri waktu agak lama (3 detik) untuk alat memproses data + frame_ack = ser.read(1) + ser.timeout = original_timeout + + if frame_ack == b'\x06': # ACK (Sukses) + frame_success = True + logging.info(f"[{port_name}] Frame {i+1} ACK diterima.") + # Beri jeda dikit sebelum kirim frame berikutnya + time.sleep(0.2) + break + + elif frame_ack == b'\x15': # NAK (Ditolak - Checksum Salah) + logging.warning(f"[{port_name}] Frame ditolak (NAK). Checksum mungkin salah.") + time.sleep(2) # Tunggu 2 detik + retry_count += 1 + + elif not frame_ack: # Timeout (Sepi) + # --- PERBAIKAN 2: BERI JEDA SAAT TIMEOUT --- + logging.warning(f"[{port_name}] Timeout (Alat diam). Menunggu sebelum retry...") + time.sleep(2) # Tunggu 2 detik agar alat recover + retry_count += 1 + + else: + logging.warning(f"[{port_name}] Respon aneh: {frame_ack}") + time.sleep(1) + retry_count += 1 + + if not frame_success: + logging.error(f"[{port_name}] Gagal kirim frame ke-{i+1}. Batal.") + all_frames_sent = False + break + # --------------------------------------------------- + # STEP 3: FINALIZE + # --------------------------------------------------- + if all_frames_sent: + ser.write(b'\x04') # EOT (End of Transmission) + logging.info(f"[{port_name}] Order {pending_order.rnoreg} SUKSES Terkirim.") + + # Update Database + setattr(pending_order, flag_col, True) + session.commit() + else: + ser.write(b'\x04') # EOT (Putus paksa karena error) + logging.error(f"[{port_name}] Pengiriman Order GAGAL.") + + else: + # Jika Handshake gagal (Dibalas NAK, atau Timeout) + logging.warning(f"[{port_name}] Handshake Gagal. Respon alat: {ack_response}") + # Jangan update flag DB, biarkan coba lagi nanti + + + session.close() + + except Exception as e: + logging.error(f"[{port_name}] Error Sending Logic: {e}") + if 'session' in locals(): session.close() + + # ========================================== + # PHASE 3: IDLE MANAGEMENT + # ========================================== + # Jika tidak ada data masuk dan tidak ada order keluar, tidur sebentar + # Ini penting agar CPU tidak 100% dan DB tidak jebol + if not has_activity: + time.sleep(1.0) + + except Exception as e: + logging.critical(f"[{port_name}] Gagal connect Serial: {e}") + time.sleep(5) + +# ========================================== +# 5. ORDER POLLER (BROADCASTER) +# ========================================== + +def order_poller(stop_event): + """Looping cek DB dan kirim order ke SEMUA alat""" + logging.info("Order Poller Berjalan...") + + while not stop_event.is_set(): + session = SessionLocal() + try: + # 1. Ambil order yang belum dikirim (flag_genexpert = False) + + orders = session.query(PaslabOrder).filter( + (PaslabOrder.flg_gxp1 == False) | + (PaslabOrder.flg_gxp2 == False) | + (PaslabOrder.flg_gxp3 == False) + ).all() + # Update daftar koneksi aktif (Thread safe) + current_connections = {} + with connection_lock: + current_connections = active_genexpert_connections.copy() + + if not orders: + time.sleep(5) + continue + for order in orders: + # Generate pesan HL7 sekali untuk order ini + hl7_msg = create_hl7_orm_message(order) + + # --- LOGIKA KIRIM BERDASARKAN FLAG --- + + # 1. Cek Target GXP 1 + if order.flg_gxp1 == False: + target_ip = TARGET_MAPPING['flg_gxp1'] # 10.10.123.73 + if target_ip in current_connections: + # Ada koneksi dari alat 1, kirim! + conn = current_connections[target_ip] + try: + send_mllp_message(conn, hl7_msg) + order.flg_gxp1 = True # Update Flag + logging.info(f"[SENT] Order {order.rnoreg} dikirim ke GXP-1 ({target_ip})") + except Exception as e: + logging.error(f"[FAIL] Gagal kirim ke GXP-1: {e}") + else: + # Alat 1 belum connect, biarkan False (pending) + pass + + # 2. Cek Target GXP 2 + if order.flg_gxp2 == False: + target_ip = TARGET_MAPPING['flg_gxp2'] # 10.10.123.74 + if target_ip in current_connections: + conn = current_connections[target_ip] + try: + send_mllp_message(conn, hl7_msg) + order.flg_gxp2 = True + logging.info(f"[SENT] Order {order.rnoreg} dikirim ke GXP-2 ({target_ip})") + except Exception as e: + logging.error(f"[FAIL] Gagal kirim ke GXP-2: {e}") + + # 3. Cek Target GXP 3 + if order.flg_gxp3 == False: + target_ip = TARGET_MAPPING['flg_gxp3'] # 10.10.120.75 + if target_ip in current_connections: + conn = current_connections[target_ip] + try: + send_mllp_message(conn, hl7_msg) + order.flg_gxp3 = True + logging.info(f"[SENT] Order {order.rnoreg} dikirim ke GXP-3 ({target_ip})") + except Exception as e: + logging.error(f"[FAIL] Gagal kirim ke GXP-3: {e}") + + # Commit perubahan flag ke database + session.commit() + except Exception as e: + logging.error(f"Error Poller: {e}") + session.rollback() + finally: + session.close() + + time.sleep(5) # Cek DB setiap 5 detik + +def manage_serial_port(config): + """Fungsi router yang memilih manajer yang tepat berdasarkan tipe alat.""" + device_type = config.get('device_type') + if device_type == 'vitek': + manage_vitek_port(config) + elif device_type in ['bd_mgit', 'bd_bactec', 'bd']: + manage_bd_port(config) + else: + logging.error(f"Tipe alat tidak diketahui: '{device_type}' untuk port {config.get('port')}. Thread dihentikan.") + +# ========================================== +# 6. FLASK API (Opsional) +# ========================================== +@app.route('/status', methods=['GET']) +def status(): + with connection_lock: + return jsonify({ + "active_connections": list(active_genexpert_connections.keys()), + "total_connected": len(active_genexpert_connections) + }) + +# ========================================== +# 7. MAIN EXECUTION +# ========================================== +if __name__ == "__main__": + print("--- MEMULAI LIS INTERFACE SYSTEM ---") + + # List untuk menampung semua thread agar bisa dimonitor + all_threads = [] + stop_event = threading.Event() + # 1. Start Thread Order Poller (Pengecek Order Baru di DB) + t_poller = threading.Thread(target=order_poller, args=(stop_event,), name="OrderPoller", daemon=True) + t_poller.start() + all_threads.append(t_poller) + + # 2. Start Thread Serial Manager (Vitek & BD) + for config in DEVICE_CONFIGS: + if config['protocol'] == 'serial': + t_serial = threading.Thread( + target=manage_serial_port, + args=(config,), + name=f"Manager-{config['device_type']}-{config['port']}", # Nama lebih jelas + daemon=True + ) + t_serial.start() + all_threads.append(t_serial) + + # 3. Start Thread TCP Server (GeneXpert) + t_tcp = threading.Thread(target=manage_tcp_server, name="Manager-TCP-GeneXpert", daemon=True) + t_tcp.start() + all_threads.append(t_tcp) + + # 4. Start Thread Flask Web Server (API) + # PENTING: Flask juga harus di thread agar tidak memblokir monitoring + t_flask = threading.Thread(target=run_flask, name="WebServer-Flask", daemon=True) + t_flask.start() + all_threads.append(t_flask) + + # 5. LOOP UTAMA (Keep-Alive & Monitoring) + # Ini sekarang bisa berjalan karena Flask sudah dipindah ke thread + try: + while True: + print(f"--- Monitoring {len(all_threads)} Threads ---") + + # Cek status setiap thread + alive_count = 0 + for t in all_threads: + if t.is_alive(): + alive_count += 1 + else: + print(f"!!! THREAD MATI: {t.name} !!!") + # Disini Anda bisa menambahkan logika restart thread jika mati + + # Jika semua mati, exit (atau restart service) + if alive_count == 0: + logging.critical("Semua thread mati. System Shutdown.") + break + + time.sleep(10) # Cek setiap 10 detik (Hemat CPU) + + except KeyboardInterrupt: + print("Mematikan Service (Ctrl+C)...") + stop_event.set() + # Thread daemon akan mati otomatis saat main exit \ No newline at end of file diff --git a/listener/requirements.txt b/listener/requirements.txt new file mode 100644 index 00000000..55263398 --- /dev/null +++ b/listener/requirements.txt @@ -0,0 +1,15 @@ +blinker==1.9.0 +click==8.1.8 +Flask==3.1.2 +greenlet==3.2.4 +hl7==0.4.5 +importlib_metadata==8.7.0 +itsdangerous==2.2.0 +Jinja2==3.1.6 +MarkupSafe==3.0.2 +psycopg2-binary==2.9.10 +pyserial==3.5 +SQLAlchemy==2.0.43 +typing_extensions==4.14.1 +Werkzeug==3.1.3 +zipp==3.23.0