Persist calendar events/attendees with dedup; link meetings to events (T6.2/T6.3/T6.6)

This commit is contained in:
iamdoubz
2026-07-01 15:07:40 -05:00
parent cd170c2907
commit 826d87c7fa
+292 -11
View File
@@ -5,7 +5,8 @@
//! summary) are regenerable; retention never touches an in-progress meeting.
use crate::models::{
ActionItem, MeetingId, MeetingListItem, MeetingStatus, SpeakerInfo, TranscriptSegment,
ActionItem, CalendarEvent, ImportedEvent, MeetingId, MeetingListItem, MeetingStatus,
Participant, SpeakerInfo, TranscriptSegment,
};
use crate::paths;
use async_trait::async_trait;
@@ -69,6 +70,18 @@ pub struct Meeting {
pub notes_markdown: String,
/// `None` until `generate_summary` has run at least once (Phase 5, FR-LLM-4).
pub summary: Option<SummaryFile>,
/// The linked calendar event, if any (T6.3/T6.6, FR-CAL-2/4) — set either
/// at `start_recording` time or later via `attach_meeting_to_event`.
pub calendar_event_id: Option<String>,
}
/// A calendar event with its attendees (T6.3/T6.4, FR-CAL-1/3) — the
/// "detail" fetch behind the pre-meeting context panel and the
/// participant-aware speaker-naming dropdown (T6.5).
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct CalendarEventDetail {
pub event: CalendarEvent,
pub participants: Vec<Participant>,
}
/// On-disk shape of `summary.json` (`docs/03-data-model.md`). Drafted action
@@ -129,6 +142,39 @@ pub trait Store: Send + Sync {
id: &MeetingId,
items: &[ActionItem],
) -> Result<(), StoreError>;
/// Persist imported calendar events + their attendees (T6.1/T6.2,
/// FR-CAL-1). Re-importing the same event (matched by `source` +
/// `raw_uid`) updates it in place rather than duplicating it. Returns
/// the number of events imported/updated.
async fn import_calendar_events(&self, events: Vec<ImportedEvent>) -> Result<u32, StoreError>;
/// Browse imported calendar events (T6.3, FR-CAL-2), optionally bounded
/// by unix-epoch start time.
async fn list_calendar_events(
&self,
from: Option<i64>,
to: Option<i64>,
) -> Result<Vec<CalendarEvent>, StoreError>;
/// A single event with its attendees (T6.3/T6.4, FR-CAL-1/3) — backs the
/// pre-meeting context panel and the speaker-naming attendee dropdown.
async fn get_calendar_event(&self, id: &str) -> Result<CalendarEventDetail, StoreError>;
/// Link a meeting (current or historical) to a calendar event (T6.3/T6.6,
/// FR-CAL-2/4). Errs if either id doesn't exist.
async fn attach_meeting_to_event(
&self,
meeting_id: &MeetingId,
event_id: &str,
) -> Result<(), StoreError>;
/// Names a speaker AND links them to a known `Participant` (T6.5/T6.6,
/// FR-SPK-4): the display name comes from the participant record, and
/// the shared `participant_id` is what gives naming "continuity" across
/// meetings — the same person, once identified anywhere, shows up
/// pre-resolvable wherever they're an attendee again.
async fn map_speaker_to_participant(
&self,
meeting_id: &MeetingId,
label: &str,
participant_id: &str,
) -> Result<(), StoreError>;
/// Full-text search across transcripts + notes (Phase 8, FR-SEARCH-1).
async fn search(&self, query: &str) -> Result<Vec<MeetingListItem>, StoreError>;
/// Startup reconcile: meetings with audio but no finalized transcript (FR-REL-1).
@@ -164,25 +210,69 @@ impl SqliteStore {
}
impl SqliteStore {
/// `participant_id: None` leaves an existing link untouched (e.g. a
/// plain free-text `rename_speaker` after `map_speaker_to_participant`
/// must not silently unlink the participant) — `COALESCE(excluded.…, …)`
/// falls back to the row's current value when the new one isn't given.
async fn upsert_speaker(
&self,
meeting_id: &MeetingId,
label: &str,
display_name: Option<&str>,
participant_id: Option<&str>,
) -> Result<(), StoreError> {
sqlx::query(
"INSERT INTO speakers (id, meeting_id, label, display_name)
VALUES (?, ?, ?, ?)
ON CONFLICT(meeting_id, label) DO UPDATE SET display_name = excluded.display_name",
"INSERT INTO speakers (id, meeting_id, label, display_name, participant_id)
VALUES (?, ?, ?, ?, ?)
ON CONFLICT(meeting_id, label) DO UPDATE SET
display_name = excluded.display_name,
participant_id = COALESCE(excluded.participant_id, participant_id)",
)
.bind(uuid::Uuid::new_v4().to_string())
.bind(meeting_id)
.bind(label)
.bind(display_name)
.bind(participant_id)
.execute(&self.pool)
.await?;
Ok(())
}
/// Finds-or-creates a `participants` row. Not a plain `ON CONFLICT`
/// upsert: `UNIQUE(name, email)` treats two NULL emails as distinct under
/// standard SQL NULL semantics, so a conflict target wouldn't reliably
/// catch a repeat import of the same no-email attendee — this looks the
/// row up explicitly instead.
async fn upsert_participant(
&self,
name: &str,
email: Option<&str>,
) -> Result<String, StoreError> {
let existing: Option<String> = match email {
Some(email) => sqlx::query("SELECT id FROM participants WHERE name = ? AND email = ?")
.bind(name)
.bind(email)
.fetch_optional(&self.pool)
.await?
.map(|r| r.get("id")),
None => sqlx::query("SELECT id FROM participants WHERE name = ? AND email IS NULL")
.bind(name)
.fetch_optional(&self.pool)
.await?
.map(|r| r.get("id")),
};
if let Some(id) = existing {
return Ok(id);
}
let id = uuid::Uuid::new_v4().to_string();
sqlx::query("INSERT INTO participants (id, name, email) VALUES (?, ?, ?)")
.bind(&id)
.bind(name)
.bind(email)
.execute(&self.pool)
.await?;
Ok(id)
}
}
/// Follows a `label -> merged_into` chain to its canonical label. Capped at 8
@@ -234,6 +324,19 @@ fn row_to_list_item(row: &sqlx::sqlite::SqliteRow) -> MeetingListItem {
}
}
fn row_to_calendar_event(row: &sqlx::sqlite::SqliteRow) -> CalendarEvent {
CalendarEvent {
id: row.get("id"),
source: row.get("source"),
subject: row.get("subject"),
organizer: row.get("organizer"),
starts_at: row.get("starts_at"),
ends_at: row.get("ends_at"),
description: row.get("description"),
raw_uid: row.get("raw_uid"),
}
}
#[async_trait]
impl Store for SqliteStore {
async fn create_meeting(&self, m: NewMeeting) -> Result<MeetingId, StoreError> {
@@ -277,8 +380,13 @@ impl Store for SqliteStore {
.await?;
for speaker in &s.speakers {
self.upsert_speaker(id, &speaker.label, speaker.display_name.as_deref())
.await?;
self.upsert_speaker(
id,
&speaker.label,
speaker.display_name.as_deref(),
speaker.participant_id.as_deref(),
)
.await?;
}
let transcript = TranscriptFile {
@@ -326,7 +434,7 @@ impl Store for SqliteStore {
async fn get_meeting(&self, id: &MeetingId) -> Result<Meeting, StoreError> {
let row = sqlx::query(
"SELECT id, title, started_at, ended_at, duration_secs, status, recorded, language, backend_used, model_used
"SELECT id, title, started_at, ended_at, duration_secs, status, recorded, language, backend_used, model_used, calendar_event_id
FROM meetings WHERE id = ?",
)
.bind(id)
@@ -335,7 +443,7 @@ impl Store for SqliteStore {
.ok_or_else(|| StoreError::NotFound(id.clone()))?;
let speaker_rows = sqlx::query(
"SELECT label, display_name, merged_into FROM speakers WHERE meeting_id = ? ORDER BY label",
"SELECT label, display_name, participant_id, merged_into FROM speakers WHERE meeting_id = ? ORDER BY label",
)
.bind(id)
.fetch_all(&self.pool)
@@ -356,7 +464,7 @@ impl Store for SqliteStore {
.map(|r| SpeakerInfo {
label: r.get("label"),
display_name: r.get("display_name"),
participant_id: None,
participant_id: r.get("participant_id"),
})
.collect();
@@ -394,6 +502,7 @@ impl Store for SqliteStore {
speakers,
notes_markdown,
summary,
calendar_event_id: row.get("calendar_event_id"),
})
}
@@ -425,7 +534,116 @@ impl Store for SqliteStore {
label: &str,
name: &str,
) -> Result<(), StoreError> {
self.upsert_speaker(id, label, Some(name)).await
self.upsert_speaker(id, label, Some(name), None).await
}
async fn list_calendar_events(
&self,
from: Option<i64>,
to: Option<i64>,
) -> Result<Vec<CalendarEvent>, StoreError> {
let rows = match (from, to) {
(Some(f), Some(t)) => {
sqlx::query(
"SELECT * FROM calendar_events WHERE starts_at >= ? AND starts_at <= ? ORDER BY starts_at",
)
.bind(f)
.bind(t)
.fetch_all(&self.pool)
.await?
}
(Some(f), None) => {
sqlx::query("SELECT * FROM calendar_events WHERE starts_at >= ? ORDER BY starts_at")
.bind(f)
.fetch_all(&self.pool)
.await?
}
(None, Some(t)) => {
sqlx::query("SELECT * FROM calendar_events WHERE starts_at <= ? ORDER BY starts_at")
.bind(t)
.fetch_all(&self.pool)
.await?
}
(None, None) => {
sqlx::query("SELECT * FROM calendar_events ORDER BY starts_at")
.fetch_all(&self.pool)
.await?
}
};
Ok(rows.iter().map(row_to_calendar_event).collect())
}
async fn get_calendar_event(&self, id: &str) -> Result<CalendarEventDetail, StoreError> {
let row = sqlx::query("SELECT * FROM calendar_events WHERE id = ?")
.bind(id)
.fetch_optional(&self.pool)
.await?
.ok_or_else(|| StoreError::NotFound(id.to_string()))?;
let event = row_to_calendar_event(&row);
let participant_rows = sqlx::query(
"SELECT p.id, p.name, p.email, cep.role FROM calendar_event_participants cep
JOIN participants p ON p.id = cep.participant_id
WHERE cep.calendar_event_id = ? ORDER BY p.name",
)
.bind(id)
.fetch_all(&self.pool)
.await?;
let participants = participant_rows
.iter()
.map(|r| Participant {
id: r.get("id"),
name: r.get("name"),
email: r.get("email"),
role: r.get("role"),
})
.collect();
Ok(CalendarEventDetail {
event,
participants,
})
}
async fn attach_meeting_to_event(
&self,
meeting_id: &MeetingId,
event_id: &str,
) -> Result<(), StoreError> {
let exists: Option<String> =
sqlx::query_scalar("SELECT id FROM calendar_events WHERE id = ?")
.bind(event_id)
.fetch_optional(&self.pool)
.await?;
if exists.is_none() {
return Err(StoreError::NotFound(format!("calendar event {event_id}")));
}
let result =
sqlx::query("UPDATE meetings SET calendar_event_id = ?, updated_at = ? WHERE id = ?")
.bind(event_id)
.bind(now_unix())
.bind(meeting_id)
.execute(&self.pool)
.await?;
if result.rows_affected() == 0 {
return Err(StoreError::NotFound(meeting_id.clone()));
}
Ok(())
}
async fn map_speaker_to_participant(
&self,
meeting_id: &MeetingId,
label: &str,
participant_id: &str,
) -> Result<(), StoreError> {
let name: String = sqlx::query_scalar("SELECT name FROM participants WHERE id = ?")
.bind(participant_id)
.fetch_optional(&self.pool)
.await?
.ok_or_else(|| StoreError::NotFound(format!("participant {participant_id}")))?;
self.upsert_speaker(meeting_id, label, Some(&name), Some(participant_id))
.await
}
async fn merge_speakers(
@@ -507,8 +725,71 @@ impl Store for SqliteStore {
Ok(())
}
async fn import_calendar_events(&self, events: Vec<ImportedEvent>) -> Result<u32, StoreError> {
let mut imported = 0u32;
for ImportedEvent { event, attendees } in events {
sqlx::query(
"INSERT INTO calendar_events (id, source, subject, organizer, starts_at, ends_at, description, raw_uid)
VALUES (?, ?, ?, ?, ?, ?, ?, ?)
ON CONFLICT(source, raw_uid) WHERE raw_uid IS NOT NULL DO UPDATE SET
subject = excluded.subject,
organizer = excluded.organizer,
starts_at = excluded.starts_at,
ends_at = excluded.ends_at,
description = excluded.description",
)
.bind(&event.id)
.bind(&event.source)
.bind(&event.subject)
.bind(&event.organizer)
.bind(event.starts_at)
.bind(event.ends_at)
.bind(&event.description)
.bind(&event.raw_uid)
.execute(&self.pool)
.await?;
// A re-import may have updated an *existing* row rather than
// inserting `event.id` — resolve the row that's actually there
// before linking attendees to it.
let row_id: String = match &event.raw_uid {
Some(raw_uid) => {
sqlx::query("SELECT id FROM calendar_events WHERE source = ? AND raw_uid = ?")
.bind(&event.source)
.bind(raw_uid)
.fetch_one(&self.pool)
.await?
.get("id")
}
None => event.id.clone(),
};
for attendee in &attendees {
let participant_id = self
.upsert_participant(&attendee.name, attendee.email.as_deref())
.await?;
sqlx::query(
"INSERT INTO calendar_event_participants (calendar_event_id, participant_id, role)
VALUES (?, ?, ?)
ON CONFLICT(calendar_event_id, participant_id) DO UPDATE SET role = excluded.role",
)
.bind(&row_id)
.bind(&participant_id)
.bind(&attendee.role)
.execute(&self.pool)
.await?;
}
imported += 1;
}
Ok(imported)
}
async fn search(&self, _query: &str) -> Result<Vec<MeetingListItem>, StoreError> {
todo!("Phase 8 — FTS5 search")
// Not wired to a command yet, but never panic on a trait method a
// future caller could reach (see commands::not_implemented).
Err(StoreError::Db(
"full-text search isn't built yet (Phase 8)".to_string(),
))
}
async fn recover_scan(&self) -> Result<Vec<MeetingId>, StoreError> {