Refuse a job kind that is enqueued with nothing to claim it
The queue coalesces, so a producer with no consumer never fails: it leaves one row per subject for ever. That is how 23,582 Thumbnail jobs accumulated unnoticed (#73), and nothing at runtime would have said so. every_queued_kind_has_a_consumer reads the shipping sources of every crate under core/, ui/, apps/ and platform/ (cfg(test) items dropped) and pairs the JobKind named at each enqueue( call with the kinds named in a fn kinds( body or a claim_next_matching( call. An enqueue that does not spell its kind is refused, since the pairing could not be checked. It guards against passing over nothing: the queue's own files and the scan must have been read. A second test runs the reader over fixed snippets so a parsing bug shows up as a failure. Run against master's scan.rs and walk.rs it names all three orphan enqueues. Refs #73
This commit is contained in:
@@ -0,0 +1,341 @@
|
|||||||
|
// TRACES: FR-PLAT-AND-3
|
||||||
|
//! No job kind is enqueued without something that claims it.
|
||||||
|
//!
|
||||||
|
//! The queue coalesces, so a producer with no consumer does not fail — it
|
||||||
|
//! just leaves a row per subject for ever. That is how the reference catalog
|
||||||
|
//! came to hold 23,582 `Thumbnail` jobs, one per photograph, re-coalesced on
|
||||||
|
//! every scan, with no handler for the kind anywhere in the tree (#73). Nothing
|
||||||
|
//! at runtime notices: the rows are cheap one at a time and invisible in the
|
||||||
|
//! interface. So the pairing is checked here, over the source, instead.
|
||||||
|
//!
|
||||||
|
//! ## What counts
|
||||||
|
//!
|
||||||
|
//! In shipping code under `core/`, `ui/`, `apps/` and `platform/` — every
|
||||||
|
//! `src/` tree, with `#[cfg(test)]` items dropped:
|
||||||
|
//!
|
||||||
|
//! - **Enqueued**: the `JobKind::X` named in the arguments of a call to
|
||||||
|
//! `enqueue(`. An enqueue whose kind is not spelled there — passed in a
|
||||||
|
//! variable — is refused outright, because this scan could not say what it
|
||||||
|
//! queues.
|
||||||
|
//! - **Claimed**: the `JobKind::X` in the body of a `fn kinds(` (what a
|
||||||
|
//! `JobHandler` declares, and all a `Runner` claims), or named in a call to
|
||||||
|
//! `claim_next_matching(`. A call to `claim_next(` claims every kind.
|
||||||
|
//! `JobKind::ALL` in either place means every kind.
|
||||||
|
//!
|
||||||
|
//! Tests and examples are left out on purpose: a unit test of the queue's
|
||||||
|
//! mechanics enqueues and claims whatever it likes, and proves nothing about
|
||||||
|
//! the app.
|
||||||
|
|
||||||
|
use std::collections::BTreeSet;
|
||||||
|
use std::fs;
|
||||||
|
use std::path::{Path, PathBuf};
|
||||||
|
|
||||||
|
/// Every `.rs` file under each `src/` of each crate in `group`.
|
||||||
|
fn crate_sources(group: &Path, out: &mut Vec<PathBuf>) {
|
||||||
|
let Ok(crates) = fs::read_dir(group) else {
|
||||||
|
return;
|
||||||
|
};
|
||||||
|
for krate in crates {
|
||||||
|
let src = krate.expect("read dir entry").path().join("src");
|
||||||
|
if src.is_dir() {
|
||||||
|
rust_files(&src, out);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn rust_files(dir: &Path, out: &mut Vec<PathBuf>) {
|
||||||
|
for entry in fs::read_dir(dir).unwrap_or_else(|e| panic!("cannot read {}: {e}", dir.display()))
|
||||||
|
{
|
||||||
|
let path = entry.expect("read dir entry").path();
|
||||||
|
if path.is_dir() {
|
||||||
|
rust_files(&path, out);
|
||||||
|
} else if path.extension().and_then(|e| e.to_str()) == Some("rs") {
|
||||||
|
out.push(path);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Blank out string literals and line comments, so neither a brace nor a
|
||||||
|
/// `JobKind::` inside prose is read as code.
|
||||||
|
fn strip_literals_and_comments(line: &str) -> String {
|
||||||
|
let mut out = String::with_capacity(line.len());
|
||||||
|
let mut chars = line.chars().peekable();
|
||||||
|
let mut in_string = false;
|
||||||
|
while let Some(c) = chars.next() {
|
||||||
|
if in_string {
|
||||||
|
match c {
|
||||||
|
'\\' => {
|
||||||
|
chars.next();
|
||||||
|
}
|
||||||
|
'"' => in_string = false,
|
||||||
|
_ => {}
|
||||||
|
}
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
match c {
|
||||||
|
'"' => in_string = true,
|
||||||
|
'/' if chars.peek() == Some(&'/') => break,
|
||||||
|
_ => out.push(c),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
out
|
||||||
|
}
|
||||||
|
|
||||||
|
/// The shipping code of a file, comments and strings blanked, with every
|
||||||
|
/// `#[cfg(test)]` item dropped. The attribute must be the whole line, so a
|
||||||
|
/// doc comment mentioning it is not mistaken for one.
|
||||||
|
fn shipping_code(text: &str) -> String {
|
||||||
|
let lines: Vec<String> = text.lines().map(strip_literals_and_comments).collect();
|
||||||
|
let mut out = String::new();
|
||||||
|
let mut i = 0;
|
||||||
|
while i < lines.len() {
|
||||||
|
if lines[i].trim() != "#[cfg(test)]" {
|
||||||
|
out.push_str(&lines[i]);
|
||||||
|
out.push('\n');
|
||||||
|
i += 1;
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
let (mut j, mut depth, mut opened) = (i + 1, 0i32, false);
|
||||||
|
while j < lines.len() {
|
||||||
|
depth += lines[j].matches('{').count() as i32;
|
||||||
|
depth -= lines[j].matches('}').count() as i32;
|
||||||
|
opened |= lines[j].contains('{');
|
||||||
|
if (opened && depth <= 0) || (!opened && lines[j].contains(';')) {
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
j += 1;
|
||||||
|
}
|
||||||
|
i = j + 1;
|
||||||
|
}
|
||||||
|
out
|
||||||
|
}
|
||||||
|
|
||||||
|
/// The text from `start` (just past an opening delimiter) to its matching
|
||||||
|
/// close.
|
||||||
|
fn balanced(code: &str, start: usize, open: char, close: char) -> &str {
|
||||||
|
let mut depth = 1;
|
||||||
|
for (i, c) in code[start..].char_indices() {
|
||||||
|
if c == open {
|
||||||
|
depth += 1;
|
||||||
|
} else if c == close {
|
||||||
|
depth -= 1;
|
||||||
|
if depth == 0 {
|
||||||
|
return &code[start..start + i];
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
&code[start..]
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Each call of `name(` in `code` that is a call rather than the function's
|
||||||
|
/// own definition or a longer name ending in it, as its argument text.
|
||||||
|
fn calls<'a>(code: &'a str, name: &str) -> Vec<&'a str> {
|
||||||
|
let needle = format!("{name}(");
|
||||||
|
let mut found = Vec::new();
|
||||||
|
for (at, _) in code.match_indices(&needle) {
|
||||||
|
let before = &code[..at];
|
||||||
|
let prev = before.chars().next_back();
|
||||||
|
if prev.is_some_and(|c| c.is_alphanumeric() || c == '_') {
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
if before.trim_end().ends_with("fn") {
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
found.push(balanced(code, at + needle.len(), '(', ')'));
|
||||||
|
}
|
||||||
|
found
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Bodies of every `fn kinds(` that has one — a trait declaration ending in
|
||||||
|
/// `;` has none.
|
||||||
|
fn kinds_bodies(code: &str) -> Vec<&str> {
|
||||||
|
let mut found = Vec::new();
|
||||||
|
for (at, _) in code.match_indices("fn kinds(") {
|
||||||
|
let rest = &code[at..];
|
||||||
|
let (Some(brace), semi) = (rest.find('{'), rest.find(';')) else {
|
||||||
|
continue;
|
||||||
|
};
|
||||||
|
if semi.is_some_and(|s| s < brace) {
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
found.push(balanced(code, at + brace + 1, '{', '}'));
|
||||||
|
}
|
||||||
|
found
|
||||||
|
}
|
||||||
|
|
||||||
|
/// The `X` of each `JobKind::X` in `text`.
|
||||||
|
fn kinds_named(text: &str) -> Vec<String> {
|
||||||
|
text.match_indices("JobKind::")
|
||||||
|
.map(|(at, m)| {
|
||||||
|
text[at + m.len()..]
|
||||||
|
.chars()
|
||||||
|
.take_while(|c| c.is_alphanumeric() || *c == '_')
|
||||||
|
.collect()
|
||||||
|
})
|
||||||
|
.collect()
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Default, Debug)]
|
||||||
|
struct Ledger {
|
||||||
|
/// Kind → where it is enqueued.
|
||||||
|
enqueued: Vec<(String, String)>,
|
||||||
|
/// Enqueue calls whose kind could not be read.
|
||||||
|
unreadable: Vec<String>,
|
||||||
|
claimed: BTreeSet<String>,
|
||||||
|
claims_everything: bool,
|
||||||
|
}
|
||||||
|
|
||||||
|
fn read(files: &[(String, String)]) -> Ledger {
|
||||||
|
let mut ledger = Ledger::default();
|
||||||
|
for (name, text) in files {
|
||||||
|
let code = shipping_code(text);
|
||||||
|
|
||||||
|
for args in calls(&code, "enqueue") {
|
||||||
|
let kinds = kinds_named(args);
|
||||||
|
if kinds.is_empty() {
|
||||||
|
ledger
|
||||||
|
.unreadable
|
||||||
|
.push(format!("{name}: enqueue({})", args.trim()));
|
||||||
|
}
|
||||||
|
for k in kinds {
|
||||||
|
ledger.enqueued.push((k, name.clone()));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
let claimed = kinds_bodies(&code)
|
||||||
|
.into_iter()
|
||||||
|
.chain(calls(&code, "claim_next_matching"));
|
||||||
|
for text in claimed {
|
||||||
|
for k in kinds_named(text) {
|
||||||
|
if k == "ALL" {
|
||||||
|
ledger.claims_everything = true;
|
||||||
|
} else {
|
||||||
|
ledger.claimed.insert(k);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if !calls(&code, "claim_next").is_empty() {
|
||||||
|
ledger.claims_everything = true;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
ledger
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn every_kind_enqueued_is_claimed_by_something() {
|
||||||
|
let repo = Path::new(env!("CARGO_MANIFEST_DIR"))
|
||||||
|
.parent()
|
||||||
|
.and_then(Path::parent)
|
||||||
|
.expect("core/dr-catalog has a grandparent");
|
||||||
|
|
||||||
|
let mut paths = Vec::new();
|
||||||
|
for group in ["core", "ui", "apps", "platform"] {
|
||||||
|
crate_sources(&repo.join(group), &mut paths);
|
||||||
|
}
|
||||||
|
let files: Vec<(String, String)> = paths
|
||||||
|
.iter()
|
||||||
|
.map(|p| {
|
||||||
|
let text = fs::read_to_string(p)
|
||||||
|
.unwrap_or_else(|e| panic!("cannot read {}: {e}", p.display()));
|
||||||
|
let name = p.strip_prefix(repo).unwrap_or(p).display().to_string();
|
||||||
|
(name, text)
|
||||||
|
})
|
||||||
|
.collect();
|
||||||
|
|
||||||
|
// A scan over nothing passes for the wrong reason. The queue's own file
|
||||||
|
// and the scan that used to feed it must both have been read, and the
|
||||||
|
// queue's definitions found in them.
|
||||||
|
for must in [
|
||||||
|
"core/dr-catalog/src/jobs.rs",
|
||||||
|
"core/dr-catalog/src/runner.rs",
|
||||||
|
"ui/dr-ui/src/library/scan.rs",
|
||||||
|
] {
|
||||||
|
assert!(
|
||||||
|
files.iter().any(|(n, _)| n == must),
|
||||||
|
"{must} was not scanned — the source walk is wrong, not the code"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
let jobs = &files
|
||||||
|
.iter()
|
||||||
|
.find(|(n, _)| n == "core/dr-catalog/src/jobs.rs")
|
||||||
|
.unwrap()
|
||||||
|
.1;
|
||||||
|
assert!(shipping_code(jobs).contains("pub fn enqueue("));
|
||||||
|
|
||||||
|
let ledger = read(&files);
|
||||||
|
|
||||||
|
assert!(
|
||||||
|
ledger.unreadable.is_empty(),
|
||||||
|
"\n\nThese enqueue calls do not name their JobKind, so this test cannot \
|
||||||
|
check that anything claims it. Spell the kind at the call:\n {}\n",
|
||||||
|
ledger.unreadable.join("\n ")
|
||||||
|
);
|
||||||
|
|
||||||
|
if ledger.claims_everything {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
let orphans: Vec<String> = ledger
|
||||||
|
.enqueued
|
||||||
|
.iter()
|
||||||
|
.filter(|(k, _)| !ledger.claimed.contains(k))
|
||||||
|
.map(|(k, at)| format!("JobKind::{k}, enqueued in {at}"))
|
||||||
|
.collect();
|
||||||
|
assert!(
|
||||||
|
orphans.is_empty(),
|
||||||
|
"\n\nEnqueued, and claimed by nothing (claimed: {:?}):\n {}\n\n\
|
||||||
|
A kind nobody claims is a row per subject that stays for ever — the \
|
||||||
|
queue coalesces, so it never fails, it only grows (#73). Register a \
|
||||||
|
JobHandler for the kind, or stop enqueueing it and add it to \
|
||||||
|
JobKind::RETIRED so the rows already queued are dropped.\n",
|
||||||
|
ledger.claimed,
|
||||||
|
orphans.join("\n ")
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
/// The reader itself, on code whose answer is known — so a parsing bug shows
|
||||||
|
/// up as this failing, rather than the real check quietly finding nothing.
|
||||||
|
#[test]
|
||||||
|
fn the_reader_sees_producers_and_consumers() {
|
||||||
|
let producer = r#"
|
||||||
|
use dr_catalog::jobs;
|
||||||
|
fn persist(tx: &Connection, id: i64) {
|
||||||
|
// jobs::enqueue(tx, JobKind::ContentHash, ...) in a comment is not a call
|
||||||
|
let _ = dr_catalog::jobs::enqueue(
|
||||||
|
tx,
|
||||||
|
JobKind::Thumbnail,
|
||||||
|
Some(id),
|
||||||
|
Priority::Background,
|
||||||
|
None,
|
||||||
|
);
|
||||||
|
jobs::enqueue(tx, kind, Some(id), Priority::Background, None)?;
|
||||||
|
}
|
||||||
|
pub fn enqueue(conn: &Connection, kind: JobKind) {}
|
||||||
|
#[cfg(test)]
|
||||||
|
mod tests {
|
||||||
|
fn t() { enqueue(&c, JobKind::FetchOriginal, None, P, None); }
|
||||||
|
}
|
||||||
|
"#;
|
||||||
|
let consumer = r#"
|
||||||
|
impl JobHandler for Faces {
|
||||||
|
fn kinds(&self) -> &[JobKind] {
|
||||||
|
&[JobKind::DetectFaces]
|
||||||
|
}
|
||||||
|
fn run(&mut self) {}
|
||||||
|
}
|
||||||
|
trait JobHandler { fn kinds(&self) -> &[JobKind]; }
|
||||||
|
fn pull(c: &Connection) { claim_next_matching(c, 0, &[JobKind::FetchPreview]); }
|
||||||
|
"#;
|
||||||
|
let ledger = read(&[
|
||||||
|
("producer.rs".into(), producer.into()),
|
||||||
|
("consumer.rs".into(), consumer.into()),
|
||||||
|
]);
|
||||||
|
|
||||||
|
let enqueued: Vec<&str> = ledger.enqueued.iter().map(|(k, _)| k.as_str()).collect();
|
||||||
|
assert_eq!(enqueued, vec!["Thumbnail"]);
|
||||||
|
assert_eq!(ledger.unreadable.len(), 1, "{:?}", ledger.unreadable);
|
||||||
|
assert_eq!(
|
||||||
|
ledger.claimed,
|
||||||
|
BTreeSet::from(["DetectFaces".to_string(), "FetchPreview".to_string()])
|
||||||
|
);
|
||||||
|
assert!(!ledger.claims_everything);
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user