Phase 5a: persistent websocket with request ids and reconnect
continuous-integration/drone/push Build encountered an error
continuous-integration/drone/push Build encountered an error
Client (ts/network_ts.ts): one WebSocket per page. Commands queue and go out one at a time with a client-chosen id; the response is matched on the echoed id (or to the in-flight request for older servers). 12 s timeout per request, exponential-backoff reconnect (0.5 s to 10 s), the in-flight request is retried after a reconnect, and log polls are de-duplicated so they cannot pile up behind a stalled connection. The old global flag that silently dropped any request made while another was in flight is gone. Server (src/webserver.go): the /data/ handler now serves any number of commands on one connection (it used to break out of its loop after the first reply without closing the socket, leaving it open and deaf; the old client papered over that by opening a new socket per request). Connection closed on exit, request id echoed in the response. Test: TestWSServesMultipleCommandsPerConnection.
This commit is contained in:
+194
-95
@@ -1,97 +1,197 @@
|
||||
// Websocket transport between the web UI and xTeVe.
|
||||
//
|
||||
// One socket stays open for the life of the page. Requests are queued and
|
||||
// sent one at a time; each carries an id that the server echoes back, and
|
||||
// the response is matched on it (or, for older servers that close after a
|
||||
// single reply, to the oldest in-flight request). Connection loss triggers
|
||||
// a reconnect with exponential backoff, and queued requests survive it.
|
||||
class WSClient {
|
||||
constructor() {
|
||||
this.socket = null;
|
||||
this.queue = [];
|
||||
this.inFlight = null;
|
||||
this.nextID = 1;
|
||||
this.reconnectDelay = 500;
|
||||
this.reconnectTimer = 0;
|
||||
this.closedByUs = false;
|
||||
this.requestTimeoutMs = 12000;
|
||||
}
|
||||
url() {
|
||||
var protocol = window.location.protocol == "https:" ? "wss://" : "ws://";
|
||||
var host = window.location.host;
|
||||
if (host == undefined || host.length < 1) {
|
||||
host = window.location.hostname;
|
||||
}
|
||||
return protocol + host + "/data/";
|
||||
}
|
||||
// send queues a command and resolves with the parsed response.
|
||||
send(payload) {
|
||||
var client = this;
|
||||
return new Promise(function (resolve, reject) {
|
||||
var req = {
|
||||
id: String(client.nextID++),
|
||||
payload: payload,
|
||||
resolve: resolve,
|
||||
reject: reject,
|
||||
timer: 0,
|
||||
sent: false
|
||||
};
|
||||
client.queue.push(req);
|
||||
client.pump();
|
||||
});
|
||||
}
|
||||
// hasQueued reports whether a command of this kind is already waiting.
|
||||
hasQueued(cmd) {
|
||||
if (this.inFlight != null && this.inFlight.payload["cmd"] == cmd) {
|
||||
return true;
|
||||
}
|
||||
for (var i = 0; i < this.queue.length; i++) {
|
||||
if (this.queue[i].payload["cmd"] == cmd) {
|
||||
return true;
|
||||
}
|
||||
}
|
||||
return false;
|
||||
}
|
||||
connect() {
|
||||
if (this.socket != null && (this.socket.readyState == WebSocket.OPEN || this.socket.readyState == WebSocket.CONNECTING)) {
|
||||
return;
|
||||
}
|
||||
var client = this;
|
||||
var socket = new WebSocket(this.url());
|
||||
this.socket = socket;
|
||||
socket.onopen = function () {
|
||||
WS_AVAILABLE = true;
|
||||
client.reconnectDelay = 500;
|
||||
client.pump();
|
||||
};
|
||||
socket.onmessage = function (e) {
|
||||
client.handleMessage(e.data);
|
||||
};
|
||||
socket.onerror = function () {
|
||||
// onclose follows; nothing to do here.
|
||||
};
|
||||
socket.onclose = function () {
|
||||
if (client.socket == socket) {
|
||||
client.socket = null;
|
||||
}
|
||||
// An in-flight request that never got its answer goes back to the
|
||||
// front of the queue and is retried on the next connection.
|
||||
if (client.inFlight != null) {
|
||||
var req = client.inFlight;
|
||||
client.inFlight = null;
|
||||
req.sent = false;
|
||||
client.queue.unshift(req);
|
||||
}
|
||||
client.scheduleReconnect();
|
||||
};
|
||||
}
|
||||
scheduleReconnect() {
|
||||
if (this.closedByUs || this.reconnectTimer != 0) {
|
||||
return;
|
||||
}
|
||||
if (this.queue.length == 0) {
|
||||
// Idle: reconnect lazily on the next request.
|
||||
return;
|
||||
}
|
||||
var client = this;
|
||||
this.reconnectTimer = window.setTimeout(function () {
|
||||
client.reconnectTimer = 0;
|
||||
client.connect();
|
||||
}, this.reconnectDelay);
|
||||
this.reconnectDelay = Math.min(this.reconnectDelay * 2, 10000);
|
||||
}
|
||||
// pump sends the next queued request if the socket is open and idle.
|
||||
pump() {
|
||||
if (this.inFlight != null || this.queue.length == 0) {
|
||||
return;
|
||||
}
|
||||
if (this.socket == null || this.socket.readyState != WebSocket.OPEN) {
|
||||
this.connect();
|
||||
return;
|
||||
}
|
||||
var req = this.queue.shift();
|
||||
this.inFlight = req;
|
||||
req.sent = true;
|
||||
req.payload["id"] = req.id;
|
||||
var client = this;
|
||||
req.timer = window.setTimeout(function () {
|
||||
if (client.inFlight == req) {
|
||||
client.inFlight = null;
|
||||
req.reject(new Error("timeout"));
|
||||
// A socket that stopped answering is not trusted any more.
|
||||
client.resetSocket();
|
||||
client.pump();
|
||||
}
|
||||
}, this.requestTimeoutMs);
|
||||
try {
|
||||
this.socket.send(JSON.stringify(req.payload));
|
||||
}
|
||||
catch (err) {
|
||||
window.clearTimeout(req.timer);
|
||||
this.inFlight = null;
|
||||
req.sent = false;
|
||||
this.queue.unshift(req);
|
||||
this.resetSocket();
|
||||
}
|
||||
}
|
||||
resetSocket() {
|
||||
var socket = this.socket;
|
||||
this.socket = null;
|
||||
if (socket != null) {
|
||||
socket.onclose = null;
|
||||
try {
|
||||
socket.close();
|
||||
}
|
||||
catch (err) { /* ignore */ }
|
||||
}
|
||||
this.scheduleReconnect();
|
||||
}
|
||||
handleMessage(raw) {
|
||||
var response;
|
||||
try {
|
||||
response = JSON.parse(raw);
|
||||
}
|
||||
catch (err) {
|
||||
return;
|
||||
}
|
||||
var req = this.inFlight;
|
||||
if (req == null) {
|
||||
return;
|
||||
}
|
||||
// Servers that echo the id let us drop stray late answers.
|
||||
if (response.hasOwnProperty("id") && response["id"] != req.id) {
|
||||
return;
|
||||
}
|
||||
window.clearTimeout(req.timer);
|
||||
this.inFlight = null;
|
||||
req.resolve(response);
|
||||
this.pump();
|
||||
}
|
||||
}
|
||||
var WS = new WSClient();
|
||||
class Server {
|
||||
constructor(cmd) {
|
||||
this.cmd = cmd;
|
||||
}
|
||||
request(data) {
|
||||
if (SERVER_CONNECTION == true) {
|
||||
var isLogUpdate = this.cmd == "updateLog";
|
||||
// Log polling must never pile up behind a stalled connection.
|
||||
if (isLogUpdate && WS.hasQueued("updateLog")) {
|
||||
return;
|
||||
}
|
||||
SERVER_CONNECTION = true;
|
||||
if (this.cmd != "updateLog") {
|
||||
if (!isLogUpdate) {
|
||||
showElement("loading", true);
|
||||
UNDO = new Object();
|
||||
setConnectionState("busy");
|
||||
}
|
||||
switch (window.location.protocol) {
|
||||
case "http:":
|
||||
this.protocol = "ws://";
|
||||
break;
|
||||
case "https:":
|
||||
this.protocol = "wss://";
|
||||
break;
|
||||
}
|
||||
var wsHost = window.location.host;
|
||||
if (wsHost == undefined || wsHost.length < 1) {
|
||||
wsHost = window.location.hostname;
|
||||
}
|
||||
// The session cookie (HttpOnly) is sent with the websocket handshake.
|
||||
var url = this.protocol + wsHost + "/data/";
|
||||
data["cmd"] = this.cmd;
|
||||
var requestCmd = data["cmd"];
|
||||
var ws = new WebSocket(url);
|
||||
var isLogUpdate = data["cmd"] == "updateLog";
|
||||
var responseReceived = false;
|
||||
var requestFinished = false;
|
||||
var timeoutMs = 12000;
|
||||
var requestTimeout;
|
||||
var finishRequest = function (state, responseSuccess = false) {
|
||||
if (requestFinished == true) {
|
||||
return;
|
||||
}
|
||||
requestFinished = true;
|
||||
SERVER_CONNECTION = false;
|
||||
window.clearTimeout(requestTimeout);
|
||||
if (responseSuccess == true) {
|
||||
if (state == "online") {
|
||||
WS_FAILURE_COUNT = 0;
|
||||
}
|
||||
}
|
||||
else {
|
||||
WS_FAILURE_COUNT++;
|
||||
}
|
||||
if (isLogUpdate == false) {
|
||||
var requestCmd = this.cmd;
|
||||
WS.send(data).then(function (response) {
|
||||
WS_FAILURE_COUNT = 0;
|
||||
if (!isLogUpdate) {
|
||||
showElement("loading", false);
|
||||
}
|
||||
if (state != "") {
|
||||
setConnectionState(state);
|
||||
}
|
||||
};
|
||||
requestTimeout = window.setTimeout(function () {
|
||||
console.log("Websocket request timed out.");
|
||||
var timeoutState = "offline";
|
||||
if (isLogUpdate == true && WS_FAILURE_COUNT < 2) {
|
||||
timeoutState = "idle";
|
||||
}
|
||||
finishRequest(timeoutState, false);
|
||||
try {
|
||||
ws.close();
|
||||
}
|
||||
catch (err) {
|
||||
console.log(err);
|
||||
}
|
||||
}, timeoutMs);
|
||||
ws.onopen = function () {
|
||||
WS_AVAILABLE = true;
|
||||
if (data["cmd"] != "updateLog") {
|
||||
setConnectionState("busy");
|
||||
}
|
||||
this.send(JSON.stringify(data));
|
||||
};
|
||||
ws.onerror = function (e) {
|
||||
console.log("No websocket connection to xTeVe could be established. Check your network configuration.");
|
||||
var errorState = "offline";
|
||||
if (isLogUpdate == true && WS_FAILURE_COUNT < 2) {
|
||||
errorState = "idle";
|
||||
}
|
||||
finishRequest(errorState, false);
|
||||
if (WS_AVAILABLE == false && isLogUpdate == false && requestCmd != "getServerConfig") {
|
||||
alert("No websocket connection to xTeVe could be established. Check your network configuration.");
|
||||
}
|
||||
};
|
||||
ws.onmessage = function (e) {
|
||||
responseReceived = true;
|
||||
finishRequest("online", true);
|
||||
var response = JSON.parse(e.data);
|
||||
setConnectionState("online");
|
||||
if (response["status"] == false) {
|
||||
setConnectionState("offline");
|
||||
alert(response["err"]);
|
||||
@@ -106,16 +206,14 @@ class Server {
|
||||
div.className = "changed";
|
||||
return;
|
||||
}
|
||||
switch (data["cmd"]) {
|
||||
switch (requestCmd) {
|
||||
case "updateLog":
|
||||
SERVER["log"] = response["log"];
|
||||
if (document.getElementById("content_log")) {
|
||||
showLogs(false);
|
||||
}
|
||||
return;
|
||||
break;
|
||||
default:
|
||||
SERVER = new Object();
|
||||
SERVER = response;
|
||||
break;
|
||||
}
|
||||
@@ -139,17 +237,20 @@ class Server {
|
||||
return;
|
||||
}
|
||||
createLayout();
|
||||
};
|
||||
ws.onclose = function () {
|
||||
if (responseReceived == true) {
|
||||
return;
|
||||
}, function (err) {
|
||||
WS_FAILURE_COUNT++;
|
||||
if (!isLogUpdate) {
|
||||
showElement("loading", false);
|
||||
}
|
||||
var closeState = "offline";
|
||||
if (isLogUpdate == true && WS_FAILURE_COUNT < 2) {
|
||||
closeState = "idle";
|
||||
var state = "offline";
|
||||
if (isLogUpdate && WS_FAILURE_COUNT < 2) {
|
||||
state = "idle";
|
||||
}
|
||||
finishRequest(closeState, false);
|
||||
};
|
||||
setConnectionState(state);
|
||||
if (WS_AVAILABLE == false && !isLogUpdate && requestCmd != "getServerConfig") {
|
||||
alert("No websocket connection to xTeVe could be established. Check your network configuration.");
|
||||
}
|
||||
});
|
||||
}
|
||||
}
|
||||
var WS_FAILURE_COUNT = 0;
|
||||
@@ -961,7 +1062,7 @@ function PageReady() {
|
||||
window.clearInterval(bootstrapTimer);
|
||||
return;
|
||||
}
|
||||
if (SERVER_CONNECTION == true) {
|
||||
if (WS.hasQueued("getServerConfig")) {
|
||||
return;
|
||||
}
|
||||
bootstrapAttempts++;
|
||||
@@ -2688,7 +2789,6 @@ var BULK_EDIT = false;
|
||||
var COLUMN_TO_SORT;
|
||||
var SEARCH_MAPPING = new Object();
|
||||
var UNDO = new Object();
|
||||
var SERVER_CONNECTION = false;
|
||||
var WS_AVAILABLE = false;
|
||||
var ACTIVE_MENU_ID = "";
|
||||
var LAST_BULK_CHECKBOX = null;
|
||||
@@ -3236,7 +3336,6 @@ function sortSelect(elem) {
|
||||
return;
|
||||
}
|
||||
function updateLog() {
|
||||
console.log("TOKEN");
|
||||
var server = new Server("updateLog");
|
||||
server.request(new Object());
|
||||
}
|
||||
|
||||
@@ -5,6 +5,10 @@ type RequestStruct struct {
|
||||
// Befehle an xTeVe
|
||||
Cmd string `json:"cmd"`
|
||||
|
||||
// Client-chosen request id, echoed in the response so the web UI can
|
||||
// match answers on a long-lived connection.
|
||||
ID string `json:"id,omitempty"`
|
||||
|
||||
// Benutzer
|
||||
DeleteUser bool `json:"deleteUser,omitempty"`
|
||||
UserData map[string]any `json:"userData,omitempty"`
|
||||
@@ -73,6 +77,8 @@ type RequestStruct struct {
|
||||
|
||||
// ResponseStruct : Antworten an den Client (WEB)
|
||||
type ResponseStruct struct {
|
||||
ID string `json:"id,omitempty"`
|
||||
|
||||
ClientInfo struct {
|
||||
ARCH string `json:"arch"`
|
||||
Branch string `json:"branch,omitempty"`
|
||||
|
||||
+7
-11
@@ -330,22 +330,18 @@ func WS(w http.ResponseWriter, r *http.Request) {
|
||||
|
||||
var newToken string
|
||||
|
||||
/*
|
||||
if r.Header.Get("Origin") != "http://"+r.Host {
|
||||
httpStatusError(w, r, 403)
|
||||
return
|
||||
}
|
||||
*/
|
||||
|
||||
// Upgrade writes its own error response (e.g. 403 for a bad Origin).
|
||||
conn, err := wsUpgrader.Upgrade(w, r, nil)
|
||||
if err != nil {
|
||||
ShowError(err, 0)
|
||||
return
|
||||
}
|
||||
defer conn.Close()
|
||||
|
||||
setGlobalDomain(r.Host)
|
||||
|
||||
// One connection serves any number of commands, one at a time, until the
|
||||
// client goes away or a write fails.
|
||||
for {
|
||||
|
||||
// Fresh structs per command: a failed command must not leak its
|
||||
@@ -360,6 +356,8 @@ func WS(w http.ResponseWriter, r *http.Request) {
|
||||
return
|
||||
}
|
||||
|
||||
response.ID = request.ID
|
||||
|
||||
if !System.ConfigurationWizard {
|
||||
|
||||
switch Settings.AuthenticationWEB {
|
||||
@@ -398,10 +396,9 @@ func WS(w http.ResponseWriter, r *http.Request) {
|
||||
response = setDefaultResponseData(response, false)
|
||||
if err = conn.WriteJSON(response); err != nil {
|
||||
ShowError(err, 1022)
|
||||
} else {
|
||||
return
|
||||
}
|
||||
return
|
||||
continue
|
||||
|
||||
case "loadFiles":
|
||||
//response.Response = Settings.Files
|
||||
@@ -572,8 +569,7 @@ func WS(w http.ResponseWriter, r *http.Request) {
|
||||
|
||||
if err = conn.WriteJSON(response); err != nil {
|
||||
ShowError(err, 1022)
|
||||
} else {
|
||||
break
|
||||
return
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -5,6 +5,7 @@ import (
|
||||
"net/http/httptest"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/gorilla/websocket"
|
||||
|
||||
@@ -137,3 +138,33 @@ func TestSetCookieTokenFlags(t *testing.T) {
|
||||
t.Error("logout must clear the cookie (negative MaxAge)")
|
||||
}
|
||||
}
|
||||
|
||||
// A browser keeps one socket open and sends commands one after another; the
|
||||
// server must answer each on the same connection and echo the request id.
|
||||
func TestWSServesMultipleCommandsPerConnection(t *testing.T) {
|
||||
srv, token := newAuthenticatedWSServer(t)
|
||||
|
||||
conn, _, err := websocket.DefaultDialer.Dial(wsURL(srv), http.Header{"Cookie": {"Token=" + token}})
|
||||
if err != nil {
|
||||
t.Fatalf("dial: %v", err)
|
||||
}
|
||||
defer conn.Close()
|
||||
|
||||
for i := 1; i <= 3; i++ {
|
||||
id := strings.Repeat("x", i)
|
||||
if err := conn.WriteJSON(map[string]any{"cmd": "noop", "id": id}); err != nil {
|
||||
t.Fatalf("write %d: %v", i, err)
|
||||
}
|
||||
conn.SetReadDeadline(time.Now().Add(2 * time.Second))
|
||||
var response ResponseStruct
|
||||
if err := conn.ReadJSON(&response); err != nil {
|
||||
t.Fatalf("read %d: %v (server stopped answering on the same connection)", i, err)
|
||||
}
|
||||
if !response.Status {
|
||||
t.Fatalf("command %d refused: %s", i, response.Error)
|
||||
}
|
||||
if response.ID != id {
|
||||
t.Fatalf("command %d: id %q not echoed, got %q", i, id, response.ID)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
+8
-8
@@ -1,7 +1,7 @@
|
||||
# xTeVe improvement checklist
|
||||
|
||||
Detailed rationale, file references and effort estimates: `tasks/improvement-plan.md`.
|
||||
Status: Phase 0 committed on branch `improvements` 2026-09-26.
|
||||
Status: Phases 0, 1, 2, 3 done on branch `improvements` (2026-09-26); Phase 5 in progress.
|
||||
|
||||
## Phase 0: Hygiene
|
||||
- [x] `.gitignore` (`.gocache/`), deleted `.gocache/`, extended `.dockerignore`
|
||||
@@ -35,12 +35,12 @@ Status: Phase 0 committed on branch `improvements` 2026-09-26.
|
||||
- [x] `streamingURLsMu` around `Data.Cache.StreamingURLS`; persisted from a snapshot
|
||||
- [x] No more `defer` inside the buffer read/redirect loops (compression done in Phase 1; imgcache in 2b)
|
||||
- [x] `thirdPartyBuffer` rewritten: `exec.CommandContext` cancelled when the last client leaves or the 20 s startup watchdog fires; `buildFFmpegArgs`/`buildVLCArgs` tested; no panic, `cmd.Start` checked, one file handle per segment. Client loop uses `r.Context()` instead of `CloseNotifier`
|
||||
- [ ] Real mutex for screen log
|
||||
- [ ] `http.Server` with timeouts; timeouts on all outbound clients
|
||||
- [ ] imgcache: download outside the lock
|
||||
- [ ] Atomic write helper (temp + fsync + rename)
|
||||
- [x] Buffer: headers set before `WriteHeader`, bogus `Content-Length:` header gone, segments flushed to the client as they arrive. Others in 2b
|
||||
- [ ] Error hygiene: no-op closures, `os.Exit`/`panic` outside main, ignored results
|
||||
- [x] `logMu` guards `WebScreenLog` and notifications; accessors `logAppend`/`webScreenLogSnapshot`/`notificationsSnapshot`; ring buffer keeps the newest N (was dropping them)
|
||||
- [x] `http.Server` with `ReadHeaderTimeout` 15 s and `IdleTimeout` 120 s (no write timeout, streams are long-lived); `providerHTTPClient` 5 min, `apiHTTPClient` 30 s, imgcache client 30 s, stream client with dial/TLS/header timeouts
|
||||
- [x] imgcache downloads outside the lock; cache URL bug fixed (was a filesystem path); query-string URLs no longer produce unservable file names
|
||||
- [x] `writeFileAtomic` (temp + fsync + rename) behind `writeByteToFile`, `saveMapToJSONFile`, `writePrivateFile` and the auth database
|
||||
- [x] Buffer: headers set before `WriteHeader`, bogus `Content-Length:` header gone, segments flushed as they arrive. xepg file-removal bug, dead range mutation, API double write, WS struct reuse, unchecked assertions in data/backup/provider/screen, notification eviction by age: all fixed with tests
|
||||
- [x] Error hygiene: no `panic`/`os.Exit`/`log.Fatal` outside `main`; SIGINT/SIGTERM handled in main via `src.Shutdown()`; fatal start-up errors exit 1
|
||||
|
||||
## Phase 3: Build, embed, CI, Docker
|
||||
- [x] `//go:embed` via `html/embed.go`; deleted `webUI.go`, `html-build.go`, `cmd/webui-gen`; ETag + Cache-Control on static assets
|
||||
@@ -58,7 +58,7 @@ Status: Phase 0 committed on branch `improvements` 2026-09-26.
|
||||
Not planned (lineup is ~170 channels, decided 2026-09-25). See plan for the reference list.
|
||||
|
||||
## Phase 5: Frontend architecture
|
||||
- [ ] Persistent websocket with request IDs, queue, backoff reconnect
|
||||
- [x] Persistent websocket: one socket per page, queued requests with ids echoed by the server, 12 s per-request timeout, exponential-backoff reconnect, log-poll dedupe. Server keeps serving on one connection and closes it on exit (upstream leaked one open socket per request). Test: `TestWSServesMultipleCommandsPerConnection`
|
||||
- [ ] Section-level re-render instead of full `createLayout()`
|
||||
- [ ] Virtualised mapping table
|
||||
- [ ] Split `menu_ts.ts`; `addEventListener`; data-driven settings rows
|
||||
|
||||
@@ -3,7 +3,6 @@ var BULK_EDIT:Boolean = false
|
||||
var COLUMN_TO_SORT:number
|
||||
var SEARCH_MAPPING = new Object()
|
||||
var UNDO = new Object()
|
||||
var SERVER_CONNECTION = false
|
||||
var WS_AVAILABLE = false
|
||||
var ACTIVE_MENU_ID:string = ""
|
||||
var LAST_BULK_CHECKBOX:HTMLInputElement = null
|
||||
@@ -761,7 +760,6 @@ function sortSelect(elem) {
|
||||
|
||||
function updateLog() {
|
||||
|
||||
console.log("TOKEN")
|
||||
var server:Server = new Server("updateLog")
|
||||
server.request(new Object())
|
||||
|
||||
|
||||
+1
-1
@@ -1023,7 +1023,7 @@ function PageReady() {
|
||||
return
|
||||
}
|
||||
|
||||
if (SERVER_CONNECTION == true) {
|
||||
if (WS.hasQueued("getServerConfig")) {
|
||||
return
|
||||
}
|
||||
|
||||
|
||||
+228
-125
@@ -1,136 +1,240 @@
|
||||
class Server {
|
||||
protocol:string
|
||||
cmd:string
|
||||
// Websocket transport between the web UI and xTeVe.
|
||||
//
|
||||
// One socket stays open for the life of the page. Requests are queued and
|
||||
// sent one at a time; each carries an id that the server echoes back, and
|
||||
// the response is matched on it (or, for older servers that close after a
|
||||
// single reply, to the oldest in-flight request). Connection loss triggers
|
||||
// a reconnect with exponential backoff, and queued requests survive it.
|
||||
|
||||
constructor(cmd:string) {
|
||||
this.cmd = cmd
|
||||
interface PendingRequest {
|
||||
id: string
|
||||
payload: Object
|
||||
resolve: (response: any) => void
|
||||
reject: (err: Error) => void
|
||||
timer: number
|
||||
sent: boolean
|
||||
}
|
||||
|
||||
class WSClient {
|
||||
private socket: WebSocket = null
|
||||
private queue: PendingRequest[] = []
|
||||
private inFlight: PendingRequest = null
|
||||
private nextID: number = 1
|
||||
private reconnectDelay: number = 500
|
||||
private reconnectTimer: number = 0
|
||||
private closedByUs: boolean = false
|
||||
|
||||
readonly requestTimeoutMs: number = 12000
|
||||
|
||||
private url(): string {
|
||||
var protocol = window.location.protocol == "https:" ? "wss://" : "ws://"
|
||||
var host = window.location.host
|
||||
if (host == undefined || host.length < 1) {
|
||||
host = window.location.hostname
|
||||
}
|
||||
return protocol + host + "/data/"
|
||||
}
|
||||
|
||||
request(data:Object):any {
|
||||
// send queues a command and resolves with the parsed response.
|
||||
send(payload: Object): Promise<any> {
|
||||
var client = this
|
||||
return new Promise(function(resolve, reject) {
|
||||
var req: PendingRequest = {
|
||||
id: String(client.nextID++),
|
||||
payload: payload,
|
||||
resolve: resolve,
|
||||
reject: reject,
|
||||
timer: 0,
|
||||
sent: false
|
||||
}
|
||||
client.queue.push(req)
|
||||
client.pump()
|
||||
})
|
||||
}
|
||||
|
||||
if (SERVER_CONNECTION == true) {
|
||||
// hasQueued reports whether a command of this kind is already waiting.
|
||||
hasQueued(cmd: string): boolean {
|
||||
if (this.inFlight != null && this.inFlight.payload["cmd"] == cmd) {
|
||||
return true
|
||||
}
|
||||
for (var i = 0; i < this.queue.length; i++) {
|
||||
if (this.queue[i].payload["cmd"] == cmd) {
|
||||
return true
|
||||
}
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
private connect(): void {
|
||||
if (this.socket != null && (this.socket.readyState == WebSocket.OPEN || this.socket.readyState == WebSocket.CONNECTING)) {
|
||||
return
|
||||
}
|
||||
|
||||
SERVER_CONNECTION = true
|
||||
|
||||
if (this.cmd != "updateLog") {
|
||||
var client = this
|
||||
var socket = new WebSocket(this.url())
|
||||
this.socket = socket
|
||||
|
||||
socket.onopen = function() {
|
||||
WS_AVAILABLE = true
|
||||
client.reconnectDelay = 500
|
||||
client.pump()
|
||||
}
|
||||
|
||||
socket.onmessage = function(e: MessageEvent) {
|
||||
client.handleMessage(e.data)
|
||||
}
|
||||
|
||||
socket.onerror = function() {
|
||||
// onclose follows; nothing to do here.
|
||||
}
|
||||
|
||||
socket.onclose = function() {
|
||||
if (client.socket == socket) {
|
||||
client.socket = null
|
||||
}
|
||||
// An in-flight request that never got its answer goes back to the
|
||||
// front of the queue and is retried on the next connection.
|
||||
if (client.inFlight != null) {
|
||||
var req = client.inFlight
|
||||
client.inFlight = null
|
||||
req.sent = false
|
||||
client.queue.unshift(req)
|
||||
}
|
||||
client.scheduleReconnect()
|
||||
}
|
||||
}
|
||||
|
||||
private scheduleReconnect(): void {
|
||||
if (this.closedByUs || this.reconnectTimer != 0) {
|
||||
return
|
||||
}
|
||||
if (this.queue.length == 0) {
|
||||
// Idle: reconnect lazily on the next request.
|
||||
return
|
||||
}
|
||||
var client = this
|
||||
this.reconnectTimer = window.setTimeout(function() {
|
||||
client.reconnectTimer = 0
|
||||
client.connect()
|
||||
}, this.reconnectDelay)
|
||||
this.reconnectDelay = Math.min(this.reconnectDelay * 2, 10000)
|
||||
}
|
||||
|
||||
// pump sends the next queued request if the socket is open and idle.
|
||||
private pump(): void {
|
||||
if (this.inFlight != null || this.queue.length == 0) {
|
||||
return
|
||||
}
|
||||
if (this.socket == null || this.socket.readyState != WebSocket.OPEN) {
|
||||
this.connect()
|
||||
return
|
||||
}
|
||||
|
||||
var req = this.queue.shift()
|
||||
this.inFlight = req
|
||||
req.sent = true
|
||||
req.payload["id"] = req.id
|
||||
|
||||
var client = this
|
||||
req.timer = window.setTimeout(function() {
|
||||
if (client.inFlight == req) {
|
||||
client.inFlight = null
|
||||
req.reject(new Error("timeout"))
|
||||
// A socket that stopped answering is not trusted any more.
|
||||
client.resetSocket()
|
||||
client.pump()
|
||||
}
|
||||
}, this.requestTimeoutMs)
|
||||
|
||||
try {
|
||||
this.socket.send(JSON.stringify(req.payload))
|
||||
} catch (err) {
|
||||
window.clearTimeout(req.timer)
|
||||
this.inFlight = null
|
||||
req.sent = false
|
||||
this.queue.unshift(req)
|
||||
this.resetSocket()
|
||||
}
|
||||
}
|
||||
|
||||
private resetSocket(): void {
|
||||
var socket = this.socket
|
||||
this.socket = null
|
||||
if (socket != null) {
|
||||
socket.onclose = null
|
||||
try { socket.close() } catch (err) { /* ignore */ }
|
||||
}
|
||||
this.scheduleReconnect()
|
||||
}
|
||||
|
||||
private handleMessage(raw: string): void {
|
||||
var response: any
|
||||
try {
|
||||
response = JSON.parse(raw)
|
||||
} catch (err) {
|
||||
return
|
||||
}
|
||||
|
||||
var req = this.inFlight
|
||||
if (req == null) {
|
||||
return
|
||||
}
|
||||
// Servers that echo the id let us drop stray late answers.
|
||||
if (response.hasOwnProperty("id") && response["id"] != req.id) {
|
||||
return
|
||||
}
|
||||
|
||||
window.clearTimeout(req.timer)
|
||||
this.inFlight = null
|
||||
req.resolve(response)
|
||||
this.pump()
|
||||
}
|
||||
}
|
||||
|
||||
var WS = new WSClient()
|
||||
|
||||
class Server {
|
||||
cmd: string
|
||||
|
||||
constructor(cmd: string) {
|
||||
this.cmd = cmd
|
||||
}
|
||||
|
||||
request(data: Object): any {
|
||||
|
||||
var isLogUpdate: boolean = this.cmd == "updateLog"
|
||||
|
||||
// Log polling must never pile up behind a stalled connection.
|
||||
if (isLogUpdate && WS.hasQueued("updateLog")) {
|
||||
return
|
||||
}
|
||||
|
||||
if (!isLogUpdate) {
|
||||
showElement("loading", true)
|
||||
UNDO = new Object()
|
||||
setConnectionState("busy")
|
||||
}
|
||||
|
||||
switch(window.location.protocol) {
|
||||
case "http:":
|
||||
this.protocol = "ws://"
|
||||
break
|
||||
case "https:":
|
||||
this.protocol = "wss://"
|
||||
break
|
||||
}
|
||||
|
||||
var wsHost:string = window.location.host
|
||||
if (wsHost == undefined || wsHost.length < 1) {
|
||||
wsHost = window.location.hostname
|
||||
}
|
||||
// The session cookie (HttpOnly) is sent with the websocket handshake.
|
||||
var url = this.protocol + wsHost + "/data/"
|
||||
|
||||
data["cmd"] = this.cmd
|
||||
var requestCmd:string = data["cmd"]
|
||||
var ws = new WebSocket(url)
|
||||
var isLogUpdate:boolean = data["cmd"] == "updateLog"
|
||||
var responseReceived:boolean = false
|
||||
var requestFinished:boolean = false
|
||||
var timeoutMs:number = 12000
|
||||
var requestTimeout:number
|
||||
var requestCmd: string = this.cmd
|
||||
|
||||
var finishRequest = function(state:string, responseSuccess:boolean = false):void {
|
||||
if (requestFinished == true) {
|
||||
return
|
||||
}
|
||||
WS.send(data).then(function(response) {
|
||||
|
||||
requestFinished = true
|
||||
SERVER_CONNECTION = false
|
||||
window.clearTimeout(requestTimeout)
|
||||
|
||||
if (responseSuccess == true) {
|
||||
if (state == "online") {
|
||||
WS_FAILURE_COUNT = 0
|
||||
}
|
||||
} else {
|
||||
WS_FAILURE_COUNT++
|
||||
}
|
||||
|
||||
if (isLogUpdate == false) {
|
||||
WS_FAILURE_COUNT = 0
|
||||
if (!isLogUpdate) {
|
||||
showElement("loading", false)
|
||||
}
|
||||
|
||||
if (state != "") {
|
||||
setConnectionState(state)
|
||||
}
|
||||
}
|
||||
|
||||
requestTimeout = window.setTimeout(function() {
|
||||
console.log("Websocket request timed out.")
|
||||
var timeoutState:string = "offline"
|
||||
if (isLogUpdate == true && WS_FAILURE_COUNT < 2) {
|
||||
timeoutState = "idle"
|
||||
}
|
||||
finishRequest(timeoutState, false)
|
||||
try {
|
||||
ws.close()
|
||||
} catch (err) {
|
||||
console.log(err)
|
||||
}
|
||||
}, timeoutMs)
|
||||
|
||||
ws.onopen = function() {
|
||||
|
||||
WS_AVAILABLE = true
|
||||
if (data["cmd"] != "updateLog") {
|
||||
setConnectionState("busy")
|
||||
}
|
||||
|
||||
this.send(JSON.stringify(data));
|
||||
|
||||
}
|
||||
|
||||
ws.onerror = function(e) {
|
||||
|
||||
console.log("No websocket connection to xTeVe could be established. Check your network configuration.")
|
||||
var errorState:string = "offline"
|
||||
if (isLogUpdate == true && WS_FAILURE_COUNT < 2) {
|
||||
errorState = "idle"
|
||||
}
|
||||
finishRequest(errorState, false)
|
||||
|
||||
if (WS_AVAILABLE == false && isLogUpdate == false && requestCmd != "getServerConfig") {
|
||||
alert("No websocket connection to xTeVe could be established. Check your network configuration.")
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
|
||||
ws.onmessage = function (e) {
|
||||
responseReceived = true
|
||||
finishRequest("online", true)
|
||||
|
||||
var response = JSON.parse(e.data);
|
||||
setConnectionState("online")
|
||||
|
||||
if (response["status"] == false) {
|
||||
setConnectionState("offline")
|
||||
|
||||
alert(response["err"])
|
||||
|
||||
if (response.hasOwnProperty("reload")) {
|
||||
location.reload()
|
||||
}
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
|
||||
if (response.hasOwnProperty("logoURL")) {
|
||||
var div = (document.getElementById("channel-icon") as HTMLInputElement)
|
||||
div.value = response["logoURL"]
|
||||
@@ -138,19 +242,17 @@ class Server {
|
||||
return
|
||||
}
|
||||
|
||||
switch (data["cmd"]) {
|
||||
switch (requestCmd) {
|
||||
case "updateLog":
|
||||
SERVER["log"] = response["log"]
|
||||
if (document.getElementById("content_log")) {
|
||||
showLogs(false)
|
||||
}
|
||||
return
|
||||
break;
|
||||
|
||||
|
||||
default:
|
||||
SERVER = new Object()
|
||||
SERVER = response
|
||||
break;
|
||||
break
|
||||
}
|
||||
|
||||
if (response.hasOwnProperty("openMenu")) {
|
||||
@@ -171,7 +273,6 @@ class Server {
|
||||
location.reload()
|
||||
}
|
||||
|
||||
|
||||
if (response.hasOwnProperty("wizard")) {
|
||||
createLayout()
|
||||
configurationWizard[response["wizard"]].createWizard()
|
||||
@@ -180,23 +281,25 @@ class Server {
|
||||
|
||||
createLayout()
|
||||
|
||||
}
|
||||
}, function(err) {
|
||||
|
||||
ws.onclose = function() {
|
||||
if (responseReceived == true) {
|
||||
return
|
||||
WS_FAILURE_COUNT++
|
||||
if (!isLogUpdate) {
|
||||
showElement("loading", false)
|
||||
}
|
||||
|
||||
var closeState:string = "offline"
|
||||
if (isLogUpdate == true && WS_FAILURE_COUNT < 2) {
|
||||
closeState = "idle"
|
||||
var state: string = "offline"
|
||||
if (isLogUpdate && WS_FAILURE_COUNT < 2) {
|
||||
state = "idle"
|
||||
}
|
||||
finishRequest(closeState, false)
|
||||
}
|
||||
|
||||
setConnectionState(state)
|
||||
|
||||
if (WS_AVAILABLE == false && !isLogUpdate && requestCmd != "getServerConfig") {
|
||||
alert("No websocket connection to xTeVe could be established. Check your network configuration.")
|
||||
}
|
||||
|
||||
})
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
var WS_FAILURE_COUNT:number = 0
|
||||
|
||||
var WS_FAILURE_COUNT: number = 0
|
||||
|
||||
Reference in New Issue
Block a user