diff --git a/src/openhuman/composio/ops.rs b/src/openhuman/composio/ops.rs index 4b8c0b833..a82b688ae 100644 --- a/src/openhuman/composio/ops.rs +++ b/src/openhuman/composio/ops.rs @@ -292,6 +292,15 @@ pub async fn composio_get_user_profile( .await .map_err(|e| format!("[composio] get_user_profile({toolkit}) failed: {e}"))?; + // Side-effect: persist profile fields into the local user_profile + // facet table so any RPC call also refreshes the local store. + let facets = super::providers::profile::persist_provider_profile(&profile); + tracing::debug!( + toolkit = %toolkit, + facets_written = facets, + "[composio] profile facets persisted from get_user_profile" + ); + Ok(RpcOutcome::new( profile, vec![format!( diff --git a/src/openhuman/composio/providers/gmail.rs b/src/openhuman/composio/providers/gmail/mod.rs similarity index 75% rename from src/openhuman/composio/providers/gmail.rs rename to src/openhuman/composio/providers/gmail/mod.rs index a3ca661d8..f46f9d634 100644 --- a/src/openhuman/composio/providers/gmail.rs +++ b/src/openhuman/composio/providers/gmail/mod.rs @@ -17,6 +17,10 @@ //! number of `execute_tool` calls per calendar day, preventing runaway //! API usage during large initial backfills. +mod sync; +#[cfg(test)] +mod tests; + use async_trait::async_trait; use serde_json::{json, Value}; @@ -151,7 +155,7 @@ impl ComposioProvider for GmailProvider { } async fn sync(&self, ctx: &ProviderContext, reason: SyncReason) -> Result { - let started_at_ms = now_ms(); + let started_at_ms = sync::now_ms(); let connection_id = ctx .connection_id .clone() @@ -181,7 +185,7 @@ impl ComposioProvider for GmailProvider { reason: reason.as_str().to_string(), items_ingested: 0, started_at_ms, - finished_at_ms: now_ms(), + finished_at_ms: sync::now_ms(), summary: "gmail sync skipped: daily budget exhausted".to_string(), details: json!({ "budget_exhausted": true }), }); @@ -212,7 +216,7 @@ impl ComposioProvider for GmailProvider { // returns newer mail. let mut query = "in:inbox -in:spam -in:trash".to_string(); if let Some(ref cursor) = state.cursor { - if let Some(date_filter) = cursor_to_gmail_after_filter(cursor) { + if let Some(date_filter) = sync::cursor_to_gmail_after_filter(cursor) { query.push_str(&format!(" after:{date_filter}")); tracing::debug!( page = page_num, @@ -252,7 +256,7 @@ impl ComposioProvider for GmailProvider { )); } - let messages = extract_messages(&resp.data); + let messages = sync::extract_messages(&resp.data); total_fetched += messages.len(); if messages.is_empty() { @@ -336,7 +340,7 @@ impl ComposioProvider for GmailProvider { } // Check for next page token. - page_token = extract_page_token(&resp.data); + page_token = sync::extract_page_token(&resp.data); if page_token.is_none() { tracing::debug!(page = page_num, "[composio:gmail] no next page token, done"); break; @@ -349,7 +353,7 @@ impl ComposioProvider for GmailProvider { } state.save(&memory).await?; - let finished_at_ms = now_ms(); + let finished_at_ms = sync::now_ms(); let summary = format!( "gmail sync ({reason}): fetched {total_fetched}, persisted {total_persisted} new, \ budget remaining {remaining}", @@ -408,144 +412,3 @@ impl ComposioProvider for GmailProvider { Ok(()) } } - -// ── helpers ──────────────────────────────────────────────────────── - -/// Walk the Composio response envelope and pull out message objects. -fn extract_messages(data: &Value) -> Vec { - let candidates = [ - data.pointer("/data/messages"), - data.pointer("/messages"), - data.pointer("/data/data/messages"), - data.pointer("/data/items"), - data.pointer("/items"), - ]; - for cand in candidates.into_iter().flatten() { - if let Some(arr) = cand.as_array() { - return arr.clone(); - } - } - Vec::new() -} - -/// Try to extract a pagination token from the API response. -fn extract_page_token(data: &Value) -> Option { - let candidates = [ - data.pointer("/data/nextPageToken"), - data.pointer("/nextPageToken"), - data.pointer("/data/data/nextPageToken"), - ]; - for cand in candidates.into_iter().flatten() { - if let Some(s) = cand.as_str() { - let trimmed = s.trim(); - if !trimmed.is_empty() { - return Some(trimmed.to_string()); - } - } - } - None -} - -/// Convert a cursor value (epoch millis or date string) into a Gmail -/// `after:YYYY/MM/DD` filter component. Returns `None` if the cursor -/// cannot be parsed. -fn cursor_to_gmail_after_filter(cursor: &str) -> Option { - // Try parsing as epoch millis first (Gmail's internalDate). - if let Ok(millis) = cursor.parse::() { - let secs = millis / 1000; - if let Some(dt) = chrono::DateTime::from_timestamp(secs, 0) { - return Some(dt.format("%Y/%m/%d").to_string()); - } - } - // Try parsing as an ISO date/datetime. - if let Ok(dt) = chrono::NaiveDate::parse_from_str(cursor, "%Y-%m-%d") { - return Some(dt.format("%Y/%m/%d").to_string()); - } - if let Ok(dt) = chrono::DateTime::parse_from_rfc3339(cursor) { - return Some(dt.format("%Y/%m/%d").to_string()); - } - None -} - -fn now_ms() -> u64 { - use std::time::{SystemTime, UNIX_EPOCH}; - SystemTime::now() - .duration_since(UNIX_EPOCH) - .map(|d| d.as_millis() as u64) - .unwrap_or(0) -} - -#[cfg(test)] -mod tests { - use super::*; - - #[test] - fn extract_messages_finds_data_messages() { - let v = json!({ - "data": { "messages": [{"id": "m1"}, {"id": "m2"}] }, - "successful": true, - }); - assert_eq!(extract_messages(&v).len(), 2); - } - - #[test] - fn extract_messages_finds_top_level_messages() { - let v = json!({ "messages": [{"id": "m1"}] }); - assert_eq!(extract_messages(&v).len(), 1); - } - - #[test] - fn extract_messages_returns_empty_when_missing() { - let v = json!({ "data": { "other": [] } }); - assert_eq!(extract_messages(&v).len(), 0); - } - - #[test] - fn extract_page_token_finds_nested() { - let v = json!({ "data": { "nextPageToken": "tok123" } }); - assert_eq!(extract_page_token(&v), Some("tok123".to_string())); - } - - #[test] - fn extract_page_token_none_when_missing() { - let v = json!({ "data": {} }); - assert_eq!(extract_page_token(&v), None); - } - - #[test] - fn cursor_to_filter_from_epoch_millis() { - // 2026-04-01 00:00:00 UTC in millis - let millis = "1774915200000"; - let filter = cursor_to_gmail_after_filter(millis); - assert!(filter.is_some()); - // Should produce a YYYY/MM/DD date. - let f = filter.unwrap(); - assert!(f.contains('/'), "Expected date with slashes, got {f}"); - } - - #[test] - fn cursor_to_filter_from_iso_date() { - assert_eq!( - cursor_to_gmail_after_filter("2026-03-15"), - Some("2026/03/15".to_string()) - ); - } - - #[test] - fn cursor_to_filter_from_rfc3339() { - let f = cursor_to_gmail_after_filter("2026-03-15T12:00:00Z"); - assert_eq!(f, Some("2026/03/15".to_string())); - } - - #[test] - fn cursor_to_filter_returns_none_for_garbage() { - assert_eq!(cursor_to_gmail_after_filter("not-a-date"), None); - } - - #[test] - fn provider_metadata_is_stable() { - let p = GmailProvider::new(); - assert_eq!(p.toolkit_slug(), "gmail"); - assert_eq!(p.sync_interval_secs(), Some(15 * 60)); - } -} diff --git a/src/openhuman/composio/providers/gmail/sync.rs b/src/openhuman/composio/providers/gmail/sync.rs new file mode 100644 index 000000000..0fb74140f --- /dev/null +++ b/src/openhuman/composio/providers/gmail/sync.rs @@ -0,0 +1,69 @@ +//! Gmail sync helpers — message extraction, pagination, cursor +//! conversion, and time utilities. + +use serde_json::Value; + +/// Walk the Composio response envelope and pull out message objects. +pub(crate) fn extract_messages(data: &Value) -> Vec { + let candidates = [ + data.pointer("/data/messages"), + data.pointer("/messages"), + data.pointer("/data/data/messages"), + data.pointer("/data/items"), + data.pointer("/items"), + ]; + for cand in candidates.into_iter().flatten() { + if let Some(arr) = cand.as_array() { + return arr.clone(); + } + } + Vec::new() +} + +/// Try to extract a pagination token from the API response. +pub(crate) fn extract_page_token(data: &Value) -> Option { + let candidates = [ + data.pointer("/data/nextPageToken"), + data.pointer("/nextPageToken"), + data.pointer("/data/data/nextPageToken"), + ]; + for cand in candidates.into_iter().flatten() { + if let Some(s) = cand.as_str() { + let trimmed = s.trim(); + if !trimmed.is_empty() { + return Some(trimmed.to_string()); + } + } + } + None +} + +/// Convert a cursor value (epoch millis or date string) into a Gmail +/// `after:YYYY/MM/DD` filter component. Returns `None` if the cursor +/// cannot be parsed. +pub(crate) fn cursor_to_gmail_after_filter(cursor: &str) -> Option { + let cursor = cursor.trim(); + // Try parsing as epoch millis first (Gmail's internalDate). + if let Ok(millis) = cursor.parse::() { + let secs = millis / 1000; + if let Some(dt) = chrono::DateTime::from_timestamp(secs, 0) { + return Some(dt.format("%Y/%m/%d").to_string()); + } + } + // Try parsing as an ISO date/datetime. + if let Ok(dt) = chrono::NaiveDate::parse_from_str(cursor, "%Y-%m-%d") { + return Some(dt.format("%Y/%m/%d").to_string()); + } + if let Ok(dt) = chrono::DateTime::parse_from_rfc3339(cursor) { + return Some(dt.format("%Y/%m/%d").to_string()); + } + None +} + +pub(crate) fn now_ms() -> u64 { + use std::time::{SystemTime, UNIX_EPOCH}; + SystemTime::now() + .duration_since(UNIX_EPOCH) + .map(|d| d.as_millis() as u64) + .unwrap_or(0) +} diff --git a/src/openhuman/composio/providers/gmail/tests.rs b/src/openhuman/composio/providers/gmail/tests.rs new file mode 100644 index 000000000..d551f3915 --- /dev/null +++ b/src/openhuman/composio/providers/gmail/tests.rs @@ -0,0 +1,75 @@ +//! Unit tests for the Gmail provider. + +use super::sync::{cursor_to_gmail_after_filter, extract_messages, extract_page_token}; +use super::GmailProvider; +use crate::openhuman::composio::providers::ComposioProvider; +use serde_json::json; + +#[test] +fn extract_messages_finds_data_messages() { + let v = json!({ + "data": { "messages": [{"id": "m1"}, {"id": "m2"}] }, + "successful": true, + }); + assert_eq!(extract_messages(&v).len(), 2); +} + +#[test] +fn extract_messages_finds_top_level_messages() { + let v = json!({ "messages": [{"id": "m1"}] }); + assert_eq!(extract_messages(&v).len(), 1); +} + +#[test] +fn extract_messages_returns_empty_when_missing() { + let v = json!({ "data": { "other": [] } }); + assert_eq!(extract_messages(&v).len(), 0); +} + +#[test] +fn extract_page_token_finds_nested() { + let v = json!({ "data": { "nextPageToken": "tok123" } }); + assert_eq!(extract_page_token(&v), Some("tok123".to_string())); +} + +#[test] +fn extract_page_token_none_when_missing() { + let v = json!({ "data": {} }); + assert_eq!(extract_page_token(&v), None); +} + +#[test] +fn cursor_to_filter_from_epoch_millis() { + // 1774915200000 ms = 2026-03-31 UTC + let millis = "1774915200000"; + assert_eq!( + cursor_to_gmail_after_filter(millis), + Some("2026/03/31".to_string()) + ); +} + +#[test] +fn cursor_to_filter_from_iso_date() { + assert_eq!( + cursor_to_gmail_after_filter("2026-03-15"), + Some("2026/03/15".to_string()) + ); +} + +#[test] +fn cursor_to_filter_from_rfc3339() { + let f = cursor_to_gmail_after_filter("2026-03-15T12:00:00Z"); + assert_eq!(f, Some("2026/03/15".to_string())); +} + +#[test] +fn cursor_to_filter_returns_none_for_garbage() { + assert_eq!(cursor_to_gmail_after_filter("not-a-date"), None); +} + +#[test] +fn provider_metadata_is_stable() { + let p = GmailProvider::new(); + assert_eq!(p.toolkit_slug(), "gmail"); + assert_eq!(p.sync_interval_secs(), Some(15 * 60)); +} diff --git a/src/openhuman/composio/providers/mod.rs b/src/openhuman/composio/providers/mod.rs index 17b27374a..cbd889db1 100644 --- a/src/openhuman/composio/providers/mod.rs +++ b/src/openhuman/composio/providers/mod.rs @@ -42,6 +42,7 @@ use super::client::{build_composio_client, ComposioClient}; pub mod gmail; pub mod notion; +pub mod profile; pub mod registry; pub mod sync_state; @@ -262,6 +263,17 @@ pub trait ComposioProvider: Send + Sync { email_domain = ?email_domain, "[composio:provider] user profile fetched" ); + + // Persist profile fields into the local user_profile + // facet table so display_name / email / avatar are + // available to the agent context and UI without a + // round-trip to the upstream provider. + let facets = profile::persist_provider_profile(&profile); + tracing::debug!( + toolkit = %toolkit, + facets_written = facets, + "[composio:provider] profile facets persisted" + ); } Err(e) => { tracing::warn!( diff --git a/src/openhuman/composio/providers/notion.rs b/src/openhuman/composio/providers/notion/mod.rs similarity index 71% rename from src/openhuman/composio/providers/notion.rs rename to src/openhuman/composio/providers/notion/mod.rs index 4d0fdd852..da1b122c1 100644 --- a/src/openhuman/composio/providers/notion.rs +++ b/src/openhuman/composio/providers/notion/mod.rs @@ -15,6 +15,10 @@ //! page are older than the cursor. //! 7. Advance the cursor and save state. +mod sync; +#[cfg(test)] +mod tests; + use async_trait::async_trait; use serde_json::{json, Value}; @@ -23,8 +27,8 @@ use super::{ pick_str, ComposioProvider, ProviderContext, ProviderUserProfile, SyncOutcome, SyncReason, }; -const ACTION_GET_ABOUT_ME: &str = "NOTION_GET_ABOUT_ME"; -const ACTION_FETCH_DATA: &str = "NOTION_FETCH_DATA"; +pub(crate) const ACTION_GET_ABOUT_ME: &str = "NOTION_GET_ABOUT_ME"; +pub(crate) const ACTION_FETCH_DATA: &str = "NOTION_FETCH_DATA"; /// Page size per API call. const PAGE_SIZE: u32 = 25; @@ -133,7 +137,7 @@ impl ComposioProvider for NotionProvider { } async fn sync(&self, ctx: &ProviderContext, reason: SyncReason) -> Result { - let started_at_ms = now_ms(); + let started_at_ms = sync::now_ms(); let connection_id = ctx .connection_id .clone() @@ -163,7 +167,7 @@ impl ComposioProvider for NotionProvider { reason: reason.as_str().to_string(), items_ingested: 0, started_at_ms, - finished_at_ms: now_ms(), + finished_at_ms: sync::now_ms(), summary: "notion sync skipped: daily budget exhausted".to_string(), details: json!({ "budget_exhausted": true }), }); @@ -219,7 +223,7 @@ impl ComposioProvider for NotionProvider { )); } - let results = extract_results(&resp.data); + let results = sync::extract_results(&resp.data); total_fetched += results.len(); if results.is_empty() { @@ -273,8 +277,8 @@ impl ComposioProvider for NotionProvider { } // Build a title from the page's properties. - let title_text = - extract_page_title(page).unwrap_or_else(|| format!("Notion page {page_id}")); + let title_text = sync::extract_page_title(page) + .unwrap_or_else(|| format!("Notion page {page_id}")); let doc_id = format!("composio-notion-page-{page_id}"); let title = format!("Notion: {title_text}"); @@ -312,7 +316,7 @@ impl ComposioProvider for NotionProvider { } // Check for next page cursor from Notion API. - notion_cursor = extract_notion_cursor(&resp.data); + notion_cursor = sync::extract_notion_cursor(&resp.data); if notion_cursor.is_none() { tracing::debug!(page = page_num, "[composio:notion] no next cursor, done"); break; @@ -325,7 +329,7 @@ impl ComposioProvider for NotionProvider { } state.save(&memory).await?; - let finished_at_ms = now_ms(); + let finished_at_ms = sync::now_ms(); let summary = format!( "notion sync ({reason}): fetched {total_fetched}, persisted {total_persisted} new, \ budget remaining {remaining}", @@ -379,148 +383,3 @@ impl ComposioProvider for NotionProvider { Ok(()) } } - -// ── helpers ──────────────────────────────────────────────────────── - -/// Walk the Composio response envelope for Notion page results. -fn extract_results(data: &Value) -> Vec { - let candidates = [ - data.pointer("/data/results"), - data.pointer("/results"), - data.pointer("/data/data/results"), - data.pointer("/data/items"), - data.pointer("/items"), - ]; - for cand in candidates.into_iter().flatten() { - if let Some(arr) = cand.as_array() { - return arr.clone(); - } - } - Vec::new() -} - -/// Extract the Notion pagination cursor (for `start_cursor` on the -/// next request). -fn extract_notion_cursor(data: &Value) -> Option { - let candidates = [ - data.pointer("/data/next_cursor"), - data.pointer("/next_cursor"), - data.pointer("/data/data/next_cursor"), - ]; - for cand in candidates.into_iter().flatten() { - if let Some(s) = cand.as_str() { - let trimmed = s.trim(); - if !trimmed.is_empty() { - return Some(trimmed.to_string()); - } - } - } - None -} - -/// Try to extract a human-readable title from a Notion page object. -/// -/// Notion pages store the title in `properties.title` or -/// `properties.Name.title[0].plain_text`. We try several shapes. -fn extract_page_title(page: &Value) -> Option { - // Try the common `properties.title.title[0].plain_text` shape. - let props = page - .get("properties") - .or_else(|| page.get("data")?.get("properties")); - if let Some(props) = props { - // Walk all properties looking for a "title" type field. - if let Some(obj) = props.as_object() { - for (_key, val) in obj { - if val.get("type").and_then(Value::as_str) == Some("title") { - if let Some(arr) = val.get("title").and_then(Value::as_array) { - let text: String = arr - .iter() - .filter_map(|t| t.get("plain_text").and_then(Value::as_str)) - .collect::>() - .join(""); - if !text.is_empty() { - return Some(text); - } - } - } - } - } - } - - // Fallback: top-level "title" field (some Composio shapes). - pick_str(page, &["title", "data.title", "name", "data.name"]) -} - -fn now_ms() -> u64 { - use std::time::{SystemTime, UNIX_EPOCH}; - SystemTime::now() - .duration_since(UNIX_EPOCH) - .map(|d| d.as_millis() as u64) - .unwrap_or(0) -} - -#[cfg(test)] -mod tests { - use super::*; - - #[test] - fn extract_results_walks_common_shapes() { - let v1 = json!({ "data": { "results": [{"id": "p1"}] } }); - let v2 = json!({ "results": [{"id": "p2"}, {"id": "p3"}] }); - let v3 = json!({ "data": {} }); - assert_eq!(extract_results(&v1).len(), 1); - assert_eq!(extract_results(&v2).len(), 2); - assert_eq!(extract_results(&v3).len(), 0); - } - - #[test] - fn extract_notion_cursor_finds_nested() { - let v = json!({ "data": { "next_cursor": "abc123" } }); - assert_eq!(extract_notion_cursor(&v), Some("abc123".to_string())); - } - - #[test] - fn extract_notion_cursor_none_when_missing() { - let v = json!({ "data": { "has_more": false } }); - assert_eq!(extract_notion_cursor(&v), None); - } - - #[test] - fn extract_page_title_from_properties() { - let page = json!({ - "id": "page-1", - "properties": { - "Name": { - "type": "title", - "title": [ - { "plain_text": "My " }, - { "plain_text": "Page Title" } - ] - } - } - }); - assert_eq!(extract_page_title(&page), Some("My Page Title".to_string())); - } - - #[test] - fn extract_page_title_fallback_to_top_level() { - let page = json!({ "title": "Fallback Title" }); - assert_eq!( - extract_page_title(&page), - Some("Fallback Title".to_string()) - ); - } - - #[test] - fn extract_page_title_returns_none_when_missing() { - let page = json!({ "id": "p1" }); - assert_eq!(extract_page_title(&page), None); - } - - #[test] - fn provider_metadata_is_stable() { - let p = NotionProvider::new(); - assert_eq!(p.toolkit_slug(), "notion"); - assert_eq!(p.sync_interval_secs(), Some(30 * 60)); - } -} diff --git a/src/openhuman/composio/providers/notion/sync.rs b/src/openhuman/composio/providers/notion/sync.rs new file mode 100644 index 000000000..1571b13c8 --- /dev/null +++ b/src/openhuman/composio/providers/notion/sync.rs @@ -0,0 +1,83 @@ +//! Notion sync helpers — result extraction, pagination cursor, +//! page title extraction, and time utilities. + +use serde_json::Value; + +use super::pick_str; + +/// Walk the Composio response envelope for Notion page results. +pub(crate) fn extract_results(data: &Value) -> Vec { + let candidates = [ + data.pointer("/data/results"), + data.pointer("/results"), + data.pointer("/data/data/results"), + data.pointer("/data/items"), + data.pointer("/items"), + ]; + for cand in candidates.into_iter().flatten() { + if let Some(arr) = cand.as_array() { + return arr.clone(); + } + } + Vec::new() +} + +/// Extract the Notion pagination cursor (for `start_cursor` on the +/// next request). +pub(crate) fn extract_notion_cursor(data: &Value) -> Option { + let candidates = [ + data.pointer("/data/next_cursor"), + data.pointer("/next_cursor"), + data.pointer("/data/data/next_cursor"), + ]; + for cand in candidates.into_iter().flatten() { + if let Some(s) = cand.as_str() { + let trimmed = s.trim(); + if !trimmed.is_empty() { + return Some(trimmed.to_string()); + } + } + } + None +} + +/// Try to extract a human-readable title from a Notion page object. +/// +/// Notion pages store the title in `properties.title` or +/// `properties.Name.title[0].plain_text`. We try several shapes. +pub(crate) fn extract_page_title(page: &Value) -> Option { + // Try the common `properties.title.title[0].plain_text` shape. + let props = page + .get("properties") + .or_else(|| page.get("data")?.get("properties")); + if let Some(props) = props { + // Walk all properties looking for a "title" type field. + if let Some(obj) = props.as_object() { + for (_key, val) in obj { + if val.get("type").and_then(Value::as_str) == Some("title") { + if let Some(arr) = val.get("title").and_then(Value::as_array) { + let text: String = arr + .iter() + .filter_map(|t| t.get("plain_text").and_then(Value::as_str)) + .collect::>() + .join(""); + if !text.is_empty() { + return Some(text); + } + } + } + } + } + } + + // Fallback: top-level "title" field (some Composio shapes). + pick_str(page, &["title", "data.title", "name", "data.name"]) +} + +pub(crate) fn now_ms() -> u64 { + use std::time::{SystemTime, UNIX_EPOCH}; + SystemTime::now() + .duration_since(UNIX_EPOCH) + .map(|d| d.as_millis() as u64) + .unwrap_or(0) +} diff --git a/src/openhuman/composio/providers/notion/tests.rs b/src/openhuman/composio/providers/notion/tests.rs new file mode 100644 index 000000000..eb706c720 --- /dev/null +++ b/src/openhuman/composio/providers/notion/tests.rs @@ -0,0 +1,67 @@ +//! Unit tests for the Notion provider. + +use super::sync::{extract_notion_cursor, extract_page_title, extract_results}; +use super::NotionProvider; +use crate::openhuman::composio::providers::ComposioProvider; +use serde_json::json; + +#[test] +fn extract_results_walks_common_shapes() { + let v1 = json!({ "data": { "results": [{"id": "p1"}] } }); + let v2 = json!({ "results": [{"id": "p2"}, {"id": "p3"}] }); + let v3 = json!({ "data": {} }); + assert_eq!(extract_results(&v1).len(), 1); + assert_eq!(extract_results(&v2).len(), 2); + assert_eq!(extract_results(&v3).len(), 0); +} + +#[test] +fn extract_notion_cursor_finds_nested() { + let v = json!({ "data": { "next_cursor": "abc123" } }); + assert_eq!(extract_notion_cursor(&v), Some("abc123".to_string())); +} + +#[test] +fn extract_notion_cursor_none_when_missing() { + let v = json!({ "data": { "has_more": false } }); + assert_eq!(extract_notion_cursor(&v), None); +} + +#[test] +fn extract_page_title_from_properties() { + let page = json!({ + "id": "page-1", + "properties": { + "Name": { + "type": "title", + "title": [ + { "plain_text": "My " }, + { "plain_text": "Page Title" } + ] + } + } + }); + assert_eq!(extract_page_title(&page), Some("My Page Title".to_string())); +} + +#[test] +fn extract_page_title_fallback_to_top_level() { + let page = json!({ "title": "Fallback Title" }); + assert_eq!( + extract_page_title(&page), + Some("Fallback Title".to_string()) + ); +} + +#[test] +fn extract_page_title_returns_none_when_missing() { + let page = json!({ "id": "p1" }); + assert_eq!(extract_page_title(&page), None); +} + +#[test] +fn provider_metadata_is_stable() { + let p = NotionProvider::new(); + assert_eq!(p.toolkit_slug(), "notion"); + assert_eq!(p.sync_interval_secs(), Some(30 * 60)); +} diff --git a/src/openhuman/composio/providers/profile.rs b/src/openhuman/composio/providers/profile.rs new file mode 100644 index 000000000..ece0300aa --- /dev/null +++ b/src/openhuman/composio/providers/profile.rs @@ -0,0 +1,219 @@ +//! Profile persistence bridge — maps [`ProviderUserProfile`] fields +//! into the local `user_profile` facet table so provider-sourced +//! identity data (display name, email, username, avatar) accumulates +//! alongside conversation-extracted preferences. +//! +//! Each non-`None` field becomes a [`FacetType::Context`] facet keyed +//! as `composio:{toolkit}:{field}`. Confidence is set to 0.95 because +//! this data comes directly from the upstream provider API — it's +//! authoritative, not inferred from conversation. +//! +//! Callers are expected to invoke [`persist_provider_profile`] after +//! every successful `fetch_user_profile` call — from +//! `on_connection_created`, periodic syncs, and the +//! `composio_get_user_profile` RPC op. + +use super::ProviderUserProfile; +use crate::openhuman::memory::store::profile::{self, FacetType}; + +/// Confidence level assigned to provider-sourced profile data. +/// +/// This is higher than conversation-inferred facets (typically 0.5–0.7) +/// because the data comes directly from the upstream provider API. +const PROVIDER_CONFIDENCE: f64 = 0.95; + +/// Persist the non-`None` fields of a [`ProviderUserProfile`] into the +/// local `user_profile` facet table. +/// +/// Returns the number of facets written (0–4). Silently returns 0 if +/// the global memory client is not yet initialised (user not signed in, +/// startup race, etc.) — callers should treat that as non-fatal. +pub fn persist_provider_profile(profile: &ProviderUserProfile) -> usize { + let Some(client) = crate::openhuman::memory::global::client_if_ready() else { + tracing::debug!( + toolkit = %profile.toolkit, + "[composio:profile] memory client not ready, skipping profile persist" + ); + return 0; + }; + let conn = client.profile_conn(); + + let now = now_secs(); + let toolkit = &profile.toolkit; + + let fields: &[(&str, &Option)] = &[ + ("display_name", &profile.display_name), + ("email", &profile.email), + ("username", &profile.username), + ("avatar_url", &profile.avatar_url), + ]; + + let mut written = 0usize; + for (field_name, value) in fields { + let Some(val) = value else { continue }; + if val.trim().is_empty() { + continue; + } + + let key = format!("composio:{toolkit}:{field_name}"); + let facet_id = format!("composio-{toolkit}-{field_name}"); + + if let Err(e) = profile::profile_upsert( + &conn, + &facet_id, + &FacetType::Context, + &key, + val, + PROVIDER_CONFIDENCE, + None, // no source segment — this comes from a provider, not a conversation + now, + ) { + tracing::warn!( + toolkit = %toolkit, + field = %field_name, + error = %e, + "[composio:profile] profile_upsert failed (non-fatal)" + ); + continue; + } + written += 1; + } + + if written > 0 { + tracing::debug!( + toolkit = %toolkit, + facets_written = written, + "[composio:profile] persisted provider profile facets" + ); + } + written +} + +fn now_secs() -> f64 { + use std::time::{SystemTime, UNIX_EPOCH}; + SystemTime::now() + .duration_since(UNIX_EPOCH) + .map(|d| d.as_secs_f64()) + .unwrap_or(0.0) +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::openhuman::memory::store::profile::{profile_load_all, PROFILE_INIT_SQL}; + use parking_lot::Mutex; + use rusqlite::Connection; + use std::sync::Arc; + + fn setup_db() -> Arc> { + let conn = Connection::open_in_memory().unwrap(); + conn.execute_batch(PROFILE_INIT_SQL).unwrap(); + Arc::new(Mutex::new(conn)) + } + + /// Directly exercise `profile_upsert` with provider-style keys to + /// verify the facet schema and conflict resolution work end-to-end. + #[test] + fn provider_fields_map_to_facets() { + let conn = setup_db(); + let now = 1000.0; + + // Simulate what persist_provider_profile does internally. + profile::profile_upsert( + &conn, + "composio-gmail-email", + &FacetType::Context, + "composio:gmail:email", + "user@example.com", + PROVIDER_CONFIDENCE, + None, + now, + ) + .unwrap(); + + profile::profile_upsert( + &conn, + "composio-gmail-display_name", + &FacetType::Context, + "composio:gmail:display_name", + "Jane Doe", + PROVIDER_CONFIDENCE, + None, + now, + ) + .unwrap(); + + let facets = profile_load_all(&conn).unwrap(); + assert_eq!(facets.len(), 2); + + let email_facet = facets.iter().find(|f| f.key == "composio:gmail:email"); + assert!(email_facet.is_some()); + let email_facet = email_facet.unwrap(); + assert_eq!(email_facet.value, "user@example.com"); + assert!((email_facet.confidence - PROVIDER_CONFIDENCE).abs() < f64::EPSILON); + assert_eq!(email_facet.evidence_count, 1); + } + + #[test] + fn repeated_persist_increments_evidence() { + let conn = setup_db(); + + // First write. + profile::profile_upsert( + &conn, + "composio-notion-email", + &FacetType::Context, + "composio:notion:email", + "user@workspace.com", + PROVIDER_CONFIDENCE, + None, + 1000.0, + ) + .unwrap(); + + // Second write — same key, same value (periodic re-sync). + profile::profile_upsert( + &conn, + "composio-notion-email-2", + &FacetType::Context, + "composio:notion:email", + "user@workspace.com", + PROVIDER_CONFIDENCE, + None, + 2000.0, + ) + .unwrap(); + + let facets = profile_load_all(&conn).unwrap(); + assert_eq!(facets.len(), 1, "duplicate key should merge into one row"); + assert_eq!(facets[0].evidence_count, 2); + } + + #[test] + fn empty_fields_are_skipped() { + let profile = ProviderUserProfile { + toolkit: "gmail".into(), + connection_id: Some("conn-1".into()), + display_name: Some("Jane".into()), + email: None, + username: Some("".into()), // empty string — should be skipped + avatar_url: Some(" ".into()), // whitespace — should be skipped + extras: serde_json::Value::Null, + }; + + // We can't call persist_provider_profile directly without the + // global memory singleton, but we can verify the filtering logic + // by checking the field iteration manually. + let fields: &[(&str, &Option)] = &[ + ("display_name", &profile.display_name), + ("email", &profile.email), + ("username", &profile.username), + ("avatar_url", &profile.avatar_url), + ]; + let non_empty_count = fields + .iter() + .filter(|(_, v)| v.as_deref().is_some_and(|s| !s.trim().is_empty())) + .count(); + assert_eq!(non_empty_count, 1, "only display_name should pass"); + } +} diff --git a/src/openhuman/memory/store/client.rs b/src/openhuman/memory/store/client.rs index e14a794e1..00c2240a5 100644 --- a/src/openhuman/memory/store/client.rs +++ b/src/openhuman/memory/store/client.rs @@ -44,6 +44,19 @@ pub struct MemoryClient { } impl MemoryClient { + /// Returns a handle to the underlying SQLite connection for direct + /// profile-facet writes via + /// [`crate::openhuman::memory::store::unified::profile::profile_upsert`]. + /// + /// Intentionally `pub(crate)` — external consumers should use the + /// higher-level `MemoryClient` API; this escape hatch exists so + /// in-crate subsystems (composio providers, archivist, learning + /// hooks) can write structured profile facets without an additional + /// round-trip through the ingestion queue. + pub(crate) fn profile_conn(&self) -> std::sync::Arc> { + std::sync::Arc::clone(&self.inner.conn) + } + /// Create a new local memory client using the default `.openhuman` directory. /// /// # Errors diff --git a/src/openhuman/memory/store/unified/profile.rs b/src/openhuman/memory/store/unified/profile.rs index c3b0f2139..e729594be 100644 --- a/src/openhuman/memory/store/unified/profile.rs +++ b/src/openhuman/memory/store/unified/profile.rs @@ -120,8 +120,8 @@ pub fn profile_upsert( (None, None) => String::new(), }; - if confidence > existing_confidence { - // Higher confidence: overwrite value + update metadata. + if confidence >= existing_confidence { + // Higher or equal confidence: overwrite value + update metadata. conn.execute( "UPDATE user_profile SET value = ?2, confidence = ?3, evidence_count = ?4,