skardi 0.6.0

High performance query engine for both offline compute and online serving
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
//! RSS/Atom subscriptions as a read-only data source (`type: rss`).
//!
//! See `docs/superpowers/specs/2026-07-22-rss-feed-support-design.md`.
#[cfg(feature = "rss")]
pub mod cache;
pub mod config;
#[cfg(feature = "rss")]
pub mod conformance;
// The compatibility corpus: committed feed documents in `fixtures/` plus the
// manifest-driven contract test over them. Test-only and additionally gated
// behind `rss`, like `testutil` below, since everything it drives
// (`parse_feed_document`, the sanitation rungs, the HTML→Markdown conversion)
// is.
#[cfg(feature = "rss")]
pub mod convert;
#[cfg(all(test, feature = "rss"))]
mod corpus;
pub mod error;
// Reads OPML files and pulls in `quick-xml`; gated so the config/error types
// above stay parseable — and `ResolvedSubscription` below stays nameable —
// in builds that omit the `rss` feature.
#[cfg(feature = "rss")]
pub mod opml;
// The fetcher's egress policy seam: an `EgressPolicy` trait plus the OSS
// `AllowAll` default, consulted before reqwest connects (see the module doc
// for why). Not `pub` — it is an internal implementation detail of the fetch
// engine (`fetch` consumes it via `super::egress`), not part of this
// provider's public surface. Gated behind `rss` alongside the rest of the
// fetch/parse engine, even though its own dependencies (reqwest, tokio) are
// already unconditional crate deps.
#[cfg(feature = "rss")]
mod egress;
// The partition-per-feed execution plan: the engine's only consumer, and the
// layer that enforces the scan deadline and the LIMIT launch gate. `pub` so
// the table provider a later task adds — in this module tree but a different
// file — can construct it.
#[cfg(feature = "rss")]
pub mod exec;
// The freshness state machine that composes every module above: it decides
// per feed whether a scan serves a cached window, revalidates it, refetches
// it, or degrades to stale rows, and it is the sole production consumer of
// `egress`, `fetch`, and `cache`.
#[cfg(feature = "rss")]
pub mod engine;
// The bounded HTTP fetcher (conditional GET, retries, egress enforcement)
// built on top of `egress`. Not `pub` for the same reason `egress` isn't:
// it is an implementation detail of the engine a later task builds on top,
// not part of this provider's public surface.
#[cfg(feature = "rss")]
mod fetch;
// Hand-rolled mock feed server the fetcher's tests drive. Test-only (never
// compiled into a release build) and additionally gated behind `rss` since
// its only consumer, `fetch`'s test module, is.
#[cfg(all(test, feature = "rss"))]
pub(crate) mod testutil;
// The parsing chain: byte-level sanitation rungs, the feed-rs parse driver
// that applies them, and the fixed Arrow schemas the providers serve. These
// were built on a parallel branch (Tasks 5-9) alongside the fetch chain
// above, which is why they land as one merge rather than task by task.
#[cfg(feature = "rss")]
pub mod parse;
#[cfg(feature = "rss")]
pub mod sanitize;
#[cfg(feature = "rss")]
pub mod schema;
// The two `TableProvider`s (`feeds` and `items`): they classify filters for
// pushdown, prune the subscription list to the feeds a scan must visit, and
// construct the `exec` plan above. `pub` so the catalog registration a later
// task adds can name them.
#[cfg(feature = "rss")]
pub mod table;

// The acceptance-criteria crosswalk: full SQL against a registered catalog
// whose feeds live on `testutil::MockFeedServer`. In-crate rather than in
// `crates/skardi/tests/` because that mock server is what binds it here:
// `testutil` is test-only `pub(crate)` and unreachable from an external
// test crate. (`register_rss_tables_with_policy` itself is `pub` — the
// egress seam an embedder calls — so it is not the constraint.)
#[cfg(all(test, feature = "rss"))]
mod integration_tests;

// One layer above `integration_tests`: the downstream composition the design
// leaves to user-space SQL — a federated join, and the two-`INSERT` archive
// that gives feed entries a history the live window does not. In-crate for the
// same reason, and additionally gated behind `chunking` because the archive's
// second statement is a `chunk()` call.
#[cfg(all(test, feature = "rss", feature = "chunking"))]
mod composition_tests;

