From 29d2cc7315b8019498cf6f534741ae6e4d9c8d34 Mon Sep 17 00:00:00 2001 From: Accusys Date: Mon, 20 Jul 2026 23:26:42 +0800 Subject: [PATCH] fix: wait for Qdrant face points before TKG build - Add Qdrant count verification before TKG build - Wait up to 30 seconds for Qdrant to have all face embeddings - Prevents race condition where TKG reads incomplete data - Add count_points method to QdrantDb --- src/core/db/qdrant_db.rs | 35 ++++++++++++++++++++++++++++++++++ src/worker/job_worker.rs | 41 ++++++++++++++++++++++++++++++++++++---- 2 files changed, 72 insertions(+), 4 deletions(-) diff --git a/src/core/db/qdrant_db.rs b/src/core/db/qdrant_db.rs index e9213bf..57a7a4c 100644 --- a/src/core/db/qdrant_db.rs +++ b/src/core/db/qdrant_db.rs @@ -884,6 +884,41 @@ impl QdrantDb { Ok(all_points) } + pub async fn count_points( + &self, + collection: &str, + filter: serde_json::Value, + ) -> Result { + let url = format!( + "{}/collections/{}/points/count", + self.base_url, collection + ); + + let body = serde_json::json!({ + "filter": filter, + "exact": true + }); + + let resp = self + .client + .post(&url) + .header("api-key", &self.api_key) + .header("Content-Type", "application/json") + .json(&body) + .send() + .await?; + + if !resp.status().is_success() { + anyhow::bail!("Qdrant count failed: {}", resp.status()); + } + + let result: serde_json::Value = resp.json().await?; + let count = result["result"]["count"] + .as_i64() + .unwrap_or(0); + Ok(count) + } + /// Update payload for points matching a filter pub async fn update_payload_by_filter( &self, diff --git a/src/worker/job_worker.rs b/src/worker/job_worker.rs index b474ce9..5ab52ff 100644 --- a/src/worker/job_worker.rs +++ b/src/worker/job_worker.rs @@ -2140,10 +2140,43 @@ impl JobWorker { traced_path, uuid ); } else { - info!( - "📝 Prerequisites met for TKG Build (face_traced.json exists): {}", - uuid - ); + // Verify Qdrant has face points before TKG build + let expected_faces = if let Ok(content) = std::fs::read_to_string(&face_json_path) { + if let Ok(face_data) = serde_json::from_str::(&content) { + face_data.get("total_faces").and_then(|t| t.as_i64()).unwrap_or(0) + } else { 0 } + } else { 0 }; + + // Wait for Qdrant to have all face points (up to 30 seconds) + let qdrant = crate::core::db::qdrant_db::QdrantDb::new(); + let mut attempts = 0; + let max_attempts = 30; + loop { + match qdrant.count_points("_faces", serde_json::json!({ + "must": [{"key": "file_uuid", "match": {"value": uuid}}] + })).await { + Ok(count) if count >= expected_faces && expected_faces > 0 => { + info!("📝 Prerequisites met for TKG Build (Qdrant has {} faces, expected {}): {}", + count, expected_faces, uuid); + break; + } + Ok(count) => { + attempts += 1; + if attempts >= max_attempts { + warn!("⚠️ TKG build proceeding with {} faces (expected {}) after 30s wait: {}", + count, expected_faces, uuid); + break; + } + debug!("⏳ TKG waiting for Qdrant ({} of {} faces): {}", count, expected_faces, uuid); + } + Err(e) => { + warn!("Qdrant count failed: {}, proceeding anyway", e); + break; + } + } + tokio::time::sleep(std::time::Duration::from_secs(1)).await; + } + let db_clone = self.db.clone(); let redis_clone = self.redis.clone(); let uuid_clone = uuid.to_string();