fix: ASRX duplication, TKG edges, trace ingest, and add pipeline progress publishing

- ASRX handler no longer stores duplicate 'asr' pre_chunks
- Pre_chunks storage made idempotent (delete-before-insert)
- Rule 1 + trace_ingest changed to query 'asrx' not 'asr'
- Trace chunks removed (dynamic from TKG/Qdrant)
- TKG scroll_face_points fixed: trace_id >= 1 (not == 1)
- TKG AsrxSegmentEntry: start/end -> start_time/end_time (match ASRX JSON)
- Unregister error handling: log instead of silent discard
- Add publish_pipeline_progress calls at each pipeline stage
  (processors, rule1, face_trace, identity_agent, TKG, rule2, completion)
This commit is contained in:
Accusys
2026-07-02 10:43:46 +08:00
parent d791d138f2
commit 3eabd45882
65 changed files with 9477 additions and 3852 deletions
+470 -70
View File
@@ -9,6 +9,7 @@ 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,
@@ -225,7 +226,7 @@ impl JobWorker {
.get_processor_results_by_job(job.id)
.await
.unwrap_or_default();
// 若有任何 processor 是 pending/skipped(未真正啟動),重新處理 job
// 若有任何 processor 是 pending/skipped/deferred(未真正啟動),重新處理 job
let has_unstarted = results.iter().any(|r| {
matches!(
r.status,
@@ -233,7 +234,21 @@ impl JobWorker {
| crate::core::db::ProcessorJobStatus::Skipped
)
});
if has_unstarted {
// Also check if there are processors without result records (deferred)
let expected_count = if job.processors.is_empty() {
crate::core::db::ProcessorType::all().len()
} else {
job.processors.len()
};
let has_deferred = results.len() < expected_count;
if has_unstarted || has_deferred {
// Call check_and_complete_job to retry deferred processors
let _ = self
.check_and_complete_job(job.id, &job.uuid, &job.processors, expected_count)
.await;
if let Err(e) = self.process_job(job.clone()).await {
error!("Failed to reprocess job {}: {}", job.uuid, e);
}
@@ -345,7 +360,16 @@ impl JobWorker {
processor_type.as_str()
));
debug!("Checking output file: {:?}", output_path);
if output_path.exists() {
// Special case: Pose processor should NOT be skipped even if pose.json exists
// because swift_face_pose creates it and pose.rs needs to interpolate
let skip_check = if *processor_type == crate::core::db::ProcessorType::Pose {
false // Always run pose.rs to check for interpolation
} else {
output_path.exists()
};
if skip_check {
info!(
"Processor {} output file exists, marking completed and skipping",
processor_type.as_str()
@@ -803,6 +827,65 @@ impl JobWorker {
}
}
// Special handling for ASRX: if ASR output exists with no_audio_track/silent_audio, skip processing
if *processor_type == crate::core::db::ProcessorType::Asrx {
let asr_output_path = format!(
"{}{}.asr.json",
crate::core::config::OUTPUT_DIR
.as_str()
.trim_end_matches('/'),
job.uuid
);
if let Ok(asr_json) = std::fs::read_to_string(&asr_output_path) {
if let Ok(asr_data) = serde_json::from_str::<serde_json::Value>(&asr_json) {
let asr_status = asr_data.get("status").and_then(|s| s.as_str());
if let Some(status) = asr_status {
if status == "no_audio_track" || status == "silent_audio" {
info!("ASRX: ASR status={}, skipping ASRX processing", status);
// Create completed result with same status
if let Err(e) = self
.db
.upsert_processor_result(
job.id,
*processor_type,
&job.uuid,
"completed",
)
.await
{
error!("Failed to create ASRX result: {}", e);
}
// Update asr_status column
let _ = sqlx::query(&format!(
"UPDATE {} SET asr_status = $1, segment_count = 0 WHERE job_id = $2 AND processor = 'asrx'",
crate::core::db::schema::table_name("processor_results")
))
.bind(status)
.bind(job.id)
.execute(self.db.pool())
.await;
let _ = self
.redis
.update_worker_processor_status(
&job.uuid,
"asrx",
"completed",
None,
0,
0,
0,
0,
0,
)
.await;
started_count += 1;
continue;
}
}
}
}
}
// Check dependencies: all dependent processors must be completed
let deps = processor_type.dependencies();
if !deps.is_empty() {
@@ -877,6 +960,7 @@ impl JobWorker {
{
error!("Failed to emit processor alert: {}", e);
}
started_count += 1;
continue;
}
}
@@ -1005,54 +1089,127 @@ impl JobWorker {
/// 檢查所有入庫步驟是否已完成(與 ingestion-status endpoint 同步邏輯)
async fn ingestion_complete(pool: &PgPool, uuid: &str, job_processors: &[String]) -> bool {
let chunk_t = schema::table_name("chunk");
let fd_t = schema::table_name("face_detections");
let pr_t = schema::table_name("processor_results");
// Only check conditions relevant to the job's processors
let has_asr_or_asrx =
job_processors.is_empty() || job_processors.iter().any(|p| p == "asrx" || p == "asr");
let has_cut = job_processors.is_empty() || job_processors.iter().any(|p| p == "cut");
let has_face = job_processors.is_empty() || job_processors.iter().any(|p| p == "face");
let rule1 = !has_asr_or_asrx
|| sqlx::query_scalar::<_, i32>(&format!(
"SELECT 1 FROM {chunk_t} WHERE file_uuid = $1 AND chunk_type = 'sentence' LIMIT 1"
// Check asr_status for ASR/ASRX - if no_audio_track or silent_audio, ingestion is complete
let asr_done: bool = if has_asr_or_asrx {
let asr_status: Option<String> = sqlx::query_scalar(&format!(
"SELECT asr_status FROM {pr_t} WHERE file_uuid = $1 AND processor IN ('asr', 'asrx') LIMIT 1"
))
.bind(uuid)
.fetch_optional(pool)
.await
.unwrap_or(None)
.unwrap_or(0)
> 0;
.unwrap_or(None);
let vector = !has_asr_or_asrx
|| sqlx::query_scalar::<_, i32>(&format!(
"SELECT 1 FROM {chunk_t} WHERE file_uuid = $1 AND chunk_type = 'sentence' AND embedding IS NOT NULL LIMIT 1"
))
.bind(uuid)
.fetch_optional(pool)
.await
.unwrap_or(None)
.unwrap_or(0)
> 0;
match asr_status.as_deref() {
Some("no_audio_track") | Some("silent_audio") => {
tracing::info!(
"[Ingestion] ASR status {} for {} - no chunks needed",
asr_status.unwrap_or_default(),
uuid
);
true
}
Some("has_transcript") => {
// Has transcript, need chunks
sqlx::query_scalar::<_, i32>(&format!(
"SELECT 1 FROM {chunk_t} WHERE file_uuid = $1 AND chunk_type = 'sentence' LIMIT 1"
))
.bind(uuid)
.fetch_optional(pool)
.await
.unwrap_or(None)
.unwrap_or(0)
> 0
}
_ => false,
}
} else {
true
};
let trace = !has_face
|| sqlx::query_scalar::<_, i64>(&format!(
"SELECT COUNT(DISTINCT trace_id) FROM {fd_t} WHERE file_uuid = $1 AND trace_id IS NOT NULL"
))
.bind(uuid)
.fetch_one(pool)
.await
.unwrap_or(0)
> 0;
// Check face_status for Face - if no_faces, ingestion is complete
let trace_done: bool = if has_face {
// Check face_traced.json file for traces directly
let output_dir = std::env::var("MOMENTRY_OUTPUT_DIR")
.unwrap_or_else(|_| "/Users/accusys/momentry/output".to_string());
let traced_path = format!("{}/{}.face_traced.json", output_dir, uuid);
let all_ok = rule1 && vector && trace;
if !all_ok {
tracing::info!(
"[Ingestion] waiting (uuid={}): rule1={} vector={} trace={}",
uuid,
rule1,
vector,
trace
tracing::info!(
"[Ingestion] Checking face traces for {}: path={}",
uuid,
traced_path
);
if std::path::Path::new(&traced_path).exists() {
if let Ok(content) = std::fs::read_to_string(&traced_path) {
if let Ok(traced_data) = serde_json::from_str::<serde_json::Value>(&content) {
if let Some(traces) = traced_data.get("traces") {
// traces can be an object (dictionary) or array
let trace_count = if traces.is_object() {
traces.as_object().map(|o| o.len()).unwrap_or(0)
} else if traces.is_array() {
traces.as_array().map(|a| a.len()).unwrap_or(0)
} else {
0
};
if trace_count > 0 {
tracing::info!(
"[Ingestion] Face traces found for {}: {} traces (from face_traced.json)",
uuid, trace_count
);
true
} else {
tracing::warn!("[Ingestion] Face traces is empty for {}", uuid);
false
}
} else {
tracing::warn!(
"[Ingestion] No 'traces' key in face_traced.json for {}",
uuid
);
false
}
} else {
tracing::warn!("[Ingestion] Failed to parse face_traced.json for {}", uuid);
false
}
} else {
tracing::warn!("[Ingestion] Failed to read face_traced.json for {}", uuid);
false
}
} else {
tracing::warn!(
"[Ingestion] face_traced.json not found for {}: {}",
uuid,
traced_path
);
false
}
} else {
tracing::info!("[Ingestion] No face processor, trace_done=true");
true
};
let all_ok = asr_done && trace_done;
tracing::info!(
"[Ingestion] all_ok={} (asr_done={}, trace_done={}) for uuid={}",
all_ok,
asr_done,
trace_done,
uuid
);
if !all_ok {
tracing::info!(
"[Ingestion] waiting (uuid={}): asr_done={} trace_done={}",
uuid,
asr_done,
trace_done
);
}
all_ok
@@ -1103,7 +1260,7 @@ vector,
.any(|r| matches!(r.status, crate::core::db::ProcessorJobStatus::Pending));
const MAX_RETRIES: i32 = 3;
if any_failed && !any_pending {
let failed_processors_to_retry: Vec<i32> = results
.iter()
@@ -1116,19 +1273,131 @@ vector,
.collect();
if !failed_processors_to_retry.is_empty() {
info!("🔄 Attempting to retry {} failed processors...", failed_processors_to_retry.len());
info!(
"🔄 Attempting to retry {} failed processors...",
failed_processors_to_retry.len()
);
for result_id in failed_processors_to_retry {
if let Ok(true) = self.db.retry_failed_processor(result_id, MAX_RETRIES).await {
if let Ok(mut conn) = self.redis.get_conn().await {
let redis_key = format!("momentry:progress:{}", uuid);
let _: Result<i32, _> = redis::AsyncCommands::del(&mut conn, &redis_key).await;
let _: Result<i32, _> =
redis::AsyncCommands::del(&mut conn, &redis_key).await;
}
}
}
}
}
// Retry deferred processors whose dependencies are now met
// Build a set of completed processor types
let completed_set: std::collections::HashSet<_> = results
.iter()
.filter(|r| matches!(r.status, ProcessorJobStatus::Completed))
.map(|r| r.processor_type)
.collect();
let mut created_deferred = false;
// Find processors in job_processors that are not in results yet
for processor_name in job_processors {
let processor_type = match crate::core::db::ProcessorType::from_db_str(processor_name) {
Some(pt) => pt,
None => continue,
};
// Skip if already has a result
if results.iter().any(|r| r.processor_type == processor_type) {
continue;
}
// Check if all dependencies are met
let deps = processor_type.dependencies();
let deps_met = deps.iter().all(|dep| completed_set.contains(dep));
if !deps_met {
continue;
}
info!(
"🔄 Deferred processor {} dependencies now met, creating result",
processor_name
);
created_deferred = true;
// Special handling for ASRX: check ASR output file
if processor_type == crate::core::db::ProcessorType::Asrx {
let asr_output_path = format!(
"{}{}.asr.json",
crate::core::config::OUTPUT_DIR
.as_str()
.trim_end_matches('/'),
uuid
);
if let Ok(asr_json) = std::fs::read_to_string(&asr_output_path) {
if let Ok(asr_data) = serde_json::from_str::<serde_json::Value>(&asr_json) {
let asr_status = asr_data.get("status").and_then(|s| s.as_str());
if let Some(status) = asr_status {
if status == "no_audio_track" || status == "silent_audio" {
info!(
"ASRX: ASR status={}, creating completed result directly",
status
);
if let Err(e) = self
.db
.upsert_processor_result(
job_id,
processor_type,
uuid,
"completed",
)
.await
{
error!("Failed to create ASRX result: {}", e);
}
let _ = sqlx::query(&format!(
"UPDATE {} SET asr_status = $1, segment_count = 0 WHERE job_id = $2 AND processor = 'asrx'",
crate::core::db::schema::table_name("processor_results")
))
.bind(status)
.bind(job_id)
.execute(self.db.pool())
.await;
let _ = self
.redis
.update_worker_processor_status(
uuid,
"asrx",
"completed",
None,
0,
0,
0,
0,
0,
)
.await;
continue;
}
}
}
}
}
// For other deferred processors, create pending result so worker can pick it up
if let Err(e) = self
.db
.upsert_processor_result(job_id, processor_type, uuid, "pending")
.await
{
error!(
"Failed to create deferred result for {}: {}",
processor_name, e
);
}
}
let any_skipped = results
.iter()
.filter(|r| job_processors.contains(&r.processor_type.as_str().to_string()))
@@ -1192,7 +1461,9 @@ vector,
} else {
info!("📝 Prerequisites met for Rule 1 Chunking. Starting ingestion...");
let db_clone = self.db.clone();
let redis_clone = self.redis.clone();
let uuid_clone = uuid.to_string();
let job_id_clone = job_id;
tokio::spawn(async move {
match db_clone.get_video_by_uuid(&uuid_clone).await {
Ok(Some(video)) => {
@@ -1217,6 +1488,9 @@ vector,
);
}
}
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;
info!("📦 Phase 1 release packaging...");
let executor =
match crate::core::processor::PythonExecutor::new() {
@@ -1240,7 +1514,10 @@ vector,
.await
{
Ok(()) => {
info!("✅ Phase 1 release packaged for {}", uuid_clone)
info!("✅ Phase 1 release packaged for {}", uuid_clone);
// Note: Job status will be updated after Rule 2 (TKG) completion
// Do not mark as completed here
}
Err(e) => error!("❌ Phase 1 release pack failed: {}", e),
}
@@ -1251,16 +1528,21 @@ vector,
Ok(None) => error!("Video not found for chunking: {}", uuid_clone),
Err(e) => error!("Failed to get video info for chunking: {}", e),
}
});
}
}
});
}
}
if all_completed {
// 🚀 P2 Trigger: Face Trace + DB Store (after Face)
if all_completed {
let mut pp = PipelineProgress::new(uuid);
pp.update_stage("processors", 1.0, "completed", None);
publish_pipeline_progress(self.redis.as_ref(), uuid, &pp).await;
// 🚀 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() {
@@ -1283,17 +1565,56 @@ if all_completed {
Ok(()) => {
info!("✅ Face trace + DB store completed for {}", uuid_clone);
// Generate trace chunks from face_detections + ASR text
info!("📝 Generating trace chunks...");
match crate::core::chunk::trace_ingest::ingest_traces(
&db_clone,
&uuid_clone,
)
.await
// Query trace count and distribution
let trace_count = match db_clone
.get_trace_count_by_file(&uuid_clone)
.await
{
Ok(n) => info!("✅ {} trace chunks created for {}", n, uuid_clone),
Err(e) => error!("❌ Trace chunk ingestion failed: {}", e),
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
);
(0, 0)
}
};
let trace_status =
crate::core::processor::TraceStatus::from_trace_count(trace_count);
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);
}
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)
@@ -1320,15 +1641,46 @@ if all_completed {
count, uuid_clone
);
// Save identity files for affected identities
let ids = sqlx::query_scalar::<_, uuid::Uuid>(
"SELECT DISTINCT i.uuid FROM identities i \
JOIN face_detections fd ON fd.identity_id = i.id \
WHERE fd.file_uuid = $1 AND fd.identity_id IS NOT NULL",
)
.bind(&uuid_clone)
.fetch_all(db_clone.pool())
.await
.unwrap_or_default();
let qdrant = crate::core::db::qdrant_db::QdrantDb::new();
let face_filter = serde_json::json!({
"must": [
{"key": "file_uuid", "match": {"value": &uuid_clone}},
{"key": "identity_id", "is_null": false}
]
});
let face_points = qdrant
.scroll_all_points("_faces", face_filter, 1000)
.await
.unwrap_or_default();
use std::collections::HashSet;
let mut identity_ids: HashSet<i32> = HashSet::new();
for p in &face_points {
if let Some(iid) = p["payload"]["identity_id"].as_i64() {
identity_ids.insert(iid as i32);
}
}
let ids: Vec<uuid::Uuid> = if !identity_ids.is_empty() {
let ids_list: Vec<i32> = identity_ids.into_iter().collect();
let id_params: Vec<String> =
ids_list.iter().map(|_| "$1".to_string()).collect();
// Use batch query: since we can't do IN with variable params via sqlx easily,
// query one by one. But typically there are few (<20) identities.
let mut result = Vec::new();
for iid in &ids_list {
if let Ok(Some(u)) = sqlx::query_scalar::<_, uuid::Uuid>(
"SELECT uuid FROM identities WHERE id = $1",
)
.bind(iid)
.fetch_optional(db_clone.pool())
.await
{
result.push(u);
}
}
result
} else {
Vec::new()
};
for id_uuid in &ids {
let us = id_uuid.to_string().replace('-', "");
if let Err(e) = crate::core::identity::storage::save_identity_file(
@@ -1374,15 +1726,58 @@ if all_completed {
if has_face && has_asr_or_asrx {
info!("📝 Prerequisites met for Identity Agent. 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).await {
Ok(()) => info!("✅ Identity Agent completed for {}", uuid_clone),
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;
}
Err(e) => error!("❌ Identity Agent failed for {}: {}", uuid_clone, e),
}
});
}
// 🚀 P4 Trigger: TKG Build (Face + ASRX) → then Rule2 ingestion
if has_face && has_asr_or_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 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!("❌ TKG build failed for {}: {}", uuid_clone, e),
}
});
}
if !Self::ingestion_complete(self.db.pool(), uuid, job_processors).await {
info!(
"Job {}: all processors done, waiting for ingestion...",
@@ -1413,6 +1808,10 @@ if all_completed {
self.redis.delete_worker_job(uuid).await?;
let mut pp = PipelineProgress::new(uuid);
pp.mark_completed();
publish_pipeline_progress(self.redis.as_ref(), uuid, &pp).await;
info!("Job {} completed successfully (ingestion done)", job_id);
} else if essential_completed && !all_completed && !any_pending && !any_skipped {
// 必要 processor 完成但部分非必要失敗 → 仍算完成(但無 pending 者才觸發)
@@ -1466,7 +1865,8 @@ if all_completed {
.await?;
}
Ok(false)
// Return true if we created deferred processors, so caller will reprocess the job
Ok(created_deferred)
}
pub async fn shutdown(&self) {
+158 -5
View File
@@ -82,6 +82,10 @@ struct ProcessorOutput {
total_frames: i32,
retry_count: i32,
pid: i32,
asr_status: Option<crate::core::processor::AsrStatus>,
segment_count: usize,
face_status: Option<crate::core::processor::FaceStatus>,
total_faces: usize,
}
#[derive(Debug, Clone)]
@@ -316,13 +320,16 @@ impl ProcessorPool {
}
// Subscribe to Redis progress pub/sub and update processor hash in real-time
let sub_db = db.clone();
let sub_redis = redis.clone();
let sub_uuid = job.uuid.clone();
let sub_processor = processor_name.clone();
let progress_handle = tokio::spawn(async move {
let cb_db = sub_db.clone();
let cb_redis = sub_redis.clone();
let cb_uuid = sub_uuid.clone();
let cb_processor = sub_processor.clone();
let last_update = std::cell::Cell::new(0i64);
if let Err(e) = sub_redis
.subscribe_and_callback(&sub_uuid, move |msg| {
tracing::info!(
@@ -338,6 +345,7 @@ impl ProcessorPool {
let r = cb_redis.clone();
let u = cb_uuid.clone();
let p = cb_processor.clone();
let p2 = p.clone();
tokio::spawn(async move {
match r
.update_worker_processor_status(
@@ -354,6 +362,46 @@ impl ProcessorPool {
Err(e) => tracing::error!("[Subscriber] FAILED {}: {}", p, e),
}
});
// Sync progress to PostgreSQL every 5 seconds
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_secs() as i64;
let elapsed = now - last_update.get();
if elapsed >= 5 {
tracing::info!(
"[Subscriber] PG sync {}: cur={} tot={} (elapsed={})",
p2,
cur,
tot,
elapsed
);
last_update.set(now);
let db_client = cb_db.clone();
let u = cb_uuid.clone();
let p = cb_processor.clone();
tokio::spawn(async move {
if let Err(e) = db_client
.update_processor_progress(
&u, &p, cur as u64, tot as u64, "running",
)
.await
{
tracing::error!(
"[Subscriber] PG progress update FAILED {}: {}",
p,
e
);
} else {
tracing::info!(
"[Subscriber] PG progress updated {}: cur={} tot={}",
p,
cur,
tot
);
}
});
}
}
})
.await
@@ -400,6 +448,32 @@ impl ProcessorPool {
error!("Failed to update processor result to completed: {}", e);
}
if let Some(ref asr_status) = output.asr_status {
if let Err(e) = db
.update_asr_status(
processor_result_id,
asr_status,
output.segment_count,
)
.await
{
error!("Failed to update ASR status: {}", e);
}
}
if let Some(ref face_status) = output.face_status {
if let Err(e) = db
.update_face_status(
processor_result_id,
face_status,
output.total_faces,
)
.await
{
error!("Failed to update FACE status: {}", e);
}
}
if let Err(e) = redis
.update_worker_processor_status(
&job.uuid,
@@ -416,6 +490,20 @@ impl ProcessorPool {
{
error!("Failed to update Redis processor status: {}", e);
}
// Also update PostgreSQL processing_status JSON
if let Err(e) = db
.update_processor_progress(
&job.uuid,
&processor_name,
output.frames_processed as u64,
output.total_frames as u64,
"completed",
)
.await
{
error!("Failed to update PostgreSQL processor status: {}", e);
}
} else {
error!(
"Processor {} output failed verification for job {}: {:?}",
@@ -569,6 +657,10 @@ impl ProcessorPool {
total_frames,
retry_count: 0,
pid: 0,
asr_status: None,
segment_count: 0,
face_status: None,
total_faces: 0,
})
}
ProcessorType::Yolo => {
@@ -612,6 +704,10 @@ impl ProcessorPool {
total_frames,
retry_count: 0,
pid: 0,
asr_status: None,
segment_count: 0,
face_status: None,
total_faces: 0,
})
}
ProcessorType::Ocr => {
@@ -655,6 +751,10 @@ impl ProcessorPool {
total_frames,
retry_count: 0,
pid: 0,
asr_status: None,
segment_count: 0,
face_status: None,
total_faces: 0,
})
}
ProcessorType::Face => {
@@ -666,9 +766,16 @@ impl ProcessorPool {
)
.await?;
let chunks_produced = result.frames.len() as i32;
let face_status = result.status.clone();
let total_faces = result.total_faces;
tracing::info!(
"FACE completed, storing {} frames for {}",
"FACE completed, status={}, {} frames, {} total faces for {}",
face_status
.as_ref()
.map(|s| s.to_string())
.unwrap_or_default(),
chunks_produced,
total_faces,
job.uuid
);
if let Err(e) = Self::store_face_chunks(db, &job.uuid, &result).await {
@@ -720,6 +827,10 @@ impl ProcessorPool {
total_frames,
retry_count: 0,
pid: 0,
asr_status: None,
segment_count: 0,
face_status,
total_faces,
})
}
ProcessorType::FaceCluster => {
@@ -741,6 +852,10 @@ impl ProcessorPool {
total_frames: 0,
retry_count: 0,
pid: 0,
asr_status: None,
segment_count: 0,
face_status: None,
total_faces: 0,
})
}
ProcessorType::Pose => {
@@ -784,6 +899,10 @@ impl ProcessorPool {
total_frames,
retry_count: 0,
pid: 0,
asr_status: None,
segment_count: 0,
face_status: None,
total_faces: 0,
})
}
ProcessorType::Hand => {
@@ -824,6 +943,10 @@ impl ProcessorPool {
total_frames,
retry_count: 0,
pid: 0,
asr_status: None,
segment_count: 0,
face_status: None,
total_faces: 0,
})
}
ProcessorType::Appearance => {
@@ -851,14 +974,24 @@ impl ProcessorPool {
total_frames,
retry_count: 0,
pid: 0,
asr_status: None,
segment_count: 0,
face_status: None,
total_faces: 0,
})
}
ProcessorType::Asr => {
let result =
processor::process_asr(video_path, output_path.to_str().unwrap(), uuid).await?;
let chunks_produced = result.segments.len() as i32;
let asr_status = result.status.clone();
let segment_count = result.segment_count;
tracing::info!(
"ASR completed, storing {} segments for {}",
"ASR completed, status={}, {} segments for {}",
asr_status
.as_ref()
.map(|s| s.to_string())
.unwrap_or_default(),
chunks_produced,
job.uuid
);
@@ -892,6 +1025,10 @@ impl ProcessorPool {
total_frames,
retry_count: 0,
pid: 0,
asr_status,
segment_count,
face_status: None,
total_faces: 0,
})
}
ProcessorType::Asrx => {
@@ -899,8 +1036,14 @@ impl ProcessorPool {
processor::process_asrx(video_path, output_path.to_str().unwrap(), uuid)
.await?;
let chunks_produced = result.segments.len() as i32;
let asr_status = result.status.clone();
let segment_count = result.segment_count;
tracing::info!(
"ASRX completed, storing {} segments for {}",
"ASRX completed, status={}, {} segments for {}",
asr_status
.as_ref()
.map(|s| s.to_string())
.unwrap_or_default(),
chunks_produced,
job.uuid
);
@@ -959,6 +1102,10 @@ impl ProcessorPool {
total_frames,
retry_count: 0,
pid: 0,
asr_status,
segment_count,
face_status: None,
total_faces: 0,
})
}
ProcessorType::Scene => {
@@ -977,6 +1124,10 @@ impl ProcessorPool {
total_frames,
retry_count: 0,
pid: 0,
asr_status: None,
segment_count: 0,
face_status: None,
total_faces: 0,
});
} else if scene_path.exists() {
tracing::info!("Scene JSON exists for {}, loading from file", job.uuid);
@@ -1025,6 +1176,10 @@ impl ProcessorPool {
total_frames,
retry_count: 0,
pid: 0,
asr_status: None,
segment_count: 0,
face_status: None,
total_faces: 0,
})
}
}
@@ -1363,8 +1518,6 @@ impl ProcessorPool {
db.store_raw_pre_chunks_batch(uuid, "asrx", &pre_chunks_to_store)
.await?;
db.store_raw_pre_chunks_batch(uuid, "asr", &pre_chunks_to_store)
.await?;
db.store_speaker_detections_batch(uuid, &speaker_detections)
.await?;
Ok(())