app.py 145 KB

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