diff --git a/.github/workspace-dep-closures.json b/.github/workspace-dep-closures.json index 724e0e73025..bc5b8f6d06c 100644 --- a/.github/workspace-dep-closures.json +++ b/.github/workspace-dep-closures.json @@ -1,6 +1,20 @@ { "//": "Generated by `cargo run -p xtask -- deps` (see tooling/xtask); do not edit by hand. Maps every Rust workspace crate to the directories of its transitive workspace dependency closure (normal + build deps, dev-deps excluded). Consumed by flake.nix to build pruned per-artifact deploy sources.", "closures": { + "activity": [ + "crates/activity", + "crates/bot_id", + "crates/channel_sender", + "crates/cowlike", + "crates/kafka_util", + "crates/macro_env", + "crates/macro_env_var", + "crates/macro_event_broker", + "crates/macro_event_topics", + "crates/macro_user_id", + "crates/model-entity", + "crates/workspace-hack" + ], "agent": [ "crates/agent", "crates/ai_toolset", @@ -27,6 +41,7 @@ "crates/agent_runtime_protocol" ], "ai_projections": [ + "crates/activity", "crates/agent", "crates/ai_projections", "crates/ai_tools", @@ -138,6 +153,7 @@ "crates/workspace-hack" ], "ai_projections_refresh_handler": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -202,6 +218,7 @@ "services/ai_projections_refresh_handler" ], "ai_tools": [ + "crates/activity", "crates/agent", "crates/ai_tools", "crates/ai_toolset", @@ -354,6 +371,7 @@ "crates/workspace-hack" ], "authentication_service": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -506,6 +524,7 @@ "crates/workspace-hack" ], "bots": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -603,6 +622,7 @@ "crates/workspace-hack" ], "call": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -665,6 +685,7 @@ "services/call_recording_preview_handler" ], "channel_bots": [ + "crates/activity", "crates/agent", "crates/ai_tools", "crates/ai_toolset", @@ -783,6 +804,7 @@ "crates/workspace-hack" ], "channels": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -834,11 +856,13 @@ "crates/workspace-hack" ], "chat": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", "crates/attachment", "crates/bot_id", + "crates/channel_sender", "crates/chat", "crates/cowlike", "crates/document_sub_type", @@ -871,6 +895,7 @@ "crates/workspace-hack" ], "comms_db_client": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -923,6 +948,7 @@ "crates/workspace-hack" ], "complete_graph": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -1171,6 +1197,7 @@ "services/contacts_service" ], "convert_service": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -1252,6 +1279,7 @@ "crates/workspace-hack" ], "crm": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -1333,6 +1361,7 @@ "services/delete_chat_handler" ], "deleted_item_poller": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -1430,6 +1459,7 @@ "services/deleted_item_poller" ], "document_cognition_service": [ + "crates/activity", "crates/agent", "crates/ai_projections", "crates/ai_tools", @@ -1555,6 +1585,7 @@ "services/unfurl_service" ], "document_storage_service": [ + "crates/activity", "crates/agent", "crates/ai_tools", "crates/ai_toolset", @@ -1761,6 +1792,7 @@ "services/document_text_extractor" ], "document_upload_finalizer_handler": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -1845,6 +1877,7 @@ "services/document_upload_finalizer_handler" ], "documents": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -1928,6 +1961,7 @@ "crates/workspace-hack" ], "docx_unzip_handler": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -2006,6 +2040,7 @@ "crates/workspace-hack" ], "email": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -2097,6 +2132,7 @@ "crates/workspace-hack" ], "email_refresh_handler": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -2159,6 +2195,7 @@ "services/email_refresh_handler" ], "email_scheduled_handler": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -2221,6 +2258,7 @@ "services/email_scheduled_handler" ], "email_service": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -2305,6 +2343,7 @@ "services/email_service" ], "email_service_client": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -2361,6 +2400,7 @@ "crates/workspace-hack" ], "email_sfs_delete_handler": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -2641,6 +2681,7 @@ "crates/workspace-hack" ], "github": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -2736,6 +2777,7 @@ "crates/workspace-hack" ], "graphql_channel": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -2811,6 +2853,7 @@ "crates/workspace-hack" ], "graphql_email": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -2867,6 +2910,7 @@ "crates/workspace-hack" ], "graphql_entity_mutation": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -3056,6 +3100,7 @@ "crates/workspace-hack" ], "graphql_properties": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -3104,6 +3149,7 @@ "crates/workspace-hack" ], "graphql_soup": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -3257,6 +3303,7 @@ "services/image_proxy_service" ], "import": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -3357,6 +3404,7 @@ "crates/workspace-hack" ], "lexical_client": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -3428,6 +3476,7 @@ "crates/workspace-hack" ], "lexical_mention_extractor": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -3778,6 +3827,7 @@ "crates/workspace-hack" ], "mcp_service": [ + "crates/activity", "crates/agent", "crates/ai_tools", "crates/ai_toolset", @@ -3894,6 +3944,7 @@ "services/mcp_service" ], "memory": [ + "crates/activity", "crates/agent", "crates/ai_tools", "crates/ai_toolset", @@ -4136,6 +4187,7 @@ "crates/workspace-hack" ], "models_search": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -4214,6 +4266,7 @@ "crates/workspace-hack" ], "models_soup": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -4283,6 +4336,7 @@ "crates/workspace-hack" ], "name_search": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -4486,6 +4540,7 @@ "services/notification_service" ], "onboarding": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -4547,6 +4602,7 @@ "crates/workspace-hack" ], "organization_retention_handler": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -4609,6 +4665,7 @@ "services/organization_retention_handler" ], "organization_retention_trigger": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -4671,6 +4728,7 @@ "services/organization_retention_trigger" ], "projects": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -4764,6 +4822,7 @@ "crates/workspace-hack" ], "properties": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -4808,6 +4867,7 @@ "crates/workspace-hack" ], "properties_service": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -4948,6 +5008,7 @@ "crates/workspace-hack" ], "scheduled_action": [ + "crates/activity", "crates/agent", "crates/ai_tools", "crates/ai_toolset", @@ -5064,6 +5125,7 @@ "services/scheduled_action" ], "search_processing_service": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -5170,6 +5232,7 @@ "services/search_processing_service" ], "search_service": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -5251,6 +5314,7 @@ "crates/workspace-hack" ], "search_service_client": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -5361,6 +5425,7 @@ "crates/workspace-hack" ], "seed_cli": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -5475,6 +5540,7 @@ "crates/workspace-hack" ], "skills": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -5556,6 +5622,7 @@ "crates/workspace-hack" ], "soup": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -5625,6 +5692,7 @@ "crates/workspace-hack" ], "soup_realtime": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -5718,6 +5786,7 @@ "crates/workspace-hack" ], "sqs_client": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -5905,6 +5974,7 @@ "crates/workspace-hack" ], "task_dedup": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -5980,6 +6050,7 @@ "crates/workspace-hack" ], "teams": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -6064,6 +6135,7 @@ "services/unfurl_service" ], "upload_extractor_lambda_handler": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -6129,6 +6201,7 @@ "services/upload_extractor_lambda_handler" ], "upload_extractor_lambda_trigger": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", @@ -6234,6 +6307,7 @@ "crates/workspace-hack" ], "webhook": [ + "crates/activity", "crates/agent", "crates/ai_toolset", "crates/ai_usage", diff --git a/.sqlx/query-29b8ca0e672d412aeeaf19c3f633fb10fb5b41ac6a0ff074eefcc71123353872.json b/.sqlx/query-29b8ca0e672d412aeeaf19c3f633fb10fb5b41ac6a0ff074eefcc71123353872.json new file mode 100644 index 00000000000..c0f966a2253 --- /dev/null +++ b/.sqlx/query-29b8ca0e672d412aeeaf19c3f633fb10fb5b41ac6a0ff074eefcc71123353872.json @@ -0,0 +1,21 @@ +{ + "db_name": "PostgreSQL", + "query": "\n INSERT INTO activity_events\n (id, actor_id, subject_id, action, action_payload,\n entity_type, entity_id, occurred_at)\n SELECT * FROM UNNEST(\n $1::uuid[], $2::text[], $3::text[], $4::text[], $5::jsonb[],\n $6::text[], $7::text[], $8::timestamptz[])\n ON CONFLICT (id) DO NOTHING\n ", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "UuidArray", + "TextArray", + "TextArray", + "TextArray", + "JsonbArray", + "TextArray", + "TextArray", + "TimestamptzArray" + ] + }, + "nullable": [] + }, + "hash": "29b8ca0e672d412aeeaf19c3f633fb10fb5b41ac6a0ff074eefcc71123353872" +} diff --git a/.sqlx/query-6d61d631a10b803cb269fa6992ef84adbf8b0ad0940ab61b8a649f74584d8b0e.json b/.sqlx/query-6d61d631a10b803cb269fa6992ef84adbf8b0ad0940ab61b8a649f74584d8b0e.json new file mode 100644 index 00000000000..53abbe286bd --- /dev/null +++ b/.sqlx/query-6d61d631a10b803cb269fa6992ef84adbf8b0ad0940ab61b8a649f74584d8b0e.json @@ -0,0 +1,15 @@ +{ + "db_name": "PostgreSQL", + "query": "\n DELETE FROM activity_events\n WHERE (entity_type, entity_id) IN\n (SELECT * FROM UNNEST($1::text[], $2::text[]))\n ", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "TextArray", + "TextArray" + ] + }, + "nullable": [] + }, + "hash": "6d61d631a10b803cb269fa6992ef84adbf8b0ad0940ab61b8a649f74584d8b0e" +} diff --git a/Cargo.lock b/Cargo.lock index 47fbeb27fab..41cdccc38cb 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -8,6 +8,30 @@ version = "0.11.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "fe438c63458706e03479442743baae6c88256498e6431708f6dfc520a26515d3" +[[package]] +name = "activity" +version = "0.1.0" +dependencies = [ + "channel_sender", + "chrono", + "kafka_util", + "macro_db_migrator", + "macro_event_broker", + "macro_user_id", + "model-entity", + "rdkafka", + "readonly", + "rootcause", + "serde", + "serde_json", + "sqlx", + "strum 0.27.2", + "tokio", + "tracing", + "uuid", + "workspace-hack", +] + [[package]] name = "adler2" version = "2.0.1" @@ -2545,6 +2569,7 @@ dependencies = [ name = "call" version = "0.1.0" dependencies = [ + "activity", "agent", "ai_toolset", "ai_usage", @@ -2763,6 +2788,7 @@ dependencies = [ name = "channels" version = "0.1.0" dependencies = [ + "activity", "ai_toolset", "anyhow", "async-trait", @@ -2830,6 +2856,7 @@ dependencies = [ name = "chat" version = "0.1.0" dependencies = [ + "activity", "agent", "ai_toolset", "anyhow", @@ -4408,6 +4435,7 @@ dependencies = [ name = "document_storage_service" version = "0.1.0" dependencies = [ + "activity", "ai_tools", "ai_usage", "analytics_client", @@ -4629,6 +4657,7 @@ dependencies = [ name = "documents" version = "0.1.0" dependencies = [ + "activity", "ai_toolset", "ai_usage", "anyhow", @@ -4930,6 +4959,7 @@ dependencies = [ name = "email" version = "0.1.0" dependencies = [ + "activity", "ai_toolset", "anyhow", "async-trait", @@ -11460,6 +11490,7 @@ dependencies = [ name = "projects" version = "0.1.0" dependencies = [ + "activity", "ai_toolset", "anyhow", "async-recursion", @@ -11529,6 +11560,7 @@ dependencies = [ name = "properties" version = "0.1.0" dependencies = [ + "activity", "ai_toolset", "anyhow", "async-trait", @@ -12042,6 +12074,17 @@ dependencies = [ "pkg-config", ] +[[package]] +name = "readonly" +version = "0.2.14" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "67a77cd6f1a55ff3cb1969c2618f73107efaa6de1daf796ad60ce8f737ae97c1" +dependencies = [ + "proc-macro2", + "quote", + "syn 3.0.3", +] + [[package]] name = "readonly_pool" version = "0.1.0" @@ -14518,6 +14561,17 @@ dependencies = [ "unicode-ident", ] +[[package]] +name = "syn" +version = "3.0.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "53e9bae58849f64dfa4f5d5ae372c8341f7305f82a3868709269343628b659a3" +dependencies = [ + "proc-macro2", + "quote", + "unicode-ident", +] + [[package]] name = "sync_service" version = "0.1.0" diff --git a/Cargo.toml b/Cargo.toml index 7887dca8fbf..cd9f4859386 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -2,6 +2,7 @@ resolver = "2" members = [ + "crates/activity", "crates/agent", "crates/agent_runtime_protocol", "crates/ai_projections", @@ -208,6 +209,7 @@ rdkafka = { version = "0.37", default-features = false, features = [ "ssl-vendored", "tokio", ] } +readonly = "0.2" recursion = "0.5.4" redis = "1.0.3" regex = "1.11.1" diff --git a/crates/activity/Cargo.toml b/crates/activity/Cargo.toml new file mode 100644 index 00000000000..70a0d3b97d6 --- /dev/null +++ b/crates/activity/Cargo.toml @@ -0,0 +1,41 @@ +[package] +edition = "2024" +name = "activity" +publish = false +version = "0.1.0" + +[features] +consumer = [ + "dep:kafka_util", + "dep:macro_event_broker", + "dep:rdkafka", + "dep:rootcause", + "dep:tokio", + "dep:tracing", + "macro_event_broker/outbound", +] +default = ["consumer", "outbound"] +outbound = ["dep:sqlx"] + +[dependencies] +channel_sender = { path = "../channel_sender" } +chrono = { workspace = true } +kafka_util = { path = "../kafka_util", optional = true } +macro_event_broker = { path = "../macro_event_broker", default-features = false, optional = true } +macro_user_id = { path = "../macro_user_id" } +model-entity = { path = "../model-entity" } +rdkafka = { workspace = true, optional = true } +readonly = { workspace = true } +rootcause = { workspace = true, optional = true } +serde = { workspace = true } +serde_json = { workspace = true } +sqlx = { workspace = true, features = ["json"], optional = true } +strum = { workspace = true } +tokio = { workspace = true, features = ["sync", "time"], optional = true } +tracing = { workspace = true, optional = true } +uuid = { workspace = true } +workspace-hack = { version = "0.1", path = "../workspace-hack" } + +[dev-dependencies] +macro_db_migrator = { path = "../macro_db_migrator" } +tokio = { workspace = true, features = ["macros", "rt-multi-thread", "test-util"] } diff --git a/crates/activity/src/domain/mod.rs b/crates/activity/src/domain/mod.rs new file mode 100644 index 00000000000..0cd0a237ec8 --- /dev/null +++ b/crates/activity/src/domain/mod.rs @@ -0,0 +1,4 @@ +//! Domain layer: the activity model and the storage port. + +pub mod models; +pub mod ports; diff --git a/crates/activity/src/domain/models.rs b/crates/activity/src/domain/models.rs new file mode 100644 index 00000000000..b873ced6072 --- /dev/null +++ b/crates/activity/src/domain/models.rs @@ -0,0 +1,345 @@ +//! The activity model — the protocol domains implement and the storage +//! shape the consumer persists. +//! +//! Ownership boundary: **this crate owns what an activity is; domains own +//! what counts as activity in their domain.** Domains construct activities +//! through two doors: +//! +//! - [`Activity::common`] for [`CommonAction`]s, which are valid on every +//! entity kind by definition — no pairing proof needed. +//! - [`Activity::from_domain`] for entity-exclusive actions, reachable only +//! through a domain-owned type implementing [`DomainActivity`]. The domain +//! crate is the authority on its own action vocabulary and its projection +//! into the durable [`Action`]. +//! +//! [`Activity`] itself cannot be literal-constructed (private field), so an +//! invalid pairing like a document with a `Messaged` action has no +//! expressible path into storage. + +#[cfg(test)] +mod test; + +use std::sync::LazyLock; + +use chrono::{DateTime, Utc}; +use macro_user_id::user_id::MacroUserIdStr; +use serde::{Deserialize, Serialize}; +use serde_json::Value; +use uuid::Uuid; + +/// A principal that can act on entities: a user or a bot, as one +/// prefix-parseable string. An agent is a bot acting with `on_behalf_of` +/// set — agent-ness is relational, not an identity kind. +pub use channel_sender::ChannelSender as Actor; +/// The entity-kind vocabulary activities are recorded against (re-exported +/// for domain mappings). +pub use model_entity::EntityType; + +/// Namespace for deriving deterministic activity ids from source event ids. +static ACTIVITY_ID_NAMESPACE: LazyLock = + LazyLock::new(|| Uuid::new_v5(&Uuid::NAMESPACE_OID, b"macro.activity_events")); + +/// Actions valid on every entity kind. +#[derive(Debug, Clone, PartialEq)] +pub enum CommonAction { + /// The entity was created. + Created, + /// The entity's content or metadata was edited. + Edited, + /// The entity was opened by its subject. + Opened, + /// The entity was soft-deleted. + Deleted, + /// A property value changed on the entity. `property` is the property + /// definition id; `from` is unset until the source event carries the + /// previous value; `to` is `None` when the value was cleared. + PropertyChanged(PropertyChange), +} + +/// Payload of [`Action::PropertyChanged`]. Serde-derived so the write and +/// (future) read codecs share one shape definition. +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +pub struct PropertyChange { + /// Property definition id. + pub property: String, + /// Previous value, when known. + pub from: Option, + /// New value; `None` when cleared. + pub to: Option, +} + +/// Payload of [`Action::ParticipantAdded`] / [`Action::ParticipantRemoved`]. +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +pub struct ParticipantChange { + /// The (un)added principal. + pub participant: Actor<'static>, +} + +/// Payload of [`Action::CallStarted`]. +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +pub struct CallStart { + /// The started call. + pub call_id: String, +} + +/// The durable action vocabulary — what the `action`/`action_payload` +/// columns hold. [`Action::to_columns`] is the storage codec, written out +/// explicitly; the future read path adds the inverse next to it. +/// +/// This union is deliberately the one place every domain's vocabulary +/// meets: storage needs a single closed vocabulary, and reads, retention +/// policy, and the API surface stay exhaustively checkable because of it. +/// Entity-exclusive variants are reachable only through a domain's +/// [`DomainActivity`] projection. +/// +/// Tags, variants, and payload fields are never renamed or repurposed — +/// stored activities are immutable and must decode forever. New payload +/// fields must tolerate absence on old rows. The stored tag is derived +/// from the variant name (strum, snake_case), so **renaming a variant is a +/// storage migration** — the pinned codec test exists to make that loud. +#[derive(Debug, Clone, PartialEq, strum::IntoStaticStr)] +#[strum(serialize_all = "snake_case")] +pub enum Action { + /// The entity was created. + Created, + /// The entity's content or metadata was edited. + Edited, + /// The entity was opened by its subject. + Opened, + /// The entity was soft-deleted. + Deleted, + /// A message was sent in the entity (channel or chat). + Messaged, + /// An email message was sent on the thread. + Sent, + /// A property value changed on the entity (see + /// [`CommonAction::PropertyChanged`]). + PropertyChanged(PropertyChange), + /// A principal was added to the entity (channel membership). + ParticipantAdded(ParticipantChange), + /// A principal was removed from the entity (channel membership). + ParticipantRemoved(ParticipantChange), + /// A call was started in the entity (channel). + CallStarted(CallStart), +} + +impl From for Action { + fn from(action: CommonAction) -> Self { + match action { + CommonAction::Created => Action::Created, + CommonAction::Edited => Action::Edited, + CommonAction::Opened => Action::Opened, + CommonAction::Deleted => Action::Deleted, + CommonAction::PropertyChanged(change) => Action::PropertyChanged(change), + } + } +} + +impl Action { + /// Whether this action is a view rather than a mutation — the only + /// classification activity queries need. + pub fn is_view(&self) -> bool { + matches!(self, Action::Opened) + } + + /// Splits the action into its `(action, action_payload)` column values: + /// the tag from the variant name (strum), the payload from the shared + /// serde structs. Exhaustive, so a new variant fails compilation until + /// its payload is decided. + pub fn to_columns(&self) -> (&'static str, Option) { + // Payload structs are plain data (strings, options, JSON values): + // serialization cannot fail. + fn payload(payload: &T) -> Option { + Some(serde_json::to_value(payload).expect("payload structs serialize infallibly")) + } + + let tag: &'static str = self.into(); + let payload = match self { + Action::Created + | Action::Edited + | Action::Opened + | Action::Deleted + | Action::Messaged + | Action::Sent => None, + Action::PropertyChanged(change) => payload(change), + Action::ParticipantAdded(change) | Action::ParticipantRemoved(change) => { + payload(change) + } + Action::CallStarted(start) => payload(start), + }; + (tag, payload) + } +} + +/// The capability a domain implements to feed activity: given one of its +/// broker events, what activity happened? +/// +/// Implemented on the domain's topic-event enum, next to it, by the crate +/// that owns it — the domain is the authority on what counts as activity +/// in its domain. Implementations match exhaustively (no `_` arm) so a new +/// event variant fails compilation until classified or explicitly dropped. +/// +/// Takes the broker `event_id` rather than the envelope so this crate's +/// model surface stays free of broker dependencies; the id feeds +/// deterministic activity ids ([`activity_id`]) and the [`event_time`] +/// fallback. If the team later wants every topic event to declare its +/// activity semantics, `TopicEvent: ActivitySource` is the one-line +/// ratchet. +pub trait ActivitySource { + /// Classifies one event into its ingest outcome. + fn ingest(&self, event_id: Uuid) -> Ingest; +} + +/// The contract a domain implements to contribute entity-exclusive +/// activity. The implementing crate is the authority on which actions its +/// entity kind supports and how they project into the durable [`Action`]. +pub trait DomainActivity { + /// The entity kind this activity applies to. + const ENTITY_TYPE: EntityType; + + /// The entity acted on. + fn entity_id(&self) -> &str; + + /// Total projection into the durable vocabulary. + fn into_action(self) -> Action; +} + +/// What ingesting one broker event asks of storage. +#[derive(Debug, Clone, PartialEq)] +pub enum Ingest { + /// Insert these activities. + Insert(Vec), + /// These entities were hard-deleted; purge their activities. + Purge(Vec<(EntityType, String)>), + /// The event carries no activity: pipeline noise, or a mutation with no + /// actor (unattributable activities are dropped until a system principal + /// exists). + Ignore, +} + +/// One activity: a principal did something to an entity at a time. +/// +/// Flat — this is (nearly) the persisted row. Its invariants (valid +/// action-for-entity pairing, `subject = on_behalf_of ?? actor`, +/// deterministic id) are established at construction and sealed by +/// `#[readonly::make]`: outside this module the fields read with normal +/// syntax but cannot be written, and literal construction is impossible — +/// [`Activity::common`] and [`Activity::from_domain`] are the only doors. +#[readonly::make] +#[derive(Debug, Clone, PartialEq)] +pub struct Activity { + /// Deterministic id: uuidv5 over (source event id, ordinal), so replays + /// of the same broker event re-derive the same activity ids. + pub id: Uuid, + /// Who mechanically acted. + pub actor: Actor<'static>, + /// Whose activity this is: `on_behalf_of ?? actor`, resolved here at + /// construction and never re-derived downstream. + pub subject_id: String, + /// The kind of entity acted on, in the soup item-type vocabulary. + pub entity_type: EntityType, + /// The entity acted on. + pub entity_id: String, + /// What they did. + pub action: Action, + /// When it happened, per the source event. + pub occurred_at: DateTime, +} + +impl Activity { + /// Builds an activity for a [`CommonAction`] — valid on every entity + /// kind by definition, so any (kind, id) pairing is accepted. + #[allow(clippy::too_many_arguments)] + pub fn common( + source_event_id: Uuid, + ordinal: u32, + actor: Actor<'static>, + on_behalf_of: Option>, + entity_type: EntityType, + entity_id: impl Into, + action: CommonAction, + occurred_at: DateTime, + ) -> Self { + Self::build( + source_event_id, + ordinal, + actor, + on_behalf_of, + entity_type, + entity_id.into(), + action.into(), + occurred_at, + ) + } + + /// Builds an activity for an entity-exclusive action, reachable only + /// through a domain-owned [`DomainActivity`] type. + pub fn from_domain( + source_event_id: Uuid, + ordinal: u32, + actor: Actor<'static>, + on_behalf_of: Option>, + entity: A, + occurred_at: DateTime, + ) -> Self { + let entity_id = entity.entity_id().to_owned(); + Self::build( + source_event_id, + ordinal, + actor, + on_behalf_of, + A::ENTITY_TYPE, + entity_id, + entity.into_action(), + occurred_at, + ) + } + + #[allow(clippy::too_many_arguments)] + fn build( + source_event_id: Uuid, + ordinal: u32, + actor: Actor<'static>, + on_behalf_of: Option>, + entity_type: EntityType, + entity_id: String, + action: Action, + occurred_at: DateTime, + ) -> Self { + let subject_id = on_behalf_of + .map(|user| user.as_ref().to_owned()) + .unwrap_or_else(|| actor.as_ref().to_owned()); + Self { + id: activity_id(source_event_id, ordinal), + actor, + subject_id, + entity_type, + entity_id, + action, + occurred_at, + } + } +} + +/// Derives the deterministic id for the `ordinal`-th activity produced by the +/// broker event `source_event_id`. +pub fn activity_id(source_event_id: Uuid, ordinal: u32) -> Uuid { + Uuid::new_v5( + &ACTIVITY_ID_NAMESPACE, + format!("{source_event_id}:{ordinal}").as_bytes(), + ) +} + +/// `occurred_at` fallback for events whose metadata carries no timestamp: +/// broker event ids are uuidv7, whose embedded time is when the source +/// domain published. Replays keep the first stored value regardless +/// (activity inserts are `ON CONFLICT DO NOTHING`). +pub fn event_time(event_id: Uuid) -> DateTime { + event_id + .get_timestamp() + .and_then(|ts| { + let (seconds, nanos) = ts.to_unix(); + DateTime::from_timestamp(i64::try_from(seconds).ok()?, nanos) + }) + .unwrap_or_else(Utc::now) +} diff --git a/crates/activity/src/domain/models/test.rs b/crates/activity/src/domain/models/test.rs new file mode 100644 index 00000000000..c41b0350755 --- /dev/null +++ b/crates/activity/src/domain/models/test.rs @@ -0,0 +1,133 @@ +use macro_user_id::user_id::MacroUserIdStr; +use serde_json::json; + +use super::*; + +fn user(id: &str) -> MacroUserIdStr<'static> { + MacroUserIdStr::try_from(id.to_string()).expect("valid user id") +} + +#[test] +fn every_action_maps_to_stable_columns() { + let participant = Actor::new_from_user(user("macro|sarah@example.com")); + // Every variant, pinned: these strings and payload shapes are the + // durable storage contract and must never change for existing tags. + let cases: Vec<(Action, &str, Option)> = vec![ + (Action::Created, "created", None), + (Action::Edited, "edited", None), + (Action::Opened, "opened", None), + (Action::Deleted, "deleted", None), + (Action::Messaged, "messaged", None), + (Action::Sent, "sent", None), + ( + Action::PropertyChanged(PropertyChange { + property: "prop-1".to_string(), + from: None, + to: Some(json!("Done")), + }), + "property_changed", + Some(json!({ "property": "prop-1", "from": null, "to": "Done" })), + ), + ( + Action::ParticipantAdded(ParticipantChange { + participant: participant.clone(), + }), + "participant_added", + Some(json!({ "participant": "macro|sarah@example.com" })), + ), + ( + Action::ParticipantRemoved(ParticipantChange { participant }), + "participant_removed", + Some(json!({ "participant": "macro|sarah@example.com" })), + ), + ( + Action::CallStarted(CallStart { + call_id: "call-1".to_string(), + }), + "call_started", + Some(json!({ "call_id": "call-1" })), + ), + ]; + + for (action, expected_tag, expected_payload) in cases { + let (tag, payload) = action.to_columns(); + assert_eq!(tag, expected_tag, "tag for {action:?}"); + assert_eq!(payload, expected_payload, "payload for {action:?}"); + } +} + +#[test] +fn only_opened_is_a_view() { + assert!(Action::Opened.is_view()); + assert!(!Action::Created.is_view()); + assert!(!Action::Edited.is_view()); + assert!(!Action::Deleted.is_view()); +} + +#[test] +fn activity_ids_are_deterministic_per_event_and_ordinal() { + let event_id = Uuid::from_u128(7); + + assert_eq!(activity_id(event_id, 0), activity_id(event_id, 0)); + assert_ne!(activity_id(event_id, 0), activity_id(event_id, 1)); + assert_ne!(activity_id(event_id, 0), activity_id(Uuid::from_u128(8), 0)); +} + +#[test] +fn subject_is_the_actor_unless_delegated() { + let direct = Activity::common( + Uuid::from_u128(1), + 0, + Actor::new_from_user(user("macro|teo@example.com")), + None, + EntityType::Document, + "doc-1", + CommonAction::Edited, + Utc::now(), + ); + assert_eq!(direct.subject_id, "macro|teo@example.com"); + assert_eq!(direct.actor.as_ref(), "macro|teo@example.com"); + + let delegated = Activity::common( + Uuid::from_u128(2), + 0, + Actor::new_from_user(user("macro|other@example.com")), + Some(user("macro|teo@example.com")), + EntityType::Document, + "doc-1", + CommonAction::Edited, + Utc::now(), + ); + assert_eq!(delegated.subject_id, "macro|teo@example.com"); + assert_eq!(delegated.actor.as_ref(), "macro|other@example.com"); +} + +#[test] +fn from_domain_projects_the_owning_kind_and_action() { + struct FixtureActivity { + id: String, + } + impl DomainActivity for FixtureActivity { + const ENTITY_TYPE: EntityType = EntityType::Channel; + fn entity_id(&self) -> &str { + &self.id + } + fn into_action(self) -> Action { + Action::Messaged + } + } + + let activity = Activity::from_domain( + Uuid::from_u128(3), + 0, + Actor::new_from_user(user("macro|teo@example.com")), + None, + FixtureActivity { + id: "chan-1".to_string(), + }, + Utc::now(), + ); + assert_eq!(activity.entity_type, EntityType::Channel); + assert_eq!(activity.entity_id, "chan-1"); + assert_eq!(activity.action, Action::Messaged); +} diff --git a/crates/activity/src/domain/ports.rs b/crates/activity/src/domain/ports.rs new file mode 100644 index 00000000000..4bf493d9759 --- /dev/null +++ b/crates/activity/src/domain/ports.rs @@ -0,0 +1,24 @@ +//! Storage port for activities. + +use model_entity::EntityType; + +use super::models::Activity; + +/// Persists activities. +pub trait ActivityRepo { + /// The adapter's error type. + type Err: std::error::Error + Send + Sync + 'static; + + /// Inserts activities idempotently: an activity whose id already exists is left + /// untouched, so at-least-once redelivery is safe. + fn insert_activities( + &self, + activities: &[Activity], + ) -> impl Future> + Send; + + /// Hard-deletes every activity for the purged entities. + fn purge_entities( + &self, + entities: &[(EntityType, String)], + ) -> impl Future> + Send; +} diff --git a/crates/activity/src/inbound/kafka_consumer.rs b/crates/activity/src/inbound/kafka_consumer.rs new file mode 100644 index 00000000000..602cd811248 --- /dev/null +++ b/crates/activity/src/inbound/kafka_consumer.rs @@ -0,0 +1,159 @@ +//! Generic Kafka consumer that materializes activities from domain events. +//! +//! This machinery knows **zero domains**: it is generic over a declared +//! event collection `C` and a host-supplied dispatcher `Fn(&C) -> Ingest`. +//! The composition root (the hosting service) declares the topics and maps +//! each decoded event to the owning domain's ingest function. +//! +//! Delivery is at least once. Malformed messages and recognized-but-inert +//! events are committed so they cannot wedge a partition. A storage failure +//! aborts the run **without committing** — committing any later record on +//! the partition would cumulatively commit past the failed one — and the +//! host supervisor restarts the consumer from the last committed offset; +//! deterministic activity ids make the replayed inserts idempotent. + +use std::future::Future; +use std::marker::PhantomData; + +use kafka_util::{GroupName, KafkaEventConsumer}; +use macro_event_broker::{KafkaConsumerAdapter, MacroEventCollection, MacroEventConsumerService}; +use rdkafka::consumer::CommitMode; +use rdkafka::message::{BorrowedMessage, Message as _}; +use rootcause::prelude::{Report, ResultExt as _}; +use tracing::Instrument as _; + +use crate::domain::{models::Ingest, ports::ActivityRepo}; + +/// Consumer group for activity materialization offsets. +struct ActivityConsumerGroup; + +impl GroupName for ActivityConsumerGroup { + const GROUP_NAME: &'static str = "activity-materializer"; +} + +/// Consumes activity-bearing topics and writes activities through the repo. +/// +/// `C` is the host's declared event collection (`declare_topics!`); `ingest` +/// is the host's dispatch from a decoded event to the owning domain's +/// mapping. +pub struct ActivityConsumer { + repo: R, + ingest: F, + _events: PhantomData C>, +} + +impl ActivityConsumer +where + R: ActivityRepo, + C: MacroEventCollection + 'static, + F: Fn(&C) -> Ingest + Send + Sync, +{ + /// Builds the consumer over an activity store and an event dispatcher. + pub fn new(repo: R, ingest: F) -> Self { + Self { + repo, + ingest, + _events: PhantomData, + } + } + + /// Applies one decoded event to storage. + async fn apply(&self, event: &C) -> Result<(), R::Err> { + match (self.ingest)(event) { + Ingest::Insert(activities) => self.repo.insert_activities(&activities).await, + Ingest::Purge(entities) => self.repo.purge_entities(&entities).await, + Ingest::Ignore => Ok(()), + } + } + + /// Runs the consumer until `shutdown` resolves. + #[tracing::instrument(skip(self, shutdown), fields(brokers), err)] + pub async fn run( + &self, + brokers: &str, + shutdown: impl Future + Send, + ) -> Result<(), Report> { + let consumer = KafkaEventConsumer::::from_env(brokers)?; + let consumer = KafkaConsumerAdapter::::new(consumer) + .subscribe::() + .context("failed to subscribe to activity topics")?; + let consumer = MacroEventConsumerService::::new(consumer); + tracing::info!( + topics = ?C::topics(), + group = ActivityConsumerGroup::GROUP_NAME, + "activity consumer listening" + ); + + let mut shutdown = std::pin::pin!(shutdown); + loop { + tokio::select! { + _ = &mut shutdown => { + tracing::info!("activity consumer shutting down"); + break; + } + result = consumer.recv() => { + let message = match result { + Ok(message) => message, + Err(e) => { + tracing::error!(error = ?e, "kafka receive error"); + continue; + } + }; + let kafka_message = message.inner(); + let span = tracing::info_span!( + "activity_source_event", + topic = kafka_message.topic(), + partition = kafka_message.partition(), + offset = kafka_message.offset(), + ); + let decoded = { + let _guard = span.enter(); + message.decode_payload() + }; + let event = match decoded { + Ok(event) => event, + Err(_) => { + // Poison record: commit so it cannot wedge the + // partition. + commit_logged(&consumer, kafka_message); + continue; + } + }; + + // A storage failure must abort the run: continuing and + // committing a later record on this partition would + // cumulatively commit past the failed one, losing it + // forever. Returning Err restarts the consumer from the + // last committed offset; deterministic activity ids make + // the replay idempotent. + self.apply(&event) + .instrument(span) + .await + .context("failed to store activities")?; + commit_logged(&consumer, kafka_message); + } + } + } + + Ok(()) + } +} + +fn commit_logged( + consumer: &MacroEventConsumerService>, + message: &BorrowedMessage<'_>, +) { + match consumer.inner().commit_message(message, CommitMode::Async) { + Ok(()) => tracing::trace!( + partition = message.partition(), + offset = message.offset(), + "committed offset" + ), + Err(error) => tracing::error!( + error = ?error, + partition = message.partition(), + offset = message.offset(), + "failed to commit offset" + ), + } +} diff --git a/crates/activity/src/inbound/mod.rs b/crates/activity/src/inbound/mod.rs new file mode 100644 index 00000000000..0bd24635fe1 --- /dev/null +++ b/crates/activity/src/inbound/mod.rs @@ -0,0 +1,3 @@ +//! Inbound adapters. + +pub mod kafka_consumer; diff --git a/crates/activity/src/lib.rs b/crates/activity/src/lib.rs new file mode 100644 index 00000000000..dcabfc38890 --- /dev/null +++ b/crates/activity/src/lib.rs @@ -0,0 +1,31 @@ +#![deny(missing_docs)] +//! The activity protocol and its storage/consumer machinery. +//! +//! One activity: a principal did something to an entity at a time. Every +//! activity surface (feeds, entity timelines, soup attribution sorts) is a +//! query over the single `activity_events` table; there are no derived +//! tables. +//! +//! Ownership boundary: this crate owns **what an activity is** — the +//! durable [`Action`](domain::models::Action) vocabulary, the +//! [`Activity`](domain::models::Activity) row shape, and the +//! [`DomainActivity`](domain::models::DomainActivity) contract. Domain +//! crates own **what counts as activity in their domain**: they depend on +//! this crate with `default-features = false` (models only) and implement +//! their own event → activity mappings. The hosting service is the +//! composition root: it declares the consumed topics and dispatches each +//! decoded event to the owning domain's mapping. +//! +//! Features: `outbound` (Postgres adapter), `consumer` (generic Kafka +//! consumer); both on by default. The models are always available. + +pub mod domain; +#[cfg(feature = "consumer")] +pub mod inbound; +#[cfg(feature = "outbound")] +pub mod outbound; + +pub use domain::models::{ + Action, Activity, ActivitySource, Actor, CallStart, CommonAction, DomainActivity, EntityType, + Ingest, ParticipantChange, PropertyChange, activity_id, event_time, +}; diff --git a/crates/activity/src/outbound/mod.rs b/crates/activity/src/outbound/mod.rs new file mode 100644 index 00000000000..1d7c4cc6bc7 --- /dev/null +++ b/crates/activity/src/outbound/mod.rs @@ -0,0 +1,3 @@ +//! Outbound adapters. + +pub mod pg_activity_repo; diff --git a/crates/activity/src/outbound/pg_activity_repo.rs b/crates/activity/src/outbound/pg_activity_repo.rs new file mode 100644 index 00000000000..02873739db5 --- /dev/null +++ b/crates/activity/src/outbound/pg_activity_repo.rs @@ -0,0 +1,100 @@ +//! Postgres adapter for the `activity_events` table (MacroDB). + +#[cfg(test)] +mod test; + +use model_entity::EntityType; +use sqlx::PgPool; + +use crate::domain::{models::Activity, ports::ActivityRepo}; + +/// Writes activities to MacroDB. +#[derive(Debug, Clone)] +pub struct PgActivityRepo { + pool: PgPool, +} + +impl PgActivityRepo { + /// Builds the adapter over a MacroDB pool. + pub fn new(pool: PgPool) -> Self { + Self { pool } + } +} + +impl ActivityRepo for PgActivityRepo { + type Err = sqlx::Error; + + async fn insert_activities(&self, activities: &[Activity]) -> Result<(), Self::Err> { + if activities.is_empty() { + return Ok(()); + } + + let mut ids = Vec::with_capacity(activities.len()); + let mut actor_ids = Vec::with_capacity(activities.len()); + let mut subject_ids = Vec::with_capacity(activities.len()); + let mut actions = Vec::with_capacity(activities.len()); + let mut payloads = Vec::with_capacity(activities.len()); + let mut entity_types = Vec::with_capacity(activities.len()); + let mut entity_ids = Vec::with_capacity(activities.len()); + let mut occurred_ats = Vec::with_capacity(activities.len()); + for activity in activities { + let (action, payload) = activity.action.to_columns(); + ids.push(activity.id); + actor_ids.push(activity.actor.as_ref().to_owned()); + subject_ids.push(activity.subject_id.clone()); + actions.push(action.to_owned()); + payloads.push(payload); + entity_types.push(activity.entity_type.as_ref().to_owned()); + entity_ids.push(activity.entity_id.clone()); + occurred_ats.push(activity.occurred_at); + } + + sqlx::query!( + r#" + INSERT INTO activity_events + (id, actor_id, subject_id, action, action_payload, + entity_type, entity_id, occurred_at) + SELECT * FROM UNNEST( + $1::uuid[], $2::text[], $3::text[], $4::text[], $5::jsonb[], + $6::text[], $7::text[], $8::timestamptz[]) + ON CONFLICT (id) DO NOTHING + "#, + &ids, + &actor_ids, + &subject_ids, + &actions, + payloads.as_slice() as &[Option], + &entity_types, + &entity_ids, + &occurred_ats, + ) + .execute(&self.pool) + .await?; + Ok(()) + } + + async fn purge_entities(&self, entities: &[(EntityType, String)]) -> Result<(), Self::Err> { + if entities.is_empty() { + return Ok(()); + } + + let entity_types: Vec = entities + .iter() + .map(|(entity_type, _)| entity_type.as_ref().to_owned()) + .collect(); + let entity_ids: Vec = entities.iter().map(|(_, id)| id.clone()).collect(); + + sqlx::query!( + r#" + DELETE FROM activity_events + WHERE (entity_type, entity_id) IN + (SELECT * FROM UNNEST($1::text[], $2::text[])) + "#, + &entity_types, + &entity_ids, + ) + .execute(&self.pool) + .await?; + Ok(()) + } +} diff --git a/crates/activity/src/outbound/pg_activity_repo/test.rs b/crates/activity/src/outbound/pg_activity_repo/test.rs new file mode 100644 index 00000000000..70abab51964 --- /dev/null +++ b/crates/activity/src/outbound/pg_activity_repo/test.rs @@ -0,0 +1,91 @@ +use chrono::Utc; +use macro_db_migrator::MACRO_DB_MIGRATIONS; +use macro_user_id::user_id::MacroUserIdStr; +use uuid::Uuid; + +use super::*; +use crate::domain::models::{Actor, CommonAction}; + +fn user(id: &str) -> MacroUserIdStr<'static> { + MacroUserIdStr::try_from(id.to_string()).expect("valid user id") +} + +fn seed(source_event: u128, action: CommonAction, entity_id: &str) -> Activity { + Activity::common( + Uuid::from_u128(source_event), + 0, + Actor::new_from_user(user("macro|actor@example.com")), + None, + model_entity::EntityType::Document, + entity_id, + action, + Utc::now(), + ) +} + +#[sqlx::test(migrator = "MACRO_DB_MIGRATIONS")] +async fn inserts_activities_with_split_action_columns(pool: PgPool) { + let repo = PgActivityRepo::new(pool.clone()); + + repo.insert_activities(&[seed(1, CommonAction::Opened, "doc-1")]) + .await + .unwrap(); + + let row = sqlx::query!( + r#" + SELECT actor_id, subject_id, action, action_payload, entity_type, entity_id + FROM activity_events + "# + ) + .fetch_one(&pool) + .await + .unwrap(); + + assert_eq!(row.actor_id, "macro|actor@example.com"); + assert_eq!(row.subject_id, "macro|actor@example.com"); + assert_eq!(row.action, "opened"); + assert_eq!(row.action_payload, None); + assert_eq!(row.entity_type, "document"); + assert_eq!(row.entity_id, "doc-1"); +} + +#[sqlx::test(migrator = "MACRO_DB_MIGRATIONS")] +async fn replayed_activities_are_absorbed_by_the_id_conflict(pool: PgPool) { + let repo = PgActivityRepo::new(pool.clone()); + let seed_activity = seed(2, CommonAction::Edited, "doc-2"); + + repo.insert_activities(std::slice::from_ref(&seed_activity)) + .await + .unwrap(); + repo.insert_activities(std::slice::from_ref(&seed_activity)) + .await + .unwrap(); + + let count = sqlx::query_scalar!(r#"SELECT COUNT(*) FROM activity_events"#) + .fetch_one(&pool) + .await + .unwrap(); + assert_eq!(count, Some(1)); +} + +#[sqlx::test(migrator = "MACRO_DB_MIGRATIONS")] +async fn purge_removes_only_that_entitys_activities(pool: PgPool) { + let repo = PgActivityRepo::new(pool.clone()); + repo.insert_activities(&[ + seed(3, CommonAction::Created, "doc-purged"), + seed(4, CommonAction::Opened, "doc-purged"), + seed(5, CommonAction::Created, "doc-kept"), + ]) + .await + .unwrap(); + + repo.purge_entities(&[(model_entity::EntityType::Document, "doc-purged".to_string())]) + .await + .unwrap(); + + let remaining = sqlx::query_scalar!(r#"SELECT entity_id FROM activity_events"#) + .fetch_all(&pool) + .await + .unwrap(); + assert_eq!(remaining, vec!["doc-kept".to_string()]); +} diff --git a/crates/call/Cargo.toml b/crates/call/Cargo.toml index 7ef5fb57667..4bd4adc963e 100644 --- a/crates/call/Cargo.toml +++ b/crates/call/Cargo.toml @@ -35,6 +35,7 @@ outbound = [ ports = ["dep:entity_access", "dep:tokio"] [dependencies] +activity = { path = "../activity", default-features = false } agent = { path = "../agent", optional = true } ai_usage = { path = "../ai_usage", optional = true } ai_toolset = { path = "../ai_toolset", optional = true } diff --git a/crates/call/src/domain/activity.rs b/crates/call/src/domain/activity.rs new file mode 100644 index 00000000000..fa48b5a7577 --- /dev/null +++ b/crates/call/src/domain/activity.rs @@ -0,0 +1,91 @@ +//! What counts as activity in the calls domain. +//! +//! Note the ownership shape: "a call started in this channel" is *call* +//! knowledge even though the activity targets a channel entity — so the +//! call crate owns the [`DomainActivity`] impl targeting +//! [`EntityType::Channel`]. Ownership follows the knowledge, not the entity +//! kind. + +#[cfg(test)] +mod test; + +use ::activity::{ + Action, Activity, ActivitySource, Actor, CallStart, CommonAction, DomainActivity, EntityType, + Ingest, +}; +use uuid::Uuid; + +use super::events::CallTopicEvent; + +/// A call started in a channel — a channel-targeted activity owned by the +/// calls domain. +#[derive(Debug, Clone, PartialEq)] +pub struct CallStartedActivity { + /// The channel the call started in. + pub channel_id: String, + /// The started call. + pub call_id: String, +} + +impl DomainActivity for CallStartedActivity { + const ENTITY_TYPE: EntityType = EntityType::Channel; + + fn entity_id(&self) -> &str { + &self.channel_id + } + + fn into_action(self) -> Action { + Action::CallStarted(CallStart { + call_id: self.call_id, + }) + } +} + +impl ActivitySource for CallTopicEvent { + /// Maps one `macro.calls` event to its ingest outcome. + /// + /// Exhaustive on purpose: a new event variant fails compilation here + /// until someone classifies it or explicitly drops it. + fn ingest(&self, event_id: Uuid) -> Ingest { + match self { + // Two activities: the call happened in the channel (timeline/feed), + // and the call itself is a new entity (soup item). + CallTopicEvent::Started(m) => { + let actor = Actor::new_from_user(m.created_by.clone()); + Ingest::Insert(vec![ + Activity::from_domain( + event_id, + 0, + actor.clone(), + None, + CallStartedActivity { + channel_id: m.channel_id.to_string(), + call_id: m.call_id.to_string(), + }, + m.created_at, + ), + Activity::common( + event_id, + 1, + actor, + None, + EntityType::Call, + m.call_id.to_string(), + CommonAction::Created, + m.created_at, + ), + ]) + } + // Call-record deletion is a hard delete; purge the call's + // activities like other hard deletes. + CallTopicEvent::RecordDeleted(m) => { + Ingest::Purge(vec![(EntityType::Call, m.call_id.to_string())]) + } + // Archival/processing pipeline, not user activity. + CallTopicEvent::RecordArchived(_) + | CallTopicEvent::RecordUpdated(_) + | CallTopicEvent::RecordSummarized(_) + | CallTopicEvent::RecordingReady(_) => Ingest::Ignore, + } + } +} diff --git a/crates/call/src/domain/activity/test.rs b/crates/call/src/domain/activity/test.rs new file mode 100644 index 00000000000..a98640ae468 --- /dev/null +++ b/crates/call/src/domain/activity/test.rs @@ -0,0 +1,64 @@ +use ::activity::Action; +use chrono::Utc; +use macro_user_id::user_id::MacroUserIdStr; +use uuid::Uuid; + +use macro_event_broker::Event; + +use super::*; +use crate::domain::events::{CallRecordDeletedMetadata, CallStartedMetadata}; + +#[test] +fn call_started_yields_channel_and_call_activities() { + let call_id = Uuid::from_u128(5); + let channel_id = Uuid::from_u128(6); + let created_at = Utc::now(); + let event = Event::with_event_id( + Uuid::now_v7(), + CallTopicEvent::Started(CallStartedMetadata { + call_id, + channel_id, + created_by: MacroUserIdStr::try_from("macro|rahul@example.com".to_string()).unwrap(), + created_at, + recording_enabled: false, + }), + ); + + let Ingest::Insert(activities) = event.event.ingest(event.event_id) else { + panic!("expected activities"); + }; + assert_eq!(activities.len(), 2); + assert_ne!(activities[0].id, activities[1].id); + + assert_eq!(activities[0].entity_type, EntityType::Channel); + assert_eq!(activities[0].entity_id, channel_id.to_string()); + assert_eq!( + activities[0].action, + Action::CallStarted(::activity::CallStart { + call_id: call_id.to_string() + }) + ); + + assert_eq!(activities[1].entity_type, EntityType::Call); + assert_eq!(activities[1].entity_id, call_id.to_string()); + assert_eq!(activities[1].action, Action::Created); + assert!(activities.iter().all(|a| a.occurred_at == created_at)); +} + +#[test] +fn record_deletion_purges_the_call() { + let call_id = Uuid::from_u128(5); + let event = Event::with_event_id( + Uuid::now_v7(), + CallTopicEvent::RecordDeleted(CallRecordDeletedMetadata { + call_id, + channel_id: Uuid::from_u128(6), + actor_user_id: None, + }), + ); + + assert_eq!( + event.event.ingest(event.event_id), + Ingest::Purge(vec![(EntityType::Call, call_id.to_string())]) + ); +} diff --git a/crates/call/src/domain/mod.rs b/crates/call/src/domain/mod.rs index 31abcf6090b..8820297c2a2 100644 --- a/crates/call/src/domain/mod.rs +++ b/crates/call/src/domain/mod.rs @@ -1,4 +1,6 @@ /// Kafka event contracts for call lifecycle events. +/// Event-to-activity mappings for this domain. +pub mod activity; pub mod events; /// Domain models for calls. diff --git a/crates/channels/Cargo.toml b/crates/channels/Cargo.toml index 91f5d18f064..b7394ecec05 100644 --- a/crates/channels/Cargo.toml +++ b/crates/channels/Cargo.toml @@ -75,6 +75,7 @@ ports = ["dep:entity_access", "dep:tokio", "entity_access/ports"] schema = ["channel_sender/schema", "dep:utoipa", "macro_user_id/schema"] [dependencies] +activity = { path = "../activity", default-features = false } anyhow = { workspace = true } attachment = { path = "../attachment", optional = true } axum = { workspace = true, optional = true } diff --git a/crates/channels/src/domain/activity.rs b/crates/channels/src/domain/activity.rs new file mode 100644 index 00000000000..543ff8a02fc --- /dev/null +++ b/crates/channels/src/domain/activity.rs @@ -0,0 +1,213 @@ +//! What counts as activity in the channels domain. +//! +//! Channels own their exclusive action vocabulary and its projection into +//! the durable [`Action`] — the activity crate never learns channel +//! semantics. + +#[cfg(test)] +mod test; + +use ::activity::{ + Action, Activity, ActivitySource, Actor, CommonAction, DomainActivity, EntityType, Ingest, + ParticipantChange, event_time, +}; +use chrono::{DateTime, Utc}; +use macro_user_id::user_id::MacroUserIdStr; +use uuid::Uuid; + +use super::broker_events::ChannelTopicEvent; + +/// Channel-exclusive actions. Common lifecycle actions go through +/// [`Activity::common`] and need no representation here. +#[derive(Debug, Clone, PartialEq)] +pub enum ChannelAction { + /// A message was posted in the channel. + Messaged, + /// A principal was added to the channel. + ParticipantAdded { + /// The added principal. + participant: Actor<'static>, + }, + /// A principal was removed from the channel. + ParticipantRemoved { + /// The removed principal. + participant: Actor<'static>, + }, +} + +/// A channel-exclusive activity: the channel it happened in, paired with an +/// action only channels support. +#[derive(Debug, Clone, PartialEq)] +pub struct ChannelActivity { + /// The channel acted on. + pub channel_id: String, + /// What happened to it. + pub action: ChannelAction, +} + +impl DomainActivity for ChannelActivity { + const ENTITY_TYPE: EntityType = EntityType::Channel; + + fn entity_id(&self) -> &str { + &self.channel_id + } + + fn into_action(self) -> Action { + match self.action { + ChannelAction::Messaged => Action::Messaged, + ChannelAction::ParticipantAdded { participant } => { + Action::ParticipantAdded(ParticipantChange { participant }) + } + ChannelAction::ParticipantRemoved { participant } => { + Action::ParticipantRemoved(ParticipantChange { participant }) + } + } + } +} + +fn exclusive( + event_id: Uuid, + ordinal: u32, + actor: Actor<'static>, + on_behalf_of: Option>, + channel_id: Uuid, + action: ChannelAction, + occurred_at: DateTime, +) -> Activity { + Activity::from_domain( + event_id, + ordinal, + actor, + on_behalf_of, + ChannelActivity { + channel_id: channel_id.to_string(), + action, + }, + occurred_at, + ) +} + +/// One activity per (un)added participant; ordinals keep replay ids stable. +fn participant_activities( + event_id: Uuid, + actor: Actor<'static>, + users: &[MacroUserIdStr<'static>], + channel_id: Uuid, + occurred_at: DateTime, + make_action: impl Fn(Actor<'static>) -> ChannelAction, +) -> Ingest { + Ingest::Insert( + users + .iter() + .enumerate() + .map(|(ordinal, user)| { + exclusive( + event_id, + u32::try_from(ordinal).unwrap_or(u32::MAX), + actor.clone(), + None, + channel_id, + make_action(Actor::new_from_user(user.clone())), + occurred_at, + ) + }) + .collect(), + ) +} + +impl ActivitySource for ChannelTopicEvent { + /// Maps one `macro.channels` event to its ingest outcome. + /// + /// Exhaustive on purpose: a new event variant fails compilation here + /// until someone classifies it or explicitly drops it. + fn ingest(&self, event_id: Uuid) -> Ingest { + let now = || event_time(event_id); + let common = + |actor: Actor<'static>, action: CommonAction, channel_id: Uuid, at: DateTime| { + Ingest::Insert(vec![Activity::common( + event_id, + 0, + actor, + None, + EntityType::Channel, + channel_id.to_string(), + action, + at, + )]) + }; + + match self { + ChannelTopicEvent::Created(m) => { + common(m.actor.clone(), CommonAction::Created, m.channel_id, now()) + } + ChannelTopicEvent::Updated(m) => common( + Actor::new_from_user(m.actor.clone()), + CommonAction::Edited, + m.channel_id, + now(), + ), + ChannelTopicEvent::Deleted(m) => { + common(m.actor.clone(), CommonAction::Deleted, m.channel_id, now()) + } + ChannelTopicEvent::MessagePosted(m) => { + // For agent (bot) messages, `triggered_by` is the user whose + // authority the message was sent under — the activity's subject. + let on_behalf_of = m.triggered_by.as_deref().and_then(|id| { + MacroUserIdStr::try_from(id.to_string()) + .inspect_err(|e| { + // Fall back to the sender as subject, but loudly: the + // triggering user's activity is being misattributed. + tracing::warn!(error=?e, triggered_by=id, "unparseable triggered_by"); + }) + .ok() + }); + Ingest::Insert(vec![exclusive( + event_id, + 0, + m.sender.clone(), + on_behalf_of, + m.channel_id, + ChannelAction::Messaged, + m.created_at, + )]) + } + ChannelTopicEvent::MessagePatched(m) => common( + m.actor.clone(), + CommonAction::Edited, + m.channel_id, + m.updated_at, + ), + // Deleting a message mutates the channel's content. + ChannelTopicEvent::MessageDeleted(m) => common( + m.actor.clone(), + CommonAction::Edited, + m.channel_id, + m.deleted_at.unwrap_or_else(now), + ), + ChannelTopicEvent::MessageAttachmentCreated(m) => { + common(m.actor.clone(), CommonAction::Edited, m.channel_id, now()) + } + ChannelTopicEvent::MessageAttachmentRemoved(m) => { + common(m.actor.clone(), CommonAction::Edited, m.channel_id, now()) + } + ChannelTopicEvent::ParticipantAdded(m) => participant_activities( + event_id, + m.added_by.clone(), + &m.added_user_ids, + m.channel_id, + now(), + |participant| ChannelAction::ParticipantAdded { participant }, + ), + ChannelTopicEvent::ParticipantRemoved(m) => participant_activities( + event_id, + Actor::new_from_user(m.removed_by.clone()), + &m.removed_user_ids, + m.channel_id, + now(), + |participant| ChannelAction::ParticipantRemoved { participant }, + ), + // Derivative of MessagePosted (the full mention list travels there). + ChannelTopicEvent::Mentioned(_) => Ingest::Ignore, + } + } +} diff --git a/crates/channels/src/domain/activity/test.rs b/crates/channels/src/domain/activity/test.rs new file mode 100644 index 00000000000..c621018d5a0 --- /dev/null +++ b/crates/channels/src/domain/activity/test.rs @@ -0,0 +1,107 @@ +use ::activity::Action; +use ::activity::EntityType; +use chrono::Utc; +use macro_user_id::user_id::MacroUserIdStr; +use uuid::Uuid; + +use macro_event_broker::Event; + +use super::*; +use crate::domain::broker_events::{ + ChannelCreatedMetadata, ChannelMessagePostedMetadata, ChannelParticipantAddedMetadata, +}; +use crate::domain::models::ChannelType; + +fn user(id: &str) -> MacroUserIdStr<'static> { + MacroUserIdStr::try_from(id.to_string()).expect("valid user id") +} + +fn envelope(event: ChannelTopicEvent) -> Event { + Event::with_event_id(Uuid::now_v7(), event) +} + +const CHANNEL_ID: Uuid = Uuid::from_u128(7); + +#[test] +fn message_posted_maps_to_messaged_with_triggered_by_as_subject() { + let created_at = Utc::now(); + let event = envelope(ChannelTopicEvent::MessagePosted( + ChannelMessagePostedMetadata { + channel_id: CHANNEL_ID, + message_id: Uuid::from_u128(8), + thread_id: None, + sender: Actor::new_from_user(user("macro|bot-like@example.com")), + triggered_by: Some("macro|teo@example.com".to_string()), + channel_type: ChannelType::Public, + content: "hi".to_string(), + mentions: vec![], + attachments: vec![], + created_at, + }, + )); + + let Ingest::Insert(activities) = event.event.ingest(event.event_id) else { + panic!("expected activities"); + }; + assert_eq!(activities.len(), 1); + assert_eq!(activities[0].action, Action::Messaged); + // Delegation: the triggering user is the subject; the sender the actor. + assert_eq!(activities[0].subject_id, "macro|teo@example.com"); + assert_eq!(activities[0].actor.as_ref(), "macro|bot-like@example.com"); + assert_eq!(activities[0].occurred_at, created_at); + assert_eq!(activities[0].entity_id, CHANNEL_ID.to_string()); + assert_eq!(activities[0].entity_type, EntityType::Channel); +} + +#[test] +fn participant_added_yields_one_activity_per_user_with_stable_ordinals() { + let event = envelope(ChannelTopicEvent::ParticipantAdded( + ChannelParticipantAddedMetadata { + channel_id: CHANNEL_ID, + channel_type: ChannelType::Public, + added_by: Actor::new_from_user(user("macro|admin@example.com")), + added_user_ids: vec![user("macro|a@example.com"), user("macro|b@example.com")], + }, + )); + + let Ingest::Insert(activities) = event.event.ingest(event.event_id) else { + panic!("expected activities"); + }; + assert_eq!(activities.len(), 2); + assert_ne!(activities[0].id, activities[1].id); + assert!( + activities + .iter() + .all(|a| a.subject_id == "macro|admin@example.com") + ); + assert_eq!( + activities[0].action, + Action::ParticipantAdded(::activity::ParticipantChange { + participant: Actor::new_from_user(user("macro|a@example.com")) + }) + ); + + // Replay derives identical ids. + let Ingest::Insert(replayed) = event.event.ingest(event.event_id) else { + panic!("expected activities"); + }; + assert_eq!(activities[0].id, replayed[0].id); + assert_eq!(activities[1].id, replayed[1].id); +} + +#[test] +fn created_maps_to_created_by_the_actor() { + let event = envelope(ChannelTopicEvent::Created(ChannelCreatedMetadata { + channel_id: CHANNEL_ID, + actor: Actor::new_from_user(user("macro|owner@example.com")), + channel_type: ChannelType::Public, + channel_name: Some("general".to_string()), + participant_user_ids: vec![user("macro|owner@example.com")], + })); + + let Ingest::Insert(activities) = event.event.ingest(event.event_id) else { + panic!("expected activities"); + }; + assert_eq!(activities[0].action, Action::Created); + assert_eq!(activities[0].subject_id, "macro|owner@example.com"); +} diff --git a/crates/channels/src/domain/mod.rs b/crates/channels/src/domain/mod.rs index c6ec8374d47..f6740496a09 100644 --- a/crates/channels/src/domain/mod.rs +++ b/crates/channels/src/domain/mod.rs @@ -1,4 +1,6 @@ /// Kafka event models for the `macro.channels` topic. +/// Event-to-activity mappings for this domain. +pub mod activity; pub mod broker_events; #[cfg(feature = "entity_mutation")] /// Unified entity-mutation capability impls. diff --git a/crates/chat/Cargo.toml b/crates/chat/Cargo.toml index 25102db7f18..e012c34e742 100644 --- a/crates/chat/Cargo.toml +++ b/crates/chat/Cargo.toml @@ -22,6 +22,7 @@ outbound = ["dep:entity_access_db_utils", "dep:macro_uuid", "dep:sqlx", "ports"] ports = [] [dependencies] +activity = { path = "../activity", default-features = false } # Domain (always included) agent = { path = "../agent" } ai_toolset = { path = "../ai_toolset" } diff --git a/crates/chat/src/domain.rs b/crates/chat/src/domain.rs index d1bb7fd4f6f..0f7585952cb 100644 --- a/crates/chat/src/domain.rs +++ b/crates/chat/src/domain.rs @@ -1,4 +1,6 @@ /// Kafka event models for the `macro.chats` topic. +/// Event-to-activity mappings for this domain. +pub mod activity; pub mod events; /// Types used by the chat domain ports. pub mod models; diff --git a/crates/chat/src/domain/activity.rs b/crates/chat/src/domain/activity.rs new file mode 100644 index 00000000000..e772df4e612 --- /dev/null +++ b/crates/chat/src/domain/activity.rs @@ -0,0 +1,131 @@ +//! What counts as activity in the chat domain. + +#[cfg(test)] +mod test; + +use ::activity::{ + Action, Activity, ActivitySource, Actor, CommonAction, DomainActivity, EntityType, Ingest, + event_time, +}; +use chrono::{DateTime, Utc}; +use uuid::Uuid; + +use super::events::{ChatMessageRole, ChatTopicEvent}; + +/// Chat-exclusive actions. Common lifecycle actions go through +/// [`Activity::common`] and need no representation here. +#[derive(Debug, Clone, PartialEq)] +pub enum ChatAction { + /// The subject sent a message in the chat. + Messaged, +} + +/// A chat-exclusive activity. +#[derive(Debug, Clone, PartialEq)] +pub struct ChatActivity { + /// The chat acted on. + pub chat_id: String, + /// What happened to it. + pub action: ChatAction, +} + +impl DomainActivity for ChatActivity { + const ENTITY_TYPE: EntityType = EntityType::Chat; + + fn entity_id(&self) -> &str { + &self.chat_id + } + + fn into_action(self) -> Action { + match self.action { + ChatAction::Messaged => Action::Messaged, + } + } +} + +impl ActivitySource for ChatTopicEvent { + /// Maps one `macro.chats` event to its ingest outcome. + /// + /// Exhaustive on purpose: a new event variant fails compilation here + /// until someone classifies it or explicitly drops it. + fn ingest(&self, event_id: Uuid) -> Ingest { + let now = || event_time(event_id); + let common = + |actor: Actor<'static>, action: CommonAction, chat_id: &str, at: DateTime| { + Ingest::Insert(vec![Activity::common( + event_id, + 0, + actor, + None, + EntityType::Chat, + chat_id, + action, + at, + )]) + }; + + match self { + ChatTopicEvent::Created(m) => common( + Actor::new_from_user(m.owner.clone()), + CommonAction::Created, + &m.chat_id, + now(), + ), + // The copy is a new chat; its creation is the activity. + ChatTopicEvent::Copied(m) => common( + Actor::new_from_user(m.owner.clone()), + CommonAction::Created, + &m.chat_id, + now(), + ), + ChatTopicEvent::Updated(m) => common( + Actor::new_from_user(m.actor_user_id.clone()), + CommonAction::Edited, + &m.chat_id, + now(), + ), + // Only the user's own prompts are their activity; assistant + // responses are a consequence, not an act by the subject. + ChatTopicEvent::MessageSent(m) => match (&m.actor_user_id, &m.role) { + (Some(actor), ChatMessageRole::User) => { + Ingest::Insert(vec![Activity::from_domain( + event_id, + 0, + Actor::new_from_user(actor.clone()), + None, + ChatActivity { + chat_id: m.chat_id.clone(), + action: ChatAction::Messaged, + }, + now(), + )]) + } + _ => Ingest::Ignore, + }, + ChatTopicEvent::Deleted(m) => match &m.actor_user_id { + Some(actor) => common( + Actor::new_from_user(actor.clone()), + CommonAction::Deleted, + &m.chat_id, + now(), + ), + None => Ingest::Ignore, + }, + // A restore is a mutation of the chat's lifecycle state. + ChatTopicEvent::Restored(m) => match &m.actor_user_id { + Some(actor) => common( + Actor::new_from_user(actor.clone()), + CommonAction::Edited, + &m.chat_id, + now(), + ), + None => Ingest::Ignore, + }, + ChatTopicEvent::PermanentlyDeleted(m) => { + Ingest::Purge(vec![(EntityType::Chat, m.chat_id.clone())]) + } + // Carries no actor. + ChatTopicEvent::MessageDeleted(_) => Ingest::Ignore, + } + } +} diff --git a/crates/chat/src/domain/activity/test.rs b/crates/chat/src/domain/activity/test.rs new file mode 100644 index 00000000000..611b8d6767b --- /dev/null +++ b/crates/chat/src/domain/activity/test.rs @@ -0,0 +1,81 @@ +use ::activity::Action; +use macro_user_id::user_id::MacroUserIdStr; +use model_entity::EntityType; +use uuid::Uuid; + +use macro_event_broker::Event; + +use super::*; +use crate::domain::events::{ + ChatCreatedMetadata, ChatMessageSentMetadata, ChatPermanentlyDeletedMetadata, +}; + +fn user(id: &str) -> MacroUserIdStr<'static> { + MacroUserIdStr::try_from(id.to_string()).expect("valid user id") +} + +fn envelope(event: ChatTopicEvent) -> Event { + Event::with_event_id(Uuid::now_v7(), event) +} + +fn message(role: ChatMessageRole, actor: Option>) -> ChatTopicEvent { + ChatTopicEvent::MessageSent(ChatMessageSentMetadata { + chat_id: "chat-1".to_string(), + message_id: "msg-1".to_string(), + role, + model: "test".to_string(), + actor_user_id: actor, + attachment_count: 0, + }) +} + +#[test] +fn created_and_user_message_map_to_activities() { + let created = envelope(ChatTopicEvent::Created(ChatCreatedMetadata { + chat_id: "chat-1".to_string(), + owner: user("macro|owner@example.com"), + name: "planning".to_string(), + project_id: None, + })); + let Ingest::Insert(activities) = created.event.ingest(created.event_id) else { + panic!("expected activities"); + }; + assert_eq!(activities[0].action, Action::Created); + assert_eq!(activities[0].entity_type, EntityType::Chat); + + let sent = envelope(message( + ChatMessageRole::User, + Some(user("macro|owner@example.com")), + )); + let Ingest::Insert(activities) = sent.event.ingest(sent.event_id) else { + panic!("expected activities"); + }; + assert_eq!(activities[0].action, Action::Messaged); +} + +#[test] +fn assistant_and_actorless_messages_are_dropped() { + let assistant = envelope(message( + ChatMessageRole::Assistant, + Some(user("macro|owner@example.com")), + )); + assert_eq!(assistant.event.ingest(assistant.event_id), Ingest::Ignore); + + let actorless = envelope(message(ChatMessageRole::User, None)); + assert_eq!(actorless.event.ingest(actorless.event_id), Ingest::Ignore); +} + +#[test] +fn permanent_delete_purges() { + let purged = envelope(ChatTopicEvent::PermanentlyDeleted( + ChatPermanentlyDeletedMetadata { + chat_id: "chat-1".to_string(), + actor_user_id: None, + project_id: None, + }, + )); + assert_eq!( + purged.event.ingest(purged.event_id), + Ingest::Purge(vec![(EntityType::Chat, "chat-1".to_string())]) + ); +} diff --git a/crates/documents/Cargo.toml b/crates/documents/Cargo.toml index d33b6bfabba..ab1ca17eb2b 100644 --- a/crates/documents/Cargo.toml +++ b/crates/documents/Cargo.toml @@ -69,6 +69,7 @@ schema = ["dep:utoipa", "macro_user_id/schema"] service = ["dep:entity_access_management", "dep:foreign_entity", "ports"] [dependencies] +activity = { path = "../activity", default-features = false } # Core anyhow = { workspace = true } chrono = { workspace = true } diff --git a/crates/documents/src/domain.rs b/crates/documents/src/domain.rs index 61b3f6727e9..1ab6ebe4af8 100644 --- a/crates/documents/src/domain.rs +++ b/crates/documents/src/domain.rs @@ -1,5 +1,7 @@ //! Domain layer: models, ports (trait interfaces), and service implementation. +/// Event-to-activity mappings for this domain. +pub mod activity; pub mod branch_name; pub mod content; /// Unified entity-mutation capability impls. diff --git a/crates/documents/src/domain/activity.rs b/crates/documents/src/domain/activity.rs new file mode 100644 index 00000000000..266c6edcaa7 --- /dev/null +++ b/crates/documents/src/domain/activity.rs @@ -0,0 +1,81 @@ +//! What counts as activity in the documents domain. +//! +//! Documents have no entity-exclusive actions yet, so every mapping goes +//! through [`Activity::common`]. + +#[cfg(test)] +mod test; + +use ::activity::{Activity, ActivitySource, Actor, CommonAction, EntityType, Ingest, event_time}; +use chrono::{DateTime, Utc}; +use uuid::Uuid; + +use super::events::DocumentTopicEvent; + +impl ActivitySource for DocumentTopicEvent { + /// Maps one `macro.documents` event to its ingest outcome. + /// + /// Exhaustive on purpose: a new event variant fails compilation here + /// until someone classifies it or explicitly drops it. + fn ingest(&self, event_id: Uuid) -> Ingest { + let single = |actor: Actor<'static>, + action: CommonAction, + document_id: &str, + occurred_at: DateTime| { + Ingest::Insert(vec![Activity::common( + event_id, + 0, + actor, + None, + EntityType::Document, + document_id, + action, + occurred_at, + )]) + }; + + match self { + DocumentTopicEvent::Created(metadata) => single( + Actor::new_from_user(metadata.owner.clone()), + CommonAction::Created, + &metadata.document_id, + metadata.created_at.unwrap_or_else(|| event_time(event_id)), + ), + DocumentTopicEvent::Updated(metadata) => match &metadata.actor_user_id { + Some(actor) => single( + Actor::new_from_user(actor.clone()), + CommonAction::Edited, + &metadata.document_id, + event_time(event_id), + ), + None => Ingest::Ignore, + }, + DocumentTopicEvent::Deleted(metadata) => match &metadata.actor_user_id { + Some(actor) => single( + Actor::new_from_user(actor.clone()), + CommonAction::Deleted, + &metadata.document_id, + event_time(event_id), + ), + None => Ingest::Ignore, + }, + // The copy is a new document; its creation is the activity. + DocumentTopicEvent::Copied(metadata) => single( + Actor::new_from_user(metadata.owner.clone()), + CommonAction::Created, + &metadata.document_id, + event_time(event_id), + ), + DocumentTopicEvent::Purged(metadata) => { + Ingest::Purge(vec![(EntityType::Document, metadata.document_id.clone())]) + } + // Extraction-pipeline noise, not user activity. + DocumentTopicEvent::ContentUploaded(_) => Ingest::Ignore, + // Carries no actor today; becomes an Edited activity once collab + // edits are attributed. + DocumentTopicEvent::SyncContentUpdated(_) => Ingest::Ignore, + // Session lifecycle (first join / last leave), no actor. + DocumentTopicEvent::Interaction(_) => Ingest::Ignore, + } + } +} diff --git a/crates/documents/src/domain/activity/test.rs b/crates/documents/src/domain/activity/test.rs new file mode 100644 index 00000000000..a7ed347d1ed --- /dev/null +++ b/crates/documents/src/domain/activity/test.rs @@ -0,0 +1,170 @@ +use ::activity::{Action, activity_id}; +use chrono::{TimeZone as _, Utc}; +use macro_user_id::user_id::MacroUserIdStr; +use model::document::FileType; +use model_entity::EntityType; +use uuid::Uuid; + +use macro_event_broker::Event; + +use super::*; +use crate::domain::events::{ + DocumentCopiedMetadata, DocumentCreatedMetadata, DocumentDeletedMetadata, + DocumentInteractionMetadata, DocumentPurgedMetadata, DocumentSyncContentUpdatedMetadata, + DocumentUpdatedMetadata, InteractionReason, +}; + +const DOCUMENT_ID: &str = "11111111-1111-1111-1111-111111111111"; + +fn user(id: &str) -> MacroUserIdStr<'static> { + MacroUserIdStr::try_from(id.to_string()).expect("valid user id") +} + +fn envelope(event: DocumentTopicEvent) -> Event { + Event::with_event_id(Uuid::now_v7(), event) +} + +fn single_activity(ingest: Ingest) -> Activity { + match ingest { + Ingest::Insert(mut activities) => { + assert_eq!(activities.len(), 1); + activities.pop().unwrap() + } + other => panic!("expected a single activity, got {other:?}"), + } +} + +#[test] +fn created_maps_to_a_created_activity_with_the_metadata_timestamp() { + let created_at = Utc.with_ymd_and_hms(2026, 8, 5, 12, 0, 0).unwrap(); + let event = envelope(DocumentTopicEvent::Created(DocumentCreatedMetadata { + document_id: DOCUMENT_ID.to_string(), + owner: user("macro|creator@example.com"), + document_name: "spec".to_string(), + file_type: Some(FileType::Md), + project_id: None, + sub_type: None, + created_at: Some(created_at), + })); + + let activity = single_activity(event.event.ingest(event.event_id)); + assert_eq!(activity.action, Action::Created); + assert_eq!(activity.subject_id, "macro|creator@example.com"); + assert_eq!(activity.entity_type, EntityType::Document); + assert_eq!(activity.entity_id, DOCUMENT_ID); + assert_eq!(activity.occurred_at, created_at); + assert_eq!(activity.id, activity_id(event.event_id, 0)); +} + +#[test] +fn attributed_update_and_delete_map_to_activities() { + let updated = envelope(DocumentTopicEvent::Updated(DocumentUpdatedMetadata { + document_id: DOCUMENT_ID.to_string(), + owner: user("macro|owner@example.com"), + actor_user_id: Some(user("macro|editor@example.com")), + document_name: Some("renamed".to_string()), + previous_project_id: None, + project_id: None, + file_type: None, + share_permission_updated: false, + })); + let activity = single_activity(updated.event.ingest(updated.event_id)); + assert_eq!(activity.action, Action::Edited); + assert_eq!(activity.subject_id, "macro|editor@example.com"); + + let deleted = envelope(DocumentTopicEvent::Deleted(DocumentDeletedMetadata { + document_id: DOCUMENT_ID.to_string(), + actor_user_id: Some(user("macro|editor@example.com")), + project_id: None, + })); + let activity = single_activity(deleted.event.ingest(deleted.event_id)); + assert_eq!(activity.action, Action::Deleted); +} + +#[test] +fn unattributable_mutations_are_dropped() { + let updated = envelope(DocumentTopicEvent::Updated(DocumentUpdatedMetadata { + document_id: DOCUMENT_ID.to_string(), + owner: user("macro|owner@example.com"), + actor_user_id: None, + document_name: None, + previous_project_id: None, + project_id: None, + file_type: None, + share_permission_updated: true, + })); + assert_eq!(updated.event.ingest(updated.event_id), Ingest::Ignore); + + let deleted = envelope(DocumentTopicEvent::Deleted(DocumentDeletedMetadata { + document_id: DOCUMENT_ID.to_string(), + actor_user_id: None, + project_id: None, + })); + assert_eq!(deleted.event.ingest(deleted.event_id), Ingest::Ignore); +} + +#[test] +fn copied_maps_to_a_created_activity_for_the_new_document() { + let event = envelope(DocumentTopicEvent::Copied(DocumentCopiedMetadata { + document_id: "22222222-2222-2222-2222-222222222222".to_string(), + source_document_id: DOCUMENT_ID.to_string(), + source_version_id: None, + owner: user("macro|copier@example.com"), + document_name: "copy".to_string(), + file_type: None, + project_id: None, + sub_type: None, + })); + + let activity = single_activity(event.event.ingest(event.event_id)); + assert_eq!(activity.action, Action::Created); + assert_eq!(activity.entity_id, "22222222-2222-2222-2222-222222222222"); +} + +#[test] +fn purge_requests_entity_deletion() { + let event = envelope(DocumentTopicEvent::Purged(DocumentPurgedMetadata { + document_id: DOCUMENT_ID.to_string(), + })); + + assert_eq!( + event.event.ingest(event.event_id), + Ingest::Purge(vec![(EntityType::Document, DOCUMENT_ID.to_string())]) + ); +} + +#[test] +fn pipeline_and_session_events_are_ignored() { + let sync = envelope(DocumentTopicEvent::SyncContentUpdated( + DocumentSyncContentUpdatedMetadata { + document_id: DOCUMENT_ID.to_string(), + file_type: FileType::Md, + document_version_id: None, + }, + )); + assert_eq!(sync.event.ingest(sync.event_id), Ingest::Ignore); + + let interaction = envelope(DocumentTopicEvent::Interaction( + DocumentInteractionMetadata { + document_id: DOCUMENT_ID.to_string(), + reason: InteractionReason::FirstJoin, + }, + )); + assert_eq!( + interaction.event.ingest(interaction.event_id), + Ingest::Ignore + ); +} + +#[test] +fn replaying_an_event_derives_identical_activity_ids() { + let event = envelope(DocumentTopicEvent::Deleted(DocumentDeletedMetadata { + document_id: DOCUMENT_ID.to_string(), + actor_user_id: Some(user("macro|editor@example.com")), + project_id: None, + })); + + let first = single_activity(event.event.ingest(event.event_id)); + let second = single_activity(event.event.ingest(event.event_id)); + assert_eq!(first.id, second.id); +} diff --git a/crates/email/Cargo.toml b/crates/email/Cargo.toml index 17a8d41d534..d307b6b6505 100644 --- a/crates/email/Cargo.toml +++ b/crates/email/Cargo.toml @@ -40,6 +40,7 @@ outbound = ["dep:sqlx", "ports"] ports = ["crm/ports", "dep:tokio", "frecency/ports"] [dependencies] +activity = { path = "../activity", default-features = false } ai_toolset = { path = "../ai_toolset", optional = true } attachment = { path = "../attachment", optional = true } anyhow = { workspace = true } diff --git a/crates/email/src/domain.rs b/crates/email/src/domain.rs index b212a7dc2d2..bdda065ab01 100644 --- a/crates/email/src/domain.rs +++ b/crates/email/src/domain.rs @@ -1,3 +1,5 @@ +/// Event-to-activity mappings for this domain. +pub mod activity; pub mod events; pub mod models; diff --git a/crates/email/src/domain/activity.rs b/crates/email/src/domain/activity.rs new file mode 100644 index 00000000000..a61d2a2f770 --- /dev/null +++ b/crates/email/src/domain/activity.rs @@ -0,0 +1,131 @@ +//! What counts as activity in the email domain. +//! +//! Only user-initiated changes are activity: inbound mail, provider syncs, +//! and link lifecycle carry no acting user and are dropped. + +#[cfg(test)] +mod test; + +use ::activity::{ + Action, Activity, ActivitySource, Actor, CommonAction, DomainActivity, EntityType, Ingest, + event_time, +}; +use macro_user_id::user_id::MacroUserIdStr; +use uuid::Uuid; + +use super::events::{EmailEventOrigin, EmailTopicEvent}; + +/// Email-thread-exclusive actions. Common lifecycle actions go through +/// [`Activity::common`] and need no representation here. +#[derive(Debug, Clone, PartialEq)] +pub enum EmailThreadAction { + /// An email message was sent on the thread. + Sent, +} + +/// An email-thread-exclusive activity. +#[derive(Debug, Clone, PartialEq)] +pub struct EmailThreadActivity { + /// The thread acted on. + pub thread_id: String, + /// What happened to it. + pub action: EmailThreadAction, +} + +impl DomainActivity for EmailThreadActivity { + const ENTITY_TYPE: EntityType = EntityType::EmailThread; + + fn entity_id(&self) -> &str { + &self.thread_id + } + + fn into_action(self) -> Action { + match self.action { + EmailThreadAction::Sent => Action::Sent, + } + } +} + +impl ActivitySource for EmailTopicEvent { + /// Maps one `macro.email` event to its ingest outcome. + /// + /// Exhaustive on purpose: a new event variant fails compilation here + /// until someone classifies it or explicitly drops it. + fn ingest(&self, event_id: Uuid) -> Ingest { + let now = || event_time(event_id); + // A thread mutation counts only when a user did it directly. + let user_edit = |actor: &Option>, + origin: &EmailEventOrigin, + thread_id: Uuid| { + match (actor, origin) { + (Some(actor), EmailEventOrigin::UserAction) => { + Ingest::Insert(vec![Activity::common( + event_id, + 0, + Actor::new_from_user(actor.clone()), + None, + EntityType::EmailThread, + thread_id.to_string(), + CommonAction::Edited, + now(), + )]) + } + _ => Ingest::Ignore, + } + }; + + match self { + // Sent counts only when the user sent from Macro; a send synced + // from another client (provider_sync) still carries the actor + // but is not an act performed here. + EmailTopicEvent::MessageSent(m) => match (&m.actor, &m.origin) { + (Some(actor), EmailEventOrigin::UserAction) => { + Ingest::Insert(vec![Activity::from_domain( + event_id, + 0, + Actor::new_from_user(actor.clone()), + None, + EmailThreadActivity { + thread_id: m.thread_id.to_string(), + action: EmailThreadAction::Sent, + }, + now(), + )]) + } + _ => Ingest::Ignore, + }, + EmailTopicEvent::ThreadArchived(m) => user_edit(&m.actor, &m.origin, m.thread_id), + EmailTopicEvent::ThreadTrashed(m) => user_edit(&m.actor, &m.origin, m.thread_id), + EmailTopicEvent::ThreadStarred(m) => user_edit(&m.actor, &m.origin, m.thread_id), + EmailTopicEvent::ThreadSpamChanged(m) => user_edit(&m.actor, &m.origin, m.thread_id), + EmailTopicEvent::ThreadLabelsUpdated(m) => user_edit(&m.actor, &m.origin, m.thread_id), + EmailTopicEvent::ThreadProjectChanged(m) => Ingest::Insert(vec![Activity::common( + event_id, + 0, + Actor::new_from_user(m.actor.clone()), + None, + EntityType::EmailThread, + m.thread_id.to_string(), + CommonAction::Edited, + now(), + )]), + // Read-state, not a content mutation (and not view telemetry either — + // opens arrive via the interaction capability). + EmailTopicEvent::ThreadRead(_) => Ingest::Ignore, + // Resolved by a later message_sent; counting both double-reports. + EmailTopicEvent::MessageSendQueued(_) | EmailTopicEvent::MessageSendCancelled(_) => { + Ingest::Ignore + } + // No acting user: inbound mail, provider syncs, link lifecycle, + // backfill/reindex pipeline. + EmailTopicEvent::MessageReceived(_) + | EmailTopicEvent::MessageDraftSynced(_) + | EmailTopicEvent::MessageDeleted(_) + | EmailTopicEvent::LinkConnected(_) + | EmailTopicEvent::LinkDisconnected(_) + | EmailTopicEvent::LinkReauthRequired(_) + | EmailTopicEvent::ThreadBackfilled(_) + | EmailTopicEvent::ThreadsReindexRequested(_) => Ingest::Ignore, + } + } +} diff --git a/crates/email/src/domain/activity/test.rs b/crates/email/src/domain/activity/test.rs new file mode 100644 index 00000000000..540658ee970 --- /dev/null +++ b/crates/email/src/domain/activity/test.rs @@ -0,0 +1,89 @@ +use ::activity::Action; +use chrono::Utc; +use macro_user_id::user_id::MacroUserIdStr; +use model_entity::EntityType; +use uuid::Uuid; + +use macro_event_broker::Event; + +use super::*; +use crate::domain::events::{MessageSentMetadata, ThreadArchivedMetadata}; + +fn user(id: &str) -> MacroUserIdStr<'static> { + MacroUserIdStr::try_from(id.to_string()).expect("valid user id") +} + +fn envelope(event: EmailTopicEvent) -> Event { + Event::with_event_id(Uuid::now_v7(), event) +} + +const THREAD_ID: Uuid = Uuid::from_u128(9); + +fn sent(actor: Option>) -> MessageSentMetadata { + sent_with_origin(actor, EmailEventOrigin::UserAction) +} + +fn sent_with_origin( + actor: Option>, + origin: EmailEventOrigin, +) -> MessageSentMetadata { + MessageSentMetadata { + link_id: Uuid::from_u128(1), + owner: user("macro|owner@example.com"), + actor, + message_id: Uuid::from_u128(2), + provider_message_id: "pm".to_string(), + thread_id: THREAD_ID, + provider_thread_id: "pt".to_string(), + subject: None, + to_emails: vec![], + cc_emails: vec![], + sent_at: Utc::now(), + origin, + } +} + +#[test] +fn user_sent_message_maps_to_sent_activity() { + let event = envelope(EmailTopicEvent::MessageSent(sent(Some(user( + "macro|teo@example.com", + ))))); + let Ingest::Insert(activities) = event.event.ingest(event.event_id) else { + panic!("expected activities"); + }; + assert_eq!(activities[0].action, Action::Sent); + assert_eq!(activities[0].entity_type, EntityType::EmailThread); + assert_eq!(activities[0].entity_id, THREAD_ID.to_string()); +} + +#[test] +fn provider_synced_send_is_dropped_even_with_an_actor() { + let event = envelope(EmailTopicEvent::MessageSent(sent_with_origin( + Some(user("macro|teo@example.com")), + EmailEventOrigin::ProviderSync, + ))); + assert_eq!(event.event.ingest(event.event_id), Ingest::Ignore); +} + +#[test] +fn provider_sync_archive_is_dropped_but_user_archive_maps() { + let archive = |origin| { + envelope(EmailTopicEvent::ThreadArchived(ThreadArchivedMetadata { + link_id: Uuid::from_u128(1), + owner: user("macro|owner@example.com"), + actor: Some(user("macro|owner@example.com")), + thread_id: THREAD_ID, + archived: true, + origin, + })) + }; + + let user_action = archive(EmailEventOrigin::UserAction); + let Ingest::Insert(activities) = user_action.event.ingest(user_action.event_id) else { + panic!("expected activities"); + }; + assert_eq!(activities[0].action, Action::Edited); + + let provider = archive(EmailEventOrigin::ProviderSync); + assert_eq!(provider.event.ingest(provider.event_id), Ingest::Ignore); +} diff --git a/crates/macro_db_client/migrations/20260805180315_create_activity_events.sql b/crates/macro_db_client/migrations/20260805180315_create_activity_events.sql new file mode 100644 index 00000000000..7487701b4ec --- /dev/null +++ b/crates/macro_db_client/migrations/20260805180315_create_activity_events.sql @@ -0,0 +1,24 @@ +-- One append-only table of activity facts: a principal did something to an +-- entity at a time. Every activity surface is a query over this table; no +-- derived tables. See the activity crate for the fact vocabulary. +CREATE TABLE activity_events ( + -- uuidv5(source event id, ordinal): one broker event may yield several + -- facts; replays re-derive the same ids, making inserts idempotent. + id UUID PRIMARY KEY, + -- Who mechanically acted, as a prefixed principal string (macro|…, bot|…). + actor_id TEXT NOT NULL, + -- Whose activity this is: on_behalf_of ?? actor, resolved at ingestion. + subject_id TEXT NOT NULL, + action TEXT NOT NULL, + -- Per-action payload; NULL for payload-free actions. + action_payload JSONB, + entity_type TEXT NOT NULL, + entity_id TEXT NOT NULL, + occurred_at TIMESTAMPTZ NOT NULL +); + +-- id doubles as the keyset-pagination tiebreaker in both indexes. +CREATE INDEX idx_activity_events_subject + ON activity_events (subject_id, occurred_at DESC, id DESC); +CREATE INDEX idx_activity_events_entity + ON activity_events (entity_type, entity_id, occurred_at DESC, id DESC); diff --git a/crates/projects/Cargo.toml b/crates/projects/Cargo.toml index 84fefa1e3b4..c288b4d3d4a 100644 --- a/crates/projects/Cargo.toml +++ b/crates/projects/Cargo.toml @@ -40,6 +40,7 @@ ports = ["dep:tokio"] service = ["dep:entity_access_management", "ports"] [dependencies] +activity = { path = "../activity", default-features = false } # Core anyhow = { workspace = true } chrono = { workspace = true } diff --git a/crates/projects/src/domain.rs b/crates/projects/src/domain.rs index 9c75d16302a..114782e7352 100644 --- a/crates/projects/src/domain.rs +++ b/crates/projects/src/domain.rs @@ -2,6 +2,8 @@ /// Unified entity-mutation capability impls. #[cfg(feature = "service")] +/// Event-to-activity mappings for this domain. +pub mod activity; pub mod entity_mutation; pub mod events; pub mod models; diff --git a/crates/projects/src/domain/activity.rs b/crates/projects/src/domain/activity.rs new file mode 100644 index 00000000000..741d9c7a12e --- /dev/null +++ b/crates/projects/src/domain/activity.rs @@ -0,0 +1,99 @@ +//! What counts as activity in the projects domain. +//! +//! Projects have no entity-exclusive actions yet, so every mapping goes +//! through [`Activity::common`]. + +#[cfg(test)] +mod test; + +use ::activity::{Activity, ActivitySource, Actor, CommonAction, EntityType, Ingest, event_time}; +use chrono::{DateTime, Utc}; +use uuid::Uuid; + +use super::events::ProjectTopicEvent; + +impl ActivitySource for ProjectTopicEvent { + /// Maps one `macro.projects` event to its ingest outcome. + /// + /// Exhaustive on purpose: a new event variant fails compilation here + /// until someone classifies it or explicitly drops it. + fn ingest(&self, event_id: Uuid) -> Ingest { + let now = || event_time(event_id); + let single = + |actor: Actor<'static>, action: CommonAction, project_id: &str, at: DateTime| { + Ingest::Insert(vec![Activity::common( + event_id, + 0, + actor, + None, + EntityType::Project, + project_id, + action, + at, + )]) + }; + + match self { + ProjectTopicEvent::Created(m) => single( + Actor::new_from_user(m.owner.clone()), + CommonAction::Created, + &m.project_id, + now(), + ), + ProjectTopicEvent::Updated(m) => match &m.actor_user_id { + Some(actor) => single( + Actor::new_from_user(actor.clone()), + CommonAction::Edited, + &m.project_id, + now(), + ), + None => Ingest::Ignore, + }, + // Activity for the root project only; the cascade lists are + // soft-deleted rows whose activities we keep. + ProjectTopicEvent::Deleted(m) => match &m.actor_user_id { + Some(actor) => single( + Actor::new_from_user(actor.clone()), + CommonAction::Deleted, + &m.project_id, + now(), + ), + None => Ingest::Ignore, + }, + ProjectTopicEvent::Restored(m) => match &m.actor_user_id { + Some(actor) => single( + Actor::new_from_user(actor.clone()), + CommonAction::Edited, + &m.project_id, + now(), + ), + None => Ingest::Ignore, + }, + // Known window: the cascade's documents/chats produce activities on + // *other* topics with no cross-topic ordering, so a redelivered + // pre-purge event can re-insert an activity for a hard-deleted + // entity after this purge ran. Entity-owned purge events re-purge on + // their own topic's replays; a periodic reconciliation sweep is the + // durable fix if the residue ever matters. + ProjectTopicEvent::PermanentlyDeleted(m) => Ingest::Purge( + std::iter::once(&m.project_id) + .chain(m.purged_project_ids.iter()) + .map(|id| (EntityType::Project, id.clone())) + .chain( + m.purged_document_ids + .iter() + .map(|id| (EntityType::Document, id.clone())), + ) + .chain( + m.purged_chat_ids + .iter() + .map(|id| (EntityType::Chat, id.clone())), + ) + .collect(), + ), + // Bulk upload; per-project creation activities would misattribute — + // the upload flow publishes Created events for contained documents. + ProjectTopicEvent::Uploaded(_) => Ingest::Ignore, + } + } +} diff --git a/crates/projects/src/domain/activity/test.rs b/crates/projects/src/domain/activity/test.rs new file mode 100644 index 00000000000..a77643f49e8 --- /dev/null +++ b/crates/projects/src/domain/activity/test.rs @@ -0,0 +1,75 @@ +use ::activity::Action; +use macro_user_id::user_id::MacroUserIdStr; +use model_entity::EntityType; +use uuid::Uuid; + +use macro_event_broker::Event; + +use super::*; +use crate::domain::events::{ProjectDeletedMetadata, ProjectPermanentlyDeletedMetadata}; + +fn user(id: &str) -> MacroUserIdStr<'static> { + MacroUserIdStr::try_from(id.to_string()).expect("valid user id") +} + +fn envelope(event: ProjectTopicEvent) -> Event { + Event::with_event_id(Uuid::now_v7(), event) +} + +#[test] +fn attributed_delete_maps_and_unattributed_is_dropped() { + let attributed = envelope(ProjectTopicEvent::Deleted(ProjectDeletedMetadata { + project_id: "proj-1".to_string(), + owner: user("macro|owner@example.com"), + actor_user_id: Some(user("macro|teo@example.com")), + parent_project_id: None, + deleted_project_ids: vec!["proj-2".to_string()], + deleted_document_ids: vec![], + deleted_chat_ids: vec![], + })); + let Ingest::Insert(activities) = attributed.event.ingest(attributed.event_id) else { + panic!("expected activities"); + }; + assert_eq!(activities.len(), 1); + assert_eq!(activities[0].action, Action::Deleted); + assert_eq!(activities[0].entity_id, "proj-1"); + + let unattributed = envelope(ProjectTopicEvent::Deleted(ProjectDeletedMetadata { + project_id: "proj-1".to_string(), + owner: user("macro|owner@example.com"), + actor_user_id: None, + parent_project_id: None, + deleted_project_ids: vec![], + deleted_document_ids: vec![], + deleted_chat_ids: vec![], + })); + assert_eq!( + unattributed.event.ingest(unattributed.event_id), + Ingest::Ignore + ); +} + +#[test] +fn permanent_delete_purges_the_whole_cascade() { + let event = envelope(ProjectTopicEvent::PermanentlyDeleted( + ProjectPermanentlyDeletedMetadata { + project_id: "proj-1".to_string(), + owner: user("macro|owner@example.com"), + actor_user_id: None, + parent_project_id: None, + purged_project_ids: vec!["proj-2".to_string()], + purged_document_ids: vec!["doc-1".to_string()], + purged_chat_ids: vec!["chat-1".to_string()], + }, + )); + + assert_eq!( + event.event.ingest(event.event_id), + Ingest::Purge(vec![ + (EntityType::Project, "proj-1".to_string()), + (EntityType::Project, "proj-2".to_string()), + (EntityType::Document, "doc-1".to_string()), + (EntityType::Chat, "chat-1".to_string()), + ]) + ); +} diff --git a/crates/properties/Cargo.toml b/crates/properties/Cargo.toml index 6a45f4c910a..ff707602b47 100644 --- a/crates/properties/Cargo.toml +++ b/crates/properties/Cargo.toml @@ -9,14 +9,12 @@ ai_tools = [ "dep:ai_toolset", "dep:async-trait", "dep:schemars", - "dep:serde_json", "ports", ] default = ["ai_tools", "inbound", "outbound", "ports"] inbound = [ "dep:axum", "dep:macro_authorization", - "dep:serde_json", "dep:tokio", "dep:utoipa", "entity_access/inbound", @@ -30,13 +28,13 @@ outbound = [ "dep:model_notifications", "dep:models_permissions", "dep:notification", - "dep:serde_json", "dep:sqlx", "ports", ] ports = [] [dependencies] +activity = { path = "../activity", default-features = false } ai_toolset = { path = "../ai_toolset", optional = true } anyhow = { workspace = true } async-trait = { workspace = true, optional = true } @@ -62,7 +60,7 @@ models_properties = { path = "../models_properties" } notification = { path = "../notification", optional = true } schemars = { workspace = true, optional = true } serde = { workspace = true } -serde_json = { workspace = true, optional = true } +serde_json = { workspace = true } sqlx = { workspace = true, optional = true } system_properties = { path = "../system_properties" } thiserror = { workspace = true } diff --git a/crates/properties/src/domain.rs b/crates/properties/src/domain.rs index 0d4561eb5c9..eee685e5367 100644 --- a/crates/properties/src/domain.rs +++ b/crates/properties/src/domain.rs @@ -1,3 +1,5 @@ +/// Event-to-activity mappings for this domain. +pub mod activity; pub mod error; pub mod events; pub mod metadata; diff --git a/crates/properties/src/domain/activity.rs b/crates/properties/src/domain/activity.rs new file mode 100644 index 00000000000..fceca8fe3bc --- /dev/null +++ b/crates/properties/src/domain/activity.rs @@ -0,0 +1,116 @@ +//! What counts as activity in the properties domain. +//! +//! Every attributed property change on an entity is activity +//! ([`CommonAction::PropertyChanged`], valid on every entity kind). +//! Property/option *definition* changes are not entity activity and are +//! dropped. + +#[cfg(test)] +mod test; + +use ::activity::{ + Activity, ActivitySource, Actor, CommonAction, EntityType, Ingest, PropertyChange, event_time, +}; +use macro_user_id::user_id::MacroUserIdStr; +use models_properties::EntityType as PropertyEntityType; +use uuid::Uuid; + +use super::events::PropertyTopicEvent; + +/// Maps the properties-domain entity vocabulary onto the soup item-type +/// vocabulary activities use. `None` drops the activity (no soup surface). +fn entity_type(property_entity: &PropertyEntityType) -> Option { + match property_entity { + // Tasks are documents with a task sub-type. + PropertyEntityType::Document | PropertyEntityType::Task => Some(EntityType::Document), + PropertyEntityType::Project => Some(EntityType::Project), + PropertyEntityType::Chat => Some(EntityType::Chat), + PropertyEntityType::Channel => Some(EntityType::Channel), + PropertyEntityType::Thread => Some(EntityType::EmailThread), + PropertyEntityType::CallRecord => Some(EntityType::Call), + PropertyEntityType::Company => Some(EntityType::CrmCompany), + PropertyEntityType::CalendarEvent | PropertyEntityType::User => None, + } +} + +fn attributed( + event_id: uuid::Uuid, + actor: &Option>, + property_entity: &PropertyEntityType, + entity_id: &str, + action: CommonAction, + occurred_at: chrono::DateTime, +) -> Ingest { + let (Some(actor), Some(entity_type)) = (actor, entity_type(property_entity)) else { + return Ingest::Ignore; + }; + Ingest::Insert(vec![Activity::common( + event_id, + 0, + Actor::new_from_user(actor.clone()), + None, + entity_type, + entity_id, + action, + occurred_at, + )]) +} + +impl ActivitySource for PropertyTopicEvent { + /// Maps one `macro.properties` event to its ingest outcome. + /// + /// Exhaustive on purpose: a new event variant fails compilation here + /// until someone classifies it or explicitly drops it. + fn ingest(&self, event_id: Uuid) -> Ingest { + match self { + PropertyTopicEvent::EntityPropertyUpdated(m) => attributed( + event_id, + &m.actor_user_id, + &m.entity_type, + &m.entity_id, + CommonAction::PropertyChanged(PropertyChange { + property: m.property_definition_id.to_string(), + // The event carries only the new value today. + from: None, + // None means cleared; a serialization failure must not + // masquerade as one, so it stores an explicit JSON null. + to: m.value.as_ref().map(|value| { + serde_json::to_value(value).unwrap_or_else(|e| { + tracing::error!(error=?e, "unserializable property value"); + serde_json::Value::Null + }) + }), + }), + m.updated_at, + ), + PropertyTopicEvent::EntityPropertyDeleted(m) => attributed( + event_id, + &m.actor_user_id, + &m.entity_type, + &m.entity_id, + CommonAction::PropertyChanged(PropertyChange { + property: m.property_definition_id.to_string(), + from: None, + to: None, + }), + event_time(event_id), + ), + // A bulk clear has no per-property detail; it is still an attributed + // mutation of the entity. + PropertyTopicEvent::EntityPropertiesCleared(m) => attributed( + event_id, + &m.actor_user_id, + &m.entity_type, + &m.entity_id, + CommonAction::Edited, + event_time(event_id), + ), + // Definition/option lifecycle is not entity activity. + PropertyTopicEvent::Created(_) + | PropertyTopicEvent::Deleted(_) + | PropertyTopicEvent::OptionCreated(_) + | PropertyTopicEvent::OptionUpdated(_) + | PropertyTopicEvent::OptionDeleted(_) => Ingest::Ignore, + } + } +} diff --git a/crates/properties/src/domain/activity/test.rs b/crates/properties/src/domain/activity/test.rs new file mode 100644 index 00000000000..3e7d79c8193 --- /dev/null +++ b/crates/properties/src/domain/activity/test.rs @@ -0,0 +1,59 @@ +use ::activity::Action; +use chrono::Utc; +use macro_user_id::user_id::MacroUserIdStr; +use model_entity::EntityType as ActivityEntityType; +use uuid::Uuid; + +use macro_event_broker::Event; + +use super::*; +use crate::domain::events::EntityPropertyUpdatedMetadata; + +fn user(id: &str) -> MacroUserIdStr<'static> { + MacroUserIdStr::try_from(id.to_string()).expect("valid user id") +} + +fn envelope(event: PropertyTopicEvent) -> Event { + Event::with_event_id(Uuid::now_v7(), event) +} + +fn update(actor: Option>) -> EntityPropertyUpdatedMetadata { + EntityPropertyUpdatedMetadata { + entity_property_id: Uuid::from_u128(4), + entity_id: "task-1".to_string(), + entity_type: PropertyEntityType::Task, + property_definition_id: Uuid::from_u128(3), + actor_user_id: actor, + value: None, + updated_at: Utc::now(), + } +} + +#[test] +fn task_property_update_maps_to_property_changed_on_the_document() { + let event = envelope(PropertyTopicEvent::EntityPropertyUpdated(update(Some( + user("macro|seamus@example.com"), + )))); + + let Ingest::Insert(activities) = event.event.ingest(event.event_id) else { + panic!("expected activities"); + }; + // Tasks are documents in the soup vocabulary. + assert_eq!(activities[0].entity_type, ActivityEntityType::Document); + assert_eq!(activities[0].subject_id, "macro|seamus@example.com"); + match &activities[0].action { + Action::PropertyChanged(change) => { + assert_eq!(change.property, Uuid::from_u128(3).to_string()); + assert_eq!(change.from, None); + // The event's value was None: cleared without deleting the row. + assert_eq!(change.to, None); + } + other => panic!("expected property_changed, got {other:?}"), + } +} + +#[test] +fn unattributed_property_update_is_dropped() { + let event = envelope(PropertyTopicEvent::EntityPropertyUpdated(update(None))); + assert_eq!(event.event.ingest(event.event_id), Ingest::Ignore); +} diff --git a/services/document_storage_service/Cargo.toml b/services/document_storage_service/Cargo.toml index 4ecacbb7c49..4df3ad1c2ce 100644 --- a/services/document_storage_service/Cargo.toml +++ b/services/document_storage_service/Cargo.toml @@ -27,6 +27,7 @@ local_auth = ["macro_authorization/local_auth"] location_check = [] [dependencies] +activity = { path = "../../crates/activity" } ai_tools = { path = "../../crates/ai_tools" } ai_usage = { path = "../../crates/ai_usage" } analytics_client = { path = "../../crates/analytics_client" } diff --git a/services/document_storage_service/src/main.rs b/services/document_storage_service/src/main.rs index 9878ec202b4..d0c073ba749 100644 --- a/services/document_storage_service/src/main.rs +++ b/services/document_storage_service/src/main.rs @@ -789,6 +789,46 @@ async fn main() -> anyhow::Result<()> { } }); + let activity_consumer_brokers = config.kafka_brokers.as_ref().to_string(); + consumer_tracker.spawn({ + let cancellation_token = consumer_cancellation_token.clone(); + let activity_repo = activity::outbound::pg_activity_repo::PgActivityRepo::new(db.clone()); + async move { + let consumer = activity::inbound::kafka_consumer::ActivityConsumer::< + _, + crate::service::activity::ActivitySourceEvent, + _, + >::new(activity_repo, crate::service::activity::ingest); + loop { + if cancellation_token.is_cancelled() { + break; + } + + tracing::info!("starting activity consumer"); + let result = consumer + .run(&activity_consumer_brokers, cancellation_token.cancelled()) + .await; + + if cancellation_token.is_cancelled() { + break; + } + + match result { + Ok(()) => tracing::error!("activity consumer exited unexpectedly"), + Err(error) => { + tracing::error!(error = ?error, "activity consumer exited unexpectedly"); + } + } + + tokio::select! { + biased; + _ = cancellation_token.cancelled() => break, + _ = tokio::time::sleep(Duration::from_secs(5)) => {} + } + } + } + }); + let call_internal_state = InternalCallRouterState::new(call_service.clone()); // Create the SQS worker for delete document processing before config is moved. diff --git a/services/document_storage_service/src/service/activity.rs b/services/document_storage_service/src/service/activity.rs new file mode 100644 index 00000000000..a9249fce990 --- /dev/null +++ b/services/document_storage_service/src/service/activity.rs @@ -0,0 +1,52 @@ +//! Composition root for the activity consumer. +//! +//! DSS is the one place that legitimately knows every domain, so it declares +//! which topics feed activity and dispatches each decoded event to the +//! owning domain's mapping. All semantics live in the domain crates; all +//! machinery lives in the `activity` crate. These arms are pure wiring. + +use activity::Ingest; +use call::domain::events::CallMacroEvent; +use channels::domain::broker_events::ChannelMacroEvent; +use chat::domain::events::ChatMacroEvent; +use documents_hex::domain::events::DocumentMacroEvent; +use email::domain::events::EmailMacroEvent; +use macro_event_broker::MacroEvent as _; +use projects_hex::domain::events::ProjectMacroEvent; +use properties::domain::events::PropertyMacroEvent; + +#[allow(clippy::enum_variant_names)] +mod source { + use super::*; + + macro_event_broker::declare_topics!( + ActivitySourceEvent: + DocumentMacroEvent, + ChannelMacroEvent, + ChatMacroEvent, + ProjectMacroEvent, + EmailMacroEvent, + PropertyMacroEvent, + CallMacroEvent, + ); +} +pub(crate) use source::ActivitySourceEvent; + +/// Dispatches one decoded event to its domain's [`ActivitySource`] impl — +/// every arm is the identical expression; all semantics live with the +/// domains. +pub(crate) fn ingest(event: &ActivitySourceEvent) -> Ingest { + fn arm(envelope: ¯o_event_broker::Event) -> Ingest { + envelope.event.ingest(envelope.event_id) + } + + match event { + ActivitySourceEvent::DocumentMacroEvent(e) => arm(e.event()), + ActivitySourceEvent::ChannelMacroEvent(e) => arm(e.event()), + ActivitySourceEvent::ChatMacroEvent(e) => arm(e.event()), + ActivitySourceEvent::ProjectMacroEvent(e) => arm(e.event()), + ActivitySourceEvent::EmailMacroEvent(e) => arm(e.event()), + ActivitySourceEvent::PropertyMacroEvent(e) => arm(e.event()), + ActivitySourceEvent::CallMacroEvent(e) => arm(e.event()), + } +} diff --git a/services/document_storage_service/src/service/mod.rs b/services/document_storage_service/src/service/mod.rs index 167fbce1726..fa41cf3b941 100644 --- a/services/document_storage_service/src/service/mod.rs +++ b/services/document_storage_service/src/service/mod.rs @@ -1,3 +1,4 @@ +pub mod activity; pub mod conn_gateway; #[cfg(feature = "delete_document_worker")] pub mod delete_document_worker;