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
2 changes: 1 addition & 1 deletion atom/compass/clock/identity.py
Original file line number Diff line number Diff line change
Expand Up @@ -41,7 +41,7 @@ def __post_init__(self) -> None:
)
if not self.name:
raise ValueError("a logical process name must not be empty")
if self.name != self.name.strip() or any(c.isspace() for c in self.name):
if any(c.isspace() for c in self.name):
raise ValueError(
f"a logical process name must not contain whitespace, got {self.name!r}; "
"names are written verbatim into the timeline record, one field per column"
Expand Down
60 changes: 7 additions & 53 deletions atom/compass/clock/lookahead.py
Original file line number Diff line number Diff line change
Expand Up @@ -18,8 +18,7 @@
the other's current time, and the whole run collapses towards a single global
event loop. So the floor buys speed, and the amount of speed it buys is the
whole of what it decides. It never buys safety, and code that treats a small
floor as a hazard has the relationship backwards. `LookaheadMatrix.serializing`
exists to report the zero links, not to reject them.
floor as a hazard has the relationship backwards.

**An undeclared pair is not a zero, and it is not skippable either.** A caller
takes a minimum over one participant's whole row, quantified over every
Expand All @@ -41,34 +40,25 @@


class LinkClass(enum.Enum):
"""The three kinds of link between logical processes, and the scale of each.

`scale_seconds` is the order of magnitude the link has been observed or
modelled at. It is documentation and a sanity reference -- a declared floor
is never derived from it -- and it is what says which link is the tight one.
"""
"""The three kinds of link between logical processes."""

#: Traffic source to engine: the modelled admission delay. The only one of
#: the three with an end-to-end measurement behind it, and the measurement
#: is path-specific; see `TRAFFIC_TO_ENGINE_FLOOR_SECONDS`.
TRAFFIC_TO_ENGINE = ("traffic_to_engine", 1.0e-2)
TRAFFIC_TO_ENGINE = "traffic_to_engine"

#: Prefill role to decode role: a router relay plus a simulated transfer of
#: the cached keys and values. Modelled at millisecond scale, which is
#: comfortable -- it is the cheap link to stretch across a node.
PREFILL_TO_DECODE = ("prefill_to_decode", 1.0e-3)
PREFILL_TO_DECODE = "prefill_to_decode"

#: Pipeline stage to pipeline stage: a modelled send and receive of the
#: intermediate tensors. Microsecond scale, and the only tight one: stages
#: belong close together for exactly this reason.
PIPELINE_STAGE_TO_STAGE = ("pipeline_stage_to_stage", 1.0e-6)

def __init__(self, label: str, scale_seconds: float) -> None:
self.label = label
self.scale_seconds = scale_seconds
PIPELINE_STAGE_TO_STAGE = "pipeline_stage_to_stage"

def __str__(self) -> str:
return self.label
return self.value

Copy link
Copy Markdown
Owner Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Non-blocking, and already true at the tip. No test covers LinkClass.__str__, and this PR rewrites its body.

Principle 6: "A declined answer with a named reason is a result."

str(link_class) has one reader: InterLpLink.__repr__. That repr is the named reason in declare's duplicate refusal (... is already declared as InterLpLink(prefill -> decode, prefill_to_decode, floor=0.001s)).

With mutate.py str_is_name, return self.name in place of return self.value:

  • at the head, 67 passed and 0 failed;
  • at the tip, the same mutant on return self.label gives 74 passed and 0 failed.

test_declaring_the_same_link_twice_is_refused matches only "already declared". So the evidence that self.label → self.value is equivalent is the probe, not a test.

My probe on node 18 agrees with the developer's: str(member) and repr(InterLpLink) are byte-identical on both sides for all 3 classes. The change is correct.

Optional one-line pin: add match="already declared as InterLpLink\(prefill -> decode, prefill_to_decode" to that test's pytest.raises.



#: Measured admission delay from the traffic source to the engine, per path, in
Expand Down Expand Up @@ -177,7 +167,7 @@ def inbound(self, target: LpId) -> tuple[InterLpLink, ...]:
not less. Zero would at least have been the conservative mistake.
"""
self._registry.require(target)
missing = self._missing_into(target)
missing = [source for source, into in self.undeclared() if into == target]
if missing:
raise KeyError(
f"{len(missing)} registered peer(s) have no declared floor into "
Expand All @@ -193,13 +183,6 @@ def inbound(self, target: LpId) -> tuple[InterLpLink, ...]:
if source != target
)

def _missing_into(self, target: LpId) -> tuple[LpId, ...]:
return tuple(
source
for source in self._registry.ids()
if source != target and (source, target) not in self._links
)

def require_complete(self) -> None:
"""Refuse unless every ordered pair of registered participants has a floor.

Expand All @@ -218,10 +201,6 @@ def require_complete(self) -> None:
"lockstep."
)

