Files
doctate/client-desktop/src/app.rs
T
Brummel c769abe3d5 feat: Add per-case dir oneliner locks
Introduces `OnelinerLocks` to serialize read-modify-write operations on
`oneliner.json`. This prevents race conditions between the manual
override API (`PUT /api/cases/{case_id}/oneliner`) and the transcription
worker's auto-regeneration process.

The manual override now acquires a lock specific to the `case_dir`
before writing. The transcription worker's `update_oneliner` function
also acquires the same lock around its post-LLM re-read-and-write
sequence. This ensures that a doctor's manual edit always takes
precedence, even if it arrives while an auto-regeneration is in
progress.

This change also includes:
- A new `routes::oneliner_override` module for the PUT handler.
- Validation for the oneliner text length and emptiness.
- Integration tests for the override functionality, covering happy path,
  error conditions, early latching, race conditions, and parallel PUTs.
2026-04-26 16:05:01 +02:00

717 lines
24 KiB
Rust
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
//! The egui application: config panel, record/stop controls, case list.
//!
//! State transitions are driven by:
//! - user clicks (new/continue/stop/open/save-config/dismiss-error)
//! - recorder events (ffmpeg finalizing → RecorderEvent::Finished)
//! - upload events (terminal failure → AppState::Error)
//!
//! Connectivity and upload-progress indicators live in the footer; they
//! never block the main state machine. See
//! `berpr-fe-den-plan-lively-journal.md` for the design rationale.
use std::path::PathBuf;
use std::sync::Arc;
use std::time::{Duration, Instant};
use doctate_client_core::{
CaseMarker, CaseStore, DEFAULT_MARKER_RETENTION, DEFAULT_ORPHAN_AUDIO_RETENTION, FooterStatus,
PendingUpload, ServerSync, SnapshotCache, SyncConfig, UploadEvent, WorkerSnapshot, WorkerState,
pick_footer_status, run_startup_cleanup, write_sidecar,
};
use doctate_common::join_url;
use doctate_common::oneliners::OnelinerState;
use doctate_common::timestamp::{extract_hhmm, now_rfc3339, recorded_at_to_filename_stem};
use eframe::egui;
use tokio::runtime::Runtime;
use tokio::sync::{mpsc, watch};
use tracing::{error, info, warn};
use uuid::Uuid;
use crate::config::Config;
use crate::recorder::{Recorder, RecorderEvent};
use crate::state::AppState;
pub struct DoctateApp {
runtime: Arc<Runtime>,
config: Option<Config>,
pending_dir: PathBuf,
state: AppState,
server_url_input: String,
api_key_input: String,
http_client: Arc<reqwest::Client>,
case_store: Arc<CaseStore>,
snapshot_rx: watch::Receiver<Arc<Vec<CaseMarker>>>,
snapshot_cache: Arc<SnapshotCache>,
recorder: Option<Recorder>,
recorder_events: Option<mpsc::UnboundedReceiver<RecorderEvent>>,
server_sync: Option<ServerSync>,
worker_rx: Option<watch::Receiver<WorkerSnapshot>>,
upload_events: Option<mpsc::UnboundedReceiver<UploadEvent>>,
/// Context for an in-flight ffmpeg flush: we left `AppState::Recording`
/// for `AppState::Idle` on Stop (so the UI stays responsive), but we
/// still need `case_id` + `recorded_at` when the recorder emits
/// `Finished`. `None` when no flush is in progress.
finalizing: Option<RecordingContext>,
}
/// Enough recording identity to finalize a pending flush or emit an
/// upload event once the recorder signals `Finished`.
#[derive(Debug, Clone)]
struct RecordingContext {
case_id: Uuid,
recorded_at: String,
}
enum UiAction {
SaveConfig,
StartNew,
Continue(Uuid),
OpenWeb(Uuid),
Stop,
DismissError,
}
impl DoctateApp {
pub fn new(
runtime: Arc<Runtime>,
pending_dir: PathBuf,
cases_dir: PathBuf,
snapshot_cache_path: PathBuf,
) -> Self {
let config = match crate::config::load() {
Ok(Some(c)) => Some(c),
Ok(None) => None,
Err(e) => {
warn!(error = %e, "config load failed");
None
}
};
let server_url_input = config
.as_ref()
.map(|c| c.server_url.clone())
.unwrap_or_default();
let api_key_input = config
.as_ref()
.map(|c| c.api_key.clone())
.unwrap_or_default();
let (store, snapshot_rx) = runtime
.block_on(async { CaseStore::open(cases_dir).await })
.expect("open case store");
let case_store = Arc::new(store);
let snapshot_cache = Arc::new(SnapshotCache::new(snapshot_cache_path));
let http_client = Arc::new(reqwest::Client::new());
// Startup cleanup — one-shot, before the worker spawns.
{
let store = case_store.clone();
let pending = pending_dir.clone();
runtime.block_on(async move {
run_startup_cleanup(
&store,
&pending,
DEFAULT_MARKER_RETENTION,
DEFAULT_ORPHAN_AUDIO_RETENTION,
)
.await;
});
}
let (state, server_sync, worker_rx, upload_events) = if let Some(cfg) = &config {
let (sync, events, rx) = spawn_worker(
&runtime,
cfg.clone(),
pending_dir.clone(),
http_client.clone(),
case_store.clone(),
snapshot_cache.clone(),
);
(AppState::Idle, Some(sync), Some(rx), Some(events))
} else {
(AppState::NotConfigured, None, None, None)
};
Self {
runtime,
config,
pending_dir,
state,
server_url_input,
api_key_input,
http_client,
case_store,
snapshot_rx,
snapshot_cache,
recorder: None,
recorder_events: None,
server_sync,
worker_rx,
upload_events,
finalizing: None,
}
}
fn on_save_config(&mut self) {
let cfg = Config {
server_url: self.server_url_input.trim().to_string(),
api_key: self.api_key_input.trim().to_string(),
oneliner_poll_interval_seconds: self
.config
.as_ref()
.and_then(|c| c.oneliner_poll_interval_seconds),
};
if let Err(e) = crate::config::save(&cfg) {
error!(error = %e, "config save failed");
self.state = AppState::Error {
message: format!("Speichern: {e}"),
};
return;
}
info!("config saved");
let store = self.case_store.clone();
let cache = self.snapshot_cache.clone();
self.runtime.block_on(async move {
if let Err(e) = store.clear_all().await {
warn!(error = %e, "case store clear failed on config change");
}
if let Err(e) = cache.clear().await {
warn!(error = %e, "snapshot cache clear failed on config change");
}
});
// Drop old worker explicitly so its abort fires before new spawn.
self.server_sync = None;
self.worker_rx = None;
self.upload_events = None;
let (sync, events, rx) = spawn_worker(
&self.runtime,
cfg.clone(),
self.pending_dir.clone(),
self.http_client.clone(),
self.case_store.clone(),
self.snapshot_cache.clone(),
);
self.server_sync = Some(sync);
self.worker_rx = Some(rx);
self.upload_events = Some(events);
self.config = Some(cfg);
self.state = AppState::Idle;
self.finalizing = None;
}
fn on_start_new(&mut self) {
self.start_recording_for(Uuid::new_v4(), /* is_continue */ false);
}
fn on_continue(&mut self, case_id: Uuid) {
self.start_recording_for(case_id, /* is_continue */ true);
}
fn start_recording_for(&mut self, case_id: Uuid, is_continue: bool) {
// Two-part guard: user-visible state AND no lingering recorder
// from a previous stop that is still flushing.
if !matches!(self.state, AppState::Idle) || self.recorder.is_some() {
return;
}
if let Err(e) = std::fs::create_dir_all(&self.pending_dir) {
self.state = AppState::Error {
message: format!("Pending-Verzeichnis: {e}"),
};
return;
}
let recorded_at = now_rfc3339();
let stem = format!("{}_{}", case_id, recorded_at_to_filename_stem(&recorded_at));
let output_path = self.pending_dir.join(format!("{stem}.m4a"));
// Persist the case marker before recording so the UI list
// reflects the new case immediately.
let store = self.case_store.clone();
let now = time::OffsetDateTime::now_utc();
self.runtime.block_on(async move {
let res = if is_continue {
store.mark_activity(case_id, now).await.map(|_| ())
} else {
store.create_local(case_id, now).await.map(|_| ())
};
if let Err(e) = res {
warn!(error = %e, "case store write failed");
}
});
let _guard = self.runtime.enter();
match Recorder::start(output_path) {
Ok((recorder, events)) => {
self.recorder = Some(recorder);
self.recorder_events = Some(events);
self.state = AppState::Recording {
case_id,
recorded_at,
started_at: Instant::now(),
};
}
Err(e) => {
self.state = AppState::Error {
message: format!("Aufnahme starten: {e}"),
};
}
}
}
fn on_stop(&mut self) {
let AppState::Recording {
case_id,
recorded_at,
..
} = &self.state
else {
return;
};
let case_id = *case_id;
let recorded_at = recorded_at.clone();
if let Some(rec) = self.recorder.take() {
rec.stop();
}
// UI is ready for the next recording immediately — the finalizing
// ffmpeg flush shows up as a footer spinner, not a blocking state.
self.finalizing = Some(RecordingContext {
case_id,
recorded_at,
});
self.state = AppState::Idle;
}
fn on_open_web(&self, case_id: Uuid) {
let Some(cfg) = &self.config else {
return;
};
let server_url = cfg.server_url.clone();
let api_key = cfg.api_key.clone();
let http = self.http_client.clone();
let return_to = format!("/web/cases/{case_id}");
let fallback_url = join_url(&server_url, &return_to);
// Fire-and-forget on the runtime. Blocking the egui UI thread on
// an HTTP round-trip would freeze the window for ~100500ms.
self.runtime.spawn(async move {
let url =
match crate::magic_link::build_magic_url(&http, &server_url, &api_key, &return_to)
.await
{
Ok(u) => u,
Err(e) => {
warn!(error = %e, "magic-link request failed; falling back to plain URL");
fallback_url
}
};
if let Err(e) = webbrowser::open(&url) {
warn!(url = %url, error = %e, "open browser failed");
}
});
}
fn drain_recorder_events(&mut self) {
let Some(rx) = self.recorder_events.as_mut() else {
return;
};
let mut events = Vec::new();
while let Ok(event) = rx.try_recv() {
events.push(event);
}
for event in events {
match event {
RecorderEvent::Started => {}
RecorderEvent::Finished { output } => {
self.handle_recording_finished(output);
}
RecorderEvent::Failed { error: err } => {
error!(error = %err, "recorder failed");
self.finalizing = None;
self.state = AppState::Error {
message: format!("Aufnahme: {err}"),
};
self.recorder_events = None;
self.recorder = None;
}
}
}
}
fn handle_recording_finished(&mut self, output: PathBuf) {
// Prefer the explicit finalizing context (Stop-flow). Fall back
// to AppState::Recording (recorder terminated before Stop — rare
// but possible).
let ctx = if let Some(ctx) = self.finalizing.take() {
ctx
} else if let AppState::Recording {
case_id,
recorded_at,
..
} = &self.state
{
let ctx = RecordingContext {
case_id: *case_id,
recorded_at: recorded_at.clone(),
};
self.state = AppState::Idle;
ctx
} else {
warn!(state = ?self.state, "recorder finished in unexpected state");
self.recorder_events = None;
self.recorder = None;
return;
};
let RecordingContext {
case_id,
recorded_at,
} = ctx;
let store = self.case_store.clone();
let now = time::OffsetDateTime::now_utc();
self.runtime.block_on(async move {
if let Err(e) = store.mark_activity(case_id, now).await {
warn!(error = %e, "mark_activity on recording finished failed");
}
});
let upload = PendingUpload {
case_id,
recorded_at,
file: output,
};
if let Err(e) = write_sidecar(&upload) {
self.state = AppState::Error {
message: format!("Sidecar schreiben: {e}"),
};
self.recorder_events = None;
self.recorder = None;
return;
}
// Kick the worker so it picks up the fresh file without waiting
// for the next poll-tick.
if let Some(sync) = &self.server_sync {
sync.notify();
} else {
warn!("no server_sync active — upload stays on disk until next launch");
}
self.recorder_events = None;
self.recorder = None;
}
fn drain_upload_events(&mut self) {
let Some(rx) = self.upload_events.as_mut() else {
return;
};
let mut events = Vec::new();
while let Ok(event) = rx.try_recv() {
events.push(event);
}
for event in events {
match event {
UploadEvent::Started(case_id) => {
info!(case_id = %case_id, "upload attempt started");
}
UploadEvent::Succeeded { case_id, status } => {
info!(case_id = %case_id, status = ?status, "upload succeeded");
}
UploadEvent::Failed {
case_id,
reason,
will_retry,
} => {
warn!(case_id = %case_id, reason = %reason, will_retry, "upload failed");
if !will_retry {
// Terminal = operator action needed. Block UI.
self.state = AppState::Error { message: reason };
}
}
}
}
}
fn render_controls(&mut self, ui: &mut egui::Ui) -> Option<UiAction> {
match &self.state {
AppState::NotConfigured => {
ui.label("Server URL:");
ui.add(egui::TextEdit::singleline(&mut self.server_url_input).desired_width(280.0));
ui.add_space(4.0);
ui.label("API Key:");
ui.add(
egui::TextEdit::singleline(&mut self.api_key_input)
.desired_width(280.0)
.password(true),
);
ui.add_space(8.0);
if ui.button("Speichern").clicked() {
return Some(UiAction::SaveConfig);
}
None
}
AppState::Idle => {
if ui.button("● Neu").clicked() {
return Some(UiAction::StartNew);
}
None
}
AppState::Recording { started_at, .. } => {
let elapsed = started_at.elapsed();
ui.label(format!(
"● REC {:02}:{:02}",
elapsed.as_secs() / 60,
elapsed.as_secs() % 60
));
ui.add_space(4.0);
if ui.button("■ Stop").clicked() {
return Some(UiAction::Stop);
}
None
}
AppState::Error { message } => {
ui.colored_label(egui::Color32::RED, message);
ui.add_space(8.0);
if ui.button("OK").clicked() {
return Some(UiAction::DismissError);
}
None
}
}
}
fn render_case_list(&self, ui: &mut egui::Ui) -> Option<UiAction> {
let snapshot = self.snapshot_rx.borrow().clone();
// Snapshot is already server-filtered — the server's `window_hours`
// per user (users.toml) is the single source of truth.
let visible: Vec<&CaseMarker> = snapshot.iter().collect();
if visible.is_empty() {
ui.add_space(8.0);
ui.colored_label(egui::Color32::GRAY, "Noch keine Fälle heute.");
return None;
}
// Pending-upload lookup for the badge. Cheap clone (Arc) from the
// worker's last published snapshot; empty set if no worker yet.
let pending_ids = self
.worker_rx
.as_ref()
.map(|rx| rx.borrow().pending_case_ids.clone())
.unwrap_or_default();
let mut action = None;
// Ready-for-new means the user-intent state is Idle AND no ffmpeg
// instance is still flushing. Upload-queue depth does NOT gate.
let idle = matches!(self.state, AppState::Idle) && self.recorder.is_none();
egui::ScrollArea::vertical().show(ui, |ui| {
for marker in visible {
ui.separator();
ui.add_space(4.0);
let time_str = extract_hhmm(&marker.last_activity_at);
let oneliner = match &marker.oneliner {
Some(
OnelinerState::Ready { text, .. } | OnelinerState::Manual { text, .. },
) => text.clone(),
Some(OnelinerState::Empty { .. }) => "∅".to_owned(),
Some(OnelinerState::Error { .. }) => "⚠".to_owned(),
None => "⏳".to_owned(),
};
let has_pending = pending_ids.contains(&marker.case_id);
let line = if has_pending {
format!("⬆ {time_str} · {oneliner}")
} else {
format!("{time_str} · {oneliner}")
};
if has_pending {
ui.colored_label(egui::Color32::from_rgb(200, 140, 0), line);
} else {
ui.label(line);
}
ui.horizontal(|ui| {
if ui
.add_enabled(idle, egui::Button::new("▶ Fortsetzen"))
.clicked()
{
action = Some(UiAction::Continue(marker.case_id));
}
if ui.button("🌐 Öffnen").clicked() {
action = Some(UiAction::OpenWeb(marker.case_id));
}
});
ui.add_space(4.0);
}
});
action
}
fn render_footer_status(&self, ui: &mut egui::Ui) {
let Some(worker_rx) = &self.worker_rx else {
ui.colored_label(egui::Color32::GRAY, "—");
return;
};
let snap = worker_rx.borrow().clone();
let last_success_age = snap.last_success.map(|i| i.elapsed());
let last_failure_age = snap.last_failure.map(|i| i.elapsed());
let threshold = self
.config
.as_ref()
.map(|c| std::cmp::max(3 * c.poll_interval(), Duration::from_secs(30)))
.unwrap_or(Duration::from_secs(30));
let status = pick_footer_status(
self.finalizing.is_some(),
snap.state,
snap.queue_len,
last_success_age,
last_failure_age,
threshold,
);
match status {
FooterStatus::Finalizing => {
ui.horizontal(|ui| {
ui.spinner();
ui.label("Speichere Aufnahme…");
});
}
FooterStatus::Uploading { count } => {
ui.horizontal(|ui| {
ui.spinner();
if count <= 1 {
ui.label("Uploading…");
} else {
ui.label(format!("Uploading ({count} offen)…"));
}
});
}
FooterStatus::Synchronizing => {
ui.horizontal(|ui| {
ui.spinner();
ui.label("Synchronisiere…");
});
}
FooterStatus::NeverConnected => {
ui.colored_label(egui::Color32::GRAY, "⚪ noch keine Verbindung");
}
FooterStatus::Offline => {
ui.colored_label(egui::Color32::DARK_RED, "🔴 Offline");
}
FooterStatus::Online => {
ui.colored_label(egui::Color32::DARK_GREEN, "🟢 Online");
}
}
}
fn handle_action(&mut self, action: UiAction) {
match action {
UiAction::SaveConfig => self.on_save_config(),
UiAction::StartNew => self.on_start_new(),
UiAction::Continue(id) => self.on_continue(id),
UiAction::OpenWeb(id) => self.on_open_web(id),
UiAction::Stop => self.on_stop(),
UiAction::DismissError => self.state = AppState::Idle,
}
}
fn worker_is_active(&self) -> bool {
let Some(rx) = &self.worker_rx else {
return false;
};
let snap = rx.borrow();
snap.state != WorkerState::Idle || snap.queue_len > 0
}
}
fn spawn_worker(
runtime: &Runtime,
config: Config,
pending_dir: PathBuf,
http: Arc<reqwest::Client>,
case_store: Arc<CaseStore>,
snapshot_cache: Arc<SnapshotCache>,
) -> (
ServerSync,
mpsc::UnboundedReceiver<UploadEvent>,
watch::Receiver<WorkerSnapshot>,
) {
let _guard = runtime.enter();
let poll_interval = config.poll_interval();
let sync_config = SyncConfig {
server_url: config.server_url,
api_key: config.api_key,
poll_interval,
};
let (sync, events_rx) =
ServerSync::spawn(http, sync_config, case_store, snapshot_cache, pending_dir);
let snapshot_rx = sync.snapshot();
(sync, events_rx, snapshot_rx)
}
impl eframe::App for DoctateApp {
fn update(&mut self, ctx: &egui::Context, _frame: &mut eframe::Frame) {
self.drain_recorder_events();
self.drain_upload_events();
let mut pending_action: Option<UiAction> = None;
egui::TopBottomPanel::top("header").show(ctx, |ui| {
ui.add_space(4.0);
ui.vertical_centered(|ui| {
ui.heading("doctate");
ui.label(self.state.label());
});
ui.add_space(4.0);
});
egui::TopBottomPanel::bottom("footer").show(ctx, |ui| {
ui.add_space(4.0);
ui.vertical_centered(|ui| {
self.render_footer_status(ui);
});
ui.add_space(4.0);
});
egui::CentralPanel::default().show(ctx, |ui| {
ui.vertical_centered(|ui| {
ui.add_space(8.0);
if let Some(a) = self.render_controls(ui) {
pending_action = Some(a);
}
ui.add_space(8.0);
});
if !matches!(self.state, AppState::NotConfigured)
&& let Some(a) = self.render_case_list(ui)
{
pending_action = Some(a);
}
});
if let Some(action) = pending_action {
self.handle_action(action);
}
// Repaint cadence: fast while visibly active (spinner, REC-timer),
// slow (1 Hz) in the quiet steady state.
let active_spinner = self.finalizing.is_some() || self.worker_is_active();
let interval = match self.state {
AppState::Recording { .. } => Duration::from_millis(250),
_ if active_spinner => Duration::from_millis(100),
_ => Duration::from_secs(1),
};
ctx.request_repaint_after(interval);
}
}