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
This commit is contained in:
@@ -884,6 +884,41 @@ impl QdrantDb {
|
|||||||
Ok(all_points)
|
Ok(all_points)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
pub async fn count_points(
|
||||||
|
&self,
|
||||||
|
collection: &str,
|
||||||
|
filter: serde_json::Value,
|
||||||
|
) -> Result<i64> {
|
||||||
|
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
|
/// Update payload for points matching a filter
|
||||||
pub async fn update_payload_by_filter(
|
pub async fn update_payload_by_filter(
|
||||||
&self,
|
&self,
|
||||||
|
|||||||
@@ -2140,10 +2140,43 @@ impl JobWorker {
|
|||||||
traced_path, uuid
|
traced_path, uuid
|
||||||
);
|
);
|
||||||
} else {
|
} else {
|
||||||
info!(
|
// Verify Qdrant has face points before TKG build
|
||||||
"📝 Prerequisites met for TKG Build (face_traced.json exists): {}",
|
let expected_faces = if let Ok(content) = std::fs::read_to_string(&face_json_path) {
|
||||||
uuid
|
if let Ok(face_data) = serde_json::from_str::<serde_json::Value>(&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 db_clone = self.db.clone();
|
||||||
let redis_clone = self.redis.clone();
|
let redis_clone = self.redis.clone();
|
||||||
let uuid_clone = uuid.to_string();
|
let uuid_clone = uuid.to_string();
|
||||||
|
|||||||
Reference in New Issue
Block a user