Skip to content
Draft
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
31 changes: 31 additions & 0 deletions AGENT.md
Original file line number Diff line number Diff line change
Expand Up @@ -482,6 +482,37 @@ there would break that gate on every build.
therefore put the subclass arm FIRST — `ErrorMapping`,
`Metrics.commitFailureResult` and `Audit.failureOutcome` all do, and
each says so.
13. **Table file format is immutable and homogeneous per table**, stored
authoritatively on `hog_table.file_format` (the single source of truth;
`write.format.default` is derived for clients at read time).
- Parquet is the default format across all tables.
- `clickhouse-mergetree-packed` is admitted for fixed-schema,
unpartitioned, unsorted append-only tables when the rollout gate
`HOGLAKE_PACKED_MERGETREE_ENABLED=true` is enabled.
- Packed tables do not support schema evolution (columns, partition specs,
or sort orders cannot be added, dropped, renamed, promoted, or commented),
deletion vectors, explicit row IDs, truncation, or Parquet compaction.
- Compaction and maintenance planners select Parquet files only. The
maintenance debt sampler excludes non-Parquet tables so packed tables
never leak permanent small-file debt into debt scores or metrics.
- The server never opens a packed part. Registration validates packed
files from metadata only (non-zero rows and bytes, no `footer_size`,
no `split_offsets`, file format equal to the table format). The two
server surfaces that DO open Parquet — the stats hydrator and
compaction (invariant 2) — select `file_format = 'parquet'` only, so a
packed object can never reach a Parquet reader server-side.
- Packed stats are counts only: `column_stats` is optional and a packed
row is always `stats_state = 'provided'` (the hydrator never claims
it, so `pending` would be permanent). Packed files carry positional
row-id ranges like any other file and never `explicit_row_ids`.
- Every client and engine must refuse a format it cannot read, typed,
before handing the object to a reader: DuckDB (`duckdb-client/`),
Hedgerow and Trino (`plugin/trino-hoglake`) refuse packed tables and
per-file non-Parquet formats in scan plans; `pyhoglake.packed` is the
only reader.
- **No rollback past V26 once a packed table exists**: pre-V26 binaries
lack format filters in their compaction candidate queries and will fail
if run against a database containing packed tables.

## Scale doctrine (read before touching a query, a loop or a lock)

Expand Down
16 changes: 16 additions & 0 deletions ci/clickhouse-local.sh
Original file line number Diff line number Diff line change
@@ -0,0 +1,16 @@
#!/usr/bin/env bash
# The pinned ClickHouse that pyhoglake's packed adapter is tested against,
# exposed as a `clickhouse` executable (PYHOGLAKE_CLICKHOUSE). Packed parts
# are version-coupled bytes, so writer and reader are pinned by digest.
# The adapter only uses temporary directories, which the container sees at
# the same path through the TMPDIR bind mount.
set -euo pipefail

image=clickhouse/clickhouse-server:26.9.8.3@sha256:230b973a00b5bac5b8925bd2898fae95a039aa6f25063fa8c507ffe9e2c7d770
tmp=${TMPDIR:-/tmp}
exec docker run --rm -i --network none \
--user "$(id -u):$(id -g)" \
--volume "$tmp:$tmp" \
--env TMPDIR="$tmp" \
--entrypoint clickhouse \
"$image" "$@"
10 changes: 10 additions & 0 deletions ci/live-python.sh
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,15 @@ ci_python=(uv run --no-project --with defusedxml==0.7.1 python)
./gradlew --no-daemon :installDist -x test -x ktlintCheck --console=plain
) 2>&1 | tee "$report_dir/build.log"
"${compose[@]}" up -d --wait --wait-timeout 90 postgres minio
# The packed MergeTree tests run real ClickHouse. Pull the pinned image up
# front so its pull time does not count against the adapter's timeouts.
clickhouse_image=$(sed -n 's/^image=//p' "$repo_dir/ci/clickhouse-local.sh")
for ((attempt = 1; attempt <= 3; attempt++)); do
docker pull --quiet "$clickhouse_image" && break
[[ $attempt -lt 3 ]] || exit 1
sleep $((attempt * 5))
done
export PYHOGLAKE_CLICKHOUSE="$repo_dir/ci/clickhouse-local.sh"