pub use config::{FeedSubscription, RssConfig};
pub use error::RssError;
#[cfg(feature = "rss")]
pub use opml::resolve_subscriptions;
// The egress-policy seam, re-exported so an embedder (Skardi Cloud, or an
// operator running Skardi as a library) can name the trait to implement a
// destination filter and inject it via `register_rss_tables_with_policy`.
// OSS ships only `AllowAll`; supplying anything stricter is the caller's.
#[cfg(feature = "rss")]
pub use egress::{AllowAll, EgressDenied, EgressPolicy, EgressReason};

#[cfg(feature = "rss")]
use std::sync::Arc;

#[cfg(feature = "rss")]
use anyhow::Result;
#[cfg(feature = "rss")]
use datafusion::catalog::{
    CatalogProvider, MemoryCatalogProvider, MemorySchemaProvider, SchemaProvider,
};
#[cfg(feature = "rss")]
use datafusion::prelude::SessionContext;

#[cfg(feature = "rss")]
use crate::sources::hierarchy::HierarchyLevel;
#[cfg(feature = "rss")]
use engine::RssEngine;
#[cfg(feature = "rss")]
use table::RssTableProvider;

/// Integer version of the `feeds`/`items` public surface. Bumped only by
/// breaking changes (column removal/rename/retype, nullability tightening,
/// enum-domain repurposing, identity/window semantics changes).
pub const RSS_SURFACE_VERSION: u32 = 1;

/// One subscription, fully resolved from either of [`RssConfig`]'s two
/// mutually exclusive input forms — an inline `feeds:` entry or an
/// `<outline>` pulled from an `opml:` file.
///
/// This is the convergence point the rest of the provider is built on:
/// every later stage — the fetcher, the TTL cache, the freshness state
/// machine, the partition-per-feed execution plan — consumes only a
/// `Vec<ResolvedSubscription>` and never looks at `RssConfig`'s input shape
/// again. It is a plain data struct with no parsing logic of its own, so
/// unlike the `opml` module (which requires the `rss` feature for
/// `quick-xml`, so a doc link to it would dangle in featureless builds) it
/// stays nameable in featureless builds — the server (or any embedder) can
/// hold it in a typed field regardless of which features a given build
/// enables.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ResolvedSubscription {
    /// Effective subscription name: an explicit `name`/`text`/`title`, or
    /// the feed's URL stripped of credentials, query, and fragment when
    /// none was given — the name is public surface (the `feed` column,
    /// log fields), and those are the URL parts that can carry a private
    /// token. Unique across the whole resolved list.
    pub name: String,
    /// Feed URL; already checked to be `http://` or `https://`.
    pub url: String,
}

/// The one schema an `rss` catalog exposes. One source is one catalog holding
/// exactly `<name>.main.feeds` and `<name>.main.items`, and `main` mirrors the
/// sqlite provider's convention, where the single schema every table lands in
/// is also spelled `main` (`sources/providers/sqlite/mod.rs:443-444`).
#[cfg(feature = "rss")]
const RSS_SCHEMA: &str = "main";
/// The per-subscription health table: one row per subscription, always.
#[cfg(feature = "rss")]
const FEEDS_TABLE: &str = "feeds";
/// The feed-entry table: one window per subscription.
#[cfg(feature = "rss")]
const ITEMS_TABLE: &str = "items";

