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, HashSet};
11use tracing::{debug, error, trace, warn};
12
13#[derive(Serialize, Deserialize, Clone)]
14pub struct Plex {
15 pub url: String,
17 pub token: String,
19 #[serde(default)]
21 pub refresh: bool,
22 #[serde(default)]
24 pub analyze: bool,
25 #[serde(default)]
28 pub empty_trash: bool,
29 pub rewrite: Option<Rewrite>,
31 #[serde(default)]
33 pub filter: PathFilter,
34 #[serde(default)]
36 pub request: Request,
37}
38
39#[derive(Deserialize, Clone, Debug)]
40#[serde(rename_all = "camelCase")]
41pub struct Media {
42 #[serde(rename = "Part")]
43 pub part: Vec<Part>,
44}
45
46#[derive(Deserialize, Clone, Debug)]
47#[serde(rename_all = "camelCase")]
48pub struct Part {
49 pub key: String,
51 pub file: String,
53 }
61
62#[derive(Deserialize, Clone, Debug)]
63#[serde(rename_all = "camelCase")]
64pub struct Metadata {
65 pub key: String,
66 #[serde(rename = "Media")]
67 pub media: Option<Vec<Media>>,
68 #[serde(rename = "type")]
69 pub t: String,
70}
71
72#[doc(hidden)]
73#[derive(Deserialize, Clone, Debug)]
74struct Location {
75 path: String,
76}
77
78#[doc(hidden)]
79#[derive(Deserialize, Clone, Debug)]
80struct Library {
81 title: String,
82 key: String,
83 refreshing: Option<bool>,
84 #[serde(rename = "scannedAt")]
85 scanned_at: Option<u64>,
86 #[serde(rename = "Location")]
87 location: Vec<Location>,
88}
89
90struct LibraryCleanup {
91 scanned_at: Option<u64>,
92 event_ids: HashSet<String>,
93 scan_failed: bool,
94}
95
96#[doc(hidden)]
97#[derive(Deserialize, Clone)]
98#[serde(rename_all = "PascalCase")]
99struct LibraryMediaContainer {
100 directory: Option<Vec<Library>>,
101 metadata: Option<Vec<Metadata>>,
102}
103
104#[doc(hidden)]
105#[derive(Deserialize, Clone)]
106#[serde(rename_all = "PascalCase")]
107struct SearchResult {
108 metadata: Option<Metadata>,
109}
110
111#[doc(hidden)]
112#[derive(Deserialize, Clone)]
113#[serde(rename_all = "PascalCase")]
114struct SearchLibraryMediaContainer {
115 #[serde(default)]
116 search_result: Vec<SearchResult>,
117}
118
119#[doc(hidden)]
120#[derive(Deserialize, Clone)]
121#[serde(rename_all = "PascalCase")]
122struct LibraryResponse {
123 media_container: LibraryMediaContainer,
124}
125
126#[doc(hidden)]
127#[derive(Deserialize, Clone)]
128#[serde(rename_all = "PascalCase")]
129struct SearchLibraryResponse {
130 media_container: SearchLibraryMediaContainer,
131}
132
133fn path_matches(part_file: &str, path: &str) -> bool {
134 let part_file = RuntimePath::new(part_file);
135 let path = RuntimePath::new(path);
136
137 if path.is_directory() {
138 part_file.starts_with(path)
139 } else {
140 part_file.equals(path)
141 }
142}
143
144fn has_matching_media(media: &[Media], path: &str) -> bool {
145 media.iter().any(|media_item| {
146 media_item
147 .part
148 .iter()
149 .any(|part| path_matches(&part.file, path))
150 })
151}
152
153fn scan_directory(path: &str) -> &str {
154 RuntimePath::new(path).parent_or_self().as_str()
155}
156
157impl Plex {
158 fn get_client(&self) -> anyhow::Result<reqwest::Client> {
159 let mut headers = header::HeaderMap::new();
160
161 headers.insert("X-Plex-Token", self.token.parse()?);
162 headers.insert("Accept", "application/json".parse()?);
163
164 self.request
165 .client_builder(headers)
166 .build()
167 .map_err(Into::into)
168 }
169
170 async fn libraries(&self) -> anyhow::Result<Vec<Library>> {
171 let client = self.get_client()?;
172 let url = get_url(&self.url)?.join("library/sections")?;
173
174 let res = client.get(url).perform().await?;
175
176 let libraries: LibraryResponse = res.json().await?;
177
178 Ok(libraries.media_container.directory.unwrap_or_default())
179 }
180
181 fn get_libraries(&self, libraries: &[Library], path: &str) -> Vec<Library> {
182 let event_path = RuntimePath::new(path);
183 let mut matches: Vec<(usize, &Library)> = vec![];
184
185 for library in libraries {
186 for location in &library.location {
187 let location_path = RuntimePath::new(&location.path);
188 if event_path.starts_with(location_path) {
189 matches.push((location_path.component_count(), library));
190 }
191 }
192 }
193
194 matches.sort_by(|(components_a, _), (components_b, _)| components_b.cmp(components_a));
196
197 matches
198 .into_iter()
199 .map(|(_, library)| library.clone())
200 .collect()
201 }
202
203 async fn get_episodes(&self, key: &str) -> anyhow::Result<LibraryResponse> {
204 let client = self.get_client()?;
205
206 let key = key.rsplit_once('/').map(|x| x.0).unwrap_or(key);
208
209 let url = get_url(&self.url)?.join(&format!("{key}/allLeaves"))?;
210
211 let res = client.get(url).perform().await?;
212
213 let lib: LibraryResponse = res.json().await?;
214
215 Ok(lib)
216 }
217
218 fn get_search_term(&self, path: &str) -> anyhow::Result<String> {
219 let parent_or_directory = RuntimePath::new(path).parent_or_self();
220 let components = parent_or_directory.normal_components().collect::<Vec<_>>();
221
222 let chosen = components
223 .iter()
224 .rev()
225 .copied()
226 .find(|component| !component.contains("Season") && !component.is_empty())
227 .map(ToString::to_string)
228 .unwrap_or_else(|| {
229 components.join(" ")
232 });
233
234 Ok(chosen
235 .split_whitespace()
236 .filter(|part| {
237 ["(", ")", "[", "]", "{", "}"]
238 .iter()
239 .all(|character| !part.contains(character))
240 })
241 .collect::<Vec<_>>()
242 .join(" "))
243 }
244
245 async fn search_items(&self, _library: &Library, path: &str) -> anyhow::Result<Vec<Metadata>> {
246 let client = self.get_client()?;
247 let mut results = vec![];
250
251 let rel_path = path.to_string();
252
253 trace!("searching for item with relative path: {}", rel_path);
254
255 let mut search_term = self.get_search_term(&rel_path)?;
256
257 while !search_term.is_empty() {
258 let mut url = get_url(&self.url)?.join("library/search")?;
259
260 url.query_pairs_mut().append_pair("includeCollections", "1");
261 url.query_pairs_mut()
262 .append_pair("includeExternalMedia", "1");
263 url.query_pairs_mut()
264 .append_pair("searchTypes", "movies,people,tv");
265 url.query_pairs_mut().append_pair("limit", "100");
266
267 trace!("searching for item with term: {}", search_term);
268
269 url.query_pairs_mut()
270 .append_pair("query", search_term.as_str());
272
273 let res = client.get(url).perform().await?;
274
275 let lib: SearchLibraryResponse = res.json().await?;
276
277 let mut metadata = lib
278 .media_container
279 .search_result
280 .into_iter()
281 .filter_map(|s| s.metadata)
282 .collect::<Vec<_>>();
283
284 metadata.sort_by(|a, b| {
286 if a.t == "episode" && b.t != "episode" {
287 std::cmp::Ordering::Less
288 } else if a.t != "episode" && b.t == "episode" {
289 std::cmp::Ordering::Greater
290 } else if a.t == "movie" && b.t != "movie" && b.t != "episode" {
291 std::cmp::Ordering::Less
292 } else if a.t != "movie" && a.t != "episode" && b.t == "movie" {
293 std::cmp::Ordering::Greater
294 } else {
295 std::cmp::Ordering::Equal
296 }
297 });
298
299 for item in &metadata {
300 if item.t == "show" {
301 let episodes = self.get_episodes(&item.key).await?;
302
303 if let Some(episode_metadata) = episodes.media_container.metadata {
304 for episode in episode_metadata {
305 if let Some(media) = &episode.media {
306 if has_matching_media(media, path) {
307 results.push(episode.clone());
308 }
309 }
310 }
311 }
312 } else if let Some(media) = &item.media {
313 if has_matching_media(media, path) {
315 results.push(item.clone());
316 }
317 }
318 }
319
320 trace!(
321 "found {} out of {} items matching search",
322 results.len(),
323 metadata.len()
324 );
325
326 if results.is_empty() {
327 let mut search_parts = search_term.split_whitespace().collect::<Vec<_>>();
328 search_parts.pop();
329 search_term = search_parts.join(" ");
330 } else {
331 break;
332 }
333 }
334
335 results.dedup_by_key(|item| item.key.clone());
337
338 Ok(results)
339 }
340
341 async fn _get_items(&self, library: &Library, path: &str) -> anyhow::Result<Vec<Metadata>> {
342 let client = self.get_client()?;
343 let url = get_url(&self.url)?.join(&format!("library/sections/{}/all", library.key))?;
344
345 let res = client.get(url).perform().await?;
346
347 let lib: LibraryResponse = res.json().await?;
348
349 let mut parts = vec![];
350
351 for item in lib.media_container.metadata.unwrap_or_default() {
353 match item.t.as_str() {
354 "show" => {
355 let episodes = self.get_episodes(&item.key).await?;
356
357 for episode in episodes.media_container.metadata.unwrap_or_default() {
358 if let Some(media) = &episode.media {
359 if has_matching_media(media, path) {
360 parts.push(episode.clone());
361 }
362 }
363 }
364 }
365 _ => {
366 if let Some(media) = &item.media {
367 if has_matching_media(media, path) {
368 parts.push(item.clone());
369 }
370 }
371 }
372 }
373 }
374
375 Ok(parts)
376 }
377
378 async fn refresh_item(&self, key: &str) -> anyhow::Result<()> {
379 let client = self.get_client()?;
380 let url = get_url(&self.url)?.join(&format!("{key}/refresh"))?;
381
382 client.put(url).perform().await.map(|_| ())
383 }
384
385 async fn analyze_item(&self, key: &str) -> anyhow::Result<()> {
386 let client = self.get_client()?;
387 let url = get_url(&self.url)?.join(&format!("{key}/analyze"))?;
388
389 client.put(url).perform().await.map(|_| ())
390 }
391
392 async fn scan(&self, ev: &ScanEvent, library: &Library) -> anyhow::Result<()> {
393 let client = self.get_client()?;
394 let mut url =
395 get_url(&self.url)?.join(&format!("library/sections/{}/refresh", library.key))?;
396
397 let ev_path = ev.get_path(&self.rewrite);
398 url.query_pairs_mut()
399 .append_pair("path", scan_directory(&ev_path));
400
401 client.get(url).perform().await.map(|_| ())
402 }
403
404 async fn empty_library_trash(
405 &self,
406 key: &str,
407 scanned_at: Option<u64>,
408 ) -> anyhow::Result<bool> {
409 let scan_finished = tokio::time::timeout(std::time::Duration::from_secs(300), async {
410 let started_at = tokio::time::Instant::now();
411 let mut observed_scan = false;
412 loop {
413 tokio::time::sleep(std::time::Duration::from_secs(1)).await;
414 let library = self
415 .libraries()
416 .await?
417 .into_iter()
418 .find(|library| library.key == key)
419 .context("scanned library no longer exists")?;
420 match library.refreshing {
421 Some(true) => observed_scan = true,
422 Some(false) => {
423 let scan_finished = scanned_at
426 .zip(library.scanned_at)
427 .is_some_and(|(before, after)| after > before);
428 if observed_scan || scan_finished {
429 return Ok::<_, anyhow::Error>(true);
430 }
431 }
432 None => anyhow::bail!("Plex did not report the library scan status"),
433 }
434 if !observed_scan && started_at.elapsed() >= std::time::Duration::from_secs(5) {
435 return Ok(false);
436 }
437 }
438 })
439 .await
440 .context("timed out waiting for Plex library scan to finish")??;
441
442 if !scan_finished {
443 return Ok(false);
444 }
445
446 let client = self.get_client()?;
447 let url = get_url(&self.url)?.join(&format!("library/sections/{key}/emptyTrash"))?;
448 client.put(url).perform().await.map(|_| true)
449 }
450}
451
452impl TargetProcess for Plex {
453 async fn process(&self, evs: &[&ScanEvent]) -> anyhow::Result<Vec<String>> {
454 let libraries = self.libraries().await.context("failed to get libraries")?;
455
456 let mut succeeded: HashMap<String, bool> = HashMap::new();
457 let mut cleanups: HashMap<String, LibraryCleanup> = HashMap::new();
458
459 for ev in evs {
460 let succeeded_entry = succeeded.entry(ev.id.clone()).or_insert(true);
461
462 let ev_path = ev.get_path(&self.rewrite);
463 let matched_libraries = self.get_libraries(&libraries, &ev_path);
464
465 if matched_libraries.is_empty() {
466 error!("no matching library for {ev_path}");
467
468 *succeeded_entry = false;
469
470 continue;
471 }
472
473 let mut processed_items = HashSet::new();
474
475 for library in matched_libraries {
476 trace!("found library '{}' for {ev_path}", library.title);
477
478 let scan_result = self.scan(ev, &library).await;
479 if self.empty_trash {
480 let cleanup =
481 cleanups
482 .entry(library.key.clone())
483 .or_insert_with(|| LibraryCleanup {
484 scanned_at: library.scanned_at,
485 event_ids: HashSet::new(),
486 scan_failed: false,
487 });
488 cleanup.event_ids.insert(ev.id.clone());
489 cleanup.scan_failed |= scan_result.is_err();
490 }
491
492 match scan_result {
493 Ok(()) => {
494 debug!("scanned '{}'", ev_path);
495
496 if self.analyze || self.refresh {
497 match self.search_items(&library, &ev_path).await {
498 Ok(items) => {
499 if items.is_empty() {
500 trace!(
501 "failed to find items for file: '{}', leaving at scan",
502 ev_path
503 );
504
505 } else {
507 trace!("found items for file '{}'", ev_path);
508
509 let mut all_success = true;
510
511 for item in items {
512 let mut item_success = true;
513
514 if processed_items.contains(&item.key) {
515 debug!(
516 "already processed item '{}' earlier, skipping",
517 item.key
518 );
519 continue;
520 }
521
522 if self.refresh {
523 match self.refresh_item(&item.key).await {
524 Ok(()) => {
525 debug!("refreshed metadata '{}'", item.key);
526 }
527 Err(e) => {
528 error!(
529 "failed to refresh metadata for '{}': {}",
530 item.key, e
531 );
532 item_success = false;
533 }
534 }
535 }
536
537 if self.analyze {
538 match self.analyze_item(&item.key).await {
539 Ok(()) => {
540 debug!("analyzed metadata '{}'", item.key);
541 }
542 Err(e) => {
543 error!(
544 "failed to analyze metadata for '{}': {}",
545 item.key, e
546 );
547 item_success = false;
548 }
549 }
550 }
551
552 if !item_success {
553 all_success = false;
554 }
555
556 processed_items.insert(item.key);
557 }
558
559 if !all_success {
560 *succeeded_entry = false;
561 }
562 }
563 }
564 Err(e) => {
565 error!("failed to get items for '{}': {:?}", ev_path, e);
566 *succeeded_entry = false;
567 }
568 };
569 }
570 }
571 Err(e) => {
572 error!("failed to scan file '{}': {}", ev_path, e);
573 *succeeded_entry = false;
574 }
575 }
576 }
577 }
578
579 for (key, cleanup) in cleanups {
580 let result = if cleanup.scan_failed {
581 Err(anyhow::anyhow!("a scan failed for this library"))
582 } else {
583 self.empty_library_trash(&key, cleanup.scanned_at).await
584 };
585
586 match result {
587 Ok(true) => debug!("emptied trash for library '{key}'"),
588 Ok(false) => warn!("skipped trash for library '{key}': Plex did not report a scan"),
589 Err(e) => {
590 error!("failed to empty trash for library '{key}': {e:#}");
591 for id in cleanup.event_ids {
592 succeeded.insert(id, false);
593 }
594 }
595 }
596 }
597
598 Ok(succeeded
599 .into_iter()
600 .filter_map(|(k, v)| if v { Some(k) } else { None })
601 .collect())
602 }
603}
604
605#[cfg(test)]
606mod tests {
607 use super::*;
608
609 fn test_plex() -> Plex {
610 Plex {
611 url: String::new(),
612 token: String::new(),
613 refresh: false,
614 analyze: false,
615 empty_trash: false,
616 rewrite: None,
617 filter: PathFilter::default(),
618 request: Request::default(),
619 }
620 }
621
622 #[test]
623 fn test_get_search_term() {
624 let plex = test_plex();
625
626 let path = "/media/TV Shows/Breaking Bad/Season 1/S01E01.mkv";
628 assert_eq!(plex.get_search_term(path).unwrap(), "Breaking Bad");
629
630 let path = "/media/Movies/The Matrix (1999) [1080p]/matrix.mkv";
632 assert_eq!(plex.get_search_term(path).unwrap(), "The Matrix");
633
634 let path = "/media/Movies/Inception/inception.mkv";
636 assert_eq!(plex.get_search_term(path).unwrap(), "Inception");
637
638 let path = "/media/TV Shows/Game of Thrones/Season 2";
640 assert_eq!(plex.get_search_term(path).unwrap(), "Game of Thrones");
641
642 let path = "/media/TV Shows/Game of Thrones";
644 assert_eq!(plex.get_search_term(path).unwrap(), "Game of Thrones");
645
646 let path = "/media/TV Shows/Doctor Who/Season 10/Season 10 Part 2/S10E12.mkv";
648 assert_eq!(plex.get_search_term(path).unwrap(), "Doctor Who");
649 }
650
651 #[test]
652 fn test_get_library() {
653 let plex = Plex {
654 url: String::new(),
655 token: String::new(),
656 refresh: false,
657 analyze: false,
658 empty_trash: false,
659 rewrite: None,
660 filter: PathFilter::default(),
661 request: Request::default(),
662 };
663
664 let libraries = [Library {
665 title: "Movies".to_string(),
666 key: "library_key_movies".to_string(),
667 refreshing: None,
668 scanned_at: None,
669 location: vec![Location {
670 path: "/media/movies".to_string(),
671 }],
672 }];
673
674 let path = "/media/movies/Inception.mkv";
675 let libraries = plex.get_libraries(&libraries, path);
676 assert!(libraries[0].key == "library_key_movies");
677
678 let nested_libraries = [
679 Library {
680 title: "Movies".to_string(),
681 key: "library_key_movies".to_string(),
682 refreshing: None,
683 scanned_at: None,
684 location: vec![Location {
685 path: "/media/movies".to_string(),
686 }],
687 },
688 Library {
689 title: "Movies".to_string(),
690 key: "library_key_movies_4k".to_string(),
691 refreshing: None,
692 scanned_at: None,
693 location: vec![Location {
694 path: "/media/movies/4k".to_string(),
695 }],
696 },
697 ];
698
699 let path = "/media/movies/4k/Inception.mkv";
700
701 let libraries = plex.get_libraries(&nested_libraries, path);
702 assert!(libraries[0].key == "library_key_movies_4k");
703 assert!(libraries[1].key == "library_key_movies");
704 }
705
706 #[test]
707 fn windows_library_matching_prefers_the_most_specific_location() {
708 let libraries = [
709 Library {
710 title: "Movies".to_string(),
711 key: "movies".to_string(),
712 refreshing: None,
713 scanned_at: None,
714 location: vec![Location {
715 path: r"\\server\media".to_string(),
716 }],
717 },
718 Library {
719 title: "4K Movies".to_string(),
720 key: "movies-4k".to_string(),
721 refreshing: None,
722 scanned_at: None,
723 location: vec![Location {
724 path: r"\\SERVER\MEDIA\4K".to_string(),
725 }],
726 },
727 ];
728
729 let matches = test_plex().get_libraries(&libraries, r"\\server\media\4k\Film\Film.mkv");
730
731 assert_eq!(
732 matches
733 .iter()
734 .map(|library| library.key.as_str())
735 .collect::<Vec<_>>(),
736 ["movies-4k", "movies"]
737 );
738 }
739
740 #[test]
741 fn unc_file_produces_a_non_empty_scan_directory_in_original_syntax() {
742 assert_eq!(
743 scan_directory(r"\\server\media\TV\Show\Season 1\S01E01.mkv"),
744 r"\\server\media\TV\Show\Season 1"
745 );
746 }
747
748 #[test]
749 fn windows_search_term_skips_season_components() {
750 assert_eq!(
751 test_plex()
752 .get_search_term(r"D:\TV Shows\Breaking Bad\Season 1\S01E01.mkv")
753 .unwrap(),
754 "Breaking Bad"
755 );
756 }
757
758 #[test]
759 fn windows_media_matching_uses_runtime_case_and_boundaries() {
760 assert!(path_matches(
761 r"D:\MEDIA\Movies\Film\Film.mkv",
762 r"d:\media\movies\film\film.MKV"
763 ));
764 assert!(path_matches(
765 r"\\server\media\Shows\Show\Episode.mkv",
766 r"\\SERVER\MEDIA\SHOWS\SHOW"
767 ));
768 assert!(!path_matches(
769 r"\\server\media-archive\Film.mkv",
770 r"\\server\media"
771 ));
772 }
773
774 #[test]
775 fn unix_search_term_without_a_non_season_parent_uses_the_directory_name() {
776 assert_eq!(
777 test_plex().get_search_term("/Season 1/file.mkv").unwrap(),
778 "Season 1"
779 );
780 }
781
782 #[test]
783 fn windows_search_term_without_a_non_season_parent_uses_the_directory_name() {
784 assert_eq!(
785 test_plex()
786 .get_search_term(r"D:\Season 1\S01E01.mkv")
787 .unwrap(),
788 "Season 1"
789 );
790
791 assert_eq!(
796 test_plex().get_search_term(r"D:\Season 1").unwrap(),
797 "Season 1"
798 );
799 }
800}