feat: independent TKG processing mechanism with operation logging
- Created src/core/tkg/ module with service, log, and models - Added tkg_operation_log table for tracking all TKG operations - TkgService: build, rebuild (with force option), delete, get_operations - Updated rebuild endpoint with force parameter - Added GET /api/v1/file/:file_uuid/tkg for operation history - Added DELETE /api/v1/file/:file_uuid/tkg for TKG deletion - TKG rebuild now always triggers Rule 2 (even with 0 edges) - Full audit trail for all TKG operations (create/update/delete/rebuild)
This commit is contained in:
+45
-38
@@ -3,7 +3,7 @@ use axum::{
|
||||
extract::{Path, Query, State},
|
||||
http::{header, StatusCode},
|
||||
response::{IntoResponse, Json, Response},
|
||||
routing::{get, post},
|
||||
routing::{delete, get, post},
|
||||
Router,
|
||||
};
|
||||
use serde::{Deserialize, Serialize};
|
||||
@@ -39,6 +39,8 @@ pub fn trace_agent_routes() -> Router<crate::api::types::AppState> {
|
||||
get(get_cooccurrence),
|
||||
)
|
||||
.route("/api/v1/file/:file_uuid/tkg/rebuild", post(rebuild_tkg))
|
||||
.route("/api/v1/file/:file_uuid/tkg", get(get_tkg_operations))
|
||||
.route("/api/v1/file/:file_uuid/tkg", delete(delete_tkg))
|
||||
.route("/api/v1/file/:file_uuid/rule2", post(ingest_rule2))
|
||||
.route(
|
||||
"/api/v1/file/:file_uuid/representative-frame",
|
||||
@@ -980,58 +982,41 @@ struct TkgRebuildResponse {
|
||||
error: Option<String>,
|
||||
}
|
||||
|
||||
#[derive(Deserialize)]
|
||||
struct RebuildParams {
|
||||
force: Option<bool>,
|
||||
}
|
||||
|
||||
async fn rebuild_tkg(
|
||||
State(state): State<crate::api::types::AppState>,
|
||||
Path(file_uuid): Path<String>,
|
||||
Query(params): Query<RebuildParams>,
|
||||
) -> Json<TkgRebuildResponse> {
|
||||
use crate::core::chunk::rule2_ingest::ingest_rule2;
|
||||
use crate::core::tkg::TkgService;
|
||||
use tracing::info;
|
||||
|
||||
let redis = crate::core::db::RedisClient::new().ok();
|
||||
let result = crate::core::processor::tkg::build_tkg(&state.db, &file_uuid, &OUTPUT_DIR, redis.map(Arc::new)).await;
|
||||
let db = state.db.clone();
|
||||
let tkg_service = TkgService::new(state.db);
|
||||
let force = params.force.unwrap_or(false);
|
||||
let result = tkg_service.rebuild(&file_uuid, &OUTPUT_DIR, force).await;
|
||||
|
||||
match result {
|
||||
Ok(r) => {
|
||||
let total_edges = r.speaker_face_edges
|
||||
+ r.mutual_gaze_edges
|
||||
+ r.face_face_edges
|
||||
+ r.co_occurrence_edges
|
||||
+ r.has_appearance_edges
|
||||
+ r.wears_edges;
|
||||
|
||||
if total_edges > 0 {
|
||||
info!(
|
||||
"[TKG] {} relationship edges found, triggering Rule 2 ingestion...",
|
||||
total_edges
|
||||
);
|
||||
match ingest_rule2(state.db.pool(), &file_uuid, None, None).await {
|
||||
Ok(count) => info!("[TKG] Rule 2 created {} relationship chunks", count),
|
||||
Err(e) => info!("[TKG] Rule 2 ingestion failed: {}", e),
|
||||
}
|
||||
// Always trigger Rule 2 (even with 0 edges)
|
||||
info!(
|
||||
"[TKG] Rebuild completed for {}: {} nodes, {} edges",
|
||||
file_uuid, r.total_nodes(), r.total_edges()
|
||||
);
|
||||
match ingest_rule2(db.pool(), &file_uuid, None, None).await {
|
||||
Ok(count) => info!("[TKG] Rule 2 created {} relationship chunks", count),
|
||||
Err(e) => info!("[TKG] Rule 2 ingestion failed: {}", e),
|
||||
}
|
||||
|
||||
Json(TkgRebuildResponse {
|
||||
success: true,
|
||||
file_uuid,
|
||||
result: Some(serde_json::json!({
|
||||
"face_track_nodes": r.face_track_nodes,
|
||||
"gaze_track_nodes": r.gaze_track_nodes,
|
||||
"lip_track_nodes": r.lip_track_nodes,
|
||||
"text_region_nodes": r.text_region_nodes,
|
||||
"appearance_trace_nodes": r.appearance_trace_nodes,
|
||||
"accessory_nodes": r.accessory_nodes,
|
||||
"object_nodes": r.object_nodes,
|
||||
"hand_nodes": r.hand_nodes,
|
||||
"speaker_nodes": r.speaker_nodes,
|
||||
"co_occurrence_edges": r.co_occurrence_edges,
|
||||
"speaker_face_edges": r.speaker_face_edges,
|
||||
"face_face_edges": r.face_face_edges,
|
||||
"mutual_gaze_edges": r.mutual_gaze_edges,
|
||||
"lip_sync_edges": r.lip_sync_edges,
|
||||
"has_appearance_edges": r.has_appearance_edges,
|
||||
"wears_edges": r.wears_edges,
|
||||
"hand_object_edges": r.hand_object_edges,
|
||||
})),
|
||||
result: Some(r.to_json()),
|
||||
error: None,
|
||||
})
|
||||
}
|
||||
@@ -1044,6 +1029,28 @@ async fn rebuild_tkg(
|
||||
}
|
||||
}
|
||||
|
||||
async fn get_tkg_operations(
|
||||
State(state): State<crate::api::types::AppState>,
|
||||
Path(file_uuid): Path<String>,
|
||||
) -> Json<Vec<crate::core::tkg::models::TkgOperationLog>> {
|
||||
use crate::core::tkg::TkgService;
|
||||
let tkg_service = TkgService::new(state.db);
|
||||
let operations = tkg_service.get_operations(&file_uuid).await.unwrap_or_default();
|
||||
Json(operations)
|
||||
}
|
||||
|
||||
async fn delete_tkg(
|
||||
State(state): State<crate::api::types::AppState>,
|
||||
Path(file_uuid): Path<String>,
|
||||
) -> Json<serde_json::Value> {
|
||||
use crate::core::tkg::TkgService;
|
||||
let tkg_service = TkgService::new(state.db);
|
||||
match tkg_service.delete_tkg_and_logs(&file_uuid).await {
|
||||
Ok(_) => Json(serde_json::json!({"success": true, "message": "TKG deleted successfully"})),
|
||||
Err(e) => Json(serde_json::json!({"success": false, "error": e.to_string()})),
|
||||
}
|
||||
}
|
||||
|
||||
// ── Representative Frame (JSON) ───────────────────────────────────
|
||||
|
||||
use crate::core::processor::tkg;
|
||||
|
||||
Reference in New Issue
Block a user