Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 4 additions & 2 deletions .github/workflows/ci.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -81,8 +81,10 @@ jobs:

integration:
# Integration suite: ducklake.write() + SchemaManager against
# in-memory DuckDB (no containers, ~10s), then the hoglake sink
# against a REAL hoglake server on the throwaway compose stack.
# in-memory DuckDB (no containers, ~10s), the Postgres-only
# maintenance ops against a throwaway Postgres (runner's initdb,
# docker fallback), then the hoglake sink against a REAL hoglake
# server on the throwaway compose stack.
runs-on: ubuntu-latest
needs: check
env:
Expand Down
23 changes: 22 additions & 1 deletion AGENT.md
Original file line number Diff line number Diff line change
@@ -1,5 +1,24 @@
# Millpond — Dead Simple Kafka to DuckLake

## Public Repository

This repository is public. Everything that lands in it is visible to anyone:
code, comments, docs, commit messages, branch names, PR titles and
descriptions, and review comments.

Never include internal or customer-identifying information, for example:

- customer, organization, team or project IDs and UUIDs
- customer names, or details that let a reader identify a customer
- per-tenant database, bucket or catalog names
- internal hostnames, service endpoints, cluster or shard names, account IDs
- secrets, tokens or credentials, even expired ones
- links to internal dashboards, incident channels or private repositories

When a real incident motivates a change, describe it generically, for
example "a production tenant" or "a large catalog". Keep identifying
diagnostics in internal channels.

## Pre-push Checklist

