From 860b73a8f408c03cad5a1745a2d255bbd77e6794 Mon Sep 17 00:00:00 2001 From: Onek8 Date: Sat, 3 Oct 2026 22:05:53 -0700 Subject: [PATCH] hxWebsocket Upstream and Network Updates --- leenkx/Sources/leenkx/network/Buffer.hx | 44 +++- leenkx/Sources/leenkx/network/Handler.hx | 2 + leenkx/Sources/leenkx/network/HttpHeader.hx | 1 + leenkx/Sources/leenkx/network/Leenkx.hx | 238 +++--------------- leenkx/Sources/leenkx/network/Log.hx | 51 ---- leenkx/Sources/leenkx/network/WebSocket.hx | 115 +++++---- .../Sources/leenkx/network/WebSocketCommon.hx | 41 ++- .../leenkx/network/WebSocketHandler.hx | 23 +- .../Sources/leenkx/network/WebSocketServer.hx | 4 - .../leenkx/network/krom/KromSecureSocket.hx | 11 +- .../Sources/leenkx/network/krom/KromSocket.hx | 29 +++ .../leenkx/network/krom/KromUdpSocket.hx | 70 ++++++ 12 files changed, 284 insertions(+), 345 deletions(-) delete mode 100644 leenkx/Sources/leenkx/network/Log.hx create mode 100644 leenkx/Sources/leenkx/network/krom/KromUdpSocket.hx diff --git a/leenkx/Sources/leenkx/network/Buffer.hx b/leenkx/Sources/leenkx/network/Buffer.hx index d3a4d02f..e003cfae 100644 --- a/leenkx/Sources/leenkx/network/Buffer.hx +++ b/leenkx/Sources/leenkx/network/Buffer.hx @@ -63,10 +63,10 @@ class Buffer { public function readUntil(delimiter:String):Bytes { var dl = delimiter.length; - for (i in 0 ... available - dl) { + for (i in 0 ... available - dl + 1) { var matched = true; for (j in 0 ... dl) { - if (peekByte(currentOffset + i + j + 1) == delimiter.charCodeAt(j)) { + if (peekByte(i + j) == delimiter.charCodeAt(j)) { continue; } matched = false; @@ -74,7 +74,8 @@ class Buffer { } if (matched) { - var bytes = readBytes(i + dl + 1); + var bytes = readBytes(i); + readBytes(dl); return bytes; } } @@ -85,7 +86,19 @@ class Buffer { public function readBytes(count:Int):Bytes { var count2 = Std.int(Math.min(count, available)); var out = Bytes.alloc(count2); - for (n in 0 ... count2) out.set(n, readByte()); + var written = 0; + while (written < count2) { + if (currentData == null || currentOffset >= currentData.length) { + currentOffset = 0; + currentData = chunks.shift(); + } + var n = Std.int(Math.min(count2 - written, currentData.length - currentOffset)); + out.blit(written, currentData, currentOffset, n); + written += n; + currentOffset += n; + available -= n; + } + length = available; return out; } @@ -115,16 +128,23 @@ class Buffer { } public function peekByte(offset:Int):Int { - if (available <= 0) throw 'No bytes available'; - var tempOffset = offset; - var tempData = chunks[0]; - if (tempData == null) { - tempData = currentData; + if (available <= 0 || offset < 0 || offset >= available) throw 'No bytes available'; + var tempData = currentData; + var chunkIndex = -1; + var tempOffset:Int; + if (tempData != null && currentOffset < tempData.length) { + tempOffset = currentOffset + offset; + } else { + tempData = chunks[0]; + chunkIndex = 0; + tempOffset = offset; } - var chunkIndex = 0; while (tempOffset >= tempData.length) { tempOffset -= tempData.length; chunkIndex++; + if (chunkIndex >= chunks.length){ + trace('No bytes available'); + } tempData = chunks[chunkIndex]; } return tempData.get(tempOffset); @@ -154,9 +174,9 @@ class Buffer { public function endsWith(e:String):Bool { var i = available - e.length; - var n = currentOffset; + var n = 0; - if (i <= 0) { + if (i < 0) { return false; } diff --git a/leenkx/Sources/leenkx/network/Handler.hx b/leenkx/Sources/leenkx/network/Handler.hx index 966384d0..f71c9770 100644 --- a/leenkx/Sources/leenkx/network/Handler.hx +++ b/leenkx/Sources/leenkx/network/Handler.hx @@ -1,6 +1,8 @@ package leenkx.network; class Handler extends WebSocketCommon { + public var validateHandshake:(HttpRequest, HttpResponse, (HttpResponse) -> Void) -> Void = null; + public function new(socket:SocketImpl) { super(socket); isClient = false; diff --git a/leenkx/Sources/leenkx/network/HttpHeader.hx b/leenkx/Sources/leenkx/network/HttpHeader.hx index cd8d382b..b4a8aca3 100644 --- a/leenkx/Sources/leenkx/network/HttpHeader.hx +++ b/leenkx/Sources/leenkx/network/HttpHeader.hx @@ -4,6 +4,7 @@ class HttpHeader { public static inline var SEC_WEBSOCKET_KEY:String = "Sec-WebSocket-Key"; public static inline var SEC_WEBSOSCKET_ACCEPT:String = "Sec-WebSocket-Accept"; public static inline var SEC_WEBSOSCKET_VERSION:String = "Sec-WebSocket-Version"; + public static inline var SEC_WEBSOCKET_PROTOCOL:String = "Sec-WebSocket-Protocol"; public static inline var UPGRADE:String = "Upgrade"; public static inline var HOST:String = "Host"; public static inline var CONNECTION:String = "Connection"; diff --git a/leenkx/Sources/leenkx/network/Leenkx.hx b/leenkx/Sources/leenkx/network/Leenkx.hx index de38fc15..65412591 100644 --- a/leenkx/Sources/leenkx/network/Leenkx.hx +++ b/leenkx/Sources/leenkx/network/Leenkx.hx @@ -1,10 +1,12 @@ package leenkx.network; +#if lnx_torrent import leenkx.network.Types; import haxe.io.Bytes; import iron.object.Object; import leenkx.system.Event; import leenkx.network.Buffer; +import leenkx.network.torrent.LeenkxClient; @:expose class Leenkx { @@ -20,24 +22,20 @@ class Leenkx { public static var onPingEvent: String = "Leenkx.onPing"; public static var onLeftEvent: String = "Leenkx.onLeft"; public static var onTimeOutEvent: String = "Leenkx.onTimeOut"; - public static var onRpcEvent: String = "Leenkx.onRpc"; - public static var onRpcResponseEvent: String = "Leenkx.onRpcResponse"; public static var onWireLeftEvent: String = "Leenkx.onWireLeft"; public static var onWireSeenEvent: String = "Leenkx.onWireSeen"; public static var onTorrentEvent: String = "Leenkx.onTorrent"; + public static var onTorrentAddedEvent: String = "Leenkx.onTorrentAdded"; public static var onTrackerEvent: String = "Leenkx.onTracker"; public static var onAnnounceEvent: String = "Leenkx.onAnnounce"; public static var onTorrentDoneEvent: String = "Leenkx.onTorrentDone"; public static var connections:Map = []; - #if js - public static var peers:js.lib.Map = new js.lib.Map(); - public static var data:js.lib.Map = new js.lib.Map(); - public static var id:js.lib.Map = new js.lib.Map(); - public static var torrent:js.lib.Map = new js.lib.Map(); - public static var file:js.lib.Map = new js.lib.Map(); - #end - public static var lxNew:Void->Void; - + public static var peers:Map = []; + public static var data:Map = []; + public static var id:Map = []; + public static var torrent:Map = []; + public static var file:Map = []; + public function new(net_Url: String, net_object: Object) { if (net_Url != null && net_object != null) { @@ -54,224 +52,65 @@ class Leenkx { } if (object != null) { - final loadEvent = Event.get(Leenkx.onLoadEvent); - final openEvent = Event.get(Leenkx.onOpenEvent); - final messageEvent = Event.get(Leenkx.onMessageEvent); - final errorEvent = Event.get(Leenkx.onErrorEvent); - final closeEvent = Event.get(Leenkx.onCloseEvent); - final seenEvent = Event.get(Leenkx.onSeenEvent); - final serverEvent = Event.get(Leenkx.onServerEvent); - final connectionsEvent = Event.get(Leenkx.onConnectionsEvent); - final pingEvent = Event.get(Leenkx.onPingEvent); - final leftEvent = Event.get(Leenkx.onLeftEvent); - final timeoutEvent = Event.get(Leenkx.onTimeOutEvent); - final rpcEvent = Event.get(Leenkx.onRpcEvent); - final rpcresponseEvent = Event.get(Leenkx.onRpcResponseEvent); - final wireleftEvent = Event.get(Leenkx.onWireLeftEvent); - final wireseenEvent = Event.get(Leenkx.onWireSeenEvent); - final torrentEvent = Event.get(Leenkx.onTorrentEvent); - final trackerEvent = Event.get(Leenkx.onTrackerEvent); - final announceEvent = Event.get(Leenkx.onAnnounceEvent); - final torrentDoneEvent = Event.get(Leenkx.onTorrentDoneEvent); + var uid = object.uid; Leenkx.connections[net_Url].onopen = function() { - if (openEvent != null) { - for (e in openEvent) { - if (e.mask == object.uid) { - e.onEvent(); - } - } - } + Event.send(Leenkx.onOpenEvent, uid); }; Leenkx.connections[net_Url].onmessage = function() { - if (messageEvent != null) { - for (e in messageEvent) { - if (e.mask == object.uid) { - e.onEvent(); - } - } - } + Event.send(Leenkx.onMessageEvent, uid); }; Leenkx.connections[net_Url].onerror = function() { - if (errorEvent != null) { - for (e in errorEvent) { - if (e.mask == object.uid) { - e.onEvent(); - } - } - } + Event.send(Leenkx.onErrorEvent, uid); }; Leenkx.connections[net_Url].onclose = function() { - if (closeEvent != null) { - for (e in closeEvent) { - if (e.mask == object.uid) { - e.onEvent(); - } - } - } + Event.send(Leenkx.onCloseEvent, uid); }; Leenkx.connections[net_Url].onseen = function() { - if (seenEvent != null) { - for (e in seenEvent) { - if (e.mask == object.uid) { - e.onEvent(); - } - } - } + Event.send(Leenkx.onSeenEvent, uid); }; Leenkx.connections[net_Url].onserver = function() { - if (serverEvent != null) { - for (e in serverEvent) { - if (e.mask == object.uid) { - e.onEvent(); - } - } - } + Event.send(Leenkx.onServerEvent, uid); }; Leenkx.connections[net_Url].onconnections = function() { - if (connectionsEvent != null) { - for (e in connectionsEvent) { - if (e.mask == object.uid) { - e.onEvent(); - } - } - } + Event.send(Leenkx.onConnectionsEvent, uid); }; Leenkx.connections[net_Url].onping = function() { - if (pingEvent != null) { - for (e in pingEvent) { - if (e.mask == object.uid) { - e.onEvent(); - } - } - } + Event.send(Leenkx.onPingEvent, uid); }; Leenkx.connections[net_Url].onleft = function() { - if (leftEvent != null) { - for (e in leftEvent) { - if (e.mask == object.uid) { - e.onEvent(); - } - } - } + Event.send(Leenkx.onLeftEvent, uid); }; Leenkx.connections[net_Url].ontimeout = function() { - if (timeoutEvent != null) { - for (e in timeoutEvent) { - if (e.mask == object.uid) { - e.onEvent(); - } - } - } - }; - Leenkx.connections[net_Url].onrpc = function() { - if (rpcEvent != null) { - for (e in rpcEvent) { - if (e.mask == object.uid) { - e.onEvent(); - } - } - } - }; - Leenkx.connections[net_Url].onrpcresponse = function() { - if (rpcresponseEvent != null) { - for (e in rpcresponseEvent) { - if (e.mask == object.uid) { - e.onEvent(); - } - } - } + Event.send(Leenkx.onTimeOutEvent, uid); }; Leenkx.connections[net_Url].onwireleft = function() { - if (wireleftEvent != null) { - for (e in wireleftEvent) { - if (e.mask == object.uid) { - e.onEvent(); - } - } - } + Event.send(Leenkx.onWireLeftEvent, uid); }; Leenkx.connections[net_Url].onwireseen = function() { - if (wireseenEvent != null) { - for (e in wireseenEvent) { - if (e.mask == object.uid) { - e.onEvent(); - } - } - } + Event.send(Leenkx.onWireSeenEvent, uid); }; Leenkx.connections[net_Url].ontorrent = function() { - if (torrentEvent != null) { - for (e in torrentEvent) { - if (e.mask == object.uid) { - e.onEvent(); - } - } - } - }; + Event.send(Leenkx.onTorrentEvent, uid); + }; + Leenkx.connections[net_Url].ontorrentadded = function() { + Event.send(Leenkx.onTorrentAddedEvent, uid); + }; Leenkx.connections[net_Url].ontracker = function() { - if (trackerEvent != null) { - for (e in trackerEvent) { - if (e.mask == object.uid) { - e.onEvent(); - } - } - } - }; + Event.send(Leenkx.onTrackerEvent, uid); + }; Leenkx.connections[net_Url].onannounce = function() { - if (announceEvent != null) { - for (e in announceEvent) { - if (e.mask == object.uid) { - e.onEvent(); - } - } - } + Event.send(Leenkx.onAnnounceEvent, uid); }; Leenkx.connections[net_Url].onload = function() { - if (loadEvent != null) { - for (e in loadEvent) { - if (e.mask == object.uid) { - e.onEvent(); - } - } - } - }; + Event.send(Leenkx.onLoadEvent, uid); + }; Leenkx.connections[net_Url].ontorrentdone = function() { - if (torrentDoneEvent != null) { - for (e in torrentDoneEvent) { - if (e.mask == object.uid) { - e.onEvent(); - } - } - } + Event.send(Leenkx.onTorrentDoneEvent, uid); }; - #if js - var lnxjs:Dynamic = js.Lib.global; var connectionUrl = net_Url; - - if (lnxjs.Leenkx != null) { - if (lnxjs.lnxNew == null) { - js.Syntax.code('globalThis.lnxNew = function(url) { return new Leenkx(url); }'); - } - Leenkx.connections[connectionUrl].onload(); - } else { - kha.Assets.loadBlobFromPath("Leenkx.js", function(b: kha.Blob) { - if (b != null) { - js.Syntax.code("(1,eval)({0})", b.toString()); - if (lnxjs.Leenkx != null && lnxjs.lnxNew == null) { - js.Syntax.code('globalThis.lnxNew = function(url) { return new Leenkx(url); }'); - } - } else { - trace("Warning: Leenkx.js blob is null - file may not be in assets"); - } - Leenkx.connections[connectionUrl].onload(); - }, function(err: kha.AssetError) { - trace("ERROR loading Leenkx.js: " + err.url + " - " + err.error); - Leenkx.connections[connectionUrl].onload(); - }); - } - #end + Leenkx.connections[connectionUrl].onload(); } } } @@ -281,11 +120,12 @@ class Leenkx { class LeenkxSocket { public var _url:String; - public var client:haxe.DynamicAccess; + public var client:LeenkxClient; + public var torrentClient:leenkx.network.torrent.TorrentClient; public var buffer:Dynamic; public var onopen:Void->Void; public var onclose:Void->Void; - public var onmessage:Void->Void; + public var onmessage:Void->Void; public var onerror:Void->Void; public var onseen:Void->Void; public var onserver:Void->Void; @@ -293,11 +133,10 @@ class LeenkxSocket { public var onping:Void->Void; public var onleft:Void->Void; public var ontimeout:Void->Void; - public var onrpc:Void->Void; - public var onrpcresponse:Void->Void; public var onwireleft:Void->Void; public var onwireseen:Void->Void; public var ontorrent:Void->Void; + public var ontorrentadded:Void->Void; public var ontracker:Void->Void; public var onannounce:Void->Void; public var onload:Void->Void; @@ -308,3 +147,4 @@ class LeenkxSocket { } } +#end diff --git a/leenkx/Sources/leenkx/network/Log.hx b/leenkx/Sources/leenkx/network/Log.hx deleted file mode 100644 index 44e1f016..00000000 --- a/leenkx/Sources/leenkx/network/Log.hx +++ /dev/null @@ -1,51 +0,0 @@ -package leenkx.network; - -class Log { - public static inline var INFO:Int = 0x000001; - public static inline var DEBUG:Int = 0x000010; - public static inline var DATA:Int = 0x000100; - - public static var mask:Int = 0; - - #if (sys || kha_krom) - public static var logFn:Dynamic->Void = function(data:Dynamic) { trace(data); }; - #elseif js - public static var logFn:Dynamic->Void = js.html.Console.log; - #end - - public static function info(data:String, id:String = null) { - if (mask & INFO != INFO) { - return; - } - - if (id != null) { - logFn('INFO :: ID-${id} :: ${data}'); - } else { - logFn('INFO :: ${data}'); - } - } - - public static function debug(data:String, id:String = null) { - if (mask & DEBUG != DEBUG) { - return; - } - - if (id != null) { - logFn('DEBUG :: ID-${id} :: ${data}'); - } else { - logFn('DEBUG :: ${data}'); - } - } - - public static function data(data:String, id:String = null) { - if (mask & DATA != DATA) { - return; - } - - if (id != null) { - logFn('DATA :: ID-${id}\n------------------------------\n${data}\n------------------------------'); - } else { - logFn('${data}'); - } - } -} diff --git a/leenkx/Sources/leenkx/network/WebSocket.hx b/leenkx/Sources/leenkx/network/WebSocket.hx index c3438f5e..acdf86b7 100644 --- a/leenkx/Sources/leenkx/network/WebSocket.hx +++ b/leenkx/Sources/leenkx/network/WebSocket.hx @@ -136,18 +136,43 @@ import haxe.io.Bytes; class WebSocket { private var _url:String; - + private var _protocols:Array = null; private var _ws:js.html.WebSocket = null; - public function new(url:String, immediateOpen=true) { + public var state(get, null): State; + function get_state() { + if (_ws == null) return Closed; + switch (_ws.readyState) { + case js.html.WebSocket.CONNECTING: + return Handshake; + case js.html.WebSocket.OPEN: + return Body; + } + + return Closed; + } + + public var protocol(get, null): String; + function get_protocol() { + if (_ws == null) + return null; + return _ws.protocol; + } + + public function new(url:String, immediateOpen=true, protocols:Array = null) { _url = url; + _protocols = protocols; if (immediateOpen) { open(); } } private function createSocket() { - return new js.html.WebSocket(_url); + if (_protocols == null) { + return new js.html.WebSocket(_url); + } + + return new js.html.WebSocket(_url, _protocols); } public function open() { @@ -292,6 +317,10 @@ class WebSocket extends WebSocketCommon { public var _host:String; public var _port:Int = 0; public var _path:String; + public var _search:String; + + private var _protocols:Array; + public var protocol(default, null):String = null; private var _processThread:Thread; private var _encodedKey:String = "wskey"; @@ -300,8 +329,9 @@ class WebSocket extends WebSocketCommon { public var additionalHeaders(get, null):Map; - public function new(url:String, immediateOpen=true) { + public function new(url:String, immediateOpen=true, protocols:Array = null) { parseUrl(url); + _protocols = protocols; super(createSocket()); @@ -310,40 +340,42 @@ class WebSocket extends WebSocketCommon { } } - inline private function parseUrl(url) + inline private function parseUrl(url:String) { - /** TO DO - FIND OUT WHAT IS BREAKING REGEX IN THE NEW PCRE2 FOR HL C - - var urlRegExp = ~/^(\w+?):\/\/([\w\.-]+)(:(\d+))?(\/.*)?$/; - - if ( ! urlRegExp.match(url)) { + var sep = url.indexOf("://"); + if (sep < 0) { throw 'Uri not matching websocket URL "${url}"'; } + _protocol = url.substr(0, sep); + var rest = url.substr(sep + 3); - _protocol = urlRegExp.matched(1); + var pathStart = rest.length; + var slashIdx = rest.indexOf("/"); + var qIdx = rest.indexOf("?"); + if (slashIdx >= 0 && slashIdx < pathStart) pathStart = slashIdx; + if (qIdx >= 0 && qIdx < pathStart) pathStart = qIdx; - _host = urlRegExp.matched(2); + var hostPort = rest.substr(0, pathStart); + var pathAndSearch = rest.substr(pathStart); - var parsedPort = Std.parseInt(urlRegExp.matched(4)); - if (parsedPort > 0 ) { - _port = parsedPort; + var colonIdx = hostPort.indexOf(":"); + if (colonIdx >= 0) { + _host = hostPort.substr(0, colonIdx); + var parsedPort = Std.parseInt(hostPort.substr(colonIdx + 1)); + if (parsedPort > 0) { + _port = parsedPort; + } + } else { + _host = hostPort; } - _path = urlRegExp.matched(5); - if (_path == null || _path.length == 0) { - _path = "/"; + + _path = pathAndSearch; + _search = null; + var searchIdx = pathAndSearch.indexOf("?"); + if (searchIdx >= 0) { + _path = pathAndSearch.substr(0, searchIdx); + _search = pathAndSearch.substr(searchIdx); } - **/ - var urlArr = url.split(":"); - if ( urlArr.length < 3) { - throw 'Uri not matching websocket URL "${url}"'; - } - _protocol = urlArr[0]; - _host = urlArr[1].substr(2, urlArr[1].length); - var parsedPort = Std.parseInt(urlArr[2].split("/")[0]); - if (parsedPort > 0 ) { - _port = parsedPort; - } - _path = urlArr[2].substr(urlArr[2].split("/")[0].length, urlArr[2].length); if (_path == null || _path.length == 0) { _path = "/"; } @@ -388,9 +420,7 @@ class WebSocket extends WebSocketCommon { #else haxe.MainLoop.addThread(function() { - Log.debug("Thread started", this.id); processLoop(this); - Log.debug("Thread ended", this.id); }); #end @@ -399,16 +429,13 @@ class WebSocket extends WebSocketCommon { } private function processThread() { - Log.debug("Thread started", this.id); var ws:WebSocket = Thread.readMessage(true); processLoop(this); - Log.debug("Thread ended", this.id); } private function processLoop(ws:WebSocket) { while (ws.state != State.Closed) { // TODO: should think about mutex ws.process(); - Sys.sleep(.01); } } @@ -422,8 +449,7 @@ class WebSocket extends WebSocketCommon { public function sendHandshake() { var httpRequest = new HttpRequest(); httpRequest.method = "GET"; - // TODO: should propably be hostname+port+path? - httpRequest.uri = _path; + httpRequest.uri = (_search != null) ? _path + _search : _path; httpRequest.httpVersion = "HTTP/1.1"; httpRequest.headers.set(HttpHeader.HOST, _host + ":" + _port); @@ -435,6 +461,9 @@ class WebSocket extends WebSocketCommon { httpRequest.headers.set(HttpHeader.CACHE_CONTROL, "no-cache"); httpRequest.headers.set(HttpHeader.ORIGIN, _socket.host().host.toString() + ":" + _socket.host().port); + if (_protocols != null) { + httpRequest.headers.set(HttpHeader.SEC_WEBSOCKET_PROTOCOL, _protocols.join(', ')); + } _encodedKey = generateWSKey(); httpRequest.headers.set(HttpHeader.SEC_WEBSOCKET_KEY, _encodedKey); @@ -487,16 +516,18 @@ class WebSocket extends WebSocketCommon { } } + var protocol = httpResponse.headers.get(HttpHeader.SEC_WEBSOCKET_PROTOCOL); + if (protocol != null) { + this.protocol = protocol; + } + _onopenCalled = false; state = State.Head; } private function generateWSKey():String { - var b = Bytes.alloc(16); - for (i in 0...16) { - b.set(i, Std.random(255)); - } - return Base64.encode(b); + return haxe.crypto.Base64.encode( + leenkx.network.torrent.Crypto.randomBytes(16)); } } diff --git a/leenkx/Sources/leenkx/network/WebSocketCommon.hx b/leenkx/Sources/leenkx/network/WebSocketCommon.hx index 6d33c399..69dbe792 100644 --- a/leenkx/Sources/leenkx/network/WebSocketCommon.hx +++ b/leenkx/Sources/leenkx/network/WebSocketCommon.hx @@ -38,7 +38,6 @@ class WebSocketCommon { public function send(data:Any) { if (Std.isOfType(data, String)) { - Log.data(data, id); sendFrame(Utf8Encoder.encode(data), OpCode.Text); } else if (Std.isOfType(data, Bytes)) { sendFrame(data, OpCode.Binary); @@ -118,7 +117,6 @@ class WebSocketCommon { } } else { var stringPayload = Utf8Encoder.decode(unmaskedMessageData); - Log.data(stringPayload, id); if (this.onmessage != null) { this.onmessage(StrMessage(stringPayload)); } @@ -138,14 +136,12 @@ class WebSocketCommon { case State.Closed: close(); case _: - trace('State not impl: ${state}'); } } public function close() { if (state != State.Closed) { try { - Log.debug("Closed", id); sendFrame(Bytes.alloc(0), OpCode.Close); state = State.Closed; _socket.close(); @@ -162,7 +158,6 @@ class WebSocketCommon { _socket.output.write(data); _socket.output.flush(); } catch (e:Dynamic) { - Log.debug(Std.string(e), id); if (onerror != null) { onerror(Std.string(e)); } @@ -195,17 +190,24 @@ class WebSocketCommon { } private static function generateMask() { - var maskData = Bytes.alloc(4); - maskData.set(0, Std.random(256)); - maskData.set(1, Std.random(256)); - maskData.set(2, Std.random(256)); - maskData.set(3, Std.random(256)); - return maskData; + return leenkx.network.torrent.Crypto.randomBytes(4); } private static function applyMask(payload:Bytes, mask:Bytes) { var maskedPayload = Bytes.alloc(payload.length); - for (n in 0 ... payload.length) maskedPayload.set(n, payload.get(n) ^ mask.get(n % mask.length)); + var m0 = mask.get(0), m1 = mask.get(1), m2 = mask.get(2), m3 = mask.get(3); + var i = 0; + while (i + 3 < payload.length) { + maskedPayload.set(i, payload.get(i) ^ m0); + maskedPayload.set(i + 1, payload.get(i + 1) ^ m1); + maskedPayload.set(i + 2, payload.get(i + 2) ^ m2); + maskedPayload.set(i + 3, payload.get(i + 3) ^ m3); + i += 4; + } + while (i < payload.length) { + maskedPayload.set(i, payload.get(i) ^ mask.get(i & 3)); + i++; + } return maskedPayload; } @@ -230,7 +232,6 @@ class WebSocketCommon { try { result = SocketImpl.select([_socket], null, null, 0.01); } catch (e:Dynamic) { - Log.debug("Error selecting socket: " + e); needClose = true; } @@ -238,12 +239,11 @@ class WebSocketCommon { if (result.read.length > 0) { try { while (true) { - var data = Bytes.alloc(1024); + var data = Bytes.alloc(65536); var read = _socket.input.readBytes(data, 0, data.length); if (read <= 0){ break; } - Log.debug("Bytes read: " + read, id); _buffer.writeBytes(data.sub(0, read)); } } catch (e:Dynamic) { @@ -259,7 +259,7 @@ class WebSocketCommon { } #else - needClose = !(e == 'Blocking' || (Std.isOfType(e, Error) && (e:Error).match(Error.Blocked))); + needClose = !(e == 'Blocking' || (Std.isOfType(e, Error) && (e:Error).match(Error.Blocked)) || (state == Handshake && e == haxe.io.Eof)); #end } @@ -273,7 +273,6 @@ class WebSocketCommon { if (needClose == true) { // dont want to send the Close frame here if (state != State.Closed) { try { - Log.debug("Closed", id); state = State.Closed; _socket.close(); } catch (e:Dynamic) { } @@ -288,8 +287,6 @@ class WebSocketCommon { public function sendHttpRequest(httpRequest:HttpRequest) { var data = httpRequest.build(); - Log.data(data, id); - try { _socket.output.write(Bytes.ofString(data)); _socket.output.flush(); @@ -304,8 +301,6 @@ class WebSocketCommon { public function sendHttpResponse(httpResponse:HttpResponse) { var data = httpResponse.build(); - Log.data(data, id); - _socket.output.write(Bytes.ofString(data)); _socket.output.flush(); } @@ -325,8 +320,6 @@ class WebSocketCommon { httpRequest.addLine(line); } - Log.data(httpRequest.toString(), id); - return httpRequest; } @@ -346,8 +339,6 @@ class WebSocketCommon { } - Log.data(httpResponse.toString(), id); - return httpResponse; } diff --git a/leenkx/Sources/leenkx/network/WebSocketHandler.hx b/leenkx/Sources/leenkx/network/WebSocketHandler.hx index 11e8f2f0..17b5d9e2 100644 --- a/leenkx/Sources/leenkx/network/WebSocketHandler.hx +++ b/leenkx/Sources/leenkx/network/WebSocketHandler.hx @@ -13,7 +13,6 @@ class WebSocketHandler extends Handler { _creationTime = Sys.time(); #end _socket.setBlocking(false); - Log.debug('New socket handler', id); } public override function handle() { @@ -23,7 +22,6 @@ class WebSocketHandler extends Handler { var currentTime = Sys.time(); #end if (this.state == State.Handshake && currentTime - _creationTime > (MAX_WAIT_TIME / 1000)) { - Log.info('No handshake detected in ${MAX_WAIT_TIME}ms, closing connection', id); this.close(); return; } @@ -70,10 +68,8 @@ class WebSocketHandler extends Handler { httpResponse.headers.set(HttpHeader.CONNECTION, "close"); httpResponse.headers.set(HttpHeader.X_WEBSOCKET_REJECT_REASON, 'Unsupported connection header: ${httpRequest.headers.get(HttpHeader.CONNECTION)}.'); } else { - Log.debug('Handshaking', id); var key = httpRequest.headers.get(HttpHeader.SEC_WEBSOCKET_KEY); var result = makeWSKeyResponse(key); - Log.debug('Handshaking key - ${result}', id); httpResponse.code = 101; httpResponse.text = "Switching Protocols"; @@ -82,14 +78,21 @@ class WebSocketHandler extends Handler { httpResponse.headers.set(HttpHeader.SEC_WEBSOSCKET_ACCEPT, result); } - sendHttpResponse(httpResponse); + function callback(httpResponse: HttpResponse) { + sendHttpResponse(httpResponse); - if (httpResponse.code == 101) { - _onopenCalled = false; - state = State.Head; - Log.debug('Connected', id); + if (httpResponse.code == 101) { + _onopenCalled = false; + state = State.Head; + } else { + close(); + } + } + + if (validateHandshake != null) { + validateHandshake(httpRequest, httpResponse, callback); } else { - close(); + callback(httpResponse); } } } diff --git a/leenkx/Sources/leenkx/network/WebSocketServer.hx b/leenkx/Sources/leenkx/network/WebSocketServer.hx index 11976eeb..ff968b15 100644 --- a/leenkx/Sources/leenkx/network/WebSocketServer.hx +++ b/leenkx/Sources/leenkx/network/WebSocketServer.hx @@ -60,7 +60,6 @@ class WebSocketServer _serverSocket.bind(new sys.net.Host(_host), _port); #end _serverSocket.listen(_maxConnections); - Log.info('Starting server - ${_host}:${_port} (maxConnections: ${_maxConnections})'); #if kha_krom kha.Scheduler.addTimeTask(function() { @@ -114,7 +113,6 @@ class WebSocketServer var handler = new T(socket); handlers.push(handler); - Log.debug("Adding to web server handler to list - total: " + handlers.length, handler.id); if (onClientAdded != null) { onClientAdded(handler); } @@ -174,7 +172,6 @@ class WebSocketServer for (h in toRemove) { handlers.remove(h); - Log.debug("Removing web server handler from list - total: " + handlers.length, h.id); if (onClientRemoved != null) { onClientRemoved(h); } @@ -186,7 +183,6 @@ class WebSocketServer while (_handlersClosed.length > 0) { var h = _handlersClosed.shift(); handlers.remove(h); - Log.debug("Removing web server handler from list - total: " + handlers.length, h.id); if (onClientRemoved != null) { onClientRemoved(h); } diff --git a/leenkx/Sources/leenkx/network/krom/KromSecureSocket.hx b/leenkx/Sources/leenkx/network/krom/KromSecureSocket.hx index 62b09a6a..523ce592 100644 --- a/leenkx/Sources/leenkx/network/krom/KromSecureSocket.hx +++ b/leenkx/Sources/leenkx/network/krom/KromSecureSocket.hx @@ -52,8 +52,15 @@ class KromSecureSocket extends KromSocket { } override public function connect(host:KromHost, port:Int):Void { - super.connect(host, port); - enableSSL(); + setBlocking(true); + try { + super.connect(host, port); + enableSSL(); + setBlocking(false); + } catch (e:Dynamic) { + setBlocking(false); + trace('Failed Krom Host'); + } } public function isSSLEnabled():Bool { diff --git a/leenkx/Sources/leenkx/network/krom/KromSocket.hx b/leenkx/Sources/leenkx/network/krom/KromSocket.hx index e7aeada4..eb1781b8 100644 --- a/leenkx/Sources/leenkx/network/krom/KromSocket.hx +++ b/leenkx/Sources/leenkx/network/krom/KromSocket.hx @@ -43,6 +43,9 @@ class KromSocket { } private function setSocketId(id:Int):Void { + if (_socketId >= 0) { + krom_socket_close(_socketId); + } _socketId = id; if (_socketId >= 0) { input = new KromSocketInput(this); @@ -247,3 +250,29 @@ class KromHost { return host; } } +// TODO: Replace all ticks timers and intervals for events +class KromPump { + static var ticks:Array Void> = []; + static var timer:haxe.Timer = null; + + public static function add(tick:Void -> Void):Void { + if (!ticks.contains(tick)) ticks.push(tick); + if (timer == null) { + timer = new haxe.Timer(10); + timer.run = runAll; + } + } + public static function remove(tick:Void -> Void):Void { + ticks.remove(tick); + if (ticks.length == 0 && timer != null) { + timer.stop(); + timer = null; + } + } + static function runAll():Void { + var snapshot = ticks.copy(); + for (tick in snapshot) { + if (ticks.contains(tick)) tick(); + } + } +} diff --git a/leenkx/Sources/leenkx/network/krom/KromUdpSocket.hx b/leenkx/Sources/leenkx/network/krom/KromUdpSocket.hx new file mode 100644 index 00000000..79505b04 --- /dev/null +++ b/leenkx/Sources/leenkx/network/krom/KromUdpSocket.hx @@ -0,0 +1,70 @@ +package leenkx.network.krom; + +import haxe.io.Bytes; + +@:native("krom_udp_create") extern function krom_udp_create():Int; +@:native("krom_udp_bind") extern function krom_udp_bind(id:Int, + addr:String, port:Int):Bool; +@:native("krom_udp_sendto") extern function krom_udp_sendto(id:Int, + host:String, port:Int, data:js.lib.ArrayBuffer):Int; +@:native("krom_udp_recvfrom") extern function krom_udp_recvfrom(id:Int, + maxLen:Int):Dynamic; +@:native("krom_socket_close") extern function krom_udp_close(id:Int):Void; +@:native("krom_socket_set_blocking") +extern function krom_udp_set_blocking(id:Int, blocking:Bool):Void; + +class KromUdpSocket { + private var _socketId:Int = -1; + + public function new() { + _socketId = krom_udp_create(); + setBlocking(false); + } + + public function bind(host:KromSocket.KromHost, port:Int):Void { + if (_socketId >= 0) { + if (!krom_udp_bind(_socketId, host.host, port)) { + trace("UDP bind failed to " + host.host + ":" + port); + } + } + } + + public function sendTo(data:Bytes, pos:Int, len:Int, + host:KromSocket.KromHost, port:Int):Int { + if (_socketId < 0) return 0; + var arrayBuffer:js.lib.ArrayBuffer = data.getData(); + var exact:js.lib.ArrayBuffer = + untyped arrayBuffer.slice(pos, pos + len); + return krom_udp_sendto(_socketId, host.host, port, exact); + } + + public function recvFrom(maxLength:Int):{data:Bytes, host:String, + port:Int} { + if (_socketId < 0) return null; + var result:Dynamic = krom_udp_recvfrom(_socketId, maxLength); + if (result == -1) { + _socketId = -1; + trace("Closed"); + return null; + } + if (result == null || Std.isOfType(result, Int)) return null; + return { + data: Bytes.ofData(result.data), + host: result.host, + port: result.port + }; + } + + public function setBlocking(blocking:Bool):Void { + if (_socketId >= 0) { + krom_udp_set_blocking(_socketId, blocking); + } + } + + public function close():Void { + if (_socketId >= 0) { + krom_udp_close(_socketId); + _socketId = -1; + } + } +}