diff --git a/src/adapters/neo4j.rs b/src/adapters/neo4j.rs index 03b9dd9..f0c5b82 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 887d1e8..8cebb12 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 { @@ -51,8 +57,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::() @@ -87,8 +93,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: {}", @@ -123,8 +129,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: {}", @@ -148,6 +154,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 }}, ] } }); @@ -160,8 +167,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: {}", @@ -186,12 +193,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": { @@ -209,8 +211,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 7fa36d2..8136632 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, @@ -107,15 +107,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", computed_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: {}", @@ -129,21 +129,13 @@ impl ObjectArchivePort for S3Adapter { }) } - fn tombstone_object( - &self, - tenant_id: &str, - object_key: &str, - reason: &str, - ) -> CoreResult<()> { - // Validate that object_key belongs to the tenant - let expected_prefix = format!("tenants/{}/", tenant_id); - if !object_key.starts_with(&expected_prefix) { - return Err(CoreError::InvalidInput(format!( - "object_key '{}' does not belong to tenant '{}'", - object_key, tenant_id - ))); + 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!( "{}/{}/{}", @@ -153,9 +145,8 @@ impl ObjectArchivePort for S3Adapter { )), self.hardening.token_env.as_deref(), ) - .header("x-amz-meta-tombstone-reason", reason) - .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: {}",