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
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
//! The Arrow surface of the `feeds` and `items` tables: the two schemas, the
//! builders that turn parsed rows into `RecordBatch`es, and the one column
//! rewrite the engine performs after a batch is built.
//!
//! Column order *is* the batch shape — the exec layer projects by index — so
//! these schemas are the single source of truth for both, and the tests pin
//! every name, type, and nullability against the plan's tables. No projection
//! or filtering happens here; that is exec's job.

use std::collections::HashMap;
use std::sync::{Arc, LazyLock};

use arrow::array::{
    ArrayRef, ListBuilder, RecordBatch, StringArray, StringBuilder, TimestampMillisecondBuilder,
    UInt16Builder, UInt32Array, UInt64Builder,
};
use arrow::datatypes::{DataType, Field, Schema, SchemaRef, TimeUnit};

use super::RSS_SURFACE_VERSION;
use super::cache::FeedStatus;
use super::parse::ItemRow;

/// Position of `items.window_status`. The engine builds a window's batch once
/// and re-labels it per serve (`fresh` / `revalidated` / `stale-error`), so this
/// index is part of the module's contract — see [`with_window_status`].
pub const WINDOW_STATUS_IDX: usize = 15;

/// Schema-metadata key carrying [`RSS_SURFACE_VERSION`], so a batch that
/// outlives its producer still names the surface generation it was built for.
const SURFACE_VERSION_KEY: &str = "skardi.rss.surface_version";

static ITEMS_SCHEMA: LazyLock<SchemaRef> = LazyLock::new(|| {
    Arc::new(
        Schema::new(vec![
            Field::new("feed", DataType::Utf8, false),
            Field::new("feed_url", DataType::Utf8, false),
            Field::new("guid", DataType::Utf8, false),
            Field::new("title", DataType::Utf8, true),
            Field::new("link", DataType::Utf8, true),
            Field::new("author", DataType::Utf8, true),
            Field::new(
                "published",
                DataType::Timestamp(TimeUnit::Millisecond, Some("UTC".into())),
                true,
            ),
            Field::new(
                "updated",
                DataType::Timestamp(TimeUnit::Millisecond, Some("UTC".into())),
                true,
            ),
            Field::new("content", DataType::Utf8, true),
            Field::new("summary", DataType::Utf8, true),
            Field::new(
                "categories",
                DataType::List(Arc::new(Field::new("item", DataType::Utf8, true))),
                true,
            ),
            Field::new("enclosure_url", DataType::Utf8, true),
            Field::new("enclosure_type", DataType::Utf8, true),
            Field::new("enclosure_length", DataType::UInt64, true),
            Field::new("position", DataType::UInt32, false),
            Field::new("window_status", DataType::Utf8, false),
            Field::new("extensions_json", DataType::Utf8, true),
        ])
        .with_metadata(surface_metadata()),
    )
});

static FEEDS_SCHEMA: LazyLock<SchemaRef> = LazyLock::new(|| {
    Arc::new(
        Schema::new(vec![
            Field::new("name", DataType::Utf8, false),
            Field::new("url", DataType::Utf8, false),
            Field::new("title", DataType::Utf8, true),
            Field::new("site_url", DataType::Utf8, true),
            Field::new("description", DataType::Utf8, true),
            Field::new(
                "last_fetch",
                DataType::Timestamp(TimeUnit::Millisecond, Some("UTC".into())),
                true,
            ),
            Field::new("last_status", DataType::Utf8, false),
            Field::new("http_status", DataType::UInt16, true),
            Field::new("last_error", DataType::Utf8, true),
            Field::new("etag", DataType::Utf8, true),
            Field::new("last_modified", DataType::Utf8, true),
            Field::new("dialect", DataType::Utf8, true),
            Field::new("dialect_declared", DataType::Utf8, true),
            Field::new("conformance_notes", DataType::Utf8, true),
            Field::new("item_count", DataType::UInt64, true),
        ])
        .with_metadata(surface_metadata()),
    )
});

fn surface_metadata() -> HashMap<String, String> {
    HashMap::from([(
        SURFACE_VERSION_KEY.to_string(),
        RSS_SURFACE_VERSION.to_string(),
    )])
}