/// Register one `type: rss` data source as the catalog `name`, exposing
/// `<name>.main.feeds` and `<name>.main.items`.
///
/// # Registration performs no network I/O
///
/// The only I/O here is reading the `opml:` file, when one is configured, via
/// [`resolve_subscriptions`]. Nothing probes a feed: the engine and the two
/// table providers are built from the resolved subscription list alone, and
/// every HTTP request happens later, inside an `items` scan. So startup cost
/// is proportional to the size of the subscription list rather than to the
/// availability of the hosts on it — a source with fifty feeds does not wait
/// on fifty upstreams to become queryable, and an unreachable host surfaces
/// as that subscription's `feeds.last_status` instead of failing the whole
/// source. `registration_is_zero_network_and_tables_queryable` pins this by
/// asserting a mock server has observed no requests after registration *and*
/// after a `feeds` scan (`feeds` is a pure state read, so it stays at zero),
/// and exactly one after an `items` scan.
///
/// # Surface version
///
/// [`RSS_SURFACE_VERSION`] is logged here. The same constant is stamped into
/// both tables' Arrow schema metadata under `skardi.rss.surface_version`
/// (`schema.rs:96-100`), so a client can read the version off a query result
/// without access to the log.
///
/// # Egress
///
/// The engine's fetcher is built with no injected policy (`None`) — the OSS
/// default: no destination filtering (the fetcher falls back to its internal
/// `AllowAll`), with system proxy variables honored. `FeedFetcher::new`
/// disables proxies exactly when a policy is injected — the switch is the
/// fact of injection — so the `None` must survive to it rather than being
/// materialized as an `AllowAll` here (names unlinked: the `fetch` module is
/// private, so a doc link to it would not resolve). A caller may inject an
/// `EgressPolicy` through [`register_rss_tables_with_policy`] (an operator's
/// own, or Skardi Cloud's at registration) to restrict egress; see
/// `docs/superpowers/specs/2026-08-03-rss-cloud-egress-design.md`.
#[cfg(feature = "rss")]
pub async fn register_rss_tables(
    session_ctx: &mut SessionContext,
    name: &str,
    config: Option<&RssConfig>,
    read_write: bool,
    hierarchy_level: HierarchyLevel,
) -> Result<()> {
    register_with_policy(session_ctx, name, config, read_write, hierarchy_level, None).await
}

/// [`register_rss_tables`] with a caller-supplied egress policy — the public
/// seam for injecting destination filtering (Skardi Cloud's reserved-range
/// policy, or an operator's own) in place of the `AllowAll` default. OSS ships
/// no non-`AllowAll` policy itself; this is the entry point an embedder calls
/// to supply one. In-crate tests use it the same way to prove a refusal
/// reaches `feeds.last_error`.
#[cfg(feature = "rss")]
pub async fn register_rss_tables_with_policy(
    session_ctx: &mut SessionContext,
    name: &str,
    config: Option<&RssConfig>,
    read_write: bool,
    hierarchy_level: HierarchyLevel,
    policy: Arc<dyn EgressPolicy>,
) -> Result<()> {
    register_with_policy(
        session_ctx,
        name,
        config,
        read_write,
        hierarchy_level,
        Some(policy),
    )
    .await
}

/// The shared body of [`register_rss_tables`] and its test seam. `policy`
/// carries the fact of injection through to `FeedFetcher::new`, which is why
/// it is an `Option` and not a defaulted `AllowAll` — see the `# Egress`
/// section on [`register_rss_tables`].
#[cfg(feature = "rss")]
async fn register_with_policy(
    session_ctx: &mut SessionContext,
    name: &str,
    config: Option<&RssConfig>,
    read_write: bool,
    hierarchy_level: HierarchyLevel,
    policy: Option<Arc<dyn EgressPolicy>>,
) -> Result<()> {
    // All invariant checks live here so every entry point — the server's
    // registration arm and the public `register_rss_tables_with_policy`
    // embedder seam — gets identical behavior; a caller may add an earlier
    // typed error, but this is the single enforcement point — the same
    // arrangement `register_open_connector_tables` uses
    // (`sources/providers/open_connector/mod.rs:144-168`).
    if hierarchy_level != HierarchyLevel::Catalog {
        return Err(RssError::CatalogHierarchyRequired {
            name: name.to_string(),
        }
        .into());
    }
    if read_write {
        return Err(RssError::ReadWriteNotSupported {
            name: name.to_string(),
        }
        .into());
    }
    let config = config.ok_or_else(|| RssError::MissingConfig {
        name: name.to_string(),
    })?;
    // `InvalidConfig` is deliberately nameless (see its doc in `error.rs`);
    // the server names the source when it wraps registration errors in
    // `ConfigError::DataSourceRegistrationFailed`, so nothing is re-stamped
    // here — the same bare call `register_open_connector_tables` makes.
    config.validate()?;

    // The one I/O step: an `opml:` path is read here, not by `validate()`.
    let subscriptions = resolve_subscriptions(name, config)?;
    let engine = Arc::new(RssEngine::new(
        name.to_string(),
        subscriptions,
        config,
        policy,
    )?);
    let subscription_count = engine.subscriptions().len();

    // Built directly rather than through `hierarchy::build_catalog`: that
    // helper's job is to drive many `TableProvider` constructions
    // concurrently and key them by `(schema, table)` name strings, and here
    // there are exactly two providers, both built synchronously from the
    // engine above. Going through it would mean re-dispatching on the two
    // names this function just wrote, with an unreachable third arm. This is
    // the shape `register_open_connector_tables` ends with
    // (`sources/providers/open_connector/mod.rs:212-274`).
    let schema_provider = Arc::new(MemorySchemaProvider::new());
    schema_provider
        .register_table(
            FEEDS_TABLE.to_string(),
            Arc::new(RssTableProvider::feeds(Arc::clone(&engine))),
        )
        .map_err(|e| {
            anyhow::anyhow!(
                "rss source '{name}': failed to register {RSS_SCHEMA}.{FEEDS_TABLE}: {e}"
            )
        })?;
    schema_provider
        .register_table(
            ITEMS_TABLE.to_string(),
            Arc::new(RssTableProvider::items(engine)),
        )
        .map_err(|e| {
            anyhow::anyhow!(
                "rss source '{name}': failed to register {RSS_SCHEMA}.{ITEMS_TABLE}: {e}"
            )
        })?;

    let catalog = Arc::new(MemoryCatalogProvider::new());
    catalog
        .register_schema(RSS_SCHEMA, schema_provider)
        .map_err(|e| {
            anyhow::anyhow!("rss source '{name}': failed to register schema '{RSS_SCHEMA}': {e}")
        })?;

    // Publishing the catalog is the last step, after every gate above has
    // passed: `register_catalog` inserts into the context's catalog list and
    // returns whatever was registered under `name` before (datafusion 52.5.0
    // `src/execution/context/mod.rs:1716-1726`, delegating to
    // `MemoryCatalogProviderList::register_catalog`, datafusion-catalog
    // 52.5.0 `src/memory/catalog.rs:54-60`), so a failed registration must
    // not have already replaced a working source's catalog.
    session_ctx.register_catalog(name, catalog);

    tracing::info!(
        source = %name,
        subscriptions = subscription_count,
        surface_version = RSS_SURFACE_VERSION,
        "RSS source registered"
    );

    Ok(())
}

