use std::borrow::Cow;
use std::collections::{BTreeMap, HashSet};
use serde::{Deserialize, Serialize};
use url::Url;
use crate::client::Response;
use crate::error::Error;
use crate::generated::services::calendars::Calendars;
use crate::generated::types::{Calendar, Recording};
use crate::http::Method;
use crate::observability::OperationInfo;
use crate::operation::Operation;
use crate::pagination::next_link;
use crate::security::is_same_origin;
use crate::services::calendars::ListedCalendar;
use crate::types::DateTime;
const TOO_FAR_BEHIND: u16 = 409;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct CalendarChangesCursor {
pub since: Option<String>,
pub version: Option<String>,
pub page: Option<String>,
pub per_page: Option<String>,
}
impl CalendarChangesCursor {
pub fn from_url(changes_url: &str) -> Result<CalendarChangesCursor, Error> {
let url = Url::parse(changes_url)
.map_err(|error| Error::usage(format!("changes URL {changes_url}: {error}")))?;
Ok(CalendarChangesCursor::from_parsed(&url))
}
fn from_parsed(url: &Url) -> CalendarChangesCursor {
CalendarChangesCursor {
since: parameter(url, "since"),
version: parameter(url, "v"),
page: parameter(url, "page"),
per_page: parameter(url, "per_page"),
}
}
fn apply(&self, operation: &mut Operation) {
operation.query_optional("since", self.since.as_ref());
operation.query_optional("v", self.version.as_ref());
operation.query_optional("page", self.page.as_ref());
operation.query_optional("per_page", self.per_page.as_ref());
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[non_exhaustive]
pub struct DeletedCalendar {
#[serde(default)]
pub id: i64,
pub deleted_at: DateTime,
}
#[derive(Debug, Clone, Default, PartialEq)]
#[non_exhaustive]
pub struct CalendarChanges {
pub added: Vec<ListedCalendar>,
pub updated: Vec<Calendar>,
pub deleted: Vec<DeletedCalendar>,
pub next_page: Option<CalendarChangesCursor>,
pub next_cursor: Option<CalendarChangesCursor>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[non_exhaustive]
pub struct DeletedRecording {
#[serde(default)]
pub id: i64,
pub deleted_at: DateTime,
#[serde(default)]
pub r#type: String,
}
#[derive(Debug, Clone, Default, PartialEq)]
#[non_exhaustive]
pub struct RecordingChanges {
pub added: BTreeMap<String, Vec<Recording>>,
pub updated: BTreeMap<String, Vec<Recording>>,
pub deleted: Vec<DeletedRecording>,
pub next_page: Option<CalendarChangesCursor>,
pub next_cursor: Option<CalendarChangesCursor>,
pub full_sync_required: bool,
}
impl Calendars<'_> {
pub async fn all_calendar_changes(
&self,
cursor: &CalendarChangesCursor,
) -> Result<CalendarChanges, Error> {
let mut all = CalendarChanges::default();
let mut cursor = cursor.clone();
for _ in 0..self.client().max_pages() {
let mut changes = self.calendar_changes(&cursor).await?;
all.added.append(&mut changes.added);
all.updated.append(&mut changes.updated);
all.deleted.append(&mut changes.deleted);
all.next_cursor = changes.next_cursor;
match changes.next_page {
None => return Ok(all),
Some(next) => cursor = next,
}
}
crate::trace::warning!(
max_pages = self.client().max_pages(),
"calendar changes pagination capped"
);
Ok(all)
}
pub async fn calendar_changes(
&self,
cursor: &CalendarChangesCursor,
) -> Result<CalendarChanges, Error> {
if cursor.since.is_none() {
return Err(Error::usage(
"a since cursor is required — start from the list's calendar_changes_url",
));
}
let mut operation = self.client().request(Method::GET, "/calendar/changes");
operation.info(changes_info("GetCalendarChanges", "calendar"));
cursor.apply(&mut operation);
operation.no_cache();
let response = self.client().execute(operation).await?;
let payload: CalendarChangesPayload = response.json()?;
let (next_page, next_cursor) = next_cursors(&response, self.client().base_url())?;
Ok(CalendarChanges {
added: payload.added,
updated: payload.updated,
deleted: payload.deleted,
next_page,
next_cursor,
})
}
pub async fn all_recording_changes(
&self,
calendar_id: i64,
cursor: &CalendarChangesCursor,
) -> Result<RecordingChanges, Error> {
let mut all = RecordingChanges::default();
let mut cursor = cursor.clone();
for _ in 0..self.client().max_pages() {
let mut changes = self.recording_changes(calendar_id, &cursor).await?;
if changes.full_sync_required {
return Ok(changes);
}
merge_recordings(&mut all.added, changes.added);
merge_recordings(&mut all.updated, changes.updated);
all.deleted.append(&mut changes.deleted);
all.next_cursor = changes.next_cursor;
match changes.next_page {
None => return Ok(all),
Some(next) => cursor = next,
}
}
crate::trace::warning!(
max_pages = self.client().max_pages(),
"recording changes pagination capped"
);
Ok(all)
}
pub async fn recording_changes(
&self,
calendar_id: i64,
cursor: &CalendarChangesCursor,
) -> Result<RecordingChanges, Error> {
if cursor.since.is_none() {
return Err(Error::usage(
"a since cursor is required — start from the calendar's recording_changes_url",
));
}
if cursor.version.is_none() {
return Err(Error::usage(
"a feed version is required — read the calendar's recording_changes_url with CalendarChangesCursor::from_url",
));
}
let mut operation = self.client().request(
Method::GET,
format!("/calendars/{calendar_id}/recording/changes"),
);
operation.info(changes_info("GetCalendarRecordingChanges", "recording"));
operation.resource_id(calendar_id);
cursor.apply(&mut operation);
operation.no_cache();
let response = match self.client().execute(operation).await {
Ok(response) => response,
Err(error) if error.http_status() == Some(TOO_FAR_BEHIND) => {
return Ok(RecordingChanges {
full_sync_required: true,
..RecordingChanges::default()
});
}
Err(error) => return Err(error),
};
let payload: RecordingChangesPayload = response.json()?;
let (next_page, next_cursor) = next_cursors(&response, self.client().base_url())?;
Ok(RecordingChanges {
added: payload.added,
updated: payload.updated,
deleted: flatten_deleted_recordings(payload.deleted),
next_page,
next_cursor,
full_sync_required: false,
})
}
}
#[derive(Debug, Default, Deserialize)]
struct CalendarChangesPayload {
#[serde(default)]
added: Vec<ListedCalendar>,
#[serde(default)]
updated: Vec<Calendar>,
#[serde(default)]
deleted: Vec<DeletedCalendar>,
}
#[derive(Debug, Default, Deserialize)]
struct RecordingChangesPayload {
#[serde(default)]
added: BTreeMap<String, Vec<Recording>>,
#[serde(default)]
updated: BTreeMap<String, Vec<Recording>>,
#[serde(default)]
deleted: BTreeMap<String, Vec<DeletedRecording>>,
}
fn changes_info(operation: &'static str, resource_type: &'static str) -> OperationInfo {
OperationInfo {
service: Cow::Borrowed("Calendars"),
operation: Cow::Borrowed(operation),
resource_type: Cow::Borrowed(resource_type),
is_mutation: false,
resource_id: None,
}
}
fn next_cursors(
response: &Response,
base_url: &Url,
) -> Result<(Option<CalendarChangesCursor>, Option<CalendarChangesCursor>), Error> {
match link_cursor(response, base_url)? {
None => Ok((None, None)),
Some(cursor) if cursor.page.is_some() => Ok((Some(cursor), None)),
Some(cursor) => Ok((None, Some(cursor))),
}
}
fn link_cursor(
response: &Response,
base_url: &Url,
) -> Result<Option<CalendarChangesCursor>, Error> {
match response.header("link").and_then(next_link) {
None => Ok(None),
Some(target) => {
let next = response.url.join(&target)?;
if is_same_origin(&next, base_url) {
Ok(Some(CalendarChangesCursor::from_parsed(&next)))
} else {
Err(Error::usage(format!(
"changes Link header points to a different origin: {next}"
)))
}
}
}
}
fn merge_recordings(
into: &mut BTreeMap<String, Vec<Recording>>,
from: BTreeMap<String, Vec<Recording>>,
) {
for (key, recordings) in from {
into.entry(key).or_default().extend(recordings);
}
}
fn flatten_deleted_recordings(
buckets: BTreeMap<String, Vec<DeletedRecording>>,
) -> Vec<DeletedRecording> {
let mut seen = HashSet::new();
buckets
.into_values()
.flatten()
.filter(|record| seen.insert(record.id))
.collect()
}
fn parameter(url: &Url, name: &str) -> Option<String> {
url.query_pairs()
.find(|(key, _)| key == name)
.map(|(_, value)| value.into_owned())
.filter(|value| !value.is_empty())
}