The internal processing pipeline of the karmamap binary, plus how
the output it produces is queried by the bundled web frontend. The Parquet
data contract — layout, schemas, encodings — lives in API.md.
A single binary, karmamap, with four pipeline stages — passes 1-3 run
under karmamap import, step 4 under karmamap prepare-update:
- Node pass: builds the mmap node cache
(node_id, day) -> positionand counts node changes intochanges/year=YYYY/nodes.parquet. - Way pass: resolves each way's node positions via the cache and counts
way changes at the distinct cells of those positions into
changes/year=YYYY/ways.parquet. - Merge pass: merges each year's
nodes.parquetandways.parquetcounts per(h3_cell, change_date)into the singlecountcolumn ofdata.parquet, sorted by(h3_cell, change_date)so row-group min/max support bbox and date-range pruning. Idempotent: an existingdata.parquetsupplies the merged total of any count whose staging file is already gone, so re-merging never zeroes it. - Step 4 (incremental cache): collapses the node cache into a second cache
holding only the last known h3 cell per node (day dropped), written to
<node-cache>.lastby default (see "Incremental cache" below). Step 4 is not part of import: runkarmamap prepare-updatefor it.
Passes 1-3 and the users-history pass run as part of karmamap import
(step 4 lives in karmamap prepare-update; see below).
Import records the snapshot's osmosis replication provenance in
manifest.json from a <base>.state.txt sidecar next to the osh, downloaded
manually with wget on the snapshot's day: upstream state.txt is always the
current state and would be too new for an older snapshot. prepare-update
records only the update stream URL, keeping the recorded sequence and
timestamp; update fetches the live state.txt from the update URL.
Pass 1 wipes changes/ (so a re-run never leaves stale partitions) and the
node cache file before rebuilding; pass 2 does not wipe the root — it adds
ways.parquet staging files next to nodes.parquet, so any pass can run
alone. The manifest is rebuilt at the end of every run.
| Case | Behavior |
|---|---|
| Node with valid coordinates | Written to the node cache + counted |
| Deleted node with a previously known position | Counted on the last known position, no cache write |
| Node with no coordinates and no previously known position | Skipped |
| Way node with an unresolved position | The node is skipped |
| Deleted way with a previously known geometry | Counted on the last known geometry |
| Deleted way with no previously known geometry | Skipped |
| Visible way with no nodes | Skipped |
| Relations | Out of scope for the change-counting passes; the users-history pass counts relation created/modified/deleted in a day's activity (relations feed the ranking only via creations) |
| Node cells of a way | Each distinct node cell counted once per way version |
| Time zone | Strict UTC |
| Source file ordering | Assumed sorted by (id, version) ascending, as documented for OSM full-history files |
The pipeline uses a single H3 resolution:
--h3-resolution(default 9): the resolution of the cells stored in the Parqueth3_cellcolumn and used across all passes. Range 0-13 (the node cache packs at most 13 H3 digits into 6 bytes).
Rows are partitioned by the calendar year of change_date, not by H3
cell, so the number of concurrently open Parquet writers is bounded by the
number of years that contain data — independent of extract size or H3
resolution.
The --node-cache argument points at a single file that pass 1 builds and
pass 2 reads (pass 1 deletes and recreates it):
record : [node_id 8B][day 2B][h3 cell 6B] (16 bytes)
file : [header 40B][block 0]...[block N-1][directory 12*N]
header : magic "OSNC", version, record_size, h3_resolution,
record count, block count, records per block, compression
Records are sorted by (node_id, day) ascending. node_id is stored
big-endian with the sign bit flipped so byte order equals numeric order;
day is the same uint16 UTC epoch-day value as change_date; the cell is
the node's H3 index at --h3-resolution (0-13) packed into 6 little-endian
bytes (h3_utils::pack_cell), with the resolution re-applied from the
header on read. The writer
enforces the ascending order while building (a decreasing node id or day
aborts loudly), so the way pass's sweep is safe. A node edited several times
in one day is stored once (last position wins).
Node deletions are not stored: resolving a deleted node returns its last
known position — an accepted approximation.
Records are grouped into blocks of 2¹⁸ (4 MiB raw) that are ZSTD-compressed
on write and appended as-is; the block and record counts are patched into
the header at finish, and a trailing directory holds each block's first
node and compressed size (offsets cumulative, recomputed in RAM at open).
The compression/records-per-block header fields act as a format stamp:
a cache written by another record layout, or a truncated file, aborts with
rebuild instructions. Worst case (incompressible data) ZSTD stores a block
raw, so the file is never meaningfully larger than the uncompressed
layout; the cost is one decompression per block on first access of each
pass.
The reader mmaps the file read-only and decompresses one block at a time into a 4 MiB cache. The per-block first keys seed the way pass's sweep cursor between batches (a binary search over the directory, then one forward-only scan); no RAM sample index is kept.
Step 4 derives a second, history-less cache from the node cache. The node
cache holds every (node_id, day) version, so its last record per node_id is
that node's last known position; the incremental cache keeps exactly one
14-byte record per node — [node_id 8B][h3 cell 6B], the day column dropped,
node_id -> h3_cell. Same block/directory scheme as the node cache
(a "INCC" magic, 2¹⁸ records per ZSTD block, trailing directory of first
key + compressed size), built by one streaming sweep of the node cache fed
into a collapsing writer.
Because the record layout differs from the node cache, blocks are never
reused across runs: the writer appends fresh blocks to <path>.tmp and swaps
it over the final path with a rename once the header and directory are
finalized, so a crash leaves either the old cache or only the tmp file,
never a torn cache. The --node-cache-last flag sets the output path
(default <node-cache>.last); the build is the whole of
karmamap prepare-update, import never runs it.
A reader (node_cache::incremental::Reader) mmaps the file and looks up a
node's last known cell via the block-directory binary search plus a linear
scan (returning 0 for nodes absent from the cache).
karmamap update applies osmosis replication diffs to an existing dataset
whose manifest.json recorded an update URL (an --update-url passed to
karmamap import, or the one recorded by karmamap prepare-update, which
never fetches state.txt). The
update stream is that recorded source URL unless --update-url is given
explicitly (which must then match), and update runs against the .last
incremental cache built by prepare-update. The starting
sequence is the recorded source sequence; each
.osc.gz diff (URL AAA/BBB/CCC.osc.gz, where N = AAA1000000 + BBB1000 +
CCC from the 3/3/3 split of the sequence) is downloaded to the diffs dir next
to the node caches — a file
already present is reused, a partial download is removed on failure — and
applied. Once every pass over a diff succeeded its file is removed, so
the diffs dir only ever holds in-flight diffs; diffs committed by
earlier runs are purged at update start. Bare update (or update 0) fetches
every diff up to the current
state.txt; update N stops after N.
Each diff invokes update-mode passes 1 and 2:
- Update node pass: counts the diff's node changes into
changes/year=YYYY/nodes.<seq>.parquetand folds created/modified positions into an in-memory overlay over the flat incremental cache (--node-cache-last, default<node-cache>.last), tracking deletions separately. - Update way pass: counts the diff's way changes into
changes/year=YYYY/ways.<seq>.parquet, resolving node refs against the overlay — post-update view for visible ways (overlay else base, deleted resolves to nothing), pre-update view for deleted ways (overlay else base = their last known geometry). Like the import way pass, refs are resolved in batches (bounded by--way-batch-mb's default) with one forward-only sweep over the base cache per batch instead of a random lookup per ref, so a diff's way pass decompresses only the cache blocks its refs touch.
Once per run the .last incremental cache is rebuilt as base + overlay minus
deletions (the incremental writer's tmp+rename swap keeps it consistent), and
the staged per-sequence counts are folded into each year's data.parquet by
merge_update_partitions — base data.parquet plus its nodes.<seq> /
ways.<seq> staging per year, one rename per year. Fetching N diffs never
rewrites the dataset N times.
Apply-once semantics: the merged data.parquet footer carries a
karmamap_source_seq key stamped with the highest sequence folded in. A
partition whose stamp is already >= the applied sequence drops any orphaned
staging and leaves the merged data untouched, so a crash after the merge
rename (or a rerun of the same diff) is a no-op. Staging carrying a sequence
later than the applied one (a crashed run that fetched further than a capped
rerun applies) is dropped rather than folded. Years holding no staging are
never rewritten.
Each update diff additionally runs the users-history scan
into users_history_update_stage/seq_<n>/; one finalize pass folds the
per-(uid, change_date) activity deltas and per-uid counter totals into
users_history.parquet (in place, sorted) and rebuilds
user_ranking.parquet from the combined existing + delta totals. The
manifest's source block is updated to reflect the highest applied sequence
and its timestamp once per run.
That merge is by file name and schema: the finalize reads the previous
users_history.parquet, user_ranking.parquet (the base of the recomputed
ranking) and suspect.parquet (the far_move_count, max_move_meters and
ranking_at_day columns, the base of the carried and frozen values) as its
three bases, and reads all of them
before it writes the first file of the run, so a base that does not match the
current schema aborts the update with nothing replaced. Every added column is
therefore a schema change: an output directory written by an older build is not
updatable, and re-importing the snapshot into a fresh output directory is the
migration. An output directory not written by this build is likewise not
updatable.
The suspect_minutes.bin minute store is the one derived file with no schema
check: a missing store is read as an empty base, so an update against a
directory whose store was written under another name silently restarts
filter-2 accumulation at that point instead of failing. The store is
self-identifying (a VMIN magic and version, no file name), so renaming it in
place preserves the folded minutes and the applied-sequence stamp; do that
before updating, or the period before the upgrade stops contributing to the
filter-2 day bits.
The suspect_tags.bin tag store is likewise self-identifying (magic TAGS,
variable-length records, no file name), and follows the same pattern.
The suspect engine (OSMPatrol filters 2 and 3, see
docs/osmpatrol-neis-2012.md) runs in the same loop. Every diff is scanned
into per-(uid, minute) modified+deleted buckets under
suspect_update_stage/counts/seq_<n>/, and modified nodes with a known
prior position are recorded by the update node pass into
suspect_update_stage/moves/seq_<n>/ as (uid, minute) rows.
The finalize three-step ordering is load-bearing:
fold_minute_countsmerges the run's staged buckets into the persisted binary storesuspect_minutes.bin(suspect_store.hpp), summing equal(uid, minute)keys — the store is the merge base for the next run and is stamped with the applied sequence so a rerun is a no-op.users_history::run_update_finalizereads it throughsuspect::flagged_daysplus this run'ssuspect::flagged_move_daysand recomputes thesuspect_flagcolumn ofusers_history.parquet(bits 0/1 base flags carried forward, ORed with this run's filter-2 and filter-3 bits; bit 2, the ranking-based filter 1, is forward-only — set only on the rows this run newly writes, see the users-history pass below). It also writessuspect.parquet, one row per flagged day with the day's total change count, the peak value each store-derived filter fired on, its far-move count with the largest of them and the ranking frozen at the day's first flag.suspect::flagged_move_daysfolds the staged node moves (> 500 m) into those same per-day bits (filter 3), a per-day far-move count and the largest distance among them; its stage is transient and removed, so every finalize is idempotent.suspect_minutes.binlikewise folds minutes (filter 2). There is no persisted move dataset beyond the carriedfar_move_count/max_move_meterscolumns ofsuspect.parquet— both filters land inusers_history.parquet'ssuspect_flagbits and only the flagged days are re-exported.
The users-history pass scores history per user and per UTC
day with cheap OSMPatrol-style heuristics. It is a single streaming scan
over the node and way history plus relation creations (no changeset metadata
is needed) and one in-memory finalize. OSM full history is
(id, version)-sorted, so the scan is a running pass with O(1) object
state, writing day-aggregates to a staged
users_history_stage/stage_*.parquet directory that finalize merges, sorts
by (uid, change_date), derives the ranking rows, and removes.
user_ranking.parquet is a pure derived view of that data: one row per
user, so it grows with new users, not new edits, and can be rebuilt from the
per-user totals without re-reading history. The non-partitioned single files
keep the uid join cheap and the numerics-only history file small.
The suspect_flag of each history row is a bit field. Bits 0 and 1
(kFlagFilter2, kFlagFilter3) stay 0 on import and are monotonically ORed
by every update finalize from the persisted minute store and the run's
move-flagged days. Bit 2 (kFlagFilter1, the paper's "new users or
ranking < 5%" filter) differs: it is not diff-based, so import writes 0 and
every update finalize sets it only on the rows it newly writes, from the
ranking it just recomputed. Base rows are carried unchanged, so a contributor
who later climbs above the threshold keeps the bit on the days already
flagged. The update finalize derives it from the same ranking::Result it
writes to user_ranking.parquet; a contributor who created nothing ranks 0,
which covers the "new users" half of the screen without a separate rule.
Bit 3 (kFlagFilter4) is a local extension, not from the original OSMPatrol
paper: it flags a day when a user's edits in a trailing 1h window span
>= 3 distinct H3 cells whose combined surface area (cell count × the
resolution's average hexagon area) is >= 20 km², gated on >= 20 true
edits in the window (the true edit count comes from suspect_minutes.bin,
so a long way counts as one edit but contributes all its cells to the spread).
suspect_cells.bin is the update-only persistent store behind this flag.
Bit 4 (kFlagFilter5) is a local extension: a trailing 1-hour window whose
modified/deleted object count is >= 100 and where one tag key covers > 90%
of those objects. The per-(uid, minute, tag_key) modified+deleted counts are
persisted in suspect_tags.bin (version 1).
The suspect engine (OSMPatrol filters 2 and 3 of Neis, Goetz & Zipf 2012)
watches the diff stream, not the full history — it runs update-only. The
diff scan (suspect::run_scan_diff) reuses the users-history
classification (visible version 1 = created, later = modified,
invisible = deleted) and counts modified + deleted objects per
(uid, minute) into suspect_update_stage/counts/seq_<n>/ (UTC minutes
since the epoch; creates are ignored).
The minute buckets are persisted as the binary block store
suspect_minutes.bin next to the node caches (see suspect_store.hpp for
the on-disk format: one 16-byte record per (uid, minute), sorted, keyed as
the update finalize's merge base). fold_minute_counts reads the store plus
the run's staged buckets, sums equal (uid, minute) keys (so a minute that
gains edits in an incoming diff amends its record), rewrites the store with a
tmp+rename swap and stamps its header with the applied sequence. A rerun of
an already-folded sequence (crash between the rename and stage cleanup) is
skipped.
suspect::flagged_days reads the store back into
(uid, day) -> (flag, max_edits_per_hour): each uid's contiguous minute
series is run through hour_spans, the trailing 60-minute window (the
minute's count plus the previous 59), and any day holding a minute whose
span is strictly above 500 is flagged, carrying the largest such span of the
day as its max_edits_per_hour metric (day = minute / 1440, so a burst
crossing midnight still lands on the day of its peak minute). The
users-history update finalize merges those flags into users_history.parquet's
suspect_flag column for the whole history. Import writes the column as 0: the
full-history scan precedes the replication stream, so no minute buckets exist
for it.
Filter 3 records modified-node moves (> 500 m) as the same updates apply:
the update node pass already tracks each node's last known H3 cell
(node_cache::incremental base plus this run's overlay), and it hands every
detected move beyond the 500 m screen to suspect::NodeMoveSink, which
stages (uid, minute) rows under
suspect_update_stage/moves/seq_<n>/ (rows that do not clear the screen
are dropped before staging). flagged_move_days folds them straight into
this run's (uid, day) -> filter-3 bit flags, per-day far-move count and
the largest distance among them, used by the users-history finalize, and
removes the stage root, so there is no persisted move dataset — filter 3
survives only as the carried/ORed day bit in users_history.parquet and the
frozen per-day far_move_count/max_move_meters columns of
suspect.parquet. The distance is measured from the prior cell
center to the new point, within one res-9 cell radius (~175 m) of the true
prior: fine for the 500 m screen, not for the paper's finer 11 m
edit-analysis flag. Filter 1 (new users / ranking < 5%) is a kFlagFilter1
bit that the users-history update finalize sets forward-only on the rows it
newly writes, from the run's recomputed ranking, instead of a diff-based
screen; base rows are never re-flagged, so the bit is monotonic.
Filter 4 (local extension, not from the original OSMPatrol paper) records
the spatial spread of H3 cells in a trailing 1h window. The diff scan
(suspect::run_scan_diff_cells) counts modified+deleted objects per
(uid, minute, h3_cell): nodes in their cell; ways once per distinct cell of
their nodes (resolved through NodeState so the geometry is correct even for
deleted ways); relations skipped. The per-diff buckets are staged under
suspect_update_stage/cells/seq_<n>/, folded by fold_cell_counts into the
persisted binary block store suspect_cells.bin (same block format as
suspect_minutes.bin but 24-byte records keyed by (uid, minute, h3_cell),
with an H3 resolution stamp in the header). suspect::flagged_cell_days
reads both suspect_cells.bin and suspect_minutes.bin in lockstep: for each
uid, it slides a 60-minute window over both streams simultaneously — the cell
stream gives the distinct cell count in the window, the minute stream gives the
true edit count for the gate — and flags the day when the window has >= 20
edits, >= 3 distinct cells, and the surface area distinct_cells * cell_area_km2(resolution) >= 20 km². The area of the largest flagging window
is kept as that day's metric. The flag lands in bit 3 of
suspect_flag and is monotonic like bits 0/1. Import writes it as 0 (no cell
buckets exist for the full-history scan). Expressing the budget as an area means
the constant 20 stays in km² when --h3-resolution changes, so it does not
need retuning; the corresponding number of distinct cells still scales with
cell size (about 190 cells at the default resolution 9, 1 at resolution 6), and
the metric counts cells touched rather than measuring the distance between them.
Since distinct_cells * cell_area is monotone in distinct_cells, exactly one
of the kFilter4MinCells and area gates binds at any given resolution.
Filter 5 (local extension, not from the original OSMPatrol paper) flags a
day when a user's modified/deleted objects in a trailing 1-hour window are
dominated by a single tag key. The diff scan (suspect::run_scan_diff_tags)
counts modified+deleted objects per (uid, minute, tag_key) on nodes,
ways, and relations (creates are excluded, matching the scope of
suspect_minutes.bin which provides the denominator). The per-diff buckets are
staged under suspect_update_stage/tags/seq_<n>/, folded by
fold_tag_counts into the persisted binary block store suspect_tags.bin
(same block format but variable-length records keyed by (uid, minute, tag_key),
max 272 bytes, version 1). suspect::flagged_tag_days reads both
suspect_tags.bin and suspect_minutes.bin in lockstep: for each uid, it slides
a 60-minute window over both streams simultaneously — the tag stream gives the
per-key counts, the minute stream gives the total modified+deleted count for the
gate — and flags the day when the window has >= 100 objects and one tag key
covers > 90% of them, keeping the sorted, de-duplicated set of keys that
cleared the ratio as that day's metric (a window may have more than one, each
just over the line). The flag lands in bit 4 of suspect_flag (0x10) and is
monotonic like the other bits. Import writes it as 0 (no tag buckets exist for
the full-history scan). Limitation: a create-only mass import (e.g. a bulk
upload of new buildings all tagged building=yes) is not flagged, because
creates are excluded from both numerator and denominator to share scope with the
minute store.
web/changes/query.js queries bbox + date range. The date range selects the
year partitions (intersected with the manifest's partition list); each
distinct file is queried once via parquetQuery, which prunes row groups on
h3_cell and change_date, then aggregates the merged count
client-side per cell and per day. The manifest's date_range (read from the
data.parquet footer stats) bounds the date pickers and histogram axis to the
exact days that hold data. The non-contiguous H3 cell set of the
viewport bbox is applied as a coarse [min, max] range filter first, then
exact membership is checked client-side. web/changes/app.js triggers a new
query on pan/zoom (debounced) and resolves its data root relative to the
page URL (../data, see README, "Serving the web frontend").
web/users/query.js looks a user up by exact username on
user_ranking.parquet (the pipeline stamps the current username per uid,
and the file is username-sorted with a uid tie-break, so the exact filter
prunes straight to the matching pages). The ranking, identity fields and
per-history totals all come from that same user row; the dataset-wide
active/max aspect stats are read once from the file's
key_value_metadata footer instead of repeated per-row columns, and the
per-aspect points are recomputed from the stored pct and the constant paper
caps — no percentile math runs in the browser.
Only the per-day activity timeline is then read from
users_history.parquet: a uid [min, max] range filter prunes the
uid-sorted file to the pages holding that user, with exact membership kept
client-side. Each day's count column already totals the six node/way change
counters plus the three relation counters; the per-uid edit total stays
node/way-only (relations are ranking-only).
web/suspect/query.js reads the 100 latest flagged days from
suspect.parquet. Because the update finalize writes that file newest-first
(change_date descending), the page issues a plain rowEnd-capped query (no
date filter): hyparquet fetches only the leading row groups' pages, so the
"100 last" table never scans the whole file. Each row is one flagged
(uid, change_date) with the username, the combined suspect_flag bits,
the day's total change count, the peak edit count, spread and tag keys of
whichever store-derived filter fired, its far-move count with the largest of
them and the ranking frozen at the day's first flag. Of the flags the table
shows one: a day carrying bit 0 (kFlagFilter2, the >500 modified/deleted
objects in one hour filter) is marked as an edit burst.