app.py 183 KB

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