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
13 changes: 8 additions & 5 deletions src/adapters/neo4j.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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)?;
Expand Down Expand Up @@ -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: {}",
Expand Down Expand Up @@ -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
Expand All @@ -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,
}),
Expand All @@ -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
Expand All @@ -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
}),
Expand Down
42 changes: 22 additions & 20 deletions src/adapters/qdrant.rs
Original file line number Diff line number Diff line change
@@ -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},
};
Expand All @@ -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 {
Expand All @@ -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::<Value>()
Expand Down Expand Up @@ -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: {}",
Expand Down Expand Up @@ -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: {}",
Expand All @@ -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 }},
]
}
});
Expand All @@ -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: {}",
Expand All @@ -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": {
Expand All @@ -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: {}",
Expand Down
47 changes: 19 additions & 28 deletions src/adapters/s3.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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: {}",
Expand All @@ -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!(
"{}/{}/{}",
Expand All @@ -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: {}",
Expand Down
Loading