systemprompt_analytics/feedback/
service.rs1use super::FeedbackFactsRepository;
7use crate::Result;
8use systemprompt_identifiers::{AnalyticsWorkerId, UserId};
9
10#[derive(Debug, Clone)]
11pub struct FactsProcessingService {
12 repository: FeedbackFactsRepository,
13}
14
15impl FactsProcessingService {
16 pub const fn new(repository: FeedbackFactsRepository) -> Self {
17 Self { repository }
18 }
19
20 pub async fn drain(
21 &self,
22 owner: &UserId,
23 worker: &AnalyticsWorkerId,
24 limit: u32,
25 ) -> Result<usize> {
26 let leases = self.repository.claim(owner, worker, limit, 60).await?;
27 let mut completed = 0;
28 for lease in leases {
29 match self.repository.apply(owner, &lease).await {
30 Ok(_) => completed += 1,
31 Err(_error) => {
32 if self.repository.retry(owner, &lease).await.is_err() {
33 tracing::warn!(change_id = %lease.change_id, "Analytics lease lost; current lease owner will resume processing");
34 }
35 },
36 }
37 }
38 Ok(completed)
39 }
40
41 pub async fn run(
42 &self,
43 owner: &UserId,
44 worker: &AnalyticsWorkerId,
45 mut shutdown: tokio::sync::watch::Receiver<bool>,
46 ) -> Result<()> {
47 let mut interval = tokio::time::interval(std::time::Duration::from_secs(1));
48 interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
49 while !*shutdown.borrow() {
50 tokio::select! {
51 changed = shutdown.changed() => { if changed.is_err() || *shutdown.borrow() { break; } },
52 _ = interval.tick() => {
53 if self.drain(owner, worker, 64).await.is_err() {
54 tracing::warn!(worker_id = %worker, "Analytics worker storage unavailable; committed changes remain pending");
55 }
56 },
57 }
58 }
59 Ok(())
60 }
61}