/// One `feeds` row, flat and owned so this module stays independent of the
/// cache's and engine's own types.
#[derive(Debug, Clone)]
pub struct FeedsRow {
    pub name: String,
    pub url: String,
    pub title: Option<String>,
    pub site_url: Option<String>,
    pub description: Option<String>,
    pub last_fetch_ms: Option<i64>,
    /// Typed rather than `&'static str` for the same reason
    /// [`with_window_status`] takes a [`FeedStatus`]: [`FeedStatus`] owns this
    /// column's domain and is the only place its strings are spelled, so a typo
    /// cannot reach the column from here.
    pub last_status: FeedStatus,
    pub http_status: Option<u16>,
    pub last_error: Option<String>,
    pub etag: Option<String>,
    pub last_modified: Option<String>,
    pub dialect: Option<String>,
    pub dialect_declared: Option<String>,
    pub conformance_notes: Option<String>,
    pub item_count: Option<u64>,
}

/// The 17-column `items` surface.
pub fn items_schema() -> SchemaRef {
    ITEMS_SCHEMA.clone()
}

/// The 15-column `feeds` surface.
pub fn feeds_schema() -> SchemaRef {
    FEEDS_SCHEMA.clone()
}

/// Encode one feed's window as an `items` batch.
///
/// `feed`/`feed_url` are constant for the batch (a partition is one feed), and
/// `position` is the row's index in the window — the feed's own ordering, which
/// is the only ordering a feed guarantees. `window_status` is filled from
/// [`FeedStatus::Fresh`] (a batch is only ever built from a document that just
/// parsed); serving a cached window re-labels it via [`with_window_status`].
pub fn build_items_batch(feed: &str, feed_url: &str, items: &[ItemRow]) -> RecordBatch {
    let rows = items.len();
    let mut guid = StringBuilder::new();
    let mut title = StringBuilder::new();
    let mut link = StringBuilder::new();
    let mut author = StringBuilder::new();
    let mut published = TimestampMillisecondBuilder::new();
    let mut updated = TimestampMillisecondBuilder::new();
    let mut content = StringBuilder::new();
    let mut summary = StringBuilder::new();
    let mut categories = ListBuilder::new(StringBuilder::new());
    let mut enclosure_url = StringBuilder::new();
    let mut enclosure_type = StringBuilder::new();
    let mut enclosure_length = UInt64Builder::new();
    let mut extensions_json = StringBuilder::new();

    for item in items {
        guid.append_value(&item.guid);
        title.append_option(item.title.as_deref());
        link.append_option(item.link.as_deref());
        author.append_option(item.author.as_deref());
        published.append_option(item.published_ms);
        updated.append_option(item.updated_ms);
        content.append_option(item.content.as_deref());
        summary.append_option(item.summary.as_deref());
        // No categories is NULL, not `[]` — absence then reads the same as every
        // other absent field on the row, which is what the column's declared
        // nullability is for.
        if item.categories.is_empty() {
            categories.append_null();
        } else {
            for category in &item.categories {
                categories.values().append_value(category);
            }
            categories.append(true);
        }
        enclosure_url.append_option(item.enclosure_url.as_deref());
        enclosure_type.append_option(item.enclosure_type.as_deref());
        enclosure_length.append_option(item.enclosure_length);
        extensions_json.append_option(item.extensions_json.as_deref());
    }

    let columns: Vec<ArrayRef> = vec![
        Arc::new(StringArray::from(vec![feed; rows])),
        Arc::new(StringArray::from(vec![feed_url; rows])),
        Arc::new(guid.finish()),
        Arc::new(title.finish()),
        Arc::new(link.finish()),
        Arc::new(author.finish()),
        Arc::new(published.finish().with_timezone("UTC")),
        Arc::new(updated.finish().with_timezone("UTC")),
        Arc::new(content.finish()),
        Arc::new(summary.finish()),
        Arc::new(categories.finish()),
        Arc::new(enclosure_url.finish()),
        Arc::new(enclosure_type.finish()),
        Arc::new(enclosure_length.finish()),
        Arc::new(UInt32Array::from_iter_values(0..rows as u32)),
        Arc::new(StringArray::from(vec![fresh_label(); rows])),
        Arc::new(extensions_json.finish()),
    ];

    RecordBatch::try_new(items_schema(), columns)
        .expect("built columns match the items schema declared above")
}

