From f8a22a7a813e2952437f63bfb4e3db4412a92ba7 Mon Sep 17 00:00:00 2001 From: Neeraj Sathish Kumar Date: Mon, 4 May 2026 12:22:05 +0530 Subject: [PATCH] fix: enforce tenant/privacy scope in remote adapters --- src/adapters/neo4j.rs | 13 ++++++++----- src/adapters/qdrant.rs | 42 ++++++++++++++++++++++-------------------- src/adapters/s3.rs | 39 ++++++++++++++++++++------------------- 3 files changed, 50 insertions(+), 44 deletions(-) diff --git a/src/adapters/neo4j.rs b/src/adapters/neo4j.rs index c8bf875..bd8c7d7 100644 --- a/src/adapters/neo4j.rs +++ b/src/adapters/neo4j.rs @@ -24,7 +24,8 @@ impl Neo4jAdapter { } let hardening = TransportHardeningProfile::baseline(Some("NEXTRAL_NEO4J_API_KEY")); validate_transport_url( - &url.replace("neo4j://", "http://").replace("bolt://", "http://"), + &url.replace("neo4j://", "http://") + .replace("bolt://", "http://"), hardening.require_tls, ) .map_err(CoreError::InvalidInput)?; @@ -64,8 +65,8 @@ impl Neo4jAdapter { self.hardening.token_env.as_deref(), ) .json(&payload) - .send() - .map_err(|error| CoreError::Io(error.to_string()))?; + .send() + .map_err(|error| CoreError::Io(error.to_string()))?; if !response.status().is_success() { return Err(CoreError::Io(format!( "neo4j request failed: {}", @@ -141,7 +142,7 @@ impl Neo4jPort for Neo4jAdapter { } let statement = format!( r#" - MATCH (n:NextralEntity {{user_id:$user_id}}) + MATCH (n:NextralEntity {{tenant_id:$tenant_id, user_id:$user_id}}) WHERE any(term in $query_entities WHERE toLower(n.name) CONTAINS toLower(term)) MATCH p=(n)-[r:NEXTRAL_RELATES_TO*1..{}]-() UNWIND relationships(p) as rel @@ -154,6 +155,7 @@ impl Neo4jPort for Neo4jAdapter { let body = self.cypher( &statement, json!({ + "tenant_id": scope.tenant_id, "user_id": scope.user_id, "query_entities": query_entities, }), @@ -171,7 +173,7 @@ impl Neo4jPort for Neo4jAdapter { fn redact_memory_edges(&self, scope: &TenantUserScope, memory_id: &str) -> CoreResult<()> { let statement = r#" - MATCH ()-[r:NEXTRAL_RELATES_TO {user_id:$user_id}]-() + MATCH ()-[r:NEXTRAL_RELATES_TO {tenant_id:$tenant_id, user_id:$user_id}]-() WHERE any(id in r.source_memory_ids WHERE id = $memory_id) SET r.source_memory_ids = [id IN r.source_memory_ids WHERE id <> $memory_id] WITH r @@ -181,6 +183,7 @@ impl Neo4jPort for Neo4jAdapter { self.cypher( statement, json!({ + "tenant_id": scope.tenant_id, "user_id": scope.user_id, "memory_id": memory_id }), diff --git a/src/adapters/qdrant.rs b/src/adapters/qdrant.rs index 2b143ee..7b47f45 100644 --- a/src/adapters/qdrant.rs +++ b/src/adapters/qdrant.rs @@ -1,5 +1,8 @@ use crate::{ - adapters::{transport::{maybe_add_bearer_auth, validate_transport_url, TransportHardeningProfile}, AdapterHealth}, + adapters::{ + transport::{maybe_add_bearer_auth, validate_transport_url, TransportHardeningProfile}, + AdapterHealth, + }, contracts::{CoreError, CoreResult}, ports::{QdrantPort, VectorPoint, VectorSearchHit, VectorSearchRequest}, }; @@ -24,9 +27,12 @@ impl QdrantAdapter { )); } let hardening = TransportHardeningProfile::baseline(Some("NEXTRAL_QDRANT_API_KEY")); - validate_transport_url(&url, hardening.require_tls) - .map_err(CoreError::InvalidInput)?; - Ok(Self { url, collection, hardening }) + validate_transport_url(&url, hardening.require_tls).map_err(CoreError::InvalidInput)?; + Ok(Self { + url, + collection, + hardening, + }) } pub fn collection_schema_json(&self) -> &'static str { @@ -47,8 +53,8 @@ impl QdrantAdapter { client.get(format!("{}/collections", self.url.trim_end_matches('/'))), self.hardening.token_env.as_deref(), ) - .send() - .map_err(|error| CoreError::Io(error.to_string()))?; + .send() + .map_err(|error| CoreError::Io(error.to_string()))?; let status = response.status(); let body = response .json::() @@ -83,8 +89,8 @@ impl QdrantPort for QdrantAdapter { self.hardening.token_env.as_deref(), ) .json(&payload) - .send() - .map_err(|error| CoreError::Io(error.to_string()))?; + .send() + .map_err(|error| CoreError::Io(error.to_string()))?; if !response.status().is_success() { return Err(CoreError::Io(format!( "qdrant ensure_collection failed: {}", @@ -119,8 +125,8 @@ impl QdrantPort for QdrantAdapter { self.hardening.token_env.as_deref(), ) .json(&payload) - .send() - .map_err(|error| CoreError::Io(error.to_string()))?; + .send() + .map_err(|error| CoreError::Io(error.to_string()))?; if !response.status().is_success() { return Err(CoreError::Io(format!( "qdrant upsert_point failed: {}", @@ -144,6 +150,7 @@ impl QdrantPort for QdrantAdapter { { "key": "tenant_id", "match": { "value": request.scope.tenant_id }}, { "key": "user_id", "match": { "value": request.scope.user_id }}, { "key": "status", "match": { "value": "active" }}, + { "key": "privacy_level", "match": { "any": request.privacy_scope }}, ] } }); @@ -156,8 +163,8 @@ impl QdrantPort for QdrantAdapter { self.hardening.token_env.as_deref(), ) .json(&payload) - .send() - .map_err(|error| CoreError::Io(error.to_string()))?; + .send() + .map_err(|error| CoreError::Io(error.to_string()))?; if !response.status().is_success() { return Err(CoreError::Io(format!( "qdrant search failed: {}", @@ -182,12 +189,7 @@ impl QdrantPort for QdrantAdapter { Ok(hits) } - fn delete_point( - &self, - collection: &str, - tenant_id: &str, - memory_id: &str, - ) -> CoreResult<()> { + fn delete_point(&self, collection: &str, tenant_id: &str, memory_id: &str) -> CoreResult<()> { let payload = json!({ "points": [memory_id], "filter": { @@ -205,8 +207,8 @@ impl QdrantPort for QdrantAdapter { self.hardening.token_env.as_deref(), ) .json(&payload) - .send() - .map_err(|error| CoreError::Io(error.to_string()))?; + .send() + .map_err(|error| CoreError::Io(error.to_string()))?; if !response.status().is_success() { return Err(CoreError::Io(format!( "qdrant delete_point failed: {}", diff --git a/src/adapters/s3.rs b/src/adapters/s3.rs index f95f569..aa5f9ec 100644 --- a/src/adapters/s3.rs +++ b/src/adapters/s3.rs @@ -68,8 +68,8 @@ impl S3Adapter { )), self.hardening.token_env.as_deref(), ) - .send() - .map_err(|error| CoreError::Io(error.to_string()))?; + .send() + .map_err(|error| CoreError::Io(error.to_string()))?; Ok(serde_json::json!({ "status": response.status().as_u16(), "bucket": self.bucket, @@ -98,15 +98,15 @@ impl ObjectArchivePort for S3Adapter { self.hardening.token_env.as_deref(), ) .body(object.bytes.clone()) - .header("x-amz-meta-tenant-id", object.tenant_id.clone()) - .header("x-amz-meta-user-id", object.user_id.clone()) - .header( - "x-amz-meta-memory-id", - object.memory_id.clone().unwrap_or_default(), - ) - .header("x-amz-meta-content-sha256", object.content_sha256.clone()) - .send() - .map_err(|error| CoreError::Io(error.to_string()))?; + .header("x-amz-meta-tenant-id", object.tenant_id.clone()) + .header("x-amz-meta-user-id", object.user_id.clone()) + .header( + "x-amz-meta-memory-id", + object.memory_id.clone().unwrap_or_default(), + ) + .header("x-amz-meta-content-sha256", object.content_sha256.clone()) + .send() + .map_err(|error| CoreError::Io(error.to_string()))?; if !response.status().is_success() { return Err(CoreError::Io(format!( "s3 put_object failed: {}", @@ -120,12 +120,13 @@ impl ObjectArchivePort for S3Adapter { }) } - fn tombstone_object( - &self, - _tenant_id: &str, - object_key: &str, - _reason: &str, - ) -> CoreResult<()> { + fn tombstone_object(&self, tenant_id: &str, object_key: &str, _reason: &str) -> CoreResult<()> { + let tenant_prefix = format!("tenants/{}/", tenant_id); + if !object_key.starts_with(&tenant_prefix) { + return Err(CoreError::InvalidInput( + "object_key is outside tenant scope".to_string(), + )); + } let response = maybe_add_bearer_auth( Client::new().delete(format!( "{}/{}/{}", @@ -135,8 +136,8 @@ impl ObjectArchivePort for S3Adapter { )), self.hardening.token_env.as_deref(), ) - .send() - .map_err(|error| CoreError::Io(error.to_string()))?; + .send() + .map_err(|error| CoreError::Io(error.to_string()))?; if !response.status().is_success() { return Err(CoreError::Io(format!( "s3 tombstone_object failed: {}",