Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 8 additions & 0 deletions crates/openhuman-core/src/config/schema/memory.rs
Original file line number Diff line number Diff line change
Expand Up @@ -148,6 +148,13 @@ pub struct MemoryConfig {
/// Written only when off, since on is the default.
#[serde(skip_serializing_if = "is_true")]
pub split_github_by_repo: bool,
/// Attribute what memory stores to who said or did it (CortexDB's
/// `observed_actor`): an assistant turn to its agent, a synced email to
Comment thread
M3gA-Mind marked this conversation as resolved.
/// its sender. Off by default, and off nothing on the wire changes. Only
Comment thread
coderabbitai[bot] marked this conversation as resolved.
/// the `cortexdb` engine honours it, and a write CortexDB refuses for it
Comment thread
M3gA-Mind marked this conversation as resolved.
/// is written again without it.
Comment thread
M3gA-Mind marked this conversation as resolved.
#[serde(skip_serializing_if = "std::ops::Not::not")]
pub observed_actor: bool,
}

/// `[memory] layout`: where memory sits on the engine.
Expand Down Expand Up @@ -218,6 +225,7 @@ impl Default for MemoryConfig {
embedding_rate_limit_per_min: DEFAULT_EMBEDDING_RATE_LIMIT_PER_MIN,
agents: BTreeMap::new(),
split_github_by_repo: true,
observed_actor: false,
}
}
}
Expand Down
14 changes: 14 additions & 0 deletions crates/openhuman-core/src/config/schema/memory_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -255,3 +255,17 @@ fn migration_drops_unmappable_disabled_and_incomplete_entries() {
assert_eq!(migrate_legacy_source(&legacy), None, "{legacy}");
}
}

#[test]
fn observed_actor_is_off_by_default_and_unwritten_until_set() {
Comment thread
M3gA-Mind marked this conversation as resolved.
let config = MemoryConfig::default();
assert!(!config.observed_actor);
let written = serde_json::to_value(&config).unwrap();
assert!(
written.get("observed_actor").is_none(),
"off is not written"
);

let on: MemoryConfig = toml::from_str("observed_actor = true").unwrap();
assert!(on.observed_actor);
}
7 changes: 6 additions & 1 deletion crates/openhuman-core/src/memory/engine.rs
Original file line number Diff line number Diff line change
Expand Up @@ -409,11 +409,16 @@ fn resolve_cortexdb(config: &Config, root: Option<&str>) -> Binding {
);
}
};
let fingerprint = format!("{CORTEXDB_ENGINE}|{endpoint}|{}", key_digest(&key));
let observed_actor = config.memory.observed_actor;
Comment thread
M3gA-Mind marked this conversation as resolved.
let fingerprint = format!(
"{CORTEXDB_ENGINE}|{endpoint}|{}|actor={observed_actor}",
key_digest(&key)
);
// A third-party endpoint: no TinyHumans attribution headers.
let (settings, layout) = rooted(
EngineSettings {
endpoint: Some(endpoint),
observed_actor,
Comment thread
M3gA-Mind marked this conversation as resolved.
..EngineSettings::default()
},
root,
Expand Down
23 changes: 23 additions & 0 deletions crates/openhuman-core/src/memory/engine_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -133,6 +133,29 @@ fn cortexdb_is_off_until_a_key_is_stored_and_rebuilds_on_a_new_key() {
assert!(!is_on(&config));
}

#[test]
fn switching_observed_actor_rebuilds_the_cortexdb_engine() {
let tmp = tempfile::tempdir().unwrap();
let mut config = config_in(&tmp);
config.memory.engine = CORTEXDB_ENGINE.to_string();
store_cortexdb_key(&config, "cdb-key-actor").unwrap();
let off = resolve(&config).engine().expect("bound");
assert!(
Arc::ptr_eq(
&off.engine,
&resolve(&config).engine().expect("bound").engine
),
"unchanged config reuses the engine"
);

config.memory.observed_actor = true;
let on = resolve(&config).engine().expect("rebound");
assert!(
!Arc::ptr_eq(&off.engine, &on.engine),
"the setting is part of the cache fingerprint"
);
}

#[test]
fn a_blank_cortexdb_key_is_refused() {
let tmp = tempfile::tempdir().unwrap();
Expand Down
27 changes: 24 additions & 3 deletions crates/openhuman-core/src/memory/sources/composio.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5,13 +5,15 @@
//! (`integrations::composio::ops::run_sync_pass`), which pages the account and
//! hands back decoded records; each record is stored as a document with
//! `source = composio:<source id>`, `tags = [toolkit, "connection:<id>", …]`,
//! its URL and its upstream timestamp. The connection tag is what deleting a
//! its URL, its upstream timestamp and, when the connector names one (an
//! email's sender), its sender as `meta.observed_actor` ([`sender_actor`]).
//! The connection tag is what deleting a
//! connection with `clear_memory` forgets by. The same conversion backs `openhuman.composio_sync`,
//! which syncs one connection on demand.

use chrono::{TimeZone, Utc};
use tinyconnectors_bus::records::ConnectorRecord;
use tinymemory_api::{DocumentBody, MemoryMeta, SourceKind, SourceRef, StoreItem};
use tinyconnectors_bus::records::{ConnectorRecord, RecordSender};
use tinymemory_api::{DocumentBody, MemoryMeta, ObservedActor, SourceKind, SourceRef, StoreItem};

use crate::config::schema::MemorySourceConfig;
use crate::config::Config;
Expand Down Expand Up @@ -97,11 +99,30 @@ pub fn record_item(
kind: SourceKind::Composio,
id: Some(source_id.to_string()),
},
observed_actor: record.sender.as_ref().and_then(sender_actor),
Comment thread
M3gA-Mind marked this conversation as resolved.
..MemoryMeta::default()
},
})
}