Always run before pushing:
Expand Down Expand Up @@ -58,7 +77,9 @@ The catalog-side orphan-recovery subcommands (`dedup-deletions`, `find-orphans`,

`main()` logs `millpond <version> (maintenance)` on startup using `importlib.metadata.version("millpond")`, mirroring `millpond/main.py`. If `MILLPOND_IMAGE` is set in the env, the image identifier is appended (`image=<value>`). The chart-side wiring is optional — without it, the package version is sufficient to identify which build is running.

`drop-partitions` and `list-droppable-partitions` implement server-side partition retention (megaduck `events_nrt` 14-day campaign and the daily cron). They forge an engine whole-file-drop commit directly in the catalog (snapshot row with `next_file_id = baseline + 1` — the F1 stats-ObjectCache bust, per-row re-verified `end_snapshot` UPDATE, RELATIVE stats decrement from `UPDATE ... RETURNING` sums, `deleted_from_table:<id>` conflict token) because the engine's DELETE cannot complete against the nrt commit rate (insert-vs-delete OCC conflict is a hard abort). These two ops are the ONLY ones that bypass the duckdb-attach `postgres_execute` path — they use a direct psycopg (libpq) connection (READ COMMITTED, real RETURNING, no wedged-duckdb teardown core dump), so `main()` skips `connect()` for them (`_DIRECT_PG_COMMANDS`). Dry-run is the default; `--execute` writes. Safety posture: leaf = one fully-resolved partition tuple = one short transaction = one snapshot; per-leaf advisory-lock re-acquisition (`pg_advisory_unlock_all` + `pg_try_advisory_lock`, refcount-safe); fresh viaduck cursor guard per leaf read from the durable `viaduck.viaduck_state` store (fail-closed on missing/empty, run-abort on consumer lag past cutoff+48h, `--no-viaduck-guard` prints what it overrides); structural floor at selection (cutoff+48h); partition-spec pinning by live `partition_id` with fail-loud aborts on old-spec/unageable/rot files (remedy: `repair-partition-values`). The stats assertion refuses to take `record_count` to 0, so a table's FINAL leaf is undroppable by design (the engine has a distinct empty-table path we deliberately don't emulate). Full design and hazard list: `DROP_PARTITIONS_PLAN.md` and `partition-drop-findings.md`.
`drop-partitions` and `list-droppable-partitions` implement server-side partition retention (megaduck `events_nrt` 14-day campaign and the daily cron). They forge an engine whole-file-drop commit directly in the catalog (snapshot row with `next_file_id = baseline + 1` — the F1 stats-ObjectCache bust, per-row re-verified `end_snapshot` UPDATE, RELATIVE stats decrement from `UPDATE ... RETURNING` sums, `deleted_from_table:<id>` conflict token) because the engine's DELETE cannot complete against the nrt commit rate (insert-vs-delete OCC conflict is a hard abort). These two ops (and `drop-orphan-inline-tables`, below) are the ONLY ones that bypass the duckdb-attach `postgres_execute` path — they use a direct psycopg (libpq) connection (READ COMMITTED, real RETURNING, no wedged-duckdb teardown core dump), so `main()` skips `connect()` for them (`_DIRECT_PG_COMMANDS`). Dry-run is the default; `--execute` writes. Safety posture: leaf = one fully-resolved partition tuple = one short transaction = one snapshot; per-leaf advisory-lock re-acquisition (`pg_advisory_unlock_all` + `pg_try_advisory_lock`, refcount-safe); fresh viaduck cursor guard per leaf read from the durable `viaduck.viaduck_state` store (fail-closed on missing/empty, run-abort on consumer lag past cutoff+48h, `--no-viaduck-guard` prints what it overrides); structural floor at selection (cutoff+48h); partition-spec pinning by live `partition_id` with fail-loud aborts on old-spec/unageable/rot files (remedy: `repair-partition-values`). The stats assertion refuses to take `record_count` to 0, so a table's FINAL leaf is undroppable by design (the engine has a distinct empty-table path we deliberately don't emulate). Full design and hazard list: `DROP_PARTITIONS_PLAN.md` and `partition-drop-findings.md`.

`drop-orphan-inline-tables` drops the Postgres tables behind DuckLake data inlining (`public.ducklake_inlined_data_<table_id>_<schema_version>`, registered in `ducklake_inlined_data_tables`) whose DuckLake table no retained snapshot can reach, and deletes their registry rows. DuckLake's GC (`DropEmptySupersededInlinedTables`) only handles superseded-and-empty tables, so DROP+CREATE churn can leak one Postgres table per dropped DuckLake table; every leaked table adds rows to `pg_class` / `pg_attribute` / `pg_type`, and DuckDB's DuckLake ATTACH scans those system catalogs in full — at ~230k leaked tables the ATTACH itself OOMs. So this op is in `_DIRECT_PG_COMMANDS` and must never need the ATTACH. The orphan predicate (`_inline_orphan_predicate`) is the range-overlap predicate of the `ducklake_unreachable_inline_tables` metric; keep them in lockstep (an integration test runs the metric's SQL against the same catalog and asserts equal counts). Per batch (default 500, capped at 2000 because each DROP holds several shared lock-table entries until commit) ONE transaction selects the next keyset page of orphans (`FOR UPDATE OF idt SKIP LOCKED`), then per present table takes `LOCK TABLE ... IN ACCESS EXCLUSIVE MODE`, checks emptiness with `SELECT 1 ... LIMIT 1` (never `reltuples`, which is a stale estimate), skips and WARNs on non-empty tables and on names that do not fullmatch `ducklake_inlined_data_<n>_<n>`, `DROP TABLE IF EXISTS`, and deletes the matching registry rows. The predicate is re-evaluated per batch. Execute mode takes the maintenance advisory lock once (`_lock_refresh`). It never runs `VACUUM FULL` (ACCESS EXCLUSIVE on system catalogs); it logs a hint instead. Gauges `maintenance_inline_orphans_{dropped_total,skipped_nonempty_total,remaining}` are unlabeled and therefore registered only for an execute run of this command, so other crons' pushes cannot zero them. Tests: `tests/integration/test_drop_orphan_inline_tables.py` runs against a throwaway real Postgres (`MILLPOND_TEST_PG_DSN`, else local `initdb`/`pg_ctl`, else docker).

`cleanup` and `cleanup-all` log a single structured throughput line on completion: `cleanup throughput: files_processed=N elapsed_s=T rate_obj_s=R queue_depth_after=A`. `files_processed` comes directly from `len(result)` — the count of rows `ducklake_cleanup_old_files` returned — rather than a queue-depth delta. A delta would be misleading whenever any other writer enqueues deletions during the call (the maintenance advisory lock by design only mutexes maintenance invocations, not arbitrary writers); `len(result)` is accurate regardless. `queue_depth_after` is queried with a single post-call snapshot and shows remaining work but is not used in the rate calculation. The line is intentionally suppressed on `cleanup --dry-run` because dry-run returns preview rows and a rate computed from those would falsely claim work was done. `cleanup-all` has no dry-run form at all — the CLI rejects `--dry-run` rather than silently no-opping (preview via `cleanup-dry-run` / `fsck-dry-run`).

Expand Down
4 changes: 2 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -310,13 +310,13 @@ Requires Docker (uses `keytool` from the Kafka container image for cert generati

The `tools/` directory ships two DuckLake-only operational binaries inside the same image as the writer:

- **`tools/ducklake_maintenance.py`** — CLI for snapshot expiry (incl. Postgres-native `expire-snapshots`), file cleanup, orphan recovery, tiered compaction, fsck, and one-shot repairs (`repair-partition-values`, `dedup-deletions`, `purge-orphan-stats`). Runs as a K8s CronJob.
- **`tools/ducklake_maintenance.py`** — CLI for snapshot expiry (incl. Postgres-native `expire-snapshots`), file cleanup, orphan recovery, tiered compaction, fsck, and one-shot repairs (`repair-partition-values`, `dedup-deletions`, `purge-orphan-stats`, `drop-orphan-inline-tables`). Runs as a K8s CronJob.
- **`tools/ducklake_metrics.py`** — Catalog-side lake-state metrics, either as a long-running Prometheus-exposition daemon or in one-shot push mode (`--once`, POSTing to `DUCKLAKE_METRICS_PUSH_URL` — the per-tenant metrics CronJob path).

`tools/justfile` (copied to `/justfile` in the image) wraps both. Recipe groups:

- `interactive` — `shell`: a DuckDB shell wired to the DuckLake with the same session setup as the subcommands
- `lifecycle` — snapshot + file lifecycle: `expire`/`expire-snapshots` (+ chain-safe `expire-7d`), `cleanup`/`cleanup-all`/`cleanup-all-safe`, orphan handling (`find-orphans`, `heal-orphans`, `delete-orphaned-files`, `fsck`), `dedup-deletions`, `purge-orphan-stats`, `repair-partition-values`, `maintain`, `checkpoint`. Destructive recipes print the target catalog and, on a TTY, demand confirmation.
- `lifecycle` — snapshot + file lifecycle: `expire`/`expire-snapshots` (+ chain-safe `expire-7d`), `cleanup`/`cleanup-all`/`cleanup-all-safe`, orphan handling (`find-orphans`, `heal-orphans`, `delete-orphaned-files`, `fsck`), `dedup-deletions`, `purge-orphan-stats`, `drop-orphan-inline-tables` (+ chain-safe `drop-orphan-inline-tables-default`), `repair-partition-values`, `maintain`, `checkpoint`. Destructive recipes print the target catalog and, on a TTY, demand confirmation.
- `compaction` — tiered `compact-to-tier-{1,2,3}` (+ dry-runs), `compact-all-tiers`, `compact-probe`, and chain-safe no-arg wrappers (`compact-all-tiers-default`) so the tenant-maintenance CronJob can run one flat recipe chain
- `bootstrap` — `bootstrap-indexes`: the DuckLake catalog btrees (compaction scans, snapshot-range reads, per-file joins) via `psql`, all `CREATE INDEX CONCURRENTLY IF NOT EXISTS`
- `metrics` — `ducklake-metrics` daemon recipes and chain-safe `metrics-once` for the per-tenant metrics CronJob
Expand Down
30 changes: 19 additions & 11 deletions docker-compose.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -37,10 +37,11 @@ services:
- /var/lib/postgresql/data

minio:
# quay.io: Docker Hub stopped serving the minio images anonymously
# ("pull access denied ... repository does not exist"). The same
# release tags stay public on quay.io.
image: quay.io/minio/minio:latest
# cgr.dev/chainguard: Docker Hub dropped minio/minio and quay.io now
# refuses anonymous pulls of it (401). The Chainguard image is public,
# has the same /usr/bin/minio entrypoint, and ships mc for the
# healthcheck.
image: cgr.dev/chainguard/minio:latest
hostname: minio
ports:
- "127.0.0.1:9000:9000"
Expand All @@ -57,16 +58,23 @@ services:
retries: 5

minio-init:
image: quay.io/minio/mc:latest
# awscli creates the bucket. quay.io now refuses anonymous pulls of
# minio/mc, and the AWS image on public.ecr.aws stays public.
image: public.ecr.aws/aws-cli/aws-cli:latest
depends_on:
minio:
condition: service_healthy
entrypoint: >
/bin/sh -c "
mc alias set local http://minio:9000 minioadmin minioadmin &&
mc mb --ignore-existing local/ducklake &&
mc anonymous set download local/ducklake
"
environment:
AWS_ACCESS_KEY_ID: minioadmin
AWS_SECRET_ACCESS_KEY: minioadmin
AWS_DEFAULT_REGION: us-east-1
entrypoint: ["/bin/sh", "-c"]
command:
- |
set -e
aws --endpoint-url http://minio:9000 s3api head-bucket --bucket ducklake 2>/dev/null \
|| aws --endpoint-url http://minio:9000 s3api create-bucket --bucket ducklake
aws --endpoint-url http://minio:9000 s3api put-bucket-policy --bucket ducklake --policy '{"Version":"2012-10-17","Statement":[{"Effect":"Allow","Principal":{"AWS":["*"]},"Action":["s3:GetObject"],"Resource":["arn:aws:s3:::ducklake/*"]}]}'

# --- Topic Setup ---

Expand Down
6 changes: 4 additions & 2 deletions justfile
Original file line number Diff line number Diff line change
Expand Up @@ -48,8 +48,10 @@ lint-fix:
test:
uv run python -m pytest tests/unit

# Run integration tests (in-memory DuckDB — no docker; the hoglake
# suite is docker-gated and has its own recipe below)
# Run integration tests (in-memory DuckDB — no docker; the Postgres-only
# maintenance tests boot a throwaway Postgres via MILLPOND_TEST_PG_DSN,
# local initdb/pg_ctl, or docker, and skip without any of them; the
# hoglake suite is docker-gated and has its own recipe below)
[group('test')]
test-integration:
uv run python -m pytest tests/integration --ignore=tests/integration/test_hoglake_integration.py
Expand Down
5 changes: 3 additions & 2 deletions tests/hoglake_stack/docker-compose.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -31,8 +31,9 @@ services:
retries: 15

minio:
# quay.io: Docker Hub stopped serving minio/minio anonymously.
image: quay.io/minio/minio:latest
# cgr.dev/chainguard: quay.io now refuses anonymous pulls of
# minio/minio (401). Same entrypoint, and mc is in the image.
image: cgr.dev/chainguard/minio:latest
command: server /data
environment:
MINIO_ROOT_USER: hoglake
Expand Down
6 changes: 5 additions & 1 deletion tests/hoglake_stack/stack.py
Original file line number Diff line number Diff line change
Expand Up @@ -68,7 +68,11 @@ def compose(*args: str, profile: str | None = None, check: bool = True) -> subpr
if profile:
cmd += ["--profile", profile]
cmd += list(args)
return subprocess.run(cmd, check=check, capture_output=True, text=True, timeout=600)
# The compose file falls back to hoglake-server:latest. Hand it the
# pinned image explicitly, otherwise a new server release changes
# what CI tests against without any commit here.
env = {**os.environ, "HOGLAKE_SERVER_IMAGE": server_image()}
return subprocess.run(cmd, check=check, capture_output=True, text=True, timeout=600, env=env)


def docker_available() -> bool:
Expand Down
Loading
Loading