From 8f4134a911e5f9669f4ce1148a464fcaa58a28cc Mon Sep 17 00:00:00 2001 From: MaD_CaT Date: Fri, 30 Jan 2026 11:11:17 +0300 Subject: [PATCH] =?UTF-8?q?=D0=94=D0=BE=D1=80=D0=B0=D0=B1=D0=BE=D1=82?= =?UTF-8?q?=D0=BA=D0=B0.=20=D0=9E=D0=B1=D0=BC=D0=B5=D0=BD=20=D1=81=205?= =?UTF-8?q?=D0=9F-28=20=D0=B2=D1=8B=D0=BD=D0=B5=D1=81=D0=B5=D0=BD=20=D0=B2?= =?UTF-8?q?=20=D0=BE=D1=82=D0=B4=D0=B5=D0=BB=D1=8C=D0=BD=D1=8B=D0=B9=20?= =?UTF-8?q?=D0=BF=D0=BE=D1=82=D0=BE=D0=BA.?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- scripts/network.gd | 91 +----------------------- scripts/tcp5p28.gd | 172 +++++++++++++++++++++++++++++++++++++++++---- 2 files changed, 159 insertions(+), 104 deletions(-) diff --git a/scripts/network.gd b/scripts/network.gd index 6691fd59..7ea36cc9 100644 --- a/scripts/network.gd +++ b/scripts/network.gd @@ -110,24 +110,6 @@ class SocketUDP extends PacketPeerUDP: push_error('%s: %s', [error_string(rc), addr]) -class SocketTCP extends TCPServer: - const RX_TIMEOUT: int = 1000 ## Время сброса счётчика принятия части пакета. - var peerstream: = PacketPeerStream.new() - var unit_name: StringName ## Уникальное имя устройства - var head_len: int = 8 ## Длина заголовка команды. - var len_place: int = 4 ## Номер начального байта с длинной блока данных. - var type_len: int = 4 ## Количество байт в размере длинны. - var rx_tick: int ## Время приёма последнего пакета. - var rx_all: bool = true ## Флаг, что пакет принят полностью. - var rx_len: int ## Количество принятых байт. - var rx_data: PackedByteArray ## Принятые данные. - var data_len: int ## Размер ожидаемых данных. - - func send_to(data: PackedByteArray): - var peer = peerstream.get_stream_peer() - peer.put_data(data) - - var poll_sockets: Array[SocketUDP] ## Сокеты для непрерывного опроса наличия новых данных var tcp_sockets: Array[Array] ## Сокеты для TCP-соединений var units: Dictionary[StringName, unit.Unit] ## Экземпляры всех сетевых устройств @@ -187,21 +169,6 @@ func create_socket_caps(unit_name: StringName) -> SocketUDP: return sock -## [param unit_name] - Уникальное имя устройства[br] -func create_socket_tcp(unit_name: StringName) -> SocketTCP: - var st = settings.UnitProfiles[unit_name][1] - var addr = st[0] - var port = st[1] - var sock: = SocketTCP.new() - sock.unit_name = unit_name - var rc: = sock.listen(port, addr) - if rc == Error.OK: - log.message(log.INFO, '\"%s\" слушает на %s:%d' % [unit_name, addr, port]) - else: - log.message(log.ERROR, '\"%s\" %s:%d \"%s\"' % [unit_name, addr, port, error_string(rc)]) - return sock - - ## [param unit_name] - Уникальное имя устройства[br] func create_serial(unit_instance: unit.Unit) -> SocketSerial: var st = settings.UnitProfiles[unit_instance.name][1] @@ -279,9 +246,11 @@ func _ready() -> void: dst_ports[unit_name] = unit_profile[3] if proto in TCP_PROTO: + var addr = settings.UnitProfiles[unit_name][1] var unit_tcp = TCP_PROTO[proto].new(unit_name) units[unit_name] = unit_tcp units_tcp[unit_name] = unit_tcp + unit_tcp.open(addr[0], addr[1]) if proto in MODBUS_PROTO: var unit_modbus = MODBUS_PROTO[proto].new(unit_name) @@ -305,7 +274,6 @@ func _ready() -> void: units_nullp[unit_name] = unit_null - tcp_sockets.append([create_socket_tcp('уарэп-5п28'), units_tcp['уарэп-5п28']]) for key: StringName in dst_ports: port_to_unit_name[dst_ports[key]] = key if logger_page: @@ -347,55 +315,6 @@ func poll_receive_udp(sock: SocketUDP) -> bool: return false -func poll_receive_tcp(sock: SocketTCP) -> bool: - if sock.is_connection_available(): # Проверить, если кто то пытается подключиться - var client: = sock.take_connection() # Принять соединение - sock.peerstream.set_stream_peer(client) # Привязать поток к новому клиенту - log.info('подключен: %s:%d к порту: %d' % [ client.get_connected_host(), client.get_connected_port(), client.get_local_port() ]) - var peer = sock.peerstream.get_stream_peer() - if peer: - while true: - var peer_len_data: = peer.get_available_bytes() - if peer_len_data <= 0: - break - get_client_data(sock, peer, peer_len_data) - if sock.unit_name in units_tcp and sock.rx_all: - units_tcp[sock.unit_name].parse(sock.rx_data, tick) - return false - - -func get_client_data(sock: SocketTCP, peer: StreamPeer, peer_len_data: int) -> void: - if tick - sock.rx_tick > sock.RX_TIMEOUT: - sock.rx_all = true - sock.rx_tick = tick - if sock.rx_all: - if peer_len_data >= sock.head_len: - sock.rx_all = false - sock.rx_len = sock.head_len - var head = peer.get_data(sock.head_len) - sock.data_len = head[1].decode_u32(sock.len_place) - var len_data_in_buf = peer.get_available_bytes() - if len_data_in_buf > 0: - sock.rx_data = head[1] - var len_rx = len_data_in_buf - if len_data_in_buf >= sock.data_len: - len_rx = sock.data_len - get_tcp_data(sock, peer, len_rx) - else: - var len_rx = (sock.head_len + sock.data_len) - sock.rx_len - if peer_len_data < len_rx: - len_rx = peer_len_data - get_tcp_data(sock, peer, len_rx) - - -func get_tcp_data(sock: SocketTCP, peer: StreamPeer, len_rx: int): - var temp_data = peer.get_data(len_rx)[1] - sock.rx_data.append_array(temp_data) - sock.rx_len += len_rx - if sock.rx_len == sock.head_len + sock.data_len: - sock.rx_all = true - - func _process(_delta: float) -> void: if Engine.is_editor_hint(): return tick = Time.get_ticks_msec() @@ -421,12 +340,6 @@ func _process(_delta: float) -> void: for unit_name: StringName in units_modbus: var unit_modbus: = units_modbus[unit_name] unit_modbus.process(tick) - for sock_unit: Array in tcp_sockets: - var sock = sock_unit[0] - var unit_tcp = sock_unit[1] - if Error.OK == unit_tcp.process(tick): - sock.send_to(unit_tcp.tx_data) - poll_receive_tcp(sock) poll_sockets.any(poll_receive_udp) diff --git a/scripts/tcp5p28.gd b/scripts/tcp5p28.gd index 1e5699bb..1a3ed439 100644 --- a/scripts/tcp5p28.gd +++ b/scripts/tcp5p28.gd @@ -3,7 +3,6 @@ class_name tcp5_p28 extends Node ## Реализация "Протокол информационного сопряжения. ## СПО АСУ изделия МП-550 с СПО 5П-28 -const ONLINE_TIMEOUT = 5000 ## Время ожидания пакета от 5П-28, мс. const BUFFER_SIZE = 2048 ## Размер приёмного буфера для датаграмм. ## Тип модуляции для запаковки в пакет для 5П-28. @@ -56,24 +55,173 @@ class TCP5P28 extends unit.Unit: 'уарэп-яу07-4в': [8, 7], 'уарэп-яу07-4к': [16, 15] } + signal data_sended() + signal disconnected(rc: Array) + signal connected(rc: Array) signal get_threats(unit: TCP5P28) ## Сигнал запрос целей. signal get_interfers(unit: TCP5P28) ## Запрос состояния выполнения целеуказания. - var tx_stack: Array ## Массив пакетов для отправки. var uarep_state: int = 1 ## Флаги состояния прибора, 0 бит, состояние ПО var uarep_line: int = 3 ## Состояние линии связи с прибором var ext_cu: bool = false + var tx_queue: Array ## Очередь пакетов для отправки + var tx_mutex: Mutex ## Мютекс для очереди передачи + var thread: Thread ## Рабочий поток + var tctl_mutex: Mutex ## Мютекс для управления потоком + var tctl_run: bool ## Управление потоком + + enum STREAM_STATE { + IDLE, + CONNECT, + WAIT, + RX_START, + RX_DONE, + TX_START, + TX_DONE, + ERROR + } + + func _init(nm) -> void: super._init(nm) ext_cu = ProjectSettings.get_setting('application/config/Внешнее управление', false) ProjectSettings.connect('settings_changed', on_setting_changed) + tx_mutex = Mutex.new() + tctl_mutex = Mutex.new() + connect('data_received', on_data_received) + connect('data_sended', on_data_sended) + connect('connected', on_connected) + connect('disconnected', on_disconnected) init_state() + func on_connected(rc: Array): + log.info('Клиент 5П-28 подключен: %s:%s' % [ rc[0], rc[1]] ) + + func on_data_sended(): pass + + func on_data_received(data: PackedByteArray): + var tick: = Time.get_ticks_msec() + parse(data, tick) + + func on_disconnected(rc: Array): + log.info('Клиент 5П-28 отключен: %s:%s' % [ rc[0], rc[1]] ) + if online: + online = false + emit_signal('line_changed', self) + tx_mutex.lock() + tx_queue.clear() + tx_mutex.unlock() + func init_state(): for key in PR_STATE.keys(): var unit_pribor = network.get_unit_instance(key) unit_pribor.connect('line_changed', Callable(self, 'pribor_line_changed').bind(key)) + func open(host: String, port: int): + tctl_run = true + thread = Thread.new() + thread.start(thread_proc.bind(host, port)) + + func close(): + tctl_mutex.lock() + tctl_run = false + tctl_mutex.unlock() + + func thread_proc(host: String, port: int): + var sock = TCPServer.new() + var stream: = StreamPeerTCP.new() + var fsm: = STREAM_STATE.CONNECT + var rx_buff: PackedByteArray + var rx_bytes: int = 0 + var head_len: int = 8 ## Длина заголовка команды. + var len_place: int = 4 ## Номер начального байта с длинной блока данных. + var status: = StreamPeerTCP.Status.STATUS_NONE + var rc: = sock.listen(port, host) + var cl_host: String + var cl_port: int + while true: + stream.poll() + var new_status = stream.get_status() + if status != new_status: + status = new_status + match status: + stream.STATUS_NONE: + call_deferred('emit_signal', 'disconnected', [cl_host, cl_port]) + fsm = STREAM_STATE.ERROR + stream.STATUS_CONNECTED: + cl_host = stream.get_connected_host() + cl_port = stream.get_connected_port() + call_deferred('emit_signal', 'connected', [cl_host, cl_port]) + stream.STATUS_ERROR: + fsm = STREAM_STATE.ERROR + + if fsm == STREAM_STATE.CONNECT: + if sock.is_connection_available(): # Проверить, если кто то пытается подключиться + stream = sock.take_connection() # Принять соединение + new_status = stream.get_status() + if new_status == stream.STATUS_CONNECTED: + fsm = STREAM_STATE.WAIT + else: + fsm = STREAM_STATE.ERROR + elif fsm == STREAM_STATE.WAIT: + if status == stream.STATUS_CONNECTED: + fsm = STREAM_STATE.RX_START + else: + OS.delay_msec(50) + fsm = STREAM_STATE.CONNECT + elif fsm == STREAM_STATE.ERROR: + stream.disconnect_from_host() + fsm = STREAM_STATE.WAIT + OS.delay_msec(1000) + elif fsm == STREAM_STATE.IDLE: + OS.delay_msec(50) + fsm = STREAM_STATE.RX_START + elif fsm == STREAM_STATE.RX_START: + rx_buff.clear() + rx_bytes = 0 + var sz = stream.get_available_bytes() + if sz <= 0: + fsm = STREAM_STATE.TX_START + continue + var rxd = stream.get_data(head_len) + if rxd[0] == Error.OK: + rx_buff.append_array(rxd[1]) + else: + fsm = STREAM_STATE.ERROR + continue + var pay: = rx_buff.decode_u32(len_place) + rx_bytes = pay + head_len + if rx_buff.size() == rx_bytes: + fsm = STREAM_STATE.RX_DONE + else: + rxd = stream.get_data(rx_bytes - rx_buff.size()) + rc = rxd[0] + if rc == Error.OK: + rx_buff.append_array(rxd[1]) + fsm = STREAM_STATE.RX_DONE + else: + fsm = STREAM_STATE.ERROR + elif fsm == STREAM_STATE.RX_DONE: + var rx_data: = rx_buff.duplicate() + call_deferred('emit_signal', 'data_received', rx_data) + fsm = STREAM_STATE.TX_START + elif fsm == STREAM_STATE.TX_START: + tx_mutex.lock() + if not tx_queue.size(): + fsm = STREAM_STATE.IDLE + tx_mutex.unlock() + continue + var data = tx_queue.pop_front() + tx_mutex.unlock() + rc = stream.put_data(data) + if rc == Error.OK: + fsm = STREAM_STATE.TX_DONE + else: + fsm = STREAM_STATE.ERROR + elif fsm == STREAM_STATE.TX_DONE: + call_deferred('emit_signal', 'data_sended') + fsm = STREAM_STATE.RX_START + func pribor_line_changed(u, key): uarep_state = tools.set_bit(uarep_state, PR_STATE[key][0], u.online) uarep_line = tools.set_bit(uarep_state, PR_STATE[key][1], u.online) @@ -95,16 +243,10 @@ class TCP5P28 extends unit.Unit: if cmd_type == CMD_TYPE_IN and ext_cu: set_interfer(data) - func process(tick: int): - if online and ((tick - rx_tick) > ONLINE_TIMEOUT): - online = false - emit_signal('line_changed', self) - if len(tx_stack): - tx_data = tx_stack[0] - tx_stack.remove_at(0) - return Error.OK - else: - return Error.ERR_UNAVAILABLE + func queue_packet(packet: PackedByteArray): + tx_mutex.lock() + tx_queue.append(packet) + tx_mutex.unlock() func pack_threats(ths: Dictionary): var data = PackedByteArray() @@ -117,7 +259,7 @@ class TCP5P28 extends unit.Unit: data.append_array(get_threats_data(th, tick0)) var data_to_send: PackedByteArray = create_header(data_len) data_to_send.append_array(data) - tx_stack.append(data_to_send) + queue_packet(data_to_send) func pack_interfers(interfers: Dictionary, ecms: Dictionary): var data = PackedByteArray() @@ -141,7 +283,7 @@ class TCP5P28 extends unit.Unit: data.append_array(cu_data) var data_to_send: PackedByteArray = create_header(data_len + 4 * interfers.size()) data_to_send.append_array(data) - tx_stack.append(data_to_send) + queue_packet(data_to_send) func create_header(data_len: int): var head_data = PackedByteArray() @@ -210,7 +352,7 @@ class TCP5P28 extends unit.Unit: data_tx.encode_u16(0, err) # Код ошибки var data_to_send: PackedByteArray = create_header(data_len) data_to_send.append_array(data_tx) - tx_stack.append(data_to_send) + queue_packet(data_to_send) func on_setting_changed(): ext_cu = ProjectSettings.get_setting('application/config/Внешнее управление', false)