memory: implement flat extract staging
This commit is contained in:
+16
-17
@@ -3172,15 +3172,15 @@ impl<C: LlmClient, St: Store> Worker<C, St> {
|
||||
});
|
||||
|
||||
let source_segment_id = self.segment_state.segment_id();
|
||||
let staging_id = if payload.is_empty() {
|
||||
String::new()
|
||||
let staging_results = if payload.is_empty() {
|
||||
Vec::new()
|
||||
} else {
|
||||
let source = memory::schema::SourceRef {
|
||||
segment_id: source_segment_id.to_string(),
|
||||
range: [start_entry as u64, end_entry as u64],
|
||||
};
|
||||
let (id, _) = match extract::write_staging(&layout, source, payload) {
|
||||
Ok(result) => result,
|
||||
match extract::write_staging(&layout, source, payload) {
|
||||
Ok(results) => results,
|
||||
Err(err) => {
|
||||
let usage = usage_capture
|
||||
.lock()
|
||||
@@ -3197,9 +3197,12 @@ impl<C: LlmClient, St: Store> Worker<C, St> {
|
||||
);
|
||||
return Err(WorkerError::ExtractStaging(err));
|
||||
}
|
||||
};
|
||||
id.to_string()
|
||||
}
|
||||
};
|
||||
let staging_id = staging_results
|
||||
.first()
|
||||
.map(|result| result.id.to_string())
|
||||
.unwrap_or_default();
|
||||
|
||||
let pointer_payload = extract::ExtractPointerPayload {
|
||||
processed_through_entry: end_entry,
|
||||
@@ -3220,16 +3223,12 @@ impl<C: LlmClient, St: Store> Worker<C, St> {
|
||||
.expect("extract_pointer poisoned") = Some(pointer_payload);
|
||||
|
||||
let mut extract_audit = extract_audit_base;
|
||||
if !staging_id.is_empty() {
|
||||
extract_audit.staging_count = 1;
|
||||
extract_audit.staging_ids.push(staging_id.clone());
|
||||
extract_audit.staging_paths.push(
|
||||
layout
|
||||
.staging_dir()
|
||||
.join(format!("{staging_id}.json"))
|
||||
.display()
|
||||
.to_string(),
|
||||
);
|
||||
extract_audit.staging_count = staging_results.len();
|
||||
for result in &staging_results {
|
||||
extract_audit.staging_ids.push(result.id.to_string());
|
||||
extract_audit
|
||||
.staging_paths
|
||||
.push(result.path.display().to_string());
|
||||
}
|
||||
let usage = usage_capture
|
||||
.lock()
|
||||
@@ -4819,7 +4818,7 @@ pub enum WorkerError {
|
||||
Skill(#[from] SkillClientError),
|
||||
|
||||
#[error("memory extract staging write failed: {0}")]
|
||||
ExtractStaging(#[source] memory::extract::StagingError),
|
||||
ExtractStaging(#[source] std::io::Error),
|
||||
|
||||
#[error("memory consolidation lock acquisition failed: {0}")]
|
||||
ConsolidationLock(#[source] memory::consolidate::LockError),
|
||||
|
||||
@@ -24,7 +24,7 @@ use llm_engine::Engine;
|
||||
use llm_engine::llm_client::event::{Event as LlmEvent, ResponseStatus, StatusEvent};
|
||||
use llm_engine::llm_client::{ClientError, LlmClient, Request};
|
||||
use memory::WorkspaceLayout;
|
||||
use memory::extract::{ExtractedPayload, write_staging};
|
||||
use memory::extract::{CandidateKind, ExtractedCandidate, ExtractedPayload, write_staging};
|
||||
use memory::schema::SourceRef;
|
||||
use session_store::FsStore;
|
||||
use session_store::{CombinedStore, FsWorkerStore};
|
||||
@@ -182,18 +182,32 @@ async fn make_worker_with(
|
||||
.unwrap()
|
||||
}
|
||||
|
||||
fn staging_payload(claim: String) -> ExtractedPayload {
|
||||
ExtractedPayload {
|
||||
candidates: vec![ExtractedCandidate {
|
||||
kind: CandidateKind::Lesson,
|
||||
claim,
|
||||
why_useful: "useful for consolidation trigger tests".into(),
|
||||
staleness: None,
|
||||
evidence_ids: Vec::new(),
|
||||
}],
|
||||
}
|
||||
}
|
||||
|
||||
fn write_n_staging(layout: &WorkspaceLayout, n: usize) -> Vec<uuid::Uuid> {
|
||||
let mut ids = Vec::new();
|
||||
for i in 0..n {
|
||||
let (id, _) = write_staging(
|
||||
let id = write_staging(
|
||||
layout,
|
||||
SourceRef {
|
||||
segment_id: format!("s-{i}"),
|
||||
range: [i as u64, i as u64],
|
||||
},
|
||||
ExtractedPayload::default(),
|
||||
staging_payload(format!("candidate-{i}")),
|
||||
)
|
||||
.unwrap();
|
||||
.unwrap()
|
||||
.remove(0)
|
||||
.id;
|
||||
ids.push(id);
|
||||
}
|
||||
ids
|
||||
|
||||
Reference in New Issue
Block a user