use std::collections::{BTreeMap, HashMap};
use std::sync::{Arc, Mutex};
use std::time::Duration;
use arrow::array::{Array, Int64Array, RecordBatch, TimestampMillisecondArray, UInt64Array};
use datafusion::prelude::SessionContext;
use tokio_rusqlite::Connection;
use super::config::{FeedSubscription, RssConfig, inline_config};
use super::egress::AllowAll;
use super::register_rss_tables_with_policy;
use super::testutil::{MockFeedServer, MockResponse, str_col, str_opt_col, total_rows};
use crate::model::chunking::ChunkingRegistry;
use crate::sources::hierarchy::HierarchyLevel;
use crate::sources::providers::sqlite::register_sqlite_tables;
const NEWS: &str = "news";
const ARCHIVE: &str = "archive";
const META: &str = "meta";
const QUERY_CEILING: Duration = Duration::from_secs(60);
const ARCHIVE_DDL: &str = "
CREATE TABLE IF NOT EXISTS news_items (
feed TEXT NOT NULL, guid TEXT NOT NULL, title TEXT, link TEXT, author TEXT,
published TIMESTAMP, content TEXT, PRIMARY KEY (feed, guid));
CREATE TABLE IF NOT EXISTS news_chunks (
feed TEXT NOT NULL, guid TEXT NOT NULL, chunk_idx INTEGER NOT NULL,
chunk_text TEXT NOT NULL, embedding BLOB, ingested_at TIMESTAMP,
PRIMARY KEY (feed, guid, chunk_idx));
";
const INSERT_ITEMS: &str = "\
INSERT INTO archive.main.news_items (feed, guid, title, link, author, published, content)
SELECT i.feed, i.guid, i.title, i.link, i.author, i.published, COALESCE(i.content, i.summary)
FROM news.main.items i
LEFT JOIN archive.main.news_items a ON a.feed = i.feed AND a.guid = i.guid
WHERE a.guid IS NULL";
fn insert_chunks(size: u32, overlap: u32, embedding: &str) -> String {
format!(
"\
INSERT INTO archive.main.news_chunks (feed, guid, chunk_idx, chunk_text, embedding, ingested_at)
SELECT s.feed, s.guid,
ROW_NUMBER() OVER (PARTITION BY s.feed, s.guid) - 1 AS chunk_idx,
s.chunk_text, {embedding} AS embedding, now() AS ingested_at
FROM (
SELECT n.feed, n.guid, UNNEST(chunk('markdown', n.content, {size}, {overlap})) AS chunk_text
FROM archive.main.news_items n
LEFT JOIN archive.main.news_chunks e ON e.feed = n.feed AND e.guid = n.guid
WHERE e.guid IS NULL AND n.content IS NOT NULL
) s"
)
}
const FIRST_SIZE: u32 = 1200;
const FIRST_OVERLAP: u32 = 120;
const REBUILD_SIZE: u32 = 600;
const REBUILD_OVERLAP: u32 = 60;
const HEALTH_REPORT: &str = "\
SELECT name, last_status, last_error, last_fetch FROM news.main.feeds
WHERE last_status IN ('error', 'never', 'stale-error')
ORDER BY name";
const PARA_HTML: &str = "<p>Filler paragraph with <strong>bold</strong> emphasis and a \
<a href=\"https://feed.example/more\">link</a> in it, written long enough that repeating it a \
handful of times pushes one entry's body well past the chunk target the first ingest uses.</p>";
const PARA_MD: &str = "Filler paragraph with **bold** emphasis and a \
[link](https://feed.example/more) in it, written long enough that repeating it a handful of \
times pushes one entry's body well past the chunk target the first ingest uses.";
const PARAS: usize = 6;
const ENTRIES: [(&str, &str, &str); 3] = [
(
"news-1",
"Alpha announcement",
"Mon, 20 Jul 2026 10:00:00 GMT",
),
(
"news-2",
"Beta announcement",
"Tue, 21 Jul 2026 11:00:00 GMT",
),
(
"news-3",
"Gamma announcement",
"Wed, 22 Jul 2026 12:00:00 GMT",
),
];
fn entry_html(title: &str) -> String {
let mut html = format!("<h2>{title}</h2>");
for _ in 0..PARAS {
html.push_str(PARA_HTML);
}
html
}
fn entry_markdown(title: &str) -> String {
let mut markdown = format!("## {title}");
for _ in 0..PARAS {
markdown.push_str("\n\n");
markdown.push_str(PARA_MD);
}
markdown
}
fn archive_feed(count: usize) -> String {
let mut doc = String::from(
r#"<?xml version="1.0" encoding="UTF-8"?>
<rss version="2.0" xmlns:content="http://purl.org/rss/1.0/modules/content/"
xmlns:dc="http://purl.org/dc/elements/1.1/">
<channel><title>World News</title><link>https://feed.example/</link>
<description>The archive fixtures' channel.</description>"#,
);
for (guid, title, pub_date) in &ENTRIES[ENTRIES.len() - count..] {
doc.push_str(&format!(
"<item><guid isPermaLink=\"false\">{guid}</guid><title>{title}</title>\
<link>https://feed.example/{guid}</link>\
<dc:creator>Ada Lovelace</dc:creator><pubDate>{pub_date}</pubDate>\
<content:encoded><![CDATA[{}]]></content:encoded></item>",
entry_html(title)
));
}
doc.push_str("</channel></rss>");
doc
}
fn simple_feed(guids: &[&str]) -> String {
let mut doc = String::from(
r#"<rss version="2.0"><channel><title>Mock Feed</title>
<link>https://feed.example/</link><description>A mock feed.</description>"#,
);
for guid in guids {
doc.push_str(&format!(
"<item><guid>{guid}</guid><title>{guid} title</title>\
<link>https://feed.example/{guid}</link>\
<description>Short summary for {guid}.</description></item>"
));
}
doc.push_str("</channel></rss>");
doc
}
#[derive(Clone)]
struct Canned {
status: u16,
body: Vec<u8>,
}
impl Canned {
fn xml(body: &str) -> Self {
Self {
status: 200,
body: body.as_bytes().to_vec(),
}
}
fn status(status: u16) -> Self {
Self {
status,
body: Vec::new(),
}
}
}
struct MockFeeds {
server: MockFeedServer,
paths: Arc<Mutex<HashMap<String, Canned>>>,
}
impl MockFeeds {
async fn start() -> Self {
let paths: Arc<Mutex<HashMap<String, Canned>>> = Arc::new(Mutex::new(HashMap::new()));
let handler_paths = Arc::clone(&paths);
let server = MockFeedServer::start(move |request| {
let paths = handler_paths
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
match paths.get(&request.path) {
Some(canned) => MockResponse::new(canned.status, canned.body.clone())
.with_header("content-type", "application/xml"),
None => MockResponse::status(404),
}
})
.await;
Self { server, paths }
}
fn serve(&self, path: &str, canned: Canned) {
self.paths
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.insert(path.to_string(), canned);
}
fn url(&self) -> String {
self.server.url()
}
fn request_count(&self) -> usize {
self.server.requests().len()
}
fn sorted_paths(&self) -> Vec<String> {
let mut paths: Vec<String> = self
.server
.requests()
.into_iter()
.map(|request| request.path)
.collect();
paths.sort();
paths
}
}
struct ArchiveDb {
_dir: tempfile::TempDir,
path: String,
conn: Connection,
}
impl ArchiveDb {
async fn create() -> Self {
let dir = tempfile::tempdir().expect("create temp dir");
let path = dir
.path()
.join("archive.db")
.to_str()
.expect("temp path is utf-8")
.to_string();
let conn = Connection::open(&path).await.expect("open archive db");
conn.call(
|conn| -> std::result::Result<(), tokio_rusqlite::rusqlite::Error> {
conn.execute_batch(ARCHIVE_DDL)
},
)
.await
.expect("create the archive schema");
Self {
_dir: dir,
path,
conn,
}
}
async fn execute(&self, sql: &'static str) -> usize {
self.conn
.call(move |conn| conn.execute(sql, []))
.await
.unwrap_or_else(|e| panic!("archive statement {sql:?}: {e}"))
}
async fn count(&self, table: &'static str) -> i64 {
self.conn
.call(move |conn| {
conn.query_row(&format!("SELECT count(*) FROM {table}"), [], |row| {
row.get(0)
})
})
.await
.unwrap_or_else(|e| panic!("counting {table}: {e}"))
}
}
struct Composition {
ctx: SessionContext,
}
impl Composition {
async fn register(
mock: &MockFeeds,
feeds: &[(&str, &str)],
archive: &ArchiveDb,
tune: impl FnOnce(&mut RssConfig),
) -> Self {
let subscriptions = feeds
.iter()
.map(|(name, path)| FeedSubscription {
url: format!("{}{path}", mock.url()),
name: Some((*name).to_string()),
})
.collect();
let mut config = inline_config(subscriptions);
config.request_timeout_seconds = 5;
config.scan_timeout_seconds = 20;
tune(&mut config);
let mut ctx = SessionContext::new();
register_rss_tables_with_policy(
&mut ctx,
NEWS,
Some(&config),
false,
HierarchyLevel::Catalog,
Arc::new(AllowAll),
)
.await
.expect("registering the rss source succeeds");
register_sqlite_tables(
&mut ctx,
ARCHIVE,
&archive.path,
None,
true,
None,
HierarchyLevel::Catalog,
)
.await
.expect("registering the writable archive succeeds");
Arc::new(ChunkingRegistry::new()).register_chunk_udf(&mut ctx);
Self { ctx }
}
async fn with_meta(mut self, db_path: &str) -> Self {
register_sqlite_tables(
&mut self.ctx,
META,
db_path,
None,
false,
None,
HierarchyLevel::Catalog,
)
.await
.expect("registering the read-only meta source succeeds");
self
}
async fn sql(&self, sql: &str) -> Vec<RecordBatch> {
tokio::time::timeout(QUERY_CEILING, async {
self.ctx
.sql(sql)
.await
.unwrap_or_else(|e| panic!("plan {sql:?}: {e}"))
.collect()
.await
.unwrap_or_else(|e| panic!("execute {sql:?}: {e}"))
})
.await
.unwrap_or_else(|_| panic!("{sql:?} did not finish within {QUERY_CEILING:?}"))
}
async fn insert(&self, sql: &str) -> u64 {
let batches = self.sql(sql).await;
assert_eq!(
total_rows(&batches),
1,
"an INSERT reports exactly one count row"
);
batches[0]
.column(0)
.as_any()
.downcast_ref::<UInt64Array>()
.expect("an INSERT's count column is UInt64")
.value(0)
}
}
fn col(batches: &[RecordBatch], name: &str) -> Vec<String> {
batches.iter().flat_map(|b| str_col(b, name)).collect()
}
fn opt_col(batches: &[RecordBatch], name: &str) -> Vec<Option<String>> {
batches.iter().flat_map(|b| str_opt_col(b, name)).collect()
}
fn i64_col(batches: &[RecordBatch], name: &str) -> Vec<i64> {
batches
.iter()
.flat_map(|batch| {
let index = batch
.schema()
.index_of(name)
.unwrap_or_else(|e| panic!("batch has no column {name:?}: {e}"));
let column = batch
.column(index)
.as_any()
.downcast_ref::<Int64Array>()
.unwrap_or_else(|| panic!("column {name:?} is not Int64"));
assert_eq!(column.null_count(), 0, "column {name:?} has NULLs");
column.values().to_vec()
})
.collect()
}
fn ts_present(batches: &[RecordBatch], name: &str) -> Vec<bool> {
batches
.iter()
.flat_map(|batch| {
let index = batch
.schema()
.index_of(name)
.unwrap_or_else(|e| panic!("batch has no column {name:?}: {e}"));
let column = batch
.column(index)
.as_any()
.downcast_ref::<TimestampMillisecondArray>()
.unwrap_or_else(|| panic!("column {name:?} is not Timestamp(Millisecond)"));
(0..column.len())
.map(|row| column.is_valid(row))
.collect::<Vec<_>>()
})
.collect()
}
#[track_caller]
fn assert_contains(actual: &str, needle: &str) {
assert!(actual.contains(needle), "expected {needle:?} in {actual:?}");
}
#[track_caller]
fn chunks_by_entry(batches: &[RecordBatch]) -> BTreeMap<(String, String), Vec<String>> {
let feeds = col(batches, "feed");
let guids = col(batches, "guid");
let indices = i64_col(batches, "chunk_idx");
let texts = col(batches, "chunk_text");
assert_eq!(feeds.len(), guids.len());
assert_eq!(feeds.len(), indices.len());
assert_eq!(feeds.len(), texts.len());
let mut grouped: BTreeMap<(String, String), Vec<(i64, String)>> = BTreeMap::new();
for row in 0..feeds.len() {
grouped
.entry((feeds[row].clone(), guids[row].clone()))
.or_default()
.push((indices[row], texts[row].clone()));
}
grouped
.into_iter()
.map(|(entry, mut rows)| {
rows.sort_by_key(|(index, _)| *index);
let seen: Vec<i64> = rows.iter().map(|(index, _)| *index).collect();
let dense: Vec<i64> = (0..rows.len() as i64).collect();
assert_eq!(
seen, dense,
"chunk_idx for {entry:?} is not a dense 0..n-1 range"
);
let texts = rows.into_iter().map(|(_, text)| text).collect();
(entry, texts)
})
.collect()
}
fn sorted(mut values: Vec<String>) -> Vec<String> {
values.sort();
values
}
#[tokio::test]
async fn federated_join_items_with_sqlite() {
let mock = MockFeeds::start().await;
for feed in ["a", "b", "c"] {
mock.serve(
&format!("/{feed}.xml"),
Canned::xml(&simple_feed(&[&format!("{feed}-1")])),
);
}
let meta_dir = tempfile::tempdir().expect("create temp dir");
let meta_path = meta_dir
.path()
.join("meta.db")
.to_str()
.expect("temp path is utf-8")
.to_string();
let meta_conn = Connection::open(&meta_path).await.expect("open meta db");
meta_conn
.call(
|conn| -> std::result::Result<(), tokio_rusqlite::rusqlite::Error> {
conn.execute_batch(
"CREATE TABLE feed_meta (feed TEXT, tier TEXT);
INSERT INTO feed_meta (feed, tier) VALUES ('a', 'primary');
INSERT INTO feed_meta (feed, tier) VALUES ('b', 'secondary');",
)
},
)
.await
.expect("seed feed_meta");
let archive = ArchiveDb::create().await;
let news = Composition::register(
&mock,
&[("a", "/a.xml"), ("b", "/b.xml"), ("c", "/c.xml")],
&archive,
|_| {},
)
.await
.with_meta(&meta_path)
.await;
let joined = news
.sql(
"SELECT i.guid, m.tier \
FROM news.main.items i JOIN meta.main.feed_meta m ON m.feed = i.feed \
ORDER BY i.guid",
)
.await;
assert_eq!(
col(&joined, "guid"),
vec!["a-1", "b-1"],
"the subscription with no metadata row is dropped by the inner join"
);
assert_eq!(col(&joined, "tier"), vec!["primary", "secondary"]);
assert_eq!(
mock.sorted_paths(),
vec!["/a.xml", "/b.xml", "/c.xml"],
"the join predicate is not on `feed`'s equality to a literal, so nothing prunes"
);
}
#[tokio::test]
async fn archive_ingest_is_idempotent_and_survives_window_roll() {
let mock = MockFeeds::start().await;
mock.serve("/world.xml", Canned::xml(&archive_feed(3)));
let archive = ArchiveDb::create().await;
let news = Composition::register(&mock, &[("world", "/world.xml")], &archive, |config| {
config.ttl_seconds = 0
})
.await;
assert_eq!(news.insert(INSERT_ITEMS).await, 3, "three entries are new");
let first_chunks = news
.insert(&insert_chunks(FIRST_SIZE, FIRST_OVERLAP, "NULL"))
.await;
assert!(
first_chunks >= 3,
"each entry's body is longer than {FIRST_SIZE} characters, so three entries owe at \
least three chunks; got {first_chunks}"
);
assert_eq!(archive.count("news_items").await, 3);
assert_eq!(archive.count("news_chunks").await, first_chunks as i64);
let archived = news
.sql("SELECT guid, content FROM archive.main.news_items ORDER BY guid")
.await;
assert_eq!(col(&archived, "guid"), vec!["news-1", "news-2", "news-3"]);
assert_eq!(
opt_col(&archived, "content")[0].as_deref(),
Some(entry_markdown(ENTRIES[0].1).as_str()),
"the archived body is the Markdown the conversion owes, not the HTML on the wire"
);
let served = news
.sql("SELECT guid, content FROM news.main.items ORDER BY guid")
.await;
assert_eq!(
opt_col(&archived, "content"),
opt_col(&served, "content"),
"the archive stores content exactly as `items` served it"
);
let stored = news
.sql("SELECT feed, guid, chunk_idx, chunk_text FROM archive.main.news_chunks")
.await;
let grouped = chunks_by_entry(&stored);
assert_eq!(grouped.len(), 3, "every entry contributed chunks");
for ((feed, guid), texts) in &grouped {
assert_eq!(feed, "world");
assert!(
texts.len() > 1,
"{guid} produced {} chunk(s); the body is no longer long enough for the \
chunk-index assertions to mean anything",
texts.len()
);
let expected = news
.sql(&format!(
"SELECT UNNEST(chunk('markdown', content, {FIRST_SIZE}, {FIRST_OVERLAP})) \
AS chunk_text FROM archive.main.news_items WHERE guid = '{guid}'"
))
.await;
assert_eq!(
sorted(texts.clone()),
sorted(col(&expected, "chunk_text")),
"the stored chunks for {guid} are not the ones chunk() produces"
);
}
let before_rerun = mock.request_count();
assert_eq!(
news.insert(INSERT_ITEMS).await,
0,
"re-running the entry INSERT wrote rows a second time"
);
assert_eq!(
news.insert(&insert_chunks(FIRST_SIZE, FIRST_OVERLAP, "NULL"))
.await,
0,
"re-running the chunk INSERT wrote rows a second time"
);
assert!(
mock.request_count() > before_rerun,
"the second run was served from cache, so it did not exercise the anti-join"
);
assert_eq!(archive.count("news_items").await, 3);
assert_eq!(archive.count("news_chunks").await, first_chunks as i64);
mock.serve("/world.xml", Canned::xml(&archive_feed(2)));
let live = news
.sql("SELECT guid FROM news.main.items ORDER BY guid")
.await;
assert_eq!(
col(&live, "guid"),
vec!["news-2", "news-3"],
"the live window did not actually shrink, so nothing below is a citability test"
);
assert_eq!(news.insert(INSERT_ITEMS).await, 0, "no entry is new");
assert_eq!(
news.insert(&insert_chunks(FIRST_SIZE, FIRST_OVERLAP, "NULL"))
.await,
0
);
assert_eq!(
archive.count("news_items").await,
3,
"an entry leaving the live window must not remove it from the archive"
);
assert_eq!(archive.count("news_chunks").await, first_chunks as i64);
let dropped = news
.sql(
"SELECT guid, title, link, published FROM archive.main.news_items \
WHERE guid = 'news-1'",
)
.await;
assert_eq!(col(&dropped, "guid"), vec!["news-1"]);
assert_eq!(
opt_col(&dropped, "title"),
vec![Some(ENTRIES[0].1.to_string())]
);
assert_eq!(
opt_col(&dropped, "link"),
vec![Some("https://feed.example/news-1".to_string())]
);
assert_eq!(
opt_col(&dropped, "published"),
vec![Some("2026-07-20T10:00:00Z".to_string())],
"the dropped entry's published time survived the window roll"
);
}
#[tokio::test]
async fn sync_closing_health_report_shape() {
let mock = MockFeeds::start().await;
for feed in ["alive", "dying"] {
mock.serve(
&format!("/{feed}.xml"),
Canned::xml(&simple_feed(&[&format!("{feed}-1")])),
);
}
let archive = ArchiveDb::create().await;
let news = Composition::register(
&mock,
&[("alive", "/alive.xml"), ("dying", "/dying.xml")],
&archive,
|config| config.ttl_seconds = 0,
)
.await;
let untouched = news.sql(HEALTH_REPORT).await;
assert_eq!(col(&untouched, "name"), vec!["alive", "dying"]);
assert_eq!(col(&untouched, "last_status"), vec!["never", "never"]);
assert_eq!(
opt_col(&untouched, "last_error"),
vec![None, None],
"a never-attempted subscription has no error to report, only a status"
);
assert_eq!(
ts_present(&untouched, "last_fetch"),
vec![false, false],
"and no as-of time"
);
assert_eq!(
mock.request_count(),
0,
"the health report reached the fetcher"
);
news.sql("SELECT guid FROM news.main.items").await;
let after_scan = mock.request_count();
let healthy = news.sql(HEALTH_REPORT).await;
assert_eq!(
total_rows(&healthy),
0,
"a healthy run's report lists nothing: {:?}",
col(&healthy, "name")
);
assert_eq!(
mock.request_count(),
after_scan,
"the health report reached the fetcher"
);
mock.serve("/dying.xml", Canned::status(500));
news.sql("SELECT guid FROM news.main.items").await;
let degraded = news.sql(HEALTH_REPORT).await;
assert_eq!(
col(°raded, "name"),
vec!["dying"],
"the healthy neighbour must not appear in the report"
);
assert_eq!(
col(°raded, "last_status"),
vec!["stale-error"],
"the feed has a cached window, so it degrades rather than going dark"
);
assert_contains(
opt_col(°raded, "last_error")[0]
.as_deref()
.expect("a degraded subscription reports why"),
"http status 500",
);
assert_eq!(
ts_present(°raded, "last_fetch"),
vec![true],
"the report carries the as-of time of the failed attempt"
);
let items = news
.sql("SELECT feed, window_status FROM news.main.items ORDER BY feed")
.await;
assert_eq!(col(&items, "feed"), vec!["alive", "dying"]);
assert_eq!(col(&items, "window_status"), vec!["fresh", "stale-error"]);
}
#[tokio::test]
async fn subscription_add_is_config_only() {
let mock = MockFeeds::start().await;
for feed in ["a", "b", "c"] {
mock.serve(
&format!("/{feed}.xml"),
Canned::xml(&simple_feed(&[&format!("{feed}-1")])),
);
}
let archive = ArchiveDb::create().await;
let v1 = Composition::register(
&mock,
&[("a", "/a.xml"), ("b", "/b.xml")],
&archive,
|config| config.ttl_seconds = 0,
)
.await;
assert_eq!(v1.insert(INSERT_ITEMS).await, 2);
assert_eq!(archive.count("news_items").await, 2);
assert_eq!(mock.sorted_paths(), vec!["/a.xml", "/b.xml"]);
drop(v1);
let after_v1 = mock.request_count();
let v2 = Composition::register(
&mock,
&[("a", "/a.xml"), ("b", "/b.xml"), ("c", "/c.xml")],
&archive,
|config| config.ttl_seconds = 0,
)
.await;
assert_eq!(
mock.request_count(),
after_v1,
"registering the new subscription performed network I/O"
);
assert_eq!(
v2.insert(INSERT_ITEMS).await,
1,
"only the added subscription's entry is new to the archive"
);
assert_eq!(archive.count("news_items").await, 3);
let archived = news_guids(&v2).await;
assert_eq!(archived, vec!["a-1", "b-1", "c-1"]);
assert!(
mock.sorted_paths().contains(&"/c.xml".to_string()),
"the new subscription's first items scan is what fetched it: {:?}",
mock.sorted_paths()
);
let feeds = v2
.sql("SELECT name, last_status, last_error FROM news.main.feeds ORDER BY name")
.await;
assert_eq!(col(&feeds, "name"), vec!["a", "b", "c"]);
assert_eq!(
col(&feeds, "last_status"),
vec!["fresh", "fresh", "fresh"],
"the added feed was fetched successfully alongside the existing two"
);
assert_eq!(opt_col(&feeds, "last_error"), vec![None, None, None]);
}
async fn news_guids(news: &Composition) -> Vec<String> {
let batches = news
.sql("SELECT guid FROM archive.main.news_items ORDER BY guid")
.await;
col(&batches, "guid")
}
#[tokio::test]
async fn parameter_change_rebuild_from_retained_content() {
let mock = MockFeeds::start().await;
mock.serve("/world.xml", Canned::xml(&archive_feed(3)));
let archive = ArchiveDb::create().await;
let news = Composition::register(&mock, &[("world", "/world.xml")], &archive, |config| {
config.ttl_seconds = 0
})
.await;
assert_eq!(news.insert(INSERT_ITEMS).await, 3);
let coarse = news
.insert(&insert_chunks(FIRST_SIZE, FIRST_OVERLAP, "NULL"))
.await;
let after_ingest = mock.request_count();
assert_eq!(after_ingest, 1, "the ingest fetched the feed exactly once");
mock.serve("/world.xml", Canned::status(500));
let deleted = archive.execute("DELETE FROM news_chunks").await;
assert_eq!(deleted, coarse as usize, "the derived table was emptied");
assert_eq!(archive.count("news_chunks").await, 0);
let fine = news
.insert(&insert_chunks(REBUILD_SIZE, REBUILD_OVERLAP, "NULL"))
.await;
assert!(
fine > coarse,
"a {REBUILD_SIZE}-character target must split the same bodies into more chunks than \
a {FIRST_SIZE}-character one did ({fine} vs {coarse})"
);
assert_eq!(archive.count("news_chunks").await, fine as i64);
assert_eq!(
mock.request_count(),
after_ingest,
"the rebuild touched the live window"
);
let stored = news
.sql("SELECT feed, guid, chunk_idx, chunk_text FROM archive.main.news_chunks")
.await;
let grouped = chunks_by_entry(&stored);
assert_eq!(grouped.len(), 3);
for ((_, guid), texts) in &grouped {
let expected = news
.sql(&format!(
"SELECT UNNEST(chunk('markdown', content, {REBUILD_SIZE}, {REBUILD_OVERLAP})) \
AS chunk_text FROM archive.main.news_items WHERE guid = '{guid}'"
))
.await;
assert_eq!(
sorted(texts.clone()),
sorted(col(&expected, "chunk_text")),
"the rebuilt chunks for {guid} are not the ones the new parameters produce"
);
}
assert_eq!(archive.count("news_items").await, 3);
let archived = news
.sql("SELECT guid, content FROM archive.main.news_items ORDER BY guid")
.await;
assert_eq!(
opt_col(&archived, "content")[0].as_deref(),
Some(entry_markdown(ENTRIES[0].1).as_str())
);
}
#[cfg(feature = "candle")]
#[tokio::test]
#[ignore = "live: requires SKARDI_TEST_EMBED_MODEL pointing at a local embedding model dir"]
async fn archive_ingest_with_candle_embeddings() {
use crate::model::candle::CandleModelRegistry;
use crate::sources::providers::sqlite::register_vec_to_binary_udf;
use arrow::array::BinaryArray;
let Ok(model) = std::env::var("SKARDI_TEST_EMBED_MODEL") else {
eprintln!("skipping: SKARDI_TEST_EMBED_MODEL not set (needs a local embedding model dir)");
return;
};
assert!(
!model.contains('\''),
"the model path is interpolated into SQL as a literal: {model:?}"
);
let mock = MockFeeds::start().await;
mock.serve("/world.xml", Canned::xml(&archive_feed(3)));
let archive = ArchiveDb::create().await;
let mut news = Composition::register(&mock, &[("world", "/world.xml")], &archive, |_| {}).await;
Arc::new(CandleModelRegistry::new()).register_candle_udf(&mut news.ctx);
register_vec_to_binary_udf(&mut news.ctx);
assert_eq!(news.insert(INSERT_ITEMS).await, 3);
let written = news
.insert(&insert_chunks(
FIRST_SIZE,
FIRST_OVERLAP,
&format!("vec_to_binary(candle('{model}', s.chunk_text))"),
))
.await;
assert!(written >= 3, "three entries owe at least three chunks");
let stored = news
.sql("SELECT embedding FROM archive.main.news_chunks")
.await;
assert_eq!(total_rows(&stored) as u64, written);
for batch in &stored {
let column = batch
.column(0)
.as_any()
.downcast_ref::<BinaryArray>()
.expect("a sqlite BLOB column reads back as Binary");
for row in 0..column.len() {
assert!(column.is_valid(row), "row {row} stored a NULL embedding");
let bytes = column.value(row);
assert!(!bytes.is_empty(), "row {row} stored an empty embedding");
assert_eq!(
bytes.len() % 4,
0,
"a packed f32 blob is a whole number of 4-byte lanes, got {} bytes",
bytes.len()
);
}
}
}