#[cfg(all(test, feature = "rss"))]
mod tests {
    use super::*;
    use crate::sources::providers::rss::config::inline_config;
    use crate::sources::providers::rss::schema::{feeds_schema, items_schema};
    use crate::sources::providers::rss::testutil::{
        MockFeedServer, MockResponse, MockResponseExt, RSS2_MINIMAL, feed_urls, str_col,
    };
    use arrow::array::RecordBatch;

    /// A config subscribing to `feeds` (`(name, path)` pairs) on `server`.
    ///
    /// `request_timeout_seconds` and `scan_timeout_seconds` are pulled well
    /// below their spec defaults (10 and 60): every server here answers
    /// immediately, so a test that starts hanging has regressed, and it
    /// should say so in seconds rather than sitting on the scan deadline for
    /// a minute.
    fn config_pointing_at(server: &MockFeedServer, feeds: &[(&str, &str)]) -> RssConfig {
        let mut config = inline_config(
            feed_urls(server, feeds)
                .into_iter()
                .map(|(name, url)| FeedSubscription {
                    url,
                    name: Some(name),
                })
                .collect(),
        );
        config.request_timeout_seconds = 5;
        config.scan_timeout_seconds = 10;
        config
    }

    /// A config that never reaches a network at all — the shape the rejection
    /// tests use, since none of them gets far enough to fetch.
    fn unreachable_config() -> RssConfig {
        inline_config(vec![FeedSubscription {
            url: "https://feed.example/f.xml".to_string(),
            name: Some("a".to_string()),
        }])
    }

    async fn query(ctx: &SessionContext, sql: &str) -> Vec<RecordBatch> {
        ctx.sql(sql)
            .await
            .unwrap_or_else(|e| panic!("plan {sql:?}: {e}"))
            .collect()
            .await
            .unwrap_or_else(|e| panic!("execute {sql:?}: {e}"))
    }

