hxWebsocket Upstream and Network Updates

This commit is contained in:
2026-10-03 22:05:53 -07:00
parent da38631584
commit 860b73a8f4
12 changed files with 284 additions and 345 deletions

View File

@ -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;
}

View File

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

View File

@ -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";

View File

@ -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,23 +22,19 @@ 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<String, leenkx.network.LeenkxSocket> = [];
#if js
public static var peers:js.lib.Map<String, String> = new js.lib.Map<String,String>();
public static var data:js.lib.Map<String, Dynamic> = new js.lib.Map<String,Dynamic>();
public static var id:js.lib.Map<String, String> = new js.lib.Map<String,String>();
public static var torrent:js.lib.Map<String, Dynamic> = new js.lib.Map<String,Dynamic>();
public static var file:js.lib.Map<String, Dynamic> = new js.lib.Map<String,Dynamic>();
#end
public static var lxNew:Void->Void;
public static var peers:Map<String, String> = [];
public static var data:Map<String, Dynamic> = [];
public static var id:Map<String, String> = [];
public static var torrent:Map<String, Dynamic> = [];
public static var file:Map<String, Dynamic> = [];
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
}
}
}
@ -281,7 +120,8 @@ class Leenkx {
class LeenkxSocket {
public var _url:String;
public var client:haxe.DynamicAccess<Dynamic>;
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;
@ -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

View File

@ -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}');
}
}
}

View File

@ -136,20 +136,45 @@ import haxe.io.Bytes;
class WebSocket {
private var _url:String;
private var _protocols:Array<String> = 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<String> = null) {
_url = url;
_protocols = protocols;
if (immediateOpen) {
open();
}
}
private function createSocket() {
if (_protocols == null) {
return new js.html.WebSocket(_url);
}
return new js.html.WebSocket(_url, _protocols);
}
public function open() {
if (_ws != null) {
throw "Socket already connected";
@ -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<String>;
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<String, String>;
public function new(url:String, immediateOpen=true) {
public function new(url:String, immediateOpen=true, protocols:Array<String> = 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 ) {
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;
}
_path = urlRegExp.matched(5);
if (_path == null || _path.length == 0) {
_path = "/";
} else {
_host = hostPort;
}
**/
var urlArr = url.split(":");
if ( urlArr.length < 3) {
throw 'Uri not matching websocket URL "${url}"';
_path = pathAndSearch;
_search = null;
var searchIdx = pathAndSearch.indexOf("?");
if (searchIdx >= 0) {
_path = pathAndSearch.substr(0, searchIdx);
_search = pathAndSearch.substr(searchIdx);
}
_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));
}
}

View File

@ -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;
}

View File

@ -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);
}
function callback(httpResponse: HttpResponse) {
sendHttpResponse(httpResponse);
if (httpResponse.code == 101) {
_onopenCalled = false;
state = State.Head;
Log.debug('Connected', id);
} else {
close();
}
}
if (validateHandshake != null) {
validateHandshake(httpRequest, httpResponse, callback);
} else {
callback(httpResponse);
}
}
}

View File

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

View File

@ -52,8 +52,15 @@ class KromSecureSocket extends KromSocket {
}
override public function connect(host:KromHost, port:Int):Void {
setBlocking(true);
try {
super.connect(host, port);
enableSSL();
setBlocking(false);
} catch (e:Dynamic) {
setBlocking(false);
trace('Failed Krom Host');
}
}
public function isSSLEnabled():Bool {

View File

@ -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 -> 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();
}
}
}

View File

@ -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;
}
}
}