Merge remote-tracking branch 'MaD_CaT/master' into workflows
# Conflicts: # scripts/prd-ctl-config.gd
This commit is contained in:
@@ -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)
|
||||
|
||||
|
||||
|
||||
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user