    #[tokio::test]
    async fn registration_is_zero_network_and_tables_queryable() {
        let server = MockFeedServer::start(|_| MockResponse::xml(RSS2_MINIMAL)).await;
        let mut ctx = SessionContext::new();
        let config = config_pointing_at(&server, &[("a", "/f.xml")]);

        // Through the test seam: production injects no policy at all, but
        // the seam exists so a test can supply one.
        register_rss_tables_with_policy(
            &mut ctx,
            "news",
            Some(&config),
            false,
            HierarchyLevel::Catalog,
            Arc::new(AllowAll),
        )
        .await
        .expect("registration succeeds");

        assert_eq!(
            server.requests().len(),
            0,
            "registration performed network I/O"
        );

        // Both tables are addressable at exactly `<name>.main.<table>`.
        let feeds = query(
            &ctx,
            "SELECT name, last_status FROM news.main.feeds ORDER BY name",
        )
        .await;
        assert_eq!(str_col(&feeds[0], "name"), vec!["a"]);
        assert_eq!(str_col(&feeds[0], "last_status"), vec!["never"]);
        assert_eq!(
            server.requests().len(),
            0,
            "feeds scan performed network I/O"
        );

        let items = query(&ctx, "SELECT guid, window_status FROM news.main.items").await;
        assert_eq!(
            items.iter().map(RecordBatch::num_rows).sum::<usize>(),
            1,
            "the one item in RSS2_MINIMAL"
        );
        assert_eq!(
            server.requests().len(),
            1,
            "an items scan is what fetches the feed"
        );
    }

    #[tokio::test]
    async fn non_catalog_hierarchy_is_rejected() {
        let mut ctx = SessionContext::new();
        let config = unreachable_config();
        let err = register_rss_tables(
            &mut ctx,
            "news",
            Some(&config),
            false,
            HierarchyLevel::Table,
        )
        .await
        .expect_err("hierarchy_level: table must be rejected");
        assert!(
            err.to_string()
                .contains("hierarchy_level must be 'catalog'"),
            "{err}"
        );
        assert!(
            ctx.catalog("news").is_none(),
            "a rejected source must leave no catalog behind"
        );
    }

    #[tokio::test]
    async fn read_write_is_rejected() {
        let mut ctx = SessionContext::new();
        let config = unreachable_config();
        let err = register_rss_tables(
            &mut ctx,
            "news",
            Some(&config),
            true,
            HierarchyLevel::Catalog,
        )
        .await
        .expect_err("read_write must be rejected");
        assert!(
            err.to_string().contains("access_mode must be read-only"),
            "{err}"
        );
        assert!(
            ctx.catalog("news").is_none(),
            "a rejected source must leave no catalog behind"
        );
    }

    #[tokio::test]
    async fn missing_config_is_rejected() {
        let mut ctx = SessionContext::new();
        let err = register_rss_tables(&mut ctx, "news", None, false, HierarchyLevel::Catalog)
            .await
            .expect_err("a source with no `rss:` block must be rejected");
        assert!(
            err.to_string()
                .contains("missing required `rss:` configuration block"),
            "{err}"
        );
        assert!(
            ctx.catalog("news").is_none(),
            "a rejected source must leave no catalog behind"
        );
    }

    #[tokio::test]
    async fn an_invalid_config_is_rejected_before_any_catalog_appears() {
        // `validate()` runs on the registration path too, not only in the
        // server's pure config check: registration is the single enforcement
        // point, so a front-end that skipped validation cannot register a
        // source the config rules forbid.
        let mut ctx = SessionContext::new();
        let mut config = unreachable_config();
        config.max_concurrent = 0;
        let err = register_rss_tables(
            &mut ctx,
            "news",
            Some(&config),
            false,
            HierarchyLevel::Catalog,
        )
        .await
        .expect_err("an invalid config must be rejected");
        // The reason itself; the source name is attached one layer up, by
        // the server's `DataSourceRegistrationFailed` wrap (`InvalidConfig`
        // is deliberately nameless — see its doc in `error.rs`).
        assert!(
            err.to_string()
                .contains("max_concurrent must be at least 1"),
            "{err}"
        );
        assert!(
            ctx.catalog("news").is_none(),
            "a rejected source must leave no catalog behind"
        );
    }

    #[test]
    fn schema_metadata_carries_surface_version() {
        // Spelled out rather than read from the constant: the key and the
        // rendered value are both wire-visible, so a rename or a silent bump
        // must fail here. `RSS_SURFACE_VERSION` is logged at registration and
        // stamped into these two schemas from the same constant, so this is
        // the query-side half of what that log line reports.
        for schema in [items_schema(), feeds_schema()] {
            assert_eq!(
                schema
                    .metadata()
                    .get("skardi.rss.surface_version")
                    .map(String::as_str),
                Some("1"),
            );
        }
        assert_eq!(RSS_SURFACE_VERSION, 1);
    }
}