233 config = sim_config.load()
234 site_tag = sim_config.require(config,
"siteTag")
235 period = float(config.get(
"simPeriodSeconds", 2.0))
236 control_host = sim_config.require(config,
"controlServer",
"bindHost")
237 control_port = sim_config.require(config,
"controlServer",
"port")
248 logging.basicConfig(level=logging.INFO, format=
"[%(levelname)s] %(message)s", stream=sys.stdout)
249 log = logging.getLogger(
"train")
255 link = _safecomm_link_from_config(log, config, TRAIN_WIRE_SIZE, sim_config)
257 log.info(
"train link: via CommServer (safeCommFreamwork path)")
260 "WEST": (sim_config.require(config,
"rbcWest",
"host"), sim_config.require(config,
"rbcWest",
"port")),
261 "EAST": (sim_config.require(config,
"rbcEast",
"host"), sim_config.require(config,
"rbcEast",
"port")),
263 link = DualHomedLink(log, targets, TRAIN_WIRE_SIZE)
264 was_connected = {site:
False for site
in SITES}
273 online_site = {
"value":
None}
301 connected_to_rbc = {}
321 HEARTBEAT_PERIOD_S = 6.0
332 _connect_msg_name =
"TRAIN_CONNECT" if "TRAIN_CONNECT" in CATALOG.message_names()
else (
333 "P0" if "P0" in CATALOG.message_names()
else None)
334 if _connect_msg_name
is None:
335 log.error(
"[IO] [%s] [internal] [CONFIG] [catalog has neither TRAIN_CONNECT nor P0 - "
336 "train connect handshake cannot be sent]", f
"train-{site_tag.lower()}")
338 def send_p0_now(target_nid):
339 """Sends the connect handshake for exactly ONE train, @target_nid -
340 no implicit "every train" fan-out any more (there is no roster to
341 fan out over - see main()'s own "no default trains" doc). The
342 legacy plain CONNECT command (handle_connect below) now passes a
343 single hardcoded nid_engine explicitly, same as connectTrain
345 if _connect_msg_name
is None:
348 s = link.sockets[site]
351 s.sendall(
_pad(CATALOG.build_envelope(_connect_msg_name, target_nid, {})))
352 log.info(
"[IO] [%s] [c-%s] [%s] [ERTMS Train connection handshake sent] [train_id=%d]",
353 f
"train-{site_tag.lower()}", site.lower(), _connect_msg_name, target_nid)
380 def _reset_train(tid):
381 state[tid] = {
"position_m": 0,
"cycle": 0}
389 last_received.pop(tid,
None)
391 def _ensure_train(tid):
392 """Lazily creates @tid's own state/train_vars entries the first
393 time it's actually named (connectTrain, or the legacy REPORT/
394 CONNECT commands) - see main()'s own "no default trains" doc.
395 Idempotent - safe to call on every touch, not just the first."""
396 state.setdefault(tid, {
"position_m": 0,
"cycle": 0})
397 train_vars.setdefault(tid, {
410 BALISE_NAMES = site_data.balise_names()
412 def get_balise_name(nid):
415 return BALISE_NAMES.get(int(nid), f
"BG_{nid}")
417 def send_position_report(tid, d_lrbg, nid_lrbg=None):
418 """Sends this train's cyclic position report as a REAL Subset-026
419 Msg 136 / Packet 0 (CATALOG's own "M136" ertms entry -
420 ab_gp_ertms_train_message_codec.c on the RBC side), not the old
421 flat scheme's own "136" kind - switched this session, on direct
422 request, for every caller of this shared function (the autonomous
423 cyclic loop below AND the REPORT/handle_report command), so there
424 is exactly one M136-sending code path, not two. GP's own real
425 Msg 136 handler (handle_msg_136(), ab_gp_proc_dispatch.c) already
426 works standalone (no prior Msg 155 session handshake required) and
427 already replies with a real Msg 24 - see that function's own doc.
428 "cycle" stays purely local bookkeeping (this sim's own REPORT-
429 command/logging counter) - real Msg 136 has no counter field of
430 its own to carry it on the wire (message_catalog.json's own "M136"
431 reference entry doc: "cycle... is a project-invented counter with
432 no real-message-136 equivalent")."""
436 st[
"position_m"] = d_lrbg
437 vars_for_train = train_vars.setdefault(tid, {})
438 vars_for_train[
"d_lrbg"] = d_lrbg
439 if nid_lrbg
is not None:
440 vars_for_train[
"nid_lrbg"] = nid_lrbg
442 cur_lrbg = vars_for_train.get(
"nid_lrbg",
_default_lrbg(tid))
443 cur_v = vars_for_train.get(
"v_train", 0)
444 cur_mode = vars_for_train.get(
"m_mode", 0)
445 cur_level = vars_for_train.get(
"m_level", 2)
455 "q_scale": vars_for_train.get(
"q_scale", ertms_codec.Q_SCALE_1M),
456 "nid_c": vars_for_train.get(
"nid_c", 0),
459 "q_dirlrbg": vars_for_train.get(
"q_dirlrbg", ertms_codec.Q_DIRLRBG_NOMINAL),
460 "q_dlrbg": vars_for_train.get(
"q_dlrbg", ertms_codec.Q_DLRBG_NOMINAL),
461 "l_doubtover": vars_for_train.get(
"l_doubtover", 0),
462 "l_doubtunder": vars_for_train.get(
"l_doubtunder", 0),
465 "m_level": cur_level,
467 link.send_all(
_pad(CATALOG.build_envelope(
"M136", tid, wire_fields)))
473 payload = dict(wire_fields)
474 payload[
"cycle"] = st[
"cycle"]
475 payload[
"nid_lrbg"] = cur_lrbg
478 def handle_report(args):
485 d_lrbg = int(args[0])
486 nid_lrbg = int(args[1])
if len(args) > 1
else None
487 payload = send_position_report(tid, d_lrbg, nid_lrbg)
488 fields_str =
" ".join(f
"{k}={v}" for k, v
in sorted(payload.items()))
489 bname = get_balise_name(payload.get(
"nid_lrbg"))
490 balise_tag = f
" balise={bname}" if bname
else ""
491 log.info(
"[IO] [%s] [c-west,c-east] [M136] [ERTMS Train position report commanded] "
492 "[train_id=%d] [P0] [%s%s]",
493 f
"train-{site_tag.lower()}", tid, fields_str, balise_tag)
496 def handle_ping(_args):
499 def handle_connect(_args):
505 connected_to_rbc[tid] =
True
506 link.ensure_connected()
510 def handle_disconnect(_args):
512 link.disconnect(site)
515 def handle_disconnect_with_online(_args):
516 site = online_site[
"value"]
518 return "ERR no site currently believed online\n"
519 link.disconnect(site)
522 def handle_disconnect_with_standby(_args):
523 site = online_site[
"value"]
525 return "ERR no site currently believed online\n"
526 standby =
"EAST" if site ==
"WEST" else "WEST"
527 link.disconnect(standby)
530 def handle_send_message(request):
531 nid_engine = request[
"nidEngine"]
532 message = request[
"message"]
533 fields = request.get(
"fields", {})
534 link.send_all(
_pad(CATALOG.build_envelope(message, nid_engine, fields)))
535 fields_str =
" ".join(f
"{k}={v}" for k, v
in sorted(fields.items()))
536 log.info(
"[IO] [%s] [c-west,c-east] [%s] [ERTMS message sent] [train_id=%d %s]",
537 f
"train-{site_tag.lower()}", message, nid_engine, fields_str)
542 if message
in (
"M150",
"M156"):
543 ertms_post_eom[nid_engine] =
True
544 next_heartbeat[nid_engine] = time.monotonic() + HEARTBEAT_PERIOD_S
545 elif message ==
"M155":
546 ertms_post_eom[nid_engine] =
False
547 return {
"status":
"OK"}
549 def handle_get_message(request):
550 nid_engine = request[
"nidEngine"]
551 message = str(request[
"message"])
552 packets = last_received.get(nid_engine, {}).get(message)
563 if packets
is None and message
in (
"24",
"15",
"M15"):
564 packets = last_received.get(nid_engine, {}).get(
"PositionReportAck")
566 return {
"status":
"ERR",
"reason": f
"no {message} received yet for nidEngine {nid_engine}"}
567 result_packets = {message: packets}
569 for pkt_name
in CATALOG.packets_for(message):
570 result_packets[pkt_name] = packets
573 if "msg_type" in packets:
574 result_packets[
"PositionReportAck"] = packets
575 result_packets[
"M24_M15"] = packets
576 result_packets[
"24"] = packets
577 result_packets[str(packets[
"msg_type"])] = packets
578 result_packets[f
"M{packets['msg_type']}"] = packets
579 return {
"status":
"OK",
"message": message,
"packets": result_packets}
581 def handle_connect_train(request):
582 nid_engine = request[
"nidEngine"]
586 if nid_engine
not in state
and len(state) >= MAX_TRAINS:
587 return {
"status":
"ERR",
588 "reason": f
"this sim already hosts {MAX_TRAINS} distinct trains (MAX_TRAINS)"}
589 _reset_train(nid_engine)
601 if "nidLrbg" in request
and request[
"nidLrbg"]
is not None:
602 train_vars[nid_engine][
"nid_lrbg"] = int(request[
"nidLrbg"])
603 if "dLrbg" in request
and request[
"dLrbg"]
is not None:
604 d_lrbg = int(request[
"dLrbg"])
605 train_vars[nid_engine][
"d_lrbg"] = d_lrbg
606 state[nid_engine][
"position_m"] = d_lrbg
607 connected_to_rbc[nid_engine] =
True
608 link.ensure_connected()
609 send_p0_now(nid_engine)
613 ertms_post_eom[nid_engine] =
False
614 log.info(
"[GENERAL] [%s] [internal] [CONNECT] [Train connected to RBC] [train_id=%d]",
615 f
"train-{site_tag.lower()}", nid_engine)
616 return {
"status":
"OK"}
618 def handle_disconnect_train(request):
619 nid_engine = request[
"nidEngine"]
621 s = link.sockets[site]
624 s.sendall(
_pad(CATALOG.build_envelope(
"TRAIN_DISCONNECT", nid_engine, {})))
627 connected_to_rbc[nid_engine] =
False
628 ertms_post_eom[nid_engine] =
False
629 _reset_train(nid_engine)
630 log.info(
"[IO] [%s] [c-west,c-east] [TRAIN_DISCONNECT] [ERTMS Train disconnected from RBC] [train_id=%d]",
631 f
"train-{site_tag.lower()}", nid_engine)
632 return {
"status":
"OK"}
634 def handle_set_train_variable(request):
635 nid_engine = request[
"nidEngine"]
636 variable = request[
"variable"]
637 value = request[
"value"]
638 train_vars.setdefault(nid_engine, {})[variable] = value
652 if variable ==
"d_lrbg":
653 state.setdefault(nid_engine, {
"position_m": 0,
"cycle": 0})[
"position_m"] = value
654 log.info(
"[GENERAL] [%s] [internal] [VAR_STORE] [Train variable set] [train_id=%d %s=%s]",
655 f
"train-{site_tag.lower()}", nid_engine, variable, value)
656 return {
"status":
"OK"}
658 def handle_send_train_message(request):
659 nid_engine = request[
"nidEngine"]
660 message = request[
"message"]
661 packet = request.get(
"packet", message)
662 stored = train_vars.get(nid_engine, {})
679 if CATALOG.is_ertms(message):
680 fields = dict(stored)
682 field_specs = CATALOG.fields_for(message, packet
if packet != message
else None)
683 fields = {name: stored[name]
for name
in field_specs
if name
in stored}
684 except Exception
as e:
685 return {
"status":
"ERR",
"reason": str(e)}
686 if message ==
"136" or message ==
"PositionReport" or message ==
"M136":
687 train_st = state.setdefault(nid_engine, {
"position_m": 0,
"cycle": 0})
688 if "cycle" not in fields:
689 train_st[
"cycle"] += 1
690 fields[
"cycle"] = train_st[
"cycle"]
691 if "nid_lrbg" not in fields:
693 if "d_lrbg" in fields:
694 train_st[
"position_m"] = fields[
"d_lrbg"]
695 if message
in (
"132",
"M132",
"146",
"M146"):
696 last_received.get(nid_engine, {}).pop(
"M24_CONFIG",
None)
697 link.send_all(
_pad(CATALOG.build_envelope(message, nid_engine, fields)))
698 fields_str =
" ".join(f
"{k}={v}" for k, v
in sorted(fields.items()))
699 msg_tag = f
"M{message}" if message.isdigit()
else message
700 log.info(
"[IO] [%s] [c-west,c-east] [%s] [ERTMS message sent from stored variables] [train_id=%d %s]",
701 f
"train-{site_tag.lower()}", msg_tag, nid_engine, fields_str)
702 return {
"status":
"OK"}
704 def handle_clear_messages(_request):
705 last_received.clear()
706 return {
"status":
"OK"}
708 control = ControlServer(
712 "REPORT": handle_report,
714 "CONNECT": handle_connect,
715 "DISCONNECT": handle_disconnect,
716 "DISCONNECTWITHONLINE": handle_disconnect_with_online,
717 "DISCONNECTWITHSTANDBY": handle_disconnect_with_standby,
720 "sendMessage": handle_send_message,
721 "getMessage": handle_get_message,
722 "connectTrain": handle_connect_train,
723 "disconnectTrain": handle_disconnect_train,
724 "setTrainVariable": handle_set_train_variable,
725 "sendTrainMessage": handle_send_train_message,
726 "clearMessages": handle_clear_messages,
728 bind_host=control_host,
733 def _stop(signum, _frame):
735 log.info(
"[GENERAL] [%s] [internal] [SIGNAL] [Train sim shutting down] [signal=%d]",
736 f
"train-{site_tag.lower()}", signum)
739 signal.signal(signal.SIGTERM, _stop)
740 signal.signal(signal.SIGINT, _stop)
742 log.info(
"[GENERAL] [%s] [internal] [INIT] [Train sim starting] "
743 "[no trains connected yet - up to %d supported, see connectTrain period=%.1fs]",
744 f
"train-{site_tag.lower()}", MAX_TRAINS, period)
748 status = StatusReporter(log, link, f
"SIM-TRAIN/{site_tag}", f
"train-{site_tag.lower()}",
749 frame_size=TRAIN_WIRE_SIZE)
750 status.emit({
"trains": sum(1
for v
in connected_to_rbc.values()
if v)})
753 link.ensure_connected()
755 now_connected = link.sockets[site]
is not None
761 if now_connected
and not was_connected[site]:
762 s = link.sockets[site]
763 for tid
in list(connected_to_rbc.keys()):
764 if not connected_to_rbc.get(tid):
767 s.sendall(
_pad(CATALOG.build_envelope(
"P0", tid, {})))
768 log.info(
"[IO] [%s] [c-%s] [P0] [ERTMS Train connection handshake sent] [train_id=%d]",
769 f
"train-{site_tag.lower()}", site.lower(), tid)
772 was_connected[site] = now_connected
774 control.poll(timeout=0.0)
775 status.tick({
"trains": sum(1
for v
in connected_to_rbc.values()
if v)})
777 now = time.monotonic()
778 if now >= next_report:
779 for tid
in list(connected_to_rbc.keys()):
780 if not connected_to_rbc.get(tid):
785 if ertms_post_eom.get(tid):
787 cur_vars = train_vars.get(tid, {})
788 cur_speed = cur_vars.get(
"v_train", 0)
790 delta_m = int(round((cur_speed * 1000.0 / 3600.0) * period))
791 new_position = state[tid][
"position_m"] + max(1, delta_m)
793 new_position = state[tid][
"position_m"]
794 payload = send_position_report(tid, new_position)
795 fields_str =
" ".join(f
"{k}={v}" for k, v
in sorted(payload.items()))
796 bname = get_balise_name(payload.get(
"nid_lrbg"))
797 balise_tag = f
" balise={bname}" if bname
else ""
805 log.info(
"[IO] [%s] [c-west,c-east] [M136] [ERTMS Train periodic position report] "
806 "[train_id=%d] [P0] [%s%s]",
807 f
"train-{site_tag.lower()}", tid, fields_str, balise_tag)
808 next_report = now + period
818 for tid
in list(connected_to_rbc.keys()):
819 if connected_to_rbc.get(tid)
and ertms_post_eom.get(tid)
and now >= next_heartbeat.get(tid, 0.0):
820 cur_vars = train_vars.get(tid, {})
821 cur_lrbg_c = cur_vars.get(
"nid_c", 0)
824 "nid_c": cur_lrbg_c,
"nid_bg": cur_lrbg_bg,
825 "prv_nid_c": cur_lrbg_c,
"prv_nid_bg": cur_lrbg_bg,
826 "d_lrbg": state[tid][
"position_m"],
827 "v_train": cur_vars.get(
"v_train", 0),
828 "m_mode": cur_vars.get(
"m_mode", 0),
829 "m_level": cur_vars.get(
"m_level", 2),
831 codec_bytes = ertms_codec.encode_msg_136_packet1(tid, packet1_report)
832 link.send_all(
_pad(CATALOG.build_ertms_envelope_raw(codec_bytes, tid)))
833 bname = get_balise_name(cur_lrbg_bg)
834 balise_tag = f
" balise={bname}" if bname
else ""
839 log.info(
"[IO] [%s] [c-west,c-east] [M136] [ERTMS post-EoM idle heartbeat sent] "
840 "[train_id=%d] [P1] [d_lrbg=%d nid_lrbg=%d%s]",
841 f
"train-{site_tag.lower()}", tid, state[tid][
"position_m"], cur_lrbg_bg, balise_tag)
842 next_heartbeat[tid] = now + HEARTBEAT_PERIOD_S
844 for site, data
in link.poll_recv(timeout=0.2):
845 decoded = CATALOG.decode_to_message(data)
847 log.warning(
"[%s] [internal] [MALFORMED] [Malformed envelope dropped] [site=%s]",
848 f
"train-{site_tag.lower()}", site)
850 nid_engine, message, fields = decoded
852 if fields.get(
"packetsPresent"):
853 last_received.setdefault(nid_engine, {})[
"M24_CONFIG"] = fields
854 last_received.setdefault(nid_engine, {})[
"M24"] = fields
855 elif "M24_CONFIG" not in last_received.get(nid_engine, {}):
856 last_received.setdefault(nid_engine, {})[
"M24"] = fields
858 last_received.setdefault(nid_engine, {})[message] = fields
859 online_site[
"value"] = site
860 fields_str =
" ".join(f
"{k}={v}" for k, v
in sorted(fields.items()))
861 if message ==
"PositionReportAck":
862 msg_type_name = f
"M{fields.get('msg_type', '24')}"
863 log.info(
"[IO] [c-%s] [%s] [%s] [ERTMS Position report ack received] [nid_engine=%d %s]",
864 site.lower(), f
"train-{site_tag.lower()}", msg_type_name, nid_engine, fields_str)
865 elif message
in (
"3",
"M3"):
866 log.info(
"[IO] [c-%s] [%s] [M3] [ERTMS Movement Authority received] [nid_engine=%d %s]",
867 site.lower(), f
"train-{site_tag.lower()}", nid_engine, fields_str)
869 link.send_all(
_pad(CATALOG.build_envelope(
"146", nid_engine, {
"ma_seq": fields[
"ma_seq"]})))
870 log.info(
"[IO] [%s] [c-%s] [M146] [ERTMS Movement Authority Ack sent] [train_id=%d ma_seq=%d]",
871 f
"train-{site_tag.lower()}", site.lower(), nid_engine, fields[
"ma_seq"])
873 msg_tag = message
if message.startswith(
"M")
else (f
"M{message}" if message.isdigit()
else message)
874 log.info(
"[IO] [c-%s] [%s] [%s] [ERTMS message received from RBC] [nid_engine=%d %s]",
875 site.lower(), f
"train-{site_tag.lower()}", msg_tag, nid_engine, fields_str)
877 log.info(
"stopped after %d cycle(s)", sum(st[
"cycle"]
for st
in state.values()))