Skip to main content

headless_lms_models/
chatbot_page_sync_statuses.rs

1use std::collections::HashMap;
2
3use crate::{course_page_markdown_content, prelude::*};
4
5#[derive(Debug, Serialize, Deserialize, PartialEq, Clone)]
6
7pub struct ChatbotPageSyncStatus {
8    pub id: Uuid,
9    pub created_at: DateTime<Utc>,
10    pub updated_at: DateTime<Utc>,
11    pub deleted_at: Option<DateTime<Utc>>,
12    pub course_id: Uuid,
13    pub page_id: Uuid,
14    pub error_message: Option<String>,
15    pub synced_page_revision_id: Option<Uuid>,
16    pub consecutive_failures: i32,
17    pub converted_markdown_content_id: Option<Uuid>,
18}
19
20pub async fn ensure_sync_statuses_exist(
21    conn: &mut PgConnection,
22    course_ids: &[Uuid],
23) -> ModelResult<HashMap<Uuid, Vec<ChatbotPageSyncStatus>>> {
24    sqlx::query!(
25        r#"
26INSERT INTO chatbot_page_sync_statuses (course_id, page_id)
27SELECT course_id,
28  id
29FROM pages
30WHERE course_id = ANY($1)
31  AND deleted_at IS NULL
32  AND hidden IS FALSE ON CONFLICT (page_id, deleted_at) DO NOTHING
33        "#,
34        course_ids
35    )
36    .execute(&mut *conn)
37    .await?;
38
39    let all_statuses = sqlx::query_as!(
40        ChatbotPageSyncStatus,
41        r#"
42SELECT *
43FROM chatbot_page_sync_statuses
44WHERE course_id = ANY($1)
45        "#,
46        course_ids
47    )
48    .fetch_all(&mut *conn)
49    .await?
50    .into_iter()
51    .fold(
52        HashMap::<Uuid, Vec<ChatbotPageSyncStatus>>::new(),
53        |mut map, status| {
54            map.entry(status.course_id).or_default().push(status);
55            map
56        },
57    );
58
59    Ok(all_statuses)
60}
61
62pub async fn save_markdown_content(
63    conn: &mut PgConnection,
64    page_id_to_history_id_md_content: HashMap<Uuid, (Uuid, String)>,
65) -> ModelResult<()> {
66    // If there are no updates to perform, return early
67    if page_id_to_history_id_md_content.is_empty() {
68        return Ok(());
69    }
70
71    let mut tx = conn.begin().await?;
72    let res = course_page_markdown_content::insert_batch(
73        &mut tx,
74        page_id_to_history_id_md_content
75            .clone()
76            .into_iter()
77            .collect(),
78    )
79    .await?;
80
81    let (md_ids, page_ids): (Vec<Uuid>, Vec<Uuid>) = page_id_to_history_id_md_content
82        .iter()
83        .filter_map(|(p_id, (ph_id, _))| {
84            let md_content = res.iter().find(|x| &x.page_history_id == ph_id);
85            md_content.map(|content| (content.id, p_id.to_owned()))
86        })
87        .collect::<Vec<(Uuid, Uuid)>>()
88        .into_iter()
89        .unzip();
90
91    sqlx::query!(
92        r#"
93UPDATE chatbot_page_sync_statuses AS cps
94SET converted_markdown_content_id = data.markdown_id
95FROM (
96    SELECT unnest($1::uuid []) AS markdown_id,
97      unnest($2::uuid []) AS page_id
98  ) AS data
99WHERE cps.page_id = data.page_id
100  AND cps.deleted_at IS NULL
101    "#,
102        &md_ids,
103        &page_ids
104    )
105    .execute(&mut *tx)
106    .await?;
107
108    tx.commit().await?;
109
110    Ok(())
111}
112
113// Given a mapping from page id to the new revision id, update the sync statuses
114pub async fn update_page_revision_ids(
115    conn: &mut PgConnection,
116    page_id_to_new_revision_id: HashMap<Uuid, Uuid>,
117) -> ModelResult<()> {
118    // If there are no updates to perform, return early
119    if page_id_to_new_revision_id.is_empty() {
120        return Ok(());
121    }
122    let (page_ids, revision_ids): (Vec<Uuid>, Vec<Uuid>) =
123        page_id_to_new_revision_id.into_iter().unzip();
124
125    sqlx::query!(
126        r#"
127UPDATE chatbot_page_sync_statuses AS cps
128SET synced_page_revision_id = data.synced_page_revision_id,
129    error_message = NULL,
130    consecutive_failures = 0
131FROM (
132    SELECT unnest($1::uuid []) AS page_id,
133      unnest($2::uuid []) AS synced_page_revision_id
134  ) AS data
135WHERE cps.page_id = data.page_id
136AND cps.deleted_at IS NULL
137    "#,
138        &page_ids,
139        &revision_ids
140    )
141    .execute(conn)
142    .await?;
143
144    Ok(())
145}
146
147pub async fn set_page_sync_error(
148    conn: &mut PgConnection,
149    page_id: Uuid,
150    error_message: &str,
151) -> ModelResult<()> {
152    sqlx::query!(
153        r#"
154UPDATE chatbot_page_sync_statuses
155SET error_message = $2,
156    consecutive_failures = consecutive_failures + 1
157WHERE page_id = $1
158AND deleted_at IS NULL
159    "#,
160        page_id,
161        error_message
162    )
163    .execute(conn)
164    .await?;
165
166    Ok(())
167}
168
169/// Clears sync statuses for the given page IDs.
170/// This is used when pages become hidden to ensure they'll be re-synced if unhidden.
171pub async fn clear_sync_statuses(conn: &mut PgConnection, page_ids: &[Uuid]) -> ModelResult<()> {
172    if page_ids.is_empty() {
173        return Ok(());
174    }
175
176    sqlx::query!(
177        r#"
178UPDATE chatbot_page_sync_statuses
179SET synced_page_revision_id = NULL,
180    error_message = NULL,
181    consecutive_failures = 0
182WHERE page_id = ANY($1)
183AND deleted_at IS NULL
184        "#,
185        page_ids
186    )
187    .execute(conn)
188    .await?;
189
190    Ok(())
191}