refactor: remove face embedding architecture - single Qdrant _faces collection
- Delete FaceEmbeddingDb module (face_embedding_db.rs) - Stub match_faces_iterative, generate_seed_embeddings, tmdb_match_handler - Remove sync_trace_embeddings, populate_face_embeddings_to_qdrant - Remove embedding from face.json output (face_processor.py) - Remove embedding from PG UPDATE (store_traced_faces.py) - Remove workspace traces staging (checkin.rs, qdrant_workspace.rs) - Fix tests: add pose_angle to Face, hand_nodes to TkgResult Disabled functions (need reimplement with _faces): - match_faces_iterative (identity agent) - generate_seed_embeddings (TMDb seeds) - tmdb_match_handler (TMDb matching) - cluster_face_embeddings, search_similar_faces - merge_traces_within_cuts
This commit is contained in:
+16
-115
@@ -14,9 +14,7 @@ struct ProcessorCleanupGuard {
|
||||
running_count: Arc<RwLock<usize>>,
|
||||
frame_count: Arc<RwLock<usize>>,
|
||||
time_count: Arc<RwLock<usize>>,
|
||||
best_effort_count: Arc<RwLock<usize>>,
|
||||
pipeline: PipelineType,
|
||||
is_best_effort: bool,
|
||||
}
|
||||
|
||||
impl Drop for ProcessorCleanupGuard {
|
||||
@@ -32,30 +30,22 @@ impl Drop for ProcessorCleanupGuard {
|
||||
*guard -= 1;
|
||||
}
|
||||
}
|
||||
if self.is_best_effort {
|
||||
if let Ok(mut guard) = self.best_effort_count.try_write() {
|
||||
if *guard > 0 {
|
||||
*guard -= 1;
|
||||
}
|
||||
}
|
||||
} else {
|
||||
match self.pipeline {
|
||||
PipelineType::Frame => {
|
||||
if let Ok(mut guard) = self.frame_count.try_write() {
|
||||
if *guard > 0 {
|
||||
*guard -= 1;
|
||||
}
|
||||
match self.pipeline {
|
||||
PipelineType::Frame => {
|
||||
if let Ok(mut guard) = self.frame_count.try_write() {
|
||||
if *guard > 0 {
|
||||
*guard -= 1;
|
||||
}
|
||||
}
|
||||
PipelineType::Time => {
|
||||
if let Ok(mut guard) = self.time_count.try_write() {
|
||||
if *guard > 0 {
|
||||
*guard -= 1;
|
||||
}
|
||||
}
|
||||
PipelineType::Time => {
|
||||
if let Ok(mut guard) = self.time_count.try_write() {
|
||||
if *guard > 0 {
|
||||
*guard -= 1;
|
||||
}
|
||||
}
|
||||
PipelineType::Cross => {} // cross pipeline not tracked in slot counts
|
||||
}
|
||||
PipelineType::Cross => {}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -106,8 +96,6 @@ pub struct ProcessorTask {
|
||||
const FRAME_SLOT_MAX: usize = 2;
|
||||
/// Time pipeline max concurrent processors (audio is heavy, run 1 at a time).
|
||||
const TIME_SLOT_MAX: usize = 1;
|
||||
/// Best-effort slot (used by low-priority processors like MediaPipe).
|
||||
const BEST_EFFORT_SLOT_MAX: usize = 1;
|
||||
|
||||
pub struct ProcessorPool {
|
||||
db: Arc<PostgresDb>,
|
||||
@@ -117,7 +105,6 @@ pub struct ProcessorPool {
|
||||
running_count: Arc<RwLock<usize>>,
|
||||
running_frame_count: Arc<RwLock<usize>>,
|
||||
running_time_count: Arc<RwLock<usize>>,
|
||||
running_best_effort_count: Arc<RwLock<usize>>,
|
||||
}
|
||||
|
||||
impl ProcessorPool {
|
||||
@@ -130,7 +117,6 @@ impl ProcessorPool {
|
||||
running_count: Arc::new(RwLock::new(0)),
|
||||
running_frame_count: Arc::new(RwLock::new(0)),
|
||||
running_time_count: Arc::new(RwLock::new(0)),
|
||||
running_best_effort_count: Arc::new(RwLock::new(0)),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -240,22 +226,16 @@ impl ProcessorPool {
|
||||
*count += 1;
|
||||
}
|
||||
// 遞增產線專屬 slot
|
||||
let is_best_effort = processor_type == ProcessorType::MediaPipe;
|
||||
if is_best_effort {
|
||||
*self.running_best_effort_count.write().await += 1;
|
||||
} else {
|
||||
match pipeline {
|
||||
PipelineType::Frame => *self.running_frame_count.write().await += 1,
|
||||
PipelineType::Time => *self.running_time_count.write().await += 1,
|
||||
PipelineType::Cross => {} // cross pipeline uses global slot
|
||||
}
|
||||
match pipeline {
|
||||
PipelineType::Frame => *self.running_frame_count.write().await += 1,
|
||||
PipelineType::Time => *self.running_time_count.write().await += 1,
|
||||
PipelineType::Cross => {} // cross pipeline uses global slot
|
||||
}
|
||||
|
||||
let running = self.running.clone();
|
||||
let running_count = self.running_count.clone();
|
||||
let running_frame_count = self.running_frame_count.clone();
|
||||
let running_time_count = self.running_time_count.clone();
|
||||
let running_best_effort_count = self.running_best_effort_count.clone();
|
||||
let child_pid: Arc<RwLock<Option<i32>>> = Arc::new(RwLock::new(None));
|
||||
running.write().await.insert(
|
||||
job_id,
|
||||
@@ -287,9 +267,7 @@ impl ProcessorPool {
|
||||
running_count: running_count.clone(),
|
||||
frame_count: running_frame_count.clone(),
|
||||
time_count: running_time_count.clone(),
|
||||
best_effort_count: running_best_effort_count.clone(),
|
||||
pipeline,
|
||||
is_best_effort,
|
||||
};
|
||||
|
||||
info!("Starting processor {} for job {}", processor_name, job.uuid);
|
||||
@@ -528,10 +506,7 @@ impl ProcessorPool {
|
||||
|
||||
// Generate output path
|
||||
let output_dir = PathBuf::from(OUTPUT_DIR.as_str());
|
||||
let suffix = match processor_type {
|
||||
ProcessorType::Story => format!("{}.story_story", job.uuid),
|
||||
_ => format!("{}.{}", job.uuid, processor_type.as_str()),
|
||||
};
|
||||
let suffix = format!("{}.{}", job.uuid, processor_type.as_str());
|
||||
let output_path = output_dir.join(format!("{}.json", suffix));
|
||||
|
||||
// Ensure output directory exists
|
||||
@@ -1052,80 +1027,6 @@ impl ProcessorPool {
|
||||
pid: 0,
|
||||
})
|
||||
}
|
||||
ProcessorType::Story => {
|
||||
let executor = crate::core::processor::PythonExecutor::new()?;
|
||||
let _ = executor
|
||||
.run(
|
||||
"parent_chunk_5w1h.py",
|
||||
&["--file-uuid", &job.uuid, "--embed"],
|
||||
uuid,
|
||||
"STORY",
|
||||
Some(std::time::Duration::from_secs(300)),
|
||||
)
|
||||
.await;
|
||||
let narratives_path = output_dir.join(format!("{}.narratives.json", job.uuid));
|
||||
let chunks_produced = if narratives_path.exists() {
|
||||
let content = std::fs::read_to_string(&narratives_path).unwrap_or_default();
|
||||
let count: i32 = serde_json::from_str::<Vec<String>>(&content)
|
||||
.map(|v| v.len() as i32)
|
||||
.unwrap_or(0);
|
||||
tracing::info!("Story generated {} narratives for {}", count, job.uuid);
|
||||
count
|
||||
} else {
|
||||
0
|
||||
};
|
||||
Ok(ProcessorOutput {
|
||||
data: serde_json::Value::Null,
|
||||
chunks_produced,
|
||||
frames_processed: total_frames,
|
||||
total_frames,
|
||||
retry_count: 0,
|
||||
pid: 0,
|
||||
})
|
||||
}
|
||||
ProcessorType::FiveW1H => {
|
||||
let executor = crate::core::processor::PythonExecutor::new()?;
|
||||
let _ = executor
|
||||
.run(
|
||||
"parent_chunk_5w1h.py",
|
||||
&["--file-uuid", &job.uuid, "--embed", "--mode", "llm"],
|
||||
uuid,
|
||||
"5W1H",
|
||||
Some(std::time::Duration::from_secs(300)),
|
||||
)
|
||||
.await;
|
||||
Ok(ProcessorOutput {
|
||||
data: serde_json::Value::Null,
|
||||
chunks_produced: 0,
|
||||
frames_processed: total_frames,
|
||||
total_frames,
|
||||
retry_count: 0,
|
||||
pid: 0,
|
||||
})
|
||||
}
|
||||
ProcessorType::MediaPipe => {
|
||||
let result = processor::process_mediapipe_v2(
|
||||
video_path,
|
||||
output_path.to_str().unwrap(),
|
||||
uuid,
|
||||
Some(&sample_frames),
|
||||
)
|
||||
.await?;
|
||||
let chunks_produced = result.frames.len() as i32;
|
||||
tracing::info!(
|
||||
"MEDIAPIPE completed, {} frames for {}",
|
||||
chunks_produced,
|
||||
job.uuid
|
||||
);
|
||||
Ok(ProcessorOutput {
|
||||
data: serde_json::to_value(result)?,
|
||||
chunks_produced,
|
||||
frames_processed: total_frames,
|
||||
total_frames,
|
||||
retry_count: 0,
|
||||
pid: 0,
|
||||
})
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user