def links(self) -> tuple[InterLpLink, ...]:
"""Every declared link, ordered by source then target."""
return tuple(self._links[key] for key in sorted(self._links))

def undeclared(self) -> tuple[tuple[LpId, LpId], ...]:
"""Ordered pairs of registered identities with no floor yet, in order.

Expand All @@ -235,28 +214,3 @@ def undeclared(self) -> tuple[tuple[LpId, LpId], ...]:
for target in ids
if source != target and (source, target) not in self._links
)

def serializing(self) -> tuple[InterLpLink, ...]:
"""The links declared at a zero floor -- correct, and slow.

Reported so a run that turns out to be serial can say which pair made it
so. Not an error, and not a warning about correctness.
"""
return tuple(link for link in self.links() if link.floor_seconds == 0.0)

def tightest(self) -> InterLpLink | None:
"""The declared link with the smallest floor: what bounds how far anything runs ahead.

Reporting, not protocol: it walks the links that exist, so on an
incomplete matrix it answers about the part that was declared. Check
`require_complete` before quoting it as the bound on a run.
"""
links = self.links()
if not links:
return None
return min(
links, key=lambda link: (link.floor_seconds, link.source, link.target)
)

def __len__(self) -> int:
return len(self._links)
11 changes: 2 additions & 9 deletions atom/compass/clock/registry.py
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,6 @@ class LpRegistry:

def __init__(self) -> None:
self._members: dict[LpId, None] = {}
self._ordered: tuple[LpId, ...] | None = None

def register(self, lp_id: LpId) -> LpId:
"""Add one logical process. Refuses a duplicate rather than ignoring it.
Expand All @@ -42,14 +41,11 @@ def register(self, lp_id: LpId) -> LpId:
if lp_id in self._members:
raise ValueError(f"{lp_id} is already registered")
self._members[lp_id] = None
self._ordered = None
return lp_id

def ids(self) -> tuple[LpId, ...]:
"""Every registered identity, in the total order, cheapest to call repeatedly."""
if self._ordered is None:
self._ordered = tuple(sorted(self._members))
return self._ordered
"""Every registered identity, in the total order."""
return tuple(sorted(self._members))

