1use axum::{Json, extract::State, http::StatusCode, response::IntoResponse};
22use chrono::{DateTime, FixedOffset, NaiveDate, TimeDelta, Utc};
23use serde::{Deserialize, Serialize};
24use sqlx::{Postgres, Transaction};
25use uuid::Uuid;
26
27use crate::{
28 app::AppState,
29 auth::AuthenticatedAgent,
30 calendar::WorkdayKind,
31 error::ApiError,
32 privacy::{Dropped, Policy, PrivacyLevel},
33 webhooks::{self, DayPayload, Event, EventKind, Webhooks},
34};
35
36const CLOSED_DAY_IS_NEWS_FOR: TimeDelta = TimeDelta::hours(24);
43
44#[derive(Debug, Deserialize)]
46pub struct DayUpload {
47 pub date: NaiveDate,
51 pub started_at: DateTime<FixedOffset>,
53 #[serde(default)]
55 pub ended_at: Option<DateTime<FixedOffset>>,
56 #[serde(default)]
57 pub pauses: Vec<PauseUpload>,
58 #[serde(default)]
59 pub tasks: Vec<TaskUpload>,
60 #[serde(default)]
70 pub tasks_are_complete: bool,
71 #[serde(default)]
79 pub kind: WorkdayKind,
80}
81
82#[derive(Debug, Deserialize)]
83pub struct PauseUpload {
84 pub started_at: DateTime<FixedOffset>,
85 #[serde(default)]
86 pub ended_at: Option<DateTime<FixedOffset>>,
87 #[serde(default)]
90 pub duration_seconds: Option<i32>,
91 #[serde(default)]
93 pub manual: bool,
94 #[serde(default)]
95 pub reason: Option<String>,
96}
97
98#[derive(Debug, Deserialize)]
99pub struct TaskUpload {
100 pub agent_task_id: i32,
103 #[serde(default)]
107 pub agent_group_id: Option<i32>,
108 pub recorded_at: DateTime<FixedOffset>,
109 pub name: String,
110 #[serde(default)]
111 pub comment: Option<String>,
112 pub completeness: i16,
114}
115
116#[derive(Debug, Serialize)]
118pub struct DayAccepted {
119 pub workday_id: Uuid,
120 pub date: NaiveDate,
121 pub kind: WorkdayKind,
125 pub pauses: usize,
126 pub tasks: usize,
127 pub deleted_tasks: u64,
130 pub privacy_level: PrivacyLevel,
134 #[serde(skip_serializing_if = "Dropped::is_empty")]
137 pub discarded: Dropped,
138}
139
140#[derive(Debug, Deserialize)]
142pub struct BatchUpload {
143 pub days: Vec<DayUpload>,
144}
145
146#[derive(Debug, Serialize)]
150pub struct BatchResult {
151 pub accepted: usize,
152 pub rejected: usize,
153 pub results: Vec<DayResult>,
154}
155
156#[derive(Debug, Serialize)]
158#[serde(tag = "status", rename_all = "lowercase")]
159pub enum DayResult {
160 Accepted {
161 #[serde(flatten)]
162 day: DayAccepted,
163 },
164 Rejected { date: NaiveDate, error: String },
168}
169
170pub async fn upload_day(State(state): State<AppState>, agent: AuthenticatedAgent, Json(day): Json<DayUpload>) -> Result<impl IntoResponse, ApiError> {
172 let policy = Policy::load(&state.pool).await?;
173 let accepted = store_day(&state.pool, agent, &day, policy, &state.webhooks, Utc::now()).await?;
174 Ok((StatusCode::OK, Json(accepted)))
175}
176
177pub async fn upload_batch(State(state): State<AppState>, agent: AuthenticatedAgent, Json(batch): Json<BatchUpload>) -> Result<impl IntoResponse, ApiError> {
183 if batch.days.len() > state.max_batch_days {
184 return Err(ApiError::new(
185 StatusCode::PAYLOAD_TOO_LARGE,
186 format!("a batch carries at most {} days; split the backlog", state.max_batch_days),
187 ));
188 }
189
190 let policy = Policy::load(&state.pool).await?;
194 let now = Utc::now();
195
196 let mut results = Vec::with_capacity(batch.days.len());
197 let mut accepted = 0;
198 let mut rejected = 0;
199
200 for day in &batch.days {
201 match store_day(&state.pool, agent, day, policy, &state.webhooks, now).await {
202 Ok(stored) => {
203 accepted += 1;
204 results.push(DayResult::Accepted { day: stored });
205 }
206 Err(error) if error.status().is_server_error() => return Err(error),
210 Err(error) => {
211 rejected += 1;
212 results.push(DayResult::Rejected {
213 date: day.date,
214 error: error.to_string(),
215 });
216 }
217 }
218 }
219
220 tracing::info!(user_id = %agent.user_id, agent_id = %agent.agent_id, accepted, rejected, "accepted a batch");
221
222 Ok((StatusCode::OK, Json(BatchResult { accepted, rejected, results })))
223}
224
225async fn store_day(
230 pool: &sqlx::PgPool,
231 agent: AuthenticatedAgent,
232 day: &DayUpload,
233 policy: Policy,
234 webhooks: &Webhooks,
235 now: DateTime<Utc>,
236) -> Result<DayAccepted, ApiError> {
237 validate(day)?;
238
239 let paused_seconds = pause_totals(&day.pauses).seconds;
243
244 let level = policy.level();
248 let (day, discarded, pause_totals) = filter(day, level);
249 let day = &day;
250
251 let mut tx = pool.begin().await?;
254
255 let was_closed: Option<bool> = sqlx::query_scalar("SELECT ended_at IS NOT NULL FROM workdays WHERE user_id = $1 AND date = $2 FOR UPDATE")
260 .bind(agent.user_id)
261 .bind(day.date)
262 .fetch_optional(&mut *tx)
263 .await?;
264
265 let workday_id = upsert_workday(&mut tx, agent.user_id, day, pause_totals).await?;
266 replace_pauses(&mut tx, workday_id, &day.pauses).await?;
267 let tasks = upsert_tasks(&mut tx, agent.user_id, day.date, &day.tasks).await?;
268 let deleted_tasks = if day.tasks_are_complete || !level.keeps_tasks() {
274 delete_missing_tasks(&mut tx, agent.user_id, day.date, &day.tasks).await?
275 } else {
276 0
277 };
278
279 if let Some(ended_at) = day.ended_at.map(|at| at.with_timezone(&Utc))
280 && was_closed != Some(true)
281 && now - ended_at <= CLOSED_DAY_IS_NEWS_FOR
282 && webhooks.anyone_hears(EventKind::DayClosed)
283 {
284 let started_at = day.started_at.with_timezone(&Utc);
285 let closed = DayPayload {
286 date: day.date,
287 kind: day.kind,
288 started_at,
289 ended_at,
290 worked_seconds: ((ended_at - started_at).num_seconds() - i64::from(paused_seconds)).max(0),
291 };
292 let person = webhooks::person(&mut tx, agent.user_id).await?;
293 webhooks::enqueue(&mut tx, webhooks, &Event::day_closed(workday_id, closed, person, webhooks, now)).await?;
294 }
295
296 tx.commit().await?;
297
298 tracing::info!(%workday_id, user_id = %agent.user_id, agent_id = %agent.agent_id, date = %day.date, pauses = day.pauses.len(), tasks, deleted_tasks, "accepted a day");
301
302 Ok(DayAccepted {
303 workday_id,
304 date: day.date,
305 kind: day.kind,
306 pauses: day.pauses.len(),
307 tasks,
308 deleted_tasks,
309 privacy_level: level,
310 discarded,
311 })
312}
313
314fn filter(day: &DayUpload, level: PrivacyLevel) -> (DayUpload, Dropped, Option<PauseTotals>) {
326 let mut discarded = Dropped::default();
327
328 let pauses = if level.keeps_pause_times() {
329 day.pauses
330 .iter()
331 .map(|pause| PauseUpload {
332 started_at: pause.started_at,
333 ended_at: pause.ended_at,
334 duration_seconds: pause.duration_seconds,
335 manual: pause.manual,
336 reason: match (&pause.reason, level.keeps_free_text()) {
337 (Some(reason), false) if !reason.is_empty() => {
338 discarded.free_text += 1;
339 None
340 }
341 (reason, true) => reason.clone(),
342 _ => None,
343 },
344 })
345 .collect()
346 } else {
347 discarded.pauses = day.pauses.len();
351 Vec::new()
352 };
353
354 let tasks = if level.keeps_tasks() {
355 day.tasks
356 .iter()
357 .map(|task| TaskUpload {
358 agent_task_id: task.agent_task_id,
359 agent_group_id: task.agent_group_id,
360 recorded_at: task.recorded_at,
361 name: task.name.clone(),
362 comment: match (&task.comment, level.keeps_free_text()) {
363 (Some(comment), false) if !comment.is_empty() => {
364 discarded.free_text += 1;
365 None
366 }
367 (comment, true) => comment.clone(),
368 _ => None,
369 },
370 completeness: task.completeness,
371 })
372 .collect()
373 } else {
374 discarded.tasks = day.tasks.len();
375 Vec::new()
376 };
377
378 let filtered = DayUpload {
384 date: day.date,
385 started_at: day.started_at,
386 ended_at: day.ended_at,
387 pauses,
388 tasks,
389 tasks_are_complete: day.tasks_are_complete,
390 kind: day.kind,
395 };
396
397 let totals = (!level.keeps_pause_times()).then(|| pause_totals(&day.pauses));
398
399 (filtered, discarded, totals)
400}
401
402#[derive(Debug, Clone, Copy, PartialEq, Eq)]
404struct PauseTotals {
405 count: i32,
406 seconds: i32,
407}
408
409fn pause_totals(pauses: &[PauseUpload]) -> PauseTotals {
416 let count = pauses.len() as i32;
417 let seconds = pauses
418 .iter()
419 .map(|pause| {
420 pause.duration_seconds.unwrap_or_else(|| {
421 pause
425 .ended_at
426 .map(|ended| (ended - pause.started_at).num_seconds().clamp(0, i32::MAX as i64) as i32)
427 .unwrap_or(0)
428 })
429 })
430 .fold(0i32, |total, seconds| total.saturating_add(seconds));
431
432 PauseTotals { count, seconds }
433}
434
435fn validate(day: &DayUpload) -> Result<(), ApiError> {
439 if let Some(ended_at) = day.ended_at
440 && ended_at < day.started_at
441 {
442 return Err(ApiError::bad_request("ended_at is before started_at"));
443 }
444
445 for (index, pause) in day.pauses.iter().enumerate() {
446 if let Some(ended_at) = pause.ended_at
447 && ended_at < pause.started_at
448 {
449 return Err(ApiError::bad_request(format!("pauses[{index}]: ended_at is before started_at")));
450 }
451 if pause.duration_seconds.is_some_and(|seconds| seconds < 0) {
452 return Err(ApiError::bad_request(format!("pauses[{index}]: duration_seconds is negative")));
453 }
454 }
455
456 for (index, task) in day.tasks.iter().enumerate() {
457 if !(0..=100).contains(&task.completeness) {
458 return Err(ApiError::bad_request(format!("tasks[{index}]: completeness must be between 0 and 100")));
459 }
460 if task.name.trim().is_empty() {
461 return Err(ApiError::bad_request(format!("tasks[{index}]: name is empty")));
462 }
463 }
464
465 Ok(())
466}
467
468async fn upsert_workday(tx: &mut Transaction<'_, Postgres>, user_id: Uuid, day: &DayUpload, pause_totals: Option<PauseTotals>) -> Result<Uuid, ApiError> {
474 let workday_id: Uuid = sqlx::query_scalar(
475 "INSERT INTO workdays (user_id, date, started_at, ended_at, paused_count, paused_seconds, kind) VALUES ($1, $2, $3, $4, $5, $6, $7)
476 ON CONFLICT (user_id, date) DO UPDATE SET
477 started_at = EXCLUDED.started_at,
478 ended_at = EXCLUDED.ended_at,
479 paused_count = EXCLUDED.paused_count,
480 paused_seconds = EXCLUDED.paused_seconds,
481 kind = EXCLUDED.kind
482 RETURNING id",
483 )
484 .bind(user_id)
485 .bind(day.date)
486 .bind(day.started_at.with_timezone(&Utc))
487 .bind(day.ended_at.map(|at| at.with_timezone(&Utc)))
488 .bind(pause_totals.map(|totals| totals.count))
489 .bind(pause_totals.map(|totals| totals.seconds))
490 .bind(day.kind)
491 .fetch_one(&mut **tx)
492 .await?;
493
494 Ok(workday_id)
495}
496
497async fn replace_pauses(tx: &mut Transaction<'_, Postgres>, workday_id: Uuid, pauses: &[PauseUpload]) -> Result<(), ApiError> {
503 sqlx::query("DELETE FROM pauses WHERE workday_id = $1")
504 .bind(workday_id)
505 .execute(&mut **tx)
506 .await?;
507
508 for pause in pauses {
509 sqlx::query("INSERT INTO pauses (workday_id, started_at, ended_at, duration_seconds, manual, reason) VALUES ($1, $2, $3, $4, $5, $6)")
510 .bind(workday_id)
511 .bind(pause.started_at.with_timezone(&Utc))
512 .bind(pause.ended_at.map(|at| at.with_timezone(&Utc)))
513 .bind(pause.duration_seconds)
514 .bind(pause.manual)
515 .bind(pause.reason.as_deref())
516 .execute(&mut **tx)
517 .await?;
518 }
519
520 Ok(())
521}
522
523async fn upsert_tasks(tx: &mut Transaction<'_, Postgres>, user_id: Uuid, date: NaiveDate, tasks: &[TaskUpload]) -> Result<usize, ApiError> {
529 for task in tasks {
530 sqlx::query(
531 "INSERT INTO tasks (user_id, agent_task_id, agent_group_id, date, recorded_at, name, comment, completeness)
532 VALUES ($1, $2, $3, $4, $5, $6, $7, $8)
533 ON CONFLICT (user_id, agent_task_id) DO UPDATE SET
534 agent_group_id = EXCLUDED.agent_group_id,
535 date = EXCLUDED.date,
536 recorded_at = EXCLUDED.recorded_at,
537 name = EXCLUDED.name,
538 comment = EXCLUDED.comment,
539 completeness = EXCLUDED.completeness",
540 )
541 .bind(user_id)
542 .bind(task.agent_task_id)
543 .bind(task.agent_group_id.unwrap_or(task.agent_task_id))
544 .bind(date)
545 .bind(task.recorded_at.with_timezone(&Utc))
546 .bind(task.name.trim())
547 .bind(task.comment.as_deref())
548 .bind(task.completeness)
549 .execute(&mut **tx)
550 .await?;
551 }
552
553 Ok(tasks.len())
554}
555
556async fn delete_missing_tasks(tx: &mut Transaction<'_, Postgres>, user_id: Uuid, date: NaiveDate, tasks: &[TaskUpload]) -> Result<u64, ApiError> {
562 let kept: Vec<i32> = tasks.iter().map(|task| task.agent_task_id).collect();
563
564 let deleted = sqlx::query("DELETE FROM tasks WHERE user_id = $1 AND date = $2 AND agent_task_id <> ALL($3)")
565 .bind(user_id)
566 .bind(date)
567 .bind(&kept)
568 .execute(&mut **tx)
569 .await?
570 .rows_affected();
571
572 Ok(deleted)
573}
574
575#[cfg(test)]
576mod tests {
577 use super::*;
578
579 fn day_json(patch: serde_json::Value) -> DayUpload {
580 let mut value = serde_json::json!({
581 "date": "2026-08-14",
582 "started_at": "2026-08-14T09:00:00-03:00",
583 "pauses": [],
584 "tasks": [],
585 });
586 let (serde_json::Value::Object(base), serde_json::Value::Object(patch)) = (&mut value, patch) else {
587 panic!("both must be objects");
588 };
589 base.extend(patch);
590 serde_json::from_value(value).expect("the fixture should deserialize")
591 }
592
593 #[test]
594 fn an_offset_is_required_on_every_instant() {
595 let bare = serde_json::json!({
598 "date": "2026-08-14",
599 "started_at": "2026-08-14T09:00:00",
600 });
601 assert!(
602 serde_json::from_value::<DayUpload>(bare).is_err(),
603 "an instant without an offset must be rejected"
604 );
605 }
606
607 #[test]
608 fn the_offset_is_preserved_as_an_instant() {
609 let day = day_json(serde_json::json!({ "started_at": "2026-08-14T09:00:00-03:00" }));
610 assert_eq!(day.started_at.with_timezone(&Utc).to_rfc3339(), "2026-08-14T12:00:00+00:00");
611 }
612
613 #[test]
614 fn a_day_may_still_be_open() {
615 let day = day_json(serde_json::json!({}));
616 assert!(day.ended_at.is_none(), "a missing ended_at means the day is still running");
617 validate(&day).expect("an open day is valid");
618 }
619
620 #[test]
621 fn a_day_cannot_end_before_it_starts() {
622 let day = day_json(serde_json::json!({ "ended_at": "2026-08-14T08:00:00-03:00" }));
623 let error = validate(&day).expect_err("a backwards day must be refused");
624 assert_eq!(error.to_string(), "ended_at is before started_at");
625 }
626
627 #[test]
628 fn impossible_pauses_and_tasks_are_named_in_the_error() {
629 let day = day_json(serde_json::json!({
630 "pauses": [
631 {"started_at": "2026-08-14T10:00:00-03:00", "ended_at": "2026-08-14T10:20:00-03:00", "duration_seconds": 1200},
632 {"started_at": "2026-08-14T12:00:00-03:00", "duration_seconds": -1},
633 ],
634 }));
635 let error = validate(&day).expect_err("a negative duration must be refused");
636 assert!(
637 error.to_string().contains("pauses[1]"),
638 "the message should point at the offending element: {error}"
639 );
640
641 let day = day_json(serde_json::json!({
642 "tasks": [{"agent_task_id": 1, "recorded_at": "2026-08-14T17:00:00-03:00", "name": "Ship it", "completeness": 101}],
643 }));
644 let error = validate(&day).expect_err("completeness above 100 must be refused");
645 assert!(
646 error.to_string().contains("tasks[0]"),
647 "the message should point at the offending element: {error}"
648 );
649 }
650
651 #[test]
652 fn a_task_group_defaults_to_the_task_itself() {
653 let day = day_json(serde_json::json!({
654 "tasks": [{"agent_task_id": 7, "recorded_at": "2026-08-14T17:00:00-03:00", "name": "Write the ingest", "completeness": 60}],
655 }));
656 let task = &day.tasks[0];
657 assert_eq!(task.agent_group_id, None, "an absent group is absent on the wire");
658 assert_eq!(task.agent_group_id.unwrap_or(task.agent_task_id), 7, "and resolves to the task itself");
659 }
660
661 #[test]
662 fn an_agent_that_says_nothing_deletes_nothing() {
663 let day = day_json(serde_json::json!({}));
667 assert!(!day.tasks_are_complete, "the authoritative set must be opt-in");
668
669 let day = day_json(serde_json::json!({ "tasks_are_complete": true }));
670 assert!(day.tasks_are_complete, "and an agent that opts in is heard");
671 }
672
673 #[test]
674 fn a_day_without_a_kind_is_a_worked_day() {
675 let day = day_json(serde_json::json!({}));
680 assert_eq!(day.kind, WorkdayKind::Work);
681
682 let day = day_json(serde_json::json!({ "kind": "vacation" }));
683 assert_eq!(day.kind, WorkdayKind::Vacation, "and an agent that says so is heard");
684 }
685
686 #[test]
687 fn a_nameless_task_is_refused() {
688 let day = day_json(serde_json::json!({
689 "tasks": [{"agent_task_id": 1, "recorded_at": "2026-08-14T17:00:00-03:00", "name": " ", "completeness": 50}],
690 }));
691 assert!(validate(&day).is_err(), "a task with a blank name carries no information");
692 }
693}