# The deferred-stats integration test requires the hydrator loop.
env HOGLAKE_JDBC_URL="jdbc:postgresql://localhost:$HOGLAKE_PG_PORT/hoglake" \
Expand All @@ -56,6 +65,7 @@ env HOGLAKE_JDBC_URL="jdbc:postgresql://localhost:$HOGLAKE_PG_PORT/hoglake" \
HOGLAKE_CLEANUP_INTERVAL_MS=0 HOGLAKE_COMPACTION_INTERVAL_MS=0 \
HOGLAKE_RETIREMENT_INTERVAL_MS=0 HOGLAKE_REINDEX_INTERVAL_MS=0 \
HOGLAKE_METRICS_INTERVAL_MS=0 HOGLAKE_MAINTENANCE_SUMMARY_INTERVAL_MS=0 \
HOGLAKE_PACKED_MERGETREE_ENABLED=true \
"$repo_dir/server/build/install/hoglake-server/bin/hoglake-server" \
> "$report_dir/server.log" 2>&1 &
server_pid=$!
Expand Down
7 changes: 7 additions & 0 deletions duckdb-client/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -112,3 +112,10 @@ They cover DDL between file preparation and commit, table recreation,
eager DDL, multi-table atomicity, and retries that preserve the request.
Run against a server with `HOGLAKE_REFUSE_BLIND_PARTITIONED_APPENDS=true`
to verify that partitioned INSERTs satisfy the strict server setting.

## Data formats