/// Who a record came from, as the memory actor `user:<address>` with their
/// name: an email address, lower-cased so one person is one actor, or a
/// phone number as the connector's dial digits. Stored as given (the engine
/// sends it only with `[memory] observed_actor` on); `None` for a blank
/// address.
fn sender_actor(sender: &RecordSender) -> Option<ObservedActor> {
let address = sender.address.trim();
(!address.is_empty()).then(|| ObservedActor {
id: format!("user:{}", address.to_ascii_lowercase()),
name: sender
.name
.as_deref()
.map(str::trim)
.filter(|name| !name.is_empty())
.map(str::to_string),
})
}

/// Files every non-empty record into `layout`'s brain, under the toolkit's
/// brain source, and forgets the previous version of each record that
/// changed or came back empty upstream ([`super::versions`]). Returns how
Expand Down
37 changes: 37 additions & 0 deletions crates/openhuman-core/src/memory/sources/composio_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -507,3 +507,40 @@ async fn a_deleted_connection_is_not_read_again() {
assert!(!pass.more_pending);
assert!(pass.failure.is_none());
}

#[test]
fn a_records_sender_becomes_its_observed_actor() {
let actor_of = |address: &str, name: Option<&str>| {
let mut rec = record("m-1", "Lunch", "see you at noon");
rec.sender = Some(RecordSender {
address: address.into(),
name: name.map(str::to_string),
});
record_item("gmail", "conn-7", "src-g", &rec)
.expect("an item")
.meta()
.observed_actor
.clone()
};
assert_eq!(
actor_of(" Priya@Acme.com ", Some(" Priya ")),
Some(ObservedActor {
id: "user:priya@acme.com".into(),
name: Some("Priya".into()),
}),
"an email address, lower-cased, with the name"
);
assert_eq!(
actor_of("+15551234567", Some(" ")),
Some(ObservedActor {
id: "user:+15551234567".into(),
name: None,
}),
"a phone number as given; a blank name is none"
);
assert_eq!(actor_of(" ", Some("Priya")), None, "no address, no actor");
Comment thread
M3gA-Mind marked this conversation as resolved.

let unsent =
record_item("gmail", "conn-7", "src-g", &record("m-2", "", "no sender")).expect("an item");
assert_eq!(unsent.meta().observed_actor, None);
}
9 changes: 9 additions & 0 deletions docs/specs/memory-v2.md
Original file line number Diff line number Diff line change
Expand Up @@ -146,6 +146,7 @@ its own schedule (`scheduled`); an engine that cannot consolidate answers
engine = "tinyhumans" # "tinyhumans" | "cortexdb"
# agent_id = "employee-7" # host binding (see "Who is acting")
# root = "team:acme"
# observed_actor = false # attribute writes to who said or did them (cortexdb only)

[memory.engines.cortexdb]
endpoint = "https://api-v1.cortexdb.ai" # key in the keychain as "memory-cortexdb"
Expand Down Expand Up @@ -180,6 +181,14 @@ The hosted engine sends the installed transport's attribution headers
(`x-sdk-name`, …) on every request; the CortexDB engine, a third party, does
not.

`observed_actor` (off by default) makes the CortexDB engine name who said or
did what it stores (CortexDB's `observed_actor`, with the memory's owner as
`subject`): an assistant turn its agent (`agent:<id>`), a synced email its
sender (`user:<address>`, with the sender's name). An email address is
trimmed and lower-cased, so one person is one actor; a phone number is kept as
the connector's dial digits (`+15551234567`). Neither is redacted. A write
CortexDB refuses for it is written again without it (tinymemory's fallback). Off, and always on the hosted engine, nothing on the wire changes.

## Agent tool: `memory`

One tool, `action` = `recall` | `fetch` | `learn` | `forget`:
Expand Down
Loading