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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 4 additions & 3 deletions omem-server/src/api/handlers/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -19,9 +19,10 @@ pub use memory::{
};
pub use profile::get_profile;
pub use sharing::{
batch_share, create_auto_share_rule, delete_auto_share_rule, list_auto_share_rules,
org_publish, org_setup, pull_memory, reshare_memory, share_all, share_all_to_user,
share_memory, share_to_user, unshare_memory,
approve_pending_share, batch_share, create_auto_share_rule, delete_auto_share_rule,
list_auto_share_rules, list_pending_shares, org_publish, org_setup, pull_memory,
reject_pending_share, reshare_memory, share_all, share_all_to_user, share_memory,
share_to_user, unshare_memory,
};
pub use spaces::{
add_member, create_space, delete_space, get_space, list_spaces, remove_member,
Expand Down
216 changes: 213 additions & 3 deletions omem-server/src/api/handlers/sharing.rs
Original file line number Diff line number Diff line change
Expand Up @@ -11,8 +11,8 @@ use crate::api::server::{normalize_space_id, personal_space_id, AppState};
use crate::domain::error::OmemError;
use crate::domain::memory::Memory;
use crate::domain::space::{
AutoShareRule, MemberRole, Provenance, SharingAction, SharingEvent, Space, SpaceMember,
SpaceType,
AutoShareRule, MemberRole, PendingShare, Provenance, SharingAction, SharingEvent, Space,
SpaceMember, SpaceType,
};
use crate::domain::tenant::AuthInfo;
use crate::store::StoreManager;
Expand Down Expand Up @@ -1357,7 +1357,19 @@ pub async fn check_auto_share(
continue;
}
if rule.require_approval {
continue;
let pending = PendingShare {
id: Uuid::new_v4().to_string(),
source_space: memory.space_id.clone(),
source_memory: memory.id.clone(),
target_space: space.id.clone(),
rule_id: rule.id.clone(),
requested_by_user: user_id.to_string(),
requested_by_agent: agent_id.to_string(),
content_preview: content_preview(&memory.content),
created_at: chrono::Utc::now().to_rfc3339(),
};
space_store.record_pending_share(&pending).await?;
break;
}

let target_store = store_manager.get_store(&space.id).await?;
Expand Down Expand Up @@ -1385,6 +1397,117 @@ pub async fn check_auto_share(
Ok(shared_to)
}

/// Materialise an approved pending share: copy the source memory (content +
/// vector) into the target space, record a Share event, and drop the queue
/// entry. Returns the new copy. Kept free of `AppState` so it stays unit-testable.
pub(crate) async fn apply_pending_share(
store_manager: &StoreManager,
space_store: &crate::store::SpaceStore,
pending: &PendingShare,
) -> Result<Memory, OmemError> {
let source_store = store_manager.get_store(&pending.source_space).await?;
let source = source_store
.get_by_id(&pending.source_memory)
.await?
.ok_or_else(|| {
OmemError::NotFound(format!(
"source memory {} no longer exists",
pending.source_memory
))
})?;
let source_vector = source_store
.get_vector_by_id(&pending.source_memory)
.await?;

let copy = make_shared_copy(
&source,
&pending.target_space,
&pending.requested_by_user,
&pending.requested_by_agent,
);
let target_store = store_manager.get_store(&pending.target_space).await?;
target_store.create(&copy, source_vector.as_deref()).await?;

let event = make_sharing_event(
SharingAction::Share,
&copy.id,
&pending.source_space,
&pending.target_space,
&pending.requested_by_user,
&pending.requested_by_agent,
&content_preview(&source.content),
);
space_store.record_sharing_event(&event).await?;
space_store.delete_pending_share(&pending.id).await?;

Ok(copy)
}

/// GET /v1/shares/pending
///
/// List shares awaiting approval in spaces the caller can write to. Poll this;
/// there is no notification system.
pub async fn list_pending_shares(
State(state): State<Arc<AppState>>,
Extension(auth): Extension<AuthInfo>,
) -> Result<Json<Vec<PendingShare>>, OmemError> {
let spaces = state
.space_store
.list_spaces_for_user(&auth.tenant_id)
.await?;
let mut pending = Vec::new();
for space in &spaces {
if verify_space_write_access(space, &auth.tenant_id).is_ok() {
pending.extend(state.space_store.list_pending_shares(&space.id).await?);
}
}
Ok(Json(pending))
}

/// POST /v1/shares/pending/{id}/approve
pub async fn approve_pending_share(
State(state): State<Arc<AppState>>,
Extension(auth): Extension<AuthInfo>,
Path(id): Path<String>,
) -> Result<Json<Memory>, OmemError> {
let pending = state
.space_store
.get_pending_share(&id)
.await?
.ok_or_else(|| OmemError::NotFound(format!("pending share {id}")))?;
let space = state
.space_store
.get_space(&pending.target_space)
.await?
.ok_or_else(|| OmemError::NotFound(format!("space {}", pending.target_space)))?;
verify_space_write_access(&space, &auth.tenant_id)?;

let copy = apply_pending_share(&state.store_manager, &state.space_store, &pending).await?;
Ok(Json(copy))
}

/// POST /v1/shares/pending/{id}/reject
pub async fn reject_pending_share(
State(state): State<Arc<AppState>>,
Extension(auth): Extension<AuthInfo>,
Path(id): Path<String>,
) -> Result<Json<serde_json::Value>, OmemError> {
let pending = state
.space_store
.get_pending_share(&id)
.await?
.ok_or_else(|| OmemError::NotFound(format!("pending share {id}")))?;
let space = state
.space_store
.get_space(&pending.target_space)
.await?
.ok_or_else(|| OmemError::NotFound(format!("space {}", pending.target_space)))?;
verify_space_write_access(&space, &auth.tenant_id)?;

state.space_store.delete_pending_share(&id).await?;
Ok(Json(serde_json::json!({ "rejected": true, "id": id })))
}

// ── Tests ────────────────────────────────────────────────────────────

#[cfg(test)]
Expand Down Expand Up @@ -1594,6 +1717,93 @@ mod tests {
assert_eq!(team_list.len(), 3);
}

#[tokio::test]
async fn test_pending_share_approval_flow() {
let env = setup().await;
let dim = env.store_manager.vector_dim() as usize;

// Team space with an auto-share rule that REQUIRES approval.
let mut team_space = make_space("team:backend", "user-001");
team_space.auto_share_rules.push(AutoShareRule {
id: "rule-1".to_string(),
source_space: "user-001".to_string(),
categories: vec!["preferences".to_string()],
tags: Vec::new(),
min_importance: 0.0,
require_approval: true,
created_at: "2025-01-01T00:00:00Z".to_string(),
});
env.space_store
.create_space(&team_space)
.await
.expect("create space");

// Source memory (with a vector) that matches the rule.
let source_store = env
.store_manager
.get_store("user-001")
.await
.expect("source store");
let mem = make_memory("prefers tabs over spaces", "user-001", "user-001");
source_store
.create(&mem, Some(&vec![0.3f32; dim]))
.await
.expect("create source");

// require_approval => ENQUEUE, do not auto-share.
let shared_to = check_auto_share(
&mem,
&env.space_store,
&env.store_manager,
"user-001",
"agent-1",
)
.await
.expect("auto share");
assert!(shared_to.is_empty(), "require_approval must not auto-share");

let team_store = env
.store_manager
.get_store("team:backend")
.await
.expect("team store");
assert_eq!(
team_store.list_all_active().await.expect("list").len(),
0,
"nothing shared before approval"
);

let pending = env
.space_store
.list_pending_shares("team:backend")
.await
.expect("list pending");
assert_eq!(pending.len(), 1);
assert_eq!(pending[0].source_memory, mem.id);
assert_eq!(pending[0].target_space, "team:backend");

// Approve => copy materialises (with vector) and the queue drains.
let copy = apply_pending_share(&env.store_manager, &env.space_store, &pending[0])
.await
.expect("apply");
assert_eq!(copy.content, "prefers tabs over spaces");

let active = team_store.list_all_active().await.expect("list");
assert_eq!(active.len(), 1);
let v = team_store
.get_vector_by_id(&active[0].id)
.await
.expect("vec")
.expect("vector present after approval");
assert!(v.iter().any(|x| *x != 0.0));
assert!(env
.space_store
.list_pending_shares("team:backend")
.await
.expect("list")
.is_empty());
}

#[tokio::test]
async fn test_auto_share_rule() {
let env = setup().await;
Expand Down
9 changes: 9 additions & 0 deletions omem-server/src/api/router.rs
Original file line number Diff line number Diff line change
Expand Up @@ -70,6 +70,15 @@ pub fn build_router(state: Arc<AppState>) -> Router {
.route("/v1/memories/{id}/pull", post(handlers::pull_memory))
.route("/v1/memories/{id}/unshare", post(handlers::unshare_memory))
.route("/v1/memories/{id}/reshare", post(handlers::reshare_memory))
.route("/v1/shares/pending", get(handlers::list_pending_shares))
.route(
"/v1/shares/pending/{id}/approve",
post(handlers::approve_pending_share),
)
.route(
"/v1/shares/pending/{id}/reject",
post(handlers::reject_pending_share),
)
.route("/v1/memories/batch-share", post(handlers::batch_share))
.route("/v1/memories/share-all", post(handlers::share_all))
.route(
Expand Down
16 changes: 16 additions & 0 deletions omem-server/src/domain/space.rs
Original file line number Diff line number Diff line change
Expand Up @@ -131,6 +131,22 @@ pub struct SharingEvent {
pub timestamp: String,
}

/// A share that matched an auto-share rule with `require_approval = true`,
/// awaiting a human decision. Held in a queue (no notifications) until an
/// approver with write access to `target_space` approves or rejects it.
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct PendingShare {
pub id: String,
pub source_space: String,
pub source_memory: String,
pub target_space: String,
pub rule_id: String,
pub requested_by_user: String,
pub requested_by_agent: String,
pub content_preview: String,
pub created_at: String,
}

#[cfg(test)]
mod tests {
use super::*;
Expand Down
Loading
Loading