refactor: cargo fmt across 27 files + behavioral fixes

Behavioral changes:
- postgres_db: reorder default processors (cut first), prevent overwriting completed status
- qdrant_db: fix scroll pagination exit condition (next.is_none())
- job_worker: add idempotency guards for face trace / TKG build; better error logging
- main.rs: add LineWriter for stdout buffering

Remaining diff is cargo fmt reformatting (line wrapping, import ordering).
This commit is contained in:
Accusys
2026-07-11 02:03:28 +08:00
parent 701727fd08
commit b98a362de5
27 changed files with 1216 additions and 546 deletions
+259 -123
View File
@@ -9,7 +9,6 @@ use tracing::{debug, error, info, warn};
use crate::api::identity_agent_api::run_identity_agent;
use crate::core::chunk::rule1_ingest;
use crate::core::config::OUTPUT_DIR;
use crate::core::progress::{publish_pipeline_progress, PipelineProgress};
use crate::core::db::qdrant_db::QdrantDb;
use crate::core::db::{
schema, MonitorJobStatus, PostgresDb, ProcessorJobStatus, RedisClient, VectorPayload,
@@ -17,6 +16,7 @@ use crate::core::db::{
};
use crate::core::embedding::Embedder;
use crate::core::processor::heuristic_scene::generate_scene_meta;
use crate::core::progress::{publish_pipeline_progress, PipelineProgress};
use crate::worker::config::WorkerConfig;
use crate::worker::processor::{ProcessorPool, ProcessorTask};
use crate::worker::resources::SystemResources;
@@ -206,7 +206,10 @@ impl JobWorker {
let should_retry = self
.check_and_complete_job(job.id, &job.uuid, &job.processors, expected_count)
.await
.unwrap_or(false);
.unwrap_or_else(|e| {
error!("check_and_complete_job failed for {}: {}", job.uuid, e);
false
});
if should_retry && self.processor_pool.can_start().await {
if let Err(e) = self.process_job(job.clone()).await {
error!("Failed to reprocess job {}: {}", job.uuid, e);
@@ -281,7 +284,10 @@ impl JobWorker {
// Check if job still exists in database (may have been deleted by unregister)
let current_job = self.db.get_monitor_job_by_uuid(&job.uuid).await?;
if current_job.is_none() {
info!("Job {} no longer exists in database (possibly unregistered), skipping", job.uuid);
info!(
"Job {} no longer exists in database (possibly unregistered), skipping",
job.uuid
);
return Ok(());
}
@@ -313,9 +319,17 @@ impl JobWorker {
.await?;
// Clear any stale PipelineProgress from previous jobs
let progress_key = format!("{}progress:{}:pipeline", crate::core::config::REDIS_KEY_PREFIX.as_str(), job.uuid);
let progress_key = format!(
"{}progress:{}:pipeline",
crate::core::config::REDIS_KEY_PREFIX.as_str(),
job.uuid
);
if let Ok(mut conn) = self.redis.get_conn().await {
let _: Option<String> = redis::cmd("DEL").arg(&progress_key).query_async(&mut conn).await.ok();
let _: Option<String> = redis::cmd("DEL")
.arg(&progress_key)
.query_async(&mut conn)
.await
.ok();
}
self.db
@@ -384,14 +398,18 @@ impl JobWorker {
if let Ok(meta) = std::fs::metadata(&tmp_path) {
// 條件 1: 檔案 > 1KB
let has_content = meta.len() > 1024;
// 條件 2: 檔案超過 120 秒未修改(確定沒人還在寫)
let is_stale = if let Ok(modified) = meta.modified() {
if let Ok(elapsed) = modified.elapsed() {
elapsed.as_secs() > 120
} else { false }
} else { false };
} else {
false
}
} else {
false
};
// 條件 3: 檢查程序是否還在跑
let proc_name = processor_type.as_str();
let process_running = std::process::Command::new("ps")
@@ -400,11 +418,11 @@ impl JobWorker {
.ok()
.and_then(|out| String::from_utf8(out.stdout).ok())
.map(|out| {
out.contains(&format!("{}_processor", proc_name)) ||
out.contains(&format!("{}_processor", proc_name))
|| out.contains(&format!("{}_processor", proc_name))
})
.unwrap_or(false);
if has_content && is_stale && !process_running {
info!(
"Found stale .tmp file ({} bytes, {}s old, process={}), renaming to .json for {}",
@@ -1272,7 +1290,8 @@ impl JobWorker {
// TKG may create 0 nodes/edges for videos with minimal content
let has_asr_or_asrx_for_tkg =
job_processors.is_empty() || job_processors.iter().any(|p| p == "asrx" || p == "asr");
let has_face_for_tkg = job_processors.is_empty() || job_processors.iter().any(|p| p == "face");
let has_face_for_tkg =
job_processors.is_empty() || job_processors.iter().any(|p| p == "face");
let tkg_done: bool = if has_asr_or_asrx_for_tkg && has_face_for_tkg {
// TKG is done if face traces are complete (TKG runs after face tracing)
@@ -1360,9 +1379,11 @@ impl JobWorker {
});
// Check for missing processors (in job_processors but not in results)
let missing_processors: Vec<String> = job_processors.iter().filter(|p| {
!results.iter().any(|r| r.processor_type.as_str() == *p)
}).cloned().collect();
let missing_processors: Vec<String> = job_processors
.iter()
.filter(|p| !results.iter().any(|r| r.processor_type.as_str() == *p))
.cloned()
.collect();
if !missing_processors.is_empty() {
info!(
@@ -1602,8 +1623,18 @@ impl JobWorker {
}
}
let mut pp = PipelineProgress::new(&uuid_clone);
pp.update_stage("rule1_ingestion", 1.0, "completed", Some(format!("{} chunks", count)));
publish_pipeline_progress(redis_clone.as_ref(), &uuid_clone, &pp).await;
pp.update_stage(
"rule1_ingestion",
1.0,
"completed",
Some(format!("{} chunks", count)),
);
publish_pipeline_progress(
redis_clone.as_ref(),
&uuid_clone,
&pp,
)
.await;
info!("📦 Phase 1 release packaging...");
let executor =
match crate::core::processor::PythonExecutor::new() {
@@ -1653,87 +1684,113 @@ impl JobWorker {
// 🚀 P2 Trigger: Face Trace + DB Store (after Face)
// Runs face_tracker.py (IoU+embedding tracking), stores trace_id + position in DB
if has_face {
info!("📝 Face completed, triggering face trace + DB store...");
let db_clone = self.db.clone();
let redis_clone = self.redis.clone();
let uuid_clone = uuid.to_string();
tokio::spawn(async move {
let executor = match crate::core::processor::PythonExecutor::new() {
Ok(ex) => ex,
Err(e) => {
error!("Failed to create PythonExecutor for face trace: {}", e);
return;
}
};
match executor
.run(
"store_traced_faces.py",
&["--file-uuid", &uuid_clone],
Some(&uuid_clone),
"TRACE_STORE",
Some(std::time::Duration::from_secs(600)),
)
.await
{
Ok(()) => {
info!("✅ Face trace + DB store completed for {}", uuid_clone);
let traced_path = format!(
"{}{}.face_traced.json",
crate::core::config::OUTPUT_DIR
.as_str()
.trim_end_matches('/'),
uuid
);
if std::path::Path::new(&traced_path).exists() {
info!("✅ Face trace already done for {}, skipping spawn", uuid);
} else {
info!("📝 Face completed, triggering face trace + DB store...");
let db_clone = self.db.clone();
let redis_clone = self.redis.clone();
let uuid_clone = uuid.to_string();
tokio::spawn(async move {
let executor = match crate::core::processor::PythonExecutor::new() {
Ok(ex) => ex,
Err(e) => {
error!("Failed to create PythonExecutor for face trace: {}", e);
return;
}
};
match executor
.run(
"store_traced_faces.py",
&["--file-uuid", &uuid_clone],
Some(&uuid_clone),
"TRACE_STORE",
Some(std::time::Duration::from_secs(600)),
)
.await
{
Ok(()) => {
info!("✅ Face trace + DB store completed for {}", uuid_clone);
// Query trace count and distribution
let trace_count = match db_clone
.get_trace_count_by_file(&uuid_clone)
.await
{
Ok(c) => c,
Err(e) => {
error!("Failed to get trace count for {}: {}", uuid_clone, e);
0
}
};
// Query trace count and distribution
let trace_count =
match db_clone.get_trace_count_by_file(&uuid_clone).await {
Ok(c) => c,
Err(e) => {
error!(
"Failed to get trace count for {}: {}",
uuid_clone, e
);
0
}
};
let (single_frame, multi_frame) = match db_clone
.get_trace_frame_count_distribution(&uuid_clone)
.await
{
Ok(dist) => dist,
Err(e) => {
error!(
"Failed to get trace distribution for {}: {}",
uuid_clone, e
let (single_frame, multi_frame) = match db_clone
.get_trace_frame_count_distribution(&uuid_clone)
.await
{
Ok(dist) => dist,
Err(e) => {
error!(
"Failed to get trace distribution for {}: {}",
uuid_clone, e
);
(0, 0)
}
};
let trace_status =
crate::core::processor::TraceStatus::from_trace_count(
trace_count,
);
(0, 0)
}
};
let trace_status =
crate::core::processor::TraceStatus::from_trace_count(trace_count);
info!(
info!(
"📊 Trace status: {} (total={}, single_frame={}, multi_frame={}) for {}",
trace_status, trace_count, single_frame, multi_frame, uuid_clone
);
// Update processor_results trace_status for Face
if let Err(e) = db_clone
.update_trace_status_for_face(
&uuid_clone,
&trace_status,
trace_count,
single_frame,
multi_frame,
)
.await
{
error!("Failed to update trace_status for {}: {}", uuid_clone, e);
}
// Update processor_results trace_status for Face
if let Err(e) = db_clone
.update_trace_status_for_face(
&uuid_clone,
&trace_status,
trace_count,
single_frame,
multi_frame,
)
.await
{
error!(
"Failed to update trace_status for {}: {}",
uuid_clone, e
);
}
let mut pp = PipelineProgress::new(&uuid_clone);
pp.update_stage("face_tracing", 1.0, "completed", Some(format!("{} traces ({} single, {} multi)", trace_count, single_frame, multi_frame)));
publish_pipeline_progress(redis_clone.as_ref(), &uuid_clone, &pp).await;
let mut pp = PipelineProgress::new(&uuid_clone);
pp.update_stage(
"face_tracing",
1.0,
"completed",
Some(format!(
"{} traces ({} single, {} multi)",
trace_count, single_frame, multi_frame
)),
);
publish_pipeline_progress(redis_clone.as_ref(), &uuid_clone, &pp)
.await;
}
Err(e) => {
error!("❌ Face trace + DB store failed for {}: {}", uuid_clone, e)
}
}
Err(e) => {
error!("❌ Face trace + DB store failed for {}: {}", uuid_clone, e)
}
}
});
});
}
}
// 🚀 P2.5 Trigger: TMDb Face Matching (after Face, if TMDb data exists)
@@ -1841,7 +1898,8 @@ impl JobWorker {
let has_seeds = {
use crate::core::db::qdrant_db::QdrantDb;
let qdrant = QdrantDb::new();
let schema = std::env::var("DATABASE_SCHEMA").unwrap_or_else(|_| "dev".to_string());
let schema =
std::env::var("DATABASE_SCHEMA").unwrap_or_else(|_| "dev".to_string());
let seeds_collection = if schema == "public" {
"momentry_public_seeds"
} else {
@@ -1850,7 +1908,7 @@ impl JobWorker {
let filter = serde_json::json!({
"must": [{"key": "file_uuid", "match": {"value": uuid}}]
});
match qdrant.scroll_all_points("_seeds", filter, 1).await {
match qdrant.scroll_all_points("_seeds", filter, 100).await {
Ok(points) => !points.is_empty(),
Err(e) => {
warn!("Failed to check _seeds for {}: {}", uuid, e);
@@ -1860,62 +1918,140 @@ impl JobWorker {
};
if has_seeds {
info!("📝 Prerequisites met for Identity Agent (has seeds). Starting analysis...");
info!(
"📝 Prerequisites met for Identity Agent (has seeds). Starting analysis..."
);
let db_clone = self.db.clone();
let redis_clone = self.redis.clone();
let uuid_clone = uuid.to_string();
tokio::spawn(async move {
match run_identity_agent(&db_clone, &uuid_clone, Some(redis_clone.clone())).await {
match run_identity_agent(&db_clone, &uuid_clone, Some(redis_clone.clone()))
.await
{
Ok(()) => {
info!("✅ Identity Agent completed for {}", uuid_clone);
let mut pp = PipelineProgress::new(&uuid_clone);
pp.update_stage("identity_agent", 1.0, "completed", None);
publish_pipeline_progress(redis_clone.as_ref(), &uuid_clone, &pp).await;
publish_pipeline_progress(redis_clone.as_ref(), &uuid_clone, &pp)
.await;
}
Err(e) => error!("❌ Identity Agent failed for {}: {}", uuid_clone, e),
}
});
} else {
info!("📝 Skipping Identity Agent for {} (no seed identities)", uuid);
info!(
"📝 Skipping Identity Agent for {} (no seed identities)",
uuid
);
}
}
// 🚀 P4 Trigger: TKG Build (Face + ASRX) → then Rule2 ingestion
// Note: build_tkg uses ON CONFLICT, so it's safe to call multiple times
if has_face && has_asrx {
info!("📝 Prerequisites met for TKG Build. Starting graph construction...");
let db_clone = self.db.clone();
let redis_clone = self.redis.clone();
let uuid_clone = uuid.to_string();
let output_dir_clone = crate::core::config::OUTPUT_DIR.clone();
tokio::spawn(async move {
match crate::core::processor::tkg::build_tkg(&db_clone, &uuid_clone, &output_dir_clone, Some(redis_clone.clone())).await {
Ok(r) => {
let total_nodes = r.face_track_nodes + r.gaze_track_nodes + r.lip_track_nodes + r.text_region_nodes + r.appearance_trace_nodes + r.accessory_nodes + r.object_nodes + r.hand_nodes + r.speaker_nodes;
let total_edges = r.co_occurrence_edges + r.speaker_face_edges + r.face_face_edges + r.mutual_gaze_edges + r.lip_sync_edges + r.has_appearance_edges + r.wears_edges + r.hand_object_edges;
info!("✅ TKG build completed for {}: {} nodes, {} edges", uuid_clone, total_nodes, total_edges);
let tkg_table = crate::core::db::schema::table_name("tkg_edges");
let tkg_done: bool = sqlx::query_scalar::<_, i32>(&format!(
"SELECT 1 FROM {tkg_table} WHERE file_uuid = $1 LIMIT 1"
))
.bind(uuid)
.fetch_optional(self.db.pool())
.await
.unwrap_or(None)
.unwrap_or(0)
> 0;
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;
if tkg_done {
info!("✅ TKG already built for {}, skipping spawn", uuid);
} else {
info!("📝 Prerequisites met for TKG Build. Starting graph construction...");
let db_clone = self.db.clone();
let redis_clone = self.redis.clone();
let uuid_clone = uuid.to_string();
let output_dir_clone = crate::core::config::OUTPUT_DIR.clone();
tokio::spawn(async move {
match crate::core::processor::tkg::build_tkg(
&db_clone,
&uuid_clone,
&output_dir_clone,
Some(redis_clone.clone()),
)
.await
{
Ok(r) => {
let total_nodes = r.face_track_nodes
+ r.gaze_track_nodes
+ r.lip_track_nodes
+ r.text_region_nodes
+ r.appearance_trace_nodes
+ r.accessory_nodes
+ r.object_nodes
+ r.hand_nodes
+ r.speaker_nodes;
let total_edges = r.co_occurrence_edges
+ r.speaker_face_edges
+ r.face_face_edges
+ r.mutual_gaze_edges
+ r.lip_sync_edges
+ r.has_appearance_edges
+ r.wears_edges
+ r.hand_object_edges;
info!(
"✅ TKG build completed for {}: {} nodes, {} edges",
uuid_clone, total_nodes, total_edges
);
// Trigger Rule 2 ingestion after TKG complete
if total_edges > 0 {
match crate::core::chunk::rule2_ingest::ingest_rule2(db_clone.pool(), &uuid_clone, None, None).await {
Ok(rule2_count) => {
info!("✅ Rule 2 ingestion completed for {}: {} relationship chunks", uuid_clone, rule2_count);
let mut pp = PipelineProgress::new(&uuid_clone);
pp.update_stage("rule2_ingestion", 1.0, "completed", Some(format!("{} chunks", rule2_count)));
publish_pipeline_progress(redis_clone.as_ref(), &uuid_clone, &pp).await;
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;
// Trigger Rule 2 ingestion after TKG complete
if total_edges > 0 {
match crate::core::chunk::rule2_ingest::ingest_rule2(
db_clone.pool(),
&uuid_clone,
None,
None,
)
.await
{
Ok(rule2_count) => {
info!("✅ Rule 2 ingestion completed for {}: {} relationship chunks", uuid_clone, rule2_count);
let mut pp = PipelineProgress::new(&uuid_clone);
pp.update_stage(
"rule2_ingestion",
1.0,
"completed",
Some(format!("{} chunks", rule2_count)),
);
publish_pipeline_progress(
redis_clone.as_ref(),
&uuid_clone,
&pp,
)
.await;
}
Err(e) => error!(
"❌ Rule 2 ingestion failed for {}: {}",
uuid_clone, e
),
}
Err(e) => error!("❌ Rule 2 ingestion failed for {}: {}", uuid_clone, e),
}
}
Err(e) => error!("❌ TKG build failed for {}: {}", uuid_clone, e),
}
Err(e) => error!("❌ TKG build failed for {}: {}", uuid_clone, e),
}
});
});
}
}
if !Self::ingestion_complete(self.db.pool(), uuid, job_processors).await {
+7 -1
View File
@@ -1505,7 +1505,13 @@ impl ProcessorPool {
"end_frame": segment.end_frame,
});
pre_chunks_to_store.push((segment.start_frame as i64, Some(segment.start_time), data, None, None));
pre_chunks_to_store.push((
segment.start_frame as i64,
Some(segment.start_time),
data,
None,
None,
));
speaker_detections.push((
segment.speaker_id.clone().unwrap_or_default(),