/// The `window_status` label [`build_items_batch`] stamps.
///
/// [`FeedStatus::Fresh`] is one of the three statuses
/// [`FeedStatus::window_status_str`] answers `Some` for, so the `expect` is
/// discharged by that function's own match arms rather than by an assumption
/// about a caller.
fn fresh_label() -> &'static str {
    FeedStatus::Fresh
        .window_status_str()
        .expect("FeedStatus::Fresh has a window_status label")
}

/// Re-label a window's freshness without rebuilding it, or `None` when
/// `status` is one of the two that serve no window at all.
///
/// Takes a [`FeedStatus`] rather than a `&str` so a misspelled label
/// (`stale_error` for `stale-error`) cannot compile: [`FeedStatus`] owns this
/// domain and is the only place the strings are spelled. `None` covers
/// [`FeedStatus::Never`] and [`FeedStatus::Error`], which have no label
/// because a feed in either state has nothing to serve — that makes
/// "relabelled with a status that isn't a window status" unrepresentable in
/// the return value instead of something this function has to invent a label
/// for.
///
/// Every other column and the schema are shared by `Arc`, so a `Some` costs
/// one `status`-wide string array — cheap enough to run on every serve of a
/// cached window.
pub fn with_window_status(batch: &RecordBatch, status: FeedStatus) -> Option<RecordBatch> {
    let label = status.window_status_str()?;
    let mut columns = batch.columns().to_vec();
    columns[WINDOW_STATUS_IDX] = Arc::new(StringArray::from(vec![label; batch.num_rows()]));
    Some(
        RecordBatch::try_new(batch.schema(), columns)
            .expect("only window_status changed, and it stayed Utf8 with the same length"),
    )
}

/// Encode feed health observations as a `feeds` batch.
pub fn build_feeds_batch(rows: &[FeedsRow]) -> RecordBatch {
    let mut name = StringBuilder::new();
    let mut url = StringBuilder::new();
    let mut title = StringBuilder::new();
    let mut site_url = StringBuilder::new();
    let mut description = StringBuilder::new();
    let mut last_fetch = TimestampMillisecondBuilder::new();
    let mut last_status = StringBuilder::new();
    let mut http_status = UInt16Builder::new();
    let mut last_error = StringBuilder::new();
    let mut etag = StringBuilder::new();
    let mut last_modified = StringBuilder::new();
    let mut dialect = StringBuilder::new();
    let mut dialect_declared = StringBuilder::new();
    let mut conformance_notes = StringBuilder::new();
    let mut item_count = UInt64Builder::new();

    for row in rows {
        name.append_value(&row.name);
        url.append_value(&row.url);
        title.append_option(row.title.as_deref());
        site_url.append_option(row.site_url.as_deref());
        description.append_option(row.description.as_deref());
        last_fetch.append_option(row.last_fetch_ms);
        last_status.append_value(row.last_status.as_str());
        http_status.append_option(row.http_status);
        last_error.append_option(row.last_error.as_deref());
        etag.append_option(row.etag.as_deref());
        last_modified.append_option(row.last_modified.as_deref());
        dialect.append_option(row.dialect.as_deref());
        dialect_declared.append_option(row.dialect_declared.as_deref());
        conformance_notes.append_option(row.conformance_notes.as_deref());
        item_count.append_option(row.item_count);
    }

    let columns: Vec<ArrayRef> = vec![
        Arc::new(name.finish()),
        Arc::new(url.finish()),
        Arc::new(title.finish()),
        Arc::new(site_url.finish()),
        Arc::new(description.finish()),
        Arc::new(last_fetch.finish().with_timezone("UTC")),
        Arc::new(last_status.finish()),
        Arc::new(http_status.finish()),
        Arc::new(last_error.finish()),
        Arc::new(etag.finish()),
        Arc::new(last_modified.finish()),
        Arc::new(dialect.finish()),
        Arc::new(dialect_declared.finish()),
        Arc::new(conformance_notes.finish()),
        Arc::new(item_count.finish()),
    ];

    RecordBatch::try_new(feeds_schema(), columns)
        .expect("built columns match the feeds schema declared above")
}

#[cfg(test)]
mod tests {
    use std::sync::Arc;

    use arrow::array::{
        Array, ListArray, StringArray, TimestampMillisecondArray, UInt16Array, UInt32Array,
        UInt64Array,
    };
    use arrow::datatypes::{DataType, Field, TimeUnit};

