use axum::{Json, extract::State, http::StatusCode, response::IntoResponse};
use chrono::{DateTime, FixedOffset, NaiveDate, TimeDelta, Utc};
use serde::{Deserialize, Serialize};
use sqlx::{Postgres, Transaction};
use uuid::Uuid;
use crate::{
app::AppState,
auth::AuthenticatedAgent,
calendar::WorkdayKind,
error::ApiError,
privacy::{Dropped, Policy, PrivacyLevel},
webhooks::{self, DayPayload, Event, EventKind, Webhooks},
};
const CLOSED_DAY_IS_NEWS_FOR: TimeDelta = TimeDelta::hours(24);
#[derive(Debug, Deserialize)]
pub struct DayUpload {
pub date: NaiveDate,
pub started_at: DateTime<FixedOffset>,
#[serde(default)]
pub ended_at: Option<DateTime<FixedOffset>>,
#[serde(default)]
pub pauses: Vec<PauseUpload>,
#[serde(default)]
pub tasks: Vec<TaskUpload>,
#[serde(default)]
pub tasks_are_complete: bool,
#[serde(default)]
pub kind: WorkdayKind,
}
#[derive(Debug, Deserialize)]
pub struct PauseUpload {
pub started_at: DateTime<FixedOffset>,
#[serde(default)]
pub ended_at: Option<DateTime<FixedOffset>>,
#[serde(default)]
pub duration_seconds: Option<i32>,
#[serde(default)]
pub manual: bool,
#[serde(default)]
pub reason: Option<String>,
}
#[derive(Debug, Deserialize)]
pub struct TaskUpload {
pub agent_task_id: i32,
#[serde(default)]
pub agent_group_id: Option<i32>,
pub recorded_at: DateTime<FixedOffset>,
pub name: String,
#[serde(default)]
pub comment: Option<String>,
pub completeness: i16,
}
#[derive(Debug, Serialize)]
pub struct DayAccepted {
pub workday_id: Uuid,
pub date: NaiveDate,
pub kind: WorkdayKind,
pub pauses: usize,
pub tasks: usize,
pub deleted_tasks: u64,
pub privacy_level: PrivacyLevel,
#[serde(skip_serializing_if = "Dropped::is_empty")]
pub discarded: Dropped,
}
#[derive(Debug, Deserialize)]
pub struct BatchUpload {
pub days: Vec<DayUpload>,
}
#[derive(Debug, Serialize)]
pub struct BatchResult {
pub accepted: usize,
pub rejected: usize,
pub results: Vec<DayResult>,
}
#[derive(Debug, Serialize)]
#[serde(tag = "status", rename_all = "lowercase")]
pub enum DayResult {
Accepted {
#[serde(flatten)]
day: DayAccepted,
},
Rejected { date: NaiveDate, error: String },
}
pub async fn upload_day(State(state): State<AppState>, agent: AuthenticatedAgent, Json(day): Json<DayUpload>) -> Result<impl IntoResponse, ApiError> {
let policy = Policy::load(&state.pool).await?;
let accepted = store_day(&state.pool, agent, &day, policy, &state.webhooks, Utc::now()).await?;
Ok((StatusCode::OK, Json(accepted)))
}
pub async fn upload_batch(State(state): State<AppState>, agent: AuthenticatedAgent, Json(batch): Json<BatchUpload>) -> Result<impl IntoResponse, ApiError> {
if batch.days.len() > state.max_batch_days {
return Err(ApiError::new(
StatusCode::PAYLOAD_TOO_LARGE,
format!("a batch carries at most {} days; split the backlog", state.max_batch_days),
));
}
let policy = Policy::load(&state.pool).await?;
let now = Utc::now();
let mut results = Vec::with_capacity(batch.days.len());
let mut accepted = 0;
let mut rejected = 0;
for day in &batch.days {
match store_day(&state.pool, agent, day, policy, &state.webhooks, now).await {
Ok(stored) => {
accepted += 1;
results.push(DayResult::Accepted { day: stored });
}
Err(error) if error.status().is_server_error() => return Err(error),
Err(error) => {
rejected += 1;
results.push(DayResult::Rejected {
date: day.date,
error: error.to_string(),
});
}
}
}
tracing::info!(user_id = %agent.user_id, agent_id = %agent.agent_id, accepted, rejected, "accepted a batch");
Ok((StatusCode::OK, Json(BatchResult { accepted, rejected, results })))
}
async fn store_day(
pool: &sqlx::PgPool,
agent: AuthenticatedAgent,
day: &DayUpload,
policy: Policy,
webhooks: &Webhooks,
now: DateTime<Utc>,
) -> Result<DayAccepted, ApiError> {
validate(day)?;
let paused_seconds = pause_totals(&day.pauses).seconds;
let level = policy.level();
let (day, discarded, pause_totals) = filter(day, level);
let day = &day;
let mut tx = pool.begin().await?;
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")
.bind(agent.user_id)
.bind(day.date)
.fetch_optional(&mut *tx)
.await?;
let workday_id = upsert_workday(&mut tx, agent.user_id, day, pause_totals).await?;
replace_pauses(&mut tx, workday_id, &day.pauses).await?;
let tasks = upsert_tasks(&mut tx, agent.user_id, day.date, &day.tasks).await?;
let deleted_tasks = if day.tasks_are_complete || !level.keeps_tasks() {
delete_missing_tasks(&mut tx, agent.user_id, day.date, &day.tasks).await?
} else {
0
};
if let Some(ended_at) = day.ended_at.map(|at| at.with_timezone(&Utc))
&& was_closed != Some(true)
&& now - ended_at <= CLOSED_DAY_IS_NEWS_FOR
&& webhooks.anyone_hears(EventKind::DayClosed)
{
let started_at = day.started_at.with_timezone(&Utc);
let closed = DayPayload {
date: day.date,
kind: day.kind,
started_at,
ended_at,
worked_seconds: ((ended_at - started_at).num_seconds() - i64::from(paused_seconds)).max(0),
};
let person = webhooks::person(&mut tx, agent.user_id).await?;
webhooks::enqueue(&mut tx, webhooks, &Event::day_closed(workday_id, closed, person, webhooks, now)).await?;
}
tx.commit().await?;
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");
Ok(DayAccepted {
workday_id,
date: day.date,
kind: day.kind,
pauses: day.pauses.len(),
tasks,
deleted_tasks,
privacy_level: level,
discarded,
})
}
fn filter(day: &DayUpload, level: PrivacyLevel) -> (DayUpload, Dropped, Option<PauseTotals>) {
let mut discarded = Dropped::default();
let pauses = if level.keeps_pause_times() {
day.pauses
.iter()
.map(|pause| PauseUpload {
started_at: pause.started_at,
ended_at: pause.ended_at,
duration_seconds: pause.duration_seconds,
manual: pause.manual,
reason: match (&pause.reason, level.keeps_free_text()) {
(Some(reason), false) if !reason.is_empty() => {
discarded.free_text += 1;
None
}
(reason, true) => reason.clone(),
_ => None,
},
})
.collect()
} else {
discarded.pauses = day.pauses.len();
Vec::new()
};
let tasks = if level.keeps_tasks() {
day.tasks
.iter()
.map(|task| TaskUpload {
agent_task_id: task.agent_task_id,
agent_group_id: task.agent_group_id,
recorded_at: task.recorded_at,
name: task.name.clone(),
comment: match (&task.comment, level.keeps_free_text()) {
(Some(comment), false) if !comment.is_empty() => {
discarded.free_text += 1;
None
}
(comment, true) => comment.clone(),
_ => None,
},
completeness: task.completeness,
})
.collect()
} else {
discarded.tasks = day.tasks.len();
Vec::new()
};
let filtered = DayUpload {
date: day.date,
started_at: day.started_at,
ended_at: day.ended_at,
pauses,
tasks,
tasks_are_complete: day.tasks_are_complete,
kind: day.kind,
};
let totals = (!level.keeps_pause_times()).then(|| pause_totals(&day.pauses));
(filtered, discarded, totals)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
struct PauseTotals {
count: i32,
seconds: i32,
}
fn pause_totals(pauses: &[PauseUpload]) -> PauseTotals {
let count = pauses.len() as i32;
let seconds = pauses
.iter()
.map(|pause| {
pause.duration_seconds.unwrap_or_else(|| {
pause
.ended_at
.map(|ended| (ended - pause.started_at).num_seconds().clamp(0, i32::MAX as i64) as i32)
.unwrap_or(0)
})
})
.fold(0i32, |total, seconds| total.saturating_add(seconds));
PauseTotals { count, seconds }
}
fn validate(day: &DayUpload) -> Result<(), ApiError> {
if let Some(ended_at) = day.ended_at
&& ended_at < day.started_at
{
return Err(ApiError::bad_request("ended_at is before started_at"));
}
for (index, pause) in day.pauses.iter().enumerate() {
if let Some(ended_at) = pause.ended_at
&& ended_at < pause.started_at
{
return Err(ApiError::bad_request(format!("pauses[{index}]: ended_at is before started_at")));
}
if pause.duration_seconds.is_some_and(|seconds| seconds < 0) {
return Err(ApiError::bad_request(format!("pauses[{index}]: duration_seconds is negative")));
}
}
for (index, task) in day.tasks.iter().enumerate() {
if !(0..=100).contains(&task.completeness) {
return Err(ApiError::bad_request(format!("tasks[{index}]: completeness must be between 0 and 100")));
}
if task.name.trim().is_empty() {
return Err(ApiError::bad_request(format!("tasks[{index}]: name is empty")));
}
}
Ok(())
}
async fn upsert_workday(tx: &mut Transaction<'_, Postgres>, user_id: Uuid, day: &DayUpload, pause_totals: Option<PauseTotals>) -> Result<Uuid, ApiError> {
let workday_id: Uuid = sqlx::query_scalar(
"INSERT INTO workdays (user_id, date, started_at, ended_at, paused_count, paused_seconds, kind) VALUES ($1, $2, $3, $4, $5, $6, $7)
ON CONFLICT (user_id, date) DO UPDATE SET
started_at = EXCLUDED.started_at,
ended_at = EXCLUDED.ended_at,
paused_count = EXCLUDED.paused_count,
paused_seconds = EXCLUDED.paused_seconds,
kind = EXCLUDED.kind
RETURNING id",
)
.bind(user_id)
.bind(day.date)
.bind(day.started_at.with_timezone(&Utc))
.bind(day.ended_at.map(|at| at.with_timezone(&Utc)))
.bind(pause_totals.map(|totals| totals.count))
.bind(pause_totals.map(|totals| totals.seconds))
.bind(day.kind)
.fetch_one(&mut **tx)
.await?;
Ok(workday_id)
}
async fn replace_pauses(tx: &mut Transaction<'_, Postgres>, workday_id: Uuid, pauses: &[PauseUpload]) -> Result<(), ApiError> {
sqlx::query("DELETE FROM pauses WHERE workday_id = $1")
.bind(workday_id)
.execute(&mut **tx)
.await?;
for pause in pauses {
sqlx::query("INSERT INTO pauses (workday_id, started_at, ended_at, duration_seconds, manual, reason) VALUES ($1, $2, $3, $4, $5, $6)")
.bind(workday_id)
.bind(pause.started_at.with_timezone(&Utc))
.bind(pause.ended_at.map(|at| at.with_timezone(&Utc)))
.bind(pause.duration_seconds)
.bind(pause.manual)
.bind(pause.reason.as_deref())
.execute(&mut **tx)
.await?;
}
Ok(())
}
async fn upsert_tasks(tx: &mut Transaction<'_, Postgres>, user_id: Uuid, date: NaiveDate, tasks: &[TaskUpload]) -> Result<usize, ApiError> {
for task in tasks {
sqlx::query(
"INSERT INTO tasks (user_id, agent_task_id, agent_group_id, date, recorded_at, name, comment, completeness)
VALUES ($1, $2, $3, $4, $5, $6, $7, $8)
ON CONFLICT (user_id, agent_task_id) DO UPDATE SET
agent_group_id = EXCLUDED.agent_group_id,
date = EXCLUDED.date,
recorded_at = EXCLUDED.recorded_at,
name = EXCLUDED.name,
comment = EXCLUDED.comment,
completeness = EXCLUDED.completeness",
)
.bind(user_id)
.bind(task.agent_task_id)
.bind(task.agent_group_id.unwrap_or(task.agent_task_id))
.bind(date)
.bind(task.recorded_at.with_timezone(&Utc))
.bind(task.name.trim())
.bind(task.comment.as_deref())
.bind(task.completeness)
.execute(&mut **tx)
.await?;
}
Ok(tasks.len())
}
async fn delete_missing_tasks(tx: &mut Transaction<'_, Postgres>, user_id: Uuid, date: NaiveDate, tasks: &[TaskUpload]) -> Result<u64, ApiError> {
let kept: Vec<i32> = tasks.iter().map(|task| task.agent_task_id).collect();
let deleted = sqlx::query("DELETE FROM tasks WHERE user_id = $1 AND date = $2 AND agent_task_id <> ALL($3)")
.bind(user_id)
.bind(date)
.bind(&kept)
.execute(&mut **tx)
.await?
.rows_affected();
Ok(deleted)
}
#[cfg(test)]
mod tests {
use super::*;
fn day_json(patch: serde_json::Value) -> DayUpload {
let mut value = serde_json::json!({
"date": "2026-08-14",
"started_at": "2026-08-14T09:00:00-03:00",
"pauses": [],
"tasks": [],
});
let (serde_json::Value::Object(base), serde_json::Value::Object(patch)) = (&mut value, patch) else {
panic!("both must be objects");
};
base.extend(patch);
serde_json::from_value(value).expect("the fixture should deserialize")
}
#[test]
fn an_offset_is_required_on_every_instant() {
let bare = serde_json::json!({
"date": "2026-08-14",
"started_at": "2026-08-14T09:00:00",
});
assert!(
serde_json::from_value::<DayUpload>(bare).is_err(),
"an instant without an offset must be rejected"
);
}
#[test]
fn the_offset_is_preserved_as_an_instant() {
let day = day_json(serde_json::json!({ "started_at": "2026-08-14T09:00:00-03:00" }));
assert_eq!(day.started_at.with_timezone(&Utc).to_rfc3339(), "2026-08-14T12:00:00+00:00");
}
#[test]
fn a_day_may_still_be_open() {
let day = day_json(serde_json::json!({}));
assert!(day.ended_at.is_none(), "a missing ended_at means the day is still running");
validate(&day).expect("an open day is valid");
}
#[test]
fn a_day_cannot_end_before_it_starts() {
let day = day_json(serde_json::json!({ "ended_at": "2026-08-14T08:00:00-03:00" }));
let error = validate(&day).expect_err("a backwards day must be refused");
assert_eq!(error.to_string(), "ended_at is before started_at");
}
#[test]
fn impossible_pauses_and_tasks_are_named_in_the_error() {
let day = day_json(serde_json::json!({
"pauses": [
{"started_at": "2026-08-14T10:00:00-03:00", "ended_at": "2026-08-14T10:20:00-03:00", "duration_seconds": 1200},
{"started_at": "2026-08-14T12:00:00-03:00", "duration_seconds": -1},
],
}));
let error = validate(&day).expect_err("a negative duration must be refused");
assert!(
error.to_string().contains("pauses[1]"),
"the message should point at the offending element: {error}"
);
let day = day_json(serde_json::json!({
"tasks": [{"agent_task_id": 1, "recorded_at": "2026-08-14T17:00:00-03:00", "name": "Ship it", "completeness": 101}],
}));
let error = validate(&day).expect_err("completeness above 100 must be refused");
assert!(
error.to_string().contains("tasks[0]"),
"the message should point at the offending element: {error}"
);
}
#[test]
fn a_task_group_defaults_to_the_task_itself() {
let day = day_json(serde_json::json!({
"tasks": [{"agent_task_id": 7, "recorded_at": "2026-08-14T17:00:00-03:00", "name": "Write the ingest", "completeness": 60}],
}));
let task = &day.tasks[0];
assert_eq!(task.agent_group_id, None, "an absent group is absent on the wire");
assert_eq!(task.agent_group_id.unwrap_or(task.agent_task_id), 7, "and resolves to the task itself");
}
#[test]
fn an_agent_that_says_nothing_deletes_nothing() {
let day = day_json(serde_json::json!({}));
assert!(!day.tasks_are_complete, "the authoritative set must be opt-in");
let day = day_json(serde_json::json!({ "tasks_are_complete": true }));
assert!(day.tasks_are_complete, "and an agent that opts in is heard");
}
#[test]
fn a_day_without_a_kind_is_a_worked_day() {
let day = day_json(serde_json::json!({}));
assert_eq!(day.kind, WorkdayKind::Work);
let day = day_json(serde_json::json!({ "kind": "vacation" }));
assert_eq!(day.kind, WorkdayKind::Vacation, "and an agent that says so is heard");
}
#[test]
fn a_nameless_task_is_refused() {
let day = day_json(serde_json::json!({
"tasks": [{"agent_task_id": 1, "recorded_at": "2026-08-14T17:00:00-03:00", "name": " ", "completeness": 50}],
}));
assert!(validate(&day).is_err(), "a task with a blank name carries no information");
}
}