1use 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
37pub const DEFAULT_LIMIT: u32 = 50;
39pub const MAX_LIMIT: u32 = 200;
41const EXAMPLES_PER_CHECK: i64 = 3;
43
44#[derive(Debug, Clone, Copy, PartialEq, Eq)]
47pub struct Quota {
48 pub limit: Option<u32>,
49 pub remaining: Option<u32>,
50}
51
52#[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 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 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 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 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 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 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 .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
475pub fn clamp_limit(limit: Option<u32>) -> u32 {
477 limit.unwrap_or(DEFAULT_LIMIT).clamp(1, MAX_LIMIT)
478}
479
480fn next_offset(offset: u32, shown: u32, total: u32) -> Option<u32> {
482 offset.checked_add(shown).filter(|next| *next < total)
483}
484
485pub 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
500pub 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
509pub 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#[derive(Debug, Clone, PartialEq)]
565pub struct IssueRows {
566 pub total: u32,
568 pub limit: u32,
569 pub offset: u32,
570 pub urls: Vec<UrlRow>,
571 pub next_offset: Option<u32>,
573}
574
575pub 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}