#!/usr/bin/env python3 """HTTP-driven python-smpplib ESME: binds named sessions against the target SMPP server and performs one action per request, answering with the result as JSON. Kept alive as one process so a session survives across requests, the way jsmpp's Java driver does (see AGENTS.md).""" import json import socket import sys import threading import time from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer from urllib.parse import parse_qs, urlparse import smpplib.client import smpplib.consts import smpplib.exceptions import smpplib.gsm import smpplib.smpp TARGET_HOST = sys.argv[1] if len(sys.argv) > 1 else "node" TARGET_PORT = int(sys.argv[2]) if len(sys.argv) > 2 else 2775 GSM_TABLE = smpplib.gsm.GSM_CHARACTER_TABLE MODES = { "receiver": "bind_receiver", "transceiver": "bind_transceiver", "transmitter": "bind_transmitter", } def as_text(value): return value.decode() if isinstance(value, bytes) else value def gsm_decode(data: bytes) -> str: """Inverse of smpplib's own gsm_encode(), through its own (vendor-specific) table - used to decode what this driver receives, so a mismatch against what was sent is smpplib's own table, not a guess at the real GSM 03.38 one.""" chars = [] i = 0 while i < len(data): byte = data[i] if byte == 0x1B and i + 1 < len(data): chars.append(GSM_TABLE[0x80 + data[i + 1]]) i += 2 else: chars.append(GSM_TABLE[byte]) i += 1 return "".join(chars) def decode_body(data: bytes, data_coding: int) -> str: if data_coding in (0, 1): return gsm_decode(data) if data_coding == 3: return data.decode("latin-1") if data_coding == 8: whole = len(data) - (len(data) % 2) return data[:whole].decode("utf-16-be") return data.hex() class Session: def __init__(self, client): self.client = client self.send_lock = threading.Lock() self.received = [] self.acks = {} self.ack_events = {} self.reader_thread = None self.reader_running = False self.reader_error = None def wait_ack(self, sequence, budget=8.0): event = self.ack_events.setdefault(sequence, threading.Event()) event.wait(budget) return self.acks.get(sequence) SESSIONS = {} SESSIONS_LOCK = threading.Lock() def session_for(name): with SESSIONS_LOCK: return SESSIONS[name] def do_bind(body): name = body["name"] mode = body["mode"] timeout_secs = float(body.get("timeoutSecs", 5)) client = smpplib.client.Client(TARGET_HOST, TARGET_PORT, timeout=timeout_secs, allow_unknown_opt_params=True) client.connect() kwargs = {"system_id": body["systemId"], "password": body["password"]} if body.get("interfaceVersion") is not None: kwargs["interface_version"] = int(body["interfaceVersion"]) getattr(client, MODES[mode])(**kwargs) session = Session(client) def on_received(pdu, **_kwargs): data = pdu.short_message or b"" session.received.append({ "dataCoding": pdu.data_coding, "esmClass": pdu.esm_class, "from": as_text(pdu.source_addr), "hex": data.hex(), "text": decode_body(data, pdu.data_coding), "to": as_text(pdu.destination_addr), }) return smpplib.consts.SMPP_ESME_ROK def on_sent(pdu, **_kwargs): session.acks[pdu.sequence] = { "messageId": as_text(getattr(pdu, "message_id", None)), "status": int(pdu.status), } session.ack_events.setdefault(pdu.sequence, threading.Event()).set() def on_error_pdu(pdu): # Overrides the default handler, which raises: a refusing status must reach # message_sent_handler like any other response, not tear down the read loop. if pdu.command == "submit_sm_resp": on_sent(pdu) client.set_message_received_handler(on_received) client.set_message_sent_handler(on_sent) client.set_error_pdu_handler(on_error_pdu) with SESSIONS_LOCK: SESSIONS[name] = session return {"ok": True} def do_start_reader(body): session = session_for(body["name"]) auto_send_enquire_link = bool(body.get("autoSendEnquireLink", True)) if session.reader_running: return {"ok": True} def run(): session.reader_running = True try: while True: session.client.read_once(auto_send_enquire_link=auto_send_enquire_link) except Exception as exc: # noqa: BLE001 - recorded, not raised: this is a driver thread session.reader_error = f"{type(exc).__name__}: {exc}" finally: session.reader_running = False session.reader_thread = threading.Thread(target=run, daemon=True) session.reader_thread.start() return {"ok": True} def encode_body(text, data_coding): if data_coding == 0: return smpplib.gsm.gsm_encode(text) if data_coding == 3: return text.encode("latin-1") if data_coding == 8: return text.encode("utf-16-be") raise ValueError(f"unsupported dataCoding {data_coding}") def do_submit(body): # Does not wait for the submit_sm_resp: a single-segment message is only answered once the # caller's own "sms" handler calls sendResp(), which the caller can only do after seeing this # call return - waiting here would deadlock exactly that handshake. Poll /ack for the result. session = session_for(body["name"]) data_coding = int(body["dataCoding"]) payload = encode_body(body.get("text", ""), data_coding) if "text" in body else b"" if body.get("extraHex"): payload += bytes.fromhex(body["extraHex"]) with session.send_lock: pdu = session.client.send_message( source_addr=body["from"], destination_addr=body["to"], short_message=payload, data_coding=data_coding, esm_class=int(body.get("esmClass", 0)), ) sequence = pdu.sequence return {"ok": True, "sequence": sequence} def do_ack(query): session = session_for(query["name"][0]) sequence = int(query["sequence"][0]) ack = session.acks.get(sequence) if ack is None: return {"found": False, "ok": True} return {"found": True, "messageId": ack["messageId"], "ok": True, "status": ack["status"]} def do_submit_long(body): session = session_for(body["name"]) data_coding = int(body["dataCoding"]) parts, encoding, esm_class = smpplib.gsm.make_parts(body["text"], encoding=data_coding, use_udhi=True) results = [] for part in parts: with session.send_lock: pdu = session.client.send_message( source_addr=body["from"], destination_addr=body["to"], short_message=part, data_coding=encoding, esm_class=esm_class, ) sequence = pdu.sequence ack = session.wait_ack(sequence, budget=8) results.append(ack) return {"ok": all(results), "parts": len(parts), "results": results} def do_enquire_link(body): session = session_for(body["name"]) with session.send_lock: pdu = smpplib.smpp.make_pdu("enquire_link", client=session.client) session.client.send_pdu(pdu) return {"ok": True} def do_received(name): session = session_for(name) return {"ok": True, "received": session.received} def do_status(name): session = session_for(name) return { "ok": True, "readerError": session.reader_error, "readerRunning": session.reader_running, "receivedCount": len(session.received), } def do_idle_silent(body): """Sleeps `seconds` sending nothing at all - no reader thread, no enquire_link - then does one read attempt to say whether the peer (our server) closed the link while it was silent.""" session = session_for(body["name"]) time.sleep(float(body["seconds"])) session.client._socket.settimeout(2) try: session.client.read_pdu() return {"closed": False, "ok": True} except socket.timeout: return {"closed": False, "ok": True} except smpplib.exceptions.ConnectionError: return {"closed": True, "ok": True} def do_unbind(body): session = session_for(body["name"]) try: session.client.unbind() except Exception: # noqa: BLE001 - best-effort teardown pass session.client.disconnect() with SESSIONS_LOCK: del SESSIONS[body["name"]] return {"ok": True} ROUTES = { "/bind": lambda body, _query: do_bind(body), "/enquireLink": lambda body, _query: do_enquire_link(body), "/idleSilent": lambda body, _query: do_idle_silent(body), "/startReader": lambda body, _query: do_start_reader(body), "/submit": lambda body, _query: do_submit(body), "/submitLong": lambda body, _query: do_submit_long(body), "/unbind": lambda body, _query: do_unbind(body), } GET_ROUTES = { "/ack": do_ack, "/received": lambda query: do_received(query["name"][0]), "/status": lambda query: do_status(query["name"][0]), } class Handler(BaseHTTPRequestHandler): def log_message(self, fmt, *args): sys.stderr.write("%s - %s\n" % (self.address_string(), fmt % args)) def _respond(self, status, payload): body = json.dumps(payload).encode() self.send_response(status) self.send_header("Content-Type", "application/json") self.send_header("Content-Length", str(len(body))) self.end_headers() self.wfile.write(body) def do_GET(self): if self.path == "/health": self._respond(200, {"ok": True}) return parsed = urlparse(self.path) handler = GET_ROUTES.get(parsed.path) if handler is None: self._respond(404, {"error": "no such route"}) return try: self._respond(200, handler(parse_qs(parsed.query))) except Exception as exc: # noqa: BLE001 - surfaced to the caller, not the process self._respond(500, {"error": f"{type(exc).__name__}: {exc}"}) def do_POST(self): handler = ROUTES.get(self.path) if handler is None: self._respond(404, {"error": "no such route"}) return length = int(self.headers.get("Content-Length", 0)) raw = self.rfile.read(length) if length else b"{}" body = json.loads(raw or b"{}") try: self._respond(200, handler(body, None)) except Exception as exc: # noqa: BLE001 - surfaced to the caller, not the process self._respond(500, {"error": f"{type(exc).__name__}: {exc}"}) def main(): server = ThreadingHTTPServer(("0.0.0.0", 8080), Handler) print(f"python-smpplib driver listening on 8080, target {TARGET_HOST}:{TARGET_PORT}", file=sys.stderr) server.serve_forever() if __name__ == "__main__": main()