headless_lms_models/
chatbot_page_sync_statuses.rs1use 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 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
113pub async fn update_page_revision_ids(
115 conn: &mut PgConnection,
116 page_id_to_new_revision_id: HashMap<Uuid, Uuid>,
117) -> ModelResult<()> {
118 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
169pub 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}