(UI; tauri-websocket): Enhance error handling and connection management in WebSocket implementation
This commit is contained in:
@@ -1,10 +1,11 @@
|
|||||||
import WebSocket, { ConnectionConfig, Message } from '@tauri-apps/plugin-websocket';
|
import WebSocket, { ConnectionConfig, Message } from '@tauri-apps/plugin-websocket';
|
||||||
import { Subject, Observable } from 'rxjs';
|
import { Subject, Observable, merge, mergeMap, throwError } from 'rxjs';
|
||||||
import { WebSocketSubject, WebSocketSubjectConfig } from 'rxjs/webSocket';
|
import { WebSocketSubject, WebSocketSubjectConfig } from 'rxjs/webSocket';
|
||||||
import { NgZone } from '@angular/core';
|
import { NgZone } from '@angular/core';
|
||||||
|
|
||||||
const LOG_PREFIX = '[tauri_ws]';
|
const LOG_PREFIX = '[tauri_ws]';
|
||||||
|
const PING_INTERVAL_MS = 10000; // Send a ping every PING_INTERVAL_MS milliseconds
|
||||||
|
const PONG_TIMEOUT_MS = 5000; // Wait PONG_TIMEOUT_MS milliseconds for a pong response
|
||||||
/**
|
/**
|
||||||
* Creates a WebSocket connection using the Tauri WebSocket API and wraps it in an RxJS WebSocketSubject-compatible interface.
|
* Creates a WebSocket connection using the Tauri WebSocket API and wraps it in an RxJS WebSocketSubject-compatible interface.
|
||||||
*
|
*
|
||||||
@@ -33,224 +34,227 @@ export function createTauriWsConnection<T>(opts: WebSocketSubjectConfig<T>, ngZo
|
|||||||
|
|
||||||
let wsConnection: WebSocket | null = null;
|
let wsConnection: WebSocket | null = null;
|
||||||
const messageSubject = new Subject<T>();
|
const messageSubject = new Subject<T>();
|
||||||
const observable$ = messageSubject.asObservable();
|
const errorSubject = new Subject<any>(); // Added for error propagation
|
||||||
|
|
||||||
// A queue for messages that need to be sent before the connection is established
|
|
||||||
const pendingMessages: T[] = [];
|
|
||||||
|
|
||||||
// Function to establish a WebSocket connection
|
// Combined stream with both messages and errors
|
||||||
const connect = (): void => {
|
const observable$ = merge(
|
||||||
WebSocket.connect(opts.url)
|
messageSubject.asObservable(),
|
||||||
.then((ws) => {
|
errorSubject.pipe(
|
||||||
wsConnection = ws;
|
mergeMap(err => throwError(() => err))
|
||||||
console.log(`${LOG_PREFIX} Connection established`);
|
)
|
||||||
|
);
|
||||||
// Run inside NgZone to ensure Angular detects this change
|
|
||||||
ngZone.run(() => {
|
|
||||||
// Notify that connection is open
|
|
||||||
opts.openObserver?.next(undefined as unknown as Event);
|
|
||||||
// Send any pending messages
|
|
||||||
while (pendingMessages.length > 0) {
|
|
||||||
const message = pendingMessages.shift();
|
|
||||||
if (message) webSocketSubject.next(message);
|
|
||||||
}
|
|
||||||
});
|
|
||||||
|
|
||||||
setupMessageListener(ws);
|
//////////////////////////////////////////////////////////////
|
||||||
})
|
// Track subscriptions
|
||||||
.catch((error: Error) => {
|
//////////////////////////////////////////////////////////////
|
||||||
console.error(`${LOG_PREFIX} Connection failed:`, error);
|
let subscriptionCount = 0;
|
||||||
reconnect();
|
|
||||||
});
|
// Wrapper with subscription tracking
|
||||||
};
|
const trackedObservable$ = new Observable<T>(subscriber => {
|
||||||
|
subscriptionCount++;
|
||||||
// Function to reconnect
|
|
||||||
let reconnectTimeout: ReturnType<typeof setTimeout> | null = null;
|
// If this is the first subscription, connect to WebSocket
|
||||||
const reconnect = () => {
|
if (subscriptionCount === 1) {
|
||||||
if (reconnectTimeout) {
|
ngZone.runOutsideAngular(() => {
|
||||||
clearTimeout(reconnectTimeout);
|
|
||||||
}
|
|
||||||
|
|
||||||
// Notify close observer
|
|
||||||
ngZone.run(() => {
|
|
||||||
opts.closeObserver?.next(undefined as unknown as CloseEvent);
|
|
||||||
})
|
|
||||||
|
|
||||||
// Remove the existing listener if it exists
|
|
||||||
removeListener();
|
|
||||||
|
|
||||||
// Close the existing connection if it exists
|
|
||||||
wsConnection?.disconnect().catch(err => console.warn(`${LOG_PREFIX} Error closing connection during reconnect:`, err));
|
|
||||||
wsConnection = null;
|
|
||||||
|
|
||||||
// Connect again after a delay
|
|
||||||
console.log(`${LOG_PREFIX} Attempting to reconnect in 2 seconds...`);
|
|
||||||
reconnectTimeout = setTimeout(() => {
|
|
||||||
reconnectTimeout = null;
|
|
||||||
connect();
|
connect();
|
||||||
}, 2000);
|
|
||||||
};
|
|
||||||
|
|
||||||
// Function to check if connection alive, and reconnect, if necessary
|
|
||||||
let pingTimeoutId: ReturnType<typeof setTimeout> | null = null;
|
|
||||||
const checkConnectionAliveOrReconnect = () => {
|
|
||||||
if (!pingTimeoutId && wsConnection) {
|
|
||||||
console.log(`${LOG_PREFIX} Verifying connection status with ping...`);
|
|
||||||
wsConnection.send( {type: 'Ping', data: [1]} ).then(() => {
|
|
||||||
pingTimeoutId = setTimeout(() => { // The timeout will be stopped if we receive a Pong response
|
|
||||||
console.error(`${LOG_PREFIX} No response to ping - connection appears dead`);
|
|
||||||
pingTimeoutId = null;
|
|
||||||
reconnect();
|
|
||||||
}, 5000);
|
|
||||||
}).catch(err => {
|
|
||||||
console.error(`${LOG_PREFIX} Failed to send ping - connection is dead:`, err);
|
|
||||||
reconnect();
|
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
}
|
|
||||||
|
|
||||||
// Function to remove the message listener
|
const subscription = observable$.subscribe({
|
||||||
|
next: value => subscriber.next(value),
|
||||||
|
error: err => subscriber.error(err),
|
||||||
|
complete: () => subscriber.complete()
|
||||||
|
});
|
||||||
|
|
||||||
|
// Cleanup function - called when unsubscribed
|
||||||
|
return () => {
|
||||||
|
subscriptionCount--;
|
||||||
|
subscription.unsubscribe();
|
||||||
|
|
||||||
|
// If this was the last subscription, close the WebSocket connection
|
||||||
|
if (subscriptionCount === 0) {
|
||||||
|
disconnect();
|
||||||
|
}
|
||||||
|
};
|
||||||
|
});
|
||||||
|
|
||||||
|
//////////////////////////////////////////////////////////////
|
||||||
|
// Function to establish a WebSocket connection
|
||||||
|
//////////////////////////////////////////////////////////////
|
||||||
let listenerRemovalFn: (() => void) | null = null; // Store the removal function for the ws listener
|
let listenerRemovalFn: (() => void) | null = null; // Store the removal function for the ws listener
|
||||||
const removeListener = () => {
|
|
||||||
if (listenerRemovalFn) {
|
const connect = (): void => {
|
||||||
try {
|
console.log(`${LOG_PREFIX} Connecting to WebSocket: ${opts.url}`);
|
||||||
listenerRemovalFn();
|
WebSocket.connect(opts.url)
|
||||||
listenerRemovalFn = null; // Clear the reference
|
.then((ws) => {
|
||||||
} catch (err) {
|
wsConnection = ws;
|
||||||
console.error(`${LOG_PREFIX} Error removing listener:`, err);
|
console.log(`${LOG_PREFIX} Connection established`);
|
||||||
}
|
opts.openObserver?.next(undefined as unknown as Event);
|
||||||
|
listenerRemovalFn = ws.addListener(messagesListener);
|
||||||
|
startHealthChecks();
|
||||||
|
})
|
||||||
|
.catch((error: Error) => {
|
||||||
|
console.error(`${LOG_PREFIX} Connection failed:`, error);
|
||||||
|
errorSubject.next(error);
|
||||||
|
});
|
||||||
|
};
|
||||||
|
|
||||||
|
const disconnect = (): void => {
|
||||||
|
stopHealthChecks();
|
||||||
|
|
||||||
|
if (listenerRemovalFn) {
|
||||||
|
try {
|
||||||
|
listenerRemovalFn();
|
||||||
|
} catch (err) {
|
||||||
|
console.error(`${LOG_PREFIX} Error removing listener:`, err);
|
||||||
}
|
}
|
||||||
|
listenerRemovalFn = null; // Clear the reference
|
||||||
|
}
|
||||||
|
|
||||||
|
const currentWs = wsConnection;
|
||||||
|
wsConnection = null;
|
||||||
|
|
||||||
|
if (!currentWs) return;
|
||||||
|
|
||||||
|
console.log(`${LOG_PREFIX} Closing WebSocket connection.`);
|
||||||
|
opts.closeObserver?.next(undefined as unknown as CloseEvent);
|
||||||
|
currentWs.disconnect().catch(err => console.warn(`${LOG_PREFIX} Error closing connection:`, err));
|
||||||
}
|
}
|
||||||
|
|
||||||
// Function to set up the message listener
|
//////////////////////////////////////////////////////////////
|
||||||
const setupMessageListener = (ws: WebSocket) => {
|
// Function to check if connection alive
|
||||||
listenerRemovalFn = ws.addListener((message: Message) => {
|
//////////////////////////////////////////////////////////////
|
||||||
// Process message inside ngZone to trigger change detection
|
let healthCheckIntervalId: ReturnType<typeof setInterval> | null = null;
|
||||||
try {
|
let pongTimeoutId: ReturnType<typeof setTimeout> | null = null;
|
||||||
switch (message.type) {
|
|
||||||
case 'Text':
|
|
||||||
try {
|
|
||||||
const deserializedMessage = deserializer({ data: message.data as string } as any);
|
|
||||||
ngZone.run(() => { messageSubject.next(deserializedMessage); }); // inside ngZone to trigger change detection
|
|
||||||
} catch (err) {
|
|
||||||
console.error(`${LOG_PREFIX} Error deserializing text message:`, err);
|
|
||||||
}
|
|
||||||
break;
|
|
||||||
|
|
||||||
case 'Binary':
|
|
||||||
try {
|
|
||||||
const uint8Array = new Uint8Array(message.data as number[]);
|
|
||||||
const deserializedMessage = deserializer({ data: uint8Array.buffer } as any);
|
|
||||||
ngZone.run(() => { messageSubject.next(deserializedMessage); }); // inside ngZone to trigger change detection
|
|
||||||
} catch (err) {
|
|
||||||
console.error(`${LOG_PREFIX} Error deserializing binary message:`, err);
|
|
||||||
}
|
|
||||||
break;
|
|
||||||
|
|
||||||
case 'Close':
|
|
||||||
console.log(`${LOG_PREFIX} Connection closed by server`);
|
|
||||||
reconnect(); // Auto-reconnect on server-initiated close
|
|
||||||
break;
|
|
||||||
|
|
||||||
case 'Ping':
|
|
||||||
break;
|
|
||||||
|
|
||||||
case 'Pong':
|
const startHealthChecks = () => {
|
||||||
console.log(`${LOG_PREFIX} Received pong response - connection is alive`);
|
stopHealthChecks(); // Ensure no multiple intervals are running
|
||||||
if (pingTimeoutId) {
|
healthCheckIntervalId = setInterval(() => {
|
||||||
// Clear the timeout since we got a response
|
if (!wsConnection) {
|
||||||
clearTimeout(pingTimeoutId);
|
stopHealthChecks();
|
||||||
pingTimeoutId = null;
|
return;
|
||||||
}
|
|
||||||
break;
|
|
||||||
|
|
||||||
// All other message types are unexpected. Proceed with reconnect.
|
|
||||||
default:
|
|
||||||
console.warn(`${LOG_PREFIX} Received unexpected message: '${message}'`);
|
|
||||||
checkConnectionAliveOrReconnect();
|
|
||||||
break;
|
|
||||||
}
|
|
||||||
} catch (error) {
|
|
||||||
console.error(`${LOG_PREFIX} Error processing message: `, error);
|
|
||||||
}
|
}
|
||||||
});
|
if (pongTimeoutId) {
|
||||||
|
// Ping already in flight, waiting for pong.
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
wsConnection.send({ type: 'Ping', data: [] })
|
||||||
|
.then(() => {
|
||||||
|
pongTimeoutId = setTimeout(() => {
|
||||||
|
console.error(`${LOG_PREFIX} No Pong received. Connection is likely dead.`);
|
||||||
|
errorSubject.next(new Error('Connection timed out'));
|
||||||
|
stopHealthChecks();
|
||||||
|
}, PONG_TIMEOUT_MS);
|
||||||
|
})
|
||||||
|
.catch(err => {
|
||||||
|
console.error(`${LOG_PREFIX} Ping send failed:`, err);
|
||||||
|
errorSubject.next(new Error(`Ping send failed: ${err}`));
|
||||||
|
stopHealthChecks();
|
||||||
|
});
|
||||||
|
}, PING_INTERVAL_MS);
|
||||||
|
};
|
||||||
|
|
||||||
|
const stopHealthChecks = () => {
|
||||||
|
if (healthCheckIntervalId) {
|
||||||
|
clearInterval(healthCheckIntervalId);
|
||||||
|
healthCheckIntervalId = null;
|
||||||
|
}
|
||||||
|
if (pongTimeoutId) {
|
||||||
|
clearTimeout(pongTimeoutId);
|
||||||
|
pongTimeoutId = null;
|
||||||
|
}
|
||||||
};
|
};
|
||||||
|
|
||||||
//////////////////////////////////////////////////////////////
|
//////////////////////////////////////////////////////////////
|
||||||
// Connect to WebSocket
|
// Messages listener
|
||||||
//////////////////////////////////////////////////////////////
|
//////////////////////////////////////////////////////////////
|
||||||
console.log(`${LOG_PREFIX} Connecting to WebSocket:`, opts.url);
|
const messagesListener = (message: Message) => {
|
||||||
|
try {
|
||||||
// Connect outside of Angular zone for better performance
|
switch (message.type) {
|
||||||
ngZone.runOutsideAngular(() => {
|
case 'Text':
|
||||||
connect();
|
try {
|
||||||
});
|
const deserializedMessage = deserializer({ data: message.data as string } as any);
|
||||||
|
ngZone.run(() => { messageSubject.next(deserializedMessage); }); // inside ngZone to trigger change detection
|
||||||
|
} catch (err) {
|
||||||
|
console.error(`${LOG_PREFIX} Error deserializing text message:`, err);
|
||||||
|
}
|
||||||
|
break;
|
||||||
|
|
||||||
|
case 'Binary':
|
||||||
|
try {
|
||||||
|
const uint8Array = new Uint8Array(message.data as number[]);
|
||||||
|
const deserializedMessage = deserializer({ data: uint8Array.buffer } as any);
|
||||||
|
ngZone.run(() => { messageSubject.next(deserializedMessage); }); // inside ngZone to trigger change detection
|
||||||
|
} catch (err) {
|
||||||
|
console.error(`${LOG_PREFIX} Error deserializing binary message:`, err);
|
||||||
|
}
|
||||||
|
break;
|
||||||
|
|
||||||
|
case 'Close':
|
||||||
|
console.warn(`${LOG_PREFIX} Connection closed by server: ${message}`);
|
||||||
|
errorSubject.next(new Error(`Connection closed by server: ${message}`));
|
||||||
|
break;
|
||||||
|
|
||||||
|
case 'Ping':
|
||||||
|
break;
|
||||||
|
|
||||||
|
case 'Pong':
|
||||||
|
// Pong received, clear the timeout.
|
||||||
|
if (pongTimeoutId) {
|
||||||
|
clearTimeout(pongTimeoutId);
|
||||||
|
pongTimeoutId = null;
|
||||||
|
}
|
||||||
|
break;
|
||||||
|
|
||||||
|
// All other message types are unexpected. Proceed with reconnect.
|
||||||
|
default:
|
||||||
|
console.warn(`${LOG_PREFIX} Received unexpected message: '${message}'`);
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
} catch (error) {
|
||||||
|
console.error(`${LOG_PREFIX} Error processing message: `, error);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
//////////////////////////////////////////////////////////////
|
//////////////////////////////////////////////////////////////
|
||||||
// RxJS WebSocketSubject-compatible implementation
|
// RxJS WebSocketSubject-compatible interface
|
||||||
//////////////////////////////////////////////////////////////
|
//////////////////////////////////////////////////////////////
|
||||||
const webSocketSubject = {
|
const webSocketSubject = {
|
||||||
// Standard Observer interface methods
|
asObservable: () => trackedObservable$,
|
||||||
next: (message: T) => {
|
next: (message: T) => {
|
||||||
// Run outside NgZone for better performance during send operations
|
// Run outside NgZone for better performance during send operations
|
||||||
ngZone.runOutsideAngular(() => {
|
ngZone.runOutsideAngular(() => {
|
||||||
// If the connection is not established yet, queue the message
|
|
||||||
if (!wsConnection) {
|
if (!wsConnection) {
|
||||||
if (pendingMessages.length >= 100) {
|
errorSubject.next(new Error('Connection not established'));
|
||||||
console.error(`${LOG_PREFIX} Too many pending messages, skipping message`);
|
return;
|
||||||
return;
|
}
|
||||||
}
|
|
||||||
pendingMessages.push(message);
|
|
||||||
console.warn(`${LOG_PREFIX} Connection not established yet, message queued`);
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
|
|
||||||
// Serialize the message using the provided serializer
|
|
||||||
let serializedMessage: any;
|
|
||||||
try {
|
try {
|
||||||
serializedMessage = serializer(message);
|
const serializedMessage = serializer(message);
|
||||||
// 'string' type is enough here, since default serializer for portmaster message returns string
|
// 'string' type is enough here, since default serializer for portmaster message returns string
|
||||||
if (typeof serializedMessage !== 'string')
|
if (typeof serializedMessage !== 'string')
|
||||||
throw new Error('Serialized message is not a string');
|
throw new Error('Serialized message is not a string');
|
||||||
|
|
||||||
|
wsConnection?.send(serializedMessage).catch((err: Error) => {
|
||||||
|
console.error(`${LOG_PREFIX} Error sending message:`, err);
|
||||||
|
errorSubject.next(err);
|
||||||
|
});
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
console.error(`${LOG_PREFIX} Error serializing message:`, error);
|
console.error(`${LOG_PREFIX} Error serializing message:`, error);
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
// Send the serialized message through the WebSocket connection
|
|
||||||
wsConnection?.send(serializedMessage).catch((err: Error) => {
|
|
||||||
console.error(`${LOG_PREFIX} Error sending message:`, err);
|
|
||||||
checkConnectionAliveOrReconnect();
|
|
||||||
});
|
|
||||||
});
|
});
|
||||||
},
|
},
|
||||||
|
|
||||||
complete: () => {
|
complete: () => {
|
||||||
if (wsConnection) {
|
if (wsConnection) {
|
||||||
console.log(`${LOG_PREFIX} Closing connection`);
|
opts.closingObserver?.next();
|
||||||
|
disconnect();
|
||||||
// Run inside NgZone to ensure Angular detects this change
|
}
|
||||||
ngZone.run(() => {
|
|
||||||
opts.closingObserver?.next();
|
|
||||||
wsConnection!.disconnect().catch((err: Error) => console.error(`${LOG_PREFIX} Error closing connection:`, err));
|
|
||||||
wsConnection = null;
|
|
||||||
opts.closeObserver?.next(undefined as unknown as CloseEvent);
|
|
||||||
});
|
|
||||||
|
|
||||||
// Remove the existing listener if it exists
|
|
||||||
removeListener();
|
|
||||||
}
|
|
||||||
|
|
||||||
messageSubject.complete();
|
messageSubject.complete();
|
||||||
},
|
},
|
||||||
|
subscribe: trackedObservable$.subscribe.bind(trackedObservable$),
|
||||||
// RxJS Observable methods required for compatibility
|
pipe: trackedObservable$.pipe.bind(trackedObservable$),
|
||||||
pipe: function(): Observable<any> {
|
|
||||||
// @ts-ignore - Ignore the parameter type mismatch
|
|
||||||
return observable$.pipe(...arguments);
|
|
||||||
},
|
|
||||||
};
|
};
|
||||||
|
|
||||||
// Cast to WebSocketSubject<T>
|
|
||||||
return webSocketSubject as unknown as WebSocketSubject<T>;
|
return webSocketSubject as unknown as WebSocketSubject<T>;
|
||||||
}
|
}
|
||||||
@@ -801,7 +801,7 @@ export class PortapiService {
|
|||||||
}
|
}
|
||||||
},
|
},
|
||||||
error: (err) => {
|
error: (err) => {
|
||||||
console.error(err, attrs);
|
console.error("[portapi] request error:", err.message || err, attrs);
|
||||||
observer.error(err);
|
observer.error(err);
|
||||||
},
|
},
|
||||||
complete: () => {
|
complete: () => {
|
||||||
|
|||||||
Reference in New Issue
Block a user