mirror of
https://github.com/tinyhumansai/openhuman.git
synced 2026-07-29 14:02:19 +00:00
feat(composio): provider folder modules + user profile persistence (#523)
* feat(composio): implement profile persistence for user data
- Introduced a new `profile` module to handle the persistence of user profile data from various providers into the local `user_profile` facet table.
- Enhanced the `composio_get_user_profile` and `fetch_user_profile` functions to call `persist_provider_profile`, ensuring that profile fields like display name, email, and avatar are stored locally for quick access.
- Added debug logging to track the number of facets written during the persistence process, improving observability of profile updates.
- This change aims to enhance user experience by reducing the need for repeated upstream API calls for frequently accessed profile information.
* style: apply cargo fmt to profile.rs and client.rs
* fix: PR review — cursor whitespace, stale profile values, test precision, log level
- Trim whitespace in cursor_to_gmail_after_filter before parsing so
leading/trailing spaces don't silently bypass the date filter.
- Change profile_upsert condition from `>` to `>=` so equal-confidence
provider refreshes replace stale values instead of only bumping
evidence_count.
- Tighten epoch millis test to assert exact date "2026/03/31" instead
of loose contains('/') check.
- Lower profile persistence log from info to debug to reduce noise in
normal connect/sync flows.
This commit is contained in:
@@ -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!(
|
||||
|
||||
+10
-147
@@ -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<SyncOutcome, String> {
|
||||
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<Value> {
|
||||
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<String> {
|
||||
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<String> {
|
||||
// Try parsing as epoch millis first (Gmail's internalDate).
|
||||
if let Ok(millis) = cursor.parse::<i64>() {
|
||||
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));
|
||||
}
|
||||
}
|
||||
@@ -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<Value> {
|
||||
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<String> {
|
||||
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<String> {
|
||||
let cursor = cursor.trim();
|
||||
// Try parsing as epoch millis first (Gmail's internalDate).
|
||||
if let Ok(millis) = cursor.parse::<i64>() {
|
||||
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)
|
||||
}
|
||||
@@ -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));
|
||||
}
|
||||
@@ -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!(
|
||||
|
||||
+13
-154
@@ -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<SyncOutcome, String> {
|
||||
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<Value> {
|
||||
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<String> {
|
||||
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<String> {
|
||||
// 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::<Vec<_>>()
|
||||
.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));
|
||||
}
|
||||
}
|
||||
@@ -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<Value> {
|
||||
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<String> {
|
||||
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<String> {
|
||||
// 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::<Vec<_>>()
|
||||
.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)
|
||||
}
|
||||
@@ -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));
|
||||
}
|
||||
@@ -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<String>)] = &[
|
||||
("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<Mutex<Connection>> {
|
||||
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<String>)] = &[
|
||||
("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");
|
||||
}
|
||||
}
|
||||
@@ -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<parking_lot::Mutex<rusqlite::Connection>> {
|
||||
std::sync::Arc::clone(&self.inner.conn)
|
||||
}
|
||||
|
||||
/// Create a new local memory client using the default `.openhuman` directory.
|
||||
///
|
||||
/// # Errors
|
||||
|
||||
@@ -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,
|
||||
|
||||
Reference in New Issue
Block a user