Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
67 changes: 40 additions & 27 deletions crates/utopia-store/src/sources.rs
Original file line number Diff line number Diff line change
Expand Up @@ -84,37 +84,50 @@ pub async fn list(pool: &PgPool, kb_id: Uuid) -> AppResult<Vec<SourceView>> {
// config 剔掉凭据:列表给 Viewer 看,哪一种连接器的密钥都不下发。
// 键在 `SOURCE_SECRET_KEYS` 一张表上——从前这里只减 `auth_header`,五种连接器
// 的密钥就这么漏出去的(#246)
// 外层别名不能与 ENTRY_SELECT 的来源别名重名;来源和当前代谓词必须留在投影内部,
// 否则 PostgreSQL 会先计算部署中的全部 observation,再丢弃与当前来源无关的行。
let rows: Vec<SourceView> = sqlx::query_as(
&format!("WITH projected AS ({}) SELECT s.id, s.kind, s.name, s.config - $2::text[] AS config, s.icon,
s.sync_interval_minutes, s.sync_cron,
s.last_sync_at, s.last_sync_status, s.last_sync_error, s.last_sync_added,
&format!("SELECT listed_source.id, listed_source.kind, listed_source.name,
listed_source.config - $2::text[] AS config, listed_source.icon,
listed_source.sync_interval_minutes, listed_source.sync_cron,
listed_source.last_sync_at, listed_source.last_sync_status,
listed_source.last_sync_error, listed_source.last_sync_added,
(SELECT count(*) FROM documents d
WHERE d.source_id = s.id AND d.deleted_at IS NULL) AS doc_count,
WHERE d.source_id = listed_source.id AND d.deleted_at IS NULL) AS doc_count,
(SELECT count(*) FROM documents d
WHERE d.source_id = s.id AND d.missing_since IS NOT NULL
WHERE d.source_id = listed_source.id AND d.missing_since IS NOT NULL
AND d.deleted_at IS NULL) AS missing_count,
CASE WHEN s.kind <> 'rss' THEN NULL
WHEN s.config->>'content_mode' IS DISTINCT FROM 'full_new_items' THEN 'disabled'
WHEN s.rss_baselined_at IS NULL THEN 'pending' ELSE 'active' END AS rss_full_content_state,
CASE WHEN s.kind='rss' THEN s.rss_generation END AS rss_full_content_generation,
(SELECT count(*)::int FROM projected e
WHERE e.source_id=s.id AND e.activation_generation=s.rss_generation AND e.state='baseline') AS rss_full_content_baseline_count,
COALESCE((SELECT count(*) FROM projected e
WHERE e.source_id = s.id AND e.activation_generation = s.rss_generation
AND e.state = 'pending'), 0) AS rss_full_content_pending_count,
COALESCE((SELECT count(*) FROM projected e
WHERE e.source_id = s.id AND e.activation_generation = s.rss_generation
AND e.state IN ('queued', 'hydrating')), 0) AS rss_full_content_queued_count,
COALESCE((SELECT count(*) FROM projected e
WHERE e.source_id = s.id AND e.activation_generation = s.rss_generation
AND e.state = 'retry_wait'), 0) AS rss_full_content_retrying_count,
COALESCE((SELECT count(*) FROM projected e
WHERE e.source_id = s.id AND e.activation_generation = s.rss_generation
AND e.state = 'complete'), 0) AS rss_full_content_complete_count,
COALESCE((SELECT count(*) FROM projected e
WHERE e.source_id = s.id AND e.activation_generation = s.rss_generation
AND e.state IN ('terminal','deleted','superseded')), 0) AS rss_full_content_terminal_count
FROM sources s WHERE s.kb_id = $1 ORDER BY s.created_at", crate::rss_full_content::ENTRY_SELECT),
CASE WHEN listed_source.kind <> 'rss' THEN NULL
WHEN listed_source.config->>'content_mode' IS DISTINCT FROM 'full_new_items' THEN 'disabled'
WHEN listed_source.rss_baselined_at IS NULL THEN 'pending' ELSE 'active' END AS rss_full_content_state,
CASE WHEN listed_source.kind='rss' THEN listed_source.rss_generation END AS rss_full_content_generation,
COALESCE(hydration.baseline_count, 0) AS rss_full_content_baseline_count,
COALESCE(hydration.pending, 0) AS rss_full_content_pending_count,
COALESCE(hydration.queued, 0) AS rss_full_content_queued_count,
COALESCE(hydration.retrying, 0) AS rss_full_content_retrying_count,
COALESCE(hydration.complete, 0) AS rss_full_content_complete_count,
COALESCE(hydration.terminal, 0) AS rss_full_content_terminal_count
FROM sources listed_source
LEFT JOIN LATERAL (
SELECT
(count(*) FILTER (WHERE projected.state = 'baseline'))::int AS baseline_count,
count(*) FILTER (WHERE projected.state = 'pending') AS pending,
count(*) FILTER (WHERE projected.state IN ('queued', 'hydrating')) AS queued,
count(*) FILTER (WHERE projected.state = 'retry_wait') AS retrying,
count(*) FILTER (WHERE projected.state = 'complete') AS complete,
count(*) FILTER (
WHERE projected.state IN ('terminal', 'deleted', 'superseded')
) AS terminal
FROM (
{}
WHERE listed_source.kind = 'rss'
AND e.source_id = listed_source.id
AND e.activation_generation = listed_source.rss_generation
) projected
) hydration ON listed_source.kind = 'rss'
WHERE listed_source.kb_id = $1 ORDER BY listed_source.created_at",
crate::rss_full_content::ENTRY_SELECT
),
)
.bind(kb_id)
.bind(SOURCE_SECRET_KEYS)
Expand Down
315 changes: 315 additions & 0 deletions crates/utopia-store/tests/rss_full_content.rs
Original file line number Diff line number Diff line change
Expand Up @@ -296,6 +296,321 @@ async fn seed_source(pool: &PgPool) -> anyhow::Result<(Uuid, Uuid, Uuid, Uuid)>
Ok((org_id, kb_id, source_id, workspace_id))
}