    use super::*;

    const FEED: &str = "news";
    const FEED_URL: &str = "https://example.com/feed.xml";

    /// The bundled semantics overlay describes exactly these two schemas'
    /// columns — no more, no fewer — and parses as a real `SemanticsFile`.
    ///
    /// Nothing else tests `docs/rss/semantics.yaml` at all, and it is the file an
    /// agent reads when planning a query, so a drift between it and the surface
    /// is a drift in what agents are told. Both directions matter: a column the
    /// overlay omits is a column with no provenance, and a column the overlay
    /// invents is a lookup that silently never matches. Two live examples this
    /// would have caught — a stale pruning rule the overlay kept after the
    /// implementation widened it, and a NULL-ability claim contradicted by a
    /// corpus fixture — were both found by review instead.
    #[test]
    fn the_bundled_semantics_overlay_matches_the_two_schemas() {
        use crate::semantics::{SEMANTICS_KIND, SemanticsFile};

        // `CARGO_MANIFEST_DIR` is `crates/skardi`; the overlay ships at the
        // repo root under `docs/`.
        let path =
            std::path::Path::new(env!("CARGO_MANIFEST_DIR")).join("../../docs/rss/semantics.yaml");
        let text = std::fs::read_to_string(&path)
            .unwrap_or_else(|e| panic!("read the bundled overlay at {}: {e}", path.display()));
        let file: SemanticsFile = serde_yaml::from_str(&text)
            .unwrap_or_else(|e| panic!("the bundled overlay must parse as a SemanticsFile: {e}"));

        assert_eq!(
            file.kind.as_deref(),
            Some(SEMANTICS_KIND),
            "the overlay must declare `kind: semantics` or the loader skips it"
        );

        for (table, schema) in [("feeds", feeds_schema()), ("items", items_schema())] {
            // The overlay's `name:` keys are fully-qualified and assume a source
            // called `news`, which the file's own header documents.
            let qualified = format!("news.main.{table}");
            let source = file
                .spec
                .sources
                .iter()
                .find(|s| s.name == qualified)
                .unwrap_or_else(|| panic!("the overlay must carry an entry for {qualified}"));
            assert!(
                source.description.is_some(),
                "{qualified} needs a table-level description"
            );

            let described: Vec<&str> = source.columns.iter().map(|c| c.name.as_str()).collect();
            let declared: Vec<&str> = schema.fields().iter().map(|f| f.name().as_str()).collect();

            for name in &declared {
                assert!(
                    described.contains(name),
                    "{qualified}.{name} is in the schema but has no description in the overlay"
                );
            }
            for name in &described {
                assert!(
                    declared.contains(name),
                    "the overlay describes {qualified}.{name}, which is not a column of that \
                     table — an agent reading it would plan a query against nothing"
                );
            }
            // Catches a duplicated column entry, which the two containment
            // checks above would each pass.
            assert_eq!(
                described.len(),
                declared.len(),
                "{qualified}: the overlay lists {} columns for a {}-column table",
                described.len(),
                declared.len()
            );
        }
    }

    /// The expected types are restated here rather than shared with the
    /// implementation on purpose: a helper used by both would let a wrong type
    /// satisfy its own assertion.
    fn ts_utc() -> DataType {
        DataType::Timestamp(TimeUnit::Millisecond, Some("UTC".into()))
    }

    fn utf8_list() -> DataType {
        DataType::List(Arc::new(Field::new("item", DataType::Utf8, true)))
    }

    fn assert_surface(schema: &SchemaRef, expected: &[(&str, DataType, bool)], count: usize) {
        assert_eq!(expected.len(), count, "the test's own expectation list");
        assert_eq!(schema.fields().len(), count, "column count");
        for (idx, (name, data_type, nullable)) in expected.iter().enumerate() {
            let field = schema.field(idx);
            assert_eq!(field.name(), name, "column {idx} name");
            assert_eq!(field.data_type(), data_type, "column {idx} `{name}` type");
            assert_eq!(
                field.is_nullable(),
                *nullable,
                "column {idx} `{name}` nullability"
            );
        }
        // Spelled out rather than taken from the constant: the key is the
        // wire-visible name, so a rename must fail here.
        assert_eq!(
            schema
                .metadata()
                .get("skardi.rss.surface_version")
                .map(String::as_str),
            Some("1"),
            "surface-version metadata"
        );
    }

