mirror of
https://github.com/MCV-Software/TWBlue.git
synced 2026-08-21 12:28:11 +02:00
fix(mastodon): close duplicate streaming connections
This commit is contained in:
@@ -473,6 +473,8 @@ class Controller(object):
|
|||||||
def logout_account(self, session_id):
|
def logout_account(self, session_id):
|
||||||
for i in sessions.sessions:
|
for i in sessions.sessions:
|
||||||
if sessions.sessions[i].session_id == session_id: session = sessions.sessions[i]
|
if sessions.sessions[i].session_id == session_id: session = sessions.sessions[i]
|
||||||
|
if hasattr(session, "stop_streaming"):
|
||||||
|
session.stop_streaming()
|
||||||
name =session.get_name()
|
name =session.get_name()
|
||||||
delete_buffers = []
|
delete_buffers = []
|
||||||
for i in self.buffers:
|
for i in self.buffers:
|
||||||
@@ -635,6 +637,8 @@ class Controller(object):
|
|||||||
for item in sessions.sessions:
|
for item in sessions.sessions:
|
||||||
if sessions.sessions[item].logged == False:
|
if sessions.sessions[item].logged == False:
|
||||||
continue
|
continue
|
||||||
|
if hasattr(sessions.sessions[item], "stop_streaming"):
|
||||||
|
sessions.sessions[item].stop_streaming()
|
||||||
sessions.sessions[item].sound.cleaner.cancel()
|
sessions.sessions[item].sound.cleaner.cancel()
|
||||||
log.debug("Saving database for " + sessions.sessions[item].session_id)
|
log.debug("Saving database for " + sessions.sessions[item].session_id)
|
||||||
sessions.sessions[item].save_persistent_data()
|
sessions.sessions[item].save_persistent_data()
|
||||||
|
|||||||
@@ -3,11 +3,11 @@ import os
|
|||||||
import paths
|
import paths
|
||||||
import time
|
import time
|
||||||
import logging
|
import logging
|
||||||
|
import threading
|
||||||
import webbrowser
|
import webbrowser
|
||||||
import wx
|
import wx
|
||||||
import mastodon
|
import mastodon
|
||||||
import demoji
|
import demoji
|
||||||
import config
|
|
||||||
import config_utils
|
import config_utils
|
||||||
import output
|
import output
|
||||||
import application
|
import application
|
||||||
@@ -36,6 +36,9 @@ class Session(base.baseSession):
|
|||||||
self.post_visibility = "public"
|
self.post_visibility = "public"
|
||||||
self.expand_spoilers = False
|
self.expand_spoilers = False
|
||||||
self.software = "mastodon"
|
self.software = "mastodon"
|
||||||
|
self.user_stream = None
|
||||||
|
self.direct_stream = None
|
||||||
|
self._stream_lock = threading.Lock()
|
||||||
pub.subscribe(self.on_status, "mastodon.status_received")
|
pub.subscribe(self.on_status, "mastodon.status_received")
|
||||||
pub.subscribe(self.on_status_updated, "mastodon.status_updated")
|
pub.subscribe(self.on_status_updated, "mastodon.status_updated")
|
||||||
pub.subscribe(self.on_notification, "mastodon.notification_received")
|
pub.subscribe(self.on_notification, "mastodon.notification_received")
|
||||||
@@ -444,19 +447,34 @@ class Session(base.baseSession):
|
|||||||
return
|
return
|
||||||
if self.software == "gotosocial":
|
if self.software == "gotosocial":
|
||||||
return
|
return
|
||||||
listener = streaming.StreamListener(session_name=self.get_name(), user_id=self.db["user_id"])
|
with self._stream_lock:
|
||||||
try:
|
if any(stream is not None and stream.is_alive() for stream in (self.user_stream, self.direct_stream)):
|
||||||
stream_healthy = self.api.stream_healthy()
|
log.debug("Streams for session {} are already running.".format(self.get_name()))
|
||||||
if stream_healthy == True:
|
return
|
||||||
self.user_stream = self.api.stream_user(listener, run_async=True, reconnect_async=True, reconnect_async_wait_sec=30)
|
listener = streaming.StreamListener(session_name=self.get_name(), user_id=self.db["user_id"])
|
||||||
self.direct_stream = self.api.stream_direct(listener, run_async=True, reconnect_async=True, reconnect_async_wait_sec=30)
|
try:
|
||||||
log.debug("Started streams for session {}.".format(self.get_name()))
|
stream_healthy = self.api.stream_healthy()
|
||||||
except Exception as e:
|
if stream_healthy == True:
|
||||||
log.exception("Detected streaming unhealthy in {} session.".format(self.get_name()))
|
self.user_stream = self.api.stream_user(listener, run_async=True, reconnect_async=True, reconnect_async_wait_sec=30)
|
||||||
|
self.direct_stream = self.api.stream_direct(listener, run_async=True, reconnect_async=True, reconnect_async_wait_sec=30)
|
||||||
|
log.debug("Started streams for session {}.".format(self.get_name()))
|
||||||
|
except Exception:
|
||||||
|
log.exception("Detected streaming unhealthy in {} session.".format(self.get_name()))
|
||||||
|
|
||||||
def stop_streaming(self):
|
def stop_streaming(self):
|
||||||
if config.app["app-settings"]["no_streaming"]:
|
with self._stream_lock:
|
||||||
return
|
streams = (self.user_stream, self.direct_stream)
|
||||||
|
self.user_stream = None
|
||||||
|
self.direct_stream = None
|
||||||
|
for stream in streams:
|
||||||
|
if stream is None:
|
||||||
|
continue
|
||||||
|
try:
|
||||||
|
stream.close()
|
||||||
|
except Exception:
|
||||||
|
log.exception("Unable to stop a stream for session {}.".format(self.get_name()))
|
||||||
|
if any(stream is not None for stream in streams):
|
||||||
|
log.debug("Stopped streams for session {}.".format(self.get_name()))
|
||||||
|
|
||||||
def check_streams(self):
|
def check_streams(self):
|
||||||
pass
|
pass
|
||||||
|
|||||||
Reference in New Issue
Block a user