async fn insert_source(
pool: &PgPool,
kb_id: Uuid,
kind: &str,
name: &str,
config: serde_json::Value,
) -> anyhow::Result<Uuid> {
let id = Uuid::now_v7();
sqlx::query(
"INSERT INTO sources (id, kb_id, kind, name, config)
VALUES ($1, $2, $3, $4, $5)",
)
.bind(id)
.bind(kb_id)
.bind(kind)
.bind(name)
.bind(config)
.execute(pool)
.await?;
Ok(id)
}

async fn job_id_for(pool: &PgPool, source_id: Uuid, external_key: &str) -> anyhow::Result<i64> {
let job_id: Option<i64> = sqlx::query_scalar(
"SELECT e.current_job_id
FROM rss_full_content_entries e
WHERE e.source_id = $1 AND e.external_key = $2",
)
.bind(source_id)
.bind(external_key)
.fetch_one(pool)
.await?;
job_id.ok_or_else(|| anyhow::anyhow!("expected a hydration job for {external_key}"))
}

#[tokio::test]
async fn source_list_exposes_scoped_rss_summary() -> anyhow::Result<()> {
let Some(url) = utopia_store::test_db::url() else {
return Ok(());
};
let pool = PgPool::connect(&url).await?;
let (org_id, kb_id, source_id, _) = seed_source(&pool).await?;
let folder_id =
insert_source(&pool, kb_id, "folder", "folder-test", serde_json::json!({})).await?;
let url_id = insert_source(
&pool,
kb_id,
"url",
"url-test",
serde_json::json!({"urls":["https://example.com"]}),
)
.await?;
sqlx::query(
"UPDATE sources
SET created_at = CASE id
WHEN $1 THEN '2020-01-01T00:00:00Z'::timestamptz
WHEN $2 THEN '2020-01-02T00:00:00Z'::timestamptz
WHEN $3 THEN '2020-01-03T00:00:00Z'::timestamptz
END
WHERE id IN ($1, $2, $3)",
)
.bind(source_id)
.bind(folder_id)
.bind(url_id)
.execute(&pool)
.await?;

let mut tx = pool.begin().await?;
utopia_store::rss_full_content::initialize_source(&mut tx, source_id).await?;
tx.commit().await?;

let listed = utopia_store::sources::list(&pool, kb_id).await?;
assert_eq!(
listed
.iter()
.map(|source| source.name.as_str())
.collect::<Vec<_>>(),
vec!["rss-full-content-test", "folder-test", "url-test"]
);
for (id, kind) in [(folder_id, "folder"), (url_id, "url")] {
let source = listed
.iter()
.find(|source| source.id == id)
.ok_or_else(|| anyhow::anyhow!("{kind} source is missing"))?;
assert!(source.rss_full_content_state.is_none());
assert!(source.rss_full_content_generation.is_none());
assert_eq!(source.rss_full_content_baseline_count, Some(0));
assert_eq!(source.rss_full_content_pending_count, 0);
assert_eq!(source.rss_full_content_queued_count, 0);
assert_eq!(source.rss_full_content_retrying_count, 0);
assert_eq!(source.rss_full_content_complete_count, 0);
assert_eq!(source.rss_full_content_terminal_count, 0);
}
let rss = listed
.iter()
.find(|source| source.id == source_id)
.ok_or_else(|| anyhow::anyhow!("RSS source is missing"))?;
assert_eq!(rss.rss_full_content_state.as_deref(), Some("pending"));
assert_eq!(
(
rss.rss_full_content_pending_count,
rss.rss_full_content_queued_count,
rss.rss_full_content_retrying_count,
rss.rss_full_content_complete_count,
rss.rss_full_content_terminal_count
),
(0, 0, 0, 0, 0)
);

utopia_store::rss_full_content::record_baseline(&pool, source_id, 1, &[]).await?;
let entries = vec![
entry("queued"),
entry("hydrating"),
entry("retrying"),
entry("complete"),
entry("superseded"),
entry("deleted"),
{
let mut value = entry("terminal");
value.has_usable_source = false;
value
},
];
utopia_store::rss_full_content::discover(&pool, source_id, 1, &entries).await?;
assert_eq!(
utopia_store::rss_full_content::claim_pending_and_enqueue(&pool, source_id, 1, 6, 5)
.await?,
6
);
utopia_store::rss_full_content::discover(&pool, source_id, 1, &[entry("pending")]).await?;
let hydrating_job = job_id_for(&pool, source_id, "hydrating").await?;
sqlx::query("UPDATE jobs SET status = 'running' WHERE id = $1")
.bind(hydrating_job)
.execute(&pool)
.await?;
let retrying_job = job_id_for(&pool, source_id, "retrying").await?;
sqlx::query("UPDATE jobs SET status = 'queued', attempts = 1 WHERE id = $1")
.bind(retrying_job)
.execute(&pool)
.await?;
let superseded_job = job_id_for(&pool, source_id, "superseded").await?;
sqlx::query("UPDATE jobs SET status = 'done' WHERE id = $1")
.bind(superseded_job)
.execute(&pool)
.await?;

utopia_store::documents::create_with_version_and_processing(
&pool,
kb_id,
"complete.md",
"text/markdown",
10,
"aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa",
Some(source_id),
None,
Some("complete"),
)
.await?;
let deleted = utopia_store::documents::create_with_version_and_processing(
&pool,
kb_id,
"deleted.md",
"text/markdown",
10,
"bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb",
Some(source_id),
None,
Some("deleted"),
)
.await?;
utopia_store::documents::delete(&pool, kb_id, deleted.id, None).await?;
utopia_store::documents::purge(&pool, kb_id, deleted.id).await?;

let listed = utopia_store::sources::list(&pool, kb_id).await?;
let rss = listed
.iter()
.find(|source| source.id == source_id)
.ok_or_else(|| anyhow::anyhow!("RSS source is missing after discovery"))?;
assert_eq!(rss.rss_full_content_state.as_deref(), Some("active"));
assert_eq!(rss.rss_full_content_pending_count, 1);
assert_eq!(
rss.rss_full_content_queued_count, 2,
"queued and hydrating share one bucket"
);
assert_eq!(rss.rss_full_content_retrying_count, 1);
assert_eq!(rss.rss_full_content_complete_count, 1);
assert_eq!(
rss.rss_full_content_terminal_count, 3,
"terminal, deleted and superseded share one bucket"
);
assert_eq!(
rss.doc_count, 1,
"purged documents stay out of the live count"
);
assert_eq!(rss.missing_count, 0);

let mut tx = pool.begin().await?;
utopia_store::rss_full_content::enable_source(&mut tx, source_id).await?;
tx.commit().await?;
utopia_store::rss_full_content::record_baseline(&pool, source_id, 2, &[]).await?;
utopia_store::rss_full_content::discover(&pool, source_id, 2, &[entry("new-generation")])
.await?;
let listed = utopia_store::sources::list(&pool, kb_id).await?;
let rss = listed
.iter()
.find(|source| source.id == source_id)
.ok_or_else(|| anyhow::anyhow!("RSS source is missing after generation change"))?;
assert_eq!(rss.rss_full_content_state.as_deref(), Some("active"));
assert_eq!(
rss.rss_full_content_pending_count, 1,
"old-generation observations are excluded"
);
assert_eq!(rss.rss_full_content_terminal_count, 0);

sqlx::query(
"UPDATE sources
SET config = jsonb_set(config, '{content_mode}', '\"feed\"')
WHERE id = $1",
)
.bind(source_id)
.execute(&pool)
.await?;
let listed = utopia_store::sources::list(&pool, kb_id).await?;
let rss = listed
.iter()
.find(|source| source.id == source_id)
.ok_or_else(|| anyhow::anyhow!("RSS source is missing after disabling"))?;
assert_eq!(rss.rss_full_content_state.as_deref(), Some("disabled"));
assert_eq!(
rss.rss_full_content_terminal_count, 1,
"disabled current rows are superseded"
);

cleanup(&pool, org_id).await?;
Ok(())
}

