feat: Add doctate-client-core crate

Introduces a new crate `doctate-client-core` to house shared client-side
logic. This includes:

- `case_store`: Manages local case marker files and merging with server
  snapshots.
- `snapshot_cache`: Placeholder for caching server responses.
- `oneliner_poller`: Placeholder for the background polling task.

The workspace configuration and `Cargo.lock` have been updated to
include the new crate.
This commit is contained in:
2026-04-18 10:03:40 +02:00
parent 94392e72d8
commit ac7157cdca
10 changed files with 647 additions and 17 deletions
Generated
+17
View File
@@ -1145,6 +1145,23 @@ dependencies = [
"libloading 0.8.9",
]
[[package]]
name = "doctate-client-core"
version = "0.1.0"
dependencies = [
"doctate-common",
"reqwest",
"serde",
"serde_json",
"tempfile",
"thiserror 1.0.69",
"time",
"tokio",
"tracing",
"uuid",
"wiremock",
]
[[package]]
name = "doctate-common"
version = "0.1.0"
+1 -1
View File
@@ -1,5 +1,5 @@
[workspace]
members = ["server", "doctate-common", "client-desktop"]
members = ["server", "doctate-common", "doctate-client-core", "client-desktop"]
resolver = "3"
[workspace.package]
+21
View File
@@ -0,0 +1,21 @@
[package]
name = "doctate-client-core"
version.workspace = true
edition.workspace = true
description = "Shared client logic for the doctate medical dictation system: case store, oneliner poller, snapshot cache"
[dependencies]
doctate-common = { path = "../doctate-common" }
tokio = { workspace = true }
reqwest = { workspace = true }
serde = { workspace = true }
serde_json = { workspace = true }
uuid = { workspace = true }
time = { workspace = true }
tracing = { workspace = true }
thiserror = "1"
[dev-dependencies]
tempfile = "3"
tokio = { workspace = true, features = ["test-util"] }
wiremock = "0.6"
+502
View File
@@ -0,0 +1,502 @@
//! 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.
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
.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(),
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);
}
#[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);
}
}
+20
View File
@@ -0,0 +1,20 @@
//! Shared client logic for the doctate dictation system.
//!
//! This crate is filesystem- and network-aware (unlike `doctate-common`,
//! which is DTO-only), but path-agnostic: callers pass concrete paths so
//! the same code works on Linux desktop, Windows desktop, Android, iOS.
//!
//! Modules:
//! - [`case_store`]: per-case marker files, merge-with-server-snapshot logic
//! - [`snapshot_cache`]: last successful `/api/oneliners` response, cached
//! on disk so the UI has data at cold start
//! - [`oneliner_poller`]: background task that polls the server and feeds
//! the case store
pub mod case_store;
pub mod oneliner_poller;
pub mod snapshot_cache;
pub use case_store::{CaseMarker, CaseStore, CaseStoreError};
pub use oneliner_poller::{OnelinerPoller, PollerConfig};
pub use snapshot_cache::{CachedSnapshot, SnapshotCache};
@@ -0,0 +1,26 @@
//! Placeholder — full implementation follows in next step.
pub struct PollerConfig {
pub server_url: String,
pub api_key: String,
pub poll_interval: std::time::Duration,
pub window_hours: u32,
}
pub struct OnelinerPoller {
_abort: tokio::task::AbortHandle,
}
impl OnelinerPoller {
/// To be implemented.
#[doc(hidden)]
pub fn from_handle(abort: tokio::task::AbortHandle) -> Self {
Self { _abort: abort }
}
}
impl Drop for OnelinerPoller {
fn drop(&mut self) {
self._abort.abort();
}
}
+24
View File
@@ -0,0 +1,24 @@
//! Placeholder — full implementation follows in next step.
use std::path::PathBuf;
use doctate_common::oneliners::OnelinersResponse;
use serde::{Deserialize, Serialize};
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct CachedSnapshot {
pub etag: String,
pub fetched_at: String,
pub response: OnelinersResponse,
}
pub struct SnapshotCache {
#[allow(dead_code)]
path: PathBuf,
}
impl SnapshotCache {
pub fn new(path: PathBuf) -> Self {
Self { path }
}
}
+2
View File
@@ -1,5 +1,6 @@
pub mod ack;
pub mod constants;
pub mod oneliners;
pub mod timestamp;
pub use ack::{AckResponse, AckStatus};
@@ -7,4 +8,5 @@ pub use constants::{
API_KEY_HEADER, CONTENT_TYPE_AUDIO_MP4, FIELD_AUDIO, FIELD_CASE_ID, FIELD_RECORDED_AT,
UPLOAD_PATH,
};
pub use oneliners::{OnelinerEntry, OnelinersResponse, ONELINERS_PATH};
pub use timestamp::{filename_stem_to_recorded_at, now_rfc3339, recorded_at_to_filename_stem};
+32
View File
@@ -0,0 +1,32 @@
//! Shared DTOs for the `/api/oneliners` endpoint.
//!
//! Server serializes [`OnelinersResponse`]; clients deserialize the same
//! type. Keeping this in a client-agnostic crate avoids drift between
//! server output and client parsing.
use serde::{Deserialize, Serialize};
/// HTTP path of the oneliners endpoint.
pub const ONELINERS_PATH: &str = "/api/oneliners";
/// Full response body of `GET /api/oneliners`.
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct OnelinersResponse {
/// RFC3339 UTC timestamp at which the server built this response.
pub as_of: String,
/// Effective window (after server-side clamp).
pub window_hours: u32,
/// Oneliners of cases whose earliest recording sits inside the window,
/// sorted newest-first by `created_at`.
pub oneliners: Vec<OnelinerEntry>,
}
/// One case summary. `oneliner` is `None` while the case is still
/// transcribing / generating; `updated_at` is `None` for the same reason.
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct OnelinerEntry {
pub case_id: String,
pub oneliner: Option<String>,
pub created_at: String,
pub updated_at: Option<String>,
}
+2 -16
View File
@@ -18,7 +18,8 @@ use axum::extract::{Query, State};
use axum::http::{header, HeaderMap, StatusCode};
use axum::response::{IntoResponse, Response};
use axum::Json;
use serde::{Deserialize, Serialize};
use doctate_common::oneliners::{OnelinerEntry, OnelinersResponse};
use serde::Deserialize;
use time::format_description::well_known::Rfc3339;
use time::{Duration, OffsetDateTime};
use tracing::warn;
@@ -35,21 +36,6 @@ pub struct OnelinersQuery {
hours: Option<u32>,
}
#[derive(Serialize)]
struct OnelinersResponse {
as_of: String,
window_hours: u32,
oneliners: Vec<OnelinerEntry>,
}
#[derive(Serialize)]
struct OnelinerEntry {
case_id: String,
oneliner: Option<String>,
created_at: String,
updated_at: Option<String>,
}
pub async fn handle_oneliners(
user: AuthenticatedUser,
State(watermark): State<OnelinerWatermark>,