The extension reads and writes Parquet tables only.
It refuses tables with `write.format.default=clickhouse-mergetree-packed`, including empty tables, before scanning or writing objects.
Scan plans also validate each file format before passing paths to the Parquet reader.
Use the Python ClickHouse packed adapter for these tables.
1 change: 1 addition & 0 deletions duckdb-client/src/include/common/hoglake_wire.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -64,6 +64,7 @@ struct HoglakeSortSpec {

struct HoglakeTableInfo {
string name;
string file_format = "parquet";
string namespace_name;
string table_uuid;
vector<HoglakeColumn> columns;
Expand Down
7 changes: 7 additions & 0 deletions duckdb-client/src/rest/hoglake_api_client.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -187,6 +187,13 @@ HoglakeTableInfo ParseTableInfo(yyjson_val *obj) {
info.name = GetString(obj, "name");
info.namespace_name = GetString(obj, "namespace");
info.table_uuid = GetString(obj, "table_uuid");
auto properties = yyjson_obj_get(obj, "properties");
if (properties && !yyjson_is_null(properties)) {
if (!yyjson_is_obj(properties)) {
throw InvalidInputException("hoglake: expected object for table properties");
}
info.file_format = GetString(properties, "write.format.default", "parquet");
}
// Table-schema numerics were the one struct the R3 sweep missed:
// record_count feeds NumericCast in GetStorageInfo.
//
Expand Down
4 changes: 4 additions & 0 deletions duckdb-client/src/storage/hoglake_insert.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -227,6 +227,10 @@ PhysicalOperator &HoglakeInsert::PlanInsert(ClientContext &context, PhysicalPlan
plan_transaction.RequireDMLAllowed(table.ParentSchema().name.GetIdentifierName(), table.GetWireInfo().name,
false /* is_delete */);
auto &wire = table.GetWireInfo();
if (wire.file_format != "parquet") {
throw NotImplementedException("hoglake: DuckDB cannot write data format '%s'; use a ClickHouse writer",
wire.file_format);
}

auto columns = wire.columns;
std::sort(columns.begin(), columns.end(),
Expand Down
6 changes: 6 additions & 0 deletions duckdb-client/src/storage/hoglake_multi_file_list.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -189,6 +189,12 @@ void HoglakeMultiFileList::LoadFileList() const {
auto &table = read_info.table;
auto ns = table.ParentSchema().name.GetIdentifierName();
files = transaction->Api().PlanScan(ns, read_info.table_name, read_info.travel);
for (const auto &file : files) {
if (file.data_file.file_format != "parquet") {
throw NotImplementedException("hoglake: DuckDB cannot read data format '%s'; use a ClickHouse reader",
file.data_file.file_format);
}
}
read_file_list = true;
}

Expand Down
4 changes: 4 additions & 0 deletions duckdb-client/src/storage/hoglake_table_entry.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,10 @@ unique_ptr<BaseStatistics> HoglakeTableEntry::GetStatistics(ClientContext &conte
}

TableFunction HoglakeTableEntry::GetScanFunction(ClientContext &context, unique_ptr<FunctionData> &bind_data) {
if (table_info.file_format != "parquet") {
throw NotImplementedException("hoglake: DuckDB cannot read data format '%s'; use a ClickHouse reader",
table_info.file_format);
}
auto function = HoglakeFunctions::GetHoglakeScanFunction(*context.db);
auto &transaction = HoglakeTransaction::Get(context, ParentCatalog());
function.function_info = HoglakeFunctionInfo::Create(*this, transaction);
Expand Down
73 changes: 73 additions & 0 deletions duckdb-client/test/fixtures/packed_format_fixture.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,73 @@
"""Create the two packed-format tables used by hoglake_packed_format.test."""

import os
import uuid

import pyarrow as pa

from pyhoglake import AlreadyExistsError, HoglakeClient, NotFoundError, S3Config

HOGLAKE_URL = os.environ.get("HOGLAKE_URL", "http://localhost:8080")
S3_ENDPOINT = os.environ.get("HOGLAKE_S3_ENDPOINT", "http://localhost:19000")
S3_ACCESS_KEY = os.environ.get("HOGLAKE_S3_ACCESS_KEY", "hoglake")
S3_SECRET_KEY = os.environ.get("HOGLAKE_S3_SECRET_KEY", "hoglake123")
CATALOG = "duckext-packed"
DATA_PATH = "s3://duckext-itest/packed/"
FORMAT = "clickhouse-mergetree-packed"


def main() -> None:
s3 = S3Config(
access_key=S3_ACCESS_KEY,
secret_key=S3_SECRET_KEY,
endpoint_override=S3_ENDPOINT,
region="us-east-1",
allow_bucket_creation=True,
)
with HoglakeClient(HOGLAKE_URL, s3=s3) as client:
s3.filesystem().create_dir("duckext-itest")
try:
catalog = client.catalog(CATALOG)
except NotFoundError:
catalog = client.create_catalog(CATALOG, DATA_PATH)
try:
namespace = catalog.create_namespace("ns1")
except AlreadyExistsError:
namespace = catalog.namespace("ns1")
for name in ("packed_empty", "packed_registered"):
try:
namespace.table(name).drop()
except NotFoundError:
pass
table = namespace.create_table(
name,
pa.schema([pa.field("id", pa.int64(), nullable=False)]),
properties={"write.format.default": FORMAT},
)
if name == "packed_registered":
catalog.commit_prepared(
{
"idempotency_key": str(uuid.uuid4()),
"read_snapshot": catalog.refresh().head_snapshot_id,
"appends": [
{
"namespace": "ns1",
"table": name,
"expected_table_uuid": table.table_uuid,
"files": [
{
"path": DATA_PATH + "data.packed",
"file_format": FORMAT,
"record_count": 1,
"file_size_bytes": 1,
"column_stats": [],
}
],
}
],
}
)


if __name__ == "__main__":
main()
2 changes: 2 additions & 0 deletions duckdb-client/test/run-live-tests.sh
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,8 @@ export DUCKEXT_S3_ENDPOINT="${DUCKEXT_S3_ENDPOINT:-localhost:9000}"
PYHOGLAKE_DIR="${PYHOGLAKE_DIR:-$HOME/src/hoglake/pyhoglake}"

echo "== fixtures (pyhoglake) =="
(cd "$PYHOGLAKE_DIR" && HOGLAKE_S3_ENDPOINT="http://$DUCKEXT_S3_ENDPOINT" \
uv run python "$HERE/test/fixtures/packed_format_fixture.py")
(cd "$PYHOGLAKE_DIR" && HOGLAKE_S3_ENDPOINT="http://$DUCKEXT_S3_ENDPOINT" \
uv run python "$HERE/test/fixtures/read_fixture.py")

Expand Down
38 changes: 38 additions & 0 deletions duckdb-client/test/sql/hoglake_packed_format.test
Original file line number Diff line number Diff line change
@@ -0,0 +1,38 @@
# Packed MergeTree is explicit unsupported input for this Parquet-only extension.
require-env HOGLAKE_URL

require hoglake

require parquet

require httpfs

require-env DUCKEXT_S3_ENDPOINT

statement ok
CREATE SECRET duckext_minio (TYPE s3, KEY_ID 'hoglake', SECRET 'hoglake123',
ENDPOINT '{DUCKEXT_S3_ENDPOINT}', USE_SSL false,
URL_STYLE 'path')

statement ok
ATTACH 'hoglake:duckext-packed' AS lake (ENDPOINT '{HOGLAKE_URL}')

statement error
SELECT * FROM lake.ns1.packed_empty
----
<REGEX>:.*DuckDB cannot read data format 'clickhouse-mergetree-packed'.*

statement error
INSERT INTO lake.ns1.packed_empty VALUES (1)
----
<REGEX>:.*DuckDB cannot write data format 'clickhouse-mergetree-packed'.*

statement error
SELECT count(*) FROM lake.ns1.packed_registered
----
<REGEX>:.*DuckDB cannot read data format 'clickhouse-mergetree-packed'.*

query I
SELECT 42
----
42
9 changes: 8 additions & 1 deletion hedgerow/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -152,7 +152,7 @@ for this table. A nonzero lag can be entirely other tables' commits
(the next cycle drains it as an empty window); use it as a
staleness/liveness signal, not a volume estimate.

## Runbook: the seven halt conditions
## Runbook: the eight halt conditions

hedgerow HALTS (exits nonzero, no retry) when continuing would be wrong.
A supervisor must NOT blindly restart these — the same condition will
Expand Down Expand Up @@ -260,6 +260,13 @@ the last error), expect up to `1 + max_window_replays` copies of the
window's rows in the destination worst-case, and restart — the window
replays once more from the committed offset.

### 8. Unsupported data format (exit 10)

Both replication modes require Parquet source and destination tables.
Packed MergeTree tables (`write.format.default=clickhouse-mergetree-packed`) are refused, including empty tables.
Each change window is checked before any file is read, rows are appended, or offsets are advanced.
Use a Parquet destination or the Python ClickHouse packed adapter instead.

## Development

```sh
Expand Down
2 changes: 2 additions & 0 deletions hedgerow/src/hedgerow/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@
PersistentFailureError,
SchemaMismatchError,
SplitBrainError,
UnsupportedFormatError,
)
from .projection import ProjectionPlan, validate_projection
from .window import Window, plan_window
Expand Down Expand Up @@ -55,6 +56,7 @@
"SchemaMismatchError",
"SourceConfig",
"SplitBrainError",
"UnsupportedFormatError",
"Window",
"__version__",
"load_config",
Expand Down
12 changes: 12 additions & 0 deletions hedgerow/src/hedgerow/daemon.py
Original file line number Diff line number Diff line change
Expand Up @@ -72,6 +72,7 @@

from .config import HedgerowConfig
from .filtering import RowFilter, build_filter
from .formats import require_parquet
from .halts import (
DataIntegrityError,
DeletesPresentError,
Expand Down Expand Up @@ -193,11 +194,17 @@ def start(self) -> None:
self._source_catalog = self._source_client.catalog(cfg.source.catalog)
self._source_ns = self._source_catalog.namespace(cfg.source.namespace)
source_table = self._source_ns.table(cfg.source.table)
require_parquet(
(source_table.properties or {}).get("write.format.default", "parquet")
)
self.source_uuid = source_table.table_uuid

dest_catalog = self._dest_client.catalog(cfg.destination.catalog)
self._dest_ns = dest_catalog.namespace(cfg.destination.namespace)
dest_table = self._dest_ns.table(cfg.destination.table)
require_parquet(
(dest_table.properties or {}).get("write.format.default", "parquet")
)
self.dest_uuid = dest_table.table_uuid

# Destination columns define the projected set; source must cover
Expand Down Expand Up @@ -251,6 +258,7 @@ def _resolve_source_table(self):
f"(pinned uuid {self.source_uuid}) no longer exists. "
"HALT: refusing to continue against a dropped table."
) from None
require_parquet((table.properties or {}).get("write.format.default", "parquet"))
if table.table_uuid != self.source_uuid:
raise IncarnationChangedError(
f"source table {cfg.catalog}/{cfg.namespace}.{cfg.table} was "
Expand All @@ -270,6 +278,7 @@ def _resolve_dest_table(self):
f"destination table {cfg.catalog}/{cfg.namespace}.{cfg.table} "
f"(pinned uuid {self.dest_uuid}) no longer exists. HALT."
) from None
require_parquet((table.properties or {}).get("write.format.default", "parquet"))
if table.table_uuid != self.dest_uuid:
raise IncarnationChangedError(
f"destination table {cfg.catalog}/{cfg.namespace}.{cfg.table} "
Expand Down Expand Up @@ -356,6 +365,9 @@ def run_once(self) -> CycleResult:
"from the source)."
)

for file in plan.files:
require_parquet(file.file_format)

rows_read = rows_appended = appends = 0
max_rows = cfg.replication.max_rows_per_append
max_append_retries = cfg.replication.max_append_retries
Expand Down
3 changes: 3 additions & 0 deletions hedgerow/src/hedgerow/discovery.py
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@
from pyhoglake.models import ChangesPlan

from .events import EventTransform
from .formats import require_parquet
from .halts import DataIntegrityError, DeletesPresentError, IncarnationChangedError
from .pending import Fragment, PendingStore, SourceFile

Expand Down Expand Up @@ -54,6 +55,8 @@ def discover_window(
raise DeletesPresentError(
"raw source has deletion vectors; refusing to skip or reinterpret pending input"
)
for file in plan.files:
require_parquet(file.file_format)
fragments = []
source_files = []
source_bytes = selected_bytes = 0
Expand Down
9 changes: 9 additions & 0 deletions hedgerow/src/hedgerow/formats.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,9 @@
from .halts import UnsupportedFormatError


def require_parquet(file_format: str) -> None:
if file_format != "parquet":
raise UnsupportedFormatError(
f"Hedgerow cannot replicate data format {file_format!r}. "
"Use Parquet source and destination tables."
)
4 changes: 4 additions & 0 deletions hedgerow/src/hedgerow/halts.py
Original file line number Diff line number Diff line change
Expand Up @@ -79,3 +79,7 @@ class PersistentFailureError(HaltError):
halting instead of retrying forever."""

exit_code = 9


class UnsupportedFormatError(HaltError):
exit_code = 10
Loading
Loading