/// The lateral aggregate must count only the listed source's own observations.
///
/// One RSS source per base is exactly the fixture that cannot catch the outer alias
/// binding back to the inner `sources s`: the source filter turns trivially true and
/// every source's entries are counted, yet a single source still reports its own
/// number. So: two RSS sources in one base and a third in another, each with a
/// different count, so a leak between any two of them shows up as a wrong number.
#[tokio::test]
async fn source_list_counts_only_the_listed_source() -> anyhow::Result<()> {
let Some(url) = utopia_store::test_db::url() else {
return Ok(());
};
let pool = PgPool::connect(&url).await?;
let (org_id, kb_id, first_id, workspace_id) = seed_source(&pool).await?;
let rss_config = serde_json::json!({
"feed_url": "https://example.com/feed",
"content_mode": "full_new_items"
});
let second_id = insert_source(&pool, kb_id, "rss", "rss-second", rss_config.clone()).await?;
let other_kb_id = Uuid::now_v7();
sqlx::query("INSERT INTO knowledge_bases (id, workspace_id, name) VALUES ($1, $2, $3)")
.bind(other_kb_id)
.bind(workspace_id)
.bind("rss-full-content-other-base")
.execute(&pool)
.await?;
let other_id = insert_source(&pool, other_kb_id, "rss", "rss-other-base", rss_config).await?;

for (source_id, keys) in [
(first_id, vec!["first-1"]),
(second_id, vec!["second-1", "second-2"]),
(other_id, vec!["other-1", "other-2", "other-3"]),
] {
let mut tx = pool.begin().await?;
utopia_store::rss_full_content::initialize_source(&mut tx, source_id).await?;
tx.commit().await?;
utopia_store::rss_full_content::record_baseline(&pool, source_id, 1, &[]).await?;
let entries: Vec<_> = keys.into_iter().map(entry).collect();
utopia_store::rss_full_content::discover(&pool, source_id, 1, &entries).await?;
}

fn pending_of(listed: &[utopia_core::models::SourceView], id: Uuid) -> anyhow::Result<i64> {
listed
.iter()
.find(|source| source.id == id)
.map(|source| source.rss_full_content_pending_count)
.ok_or_else(|| anyhow::anyhow!("source {id} is missing from its base's list"))
}

let listed = utopia_store::sources::list(&pool, kb_id).await?;
assert_eq!(listed.len(), 2, "the first base lists its own two sources");
assert_eq!(
pending_of(&listed, first_id)?,
1,
"the first source counts only its own observation"
);
assert_eq!(
pending_of(&listed, second_id)?,
2,
"a second source in the same base does not inherit the first's rows"
);
for source in &listed {
assert_eq!(source.rss_full_content_complete_count, 0);
assert_eq!(source.rss_full_content_terminal_count, 0);
}

let listed = utopia_store::sources::list(&pool, other_kb_id).await?;
assert_eq!(listed.len(), 1, "the other base lists only its own source");
assert_eq!(
pending_of(&listed, other_id)?,
3,
"a source in another base sees none of the first base's rows"
);

cleanup(&pool, org_id).await?;
Ok(())
}

fn entry(key: impl Into<String>) -> utopia_store::rss_full_content::NewEntry {
utopia_store::rss_full_content::NewEntry {
external_key: key.into(),
Expand Down
Loading
Loading