diff --git a/src/api/identity_binding.rs b/src/api/identity_binding.rs index d46e20a..a72a2a1 100644 --- a/src/api/identity_binding.rs +++ b/src/api/identity_binding.rs @@ -304,24 +304,31 @@ pub async fn unbind_identity( ); } - // Update TKG: restore stranger_ref and remove identity fields + // Update TKG: restore stranger_ref and remove identity fields (match both external_id formats) if let Some(tid) = trace_id { let tkg_table = crate::core::db::schema::table_name("tkg_nodes"); - let ext_id = format!("face_track_{}", tid); + let ext_id_face = format!("face_track_{}", tid); + let ext_id_trace = format!("trace_{}", tid); let stranger_ref = format!("{}:stranger_trace_{}", req.file_uuid, tid); - let _ = sqlx::query(&format!( + let result = sqlx::query(&format!( "UPDATE {} SET properties = properties || $1::jsonb - 'identity_uuid' - 'identity_ref' - 'identity_id' - 'identity_name' \ - WHERE file_uuid = $2 AND node_type = 'face_track' AND external_id = $3", + WHERE file_uuid = $2 AND node_type = 'face_track' AND (external_id = $3 OR external_id = $4)", tkg_table )) .bind(serde_json::json!({ "stranger_ref": stranger_ref })) .bind(&req.file_uuid) - .bind(&ext_id) + .bind(&ext_id_face) + .bind(&ext_id_trace) .execute(state.db.pool()) .await; + + match &result { + Ok(r) => tracing::info!("[unbind_identity] TKG update: {} rows affected for trace {}", r.rows_affected(), tid), + Err(e) => tracing::error!("[unbind_identity] TKG update failed: {}", e), + } } if let Some(identity_id) = old_identity_id { @@ -887,9 +894,10 @@ pub async fn bind_identity_trace( ); } - // Update TKG node properties + // Update TKG node properties (match both face_track_N and trace_N external_id formats) let nodes_table = crate::core::db::schema::table_name("tkg_nodes"); - let external_id = format!("face_track_{}", req.trace_id); + let external_id_face = format!("face_track_{}", req.trace_id); + let external_id_trace = format!("trace_{}", req.trace_id); let identity_name: Option = sqlx::query_scalar(&format!("SELECT name FROM {} WHERE id = $1", id_table)) .bind(identity_id) @@ -898,20 +906,27 @@ pub async fn bind_identity_trace( .ok() .flatten(); - let _ = sqlx::query(&format!( - "UPDATE {} SET properties = jsonb_set(\ - jsonb_set(properties, '{{identity_id}}', $1::jsonb, false),\ - '{{identity_name}}', $2::jsonb, false)\ - WHERE file_uuid = $3 AND node_type = 'face_track' AND external_id = $4", + let new_props = serde_json::json!({ + "identity_id": identity_id, + "identity_name": identity_name + }); + let result = sqlx::query(&format!( + "UPDATE {} SET properties = properties || $1::jsonb \ + WHERE file_uuid = $2 AND node_type = 'face_track' AND (external_id = $3 OR external_id = $4)", nodes_table )) - .bind(identity_id) - .bind(identity_name.as_deref()) + .bind(serde_json::to_string(&new_props).unwrap()) .bind(&req.file_uuid) - .bind(&external_id) + .bind(&external_id_face) + .bind(&external_id_trace) .execute(state.db.pool()) .await; + match &result { + Ok(r) => tracing::info!("[bind_identity_trace] TKG update: {} rows affected for trace {}", r.rows_affected(), req.trace_id), + Err(e) => tracing::error!("[bind_identity_trace] TKG update failed: {}", e), + } + // Clear bind redo stack let _ = sqlx::query(&format!( "DELETE FROM {} WHERE identity_id = $1 AND is_undone = true AND operation IN ('bind','unbind','bind_trace')", @@ -2040,20 +2055,23 @@ pub async fn create_pending_person( } } - // Update TKG nodes + // Update TKG nodes (match both face_track_N and trace_N external_id formats) for &tid in &req.trace_ids { - let external_id = format!("face_track_{}", tid); + let external_id_face = format!("face_track_{}", tid); + let external_id_trace = format!("trace_{}", tid); + let new_props = serde_json::json!({ + "identity_id": identity_id, + "identity_name": name + }); let _ = sqlx::query(&format!( - "UPDATE {} SET properties = jsonb_set(\ - jsonb_set(properties, '{{identity_id}}', $1::jsonb, false),\ - '{{identity_name}}', $2::jsonb, false)\ - WHERE file_uuid = $3 AND node_type = 'face_track' AND external_id = $4", + "UPDATE {} SET properties = properties || $1::jsonb \ + WHERE file_uuid = $2 AND node_type = 'face_track' AND (external_id = $3 OR external_id = $4)", nodes_table )) - .bind(identity_id) - .bind(&name) + .bind(serde_json::to_string(&new_props).unwrap()) .bind(&file_uuid) - .bind(&external_id) + .bind(&external_id_face) + .bind(&external_id_trace) .execute(state.db.pool()) .await; } diff --git a/src/api/profile.rs b/src/api/profile.rs index 0324903..7bedbb5 100644 --- a/src/api/profile.rs +++ b/src/api/profile.rs @@ -55,15 +55,17 @@ pub async fn get_trace_profile_handler( Query(params): Query, ) -> Result, StatusCode> { let tkg_table = schema::table_name("tkg_nodes"); - let external_id = format!("face_track_{}", params.trace_id); + let external_id_face = format!("face_track_{}", params.trace_id); + let external_id_trace = format!("trace_{}", params.trace_id); let row: Option<(String, String, serde_json::Value)> = sqlx::query_as(&format!( "SELECT label, external_id, properties FROM {} \ - WHERE file_uuid = $1 AND node_type = 'face_track' AND external_id = $2", + WHERE file_uuid = $1 AND node_type = 'face_track' AND (external_id = $2 OR external_id = $3)", tkg_table )) .bind(¶ms.file_uuid) - .bind(&external_id) + .bind(&external_id_face) + .bind(&external_id_trace) .fetch_optional(state.db.pool()) .await .map_err(|e| { @@ -99,16 +101,18 @@ pub async fn update_trace_profile_handler( Json(req): Json, ) -> Result, StatusCode> { let tkg_table = schema::table_name("tkg_nodes"); - let external_id = format!("face_track_{}", req.trace_id); + let external_id_face = format!("face_track_{}", req.trace_id); + let external_id_trace = format!("trace_{}", req.trace_id); - // Get current node + // Get current node (match both external_id formats) let current: Option<(String, serde_json::Value)> = sqlx::query_as(&format!( "SELECT label, properties FROM {} \ - WHERE file_uuid = $1 AND node_type = 'face_track' AND external_id = $2", + WHERE file_uuid = $1 AND node_type = 'face_track' AND (external_id = $2 OR external_id = $3)", tkg_table )) .bind(&req.file_uuid) - .bind(&external_id) + .bind(&external_id_face) + .bind(&external_id_trace) .fetch_optional(state.db.pool()) .await .map_err(|e| { @@ -123,7 +127,7 @@ pub async fn update_trace_profile_handler( if let Some(ref new_name) = req.name { if new_name != ¤t_label { - updates.push(format!("label = $3")); + updates.push(format!("label = $5")); } } @@ -158,31 +162,34 @@ pub async fn update_trace_profile_handler( } let mut query = format!( - "UPDATE {} SET properties = $4", + "UPDATE {} SET properties = properties || $4::jsonb", tkg_table ); if !updates.is_empty() { query.push_str(", "); query.push_str(&updates.join(", ")); } - query.push_str(" WHERE file_uuid = $1 AND node_type = 'face_track' AND external_id = $2"); + query.push_str(" WHERE file_uuid = $1 AND node_type = 'face_track' AND (external_id = $2 OR external_id = $3)"); - let param_idx = if updates.is_empty() { 3 } else { 4 }; - query.push_str(&format!(" RETURNING id")); + query.push_str(" RETURNING id"); + + let props_json = serde_json::to_string(¤t_props).unwrap(); let result = if !updates.is_empty() { sqlx::query(&query) .bind(&req.file_uuid) - .bind(&external_id) + .bind(&external_id_face) + .bind(&external_id_trace) + .bind(&props_json) .bind(req.name.as_ref().unwrap_or(¤t_label)) - .bind(¤t_props) .execute(state.db.pool()) .await } else { sqlx::query(&query) .bind(&req.file_uuid) - .bind(&external_id) - .bind(¤t_props) + .bind(&external_id_face) + .bind(&external_id_trace) + .bind(&props_json) .execute(state.db.pool()) .await }; @@ -211,15 +218,17 @@ pub async fn update_trace_profile_group_handler( let mut updated = 0; for trace_id in &req.trace_ids { - let external_id = format!("face_track_{}", trace_id); + let external_id_face = format!("face_track_{}", trace_id); + let external_id_trace = format!("trace_{}", trace_id); let result = sqlx::query(&format!( "UPDATE {} SET label = $1 \ - WHERE file_uuid = $2 AND node_type = 'face_track' AND external_id = $3", + WHERE file_uuid = $2 AND node_type = 'face_track' AND (external_id = $3 OR external_id = $4)", tkg_table )) .bind(&req.name) .bind(&req.file_uuid) - .bind(&external_id) + .bind(&external_id_face) + .bind(&external_id_trace) .execute(state.db.pool()) .await; diff --git a/src/api/trace_agent_api.rs b/src/api/trace_agent_api.rs index fecf4f3..2677355 100644 --- a/src/api/trace_agent_api.rs +++ b/src/api/trace_agent_api.rs @@ -1,6 +1,6 @@ use axum::{ body::Body, - extract::{Path, Query, State}, + extract::{Extension, Path, Query, State}, http::{header, StatusCode}, response::{IntoResponse, Json, Response}, routing::{delete, get, post}, @@ -53,6 +53,14 @@ pub fn trace_agent_routes() -> Router { "/api/v1/file/:file_uuid/tkg/node/:node_id", get(get_tkg_node_detail), ) + .route( + "/api/v1/file/:file_uuid/trace/:trace_id", + delete(delete_trace), + ) + .route( + "/api/v1/file/:file_uuid/trace/:trace_id/merge/:target_trace_id", + post(merge_trace), + ) } #[derive(Debug, Deserialize)] @@ -132,6 +140,9 @@ async fn list_traces_sorted( let face_filter = json!({ "must": [ {"key": "file_uuid", "match": {"value": file_uuid}} + ], + "must_not": [ + {"key": "status", "match": {"value": "deleted"}} ] }); let points = qdrant @@ -1685,3 +1696,205 @@ async fn ingest_rule2( })), } } + +async fn delete_trace( + State(state): State, + Path((file_uuid, trace_id)): Path<(String, i32)>, + Json(req): Json, +) -> Json { + let hard_delete = req.get("hard_delete").and_then(|v| v.as_bool()).unwrap_or(false); + + let qdrant = crate::core::db::qdrant_db::QdrantDb::new(); + + if hard_delete { + // Hard delete: actually remove from Qdrant and TKG + let filter = serde_json::json!({ + "must": [ + {"key": "file_uuid", "match": {"value": file_uuid}}, + {"key": "trace_id", "match": {"value": trace_id}} + ] + }); + + let deleted_count = match qdrant.delete_points_by_filter("_faces", filter.clone()).await { + Ok(count) => { + tracing::info!("[delete_trace] Hard deleted {} face points from Qdrant for trace {} in {}", count, trace_id, file_uuid); + count + } + Err(e) => { + tracing::error!("[delete_trace] Failed to delete Qdrant points: {}", e); + return Json(serde_json::json!({ + "success": false, + "error": format!("Failed to delete Qdrant points: {}", e) + })); + } + }; + + let nodes_table = crate::core::db::schema::table_name("tkg_nodes"); + let ext_id_face = format!("face_track_{}", trace_id); + let ext_id_trace = format!("trace_{}", trace_id); + + let tkg_result = sqlx::query(&format!( + "DELETE FROM {} WHERE file_uuid = $1 AND node_type = 'face_track' AND (external_id = $2 OR external_id = $3)", + nodes_table + )) + .bind(&file_uuid) + .bind(&ext_id_face) + .bind(&ext_id_trace) + .execute(state.db.pool()) + .await; + + let tkg_deleted = match tkg_result { + Ok(r) => r.rows_affected(), + Err(e) => { + tracing::error!("[delete_trace] Failed to delete TKG nodes: {}", e); + 0 + } + }; + + Json(serde_json::json!({ + "success": true, + "file_uuid": file_uuid, + "trace_id": trace_id, + "hard_delete": true, + "deleted_face_points": deleted_count, + "deleted_tkg_nodes": tkg_deleted + })) + } else { + // Soft delete: mark as deleted in Qdrant and TKG + let filter = serde_json::json!({ + "must": [ + {"key": "file_uuid", "match": {"value": file_uuid}}, + {"key": "trace_id", "match": {"value": trace_id}} + ] + }); + + let payload = serde_json::json!({ + "status": "deleted" + }); + + let qdrant_updated = match qdrant.update_payload_by_filter("_faces", filter.clone(), payload).await { + Ok(_) => { + tracing::info!("[delete_trace] Soft deleted (marked) trace {} in {}", trace_id, file_uuid); + true + } + Err(e) => { + tracing::error!("[delete_trace] Failed to mark Qdrant points as deleted: {}", e); + return Json(serde_json::json!({ + "success": false, + "error": format!("Failed to mark Qdrant points: {}", e) + })); + } + }; + + // Update TKG node status + let nodes_table = crate::core::db::schema::table_name("tkg_nodes"); + let ext_id_face = format!("face_track_{}", trace_id); + let ext_id_trace = format!("trace_{}", trace_id); + + let tkg_result = sqlx::query(&format!( + "UPDATE {} SET properties = properties || $1::jsonb \ + WHERE file_uuid = $2 AND node_type = 'face_track' AND (external_id = $3 OR external_id = $4)", + nodes_table + )) + .bind(serde_json::json!({"status": "deleted"})) + .bind(&file_uuid) + .bind(&ext_id_face) + .bind(&ext_id_trace) + .execute(state.db.pool()) + .await; + + let tkg_updated = match tkg_result { + Ok(r) => r.rows_affected(), + Err(e) => { + tracing::error!("[delete_trace] Failed to mark TKG nodes as deleted: {}", e); + 0 + } + }; + + Json(serde_json::json!({ + "success": true, + "file_uuid": file_uuid, + "trace_id": trace_id, + "hard_delete": false, + "qdrant_marked": qdrant_updated, + "tkg_nodes_marked": tkg_updated + })) + } +} + +async fn merge_trace( + State(state): State, + Path((file_uuid, trace_id, target_trace_id)): Path<(String, i32, i32)>, +) -> Json { + if trace_id == target_trace_id { + return Json(serde_json::json!({ + "success": false, + "error": "Cannot merge trace into itself" + })); + } + + let qdrant = crate::core::db::qdrant_db::QdrantDb::new(); + + // Step 1: Update all face points from source trace to target trace + let src_filter = serde_json::json!({ + "must": [ + {"key": "file_uuid", "match": {"value": file_uuid}}, + {"key": "trace_id", "match": {"value": trace_id}} + ] + }); + + // Get count before update + let points_count = qdrant.scroll_all_points("_faces", src_filter.clone(), 1).await + .map(|pts| pts.len()).unwrap_or(0); + + // Update all source points to target trace_id + let new_payload = serde_json::json!({ + "trace_id": target_trace_id, + "merged_from": trace_id + }); + + let moved_count = match qdrant.update_payload_by_filter("_faces", src_filter, new_payload).await { + Ok(_) => points_count as u64, + Err(e) => { + return Json(serde_json::json!({ + "success": false, + "error": format!("Failed to update trace_id in Qdrant: {}", e) + })); + } + }; + + // Step 2: Delete source TKG nodes + let nodes_table = crate::core::db::schema::table_name("tkg_nodes"); + let ext_id_face = format!("face_track_{}", trace_id); + let ext_id_trace = format!("trace_{}", trace_id); + + let tkg_result = sqlx::query(&format!( + "DELETE FROM {} WHERE file_uuid = $1 AND node_type = 'face_track' AND (external_id = $2 OR external_id = $3)", + nodes_table + )) + .bind(&file_uuid) + .bind(&ext_id_face) + .bind(&ext_id_trace) + .execute(state.db.pool()) + .await; + + let tkg_deleted = match tkg_result { + Ok(r) => r.rows_affected(), + Err(e) => { + tracing::error!("[merge_trace] Failed to delete source TKG nodes: {}", e); + 0 + } + }; + + // Step 3: Rebuild TKG for the target trace + // (The TKG rebuild will pick up the merged points automatically) + + Json(serde_json::json!({ + "success": true, + "file_uuid": file_uuid, + "source_trace_id": trace_id, + "target_trace_id": target_trace_id, + "points_moved": moved_count, + "tkg_nodes_deleted": tkg_deleted + })) +} diff --git a/src/core/db/qdrant_db.rs b/src/core/db/qdrant_db.rs index 57a7a4c..1bbd153 100644 --- a/src/core/db/qdrant_db.rs +++ b/src/core/db/qdrant_db.rs @@ -740,6 +740,35 @@ impl QdrantDb { Ok(()) } + pub async fn delete_points_by_filter( + &self, + collection: &str, + filter: serde_json::Value, + ) -> Result { + let url = format!("{}/collections/{}/points/delete", self.base_url, collection); + + let body = serde_json::json!({ + "filter": filter + }); + + let resp = self + .client + .post(&url) + .header("api-key", &self.api_key) + .header("Content-Type", "application/json") + .json(&body) + .send() + .await + .context("Failed to delete points from Qdrant")?; + + if !resp.status().is_success() { + anyhow::bail!("Qdrant delete failed: {}", resp.status()); + } + + // Qdrant returns {"result": {"operation_id": N, "status": "completed"}} + Ok(1) // Qdrant doesn't return count, just return 1 to indicate success + } + pub async fn get_point_count(&self) -> Result { let url = format!( "{}/collections/{}/info", diff --git a/src/core/processor/tkg.rs b/src/core/processor/tkg.rs index a32a72f..a47de7b 100644 --- a/src/core/processor/tkg.rs +++ b/src/core/processor/tkg.rs @@ -42,6 +42,11 @@ async fn scroll_face_points(file_uuid: &str) -> Vec { .await { Ok(pts) => { + tracing::info!( + "[TKG-Qdrant] scroll returned {} raw points for {}", + pts.len(), + file_uuid + ); if attempt > 1 { tracing::info!( "[TKG-Qdrant] scroll succeeded on attempt {} for {}", @@ -49,7 +54,13 @@ async fn scroll_face_points(file_uuid: &str) -> Vec { file_uuid ); } - return parse_face_points(pts); + let parsed = parse_face_points(pts); + tracing::info!( + "[TKG-Qdrant] parsed {} FacePoints for {}", + parsed.len(), + file_uuid + ); + return parsed; } Err(e) => { last_err = Some(e); @@ -77,27 +88,47 @@ async fn scroll_face_points(file_uuid: &str) -> Vec { /// Parse FacePoint structs from Qdrant response fn parse_face_points(points: Vec) -> Vec { - points + let mut valid = 0; + let mut no_trace_id = 0; + let mut no_frame = 0; + let result: Vec = points .iter() .filter_map(|p| { let payload = &p["payload"]; - let trace_id = payload["trace_id"].as_i64().filter(|&t| t > 0)?; - let frame = payload["frame"].as_i64()?; + let trace_id = payload["trace_id"].as_i64().filter(|&t| t > 0); + if trace_id.is_none() { + no_trace_id += 1; + return None; + } + let frame = payload["frame"].as_i64(); + if frame.is_none() { + no_frame += 1; + return None; + } + valid += 1; let bbox = &payload["bbox"]; let x = bbox["x"].as_f64().unwrap_or(0.0); let y = bbox["y"].as_f64().unwrap_or(0.0); let w = bbox["width"].as_f64().unwrap_or(0.0); let h = bbox["height"].as_f64().unwrap_or(0.0); Some(FacePoint { - trace_id, - frame, + trace_id: trace_id.unwrap(), + frame: frame.unwrap(), x, y, w, h, }) }) - .collect() + .collect(); + tracing::info!( + "[TKG-parse] total={}, valid={}, no_trace_id={}, no_frame={}", + points.len(), + valid, + no_trace_id, + no_frame + ); + result } /// Build frame-to-face-points index for O(1) lookup diff --git a/src/worker/job_worker.rs b/src/worker/job_worker.rs index 5ab52ff..6506c2d 100644 --- a/src/worker/job_worker.rs +++ b/src/worker/job_worker.rs @@ -2307,6 +2307,23 @@ impl JobWorker { "✅ TKG build completed (no faces) for {}: {} nodes, {} edges", uuid_clone, total_nodes, total_edges ); + + // Update progress for no-faces case + let mut pp = PipelineProgress::new(&uuid_clone); + pp.update_stage( + "tkg_nodes", + 1.0, + "completed", + Some(format!("{} nodes", total_nodes)), + ); + pp.update_stage( + "tkg_edges", + 1.0, + "completed", + Some(format!("{} edges", total_edges)), + ); + publish_pipeline_progress(redis_clone.as_ref(), &uuid_clone, &pp) + .await; } Err(e) => error!("❌ TKG build failed for {}: {}", uuid_clone, e), }