Skip to content

Commit 227f9b9

Browse files
committed
fix(ingest): persist staged content pointers
1 parent cdc8978 commit 227f9b9

2 files changed

Lines changed: 15 additions & 5 deletions

File tree

src/memory/ingest/pipeline.rs

Lines changed: 8 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -17,7 +17,7 @@ use chrono::Utc;
1717

1818
use crate::memory::chunks::{
1919
self, chunk_markdown, claim_source_ingest_tx, delete_source_ingest, get_chunk_lifecycle_status,
20-
is_source_ingested, set_chunk_lifecycle_status, set_chunk_raw_refs, upsert_chunks,
20+
is_source_ingested, set_chunk_lifecycle_status, set_chunk_raw_refs, upsert_staged_chunks_tx,
2121
with_connection, ChunkerInput, ChunkerOptions, SourceKind, CHUNK_STATUS_PENDING_EXTRACTION,
2222
};
2323
use crate::memory::config::MemoryConfig;
@@ -146,7 +146,7 @@ async fn persist_score_enqueue(
146146
) -> Result<IngestSummary> {
147147
// 3. Write each chunk body to the content store (atomic write + sha256).
148148
let content_root = chunks::content_root(config);
149-
content::stage_chunks(&content_root, &chunks)
149+
let staged = content::stage_chunks(&content_root, &chunks)
150150
.map_err(|e| anyhow!("stage_chunks failed: {e}"))?;
151151

152152
// 4. Snapshot each chunk's CURRENT lifecycle BEFORE the upsert. A chunk that
@@ -159,7 +159,12 @@ async fn persist_score_enqueue(
159159
}
160160

161161
// 5. Persist chunk rows (idempotent on deterministic chunk id).
162-
let chunks_written = upsert_chunks(config, &chunks)?;
162+
let chunks_written = with_connection(config, |connection| {
163+
let transaction = connection.unchecked_transaction()?;
164+
let count = upsert_staged_chunks_tx(&transaction, &staged)?;
165+
transaction.commit()?;
166+
Ok(count)
167+
})?;
163168

164169
// 5b. Raw-archive-backed bodies: attach refs so a worker can resolve them.
165170
if let Some(refs) = opts.raw_refs.as_ref() {

src/memory/ingest/pipeline_tests.rs

Lines changed: 7 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -6,8 +6,8 @@ use tempfile::TempDir;
66

77
use super::{ingest_chat, ingest_document, ingest_document_versioned, ingest_email_with_raw_refs};
88
use crate::memory::chunks::{
9-
count_chunks, get_chunk_lifecycle_status, is_source_ingested, SourceKind,
10-
CHUNK_STATUS_PENDING_EXTRACTION,
9+
count_chunks, get_chunk_content_pointers, get_chunk_lifecycle_status, is_source_ingested,
10+
SourceKind, CHUNK_STATUS_PENDING_EXTRACTION,
1111
};
1212
use crate::memory::chunks::{get_chunk_raw_refs, RawRef};
1313
use crate::memory::config::MemoryConfig;
@@ -99,6 +99,11 @@ async fn ingest_chat_writes_chunks_and_enqueues_extract_jobs() {
9999
assert!(out.chunks_written >= 1);
100100
assert_eq!(count_chunks(&cfg).unwrap(), out.chunks_written as u64);
101101
assert_eq!(out.chunk_ids.len(), out.chunks_written);
102+
let pointers = get_chunk_content_pointers(&cfg, &out.chunk_ids[0])
103+
.unwrap()
104+
.unwrap();
105+
assert!(!pointers.0.is_empty());
106+
assert!(!pointers.1.is_empty());
102107

103108
// Every scheduled chunk got an extract job and is parked at pending.
104109
assert_eq!(out.extract_jobs_enqueued, out.chunk_ids.len());

0 commit comments

Comments
 (0)