diff --git a/README.md b/README.md index e866aced..7edecd3b 100644 --- a/README.md +++ b/README.md @@ -63,7 +63,7 @@ else on request, so a host pays only for what it uses. | `documents` | `documents` | Format sniffing and conversion to markdown, producing `StoreItem::Document` | | `documents-office` | `documents::OfficeConverter` | PDF, DOCX, PPTX and XLSX to markdown (implies `documents`) | | `brain` | `brain` | Files into `tinymemory_tools::BrainDocument`s, by the source type their format implies (implies `documents`) | -| `sources` | `sources` | Folder, file and conversation readers, Composio normalisers (implies `documents`) | +| `sources` | `sources` | Folder, file and conversation readers (implies `documents`) | | `sources-network` | `sources::fetch` and the network readers | GitHub, RSS and web-page readers and `fetch_url`, behind the SSRF guard (implies `sources`) | | `safety` | `safety` | Secret and PII scrubbing of a `StoreItem` | | `legacy-import` | `import` | Migrating a v1 (embedded TinyCortex) workspace into any engine | diff --git a/crates/tinymemory-api/src/meta/mod.rs b/crates/tinymemory-api/src/meta/mod.rs index c7736efb..ae32e35f 100644 --- a/crates/tinymemory-api/src/meta/mod.rs +++ b/crates/tinymemory-api/src/meta/mod.rs @@ -156,26 +156,27 @@ pub enum SourceKind { Github, /// An RSS or Atom feed. Rss, - /// A Composio toolkit payload (Gmail, Slack, Notion, ...). - Composio, /// A host conversation. Conversation, /// Written by an agent directly; the default. #[default] Agent, /// Imported from a legacy store. + /// + /// The retired `composio` kind decodes to this variant, so items stored + /// under it stay readable. + #[serde(alias = "composio")] Import, } impl SourceKind { /// Every kind, in declaration order. - pub const ALL: [Self; 9] = [ + pub const ALL: [Self; 8] = [ Self::Folder, Self::File, Self::Link, Self::Github, Self::Rss, - Self::Composio, Self::Conversation, Self::Agent, Self::Import, @@ -190,7 +191,6 @@ impl SourceKind { Self::Link => "link", Self::Github => "github", Self::Rss => "rss", - Self::Composio => "composio", Self::Conversation => "conversation", Self::Agent => "agent", Self::Import => "import", diff --git a/crates/tinymemory-api/src/meta/mod_tests.rs b/crates/tinymemory-api/src/meta/mod_tests.rs index def04bdf..8299c7ad 100644 --- a/crates/tinymemory-api/src/meta/mod_tests.rs +++ b/crates/tinymemory-api/src/meta/mod_tests.rs @@ -43,3 +43,14 @@ fn source_kind_wire_strings_match_serde() { assert_eq!(json, serde_json::json!(kind.as_str())); } } + +#[test] +fn the_retired_composio_source_kind_decodes_as_import() { + let kind: SourceKind = + serde_json::from_value(serde_json::json!("composio")).expect("legacy kind decodes"); + assert_eq!(kind, SourceKind::Import); + assert_eq!( + serde_json::to_value(SourceKind::Import).expect("serialise"), + serde_json::json!("import") + ); +} diff --git a/crates/tinymemory-integrations/Cargo.toml b/crates/tinymemory-integrations/Cargo.toml index 69c651b7..c8e3b2ec 100644 --- a/crates/tinymemory-integrations/Cargo.toml +++ b/crates/tinymemory-integrations/Cargo.toml @@ -18,7 +18,7 @@ tinymemory-api = { path = "../tinymemory-api" } # `MemoryEngine`, `BearerSource`, `DocumentConverter` and `SourceReader` are # object-safe async traits. async-trait = { version = "0.1", optional = true } -# Every wire body, envelope, cursor, Composio payload and checkpoint is JSON. +# Every wire body, envelope, cursor and checkpoint is JSON. serde = { version = "1", features = ["derive"], optional = true } serde_json = { version = "1", optional = true } # The typed errors of `documents`, `sources` and `import`. @@ -122,7 +122,7 @@ documents-office = ["documents", "dep:pdf-extract", "dep:calamine", "dep:quick-x # placed under the source type their format implies. brain = ["documents", "dep:tinymemory-tools"] # `sources`: readers that turn folders, files and conversations into -# `StoreItem`s, and the Composio payload normalisers. Links no HTTP stack. +# `StoreItem`s. Links no HTTP stack. sources = ["documents", "dep:schemars", "dep:regex", "dep:walkdir", "dep:chrono", "dep:log", "dep:tracing"] # The readers that fetch over the network — GitHub, RSS, web pages — and # `sources::fetch::fetch_url`, all behind the shared SSRF guard. diff --git a/crates/tinymemory-integrations/README.md b/crates/tinymemory-integrations/README.md index 6d02b39a..075f7c80 100644 --- a/crates/tinymemory-integrations/README.md +++ b/crates/tinymemory-integrations/README.md @@ -17,7 +17,7 @@ wants one integration enables one feature and links nothing else. | `cortex`, `registry`, `config` | `cortex` (default) | `CortexEngine` over two wires (`cortexdb`, `tinyhumans`); `list_engines`, `build_engine`, `EngineCredential`; `MemoryConfig` | [`src/cortex/README.md`](src/cortex/README.md) | | `documents` | `documents` | Format sniffing and conversion to markdown, producing `StoreItem::Document`. No I/O. | [`src/documents/README.md`](src/documents/README.md) | | `documents::OfficeConverter` | `documents-office` | PDF, DOCX, PPTX and XLSX to markdown, in process | (same) | -| `sources` | `sources` | Readers for folders, files and conversations; Composio payload normalisers; `collect_items`. Links no HTTP stack. | [`src/sources/README.md`](src/sources/README.md) | +| `sources` | `sources` | Readers for folders, files and conversations; `collect_items`. Links no HTTP stack. | [`src/sources/README.md`](src/sources/README.md) | | `sources::fetch`, GitHub, RSS and web-page readers | `sources-network` | The network readers and `fetch_url`, all behind the SSRF guard | (same) | | `safety` | `safety` | Secret and PII scrubbing of a `StoreItem` before it is stored | [`src/safety/README.md`](src/safety/README.md) | | `import` | `legacy-import` | Reads a v1 (embedded TinyCortex) workspace and migrates it into any engine, resumably | [`src/import/README.md`](src/import/README.md) | diff --git a/crates/tinymemory-integrations/src/lib.rs b/crates/tinymemory-integrations/src/lib.rs index f320a470..33a19e6b 100644 --- a/crates/tinymemory-integrations/src/lib.rs +++ b/crates/tinymemory-integrations/src/lib.rs @@ -8,7 +8,7 @@ //! | [`cortex`], [`registry`], [`config`] | `cortex` (default) | The CortexDB engine over its two wires, and building one from configuration | //! | `documents` | `documents`, `documents-office` | Format sniffing and conversion to markdown, emitting `StoreItem::Document` | //! | `brain` | `brain` | Converting files into the brain documents `tinymemory_tools::Brain` ingests, by source type | -//! | `sources` | `sources`, `sources-network` | Readers turning folders, files, links, GitHub, RSS, Composio payloads and conversations into `StoreItem`s | +//! | `sources` | `sources`, `sources-network` | Readers turning folders, files, links, GitHub, RSS, and conversations into `StoreItem`s | //! | `safety` | `safety` | Secret and PII scrubbing for a `StoreItem` before it is stored | //! | `import` | `legacy-import` | Migrating a legacy v1 (embedded TinyCortex) workspace into any engine | //! diff --git a/crates/tinymemory-integrations/src/sources/README.md b/crates/tinymemory-integrations/src/sources/README.md index dca6befa..0e151d6f 100644 --- a/crates/tinymemory-integrations/src/sources/README.md +++ b/crates/tinymemory-integrations/src/sources/README.md @@ -2,7 +2,7 @@ `tinymemory_integrations::sources` (features `sources` and `sources-network`): readers that turn a source into `StoreItem`s — a folder, a single file, a web -page, a GitHub repository, an RSS feed, a Composio toolkit payload, or the +page, a GitHub repository, an RSS feed, or the host's local conversation threads. Conversion to markdown and language detection come from the sibling [`documents`](../documents/README.md) module. Architecture overview: @@ -21,7 +21,6 @@ checks it with `MemorySourceEntry::validate`, nothing more. | `fetch` | one URL into a `RawDocument` or a link item (`sources-network`); the RSS and web-page readers fetch through it with their own body caps | | `fetch::ssrf` | the SSRF guard: scheme and host policy, one address classifier for literal and resolved addresses, a public-only DNS resolver, per-hop redirect checks, and a capped body reader | | `items` | reader output to `StoreItem`s with `MemoryMeta` filled per kind; `collect_items` drives a reader end to end | -| `composio` | toolkit normalisers (Gmail, Slack, GitHub, Linear, Notion, ClickUp), the `fields::pick_str` lookup they share, `normalise_payload` and `payload_items`. `readers::composio::ComposioReader` is only a placeholder reader | | `error` | the module `Error`, mapped onto `tinymemory_api::Error` | ## Kinds and metadata @@ -35,7 +34,6 @@ Every item's `meta.source` is `SourceRef { kind, id: Some(entry.id) }`. | `web_page` | `Link` | document | `url` | | `github_repo` | `Github` | document | `repo` (`owner/name`), `commit` (commits), `url` (issues, PRs), `observed_at` | | `rss_feed` | `Rss` | document | `url` (the entry's link), `observed_at` (published) | -| `composio` | `Composio` | document | `tags = [toolkit]`; payloads add `url`, `observed_at`, `thread_id`, `repo` | | `conversation` | `Conversation` | conversation | `workspace`, `thread_id`, `turns`, `observed_at` (last turn) | ## Folder selection @@ -74,7 +72,7 @@ named and numeric entities decode the same way everywhere. ## Features -- `sources` — the local readers, `items`, `composio` and `types`; implies +- `sources` — the local readers, `items` and `types`; implies `documents`. Links no HTTP stack. - `sources-network` — adds the GitHub, RSS and web-page readers, `fetch`, and the SSRF guard. diff --git a/crates/tinymemory-integrations/src/sources/composio/clickup/mod.rs b/crates/tinymemory-integrations/src/sources/composio/clickup/mod.rs deleted file mode 100644 index b4ce64e5..00000000 --- a/crates/tinymemory-integrations/src/sources/composio/clickup/mod.rs +++ /dev/null @@ -1,69 +0,0 @@ -//! ClickUp host normalization helpers — result extraction and task-title and -//! timestamp extraction. -//! -//! ClickUp's REST API (and therefore Composio's wrapping of it) returns -//! task lists in a small handful of shapes depending on which endpoint -//! is called. The functions here walk the union of common shapes so the -//! provider doesn't have to branch per Composio envelope variant. - -use serde_json::Value; - -use super::fields::pick_str; - -/// Walk the Composio response envelope for ClickUp task list results. -/// -/// ClickUp's "filtered team tasks" endpoint returns `{ "tasks": [...] }` -/// at the top level; Composio re-wraps the upstream payload under -/// `data` or `data.data` depending on the action. We probe each shape -/// in order and return the first array we find. -pub fn extract_tasks(data: &Value) -> Vec { - let candidates = [ - data.pointer("/data/tasks"), - data.pointer("/tasks"), - data.pointer("/data/data/tasks"), - data.pointer("/data/results"), - data.pointer("/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 a human-readable title from a ClickUp task object. -/// -/// ClickUp tasks store the name at `name` (or `data.name` after Composio -/// envelope wrapping). When the name is missing we fall back to the -/// task ID so chunks remain identifiable. -pub fn extract_task_name(task: &Value) -> Option { - pick_str(task, &["name", "data.name", "title", "data.title"]) -} - -/// Extract a stable cursor timestamp (milliseconds since epoch as a -/// string) from a ClickUp task object. -/// -/// The ClickUp API returns `date_updated` as a stringified epoch ms -/// (e.g. `"1733412345678"`); we keep it as a string so lexicographic -/// comparison against the stored cursor remains valid as long as the -/// length doesn't change (it won't until year 33658). -pub fn extract_task_updated(task: &Value) -> Option { - pick_str( - task, - &[ - "date_updated", - "data.date_updated", - "updated_at", - "data.updated_at", - "dateUpdated", - "data.dateUpdated", - ], - ) -} - -#[cfg(test)] -#[path = "mod_tests.rs"] -mod tests; diff --git a/crates/tinymemory-integrations/src/sources/composio/clickup/mod_tests.rs b/crates/tinymemory-integrations/src/sources/composio/clickup/mod_tests.rs deleted file mode 100644 index 5a520401..00000000 --- a/crates/tinymemory-integrations/src/sources/composio/clickup/mod_tests.rs +++ /dev/null @@ -1,58 +0,0 @@ -//! Tests for the ClickUp normaliser. - -use super::*; -use serde_json::json; - -#[test] -fn extract_tasks_from_data_tasks() { - let data = json!({ "data": { "tasks": [{"id": "t1"}] } }); - assert_eq!(extract_tasks(&data).len(), 1); -} - -#[test] -fn extract_tasks_from_top_level_tasks() { - let data = json!({ "tasks": [{"id": "a"}, {"id": "b"}] }); - assert_eq!(extract_tasks(&data).len(), 2); -} - -#[test] -fn extract_tasks_empty_when_missing() { - let data = json!({ "foo": "bar" }); - assert!(extract_tasks(&data).is_empty()); -} - -#[test] -fn extract_task_name_from_top_level() { - let task = json!({ "id": "t1", "name": "Build feature X" }); - assert_eq!(extract_task_name(&task), Some("Build feature X".into())); -} - -#[test] -fn extract_task_name_falls_back_to_data_name() { - let task = json!({ "data": { "name": "Wrapped" } }); - assert_eq!(extract_task_name(&task), Some("Wrapped".into())); -} - -#[test] -fn extract_task_name_none_when_missing() { - let task = json!({ "id": "t1" }); - assert!(extract_task_name(&task).is_none()); -} - -#[test] -fn extract_task_updated_handles_string_form() { - let task = json!({ "date_updated": "1733412345678" }); - assert_eq!( - extract_task_updated(&task), - Some("1733412345678".to_string()) - ); -} - -#[test] -fn extract_task_updated_handles_nested_data() { - let task = json!({ "data": { "dateUpdated": "1700000000000" } }); - assert_eq!( - extract_task_updated(&task), - Some("1700000000000".to_string()) - ); -} diff --git a/crates/tinymemory-integrations/src/sources/composio/documents/mod.rs b/crates/tinymemory-integrations/src/sources/composio/documents/mod.rs deleted file mode 100644 index 04753e4b..00000000 --- a/crates/tinymemory-integrations/src/sources/composio/documents/mod.rs +++ /dev/null @@ -1,343 +0,0 @@ -//! One Composio response in, documents and `StoreItem`s out. -//! -//! [`normalise_payload`] dispatches on the toolkit slug to the matching -//! normaliser and reads each record's title, body, link and timestamp. -//! Toolkits without a dedicated normaliser fall back to a generic walk that -//! keeps each record as fenced JSON, so a new toolkit is ingested (verbosely) -//! rather than dropped. [`payload_items`] wraps the documents as -//! `StoreItem::Document`s. - -use crate::documents::{DocumentFormat, markdown_from_text}; -use chrono::{DateTime, TimeZone, Utc}; -use serde_json::Value; -use tinymemory_api::{DocumentBody, MemoryMeta, SourceKind, StoreItem}; - -use super::fields::pick_str; -use super::{clickup, github, gmail_post_process, linear, notion}; - -/// One record of a Composio payload, normalised: an email, a message, an -/// issue, a task or a page. -#[derive(Debug, Clone, PartialEq)] -pub struct ComposioDocument { - /// The provider's id for the record, when it has one. - pub id: Option, - /// A human-readable title. - pub title: Option, - /// The record's text as markdown. Never empty. - pub body: String, - /// A link back to the record in the provider's UI. - pub url: Option, - /// When the record was last updated or sent. - pub observed_at: Option>, - /// The email or message thread the record belongs to. - pub thread_id: Option, - /// The repository a GitHub record belongs to, as `owner/name`. - pub repo: Option, -} - -impl ComposioDocument { - /// A record with only a body. - fn with_body(body: String) -> Self { - Self { - id: None, - title: None, - body, - url: None, - observed_at: None, - thread_id: None, - repo: None, - } - } - - /// Wrap this record as a [`StoreItem::Document`] from `toolkit`, read - /// through the connection or source `source_id`. - #[must_use] - pub fn into_store_item(self, toolkit: &str, source_id: &str) -> StoreItem { - let mut meta = MemoryMeta::from_source(SourceKind::Composio, Some(source_id.to_string())); - meta.url = self.url; - meta.observed_at = self.observed_at; - meta.thread_id = self.thread_id; - meta.repo = self.repo; - meta.tags = vec![toolkit.to_string()]; - StoreItem::Document { - title: self.title, - body: DocumentBody::Text(self.body), - mime: Some(DocumentFormat::Markdown.mime().to_string()), - meta, - } - } -} - -/// Normalise one Composio response from `toolkit` into documents. -/// -/// `data` is the action's response; for Gmail and Slack, run -/// [`gmail_post_process::post_process`] / [`super::slack_post_process::post_process`] -/// on it first so it carries the slim `messages[]` shape. Records with no text -/// at all are skipped. -#[must_use] -pub fn normalise_payload(toolkit: &str, data: &Value) -> Vec { - let documents: Vec = match toolkit.to_ascii_lowercase().as_str() { - "gmail" => array_at(data, &["/messages", "/data/messages"]) - .iter() - .filter_map(gmail_message) - .collect(), - "slack" => array_at(data, &["/messages", "/data/messages"]) - .iter() - .filter_map(slack_message) - .collect(), - "github" => github::extract_issues(data) - .iter() - .filter_map(github_issue) - .collect(), - "linear" => linear::extract_issues(data) - .iter() - .filter_map(linear_issue) - .collect(), - "notion" => notion_pages(data), - "clickup" => clickup::extract_tasks(data) - .iter() - .filter_map(clickup_task) - .collect(), - _ => generic_records(data), - }; - log::debug!( - "[memory_sources:composio] normalised toolkit={toolkit} documents={}", - documents.len() - ); - documents -} - -/// Normalise one Composio response and wrap every record as a -/// [`StoreItem::Document`] with `source.kind = Composio`, -/// `source.id = source_id` and `tags = [toolkit]`. -#[must_use] -pub fn payload_items(toolkit: &str, source_id: &str, data: &Value) -> Vec { - normalise_payload(toolkit, data) - .into_iter() - .map(|document| document.into_store_item(toolkit, source_id)) - .collect() -} - -/// The first array found at any of `pointers`. -fn array_at<'a>(data: &'a Value, pointers: &[&str]) -> &'a [Value] { - pointers - .iter() - .find_map(|pointer| data.pointer(pointer).and_then(Value::as_array)) - .map_or(&[], Vec::as_slice) -} - -/// A body as markdown: HTML (as sniffed) is converted, anything else kept. -fn to_markdown(text: &str) -> String { - let format = DocumentFormat::sniff(text.as_bytes(), None, None); - markdown_from_text(text, format).trim().to_string() -} - -/// Build a document from its parts, falling back to the title as the body; -/// `None` when there is no text at all. -fn document(title: Option, body: Option) -> Option { - let body = body - .map(|body| to_markdown(&body)) - .filter(|body| !body.is_empty()) - .or_else(|| title.clone())?; - let mut document = ComposioDocument::with_body(body); - document.title = title; - Some(document) -} - -/// Parse an ISO 8601 / RFC 3339 / RFC 2822 timestamp. -fn parse_time(text: &str) -> Option> { - gmail_post_process::parse_email_date(text) -} - -/// Parse an epoch-milliseconds string (ClickUp's `date_updated`). -fn parse_epoch_ms(text: &str) -> Option> { - let millis = text.trim().parse::().ok()?; - Utc.timestamp_millis_opt(millis).single() -} - -/// Parse a Slack `ts` (`"1712345678.123456"`, epoch seconds with a fraction). -fn parse_slack_ts(text: &str) -> Option> { - let (seconds, fraction) = text.split_once('.').unwrap_or((text, "0")); - let seconds = seconds.parse::().ok()?; - let micros = format!("{fraction:0<6}").get(..6)?.parse::().ok()?; - Utc.timestamp_opt(seconds, micros * 1_000).single() -} - -/// A string field, or a number rendered as a string. -fn scalar(value: &Value, key: &str) -> Option { - match value.get(key)? { - Value::String(text) if !text.trim().is_empty() => Some(text.trim().to_string()), - Value::Number(number) => Some(number.to_string()), - _ => None, - } -} - -/// A post-processed Gmail message: headers above the body. -fn gmail_message(message: &Value) -> Option { - let subject = pick_str(message, &["subject"]); - let markdown = pick_str(message, &["markdown", "messageText"]); - let mut header = String::new(); - for (label, key) in [("From", "from"), ("To", "to"), ("Date", "date")] { - if let Some(value) = pick_str(message, &[key]) { - header.push_str(&format!("{label}: {value}\n")); - } - } - let body = match (markdown, header.is_empty()) { - (Some(markdown), false) => Some(format!("{header}\n{markdown}")), - (Some(markdown), true) => Some(markdown), - (None, _) => None, - }; - let mut document = document(subject, body)?; - document.id = pick_str(message, &["id", "messageId"]); - document.thread_id = pick_str(message, &["threadId", "thread_id"]); - document.observed_at = pick_str(message, &["date"]).as_deref().and_then(parse_time); - Some(document) -} - -/// A post-processed Slack message. -fn slack_message(message: &Value) -> Option { - let user = pick_str(message, &["user"]); - let channel = pick_str(message, &["channel_id"]); - let title = match (&user, &channel) { - (Some(user), Some(channel)) => Some(format!("Slack message from {user} in {channel}")), - (Some(user), None) => Some(format!("Slack message from {user}")), - (None, _) => Some("Slack message".to_string()), - }; - let text = pick_str(message, &["text"])?; - let mut document = document(title, Some(text))?; - let ts = pick_str(message, &["ts"]); - document.id = ts.clone(); - document.url = pick_str(message, &["permalink"]); - document.thread_id = pick_str(message, &["thread_ts"]); - document.observed_at = ts.as_deref().and_then(parse_slack_ts); - Some(document) -} - -/// A GitHub issue or pull request from a search response. -fn github_issue(issue: &Value) -> Option { - let mut document = document( - github::extract_issue_title(issue), - pick_str(issue, &["body", "data.body"]), - )?; - document.id = github::extract_issue_id(issue); - document.url = pick_str(issue, &["html_url", "data.html_url"]); - document.repo = document.url.as_deref().and_then(github_repo); - document.observed_at = github::extract_issue_updated_at(issue) - .as_deref() - .and_then(parse_time); - Some(document) -} - -/// `owner/name` from a `https://github.com/owner/name/...` link. -fn github_repo(url: &str) -> Option { - let rest = url.split_once("github.com/")?.1; - let mut parts = rest.split('/'); - let owner = parts.next().filter(|part| !part.is_empty())?; - let name = parts.next().filter(|part| !part.is_empty())?; - Some(format!("{owner}/{name}")) -} - -/// A Linear issue. -fn linear_issue(issue: &Value) -> Option { - let mut document = document( - linear::extract_issue_title(issue), - pick_str(issue, &["description", "data.description"]), - )?; - document.id = pick_str(issue, &["identifier", "id", "data.identifier", "data.id"]); - document.url = pick_str(issue, &["url", "data.url"]); - document.observed_at = linear::extract_issue_updated(issue) - .as_deref() - .and_then(parse_time); - Some(document) -} - -/// Notion pages from a search response, or the one page a -/// `NOTION_GET_PAGE_MARKDOWN` response carries. -fn notion_pages(data: &Value) -> Vec { - let results = notion::extract_results(data); - if results.is_empty() { - return notion::extract_page_markdown(data) - .and_then(|markdown| { - let title = notion::extract_page_title(data); - let mut document = document(title, Some(markdown))?; - document.id = pick_str(data, &["id", "data.id", "page_id", "data.page_id"]); - document.url = pick_str(data, &["url", "data.url"]); - Some(document) - }) - .into_iter() - .collect(); - } - results - .iter() - .filter_map(|page| { - let mut document = document( - notion::extract_page_title(page), - notion::extract_page_markdown(page), - )?; - document.id = pick_str(page, &["id", "data.id"]); - document.url = pick_str(page, &["url", "data.url"]); - document.observed_at = pick_str(page, &["last_edited_time", "data.last_edited_time"]) - .as_deref() - .and_then(parse_time); - Some(document) - }) - .collect() -} - -/// A ClickUp task. -fn clickup_task(task: &Value) -> Option { - let mut document = document( - clickup::extract_task_name(task), - pick_str( - task, - &["markdown_description", "description", "text_content"], - ), - )?; - document.id = scalar(task, "id"); - document.url = pick_str(task, &["url", "data.url"]); - document.observed_at = clickup::extract_task_updated(task) - .as_deref() - .and_then(|text| parse_epoch_ms(text).or_else(|| parse_time(text))); - Some(document) -} - -/// Any other toolkit: each record in the first list found (or the whole -/// payload) as fenced JSON, titled and linked when it says how. -fn generic_records(data: &Value) -> Vec { - let records = array_at( - data, - &[ - "/data/items", - "/items", - "/data/results", - "/results", - "/data/data", - "/data", - ], - ); - let records: Vec<&Value> = if records.is_empty() { - vec![data] - } else { - records.iter().collect() - }; - records - .into_iter() - .filter(|record| !record.is_null()) - .filter_map(|record| { - let json = serde_json::to_string_pretty(record).ok()?; - let title = pick_str(record, &["title", "name", "subject"]); - let mut document = ComposioDocument::with_body(format!("```json\n{json}\n```")); - document.title = title; - document.id = scalar(record, "id"); - document.url = pick_str(record, &["url", "html_url", "permalink", "link"]); - document.observed_at = pick_str(record, &["updated_at", "updatedAt", "created_at"]) - .as_deref() - .and_then(parse_time); - Some(document) - }) - .collect() -} - -#[cfg(test)] -#[path = "mod_tests.rs"] -mod tests; diff --git a/crates/tinymemory-integrations/src/sources/composio/documents/mod_tests.rs b/crates/tinymemory-integrations/src/sources/composio/documents/mod_tests.rs deleted file mode 100644 index 4d689128..00000000 --- a/crates/tinymemory-integrations/src/sources/composio/documents/mod_tests.rs +++ /dev/null @@ -1,197 +0,0 @@ -//! Tests for turning Composio payloads into documents and `StoreItem`s. - -use super::*; -use serde_json::json; -use tinymemory_api::ItemKind; - -fn text_of(item: &StoreItem) -> (&Option, &str, &MemoryMeta) { - match item { - StoreItem::Document { - title, - body: DocumentBody::Text(text), - meta, - .. - } => (title, text.as_str(), meta), - other => panic!("expected a text document, got {other:?}"), - } -} - -#[test] -fn a_github_search_becomes_items_with_url_repo_time_and_toolkit_tag() { - let data = json!({ - "data": { "items": [{ - "id": 42, - "title": "Fix the build", - "body": "The build is **red**.", - "html_url": "https://github.com/acme/widgets/issues/7", - "updated_at": "2024-05-21T15:30:00Z" - }]} - }); - let items = payload_items("github", "conn_gh", &data); - assert_eq!(items.len(), 1); - let item = &items[0]; - assert_eq!(item.kind(), ItemKind::Document); - item.validate().unwrap(); - - let (title, body, meta) = text_of(item); - assert_eq!( - title.as_deref(), - Some("GitHub: acme/widgets#7: Fix the build") - ); - assert_eq!(body, "The build is **red**."); - assert_eq!(meta.source.kind, SourceKind::Composio); - assert_eq!(meta.source.id.as_deref(), Some("conn_gh")); - assert_eq!(meta.tags, vec!["github".to_string()]); - assert_eq!( - meta.url.as_deref(), - Some("https://github.com/acme/widgets/issues/7") - ); - assert_eq!(meta.repo.as_deref(), Some("acme/widgets")); - assert_eq!( - meta.observed_at, - Some(Utc.with_ymd_and_hms(2024, 5, 21, 15, 30, 0).unwrap()) - ); -} - -#[test] -fn post_processed_gmail_messages_carry_headers_thread_and_date() { - let data = json!({ "messages": [{ - "id": "m1", - "threadId": "t1", - "subject": "Lunch?", - "from": "Ann ", - "to": "me@example.com", - "date": "Tue, 21 May 2024 12:00:00 +0000", - "markdown": "Tacos at noon." - }]}); - let documents = normalise_payload("gmail", &data); - assert_eq!(documents.len(), 1); - let document = &documents[0]; - assert_eq!(document.title.as_deref(), Some("Lunch?")); - assert!(document.body.starts_with("From: Ann \n")); - assert!(document.body.ends_with("Tacos at noon.")); - assert_eq!(document.thread_id.as_deref(), Some("t1")); - assert_eq!(document.id.as_deref(), Some("m1")); - assert_eq!( - document.observed_at, - Some(Utc.with_ymd_and_hms(2024, 5, 21, 12, 0, 0).unwrap()) - ); -} - -#[test] -fn slack_messages_take_permalink_and_ts_time() { - let data = json!({ "messages": [{ - "ts": "1716300000.000100", - "user": "U1", - "channel_id": "C1", - "text": "deploy done", - "permalink": "https://acme.slack.com/archives/C1/p1716300000000100" - }]}); - let items = payload_items("slack", "src_slack", &data); - let (title, body, meta) = text_of(&items[0]); - assert_eq!(title.as_deref(), Some("Slack message from U1 in C1")); - assert_eq!(body, "deploy done"); - assert_eq!( - meta.url.as_deref(), - Some("https://acme.slack.com/archives/C1/p1716300000000100") - ); - assert_eq!( - meta.observed_at.map(|at| at.timestamp()), - Some(1_716_300_000) - ); - assert_eq!(meta.tags, vec!["slack".to_string()]); -} - -#[test] -fn linear_notion_and_clickup_records_are_normalised() { - let linear = json!({ "nodes": [{ - "identifier": "ENG-1", - "title": "Ship v2", - "description": "All of it.", - "url": "https://linear.app/acme/issue/ENG-1", - "updatedAt": "2024-01-02T03:04:05Z" - }]}); - let documents = normalise_payload("linear", &linear); - assert_eq!(documents[0].id.as_deref(), Some("ENG-1")); - assert_eq!(documents[0].body, "All of it."); - assert!(documents[0].observed_at.is_some()); - - let notion = json!({ "results": [{ - "id": "p1", - "url": "https://notion.so/p1", - "last_edited_time": "2024-01-02T03:04:05.000Z", - "properties": { "Name": { "type": "title", "title": [{ "plain_text": "Roadmap" }] } } - }]}); - let documents = normalise_payload("notion", ¬ion); - assert_eq!(documents[0].title.as_deref(), Some("Roadmap")); - assert_eq!( - documents[0].body, "Roadmap", - "a page without body keeps its title" - ); - assert_eq!(documents[0].url.as_deref(), Some("https://notion.so/p1")); - - let page_markdown = json!({ "data": { "markdown": "# Plan\n\nDo it." }, "id": "p2" }); - let documents = normalise_payload("notion", &page_markdown); - assert_eq!(documents.len(), 1); - assert_eq!(documents[0].body, "# Plan\n\nDo it."); - - let clickup = json!({ "tasks": [{ - "id": "t9", - "name": "Write docs", - "description": "

The README

", - "url": "https://app.clickup.com/t/t9", - "date_updated": "1700000000000" - }]}); - let documents = normalise_payload("clickup", &clickup); - assert_eq!(documents[0].id.as_deref(), Some("t9")); - assert_eq!( - documents[0].observed_at.map(|at| at.timestamp()), - Some(1_700_000_000) - ); -} - -#[test] -fn html_bodies_are_converted_to_markdown() { - let data = json!({ "nodes": [{ - "title": "Page", - "description": "

Hi

There

" - }]}); - let documents = normalise_payload("linear", &data); - assert_eq!(documents[0].body, "## Hi\n\nThere"); -} - -#[test] -fn records_without_any_text_are_skipped() { - let data = json!({ "messages": [{ "ts": "1.0", "text": " " }] }); - assert!(normalise_payload("slack", &data).is_empty()); - assert!(normalise_payload("github", &json!({ "items": [{}] })).is_empty()); -} - -#[test] -fn an_unknown_toolkit_keeps_each_record_as_fenced_json() { - let data = json!({ "data": { "items": [ - { "id": 1, "name": "Alpha", "url": "https://example.com/1", "updated_at": "2024-01-01T00:00:00Z" }, - { "id": 2 } - ]}}); - let items = payload_items("hubspot", "conn_hs", &data); - assert_eq!(items.len(), 2); - let (title, body, meta) = text_of(&items[0]); - assert_eq!(title.as_deref(), Some("Alpha")); - assert!(body.starts_with("```json\n")); - assert!(body.contains("\"name\": \"Alpha\"")); - assert_eq!(meta.url.as_deref(), Some("https://example.com/1")); - assert_eq!(meta.tags, vec!["hubspot".to_string()]); - - let single = normalise_payload("hubspot", &json!({ "ok": true })); - assert_eq!(single.len(), 1); -} - -#[test] -fn slack_timestamps_parse_with_and_without_a_fraction() { - assert_eq!(parse_slack_ts("10").map(|at| at.timestamp()), Some(10)); - assert_eq!( - parse_slack_ts("10.5").map(|at| at.timestamp_subsec_micros()), - Some(500_000) - ); - assert_eq!(parse_slack_ts("x"), None); -} diff --git a/crates/tinymemory-integrations/src/sources/composio/fields/mod.rs b/crates/tinymemory-integrations/src/sources/composio/fields/mod.rs deleted file mode 100644 index 4c257085..00000000 --- a/crates/tinymemory-integrations/src/sources/composio/fields/mod.rs +++ /dev/null @@ -1,44 +0,0 @@ -//! Field lookup shared by the Composio normalisers: pull a string out of a -//! payload by trying several dotted paths, because Composio wraps the same -//! upstream field at different depths depending on the action and version. - -/// Walk a JSON object using a list of dotted-path candidates and return the -/// first non-empty **string** match, trimmed. -/// -/// Each path is split on `.` and followed with `Value::get`, so it only -/// descends through objects — it never indexes into an array. A leaf that is -/// not a string (a number, a bool) is rejected rather than coerced, so a -/// payload whose `id` is `42` rather than `"42"` yields `None` here. That -/// differs from the private `scalar` lookup in the `documents` mapping, which -/// renders numbers; the normalisers were written against the -/// reject-non-strings behaviour and `pick_str_rejects_non_string_values` pins -/// it. -pub fn pick_str(value: &serde_json::Value, paths: &[&str]) -> Option { - for path in paths { - let mut cur = value; - let mut ok = true; - for segment in path.split('.') { - match cur.get(segment) { - Some(next) => cur = next, - None => { - ok = false; - break; - } - } - } - if !ok { - continue; - } - if let Some(s) = cur.as_str() { - let trimmed = s.trim(); - if !trimmed.is_empty() { - return Some(trimmed.to_string()); - } - } - } - None -} - -#[cfg(test)] -#[path = "mod_tests.rs"] -mod tests; diff --git a/crates/tinymemory-integrations/src/sources/composio/fields/mod_tests.rs b/crates/tinymemory-integrations/src/sources/composio/fields/mod_tests.rs deleted file mode 100644 index 330ec404..00000000 --- a/crates/tinymemory-integrations/src/sources/composio/fields/mod_tests.rs +++ /dev/null @@ -1,37 +0,0 @@ -//! Tests for the shared Composio field lookup. - -use super::*; -use serde_json::json; - -#[test] -fn pick_str_finds_first_non_empty_match() { - let v = json!({"data": {"user": {"name": "Ada", "email": "ada@example.com"}}}); - assert_eq!( - pick_str(&v, &["data.user.name", "data.user.email"]), - Some("Ada".into()) - ); - assert_eq!( - pick_str(&v, &["data.missing", "data.user.email"]), - Some("ada@example.com".into()) - ); - assert_eq!(pick_str(&v, &["nope.nope"]), None); -} - -#[test] -fn pick_str_respects_path_order() { - let v = json!({"a": "first", "b": "second"}); - assert_eq!(pick_str(&v, &["a", "b"]), Some("first".into())); - assert_eq!(pick_str(&v, &["b", "a"]), Some("second".into())); -} - -/// The drift guard for the behaviour documented on [`pick_str`]. If this -/// ever starts returning `Some("42")`, the normalisers' emitted ids have -/// changed. -#[test] -fn pick_str_rejects_non_string_values() { - let v = json!({"count": 42, "flag": true, "empty": "", "whitespace": " "}); - assert_eq!(pick_str(&v, &["count"]), None); - assert_eq!(pick_str(&v, &["flag"]), None); - assert_eq!(pick_str(&v, &["empty"]), None); - assert_eq!(pick_str(&v, &["whitespace"]), None); -} diff --git a/crates/tinymemory-integrations/src/sources/composio/github/mod.rs b/crates/tinymemory-integrations/src/sources/composio/github/mod.rs deleted file mode 100644 index 12e7a615..00000000 --- a/crates/tinymemory-integrations/src/sources/composio/github/mod.rs +++ /dev/null @@ -1,116 +0,0 @@ -//! GitHub host normalization helpers — issue extraction and issue id, title and -//! timestamp helpers. -//! -//! GitHub's REST API (proxied through Composio) returns search results in a -//! small number of shapes. The functions here -//! walk the union of common Composio envelope variants so the provider stays -//! clean and branch-free. - -use serde_json::Value; - -use super::fields::pick_str; - -/// Walk the Composio response envelope for GitHub search issue results. -/// -/// `GITHUB_SEARCH_ISSUES_AND_PULL_REQUESTS` wraps GitHub's `GET /search/issues` response, which -/// returns `{"total_count": N, "items": [...]}`. Composio may re-wrap this under -/// `data` or `data.data`; we probe each shape in order. -pub fn extract_issues(data: &Value) -> Vec { - let candidates = [ - data.pointer("/data/items"), - data.pointer("/items"), - data.pointer("/data/data/items"), - data.pointer("/data/results"), - data.pointer("/results"), - ]; - for cand in candidates.into_iter().flatten() { - if let Some(arr) = cand.as_array() { - return arr.clone(); - } - } - Vec::new() -} - -/// Extract a stable, globally unique identifier for a GitHub issue or PR. -/// -/// GitHub's internal `id` field is a large integer unique across all issues -/// and PRs on github.com. We convert it to a string for use as a sync key. -/// Falls back to composing from `html_url` path if `id` is absent. -pub fn extract_issue_id(issue: &Value) -> Option { - // Primary: numeric internal GitHub ID. - if let Some(id) = issue.get("id").or_else(|| issue.pointer("/data/id")) { - if let Some(n) = id.as_u64() { - return Some(n.to_string()); - } - if let Some(s) = id.as_str() { - let trimmed = s.trim(); - if !trimmed.is_empty() { - return Some(trimmed.to_string()); - } - } - } - // Fallback: parse owner/repo/number from html_url path segments. - // URL shape: https://github.com/{owner}/{repo}/issues/{number} - if let Some(url) = pick_str(issue, &["html_url", "data.html_url", "url", "data.url"]) - && let Some(slug) = github_url_to_slug(&url) - { - return Some(slug); - } - None -} - -/// Build a human-readable document title for a GitHub issue/PR. -/// -/// Format: `GitHub: {owner}/{repo}#{number}: {title}`. -/// Falls back to just the title or a placeholder when fields are missing. -pub fn extract_issue_title(issue: &Value) -> Option { - let title = pick_str(issue, &["title", "data.title"])?; - - // Best-effort: extract owner/repo#N from html_url for the prefix. - let prefix = pick_str(issue, &["html_url", "data.html_url"]) - .and_then(|url| github_url_to_slug(&url)) - .unwrap_or_default(); - - if prefix.is_empty() { - Some(title) - } else { - Some(format!("GitHub: {prefix}: {title}")) - } -} - -/// Parse `https://github.com/{owner}/{repo}/issues/{number}` (or `/pull/`) -/// into `"{owner}/{repo}#{number}"`. Returns `None` for unrecognised shapes. -fn github_url_to_slug(url: &str) -> Option { - let segs: Vec<&str> = url.trim_end_matches('/').split('/').collect(); - // Minimum: ["https:", "", "github.com", owner, repo, "issues", number] - if segs.len() >= 7 { - let number = segs[segs.len() - 1]; - let _kind = segs[segs.len() - 2]; // "issues" or "pull" — ignored - let repo = segs[segs.len() - 3]; - let owner = segs[segs.len() - 4]; - if !owner.is_empty() && !repo.is_empty() && !number.is_empty() { - return Some(format!("{owner}/{repo}#{number}")); - } - } - None -} - -/// Extract the `updated_at` ISO 8601 timestamp from a GitHub issue. -/// -/// GitHub returns `updated_at` as `"2024-05-21T15:30:00Z"`. ISO 8601 strings -/// sort lexicographically, so we use them directly as the sync cursor. -pub fn extract_issue_updated_at(issue: &Value) -> Option { - pick_str( - issue, - &[ - "updated_at", - "data.updated_at", - "updatedAt", - "data.updatedAt", - ], - ) -} - -#[cfg(test)] -#[path = "mod_tests.rs"] -mod tests; diff --git a/crates/tinymemory-integrations/src/sources/composio/github/mod_tests.rs b/crates/tinymemory-integrations/src/sources/composio/github/mod_tests.rs deleted file mode 100644 index 0b7322f1..00000000 --- a/crates/tinymemory-integrations/src/sources/composio/github/mod_tests.rs +++ /dev/null @@ -1,97 +0,0 @@ -//! Tests for the GitHub normaliser. - -use super::*; -use serde_json::json; - -#[test] -fn extract_issues_from_data_items() { - let data = json!({ "data": { "items": [{"id": 1}] } }); - assert_eq!(extract_issues(&data).len(), 1); -} - -#[test] -fn extract_issues_from_top_level_items() { - let data = json!({ "items": [{"id": 1}, {"id": 2}] }); - assert_eq!(extract_issues(&data).len(), 2); -} - -#[test] -fn extract_issues_empty_when_missing() { - let data = json!({ "foo": "bar" }); - assert!(extract_issues(&data).is_empty()); -} - -#[test] -fn extract_issue_id_from_numeric_field() { - let issue = json!({ "id": 123456789u64, "title": "Fix bug" }); - assert_eq!(extract_issue_id(&issue), Some("123456789".to_string())); -} - -#[test] -fn extract_issue_id_from_wrapped_data() { - let issue = json!({ "data": { "id": 99u64 } }); - assert_eq!(extract_issue_id(&issue), Some("99".to_string())); -} - -#[test] -fn extract_issue_id_falls_back_to_html_url() { - let issue = json!({ - "html_url": "https://github.com/owner/repo/issues/42" - }); - assert_eq!(extract_issue_id(&issue), Some("owner/repo#42".to_string())); -} - -#[test] -fn extract_issue_id_none_when_missing() { - let issue = json!({ "title": "No ID here" }); - assert!(extract_issue_id(&issue).is_none()); -} - -#[test] -fn extract_issue_title_builds_prefixed_title() { - let issue = json!({ - "id": 1u64, - "title": "Fix race condition", - "html_url": "https://github.com/acme/core/issues/99" - }); - assert_eq!( - extract_issue_title(&issue), - Some("GitHub: acme/core#99: Fix race condition".to_string()) - ); -} - -#[test] -fn extract_issue_title_returns_raw_title_when_no_url() { - let issue = json!({ "title": "Bare title" }); - assert_eq!(extract_issue_title(&issue), Some("Bare title".to_string())); -} - -#[test] -fn extract_issue_title_none_when_missing() { - let issue = json!({ "id": 1u64 }); - assert!(extract_issue_title(&issue).is_none()); -} - -#[test] -fn extract_issue_updated_at_from_top_level() { - let issue = json!({ "updated_at": "2024-05-21T15:30:00Z" }); - assert_eq!( - extract_issue_updated_at(&issue), - Some("2024-05-21T15:30:00Z".to_string()) - ); -} - -#[test] -fn extract_issue_updated_at_from_data_wrapper() { - let issue = json!({ "data": { "updated_at": "2023-01-01T00:00:00Z" } }); - assert_eq!( - extract_issue_updated_at(&issue), - Some("2023-01-01T00:00:00Z".to_string()) - ); -} - -#[test] -fn extract_issue_updated_at_none_when_missing() { - let issue = json!({ "id": 1u64 }); - assert!(extract_issue_updated_at(&issue).is_none()); -} diff --git a/crates/tinymemory-integrations/src/sources/composio/gmail_post_process/mod.rs b/crates/tinymemory-integrations/src/sources/composio/gmail_post_process/mod.rs deleted file mode 100644 index c6c78ae7..00000000 --- a/crates/tinymemory-integrations/src/sources/composio/gmail_post_process/mod.rs +++ /dev/null @@ -1,498 +0,0 @@ -//! Gmail-specific post-processing of Composio action responses. -//! -//! The upstream `GMAIL_FETCH_EMAILS` payload is extremely verbose -//! (full MIME tree under `payload.parts[]`, 50+ `Received:` headers, -//! display-layer noise the model never uses). This module rewrites -//! it into a slim envelope per message: -//! -//! ```json -//! { -//! "messages": [ -//! { -//! "id": "…", -//! "threadId": "…", -//! "subject": "…", -//! "from": "…", -//! "to": "…", -//! "date": "…", -//! "labels": ["INBOX", "UNREAD"], -//! "markdown": "…body…", -//! "attachments": [ { "filename": "...", "mimeType": "..." } ] -//! } -//! ], -//! "nextPageToken": "…", -//! "resultSizeEstimate": 201 -//! } -//! ``` -//! -//! ## Body source -//! -//! Composio's backend ships a -//! `markdownFormatted` field on the response envelope — one string -//! per tool call, pre-rendered with HTML stripped, URLs shortened, -//! footers removed, whitespace normalised. We split it per message -//! along `\n---\n` boundaries (with `## ` heading fallbacks) and -//! pin each slice to the corresponding entry in `messages[]` via -//! [`apply_response_level_markdown`]. The reshape's -//! `extract_markdown_body` then prefers that pinned field over -//! falling back to the upstream `messageText`. -//! -//! No in-house HTML→markdown conversion lives here anymore — the -//! backend does the cleaning. If `markdownFormatted` is absent for -//! a given response we fall through to whatever plain text the -//! upstream provided in `messageText`. -//! -//! Callers that need the raw Composio shape can pass `raw_html: -//! true` (or `rawHtml: true`) in the action arguments — this -//! short-circuits the reshape entirely. -//! -//! Only `GMAIL_FETCH_EMAILS` is reshaped today; other Gmail action -//! responses are passed through unchanged. When we add envelopes for -//! more slugs they should live in this file, branched from -//! [`post_process`]. - -use serde_json::{Map, Value, json}; - -/// Entry point a host calls on each Gmail action response (the slug names -/// the action) before handing it to `normalise_payload`. -/// -/// Dispatches on the Composio action slug. Unknown Gmail slugs fall -/// through to a no-op. -pub fn post_process(slug: &str, arguments: Option<&Value>, data: &mut Value) { - if is_raw_html_flag_set(arguments) { - tracing::debug!( - slug, - "[composio:gmail][post-process] raw_html flag set, passing through" - ); - return; - } - if slug == "GMAIL_FETCH_EMAILS" { - reshape_fetch_emails(data) - } -} - -/// Stash per-message slices of the response-level `markdownFormatted` -/// onto the corresponding entries inside `data.messages[]`. -/// -/// The Composio backend (tinyhumansai/backend#683) ships ONE -/// `markdownFormatted` string per tool call covering all messages — -/// already URL-shortened, footer-stripped, and whitespace-normalised. -/// To get per-email files in the raw archive we split that string -/// along section boundaries (`## ` headings or `---` rules) and pin -/// each slice to the message at the same index. `extract_markdown_body` -/// then prefers `msg.markdownFormatted` over re-decoding the MIME -/// tree. -/// -/// **Must be called BEFORE [`post_process`]** because `post_process` -/// reshapes `data` into the slim envelope; once `messages[]` carries -/// our slim shape the upstream message ordering is already locked in -/// but we may have lost original ordering signals if any. -/// -/// No-op when the slice count doesn't match `messages.len()` — we -/// can't safely align segments to messages without an exact match, -/// so we let `extract_markdown_body` fall through to its MIME path. -pub fn apply_response_level_markdown(data: &mut Value, top_md: &str) { - let trimmed = top_md.trim(); - if trimmed.is_empty() { - return; - } - // Presence is checked immutably first, then fetched mutably. The original - // form re-fetched with `unwrap()` after a mutable probe, which is sound but - // relies on the reader to see why; this crate lints against `unwrap`, and the - // immutable probe expresses the same reasoning to the compiler. - let container = if data.get("messages").is_some() { - data - } else if data.get("data").and_then(Value::as_object).is_some() { - match data.get_mut("data") { - Some(inner) => inner, - None => return, - } - } else { - tracing::debug!( - "[composio:gmail][post-process] apply_response_level_markdown: \ - no messages container in response — skipping" - ); - return; - }; - let Some(messages) = container.get_mut("messages").and_then(|v| v.as_array_mut()) else { - return; - }; - let count = messages.len(); - if count == 0 { - return; - } - // Clone hints out of the messages array so the slice borrows - // don't conflict with the upcoming `messages.iter_mut()` mutation. - let hints: Vec = messages.clone(); - let Some(slices) = split_response_markdown_per_message_with_hint(trimmed, count, Some(&hints)) - else { - tracing::debug!( - messages = count, - md_len = trimmed.len(), - "[composio:gmail][post-process] could not split response-level markdownFormatted \ - into {count} slices — falling back to per-message MIME decode" - ); - return; - }; - for (msg, slice) in messages.iter_mut().zip(slices) { - if let Some(obj) = msg.as_object_mut() { - obj.insert("markdownFormatted".to_string(), Value::String(slice)); - } - } - tracing::debug!( - messages = count, - "[composio:gmail][post-process] stashed per-message markdownFormatted slices" - ); -} - -/// Split a top-level `markdownFormatted` string into per-message -/// segments. Returns `Some(slices)` only when the split yields -/// exactly `expected_count` entries — otherwise the format isn't one -/// of the patterns we know about and we let the caller fall back. -/// -/// Primary boundary is the `\n---\n` horizontal rule the backend -/// emits between messages (confirmed against real -/// `GMAIL_FETCH_EMAILS` output). H2/H3 headings are kept as -/// fallbacks for older renderings. The preamble (`# Inbox (N -/// messages)`-style intro, if present) is dropped — we accept -/// either `expected` parts (no preamble) or `expected + 1` -/// (preamble + N messages). -/// -/// `messages_hint` is the slim message array from the same response -/// — when present we use the per-message `subject` field to verify -/// each segment really does belong to the message at the same index. -/// Mismatches force a fallback so we never write a wrong-message body -/// to the raw archive. -/// -/// The hint is what makes the split reliable: the blob's own section headings -/// are backend-rendered and have changed shape between versions, so matching -/// on them alone silently mis-attributed bodies. -pub fn split_response_markdown_per_message_with_hint( - md: &str, - expected_count: usize, - messages_hint: Option<&[Value]>, -) -> Option> { - if expected_count == 0 { - return None; - } - if expected_count == 1 { - return Some(vec![md.to_string()]); - } - - // Boundary patterns to try, in priority order. `\n---\n` is the - // confirmed marker; the heading variants stay as belt-and-braces - // for older / variant backend renderings. - let candidates: &[(&str, &str)] = &[ - ("\n---\n", "---\n"), - ("\n\n## ", "## "), - ("\n\n### ", "### "), - ("\n\n# ", "# "), - ("\n***\n", "***\n"), - ]; - - for (sep, prefix) in candidates { - let parts: Vec<&str> = md.split(sep).collect(); - let (drop_preamble, prepend_first) = if parts.len() == expected_count { - (false, false) // no preamble; first segment had no prefix - } else if parts.len() == expected_count + 1 { - (true, true) // preamble dropped; every kept segment had a prefix - } else { - continue; - }; - let segments: Vec = parts - .into_iter() - .skip(if drop_preamble { 1 } else { 0 }) - .enumerate() - .map(|(i, s)| { - if i == 0 && !prepend_first { - s.to_string() - } else { - format!("{prefix}{s}") - } - }) - .collect(); - - // Validate alignment against the JSON message array: every - // segment whose corresponding message has a non-empty subject - // must mention that subject somewhere in its body. If a single - // pair fails, we treat the split as unreliable and try the - // next pattern. Empty / null subjects skip validation (e.g. - // notification mails where the subject is ""). - if let Some(hints) = messages_hint - && !validate_segments_against_hints(&segments, hints) - { - tracing::debug!( - expected = expected_count, - sep = sep, - "[composio:gmail][post-process] split candidate failed subject check" - ); - continue; - } - return Some(segments); - } - None -} - -/// True if every (segment, message) pair where the message has a -/// non-empty subject contains that subject somewhere in the segment -/// (case-insensitive substring match — a defensive heuristic, not a -/// strict equality check, since the backend may format subjects -/// inside markdown links or with surrounding decoration). -fn validate_segments_against_hints(segments: &[String], hints: &[Value]) -> bool { - if segments.len() != hints.len() { - return false; - } - for (seg, hint) in segments.iter().zip(hints.iter()) { - let subject = hint - .get("subject") - .and_then(|v| v.as_str()) - .unwrap_or("") - .trim(); - if subject.is_empty() { - continue; - } - if !seg - .to_ascii_lowercase() - .contains(&subject.to_ascii_lowercase()) - { - return false; - } - } - true -} - -/// Returns true when the caller explicitly set `raw_html: true` (or the -/// camelCase `rawHtml: true`) in the `arguments` object. -fn is_raw_html_flag_set(arguments: Option<&Value>) -> bool { - let Some(obj) = arguments.and_then(|v| v.as_object()) else { - return false; - }; - obj.get("raw_html") - .or_else(|| obj.get("rawHtml")) - .and_then(|v| v.as_bool()) - .unwrap_or(false) -} - -/// Rewrite a `GMAIL_FETCH_EMAILS` `data` object in place into the slim -/// envelope documented at the module level. -/// -/// The Composio response can be shaped either as `{ messages, nextPageToken, ... }` -/// directly, or wrapped one level deeper under `{ data: { messages: … } }` -/// depending on backend version; we handle both. -fn reshape_fetch_emails(data: &mut Value) { - // Unwrap an optional `data:` envelope so downstream logic only has - // to deal with one shape. - let container = if data.get("messages").is_some() { - data - } else if data.get("data").and_then(Value::as_object).is_some() { - match data.get_mut("data") { - Some(inner) => inner, - None => return, - } - } else { - return; - }; - - let Some(obj) = container.as_object_mut() else { - return; - }; - - let raw_messages = obj - .remove("messages") - .and_then(|v| match v { - Value::Array(arr) => Some(arr), - _ => None, - }) - .unwrap_or_default(); - let next_page_token = obj.remove("nextPageToken").unwrap_or(Value::Null); - let result_size_estimate = obj.remove("resultSizeEstimate").unwrap_or(Value::Null); - - let messages: Vec = raw_messages.into_iter().map(reshape_message).collect(); - - let mut envelope = Map::new(); - envelope.insert("messages".into(), Value::Array(messages)); - if !next_page_token.is_null() { - envelope.insert("nextPageToken".into(), next_page_token); - } - if !result_size_estimate.is_null() { - envelope.insert("resultSizeEstimate".into(), result_size_estimate); - } - - *container = Value::Object(envelope); -} - -/// Parse an RFC 3339 or RFC 2822 date string into a UTC `DateTime`. -pub fn parse_email_date(date_str: &str) -> Option> { - date_str - .parse::>() - .or_else(|_| { - chrono::DateTime::parse_from_rfc2822(date_str).map(|d| d.with_timezone(&chrono::Utc)) - }) - .ok() -} - -const EMAIL_LOCAL_TIME_FMT: &str = "%Y-%m-%d %I:%M %p %:z"; - -/// Format a UTC `DateTime` in the given timezone. Returns `None` when the -/// formatted result is identical to the UTC rendering (no-op for UTC hosts). -pub fn format_at_tz( - utc: chrono::DateTime, - tz: &Tz, -) -> Option -where - Tz::Offset: std::fmt::Display, -{ - let local_dt = utc.with_timezone(tz); - let formatted = local_dt.format(EMAIL_LOCAL_TIME_FMT).to_string(); - - let utc_formatted = utc.format(EMAIL_LOCAL_TIME_FMT).to_string(); - if formatted == utc_formatted { - return None; - } - Some(formatted) -} - -/// Convert a UTC email timestamp string to a human-readable local-time string. -/// -/// Accepts RFC 3339 (`"2026-05-31T10:33:00Z"`) or RFC 2822 -/// (`"Sat, 31 May 2026 10:33:00 +0000"`) input. Returns a formatted string -/// in the host's local timezone, e.g. `"2026-05-31 05:33 AM -05:00"`, -/// so the agent can present local times without UTC arithmetic. -/// -/// The raw `date` field is always preserved alongside this field so -/// internal sorting, deduplication, and debugging remain UTC-based. -/// -/// Returns `None` when the input cannot be parsed or the output format -/// would be identical to the UTC input (no-op for UTC hosts). -pub fn format_email_local_time(date_str: &str) -> Option { - let utc = parse_email_date(date_str)?; - format_at_tz(utc, &chrono::Local) -} - -/// Map one raw Composio message object to its slim counterpart. -/// -/// Body source picked by [`extract_markdown_body`]: -/// 1. The per-message `markdownFormatted` slice pinned by -/// [`apply_response_level_markdown`] (preferred — backend-rendered). -/// 2. The upstream `messageText` plaintext (fallback). -/// 3. Empty string. -fn reshape_message(raw: Value) -> Value { - let Value::Object(obj) = raw else { - return raw; - }; - - let id = obj.get("messageId").cloned().unwrap_or(Value::Null); - let thread_id = obj.get("threadId").cloned().unwrap_or(Value::Null); - let subject = obj.get("subject").cloned().unwrap_or(Value::Null); - let sender = obj.get("sender").cloned().unwrap_or(Value::Null); - let to = obj.get("to").cloned().unwrap_or(Value::Null); - let date = obj - .get("messageTimestamp") - .cloned() - .or_else(|| pick_header(&obj, "Date")) - .unwrap_or(Value::Null); - let labels = obj - .get("labelIds") - .cloned() - .unwrap_or_else(|| Value::Array(Vec::new())); - let list_unsubscribe = pick_header(&obj, "List-Unsubscribe").unwrap_or(Value::Null); - - let markdown = extract_markdown_body(&obj); - let attachments = extract_attachments(&obj); - - // Compute a local-time representation of the UTC `date` so the agent - // presents times in the user's timezone rather than quoting raw UTC. - let date_local = date.as_str().and_then(format_email_local_time); - - let mut out = Map::new(); - out.insert("id".into(), id); - out.insert("threadId".into(), thread_id); - out.insert("subject".into(), subject); - out.insert("from".into(), sender); - out.insert("to".into(), to); - out.insert("date".into(), date); - if let Some(local) = date_local { - out.insert("date_local".into(), Value::String(local)); - } - out.insert("labels".into(), labels); - if !list_unsubscribe.is_null() { - out.insert("list_unsubscribe".into(), list_unsubscribe); - } - out.insert("markdown".into(), Value::String(markdown)); - if !attachments.is_empty() { - out.insert("attachments".into(), Value::Array(attachments)); - } - Value::Object(out) -} - -/// Find a header value by (case-insensitive) name in the Composio -/// `payload.headers[]` array. Returns `Some(Value::String)` on hit. -fn pick_header(msg: &Map, name: &str) -> Option { - let headers = msg.get("payload")?.get("headers")?.as_array()?; - for h in headers { - let hn = h.get("name").and_then(|v| v.as_str()).unwrap_or(""); - if hn.eq_ignore_ascii_case(name) - && let Some(v) = h.get("value").and_then(|v| v.as_str()) - { - return Some(Value::String(v.to_string())); - } - } - None -} - -/// Pick a body for the slim envelope. -/// -/// We trust the Composio backend's pre-rendered `markdownFormatted` -/// (set per-message by [`apply_response_level_markdown`] from the -/// response-level field). When that's absent we fall back to the -/// upstream's plain-text `messageText` verbatim — no in-house -/// HTML→markdown decoding lives here anymore. The backend already -/// strips HTML, shortens URLs, and normalises whitespace; running -/// our own pipeline on top duplicated work and corrupted some -/// renderings. -fn extract_markdown_body(msg: &Map) -> String { - if let Some(formatted) = msg - .get("markdownFormatted") - .or_else(|| msg.get("markdown_formatted")) - .and_then(|v| v.as_str()) - .map(str::trim) - .filter(|s| !s.is_empty()) - { - return formatted.to_string(); - } - if let Some(text) = msg - .get("messageText") - .and_then(|v| v.as_str()) - .map(str::trim) - .filter(|s| !s.is_empty()) - { - return text.to_string(); - } - String::new() -} - -/// Pull a minimal attachments descriptor from the Composio -/// `attachmentList` array. -fn extract_attachments(msg: &Map) -> Vec { - if let Some(list) = msg.get("attachmentList").and_then(|v| v.as_array()) { - return list - .iter() - .filter_map(|a| { - let filename = a.get("filename").and_then(|v| v.as_str())?; - if filename.is_empty() { - return None; - } - let mime = a - .get("mimeType") - .and_then(|v| v.as_str()) - .unwrap_or_default(); - Some(json!({ "filename": filename, "mimeType": mime })) - }) - .collect(); - } - Vec::new() -} - -#[cfg(test)] -#[path = "mod_tests.rs"] -mod tests; diff --git a/crates/tinymemory-integrations/src/sources/composio/gmail_post_process/mod_tests.rs b/crates/tinymemory-integrations/src/sources/composio/gmail_post_process/mod_tests.rs deleted file mode 100644 index dd1b44af..00000000 --- a/crates/tinymemory-integrations/src/sources/composio/gmail_post_process/mod_tests.rs +++ /dev/null @@ -1,356 +0,0 @@ -//! Tests for the Gmail post-processor. - -use super::*; -use serde_json::json; - -fn fixture_with_backend_markdown() -> Value { - json!({ - "messages": [ - { - "messageId": "m1", - "threadId": "t1", - "subject": "Hello", - "sender": "a@x.com", - "to": "b@y.com", - "messageTimestamp": "2026-04-17T12:00:00Z", - "labelIds": ["INBOX", "UNREAD"], - // Pre-rendered slice (set by `apply_response_level_markdown` - // in production; inline here for the reshape test). - "markdownFormatted": "# Hello\n\nbody copy", - "messageText": "fallback should not be used", - "display_url": "ignore-me", - "preview": { "body": "Hi plain", "subject": "Hello" }, - "attachmentList": [ - { "filename": "report.pdf", "mimeType": "application/pdf", "size": 12345 }, - { "filename": "", "mimeType": "text/html" } - ], - "payload": {} - } - ], - "nextPageToken": "tok-1", - "resultSizeEstimate": 42 - }) -} - -#[test] -fn reshape_emits_slim_envelope() { - let mut v = fixture_with_backend_markdown(); - post_process("GMAIL_FETCH_EMAILS", None, &mut v); - - assert_eq!(v["nextPageToken"], "tok-1"); - assert_eq!(v["resultSizeEstimate"], 42); - - let msgs = v["messages"].as_array().unwrap(); - assert_eq!(msgs.len(), 1); - let m = &msgs[0]; - - assert_eq!(m["id"], "m1"); - assert_eq!(m["threadId"], "t1"); - assert_eq!(m["subject"], "Hello"); - assert_eq!(m["from"], "a@x.com"); - assert_eq!(m["to"], "b@y.com"); - assert_eq!(m["date"], "2026-04-17T12:00:00Z"); - assert_eq!(m["labels"], json!(["INBOX", "UNREAD"])); - - let md = m["markdown"].as_str().unwrap(); - assert_eq!(md, "# Hello\n\nbody copy"); - - // Noise fields removed. - assert!(m.get("display_url").is_none()); - assert!(m.get("preview").is_none()); - assert!(m.get("payload").is_none()); - assert!(m.get("messageText").is_none()); - - // Attachments: empty filename entry is filtered. - let atts = m["attachments"].as_array().unwrap(); - assert_eq!(atts.len(), 1); - assert_eq!(atts[0]["filename"], "report.pdf"); - assert_eq!(atts[0]["mimeType"], "application/pdf"); -} - -#[test] -fn raw_html_flag_passes_through_unchanged() { - let mut v = fixture_with_backend_markdown(); - let original = v.clone(); - let args = json!({ "raw_html": true }); - post_process("GMAIL_FETCH_EMAILS", Some(&args), &mut v); - assert_eq!( - v, original, - "raw_html=true must preserve the Composio shape" - ); -} - -#[test] -fn camel_case_raw_html_also_recognized() { - let mut v = fixture_with_backend_markdown(); - let original = v.clone(); - let args = json!({ "rawHtml": true }); - post_process("GMAIL_FETCH_EMAILS", Some(&args), &mut v); - assert_eq!(v, original); -} - -#[test] -fn falls_back_to_message_text_when_no_backend_markdown() { - let mut v = json!({ - "messages": [{ - "messageId": "m1", - "threadId": "t1", - "subject": "s", - "sender": "a@x.com", - "to": "b@y.com", - "messageTimestamp": "2026-04-17", - "labelIds": [], - "messageText": " plain body text ", - "payload": {} - }], - "nextPageToken": null - }); - post_process("GMAIL_FETCH_EMAILS", None, &mut v); - let md = v["messages"][0]["markdown"].as_str().unwrap(); - assert_eq!(md, "plain body text"); - assert!(v.get("nextPageToken").is_none(), "null tokens dropped"); -} - -#[test] -fn unwraps_data_envelope() { - let mut v = json!({ - "data": { - "messages": [{ - "messageId": "m1", - "threadId": "t1", - "subject": "s", - "sender": "a@x.com", - "to": "b@y.com", - "messageTimestamp": "2026-04-17", - "labelIds": [], - "messageText": "body", - "payload": {} - }] - } - }); - post_process("GMAIL_FETCH_EMAILS", None, &mut v); - // Reshape writes into `data` in place. - let msgs = v["data"]["messages"].as_array().unwrap(); - assert_eq!(msgs.len(), 1); - assert_eq!(msgs[0]["markdown"], "body"); -} - -#[test] -fn non_fetch_slug_is_noop() { - let mut v = json!({ "messages": [{ "messageId": "m1", "messageText": "x" }] }); - let original = v.clone(); - post_process("GMAIL_SEND_EMAIL", None, &mut v); - assert_eq!(v, original); -} - -#[test] -fn prefers_backend_markdown_formatted_when_present() { - // Composio backend (tinyhumansai/backend#683 +) ships - // `markdownFormatted` already URL-shortened + footer-stripped - // per message (after `apply_response_level_markdown` slices the - // response-level field). When present, our post-processor must - // use it verbatim instead of falling back to `messageText`. - let mut v = json!({ - "messages": [{ - "messageId": "m1", - "threadId": "t1", - "subject": "s", - "sender": "a@x.com", - "to": "b@y.com", - "messageTimestamp": "2026-04-17", - "labelIds": [], - "markdownFormatted": "# Already nice\n\nShort URL: https://gh.io/abc", - "messageText": "fallback should not be used", - "payload": {} - }] - }); - post_process("GMAIL_FETCH_EMAILS", None, &mut v); - let md = v["messages"][0]["markdown"].as_str().unwrap(); - assert_eq!(md, "# Already nice\n\nShort URL: https://gh.io/abc"); -} - -#[test] -fn empty_markdown_formatted_falls_through_to_message_text() { - let mut v = json!({ - "messages": [{ - "messageId": "m1", - "threadId": "t1", - "subject": "s", - "sender": "a@x.com", - "to": "b@y.com", - "messageTimestamp": "2026-04-17", - "labelIds": [], - "markdownFormatted": " \n \n", - "messageText": "real body", - "payload": {} - }] - }); - post_process("GMAIL_FETCH_EMAILS", None, &mut v); - let md = v["messages"][0]["markdown"].as_str().unwrap(); - assert!(md.contains("real body")); -} - -// ── split_response_markdown_per_message_with_hint ─────────────────────── - -#[test] -fn split_response_markdown_uses_horizontal_rule_marker() { - // The confirmed backend marker is `\n---\n`. Three messages → - // expect three slices when there's no preamble. - let md = "## Alice's update\n\nbody A with https://gh.io/abc\n---\n## Bob's reply\n\nbody B\n---\n## Carol\n\nbody C"; - let slices = super::split_response_markdown_per_message_with_hint(md, 3, None).unwrap(); - assert_eq!(slices.len(), 3); - assert!(slices[0].contains("Alice's update")); - assert!(slices[1].contains("Bob's reply")); - assert!(slices[2].contains("Carol")); - // The `---\n` prefix is preserved on every-but-the-first segment - // so the section break survives the round-trip. - assert!(slices[1].starts_with("---\n")); - assert!(slices[2].starts_with("---\n")); -} - -#[test] -fn split_response_markdown_drops_preamble() { - // When a preamble like `# Inbox` precedes the first marker, we - // see N+1 parts after split — the preamble must be dropped. - let md = "# Inbox (2 messages)\n---\n## A\n\nbody A\n---\n## B\n\nbody B"; - let slices = super::split_response_markdown_per_message_with_hint(md, 2, None).unwrap(); - assert_eq!(slices.len(), 2); - assert!(slices[0].contains("body A")); - assert!(slices[1].contains("body B")); - // Both segments should carry the prefix when preamble was dropped. - assert!(slices[0].starts_with("---\n")); - assert!(slices[1].starts_with("---\n")); -} - -#[test] -fn split_response_markdown_falls_back_to_h2_marker() { - // No `---` rules — backend used h2 headings as boundaries. - let md = "## Alice\n\nbody A\n\n## Bob\n\nbody B"; - let slices = super::split_response_markdown_per_message_with_hint(md, 2, None).unwrap(); - assert_eq!(slices.len(), 2); - assert!(slices[0].contains("body A")); - assert!(slices[1].contains("body B")); -} - -#[test] -fn split_response_markdown_returns_none_on_count_mismatch() { - let md = "## only one section here"; - assert!(super::split_response_markdown_per_message_with_hint(md, 3, None).is_none()); -} - -#[test] -fn split_response_markdown_single_message_returns_whole_input() { - let md = "## solo\n\nthe whole body"; - let slices = super::split_response_markdown_per_message_with_hint(md, 1, None).unwrap(); - assert_eq!(slices, vec![md.to_string()]); -} - -#[test] -fn split_with_hint_rejects_when_subjects_dont_match() { - let md = "## Foo\nbody1\n---\n## Bar\nbody2"; - let hints = vec![ - json!({"subject": "Completely different subject A"}), - json!({"subject": "Completely different subject B"}), - ]; - let out = super::split_response_markdown_per_message_with_hint(md, 2, Some(&hints)); - assert!(out.is_none(), "subject mismatch must force fallback"); -} - -#[test] -fn split_with_hint_accepts_when_subjects_match() { - let md = "## Welcome to Gmail\nbody1\n---\n## Your invoice\nbody2"; - let hints = vec![ - json!({"subject": "Welcome to Gmail"}), - json!({"subject": "Your invoice"}), - ]; - let slices = super::split_response_markdown_per_message_with_hint(md, 2, Some(&hints)).unwrap(); - assert_eq!(slices.len(), 2); - assert!(slices[0].contains("Welcome to Gmail")); - assert!(slices[1].contains("Your invoice")); -} - -#[test] -fn split_with_hint_skips_messages_with_blank_subject() { - let md = "## A\nbody1\n---\n## B\nbody2"; - let hints = vec![json!({"subject": "A"}), json!({"subject": ""})]; - let slices = super::split_response_markdown_per_message_with_hint(md, 2, Some(&hints)).unwrap(); - assert_eq!(slices.len(), 2); -} - -// ── format_email_local_time ────────────────────────────────────────────────── - -#[test] -fn format_email_local_time_returns_none_for_unparseable_date() { - assert!(super::format_email_local_time("not-a-date").is_none()); - assert!(super::format_email_local_time("").is_none()); -} - -#[test] -fn format_email_local_time_preserves_utc_raw_date_in_reshape() { - let mut v = json!({ - "messages": [{ - "messageId": "m1", - "threadId": "t1", - "subject": "Test", - "sender": "a@example.com", - "to": "b@example.com", - "messageTimestamp": "2026-05-31T10:33:00Z", - "labelIds": [], - "messageText": "body", - "payload": {} - }] - }); - post_process("GMAIL_FETCH_EMAILS", None, &mut v); - let msg = &v["messages"][0]; - assert_eq!(msg["date"], "2026-05-31T10:33:00Z"); -} - -#[test] -fn parse_email_date_accepts_rfc3339_and_rfc2822() { - assert!(super::parse_email_date("2026-05-31T10:33:00Z").is_some()); - assert!(super::parse_email_date("Sun, 31 May 2026 10:33:00 +0000").is_some()); - assert!(super::parse_email_date("not-a-date").is_none()); -} - -#[test] -fn format_at_tz_deterministic_with_fixed_offset() { - use chrono::FixedOffset; - - let utc = super::parse_email_date("2026-05-31T10:33:00Z").unwrap(); - - let est = FixedOffset::west_opt(5 * 3600).unwrap(); - let result = super::format_at_tz(utc, &est).unwrap(); - assert_eq!(result, "2026-05-31 05:33 AM -05:00"); - - let ist = FixedOffset::east_opt(5 * 3600 + 1800).unwrap(); - let result = super::format_at_tz(utc, &ist).unwrap(); - assert_eq!(result, "2026-05-31 04:03 PM +05:30"); -} - -#[test] -fn format_at_tz_returns_none_for_utc() { - let utc = super::parse_email_date("2026-05-31T10:33:00Z").unwrap(); - let utc_tz = chrono::FixedOffset::east_opt(0).unwrap(); - assert!(super::format_at_tz(utc, &utc_tz).is_none()); -} - -#[test] -fn apply_response_level_markdown_stashes_per_message_field() { - let mut data = json!({ - "messages": [ - {"messageId": "m1", "subject": "Hello"}, - {"messageId": "m2", "subject": "World"}, - ] - }); - let top_md = "## Hello\nbody A — link https://gh.io/abc\n---\n## World\nbody B"; - super::apply_response_level_markdown(&mut data, top_md); - let m1 = data["messages"][0]["markdownFormatted"].as_str().unwrap(); - let m2 = data["messages"][1]["markdownFormatted"].as_str().unwrap(); - assert!(m1.contains("Hello")); - assert!( - m1.contains("https://gh.io/abc"), - "shortened URL must survive" - ); - assert!(m2.contains("World")); - assert!(!m1.contains("World"), "no cross-message bleed"); -} diff --git a/crates/tinymemory-integrations/src/sources/composio/linear/mod.rs b/crates/tinymemory-integrations/src/sources/composio/linear/mod.rs deleted file mode 100644 index ea52092a..00000000 --- a/crates/tinymemory-integrations/src/sources/composio/linear/mod.rs +++ /dev/null @@ -1,79 +0,0 @@ -//! Linear host normalization helpers — result extraction and issue-title and -//! timestamp extraction. -//! -//! Linear's GraphQL API (and therefore Composio's wrapping of it) returns -//! connection-style lists (`{ nodes: [...], pageInfo: {...} }`) at the top -//! level or nested under `data`. The functions here walk the union of -//! common shapes so the provider does not have to branch per Composio -//! envelope variant. - -use serde_json::Value; - -use super::fields::pick_str; - -/// Walk the Composio response envelope for Linear issue list results. -/// -/// Linear's list endpoints return `{ nodes: [...] }` or -/// `{ issues: { nodes: [...] } }` shapes; Composio may re-wrap the -/// upstream payload under `data` or `data.data`. We probe each shape -/// in order and return the first array we find. -pub fn extract_issues(data: &Value) -> Vec { - let candidates = [ - data.pointer("/data/nodes"), - data.pointer("/nodes"), - data.pointer("/data/issues/nodes"), - data.pointer("/issues/nodes"), - data.pointer("/data/data/nodes"), - data.pointer("/data/data/issues/nodes"), - data.pointer("/data/results"), - data.pointer("/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 a human-readable title from a Linear issue object. -/// -/// Linear issues store the name at `title` (or `data.title` after -/// Composio envelope wrapping). Falls back to `name` / `identifier` -/// so the chunk remains identifiable even for unusual response shapes. -pub fn extract_issue_title(issue: &Value) -> Option { - pick_str( - issue, - &[ - "title", - "data.title", - "name", - "data.name", - "identifier", - "data.identifier", - ], - ) -} - -/// Extract a stable cursor timestamp from a Linear issue object. -/// -/// Linear uses ISO-8601 strings for timestamps (`updatedAt`). We keep -/// the value as a string so lexicographic comparison against the stored -/// cursor is valid. -pub fn extract_issue_updated(issue: &Value) -> Option { - pick_str( - issue, - &[ - "updatedAt", - "data.updatedAt", - "updated_at", - "data.updated_at", - ], - ) -} - -#[cfg(test)] -#[path = "mod_tests.rs"] -mod tests; diff --git a/crates/tinymemory-integrations/src/sources/composio/linear/mod_tests.rs b/crates/tinymemory-integrations/src/sources/composio/linear/mod_tests.rs deleted file mode 100644 index c0acff63..00000000 --- a/crates/tinymemory-integrations/src/sources/composio/linear/mod_tests.rs +++ /dev/null @@ -1,93 +0,0 @@ -//! Tests for the Linear normaliser. - -use super::*; -use serde_json::json; - -// ── extract_issues ─────────────────────────────────────────────── - -#[test] -fn extract_issues_from_data_nodes() { - let data = json!({ "data": { "nodes": [{"id": "i1"}, {"id": "i2"}] } }); - assert_eq!(extract_issues(&data).len(), 2); -} - -#[test] -fn extract_issues_from_top_level_nodes() { - let data = json!({ "nodes": [{"id": "i3"}] }); - assert_eq!(extract_issues(&data).len(), 1); -} - -#[test] -fn extract_issues_from_data_issues_nodes() { - let data = - json!({ "data": { "issues": { "nodes": [{"id": "i4"}, {"id": "i5"}, {"id": "i6"}] } } }); - assert_eq!(extract_issues(&data).len(), 3); -} - -#[test] -fn extract_issues_from_top_level_issues_nodes() { - let data = json!({ "issues": { "nodes": [{"id": "i7"}] } }); - assert_eq!(extract_issues(&data).len(), 1); -} - -#[test] -fn extract_issues_from_doubly_nested_issues_nodes() { - let data = - json!({ "data": { "data": { "issues": { "nodes": [{"id": "i8"}, {"id": "i9"}] } } } }); - assert_eq!(extract_issues(&data).len(), 2); -} - -#[test] -fn extract_issues_from_results() { - let data = json!({ "results": [{"id": "i7"}] }); - assert_eq!(extract_issues(&data).len(), 1); -} - -#[test] -fn extract_issues_empty_when_missing() { - let data = json!({ "foo": "bar" }); - assert!(extract_issues(&data).is_empty()); -} - -// ── extract_issue_title ────────────────────────────────────────── - -#[test] -fn extract_issue_title_from_title_field() { - let issue = json!({ "id": "i1", "title": "Fix the login bug" }); - assert_eq!( - extract_issue_title(&issue), - Some("Fix the login bug".into()) - ); -} - -#[test] -fn extract_issue_title_falls_back_to_wrapped_data() { - let issue = json!({ "data": { "title": "Wrapped issue" } }); - assert_eq!(extract_issue_title(&issue), Some("Wrapped issue".into())); -} - -#[test] -fn extract_issue_title_falls_back_to_identifier() { - let issue = json!({ "identifier": "ENG-42" }); - assert_eq!(extract_issue_title(&issue), Some("ENG-42".into())); -} - -// ── extract_issue_updated ──────────────────────────────────────── - -#[test] -fn extract_issue_updated_from_updated_at() { - let issue = json!({ "updatedAt": "2026-03-01T12:00:00.000Z" }); - assert_eq!( - extract_issue_updated(&issue), - Some("2026-03-01T12:00:00.000Z".to_string()) - ); -} - -#[test] -fn extract_issue_updated_falls_back_to_snake_case() { - let issue = json!({ "data": { "updated_at": "2026-01-15T08:30:00.000Z" } }); - assert_eq!( - extract_issue_updated(&issue), - Some("2026-01-15T08:30:00.000Z".to_string()) - ); -} diff --git a/crates/tinymemory-integrations/src/sources/composio/mod.rs b/crates/tinymemory-integrations/src/sources/composio/mod.rs deleted file mode 100644 index 4df32030..00000000 --- a/crates/tinymemory-integrations/src/sources/composio/mod.rs +++ /dev/null @@ -1,35 +0,0 @@ -//! Composio toolkit payloads: normalisers and the mapping to `StoreItem`s. -//! -//! A host runs Composio actions with its own credentials and hands the raw -//! responses here. Two layers turn them into memory: -//! -//! 1. **Normalisers**, one module per toolkit, are pure -//! `serde_json::Value` transforms: they walk Composio's envelope variants -//! and pull out the tasks, issues, pages or messages -//! ([`clickup`], [`github`], [`linear`], [`notion`]), or rewrite a verbose -//! response into a slim one in place ([`gmail_post_process`], -//! [`slack_post_process`]). [`fields`] holds the path lookup they share. -//! 2. [`normalise_payload`] turns one (post-processed) response into -//! [`ComposioDocument`]s, and [`payload_items`] turns those into -//! [`StoreItem::Document`](tinymemory_api::StoreItem::Document)s with -//! `source.kind = Composio`, `source.id` = the connection or source id, -//! `url` and `observed_at` where the payload has them, and -//! `tags = [toolkit]`. -//! -//! Nothing here holds a credential, opens a socket or decides when to sync. -//! -//! One caveat on "pure": [`gmail_post_process::format_email_local_time`] -//! renders in `chrono::Local`, so it reads the host's timezone. The raw UTC -//! fields are preserved alongside, so ordering and identity stay UTC-based. - -pub mod clickup; -pub mod fields; -pub mod github; -pub mod gmail_post_process; -pub mod linear; -pub mod notion; -pub mod slack_post_process; - -mod documents; - -pub use documents::{ComposioDocument, normalise_payload, payload_items}; diff --git a/crates/tinymemory-integrations/src/sources/composio/notion/mod.rs b/crates/tinymemory-integrations/src/sources/composio/notion/mod.rs deleted file mode 100644 index 00c1df50..00000000 --- a/crates/tinymemory-integrations/src/sources/composio/notion/mod.rs +++ /dev/null @@ -1,88 +0,0 @@ -//! Notion host normalization helpers — result extraction, page markdown and -//! page title extraction. - -use serde_json::Value; - -use super::fields::pick_str; - -/// Walk the Composio response envelope for Notion page results. -pub 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 rendered page body markdown from a `NOTION_GET_PAGE_MARKDOWN` -/// response. Composio wraps action output in varying envelope shapes, so we -/// try the common locations tolerantly and return the first non-empty string. -/// Returns `None` if no markdown field is found (caller falls back to the -/// metadata-only body and logs the raw shape for diagnosis). -pub fn extract_page_markdown(data: &Value) -> Option { - const PATHS: &[&str] = &[ - "/markdown", - "/data/markdown", - "/data/response_data/markdown", - "/response_data/markdown", - "/data/content", - "/content", - "/data/markdown_content", - "/markdown_content", - "/text", - "/data/text", - ]; - for p in PATHS { - if let Some(s) = data.pointer(p).and_then(Value::as_str) - && !s.trim().is_empty() - { - return Some(s.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 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") - && 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"]) -} - -#[cfg(test)] -#[path = "mod_tests.rs"] -mod tests; diff --git a/crates/tinymemory-integrations/src/sources/composio/notion/mod_tests.rs b/crates/tinymemory-integrations/src/sources/composio/notion/mod_tests.rs deleted file mode 100644 index ab2e79bc..00000000 --- a/crates/tinymemory-integrations/src/sources/composio/notion/mod_tests.rs +++ /dev/null @@ -1,111 +0,0 @@ -//! Tests for the Notion normaliser. - -use super::*; -use serde_json::json; - -#[test] -fn extract_results_from_data_results() { - let data = json!({"data": {"results": [{"id": "page1"}]}}); - let results = extract_results(&data); - assert_eq!(results.len(), 1); -} - -#[test] -fn extract_page_markdown_reads_top_level_field() { - // Matches the live GET_PAGE_MARKDOWN envelope observed empirically: - // {id, markdown, object, request_id, truncated, unknown_block_ids}. - let data = json!({ - "id": "p1", - "markdown": "# Heading\n\nbody text", - "object": "page", - "truncated": false, - }); - assert_eq!( - extract_page_markdown(&data).as_deref(), - Some("# Heading\n\nbody text") - ); -} - -#[test] -fn extract_page_markdown_reads_nested_envelope() { - let data = json!({ "data": { "markdown": "nested body" } }); - assert_eq!(extract_page_markdown(&data).as_deref(), Some("nested body")); -} - -#[test] -fn extract_page_markdown_none_for_empty_or_missing() { - // Empty markdown (a DB row with no body blocks) → None → metadata-only. - assert_eq!(extract_page_markdown(&json!({ "markdown": "" })), None); - assert_eq!(extract_page_markdown(&json!({ "markdown": " " })), None); - // No markdown field at all → None. - assert_eq!(extract_page_markdown(&json!({ "id": "p1" })), None); -} - -#[test] -fn extract_results_from_top_level() { - let data = json!({"results": [{"id": "a"}, {"id": "b"}]}); - let results = extract_results(&data); - assert_eq!(results.len(), 2); -} - -#[test] -fn extract_results_from_data_items() { - let data = json!({"data": {"items": [{"id": "x"}]}}); - let results = extract_results(&data); - assert_eq!(results.len(), 1); -} - -#[test] -fn extract_results_empty_when_no_match() { - let data = json!({"foo": "bar"}); - assert!(extract_results(&data).is_empty()); -} - -#[test] -fn extract_page_title_from_properties_title_type() { - let page = json!({ - "properties": { - "Name": { - "type": "title", - "title": [{"plain_text": "Hello"}, {"plain_text": " World"}] - } - } - }); - assert_eq!(extract_page_title(&page), Some("Hello World".into())); -} - -#[test] -fn extract_page_title_from_nested_data_properties() { - let page = json!({ - "data": { - "properties": { - "Title": { - "type": "title", - "title": [{"plain_text": "My Page"}] - } - } - } - }); - assert_eq!(extract_page_title(&page), Some("My Page".into())); -} - -#[test] -fn extract_page_title_fallback_to_top_level_title() { - let page = json!({"title": "Fallback Title"}); - assert_eq!(extract_page_title(&page), Some("Fallback Title".into())); -} - -#[test] -fn extract_page_title_none_when_empty() { - let page = json!({"properties": {"Name": {"type": "title", "title": []}}}); - // Empty title array means no text - assert!( - extract_page_title(&page).is_none() || extract_page_title(&page) == Some(String::new()) - ); -} - -#[test] -fn extract_page_title_none_when_no_title_field() { - let page = json!({"id": "123"}); - assert!(extract_page_title(&page).is_none()); -} diff --git a/crates/tinymemory-integrations/src/sources/composio/slack_post_process/mod.rs b/crates/tinymemory-integrations/src/sources/composio/slack_post_process/mod.rs deleted file mode 100644 index 764b885a..00000000 --- a/crates/tinymemory-integrations/src/sources/composio/slack_post_process/mod.rs +++ /dev/null @@ -1,324 +0,0 @@ -//! Slack-specific post-processing of Composio action responses. -//! -//! Composio's Slack responses are verbose API envelopes. This module -//! rewrites each supported action's response into a slim, stable shape -//! that the ingest pipeline and enrichers can consume without walking -//! Composio's unstable nested envelopes. -//! -//! ## Supported slugs -//! -//! - `SLACK_FETCH_CONVERSATION_HISTORY` — reshapes into top-level -//! `messages[]` with `{ ts, user, text, thread_ts, channel_id }`. -//! Empty-text messages are dropped. `channel_id` is absent here (it's -//! in the request, not the response); the caller injects it via the -//! enricher in the host's `SlackSyncPipeline`. -//! -//! - `SLACK_LIST_CONVERSATIONS` — reshapes into top-level `channels[]` -//! with `{ id, name, is_private }` per channel. Entries with an empty -//! id are dropped. -//! -//! - `SLACK_SEARCH_MESSAGES` — reshapes `messages.matches[]` (possibly -//! nested) into top-level `messages[]` with `{ ts, user, text, -//! thread_ts, channel_id }`. `channel_id` is pulled from each match's -//! `channel.id` field. `paging.pages` is preserved at top-level for -//! caller pagination. -//! -//! ## Design note: user-id resolution is NOT here -//! -//! `SlackUsers` is a per-sync cache built from a separate API call — -//! not a function of any individual response. Resolving user ids -//! happens in the host's `SlackSyncPipeline` (the enricher layer), keeping -//! this module purely data-shape–oriented. -//! This matches Gmail's pattern of "post_process is data-only". -//! -//! Unknown slugs are silently no-ops so new Composio actions don't -//! break the provider. - -use serde_json::{Map, Value}; - -/// Entry point a host calls on each Slack action response (the slug names -/// the action) before handing it to `normalise_payload`. -/// -/// Dispatches on the Composio action slug and rewrites `data` in place. -/// Unknown slugs are silently ignored. -pub fn post_process(slug: &str, _arguments: Option<&Value>, data: &mut Value) { - log::debug!("[composio:slack][post-process] slug={slug}"); - match slug { - "SLACK_FETCH_CONVERSATION_HISTORY" => reshape_fetch_history(data), - "SLACK_LIST_CONVERSATIONS" => reshape_list_conversations(data), - "SLACK_SEARCH_MESSAGES" => reshape_search_messages(data), - _ => { - log::debug!("[composio:slack][post-process] unknown slug={slug}, passing through"); - } - } -} - -// ─── SLACK_FETCH_CONVERSATION_HISTORY ────────────────────────────────────── - -/// Rewrite a `SLACK_FETCH_CONVERSATION_HISTORY` response in place. -/// -/// Walks possible nested envelopes (`/data/messages`, `/messages`, -/// `/data/data/messages`) to find the raw messages array, drops messages -/// with empty `text`, and emits a slim `{ ts, user, text, thread_ts }` -/// shape under a top-level `messages[]` key. The consumed nested array is -/// removed from the payload so the raw verbose rows don't linger alongside -/// the slim copy. The caller injects `channel_id` via -/// the host's Slack sync pipeline. -fn reshape_fetch_history(data: &mut Value) { - let arr = take_array( - data, - &["/data/messages", "/messages", "/data/data/messages"], - 0, - ); - let slim: Vec = arr.into_iter().filter_map(slim_history_message).collect(); - with_object(data, |obj| { - obj.insert("messages".to_string(), Value::Array(slim)); - }); - log::debug!("[composio:slack][post-process] SLACK_FETCH_CONVERSATION_HISTORY reshaped"); -} - -fn slim_history_message(raw: Value) -> Option { - let text = raw - .get("text") - .and_then(|v| v.as_str()) - .unwrap_or("") - .trim(); - if text.is_empty() { - return None; - } - let mut out = Map::new(); - // `ts` is required: without it a caller can neither cursor nor archive. - out.insert("ts".into(), raw.get("ts")?.clone()); - if let Some(user) = raw.get("user").or_else(|| raw.get("bot_id")) { - out.insert("user".into(), user.clone()); - } - out.insert("text".into(), Value::String(text.to_string())); - if let Some(thread_ts) = raw.get("thread_ts") { - out.insert("thread_ts".into(), thread_ts.clone()); - } - if let Some(permalink) = raw.get("permalink") { - out.insert("permalink".into(), permalink.clone()); - } - Some(Value::Object(out)) -} - -/// Find the first array at any of `candidates`, remove that field (plus -/// `envelope_depth` ancestor object envelopes) from `data`, and return the -/// array. Removing the consumed nested payload keeps the reshaped output from -/// carrying duplicate raw rows. -fn take_array(data: &mut Value, candidates: &[&str], envelope_depth: usize) -> Vec { - for path in candidates { - let arr = match data.pointer(path).and_then(|v| v.as_array().cloned()) { - Some(a) => a, - None => continue, - }; - let mut remove_path = path.to_string(); - for _ in 0..envelope_depth { - remove_path = match remove_path.rsplit_once('/') { - Some((parent, _)) => parent.to_string(), - None => break, - }; - } - remove_nested(data, &remove_path); - return arr; - } - Vec::new() -} - -/// Remove the field at `path` from `data`, pruning any ancestor object that -/// the removal left empty so a consumed `data` envelope disappears entirely -/// instead of lingering as `{}`. -fn remove_nested(data: &mut Value, path: &str) { - let segments: Vec<&str> = path - .trim_start_matches('/') - .split('/') - .filter(|s| !s.is_empty()) - .collect(); - if segments.is_empty() { - return; - } - - // Remove the leaf field. - let mut current = &mut *data; - for seg in &segments[..segments.len() - 1] { - current = match current.get_mut(*seg) { - Some(next) => next, - None => return, - }; - } - if let Value::Object(map) = current { - map.remove(segments[segments.len() - 1]); - } - - // Prune empty object ancestors, deepest first. - for depth in (0..segments.len().saturating_sub(1)).rev() { - // Re-walk to the object at `segments[..=depth]`. - let mut ancestor = &mut *data; - for seg in &segments[..=depth] { - ancestor = match ancestor.get_mut(*seg) { - Some(next) => next, - None => return, - }; - } - if !matches!(ancestor, Value::Object(m) if m.is_empty()) { - break; - } - // Remove it from its parent (`segments[..depth]`). For `depth == 0` - // the parent is the top-level object, so an emptied `data` envelope - // key disappears entirely. - let mut parent = &mut *data; - for seg in &segments[..depth] { - parent = match parent.get_mut(*seg) { - Some(next) => next, - None => return, - }; - } - if let Value::Object(map) = parent { - map.remove(segments[depth]); - } - } -} - -// ─── SLACK_LIST_CONVERSATIONS ─────────────────────────────────────────────── - -/// Rewrite a `SLACK_LIST_CONVERSATIONS` response in place. -/// -/// Reshapes into a top-level `channels[]` with `{ id, name, is_private }` -/// per channel; entries with an empty id are dropped. -fn reshape_list_conversations(data: &mut Value) { - let arr = take_array( - data, - &[ - "/data/channels", - "/channels", - "/data/data/channels", - "/data/conversations", - "/conversations", - ], - 0, - ); - - let slim: Vec = arr.into_iter().filter_map(slim_channel).collect(); - with_object(data, |obj| { - obj.insert("channels".to_string(), Value::Array(slim)); - }); - log::debug!("[composio:slack][post-process] SLACK_LIST_CONVERSATIONS reshaped"); -} - -fn slim_channel(raw: Value) -> Option { - let id = raw.get("id").and_then(|v| v.as_str()).unwrap_or("").trim(); - if id.is_empty() { - return None; - } - let name = raw - .get("name") - .and_then(|v| v.as_str()) - .unwrap_or(id) - .trim(); - let is_private = raw - .get("is_private") - .and_then(|v| v.as_bool()) - .unwrap_or(false); - Some(Value::Object({ - let mut m = Map::new(); - m.insert("id".into(), Value::String(id.to_string())); - m.insert("name".into(), Value::String(name.to_string())); - m.insert("is_private".into(), Value::Bool(is_private)); - m - })) -} - -// ─── SLACK_SEARCH_MESSAGES ────────────────────────────────────────────────── - -/// Rewrite a `SLACK_SEARCH_MESSAGES` response in place. -/// -/// Reshapes `messages.matches[]` (possibly nested under one or two -/// `data` envelopes) into top-level `messages[]`. `channel_id` is pulled -/// from each match's `channel.id` field. `paging.pages` is preserved at -/// top-level under `pages` for the caller to drive pagination. -fn reshape_search_messages(data: &mut Value) { - // Preserve paging info before mutating data (take_array below removes the - // envelope that carries it). - let pages = [ - data.pointer("/data/messages/paging/pages"), - data.pointer("/messages/paging/pages"), - data.pointer("/data/data/messages/paging/pages"), - ] - .into_iter() - .flatten() - .find_map(|v| v.as_u64()) - .unwrap_or(1); - - // Envelope depth 1 removes the `messages` object (matches + paging) that - // held the consumed rows, not just the `matches` array. - let arr = take_array( - data, - &[ - "/data/messages/matches", - "/messages/matches", - "/data/data/messages/matches", - ], - 1, - ); - - let slim: Vec = arr.into_iter().filter_map(slim_search_match).collect(); - with_object(data, |obj| { - obj.insert("messages".to_string(), Value::Array(slim)); - obj.insert("pages".to_string(), Value::Number(pages.into())); - }); - log::debug!("[composio:slack][post-process] SLACK_SEARCH_MESSAGES reshaped"); -} - -fn slim_search_match(raw: Value) -> Option { - let text = raw - .get("text") - .and_then(|v| v.as_str()) - .unwrap_or("") - .trim(); - if text.is_empty() { - return None; - } - let ts = raw.get("ts")?; - let channel_id = raw - .pointer("/channel/id") - .and_then(|v| v.as_str()) - .unwrap_or("") - .trim(); - - let mut out = Map::new(); - out.insert("ts".into(), ts.clone()); - if let Some(user) = raw.get("user").or_else(|| raw.get("bot_id")) { - out.insert("user".into(), user.clone()); - } - out.insert("text".into(), Value::String(text.to_string())); - if let Some(thread_ts) = raw.get("thread_ts") { - out.insert("thread_ts".into(), thread_ts.clone()); - } - if !channel_id.is_empty() { - out.insert("channel_id".into(), Value::String(channel_id.to_string())); - } - if let Some(permalink) = raw.get("permalink") { - out.insert("permalink".into(), permalink.clone()); - } - Some(Value::Object(out)) -} - -// ─── Helpers ──────────────────────────────────────────────────────────────── - -/// Edit `data` as a JSON object, replacing it with an empty object first if -/// it is not one. -/// -/// Takes the value out, edits the map, and puts it back, so there is no -/// "re-borrow as an object" step that would need an `expect`. -fn with_object(data: &mut Value, edit: impl FnOnce(&mut Map)) { - let mut map = match std::mem::take(data) { - Value::Object(map) => map, - _ => Map::new(), - }; - edit(&mut map); - *data = Value::Object(map); -} - -#[cfg(test)] -#[path = "mod_tests.rs"] -mod tests; diff --git a/crates/tinymemory-integrations/src/sources/composio/slack_post_process/mod_tests.rs b/crates/tinymemory-integrations/src/sources/composio/slack_post_process/mod_tests.rs deleted file mode 100644 index 73a6702b..00000000 --- a/crates/tinymemory-integrations/src/sources/composio/slack_post_process/mod_tests.rs +++ /dev/null @@ -1,306 +0,0 @@ -//! Tests for the Slack post-processor. - -use super::*; -use serde_json::json; - -// ─── SLACK_FETCH_CONVERSATION_HISTORY ───────────────────────────────────── - -#[test] -fn history_reshapes_top_level_messages() { - let mut data = json!({ - "messages": [ - { "ts": "1714003200.000100", "user": "U1", "text": "hello" }, - { "ts": "1714003300.000200", "user": "U2", "text": "world", "thread_ts": "1714003200.0" }, - { "ts": "1714003400.000300", "user": "U3", "text": " " }, // dropped: empty text - ], - "response_metadata": { "next_cursor": "abc" } - }); - post_process("SLACK_FETCH_CONVERSATION_HISTORY", None, &mut data); - - let msgs = data["messages"].as_array().unwrap(); - assert_eq!(msgs.len(), 2, "empty-text message must be dropped"); - assert_eq!(msgs[0]["ts"], "1714003200.000100"); - assert_eq!(msgs[0]["user"], "U1"); - assert_eq!(msgs[0]["text"], "hello"); - assert!(msgs[0].get("thread_ts").is_none()); - assert_eq!(msgs[1]["thread_ts"], "1714003200.0"); -} - -#[test] -fn history_reshapes_nested_data_envelope() { - let mut data = json!({ - "data": { - "messages": [ - { "ts": "1714003200.0", "user": "U1", "text": "hi" } - ] - } - }); - post_process("SLACK_FETCH_CONVERSATION_HISTORY", None, &mut data); - let msgs = data["messages"].as_array().unwrap(); - assert_eq!(msgs.len(), 1); - assert_eq!(msgs[0]["text"], "hi"); -} - -#[test] -fn history_reshapes_doubly_nested_envelope() { - let mut data = json!({ - "data": { - "data": { - "messages": [ - { "ts": "1714003200.0", "user": "U1", "text": "deep" } - ] - } - } - }); - post_process("SLACK_FETCH_CONVERSATION_HISTORY", None, &mut data); - let msgs = data["messages"].as_array().unwrap(); - assert_eq!(msgs.len(), 1); - assert_eq!(msgs[0]["text"], "deep"); -} - -#[test] -fn history_drops_message_without_ts() { - let mut data = json!({ - "messages": [ - { "user": "U1", "text": "no timestamp" }, - { "ts": "1714003200.0", "user": "U2", "text": "has ts" }, - ] - }); - post_process("SLACK_FETCH_CONVERSATION_HISTORY", None, &mut data); - let msgs = data["messages"].as_array().unwrap(); - assert_eq!(msgs.len(), 1); - assert_eq!(msgs[0]["text"], "has ts"); -} - -#[test] -fn history_removes_nested_envelope_after_reshape() { - let mut data = json!({ - "data": { - "messages": [ - { "ts": "1714003200.0", "user": "U1", "text": "hi" } - ] - } - }); - post_process("SLACK_FETCH_CONVERSATION_HISTORY", None, &mut data); - - let msgs = data["messages"].as_array().unwrap(); - assert_eq!(msgs.len(), 1); - assert_eq!(msgs[0]["text"], "hi"); - assert!( - data.pointer("/data").is_none(), - "consumed `data.messages` envelope must be removed, got: {data}" - ); -} - -// ─── SLACK_LIST_CONVERSATIONS ───────────────────────────────────────────── - -#[test] -fn list_conversations_reshapes_channels() { - let mut data = json!({ - "data": { - "channels": [ - { "id": "C1", "name": "eng", "is_private": false, "extra": "noise" }, - { "id": "G1", "name": "ops", "is_private": true }, - { "id": "", "name": "empty-id" }, // dropped - ] - } - }); - post_process("SLACK_LIST_CONVERSATIONS", None, &mut data); - let channels = data["channels"].as_array().unwrap(); - assert_eq!(channels.len(), 2, "empty-id entry must be dropped"); - assert_eq!(channels[0]["id"], "C1"); - assert_eq!(channels[0]["name"], "eng"); - assert_eq!(channels[0]["is_private"], false); - assert!( - channels[0].get("extra").is_none(), - "noise fields must be removed" - ); - assert_eq!(channels[1]["id"], "G1"); - assert_eq!(channels[1]["is_private"], true); -} - -#[test] -fn list_conversations_falls_back_to_conversations_key() { - let mut data = json!({ - "conversations": [ - { "id": "C2", "name": "dev", "is_private": false } - ] - }); - post_process("SLACK_LIST_CONVERSATIONS", None, &mut data); - let channels = data["channels"].as_array().unwrap(); - assert_eq!(channels.len(), 1); - assert_eq!(channels[0]["id"], "C2"); - assert!( - data.pointer("/conversations").is_none(), - "consumed `conversations` field must be removed" - ); -} - -// ─── SLACK_SEARCH_MESSAGES ──────────────────────────────────────────────── - -#[test] -fn search_messages_reshapes_matches() { - let mut data = json!({ - "messages": { - "matches": [ - { - "ts": "1714003200.0", - "user": "U1", - "text": "hello from search", - "channel": { "id": "C1" } - }, - { - "ts": "1714003300.0", - "user": "U2", - "text": " ", // dropped: whitespace only - "channel": { "id": "C1" } - }, - ], - "paging": { "pages": 3 } - } - }); - post_process("SLACK_SEARCH_MESSAGES", None, &mut data); - let msgs = data["messages"].as_array().unwrap(); - assert_eq!(msgs.len(), 1, "empty-text match must be dropped"); - assert_eq!(msgs[0]["ts"], "1714003200.0"); - assert_eq!(msgs[0]["text"], "hello from search"); - assert_eq!(msgs[0]["channel_id"], "C1"); - assert_eq!(data["pages"], 3, "paging.pages must be preserved"); -} - -#[test] -fn search_messages_nested_data_envelope() { - let mut data = json!({ - "data": { - "messages": { - "matches": [ - { "ts": "1714003200.0", "user": "U1", "text": "nested", "channel": { "id": "C2" } } - ], - "paging": { "pages": 1 } - } - } - }); - post_process("SLACK_SEARCH_MESSAGES", None, &mut data); - let msgs = data["messages"].as_array().unwrap(); - assert_eq!(msgs.len(), 1); - assert_eq!(msgs[0]["channel_id"], "C2"); - assert_eq!(data["pages"], 1_u64); -} - -#[test] -fn search_messages_no_matches_emits_empty_array() { - let mut data = json!({ "messages": { "matches": [] } }); - post_process("SLACK_SEARCH_MESSAGES", None, &mut data); - let msgs = data["messages"].as_array().unwrap(); - assert!(msgs.is_empty()); -} - -#[test] -fn search_messages_removes_nested_envelope_after_reshape() { - let mut data = json!({ - "data": { - "messages": { - "matches": [ - { "ts": "1714003200.0", "user": "U1", "text": "nested", "channel": { "id": "C2" } } - ], - "paging": { "pages": 1 } - } - } - }); - post_process("SLACK_SEARCH_MESSAGES", None, &mut data); - - let msgs = data["messages"].as_array().unwrap(); - assert_eq!(msgs.len(), 1); - assert_eq!(msgs[0]["channel_id"], "C2"); - assert_eq!(data["pages"], 1_u64); - assert!( - data.pointer("/data").is_none(), - "consumed `data.messages` envelope must be removed, got: {data}" - ); -} - -#[test] -fn search_messages_doubly_nested_paging_preserved() { - let mut data = json!({ - "data": { - "data": { - "messages": { - "matches": [ - { "ts": "1714003200.0", "user": "U1", "text": "deep", "channel": { "id": "C3" } } - ], - "paging": { "pages": 4 } - } - } - } - }); - post_process("SLACK_SEARCH_MESSAGES", None, &mut data); - - let msgs = data["messages"].as_array().unwrap(); - assert_eq!(msgs.len(), 1); - assert_eq!(msgs[0]["text"], "deep"); - assert_eq!( - data["pages"], 4_u64, - "doubly-nested paging must be preserved" - ); - assert!( - data.pointer("/data").is_none(), - "consumed `data.data.messages` envelope must be removed, got: {data}" - ); -} - -// ─── Unknown slug ───────────────────────────────────────────────────────── - -#[test] -fn unknown_slug_is_noop() { - let mut data = json!({ "foo": "bar" }); - let original = data.clone(); - post_process("SLACK_SEND_MESSAGE", None, &mut data); - assert_eq!(data, original, "unknown slug must not mutate data"); -} - -#[test] -fn malformed_history_payload_becomes_an_empty_stable_shape() { - for mut data in [json!(null), json!([]), json!({"messages": "wrong"})] { - post_process("SLACK_FETCH_CONVERSATION_HISTORY", None, &mut data); - assert_eq!(data, json!({"messages": []})); - } -} - -#[test] -fn malformed_channel_rows_are_dropped_and_defaults_are_stable() { - let mut data = json!({ - "channels": [ - null, - "not an object", - {"id": 7, "name": "numeric id"}, - {"id": " C1 ", "name": 99, "is_private": "yes"} - ] - }); - post_process("SLACK_LIST_CONVERSATIONS", None, &mut data); - assert_eq!( - data, - json!({"channels": [{"id":"C1", "name":"C1", "is_private":false}]}) - ); -} - -#[test] -fn malformed_search_rows_and_page_counts_fall_back_safely() { - let mut data = json!({ - "messages": { - "matches": [ - null, - {"ts":"1.0", "text":"no channel"}, - {"channel":{"id":"C1"}, "text":"no timestamp"}, - {"ts":"2.0", "channel":{"id":"C1"}, "text":" valid "} - ], - "paging": {"pages":"many"} - } - }); - post_process("SLACK_SEARCH_MESSAGES", None, &mut data); - assert_eq!(data["pages"], 1); - let messages = data["messages"].as_array().unwrap(); - assert_eq!(messages.len(), 2); - assert_eq!(messages[0]["text"], "no channel"); - assert!(messages[0].get("channel_id").is_none()); - assert_eq!(messages[1]["text"], "valid"); -} diff --git a/crates/tinymemory-integrations/src/sources/items/mod.rs b/crates/tinymemory-integrations/src/sources/items/mod.rs index 40205a2e..9bbd5f89 100644 --- a/crates/tinymemory-integrations/src/sources/items/mod.rs +++ b/crates/tinymemory-integrations/src/sources/items/mod.rs @@ -11,7 +11,6 @@ //! | github | document | `repo` (`owner/name`), `commit` (commit items), `url` (issues and PRs), `observed_at` | //! | link | document | `url` | //! | rss | document | `url` (the entry's link), `observed_at` (published) | -//! | composio | document | `tags = [toolkit]` (payloads: see [`crate::sources::composio`]) | //! | conversation | conversation | `workspace`, `thread_id`, `turns`, `observed_at` (last turn) | //! //! Every document body is markdown, converted through `crate::documents`: @@ -38,23 +37,16 @@ use crate::sources::readers::local_file::LocalFile; use crate::sources::types::{ContentType, MemorySourceEntry, SourceContent, SourceKind}; /// Metadata naming `entry` as the source: `source.kind` is the entry's kind -/// mapped onto the contract, `source.id` its id. A Composio entry also gets -/// its toolkit as a tag. +/// mapped onto the contract, `source.id` its id. #[must_use] pub fn base_meta(entry: &MemorySourceEntry) -> MemoryMeta { - let mut meta = MemoryMeta { + MemoryMeta { source: SourceRef { kind: entry.kind.api_kind(), id: Some(entry.id.clone()), }, ..MemoryMeta::default() - }; - if entry.kind == SourceKind::Composio - && let Some(toolkit) = entry.toolkit.as_deref().filter(|t| !t.is_empty()) - { - meta.tags = vec![toolkit.to_string()]; } - meta } /// Convert a local file and wrap it as a document. @@ -197,7 +189,6 @@ pub fn content_item( meta.language = language_for_path(&content.id).map(str::to_string); } SourceKind::Conversation => meta.thread_id = Some(content.id.clone()), - SourceKind::Composio => {} } Ok(StoreItem::Document { diff --git a/crates/tinymemory-integrations/src/sources/items/mod_tests.rs b/crates/tinymemory-integrations/src/sources/items/mod_tests.rs index b050d8ea..e11cae90 100644 --- a/crates/tinymemory-integrations/src/sources/items/mod_tests.rs +++ b/crates/tinymemory-integrations/src/sources/items/mod_tests.rs @@ -54,17 +54,12 @@ fn base_meta_names_the_source_by_contract_kind_and_entry_id() { (SourceKind::WebPage, Api::Link), (SourceKind::GithubRepo, Api::Github), (SourceKind::RssFeed, Api::Rss), - (SourceKind::Composio, Api::Composio), (SourceKind::Conversation, Api::Conversation), ] { let meta = base_meta(&entry(kind)); assert_eq!(meta.source.kind, api); assert_eq!(meta.source.id.as_deref(), Some("src_test")); } - - let mut composio = entry(SourceKind::Composio); - composio.toolkit = Some("gmail".into()); - assert_eq!(base_meta(&composio).tags, vec!["gmail".to_string()]); } #[test] @@ -168,7 +163,7 @@ fn link_items_take_the_page_url() { } #[test] -fn reader_content_for_local_kinds_and_composio_fills_what_it_can() { +fn reader_content_for_local_kinds_fills_what_it_can() { let mut folder = entry(SourceKind::Folder); folder.path = Some("/notes".into()); let item = content_item( @@ -212,21 +207,6 @@ fn reader_content_for_local_kinds_and_composio_fills_what_it_can() { ) .unwrap(); assert_eq!(conversation.meta().thread_id.as_deref(), Some("t1")); - - let mut composio = entry(SourceKind::Composio); - composio.toolkit = Some("slack".into()); - let item = content_item( - &composio, - content( - "c1", - "sync data", - ContentType::Plaintext, - serde_json::json!({}), - ), - None, - ) - .unwrap(); - assert_eq!(item.meta().tags, vec!["slack".to_string()]); } #[test] diff --git a/crates/tinymemory-integrations/src/sources/mod.rs b/crates/tinymemory-integrations/src/sources/mod.rs index dad654b8..10668fc9 100644 --- a/crates/tinymemory-integrations/src/sources/mod.rs +++ b/crates/tinymemory-integrations/src/sources/mod.rs @@ -1,5 +1,5 @@ //! Source readers for TinyMemory: turn a folder, a file, a web page, a GitHub -//! repository, an RSS feed, a Composio toolkit payload or a local conversation +//! repository, an RSS feed or a local conversation //! into [`StoreItem`](tinymemory_api::StoreItem)s. //! //! - **Configuration** — what a source *is* ([`MemorySourceEntry`], keyed by @@ -13,8 +13,6 @@ //! - **Items** — [`items`] maps reader output to `StoreItem`s with //! [`MemoryMeta`](tinymemory_api::MemoryMeta) filled per kind; every text //! body is converted to markdown through [`crate::documents`]. -//! - **Composio** — [`composio`] normalises toolkit payloads (Gmail, Slack, -//! GitHub, Linear, Notion, ClickUp) and maps them to items. //! //! Scheduling, credentials and egress budgets stay with the host: this module //! reads when asked. @@ -63,7 +61,6 @@ //! `readers::reader_for_request` and the SSRF guard. Without it, a host that //! only reads local sources links no HTTP stack. -pub mod composio; pub mod error; #[cfg(feature = "sources-network")] pub mod fetch; diff --git a/crates/tinymemory-integrations/src/sources/readers/composio/mod.rs b/crates/tinymemory-integrations/src/sources/readers/composio/mod.rs deleted file mode 100644 index 5d0c3d84..00000000 --- a/crates/tinymemory-integrations/src/sources/readers/composio/mod.rs +++ /dev/null @@ -1,78 +0,0 @@ -//! Composio source reader — a placeholder over the provider pipeline. -//! -//! Composio data does not arrive item by item: the host runs toolkit actions -//! with its credentials and hands the responses to [`crate::sources::composio`], which -//! normalises them and maps them to `StoreItem`s. For a Composio source, -//! `list_items` returns the connection as one sync target and `read_item` -//! describes that pipeline. The reader exists so `reader_for_request` can hand -//! out a reader for every source kind uniformly. - -use std::path::Path; - -use async_trait::async_trait; - -use super::SourceReader; -use crate::sources::error::Result; -use crate::sources::types::{ - ContentType, MemorySourceEntry, SourceContent, SourceItem, SourceKind, -}; - -/// Lists a Composio connection as a single sync target. -/// -/// Composio data arrives through the provider sync pipeline rather than -/// item-by-item, so `read_item` returns a description of that rather than -/// content. The reader exists so `reader_for_request` can serve every source -/// kind uniformly. -#[derive(Debug, Clone, Copy, Default)] -pub struct ComposioReader; - -#[async_trait] -impl SourceReader for ComposioReader { - fn kind(&self) -> SourceKind { - SourceKind::Composio - } - - async fn list_items( - &self, - source: &MemorySourceEntry, - _workspace: &Path, - ) -> Result> { - let toolkit = source.toolkit.as_deref().unwrap_or("unknown"); - let connection_id = source.connection_id.as_deref().unwrap_or("unknown"); - - log::debug!( - "[memory_sources:composio] list_items toolkit={toolkit} connection_id={connection_id}" - ); - - Ok(vec![SourceItem { - id: connection_id.to_string(), - title: format!("{toolkit} connection"), - updated_at_ms: None, - }]) - } - - async fn read_item( - &self, - source: &MemorySourceEntry, - item_id: &str, - _workspace: &Path, - ) -> Result { - let toolkit = source.toolkit.as_deref().unwrap_or("unknown"); - Ok(SourceContent { - id: item_id.to_string(), - title: format!("{toolkit} sync data"), - body: format!( - "Composio {toolkit} data is synced via the provider sync pipeline, not read item-by-item." - ), - content_type: ContentType::Plaintext, - metadata: serde_json::json!({ - "toolkit": toolkit, - "connection_id": source.connection_id, - }), - }) - } -} - -#[cfg(test)] -#[path = "mod_tests.rs"] -mod tests; diff --git a/crates/tinymemory-integrations/src/sources/readers/composio/mod_tests.rs b/crates/tinymemory-integrations/src/sources/readers/composio/mod_tests.rs deleted file mode 100644 index d139ff15..00000000 --- a/crates/tinymemory-integrations/src/sources/readers/composio/mod_tests.rs +++ /dev/null @@ -1,51 +0,0 @@ -//! Tests for the surrounding module. - -use super::*; -use std::path::Path; - -fn test_source() -> MemorySourceEntry { - MemorySourceEntry { - id: "src_1".into(), - kind: SourceKind::Composio, - label: "Gmail".into(), - enabled: true, - toolkit: Some("gmail".into()), - connection_id: Some("cmp_123".into()), - path: None, - glob: None, - url: None, - branch: None, - paths: Vec::new(), - max_items: None, - max_commits: None, - max_issues: None, - max_prs: None, - selector: None, - max_tokens_per_sync: None, - max_cost_per_sync_usd: None, - sync_depth_days: None, - } -} - -#[tokio::test] -async fn list_items_returns_connection_as_item() { - let reader = ComposioReader; - let items = reader - .list_items(&test_source(), Path::new(".")) - .await - .unwrap(); - assert_eq!(items.len(), 1); - assert_eq!(items[0].id, "cmp_123"); -} - -#[tokio::test] -async fn read_item_describes_the_provider_pipeline() { - let content = ComposioReader - .read_item(&test_source(), "cmp_123", Path::new(".")) - .await - .unwrap(); - assert_eq!(content.title, "gmail sync data"); - assert!(content.body.contains("provider sync pipeline")); - assert_eq!(content.metadata["toolkit"], "gmail"); - assert_eq!(content.metadata["connection_id"], "cmp_123"); -} diff --git a/crates/tinymemory-integrations/src/sources/readers/conversation/mod_tests.rs b/crates/tinymemory-integrations/src/sources/readers/conversation/mod_tests.rs index 8f49976f..98961bde 100644 --- a/crates/tinymemory-integrations/src/sources/readers/conversation/mod_tests.rs +++ b/crates/tinymemory-integrations/src/sources/readers/conversation/mod_tests.rs @@ -11,8 +11,6 @@ fn conversation_source() -> MemorySourceEntry { kind: SourceKind::Conversation, label: "Conversations".into(), enabled: true, - toolkit: None, - connection_id: None, path: None, glob: None, url: None, diff --git a/crates/tinymemory-integrations/src/sources/readers/folder/mod_tests.rs b/crates/tinymemory-integrations/src/sources/readers/folder/mod_tests.rs index 32060ca7..970e96d3 100644 --- a/crates/tinymemory-integrations/src/sources/readers/folder/mod_tests.rs +++ b/crates/tinymemory-integrations/src/sources/readers/folder/mod_tests.rs @@ -11,8 +11,6 @@ fn folder_source(path: &str) -> MemorySourceEntry { kind: SourceKind::Folder, label: "Test folder".into(), enabled: true, - toolkit: None, - connection_id: None, path: Some(path.into()), glob: None, url: None, diff --git a/crates/tinymemory-integrations/src/sources/readers/github/mod_tests.rs b/crates/tinymemory-integrations/src/sources/readers/github/mod_tests.rs index 66168524..2b00d73c 100644 --- a/crates/tinymemory-integrations/src/sources/readers/github/mod_tests.rs +++ b/crates/tinymemory-integrations/src/sources/readers/github/mod_tests.rs @@ -10,8 +10,6 @@ fn github_source(url: Option<&str>) -> MemorySourceEntry { kind: SourceKind::GithubRepo, label: "GitHub".into(), enabled: true, - toolkit: None, - connection_id: None, path: None, glob: None, url: url.map(str::to_string), diff --git a/crates/tinymemory-integrations/src/sources/readers/mod.rs b/crates/tinymemory-integrations/src/sources/readers/mod.rs index 503a21ee..e74907af 100644 --- a/crates/tinymemory-integrations/src/sources/readers/mod.rs +++ b/crates/tinymemory-integrations/src/sources/readers/mod.rs @@ -23,14 +23,9 @@ //! on a timer. A `None` from [`reader_for`] therefore means "route this through //! the host's sync runner", which keeps the host in charge of the network. //! -//! `composio` is represented by a placeholder reader -//! ([`composio::ComposioReader`]): its data arrives through the credentialed -//! provider pipeline, and [`crate::sources::composio`] turns those payloads into items. -//! //! A host servicing an *explicit user request* (not a timer) that wants one //! reader for any kind uses `reader_for_request`. -pub mod composio; pub mod conversation; pub mod file; pub mod folder; @@ -128,8 +123,7 @@ pub fn is_locally_readable(kind: &SourceKind) -> bool { /// Get the reader for a source kind that is safe to drive on a timer. /// /// Returns `Some` for [`SourceKind::Folder`], [`SourceKind::File`] and -/// [`SourceKind::Conversation`]. Network-backed kinds (`composio`, -/// `github_repo`, `rss_feed`, `web_page`) return `None` so the caller defers to +/// [`SourceKind::Conversation`]. Network-backed kinds (/// `github_repo`, `rss_feed`, `web_page`) return `None` so the caller defers to /// the host's sync runner, which constructs those readers once it has /// authorized the fetch. #[must_use] @@ -138,10 +132,7 @@ pub fn reader_for(kind: &SourceKind) -> Option> { SourceKind::Folder => Some(Box::new(folder::FolderReader)), SourceKind::File => Some(Box::new(file::FileReader)), SourceKind::Conversation => Some(Box::new(conversation::ConversationReader)), - SourceKind::Composio - | SourceKind::GithubRepo - | SourceKind::RssFeed - | SourceKind::WebPage => None, + SourceKind::GithubRepo | SourceKind::RssFeed | SourceKind::WebPage => None, } } @@ -156,7 +147,6 @@ pub fn reader_for(kind: &SourceKind) -> Option> { #[must_use] pub fn reader_for_request(kind: &SourceKind) -> Box { match kind { - SourceKind::Composio => Box::new(composio::ComposioReader), SourceKind::Conversation => Box::new(conversation::ConversationReader), SourceKind::Folder => Box::new(folder::FolderReader), SourceKind::File => Box::new(file::FileReader), diff --git a/crates/tinymemory-integrations/src/sources/readers/rss/mod_tests.rs b/crates/tinymemory-integrations/src/sources/readers/rss/mod_tests.rs index cf6fc2ca..a6dbdf41 100644 --- a/crates/tinymemory-integrations/src/sources/readers/rss/mod_tests.rs +++ b/crates/tinymemory-integrations/src/sources/readers/rss/mod_tests.rs @@ -38,8 +38,6 @@ fn rss_source(url: Option<&str>, max_items: Option) -> MemorySourceEntry { label: "Feed".into(), kind: SourceKind::RssFeed, enabled: true, - toolkit: None, - connection_id: None, path: None, glob: None, url: url.map(str::to_string), diff --git a/crates/tinymemory-integrations/src/sources/readers/web_page/mod_tests.rs b/crates/tinymemory-integrations/src/sources/readers/web_page/mod_tests.rs index 39c0baba..30a695af 100644 --- a/crates/tinymemory-integrations/src/sources/readers/web_page/mod_tests.rs +++ b/crates/tinymemory-integrations/src/sources/readers/web_page/mod_tests.rs @@ -9,8 +9,6 @@ fn web_source(url: Option<&str>, selector: Option<&str>) -> MemorySourceEntry { kind: SourceKind::WebPage, label: "Reference page".into(), enabled: true, - toolkit: None, - connection_id: None, path: None, glob: None, url: url.map(str::to_string), diff --git a/crates/tinymemory-integrations/src/sources/types/mod.rs b/crates/tinymemory-integrations/src/sources/types/mod.rs index 156f2589..6e974d33 100644 --- a/crates/tinymemory-integrations/src/sources/types/mod.rs +++ b/crates/tinymemory-integrations/src/sources/types/mod.rs @@ -31,9 +31,6 @@ pub(crate) fn default_true() -> bool { #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)] #[serde(rename_all = "snake_case")] pub enum SourceKind { - /// A Composio OAuth connector (Gmail, Slack, Notion, …). Network-backed; - /// the live fetch is owned by the host, not this module. - Composio, /// Local agent conversation transcripts stored in the workspace. Conversation, /// A local folder of files matched by an optional glob. @@ -50,8 +47,7 @@ pub enum SourceKind { impl SourceKind { /// Every kind, in declaration order. - pub const ALL: [Self; 7] = [ - Self::Composio, + pub const ALL: [Self; 6] = [ Self::Conversation, Self::Folder, Self::File, @@ -64,7 +60,6 @@ impl SourceKind { #[must_use] pub fn as_str(&self) -> &'static str { match self { - SourceKind::Composio => "composio", SourceKind::Conversation => "conversation", SourceKind::Folder => "folder", SourceKind::File => "file", @@ -81,7 +76,6 @@ impl SourceKind { pub fn api_kind(&self) -> tinymemory_api::SourceKind { use tinymemory_api::SourceKind as Api; match self { - SourceKind::Composio => Api::Composio, SourceKind::Conversation => Api::Conversation, SourceKind::Folder => Api::Folder, SourceKind::File => Api::File, @@ -109,14 +103,6 @@ pub struct MemorySourceEntry { #[serde(default = "default_true")] pub enabled: bool, - // ── Composio ── - /// Composio toolkit slug (e.g. `gmail`). Required for `composio`. - #[serde(default, skip_serializing_if = "Option::is_none")] - pub toolkit: Option, - /// Composio connection id. Required for `composio`. - #[serde(default, skip_serializing_if = "Option::is_none")] - pub connection_id: Option, - // ── Folder / File ── /// Filesystem path of the folder or file to read. Required for `folder` /// and `file`; a relative path is anchored on the workspace. @@ -180,8 +166,6 @@ impl MemorySourceEntry { kind, label: label.into(), enabled: true, - toolkit: None, - connection_id: None, path: None, glob: None, url: None, @@ -201,8 +185,7 @@ impl MemorySourceEntry { /// Validate the fields this entry's [`SourceKind`] requires. /// /// `id` and `label` are required for every kind, and `id` must not contain - /// `:` or control characters. Composio needs `toolkit` and - /// `connection_id`; folders and files need `path`; GitHub repositories, RSS + /// `:` or control characters. Folders and files need `path`; GitHub repositories, RSS /// feeds and web pages need `url`. An empty string counts as missing. /// /// # Errors @@ -221,10 +204,6 @@ impl MemorySourceEntry { return Err(Error::Invalid("label is required".to_string())); } match self.kind { - SourceKind::Composio => { - require_field(&self.toolkit, "toolkit")?; - require_field(&self.connection_id, "connection_id") - } SourceKind::Conversation => Ok(()), SourceKind::Folder | SourceKind::File => require_field(&self.path, "path"), SourceKind::GithubRepo | SourceKind::RssFeed | SourceKind::WebPage => { diff --git a/crates/tinymemory-integrations/src/sources/types/mod_tests.rs b/crates/tinymemory-integrations/src/sources/types/mod_tests.rs index 2ea2171b..7ac3adb0 100644 --- a/crates/tinymemory-integrations/src/sources/types/mod_tests.rs +++ b/crates/tinymemory-integrations/src/sources/types/mod_tests.rs @@ -5,7 +5,6 @@ use super::*; #[test] fn source_kind_round_trips_via_serde() { for kind in [ - SourceKind::Composio, SourceKind::Conversation, SourceKind::Folder, SourceKind::File, @@ -21,7 +20,6 @@ fn source_kind_round_trips_via_serde() { #[test] fn source_kind_as_str_matches_wire_strings() { - assert_eq!(SourceKind::Composio.as_str(), "composio"); assert_eq!(SourceKind::Conversation.as_str(), "conversation"); assert_eq!(SourceKind::Folder.as_str(), "folder"); assert_eq!(SourceKind::GithubRepo.as_str(), "github_repo"); @@ -30,26 +28,6 @@ fn source_kind_as_str_matches_wire_strings() { assert_eq!(SourceKind::WebPage.as_str(), "web_page"); } -#[test] -fn validate_composio_requires_toolkit_and_connection_id() { - let entry = MemorySourceEntry { - id: "src_1".into(), - kind: SourceKind::Composio, - label: "Gmail".into(), - enabled: true, - toolkit: Some("gmail".into()), - connection_id: None, - ..default_entry() - }; - assert!(entry.validate().is_err()); - - let valid = MemorySourceEntry { - connection_id: Some("cmp_123".into()), - ..entry - }; - assert!(valid.validate().is_ok()); -} - #[test] fn validate_folder_requires_path() { let entry = MemorySourceEntry { @@ -89,8 +67,10 @@ fn validate_file_requires_path() { #[test] fn the_removed_twitter_query_kind_no_longer_decodes() { - let decoded = serde_json::from_str::("\"twitter_query\""); - assert!(decoded.is_err()); + for removed in ["twitter_query", "composio"] { + let decoded = serde_json::from_str::(&format!("\"{removed}\"")); + assert!(decoded.is_err(), "{removed} must not decode"); + } } #[test] @@ -100,7 +80,6 @@ fn every_config_kind_maps_onto_a_contract_source_kind() { assert_eq!( mapped, vec![ - Api::Composio, Api::Conversation, Api::Folder, Api::File, @@ -245,8 +224,6 @@ pub(super) fn default_entry() -> MemorySourceEntry { kind: SourceKind::Folder, label: String::new(), enabled: true, - toolkit: None, - connection_id: None, path: None, glob: None, url: None, @@ -277,8 +254,6 @@ fn source_entry_wire_format_is_pinned() { kind: SourceKind::GithubRepo, label: "Pinned".into(), enabled: false, - toolkit: Some("gmail".into()), - connection_id: Some("conn-1".into()), path: Some("/notes".into()), glob: Some("**/*.md".into()), url: Some("https://github.com/tinyhumansai/tinymemory".into()), @@ -301,8 +276,6 @@ fn source_entry_wire_format_is_pinned() { "kind": "github_repo", "label": "Pinned", "enabled": false, - "toolkit": "gmail", - "connection_id": "conn-1", "path": "/notes", "glob": "**/*.md", "url": "https://github.com/tinyhumansai/tinymemory", diff --git a/crates/tinymemory-integrations/tests/reader_dispatch.rs b/crates/tinymemory-integrations/tests/reader_dispatch.rs index 6cb12752..946faac7 100644 --- a/crates/tinymemory-integrations/tests/reader_dispatch.rs +++ b/crates/tinymemory-integrations/tests/reader_dispatch.rs @@ -17,7 +17,6 @@ fn timer_dispatch_constructs_only_readers_that_never_need_network() { } for kind in [ - SourceKind::Composio, SourceKind::GithubRepo, SourceKind::RssFeed, SourceKind::WebPage, diff --git a/crates/tinymemory-tools/src/layout/source.rs b/crates/tinymemory-tools/src/layout/source.rs index dcf203d2..4acc6932 100644 --- a/crates/tinymemory-tools/src/layout/source.rs +++ b/crates/tinymemory-tools/src/layout/source.rs @@ -60,10 +60,9 @@ impl BrainSource { pub fn source_kind(&self) -> SourceKind { match self { Self::Files | Self::Pdf | Self::Markdown => SourceKind::File, - Self::Notion => SourceKind::Composio, Self::Github => SourceKind::Github, Self::Web => SourceKind::Link, - Self::Other(_) => SourceKind::Import, + Self::Notion | Self::Other(_) => SourceKind::Import, } } } diff --git a/crates/tinymemory-tools/tests/fixtures/tool_contracts.json b/crates/tinymemory-tools/tests/fixtures/tool_contracts.json index fab2afbb..cb1eacc5 100644 --- a/crates/tinymemory-tools/tests/fixtures/tool_contracts.json +++ b/crates/tinymemory-tools/tests/fixtures/tool_contracts.json @@ -56,7 +56,6 @@ "link", "github", "rss", - "composio", "conversation", "agent", "import" @@ -191,7 +190,6 @@ "link", "github", "rss", - "composio", "conversation", "agent", "import" @@ -332,7 +330,6 @@ "link", "github", "rss", - "composio", "conversation", "agent", "import" @@ -472,7 +469,6 @@ "link", "github", "rss", - "composio", "conversation", "agent", "import" @@ -685,7 +681,6 @@ "link", "github", "rss", - "composio", "conversation", "agent", "import" diff --git a/docs/architecture/api-items.md b/docs/architecture/api-items.md index 2909c584..f62ef289 100644 --- a/docs/architecture/api-items.md +++ b/docs/architecture/api-items.md @@ -165,7 +165,7 @@ namespace are omitted on serialisation. | `observed_actor` | `Option` | `{ id, name? }`: who said or did it when not the memory's owner, such as an email's sender or a channel's (`id` is `type:id`, `user:priya@acme.com`, `user:+15551234567`). Excluded from the fingerprint; the CortexDB engine sends it only with attribution on. | `MemoryMeta::from_source(kind, id)` sets only the source. `SourceKind` is -`folder | file | link | github | rss | composio | conversation | agent | +`folder | file | link | github | rss | conversation | agent | import` (`SourceKind::ALL`, `as_str`); `agent` is the default. ## `MetaFilter` diff --git a/docs/architecture/integrations-sources.md b/docs/architecture/integrations-sources.md index 98be38d1..1b0ce235 100644 --- a/docs/architecture/integrations-sources.md +++ b/docs/architecture/integrations-sources.md @@ -22,7 +22,6 @@ optional fields are required; `validate()` checks them. | `web_page` | `url` | `selector` | `Link` | | `github_repo` | `url` | `branch`, `paths`, `max_commits`, `max_issues`, `max_prs` (default 1000 each) | `Github` | | `rss_feed` | `url` | `max_items` (default 50) | `Rss` | -| `composio` | `toolkit`, `connection_id` | | `Composio` | Every entry also needs a non-blank `id` (no `:` or control characters) and a non-empty `label`, and carries `enabled` and optional sync-budget fields @@ -52,7 +51,6 @@ for a host servicing an explicit user request, never a polling loop. | `WebPageReader` | one: the page URL | With a CSS `selector`, only the text of matching elements (plain text; only the last compound of a descendant chain is honoured). Otherwise the whole page as markdown. 10 MiB body cap. | | `RssReader` | one per feed entry (RSS or Atom), up to `max_items` | The parsed feed is cached for 60 seconds so a list-then-read pass downloads it once. 5 MiB cap; non-UTF-8 bodies are refused. | | `GithubReader` | `commit:`, `issue:`, `pr:` | See below. | -| `ComposioReader` | one: the connection | A placeholder; see Composio. | ### local_file and ensure_within_base @@ -96,7 +94,6 @@ are rejected). It combines three transports: | github | document | `repo` (`owner/name`), `commit` (commits), `url` (issues and PRs), `observed_at` | | web page | document | `url` | | rss | document | `url` (entry link), `observed_at` (published) | -| composio | document | `tags = [toolkit]`; payloads add `url`, `observed_at`, `thread_id`, `repo` | | conversation | conversation | `workspace`, `thread_id`, `turns` (`0..=n-1`), `observed_at` (last turn, else mtime) | `collect_items(reader, entry, workspace, converter)` lists and reads every @@ -105,48 +102,6 @@ with its error and the pass continues; only a listing failure is an `Err`. Other entry points: `file_item` (a path with no source), `conversation_item`, `content_item` and `items::local_file_item`. -## Composio - -Composio data does not arrive item by item. A host runs toolkit actions with -its own credentials and hands the raw responses to `sources::composio`, which -holds no credential, opens no socket and decides nothing about when to sync. - -1. **Normalisers**, pure `serde_json::Value` transforms, one module per - toolkit. They walk Composio's envelope variants (top level, under `data`, - under `data.data`) and return the first array found: `clickup` - (`extract_tasks`), `github` (`extract_issues`), `linear` - (`extract_issues`), `notion` (`extract_results`, `extract_page_markdown`), - each with title, id and updated-time helpers. `fields::pick_str` is the - shared lookup: it tries dotted paths, descends only through objects, and - rejects non-string leaves. -2. **Post-processors the host must call** for two toolkits, because their raw - responses are too verbose: - - `gmail_post_process::post_process(slug, arguments, &mut data)` rewrites a - `GMAIL_FETCH_EMAILS` response into slim `messages[]` (other Gmail slugs - pass through; `raw_html: true` in the arguments skips the reshape). If the - response carries a response-level `markdownFormatted` string, call - `apply_response_level_markdown(&mut data, markdown)` **before** - `post_process`; it is a no-op unless the split count matches the message - count. `format_email_local_time` renders in the host's local timezone; - the raw UTC fields are preserved. - - `slack_post_process::post_process(slug, arguments, &mut data)` reshapes - `SLACK_FETCH_CONVERSATION_HISTORY`, `SLACK_LIST_CONVERSATIONS` and - `SLACK_SEARCH_MESSAGES`; unknown slugs are no-ops. `channel_id` for history - is injected by the host (it is in the request, not the response), and user - ids are resolved by the host. -3. `normalise_payload(toolkit, &data)` returns `ComposioDocument`s (id, title, - markdown body, url, `observed_at`, `thread_id`, `repo`), dispatching on the - case-insensitive toolkit slug: `gmail`, `slack`, `github`, `linear`, - `notion`, `clickup`; any other toolkit falls back to each record as fenced - JSON, so a new toolkit is ingested verbosely rather than dropped. Records - with no text are skipped. `payload_items(toolkit, source_id, &data)` wraps - them as `StoreItem::Document` with `source.kind = Composio`, - `source.id = source_id` and `tags = [toolkit]`. - -`readers::composio::ComposioReader` is only a placeholder so -`reader_for_request` can serve every kind: `list_items` returns the connection -as one sync target. - ## Fetching and the SSRF guard (sources-network) `fetch::fetch_url(url)` fetches one URL into a `RawDocument`: the `Content-Type` diff --git a/docs/architecture/overview.md b/docs/architecture/overview.md index 94e6916e..d4c3589f 100644 --- a/docs/architecture/overview.md +++ b/docs/architecture/overview.md @@ -55,7 +55,7 @@ alone, so the contract is the only coupling point. | `tinymemory-integrations` | `cortex` (default) | `cortex::CortexEngine` (both wires), `registry` (`list_engines`, `build_engine`), `config` (`MemoryConfig`) | | | `documents` | Format sniffing and conversion to markdown | | | `documents-office` | PDF, DOCX, PPTX, XLSX conversion (implies `documents`) | -| | `sources` | Readers for folders, files, conversations; Composio normalisers (implies `documents`) | +| | `sources` | Readers for folders, files, conversations (implies `documents`) | | | `sources-network` | GitHub, RSS and web-page readers behind the SSRF guard (implies `sources`) | | | `safety` | Secret and PII scrubbing of a `StoreItem` | | | `legacy-import` | Reading a v1 workspace; `import::migrate` and `migrate_with` | @@ -72,7 +72,7 @@ source reader ──▶ documents ──▶ safety ──▶ engine.store ── ``` 1. **Source reader** (`sources`): lists a configured source (folder, file, - link, GitHub, RSS, Composio payload, conversation) and reads each entry. + link, GitHub, RSS, conversation) and reads each entry. `sources::collect_items` does this for one source; one bad entry lands in `Collected::skipped` instead of aborting the pass. 2. **Documents conversion** (`documents`): sniffs the format and converts the diff --git a/docs/specs/memory-v2.md b/docs/specs/memory-v2.md index 90a127f2..9bc19b80 100644 --- a/docs/specs/memory-v2.md +++ b/docs/specs/memory-v2.md @@ -30,7 +30,7 @@ Three crates, one directory each under `crates/`. Their relationships are in | --- | --- | | `tinymemory-api` | The contract: `MemoryEngine`, request/response types, `MemoryMeta`, `MetaFilter`, namespaces, `EngineDescriptor`, `Error`. No I/O. Feature `conformance`: the behavioural suite every engine must pass, plus a reference in-memory engine. | | `tinymemory-tools` | The agent tool spec `MemoryTools` (seven tools, host-fixed namespace and reach) and `context`, the `ContextCompiler` that builds `context.md` from an engine. | -| `tinymemory-integrations` | Everything that talks to the outside world, each behind a feature: `cortex` (default; the CortexDB engine, registered twice as `cortexdb` and `tinyhumans`, plus the registry `list_engines`/`build_engine` and `MemoryConfig`), `documents` and `documents-office` (format sniffing and conversion to markdown, emitting `StoreItem::Document`; PDF/DOCX/PPTX/XLSX via `OfficeConverter` or a host `DocumentConverter`), `sources` and `sources-network` (readers for folder, file, link, GitHub, RSS, Composio payloads and conversations, with the SSRF guard), `safety` (secret/PII scrubbing applied to every item before `store`), and `legacy-import` (reads a v1 TinyCortex workspace and migrates it into any engine). | +| `tinymemory-integrations` | Everything that talks to the outside world, each behind a feature: `cortex` (default; the CortexDB engine, registered twice as `cortexdb` and `tinyhumans`, plus the registry `list_engines`/`build_engine` and `MemoryConfig`), `documents` and `documents-office` (format sniffing and conversion to markdown, emitting `StoreItem::Document`; PDF/DOCX/PPTX/XLSX via `OfficeConverter` or a host `DocumentConverter`), `sources` and `sources-network` (readers for folder, file, link, GitHub, RSS and conversations, with the SSRF guard), `safety` (secret/PII scrubbing applied to every item before `store`), and `legacy-import` (reads a v1 TinyCortex workspace and migrates it into any engine). | Deleted: the earlier `tinymemory` facade, `tinymemory-cortex`, `tinymemory-documents`, `tinymemory-sources`, `tinymemory-safety`, @@ -91,7 +91,7 @@ pub struct MemoryMeta { pub tags: Vec, pub observed_at: Option>, } -pub enum SourceKind { Folder, File, Link, Github, Rss, Composio, Conversation, Agent, Import } +pub enum SourceKind { Folder, File, Link, Github, Rss, Conversation, Agent, Import } ``` `MetaFilter` has the same optional fields (each an exact match, `folder` and