Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
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
34 changes: 26 additions & 8 deletions AGENT.md

Large diffs are not rendered by default.

36 changes: 23 additions & 13 deletions README.md

Large diffs are not rendered by default.

98 changes: 88 additions & 10 deletions millpond/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -777,11 +777,35 @@ def _bad(entry: str, why: str) -> RuntimeError:
_BACKOFF_JITTER = 0.25


# Local work every hoglake attempt does before it issues its commit, and
# therefore outside the request timeout the model used to be built
# entirely from: encode the fanout's parquet and upload it. It is charged
# PER ATTEMPT because a rebuild really does redo it — which it did not
# used to, because a refused commit was replayed rather than rebuilt, and
# a `ddl_since_read_snapshot` refusal is now a routine per-attempt cost
# while a producer rollout adds columns under a running fleet.
#
# BASIS, so the number is auditable rather than round: the local-stack
# integration suite flushes 9 rows over 3 partitions and the slowest
# observed `_prepare` in those runs is well under a second, which says
# nothing about production. The figure that does is the hoglake-side
# measurement the upload path was built from — ~4.3 s to encode and
# ~5 waves of round trips to upload the prod-us events writer's ~105K
# rows over ~271 partitions, call it ~10 s with the object store
# misbehaving — doubled for headroom. It is NOT measured on this side;
# it is a deliberate over-estimate whose only job is to stop the default
# budget from sitting inside the measurement error of the liveness
# deadline. Re-derive it from a real prod flush profile before trusting
# it to more than that.
_PREPARE_ALLOWANCE_S = 20.0


def _hoglake_worst_case_flush_s(max_retries: int, timeout_s: float) -> float:
"""Longest a single sink.write() can take: every attempt spending its
full request timeout, with the whole backoff ladder between them."""
full request timeout plus its own encode-and-upload, with the whole
backoff ladder between them."""
ladder = sum(min(_BACKOFF_BASE_S * (2**attempt), _BACKOFF_CAP_S) for attempt in range(max(0, max_retries - 1)))
return max_retries * timeout_s + ladder * (1 + _BACKOFF_JITTER)
return max_retries * (timeout_s + _PREPARE_ALLOWANCE_S) + ladder * (1 + _BACKOFF_JITTER)
Comment thread
jghoman marked this conversation as resolved.


def _check_hoglake_liveness_budget(max_retries: int, timeout_s: float) -> None:
Expand All @@ -804,6 +828,47 @@ def _check_hoglake_liveness_budget(max_retries: int, timeout_s: float) -> None:
)


# pyhoglake's upload fan-out knob. Not a millpond variable, which is
# exactly why it is checked here: pyhoglake reads it inside
# `resolve_concurrency`, on the prepared-append path, and refuses a bad
# value as a client-side `ValidationError` — which millpond's
# `is_retryable` correctly calls PERMANENT. So a typo'd fan-out is a pod
# that starts clean, passes its probes, takes its partitions, and then
# crashes on its first flush with a 422-shaped error about a value no
# millpond code mentions. Every other knob is validated at load; this one
# was only not because it belongs to someone else.
_UPLOAD_CONCURRENCY_ENV = "PYHOGLAKE_UPLOAD_CONCURRENCY"


def _check_upload_concurrency() -> None:
"""Refuse a PYHOGLAKE_UPLOAD_CONCURRENCY pyhoglake would refuse.

Parsed exactly as `pyhoglake.upload.resolve_concurrency` parses it —
strip, `int()`, reject below 1 — so this cannot refuse a value
pyhoglake accepts or accept one it refuses. Unset is the library's
default (64) and is left alone; an upper bound is deliberately NOT
imposed, because the ceiling is the object store's and pyhoglake
caps the fan-out at the group count anyway.
"""
raw = os.environ.get(_UPLOAD_CONCURRENCY_ENV, "").strip()
if not raw:
return
try:
value = int(raw)
except ValueError:
raise RuntimeError(
f"{_UPLOAD_CONCURRENCY_ENV} must be a positive integer, got {raw!r}. pyhoglake reads it "
f"on the upload path and refuses a bad value permanently, so this would crash the first "
f"flush rather than the pod start."
) from None
if value < 1:
raise RuntimeError(
f"{_UPLOAD_CONCURRENCY_ENV} must be a positive integer, got {value}. pyhoglake reads it "
f"on the upload path and refuses a bad value permanently, so this would crash the first "
f"flush rather than the pod start."
)


_S3_URI = re.compile(r"^s3://[a-z0-9][a-z0-9.\-]{1,61}[a-z0-9](/.*)?$")


Expand Down Expand Up @@ -897,16 +962,29 @@ def _load_hoglake_fields() -> dict:
# failure instead — the one outcome that has to hold a prepared
# payload and resend it blind.
#
# The defaults (8 attempts x 45s, plus a jittered ladder main.py caps
# at 30s a step) come to ~474s of the 480s budget. That is the
# all-eight-attempts-black-hole case and it is deliberately close to
# the line: a catalog that has not answered a single request in eight
# minutes is one this pod should be dying over. What the check
# prevents is the same arithmetic going unnoticed when an operator
# raises either knob.
max_retries = _positive_int("HOGLAKE_MAX_RETRY_COUNT", 8)
# The model charges each attempt its request timeout AND a
# _PREPARE_ALLOWANCE_S of local encode-and-upload, because on the
# hoglake path a retry REBUILDS the flush rather than replaying it
# whenever the server refuses the commit as `ddl_since_read_snapshot`
# — which a producer rollout adding columns under a running fleet
# makes a routine per-attempt cost, not a rarity. That term is what
# moved the default: 8 attempts x (45s + 20s) plus the jittered
# ladder is ~634s, well past the 480s deadline, so the default is 6
# attempts (~429s, leaving ~51s of margin). The previous default of 8
# fit only because the model did not charge for the work outside the
# request, and "deliberately close to the line" is not a thing to be
# about a SIGKILL: a flush that is killed mid-commit is the one
# outcome with no error in the pod's own logs.
#
# Six is still a real ladder against the case the ladder exists for
# (a convoyed catalog answering 503 + Retry-After, which main.py
# honours as a floor), and an operator who wants more has to lower
# the timeout to buy it — which is the trade being made explicit
# rather than discovered.
max_retries = _positive_int("HOGLAKE_MAX_RETRY_COUNT", 6)
timeout_s = _positive_float("HOGLAKE_REQUEST_TIMEOUT_S", 45.0)
_check_hoglake_liveness_budget(max_retries, timeout_s)
_check_upload_concurrency()

partition_raw = os.environ.get("HOGLAKE_PARTITION_BY", "").strip()
if not partition_raw and os.environ.get("DUCKLAKE_PARTITION_BY", "").strip():
Expand Down
Loading
Loading