1use super::{Request, RequestBuilderPerform};
2use crate::settings::path_filter::PathFilter;
3use crate::settings::rewrite::Rewrite;
4use crate::settings::targets::TargetProcess;
5use anyhow::Context;
6use autopulse_database::models::ScanEvent;
7use autopulse_utils::{get_url, RuntimePath};
8use reqwest::header;
9use serde::{Deserialize, Serialize};
10use std::{collections::HashMap, fmt::Display, io::Cursor};
11use struson::{
12 json_path,
13 reader::{JsonReader, JsonStreamReader},
14};
15use tokio::sync::mpsc::UnboundedReceiver;
16use tracing::{debug, error};
17
18#[doc(hidden)]
19const fn default_true() -> bool {
20 true
21}
22
23#[derive(Serialize, Clone, Deserialize)]
24pub struct Emby {
25 pub url: String,
27 pub token: String,
29 #[serde(default)]
31 pub metadata_refresh_mode: EmbyMetadataRefreshMode,
32 #[serde(default = "default_true")]
34 pub refresh_metadata: bool,
35 pub rewrite: Option<Rewrite>,
37 #[serde(default)]
39 pub filter: PathFilter,
40 #[serde(default)]
42 pub request: Request,
43 #[serde(default)]
45 pub path_match: PathMatch,
46}
47
48#[derive(Serialize, Clone, Deserialize, Default)]
50#[serde(rename_all = "snake_case")]
51pub enum PathMatch {
52 #[default]
54 CaseSensitive,
55 CaseInsensitive,
57}
58
59#[derive(Serialize, Clone, Deserialize)]
61#[serde(rename_all = "snake_case")]
62#[derive(Default)]
63pub enum EmbyMetadataRefreshMode {
64 None,
66 ValidationOnly,
68 Default,
70 #[default]
72 FullRefresh,
73}
74
75impl Display for EmbyMetadataRefreshMode {
76 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
77 let mode = match self {
78 Self::None => "None",
79 Self::ValidationOnly => "ValidationOnly",
80 Self::Default => "Default",
81 Self::FullRefresh => "FullRefresh",
82 };
83
84 write!(f, "{mode}")
85 }
86}
87
88#[derive(Deserialize, Clone, Eq, PartialEq, Hash)]
89#[serde(rename_all = "PascalCase")]
90#[doc(hidden)]
91struct Library {
92 #[allow(dead_code)]
93 name: String,
94 locations: Vec<String>,
95 item_id: String,
96 collection_type: Option<String>,
97}
98
99#[derive(Serialize, Clone)]
100#[serde(rename_all = "PascalCase")]
101#[doc(hidden)]
102struct UpdateRequest {
103 path: String,
104 update_type: String,
105}
106
107#[derive(Serialize, Clone)]
108#[serde(rename_all = "PascalCase")]
109#[doc(hidden)]
110struct ScanPayload {
111 updates: Vec<UpdateRequest>,
112}
113
114#[derive(Deserialize, Clone)]
115#[serde(rename_all = "PascalCase")]
116#[doc(hidden)]
117struct Item {
118 id: String,
119 path: Option<String>,
120}
121
122impl Emby {
123 fn get_client(&self) -> anyhow::Result<reqwest::Client> {
124 let mut headers = header::HeaderMap::new();
125
126 headers.insert("X-Emby-Token", self.token.parse()?);
127 headers.insert(
128 "Authorization",
129 format!("MediaBrowser Token=\"{}\"", self.token).parse()?,
130 );
131 headers.insert("Accept", "application/json".parse()?);
132
133 self.request
134 .client_builder(headers)
135 .build()
136 .map_err(Into::into)
137 }
138
139 async fn libraries(&self) -> anyhow::Result<Vec<Library>> {
140 let client = self.get_client()?;
141 let url = get_url(&self.url)?.join("Library/VirtualFolders")?;
142
143 let res = client.get(url).perform().await?;
144
145 Ok(res.json().await?)
146 }
147
148 fn get_libraries(&self, libraries: &[Library], path: &str) -> Vec<Library> {
149 let mut matched: Vec<Library> = vec![];
150
151 for library in libraries {
152 for location in &library.locations {
153 if self.path_prefix_matches(location, path) {
154 matched.push(library.clone());
155 break; }
157 }
158 }
159
160 matched
161 }
162
163 fn path_prefix_matches(&self, location: &str, ev_path: &str) -> bool {
164 let location = RuntimePath::new(location);
165 let event = RuntimePath::new(ev_path);
166
167 match self.path_match {
168 PathMatch::CaseSensitive => event.starts_with_case_sensitive(location),
169 PathMatch::CaseInsensitive => event.starts_with_ascii_case_insensitive(location),
170 }
171 }
172
173 async fn _get_item(&self, library: &Library, path: &str) -> anyhow::Result<Option<Item>> {
174 let client = self.get_client()?;
175 let mut url = get_url(&self.url)?.join("Items")?;
176
177 url.query_pairs_mut().append_pair("Recursive", "true");
178 url.query_pairs_mut().append_pair("Fields", "Path");
179 url.query_pairs_mut().append_pair("EnableImages", "false");
180 if let Some(collection_type) = &library.collection_type {
181 url.query_pairs_mut().append_pair(
182 "IncludeItemTypes",
183 match collection_type.as_str() {
184 "tvshows" => "Episode",
185 "books" => "Book",
186 "music" => "Audio",
187 "movie" => "VideoFile,Movie",
188 _ => "",
189 },
190 );
191 }
192 url.query_pairs_mut()
193 .append_pair("ParentId", &library.item_id);
194 url.query_pairs_mut()
195 .append_pair("EnableTotalRecordCount", "false");
196
197 let res = client.get(url).perform().await?;
198
199 let bytes = res.bytes().await?;
201
202 let mut json_reader = JsonStreamReader::new(Cursor::new(bytes));
203
204 json_reader.seek_to(&json_path!["Items"])?;
205 json_reader.begin_array()?;
206
207 while json_reader.has_next()? {
208 let item: Item = json_reader.deserialize_next()?;
209
210 if item.path == Some(path.to_owned()) {
211 return Ok(Some(item));
212 }
213 }
214
215 Ok(None)
216 }
217
218 fn fetch_items(
219 &self,
220 library: &Library,
221 ) -> anyhow::Result<(
222 UnboundedReceiver<Item>,
223 tokio::task::JoinHandle<anyhow::Result<()>>,
224 )> {
225 let (tx, rx) = tokio::sync::mpsc::unbounded_channel();
226 let limit = 1000;
227
228 let client = self.get_client()?;
229 let mut url = get_url(&self.url)?.join("Items")?;
230
231 url.query_pairs_mut().append_pair("Recursive", "true");
232 url.query_pairs_mut().append_pair("Fields", "Path");
233 url.query_pairs_mut().append_pair("EnableImages", "false");
234 url.query_pairs_mut()
235 .append_pair("ParentId", &library.item_id);
236 url.query_pairs_mut()
237 .append_pair("EnableTotalRecordCount", "false");
238 url.query_pairs_mut()
239 .append_pair("Limit", &limit.to_string());
240 if let Some(collection_type) = &library.collection_type {
241 url.query_pairs_mut().append_pair(
242 "IncludeItemTypes",
243 match collection_type.as_str() {
244 "tvshows" => "Episode",
245 "books" => "Book",
246 "music" => "Audio",
247 "movie" => "VideoFile,Movie",
248 _ => "",
249 },
250 );
251 }
252
253 let handle = tokio::spawn(async move {
254 let mut page = 0;
255
256 loop {
257 let mut page_url = url.clone();
258 page_url
259 .query_pairs_mut()
260 .append_pair("StartIndex", &(page * limit).to_string());
261
262 let res = client.get(page_url).perform().await?;
263
264 let bytes = res.bytes().await?;
265
266 let mut json_reader = JsonStreamReader::new(Cursor::new(bytes));
267
268 json_reader.seek_to(&json_path!["Items"])?;
269 json_reader.begin_array()?;
270
271 let mut found_items_count = 0;
272
273 while json_reader.has_next()? {
274 let item: Item = json_reader.deserialize_next()?;
275
276 tx.send(item)?;
277
278 found_items_count += 1;
279 }
280
281 if found_items_count < limit {
282 break;
283 }
284
285 page += 1;
286 }
287
288 drop(tx);
289
290 Ok(())
291 });
292
293 Ok((rx, handle))
294 }
295
296 async fn get_items<'a>(
297 &self,
298 library: &Library,
299 events: Vec<&'a ScanEvent>,
300 ) -> anyhow::Result<(Vec<(&'a ScanEvent, Item)>, Vec<&'a ScanEvent>)> {
301 let (mut rx, handle) = self.fetch_items(library)?;
302
303 let mut found_in_library = Vec::new();
304 let mut not_found_in_library = events.clone();
305
306 while let Some(item) = rx.recv().await {
307 if let Some(ev) = events
308 .iter()
309 .find(|ev| item.path == Some(ev.get_path(&self.rewrite)))
310 {
311 found_in_library.push((*ev, item.clone()));
312 not_found_in_library.retain(|&e| e.id != ev.id);
313
314 if not_found_in_library.is_empty() {
315 break;
316 }
317 }
318 }
319
320 handle.abort();
321
322 Ok((found_in_library, not_found_in_library))
323 }
324
325 async fn scan(&self, ev: &[&ScanEvent]) -> anyhow::Result<()> {
327 let client = self.get_client()?;
328 let url = get_url(&self.url)?.join("Library/Media/Updated")?;
329
330 let updates = ev
331 .iter()
332 .map(|ev| UpdateRequest {
333 path: ev.get_path(&self.rewrite),
334 update_type: "Modified".to_string(),
335 })
336 .collect();
337
338 let body = ScanPayload { updates };
339
340 client
341 .post(url)
342 .header("Content-Type", "application/json")
343 .json(&body)
344 .perform()
345 .await
346 .map(|_| ())
347 }
348
349 async fn refresh_item(&self, item: &Item) -> anyhow::Result<()> {
350 let client = self.get_client()?;
351 let mut url = get_url(&self.url)?.join(&format!("Items/{}/Refresh", item.id))?;
352
353 url.query_pairs_mut().append_pair(
354 "MetadataRefreshMode",
355 &self.metadata_refresh_mode.to_string(),
356 );
357 url.query_pairs_mut()
358 .append_pair("ImageRefreshMode", &self.metadata_refresh_mode.to_string());
359 url.query_pairs_mut()
360 .append_pair("ReplaceAllMetadata", "true");
361 url.query_pairs_mut().append_pair("Recursive", "true");
362
363 url.query_pairs_mut()
365 .append_pair("ReplaceAllImages", "false");
366 url.query_pairs_mut()
367 .append_pair("RegenerateTrickplay", "false");
368
369 client.post(url).perform().await.map(|_| ())
370 }
371}
372
373impl TargetProcess for Emby {
374 async fn process(&self, evs: &[&ScanEvent]) -> anyhow::Result<Vec<String>> {
375 let libraries = self
376 .libraries()
377 .await
378 .context("failed to fetch libraries")?;
379
380 let mut succeeded: HashMap<String, bool> = HashMap::new();
381
382 let mut to_find = HashMap::new();
383 let mut to_refresh = Vec::new();
384 let mut to_scan = Vec::new();
385
386 if self.refresh_metadata {
387 for ev in evs {
388 let ev_path = ev.get_path(&self.rewrite);
389
390 let matched_libraries = self.get_libraries(&libraries, &ev_path);
391
392 if matched_libraries.is_empty() {
393 let known: Vec<&str> = libraries
394 .iter()
395 .flat_map(|l| l.locations.iter().map(String::as_str))
396 .collect();
397 error!(
398 "failed to find library for file '{ev_path}'. Known locations: {known:?}"
399 );
400 continue;
401 }
402
403 for library in matched_libraries {
404 to_find.entry(library).or_insert_with(Vec::new).push(*ev);
405 }
406 }
407
408 for (library, library_events) in to_find {
409 let (found_in_library, not_found_in_library) = self
410 .get_items(&library, library_events.clone())
411 .await
412 .with_context(|| {
413 format!(
414 "failed to fetch items for library: {}",
415 library.name.clone()
416 )
417 })?;
418
419 to_refresh.extend(found_in_library);
420 to_scan.extend(not_found_in_library);
421 }
422
423 for (ev, item) in to_refresh {
424 match self.refresh_item(&item).await {
425 Ok(()) => {
426 debug!("refreshed item: {}", item.id);
427 *succeeded.entry(ev.id.clone()).or_insert(true) &= true;
428 }
429 Err(e) => {
430 error!("failed to refresh item: {}", e);
431 succeeded.insert(ev.id.clone(), false);
432 }
433 }
434 }
435 } else {
436 to_scan.extend(evs.iter().copied());
437 }
438
439 if !to_scan.is_empty() {
440 match self.scan(&to_scan).await {
441 Ok(()) => {
442 for ev in &to_scan {
443 debug!("scanned file: {}", ev.file_path);
444
445 *succeeded.entry(ev.id.clone()).or_insert(true) &= true;
446 }
447 }
448 Err(e) => {
449 error!("failed to scan items: {}", e);
450
451 for ev in &to_scan {
452 succeeded.insert(ev.id.clone(), false);
453 }
454 }
455 }
456 }
457
458 Ok(succeeded
459 .iter()
460 .filter_map(|(k, v)| if *v { Some(k.clone()) } else { None })
461 .collect())
462 }
463}
464
465#[cfg(test)]
466mod tests {
467 use super::*;
468
469 fn lib(name: &str, paths: &[&str]) -> Library {
470 Library {
471 name: name.to_string(),
472 locations: paths.iter().map(|s| (*s).to_string()).collect(),
473 item_id: format!("id-{name}"),
474 collection_type: None,
475 }
476 }
477
478 fn target() -> Emby {
479 Emby {
480 url: "http://x".to_string(),
481 token: "t".to_string(),
482 metadata_refresh_mode: EmbyMetadataRefreshMode::default(),
483 refresh_metadata: true,
484 rewrite: None,
485 filter: PathFilter::default(),
486 request: Request::default(),
487 path_match: PathMatch::default(),
488 }
489 }
490
491 #[test]
492 fn matches_when_library_location_is_exact_prefix() {
493 let libs = vec![lib("TV", &["/media/TV"])];
494 let m = target().get_libraries(&libs, "/media/TV/Show/S01E01.mkv");
495 assert_eq!(m.len(), 1);
496 }
497
498 #[test]
499 fn matches_when_library_location_has_trailing_slash() {
500 let libs = vec![lib("TV", &["/media/TV/"])];
501 let m = target().get_libraries(&libs, "/media/TV/Show/S01E01.mkv");
502 assert_eq!(
503 m.len(),
504 1,
505 "trailing slash on library location must not break match"
506 );
507 }
508
509 #[test]
510 fn no_match_when_path_outside_library() {
511 let libs = vec![lib("TV", &["/data/TV"])];
512 let m = target().get_libraries(&libs, "/media/TV/Show.mkv");
513 assert_eq!(m.len(), 0);
514 }
515
516 #[test]
517 fn case_insensitive_matching_when_opted_in() {
518 let mut t = target();
519 t.path_match = PathMatch::CaseInsensitive;
520 let libs = vec![lib("TV", &["/Media/TV"])];
521 let m = t.get_libraries(&libs, "/media/tv/Show.mkv");
522 assert_eq!(
523 m.len(),
524 1,
525 "case-insensitive must match across case differences"
526 );
527 }
528
529 #[test]
530 fn case_sensitive_matching_by_default() {
531 let libs = vec![lib("TV", &["/Media/TV"])];
532 let m = target().get_libraries(&libs, "/media/tv/Show.mkv");
533 assert_eq!(m.len(), 0);
534 }
535
536 #[test]
537 fn case_sensitive_mode_parses_windows_paths_but_preserves_case_policy() {
538 let libs = vec![lib("TV", &[r"C:\Media\TV\"])];
539 assert_eq!(
540 target()
541 .get_libraries(&libs, r"C:\Media\TV\Show\S01E01.mkv")
542 .len(),
543 1
544 );
545 assert_eq!(
546 target()
547 .get_libraries(&libs, r"c:\media\tv\Show\S01E01.mkv")
548 .len(),
549 0
550 );
551 }
552
553 #[test]
554 fn case_insensitive_mode_folds_windows_ascii_case() {
555 let mut emby = target();
556 emby.path_match = PathMatch::CaseInsensitive;
557 let libs = vec![lib("TV", &[r"\\SERVER\MEDIA\TV"])];
558
559 assert_eq!(
560 emby.get_libraries(&libs, r"\\server\media\tv\Show\S01E01.mkv")
561 .len(),
562 1
563 );
564 }
565
566 #[test]
567 fn library_prefix_requires_a_component_boundary() {
568 let libs = vec![lib("TV", &[r"\\server\media"])];
569 assert!(target()
570 .get_libraries(&libs, r"\\server\media-archive\Film.mkv")
571 .is_empty());
572 }
573
574 #[test]
575 fn library_prefix_does_not_cross_runtime_flavors() {
576 let mut emby = target();
577 emby.path_match = PathMatch::CaseInsensitive;
578 let libs = vec![lib("TV", &["/media/TV"])];
579 assert!(emby
580 .get_libraries(&libs, r"C:\media\TV\Show\S01E01.mkv")
581 .is_empty());
582 }
583}