178 lines
6.6 KiB
Rust
178 lines
6.6 KiB
Rust
use crate::*;
|
|
use module_ai::{AI_RESULT_REF_ID_PREFIX, AI_TEXT_CHUNK_ID_PREFIX};
|
|
|
|
pub(crate) fn build_ai_task_row(snapshot: &AiTaskSnapshot) -> AiTask {
|
|
AiTask {
|
|
task_id: snapshot.task_id.clone(),
|
|
task_kind: snapshot.task_kind,
|
|
owner_user_id: snapshot.owner_user_id.clone(),
|
|
request_label: snapshot.request_label.clone(),
|
|
source_module: snapshot.source_module.clone(),
|
|
source_entity_id: snapshot.source_entity_id.clone(),
|
|
request_payload_json: snapshot.request_payload_json.clone(),
|
|
status: snapshot.status,
|
|
failure_message: snapshot.failure_message.clone(),
|
|
latest_text_output: snapshot.latest_text_output.clone(),
|
|
latest_structured_payload_json: snapshot.latest_structured_payload_json.clone(),
|
|
version: snapshot.version,
|
|
created_at: Timestamp::from_micros_since_unix_epoch(snapshot.created_at_micros),
|
|
started_at: snapshot
|
|
.started_at_micros
|
|
.map(Timestamp::from_micros_since_unix_epoch),
|
|
completed_at: snapshot
|
|
.completed_at_micros
|
|
.map(Timestamp::from_micros_since_unix_epoch),
|
|
updated_at: Timestamp::from_micros_since_unix_epoch(snapshot.updated_at_micros),
|
|
}
|
|
}
|
|
|
|
pub(crate) fn build_ai_task_snapshot_from_row(
|
|
ctx: &ReducerContext,
|
|
row: &AiTask,
|
|
) -> AiTaskSnapshot {
|
|
let mut stages = ctx
|
|
.db
|
|
.ai_task_stage()
|
|
.iter()
|
|
.filter(|stage| stage.task_id == row.task_id)
|
|
.map(|stage| build_ai_task_stage_snapshot_from_row(&stage))
|
|
.collect::<Vec<_>>();
|
|
stages.sort_by_key(|stage| stage.order);
|
|
|
|
let mut result_references = ctx
|
|
.db
|
|
.ai_result_reference()
|
|
.iter()
|
|
.filter(|reference| reference.task_id == row.task_id)
|
|
.map(|reference| build_ai_result_reference_snapshot_from_row(&reference))
|
|
.collect::<Vec<_>>();
|
|
result_references.sort_by_key(|reference| reference.created_at_micros);
|
|
|
|
AiTaskSnapshot {
|
|
task_id: row.task_id.clone(),
|
|
task_kind: row.task_kind,
|
|
owner_user_id: row.owner_user_id.clone(),
|
|
request_label: row.request_label.clone(),
|
|
source_module: row.source_module.clone(),
|
|
source_entity_id: row.source_entity_id.clone(),
|
|
request_payload_json: row.request_payload_json.clone(),
|
|
status: row.status,
|
|
failure_message: row.failure_message.clone(),
|
|
stages,
|
|
result_references,
|
|
latest_text_output: row.latest_text_output.clone(),
|
|
latest_structured_payload_json: row.latest_structured_payload_json.clone(),
|
|
version: row.version,
|
|
created_at_micros: row.created_at.to_micros_since_unix_epoch(),
|
|
started_at_micros: row
|
|
.started_at
|
|
.map(|value| value.to_micros_since_unix_epoch()),
|
|
completed_at_micros: row
|
|
.completed_at
|
|
.map(|value| value.to_micros_since_unix_epoch()),
|
|
updated_at_micros: row.updated_at.to_micros_since_unix_epoch(),
|
|
}
|
|
}
|
|
|
|
pub(crate) fn build_ai_task_stage_row(
|
|
task_id: &str,
|
|
snapshot: &AiTaskStageSnapshot,
|
|
) -> AiTaskStage {
|
|
AiTaskStage {
|
|
task_stage_id: generate_ai_task_stage_id(task_id, snapshot.stage_kind),
|
|
task_id: task_id.to_string(),
|
|
stage_kind: snapshot.stage_kind,
|
|
label: snapshot.label.clone(),
|
|
detail: snapshot.detail.clone(),
|
|
stage_order: snapshot.order,
|
|
status: snapshot.status,
|
|
text_output: snapshot.text_output.clone(),
|
|
structured_payload_json: snapshot.structured_payload_json.clone(),
|
|
warning_messages: snapshot.warning_messages.clone(),
|
|
started_at: snapshot
|
|
.started_at_micros
|
|
.map(Timestamp::from_micros_since_unix_epoch),
|
|
completed_at: snapshot
|
|
.completed_at_micros
|
|
.map(Timestamp::from_micros_since_unix_epoch),
|
|
}
|
|
}
|
|
|
|
pub(crate) fn build_ai_task_stage_snapshot_from_row(row: &AiTaskStage) -> AiTaskStageSnapshot {
|
|
AiTaskStageSnapshot {
|
|
stage_kind: row.stage_kind,
|
|
label: row.label.clone(),
|
|
detail: row.detail.clone(),
|
|
order: row.stage_order,
|
|
status: row.status,
|
|
text_output: row.text_output.clone(),
|
|
structured_payload_json: row.structured_payload_json.clone(),
|
|
warning_messages: row.warning_messages.clone(),
|
|
started_at_micros: row
|
|
.started_at
|
|
.map(|value| value.to_micros_since_unix_epoch()),
|
|
completed_at_micros: row
|
|
.completed_at
|
|
.map(|value| value.to_micros_since_unix_epoch()),
|
|
}
|
|
}
|
|
|
|
pub(crate) fn build_ai_text_chunk_row(snapshot: &AiTextChunkSnapshot) -> AiTextChunk {
|
|
AiTextChunk {
|
|
text_chunk_row_id: format!(
|
|
"{}{}_{}_{}",
|
|
AI_TEXT_CHUNK_ID_PREFIX,
|
|
snapshot.task_id,
|
|
snapshot.stage_kind.as_str(),
|
|
snapshot.sequence
|
|
),
|
|
chunk_id: snapshot.chunk_id.clone(),
|
|
task_id: snapshot.task_id.clone(),
|
|
stage_kind: snapshot.stage_kind,
|
|
sequence: snapshot.sequence,
|
|
delta_text: snapshot.delta_text.clone(),
|
|
created_at: Timestamp::from_micros_since_unix_epoch(snapshot.created_at_micros),
|
|
}
|
|
}
|
|
|
|
pub(crate) fn build_ai_text_chunk_snapshot_from_row(row: &AiTextChunk) -> AiTextChunkSnapshot {
|
|
AiTextChunkSnapshot {
|
|
chunk_id: row.chunk_id.clone(),
|
|
task_id: row.task_id.clone(),
|
|
stage_kind: row.stage_kind,
|
|
sequence: row.sequence,
|
|
delta_text: row.delta_text.clone(),
|
|
created_at_micros: row.created_at.to_micros_since_unix_epoch(),
|
|
}
|
|
}
|
|
|
|
pub(crate) fn build_ai_result_reference_row(
|
|
snapshot: &AiResultReferenceSnapshot,
|
|
) -> AiResultReference {
|
|
AiResultReference {
|
|
result_reference_row_id: format!(
|
|
"{}{}_{}",
|
|
AI_RESULT_REF_ID_PREFIX, snapshot.task_id, snapshot.result_ref_id
|
|
),
|
|
result_ref_id: snapshot.result_ref_id.clone(),
|
|
task_id: snapshot.task_id.clone(),
|
|
reference_kind: snapshot.reference_kind,
|
|
reference_id: snapshot.reference_id.clone(),
|
|
label: snapshot.label.clone(),
|
|
created_at: Timestamp::from_micros_since_unix_epoch(snapshot.created_at_micros),
|
|
}
|
|
}
|
|
|
|
pub(crate) fn build_ai_result_reference_snapshot_from_row(
|
|
row: &AiResultReference,
|
|
) -> AiResultReferenceSnapshot {
|
|
AiResultReferenceSnapshot {
|
|
result_ref_id: row.result_ref_id.clone(),
|
|
task_id: row.task_id.clone(),
|
|
reference_kind: row.reference_kind,
|
|
reference_id: row.reference_id.clone(),
|
|
label: row.label.clone(),
|
|
created_at_micros: row.created_at.to_micros_since_unix_epoch(),
|
|
}
|
|
}
|