diff --git a/server/bundles/io.cloudbeaver.server/src/io/cloudbeaver/server/websockets/CBJettyWebSocketManager.java b/server/bundles/io.cloudbeaver.server/src/io/cloudbeaver/server/websockets/CBJettyWebSocketManager.java index 806e8329a5..e73700ff82 100644 --- a/server/bundles/io.cloudbeaver.server/src/io/cloudbeaver/server/websockets/CBJettyWebSocketManager.java +++ b/server/bundles/io.cloudbeaver.server/src/io/cloudbeaver/server/websockets/CBJettyWebSocketManager.java @@ -51,7 +51,6 @@ public class CBJettyWebSocketManager implements JettyWebSocketCreator { return null; } if (webSession == null) { - log.error("CloudBeaver session not found"); return null; } var webSessionId = webSession.getSessionId(); @@ -68,13 +67,17 @@ public class CBJettyWebSocketManager implements JettyWebSocketCreator { @Nullable private BaseWebSession resolveWebSession(JettyServerUpgradeRequest request) throws DBException { if (request.getHttpServletRequest().getSession() == null) { + log.debug("CloudBeaver web session not exist, try to create headless session"); return webSessionManager.getHeadlessSession(request.getHttpServletRequest(), true); } var webSessionId = request.getHttpServletRequest().getSession().getId(); var webSession = webSessionManager.getSession(webSessionId); - return webSession != null - ? webSession - : webSessionManager.getHeadlessSession(request.getHttpServletRequest(), true); + if (webSession != null) { + return webSession; + } + log.error("CloudBeaver session not found with id " + webSessionId + ", try to create headless session"); + + return webSessionManager.getHeadlessSession(request.getHttpServletRequest(), true); } public void sendPing() { diff --git a/webapp/packages/core-root/src/SessionEventSource.ts b/webapp/packages/core-root/src/SessionEventSource.ts index c03b57efdc..f8cba93b6b 100644 --- a/webapp/packages/core-root/src/SessionEventSource.ts +++ b/webapp/packages/core-root/src/SessionEventSource.ts @@ -6,7 +6,7 @@ * you may not use this file except in compliance with the License. */ -import { catchError, debounceTime, filter, map, merge, Observable, retry, RetryConfig, Subject } from 'rxjs'; +import { catchError, debounceTime, filter, interval, map, merge, Observable, repeat, retry, share, Subject, throwError } from 'rxjs'; import { webSocket, WebSocketSubject } from 'rxjs/webSocket'; import { injectable } from '@cloudbeaver/core-di'; @@ -20,7 +20,9 @@ import { ServiceError, } from '@cloudbeaver/core-sdk'; +import { NetworkStateService } from './NetworkStateService'; import type { IBaseServerEvent, IServerEventCallback, IServerEventEmitter, Subscription } from './ServerEventEmitter/IServerEventEmitter'; +import { SessionExpireService } from './SessionExpireService'; export { ServerEventId, SessionEventTopic, ClientEventId }; @@ -37,11 +39,7 @@ export interface ITopicSubEvent extends ISessionEvent { topicId: SessionEventTopic; } -const retryInterval = 5000; - -const retryConfig: RetryConfig = { - delay: retryInterval, -}; +const RETRY_INTERVAL = 30 * 1000; @injectable() export class SessionEventSource @@ -55,8 +53,11 @@ implements IServerEventEmitter; private readonly subject: WebSocketSubject; private readonly oldEventsSubject: Subject; + private readonly retryTimer: Observable; constructor( + private readonly networkStateService: NetworkStateService, + private readonly sessionExpireService: SessionExpireService, private readonly environmentService: EnvironmentService, private readonly graphQLService: GraphQLService ) { @@ -65,6 +66,10 @@ implements IServerEventEmitter !this.sessionExpireService.expired && networkStateService.state) + ); this.subject = webSocket({ url: environmentService.wsEndpoint, closeObserver: this.closeSubject, @@ -97,7 +102,7 @@ implements IServerEventEmitter event.id === id), map(mapTo) ) @@ -115,7 +120,7 @@ implements IServerEventEmitter):Observable => source.pipe( + share(), + catchError(this.errorHandler), + retry({ delay: () => this.retryTimer }), + repeat({ delay: () => this.retryTimer }), + ); + } + private errorHandler(error: any, caught: Observable): Observable { this.errorSubject.next(new ServiceError('WebSocket connection error')); - return caught; + return throwError(() => error); } }