Доработка. Обмен с 5П-28 вынесен в отдельный поток.

This commit is contained in:
MaD_CaT
2026-01-30 11:11:17 +03:00
parent 6950b09879
commit 8f4134a911
2 changed files with 159 additions and 104 deletions

View File

@@ -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)

View File

@@ -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)