Skip to main content

autopulse_service/settings/targets/
emby.rs

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    /// URL to the Jellyfin/Emby server
26    pub url: String,
27    /// API token for the Jellyfin/Emby server
28    pub token: String,
29    /// Metadata refresh mode (default: `FullRefresh`)
30    #[serde(default)]
31    pub metadata_refresh_mode: EmbyMetadataRefreshMode,
32    /// Whether to try to refresh metadata for the item instead of scan (default: true)
33    #[serde(default = "default_true")]
34    pub refresh_metadata: bool,
35    /// Rewrite path for the file
36    pub rewrite: Option<Rewrite>,
37    /// Path filter matched against the target-rewritten path.
38    #[serde(default)]
39    pub filter: PathFilter,
40    /// HTTP request options
41    #[serde(default)]
42    pub request: Request,
43    /// How library locations are compared to incoming paths. Default `case_sensitive`.
44    #[serde(default)]
45    pub path_match: PathMatch,
46}
47
48/// How library locations are compared to incoming paths.
49#[derive(Serialize, Clone, Deserialize, Default)]
50#[serde(rename_all = "snake_case")]
51pub enum PathMatch {
52    /// Exact byte-for-byte prefix match (default; safe on Linux).
53    #[default]
54    CaseSensitive,
55    /// Lowercase both sides before comparing. Useful on Windows / UNC paths.
56    CaseInsensitive,
57}
58
59/// Metadata refresh mode for Jellyfin/Emby
60#[derive(Serialize, Clone, Deserialize)]
61#[serde(rename_all = "snake_case")]
62#[derive(Default)]
63pub enum EmbyMetadataRefreshMode {
64    /// `none`
65    None,
66    /// `validation_only`
67    ValidationOnly,
68    /// `default`
69    Default,
70    /// `full_refresh`
71    #[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; // one location is enough to match this library
156                }
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        // Possibly unneeded unless we can use streams
200        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    // not as effective as refreshing the item, but good enough
326    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        // TODO: Possible options in future?
364        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}