245 lines
9.6 KiB
GDScript
245 lines
9.6 KiB
GDScript
class_name rep_network extends Node
|
||
|
||
|
||
const PROTOCOL_TYPES: = {'yau07tx': yau07.YaU07, 'json-capsrpb': capsrpb.CapsRpb, 'udp-json-mbcs': sch3.Sch3}
|
||
const TCP_PROTO: = {'5p28': tcp5p28.TCP5P28}
|
||
const SERIAL_PROTO: = {'uart': spt25.SPT25 }
|
||
|
||
var tick: int = 0
|
||
|
||
class Serial extends SerialPort:
|
||
func send_to(data: PackedByteArray): return write_raw(data)
|
||
func _to_string() -> String: return 'последовательный порт \"%s\" @ %d (%s)' % [self.port, self.baudrate, 'открыт' if is_open() else 'закрыт']
|
||
|
||
func try_open(parsing_proc: Callable):
|
||
if is_open():
|
||
return Error.OK
|
||
var rc = open(port)
|
||
if rc == Error.OK:
|
||
data_received.connect(parsing_proc)
|
||
start_monitoring(timeout)
|
||
return rc
|
||
|
||
|
||
class Socket extends PacketPeerUDP:
|
||
func send_to(addr, data):
|
||
set_dest_address(addr[0], addr[1])
|
||
put_packet(data)
|
||
|
||
|
||
class SocketTCP extends TCPServer:
|
||
const RX_TIMEOUT: int = 1000 ## Время сброса счётчика принятия части пакета.
|
||
var peerstream: = PacketPeerStream.new()
|
||
var unit_key: Array ## Ключ юнита.
|
||
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):
|
||
var peer = peerstream.get_stream_peer()
|
||
peer.put_data(data)
|
||
|
||
|
||
var poll_sockets: Array[Socket]
|
||
var tcp_sockets: Array
|
||
var units: Dictionary
|
||
var units_udp: Dictionary
|
||
var units_tcp: Dictionary
|
||
var units_serial: Dictionary
|
||
var serials: Dictionary
|
||
var dst_ports: Dictionary
|
||
var unit_keys: Dictionary
|
||
var sock_unicast: Socket
|
||
var sock_capsrpb: Socket
|
||
var logger_page: Node
|
||
|
||
|
||
func on_serial_data(data, unit): unit.parse(data, tick)
|
||
|
||
|
||
func create_socket(nm) -> Socket:
|
||
var addr = settings.UnitProfiles[nm][1]
|
||
var port = settings.UnitProfiles[nm][2][0]
|
||
var broad = settings.UnitProfiles[nm][3]
|
||
var bind = settings.UnitProfiles[nm][4]
|
||
var sock: = Socket.new()
|
||
sock.set_broadcast_enabled(broad)
|
||
var rc = Error.OK if not bind else sock.bind(port, addr)
|
||
var errlevel = log.INFO if rc == Error.OK else log.ERROR
|
||
log.message(errlevel, '%s: сокет:%s:%d привязанный:%s %s' % [nm, addr, port, ['нет', 'да'][int(bind)], error_string(rc)])
|
||
return sock
|
||
|
||
|
||
func create_tcpsocket(nm) -> SocketTCP:
|
||
var addr = settings.UnitProfiles[nm][1]
|
||
var port = settings.UnitProfiles[nm][2][0]
|
||
var sock = SocketTCP.new()
|
||
sock.unit_key = settings.get_unit_key(nm)
|
||
sock.listen(port, addr)
|
||
log.message(log.INFO, '%s: сокет:%s:%d привязанный:' % [nm, addr, port])
|
||
return sock
|
||
|
||
|
||
func create_serial(nm) -> Serial:
|
||
var port = settings.UnitProfiles[nm][1]
|
||
var baud = settings.UnitProfiles[nm][2][0]
|
||
var sp: = Serial.new()
|
||
sp.port = port
|
||
sp.baudrate = baud
|
||
sp.timeout = 20000
|
||
return sp
|
||
|
||
|
||
func init_sockets():
|
||
poll_sockets.append(create_socket('уарэп-яу07-частный'))
|
||
poll_sockets.append(create_socket('уарэп-яу07-общий'))
|
||
poll_sockets.append(create_socket('уарэп-капсрпб'))
|
||
poll_sockets.append(create_socket('уарэп-щ3'))
|
||
sock_unicast = poll_sockets[0]
|
||
sock_capsrpb = poll_sockets[2]
|
||
log.info('Модуль работы с сетью готов')
|
||
|
||
|
||
func _ready() -> void:
|
||
var tmr_wsa = Timer.new()
|
||
add_child(tmr_wsa)
|
||
tmr_wsa.one_shot = true
|
||
tmr_wsa.connect('timeout', init_sockets)
|
||
tmr_wsa.start(5.0)
|
||
|
||
for unit_name in settings.UnitProfiles:
|
||
var st = settings.UnitProfiles[unit_name]
|
||
var proto = st[0]
|
||
if proto in PROTOCOL_TYPES:
|
||
var unit = PROTOCOL_TYPES[proto].new(unit_name)
|
||
var unit_key = settings.get_unit_key(unit_name)
|
||
units[unit_key] = unit
|
||
units_udp[unit_key] = unit
|
||
if logger_page:
|
||
unit.connect('line_changed', Callable(logger_page, 'on_line_changed').bind(unit_key))
|
||
unit.connect('command_fail', Callable(logger_page, 'on_command_fail').bind(unit_key))
|
||
if len(st[2]) > 1:
|
||
dst_ports[unit_key] = st[2][1]
|
||
else:
|
||
dst_ports[unit_key] = st[2][0]
|
||
|
||
if proto in TCP_PROTO:
|
||
var unit = TCP_PROTO[proto].new(unit_name)
|
||
var unit_key = settings.get_unit_key(unit_name)
|
||
units_tcp[unit_key] = unit
|
||
if logger_page:
|
||
unit.connect('line_changed', Callable(logger_page, 'on_line_changed').bind(unit_key))
|
||
unit.connect('command_fail', Callable(logger_page, 'on_command_fail').bind(unit_key))
|
||
|
||
if proto in SERIAL_PROTO:
|
||
var unit = SERIAL_PROTO[proto].new(unit_name)
|
||
var unit_key = settings.get_unit_key(unit_name)
|
||
if logger_page:
|
||
unit.connect('line_changed', Callable(logger_page, 'on_line_changed').bind(unit_key))
|
||
unit.connect('command_fail', Callable(logger_page, 'on_command_fail').bind(unit_key))
|
||
var serial = create_serial(unit_name)
|
||
units_serial[unit_key] = unit
|
||
units[unit_key] = unit
|
||
serials[unit_key] = serial
|
||
var rc = serial.try_open(on_serial_data.bind(unit))
|
||
if rc != Error.OK:
|
||
log.message(log.ERROR, 'Невозможно открыть %s для \"%s\"' % [serial, unit_name])
|
||
else:
|
||
log.message(log.INFO, '%s для \"%s\"' % [serial, unit_name])
|
||
|
||
tcp_sockets.append([create_tcpsocket('уарэп-5п28'), units_tcp[settings.get_unit_key('уарэп-5п28')]])
|
||
|
||
for key in units: log.info('%s %s:%d' % [units[key].name, key[0], key[1]])
|
||
for key in units_tcp: log.info('%s %s:%d' % [units_tcp[key].name, key[0], key[1]])
|
||
for key in dst_ports: unit_keys[dst_ports[key]] = key
|
||
|
||
|
||
func poll_receive(sock: Socket) -> bool:
|
||
while sock.get_available_packet_count() > 0:
|
||
var data: = sock.get_packet()
|
||
var addr: = sock.get_packet_ip()
|
||
var port: = sock.get_packet_port()
|
||
var addr_port: = [addr, port]
|
||
if addr_port in units:
|
||
units[addr_port].parse(data, tick)
|
||
return false
|
||
if port in unit_keys:
|
||
units[unit_keys[port]].parse(data, tick)
|
||
return false
|
||
|
||
|
||
func poll_receive_tcp(sock: SocketTCP) -> bool:
|
||
if sock.is_connection_available(): # check if someone's trying to connect
|
||
var client = sock.take_connection() # accept connection
|
||
sock.peerstream.set_stream_peer(client) # bind peerstream to new client
|
||
var peer = sock.peerstream.get_stream_peer()
|
||
if peer:
|
||
var peer_len_data = peer.get_available_bytes()
|
||
if peer_len_data > 0:
|
||
get_client_data(sock, peer_len_data, peer)
|
||
if sock.unit_key in units_tcp and sock.rx_all: units_tcp[sock.unit_key].parse(sock.rx_data, tick)
|
||
return false
|
||
|
||
|
||
func get_client_data(sock: SocketTCP, peer_len_data, peer):
|
||
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, len_rx):
|
||
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:
|
||
tick = Time.get_ticks_msec()
|
||
for addr in units_udp:
|
||
var unit = units[addr]
|
||
match unit.process(tick):
|
||
Error.OK:
|
||
if sock_unicast:
|
||
sock_unicast.send_to([addr[0], dst_ports[addr]], unit.tx_data.slice(0, unit.tx_len))
|
||
Error.FAILED:
|
||
emit_signal('socket_error', 'ошибка: %s %s:%d' % [unit, addr[0], addr[1]])
|
||
for key in units_serial:
|
||
var unit = units_serial[key]
|
||
var serial = serials[key]
|
||
match unit.process(tick):
|
||
Error.OK: serial.send_to(unit.tx_data)
|
||
Error.FAILED: emit_signal('socket_error', 'ошибка: \"%s\"' % unit.name)
|
||
for sock_unit in tcp_sockets:
|
||
var sock = sock_unit[0]
|
||
var unit = sock_unit[1]
|
||
match unit.process(tick):
|
||
Error.OK: sock.send_to(unit.tx_data)
|
||
Error.FAILED: emit_signal('socket_error', 'ошибка: %s' % unit)
|
||
poll_receive_tcp(sock)
|
||
poll_sockets.any(func(sock): poll_receive(sock))
|