feat: add Vision LLM integration (CLIP + Qwen3-VL cascade)
- Add Qwen3-VL dynamic management (start/stop/status CLI) - Add CLIP + Qwen3-VL cascade detection strategy - Add Vision CLI commands (vision start/stop/status, detect) - Add cascade_vision processor module - Add clip processor module - Add qwen_vl_manager module Changes: - scripts/start_qwen3vl.sh, stop_qwen3vl.sh: Qwen3-VL management scripts - src/core/vision/: Qwen3-VL manager module - src/core/processor/cascade_vision.rs: CLIP + Qwen3-VL cascade logic - src/core/processor/clip.rs: CLIP classification and detection - src/api/clip_api.rs: CLIP API endpoints - src/cli/vision.rs: Vision CLI implementation - src/cli/args.rs: Add Vision and Detect commands - src/main.rs: Integrate Vision CLI - src/core/mod.rs: Add vision module - src/core/processor/mod.rs: Add cascade_vision module
This commit is contained in:
+30
-21
@@ -17,8 +17,8 @@ pub async fn store_asrx_chunks(db: &PostgresDb, uuid: &str) -> Result<()> {
|
||||
|
||||
let json_str = std::fs::read_to_string(&asrx_path)
|
||||
.with_context(|| format!("ASRX file not found: {:?}", asrx_path))?;
|
||||
let result: AsrxResult = serde_json::from_str(&json_str)
|
||||
.context("Failed to parse ASRX JSON")?;
|
||||
let result: AsrxResult =
|
||||
serde_json::from_str(&json_str).context("Failed to parse ASRX JSON")?;
|
||||
|
||||
let segments_count = result.segments.len();
|
||||
let mut pre_chunks = Vec::new();
|
||||
@@ -41,21 +41,26 @@ pub async fn store_asrx_chunks(db: &PostgresDb, uuid: &str) -> Result<()> {
|
||||
));
|
||||
}
|
||||
|
||||
db.store_raw_pre_chunks_batch(uuid, "asrx", &pre_chunks).await?;
|
||||
db.store_raw_pre_chunks_batch(uuid, "asr", &pre_chunks).await?;
|
||||
db.store_speaker_detections_batch(uuid, &speaker_detections).await?;
|
||||
db.store_raw_pre_chunks_batch(uuid, "asrx", &pre_chunks)
|
||||
.await?;
|
||||
db.store_raw_pre_chunks_batch(uuid, "asr", &pre_chunks)
|
||||
.await?;
|
||||
db.store_speaker_detections_batch(uuid, &speaker_detections)
|
||||
.await?;
|
||||
|
||||
println!("Stored {} ASRX pre-chunks for {}", segments_count, uuid);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub async fn execute_rule1(db: &PostgresDb, uuid: &str) -> Result<usize> {
|
||||
let video = db.get_video_by_uuid(uuid)
|
||||
let video = db
|
||||
.get_video_by_uuid(uuid)
|
||||
.await?
|
||||
.context("Video not found")?;
|
||||
let fps = video.fps;
|
||||
|
||||
let count = rule1_ingest::execute_rule1(db, uuid, fps).await
|
||||
let count = rule1_ingest::execute_rule1(db, uuid, fps)
|
||||
.await
|
||||
.context("Rule 1 ingestion failed")?;
|
||||
|
||||
println!("Rule 1 completed: {} chunks inserted for {}", count, uuid);
|
||||
@@ -68,17 +73,15 @@ pub async fn vectorize_chunks(uuid: &str) -> Result<()> {
|
||||
let embedder = Embedder::new("embeddinggemma-300m".to_string());
|
||||
|
||||
let chunk_table = schema::table_name("chunk");
|
||||
let rows = sqlx::query_as::<_, (String, String, String, i64, i64, f64, f64, String)>(
|
||||
&format!(
|
||||
"SELECT chunk_id, chunk_type, text_content, start_frame, end_frame, \
|
||||
let rows = sqlx::query_as::<_, (String, String, String, i64, i64, f64, f64, String)>(&format!(
|
||||
"SELECT chunk_id, chunk_type, text_content, start_frame, end_frame, \
|
||||
start_time, end_time, content::text \
|
||||
FROM {} WHERE file_uuid = $1 AND chunk_type = 'sentence' \
|
||||
AND embedding IS NULL \
|
||||
AND (text_content IS NOT NULL AND text_content != '') \
|
||||
ORDER BY id",
|
||||
chunk_table
|
||||
),
|
||||
)
|
||||
chunk_table
|
||||
))
|
||||
.bind(uuid)
|
||||
.fetch_all(db.pool())
|
||||
.await?;
|
||||
@@ -91,7 +94,9 @@ pub async fn vectorize_chunks(uuid: &str) -> Result<()> {
|
||||
let total = rows.len();
|
||||
let mut stored = 0usize;
|
||||
|
||||
for (chunk_id, _chunk_type, text, start_frame, end_frame, start_time, end_time, _content_str) in &rows {
|
||||
for (chunk_id, _chunk_type, text, start_frame, end_frame, start_time, end_time, _content_str) in
|
||||
&rows
|
||||
{
|
||||
if text.is_empty() {
|
||||
continue;
|
||||
}
|
||||
@@ -127,13 +132,15 @@ pub async fn vectorize_chunks(uuid: &str) -> Result<()> {
|
||||
}
|
||||
}
|
||||
|
||||
println!("Vectorization complete: {}/{} vectors for {}", stored, total, uuid);
|
||||
println!(
|
||||
"Vectorization complete: {}/{} vectors for {}",
|
||||
stored, total, uuid
|
||||
);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub async fn run_phase1(uuid: &str) -> Result<()> {
|
||||
let executor = PythonExecutor::new()
|
||||
.context("Failed to create PythonExecutor")?;
|
||||
let executor = PythonExecutor::new().context("Failed to create PythonExecutor")?;
|
||||
|
||||
executor
|
||||
.run(
|
||||
@@ -154,15 +161,17 @@ pub async fn mark_complete(db: &PostgresDb, uuid: &str) -> Result<()> {
|
||||
use crate::core::db::MonitorJobStatus;
|
||||
use crate::core::db::VideoStatus;
|
||||
|
||||
let job_id = sqlx::query_scalar::<_, i32>(
|
||||
&format!("SELECT id FROM {} WHERE uuid = $1 LIMIT 1", schema::table_name("monitor_jobs")),
|
||||
)
|
||||
let job_id = sqlx::query_scalar::<_, i32>(&format!(
|
||||
"SELECT id FROM {} WHERE uuid = $1 LIMIT 1",
|
||||
schema::table_name("monitor_jobs")
|
||||
))
|
||||
.bind(uuid)
|
||||
.fetch_optional(db.pool())
|
||||
.await?;
|
||||
|
||||
if let Some(job_id) = job_id {
|
||||
db.update_job_status(job_id, MonitorJobStatus::Completed).await?;
|
||||
db.update_job_status(job_id, MonitorJobStatus::Completed)
|
||||
.await?;
|
||||
println!("Job {} marked as completed", job_id);
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user