def require(self, lp_id: LpId) -> LpId:
"""Return `lp_id`, or refuse with the names that are registered.
Expand All @@ -76,6 +72,3 @@ def __iter__(self):

def __len__(self) -> int:
return len(self._members)

Copy link
Copy Markdown
Owner Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Non-blocking. The F7b ruling has no test that fires. This is at __contains__, L67 above; the anchor is here because L67 is outside the diff.

Principle 6: "Refuse rather than fall back. A declined answer with a named reason is a result. A guessed one is a defect."
AI_DEV_RULES, gate 4: "A check counts only once someone has seen it fire."

The coordinator skipped F7b so that [] in registry keeps raising TypeError rather than returning a silent False. The head keeps that behaviour: the probe on node 18 at the merged tree d58491587 (its clock files are byte-identical to 5f9ee034e) gives [] in registry → TypeError: unhashable type: 'list'.

No test holds it, though. With mutate.py contains_via_iter, which deletes LpRegistry.__contains__ and changes nothing else, the three clock files give 67 passed, 0 failed (test_clock_lp_identity.py, test_clock_order_across_processes.py, test_backend_kv_geometry.py). test_membership_and_size still passes, because the __iter__ fallback gives the same answer for every LpId.

So the next ponytail pass will measure F7b as free and propose it again. Two lines in test_membership_and_size would pin the ruling:

    with pytest.raises(TypeError):
        [] in registry

This is not blocking: the PR follows the ruling, and the behaviour is intact at the head.


def __repr__(self) -> str:
return f"LpRegistry({', '.join(str(lp_id) for lp_id in self.ids())})"
74 changes: 4 additions & 70 deletions tests/compass/test_clock_lp_identity.py
Original file line number Diff line number Diff line change
Expand Up @@ -44,7 +44,6 @@

from atom.compass import clock
from atom.compass.clock import (
TRAFFIC_TO_ENGINE_FLOOR_SECONDS,
InterLpLink,
LinkClass,
LookaheadMatrix,
Expand Down Expand Up @@ -173,20 +172,12 @@ def test_an_undeclared_pair_is_refused_rather_than_read_as_zero():


def test_a_zero_floor_is_a_declaration_and_not_an_error():
# Zero serializes the pair. It stays correct, so it is accepted and reported
# rather than rejected -- the floor buys speed, never safety.
# Zero serializes the pair. It stays correct, so it is accepted rather than
# rejected -- the floor buys speed, never safety.
registry = _registry(PREFILL, DECODE)
matrix = LookaheadMatrix(registry)
link = matrix.declare(PREFILL, DECODE, LinkClass.PREFILL_TO_DECODE, 0.0)
matrix.declare(PREFILL, DECODE, LinkClass.PREFILL_TO_DECODE, 0.0)
assert matrix.lookahead(PREFILL, DECODE) == 0.0
assert matrix.serializing() == (link,)


def test_a_nonzero_floor_is_not_reported_as_serializing():
registry = _registry(PREFILL, DECODE)
matrix = LookaheadMatrix(registry)
matrix.declare(PREFILL, DECODE, LinkClass.PREFILL_TO_DECODE, 1.0e-3)
assert matrix.serializing() == ()


@pytest.mark.parametrize("bad", [-1.0e-9, float("nan"), float("inf")])
Expand Down Expand Up @@ -278,7 +269,7 @@ def test_require_complete_accepts_a_matrix_with_every_pair_declared():


def test_a_declared_floor_cannot_be_rewritten_through_a_handed_out_link():
# `links`, `inbound` and `declare` all return the object itself. If it were
# `inbound` and `declare` both return the object itself. If it were
# writable, every refusal in `declare` -- negative, NaN, infinite, already
# declared -- would be reachable around.
registry = _registry(PREFILL, DECODE)
Expand All @@ -289,17 +280,6 @@ def test_a_declared_floor_cannot_be_rewritten_through_a_handed_out_link():
assert matrix.lookahead(PREFILL, DECODE) == pytest.approx(2.0e-3)


def test_inbound_links_come_back_in_the_total_order_whatever_order_they_were_declared():
# This is the shape the grant rule reads: every floor into one participant.
# It has to be ordered by identity, because a sum or a min taken over it in
# declaration order would depend on set-up.
registry = _registry(TRAFFIC, PREFILL, DECODE)
matrix = LookaheadMatrix(registry)
matrix.declare(TRAFFIC, DECODE, LinkClass.TRAFFIC_TO_ENGINE, 9.0e-3)
matrix.declare(PREFILL, DECODE, LinkClass.PREFILL_TO_DECODE, 1.0e-3)
assert tuple(link.source for link in matrix.inbound(DECODE)) == (PREFILL, TRAFFIC)


def test_inserting_a_participant_moves_nothing_that_was_already_declared():
# `pp-stage-1` sorts between `decode` and `traffic-source`, so under any
# index-based addressing it would displace one of them. Under identity the
Expand All @@ -318,52 +298,6 @@ def test_inserting_a_participant_moves_nothing_that_was_already_declared():
assert tuple(link.source for link in matrix.inbound(DECODE)) == (inserted, TRAFFIC)


def test_the_tightest_link_is_the_one_that_bounds_the_run():
registry = _registry(TRAFFIC, PREFILL, DECODE)
matrix = LookaheadMatrix(registry)
matrix.declare(TRAFFIC, PREFILL, LinkClass.TRAFFIC_TO_ENGINE, 9.0e-3)
matrix.declare(PREFILL, DECODE, LinkClass.PREFILL_TO_DECODE, 1.0e-3)
assert matrix.tightest().link_class is LinkClass.PREFILL_TO_DECODE
assert LookaheadMatrix(LpRegistry()).tightest() is None


def test_links_are_listed_in_a_stable_order():
registry = _registry(TRAFFIC, PREFILL, DECODE)
matrix = LookaheadMatrix(registry)
matrix.declare(TRAFFIC, PREFILL, LinkClass.TRAFFIC_TO_ENGINE, 9.0e-3)
matrix.declare(PREFILL, DECODE, LinkClass.PREFILL_TO_DECODE, 1.0e-3)
assert [(str(link.source), str(link.target)) for link in matrix.links()] == [
("prefill", "decode"),
("traffic-source", "prefill"),
]
assert len(matrix) == 2


# --- the configured scales ---------------------------------------------------


def test_the_three_link_classes_carry_the_scales_that_were_modelled():
# Recorded so the tight link is a property of the class rather than folklore:
# stage-to-stage is three orders below the admission delay, which is why
# stages belong close together and role boundaries are the cheap ones to
# stretch.
scales = {link_class: link_class.scale_seconds for link_class in LinkClass}
assert (
scales[LinkClass.PIPELINE_STAGE_TO_STAGE]
< scales[LinkClass.PREFILL_TO_DECODE]
< scales[LinkClass.TRAFFIC_TO_ENGINE]
)
assert min(scales, key=scales.get) is LinkClass.PIPELINE_STAGE_TO_STAGE


def test_the_admission_delay_is_kept_per_path():
# One number for both paths would misprice whichever was not measured.
assert TRAFFIC_TO_ENGINE_FLOOR_SECONDS == {
"offline_batch": 13.0e-3,
"serving": 9.0e-3,
}


# --- what the package is allowed to name -------------------------------------

# A participant knows a peer by name. Where that peer runs, how it is reached,
Expand Down
6 changes: 0 additions & 6 deletions tests/compass/test_clock_order_across_processes.py
Original file line number Diff line number Diff line change
Expand Up @@ -70,7 +70,6 @@
print("seed", os.environ.get("PYTHONHASHSEED", "<unset>"))
print("hash", hash(names[0]))
print("module", atom.__file__)
print("arrival", " ".join(arrival))
print("registry", " ".join(str(lp_id) for lp_id in registry.ids()))

Copy link
Copy Markdown
Owner Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Non-blocking. ponytail-review.

test_clock_order_across_processes.py:L22-24, L63-64, L67, L74-77, L82, L104: delete: the arrival rotation, which no remaining test checks and no mutant needs. The children register namesdirectly, and_child(seed) drops the rotate argument. About −12 lines.

Once T2 removed test_the_children_really_did_register_in_different_orders, nothing reads the rotation. Nothing it does is needed for detection either. Measured on node 18 over the three clock files at the merged tree d58491587 (its clock files are byte-identical to 5f9ee034e):

Mutant Failures
ids_insertion (tuple(self._members)), rotation on 4 failed
ids_insertion with the child's rotate = 0 the same 4 failed
ids_hash with rotate = 0 the same 4 failed

The 4 failures are the same tests each time, test_the_total_order_is_the_same_in_every_process among them. NAMES is not sorted, so orders[0] == " ".join(sorted(NAMES)) catches an insertion-order registry without any rotation. In-process, test_the_order_does_not_depend_on_the_order_of_registration catches it too.

Removing the rotation also removes the reason for the docstring's "built from the unrotated list on purpose" clause (L16-18) and for the comment at L74-77.

This is a suggestion, not a condition. The rotation is harmless, and removing it goes beyond the audit's 12 findings.

# Built from `names`, not from `arrival`: the seed is then the only thing that
# differs between children, so a difference in this line is attributable to it
Expand Down Expand Up @@ -130,11 +129,6 @@ def test_every_child_ran_against_the_tree_under_test(runs):
)


def test_the_children_really_did_register_in_different_orders(runs):
arrivals = [run["arrival"] for run in runs]
assert len(dict.fromkeys(arrivals)) == len(SEEDS), arrivals


def test_the_total_order_is_the_same_in_every_process(runs):
expected = " ".join(sorted(NAMES))
orders = [run["registry"] for run in runs]
Expand Down