Skip to main content

synd_persistence/sqlite/feed_registry/crawl/state/
mod.rs

1use 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/// `crawl_state` columns shared by the state load and the scheduler's due
106/// input queries.
107#[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;