Files
doctate/doctate-client-core/src/case_store.rs
T
Brummel cdfa5ae90e Refactor timestamp handling for activity sorting
Introduce `last_recording_at` to `OnelinerEntry` and use it as the
primary sort key for cases. This ensures that cases with recent
dictation activity are prioritized, even if their initial creation date
is older.

The logic for determining a case's activity has been updated to consider
the most recent `.m4a` file's modification time (`last_recording_at`),
falling back to `updated_at` or `created_at` if necessary.

This change also refactors the timestamp formatting and handling within
the web interface to correctly display and sort cases based on their
actual last activity time, improving user experience and data relevance.
The `now_rfc3339` function is updated to strip sub-second precision for
consistent filename generation.
2026-04-18 11:33:35 +02:00

563 lines
21 KiB
Rust

//! Per-case marker store for native clients.
//!
//! Each known case (either created locally by this client or seen in a
//! server snapshot) has one JSON file: `{dir}/{case_id}.json`. The store
//! keeps an in-memory copy and broadcasts a sorted snapshot on every
//! write via [`tokio::sync::watch`].
//!
//! The merge-with-server-snapshot logic only **adds or updates** markers.
//! It never deletes: `/api/oneliners` is a sliding-window view, not a
//! full case inventory, so "missing from snapshot" is ambiguous (deleted
//! vs. fell out of the window). Cleanup is a separate, later concern.
use std::collections::HashMap;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use doctate_common::oneliners::OnelinersResponse;
use serde::{Deserialize, Serialize};
use thiserror::Error;
use time::format_description::well_known::Rfc3339;
use time::OffsetDateTime;
use tokio::sync::{watch, Mutex};
use tracing::{debug, warn};
use uuid::Uuid;
/// On-disk representation of a single case known to this client.
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
pub struct CaseMarker {
pub case_id: Uuid,
/// RFC3339 UTC; set once at creation or on first server-snapshot merge.
pub created_at: String,
/// RFC3339 UTC; bumped on every new recording + on server-snapshot
/// updates. Drives the UI sort order (newest activity first).
pub last_activity_at: String,
/// `true` once this case has been seen in any `/api/oneliners`
/// response. Reserved for future delete-propagation logic; today only
/// useful for diagnostics.
pub synced_to_server: bool,
/// `None` while transcription / oneliner generation is pending; set
/// from the server snapshot when available.
pub oneliner: Option<String>,
}
#[derive(Debug, Error)]
pub enum CaseStoreError {
#[error("I/O error at {path}: {source}")]
Io {
path: PathBuf,
#[source]
source: std::io::Error,
},
#[error("serialize marker: {0}")]
Serialize(#[from] serde_json::Error),
#[error("format timestamp: {0}")]
Format(#[from] time::error::Format),
}
type Snapshot = Arc<Vec<CaseMarker>>;
/// Shared case store. Cheap to `clone()`; clones share the same state.
#[derive(Clone)]
pub struct CaseStore {
dir: PathBuf,
state: Arc<Mutex<HashMap<Uuid, CaseMarker>>>,
tx: watch::Sender<Snapshot>,
}
impl CaseStore {
/// Open or create the store at `dir`. Reads every `*.json` file in
/// the directory as an initial marker. Returns the store plus a
/// receiver that yields sorted snapshots (newest `last_activity_at`
/// first) on every mutation.
pub async fn open(dir: PathBuf) -> Result<(Self, watch::Receiver<Snapshot>), CaseStoreError> {
tokio::fs::create_dir_all(&dir)
.await
.map_err(|e| CaseStoreError::Io {
path: dir.clone(),
source: e,
})?;
let initial = read_all_markers(&dir).await?;
let map: HashMap<Uuid, CaseMarker> =
initial.into_iter().map(|m| (m.case_id, m)).collect();
let snapshot = sorted_snapshot(&map);
let (tx, rx) = watch::channel(snapshot);
let store = Self {
dir,
state: Arc::new(Mutex::new(map)),
tx,
};
Ok((store, rx))
}
/// Append a new case the client itself just started (no server
/// knowledge yet). Writes the marker atomically and broadcasts the
/// new snapshot. If a marker with the same `case_id` already exists,
/// it is overwritten — callers must check before if that matters.
pub async fn create_local(
&self,
case_id: Uuid,
at: OffsetDateTime,
) -> Result<CaseMarker, CaseStoreError> {
let ts = format_utc(at)?;
let marker = CaseMarker {
case_id,
created_at: ts.clone(),
last_activity_at: ts,
synced_to_server: false,
oneliner: None,
};
let mut state = self.state.lock().await;
write_marker_file(&self.dir, &marker).await?;
state.insert(case_id, marker.clone());
self.broadcast(&state);
Ok(marker)
}
/// Record new activity on an existing case (e.g. a follow-up
/// recording). Creates the marker if missing — callers may invoke
/// this after an upload without caring whether `create_local` ran
/// first. `synced_to_server` is preserved.
pub async fn mark_activity(
&self,
case_id: Uuid,
at: OffsetDateTime,
) -> Result<(), CaseStoreError> {
let ts = format_utc(at)?;
let mut state = self.state.lock().await;
let marker = state.entry(case_id).or_insert_with(|| CaseMarker {
case_id,
created_at: ts.clone(),
last_activity_at: ts.clone(),
synced_to_server: false,
oneliner: None,
});
marker.last_activity_at = ts;
write_marker_file(&self.dir, marker).await?;
self.broadcast(&state);
Ok(())
}
/// Apply a server `/api/oneliners` response. For every entry in the
/// snapshot: insert (with `synced_to_server = true`) or update the
/// local marker (oneliner + last_activity_at from server timestamps,
/// taking the newer of local vs. server). Never deletes.
///
/// The `last_activity_at` seed from the server prefers
/// `last_recording_at` (most recent m4a mtime — the authoritative
/// "when was this case last worked on?") over `updated_at` (oneliner
/// regeneration time, which collapses to the recovery-scan moment
/// after a restart). Falls back to `created_at` if neither is set.
pub async fn merge_server_snapshot(
&self,
snapshot: &OnelinersResponse,
) -> Result<(), CaseStoreError> {
let mut state = self.state.lock().await;
for entry in &snapshot.oneliners {
let Ok(case_id) = Uuid::parse_str(&entry.case_id) else {
warn!(case_id = %entry.case_id, "server snapshot entry with invalid case_id, skipped");
continue;
};
let server_activity = entry
.last_recording_at
.clone()
.or_else(|| entry.updated_at.clone())
.unwrap_or_else(|| entry.created_at.clone());
let updated_marker = match state.get(&case_id) {
Some(existing) => {
let mut m = existing.clone();
m.synced_to_server = true;
m.oneliner = entry.oneliner.clone();
// Keep whichever `last_activity_at` is newer — the
// client may have recorded an addendum the server
// hasn't processed yet.
if server_activity > m.last_activity_at {
m.last_activity_at = server_activity;
}
m
}
None => CaseMarker {
case_id,
created_at: entry.created_at.clone(),
last_activity_at: server_activity,
synced_to_server: true,
oneliner: entry.oneliner.clone(),
},
};
write_marker_file(&self.dir, &updated_marker).await?;
state.insert(case_id, updated_marker);
}
self.broadcast(&state);
Ok(())
}
/// Remove every marker from memory and disk. Used on config change
/// (new API key = potentially different user).
pub async fn clear_all(&self) -> Result<(), CaseStoreError> {
let mut state = self.state.lock().await;
let mut entries =
tokio::fs::read_dir(&self.dir)
.await
.map_err(|e| CaseStoreError::Io {
path: self.dir.clone(),
source: e,
})?;
while let Ok(Some(entry)) = entries.next_entry().await {
let path = entry.path();
if path.extension().and_then(|s| s.to_str()) != Some("json") {
continue;
}
if let Err(e) = tokio::fs::remove_file(&path).await {
warn!(path = %path.display(), error = %e, "failed to remove marker during clear");
}
}
state.clear();
self.broadcast(&state);
Ok(())
}
/// Current in-memory list, for tests and inspection.
pub async fn list(&self) -> Vec<CaseMarker> {
let state = self.state.lock().await;
sorted_vec(&state)
}
fn broadcast(&self, state: &HashMap<Uuid, CaseMarker>) {
let snapshot = sorted_snapshot(state);
// receiver-count=0 is fine (nobody listening); ignore.
let _ = self.tx.send(snapshot);
}
}
fn format_utc(at: OffsetDateTime) -> Result<String, time::error::Format> {
at.to_offset(time::UtcOffset::UTC).format(&Rfc3339)
}
fn sorted_snapshot(state: &HashMap<Uuid, CaseMarker>) -> Snapshot {
Arc::new(sorted_vec(state))
}
fn sorted_vec(state: &HashMap<Uuid, CaseMarker>) -> Vec<CaseMarker> {
let mut v: Vec<CaseMarker> = state.values().cloned().collect();
v.sort_by(|a, b| b.last_activity_at.cmp(&a.last_activity_at));
v
}
async fn read_all_markers(dir: &Path) -> Result<Vec<CaseMarker>, CaseStoreError> {
let mut out = Vec::new();
let mut entries = tokio::fs::read_dir(dir).await.map_err(|e| CaseStoreError::Io {
path: dir.to_owned(),
source: e,
})?;
while let Ok(Some(entry)) = entries.next_entry().await {
let path = entry.path();
if path.extension().and_then(|s| s.to_str()) != Some("json") {
continue;
}
match read_marker_file(&path).await {
Ok(m) => out.push(m),
Err(e) => {
// A single unreadable file should not poison the store;
// log and skip so recovery of the rest proceeds.
warn!(path = %path.display(), error = %e, "skipping unreadable marker");
}
}
}
Ok(out)
}
async fn read_marker_file(path: &Path) -> Result<CaseMarker, CaseStoreError> {
let bytes = tokio::fs::read(path).await.map_err(|e| CaseStoreError::Io {
path: path.to_owned(),
source: e,
})?;
Ok(serde_json::from_slice(&bytes)?)
}
/// Write a marker atomically: temp-file + rename. The rename on the same
/// filesystem is atomic on POSIX and NTFS, so readers never see a
/// partially-written file.
async fn write_marker_file(dir: &Path, marker: &CaseMarker) -> Result<(), CaseStoreError> {
let bytes = serde_json::to_vec_pretty(marker)?;
let final_path = dir.join(format!("{}.json", marker.case_id));
let tmp_path = dir.join(format!("{}.json.tmp", marker.case_id));
tokio::fs::write(&tmp_path, &bytes)
.await
.map_err(|e| CaseStoreError::Io {
path: tmp_path.clone(),
source: e,
})?;
tokio::fs::rename(&tmp_path, &final_path)
.await
.map_err(|e| CaseStoreError::Io {
path: final_path.clone(),
source: e,
})?;
debug!(path = %final_path.display(), "marker written");
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
use doctate_common::oneliners::OnelinerEntry;
use tempfile::TempDir;
fn t(s: &str) -> OffsetDateTime {
OffsetDateTime::parse(s, &Rfc3339).unwrap()
}
fn server_entry(case_id: Uuid, oneliner: Option<&str>, created: &str) -> OnelinerEntry {
OnelinerEntry {
case_id: case_id.to_string(),
oneliner: oneliner.map(str::to_owned),
created_at: created.to_owned(),
last_recording_at: Some(created.to_owned()),
updated_at: Some(created.to_owned()),
}
}
fn server_snapshot(entries: Vec<OnelinerEntry>) -> OnelinersResponse {
OnelinersResponse {
as_of: "2026-04-18T10:00:00Z".into(),
window_hours: 72,
oneliners: entries,
}
}
#[tokio::test]
async fn create_local_persists_and_reloads() {
let tmp = TempDir::new().unwrap();
let (store, mut rx) = CaseStore::open(tmp.path().to_owned()).await.unwrap();
assert!(rx.borrow().is_empty());
let id = Uuid::new_v4();
let marker = store.create_local(id, t("2026-04-18T10:32:00Z")).await.unwrap();
assert_eq!(marker.case_id, id);
assert!(!marker.synced_to_server);
assert!(marker.oneliner.is_none());
rx.changed().await.unwrap();
let snap = rx.borrow().clone();
assert_eq!(snap.len(), 1);
assert_eq!(snap[0].case_id, id);
// Reopen: state must survive.
drop(store);
let (store2, _rx2) = CaseStore::open(tmp.path().to_owned()).await.unwrap();
assert_eq!(store2.list().await.len(), 1);
}
#[tokio::test]
async fn mark_activity_updates_timestamp() {
let tmp = TempDir::new().unwrap();
let (store, _rx) = CaseStore::open(tmp.path().to_owned()).await.unwrap();
let id = Uuid::new_v4();
store.create_local(id, t("2026-04-18T10:00:00Z")).await.unwrap();
store.mark_activity(id, t("2026-04-18T10:45:00Z")).await.unwrap();
let list = store.list().await;
assert_eq!(list.len(), 1);
assert_eq!(list[0].last_activity_at, "2026-04-18T10:45:00Z");
// `created_at` must not change.
assert_eq!(list[0].created_at, "2026-04-18T10:00:00Z");
}
#[tokio::test]
async fn mark_activity_creates_if_missing() {
let tmp = TempDir::new().unwrap();
let (store, _rx) = CaseStore::open(tmp.path().to_owned()).await.unwrap();
let id = Uuid::new_v4();
store.mark_activity(id, t("2026-04-18T10:00:00Z")).await.unwrap();
assert_eq!(store.list().await.len(), 1);
}
#[tokio::test]
async fn merge_adds_unknown_server_case() {
let tmp = TempDir::new().unwrap();
let (store, _rx) = CaseStore::open(tmp.path().to_owned()).await.unwrap();
let id = Uuid::new_v4();
let snap = server_snapshot(vec![server_entry(
id,
Some("Kniegelenk re."),
"2026-04-18T09:30:00Z",
)]);
store.merge_server_snapshot(&snap).await.unwrap();
let list = store.list().await;
assert_eq!(list.len(), 1);
assert_eq!(list[0].case_id, id);
assert_eq!(list[0].oneliner.as_deref(), Some("Kniegelenk re."));
assert!(list[0].synced_to_server);
}
#[tokio::test]
async fn merge_updates_existing_local_marker() {
let tmp = TempDir::new().unwrap();
let (store, _rx) = CaseStore::open(tmp.path().to_owned()).await.unwrap();
let id = Uuid::new_v4();
store.create_local(id, t("2026-04-18T10:00:00Z")).await.unwrap();
let snap = server_snapshot(vec![server_entry(
id,
Some("Hypertonie"),
"2026-04-18T10:00:00Z",
)]);
store.merge_server_snapshot(&snap).await.unwrap();
let list = store.list().await;
assert_eq!(list.len(), 1);
assert!(list[0].synced_to_server);
assert_eq!(list[0].oneliner.as_deref(), Some("Hypertonie"));
}
#[tokio::test]
async fn merge_keeps_local_when_missing_from_server() {
// User recorded locally, upload not yet through. Server snapshot
// is empty — marker must survive.
let tmp = TempDir::new().unwrap();
let (store, _rx) = CaseStore::open(tmp.path().to_owned()).await.unwrap();
let id = Uuid::new_v4();
store.create_local(id, t("2026-04-18T10:00:00Z")).await.unwrap();
let snap = server_snapshot(vec![]);
store.merge_server_snapshot(&snap).await.unwrap();
assert_eq!(store.list().await.len(), 1);
}
#[tokio::test]
async fn merge_keeps_synced_marker_when_missing_from_server() {
// Case was synced earlier, now out of the server's time window.
// Must not be deleted.
let tmp = TempDir::new().unwrap();
let (store, _rx) = CaseStore::open(tmp.path().to_owned()).await.unwrap();
let id = Uuid::new_v4();
store.create_local(id, t("2026-04-18T10:00:00Z")).await.unwrap();
let synced_snap = server_snapshot(vec![server_entry(
id,
Some("old"),
"2026-04-18T10:00:00Z",
)]);
store.merge_server_snapshot(&synced_snap).await.unwrap();
// Next poll: server window moved, case fell out.
store.merge_server_snapshot(&server_snapshot(vec![])).await.unwrap();
let list = store.list().await;
assert_eq!(list.len(), 1);
assert!(list[0].synced_to_server);
assert_eq!(list[0].oneliner.as_deref(), Some("old"));
}
#[tokio::test]
async fn merge_prefers_newer_last_activity() {
let tmp = TempDir::new().unwrap();
let (store, _rx) = CaseStore::open(tmp.path().to_owned()).await.unwrap();
let id = Uuid::new_v4();
store.create_local(id, t("2026-04-18T09:00:00Z")).await.unwrap();
store.mark_activity(id, t("2026-04-18T11:00:00Z")).await.unwrap();
// Server says last activity was 10:00 (stale) — local 11:00 must win.
let snap = server_snapshot(vec![server_entry(
id,
Some("x"),
"2026-04-18T10:00:00Z",
)]);
store.merge_server_snapshot(&snap).await.unwrap();
let list = store.list().await;
assert_eq!(list[0].last_activity_at, "2026-04-18T11:00:00Z");
}
#[tokio::test]
async fn clear_all_removes_everything() {
let tmp = TempDir::new().unwrap();
let (store, _rx) = CaseStore::open(tmp.path().to_owned()).await.unwrap();
store.create_local(Uuid::new_v4(), t("2026-04-18T10:00:00Z")).await.unwrap();
store.create_local(Uuid::new_v4(), t("2026-04-18T11:00:00Z")).await.unwrap();
assert_eq!(store.list().await.len(), 2);
store.clear_all().await.unwrap();
assert!(store.list().await.is_empty());
// Marker files on disk are gone.
let mut count = 0;
let mut entries = tokio::fs::read_dir(tmp.path()).await.unwrap();
while let Ok(Some(e)) = entries.next_entry().await {
if e.path().extension().and_then(|s| s.to_str()) == Some("json") {
count += 1;
}
}
assert_eq!(count, 0);
}
/// The client must seed `last_activity_at` from `last_recording_at`
/// (the authoritative dictation-activity timestamp), not from
/// `updated_at` (oneliner regeneration time, which collapses to a
/// narrow window after a startup recovery scan and breaks ordering).
#[tokio::test]
async fn merge_prefers_last_recording_over_updated_at() {
let tmp = TempDir::new().unwrap();
let (store, _rx) = CaseStore::open(tmp.path().to_owned()).await.unwrap();
let id = Uuid::new_v4();
// Server: case created yesterday, last dictation 30 min ago,
// oneliner regenerated just now in a recovery scan.
let entry = OnelinerEntry {
case_id: id.to_string(),
oneliner: Some("x".into()),
created_at: "2026-04-17T09:00:00Z".into(),
last_recording_at: Some("2026-04-18T09:30:00Z".into()),
updated_at: Some("2026-04-18T10:00:00Z".into()),
};
store
.merge_server_snapshot(&server_snapshot(vec![entry]))
.await
.unwrap();
let list = store.list().await;
assert_eq!(list[0].last_activity_at, "2026-04-18T09:30:00Z");
}
/// Backward-compat: if the server response omits `last_recording_at`
/// (older server or future schema shift), the client falls back to
/// `updated_at` then `created_at`.
#[tokio::test]
async fn merge_falls_back_to_updated_at_when_last_recording_missing() {
let tmp = TempDir::new().unwrap();
let (store, _rx) = CaseStore::open(tmp.path().to_owned()).await.unwrap();
let id = Uuid::new_v4();
let entry = OnelinerEntry {
case_id: id.to_string(),
oneliner: Some("x".into()),
created_at: "2026-04-17T09:00:00Z".into(),
last_recording_at: None,
updated_at: Some("2026-04-18T10:00:00Z".into()),
};
store
.merge_server_snapshot(&server_snapshot(vec![entry]))
.await
.unwrap();
let list = store.list().await;
assert_eq!(list[0].last_activity_at, "2026-04-18T10:00:00Z");
}
#[tokio::test]
async fn sorted_newest_first() {
let tmp = TempDir::new().unwrap();
let (store, _rx) = CaseStore::open(tmp.path().to_owned()).await.unwrap();
let older = Uuid::new_v4();
let newer = Uuid::new_v4();
store.create_local(older, t("2026-04-18T08:00:00Z")).await.unwrap();
store.create_local(newer, t("2026-04-18T10:00:00Z")).await.unwrap();
let list = store.list().await;
assert_eq!(list[0].case_id, newer);
assert_eq!(list[1].case_id, older);
}
}