fix(checkpoints): derive the tag wait from the VPG's cadence #12
+25
-23
@@ -37,14 +37,19 @@ venv that has this package installed.
|
|||||||
The hook is synchronous: the host waits. That is the point, because the
|
The hook is synchronous: the host waits. That is the point, because the
|
||||||
checkpoint has to exist before the change does.
|
checkpoint has to exist before the change does.
|
||||||
|
|
||||||
Budget for the slow path, not the fast one. A successful tag took about 4s
|
Budget for the slow path, not the fast one. How long the guard takes is set by
|
||||||
against a healthy vSphere-protected VPG, but a **refusal took 63s**, because
|
the VPG's checkpoint cadence, which is set in turn by its protected site:
|
||||||
`wait_for_tag` spends 45s before giving up.
|
|
||||||
|
|
||||||
So:
|
| protected at | cadence | guard takes |
|
||||||
|
|---|---|---|
|
||||||
|
| vSphere | 5s | ~7s |
|
||||||
|
| Azure | 60s | ~40s |
|
||||||
|
| AWS | 630s | ~111s |
|
||||||
|
|
||||||
- `timeout` in settings.json: 180 (seconds)
|
The tag wait is derived from that cadence and capped at 300s, so:
|
||||||
- `ZERTO_HOOK_GUARD_TIMEOUT`: 150 (seconds), kept under it
|
|
||||||
|
- `timeout` in settings.json: 360 (seconds)
|
||||||
|
- `ZERTO_HOOK_GUARD_TIMEOUT`: 330 (seconds), kept under it
|
||||||
|
|
||||||
If the host's timeout fires first it cancels the hook and **discards its
|
If the host's timeout fires first it cancels the hook and **discards its
|
||||||
output**, and the tool call carries on through the normal permission flow. A
|
output**, and the tool call carries on through the normal permission flow. A
|
||||||
@@ -68,26 +73,23 @@ likely to be running unattended.
|
|||||||
|
|
||||||
`~/.zerto-guard-hook.log`, or `ZERTO_HOOK_LOG`. One line per decision.
|
`~/.zerto-guard-hook.log`, or `ZERTO_HOOK_LOG`. One line per decision.
|
||||||
|
|
||||||
## Known issue: false denials on cloud-protected VPGs
|
## Why the timeouts are derived, not fixed
|
||||||
|
|
||||||
The `win2019-1` denial above was a **false negative**, and it is worth
|
This hook used to deny every change to a cloud-protected VM.
|
||||||
understanding before relying on this hook in an estate with cloud-protected
|
|
||||||
workloads.
|
|
||||||
|
|
||||||
`wait_for_tag` gives up after a hardcoded 45s. That is generous for a
|
`wait_for_tag` gave up after a hardcoded 45s. That is generous on a
|
||||||
vSphere-protected VPG, which checkpoints every 5s and surfaces a tag in about
|
vSphere-protected VPG, which checkpoints every 5s and surfaces a tag in about
|
||||||
4s. It is far too short elsewhere: a tag takes ~34s to appear on an
|
4s, and impossible on an AWS-protected one, where a tag takes ~128s because the
|
||||||
Azure-protected VPG and ~128s on an AWS-protected one, because journal cadence
|
journal only checkpoints every 630s.
|
||||||
is set by the protected site (5s vSphere, 60s Azure, 630s AWS).
|
|
||||||
|
|
||||||
So the hook denied the change, and the checkpoint landed anyway. It is in the
|
So the guard reported "no checkpoint, refusing the change" while Zerto was in the
|
||||||
journal as `cp 56`, timestamped a minute after the hook reported failure. The
|
middle of creating one. The checkpoint landed a minute later, in the journal,
|
||||||
guard told the agent there was no rewind point while Zerto was in the middle of
|
after the agent had already been told there was no rewind point.
|
||||||
creating one.
|
|
||||||
|
|
||||||
That failure mode is worse than the one the hook guards against, because it is
|
That is a worse failure than the one this hook exists to prevent. It is silent,
|
||||||
silent and looks correct: legitimate work is refused on every cloud-protected VM
|
it looks correct in the log, and it blocks legitimate work on every
|
||||||
while the log reads like the guard is doing its job.
|
cloud-protected VM in the estate.
|
||||||
|
|
||||||
Until `wait_for_tag` becomes cadence-aware, scope this hook's matcher to
|
Both budgets are now derived from the VPG's measured cadence rather than
|
||||||
vSphere-protected workloads.
|
guessed, which is why the numbers above differ by a factor of fifteen between
|
||||||
|
platforms.
|
||||||
|
|||||||
@@ -8,7 +8,7 @@
|
|||||||
{
|
{
|
||||||
"type": "command",
|
"type": "command",
|
||||||
"command": "/home/you/zerto-ai-rewind/.venv/bin/python /home/you/zerto-ai-rewind/hooks/zerto_guard_hook.py",
|
"command": "/home/you/zerto-ai-rewind/.venv/bin/python /home/you/zerto-ai-rewind/hooks/zerto_guard_hook.py",
|
||||||
"timeout": 180
|
"timeout": 360
|
||||||
}
|
}
|
||||||
]
|
]
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -24,7 +24,7 @@ tagged checkpoint and waits for the Zerto task to reach Completed.
|
|||||||
"matcher": "mcp__.*",
|
"matcher": "mcp__.*",
|
||||||
"hooks": [{"type": "command",
|
"hooks": [{"type": "command",
|
||||||
"command": "/path/to/.venv/bin/python /path/to/hooks/zerto_guard_hook.py",
|
"command": "/path/to/.venv/bin/python /path/to/hooks/zerto_guard_hook.py",
|
||||||
"timeout": 180}]}]}}
|
"timeout": 360}]}]}}
|
||||||
"""
|
"""
|
||||||
|
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
@@ -40,7 +40,11 @@ LOG = os.environ.get("ZERTO_HOOK_LOG", os.path.expanduser("~/.zerto-guard-hook.l
|
|||||||
# Seconds the hook will wait for the checkpoint. Must stay under the hook
|
# Seconds the hook will wait for the checkpoint. Must stay under the hook
|
||||||
# timeout configured in settings.json, or the host cancels us and the tool
|
# timeout configured in settings.json, or the host cancels us and the tool
|
||||||
# call proceeds unguarded through the normal permission flow.
|
# call proceeds unguarded through the normal permission flow.
|
||||||
GUARD_TIMEOUT_S = float(os.environ.get("ZERTO_HOOK_GUARD_TIMEOUT", "150"))
|
# Must exceed the largest tag wait the guard can take. That is now derived from
|
||||||
|
# the VPG's checkpoint cadence and capped at 300s (MAX_TAG_TIMEOUT_S), because a
|
||||||
|
# tag takes ~128s to surface on an AWS-protected VPG. Too small a budget here
|
||||||
|
# just moves the false denial from the guard into the hook.
|
||||||
|
GUARD_TIMEOUT_S = float(os.environ.get("ZERTO_HOOK_GUARD_TIMEOUT", "330"))
|
||||||
UNKNOWN_DECISION = os.environ.get("ZERTO_HOOK_UNKNOWN", "prompt") # prompt | allow | deny
|
UNKNOWN_DECISION = os.environ.get("ZERTO_HOOK_UNKNOWN", "prompt") # prompt | allow | deny
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -4,6 +4,7 @@ from __future__ import annotations
|
|||||||
|
|
||||||
import asyncio
|
import asyncio
|
||||||
import re
|
import re
|
||||||
|
import statistics
|
||||||
from datetime import UTC, datetime
|
from datetime import UTC, datetime
|
||||||
from typing import Any
|
from typing import Any
|
||||||
|
|
||||||
@@ -23,6 +24,15 @@ _WS = re.compile(r"\s+")
|
|||||||
# so the name stays readable in the Zerto UI checkpoint list.
|
# so the name stays readable in the Zerto UI checkpoint list.
|
||||||
TAG_MAX_LEN = 250
|
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:
|
def _clean(value: Any, limit: int) -> str:
|
||||||
"""One field of a checkpoint name: single-line, no separator collisions."""
|
"""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 ""
|
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(
|
async def wait_for_tag(
|
||||||
client: ZertoClient,
|
client: ZertoClient,
|
||||||
vpg_identifier: str,
|
vpg_identifier: str,
|
||||||
tag: str,
|
tag: str,
|
||||||
*,
|
*,
|
||||||
timeout_s: float = 45.0,
|
timeout_s: float | None = None,
|
||||||
interval_s: float = 1.5,
|
interval_s: float | None = None,
|
||||||
) -> dict[str, Any]:
|
) -> 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
|
deadline = asyncio.get_event_loop().time() + timeout_s
|
||||||
last: list[dict[str, Any]] = []
|
|
||||||
while asyncio.get_event_loop().time() < deadline:
|
while asyncio.get_event_loop().time() < deadline:
|
||||||
last = await client.list_checkpoints(vpg_identifier)
|
await asyncio.sleep(interval_s)
|
||||||
for row in last:
|
rows = await client.list_checkpoints(vpg_identifier)
|
||||||
|
for row in rows:
|
||||||
if checkpoint_tag(row) == tag:
|
if checkpoint_tag(row) == tag:
|
||||||
return row
|
return row
|
||||||
await asyncio.sleep(interval_s)
|
measured = f"{cadence:.0f}s" if cadence else "unknown"
|
||||||
raise ZertoError(
|
raise ZertoError(
|
||||||
f"Tagged checkpoint {tag!r} did not appear on VPG {vpg_identifier} "
|
f"Tagged checkpoint {tag!r} did not appear on VPG {vpg_identifier} "
|
||||||
f"within {timeout_s:.0f}s. Do not mutate. "
|
f"within {timeout_s:.0f}s (this VPG checkpoints about every {measured}). "
|
||||||
"On a cloud-protected VPG the tag routinely takes longer than this to appear "
|
"Do not mutate. Check the Zerto task before assuming the insert failed: "
|
||||||
"(measured ~34s on Azure, ~128s on AWS), so this timeout may simply be too "
|
"a completed task with no visible checkpoint means the wait was short, "
|
||||||
"short rather than the insert having failed. Check the Zerto task before "
|
"not that the insert was rejected."
|
||||||
"assuming it did not land."
|
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -1,5 +1,7 @@
|
|||||||
from datetime import UTC, datetime
|
from datetime import UTC, datetime
|
||||||
|
|
||||||
|
import pytest
|
||||||
|
|
||||||
from zerto_rewind_mcp.checkpoints import (
|
from zerto_rewind_mcp.checkpoints import (
|
||||||
TAG_MAX_LEN,
|
TAG_MAX_LEN,
|
||||||
checkpoint_id,
|
checkpoint_id,
|
||||||
@@ -50,3 +52,93 @@ def test_checkpoint_row_keys():
|
|||||||
row2 = {"checkpointId": "cp-2", "tag": "t"}
|
row2 = {"checkpointId": "cp-2", "tag": "t"}
|
||||||
assert checkpoint_id(row2) == "cp-2"
|
assert checkpoint_id(row2) == "cp-2"
|
||||||
assert checkpoint_tag(row2) == "t"
|
assert checkpoint_tag(row2) == "t"
|
||||||
|
|
||||||
|
|
||||||
|
def _rows(*offsets_seconds):
|
||||||
|
from datetime import timedelta
|
||||||
|
|
||||||
|
base = datetime(2026, 9, 22, 12, 0, 0, tzinfo=UTC)
|
||||||
|
return [
|
||||||
|
{"TimeStamp": (base + timedelta(seconds=o)).isoformat().replace("+00:00", "Z")}
|
||||||
|
for o in offsets_seconds
|
||||||
|
]
|
||||||
|
|
||||||
|
|
||||||
|
def test_cadence_measures_the_median_gap():
|
||||||
|
from zerto_rewind_mcp.checkpoints import cadence_seconds
|
||||||
|
|
||||||
|
assert cadence_seconds(_rows(0, 5, 10, 15, 20)) == 5.0
|
||||||
|
assert cadence_seconds(_rows(0, 60, 120, 180)) == 60.0
|
||||||
|
# one irregular gap must not drag the answer around
|
||||||
|
assert cadence_seconds(_rows(0, 5, 10, 400, 405, 410)) == 5.0
|
||||||
|
|
||||||
|
|
||||||
|
def test_cadence_is_none_when_unmeasurable():
|
||||||
|
from zerto_rewind_mcp.checkpoints import cadence_seconds
|
||||||
|
|
||||||
|
assert cadence_seconds([]) is None
|
||||||
|
assert cadence_seconds([{"TimeStamp": "not-a-date"}]) is None
|
||||||
|
assert cadence_seconds(_rows(0)) is None # one checkpoint gives no gap
|
||||||
|
|
||||||
|
|
||||||
|
def test_tag_wait_budget_covers_every_measured_platform():
|
||||||
|
"""The three cadences measured in one estate, and what each actually needed."""
|
||||||
|
from zerto_rewind_mcp.checkpoints import tag_wait_budget
|
||||||
|
|
||||||
|
for cadence, observed_visibility in ((5.0, 4.0), (60.0, 34.0), (630.0, 128.0)):
|
||||||
|
budget, interval = tag_wait_budget(cadence)
|
||||||
|
assert budget > observed_visibility, (
|
||||||
|
f"cadence {cadence}s budgets {budget}s but the tag took {observed_visibility}s"
|
||||||
|
)
|
||||||
|
assert interval >= 1.5
|
||||||
|
|
||||||
|
|
||||||
|
def test_tag_wait_budget_is_clamped_at_both_ends():
|
||||||
|
from zerto_rewind_mcp.checkpoints import (
|
||||||
|
DEFAULT_TAG_TIMEOUT_S,
|
||||||
|
MAX_TAG_TIMEOUT_S,
|
||||||
|
MIN_TAG_TIMEOUT_S,
|
||||||
|
tag_wait_budget,
|
||||||
|
)
|
||||||
|
|
||||||
|
assert tag_wait_budget(0.1)[0] == MIN_TAG_TIMEOUT_S # absurdly fast VPG
|
||||||
|
assert tag_wait_budget(100_000)[0] == MAX_TAG_TIMEOUT_S # absurdly slow one
|
||||||
|
assert tag_wait_budget(None)[0] == DEFAULT_TAG_TIMEOUT_S
|
||||||
|
assert tag_wait_budget(None)[1] == 1.5
|
||||||
|
|
||||||
|
|
||||||
|
def test_wait_for_tag_returns_immediately_when_already_present():
|
||||||
|
import asyncio
|
||||||
|
|
||||||
|
from zerto_rewind_mcp.checkpoints import wait_for_tag
|
||||||
|
|
||||||
|
class Client:
|
||||||
|
def __init__(self):
|
||||||
|
self.calls = 0
|
||||||
|
|
||||||
|
async def list_checkpoints(self, vpg):
|
||||||
|
self.calls += 1
|
||||||
|
return [{"Tag": "ai:x", "CheckpointId": "7"}]
|
||||||
|
|
||||||
|
c = Client()
|
||||||
|
row = asyncio.run(wait_for_tag(c, "vpg", "ai:x"))
|
||||||
|
assert row["CheckpointId"] == "7"
|
||||||
|
assert c.calls == 1 # no sleep, no second poll
|
||||||
|
|
||||||
|
|
||||||
|
def test_wait_for_tag_error_names_the_measured_cadence():
|
||||||
|
import asyncio
|
||||||
|
|
||||||
|
from zerto_rewind_mcp.checkpoints import wait_for_tag
|
||||||
|
from zerto_rewind_mcp.client import ZertoError
|
||||||
|
|
||||||
|
class Client:
|
||||||
|
async def list_checkpoints(self, vpg):
|
||||||
|
return _rows(0, 60, 120, 180) # 60s cadence, tag never appears
|
||||||
|
|
||||||
|
with pytest.raises(ZertoError) as err:
|
||||||
|
# explicit tiny timeout so the test does not actually wait 150s
|
||||||
|
asyncio.run(wait_for_tag(Client(), "vpg", "ai:missing", timeout_s=0.01, interval_s=0.01))
|
||||||
|
msg = str(err.value)
|
||||||
|
assert "about every 60s" in msg
|
||||||
|
assert "was short, not that the insert was rejected" in msg
|
||||||
|
|||||||
Reference in New Issue
Block a user