fix(checkpoints): derive the tag wait from the VPG's cadence (#12)
This commit was merged in pull request #12.
This commit is contained in:
@@ -4,6 +4,7 @@ from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import re
|
||||
import statistics
|
||||
from datetime import UTC, datetime
|
||||
from typing import Any
|
||||
|
||||
@@ -23,6 +24,15 @@ _WS = re.compile(r"\s+")
|
||||
# so the name stays readable in the Zerto UI checkpoint list.
|
||||
TAG_MAX_LEN = 250
|
||||
|
||||
# Tag-wait budget. Checkpoint cadence is set by the VPG's protected site and
|
||||
# measured 5s (vSphere), 60s (Azure) and 630s (AWS) in one estate, so the wait
|
||||
# has to be derived rather than fixed.
|
||||
DEFAULT_TAG_TIMEOUT_S = 45.0 # cadence unmeasurable; the old constant
|
||||
MIN_TAG_TIMEOUT_S = 45.0
|
||||
MAX_TAG_TIMEOUT_S = 300.0 # AWS needed 128s; this leaves real headroom
|
||||
MIN_POLL_INTERVAL_S = 1.5
|
||||
MAX_POLL_INTERVAL_S = 15.0
|
||||
|
||||
|
||||
def _clean(value: Any, limit: int) -> str:
|
||||
"""One field of a checkpoint name: single-line, no separator collisions."""
|
||||
@@ -74,29 +84,81 @@ def checkpoint_id(row: dict[str, Any]) -> str:
|
||||
return str(value) if value is not None else ""
|
||||
|
||||
|
||||
def checkpoint_gaps(rows: list[dict[str, Any]], sample: int = 20) -> list[float]:
|
||||
"""Seconds between consecutive checkpoints, newest `sample` of them."""
|
||||
stamps: list[datetime] = []
|
||||
for row in rows[-sample:]:
|
||||
raw = str(pick(row, "TimeStamp", "Timestamp", "timestamp") or "")
|
||||
try:
|
||||
stamps.append(datetime.fromisoformat(raw))
|
||||
except ValueError:
|
||||
continue
|
||||
return [
|
||||
(stamps[i + 1] - stamps[i]).total_seconds()
|
||||
for i in range(len(stamps) - 1)
|
||||
if stamps[i + 1] >= stamps[i]
|
||||
]
|
||||
|
||||
|
||||
def cadence_seconds(rows: list[dict[str, Any]], sample: int = 20) -> float | None:
|
||||
"""How often this VPG writes a checkpoint. None when it cannot be measured."""
|
||||
gaps = checkpoint_gaps(rows, sample)
|
||||
return statistics.median(gaps) if gaps else None
|
||||
|
||||
|
||||
def tag_wait_budget(cadence: float | None) -> tuple[float, float]:
|
||||
"""How long to wait for a tag, and how often to look, given the cadence.
|
||||
|
||||
Cadence is set by the VPG's PROTECTED site, and the spread is enormous:
|
||||
measured 5s on vSphere, 60s on Azure, 630s on AWS. A single constant cannot
|
||||
serve all three. The old fixed 45s was generous for vSphere and impossible
|
||||
for AWS, where a tag took 128s to surface, so the guard reported "no
|
||||
checkpoint" while Zerto was still creating one and the change was refused
|
||||
for no reason.
|
||||
|
||||
Visibility does not scale linearly with cadence (the insert makes its own
|
||||
off-cadence checkpoint), so this is 2x cadence plus headroom, clamped.
|
||||
"""
|
||||
if cadence is None or cadence <= 0:
|
||||
return DEFAULT_TAG_TIMEOUT_S, MIN_POLL_INTERVAL_S
|
||||
timeout = min(max(2 * cadence + 30, MIN_TAG_TIMEOUT_S), MAX_TAG_TIMEOUT_S)
|
||||
interval = min(max(cadence / 10, MIN_POLL_INTERVAL_S), MAX_POLL_INTERVAL_S)
|
||||
return timeout, interval
|
||||
|
||||
|
||||
async def wait_for_tag(
|
||||
client: ZertoClient,
|
||||
vpg_identifier: str,
|
||||
tag: str,
|
||||
*,
|
||||
timeout_s: float = 45.0,
|
||||
interval_s: float = 1.5,
|
||||
timeout_s: float | None = None,
|
||||
interval_s: float | None = None,
|
||||
) -> dict[str, Any]:
|
||||
"""Wait until the tag is listed. Budget derived from the VPG's own cadence."""
|
||||
rows = await client.list_checkpoints(vpg_identifier)
|
||||
for row in rows:
|
||||
if checkpoint_tag(row) == tag:
|
||||
return row
|
||||
|
||||
cadence = cadence_seconds(rows)
|
||||
budget, poll = tag_wait_budget(cadence)
|
||||
timeout_s = budget if timeout_s is None else timeout_s
|
||||
interval_s = poll if interval_s is None else interval_s
|
||||
|
||||
deadline = asyncio.get_event_loop().time() + timeout_s
|
||||
last: list[dict[str, Any]] = []
|
||||
while asyncio.get_event_loop().time() < deadline:
|
||||
last = await client.list_checkpoints(vpg_identifier)
|
||||
for row in last:
|
||||
await asyncio.sleep(interval_s)
|
||||
rows = await client.list_checkpoints(vpg_identifier)
|
||||
for row in rows:
|
||||
if checkpoint_tag(row) == tag:
|
||||
return row
|
||||
await asyncio.sleep(interval_s)
|
||||
measured = f"{cadence:.0f}s" if cadence else "unknown"
|
||||
raise ZertoError(
|
||||
f"Tagged checkpoint {tag!r} did not appear on VPG {vpg_identifier} "
|
||||
f"within {timeout_s:.0f}s. Do not mutate. "
|
||||
"On a cloud-protected VPG the tag routinely takes longer than this to appear "
|
||||
"(measured ~34s on Azure, ~128s on AWS), so this timeout may simply be too "
|
||||
"short rather than the insert having failed. Check the Zerto task before "
|
||||
"assuming it did not land."
|
||||
f"within {timeout_s:.0f}s (this VPG checkpoints about every {measured}). "
|
||||
"Do not mutate. Check the Zerto task before assuming the insert failed: "
|
||||
"a completed task with no visible checkpoint means the wait was short, "
|
||||
"not that the insert was rejected."
|
||||
)
|
||||
|
||||
|
||||
|
||||
Reference in New Issue
Block a user