Skip to content

Commit e8d508d

Browse files
committed
SSE
1 parent 8268d10 commit e8d508d

13 files changed

Lines changed: 240 additions & 1 deletion

File tree

Cargo.lock

Lines changed: 23 additions & 0 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

Cargo.toml

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -53,6 +53,7 @@ strum_macros = "0.26.1"
5353
tokio-util = { version = "0.7.11", features = ["compat", "io", "io-util"] }
5454
time = "0.3.34"
5555
tokio-stream = { version = "0.1.14", features = ["io-util"] }
56+
async-stream = "0.3"
5657
socketioxide = { version = "0.13.1", features = ["extensions","state", "v4"] }
5758
async-trait = "0.1.77"
5859
mime_guess = "2.0.4"

src/main.rs

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -172,6 +172,7 @@ async fn app() -> Result<Router> {
172172
.nest("/credentials", routes::credentials::routes(mc.clone()))
173173
.nest("/backups", routes::backups::routes(mc.clone()))
174174
.nest("/plugins", routes::plugins::routes(mc.clone()))
175+
.nest("/sse", routes::sse::routes(mc.clone()))
175176
.fallback(fallback)
176177
.layer(middleware::from_fn(mw_range::mw_range))
177178
//.layer(middleware::map_response(main_response_mapper))

src/model/backups.rs

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@ use tokio_util::io::{ReaderStream, StreamReader};
1919
use crate::{domain::{backup::{self, Backup, BackupError, BackupFile, BackupFileProgress, BackupMessage, BackupProcessStatus, BackupStatus, BackupWithStatus}, library::{LibraryRole, ServerLibrary}, media::{self, Media, MediaForUpdate, DEFAULT_MIME}, progress::{RsProgress, RsProgressType}}, error::{RsError, RsResult}, model::libraries::ServerLibraryForAdd, plugins::sources::{async_reader_progress::ProgressReader, error::SourcesError, AsyncReadPinBox, FileStreamResult, SourceRead}, routes::mw_range::RangeDefinition, tools::{clock::now, encryption::{ceil_to_multiple_of_16, derive_key, estimated_encrypted_size, random_iv, AesTokioDecryptStream, AesTokioEncryptStream}, log::{log_error, log_info}}};
2020

2121
use super::{error::{Error, Result}, medias::{MediaFileQuery, MediaQuery, MediaSource}, store::sql::backups::BackupInfos, users::{ConnectedUser, UserRole}, ModelController};
22+
use crate::routes::sse::SseEvent;
2223

2324
#[derive(Debug, Serialize, Deserialize, Clone)]
2425
pub struct BackupForAdd {
@@ -54,6 +55,7 @@ pub struct BackupForUpdate {
5455
impl ModelController {
5556

5657
pub fn send_backup_status(&self, message: BackupMessage) {
58+
self.broadcast_sse(SseEvent::Backups(message.clone()));
5759
self.for_connected_users(&message, |user, socket, message| {
5860
if let Some(library) = &message.backup.backup.library {
5961
let r = user.check_library_role(library, LibraryRole::Admin);
@@ -69,6 +71,7 @@ impl ModelController {
6971
});
7072
}
7173
pub fn send_backup_file_status(&self, message: BackupFileProgress) {
74+
self.broadcast_sse(SseEvent::BackupsFiles(message.clone()));
7275
self.for_connected_users(&message, |user, socket, message| {
7376
let r = user.check_role(&UserRole::Admin);
7477
if r.is_ok() {

src/model/episodes.rs

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@ use tokio_util::io::StreamReader;
1717
use crate::{domain::{deleted::RsDeleted, episode::{self, Episode, EpisodeWithAction, EpisodeWithShow, EpisodesMessage}, library::LibraryRole, people::{PeopleMessage, Person}, serie::{self, Serie, SeriesMessage}, ElementAction}, error::RsResult, plugins::{medias::imdb::ImdbContext, sources::{error::SourcesError, AsyncReadPinBox, FileStreamResult, Source}}, tools::{array_tools::Dedup, clock::now, image_tools::{resize_image_reader, ImageSize}, log::log_info}};
1818

1919
use super::{error::{Error, Result}, medias::{RsSort, RsSortOrder}, store::sql::SqlOrder, users::{ConnectedUser, HistoryQuery}, ModelController};
20+
use crate::routes::sse::SseEvent;
2021

2122

2223
#[derive(Debug, Serialize, Deserialize, Clone, Default)]
@@ -203,6 +204,7 @@ impl ModelController {
203204

204205

205206
pub fn send_episode(&self, message: EpisodesMessage) {
207+
self.broadcast_sse(SseEvent::Episodes(message.clone()));
206208
self.for_connected_users(&message, |user, socket, message| {
207209
let r = user.check_library_role(&message.library, LibraryRole::Read);
208210
if r.is_ok() {

src/model/medias.rs

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,7 @@ use crate::{domain::{MediaElement, deleted::RsDeleted, library::LibraryType, med
2727
use crate::{domain::{library::LibraryRole, media::{FileType, GroupMediaDownload, Media, MediaDownloadUrl, MediaForAdd, MediaForInsert, MediaForUpdate, MediaItemReference, MediaWithAction, MediasMessage, ProgressMessage}, progress::{RsProgress, RsProgressType}, ElementAction}, error::RsResult, plugins::{get_plugin_fodler, sources::{async_reader_progress::ProgressReader, error::SourcesError, AsyncReadPinBox, FileStreamResult, SourceRead}}, routes::mw_range::RangeDefinition, server::get_server_port, tools::{auth::{sign_local, ClaimsLocal}, file_tools::{file_type_from_mime, get_extension_from_mime}, image_tools::{self, resize_image_reader, ImageSize}, log::{log_error, log_warn, log_info, LogServiceType}, prediction::{predict_net, preload_model, PredictionTagResult}, video_tools::{self, probe_video, VideoTime}}};
2828

2929
use super::{error::{Error, Result}, plugins::PluginQuery, store::{self, sql::library::medias::MediaBackup}, users::ConnectedUser, ModelController, VideoConvertQueueElement};
30+
use crate::routes::sse::SseEvent;
3031

3132
pub const CRYPTO_HEADER_SIZE: u64 = 16 + 4 + 4 + 32 + 256;
3233

@@ -258,6 +259,7 @@ impl ModelController {
258259

259260

260261
pub fn send_media(&self, message: MediasMessage) {
262+
self.broadcast_sse(SseEvent::Medias(message.clone()));
261263
self.for_connected_users(&message, |user, socket, message| {
262264
let r = user.check_library_role(&message.library, LibraryRole::Read);
263265
if r.is_ok() {
@@ -267,6 +269,7 @@ impl ModelController {
267269
}
268270

269271
pub fn send_progress(&self, message: ProgressMessage) {
272+
self.broadcast_sse(SseEvent::MediasProgress(message.clone()));
270273
self.for_connected_users(&message, |user, socket, message| {
271274
let r = user.check_library_role(&message.library, LibraryRole::Read);
272275
if r.is_ok() {
@@ -276,6 +279,7 @@ impl ModelController {
276279
}
277280

278281
pub fn send_convert_progress(&self, message: ConvertMessage) {
282+
self.broadcast_sse(SseEvent::ConvertProgress(message.clone()));
279283
self.for_connected_users(&message, |user, socket, message| {
280284
let r = user.check_library_role(&message.library, LibraryRole::Read);
281285
if r.is_ok() {

src/model/mod.rs

Lines changed: 15 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -73,9 +73,11 @@ use socketioxide::{extract::SocketRef, SocketIo};
7373
use tokio::{
7474
fs::{self, remove_file, File},
7575
io::{copy, AsyncRead, BufReader},
76-
sync::RwLock,
76+
sync::{broadcast, RwLock},
7777
};
7878

79+
use crate::routes::sse::SseEvent;
80+
7981
#[derive(Debug, Serialize, Deserialize, Clone)]
8082
pub struct VideoConvertQueueElement {
8183
request: VideoConvertRequest,
@@ -137,6 +139,9 @@ pub struct ModelController {
137139
pub backup_processes: Arc<RwLock<Vec<BackupProcessStatus>>>,
138140

139141
pub chache_libraries: Arc<RwLock<HashMap<String, ServerLibrary>>>,
142+
143+
/// Broadcast channel for SSE events
144+
pub sse_tx: broadcast::Sender<SseEvent>,
140145
}
141146

142147
// Constructor
@@ -145,6 +150,7 @@ impl ModelController {
145150
let tmdb = TmdbContext::new("4a01db3a73eed5cf17e9c7c27fd9d008".to_string()).await?;
146151
let fanart = FanArtContext::new("a6eb2f1acb7b54550e498a9b37a574fa".to_string());
147152
let scheduler = RsScheduler::new();
153+
let (sse_tx, _) = broadcast::channel::<SseEvent>(1024);
148154

149155
let mc = Self {
150156
store: Arc::new(store),
@@ -164,6 +170,7 @@ impl ModelController {
164170
convert_current_process: Arc::new(RwLock::new(None)),
165171

166172
backup_processes: Arc::new(RwLock::new(vec![])),
173+
sse_tx,
167174
};
168175

169176
let pm_forload = mc.plugin_manager.clone();
@@ -550,7 +557,13 @@ impl ModelController {
550557
}
551558
}
552559

560+
/// Broadcasts an event to all SSE subscribers
561+
pub fn broadcast_sse(&self, event: SseEvent) {
562+
let _ = self.sse_tx.send(event);
563+
}
564+
553565
pub fn send_library(&self, message: LibraryMessage) {
566+
self.broadcast_sse(SseEvent::Library(message.clone()));
554567
self.for_connected_users(&message, |user, socket, message| {
555568
if let Some(message) = message.for_socket(user) {
556569
let _ = socket.emit("library", message);
@@ -559,6 +572,7 @@ impl ModelController {
559572
}
560573

561574
pub fn send_library_status(&self, message: LibraryStatusMessage) {
575+
self.broadcast_sse(SseEvent::LibraryStatus(message.clone()));
562576
self.for_connected_users(&message, |user, socket, message| {
563577
if user
564578
.check_library_role(&message.library, LibraryRole::Admin)

src/model/movies.rs

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@ use tokio_util::io::StreamReader;
1717
use crate::{domain::{deleted::RsDeleted, library::LibraryRole, movie::{Movie, MovieForUpdate, MovieWithAction, MoviesMessage}, people::{PeopleMessage, Person}, ElementAction, MediaElement}, error::RsResult, plugins::{medias::imdb::ImdbContext, sources::{error::SourcesError, path_provider::PathProvider, AsyncReadPinBox, FileStreamResult, Source}}, server::get_server_folder_path_array, tools::{image_tools::{convert_image_reader, resize_image_reader, ImageSize}, log::log_info}};
1818

1919
use super::{error::{Error, Result}, store::sql::SqlOrder, users::{ConnectedUser, HistoryQuery}, ModelController};
20+
use crate::routes::sse::SseEvent;
2021

2122

2223

@@ -237,6 +238,7 @@ impl ModelController {
237238

238239

239240
pub fn send_movie(&self, message: MoviesMessage) {
241+
self.broadcast_sse(SseEvent::Movies(message.clone()));
240242
self.for_connected_users(&message, |user, socket, message| {
241243
let r = user.check_library_role(&message.library, LibraryRole::Read);
242244
if r.is_ok() {

src/model/people.rs

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -48,6 +48,7 @@ use super::{
4848
users::ConnectedUser,
4949
ModelController,
5050
};
51+
use crate::routes::sse::SseEvent;
5152

5253
// Default face recognition threshold when library doesn't specify one
5354
pub const DEFAULT_FACE_THRESHOLD: f32 = 0.4;
@@ -273,6 +274,7 @@ impl ModelController {
273274
}
274275

275276
pub fn send_people(&self, message: PeopleMessage) {
277+
self.broadcast_sse(SseEvent::People(message.clone()));
276278
self.for_connected_users(&message, |user, socket, message| {
277279
let r = user.check_library_role(&message.library, LibraryRole::Read);
278280
if r.is_ok() {

src/model/series.rs

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@ use tokio_util::io::StreamReader;
1717
use crate::{domain::{deleted::RsDeleted, library::LibraryRole, people::{PeopleMessage, Person}, serie::{Serie, SerieStatus, SerieWithAction, SeriesMessage}, ElementAction, MediaElement}, error::RsResult, plugins::{medias::imdb::ImdbContext, sources::{error::SourcesError, path_provider::PathProvider, AsyncReadPinBox, FileStreamResult, Source}}, server::get_server_folder_path_array, tools::{image_tools::{convert_image_reader, resize_image_reader, ImageSize}, log::log_info}};
1818

1919
use super::{episodes::{EpisodeForUpdate, EpisodeQuery}, error::{Error, Result}, medias::{MediaQuery, RsSort}, store::sql::SqlOrder, users::ConnectedUser, ModelController};
20+
use crate::routes::sse::SseEvent;
2021

2122

2223
impl FromSql for SerieStatus {
@@ -206,6 +207,7 @@ impl ModelController {
206207

207208

208209
pub fn send_serie(&self, message: SeriesMessage) {
210+
self.broadcast_sse(SseEvent::Series(message.clone()));
209211
self.for_connected_users(&message, |user, socket, message| {
210212
let r = user.check_library_role(&message.library, LibraryRole::Read);
211213
if r.is_ok() {

0 commit comments

Comments
 (0)