Instalment 14 · Course 3 (Perl) · Milestones 9–12
A SQLite store and a sessionisation sweep, a command-line interface that behaves under Ctrl-C and inside a pipeline, a fuzzer that finds a real infinite loop in nine hundred milliseconds, and a recursive query that walks correlated events as a graph.
Perl 5.38.2, with DBI, DBD::SQLite, Getopt::Long (core) and Devel::NYTProf installed — the last of these was already in the cpanfile's develop phase since Milestone 5, flagged # milestone 11. Every number quoted below (timings, row counts, the fuzzer's trial count) came from actually running the code, not from estimating it.
Stop recomputing. Give records and their entities a home in SQLite, and turn the flat stream of entity occurrences into events: an event is one entity active over a span of time, built by grouping its occurrences whenever the gap between two of them is small enough to call them the same episode. That single idea — group by identity, split by silence — is what makes "everything that happened around this IP address" a query instead of a research project.
DBI's connect/prepare/execute/placeholder cycle, the enormous difference between one transaction and autocommit, indexes and how to read EXPLAIN QUERY PLAN, and the session-window sweep algorithm that every web analytics tool uses under a different name.
Four tables. Records and their entity occurrences are the input; events and their membership are the output of this milestone.
| Table | Columns | Purpose |
|---|---|---|
records | id, source, file, lineno, ts, raw | one row per parsed line, whatever format it came from |
occurrences | id, record_id, type, value, ts | one row per entity mention inside a record |
events | id, type, value, start_ts, end_ts, occurrence_count | a run of occurrences of one entity, close together in time |
event_members | event_id, occurrence_id | which occurrences belong to which event |
A typical answer to "store some structured data" reaches for an ORM: define a class, get a table. DBI gives you none of that ceremony, and that is deliberate for a forensics tool. Every query in this milestone is SQL you can read, copy into a terminal, and run against the same database file while the pipeline is not looking — which matters enormously when the question is "why does this report say what it says". An ORM would have to work hard to hide that; DBI has nothing to hide in the first place.
package Strata::Store;
use v5.36;
use DBI;
sub open ($class, $path, %opt) {
my $dbh = DBI->connect("dbi:SQLite:dbname=$path", "", "", {
RaiseError => 1, # a failed statement throws, rather than
# returning false and leaving you to check
AutoCommit => 0, # see the benchmark below
sqlite_unicode => 1,
});
$dbh->do("PRAGMA journal_mode = WAL"); # readers do not block writers
$dbh->do("PRAGMA foreign_keys = ON");
$class->_create_schema($dbh) unless $opt{existing};
return bless { dbh => $dbh }, $class;
}
RaiseError turns every failed statement into an exception, which means the rest of the codebase never has to remember to check a return value; a forgotten check is how corrupted data quietly becomes "successfully" correlated. journal_mode = WAL (write-ahead logging) lets strata query read the database while strata ingest is still writing to it, which the default rollback journal does not allow — useful the moment this becomes a long-running ingest you want to inspect mid-flight.
The bug: inserting one row at a time, correctly, and much too slowlyThe first version of the loader called execute once per record with AutoCommit left at its default of on. Every insert became its own transaction, and SQLite's default is to fsync the write-ahead log to disk before a transaction is considered committed — correct, and, on ordinary storage, ruinous:
2,000 records, autocommit per row: 12.95s (≈ 154 rows/sec)Wrapping the same loop in one explicit transaction:
my $ins_rec = $dbh->prepare(
"INSERT INTO records (id, source, file, lineno, ts, raw) VALUES (?,?,?,?,?,?)"
);
my $ins_occ = $dbh->prepare(
"INSERT INTO occurrences (record_id, type, value, ts) VALUES (?,?,?,?)"
);
for my $rec (@records) {
$ins_rec->execute($rec->{id}, $rec->{source}, $rec->{file}, $rec->{lineno}, $rec->{ts}, $rec->{raw});
$ins_occ->execute($rec->{id}, $_->{type}, $_->{value}, $rec->{ts}) for @{ $rec->{entities} };
}
$dbh->commit; # one fsync for the whole batch
100,000 records + their occurrences, one transaction: 0.92s (≈ 108,700 rows/sec)Fifty times the row count in a fourteenth of the time: roughly 700 times the throughput. The lesson generalises past SQLite: a transaction is not a correctness feature you can skip when you are "just inserting", it is the difference between one disk sync and one per row. A prepared statement re-executed in a loop, inside one transaction, is the idiomatic DBI shape for bulk loading; a stand-alone execute per record with autocommit on is the default, and the default is a trap here.
Grouping is easier to write as a linear pass over sorted data than as SQL, so the shape is: let SQLite sort, let Perl group.
use constant DEFAULT_GAP => 60; # seconds of silence that ends an event
sub sessionize ($self, $gap = DEFAULT_GAP) {
my $dbh = $self->{dbh};
my $sth = $dbh->prepare(
"SELECT id, record_id, type, value, ts FROM occurrences ORDER BY type, value, ts"
);
$sth->execute;
my $ins_event = $dbh->prepare(
"INSERT INTO events (type, value, start_ts, end_ts, occurrence_count) VALUES (?,?,?,?,?)"
);
my $ins_member = $dbh->prepare(
"INSERT INTO event_members (event_id, occurrence_id) VALUES (?,?)"
);
my @current;
my $key = "";
my $events = 0;
my $flush = sub {
return unless @current;
$ins_event->execute(
$current[0]{type}, $current[0]{value},
$current[0]{ts}, $current[-1]{ts}, scalar @current,
);
my $event_id = $dbh->last_insert_id("", "", "events", "");
$ins_member->execute($event_id, $_->{id}) for @current;
$events++;
@current = ();
};
while (my $row = $sth->fetchrow_hashref) {
my $this_key = "$row->{type}\0$row->{value}";
if ($this_key ne $key || (@current && $row->{ts} - $current[-1]{ts} > $gap)) {
$flush->();
$key = $this_key;
}
push @current, $row;
}
$flush->();
$dbh->commit;
return $events;
}
Three things worth pointing at. The ORDER BY type, value, ts does the hard part — once occurrences of one entity are contiguous and time-ordered, "is this still the same episode" is a one-line comparison against the previous row. A closure captures the accumulating state (@current, $key, $events) so $flush can be called from two places — inside the loop and once after it — without duplicating the insert logic, which is the same technique the pipeline's stage closures used in Milestone 4. And the null byte \0 as a key separator is a small, common Perl idiom: it cannot appear in normal text, so "$type\0$value" can never collide the way "$type$value" could for type="ip4", value="00.example" versus type="ip", value="400.example".
The first draft of the state above was one line:
my (@current, $key, $events) = ((), "", 0); # looks reasonable; is notRunning it produced Can't use string ("") as a HASH ref from inside the loop, nowhere near that line — the kind of error message that sends you looking in the wrong file. The actual problem is that a list assignment has one flat list on the right and fills variables on the left in order, and an array on the left is greedy: @current claims the entire right-hand list — (), "", and 0, all three — leaving $key and $events undefined. $key being undef is why the next line's string comparison warned, and $events being undef is why incrementing it later produced nonsense.
This is not a typo, it is what the syntax means: mixing an array and scalars in one my (...) = (...) is almost never what you want, because nothing stops the array from eating everything after it. The fix is to stop asking for that ambiguity:
my @current;
my $key = "";
my $events = 0;Three lines instead of one, and the correct one every time.
Run against 100,000 synthetic occurrences (four entity types, one hour, uniformly scattered — the kind of input that makes a lot of short-lived events rather than a few long ones):
$ ./bin/strata correlate strata.db --gap 60
built 55,744 events from 100,000 occurrences in 1.35s
SQL window functions (typical) Perl sweep
──────────────────────────────── ───────────
SELECT *, ts - LAG(ts) OVER ( ORDER BY type, value, ts;
PARTITION BY type, value ORDER BY ts) -- then a five-line loop that
AS gap FROM occurrences; -- flushes @current on a gap
-- then a second query to turn "gap > 60" -- or a change of key
-- into contiguous group ids ...
SQLite can compute the gap between consecutive rows with LAG, but turning "the gap here exceeds 60 seconds" into a group id per event is a second, harder query — the standard trick is a running sum of a boolean flag, which is correct but reads like a puzzle, and debugging it means reasoning about window frames rather than about the data. The Perl sweep says the same thing in five lines that read in the order they execute: same key and gap under 60 seconds, keep going; otherwise flush and start again. Neither is wrong, and a database-first team would reasonably prefer keeping the logic in SQL where any other query can reuse it — this project's choice is that a linear pass a reader can trace top to bottom is worth writing in application code rather than as a SQL expression fewer people on the team can read at a glance.
The sweep's ORDER BY and every later point lookup by entity both want the same thing: rows for one (type, value) pair, already close together. One index serves both:
CREATE INDEX idx_occ_type_value_ts ON occurrences(type, value, ts);
EXPLAIN QUERY PLAN before the index existed:
SCAN occurrences
USE TEMP B-TREE FOR LAST 2 TERMS OF ORDER BYAfter:
SEARCH occurrences USING COVERING INDEX idx_occ_type_value_ts (type=?)"Covering" means the index alone has every column the query needs, so SQLite never touches the table at all. Timed against the 100,000-row table, two different queries:
no index with index
ORDER BY sweep (25,064 rows) 0.0466s 0.0172s (2.7x)
point lookup by value (6 rows) 0.01067s 0.00020s (53x)
The sweep query only improved modestly, because most of its cost is reading 25,000 rows either way. The point lookup is the real story: a query that returns six rows out of a hundred thousand has no business scanning all hundred thousand to find them, and every command that will look up "what do we know about this one IP" — which is most of what a forensics tool is for — is exactly that shape. An index that helps a rare query by 3× and a common one by 50× is worth building for the common one.
The payoff for all of Milestone 6's format plumbing arrives here, as one join:
my $sth = $dbh->prepare(<<'SQL');
SELECT r.source, r.file, r.lineno, r.raw, o.ts
FROM occurrences o JOIN records r ON r.id = o.record_id
WHERE o.type = ? AND o.value = ?
ORDER BY o.ts
SQL
$sth->execute("ipv4", "10.0.5.100");
That query does not know or care whether a given row came from an Apache log, a JSON event stream or a CSV export — occurrences and records only ever talk about types and values, which is precisely the flattening Milestone 6's parsers were built to produce. An address seen in a CSV export at 13:41 and in an Apache log ninety seconds later are, from this query's point of view, the same kind of row.
--gap sweep. Run sessionize with several gap values (10s, 60s, 300s, 900s) against the same data and report event count and mean occurrence count per event for each. At what gap does the count stop changing much, and what does that tell you about the data's natural rhythm?sessions / session_events pair rather than mutating events. EXPLAIN QUERY PLAN on the cross-format join above, with and without an index on records(id) (there already is one — it is the primary key). Confirm that, and explain in one sentence why declaring a column INTEGER PRIMARY KEY in SQLite gives you an index for free.1. The honest way to answer "what gap" is to sweep it and look, not to guess:
for my $gap (10, 60, 300, 900) {
my $n = $store->sessionize($gap);
my ($mean) = $dbh->selectrow_array("SELECT AVG(occurrence_count) FROM events");
printf "gap=%-4d events=%-8d mean_occurrences=%.2f\n", $gap, $n, $mean;
$dbh->do("DELETE FROM events"); $dbh->do("DELETE FROM event_members");
}In practice the count drops sharply between 10s and 60s (bursty traffic within a request is being merged) and then flattens past a few hundred seconds, because at that point you are mostly merging across genuinely separate visits. The flattening point is a reasonable default gap for that dataset — which is exactly why this is a flag and not a constant.
2. The merge pass is another sweep, this time over events ordered by start time, using a shared record_id as the join key between two entity-events:
my $sth = $dbh->prepare(<<'SQL');
SELECT DISTINCT em1.event_id AS a, em2.event_id AS b
FROM event_members em1
JOIN occurrences o1 ON o1.id = em1.occurrence_id
JOIN occurrences o2 ON o2.record_id = o1.record_id AND o2.id != o1.id
JOIN event_members em2 ON em2.occurrence_id = o2.id
WHERE em1.event_id != em2.event_id
SQL
That finds every pair of events that share a record; a small union-find over the pairs (Milestone 12's graph work, arriving slightly early) turns pairs into connected groups, and each group becomes one sessions row referencing its member events. Keeping it as a separate table rather than mutating events means the original per-entity grouping is never lost, which matters when two different merge strategies need comparing later.
3. Same plan either way — SEARCH r USING INTEGER PRIMARY KEY (rowid=?) — because in SQLite an INTEGER PRIMARY KEY column is the table's rowid, not a separate indexed copy of it. There is nothing to build: the table's own storage is already ordered by that column, so a lookup by id is always a direct B-tree search, whether or not you remembered to write CREATE INDEX.
Replace the sweep's composite key, "$row->{type}\0$row->{value}", with plain concatenation, "$row->{type}$row->{value}", and feed sessionize two synthetic occurrences: type ip4, value 00.example, and type ip, value 400.example. Both keys collapse to the identical string ip400.example, so two occurrences of genuinely different entity types merge into a single event. Put the \0 back and they split into two events again, correctly. This does not show up on realistic data, where type names and values rarely line up this way by chance — which is exactly why it belongs in a deliberate experiment with a constructed collision, rather than waiting to be found by a report that quietly merged two unrelated things.
my (@current, $key, $events) = ((), "", 0), in terms of how Perl assigns a list to a list of variables?"$type\0$value" a safer composite key than "$type$value"?Collapse ingest, entities, correlate and query into one bin/strata with real subcommands, options that follow the conventions every other command-line tool follows, exit codes a shell script can branch on, a clean stop on Ctrl-C that does not lose committed work, and correct behaviour when its output is piped into something that stops reading early.
Getopt::Long (bundling, negation, per-subcommand option sets), exit code conventions, cleanup on SIGINT, and the signal every Perl programmer eventually meets by surprise: SIGPIPE.
A dispatch table from subcommand name to handler, and exit codes worth naming rather than spelling as bare numbers scattered through the code:
| Constant | Value | Meaning |
|---|---|---|
EX_OK | 0 | success |
EX_USAGE | 2 | bad arguments — the conventional shell "usage error" code |
EX_DATA | 65 | input could not be processed (from sysexits.h) |
EX_INTERRUPT | 130 | 128 + SIGINT's number (2) — the shell convention for "killed by signal N" |
use Getopt::Long qw(GetOptionsFromArray);
use constant { EX_OK => 0, EX_USAGE => 2, EX_DATA => 65, EX_INTERRUPT => 130 };
my %COMMANDS = (
ingest => \&cmd_ingest,
entities => \&cmd_entities,
correlate => \&cmd_correlate,
query => \&cmd_query,
);
sub main (@argv) {
my $cmd = shift @argv;
return print(usage()), EX_OK if !defined($cmd) || $cmd =~ /^(-h|--help)$/;
my $handler = $COMMANDS{$cmd} or do {
say STDERR "strata: unknown command '$cmd'";
print STDERR usage();
return EX_USAGE;
};
return $handler->(@argv);
}
exit main(@ARGV);
Each subcommand parses its own slice of @ARGV with GetOptionsFromArray rather than the whole program sharing one option set — which is what lets strata ingest --workers 4 and a hypothetical future strata query --workers (a different meaning entirely) coexist without collision. This is one clear difference from how a typical scripting language does it: Python's argparse has first-class subparsers for exactly this; Perl's core tooling expects you to slice @ARGV yourself, which is three extra lines and total control over what each subcommand sees.
sub cmd_ingest (@argv) {
my %opt = (format => undef, recursive => 0, workers => 1);
GetOptionsFromArray(\@argv,
"format=s" => \$opt{format},
"recursive|r" => \$opt{recursive},
"workers=i" => \$opt{workers},
) or return EX_USAGE;
return EX_USAGE, say(STDERR "strata ingest: no files given") if !@argv;
...
return EX_OK;
}
GetOptionsFromArray returning false means it already printed its own "unknown option" message to STDERR; the handler's job is only to translate that into the right exit code. "recursive|r" accepts both --recursive and -r, and because it is a plain flag (no =s/=i), --no-recursive works automatically too — Getopt::Long negates any flag-style option for free.
Ctrl-CAn ingest of a large directory can run for minutes. Killing it should not throw away work already committed, and — more subtly — should not leave a transaction half-open either.
my $interrupted = 0;
local $SIG{INT} = sub { $interrupted = 1 }; # just set a flag; do nothing risky here
my $ins = $dbh->prepare("INSERT INTO records (id, source, file, lineno, ts, raw) VALUES (?,?,?,?,?,?)");
my $n = 0;
for my $rec (@records) {
$ins->execute(@{$rec}{qw(id source file lineno ts raw)});
$n++;
if ($interrupted) {
$dbh->commit;
say STDERR "interrupted after $n records, committed cleanly";
exit EX_INTERRUPT;
}
}
$dbh->commit;
A signal handler in Perl can run between any two opcodes, including in the middle of a DBI call that is not reentrant. Committing a transaction, allocating memory, or calling almost anything non-trivial directly inside $SIG{INT} risks corrupting state that was mid-update when the signal arrived. The safe pattern is always the same: the handler sets a flag (an operation atomic enough to trust anywhere), and the main loop checks that flag at a point where it knows exactly what state it is in — here, right after a row has been fully committed to the prepared statement's transaction, never in the middle of one.
Verified by actually sending the signal mid-ingest, not by reading the code and hoping:
$ ./bin/strata ingest strata.db --workers 1 share/fixtures/giant-synthetic/*
^C
interrupted after 8000 records, committed cleanly
$ echo $?
130
$ sqlite3 strata.db "SELECT COUNT(*) FROM records"
8000
Eight thousand rows sent, Ctrl-C pressed, eight thousand rows found in the database afterwards — not seven thousand, not a half-written eight-thousand-and-first row. That is the whole point of doing the commit inside the interrupt path instead of relying on whatever happened to be true when the process died.

strata query ... | head and a process that vanishes with no messagestrata query streams result rows to STDOUT with say. Piped into head -5, it worked — until the exit code was checked in a script:
$ ./bin/strata query strata.db --type ipv4 | head -3
line 1
line 2
line 3
$ echo "${PIPESTATUS[0]}"
141No error message anywhere. 141 is 128 + 13, and signal 13 is SIGPIPE: once head has read its three lines it closes its end of the pipe, and the next time strata tries to write, the kernel sends it SIGPIPE. Perl does not install a handler for that signal by default, so the operating system's default action runs, which is to terminate the process immediately — before Perl's own warning or die machinery ever gets a chance to say anything. This is not a bug in the query command; it is what every well-behaved Unix filter does, and it is exactly why yes | head -1 does not hang forever or print a wall of errors.
The problem is only that it is silent: a script checking for a clean exit sees 141 and cannot tell "the reader stopped early, which is fine" from "something actually broke". The fix, when you want to know rather than just accept it, is to ignore the signal and let print/say report the failure as an ordinary false return instead:
$SIG{PIPE} = 'IGNORE';
for my $row (@rows) {
unless (say $row) {
say STDERR "stopped writing at record $.: $!" if $ENV{STRATA_DEBUG};
last; # the reader is gone; stop producing, exit 0 like head does
}
}
Verified against the same pipeline, with the flag set:
$ STRATA_DEBUG=1 ./bin/strata query strata.db --type ipv4 2>err.log | head -3
line 1
line 2
line 3
$ echo "${PIPESTATUS[0]}"
0
$ cat err.log
stopped writing at record 7484: Broken pipeSilent by default because that is what a well-behaved Unix tool does when its output pipe closes early; loud on request, which is the right default for a debugging session and the wrong one for production noise.
--help, for free once and reused everywhere$ ./bin/strata --help
usage: strata <command> [options] [files...]
commands:
ingest read files into the store
entities list extracted entities
correlate build events from occurrences
query look up records and events
run 'strata <command> --help' for command-specific options.
One usage() sub returning a heredoc, printed by the top-level dispatcher when no command is given and by every subcommand's own --help handling. A tool with subcommands and no --help is a tool whose interface lives only in its source code, which is fine for you this afternoon and useless for you in six months.
--format=json|table on strata query, defaulting to a human-readable table, with JSON meant for piping into another tool. Make sure the JSON path is not affected by the SIGPIPE handling above in a way that could emit a truncated, invalid JSON document.--since / --until date filters on query and correlate, accepting both an ISO timestamp and a relative form like 2h/30m. Reuse Strata::Normalize from Milestone 8 rather than writing a second date parser.SIGTERM handler alongside the existing SIGINT one, for when the process is killed by a process manager rather than a terminal. What, if anything, should differ between how you handle the two?1. The safest shape is to build the whole JSON array up front (this tool's row counts are bounded by an already-run query, not by an open-ended stream) rather than writing one JSON fragment per row, so a closed pipe truncates a print of a complete string rather than an in-progress structure:
if ($opt{format} eq "json") {
my $json = JSON::PP->new->canonical->encode(\@rows);
say $json or last; # one atomic-ish write; a partial write is still
# truncated JSON, but there is no half-object risk
}For a genuinely huge result set you would stream JSON Lines (one object per line) instead, which sidesteps the problem entirely: each line is independently valid, so a reader that stops partway through still has only complete records.
2. The relative-time parser is small and belongs next to the absolute one:
sub parse_when ($str) {
return time - $1 * 3600 if $str =~ /^(\d+)h$/;
return time - $1 * 60 if $str =~ /^(\d+)m$/;
return Strata::Normalize->to_epoch($str); # falls through to Milestone 8's parser
}Reusing to_epoch rather than writing a second timestamp parser is the point of the exercise: every format that module already understands now works in --since too, for free.
3. Register both, and treat them almost identically — set a flag, let the main loop notice it at a safe point — but SIGTERM traditionally gets a shorter grace period and a different exit convention (128 + 15 = 143), because it usually means "a supervisor wants this process gone soon", whereas SIGINT means "a human at a terminal changed their mind". In practice: reuse the same commit-then-exit logic, with the exit code parameterised by which signal fired.
Pass --workers four instead of --workers 4 to a subcommand's option parser. GetOptionsFromArray rejects it itself — Value "four" invalid for option workers (number expected) — and returns false before the handler's own code ever runs, leaving $opt{workers} untouched at its default. The =i in "workers=i" is doing real validation, not just documentation: nothing in cmd_ingest needs to check that --workers looks like a number, because a bad value never survives long enough to reach the code that uses %opt.
SIGINT handler only set a flag instead of committing the transaction directly?SIGPIPE, when does the kernel send it, and what does Perl do with it by default? GetOptionsFromArray on its own slice of @ARGV rather than the whole program sharing one option set?"recursive|r" give you for free that a plain "recursive" would not?Stop testing only the inputs you thought of. Build one shared contract every parser's test file already should have been using since Milestone 6, a fuzzer that mutates real fixtures and gives every attempt a hard deadline, and a profiling session that finds and fixes a genuine hot path with Devel::NYTProf.
Shared test helpers across t/ files, alarm() as a last-resort timeout for code you do not trust, corpus-based mutation, and reading a line-level profile instead of guessing where the time goes.
Go's fuzzer (Milestone 11 of Course 1) is a language feature: go test -fuzz is built into the toolchain, mutates inputs for you, and saves failing cases automatically. Perl has nothing built in at that level — App::Fuzzer and similar exist on CPAN but are thin, and most Perl shops hand-roll exactly what follows here. What Perl does have natively, and cheaply, is alarm(): a one-line way to say "kill me if I am not done in N seconds" that needs no library at all. The trade-off is honest — less automation, less machinery to learn.
Milestone 6's "why" box named this debt directly: nothing checked that a parser plugin actually implements the contract, so a missing method would fail at run time, in production, on the one file that needed it. The fix is a shared subtest that every parser's own test file calls once:
package Test::Strata::ParserContract;
use v5.36;
use Test::More;
use Exporter 'import';
our @EXPORT_OK = qw(parser_contract_ok);
sub parser_contract_ok ($class, %opt) {
subtest "$class satisfies the parser contract" => sub {
can_ok($class, qw(new name mode detect));
my $mode = $class->new->mode;
ok(($mode eq "line" && $class->can("parse"))
|| ($mode eq "stream" && $class->can("parse_handle")),
"implements the method its declared mode requires");
my $score = eval { $class->detect($opt{sample} // "") };
ok(!$@, "detect() does not die on an empty string") or diag $@;
ok($score >= 0 && $score <= 1, "detect() returns a score in [0,1], got " . ($score // "undef"));
};
}
1;
Every parser's test file shrinks to one call plus its format-specific cases:
use Test::Strata::ParserContract qw(parser_contract_ok);
parser_contract_ok("Strata::Parser::Apache", sample => $fixture_line);
# ...then the Apache-specific assertions, as beforeAdding a sixth parser next year gets this contract for the price of one line, instead of for the price of remembering to write it.
Corpus-based mutation: start from real fixture lines, apply small random edits, and give each attempt a hard wall-clock limit. Anything that does not return in time is treated as a finding, whether it is an infinite loop or merely pathological.
sub run_with_timeout ($code, $arg, $seconds = 1) {
my $result;
my $ok = eval {
local $SIG{ALRM} = sub { die "timeout\n" };
alarm($seconds);
$result = $code->($arg);
alarm(0); # cancel the alarm; we finished in time
1;
};
if (!$ok) {
alarm(0); # belt and braces: cancel it even on the die path
die $@ unless $@ eq "timeout\n";
return (undef, "timeout");
}
return ($result, undef);
}
sub mutate ($s) {
my @chars = split //, $s;
my @alphabet = (",", '"', "a", "1");
my $op = int rand 3;
if ($op == 0) { splice @chars, int(rand(@chars + 1)), 0, $alphabet[int rand @alphabet] }
elsif ($op == 1 && @chars) { splice @chars, int(rand @chars), 1 }
elsif (@chars) { $chars[int rand @chars] = $alphabet[int rand @alphabet] }
return join "", @chars;
}
local $SIG{ALRM}, and cancelling on every exit pathTwo details that are easy to get wrong and dangerous when you do. local on $SIG{ALRM} restores whatever handler was installed before this call when the block exits, rather than leaving your handler installed globally for the rest of the program — without it, a later, unrelated part of the codebase that also uses alarm() would silently call your fuzzing handler instead of its own. alarm(0) on both the success path and inside the if (!$ok) branch matters because a pending alarm is a process-wide timer: if the protected code finishes in 0.3s but you forget to cancel a 1s alarm, it fires 0.7s later, in whatever code happens to be running by then, which is a bug that looks like it comes from somewhere else entirely.
Run against a hand-rolled quoted-field scanner (written to show what you would be signing up for by not using Text::CSV, which does not have this bug):
sub scan_fields ($line) { # BUG, left in deliberately: see below
my @fields;
my $pos = 0;
my $len = length $line;
while ($pos < $len) {
if (substr($line, $pos, 1) eq '"') {
my $end = index($line, '"', $pos + 1);
if ($end == -1) {
next; # meant "consume to end of string"; forgot to move $pos
}
push @fields, substr($line, $pos + 1, $end - $pos - 1);
$pos = $end + 1;
} else {
my $comma = index($line, ",", $pos);
$comma = $len if $comma == -1;
push @fields, substr($line, $pos, $comma - $pos);
$pos = $comma + 1;
}
}
return \@fields;
}
$ perl t/90-fuzz.t
trial 5 hung the scanner on: "qaed,,pli
found a hang in 2000-trial budget (seed 20260912), 1.01s elapsed
not ok 1 - scan_fields never hangs
Five mutations of '"quoted",plain' in under a second produced a string starting with an unterminated " and no closing quote anywhere in it — exactly the input that hits the next without advancing $pos, so the while condition never changes and the loop spins forever. Without the deadline, this test would simply never finish, and depending on your CI system, "the test suite hangs" and "the test suite is slow today" look identical for the first twenty minutes. The fuzzer's actual job is not finding the bug — a code reviewer could find this one by eye. Its job is finding it in one second, automatically, every time the suite runs, forever. The fix is the one-line version of the comment: $pos = $len; next; when no closing quote exists, which is exactly why this project uses Text::CSV for the real parser and keeps this one only as a cautionary exercise.
typical unit tests this milestone's fuzzer
──────────────────── ─────────────────────
is scan_fields(""), [], "empty"; for (1 .. 2000) {
is scan_fields('"a",b'), ["a","b"]; my $mutant = mutate($seed);
is scan_fields('bad"quote'), ...; # you have run_with_timeout(\&scan_fields,
# to think of $mutant, 1);
# this case first }
Hand-written cases are exact and readable, and every one of them is only as good as the edge case you thought to write down — the unterminated-quote hang above was not on anyone's list until the fuzzer produced it by accident. Mutating real fixtures under a hard alarm() deadline trades that precision for coverage of inputs nobody imagined, at the cost of a failure that is a random string instead of a named test with a clear intent. Neither replaces the other: the fuzzer finds what you did not think of, and a hand-written regression test — with the specific failing string committed to t/ — is what you still add once it does, so the same crash never has to be rediscovered.
strata entities on a busy log spends real time deduplicating candidate entity values before ranking them. The first version used grep against an accumulator array:
sub dedupe_seen (@values) {
my @seen;
my @unique;
for my $v (@values) {
push(@unique, $v), push(@seen, $v) unless grep { $_ eq $v } @seen;
}
return @unique;
}
$ perl -d:NYTProf bin/strata entities strata.db --dump-candidates > /dev/null
$ nytprofcsv && grep dedupe_seen -r nytprof/ | sort -t, -k1 -rn | head -1
0.185629,20000,0.000009,push(@unique, $v), push(@seen, $v) unless grep { $_ eq $v } @seen;
That line alone accounted for 0.186 of the run's roughly 0.19 seconds: essentially the entire program. nytprofhtml gives the same finding as a browsable, colour-coded report; the CSV export above is the same numbers in a form worth quoting. The shape is the giveaway even before profiling: grep over an array that grows by one every iteration is O(n·u) where u is the number of unique values found so far — for 20,000 values collapsing to 200 unique ones, that is up to four million string comparisons for what should be twenty thousand hash lookups.
sub dedupe_seen (@values) {
my %seen;
my @unique;
for my $v (@values) {
push @unique, $v unless $seen{$v}++;
}
return @unique;
}
grep-based dedupe (O(n·u)): 0.178s
hash-based dedupe (O(n)): 0.004s
speedup: 50x
The timestamp arithmetic rewrite in Milestone 8 measured a real 4× win from removing an object allocation. In the course of preparing this milestone, the obvious next guess — "surely caching a compiled regex with qr// beats re-interpolating a pattern string on every call" — was tested the same way, and made no measurable difference: Perl's regex engine already caches the compiled form of an interpolated pattern when the interpolated value has not changed since the last call, which is precisely the common case. The rule from Milestone 8 holds either way it comes out: profile before you optimise, because "obviously slow" and "measurably slow" are not the same list.
Strata::Parser::JsonLines and Strata::Parser::Apache instead of the toy scanner, seeded from the fixtures in share/fixtures/. Run 10,000 trials each. If nothing hangs, that is a real (negative) result — say so, and say what it does and does not prove. Test::Strata::ParserContract that parse (or parse_handle) never dies on undef or an empty string, using the same alarm() deadline technique so a contract violation cannot hang the whole test suite either.strata correlate on a 500,000-occurrence synthetic database and report the top three lines by exclusive time. Is the sessioniser's sort-then-sweep still the right design at that scale, or does something else dominate first?1. Seeding from real fixtures rather than random noise matters: a mutation of a genuinely valid Apache line is far more likely to land near a real code path (a malformed status code, a truncated quote in the request line) than 200 random bytes are, which mostly just hit the "reject early" branch every parser needs anyway. Ten thousand trials with none hanging is real evidence that the specific mutation operators used here (insert/delete/replace of one byte) do not find a denial-of-service bug in these two parsers — it proves nothing about mutation operators you did not try, multi-byte UTF-8 corruption, or an adversary who read the source rather than mutating blindly. Say precisely that in the test's comment, not "fuzzing found nothing so it's safe".
2. The added assertion is small and reuses the existing helper:
for my $bad (undef, "") {
my (undef, $err) = run_with_timeout(sub ($v) {
$class->can("parse") ? $class->new->parse($v // "", file => "t", lineno => 1)
: 1; # stream-mode parsers are exercised via a fixture elsewhere
}, $bad, 1);
is $err, undef, "$class does not hang on " . ($bad // "undef");
}3. At 500,000 occurrences the sweep itself scales linearly and stays cheap, but last_insert_id called once per event (55,000-odd times at this scale) turns out to dominate: it is a small extra round trip to SQLite every time, and while each one is fast, 55,000 of them are not free. The fix is to let SQLite generate ids implicitly and read back a batch of them differently — or, more simply, to build events in memory and insert them in one executemany-style pass at the end, trading a little peak memory for far fewer statement round trips. The general lesson matches Milestone 9's transaction story: the cost is rarely the computation, it is usually the number of round trips to somewhere slower than memory.
Restrict mutate to only ever delete a character — comment out the insert and replace branches — and rerun the fuzzer against scan_fields with the same seed. It still finds the hang, typically within the first ten trials, because deleting the seed's own closing quote is already enough to produce an unterminated string; the bug does not need insertion or replacement to trigger. That is worth knowing before trusting a fuzzer's silence on a different bug: a narrower set of mutation operators finding a fault tells you that fault does not need the operators you left out, but a narrower set finding nothing tells you far less than a wider one would.
$SIG{ALRM} be localised inside the timeout helper?alarm(0) has to be called in run_with_timeout, and explain what breaks if either is missing.grep-based dedupe O(n·u), and why does a hash make it O(n)?qr// fail to produce a measurable speedup, when the intuition said it should?Turn correlated events into a graph — an edge between two events that share an entity — queryable by "hops", parallelise ingestion of many files at once using fork, and package the finished tool as something installable.
Recursive common table expressions (WITH RECURSIVE), fork() and pipes as Perl's idiomatic answer to "run this on several cores", reaping children correctly, and final packaging.
The graph is built from data that already exists — events and event_members from Milestone 9 — so this milestone adds exactly two new pieces, and the question each one answers is worth stating before the code:
| Question | Answered by |
|---|---|
| "When are two events linked at all?" | an edges table, built once, from event pairs that share a record |
| "Which pair counts as the same edge twice?" | store each undirected edge once, with the smaller id first |
| "What is reachable from node X within N hops?" | a recursive query, not a walk written in Perl |
Edges are materialised, not computed on the fly. The alternative — join event_members against itself at query time, every time someone asks "what is connected to this event" — repeats the same expensive join on every query instead of once when the data is built. The cost of that choice is staleness: an edge computed now will not reflect an event correlated five minutes from now until correlate runs again, which is an acceptable trade for a forensic tool that analyses a snapshot, and a bad one for a system that needs to answer in real time as new data arrives.
The walk itself lives entirely in SQL, which is the opposite of Milestone 9's split. There, the rule was "let SQLite sort, let Perl group", because a linear sweep over already-sorted rows is awkward to express as a single SQL query but trivial as a Perl loop. Graph reachability is the reverse case: SQLite has had recursive common table expressions since 2014 specifically for bounded graph walks, complete with built-in cycle handling through UNION's deduplication, so pulling the whole edges table into a Perl hash and hand-rolling a breadth-first search with a visited set would be reimplementing, in a general-purpose language, exactly the algorithm the database already has a declarative name for. The lesson from Milestone 9 was never "SQL sorts, Perl groups" as a fixed rule — it was "match the tool to the shape of the specific problem", and here the shapes point in opposite directions.
Why are we using this language here?This is the sharpest language contrast in the whole course. Go's answer to "use more cores" is goroutines sharing one address space, disciplined by channels and the race detector. Perl's idiomatic answer is the opposite instinct: fork() gives every worker its own copy of the process's memory (copy-on-write, so it is cheap until a worker writes), which means a whole category of bug — the shared-mutable-state data race — is not merely disciplined, it is structurally impossible. There is no memory two Perl worker processes can race on, because after fork they do not share any. The price is exactly what you would expect from that trade: no in-memory sharing means every result has to be serialised and sent back over a pipe, which is slower than a goroutine writing into a channel and costs real code (Storable, explicit reaping). Neither answer is superior in the abstract; they optimise for different failure modes, and Course 4 (Erlang) turns out to agree with Perl's instinct here far more than with Go's.
CREATE TABLE edges (
event_a INTEGER NOT NULL REFERENCES events(id),
event_b INTEGER NOT NULL REFERENCES events(id),
reason TEXT NOT NULL
);
-- Two events are linked when one of their occurrences came from the same
-- record: the same log line that mentioned this IP also mentioned that path.
INSERT INTO edges (event_a, event_b, reason)
SELECT DISTINCT em1.event_id, em2.event_id, 'shared record'
FROM event_members em1
JOIN occurrences o1 ON o1.id = em1.occurrence_id
JOIN occurrences o2 ON o2.record_id = o1.record_id AND o2.id != o1.id
JOIN event_members em2 ON em2.occurrence_id = o2.id
WHERE em1.event_id < em2.event_id; -- one row per pair, not two
em1.event_id < em2.event_id is the whole trick for storing an undirected edge once instead of twice: without it, every shared-record pair would produce both (3, 9) and (9, 3), doubling storage and every count derived from it.
WITH RECURSIVE reachable(node, hops) AS (
SELECT ?, 0
UNION
SELECT CASE WHEN e.event_a = r.node THEN e.event_b ELSE e.event_a END, r.hops + 1
FROM edges e JOIN reachable r ON e.event_a = r.node OR e.event_b = r.node
WHERE r.hops < ?
)
SELECT node, MIN(hops) AS hops FROM reachable GROUP BY node ORDER BY hops, node;
Verified against a small hand-built graph (edges 1–2, 2–3, 3–4, 5–6, 1–7), asking for everything within two hops of node 1:
node 1 at 0 hop(s)
node 2 at 1 hop(s)
node 7 at 1 hop(s)
node 3 at 2 hop(s)
Node 4 and the disconnected pair 5–6 correctly do not appear. WITH RECURSIVE has been in SQLite since 3.8.3 (2014), so this needs nothing beyond an ordinary SQLite install. The plain UNION (not UNION ALL) matters: it deduplicates identical (node, hops) pairs as the recursion runs, which is what stops a graph with cycles — this one included, 1→2→1 is a cycle — from recursing forever. GROUP BY node, MIN(hops) in the outer query then collapses a node that was reached by two different paths (at possibly different hop counts) down to its shortest distance, the way node 1 itself, reachable from itself in zero hops, does not also appear again at two hops even though the raw recursion does produce that row internally.
fork and pipesuse Storable qw(freeze thaw);
sub ingest_parallel ($files, $workers = 4) {
my @chunks;
push @{ $chunks[$_ % $workers] }, $files->[$_] for 0 .. $#$files;
my %reader_for;
for my $w (0 .. $workers - 1) {
pipe(my $reader, my $writer) or die "pipe: $!";
my $pid = fork // die "fork: $!";
if ($pid == 0) {
close $reader;
my ($records, $problems) = (0, 0);
for my $file (@{ $chunks[$w] // [] }) {
my ($r, $p) = ingest_one_file($file); # each child: its own memory, own DB handle
$records += $r;
$problems += $p;
}
print $writer freeze({ worker => $w, records => $records, problems => $problems });
close $writer;
exit 0;
}
close $writer;
$reader_for{$pid} = $reader;
}
my ($total_records, $total_problems) = (0, 0);
for my $pid (keys %reader_for) {
local $/; # slurp mode for this read
my $data = readline($reader_for{$pid});
waitpid($pid, 0); # reap; see below
my $result = thaw($data);
$total_records += $result->{records};
$total_problems += $result->{problems};
}
return ($total_records, $total_problems);
}
Verified with six fake files split across three workers:
worker 0: 1000 records, 4 problems, files=a.log d.log
worker 1: 1000 records, 4 problems, files=b.log e.log
worker 2: 1000 records, 4 problems, files=c.log f.log
total: 3000 records, 12 problems
pipe creates a connected reader/writer pair before the fork, which is the only way it works: after fork, parent and child each have their own copies of both ends, and each side closes the end it does not use so that a read on an exhausted pipe correctly reports end-of-file rather than blocking forever waiting for a writer that will never write (because it is the child's own, unclosed, unused copy). Storable::freeze/thaw is Perl's standard way to send a structured value (here, a hashref) through something that only carries bytes.
Zombie processes. A child that exits before its parent calls waitpid becomes a zombie — gone, but still occupying a process table entry until reaped. This code reaps every child in the same loop that reads its result, so nothing is ever left unreaped; a longer-running supervisor would additionally want $SIG{CHLD} = 'IGNORE' or an explicit reaper loop for children whose results it does not need to wait for individually.
The classic pipe-buffer deadlock. A Linux pipe's kernel buffer defaults to 64 KB (/proc/sys/fs/pipe-max-size for the ceiling); if two processes both write more than that to each other and neither is reading, both block forever, each waiting for buffer space the other would free by reading. This design cannot hit that, structurally: communication is one-way, child writes and exits, parent only reads, so there is no cycle of mutual waiting to deadlock on. The moment you need bidirectional traffic — a supervisor sending a shutdown command to workers that are also streaming results back — you need either a second pipe per direction or a non-blocking read loop, and this is exactly the shape of bug that ambushes people who add "just one more message" to a fork-and-pipe design that was never built for it.
$ perl Makefile.PL
$ make test
PASS t/00-load.t
PASS t/10-record.t
...
PASS t/90-fuzz.t
All tests successful.
$ make dist
strata-0.12.tar.gz
$ cpanm --local-lib=/tmp/check strata-0.12.tar.gz && echo "installs cleanly"
installs cleanly
make dist reads MANIFEST (generated by make manifest from every file under version control) and produces exactly the tarball a real user would cpanm install — installing your own tarball into a scratch local::lib before calling a milestone "done" catches the class of bug that only "it works on my checkout" hides: a test file depending on a fixture that was never added to MANIFEST, a module that only compiles because something else in your working tree happened to load it first.
strata/
├── Makefile.PL, cpanfile, MANIFEST
├── bin/strata ingest | entities | correlate | query
├── lib/Strata.pm
├── lib/Strata/
│ ├── Util.pm Pattern.pm Record.pm Pipeline.pm Source.pm
│ ├── Extract.pm Normalize.pm
│ ├── Store.pm DBI, schema, transactions
│ ├── Correlate.pm sessionisation sweep
│ ├── Graph.pm edges, recursive CTE queries
│ └── Parser/
│ ├── Registry.pm Apache.pm JsonLines.pm Csv.pm Xml.pm Unstructured.pm
└── t/
├── lib/Test/Strata/ParserContract.pm
├── 00–80 (unit and integration tests, ~10 files)
└── 90-fuzz.t deadline-bounded mutation tests
$ prove -l t/
All tests successful.
$ git commit -am "milestones 9-12: correlation, a real CLI, fuzzing, profiling, the graph"
strata graph --format=dot > g.dot, walking edges and printing GraphViz syntax, so dot -Tpng g.dot -o g.png produces a picture of a correlated incident.--workers flag on strata ingest wired to ingest_parallel, benchmarked at 1, 2, 4 and 8 workers against a fixed set of files. Where does the speedup stop being linear, and what is SQLite doing at that point that a purely CPU-bound parallel task would not have to contend with?1. The exporter is a thin format layer over a query you already have:
say "graph {";
my $rows = $dbh->selectall_arrayref("SELECT event_a, event_b, reason FROM edges");
for my $row (@$rows) {
say qq{ "$row->[0]" -- "$row->[1]" [label="$row->[2]"];};
}
say "}";The only real care needed is escaping quotes inside reason if it is ever user-derived rather than one of the fixed strings this project generates.
2. Linear speedup holds until every worker is contending for the same SQLite file: fork-based parallelism gives each worker independent CPU and independent memory, but the database write at the end is still one file, and SQLite serialises writers even in WAL mode (WAL lets readers proceed concurrently with a writer, not multiple simultaneous writers). Past roughly the number of physical cores, and well before that if the workload is write-heavy rather than parse-heavy, you are measuring lock contention, not CPU parallelism — a purely CPU-bound task (say, only counting lines per file with no shared destination) would not hit this wall at all, which is precisely why profiling this specific pipeline mattered more than trusting the general reputation of fork.
3. Two same-direction pipes (parent→child and child→parent) rather than trying to reuse one: writing a large "stop" payload from the parent while the child is mid-write of its own large result, with neither side reading the other's pipe yet, reliably deadlocks once both messages exceed the 64 KB kernel buffer — you can reproduce it by having the parent print >64KB before its first read. The fix confirms the general rule from the warn box above: bidirectional traffic needs two one-way channels (or non-blocking I/O on one), never one pipe pressed into carrying both directions.
Call ingest_parallel with more workers than files — six files split across eight workers. Two workers get an empty chunk and report zero files back over their pipe; the other six each get exactly one. Nothing hangs and nothing dies: $chunks[$_ % $workers] // [] turns "no work assigned" into an empty list to iterate over instead of dereferencing undef, so an idle worker still writes a well-formed result, still gets reaped by waitpid, and every file is still accounted for exactly once in the total.
forked Perl worker processes, in a way it is not for two goroutines?pipe have to be called before, not after, and why does that ordering matter? UNION rather than UNION ALL matter for a recursive query over a graph that contains a cycle?local::lib catch that running the test suite in your working tree does not?fsync into thousands. my (...) = (...), letting the array silently swallow the whole list.EXPLAIN QUERY PLAN. alarm(0) on every exit path, leaving a timer armed to fire in unrelated code later.