    fn strings(batch: &RecordBatch, idx: usize) -> Vec<Option<&str>> {
        batch
            .column(idx)
            .as_any()
            .downcast_ref::<StringArray>()
            .unwrap_or_else(|| panic!("column {idx} is not Utf8"))
            .iter()
            .collect()
    }

    fn two_items() -> Vec<ItemRow> {
        vec![
            ItemRow {
                guid: "urn:uuid:1".into(),
                title: Some("First post".into()),
                link: Some("https://example.com/1".into()),
                author: Some("Ada".into()),
                published_ms: Some(1_700_000_000_000),
                updated_ms: Some(1_700_000_060_000),
                content: Some("# body".into()),
                summary: Some("teaser".into()),
                categories: vec!["rust".into(), "arrow".into()],
                enclosure_url: Some("https://example.com/1.mp3".into()),
                enclosure_type: Some("audio/mpeg".into()),
                enclosure_length: Some(1024),
                extensions_json: Some(r#"{"rights":"CC0"}"#.into()),
            },
            // Every optional absent: the null-propagation half of the round trip.
            ItemRow {
                guid: "urn:uuid:2".into(),
                title: None,
                link: None,
                author: None,
                published_ms: None,
                updated_ms: None,
                content: None,
                summary: None,
                categories: Vec::new(),
                enclosure_url: None,
                enclosure_type: None,
                enclosure_length: None,
                extensions_json: None,
            },
        ]
    }

    fn never_row() -> FeedsRow {
        FeedsRow {
            name: FEED.into(),
            url: FEED_URL.into(),
            title: None,
            site_url: None,
            description: None,
            last_fetch_ms: None,
            last_status: FeedStatus::Never,
            http_status: None,
            last_error: None,
            etag: None,
            last_modified: None,
            dialect: None,
            dialect_declared: None,
            conformance_notes: None,
            item_count: None,
        }
    }

    #[test]
    fn items_schema_matches_spec() {
        let schema = items_schema();
        let expected: Vec<(&str, DataType, bool)> = vec![
            ("feed", DataType::Utf8, false),
            ("feed_url", DataType::Utf8, false),
            ("guid", DataType::Utf8, false),
            ("title", DataType::Utf8, true),
            ("link", DataType::Utf8, true),
            ("author", DataType::Utf8, true),
            ("published", ts_utc(), true),
            ("updated", ts_utc(), true),
            ("content", DataType::Utf8, true),
            ("summary", DataType::Utf8, true),
            ("categories", utf8_list(), true),
            ("enclosure_url", DataType::Utf8, true),
            ("enclosure_type", DataType::Utf8, true),
            ("enclosure_length", DataType::UInt64, true),
            ("position", DataType::UInt32, false),
            ("window_status", DataType::Utf8, false),
            ("extensions_json", DataType::Utf8, true),
        ];
        assert_surface(&schema, &expected, 17);
        // The engine rewrites this column by index; a reorder must fail here
        // rather than silently mislabel every row's freshness.
        assert_eq!(schema.field(WINDOW_STATUS_IDX).name(), "window_status");
    }

    #[test]
    fn feeds_schema_matches_spec() {
        let expected: Vec<(&str, DataType, bool)> = vec![
            ("name", DataType::Utf8, false),
            ("url", DataType::Utf8, false),
            ("title", DataType::Utf8, true),
            ("site_url", DataType::Utf8, true),
            ("description", DataType::Utf8, true),
            ("last_fetch", ts_utc(), true),
            ("last_status", DataType::Utf8, false),
            ("http_status", DataType::UInt16, true),
            ("last_error", DataType::Utf8, true),
            ("etag", DataType::Utf8, true),
            ("last_modified", DataType::Utf8, true),
            ("dialect", DataType::Utf8, true),
            ("dialect_declared", DataType::Utf8, true),
            ("conformance_notes", DataType::Utf8, true),
            ("item_count", DataType::UInt64, true),
        ];
        assert_surface(&feeds_schema(), &expected, 15);
    }

    #[test]
    fn items_batch_round_trips_rows() {
        let batch = build_items_batch(FEED, FEED_URL, &two_items());
        assert_eq!(batch.schema(), items_schema());
        assert_eq!(batch.num_rows(), 2);

        assert_eq!(strings(&batch, 0), vec![Some(FEED), Some(FEED)]);
        assert_eq!(strings(&batch, 1), vec![Some(FEED_URL), Some(FEED_URL)]);
        assert_eq!(
            strings(&batch, 2),
            vec![Some("urn:uuid:1"), Some("urn:uuid:2")]
        );
        assert_eq!(strings(&batch, 3), vec![Some("First post"), None]);
        assert_eq!(
            strings(&batch, 4),
            vec![Some("https://example.com/1"), None]
        );
        assert_eq!(strings(&batch, 5), vec![Some("Ada"), None]);

        let published = batch
            .column(6)
            .as_any()
            .downcast_ref::<TimestampMillisecondArray>()
            .expect("published is Timestamp(ms)");
        assert_eq!(
            published.iter().collect::<Vec<_>>(),
            vec![Some(1_700_000_000_000), None]
        );
        assert_eq!(
            published.data_type(),
            &ts_utc(),
            "the array must carry the column's timezone, not a bare Timestamp"
        );
        let updated = batch
            .column(7)
            .as_any()
            .downcast_ref::<TimestampMillisecondArray>()
            .expect("updated is Timestamp(ms)");
        assert_eq!(
            updated.iter().collect::<Vec<_>>(),
            vec![Some(1_700_000_060_000), None]
        );

        assert_eq!(strings(&batch, 8), vec![Some("# body"), None]);
        assert_eq!(strings(&batch, 9), vec![Some("teaser"), None]);

        let categories = batch
            .column(10)
            .as_any()
            .downcast_ref::<ListArray>()
            .expect("categories is a List");
        let first = categories.value(0);
        assert_eq!(
            first
                .as_any()
                .downcast_ref::<StringArray>()
                .expect("list items are Utf8")
                .iter()
                .collect::<Vec<_>>(),
            vec![Some("rust"), Some("arrow")]
        );
        // NULL rather than `[]`: every other absent field is NULL, and that
        // consistency is what makes the column's declared nullability mean
        // something.
        assert!(
            categories.is_null(1),
            "an item with no categories is NULL, not an empty list"
        );

        assert_eq!(
            strings(&batch, 11),
            vec![Some("https://example.com/1.mp3"), None]
        );
        assert_eq!(strings(&batch, 12), vec![Some("audio/mpeg"), None]);
        let length = batch
            .column(13)
            .as_any()
            .downcast_ref::<UInt64Array>()
            .expect("enclosure_length is UInt64");
        assert_eq!(length.iter().collect::<Vec<_>>(), vec![Some(1024), None]);

        let position = batch
            .column(14)
            .as_any()
            .downcast_ref::<UInt32Array>()
            .expect("position is UInt32");
        assert_eq!(
            position.iter().collect::<Vec<_>>(),
            vec![Some(0), Some(1)],
            "position is the row's index in the window"
        );
        assert_eq!(
            strings(&batch, WINDOW_STATUS_IDX),
            vec![Some("fresh"), Some("fresh")]
        );
        assert_eq!(strings(&batch, 16), vec![Some(r#"{"rights":"CC0"}"#), None]);
    }

    #[test]
    fn with_window_status_swaps_only_column_15() {
        let batch = build_items_batch(FEED, FEED_URL, &two_items());
        let relabelled = with_window_status(&batch, FeedStatus::StaleError)
            .expect("stale-error is a window status");

        for idx in 0..batch.num_columns() {
            if idx == WINDOW_STATUS_IDX {
                continue;
            }
            assert!(
                Arc::ptr_eq(batch.column(idx), relabelled.column(idx)),
                "column {idx} was rebuilt instead of shared"
            );
        }
        assert!(
            Arc::ptr_eq(&batch.schema(), &relabelled.schema()),
            "the schema was rebuilt instead of shared"
        );
        assert_eq!(
            strings(&relabelled, WINDOW_STATUS_IDX),
            vec![Some("stale-error"), Some("stale-error")]
        );
        assert_eq!(
            strings(&batch, WINDOW_STATUS_IDX),
            vec![Some("fresh"), Some("fresh")],
            "the source batch must be untouched — it is the cached window"
        );
    }

    #[test]
    fn feeds_batch_round_trips() {
        let observed = FeedsRow {
            name: "blog".into(),
            url: "https://blog.example.com/atom.xml".into(),
            title: Some("The Blog".into()),
            site_url: Some("https://blog.example.com/".into()),
            description: Some("posts".into()),
            last_fetch_ms: Some(1_700_000_000_000),
            last_status: FeedStatus::Revalidated,
            http_status: Some(304),
            last_error: None,
            etag: Some("\"abc\"".into()),
            last_modified: Some("Wed, 21 Oct 2015 07:28:00 GMT".into()),
            dialect: Some("atom".into()),
            dialect_declared: Some("atom".into()),
            conformance_notes: Some("sanitation: stripped-control-chars".into()),
            item_count: Some(7),
        };
        let batch = build_feeds_batch(&[never_row(), observed]);
        assert_eq!(batch.schema(), feeds_schema());
        assert_eq!(batch.num_rows(), 2);

        // Row 0 — never fetched: only identity and status are known.
        assert_eq!(strings(&batch, 0), vec![Some(FEED), Some("blog")]);
        assert_eq!(strings(&batch, 1)[0], Some(FEED_URL));
        assert_eq!(strings(&batch, 6), vec![Some("never"), Some("revalidated")]);
        for idx in [2, 3, 4, 5, 7, 8, 9, 10, 11, 12, 13, 14] {
            assert!(
                batch.column(idx).is_null(0),
                "column {idx} `{}` must be NULL before the first fetch",
                batch.schema().field(idx).name()
            );
        }

        // Row 1 — everything the engine can observe, round-tripped.
        assert_eq!(strings(&batch, 2)[1], Some("The Blog"));
        assert_eq!(strings(&batch, 3)[1], Some("https://blog.example.com/"));
        assert_eq!(strings(&batch, 4)[1], Some("posts"));
        let last_fetch = batch
            .column(5)
            .as_any()
            .downcast_ref::<TimestampMillisecondArray>()
            .expect("last_fetch is Timestamp(ms)");
        assert_eq!(last_fetch.data_type(), &ts_utc());
        assert_eq!(last_fetch.value(1), 1_700_000_000_000);
        let http_status = batch
            .column(7)
            .as_any()
            .downcast_ref::<UInt16Array>()
            .expect("http_status is UInt16");
        assert_eq!(http_status.value(1), 304);
        assert_eq!(
            strings(&batch, 8)[1],
            None,
            "a revalidated feed carries no error"
        );
        assert_eq!(strings(&batch, 9)[1], Some("\"abc\""));
        assert_eq!(
            strings(&batch, 10)[1],
            Some("Wed, 21 Oct 2015 07:28:00 GMT")
        );
        assert_eq!(strings(&batch, 11)[1], Some("atom"));
        assert_eq!(strings(&batch, 12)[1], Some("atom"));
        assert_eq!(
            strings(&batch, 13)[1],
            Some("sanitation: stripped-control-chars")
        );
        let item_count = batch
            .column(14)
            .as_any()
            .downcast_ref::<UInt64Array>()
            .expect("item_count is UInt64");
        assert_eq!(item_count.value(1), 7);
    }

    #[test]
    fn empty_items_batch_has_zero_rows_17_cols() {
        let batch = build_items_batch(FEED, FEED_URL, &[]);
        assert_eq!(batch.num_rows(), 0);
        assert_eq!(batch.num_columns(), 17);
        assert_eq!(batch.schema(), items_schema());
        // A feed can legitimately serve an empty window; re-labelling one must
        // not be the path that panics.
        let relabelled = with_window_status(&batch, FeedStatus::Revalidated)
            .expect("revalidated is a window status");
        assert_eq!(relabelled.num_rows(), 0);
        assert_eq!(relabelled.num_columns(), 17);
    }

    /// The two statuses that serve zero rows have no `window_status` label, so
    /// relabelling with one answers `None` rather than writing a string
    /// outside the column's domain.
    #[test]
    fn with_window_status_refuses_the_zero_row_statuses() {
        let batch = build_items_batch(FEED, FEED_URL, &two_items());
        assert!(with_window_status(&batch, FeedStatus::Never).is_none());
        assert!(with_window_status(&batch, FeedStatus::Error).is_none());
    }
}