synd_persistence/sqlite/feed_registry/crawl/state/
mod.rs1use chrono::{DateTime, Utc};
2use sqlx::{Sqlite, Transaction};
3use synd_feed::{
4 feed::service::{FeedConditionalFetch, FeedHttpStatus},
5 types::FeedUrl,
6};
7use synd_registry::{
8 RegistryDbResult,
9 crawl::state::{
10 CrawlHealth, CrawlState, CrawlStateError, FailureStreak, LastCrawlResult,
11 UpsertCrawlStateCommand,
12 },
13 db::CrawlStateDb,
14};
15
16use super::super::{
17 SqliteRegistryTx, codec,
18 error::{IntoDbResult, SqliteError, SqliteResult},
19 feed,
20};
21
22async fn load(
23 tx: &mut Transaction<'_, Sqlite>,
24 feed_url: &FeedUrl,
25) -> SqliteResult<Option<CrawlState>> {
26 let row = sqlx::query_as::<_, CrawlStateRow>(
27 r#"
28 SELECT
29 cs.last_started_at,
30 cs.last_finished_at,
31 cs.last_http_status,
32 cs.last_error_kind,
33 cs.failure_streak,
34 cs.retry_after,
35 cs.etag,
36 cs.last_modified
37 FROM crawl_state AS cs
38 INNER JOIN feed AS f
39 ON f.pk = cs.feed_pk
40 WHERE f.url = ?
41 "#,
42 )
43 .bind(feed_url.as_str())
44 .fetch_optional(&mut **tx)
45 .await?;
46
47 row.map(|row| row.into_state(feed_url)).transpose()
48}
49
50async fn upsert(
51 tx: &mut Transaction<'_, Sqlite>,
52 command: UpsertCrawlStateCommand,
53) -> SqliteResult<()> {
54 let feed_pk = feed::resolve_pk(tx, &command.feed_url).await?;
55 let last_http_status = command
56 .last
57 .http_status
58 .map(|status| i64::from(status.as_u16()));
59 let last_error_kind = command
60 .last
61 .error
62 .map(|error| codec::encode_crawl_state_error_kind(error.kind));
63 let failure_streak = encode_u64(command.health.failure_streak.value(), "failure streak")?;
64
65 sqlx::query(
66 r#"
67 INSERT INTO crawl_state (
68 feed_pk,
69 last_started_at,
70 last_finished_at,
71 last_http_status,
72 last_error_kind,
73 failure_streak,
74 retry_after,
75 etag,
76 last_modified
77 )
78 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)
79 ON CONFLICT(feed_pk) DO UPDATE SET
80 last_started_at = excluded.last_started_at,
81 last_finished_at = excluded.last_finished_at,
82 last_http_status = excluded.last_http_status,
83 last_error_kind = excluded.last_error_kind,
84 failure_streak = excluded.failure_streak,
85 retry_after = excluded.retry_after,
86 etag = excluded.etag,
87 last_modified = excluded.last_modified
88 "#,
89 )
90 .bind(feed_pk)
91 .bind(command.last.started_at)
92 .bind(command.last.finished_at)
93 .bind(last_http_status)
94 .bind(last_error_kind)
95 .bind(failure_streak)
96 .bind(command.last.retry_after)
97 .bind(command.conditional.etag.as_deref())
98 .bind(command.conditional.last_modified.as_deref())
99 .execute(&mut **tx)
100 .await?;
101
102 Ok(())
103}
104
105#[derive(sqlx::FromRow)]
108pub(in crate::sqlite::feed_registry) struct CrawlStateRow {
109 pub(in crate::sqlite::feed_registry) last_started_at: DateTime<Utc>,
110 pub(in crate::sqlite::feed_registry) last_finished_at: DateTime<Utc>,
111 pub(in crate::sqlite::feed_registry) last_http_status: Option<i64>,
112 pub(in crate::sqlite::feed_registry) last_error_kind: Option<String>,
113 pub(in crate::sqlite::feed_registry) failure_streak: i64,
114 pub(in crate::sqlite::feed_registry) retry_after: Option<DateTime<Utc>>,
115 pub(in crate::sqlite::feed_registry) etag: Option<String>,
116 pub(in crate::sqlite::feed_registry) last_modified: Option<String>,
117}
118
119impl CrawlStateRow {
120 pub(in crate::sqlite::feed_registry) fn into_state(
121 self,
122 feed_url: &FeedUrl,
123 ) -> SqliteResult<CrawlState> {
124 let http_status = self
125 .last_http_status
126 .map(|status| {
127 u16::try_from(status).map(FeedHttpStatus::new).map_err(|_| {
128 SqliteError::decode_message(format!(
129 "crawl state http status out of range: {status}"
130 ))
131 })
132 })
133 .transpose()?;
134 let error = self
135 .last_error_kind
136 .as_deref()
137 .map(codec::decode_crawl_state_error_kind)
138 .transpose()?
139 .map(|kind| CrawlStateError { kind });
140 let failure_streak = u64::try_from(self.failure_streak).map_err(|_| {
141 SqliteError::decode_message(format!(
142 "crawl state failure streak must be non-negative: {}",
143 self.failure_streak
144 ))
145 })?;
146
147 let last = LastCrawlResult {
148 started_at: self.last_started_at,
149 finished_at: self.last_finished_at,
150 http_status,
151 error,
152 retry_after: self.retry_after,
153 };
154
155 Ok(CrawlState {
156 feed_url: feed_url.clone(),
157 last,
158 health: CrawlHealth {
159 failure_streak: FailureStreak::new(failure_streak),
160 },
161 conditional: FeedConditionalFetch {
162 etag: self.etag,
163 last_modified: self.last_modified,
164 },
165 })
166 }
167}
168
169fn encode_u64(value: u64, field: &'static str) -> SqliteResult<i64> {
170 i64::try_from(value)
171 .map_err(|_| SqliteError::decode_message(format!("{field} exceeds SQLite INTEGER range")))
172}
173
174impl CrawlStateDb for SqliteRegistryTx<'_> {
175 async fn load_crawl_state(
176 &mut self,
177 feed_url: &FeedUrl,
178 ) -> RegistryDbResult<Option<CrawlState>> {
179 load(&mut self.tx, feed_url).await.db()
180 }
181
182 async fn upsert_crawl_state(
183 &mut self,
184 command: UpsertCrawlStateCommand,
185 ) -> RegistryDbResult<()> {
186 upsert(&mut self.tx, command).await.db()
187 }
188}
189
190#[cfg(test)]
191mod tests;