Skip to main content

codoseo_web/agent/
service.rs

1//! The agent API's one service layer. The REST routes (`/api/v1`) and the cloud MCP tools both
2//! call these methods, so they return the same JSON ([`codoseo_mcp::cloud::types`]) and charge
3//! the same daily quota.
4//!
5//! Every keyed method takes the already-authenticated [`ApiCaller`], counts one call against
6//! the account's allowance first (a refused call counts nothing), and only then looks at its
7//! arguments. Arguments arrive as the strings a caller sent (ids, check slugs, severities) and
8//! are validated here, after the charge, so a malformed call costs the same as any other.
9//! Another account's site is the same `NotFound` as an unknown id.
10
11use std::collections::HashMap;
12
13use codoseo_checks::def;
14use codoseo_core::check::{CheckId, IssueBits, Severity};
15use codoseo_core::plan::PlanLimits;
16use codoseo_mcp::cloud::types::{
17    ActiveCrawl, ChangeInfo, ChangesPage, CrawlHealth, CrawlQueued, IssueUrlsPage, PageInfo,
18    PageIssue, RedirectHop, SiteHealth, SiteInfo, Usage, clip_change_text,
19};
20use codoseo_mcp::local::{stop_reason_code, stop_reason_words};
21use codoseo_mcp::types::{FailingCheck, MAX_FAILING_CHECKS, UrlRow, rank_failing};
22use codoseo_store::api_keys::{self, Charge};
23use codoseo_store::crawls::{self, Crawl, CrawlStatus, ManualOutcome, ManualWindow};
24use codoseo_store::explorer::{self, PageFilter};
25use codoseo_store::reports;
26use codoseo_store::sites::{self, Site};
27use sqlx::PgPool;
28use time::OffsetDateTime;
29use url::Url;
30use uuid::Uuid;
31
32use super::auth::ApiCaller;
33use super::error::AgentError;
34use crate::crawl_policy::{limit_message, manual_priority};
35use crate::state::AppState;
36
37/// Rows a list returns when the caller doesn't say.
38pub const DEFAULT_LIMIT: u32 = 50;
39/// The most rows one call returns.
40pub const MAX_LIMIT: u32 = 200;
41/// Example URLs listed per failing check.
42const EXAMPLES_PER_CHECK: i64 = 3;
43
44/// Where the caller stands against today's allowance after a call. `None` fields mean the plan
45/// has no limit (self-hosted).
46#[derive(Debug, Clone, Copy, PartialEq, Eq)]
47pub struct Quota {
48    pub limit: Option<u32>,
49    pub remaining: Option<u32>,
50}
51
52/// A keyed call's result, with the quota it left behind (`None` when the charge itself failed
53/// on a database error). The REST layer puts the quota in headers; MCP only needs the outcome.
54#[derive(Debug)]
55pub struct Reply<T> {
56    pub quota: Option<Quota>,
57    pub outcome: Result<T, AgentError>,
58}
59
60impl<T> Reply<T> {
61    pub fn into_result(self) -> Result<T, AgentError> {
62        self.outcome
63    }
64}
65
66pub struct AgentService<'a> {
67    state: &'a AppState,
68}
69
70impl<'a> AgentService<'a> {
71    pub fn new(state: &'a AppState) -> AgentService<'a> {
72        AgentService { state }
73    }
74
75    /// Counts one call, then runs `work`. `work` is a future that hasn't started: nothing in it
76    /// runs when the charge is refused.
77    async fn metered<T>(
78        &self,
79        caller: &ApiCaller,
80        work: impl Future<Output = Result<T, AgentError>>,
81    ) -> Reply<T> {
82        let limit = PlanLimits::for_plan(caller.account.plan).api_calls_per_day;
83        match api_keys::charge(&self.state.pool, caller.account.id, limit).await {
84            Err(e) => Reply {
85                quota: None,
86                outcome: Err(e.into()),
87            },
88            Ok(Charge::OverQuota { limit }) => {
89                // The reset comes from Postgres' clock, the one the day was counted by.
90                let outcome = match api_keys::day_end(&self.state.pool, None).await {
91                    Ok(end) => Err(AgentError::QuotaExceeded {
92                        limit,
93                        retry_after_secs: u64::try_from(end.seconds_left).unwrap_or(1),
94                    }),
95                    Err(e) => Err(e.into()),
96                };
97                Reply {
98                    quota: Some(Quota {
99                        limit: Some(limit),
100                        remaining: Some(0),
101                    }),
102                    outcome,
103                }
104            }
105            Ok(Charge::Ok { used, limit }) => Reply {
106                quota: Some(Quota {
107                    limit,
108                    remaining: limit.map(|l| l.saturating_sub(clamp_u32(used))),
109                }),
110                outcome: work.await,
111            },
112        }
113    }
114
115    /// Counts a call that the caller's own arguments already made unanswerable (a query string
116    /// that doesn't parse) and answers with `error`: every authenticated request costs one.
117    pub async fn refuse<T>(&self, caller: &ApiCaller, error: AgentError) -> Reply<T> {
118        self.metered(caller, async { Err(error) }).await
119    }
120
121    pub async fn list_sites(&self, caller: &ApiCaller) -> Reply<Vec<SiteInfo>> {
122        self.metered(caller, self.list_sites_work(caller)).await
123    }
124
125    pub async fn site_health(&self, caller: &ApiCaller, site: &str) -> Reply<SiteHealth> {
126        self.metered(caller, self.site_health_work(caller, site))
127            .await
128    }
129
130    pub async fn issue_urls(
131        &self,
132        caller: &ApiCaller,
133        site: &str,
134        check: &str,
135        limit: Option<u32>,
136        offset: Option<u32>,
137    ) -> Reply<IssueUrlsPage> {
138        self.metered(
139            caller,
140            self.issue_urls_work(caller, site, check, limit, offset),
141        )
142        .await
143    }
144
145    pub async fn page(&self, caller: &ApiCaller, site: &str, url: &str) -> Reply<PageInfo> {
146        self.metered(caller, self.page_work(caller, site, url))
147            .await
148    }
149
150    pub async fn changes(
151        &self,
152        caller: &ApiCaller,
153        site: &str,
154        severity: Option<&str>,
155        limit: Option<u32>,
156        offset: Option<u32>,
157    ) -> Reply<ChangesPage> {
158        self.metered(
159            caller,
160            self.changes_work(caller, site, severity, limit, offset),
161        )
162        .await
163    }
164
165    /// Queues a manual crawl exactly as the Run crawl button does: the plan's allowance and
166    /// lane, one crawl at a time.
167    pub async fn run_crawl(&self, caller: &ApiCaller, site: &str) -> Reply<CrawlQueued> {
168        self.metered(caller, self.run_crawl_work(caller, site))
169            .await
170    }
171
172    /// Today's calls and the allowance. Free: it isn't counted.
173    pub async fn usage(&self, caller: &ApiCaller) -> Reply<Usage> {
174        let limit = PlanLimits::for_plan(caller.account.plan).api_calls_per_day;
175        let pool = &self.state.pool;
176        let outcome = async {
177            let calls = clamp_u32(api_keys::usage_today(pool, caller.account.id).await?);
178            let end = api_keys::day_end(pool, None).await?;
179            Ok::<_, AgentError>(Usage {
180                calls_today: calls,
181                limit,
182                remaining: limit.map(|l| l.saturating_sub(calls)),
183                resets_at: end.resets_at,
184            })
185        }
186        .await;
187        Reply {
188            quota: outcome.as_ref().ok().map(|u| Quota {
189                limit: u.limit,
190                remaining: u.remaining,
191            }),
192            outcome,
193        }
194    }
195
196    /// One of the caller's sites; another account's looks exactly like an unknown id.
197    async fn site(&self, caller: &ApiCaller, raw: &str) -> Result<Site, AgentError> {
198        let id: Uuid = raw
199            .trim()
200            .parse()
201            .map_err(|_| AgentError::site_not_found())?;
202        sites::get_for_account(&self.state.pool, caller.account.id, id)
203            .await?
204            .ok_or_else(AgentError::site_not_found)
205    }
206
207    async fn list_sites_work(&self, caller: &ApiCaller) -> Result<Vec<SiteInfo>, AgentError> {
208        let pool = &self.state.pool;
209        let list = sites::list_for_account(pool, caller.account.id).await?;
210        let latest: HashMap<Uuid, (Option<i16>, Option<OffsetDateTime>)> =
211            crawls::latest_done_for_account(pool, caller.account.id)
212                .await?
213                .into_iter()
214                .map(|(id, score, at)| (id, (score, at)))
215                .collect();
216        Ok(list
217            .into_iter()
218            .map(|s| {
219                let (score, at) = latest.get(&s.id).copied().unwrap_or((None, None));
220                SiteInfo {
221                    id: s.id,
222                    domain: s.domain,
223                    start_url: s.start_url,
224                    monitoring_active: s.monitoring_active,
225                    schedule: s.schedule,
226                    health_score: score.and_then(|n| u8::try_from(n).ok()),
227                    last_crawled_at: at,
228                }
229            })
230            .collect())
231    }
232
233    async fn site_health_work(
234        &self,
235        caller: &ApiCaller,
236        site: &str,
237    ) -> Result<SiteHealth, AgentError> {
238        let pool = &self.state.pool;
239        let site = self.site(caller, site).await?;
240        let latest = match crawls::latest_done(pool, site.id).await? {
241            Some(crawl) => Some(crawl_health(pool, &crawl).await?),
242            None => None,
243        };
244        let active = crawls::active(pool, site.id).await?.map(|c| ActiveCrawl {
245            crawl_id: c.id,
246            number: c.number,
247            status: if c.status == CrawlStatus::Running {
248                "running"
249            } else {
250                "queued"
251            }
252            .to_owned(),
253            pages_done: c.progress().map(|p| p.pages_done),
254        });
255        let audit_url = self
256            .state
257            .config
258            .base_url
259            .join(&format!("s/{}/audit", site.id))
260            .map_or_else(|_| format!("/s/{}/audit", site.id), |u| u.to_string());
261        Ok(SiteHealth {
262            id: site.id,
263            next_crawl_at: sites::next_crawl_at(pool, site.id).await?,
264            domain: site.domain,
265            start_url: site.start_url,
266            monitoring_active: site.monitoring_active,
267            schedule: site.schedule,
268            latest_crawl: latest,
269            active_crawl: active,
270            audit_url,
271        })
272    }
273
274    async fn issue_urls_work(
275        &self,
276        caller: &ApiCaller,
277        site: &str,
278        check: &str,
279        limit: Option<u32>,
280        offset: Option<u32>,
281    ) -> Result<IssueUrlsPage, AgentError> {
282        let pool = &self.state.pool;
283        let site = self.site(caller, site).await?;
284        let check = parse_check(check)?;
285        let mut page = IssueUrlsPage {
286            site_id: site.id,
287            check,
288            title: def(check).title.to_owned(),
289            crawl_number: None,
290            total: 0,
291            limit: clamp_limit(limit),
292            offset: offset.unwrap_or(0),
293            urls: Vec::new(),
294            next_offset: None,
295        };
296        let Some(crawl) = crawls::latest_done(pool, site.id).await? else {
297            return Ok(page);
298        };
299        let rows = issue_page(pool, crawl.id, check, limit, offset).await?;
300        page.crawl_number = Some(crawl.number);
301        page.total = rows.total;
302        page.urls = rows.urls;
303        page.next_offset = rows.next_offset;
304        Ok(page)
305    }
306
307    async fn page_work(
308        &self,
309        caller: &ApiCaller,
310        site: &str,
311        url: &str,
312    ) -> Result<PageInfo, AgentError> {
313        let pool = &self.state.pool;
314        let site = self.site(caller, site).await?;
315        if url.trim().is_empty() {
316            return Err(AgentError::BadRequest("The url is required.".to_owned()));
317        }
318        let base = Url::parse(&site.start_url)
319            .map_err(|e| AgentError::Internal(format!("site start url: {e}")))?;
320        let url = codoseo_core::url::normalize(&base, url)
321            .ok_or_else(|| AgentError::BadRequest("That is not a valid page URL.".to_owned()))?;
322        let crawl = crawls::latest_done(pool, site.id).await?.ok_or_else(|| {
323            AgentError::NotFound("This site has no finished crawl yet.".to_owned())
324        })?;
325        let page = explorer::page(pool, site.id, crawl.id, codoseo_core::url::url_hash(&url))
326            .await?
327            .ok_or_else(|| {
328                AgentError::NotFound("That URL is not in the latest crawl of this site.".to_owned())
329            })?;
330        let mut issues: Vec<PageIssue> = IssueBits(page.issues)
331            .iter()
332            .filter_map(CheckId::from_bit)
333            .map(|check| PageIssue {
334                check,
335                title: def(check).title.to_owned(),
336                severity: def(check).severity,
337            })
338            .collect();
339        issues.sort_by_key(|i| (i.severity, i.check));
340        Ok(PageInfo {
341            site_id: site.id,
342            crawl_number: crawl.number,
343            url: page.url,
344            status: page.status,
345            redirect_chain: page
346                .redirect_chain
347                .into_iter()
348                .map(|(status, url)| RedirectHop { status, url })
349                .collect(),
350            redirect_target: page.redirect_target,
351            response_ms: page.response_ms,
352            size_bytes: page.size_bytes,
353            content_type: page.content_type,
354            depth: page.depth,
355            in_sitemap: page.in_sitemap,
356            indexability: page.indexability,
357            title: page.title,
358            meta_description: page.meta_description,
359            meta_robots: page.meta_robots,
360            x_robots_tag: page.x_robots_tag,
361            canonical: page.canonical,
362            h1: page.h1,
363            h2: page.h2,
364            word_count: page.word_count,
365            inlinks: page.inlinks,
366            outlinks_internal: page.outlinks_internal,
367            outlinks_external: page.outlinks_external,
368            issues,
369        })
370    }
371
372    async fn changes_work(
373        &self,
374        caller: &ApiCaller,
375        site: &str,
376        severity: Option<&str>,
377        limit: Option<u32>,
378        offset: Option<u32>,
379    ) -> Result<ChangesPage, AgentError> {
380        let pool = &self.state.pool;
381        let site = self.site(caller, site).await?;
382        let severity = severity
383            .map(str::trim)
384            .filter(|s| !s.is_empty())
385            .map(parse_severity)
386            .transpose()?;
387        let (limit, offset) = (clamp_limit(limit), offset.unwrap_or(0));
388        let mut out = ChangesPage {
389            site_id: site.id,
390            crawl_number: None,
391            crawl_finished_at: None,
392            total: 0,
393            limit,
394            offset,
395            changes: Vec::new(),
396            next_offset: None,
397        };
398        let Some(crawl) = crawls::latest_done(pool, site.id).await? else {
399            return Ok(out);
400        };
401        let counts = reports::change_kind_counts(pool, crawl.id).await?;
402        let rows = reports::changes_for_crawl_at(
403            pool,
404            crawl.id,
405            severity,
406            i64::from(offset),
407            i64::from(limit),
408        )
409        .await?;
410        out.crawl_number = Some(crawl.number);
411        out.crawl_finished_at = crawl.finished_at;
412        out.total = clamp_u32(severity.map_or_else(|| counts.total(), |s| counts.severity(s)));
413        let shown = u32::try_from(rows.len()).unwrap_or(u32::MAX);
414        out.changes = rows
415            .into_iter()
416            .map(|c| ChangeInfo {
417                kind: c.kind,
418                severity: c.severity,
419                url: c.url,
420                before: clip_change_text(&c.before),
421                after: clip_change_text(&c.after),
422            })
423            .collect();
424        out.next_offset = next_offset(offset, shown, out.total);
425        Ok(out)
426    }
427
428    async fn run_crawl_work(
429        &self,
430        caller: &ApiCaller,
431        site: &str,
432    ) -> Result<CrawlQueued, AgentError> {
433        let site = self.site(caller, site).await?;
434        let plan = caller.account.plan;
435        let allowance = PlanLimits::for_plan(plan).manual_crawls;
436        let outcome = crawls::enqueue_manual_checked(
437            &self.state.pool,
438            site.id,
439            &site.domain,
440            manual_priority(plan),
441            ManualWindow::for_allowance(allowance),
442        )
443        .await
444        // The site was deleted between the lookup and the insert.
445        .map_err(|e| match e {
446            sqlx::Error::RowNotFound => AgentError::site_not_found(),
447            e => e.into(),
448        })?;
449        match outcome {
450            ManualOutcome::Queued { id, number } => Ok(CrawlQueued {
451                site_id: site.id,
452                crawl_id: id,
453                number,
454                status: "queued".to_owned(),
455            }),
456            ManualOutcome::Busy(status) => {
457                let what = if status == CrawlStatus::Running {
458                    "running"
459                } else {
460                    "queued"
461                };
462                Err(AgentError::CrawlInProgress(format!(
463                    "A crawl is already {what} for this site."
464                )))
465            }
466            ManualOutcome::LimitReached { frees_at } => Err(AgentError::PlanLimit(limit_message(
467                plan,
468                allowance,
469                frees_at - OffsetDateTime::now_utc(),
470            ))),
471        }
472    }
473}
474
475/// How many rows a list call returns: `limit` if given, at most [`MAX_LIMIT`], at least 1.
476pub fn clamp_limit(limit: Option<u32>) -> u32 {
477    limit.unwrap_or(DEFAULT_LIMIT).clamp(1, MAX_LIMIT)
478}
479
480/// The offset of the next page, when rows remain after `offset + shown`.
481fn next_offset(offset: u32, shown: u32, total: u32) -> Option<u32> {
482    offset.checked_add(shown).filter(|next| *next < total)
483}
484
485/// A check by slug.
486pub fn parse_check(slug: &str) -> Result<CheckId, AgentError> {
487    CheckId::from_slug(slug.trim()).ok_or_else(|| {
488        let valid: Vec<&str> = CheckId::ALL.iter().map(|c| c.slug()).collect();
489        AgentError::BadRequest(format!(
490            "Unknown check \"{slug}\". The check slugs are: {}.",
491            valid.join(", ")
492        ))
493    })
494}
495
496fn clamp_u32(n: i64) -> u32 {
497    u32::try_from(n.max(0)).unwrap_or(u32::MAX)
498}
499
500/// `critical`, `warning` or `notice` (lowercase, as the web app's filters write them).
501pub fn parse_severity(s: &str) -> Result<Severity, AgentError> {
502    Severity::from_slug(s.trim()).ok_or_else(|| {
503        AgentError::BadRequest(format!(
504            "Unknown severity \"{s}\". Use critical, warning or notice."
505        ))
506    })
507}
508
509/// A finished crawl's score, counts and failing checks (most severe first, with example URLs
510/// fetched in one query). It reads the crawl only, so it serves any crawl id the caller is
511/// already allowed to see, a site's or a quick audit's.
512pub async fn crawl_health(pool: &PgPool, crawl: &Crawl) -> Result<CrawlHealth, AgentError> {
513    let summary = crawl.summary();
514    let mut failing = rank_failing(
515        summary
516            .iter()
517            .flat_map(|s| s.counts.iter())
518            .filter_map(|(slug, n)| Some((CheckId::from_slug(slug)?, *n))),
519    );
520    let more = failing.len().saturating_sub(MAX_FAILING_CHECKS);
521    failing.truncate(MAX_FAILING_CHECKS);
522    let ids: Vec<CheckId> = failing.iter().map(|&(id, _)| id).collect();
523    let mut examples = explorer::example_urls(pool, crawl.id, &ids, EXAMPLES_PER_CHECK).await?;
524    let failing_checks = failing
525        .into_iter()
526        .map(|(check, count)| {
527            let d = def(check);
528            FailingCheck {
529                check,
530                title: d.title.to_owned(),
531                severity: d.severity,
532                count,
533                example_urls: examples
534                    .remove(&check)
535                    .unwrap_or_default()
536                    .iter()
537                    .filter_map(|u| Url::parse(u).ok())
538                    .collect(),
539            }
540        })
541        .collect();
542    Ok(CrawlHealth {
543        crawl_id: crawl.id,
544        number: crawl.number,
545        finished_at: crawl.finished_at,
546        health_score: crawl.health_score.and_then(|n| u8::try_from(n).ok()),
547        checks_passed: crawl.checks_passed.and_then(|n| u16::try_from(n).ok()),
548        checks_total: crawl.checks_total.and_then(|n| u16::try_from(n).ok()),
549        pages_crawled: summary.as_ref().map_or(0, |s| s.report_summary.pages),
550        stop_reason: summary.as_ref().map_or_else(
551            || "unknown".to_owned(),
552            |s| stop_reason_words(&s.stop_reason),
553        ),
554        stop_code: summary
555            .as_ref()
556            .map_or("unknown", |s| stop_reason_code(&s.stop_reason))
557            .to_owned(),
558        failing_checks,
559        more_failing_checks: u16::try_from(more).unwrap_or(u16::MAX),
560    })
561}
562
563/// One page of the pages of a crawl that fail `check`.
564#[derive(Debug, Clone, PartialEq)]
565pub struct IssueRows {
566    /// Pages failing the check in the crawl.
567    pub total: u32,
568    pub limit: u32,
569    pub offset: u32,
570    pub urls: Vec<UrlRow>,
571    /// Pass as `offset` for the next page; none on the last page.
572    pub next_offset: Option<u32>,
573}
574
575/// The pages of `crawl_id` that fail `check`, `limit` (clamped) from `offset`, in crawl order.
576pub async fn issue_page(
577    pool: &PgPool,
578    crawl_id: Uuid,
579    check: CheckId,
580    limit: Option<u32>,
581    offset: Option<u32>,
582) -> Result<IssueRows, AgentError> {
583    let (limit, offset) = (clamp_limit(limit), offset.unwrap_or(0));
584    let filter = PageFilter::Check(check);
585    let (matching, _) = explorer::match_count(pool, crawl_id, filter, "").await?;
586    let rows =
587        explorer::rows_at(pool, crawl_id, filter, i64::from(offset), i64::from(limit)).await?;
588    let total = clamp_u32(matching);
589    let shown = u32::try_from(rows.len()).unwrap_or(u32::MAX);
590    let urls = rows
591        .into_iter()
592        .filter_map(|r| {
593            Some(UrlRow {
594                url: Url::parse(&r.url).ok()?,
595                status: r.status,
596                title: r.title,
597                indexability: r.indexability,
598            })
599        })
600        .collect();
601    Ok(IssueRows {
602        total,
603        limit,
604        offset,
605        urls,
606        next_offset: next_offset(offset, shown, total),
607    })
608}
609
610#[cfg(test)]
611mod tests {
612    use super::*;
613
614    #[test]
615    fn limits_are_clamped() {
616        assert_eq!(clamp_limit(None), 50);
617        assert_eq!(clamp_limit(Some(0)), 1);
618        assert_eq!(clamp_limit(Some(10_000)), 200);
619    }
620
621    #[test]
622    fn severities_are_lowercase_slugs() {
623        assert_eq!(parse_severity("critical").unwrap(), Severity::Critical);
624        assert_eq!(parse_severity(" notice ").unwrap(), Severity::Notice);
625        assert!(parse_severity("Critical").is_err());
626        assert!(parse_severity("severe").is_err());
627    }
628
629    #[test]
630    fn the_next_page_starts_where_this_one_ends_until_the_total_is_reached() {
631        assert_eq!(next_offset(0, 50, 120), Some(50));
632        assert_eq!(next_offset(100, 20, 120), None);
633        assert_eq!(next_offset(0, 0, 0), None);
634        assert_eq!(next_offset(u32::MAX, 5, u32::MAX), None);
635    }
636}