app.py 146 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480148114821483148414851486148714881489149014911492149314941495149614971498149915001501150215031504150515061507150815091510151115121513151415151516151715181519152015211522152315241525152615271528152915301531153215331534153515361537153815391540154115421543154415451546154715481549155015511552155315541555155615571558155915601561156215631564156515661567156815691570157115721573157415751576157715781579158015811582158315841585158615871588158915901591159215931594159515961597159815991600160116021603160416051606160716081609161016111612161316141615161616171618161916201621162216231624162516261627162816291630163116321633163416351636163716381639164016411642164316441645164616471648164916501651165216531654165516561657165816591660166116621663166416651666166716681669167016711672167316741675167616771678167916801681168216831684168516861687168816891690169116921693169416951696169716981699170017011702170317041705170617071708170917101711171217131714171517161717171817191720172117221723172417251726172717281729173017311732173317341735173617371738173917401741174217431744174517461747174817491750175117521753175417551756175717581759176017611762176317641765176617671768176917701771177217731774177517761777177817791780178117821783178417851786178717881789179017911792179317941795179617971798179918001801180218031804180518061807180818091810181118121813181418151816181718181819182018211822182318241825182618271828182918301831183218331834183518361837183818391840184118421843184418451846184718481849185018511852185318541855185618571858185918601861186218631864186518661867186818691870187118721873187418751876187718781879188018811882188318841885188618871888188918901891189218931894189518961897189818991900190119021903190419051906190719081909191019111912191319141915191619171918191919201921192219231924192519261927192819291930193119321933193419351936193719381939194019411942194319441945194619471948194919501951195219531954195519561957195819591960196119621963196419651966196719681969197019711972197319741975197619771978197919801981198219831984198519861987198819891990199119921993199419951996199719981999200020012002200320042005200620072008200920102011201220132014201520162017201820192020202120222023202420252026202720282029203020312032203320342035203620372038203920402041204220432044204520462047204820492050205120522053205420552056205720582059206020612062206320642065206620672068206920702071207220732074207520762077207820792080208120822083208420852086208720882089209020912092209320942095209620972098209921002101210221032104210521062107210821092110211121122113211421152116211721182119212021212122212321242125212621272128212921302131213221332134213521362137213821392140214121422143214421452146214721482149215021512152215321542155215621572158215921602161216221632164216521662167216821692170217121722173217421752176217721782179218021812182218321842185218621872188218921902191219221932194219521962197219821992200220122022203220422052206220722082209221022112212221322142215221622172218221922202221222222232224222522262227222822292230223122322233223422352236223722382239224022412242224322442245224622472248224922502251225222532254225522562257225822592260226122622263226422652266226722682269227022712272227322742275227622772278227922802281228222832284228522862287228822892290229122922293229422952296229722982299230023012302230323042305230623072308230923102311231223132314231523162317231823192320232123222323232423252326232723282329233023312332233323342335233623372338233923402341234223432344234523462347234823492350235123522353235423552356235723582359236023612362236323642365236623672368236923702371237223732374237523762377237823792380238123822383238423852386238723882389239023912392239323942395239623972398239924002401240224032404240524062407240824092410241124122413241424152416241724182419242024212422242324242425242624272428242924302431243224332434243524362437243824392440244124422443244424452446244724482449245024512452245324542455245624572458245924602461246224632464246524662467246824692470247124722473247424752476247724782479248024812482248324842485248624872488248924902491249224932494249524962497249824992500250125022503250425052506250725082509251025112512251325142515251625172518251925202521252225232524252525262527252825292530253125322533253425352536253725382539254025412542254325442545254625472548254925502551255225532554255525562557255825592560256125622563256425652566256725682569257025712572257325742575257625772578257925802581258225832584258525862587258825892590259125922593259425952596259725982599260026012602260326042605260626072608260926102611261226132614261526162617261826192620262126222623262426252626262726282629263026312632263326342635263626372638263926402641264226432644264526462647264826492650265126522653265426552656265726582659266026612662266326642665266626672668266926702671267226732674267526762677267826792680268126822683268426852686268726882689269026912692269326942695269626972698269927002701270227032704270527062707270827092710271127122713271427152716271727182719272027212722272327242725272627272728272927302731273227332734273527362737273827392740274127422743274427452746274727482749275027512752275327542755275627572758275927602761276227632764276527662767276827692770277127722773277427752776277727782779278027812782278327842785278627872788278927902791279227932794279527962797279827992800280128022803280428052806280728082809281028112812281328142815281628172818281928202821282228232824282528262827282828292830283128322833283428352836283728382839284028412842284328442845284628472848284928502851285228532854285528562857285828592860286128622863286428652866286728682869287028712872287328742875287628772878287928802881288228832884288528862887288828892890289128922893289428952896289728982899290029012902290329042905290629072908290929102911291229132914291529162917291829192920292129222923292429252926292729282929293029312932293329342935293629372938293929402941294229432944294529462947294829492950295129522953295429552956295729582959296029612962296329642965296629672968296929702971297229732974297529762977297829792980298129822983298429852986298729882989299029912992299329942995299629972998299930003001300230033004300530063007300830093010301130123013301430153016301730183019302030213022302330243025302630273028302930303031303230333034303530363037303830393040304130423043304430453046304730483049305030513052305330543055305630573058305930603061306230633064306530663067306830693070307130723073307430753076307730783079308030813082308330843085308630873088308930903091309230933094309530963097309830993100310131023103310431053106310731083109311031113112311331143115311631173118311931203121312231233124312531263127312831293130313131323133313431353136313731383139314031413142314331443145314631473148314931503151315231533154315531563157315831593160316131623163316431653166316731683169317031713172317331743175317631773178317931803181318231833184318531863187318831893190319131923193319431953196319731983199320032013202320332043205320632073208320932103211321232133214321532163217321832193220322132223223322432253226322732283229323032313232323332343235323632373238323932403241324232433244324532463247324832493250325132523253325432553256325732583259326032613262326332643265326632673268326932703271327232733274327532763277327832793280328132823283328432853286328732883289329032913292329332943295329632973298329933003301330233033304330533063307330833093310331133123313331433153316331733183319332033213322332333243325332633273328332933303331333233333334333533363337333833393340334133423343334433453346334733483349335033513352335333543355335633573358335933603361336233633364336533663367336833693370337133723373337433753376337733783379338033813382338333843385338633873388338933903391339233933394339533963397339833993400340134023403340434053406340734083409341034113412341334143415341634173418341934203421342234233424342534263427342834293430343134323433343434353436343734383439344034413442344334443445344634473448344934503451345234533454345534563457345834593460346134623463346434653466346734683469347034713472347334743475347634773478347934803481348234833484348534863487348834893490349134923493349434953496349734983499350035013502350335043505350635073508350935103511351235133514351535163517351835193520352135223523352435253526352735283529353035313532353335343535353635373538353935403541354235433544354535463547354835493550355135523553355435553556355735583559356035613562356335643565356635673568356935703571357235733574357535763577357835793580358135823583358435853586358735883589359035913592359335943595359635973598359936003601360236033604360536063607360836093610361136123613361436153616361736183619362036213622362336243625362636273628362936303631363236333634363536363637363836393640364136423643364436453646364736483649365036513652365336543655365636573658365936603661366236633664366536663667366836693670367136723673367436753676367736783679368036813682368336843685368636873688368936903691369236933694369536963697369836993700370137023703370437053706370737083709371037113712371337143715371637173718371937203721372237233724372537263727372837293730373137323733373437353736373737383739374037413742374337443745374637473748374937503751375237533754375537563757375837593760376137623763376437653766376737683769377037713772377337743775377637773778377937803781378237833784378537863787378837893790379137923793379437953796379737983799380038013802380338043805380638073808380938103811381238133814381538163817
  1. from flask import Flask, jsonify, request, abort
  2. from flask_cors import CORS
  3. from flask_socketio import SocketIO, emit
  4. import threading
  5. import queue
  6. import time
  7. import json
  8. import os
  9. import logging
  10. import uuid
  11. # 配置文件路径
  12. SERIAL_CONFIG_FILE = os.path.join(os.path.dirname(os.path.abspath(__file__)), 'serial_config.json')
  13. # 导入配置
  14. from config import (
  15. MAX_BUFFER_SIZE,
  16. FLASK_SECRET_KEY,
  17. FLASK_DEBUG,
  18. FLASK_HOST,
  19. FLASK_PORT,
  20. LOG_LEVEL,
  21. LOG_FORMAT,
  22. LOG_FILE,
  23. SOCKETIO_ASYNC_MODE,
  24. SOCKETIO_ALLOWED_ORIGINS,
  25. DEFAULT_FORWARD_SERIAL_TO_MQTT,
  26. DEFAULT_FORWARD_MQTT_TO_SERIAL,
  27. DEFAULT_MQTT_PUBLISH_TOPIC,
  28. SOCKETIO_NAMESPACE_DATA,
  29. SOCKETIO_NAMESPACE_STATUS,
  30. SOCKETIO_NAMESPACE_CONTROL,
  31. ERROR_CODES,
  32. ERROR_MESSAGES,
  33. DEFAULT_CUSTOMER_ID,
  34. DEFAULT_DTU_ID,
  35. DTU_HEARTBEAT_INTERVAL,
  36. DTU_OFFLINE_TIMEOUT,
  37. DEFAULT_MQTT_TOPIC_PREFIX,
  38. DTU_REGISTRATION_TOPIC,
  39. DTU_STATUS_TOPIC,
  40. DTU_CONTROL_TOPIC,
  41. DTU_RESPONSE_TOPIC,
  42. DTU_EVENT_TOPIC,
  43. DTU_ALARM_TOPIC,
  44. DTU_BROADCAST_TOPIC,
  45. MAX_PANELS,
  46. MAX_PORTS_PER_PANEL,
  47. DHT11_ENABLED,
  48. DHT11_GPIO_PIN,
  49. DHT11_GPIO_CHIP,
  50. DHT11_POLL_INTERVAL,
  51. DHT11_SIMULATE,
  52. NEXTION_ENABLED,
  53. NEXTION_SCREEN_PORT,
  54. NEXTION_SCREEN_BAUD,
  55. NEXTION_PANEL_NAME_PREFIX,
  56. NEXTION_CMD_DELAY_MS
  57. )
  58. # 配置日志
  59. logging_config = {
  60. 'level': getattr(logging, LOG_LEVEL),
  61. 'format': LOG_FORMAT
  62. }
  63. if LOG_FILE:
  64. logging_config['filename'] = LOG_FILE
  65. logging.basicConfig(**logging_config)
  66. logger = logging.getLogger('serial_mqtt_gateway')
  67. from modules.serial_port import SerialPort
  68. from modules.mqtt_client import MQTTClient
  69. from modules.network_config import network_manager
  70. from modules.modbus_rtu import ModbusRTUClient, ANTENNA_ADDRESSES, AddressConfigProtocol, build_broadcast_query, build_confirm_address, build_assign_address
  71. from modules.dht11_sensor import DHT11Sensor
  72. from modules.nextion_display import NextionDisplay
  73. app = Flask(__name__)
  74. app.config['SECRET_KEY'] = FLASK_SECRET_KEY
  75. # 配置CORS以允许nginx代理的前端访问
  76. CORS(app, resources={r"/api/*": {"origins": "*"}, r"/socket.io/*": {"origins": "*"}})
  77. # 初始化SocketIO
  78. socketio = SocketIO(
  79. app,
  80. cors_allowed_origins=SOCKETIO_ALLOWED_ORIGINS,
  81. async_mode=SOCKETIO_ASYNC_MODE,
  82. manage_session=False, # 禁用会话管理以提高性能
  83. ping_timeout=30, # 心跳超时时间
  84. ping_interval=25, # 心跳间隔
  85. logger=FLASK_DEBUG, # 根据Flask调试模式决定是否记录SocketIO日志
  86. engineio_logger=FLASK_DEBUG
  87. )
  88. # 初始化串口和MQTT客户端
  89. serial_client = SerialPort()
  90. mqtt_client = MQTTClient()
  91. modbus_client = ModbusRTUClient(serial_client)
  92. address_config = AddressConfigProtocol(serial_client)
  93. # Nextion 串口屏实例
  94. screen_display = NextionDisplay()
  95. screen_current_panel_index = 0
  96. screen_lock = threading.Lock()
  97. screen_refresh_queue = queue.Queue()
  98. screen_refresh_thread = None
  99. # 初始化 DHT11 传感器(后续在 __main__ 中设置回调并启动)
  100. dht11_sensor = None
  101. if DHT11_ENABLED and (DHT11_GPIO_PIN is not None or DHT11_SIMULATE):
  102. dht11_sensor = DHT11Sensor(
  103. gpio_pin=DHT11_GPIO_PIN,
  104. gpio_chip=DHT11_GPIO_CHIP,
  105. poll_interval=DHT11_POLL_INTERVAL,
  106. simulate=DHT11_SIMULATE
  107. )
  108. elif DHT11_ENABLED and DHT11_GPIO_PIN is None and not DHT11_SIMULATE:
  109. logger.warning("DHT11 已启用但未配置 GPIO pin (DHT11_GPIO_PIN),跳过本地传感器启动")
  110. # 转发标志
  111. forward_serial_to_mqtt = DEFAULT_FORWARD_SERIAL_TO_MQTT
  112. forward_mqtt_to_serial = DEFAULT_FORWARD_MQTT_TO_SERIAL
  113. mqtt_publish_topic = DEFAULT_MQTT_PUBLISH_TOPIC
  114. # DTU MQTT协议配置(运行时可修改)
  115. dtu_config = {
  116. 'topic_prefix': DEFAULT_MQTT_TOPIC_PREFIX,
  117. 'customer_id': DEFAULT_CUSTOMER_ID,
  118. 'dtu_id': DEFAULT_DTU_ID,
  119. 'firmware_version': 'v1.0.0',
  120. 'hardware_version': 'v1.0',
  121. 'heartbeat_interval': DTU_HEARTBEAT_INTERVAL,
  122. 'enabled': True
  123. }
  124. DTU_CONFIG_FILE = os.path.join(os.path.dirname(os.path.abspath(__file__)), 'dtu_config.json')
  125. def load_dtu_config():
  126. global dtu_config
  127. try:
  128. if os.path.exists(DTU_CONFIG_FILE):
  129. with open(DTU_CONFIG_FILE, 'r') as f:
  130. saved = json.load(f)
  131. dtu_config.update(saved)
  132. except Exception as e:
  133. logger.warning(f"加载DTU配置失败: {e}")
  134. def save_dtu_config():
  135. try:
  136. with open(DTU_CONFIG_FILE, 'w') as f:
  137. json.dump(dtu_config, f, indent=2)
  138. except Exception as e:
  139. logger.warning(f"保存DTU配置失败: {e}")
  140. # 端口状态追踪(用于事件检测)
  141. port_state = {} # {panel_id: {port_id: {'last_uid': str, 'expected_uid': str, 'alarm_count': int}}}
  142. def get_sorted_panels():
  143. """按 position 排序返回 panel 列表。"""
  144. return sorted(panel_config.items(), key=lambda x: x[1].get('position', 0))
  145. def _port_status_to_pic(port_state_entry):
  146. """将端口状态推导为 Nextion pic 编号。"""
  147. if not port_state_entry:
  148. return 0
  149. last_uid = port_state_entry.get('last_uid')
  150. expected_uid = port_state_entry.get('expected_uid')
  151. last_polled = port_state_entry.get('last_polled_at')
  152. if not last_polled:
  153. return 0 # UNKNOWN
  154. if last_uid is None:
  155. return 1 # DISCONNECTED
  156. if expected_uid and last_uid != expected_uid:
  157. return 2 # ILLEGAL
  158. return 3 # CONNECTED
  159. def _get_display_ip():
  160. """获取用于屏幕显示的本机 IP 地址。"""
  161. try:
  162. net_status = network_manager.get_network_status()
  163. for iface, data in net_status.get('interfaces', {}).items():
  164. if iface in ('lo', 'docker0') or iface.startswith(('br-', 'veth')):
  165. continue
  166. for addr in data.get('ip_addresses', []):
  167. if '.' in addr and not addr.startswith('127.'):
  168. return addr
  169. # 回退:尝试网络配置中的静态 IP
  170. net_config = network_manager.get_network_config()
  171. static_ip = net_config.get('static_config', {}).get('ip_address')
  172. if static_ip:
  173. return static_ip
  174. except Exception as e:
  175. logger.warning(f"获取显示 IP 失败: {e}")
  176. return ''
  177. def _screen_refresh_worker():
  178. """在后台线程中异步刷新 Nextion 屏幕。"""
  179. while True:
  180. try:
  181. item = screen_refresh_queue.get()
  182. if item is None:
  183. break
  184. panel_id = item[0]
  185. refresh_screen_panel(panel_id)
  186. except Exception as e:
  187. logger.error(f"屏幕刷新工作线程异常: {e}")
  188. finally:
  189. try:
  190. screen_refresh_queue.task_done()
  191. except Exception:
  192. pass
  193. def refresh_screen_panel(panel_id):
  194. """将指定 panel 的数据下发到 Nextion 屏幕。"""
  195. global screen_current_panel_index
  196. with screen_lock:
  197. panels = get_sorted_panels()
  198. if not panels:
  199. return False, '没有可用的 panel'
  200. panel_ids = [p[0] for p in panels]
  201. if panel_id not in panel_ids:
  202. return False, f'panel 不存在: {panel_id}'
  203. if not screen_display.get_status().get('connected'):
  204. return False, '屏幕串口未连接'
  205. ps = port_state.get(panel_id, {})
  206. now = time.time()
  207. online_count = sum(1 for pid in panel_ids if device_last_seen.get(pid, 0) > now - 60)
  208. idx = panel_ids.index(panel_id)
  209. screen_current_panel_index = idx
  210. mqtt_st = mqtt_client.get_status()
  211. mqtt_connected = mqtt_st.get('connected', False) if isinstance(mqtt_st, dict) else bool(mqtt_st)
  212. delay_ms = NEXTION_CMD_DELAY_MS
  213. # 获取本机 IP 地址显示在 t1.txt
  214. ip_address = _get_display_ip()
  215. screen_display.send_cmd(f't2.txt="{NEXTION_PANEL_NAME_PREFIX}{idx + 1}"', delay_ms=delay_ms)
  216. screen_display.send_cmd(f't1.txt="{ip_address}"', delay_ms=delay_ms)
  217. screen_display.send_cmd(f't4.txt="{dtu_config.get("firmware_version", "v1.0.0")}"', delay_ms=delay_ms)
  218. mqtt_str = '已连接' if mqtt_connected else '未连接'
  219. screen_display.send_cmd(f'g0.txt="MQTT:{mqtt_str} 当前终端数:{online_count}"', delay_ms=delay_ms)
  220. screen_display.send_cmd(f'g0.pco={0 if mqtt_connected else 63488}', delay_ms=delay_ms)
  221. for port_id in range(1, 25):
  222. pic = _port_status_to_pic(ps.get(port_id, {}))
  223. screen_display.send_cmd(f'p{port_id - 1}.pic={pic}', delay_ms=delay_ms)
  224. logger.info(f"屏幕已刷新: {panel_id} (终端{idx + 1})")
  225. return True, '屏幕已刷新'
  226. def switch_screen_panel(direction):
  227. """切换当前显示的 panel。"""
  228. global screen_current_panel_index
  229. with screen_lock:
  230. panels = get_sorted_panels()
  231. if not panels:
  232. return False, '没有可用的 panel'
  233. screen_current_panel_index = (screen_current_panel_index + direction) % len(panels)
  234. panel_id = panels[screen_current_panel_index][0]
  235. logger.info(f'屏幕切换: 终端{screen_current_panel_index + 1} -> {panel_id}')
  236. return refresh_screen_panel(panel_id)
  237. def screen_button_event_handler(page_id, comp_id, event):
  238. """处理屏幕按钮事件。"""
  239. if event != 0x01:
  240. return
  241. if page_id == 0x00 and comp_id == 0x05:
  242. logger.info('串口屏事件: 上一终端')
  243. switch_screen_panel(-1)
  244. elif page_id == 0x00 and comp_id == 0x06:
  245. logger.info('串口屏事件: 下一终端')
  246. switch_screen_panel(1)
  247. # 面板配置(从地址配置模块加载)
  248. panel_config = {} # {panel_id: {'address': int, 'position': int, 'panel_uid': str}}
  249. # 数据存储缓冲区
  250. serial_data_buffer = []
  251. mqtt_data_buffer = []
  252. # 状态标志
  253. serial_status = False
  254. mqtt_status = False
  255. # 串口配置保存和加载函数
  256. def save_serial_config(port, baudrate=9600, timeout=1.0, bytesize=8, parity='N', stopbits=1):
  257. """保存串口配置"""
  258. config = {
  259. 'port': port,
  260. 'baudrate': baudrate,
  261. 'timeout': timeout,
  262. 'bytesize': bytesize,
  263. 'parity': parity,
  264. 'stopbits': stopbits
  265. }
  266. try:
  267. with open(SERIAL_CONFIG_FILE, 'w') as f:
  268. json.dump(config, f)
  269. logger.info(f"串口配置已保存: {port} @ {baudrate}")
  270. except Exception as e:
  271. logger.error(f"保存串口配置失败: {e}")
  272. def load_serial_config():
  273. """加载串口配置"""
  274. if os.path.exists(SERIAL_CONFIG_FILE):
  275. try:
  276. with open(SERIAL_CONFIG_FILE, 'r') as f:
  277. config = json.load(f)
  278. logger.info(f"已加载串口配置: {config.get('port')} @ {config.get('baudrate')}")
  279. return config
  280. except Exception as e:
  281. logger.error(f"加载串口配置失败: {e}")
  282. return None
  283. def auto_connect_serial():
  284. """自动连接上次使用的串口"""
  285. config = load_serial_config()
  286. if config and config.get('port'):
  287. port = config.get('port')
  288. baudrate = config.get('baudrate', 9600)
  289. timeout = config.get('timeout', 1.0)
  290. bytesize = config.get('bytesize', 8)
  291. parity = config.get('parity', 'N')
  292. stopbits = config.get('stopbits', 1)
  293. logger.info(f"尝试自动连接串口: {port} @ {baudrate}")
  294. success, message = serial_client.connect(port, baudrate=baudrate, timeout=timeout, bytesize=bytesize, parity=parity, stopbits=stopbits)
  295. if success:
  296. logger.info(f"自动连接串口成功: {port}")
  297. return True
  298. else:
  299. logger.warning(f"自动连接串口失败: {message}")
  300. return False
  301. # 客户端连接管理
  302. connected_clients = {
  303. 'data': set(),
  304. 'status': set(),
  305. 'control': set()
  306. }
  307. # 设置回调函数
  308. def serial_data_handler(data):
  309. """处理串口接收的数据"""
  310. try:
  311. timestamp = time.strftime('%Y-%m-%d %H:%M:%S')
  312. serial_data_buffer.append({
  313. 'timestamp': timestamp,
  314. 'data': data,
  315. 'direction': 'in'
  316. })
  317. if len(serial_data_buffer) > MAX_BUFFER_SIZE:
  318. serial_data_buffer.pop(0)
  319. socketio.emit('serial_data', {
  320. 'timestamp': timestamp,
  321. 'data': data,
  322. 'direction': 'in'
  323. }, namespace=SOCKETIO_NAMESPACE_DATA)
  324. if forward_serial_to_mqtt and mqtt_client.get_status():
  325. success, msg = mqtt_client.publish(mqtt_publish_topic, data)
  326. if not success:
  327. logger.warning(f"串口数据转发到MQTT失败: {msg}")
  328. except Exception as e:
  329. logger.error(f"处理串口数据时出错: {str(e)}")
  330. def serial_send_handler(data):
  331. """处理串口发送的数据"""
  332. try:
  333. timestamp = time.strftime('%Y-%m-%d %H:%M:%S')
  334. serial_data_buffer.append({
  335. 'timestamp': timestamp,
  336. 'data': data,
  337. 'direction': 'out'
  338. })
  339. if len(serial_data_buffer) > MAX_BUFFER_SIZE:
  340. serial_data_buffer.pop(0)
  341. socketio.emit('serial_data', {
  342. 'timestamp': timestamp,
  343. 'data': data,
  344. 'direction': 'out'
  345. }, namespace=SOCKETIO_NAMESPACE_DATA)
  346. except Exception as e:
  347. logger.error(f"处理串口发送数据时出错: {str(e)}")
  348. def serial_status_handler(status):
  349. """处理串口状态变化"""
  350. try:
  351. global serial_status
  352. serial_status = status
  353. # 通过WebSocket广播状态变化
  354. socketio.emit('serial_status', {
  355. 'connected': status
  356. }, namespace=SOCKETIO_NAMESPACE_STATUS)
  357. logger.info(f"串口状态更新: {'已连接' if status else '已断开'}")
  358. except Exception as e:
  359. logger.error(f"处理串口状态时出错: {str(e)}")
  360. def mqtt_data_handler(data):
  361. """处理MQTT接收的数据"""
  362. try:
  363. # 添加到缓冲区
  364. timestamp = time.strftime('%Y-%m-%d %H:%M:%S')
  365. mqtt_data_buffer.append({
  366. 'timestamp': timestamp,
  367. 'topic': data['topic'],
  368. 'payload': data['payload']
  369. })
  370. # 保持缓冲区大小
  371. if len(mqtt_data_buffer) > MAX_BUFFER_SIZE:
  372. mqtt_data_buffer.pop(0)
  373. # 通过WebSocket广播数据
  374. socketio.emit('mqtt_data', {
  375. 'timestamp': timestamp,
  376. 'topic': data['topic'],
  377. 'payload': data['payload']
  378. }, namespace=SOCKETIO_NAMESPACE_DATA)
  379. # 如果启用了转发且串口已连接,转发数据到串口
  380. if forward_mqtt_to_serial and serial_client.get_status():
  381. success, msg = serial_client.send_data(data['payload'])
  382. if not success:
  383. logger.warning(f"MQTT数据转发到串口失败: {msg}")
  384. except Exception as e:
  385. logger.error(f"处理MQTT数据时出错: {str(e)}")
  386. def mqtt_status_handler(status):
  387. """处理MQTT状态变化"""
  388. try:
  389. global mqtt_status
  390. mqtt_status = status
  391. # 通过WebSocket广播状态变化
  392. socketio.emit('mqtt_status', {
  393. 'connected': status
  394. }, namespace=SOCKETIO_NAMESPACE_STATUS)
  395. logger.info(f"MQTT状态更新: {'已连接' if status else '已断开'}")
  396. # 如果MQTT连接成功且启用了DTU协议,发送注册消息和订阅控制主题
  397. if status and dtu_config.get('enabled'):
  398. socketio.sleep(1) # 等待连接稳定
  399. # 发送DTU注册消息
  400. dtu_register()
  401. # 订阅控制主题
  402. control_topic = build_dtu_topic(dtu_config['customer_id'], 'dtu', dtu_config['dtu_id'], 'control')
  403. mqtt_client.subscribe(control_topic, qos=1)
  404. logger.info(f"已订阅控制主题: {control_topic}")
  405. # 订阅广播主题 (DTU发现 7.14 和批量配置 7.15)
  406. discover_topic = build_dtu_topic('broadcast', 'dtu', 'discover')
  407. config_topic = build_dtu_topic('broadcast', 'dtu', 'config')
  408. mqtt_client.subscribe(discover_topic, qos=1)
  409. mqtt_client.subscribe(config_topic, qos=1)
  410. logger.info(f"已订阅广播主题: {discover_topic}, {config_topic}")
  411. # MQTT 状态变化时刷新当前屏幕显示
  412. panels = get_sorted_panels()
  413. if panels and screen_display.get_status().get('connected'):
  414. try:
  415. with screen_lock:
  416. current_idx = screen_current_panel_index
  417. if 0 <= current_idx < len(panels):
  418. screen_refresh_queue.put((panels[current_idx][0],))
  419. except Exception as e:
  420. logger.error(f"MQTT 状态变化刷新屏幕失败: {e}")
  421. except Exception as e:
  422. logger.error(f"处理MQTT状态时出错: {str(e)}")
  423. # ========== DTU MQTT协议处理函数 ==========
  424. def build_dtu_topic(*parts):
  425. """构建DTU MQTT主题"""
  426. prefix = dtu_config.get('topic_prefix', '线架系统')
  427. return '/'.join([prefix] + list(parts))
  428. def dtu_register(discovery_request_id=None):
  429. """发送DTU注册消息"""
  430. if not mqtt_client.get_status():
  431. logger.warning("MQTT未连接,无法发送注册消息")
  432. return False
  433. try:
  434. # 加载面板配置
  435. devices = address_config.get_stored_devices()
  436. # 优先使用 dtu_config 中预设的 panel_id 映射 (uid_hex -> panel_id)
  437. panel_id_map = dtu_config.get('panel_id_mapping', {}) or {}
  438. panels = []
  439. panel_id = 1
  440. for uid_hex, addr in devices.items():
  441. custom_id = panel_id_map.get(uid_hex) or panel_id_map.get(addr)
  442. if custom_id:
  443. pid = str(custom_id)
  444. else:
  445. pid = f"PANEL_{dtu_config['dtu_id']}_{addr}"
  446. panels.append({
  447. 'panel_id': pid,
  448. 'address': addr,
  449. 'position': panel_id
  450. })
  451. panel_id += 1
  452. # 检测网络类型
  453. network_type = 'ethernet'
  454. signal_strength = None
  455. try:
  456. interfaces = network_manager.get_status().get('interfaces', {})
  457. for name, info in interfaces.items():
  458. if name in ('lo', 'docker0') or name.startswith(('br-', 'veth')):
  459. continue
  460. if info.get('status') != 'UP':
  461. continue
  462. if name.startswith(('wlan', 'wlp')):
  463. network_type = 'wifi'
  464. # 读取信号强度 (dBm, 越接近 0 越强)
  465. try:
  466. import re
  467. with open('/proc/net/wireless') as f:
  468. for line in f.readlines()[2:]:
  469. m = re.match(r'\s*(\S+):.*\s(\d+)\.\s+(-?\d+)\.', line)
  470. if m:
  471. signal_strength = int(m.group(3))
  472. break
  473. except Exception:
  474. pass
  475. elif name.startswith(('wwan', 'ppp', 'usb', 'eth')):
  476. # 移动网卡
  477. if name.startswith(('wwan', 'ppp', 'usb')):
  478. network_type = 'cellular'
  479. break
  480. except Exception:
  481. pass
  482. payload = {
  483. 'msg_id': f"reg_{int(time.time() * 1000)}_{uuid.uuid4().hex[:6]}",
  484. 'timestamp': int(time.time() * 1000),
  485. 'dtu_id': dtu_config['dtu_id'],
  486. 'type': 'REGISTER',
  487. 'payload': {
  488. 'firmware_version': dtu_config.get('firmware_version', 'v1.0.0'),
  489. 'hardware_version': dtu_config.get('hardware_version', 'v1.0'),
  490. 'panel_count': len(panels),
  491. 'network_type': network_type,
  492. 'signal_strength': signal_strength,
  493. 'uptime': int(time.time() * 1000),
  494. 'panels': panels,
  495. 'discovery_request_id': discovery_request_id
  496. }
  497. }
  498. if discovery_request_id:
  499. payload['payload']['discovery_request_id'] = discovery_request_id
  500. topic = build_dtu_topic(dtu_config['customer_id'], 'dtu', dtu_config['dtu_id'], 'register')
  501. success, msg = mqtt_client.publish(topic, json.dumps(payload), qos=1)
  502. if success:
  503. logger.info(f"DTU注册消息已发送: {topic}")
  504. else:
  505. logger.error(f"DTU注册消息发送失败: {msg}")
  506. return success
  507. except Exception as e:
  508. logger.error(f"发送DTU注册消息失败: {str(e)}")
  509. return False
  510. _discover_dedup = {}
  511. def handle_broadcast_discover(payload):
  512. """处理广播发现消息 (7.14)"""
  513. try:
  514. if not isinstance(payload, dict) or payload.get('type') != 'DISCOVER':
  515. return
  516. p = payload.get('payload', {})
  517. cust_id = p.get('customer_id', '')
  518. if cust_id and cust_id != dtu_config.get('customer_id'):
  519. return
  520. request_id = p.get('request_id', '')
  521. now = time.time()
  522. if request_id and request_id in _discover_dedup:
  523. if now - _discover_dedup[request_id] < 60:
  524. return
  525. _discover_dedup[request_id] = now
  526. # 清理超过 60s 的旧条目,防止内存泄漏
  527. if len(_discover_dedup) > 256:
  528. for k in list(_discover_dedup.keys()):
  529. if now - _discover_dedup[k] > 60:
  530. _discover_dedup.pop(k, None)
  531. logger.info(f"收到DISCOVER广播,响应注册消息 (request_id={request_id})")
  532. import random
  533. delay = random.uniform(0, 3.0)
  534. time.sleep(delay)
  535. dtu_register(discovery_request_id=request_id)
  536. dtu_publish_status()
  537. except Exception as e:
  538. logger.error(f"处理广播发现消息失败: {str(e)}")
  539. def handle_broadcast_config(payload):
  540. """处理广播批量配置消息 (7.15)"""
  541. try:
  542. if isinstance(payload, dict) and payload.get('type') == 'CONFIG':
  543. cfg = payload.get('payload', {})
  544. cust_id = cfg.get('customer_id', '')
  545. if cust_id and cust_id != dtu_config.get('customer_id'):
  546. return
  547. config_version = cfg.get('config_version', '')
  548. force_apply = cfg.get('force_apply', False)
  549. items = cfg.get('items', {})
  550. local_version = dtu_config.get('config_version', '')
  551. incoming_msg_id = payload.get('msg_id', '')
  552. if local_version == config_version and not force_apply:
  553. logger.info(f"配置版本 {config_version} 已存在,跳过")
  554. _send_config_response(config_version, [], True, 1011, incoming_msg_id)
  555. return
  556. valid_ranges = {
  557. 'heartbeat_interval': (10, 3600),
  558. 'poll_interval_ms': (100, 10000),
  559. 'modbus_timeout_ms': (200, 5000),
  560. 'mqtt_keepalive': (10, 300),
  561. 'event_debounce_ms': (0, 5000),
  562. }
  563. applied = []
  564. for key, val in items.items():
  565. if key in valid_ranges:
  566. lo, hi = valid_ranges[key]
  567. if not (lo <= val <= hi):
  568. _send_config_response(config_version, applied, False, 1012, incoming_msg_id)
  569. return
  570. if key == 'heartbeat_interval':
  571. dtu_config['heartbeat_interval'] = val
  572. elif key == 'mqtt_keepalive':
  573. dtu_config['mqtt_keepalive'] = val
  574. elif key == 'event_debounce_ms':
  575. dtu_config['event_debounce_ms'] = val
  576. elif key == 'poll_interval_ms':
  577. dtu_config['poll_interval_ms'] = val
  578. elif key == 'modbus_timeout_ms':
  579. dtu_config['modbus_timeout_ms'] = val
  580. applied.append(key)
  581. elif key in ('alarm_enabled',):
  582. dtu_config['alarm_enabled'] = bool(val)
  583. applied.append(key)
  584. else:
  585. _send_config_response(config_version, applied, False, 1013, incoming_msg_id)
  586. return
  587. if applied:
  588. dtu_config['config_version'] = config_version
  589. try:
  590. save_dtu_config()
  591. except Exception as e:
  592. logger.error(f"配置写入持久化失败: {e}")
  593. _send_config_response(config_version, applied, False, 1014, incoming_msg_id)
  594. return
  595. _send_config_response(config_version, applied, True, 0, incoming_msg_id)
  596. except Exception as e:
  597. logger.error(f"处理广播配置消息失败: {str(e)}")
  598. def _send_config_response(config_version, applied_items, success, error_code, original_msg_id=''):
  599. """发送批量配置响应 (随机 0-2000ms 延迟避免广播风暴)"""
  600. import time as _t, json as _j, random as _r
  601. delay = _r.uniform(0, 2.0)
  602. def _delayed_publish():
  603. _t.sleep(delay)
  604. topic = build_dtu_topic(dtu_config['customer_id'], 'dtu', dtu_config['dtu_id'], 'response')
  605. payload = {
  606. 'msg_id': f"rsp_cfg_{int(_t.time() * 1000)}",
  607. 'timestamp': int(_t.time() * 1000),
  608. 'dtu_id': dtu_config['dtu_id'],
  609. 'type': 'RESPONSE',
  610. 'payload': {
  611. 'original_msg_id': original_msg_id or f"cfg_{config_version}",
  612. 'command': 'BROADCAST_CONFIG',
  613. 'success': success,
  614. 'result': {
  615. 'applied_version': config_version,
  616. 'applied_items': applied_items,
  617. 'restart_required': False
  618. } if success else None,
  619. 'error_code': error_code,
  620. 'error_message': get_ota_error_message(error_code) if error_code and error_code >= 1000 else None
  621. }
  622. }
  623. mqtt_client.publish(topic, _j.dumps(payload), qos=1)
  624. # 异步执行避免阻塞 MQTT 主循环
  625. threading.Thread(target=_delayed_publish, daemon=True).start()
  626. def dtu_publish_status(force=False):
  627. """发送DTU状态/心跳消息
  628. Args:
  629. force: 保留参数。事件/告警/查询触发时传入 True,用于立即上报(区别于定时心跳)
  630. """
  631. if not mqtt_client.get_status() or not dtu_config.get('enabled'):
  632. return False
  633. # 调试日志: 区分定时 vs 触发上报
  634. if force:
  635. logger.debug("dtu_publish_status(force=True) - 事件触发立即上报")
  636. try:
  637. # panel_online 计数: 在最近 60s 内所有 24 个端口都被成功轮询过
  638. # 严格语义: 全部 24 端口在最近一轮 poll 周期内被 poll 成功 → online
  639. now_ms = int(time.time() * 1000)
  640. threshold_ms = 60_000 # 60 秒未轮询视为 offline
  641. poll_ok = 0
  642. for pid, cfg in panel_config.items():
  643. ps = port_state.get(pid, {})
  644. expected_ports = cfg.get('port_count', 24)
  645. polled_recently = sum(
  646. 1 for p in ps.values()
  647. if p.get('last_polled_at') and (now_ms - p.get('last_polled_at', 0)) < threshold_ms
  648. )
  649. if polled_recently >= expected_ports:
  650. poll_ok += 1
  651. panel_online = poll_ok
  652. panel_offline = max(0, len(panel_config) - panel_online)
  653. # 检查串口状态
  654. serial_st = serial_client.get_status()
  655. rs485_status = 'normal' if (isinstance(serial_st, dict) and serial_st.get('connected', False)) else 'error'
  656. # 网络状态: 检查实际接口 UP 状态
  657. # 规则: 有线 (eth*/wlan*) UP → online; 仅 wifi 信号弱 (无 link/低质量) → weak; 否则 offline
  658. network_status = 'offline'
  659. try:
  660. interfaces = network_manager.get_status().get('interfaces', {})
  661. real_ifs = {i: s for i, s in interfaces.items()
  662. if s.get('status') == 'UP'
  663. and i not in ('lo', 'docker0')
  664. and not i.startswith(('br-', 'veth'))}
  665. if real_ifs:
  666. # 优先有线连接
  667. if any(i.startswith(('eth', 'en')) for i in real_ifs):
  668. network_status = 'online'
  669. else:
  670. # 仅无线: 检查信号强度 (Linux: /proc/net/wireless)
  671. try:
  672. import re
  673. with open('/proc/net/wireless') as f:
  674. lines = f.readlines()
  675. quality = None
  676. for line in lines[2:]: # 跳过表头
  677. m = re.match(r'\s*(\S+):.*\s(\d+)\.\s+(\d+)\.\s+(\d+)', line)
  678. if m:
  679. # link quality (字段 2, 0-70); signal level (字段 3, dBm, 越低越差)
  680. link = int(m.group(2))
  681. sig = int(m.group(3))
  682. # 估算: link < 30 或 sig < -75 dBm 视为弱
  683. if quality is None or link < quality:
  684. quality = link
  685. network_status = 'weak' if (quality is not None and quality < 30) else 'online'
  686. except Exception:
  687. network_status = 'online'
  688. except Exception:
  689. network_status = 'offline'
  690. # 屏幕连接状态:从 Nextion 屏幕实例读取
  691. screen_connected = screen_display.get_status().get('connected', False)
  692. # CPU/内存使用率 (Linux 才有意义)
  693. cpu_usage = None
  694. memory_usage = None
  695. dtu_temperature = None
  696. try:
  697. import psutil
  698. cpu_usage = psutil.cpu_percent(interval=None)
  699. memory_usage = psutil.virtual_memory().percent
  700. # 主板温度: 优先 psutil.sensors_temperatures
  701. try:
  702. temps = psutil.sensors_temperatures(fahrenheit=False) or {}
  703. for chip_name, entries in temps.items():
  704. for entry in entries:
  705. if entry.current and entry.current > 0:
  706. dtu_temperature = round(entry.current, 1)
  707. break
  708. if dtu_temperature is not None:
  709. break
  710. except Exception:
  711. pass
  712. # 回退: 直接读 /sys/class/thermal/thermal_zone*/temp
  713. if dtu_temperature is None:
  714. import glob as _glob
  715. for tz in sorted(_glob.glob('/sys/class/thermal/thermal_zone*/temp')):
  716. try:
  717. with open(tz) as f:
  718. v = int(f.read().strip())
  719. if v > 0:
  720. dtu_temperature = round(v / 1000.0, 1)
  721. break
  722. except Exception:
  723. continue
  724. except Exception:
  725. pass
  726. # 优先使用外部上报的 dtu_temperature (来自 MQTT/RS485 传感器)
  727. dtu_temperature = env_sensor_data.get('dtu_temperature') or dtu_temperature
  728. payload = {
  729. 'msg_id': f"hbt_{int(time.time() * 1000)}_{uuid.uuid4().hex[:6]}",
  730. 'timestamp': int(time.time() * 1000),
  731. 'dtu_id': dtu_config['dtu_id'],
  732. 'type': 'STATUS',
  733. 'payload': {
  734. 'cpu_usage': cpu_usage,
  735. 'memory_usage': memory_usage,
  736. 'temperature': dtu_temperature,
  737. 'network_status': network_status,
  738. 'panel_online': panel_online,
  739. 'panel_offline': panel_offline,
  740. 'screen_connected': screen_connected,
  741. 'mqtt_connected': mqtt_client.get_status().get('connected', False) if isinstance(mqtt_client.get_status(), dict) else bool(mqtt_client.get_status()),
  742. 'rs485_status': rs485_status
  743. }
  744. }
  745. topic = build_dtu_topic(dtu_config['customer_id'], 'dtu', dtu_config['dtu_id'], 'status')
  746. success, _ = mqtt_client.publish(topic, json.dumps(payload), qos=1)
  747. return success
  748. except Exception as e:
  749. logger.error(f"发送DTU状态消息失败: {str(e)}")
  750. return False
  751. def dtu_publish_event(panel_id, port_id, event_type, jumper_uid, previous_jumper_uid=None):
  752. """发送端口事件消息"""
  753. if not mqtt_client.get_status() or not dtu_config.get('enabled'):
  754. return False
  755. try:
  756. payload = {
  757. 'msg_id': f"evt_{int(time.time() * 1000)}",
  758. 'timestamp': int(time.time() * 1000),
  759. 'dtu_id': dtu_config['dtu_id'],
  760. 'type': 'EVENT',
  761. 'payload': {
  762. 'panel_id': panel_id,
  763. 'port_id': port_id,
  764. 'event_type': event_type,
  765. 'event_id': f"evt_{uuid.uuid4().hex[:8]}",
  766. 'jumper_uid': jumper_uid,
  767. 'previous_jumper_uid': previous_jumper_uid
  768. }
  769. }
  770. topic = build_dtu_topic(dtu_config['customer_id'], 'patchpanel', dtu_config['dtu_id'], panel_id, 'event')
  771. success, _ = mqtt_client.publish(topic, json.dumps(payload), qos=1)
  772. if success:
  773. logger.info(f"端口事件已发送: {panel_id}:{port_id} - {event_type}")
  774. # 记录到事件历史
  775. port_event_history.append({
  776. 'timestamp': int(time.time() * 1000),
  777. 'panel_id': panel_id,
  778. 'port_id': port_id,
  779. 'event_type': event_type,
  780. 'jumper_uid': jumper_uid,
  781. 'previous_jumper_uid': previous_jumper_uid
  782. })
  783. # 限制历史记录数量
  784. if len(port_event_history) > 1000:
  785. port_event_history[:] = port_event_history[-1000:]
  786. # 事件驱动: 立即上报一次 status (符合 7.4 "状态变化时立即上报")
  787. try:
  788. dtu_publish_status(force=True)
  789. except Exception:
  790. pass
  791. return success
  792. except Exception as e:
  793. logger.error(f"发送端口事件失败: {str(e)}")
  794. return False
  795. def dtu_publish_alarm(panel_id, port_id, alarm_type, expected_jumper_uid, actual_jumper_uid, severity='WARNING'):
  796. """发送非法告警消息"""
  797. if not mqtt_client.get_status() or not dtu_config.get('enabled'):
  798. return False
  799. try:
  800. description = f"端口{port_id}期望跳线{expected_jumper_uid},实际{'未读到' if not actual_jumper_uid else actual_jumper_uid}"
  801. payload = {
  802. 'msg_id': f"alm_{int(time.time() * 1000)}",
  803. 'timestamp': int(time.time() * 1000),
  804. 'dtu_id': dtu_config['dtu_id'],
  805. 'type': 'ALARM',
  806. 'payload': {
  807. 'panel_id': panel_id,
  808. 'port_id': port_id,
  809. 'alarm_type': alarm_type,
  810. 'severity': severity,
  811. 'expected_jumper_uid': expected_jumper_uid,
  812. 'actual_jumper_uid': actual_jumper_uid,
  813. 'description': description
  814. }
  815. }
  816. topic = build_dtu_topic(dtu_config['customer_id'], 'patchpanel', dtu_config['dtu_id'], panel_id, 'alarm')
  817. success, _ = mqtt_client.publish(topic, json.dumps(payload), qos=1)
  818. if success:
  819. logger.warning(f"非法告警已发送: {panel_id}:{port_id} - {alarm_type}")
  820. # 记录到事件历史
  821. port_event_history.append({
  822. 'timestamp': int(time.time() * 1000),
  823. 'panel_id': panel_id,
  824. 'port_id': port_id,
  825. 'event_type': 'ALARM',
  826. 'alarm_type': alarm_type,
  827. 'expected_jumper_uid': expected_jumper_uid,
  828. 'actual_jumper_uid': actual_jumper_uid
  829. })
  830. # 限制历史记录数量
  831. if len(port_event_history) > 1000:
  832. port_event_history[:] = port_event_history[-1000:]
  833. # 事件驱动: 告警后立即上报一次 status
  834. try:
  835. dtu_publish_status(force=True)
  836. except Exception:
  837. pass
  838. return success
  839. except Exception as e:
  840. logger.error(f"发送非法告警失败: {str(e)}")
  841. return False
  842. def dtu_publish_panel_status(panel_id, address, ports_data):
  843. """发送配线架状态消息 (7.7)"""
  844. if not mqtt_client.get_status() or not dtu_config.get('enabled'):
  845. return False
  846. try:
  847. topic = build_dtu_topic(dtu_config['customer_id'], 'patchpanel', dtu_config['dtu_id'], panel_id, 'status')
  848. payload = {
  849. 'msg_id': f"pst_{int(time.time() * 1000)}",
  850. 'timestamp': int(time.time() * 1000),
  851. 'dtu_id': dtu_config['dtu_id'],
  852. 'type': 'STATUS',
  853. 'payload': {
  854. 'panel_id': panel_id,
  855. 'position': panel_config.get(panel_id, {}).get('position', 0),
  856. 'address': address,
  857. 'online': True,
  858. 'last_poll_time': int(time.time() * 1000),
  859. 'ports': ports_data
  860. }
  861. }
  862. mqtt_client.publish(topic, json.dumps(payload), qos=0)
  863. return True
  864. except Exception as e:
  865. logger.error(f"发送面板状态失败: {str(e)}")
  866. return False
  867. def dtu_publish_jumper_status():
  868. """发送跳线状态汇总 (7.8) - STATUS envelope with panels[] array"""
  869. if not mqtt_client.get_status() or not dtu_config.get('enabled'):
  870. return False
  871. try:
  872. panels_arr = []
  873. for panel_id, cfg in panel_config.items():
  874. ports_arr = []
  875. ps = port_state.get(panel_id, {})
  876. for port_id in range(1, 25):
  877. p = ps.get(port_id, {})
  878. last_uid = p.get('last_uid')
  879. expected_uid = p.get('expected_uid')
  880. last_polled = p.get('last_polled_at')
  881. if not last_polled:
  882. status = 'UNKNOWN'
  883. elif last_uid is None:
  884. status = 'DISCONNECTED'
  885. if expected_uid:
  886. status = 'ILLEGAL'
  887. elif expected_uid and last_uid != expected_uid:
  888. status = 'ILLEGAL'
  889. else:
  890. status = 'CONNECTED'
  891. ports_arr.append({
  892. 'port_id': port_id,
  893. 'status': status,
  894. 'jumper_uid': last_uid
  895. })
  896. panels_arr.append({
  897. 'panel_id': panel_id,
  898. 'position': cfg.get('position', 0),
  899. 'ports': ports_arr
  900. })
  901. topic = build_dtu_topic(dtu_config['customer_id'], 'jumper', dtu_config['dtu_id'], 'status')
  902. payload = {
  903. 'msg_id': f"jmp_{int(time.time() * 1000)}",
  904. 'timestamp': int(time.time() * 1000),
  905. 'dtu_id': dtu_config['dtu_id'],
  906. 'type': 'STATUS',
  907. 'payload': {'panels': panels_arr}
  908. }
  909. mqtt_client.publish(topic, json.dumps(payload), qos=0)
  910. return True
  911. except Exception as e:
  912. logger.error(f"发送跳线状态汇总失败: {str(e)}")
  913. return False
  914. _env_sensor_data = {
  915. 'temperature': None,
  916. 'humidity': None,
  917. 'sensor_update_time': None
  918. }
  919. def update_env_sensor_data(temperature, humidity, dtu_temperature=None, sensor_update_time=None):
  920. """更新环境传感器数据 - 委托给完整实现(带历史和告警)"""
  921. update_env_sensor_data_full(temperature, humidity, dtu_temperature, sensor_update_time)
  922. def dtu_publish_env_sensor():
  923. """发送环境传感器数据 (7.9) - STATUS envelope"""
  924. if not mqtt_client.get_status() or not dtu_config.get('enabled'):
  925. return False
  926. try:
  927. topic = build_dtu_topic(dtu_config['customer_id'], 'env', dtu_config['dtu_id'], 'sensor')
  928. payload = {
  929. 'msg_id': f"sen_{int(time.time() * 1000)}",
  930. 'timestamp': int(time.time() * 1000),
  931. 'dtu_id': dtu_config['dtu_id'],
  932. 'type': 'STATUS',
  933. 'payload': {
  934. 'temperature': _env_sensor_data.get('temperature'),
  935. 'humidity': _env_sensor_data.get('humidity'),
  936. 'sensor_update_time': _env_sensor_data.get('sensor_update_time')
  937. }
  938. }
  939. mqtt_client.publish(topic, json.dumps(payload), qos=0)
  940. return True
  941. except Exception as e:
  942. logger.error(f"发送环境传感器数据失败: {str(e)}")
  943. return False
  944. def dtu_handle_control(topic, payload):
  945. """处理下行控制指令"""
  946. try:
  947. if not dtu_config.get('enabled'):
  948. return
  949. command = payload.get('payload', {}).get('command')
  950. target = payload.get('payload', {}).get('target')
  951. params = payload.get('payload', {}).get('params', {})
  952. logger.info(f"收到控制指令: command={command}, target={target}")
  953. response_payload = {
  954. 'msg_id': f"rsp_{int(time.time() * 1000)}_{uuid.uuid4().hex[:6]}",
  955. 'timestamp': int(time.time() * 1000),
  956. 'dtu_id': dtu_config['dtu_id'],
  957. 'type': 'RESPONSE',
  958. 'payload': {
  959. 'original_msg_id': payload.get('msg_id'),
  960. 'command': command,
  961. 'success': True,
  962. 'result': {},
  963. 'error_code': 0,
  964. 'error_message': None
  965. }
  966. }
  967. topic_response_dtu = build_dtu_topic(
  968. dtu_config['customer_id'], 'dtu', dtu_config['dtu_id'], 'response')
  969. # 处理各命令
  970. if command == 'SET_PORT_LED':
  971. # 设置端口LED
  972. port_id = params.get('port_id')
  973. led_mode = params.get('led_mode', 'OFF')
  974. # 参数校验
  975. if port_id is None or not (1 <= int(port_id) <= 24):
  976. response_payload['payload']['success'] = False
  977. response_payload['payload']['error_code'] = 1003
  978. response_payload['payload']['error_message'] = f"端口号非法: {port_id}"
  979. mqtt_client.publish(topic_response_dtu, json.dumps(response_payload), qos=1)
  980. return
  981. if led_mode not in ('OFF', 'BLINK_RED', 'BLINK_GREEN', 'BLINK_BLUE'):
  982. response_payload['payload']['success'] = False
  983. response_payload['payload']['error_code'] = 1004
  984. response_payload['payload']['error_message'] = f"LED模式非法: {led_mode}"
  985. mqtt_client.publish(topic_response_dtu, json.dumps(response_payload), qos=1)
  986. return
  987. # LED模式映射
  988. led_mode_map = {'OFF': 0, 'BLINK_RED': 1, 'BLINK_GREEN': 2, 'BLINK_BLUE': 3}
  989. color = led_mode_map.get(led_mode, 0)
  990. # 查找设备地址
  991. device_address = 1 # 默认
  992. for panel_id, cfg in panel_config.items():
  993. if panel_id == target:
  994. device_address = cfg.get('address', 1)
  995. break
  996. result = modbus_client.set_rgb_led(device_address, port_id, color)
  997. response_payload['payload']['success'] = 'error' not in result
  998. response_payload['payload']['result'] = result
  999. if 'error' in result:
  1000. response_payload['payload']['error_code'] = 1005
  1001. response_payload['payload']['error_message'] = result['error']
  1002. elif command == 'QUERY_DTU_STATUS':
  1003. # 查询DTU状态
  1004. dtu_publish_status(force=True)
  1005. elif command == 'SYNC_PORT_MAPPING':
  1006. # 同步单端口期望映射
  1007. port_id = params.get('port_id')
  1008. jumper_uid = params.get('jumper_uid')
  1009. if port_id is None or target is None or target not in panel_config:
  1010. response_payload['payload']['success'] = False
  1011. response_payload['payload']['error_code'] = 1006 if port_id is None else 1002
  1012. response_payload['payload']['error_message'] = "参数缺失" if port_id is None else f"目标面板不存在: {target}"
  1013. mqtt_client.publish(topic_response_dtu, json.dumps(response_payload), qos=1)
  1014. return
  1015. if target not in port_state:
  1016. port_state[target] = {}
  1017. if port_id not in port_state[target]:
  1018. port_state[target][port_id] = {'last_uid': None, 'expected_uid': None, 'alarm_count': 0}
  1019. port_state[target][port_id]['expected_uid'] = jumper_uid
  1020. save_dtu_config()
  1021. elif command == 'SYNC_ALL_MAPPING':
  1022. # 批量同步期望映射
  1023. mappings = params.get('mappings', [])
  1024. if not isinstance(mappings, list):
  1025. response_payload['payload']['success'] = False
  1026. response_payload['payload']['error_code'] = 1006
  1027. response_payload['payload']['error_message'] = "mappings 字段格式错误"
  1028. mqtt_client.publish(topic_response_dtu, json.dumps(response_payload), qos=1)
  1029. return
  1030. try:
  1031. for mapping in mappings:
  1032. panel_id = mapping.get('panel_id')
  1033. port_id = mapping.get('port_id')
  1034. jumper_uid = mapping.get('jumper_uid')
  1035. if panel_id not in panel_config:
  1036. raise ValueError(f"目标面板不存在: {panel_id}")
  1037. if panel_id not in port_state:
  1038. port_state[panel_id] = {}
  1039. if port_id not in port_state[panel_id]:
  1040. port_state[panel_id][port_id] = {'last_uid': None, 'expected_uid': None, 'alarm_count': 0}
  1041. port_state[panel_id][port_id]['expected_uid'] = jumper_uid
  1042. save_dtu_config()
  1043. except Exception as e:
  1044. response_payload['payload']['success'] = False
  1045. response_payload['payload']['error_code'] = 1007
  1046. response_payload['payload']['error_message'] = f"同步映射失败: {e}"
  1047. mqtt_client.publish(topic_response_dtu, json.dumps(response_payload), qos=1)
  1048. return
  1049. elif command == 'REBOOT':
  1050. # 重启DTU: 1s 延迟后执行 reboot
  1051. import subprocess
  1052. logger.warning(f"MQTT REBOOT 命令收到,1 秒后执行系统重启")
  1053. threading.Thread(target=lambda: (
  1054. time.sleep(1),
  1055. subprocess.run(['reboot'], capture_output=True)
  1056. ), daemon=True).start()
  1057. response_payload['payload']['result'] = {'message': 'Reboot scheduled in 1s'}
  1058. elif command == 'READ_PANEL_STATUS':
  1059. # 读取面板状态
  1060. # 读取面板的多个寄存器获取状态信息
  1061. # 寄存器地址定义: 0x0000=运行状态, 0x0001=LED控制, 0x0002-0x0009=天线卡状态
  1062. device_address = 1
  1063. for panel_id, cfg in panel_config.items():
  1064. if panel_id == target:
  1065. device_address = cfg.get('address', 1)
  1066. break
  1067. # 读取面板状态寄存器 (地址0x0000开始,读取10个寄存器)
  1068. status_result = modbus_client.read_holding_registers(device_address, 0x0000, 10)
  1069. # 解析状态
  1070. registers = status_result.get('registers', [])
  1071. panel_status = {
  1072. 'device_address': device_address,
  1073. 'run_status': registers[0] if len(registers) > 0 else None,
  1074. 'led_control': registers[1] if len(registers) > 1 else None,
  1075. 'antenna_status': registers[2:10] if len(registers) >= 10 else [],
  1076. 'raw_registers': registers
  1077. }
  1078. response_payload['payload']['success'] = 'error' not in status_result
  1079. response_payload['payload']['result'] = panel_status
  1080. elif command == 'QUERY_JUMPER_STATUS':
  1081. # 查询跳线器状态汇总 (7.8) - 触发 STATUS 上行至 jumper/{dtu_id}/status
  1082. dtu_publish_jumper_status()
  1083. return # dtu_publish_jumper_status 已发送响应
  1084. elif command == 'QUERY_ENV_SENSOR':
  1085. # 查询环境传感器 (7.9) - 触发 STATUS 上行至 env/{dtu_id}/sensor
  1086. dtu_publish_env_sensor()
  1087. return # dtu_publish_env_sensor 已发送响应
  1088. elif command == 'OTA_UPGRADE':
  1089. params = payload.get('payload', {}).get('params', {})
  1090. target_version = params.get('firmware_version')
  1091. force_upgrade = params.get('force_upgrade', False)
  1092. current_version = dtu_config.get('firmware_version', 'v1.0.0')
  1093. incoming_msg_id = payload.get('msg_id', '')
  1094. if target_version == current_version and not force_upgrade:
  1095. response_payload['payload']['success'] = False
  1096. response_payload['payload']['error_code'] = 1021
  1097. response_payload['payload']['error_message'] = f"已是目标版本: {target_version}"
  1098. mqtt_client.publish(topic_response_dtu, json.dumps(response_payload), qos=1)
  1099. return
  1100. if params.get('firmware_url'):
  1101. import threading
  1102. threading.Thread(target=lambda: _run_ota({**params, '_incoming_msg_id': incoming_msg_id}), daemon=True).start()
  1103. response_payload['payload']['success'] = True
  1104. response_payload['payload']['result'] = {
  1105. 'firmware_version': target_version,
  1106. 'ota_status': 'DOWNLOADING',
  1107. 'restart_required': False
  1108. }
  1109. elif command == 'OTA_CANCEL':
  1110. # 取消OTA升级
  1111. # 向OTA控制寄存器写入取消命令 (0x0000=取消)
  1112. device_address = 1
  1113. for panel_id, cfg in panel_config.items():
  1114. if panel_id == target:
  1115. device_address = cfg.get('address', 1)
  1116. break
  1117. # OTA控制寄存器地址: 0x0300
  1118. # 值: 0x0000=取消升级
  1119. cancel_result = modbus_client.write_single_register(device_address, 0x0300, 0x0000)
  1120. response_payload['payload']['success'] = 'error' not in cancel_result
  1121. response_payload['payload']['result'] = {
  1122. 'message': 'OTA upgrade cancelled' if 'error' not in cancel_result else 'Failed to cancel OTA',
  1123. 'raw_result': cancel_result
  1124. }
  1125. else:
  1126. # 未知命令
  1127. response_payload['payload']['success'] = False
  1128. response_payload['payload']['error_code'] = 1001
  1129. response_payload['payload']['error_message'] = f"未知命令: {command}"
  1130. mqtt_client.publish(topic_response_dtu, json.dumps(response_payload), qos=1)
  1131. return
  1132. # 发送响应(QUERY_DTU_STATUS 已由 dtu_publish_status 发送;QUERY_JUMPER_STATUS/QUERY_ENV_SENSOR 已在分支中发布)
  1133. if command == 'QUERY_DTU_STATUS':
  1134. pass
  1135. else:
  1136. cmd = command
  1137. if cmd == 'READ_PANEL_STATUS':
  1138. if target and target != 'all' and target in panel_config:
  1139. # 单个面板: 使用 panel_id 作为 topic segment
  1140. topic = build_dtu_topic(dtu_config['customer_id'], 'patchpanel',
  1141. dtu_config['dtu_id'], target, 'status')
  1142. mqtt_client.publish(topic, json.dumps(response_payload), qos=0)
  1143. else:
  1144. # target=all: 遍历所有面板发布到各自的 topic
  1145. for pid in panel_config.keys():
  1146. topic = build_dtu_topic(dtu_config['customer_id'], 'patchpanel',
  1147. dtu_config['dtu_id'], pid, 'status')
  1148. mqtt_client.publish(topic, json.dumps(response_payload), qos=0)
  1149. else:
  1150. topic = build_dtu_topic(dtu_config['customer_id'], 'dtu', dtu_config['dtu_id'], 'response')
  1151. p = response_payload['payload']
  1152. if not p.get('success', True) and p.get('error_code', 0) == 0:
  1153. p['error_code'] = 9999
  1154. p['error_message'] = p.get('error_message') or '未知错误'
  1155. mqtt_client.publish(topic, json.dumps(response_payload), qos=1)
  1156. except Exception as e:
  1157. logger.error(f"处理控制指令失败: {str(e)}")
  1158. # 启动DTU心跳定时器
  1159. def start_dtu_heartbeat():
  1160. """启动DTU心跳定时任务"""
  1161. def heartbeat_task():
  1162. env_tick = 0
  1163. while True:
  1164. socketio.sleep(dtu_config.get('heartbeat_interval', DTU_HEARTBEAT_INTERVAL))
  1165. if mqtt_client.get_status() and dtu_config.get('enabled'):
  1166. dtu_publish_status()
  1167. # 每 2 个心跳周期 (120s) 主动推送一次环境传感器数据 (7.9)
  1168. env_tick += 1
  1169. if env_tick >= 2:
  1170. env_tick = 0
  1171. dtu_publish_env_sensor()
  1172. socketio.start_background_task(target=heartbeat_task)
  1173. # 修改mqtt_data_handler以处理控制指令
  1174. def mqtt_data_handler_extended(data):
  1175. """处理MQTT接收的数据(扩展版,含DTU协议)"""
  1176. try:
  1177. topic = data.get('topic', '')
  1178. payload_str = data.get('payload', '')
  1179. # 尝试解析JSON
  1180. try:
  1181. payload = json.loads(payload_str) if isinstance(payload_str, str) else payload_str
  1182. except:
  1183. payload = payload_str
  1184. # 添加到缓冲区
  1185. timestamp = time.strftime('%Y-%m-%d %H:%M:%S')
  1186. mqtt_data_buffer.append({
  1187. 'timestamp': timestamp,
  1188. 'topic': topic,
  1189. 'payload': payload_str
  1190. })
  1191. if len(mqtt_data_buffer) > MAX_BUFFER_SIZE:
  1192. mqtt_data_buffer.pop(0)
  1193. # 通过WebSocket广播
  1194. socketio.emit('mqtt_data', {
  1195. 'timestamp': timestamp,
  1196. 'topic': topic,
  1197. 'payload': payload_str
  1198. }, namespace=SOCKETIO_NAMESPACE_DATA)
  1199. # 检查是否是控制指令主题
  1200. expected_control_topic = build_dtu_topic(dtu_config['customer_id'], 'dtu', dtu_config['dtu_id'], 'control')
  1201. if topic == expected_control_topic:
  1202. dtu_handle_control(topic, payload)
  1203. # 检查是否是状态主题,用于处理OTA状态上报
  1204. expected_status_topic = build_dtu_topic(dtu_config['customer_id'], 'dtu', dtu_config['dtu_id'], 'status')
  1205. if topic == expected_status_topic and isinstance(payload, dict):
  1206. if 'ota_status' in payload or 'firmware_version' in payload:
  1207. handle_ota_status(dtu_config['dtu_id'], payload)
  1208. # 检查是否是环境传感器数据主题 (兼容旧 sensor/.../data 路径,spec 7.9 为 env/.../sensor)
  1209. sensor_topic = build_dtu_topic(dtu_config['customer_id'], 'sensor', dtu_config['dtu_id'], 'data')
  1210. env_topic = build_dtu_topic(dtu_config['customer_id'], 'env', dtu_config['dtu_id'], 'sensor')
  1211. if topic in (sensor_topic, env_topic) and isinstance(payload, dict):
  1212. if 'temperature' in payload or 'humidity' in payload:
  1213. update_env_sensor_data(
  1214. payload.get('temperature'),
  1215. payload.get('humidity'),
  1216. payload.get('dtu_temperature'),
  1217. payload.get('sensor_update_time')
  1218. )
  1219. # 检查是否是广播主题 - DTU发现 (7.14)
  1220. discover_topic = build_dtu_topic('broadcast', 'dtu', 'discover')
  1221. if topic == discover_topic:
  1222. handle_broadcast_discover(payload)
  1223. # 检查是否是广播主题 - 批量配置 (7.15)
  1224. config_topic = build_dtu_topic('broadcast', 'dtu', 'config')
  1225. if topic == config_topic:
  1226. handle_broadcast_config(payload)
  1227. # 转发到串口(如果启用)
  1228. if forward_mqtt_to_serial and serial_client.get_status():
  1229. success, msg = serial_client.send_data(payload_str)
  1230. if not success:
  1231. logger.warning(f"MQTT数据转发到串口失败: {msg}")
  1232. except Exception as e:
  1233. logger.error(f"处理MQTT数据时出错: {str(e)}")
  1234. # WebSocket事件处理
  1235. @socketio.on('connect', namespace=SOCKETIO_NAMESPACE_DATA)
  1236. def handle_data_connect():
  1237. """处理数据命名空间的连接"""
  1238. try:
  1239. client_id = str(uuid.uuid4())
  1240. connected_clients['data'].add(client_id)
  1241. logger.info(f'客户端已连接到数据命名空间,当前连接数: {len(connected_clients["data"])}')
  1242. # 发送当前的缓冲区数据,限制发送的历史记录数量
  1243. max_history = 100 # 限制发送的历史记录数量以提高性能
  1244. emit('serial_data_history', {'data': serial_data_buffer[-max_history:]})
  1245. emit('mqtt_data_history', {'data': mqtt_data_buffer[-max_history:]})
  1246. # 存储客户端ID以便断开连接时使用
  1247. socketio.start_background_task(target=lambda: None) # 确保上下文可用
  1248. except Exception as e:
  1249. logger.error(f"处理数据命名空间连接时出错: {str(e)}")
  1250. @socketio.on('connect', namespace=SOCKETIO_NAMESPACE_STATUS)
  1251. def handle_status_connect():
  1252. """处理状态命名空间的连接"""
  1253. try:
  1254. client_id = str(uuid.uuid4())
  1255. connected_clients['status'].add(client_id)
  1256. logger.info(f'客户端已连接到状态命名空间,当前连接数: {len(connected_clients["status"])}')
  1257. # 发送当前状态
  1258. emit('serial_status', {'connected': serial_status})
  1259. emit('mqtt_status', {'connected': mqtt_status})
  1260. emit('forward_status', {
  1261. 'serial_to_mqtt': forward_serial_to_mqtt,
  1262. 'mqtt_to_serial': forward_mqtt_to_serial,
  1263. 'publish_topic': mqtt_publish_topic
  1264. })
  1265. except Exception as e:
  1266. logger.error(f"处理状态命名空间连接时出错: {str(e)}")
  1267. @socketio.on('connect', namespace=SOCKETIO_NAMESPACE_CONTROL)
  1268. def handle_control_connect():
  1269. """处理控制命名空间的连接"""
  1270. try:
  1271. client_id = str(uuid.uuid4())
  1272. connected_clients['control'].add(client_id)
  1273. logger.info(f'客户端已连接到控制命名空间,当前连接数: {len(connected_clients["control"])}')
  1274. except Exception as e:
  1275. logger.error(f"处理控制命名空间连接时出错: {str(e)}")
  1276. @socketio.on('disconnect', namespace=SOCKETIO_NAMESPACE_DATA)
  1277. def handle_data_disconnect():
  1278. """处理数据命名空间的断开连接"""
  1279. try:
  1280. # 清理客户端连接记录
  1281. # 在实际应用中,可能需要更复杂的逻辑来追踪具体哪个客户端断开了连接
  1282. if len(connected_clients['data']) > 0:
  1283. # 这里简化处理,实际应该维护session到client_id的映射
  1284. connected_clients['data'].pop() # 注意:这是一个简化的实现
  1285. logger.info(f'客户端已断开数据命名空间的连接,当前连接数: {len(connected_clients["data"])}')
  1286. except Exception as e:
  1287. logger.error(f"处理数据命名空间断开连接时出错: {str(e)}")
  1288. @socketio.on('disconnect', namespace=SOCKETIO_NAMESPACE_STATUS)
  1289. def handle_status_disconnect():
  1290. """处理状态命名空间的断开连接"""
  1291. try:
  1292. # 清理客户端连接记录
  1293. if len(connected_clients['status']) > 0:
  1294. connected_clients['status'].pop()
  1295. logger.info(f'客户端已断开状态命名空间的连接,当前连接数: {len(connected_clients["status"])}')
  1296. except Exception as e:
  1297. logger.error(f"处理状态命名空间断开连接时出错: {str(e)}")
  1298. @socketio.on('disconnect', namespace=SOCKETIO_NAMESPACE_CONTROL)
  1299. def handle_control_disconnect():
  1300. """处理控制命名空间的断开连接"""
  1301. try:
  1302. # 清理客户端连接记录
  1303. if len(connected_clients['control']) > 0:
  1304. connected_clients['control'].pop()
  1305. logger.info(f'客户端已断开控制命名空间的连接,当前连接数: {len(connected_clients["control"])}')
  1306. except Exception as e:
  1307. logger.error(f"处理控制命名空间断开连接时出错: {str(e)}")
  1308. @socketio.on('serial_send', namespace=SOCKETIO_NAMESPACE_CONTROL)
  1309. def handle_serial_send(data):
  1310. """通过WebSocket处理串口发送数据请求"""
  1311. try:
  1312. message = data.get('message', '')
  1313. if not message:
  1314. emit('serial_send_response', {
  1315. 'success': False,
  1316. 'message': '消息内容不能为空',
  1317. 'error_code': ERROR_CODES['CONFIG_ERROR']
  1318. })
  1319. return
  1320. if not (isinstance(serial_client.get_status(), dict) and serial_client.get_status().get("connected", False)):
  1321. emit('serial_send_response', {
  1322. 'success': False,
  1323. 'message': '串口未连接',
  1324. 'error_code': ERROR_CODES['SERIAL_CONNECTION_ERROR']
  1325. })
  1326. return
  1327. success, msg = serial_client.send_data(message)
  1328. emit('serial_send_response', {
  1329. 'success': success,
  1330. 'message': msg,
  1331. 'error_code': ERROR_CODES['SUCCESS'] if success else ERROR_CODES['SERIAL_SEND_ERROR']
  1332. })
  1333. if success:
  1334. logger.info(f"通过WebSocket发送串口数据成功: {message[:50]}..." if len(message) > 50 else message)
  1335. except Exception as e:
  1336. error_msg = f"处理串口发送请求时出错: {str(e)}"
  1337. logger.error(error_msg)
  1338. emit('serial_send_response', {
  1339. 'success': False,
  1340. 'message': error_msg,
  1341. 'error_code': ERROR_CODES['UNKNOWN_ERROR']
  1342. })
  1343. @socketio.on('mqtt_publish', namespace=SOCKETIO_NAMESPACE_CONTROL)
  1344. def handle_mqtt_publish(data):
  1345. """通过WebSocket处理MQTT发布数据请求"""
  1346. try:
  1347. topic = data.get('topic', '')
  1348. message = data.get('message', '')
  1349. if not topic or not message:
  1350. emit('mqtt_publish_response', {
  1351. 'success': False,
  1352. 'message': '主题或消息内容不能为空',
  1353. 'error_code': ERROR_CODES['CONFIG_ERROR']
  1354. })
  1355. return
  1356. if not mqtt_client.get_status():
  1357. emit('mqtt_publish_response', {
  1358. 'success': False,
  1359. 'message': 'MQTT未连接',
  1360. 'error_code': ERROR_CODES['MQTT_CONNECTION_ERROR']
  1361. })
  1362. return
  1363. success, msg = mqtt_client.publish(topic, message)
  1364. emit('mqtt_publish_response', {
  1365. 'success': success,
  1366. 'message': msg,
  1367. 'error_code': ERROR_CODES['SUCCESS'] if success else ERROR_CODES['MQTT_PUBLISH_ERROR']
  1368. })
  1369. if success:
  1370. logger.info(f"通过WebSocket发布MQTT消息成功: 主题={topic}, 消息={message[:50]}..." if len(message) > 50 else message)
  1371. except Exception as e:
  1372. error_msg = f"处理MQTT发布请求时出错: {str(e)}"
  1373. logger.error(error_msg)
  1374. emit('mqtt_publish_response', {
  1375. 'success': False,
  1376. 'message': error_msg,
  1377. 'error_code': ERROR_CODES['UNKNOWN_ERROR']
  1378. })
  1379. @socketio.on('update_forward_config', namespace=SOCKETIO_NAMESPACE_CONTROL)
  1380. def handle_update_forward_config(data):
  1381. """通过WebSocket更新转发配置"""
  1382. global forward_serial_to_mqtt, forward_mqtt_to_serial, mqtt_publish_topic
  1383. try:
  1384. # 更新转发标志
  1385. if 'serial_to_mqtt' in data:
  1386. forward_serial_to_mqtt = bool(data['serial_to_mqtt'])
  1387. if 'mqtt_to_serial' in data:
  1388. forward_mqtt_to_serial = bool(data['mqtt_to_serial'])
  1389. if 'publish_topic' in data:
  1390. new_topic = str(data['publish_topic'])
  1391. if not new_topic.strip():
  1392. raise ValueError("发布主题不能为空")
  1393. mqtt_publish_topic = new_topic
  1394. # 广播配置更新
  1395. socketio.emit('forward_status', {
  1396. 'serial_to_mqtt': forward_serial_to_mqtt,
  1397. 'mqtt_to_serial': forward_mqtt_to_serial,
  1398. 'publish_topic': mqtt_publish_topic
  1399. }, namespace=SOCKETIO_NAMESPACE_STATUS)
  1400. emit('update_forward_config_response', {
  1401. 'success': True,
  1402. 'message': '转发配置已更新',
  1403. 'error_code': ERROR_CODES['SUCCESS']
  1404. })
  1405. logger.info(f"转发配置已更新 - 串口到MQTT: {forward_serial_to_mqtt}, MQTT到串口: {forward_mqtt_to_serial}, 发布主题: {mqtt_publish_topic}")
  1406. except ValueError as e:
  1407. error_msg = str(e)
  1408. logger.warning(f"转发配置更新失败: {error_msg}")
  1409. emit('update_forward_config_response', {
  1410. 'success': False,
  1411. 'message': error_msg,
  1412. 'error_code': ERROR_CODES['CONFIG_ERROR']
  1413. })
  1414. except Exception as e:
  1415. error_msg = f'更新转发配置失败: {str(e)}'
  1416. logger.error(error_msg)
  1417. emit('update_forward_config_response', {
  1418. 'success': False,
  1419. 'message': error_msg,
  1420. 'error_code': ERROR_CODES['UNKNOWN_ERROR']
  1421. })
  1422. @socketio.on('clear_data_buffer', namespace=SOCKETIO_NAMESPACE_CONTROL)
  1423. def handle_clear_data_buffer(data):
  1424. """通过WebSocket清空数据缓冲区"""
  1425. try:
  1426. buffer_type = data.get('type', '')
  1427. if buffer_type == 'serial':
  1428. serial_data_buffer.clear()
  1429. emit('clear_data_buffer_response', {
  1430. 'success': True,
  1431. 'message': '串口数据缓冲区已清空',
  1432. 'error_code': ERROR_CODES['SUCCESS']
  1433. })
  1434. logger.info("串口数据缓冲区已清空")
  1435. elif buffer_type == 'mqtt':
  1436. mqtt_data_buffer.clear()
  1437. emit('clear_data_buffer_response', {
  1438. 'success': True,
  1439. 'message': 'MQTT数据缓冲区已清空',
  1440. 'error_code': ERROR_CODES['SUCCESS']
  1441. })
  1442. logger.info("MQTT数据缓冲区已清空")
  1443. elif buffer_type == 'all':
  1444. serial_data_buffer.clear()
  1445. mqtt_data_buffer.clear()
  1446. emit('clear_data_buffer_response', {
  1447. 'success': True,
  1448. 'message': '所有数据缓冲区已清空',
  1449. 'error_code': ERROR_CODES['SUCCESS']
  1450. })
  1451. logger.info("所有数据缓冲区已清空")
  1452. else:
  1453. emit('clear_data_buffer_response', {
  1454. 'success': False,
  1455. 'message': '无效的缓冲区类型',
  1456. 'error_code': ERROR_CODES['CONFIG_ERROR']
  1457. })
  1458. logger.warning(f"清空缓冲区失败: 无效的缓冲区类型 '{buffer_type}'")
  1459. except Exception as e:
  1460. error_msg = f'清空缓冲区失败: {str(e)}'
  1461. logger.error(error_msg)
  1462. emit('clear_data_buffer_response', {
  1463. 'success': False,
  1464. 'message': error_msg,
  1465. 'error_code': ERROR_CODES['UNKNOWN_ERROR']
  1466. })
  1467. # 设置回调
  1468. serial_client.set_data_callback(serial_data_handler)
  1469. serial_client.set_send_callback(serial_send_handler)
  1470. serial_client.set_status_callback(serial_status_handler)
  1471. mqtt_client.set_data_callback(mqtt_data_handler_extended)
  1472. mqtt_client.set_status_callback(mqtt_status_handler)
  1473. # 设备配置文件路径
  1474. DEVICE_CONFIG_FILE = '/root/dzxj_dtu/devices.json'
  1475. confirm_loop_running = False
  1476. def save_device_config():
  1477. """保存设备配置到文件"""
  1478. try:
  1479. filepath = DEVICE_CONFIG_FILE
  1480. success = address_config.save_config(filepath)
  1481. if success:
  1482. logger.info(f"设备配置已保存到: {filepath}")
  1483. except Exception as e:
  1484. logger.error(f"保存设备配置失败: {str(e)}")
  1485. def load_device_config():
  1486. """加载设备配置"""
  1487. try:
  1488. filepath = DEVICE_CONFIG_FILE
  1489. if os.path.exists(filepath):
  1490. success = address_config.load_config(filepath)
  1491. if success:
  1492. devices = address_config.get_stored_devices()
  1493. logger.info(f"已加载 {len(devices)} 个设备配置")
  1494. return devices
  1495. except Exception as e:
  1496. logger.error(f"加载设备配置失败: {str(e)}")
  1497. return {}
  1498. def confirm_loop():
  1499. """后台确认线程:每10秒对所有已存储设备发送 confirm_address"""
  1500. global confirm_loop_running
  1501. confirm_loop_running = True
  1502. logger.info("启动设备确认线程 (间隔10秒)")
  1503. while confirm_loop_running:
  1504. try:
  1505. devices = address_config.get_stored_devices()
  1506. if not devices:
  1507. time.sleep(10)
  1508. continue
  1509. if not serial_client.get_status():
  1510. time.sleep(10)
  1511. continue
  1512. _st = serial_client.get_status()
  1513. if not (isinstance(_st, dict) and _st.get('connected', False)):
  1514. time.sleep(10)
  1515. continue
  1516. for uid_hex, addr in devices.items():
  1517. try:
  1518. uid_bytes = bytes.fromhex(uid_hex)
  1519. cmd = build_confirm_address(addr, uid_bytes)
  1520. success, msg = serial_client.send_raw(cmd)
  1521. if success:
  1522. logger.debug(f"确认设备: 地址={addr}, UID={uid_hex[:16]}...")
  1523. else:
  1524. logger.warning(f"确认失败 地址={addr}: {msg}")
  1525. except Exception as e:
  1526. logger.error(f"确认异常 地址={addr}: {str(e)}")
  1527. time.sleep(0.1)
  1528. except Exception as e:
  1529. logger.error(f"确认线程异常: {str(e)}")
  1530. time.sleep(10)
  1531. # API路由
  1532. # 移除静态文件服务,前端由nginx提供服务
  1533. # 根路径路由
  1534. @app.route('/')
  1535. def index():
  1536. return send_from_directory(STATIC_FOLDER, 'index.html')
  1537. # 404错误处理器 - 解决SPA路由刷新问题
  1538. @app.errorhandler(404)
  1539. def page_not_found(e):
  1540. """处理所有404错误,对于非API路径返回index.html"""
  1541. path = request.path
  1542. # 检查是否是API请求
  1543. if path.startswith('/api/'):
  1544. # 对于API请求,返回404错误
  1545. return jsonify({
  1546. 'success': False,
  1547. 'message': 'API endpoint not found'
  1548. }), 404
  1549. # 对于所有非API路径,返回index.html让前端路由处理
  1550. return send_from_directory(STATIC_FOLDER, 'index.html'), 200
  1551. @app.route('/api/serial/ports', methods=['GET'])
  1552. def get_serial_ports():
  1553. """获取可用串口列表"""
  1554. ports = serial_client.list_ports()
  1555. return jsonify({
  1556. 'success': True,
  1557. 'ports': ports
  1558. })
  1559. @app.route('/api/serial/connect', methods=['POST'])
  1560. def serial_connect():
  1561. """连接串口"""
  1562. try:
  1563. data = request.json
  1564. port = data.get('port')
  1565. baudrate = data.get('baudrate', 9600)
  1566. bytesize = data.get('bytesize', 8)
  1567. parity = data.get('parity', 'N')
  1568. stopbits = data.get('stopbits', 1)
  1569. timeout = data.get('timeout', 0.1)
  1570. if not port:
  1571. logger.warning("连接串口请求缺少串口名称")
  1572. return jsonify({
  1573. 'success': False,
  1574. 'message': '串口名称不能为空',
  1575. 'error_code': ERROR_CODES['CONFIG_ERROR']
  1576. }), 400
  1577. # 先断开之前的连接
  1578. _status = serial_client.get_status()
  1579. if isinstance(_status, dict) and _status.get('connected', False):
  1580. _port = serial_client.current_config.port if serial_client.current_config else 'unknown'
  1581. logger.info(f"断开现有串口连接: {_port}")
  1582. serial_client.disconnect()
  1583. # 连接新的串口
  1584. logger.info(f"尝试连接串口: {port}, 波特率: {baudrate}")
  1585. success, message = serial_client.connect(port, baudrate=baudrate, timeout=timeout, bytesize=bytesize, parity=parity, stopbits=stopbits)
  1586. status_code = 200 if success else 400
  1587. error_code = ERROR_CODES['SUCCESS'] if success else ERROR_CODES['SERIAL_CONNECTION_ERROR']
  1588. response = {
  1589. 'success': success,
  1590. 'message': message,
  1591. 'error_code': error_code
  1592. }
  1593. if success:
  1594. logger.info(f"串口连接成功: {port}")
  1595. # 保存串口配置
  1596. save_serial_config(port, baudrate, timeout, bytesize, parity, stopbits)
  1597. else:
  1598. logger.error(f"串口连接失败: {message}")
  1599. return jsonify(response), status_code
  1600. except Exception as e:
  1601. error_msg = f'连接串口时出错: {str(e)}'
  1602. logger.exception(error_msg) # 使用exception记录完整堆栈
  1603. return jsonify({
  1604. 'success': False,
  1605. 'message': error_msg,
  1606. 'error_code': ERROR_CODES['UNKNOWN_ERROR']
  1607. }), 500
  1608. @app.route('/api/serial/disconnect', methods=['POST'])
  1609. def serial_disconnect():
  1610. """断开串口连接"""
  1611. success, message = serial_client.disconnect()
  1612. return jsonify({
  1613. 'success': success,
  1614. 'message': message
  1615. })
  1616. @app.route('/api/serial/status', methods=['GET'])
  1617. def serial_get_status():
  1618. """获取串口状态"""
  1619. _st = serial_client.get_status()
  1620. saved_config = {}
  1621. try:
  1622. with open(SERIAL_CONFIG_FILE, 'r') as f:
  1623. saved_config = json.load(f)
  1624. except (FileNotFoundError, json.JSONDecodeError):
  1625. pass
  1626. result = {'saved_config': saved_config}
  1627. if isinstance(_st, dict):
  1628. result['connected'] = _st.get('connected', False)
  1629. if result['connected']:
  1630. cfg = _st.get('config')
  1631. result['port'] = cfg.port if cfg else None
  1632. else:
  1633. result['port'] = None
  1634. else:
  1635. result['connected'] = bool(_st)
  1636. result['port'] = None
  1637. return jsonify(result)
  1638. @app.route('/api/serial/send', methods=['POST'])
  1639. def serial_send():
  1640. """发送数据到串口"""
  1641. data = request.json
  1642. message = data.get('message')
  1643. if not message:
  1644. return jsonify({
  1645. 'success': False,
  1646. 'message': '消息内容不能为空'
  1647. })
  1648. success, message = serial_client.send_data(message)
  1649. return jsonify({
  1650. 'success': success,
  1651. 'message': message
  1652. })
  1653. @app.route('/api/mqtt/connect', methods=['POST'])
  1654. def mqtt_connect():
  1655. """连接MQTT服务器"""
  1656. try:
  1657. data = request.json
  1658. host = data.get('broker') or data.get('host', 'localhost')
  1659. port = data.get('port', 1883)
  1660. client_id = data.get('client_id', f'serial_gateway_{int(time.time())}')
  1661. username = data.get('username')
  1662. password = data.get('password')
  1663. keepalive = data.get('keepalive', 60)
  1664. # 先断开之前的连接
  1665. if mqtt_client.get_status():
  1666. logger.info(f"断开现有MQTT连接: {mqtt_client.host}:{mqtt_client.port}")
  1667. mqtt_client.disconnect()
  1668. # 连接新的MQTT服务器
  1669. logger.info(f"尝试连接MQTT服务器: {host}:{port}, 客户端ID: {client_id}")
  1670. will_topic = build_dtu_topic(dtu_config['customer_id'], 'dtu', dtu_config['dtu_id'], 'status')
  1671. will_payload = json.dumps({
  1672. 'dtu_id': dtu_config['dtu_id'], 'type': 'STATUS',
  1673. 'payload': {'online': False, 'reason': 'CONNECTION_LOST'}
  1674. })
  1675. success, message = mqtt_client.connect(
  1676. host=host, port=port, client_id=client_id,
  1677. username=username, password=password,
  1678. keepalive=keepalive,
  1679. will_topic=will_topic, will_payload=will_payload, will_qos=1, will_retain=False
  1680. )
  1681. # 如果连接成功,订阅主题
  1682. if success and 'topics' in data:
  1683. mqtt_client.subscribe(data['topics'])
  1684. status_code = 200 if success else 400
  1685. error_code = ERROR_CODES['SUCCESS'] if success else ERROR_CODES['MQTT_CONNECTION_ERROR']
  1686. response = {
  1687. 'success': success,
  1688. 'message': message,
  1689. 'error_code': error_code
  1690. }
  1691. if success:
  1692. logger.info(f"MQTT服务器连接成功: {host}:{port}")
  1693. else:
  1694. logger.error(f"MQTT服务器连接失败: {message}")
  1695. return jsonify(response), status_code
  1696. except Exception as e:
  1697. error_msg = f'连接MQTT服务器时出错: {str(e)}'
  1698. logger.exception(error_msg)
  1699. return jsonify({
  1700. 'success': False,
  1701. 'message': error_msg,
  1702. 'error_code': ERROR_CODES['UNKNOWN_ERROR']
  1703. }), 500
  1704. @app.route('/api/mqtt/disconnect', methods=['POST'])
  1705. def mqtt_disconnect():
  1706. """断开MQTT连接"""
  1707. success, message = mqtt_client.disconnect()
  1708. return jsonify({
  1709. 'success': success,
  1710. 'message': message
  1711. })
  1712. @app.route('/api/mqtt/status', methods=['GET'])
  1713. def mqtt_get_status():
  1714. """获取MQTT状态"""
  1715. _st = mqtt_client.get_status()
  1716. return jsonify({
  1717. 'connected': _st.get('connected', False) if isinstance(_st, dict) else bool(_st)
  1718. })
  1719. @app.route('/api/mqtt/publish', methods=['POST'])
  1720. def mqtt_publish():
  1721. """发布MQTT消息"""
  1722. data = request.json
  1723. topic = data.get('topic')
  1724. message = data.get('message')
  1725. if not topic or not message:
  1726. return jsonify({
  1727. 'success': False,
  1728. 'message': '主题和消息内容不能为空'
  1729. })
  1730. success, message = mqtt_client.publish(topic, message)
  1731. return jsonify({
  1732. 'success': success,
  1733. 'message': message
  1734. })
  1735. @app.route('/api/mqtt/subscribe', methods=['POST'])
  1736. def mqtt_subscribe():
  1737. """订阅MQTT主题"""
  1738. data = request.json
  1739. topics = data.get('topics', [])
  1740. if not topics:
  1741. return jsonify({
  1742. 'success': False,
  1743. 'message': '请至少订阅一个主题'
  1744. })
  1745. success, message = mqtt_client.subscribe(topics)
  1746. return jsonify({
  1747. 'success': success,
  1748. 'message': message
  1749. })
  1750. @app.route('/api/mqtt/broadcast_discover', methods=['POST'])
  1751. def mqtt_broadcast_discover():
  1752. """发送DTU发现广播 (7.14)
  1753. 服务器通过广播主题发送DISCOVER消息,
  1754. 触发所有在线DTU进行注册响应
  1755. """
  1756. data = request.json or {}
  1757. customer_id = data.get('customer_id', dtu_config.get('customer_id', 'default'))
  1758. discover_topic = build_dtu_topic('broadcast', 'dtu', 'discover')
  1759. request_id = f"dsc_{int(time.time() * 1000)}"
  1760. payload = {
  1761. 'msg_id': request_id,
  1762. 'timestamp': int(time.time() * 1000),
  1763. 'dtu_id': 'PLATFORM',
  1764. 'type': 'DISCOVER',
  1765. 'payload': {
  1766. 'request_id': request_id,
  1767. 'customer_id': customer_id
  1768. }
  1769. }
  1770. success, message = mqtt_client.publish(discover_topic, json.dumps(payload), qos=1)
  1771. logger.info(f"发送DTU发现广播: topic={discover_topic}, payload={payload}")
  1772. return jsonify({
  1773. 'success': success,
  1774. 'message': 'DTU发现广播已发送' if success else message,
  1775. 'topic': discover_topic
  1776. })
  1777. @app.route('/api/mqtt/broadcast_config', methods=['POST'])
  1778. def mqtt_broadcast_config():
  1779. """发送DTU批量配置广播 (7.15)
  1780. 服务器通过广播主题发送CONFIG消息,
  1781. 用于批量配置所有在线DTU
  1782. """
  1783. data = request.json or {}
  1784. customer_id = data.get('customer_id', dtu_config.get('customer_id', 'default'))
  1785. config_items = data.get('items', {})
  1786. config_version = data.get('config_version', f"cfg_{int(time.time())}")
  1787. force_apply = data.get('force_apply', False)
  1788. config_topic = build_dtu_topic('broadcast', 'dtu', 'config')
  1789. payload = {
  1790. 'msg_id': f"cfg_{int(time.time() * 1000)}",
  1791. 'timestamp': int(time.time() * 1000),
  1792. 'dtu_id': 'PLATFORM',
  1793. 'type': 'CONFIG',
  1794. 'payload': {
  1795. 'customer_id': customer_id,
  1796. 'config_version': config_version,
  1797. 'force_apply': force_apply,
  1798. 'items': config_items
  1799. }
  1800. }
  1801. success, message = mqtt_client.publish(config_topic, json.dumps(payload), qos=1)
  1802. logger.info(f"发送DTU批量配置广播: topic={config_topic}, payload={payload}")
  1803. return jsonify({
  1804. 'success': success,
  1805. 'message': 'DTU批量配置广播已发送' if success else message,
  1806. 'topic': config_topic,
  1807. 'config': {
  1808. 'config_version': config_version,
  1809. 'items': config_items,
  1810. 'force_apply': force_apply
  1811. }
  1812. })
  1813. @app.route('/api/data/serial', methods=['GET'])
  1814. def get_serial_data():
  1815. """获取串口数据"""
  1816. return jsonify({
  1817. 'data': serial_data_buffer
  1818. })
  1819. @app.route('/api/data/mqtt', methods=['GET'])
  1820. def get_mqtt_data():
  1821. """获取MQTT数据"""
  1822. return jsonify({
  1823. 'data': mqtt_data_buffer
  1824. })
  1825. @app.route('/api/forward/config', methods=['POST'])
  1826. def set_forward_config():
  1827. """设置转发配置"""
  1828. global forward_serial_to_mqtt, forward_mqtt_to_serial, mqtt_publish_topic
  1829. data = request.json
  1830. forward_serial_to_mqtt = data.get('serial_to_mqtt', False)
  1831. forward_mqtt_to_serial = data.get('mqtt_to_serial', False)
  1832. if 'publish_topic' in data:
  1833. mqtt_publish_topic = data['publish_topic']
  1834. return jsonify({
  1835. 'success': True,
  1836. 'message': '转发配置已更新',
  1837. 'config': {
  1838. 'serial_to_mqtt': forward_serial_to_mqtt,
  1839. 'mqtt_to_serial': forward_mqtt_to_serial,
  1840. 'publish_topic': mqtt_publish_topic
  1841. }
  1842. })
  1843. @app.route('/api/forward/status', methods=['GET'])
  1844. def get_forward_status():
  1845. """获取转发状态"""
  1846. return jsonify({
  1847. 'serial_to_mqtt': forward_serial_to_mqtt,
  1848. 'mqtt_to_serial': forward_mqtt_to_serial,
  1849. 'publish_topic': mqtt_publish_topic
  1850. })
  1851. # 健康检查端点
  1852. @app.route('/api/health', methods=['GET'])
  1853. def health_check():
  1854. """健康检查端点"""
  1855. try:
  1856. # 获取客户端连接数量
  1857. client_counts = {
  1858. 'data': len(connected_clients['data']),
  1859. 'status': len(connected_clients['status']),
  1860. 'control': len(connected_clients['control'])
  1861. }
  1862. # 获取缓冲区大小
  1863. buffer_sizes = {
  1864. 'serial': len(serial_data_buffer),
  1865. 'mqtt': len(mqtt_data_buffer)
  1866. }
  1867. # 执行系统负载检查
  1868. # 注意:这只是一个简化的负载检查,实际应用中可能需要更复杂的监控
  1869. is_healthy = True
  1870. load_warnings = []
  1871. # 检查缓冲区是否过大
  1872. if buffer_sizes['serial'] > MAX_BUFFER_SIZE * 0.8:
  1873. is_healthy = False
  1874. load_warnings.append(f"串口缓冲区接近最大容量: {buffer_sizes['serial']}/{MAX_BUFFER_SIZE}")
  1875. if buffer_sizes['mqtt'] > MAX_BUFFER_SIZE * 0.8:
  1876. is_healthy = False
  1877. load_warnings.append(f"MQTT缓冲区接近最大容量: {buffer_sizes['mqtt']}/{MAX_BUFFER_SIZE}")
  1878. # 检查WebSocket连接数是否过多
  1879. total_clients = sum(client_counts.values())
  1880. if total_clients > 100: # 设置合理的阈值
  1881. is_healthy = False
  1882. load_warnings.append(f"WebSocket连接数过多: {total_clients}")
  1883. # 尝试获取网络状态信息(使用try-except包装,防止网络模块出错导致健康检查失败)
  1884. network_info = None
  1885. try:
  1886. network_info = network_manager.get_network_status()
  1887. except Exception as e:
  1888. logger.warning(f"获取网络状态时出错: {str(e)}")
  1889. # 不影响整体健康检查,只添加警告
  1890. load_warnings.append(f"网络状态获取失败: {str(e)}")
  1891. response = {
  1892. 'status': 'healthy' if is_healthy else 'warning',
  1893. 'timestamp': time.strftime('%Y-%m-%d %H:%M:%S'),
  1894. 'services': {
  1895. 'serial': serial_status,
  1896. 'mqtt': mqtt_status,
  1897. 'websocket': True
  1898. },
  1899. 'client_counts': client_counts,
  1900. 'buffer_sizes': buffer_sizes,
  1901. 'warnings': load_warnings
  1902. }
  1903. # 仅当获取到网络信息时添加
  1904. if network_info:
  1905. response['network'] = network_info
  1906. return jsonify(response), 200
  1907. except Exception as e:
  1908. # 捕获所有异常,确保健康检查不会返回500错误
  1909. logger.error(f"健康检查端点出错: {str(e)}")
  1910. # 返回一个基础的健康状态,至少显示服务在运行
  1911. return jsonify({
  1912. 'status': 'error',
  1913. 'timestamp': time.strftime('%Y-%m-%d %H:%M:%S'),
  1914. 'error': str(e),
  1915. 'services': {
  1916. 'api': True, # API服务本身是运行的
  1917. 'serial': None,
  1918. 'mqtt': None,
  1919. 'websocket': None
  1920. }
  1921. }), 200 # 仍然返回200,避免健康检查导致的连锁反应
  1922. @app.route('/api/screen/status', methods=['GET'])
  1923. def get_screen_status():
  1924. """获取 Nextion 屏幕状态。"""
  1925. with screen_lock:
  1926. panels = get_sorted_panels()
  1927. current_index = screen_current_panel_index
  1928. current_panel_id = None
  1929. if panels and 0 <= current_index < len(panels):
  1930. current_panel_id = panels[current_index][0]
  1931. panel_count = len(panels)
  1932. status = screen_display.get_status()
  1933. return jsonify({
  1934. 'success': True,
  1935. 'data': {
  1936. 'connected': status.get('connected', False),
  1937. 'port': status.get('port'),
  1938. 'baudrate': status.get('baudrate'),
  1939. 'current_panel_index': current_index,
  1940. 'current_panel_id': current_panel_id,
  1941. 'panel_count': panel_count
  1942. }
  1943. })
  1944. @app.route('/api/screen/refresh', methods=['POST'])
  1945. def refresh_screen():
  1946. """手动刷新当前 panel 到屏幕。"""
  1947. with screen_lock:
  1948. panels = get_sorted_panels()
  1949. if not panels:
  1950. return jsonify({'success': False, 'message': '没有可用的 panel'}), 400
  1951. if not (0 <= screen_current_panel_index < len(panels)):
  1952. return jsonify({'success': False, 'message': '当前 panel 索引无效'}), 400
  1953. panel_id = panels[screen_current_panel_index][0]
  1954. success, message = refresh_screen_panel(panel_id)
  1955. if success:
  1956. return jsonify({'success': True, 'message': message})
  1957. return jsonify({'success': False, 'message': message}), 400
  1958. # 网络配置相关API
  1959. @app.route('/api/network/config', methods=['GET'])
  1960. def get_network_config():
  1961. """获取当前网络配置"""
  1962. try:
  1963. config = network_manager.get_network_config()
  1964. return jsonify(config)
  1965. except Exception as e:
  1966. logger.error(f'获取网络配置失败: {str(e)}')
  1967. return jsonify({'success': False, 'message': f'获取网络配置失败: {str(e)}'}), 500
  1968. @app.route('/api/network/config', methods=['POST'])
  1969. def update_network_config():
  1970. """更新网络配置"""
  1971. try:
  1972. config_data = request.json
  1973. result = network_manager.update_network_config(config_data)
  1974. if result['success']:
  1975. return jsonify(result)
  1976. else:
  1977. return jsonify(result), 400
  1978. except Exception as e:
  1979. logger.error(f'更新网络配置失败: {str(e)}')
  1980. return jsonify({'success': False, 'message': f'更新网络配置失败: {str(e)}'}), 500
  1981. @app.route('/api/network/status', methods=['GET'])
  1982. def get_network_status():
  1983. """获取网络状态信息"""
  1984. try:
  1985. status = network_manager.get_network_status()
  1986. return jsonify(status)
  1987. except Exception as e:
  1988. logger.error(f'获取网络状态失败: {str(e)}')
  1989. return jsonify({'success': False, 'message': f'获取网络状态失败: {str(e)}'}), 500
  1990. @app.route('/api/network/restart', methods=['POST'])
  1991. def restart_network_service():
  1992. """重启网络服务以应用新配置"""
  1993. try:
  1994. result = network_manager.restart_network_service()
  1995. if result['success']:
  1996. return jsonify(result)
  1997. else:
  1998. return jsonify(result), 500
  1999. except Exception as e:
  2000. logger.error(f'重启网络服务失败: {str(e)}')
  2001. return jsonify({'success': False, 'message': f'重启网络服务失败: {str(e)}'}), 500
  2002. # Modbus RTU API
  2003. @app.route('/api/modbus/antenna_addresses', methods=['GET'])
  2004. def get_antenna_addresses():
  2005. """获取天线地址映射表"""
  2006. return jsonify({
  2007. 'success': True,
  2008. 'antennas': {str(k): f"0x{v:04x}" for k, v in ANTENNA_ADDRESSES.items()}
  2009. })
  2010. @app.route('/api/modbus/read_antenna', methods=['POST'])
  2011. def modbus_read_antenna():
  2012. """读取指定天线的卡号
  2013. 请求参数:
  2014. {
  2015. "device_address": 1, // 设备地址 (1-247)
  2016. "antenna": 1, // 天线编号 (1-24)
  2017. "timeout": 1.0 // 可选,超时时间(秒)
  2018. }
  2019. """
  2020. try:
  2021. data = request.json
  2022. device_address = data.get('device_address', 1)
  2023. antenna = data.get('antenna', 1)
  2024. timeout = data.get('timeout')
  2025. if antenna < 1 or antenna > 24:
  2026. return jsonify({
  2027. 'success': False,
  2028. 'message': f'无效的天线编号: {antenna}, 必须是1-24'
  2029. }), 400
  2030. if not (isinstance(serial_client.get_status(), dict) and serial_client.get_status().get("connected", False)):
  2031. return jsonify({
  2032. 'success': False,
  2033. 'message': '串口未连接'
  2034. }), 400
  2035. result = modbus_client.read_antenna_card(
  2036. device_address=device_address,
  2037. antenna_num=antenna,
  2038. timeout=timeout
  2039. )
  2040. if 'error' in result:
  2041. return jsonify({
  2042. 'success': False,
  2043. 'message': result['error'],
  2044. 'raw_data': result.get('raw_data', '')
  2045. }), 400
  2046. return jsonify({
  2047. 'success': True,
  2048. 'data': result
  2049. })
  2050. except Exception as e:
  2051. logger.error(f'读取天线数据失败: {str(e)}')
  2052. return jsonify({
  2053. 'success': False,
  2054. 'message': str(e)
  2055. }), 500
  2056. @app.route('/api/modbus/read_registers', methods=['POST'])
  2057. def modbus_read_registers():
  2058. """读取保持寄存器
  2059. 请求参数:
  2060. {
  2061. "device_address": 1,
  2062. "start_address": 2, // 起始地址 (十六进制如0x0002或十进制如2)
  2063. "quantity": 4, // 寄存器数量
  2064. "timeout": 1.0
  2065. }
  2066. """
  2067. try:
  2068. data = request.json
  2069. device_address = data.get('device_address', 1)
  2070. start_address = data.get('start_address', 0)
  2071. quantity = data.get('quantity', 1)
  2072. timeout = data.get('timeout')
  2073. # 支持十六进制字符串
  2074. if isinstance(start_address, str):
  2075. start_address = int(start_address, 16)
  2076. if isinstance(device_address, str):
  2077. device_address = int(device_address, 16)
  2078. if not (isinstance(serial_client.get_status(), dict) and serial_client.get_status().get("connected", False)):
  2079. return jsonify({
  2080. 'success': False,
  2081. 'message': '串口未连接'
  2082. }), 400
  2083. result = modbus_client.read_holding_registers(
  2084. device_address=device_address,
  2085. start_address=start_address,
  2086. quantity=quantity,
  2087. timeout=timeout
  2088. )
  2089. if 'error' in result:
  2090. return jsonify({
  2091. 'success': False,
  2092. 'message': result['error'],
  2093. 'raw_data': result.get('raw_data', '')
  2094. }), 400
  2095. return jsonify({
  2096. 'success': True,
  2097. 'data': result
  2098. })
  2099. except Exception as e:
  2100. logger.error(f'读取寄存器失败: {str(e)}')
  2101. return jsonify({
  2102. 'success': False,
  2103. 'message': str(e)
  2104. }), 500
  2105. @app.route('/api/modbus/write_register', methods=['POST'])
  2106. def modbus_write_register():
  2107. """写单个寄存器
  2108. 请求参数:
  2109. {
  2110. "device_address": 1,
  2111. "register_address": 1, // 寄存器地址
  2112. "value": 256, // 写入的值
  2113. "timeout": 1.0
  2114. }
  2115. """
  2116. try:
  2117. data = request.json
  2118. device_address = data.get('device_address', 1)
  2119. register_address = data.get('register_address', 1)
  2120. value = data.get('value', 0)
  2121. timeout = data.get('timeout')
  2122. if not (isinstance(serial_client.get_status(), dict) and serial_client.get_status().get("connected", False)):
  2123. return jsonify({
  2124. 'success': False,
  2125. 'message':'串口未连接'
  2126. }), 400
  2127. result = modbus_client.write_single_register(
  2128. device_address=device_address,
  2129. register_address=register_address,
  2130. value=value,
  2131. timeout=timeout
  2132. )
  2133. if 'error' in result:
  2134. return jsonify({
  2135. 'success': False,
  2136. 'message': result['error'],
  2137. 'raw_data': result.get('raw_data', '')
  2138. }), 400
  2139. return jsonify({
  2140. 'success': True,
  2141. 'data': result
  2142. })
  2143. except Exception as e:
  2144. logger.error(f'写寄存器失败: {str(e)}')
  2145. return jsonify({
  2146. 'success': False,
  2147. 'message': str(e)
  2148. }), 500
  2149. @app.route('/api/modbus/set_rgb_led', methods=['POST'])
  2150. def modbus_set_rgb_led():
  2151. """设置RGB灯状态
  2152. 请求参数:
  2153. {
  2154. "device_address": 1,
  2155. "led_number": 1, // 灯编号 (1-24)
  2156. "color": 1, // 颜色: 0=灭, 1=红灯, 2=绿灯, 3=蓝灯
  2157. "timeout": 1.0
  2158. }
  2159. """
  2160. try:
  2161. data = request.json
  2162. device_address = data.get('device_address', 1)
  2163. led_number = data.get('led_number', 1)
  2164. color = data.get('color', 0)
  2165. timeout = data.get('timeout')
  2166. if not (isinstance(serial_client.get_status(), dict) and serial_client.get_status().get("connected", False)):
  2167. return jsonify({
  2168. 'success': False,
  2169. 'message': '串口未连接'
  2170. }), 400
  2171. result = modbus_client.set_rgb_led(
  2172. device_address=device_address,
  2173. led_number=led_number,
  2174. color=color,
  2175. timeout=timeout
  2176. )
  2177. if 'error' in result:
  2178. return jsonify({
  2179. 'success': False,
  2180. 'message': result['error'],
  2181. 'raw_data': result.get('raw_data', '')
  2182. }), 400
  2183. # 更新 LED 状态追踪
  2184. led_states[led_number] = color
  2185. return jsonify({
  2186. 'success': True,
  2187. 'data': result
  2188. })
  2189. except Exception as e:
  2190. logger.error(f'设置RGB灯失败: {str(e)}')
  2191. return jsonify({
  2192. 'success': False,
  2193. 'message': str(e)
  2194. }), 500
  2195. @app.route('/api/modbus/scan', methods=['POST'])
  2196. def modbus_scan_devices():
  2197. """扫描在线设备
  2198. 请求参数:
  2199. {
  2200. "max_address": 247, // 最大设备地址
  2201. "timeout": 0.2 // 单个设备超时时间
  2202. }
  2203. """
  2204. try:
  2205. data = request.json or {}
  2206. max_address = data.get('max_address', 247)
  2207. timeout = data.get('timeout', 0.2)
  2208. if not (isinstance(serial_client.get_status(), dict) and serial_client.get_status().get("connected", False)):
  2209. return jsonify({
  2210. 'success': False,
  2211. 'message': '串口未连接'
  2212. }), 400
  2213. # 异步扫描可能更好,但这里先同步实现
  2214. result = modbus_client.scan_devices(max_address)
  2215. return jsonify({
  2216. 'success': True,
  2217. 'devices': result
  2218. })
  2219. except Exception as e:
  2220. logger.error(f'扫描设备失败: {str(e)}')
  2221. return jsonify({
  2222. 'success': False,
  2223. 'message': str(e)
  2224. }), 500
  2225. # ========== 地址配置协议 API ==========
  2226. @app.route('/api/modbus/broadcast_query', methods=['POST'])
  2227. def modbus_broadcast_query():
  2228. """发送广播查询指令"""
  2229. try:
  2230. data = request.json or {}
  2231. timeout = data.get('timeout')
  2232. if not (isinstance(serial_client.get_status(), dict) and serial_client.get_status().get("connected", False)):
  2233. return jsonify({'success': False, 'message': '串口未连接'}), 400
  2234. responses = address_config.broadcast_query(timeout)
  2235. assign_results = []
  2236. if responses:
  2237. assign_results = address_config.process_responses(responses)
  2238. for r in assign_results:
  2239. uid_bytes = r.get('uid', r.get('request', ''))
  2240. save_device_config()
  2241. return jsonify({
  2242. 'success': True,
  2243. 'responses': responses,
  2244. 'assign_results': assign_results,
  2245. 'count': len(responses)
  2246. })
  2247. except Exception as e:
  2248. logger.error(f'广播查询失败: {str(e)}')
  2249. return jsonify({'success': False, 'message': str(e)}), 500
  2250. @app.route('/api/modbus/auto_configure', methods=['POST'])
  2251. def modbus_auto_configure():
  2252. """自动配置设备地址"""
  2253. try:
  2254. data = request.json or {}
  2255. timeout = data.get('timeout')
  2256. if not (isinstance(serial_client.get_status(), dict) and serial_client.get_status().get("connected", False)):
  2257. return jsonify({'success': False, 'message': '串口未连接'}), 400
  2258. result = address_config.auto_configure(timeout)
  2259. if result.get('success') and (result.get('discovered', 0) > 0 or result.get('confirmed', 0) > 0 or result.get('assigned', 0) > 0):
  2260. save_device_config()
  2261. return jsonify(result)
  2262. except Exception as e:
  2263. logger.error(f'自动配置失败: {str(e)}')
  2264. return jsonify({'success': False, 'message': str(e)}), 500
  2265. @app.route('/api/modbus/stored_devices', methods=['GET'])
  2266. def get_stored_devices():
  2267. """获取已存储的设备列表"""
  2268. return jsonify({'success': True, 'devices': address_config.get_stored_devices()})
  2269. @app.route('/api/modbus/stored_devices', methods=['POST'])
  2270. def add_stored_device():
  2271. """添加已存储的设备"""
  2272. try:
  2273. data = request.json
  2274. uid = data.get('uid', '').lower()
  2275. address = data.get('address', 1)
  2276. if len(uid) != 24:
  2277. return jsonify({'success': False, 'message': 'UID长度必须是24个十六进制字符'}), 400
  2278. address_config.add_stored_device(uid, address)
  2279. save_device_config()
  2280. return jsonify({'success': True, 'message': f'已添加设备: UID={uid}, 地址={address}'})
  2281. except Exception as e:
  2282. return jsonify({'success': False, 'message': str(e)}), 500
  2283. @app.route('/api/modbus/load_config', methods=['POST'])
  2284. def modbus_load_config():
  2285. """从文件加载设备配置"""
  2286. try:
  2287. data = request.json
  2288. filepath = data.get('filepath', '/tmp/modbus_devices.json')
  2289. success = address_config.load_config(filepath)
  2290. return jsonify({'success': success, 'devices': address_config.get_stored_devices()})
  2291. except Exception as e:
  2292. return jsonify({'success': False, 'message': str(e)}), 500
  2293. @app.route('/api/modbus/save_config', methods=['POST'])
  2294. def modbus_save_config():
  2295. """保存设备配置到文件"""
  2296. try:
  2297. data = request.json
  2298. filepath = data.get('filepath', '/tmp/modbus_devices.json')
  2299. success = address_config.save_config(filepath)
  2300. return jsonify({'success': success, 'message': f'配置已保存到: {filepath}' if success else '保存失败'})
  2301. except Exception as e:
  2302. return jsonify({'success': False, 'message': str(e)}), 500
  2303. @app.route('/api/modbus/confirm_devices', methods=['POST'])
  2304. def confirm_devices():
  2305. """手动触发对所有已存储设备的确认"""
  2306. try:
  2307. devices = address_config.get_stored_devices()
  2308. if not devices:
  2309. return jsonify({'success': False, 'message': '没有已存储的设备'}), 400
  2310. if not (isinstance(serial_client.get_status(), dict) and serial_client.get_status().get("connected", False)):
  2311. return jsonify({'success': False, 'message': '串口未连接'}), 400
  2312. results = []
  2313. for uid_hex, addr in devices.items():
  2314. try:
  2315. uid_bytes = bytes.fromhex(uid_hex)
  2316. cmd = build_confirm_address(addr, uid_bytes)
  2317. success, msg = serial_client.send_raw(cmd)
  2318. logger.info(f"确认设备 地址={addr}: {'成功' if success else '失败 ' + msg}")
  2319. results.append({'address': addr, 'uid': uid_hex, 'success': success, 'message': msg})
  2320. except Exception as e:
  2321. results.append({'address': addr, 'uid': uid_hex, 'success': False, 'message': str(e)})
  2322. time.sleep(0.1)
  2323. return jsonify({'success': True, 'results': results, 'count': len(results)})
  2324. except Exception as e:
  2325. return jsonify({'success': False, 'message': str(e)}), 500
  2326. # 端口状态追踪 - 用于存储每个端口的事件状态
  2327. port_event_history = [] # 存储最近的端口事件
  2328. @app.route('/api/port/status', methods=['GET'])
  2329. def get_port_status():
  2330. """获取所有端口状态"""
  2331. return jsonify({'success': True, 'port_state': port_state})
  2332. @app.route('/api/port/events', methods=['GET'])
  2333. def get_port_events():
  2334. """获取端口事件历史"""
  2335. limit = request.args.get('limit', 100, type=int)
  2336. return jsonify({
  2337. 'success': True,
  2338. 'events': port_event_history[-limit:]
  2339. })
  2340. @app.route('/api/port/clear_events', methods=['POST'])
  2341. def clear_port_events():
  2342. """清除端口事件历史"""
  2343. global port_event_history
  2344. port_event_history = []
  2345. return jsonify({'success': True, 'message': '事件历史已清除'})
  2346. # 设备最后响应时间跟踪
  2347. device_last_seen = {}
  2348. @app.route('/api/panel/status', methods=['GET'])
  2349. def get_panel_status():
  2350. """获取所有面板状态"""
  2351. now = time.time()
  2352. PANEL_OFFLINE_TIMEOUT = 60
  2353. panel_status = {}
  2354. for panel_id, ports in port_state.items():
  2355. port_count = len(ports)
  2356. alarm_count = sum(1 for p in ports.values() if p.get('alarm_count', 0) > 0)
  2357. connected_count = sum(1 for p in ports.values() if p.get('last_uid'))
  2358. panel_status[panel_id] = {
  2359. 'panel_id': panel_id,
  2360. 'port_count': port_count,
  2361. 'connected_count': connected_count,
  2362. 'alarm_count': alarm_count,
  2363. 'status': 'online' if connected_count > 0 else 'offline'
  2364. }
  2365. for panel_id, cfg in panel_config.items():
  2366. last_seen = device_last_seen.get(panel_id, 0)
  2367. is_online = (now - last_seen) < PANEL_OFFLINE_TIMEOUT
  2368. if panel_id in panel_status:
  2369. panel_status[panel_id].update({
  2370. 'address': cfg.get('address'),
  2371. 'position': cfg.get('position'),
  2372. 'panel_uid': cfg.get('panel_uid'),
  2373. 'status': 'online' if is_online else 'offline'
  2374. })
  2375. else:
  2376. panel_status[panel_id] = {
  2377. 'panel_id': panel_id,
  2378. 'address': cfg.get('address'),
  2379. 'position': cfg.get('position'),
  2380. 'panel_uid': cfg.get('panel_uid'),
  2381. 'port_count': 0,
  2382. 'connected_count': 0,
  2383. 'alarm_count': 0,
  2384. 'status': 'online' if is_online else 'offline'
  2385. }
  2386. return jsonify({'success': True, 'panels': panel_status})
  2387. # LED 状态追踪 (内存中跟踪每个灯的最后设置状态)
  2388. led_states = {i: 0 for i in range(1, 25)} # 0=off, 1=red, 2=green, 3=blue
  2389. @app.route('/api/modbus/led_status', methods=['GET'])
  2390. def get_led_status():
  2391. """获取所有 LED 状态"""
  2392. return jsonify({'success': True, 'leds': led_states})
  2393. @app.route('/api/modbus/set_all_leds', methods=['POST'])
  2394. def set_all_leds():
  2395. """批量设置 LED
  2396. 请求: {"device_address": 1, "color": 1}
  2397. 设置所有 24 个 LED 到指定颜色
  2398. """
  2399. try:
  2400. data = request.json
  2401. device_address = data.get('device_address', 1)
  2402. color = data.get('color', 0)
  2403. timeout = data.get('timeout')
  2404. if not (isinstance(serial_client.get_status(), dict) and serial_client.get_status().get("connected", False)):
  2405. return jsonify({'success': False, 'message': '串口未连接'}), 400
  2406. results = []
  2407. for led_num in range(1, 25):
  2408. result = modbus_client.set_rgb_led(
  2409. device_address=device_address,
  2410. led_number=led_num,
  2411. color=color,
  2412. timeout=timeout
  2413. )
  2414. if 'error' in result:
  2415. results.append({'led': led_num, 'success': False, 'error': result['error']})
  2416. else:
  2417. results.append({'led': led_num, 'success': True})
  2418. led_states[led_num] = color
  2419. return jsonify({'success': True, 'results': results})
  2420. except Exception as e:
  2421. logger.error(f'批量设置LED失败: {str(e)}')
  2422. return jsonify({'success': False, 'message': str(e)}), 500
  2423. # 在 set_rgb_led 中更新 LED 状态追踪
  2424. # 不再需要静态文件目录,前端由nginx提供服务
  2425. # ========== DTU MQTT协议配置API ==========
  2426. @app.route('/api/dtu/config', methods=['GET'])
  2427. def get_dtu_config():
  2428. """获取DTU配置"""
  2429. return jsonify({
  2430. 'success': True,
  2431. 'data': dtu_config
  2432. })
  2433. @app.route('/api/dtu/config', methods=['POST'])
  2434. def update_dtu_config():
  2435. """更新DTU配置"""
  2436. try:
  2437. data = request.json
  2438. old_config = dtu_config.copy()
  2439. # 更新配置项
  2440. if 'topic_prefix' in data:
  2441. dtu_config['topic_prefix'] = data['topic_prefix']
  2442. if 'customer_id' in data:
  2443. dtu_config['customer_id'] = data['customer_id']
  2444. if 'dtu_id' in data:
  2445. dtu_config['dtu_id'] = data['dtu_id']
  2446. if 'firmware_version' in data:
  2447. dtu_config['firmware_version'] = data['firmware_version']
  2448. if 'hardware_version' in data:
  2449. dtu_config['hardware_version'] = data['hardware_version']
  2450. if 'heartbeat_interval' in data:
  2451. dtu_config['heartbeat_interval'] = data['heartbeat_interval']
  2452. if 'enabled' in data:
  2453. dtu_config['enabled'] = data['enabled']
  2454. # 如果MQTT已连接且主题配置发生变化,重新订阅
  2455. if mqtt_status and dtu_config.get('enabled'):
  2456. old_control_topic = build_dtu_topic(
  2457. old_config.get('customer_id', DEFAULT_CUSTOMER_ID),
  2458. 'dtu',
  2459. old_config.get('dtu_id', DEFAULT_DTU_ID),
  2460. 'control'
  2461. )
  2462. new_control_topic = build_dtu_topic(
  2463. dtu_config['customer_id'],
  2464. 'dtu',
  2465. dtu_config['dtu_id'],
  2466. 'control'
  2467. )
  2468. if old_control_topic != new_control_topic:
  2469. # 取消旧订阅,订阅新主题
  2470. mqtt_client.unsubscribe(old_control_topic)
  2471. mqtt_client.subscribe(new_control_topic)
  2472. logger.info(f"控制主题已更新: {old_control_topic} -> {new_control_topic}")
  2473. # 重新发送注册消息
  2474. dtu_register()
  2475. return jsonify({
  2476. 'success': True,
  2477. 'message': 'DTU配置已更新',
  2478. 'data': dtu_config
  2479. })
  2480. except Exception as e:
  2481. logger.error(f"更新DTU配置失败: {str(e)}")
  2482. return jsonify({'success': False, 'message': str(e)}), 500
  2483. @app.route('/api/dtu/register', methods=['POST'])
  2484. def manual_dtu_register():
  2485. """手动触发DTU注册"""
  2486. try:
  2487. if not mqtt_status:
  2488. return jsonify({'success': False, 'message': 'MQTT未连接'}), 400
  2489. success = dtu_register()
  2490. if success:
  2491. return jsonify({'success': True, 'message': '注册消息已发送'})
  2492. else:
  2493. return jsonify({'success': False, 'message': '注册消息发送失败'}), 500
  2494. except Exception as e:
  2495. logger.error(f"手动触发DTU注册失败: {str(e)}")
  2496. return jsonify({'success': False, 'message': str(e)}), 500
  2497. def _get_cpu_usage():
  2498. try:
  2499. import psutil
  2500. return psutil.cpu_percent(interval=None)
  2501. except Exception:
  2502. return None
  2503. def _get_memory_usage():
  2504. try:
  2505. import psutil
  2506. return psutil.virtual_memory().percent
  2507. except Exception:
  2508. return None
  2509. @app.route('/api/dtu/status', methods=['GET'])
  2510. def get_dtu_status():
  2511. """获取DTU状态"""
  2512. try:
  2513. # 获取串口状态
  2514. serial_st = serial_client.get_status()
  2515. # 获取面板状态
  2516. devices = address_config.get_stored_devices()
  2517. # 主板温度: 优先外部上报,否则本地读取
  2518. dtu_temperature = env_sensor_data.get('dtu_temperature')
  2519. if dtu_temperature is None:
  2520. try:
  2521. import psutil
  2522. temps = psutil.sensors_temperatures(fahrenheit=False) or {}
  2523. for chip_name, entries in temps.items():
  2524. for entry in entries:
  2525. if entry.current and entry.current > 0:
  2526. dtu_temperature = round(entry.current, 1)
  2527. break
  2528. if dtu_temperature is not None:
  2529. break
  2530. except Exception:
  2531. pass
  2532. if dtu_temperature is None:
  2533. import glob as _glob
  2534. for tz in sorted(_glob.glob('/sys/class/thermal/thermal_zone*/temp')):
  2535. try:
  2536. with open(tz) as f:
  2537. v = int(f.read().strip())
  2538. if v > 0:
  2539. dtu_temperature = round(v / 1000.0, 1)
  2540. break
  2541. except Exception:
  2542. continue
  2543. status = {
  2544. 'dtu_id': dtu_config.get('dtu_id'),
  2545. 'mqtt_connected': mqtt_status,
  2546. 'serial_connected': isinstance(serial_st, dict) and serial_st.get('connected', False),
  2547. 'dtu_enabled': dtu_config.get('enabled', True),
  2548. 'topic_prefix': dtu_config.get('topic_prefix'),
  2549. 'customer_id': dtu_config.get('customer_id'),
  2550. 'panel_count': len(devices),
  2551. 'cpu_usage': _get_cpu_usage(),
  2552. 'memory_usage': _get_memory_usage(),
  2553. 'temperature': dtu_temperature,
  2554. 'dtu_temperature': dtu_temperature,
  2555. 'firmware_version': dtu_config.get('firmware_version', 'v1.0.0'),
  2556. 'topics': {
  2557. 'register': build_dtu_topic(dtu_config['customer_id'], 'dtu', dtu_config['dtu_id'], 'register'),
  2558. 'status': build_dtu_topic(dtu_config['customer_id'], 'dtu', dtu_config['dtu_id'], 'status'),
  2559. 'control': build_dtu_topic(dtu_config['customer_id'], 'dtu', dtu_config['dtu_id'], 'control'),
  2560. 'response': build_dtu_topic(dtu_config['customer_id'], 'dtu', dtu_config['dtu_id'], 'response'),
  2561. 'event': build_dtu_topic(dtu_config['customer_id'], 'dtu', dtu_config['dtu_id'], 'event'),
  2562. 'alarm': build_dtu_topic(dtu_config['customer_id'], 'dtu', dtu_config['dtu_id'], 'alarm')
  2563. }
  2564. }
  2565. return jsonify({'success': True, 'data': status})
  2566. except Exception as e:
  2567. logger.error(f"获取DTU状态失败: {str(e)}")
  2568. return jsonify({'success': False, 'message': str(e)}), 500
  2569. @app.route('/api/dtu/control', methods=['POST'])
  2570. def dtu_control():
  2571. """发送DTU控制命令(HTTP 入口,复用 dtu_handle_control 逻辑)"""
  2572. try:
  2573. data = request.json or {}
  2574. # 兼容两种入参格式:
  2575. # 1) {"command": "...", "target": "...", "params": {...}} 简洁格式
  2576. # 2) {"payload": {"command": "...", "target": "...", "params": {...}}} MQTT envelope 格式
  2577. inner = data.get('payload') if isinstance(data.get('payload'), dict) else data
  2578. command = inner.get('command')
  2579. # REBOOT 走快速通道立即执行
  2580. if command == 'REBOOT':
  2581. import subprocess
  2582. logger.warning("执行系统重启命令")
  2583. threading.Thread(target=lambda: (
  2584. time.sleep(1),
  2585. subprocess.run(['reboot'], capture_output=True)
  2586. ), daemon=True).start()
  2587. return jsonify({'success': True, 'message': '系统正在重启...'})
  2588. # 其他命令委托给 MQTT 处理器(构造 MQTT envelope 后调用)
  2589. mqtt_envelope = {
  2590. 'msg_id': data.get('msg_id', f"http_{int(time.time() * 1000)}_{uuid.uuid4().hex[:6]}"),
  2591. 'timestamp': data.get('timestamp', int(time.time() * 1000)),
  2592. 'dtu_id': data.get('dtu_id', dtu_config.get('dtu_id')),
  2593. 'type': 'CONTROL',
  2594. 'payload': {
  2595. 'command': command,
  2596. 'target': inner.get('target'),
  2597. 'params': inner.get('params', {})
  2598. }
  2599. }
  2600. # 调用 MQTT 处理器(其内部会发送响应到 dtu/.../response 主题)
  2601. dtu_handle_control('http/control', mqtt_envelope)
  2602. # 如果命令是 QUERY_JUMPER_STATUS / QUERY_ENV_SENSOR / QUERY_DTU_STATUS,
  2603. # 实际响应已发布到对应主题,HTTP 端点仅返回成功标记
  2604. if command in ('QUERY_JUMPER_STATUS', 'QUERY_ENV_SENSOR', 'QUERY_DTU_STATUS'):
  2605. return jsonify({
  2606. 'success': True,
  2607. 'message': f'{command} 已发布到 MQTT 主题'
  2608. })
  2609. # 其他命令: 从 MQTT 主题同步拉取最新响应(简化处理:直接返回已发布标记)
  2610. return jsonify({
  2611. 'success': True,
  2612. 'message': f'{command} 命令已下发,等待 MQTT 响应'
  2613. })
  2614. except Exception as e:
  2615. logger.error(f"DTU控制命令失败: {str(e)}")
  2616. return jsonify({'success': False, 'message': str(e)}), 500
  2617. # OTA状态存储
  2618. ota_status = {
  2619. 'status': 'IDLE', # IDLE, DOWNLOADING, VERIFYING, FLASHING, SUCCESS, FAILED
  2620. 'progress': 0,
  2621. 'firmware_version': None,
  2622. 'target_version': None,
  2623. 'error_code': None,
  2624. 'error_message': None,
  2625. 'last_update': None
  2626. }
  2627. @app.route('/api/dtu/ota_detect', methods=['POST'])
  2628. def ota_detect():
  2629. """检测固件包信息(自动下载并解析 manifest)"""
  2630. try:
  2631. url = request.json.get('url', '')
  2632. if not url:
  2633. return jsonify({'success': False, 'message': '请提供固件 URL'}), 400
  2634. import urllib.request, tempfile, tarfile, json, hashlib
  2635. tmp = tempfile.mktemp(suffix='.tar.gz')
  2636. try:
  2637. urllib.request.urlretrieve(url, tmp)
  2638. except Exception as e:
  2639. return jsonify({'success': False, 'message': f'下载失败: {str(e)}'}), 400
  2640. file_size = os.path.getsize(tmp)
  2641. h = hashlib.md5()
  2642. with open(tmp, 'rb') as f:
  2643. for chunk in iter(lambda: f.read(65536), b''):
  2644. h.update(chunk)
  2645. md5sum = h.hexdigest()
  2646. version = ''
  2647. try:
  2648. with tarfile.open(tmp, 'r:gz') as tar:
  2649. m = tar.extractfile('firmware/firmware.json')
  2650. if m:
  2651. manifest = json.loads(m.read())
  2652. version = manifest.get('version', '')
  2653. except Exception as e:
  2654. logger.warning(f"解析 firmware.json 失败: {e}")
  2655. os.remove(tmp)
  2656. return jsonify({
  2657. 'success': True,
  2658. 'data': {
  2659. 'file_size': file_size,
  2660. 'checksum': md5sum,
  2661. 'checksum_type': 'MD5',
  2662. 'firmware_version': version
  2663. }
  2664. })
  2665. except Exception as e:
  2666. logger.error(f"检测固件失败: {str(e)}")
  2667. return jsonify({'success': False, 'message': str(e)}), 500
  2668. @app.route('/api/dtu/ota_status', methods=['GET'])
  2669. def get_ota_status():
  2670. """获取OTA升级状态"""
  2671. try:
  2672. return jsonify({
  2673. 'success': True,
  2674. 'data': {
  2675. 'current_firmware': dtu_config.get('firmware_version', 'v1.0.0'),
  2676. 'ota_status': ota_status.get('status', 'IDLE'),
  2677. 'ota_progress': ota_status.get('progress', 0),
  2678. 'target_version': ota_status.get('target_version'),
  2679. 'error_code': ota_status.get('error_code'),
  2680. 'error_message': ota_status.get('error_message'),
  2681. 'last_update': ota_status.get('last_update')
  2682. }
  2683. })
  2684. except Exception as e:
  2685. logger.error(f"获取OTA状态失败: {str(e)}")
  2686. return jsonify({'success': False, 'message': str(e)}), 500
  2687. def handle_ota_status(dtu_id, payload):
  2688. """处理MQTT上报的OTA状态"""
  2689. ota_status['status'] = payload.get('ota_status', 'IDLE')
  2690. ota_status['progress'] = payload.get('ota_progress', 0)
  2691. ota_status['firmware_version'] = payload.get('firmware_version')
  2692. ota_status['last_update'] = time.strftime('%Y-%m-%d %H:%M:%S')
  2693. error_code = payload.get('error_code')
  2694. if error_code and error_code != 0:
  2695. ota_status['error_code'] = error_code
  2696. ota_status['error_message'] = payload.get('error_message', get_ota_error_message(error_code))
  2697. ota_status['status'] = 'FAILED'
  2698. elif ota_status['status'] == 'SUCCESS':
  2699. ota_status['error_code'] = None
  2700. ota_status['error_message'] = None
  2701. if ota_status.get('firmware_version'):
  2702. dtu_config['firmware_version'] = ota_status['firmware_version']
  2703. save_dtu_config()
  2704. logger.info(f"MQTT OTA状态更新: {ota_status['status']}, 进度: {ota_status['progress']}%")
  2705. def get_ota_error_message(error_code):
  2706. error_messages = {1021: '已是目标版本', 1022: '校验失败', 1023: '下载失败', 1024: '写入失败', 1025: '存储空间不足'}
  2707. return error_messages.get(error_code, f'未知错误码: {error_code}')
  2708. def _publish_ota_progress():
  2709. """通过MQTT发布OTA进度"""
  2710. try:
  2711. if not mqtt_client.get_status():
  2712. return
  2713. topic = build_dtu_topic(dtu_config['customer_id'], 'dtu', dtu_config['dtu_id'], 'status')
  2714. payload = {
  2715. 'msg_id': f"ota_status_{int(time.time() * 1000)}",
  2716. 'timestamp': int(time.time() * 1000),
  2717. 'dtu_id': dtu_config['dtu_id'],
  2718. 'type': 'STATUS',
  2719. 'payload': {
  2720. 'ota_status': ota_status.get('status', 'IDLE'),
  2721. 'ota_progress': ota_status.get('progress', 0),
  2722. 'firmware_version': ota_status.get('target_version', ''),
  2723. 'error_code': ota_status.get('error_code', 0),
  2724. 'error_message': ota_status.get('error_message')
  2725. }
  2726. }
  2727. mqtt_client.publish(topic, json.dumps(payload), qos=1)
  2728. except Exception:
  2729. pass
  2730. def _send_ota_error_response(error_code, original_msg_id=''):
  2731. """通过MQTT发送OTA错误响应(1021-1025)"""
  2732. try:
  2733. if not mqtt_client.get_status():
  2734. return
  2735. topic = build_dtu_topic(dtu_config['customer_id'], 'dtu', dtu_config['dtu_id'], 'response')
  2736. payload = {
  2737. 'msg_id': f"rsp_ota_{int(time.time() * 1000)}_{uuid.uuid4().hex[:6]}",
  2738. 'timestamp': int(time.time() * 1000),
  2739. 'dtu_id': dtu_config['dtu_id'],
  2740. 'type': 'RESPONSE',
  2741. 'payload': {
  2742. 'original_msg_id': original_msg_id,
  2743. 'command': 'OTA_UPGRADE',
  2744. 'success': False,
  2745. 'result': None,
  2746. 'error_code': error_code,
  2747. 'error_message': get_ota_error_message(error_code)
  2748. }
  2749. }
  2750. mqtt_client.publish(topic, json.dumps(payload), qos=1)
  2751. except Exception:
  2752. pass
  2753. def _run_ota(params):
  2754. """执行OTA升级(后台线程,供HTTP和MQTT共用)"""
  2755. import subprocess, os, hashlib, shutil, tarfile, tempfile, urllib.request
  2756. url = params['firmware_url']
  2757. version = params['firmware_version']
  2758. file_size = params['file_size']
  2759. checksum = params['checksum']
  2760. checksum_type = params.get('checksum_type', 'MD5').upper()
  2761. force = params.get('force_upgrade', False)
  2762. ota_status['status'] = 'DOWNLOADING'
  2763. ota_status['progress'] = 0
  2764. ota_status['target_version'] = version
  2765. ota_status['error_code'] = None
  2766. ota_status['error_message'] = None
  2767. ota_status['last_update'] = time.strftime('%Y-%m-%d %H:%M:%S')
  2768. _publish_ota_progress()
  2769. try:
  2770. tmp_dir = tempfile.mkdtemp(prefix='ota_')
  2771. fw_path = os.path.join(tmp_dir, 'firmware.tar.gz')
  2772. ota_status['progress'] = 5
  2773. ota_status['last_update'] = time.strftime('%Y-%m-%d %H:%M:%S')
  2774. logger.info(f"OTA: 开始下载固件 {url}")
  2775. req = urllib.request.Request(url, headers={'User-Agent': 'OTA-Updater'})
  2776. # 进度上报节流: 每 10% 或每 5 秒
  2777. last_published_pct = 0
  2778. last_publish_time = time.time()
  2779. with urllib.request.urlopen(req, timeout=120) as resp:
  2780. with open(fw_path, 'wb') as f:
  2781. total = int(resp.headers.get('Content-Length', 0))
  2782. downloaded = 0
  2783. while True:
  2784. chunk = resp.read(65536)
  2785. if not chunk: break
  2786. f.write(chunk)
  2787. downloaded += len(chunk)
  2788. if total:
  2789. pct = 5 + int(downloaded / total * 30)
  2790. ota_status['progress'] = min(pct, 35)
  2791. now = time.time()
  2792. if pct - last_published_pct >= 10 or (now - last_publish_time) >= 5:
  2793. last_published_pct = pct
  2794. last_publish_time = now
  2795. ota_status['last_update'] = time.strftime('%Y-%m-%d %H:%M:%S')
  2796. _publish_ota_progress()
  2797. dl_size = os.path.getsize(fw_path)
  2798. ota_status['progress'] = 40
  2799. ota_status['last_update'] = time.strftime('%Y-%m-%d %H:%M:%S')
  2800. _publish_ota_progress()
  2801. if abs(dl_size - file_size) > 1024:
  2802. raise Exception(f"文件大小不匹配: 预期{file_size}, 实际{dl_size}")
  2803. ota_status['status'] = 'VERIFYING'
  2804. ota_status['progress'] = 50
  2805. ota_status['last_update'] = time.strftime('%Y-%m-%d %H:%M:%S')
  2806. _publish_ota_progress()
  2807. h = hashlib.new(checksum_type)
  2808. with open(fw_path, 'rb') as f:
  2809. for chunk in iter(lambda: f.read(65536), b''): h.update(chunk)
  2810. actual_checksum = h.hexdigest().lower()
  2811. if actual_checksum != checksum.lower():
  2812. raise Exception(f"校验和不匹配: 预期{checksum}, 实际{actual_checksum}")
  2813. ota_status['progress'] = 60
  2814. ota_status['last_update'] = time.strftime('%Y-%m-%d %H:%M:%S')
  2815. _publish_ota_progress()
  2816. extract_dir = os.path.join(tmp_dir, 'firmware')
  2817. os.makedirs(extract_dir, exist_ok=True)
  2818. with tarfile.open(fw_path, 'r:gz') as tar: tar.extractall(extract_dir)
  2819. manifest_path = os.path.join(extract_dir, 'firmware', 'firmware.json')
  2820. if not os.path.exists(manifest_path):
  2821. raise Exception("固件包缺少 firmware.json")
  2822. with open(manifest_path, 'r') as f: manifest = json.load(f)
  2823. fw_version = manifest.get('version', version)
  2824. if not force:
  2825. current_ver = dtu_config.get('firmware_version', 'v0.0.0')
  2826. if fw_version == current_ver:
  2827. raise Exception(f'已是目标版本 {current_ver}')
  2828. ota_status['status'] = 'FLASHING'
  2829. ota_status['progress'] = 70
  2830. ota_status['last_update'] = time.strftime('%Y-%m-%d %H:%M:%S')
  2831. _publish_ota_progress()
  2832. fw_root = os.path.join(extract_dir, 'firmware')
  2833. project_dir = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
  2834. excludes = {'__pycache__', '.git', 'log', 'venv'}
  2835. for root, dirs, files in os.walk(fw_root):
  2836. rel = os.path.relpath(root, fw_root)
  2837. if rel == '.': rel = ''
  2838. parts = rel.split(os.sep) if rel else []
  2839. if parts and parts[0] in excludes: continue
  2840. for fname in files:
  2841. if fname == 'firmware.json': continue
  2842. src = os.path.join(root, fname)
  2843. dst = os.path.join(project_dir, rel, fname)
  2844. os.makedirs(os.path.dirname(dst), exist_ok=True)
  2845. shutil.copy2(src, dst)
  2846. ota_status['progress'] = 90
  2847. ota_status['last_update'] = time.strftime('%Y-%m-%d %H:%M:%S')
  2848. dtu_config['firmware_version'] = fw_version
  2849. save_dtu_config()
  2850. ota_status['status'] = 'SUCCESS'
  2851. ota_status['progress'] = 100
  2852. ota_status['firmware_version'] = fw_version
  2853. ota_status['last_update'] = time.strftime('%Y-%m-%d %H:%M:%S')
  2854. _publish_ota_progress()
  2855. shutil.rmtree(tmp_dir, ignore_errors=True)
  2856. for i in range(10, 0, -1):
  2857. ota_status['last_update'] = time.strftime('%Y-%m-%d %H:%M:%S')
  2858. time.sleep(1)
  2859. logger.info("OTA: 重启服务...")
  2860. subprocess.Popen(
  2861. [subprocess.sys.executable, '-m', 'flask', 'run', '--host=0.0.0.0', '--port=5001'],
  2862. cwd=os.path.dirname(os.path.abspath(__file__)),
  2863. stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL
  2864. )
  2865. os._exit(0)
  2866. except Exception as e:
  2867. ota_status['status'] = 'FAILED'
  2868. ota_status['error_message'] = str(e)
  2869. ota_status['last_update'] = time.strftime('%Y-%m-%d %H:%M:%S')
  2870. # 映射错误到标准错误码 (1021-1025)
  2871. msg = str(e)
  2872. if '校验和' in msg or 'firmware.json' in msg:
  2873. ota_status['error_code'] = 1022
  2874. elif '已是目标版本' in msg:
  2875. ota_status['error_code'] = 1021
  2876. elif '下载' in msg or 'urlopen' in msg or 'Connection' in msg or 'timeout' in msg:
  2877. ota_status['error_code'] = 1023
  2878. elif '磁盘' in msg or '空间' in msg or 'No space' in msg:
  2879. ota_status['error_code'] = 1025
  2880. elif '写入' in msg or 'flash' in msg.lower() or 'shutil' in msg:
  2881. ota_status['error_code'] = 1024
  2882. else:
  2883. ota_status['error_code'] = 9999
  2884. _publish_ota_progress()
  2885. # 通过MQTT发送错误响应(同时适配HTTP和MQTT触发)
  2886. _send_ota_error_response(ota_status['error_code'], params.get('_incoming_msg_id', ''))
  2887. logger.error(f"OTA: 升级失败 - {str(e)}")
  2888. if 'tmp_dir' in dir() and tmp_dir and os.path.exists(tmp_dir):
  2889. shutil.rmtree(tmp_dir, ignore_errors=True)
  2890. @app.route('/api/dtu/ota_upgrade', methods=['POST'])
  2891. def trigger_ota_upgrade():
  2892. """触发OTA升级"""
  2893. try:
  2894. data = request.json
  2895. if not data:
  2896. return jsonify({'success': False, 'message': '请求体不能为空'}), 400
  2897. required_fields = ['firmware_url', 'firmware_version', 'file_size', 'checksum', 'checksum_type']
  2898. missing = [f for f in required_fields if not data.get(f)]
  2899. if missing:
  2900. return jsonify({'success': False, 'message': f'缺少必填字段: {", ".join(missing)}'}), 400
  2901. # 已是目标版本且未强制: 1021 (同时通过 MQTT 上报响应)
  2902. force_upgrade = data.get('force_upgrade', False)
  2903. current_version = dtu_config.get('firmware_version', 'v1.0.0')
  2904. if data['firmware_version'] == current_version and not force_upgrade:
  2905. # 通过 MQTT 上报 1021 响应 (与其他错误码一致)
  2906. try:
  2907. _send_ota_error_response(1021, data.get('msg_id', ''))
  2908. except Exception:
  2909. pass
  2910. return jsonify({
  2911. 'success': False,
  2912. 'error_code': 1021,
  2913. 'message': f'已是目标版本: {current_version}'
  2914. }), 400
  2915. # 注入 _incoming_msg_id 使 MQTT 错误响应可回填
  2916. ota_params = {**data, '_incoming_msg_id': data.get('msg_id', '')}
  2917. t = threading.Thread(target=_run_ota, args=(ota_params,), daemon=True)
  2918. t.start()
  2919. return jsonify({
  2920. 'success': True,
  2921. 'message': 'OTA升级已启动',
  2922. 'data': {
  2923. 'target_version': data['firmware_version'],
  2924. 'ota_status': 'DOWNLOADING'
  2925. }
  2926. })
  2927. except Exception as e:
  2928. logger.error(f"触发OTA升级失败: {str(e)}")
  2929. return jsonify({'success': False, 'message': str(e)}), 500
  2930. # ==================== 环境传感器 API ====================
  2931. # 环境传感器数据存储
  2932. env_sensor_data = {
  2933. 'temperature': None, # 环境温度 (℃)
  2934. 'humidity': None, # 环境湿度 (%)
  2935. 'dtu_temperature': None, # DTU主板温度 (℃)
  2936. 'update_time': None, # 数据更新时间
  2937. 'sensor_update_time': None, # 传感器更新时间
  2938. 'connected': False # 传感器连接状态
  2939. }
  2940. # 环境传感器告警阈值
  2941. env_sensor_threshold = {
  2942. 'temp_high': 45.0, # 环境温度上限 (℃)
  2943. 'temp_low': -10.0, # 环境温度下限 (℃)
  2944. 'humidity_high': 80.0, # 环境湿度上限 (%)
  2945. 'humidity_low': 20.0, # 环境湿度下限 (%)
  2946. 'dtu_temp_high': 70.0 # DTU主板温度上限 (℃)
  2947. }
  2948. # 环境传感器历史数据 (保留最近1000条)
  2949. env_sensor_history = []
  2950. @app.route('/api/sensor/env', methods=['GET'])
  2951. def get_env_sensor_data():
  2952. """获取环境传感器当前数据"""
  2953. data = dict(env_sensor_data)
  2954. # DTU 主板温度回退: 若外部未上报,从 sysfs 本地读取
  2955. if data.get('dtu_temperature') is None:
  2956. try:
  2957. import psutil
  2958. temps = psutil.sensors_temperatures(fahrenheit=False) or {}
  2959. for chip_name, entries in temps.items():
  2960. for entry in entries:
  2961. if entry.current and entry.current > 0:
  2962. data['dtu_temperature'] = round(entry.current, 1)
  2963. break
  2964. if data['dtu_temperature'] is not None:
  2965. break
  2966. except Exception:
  2967. pass
  2968. if data.get('dtu_temperature') is None:
  2969. import glob as _glob
  2970. for tz in sorted(_glob.glob('/sys/class/thermal/thermal_zone*/temp')):
  2971. try:
  2972. with open(tz) as f:
  2973. v = int(f.read().strip())
  2974. if v > 0:
  2975. data['dtu_temperature'] = round(v / 1000.0, 1)
  2976. data['update_time'] = data.get('update_time') or __import__('datetime').datetime.now().strftime('%Y-%m-%d %H:%M:%S')
  2977. break
  2978. except Exception:
  2979. continue
  2980. # 附加 CPU/内存使用率(始终本地读取)
  2981. data['cpu_usage'] = _get_cpu_usage()
  2982. data['memory_usage'] = _get_memory_usage()
  2983. return jsonify({
  2984. 'success': True,
  2985. 'data': data
  2986. })
  2987. @app.route('/api/sensor/threshold', methods=['GET', 'POST'])
  2988. def handle_env_sensor_threshold():
  2989. """获取或设置环境传感器告警阈值"""
  2990. if request.method == 'GET':
  2991. return jsonify({
  2992. 'success': True,
  2993. 'data': env_sensor_threshold
  2994. })
  2995. else:
  2996. try:
  2997. data = request.get_json()
  2998. for key in env_sensor_threshold.keys():
  2999. if key in data:
  3000. env_sensor_threshold[key] = float(data[key])
  3001. logger.info(f"更新环境传感器告警阈值: {env_sensor_threshold}")
  3002. return jsonify({
  3003. 'success': True,
  3004. 'message': '阈值设置成功',
  3005. 'data': env_sensor_threshold
  3006. })
  3007. except Exception as e:
  3008. logger.error(f"设置告警阈值失败: {str(e)}")
  3009. return jsonify({'success': False, 'message': str(e)}), 400
  3010. @app.route('/api/sensor/history', methods=['GET'])
  3011. def get_env_sensor_history():
  3012. """获取环境传感器历史数据"""
  3013. try:
  3014. import datetime
  3015. # 获取查询参数
  3016. range_param = request.args.get('range', '24h')
  3017. start_time = request.args.get('start_time')
  3018. end_time = request.args.get('end_time')
  3019. limit = request.args.get('limit', 100, type=int)
  3020. filtered_history = env_sensor_history
  3021. # 解析range参数
  3022. if not start_time and not end_time:
  3023. now = datetime.datetime.now()
  3024. if range_param == '1h':
  3025. start_time = (now - datetime.timedelta(hours=1)).strftime('%Y-%m-%d %H:%M:%S')
  3026. elif range_param == '6h':
  3027. start_time = (now - datetime.timedelta(hours=6)).strftime('%Y-%m-%d %H:%M:%S')
  3028. elif range_param == '24h':
  3029. start_time = (now - datetime.timedelta(hours=24)).strftime('%Y-%m-%d %H:%M:%S')
  3030. elif range_param == '7d':
  3031. start_time = (now - datetime.timedelta(days=7)).strftime('%Y-%m-%d %H:%M:%S')
  3032. elif range_param == '30d':
  3033. start_time = (now - datetime.timedelta(days=30)).strftime('%Y-%m-%d %H:%M:%S')
  3034. # 按时间过滤
  3035. if start_time:
  3036. filtered_history = [h for h in filtered_history if h.get('update_time') >= start_time]
  3037. if end_time:
  3038. filtered_history = [h for h in filtered_history if h.get('update_time') <= end_time]
  3039. # 按时间排序 (新到旧)
  3040. filtered_history = sorted(filtered_history, key=lambda x: x.get('update_time', ''), reverse=True)
  3041. # 限制返回数量
  3042. filtered_history = filtered_history[:limit]
  3043. return jsonify({
  3044. 'success': True,
  3045. 'data': filtered_history,
  3046. 'total': len(filtered_history)
  3047. })
  3048. except Exception as e:
  3049. logger.error(f"获取历史数据失败: {str(e)}")
  3050. return jsonify({'success': False, 'message': str(e)}), 500
  3051. def update_env_sensor_data_full(temperature, humidity, dtu_temperature, sensor_update_time):
  3052. """更新环境传感器数据 (由MQTT消息触发) - 完整版,含历史和告警"""
  3053. import datetime
  3054. now = datetime.datetime.now().strftime('%Y-%m-%d %H:%M:%S')
  3055. env_sensor_data['temperature'] = temperature
  3056. env_sensor_data['humidity'] = humidity
  3057. env_sensor_data['dtu_temperature'] = dtu_temperature
  3058. env_sensor_data['sensor_update_time'] = sensor_update_time
  3059. env_sensor_data['update_time'] = now
  3060. env_sensor_data['connected'] = True
  3061. # 同步到 MQTT 上行所用的 _env_sensor_data
  3062. _env_sensor_data['temperature'] = temperature
  3063. _env_sensor_data['humidity'] = humidity
  3064. _env_sensor_data['sensor_update_time'] = sensor_update_time
  3065. # 添加到历史记录
  3066. history_record = {
  3067. 'temperature': temperature,
  3068. 'humidity': humidity,
  3069. 'dtu_temperature': dtu_temperature,
  3070. 'sensor_update_time': sensor_update_time,
  3071. 'update_time': now
  3072. }
  3073. env_sensor_history.append(history_record)
  3074. # 保留最近1000条
  3075. if len(env_sensor_history) > 1000:
  3076. env_sensor_history[:] = env_sensor_history[-1000:]
  3077. # 检查告警
  3078. check_env_sensor_alarms(temperature, humidity, dtu_temperature)
  3079. logger.info(f"环境传感器数据更新: 温度={temperature}℃, 湿度={humidity}%, DTU温度={dtu_temperature}℃")
  3080. def dht11_data_callback(temperature, humidity):
  3081. """DHT11 本地 GPIO 传感器数据回调。
  3082. 将读取到的温湿度更新到环境传感器数据区,并通过 MQTT 上报。
  3083. 保留已有的 dtu_temperature(主板温度),仅更新环境温湿度字段。
  3084. """
  3085. try:
  3086. sensor_update_time = int(time.time() * 1000)
  3087. dtu_temperature = env_sensor_data.get('dtu_temperature')
  3088. update_env_sensor_data_full(temperature, humidity, dtu_temperature, sensor_update_time)
  3089. # 立即通过 MQTT 上报环境传感器数据
  3090. try:
  3091. dtu_publish_env_sensor()
  3092. except Exception as pub_err:
  3093. logger.warning(f"DHT11 数据 MQTT 上报失败: {pub_err}")
  3094. logger.info(f"DHT11 本地传感器数据已处理: 温度={temperature}°C, 湿度={humidity}%")
  3095. except Exception as e:
  3096. logger.error(f"处理 DHT11 本地传感器数据失败: {e}")
  3097. def check_env_sensor_alarms(temperature, humidity, dtu_temperature):
  3098. """检查环境传感器告警"""
  3099. alarms = []
  3100. if temperature is not None:
  3101. if temperature > env_sensor_threshold['temp_high']:
  3102. alarms.append({
  3103. 'type': 'temperature_high',
  3104. 'message': f'环境温度过高: {temperature}℃ (阈值: {env_sensor_threshold["temp_high"]}℃)',
  3105. 'level': 'warning'
  3106. })
  3107. elif temperature < env_sensor_threshold['temp_low']:
  3108. alarms.append({
  3109. 'type': 'temperature_low',
  3110. 'message': f'环境温度过低: {temperature}℃ (阈值: {env_sensor_threshold["temp_low"]}℃)',
  3111. 'level': 'warning'
  3112. })
  3113. if humidity is not None:
  3114. if humidity > env_sensor_threshold['humidity_high']:
  3115. alarms.append({
  3116. 'type': 'humidity_high',
  3117. 'message': f'环境湿度过高: {humidity}% (阈值: {env_sensor_threshold["humidity_high"]}%)',
  3118. 'level': 'warning'
  3119. })
  3120. elif humidity < env_sensor_threshold['humidity_low']:
  3121. alarms.append({
  3122. 'type': 'humidity_low',
  3123. 'message': f'环境湿度过低: {humidity}% (阈值: {env_sensor_threshold["humidity_low"]}%)',
  3124. 'level': 'warning'
  3125. })
  3126. if dtu_temperature is not None:
  3127. if dtu_temperature > env_sensor_threshold['dtu_temp_high']:
  3128. alarms.append({
  3129. 'type': 'dtu_temp_high',
  3130. 'message': f'DTU主板温度过高: {dtu_temperature}℃ (阈值: {env_sensor_threshold["dtu_temp_high"]}℃)',
  3131. 'level': 'critical'
  3132. })
  3133. # 发送告警通知
  3134. for alarm in alarms:
  3135. logger.warning(f"环境传感器告警: {alarm['message']}")
  3136. # 可以通过WebSocket发送告警
  3137. socketio.emit('env_sensor_alarm', alarm)
  3138. if __name__ == '__main__':
  3139. try:
  3140. # 启动前的初始化工作
  3141. logger.info('启动串口-MQTT网关服务...')
  3142. logger.info(f"配置信息: 主机={FLASK_HOST}, 端口={FLASK_PORT}, 调试模式={FLASK_DEBUG}")
  3143. # 加载DTU配置
  3144. load_dtu_config()
  3145. # 加载设备配置
  3146. loaded = load_device_config()
  3147. logger.info(f"已加载 {len(loaded)} 个设备配置")
  3148. panel_config.clear()
  3149. for uid_hex, addr in loaded.items():
  3150. panel_id = f"PANEL_{dtu_config.get('dtu_id', 'DTU001')}_{addr}"
  3151. panel_config[panel_id] = {'address': addr, 'position': addr, 'panel_uid': uid_hex}
  3152. logger.info(f"已加载 {len(panel_config)} 个面板配置")
  3153. # 启动自动发现循环
  3154. def auto_discover_loop():
  3155. while True:
  3156. time.sleep(30)
  3157. try:
  3158. _st = serial_client.get_status()
  3159. if not (isinstance(_st, dict) and _st.get('connected', False)):
  3160. continue
  3161. if not dtu_config.get('enabled'):
  3162. continue
  3163. result = address_config.auto_configure(timeout=2.0)
  3164. if result.get('discovered', 0) > 0:
  3165. logger.info(f"自动发现: {result.get('discovered')} 个设备")
  3166. save_device_config()
  3167. now = time.time()
  3168. for uid in address_config.get_stored_devices():
  3169. device_last_seen[uid] = now
  3170. for pid in panel_config:
  3171. device_last_seen[pid] = now
  3172. except Exception as e:
  3173. logger.error(f"自动发现异常: {str(e)}")
  3174. import threading
  3175. t = threading.Thread(target=auto_discover_loop, daemon=True)
  3176. t.start()
  3177. logger.info("启动自动发现线程 (间隔30秒)")
  3178. # 启动DTU心跳定时器
  3179. if dtu_config.get('enabled'):
  3180. start_dtu_heartbeat()
  3181. logger.info("启动DTU心跳定时器")
  3182. # 启动端口轮询循环
  3183. def port_poll_loop():
  3184. while True:
  3185. interval = dtu_config.get('poll_interval_ms', 5000) / 1000.0
  3186. time.sleep(max(1, interval))
  3187. try:
  3188. _st = serial_client.get_status()
  3189. if not (isinstance(_st, dict) and _st.get('connected', False)):
  3190. continue
  3191. if not dtu_config.get('enabled'):
  3192. continue
  3193. if not panel_config:
  3194. continue
  3195. online = 0
  3196. offline = 0
  3197. for panel_id, cfg in panel_config.items():
  3198. addr = cfg.get('address', 1)
  3199. panel_ok = False
  3200. panel_ports = []
  3201. for port_id in range(1, 25):
  3202. try:
  3203. result = modbus_client.read_antenna_card(addr, port_id, timeout=1.0)
  3204. except Exception:
  3205. result = {'error': 'exception'}
  3206. if 'error' in result:
  3207. panel_ports.append({'port_id': port_id, 'status': 'UNKNOWN', 'jumper_uid': None})
  3208. continue
  3209. panel_ok = True
  3210. card_str = result.get('card_number_str', '')
  3211. uid = card_str.upper() if card_str and card_str != '0000000000000000' else ''
  3212. if panel_id not in port_state:
  3213. port_state[panel_id] = {}
  3214. if port_id not in port_state[panel_id]:
  3215. port_state[panel_id][port_id] = {'last_uid': None, 'expected_uid': None, 'alarm_count': 0, 'last_polled_at': int(time.time() * 1000)}
  3216. ps = port_state[panel_id][port_id]
  3217. ps['last_polled_at'] = int(time.time() * 1000)
  3218. last_uid = ps.get('last_uid')
  3219. if uid and uid != last_uid:
  3220. dtu_publish_event(panel_id, port_id, 'CONNECT' if not last_uid else 'MOVE', uid, last_uid)
  3221. ps['last_uid'] = uid
  3222. elif not uid and last_uid:
  3223. dtu_publish_event(panel_id, port_id, 'DISCONNECT', None, last_uid)
  3224. ps['last_uid'] = None
  3225. ps_exp = ps.get('expected_uid')
  3226. if dtu_config.get('alarm_enabled', True):
  3227. if ps_exp and uid and uid != ps_exp:
  3228. ps['alarm_count'] = ps.get('alarm_count', 0) + 1
  3229. sev = 'CRITICAL' if ps['alarm_count'] >= 3 else 'WARNING'
  3230. dtu_publish_alarm(panel_id, port_id, 'ILLEGAL_CONNECT', ps_exp, uid, sev)
  3231. elif ps_exp and not uid:
  3232. ps['alarm_count'] = ps.get('alarm_count', 0) + 1
  3233. sev = 'CRITICAL' if ps['alarm_count'] >= 3 else 'WARNING'
  3234. dtu_publish_alarm(panel_id, port_id, 'ILLEGAL_DISCONNECT', ps_exp, None, sev)
  3235. elif ps_exp and uid == ps_exp:
  3236. ps['alarm_count'] = 0
  3237. if uid:
  3238. port_status = 'ILLEGAL' if (ps_exp and uid != ps_exp) else 'CONNECTED'
  3239. else:
  3240. port_status = 'DISCONNECTED'
  3241. panel_ports.append({'port_id': port_id, 'status': port_status, 'jumper_uid': uid or None})
  3242. if panel_ok:
  3243. online += 1
  3244. device_last_seen[panel_id] = time.time()
  3245. device_last_seen[cfg.get('panel_uid', '')] = time.time()
  3246. dtu_publish_panel_status(panel_id, addr, panel_ports)
  3247. # 如果该 panel 是当前显示的 panel,刷新屏幕
  3248. panels = get_sorted_panels()
  3249. trigger_panel_id = None
  3250. with screen_lock:
  3251. if panels and 0 <= screen_current_panel_index < len(panels) and panels[screen_current_panel_index][0] == panel_id:
  3252. trigger_panel_id = panel_id
  3253. if trigger_panel_id:
  3254. try:
  3255. screen_refresh_queue.put((trigger_panel_id,))
  3256. except Exception as e:
  3257. logger.error(f"刷新屏幕失败: {e}")
  3258. else:
  3259. offline += 1
  3260. dtu_publish_status(force=True)
  3261. except Exception as e:
  3262. logger.error(f"端口轮询异常: {str(e)}")
  3263. t2 = threading.Thread(target=port_poll_loop, daemon=True)
  3264. t2.start()
  3265. logger.info("启动端口轮询线程 (间隔5秒)")
  3266. # 启动 DHT11 本地传感器读取线程
  3267. if dht11_sensor is not None:
  3268. dht11_sensor.on_data = dht11_data_callback
  3269. dht11_sensor.start()
  3270. # 自动连接上次使用的串口
  3271. logger.info("尝试自动连接串口...")
  3272. auto_connect_serial()
  3273. # 连接 Nextion 串口屏
  3274. if NEXTION_ENABLED:
  3275. try:
  3276. ok, msg = screen_display.connect(NEXTION_SCREEN_PORT, NEXTION_SCREEN_BAUD)
  3277. if ok:
  3278. screen_display.set_button_callback(screen_button_event_handler)
  3279. logger.info("Nextion 屏幕连接成功")
  3280. panels = get_sorted_panels()
  3281. if panels:
  3282. refresh_screen_panel(panels[0][0])
  3283. else:
  3284. logger.warning(f"Nextion 屏幕连接失败: {msg}")
  3285. except Exception as e:
  3286. logger.error(f"Nextion 屏幕初始化异常: {e}")
  3287. # 启动 Nextion 屏幕异步刷新工作线程
  3288. screen_refresh_thread = threading.Thread(target=_screen_refresh_worker, daemon=True)
  3289. screen_refresh_thread.start()
  3290. logger.info("启动 Nextion 屏幕异步刷新线程")
  3291. # 启动服务
  3292. socketio.run(
  3293. app,
  3294. host=FLASK_HOST,
  3295. port=FLASK_PORT,
  3296. debug=FLASK_DEBUG,
  3297. use_reloader=False, # 禁用重载器以避免重复初始化问题
  3298. log_output=False, # 禁用Flask的日志输出,使用我们自己的日志配置
  3299. allow_unsafe_werkzeug=True
  3300. )
  3301. except KeyboardInterrupt:
  3302. # 优雅退出
  3303. logger.info('正在关闭应用...')
  3304. try:
  3305. if isinstance(serial_client.get_status(), dict) and serial_client.get_status().get("connected", False):
  3306. serial_client.disconnect()
  3307. logger.info('串口连接已断开')
  3308. if mqtt_client.get_status():
  3309. mqtt_client.disconnect()
  3310. logger.info('MQTT连接已断开')
  3311. if dht11_sensor is not None:
  3312. dht11_sensor.stop()
  3313. logger.info('DHT11 传感器线程已停止')
  3314. except Exception as e:
  3315. logger.error(f'关闭连接时出错: {str(e)}')
  3316. # 清理WebSocket连接
  3317. for client_type in connected_clients:
  3318. connected_clients[client_type].clear()
  3319. logger.info('应用已安全关闭')
  3320. except Exception as e:
  3321. logger.exception(f'应用启动失败') # 使用exception记录完整堆栈
  3322. # 确保资源被释放
  3323. try:
  3324. serial_client.disconnect()
  3325. mqtt_client.disconnect()
  3326. except:
  3327. pass