* tests: read checked-in files as UTF-8 instead of the platform default Path.read_text() with no encoding uses locale.getpreferredencoding(), which is UTF-8 on the Linux runners and cp1252 on a stock Windows install. Nine module-level reads of checked-in source files were relying on that default. studio/backend/routes/inference.py carries the DeepSeek tool-call token regexes, so it holds U+FF5C and U+2581. Under cp1252 that read raised UnicodeDecodeError on byte 0x81 at position 97806, and because the reads run at import time it took test_cancel_atomicity.py and test_cancel_id_wiring.py out at collection, not as failures. Green on CI, permanently broken for a Windows contributor running the suite locally. Adds a guard: at module scope there is no tmp_path fixture, so a bare read_text()/write_text()/open() there is always touching a checked-in file. That makes the rule mechanical enough to enforce with no allowlist, while staying quiet about temp-dir I/O inside test bodies where the platform default is harmless. The repo already spells this correctly in 464 other places; this only stops the stragglers coming back. * tests: cover import-time helper reads and keep the guard py3.9-safe Follows up on the Codex review: - add `from __future__ import annotations`, since `str | None` in `_offender` is evaluated at import on Python 3.9 and pyproject declares requires-python ">=3.9,<3.15". - widen the guard from module scope to import time. Class bodies and the bodies of module-level helpers called from an executing statement run during collection too, so `CODE = _extract_mixed_precision_code()` was the same hazard as an inline read. `if __name__ == "__main__":` blocks are skipped: pytest never executes them. - scan studio/backend/tests/ as well as tests/. Both trees are collected on Windows by separate CI jobs, and the offender that started this, test_tool_xml_strip.py reading routes/inference.py, lives there. Widening it surfaced seven more import-time reads of checked-in sources; all now name utf-8. * [pre-commit.ci] auto fixes from pre-commit.com hooks for more information, see https://pre-commit.ci * Harden the import-time encoding guard for PR #7438 Close the detector gaps raised in review, all of which I reproduced against the actual AST before changing anything. False negatives (the guard let a real hazard through): - _is_main_guard ignored the comparison operator, so if __name__ != "__main__" counted as script-only even though its body runs at import. - The else arm of a main guard was discarded with the rest of the If node. - Decorators and argument defaults on a module-level def were skipped with the body, though both are evaluated when the def executes. - Path.open() in text mode was invisible; only builtin open() was matched. - encoding = None and encoding = "locale" both re-select the platform default, but the keyword merely being present counted as pinned. False positives (the guard would have blocked a compliant contributor): - A non-literal mode fell through to the "r" default, so open(p, mode) was flagged even when mode is "rb", where adding encoding= is a ValueError and there is no edit that satisfies the rule. - Same for open(*args) and a **kwargs splat, which hide the mode and can hide an encoding. - Lambda bodies and comprehension elements were walked even though neither runs at definition. Verified: still reports the same 22 offenders on unpatched main, green on this branch and on the tree merged with latest main (557 files), and an adversarial corpus of 33 cases now scores zero false positives and zero false negatives. Also corrected two docstring claims: neither collecting job runs on Windows, and the read is governed by locale.getencoding(). * Walk eager comprehensions and treat io.open as the builtin Two regressions from the previous commit, both reproduced against the AST before changing anything. Lumping list, set and dict comprehensions in with generator expressions was wrong. Only a genexp is lazy; the other three run their element expression, their filters and their nested iterators immediately, so CONTENTS = [p.read_text() for p in PATHS] at module scope is an import-time read the guard was silently missing. Comprehensions are now walked in full and only the genexp keeps the outermost-iterable-only treatment. io was also in the not-a-path-opener list, but io.open is the builtin, with the same mode position and the same platform default. io.open(CHECKED_IN_FILE) is exactly the hazard this guard exists for, so it is matched now, with binary modes and a pinned encoding still exempt. tarfile.open and fitz.open stay exempt since neither has an encoding to name. Verified: 13 targeted cases covering all five eager comprehension forms and io.open in text, binary and pinned shapes all classify correctly; still 22 offenders on unpatched main; green on this branch and on the tree merged with latest main. * Close three more walker gaps in the import-time guard All three reproduced against the AST first. A generator expression handed straight to a call is consumed there, so DATA = "".join(p.read_text() for p in paths) runs its element at import. Only an unconsumed genexp bound to a name stays lazy, so the walker now follows the consumed ones in full and keeps the outermost-iterable-only treatment for the rest. if "__main__" == __name__ is an equivalent and accepted spelling of the main guard, but requiring __name__ on the left meant its body was treated as import-time code. That is a false positive on a block pytest never runs, so both operand orders are recognised now. The helper table was built from module-level defs only, so a def in a class body invoked while the class is constructed was never followed, contradicting the walker's stated coverage of class bodies. Helpers are now collected from the module body and from class bodies at any nesting. Verified: 15 targeted cases including all three fixes and the earlier ones still classify correctly; still 22 offenders on unpatched main; green on this branch and on the tree merged with latest main. * Handle positional read_text encodings, lazy generators and nested helpers * Guard reads reached from test bodies, unbound Path calls and __file__ paths * [pre-commit.ci] auto fixes from pre-commit.com hooks for more information, see https://pre-commit.ci * Follow derived paths, skip lazy generator helpers, cover compressed openers * Guard the CLI tests, helper parameters and unbound Path arguments * Discover test roots and follow literal, in-place and tuple-derived paths * Identify module openers by import, unwrap starred paths, pin subprocess snippets * Resolve import origins, seed helper locals, follow named generators and parametrize * Scope imports lexically, list tracked test files, bind unpacked names * Resolve aliased openers, keyword-only params, destructured targets, next() * Pin the encoding on subprocess snippets, workflow lint and CLI output for PR #7438 * Harden the CLI encoding guard against detached streams for PR #7438 * Tighten the encoding guard's path and scope analysis for PR #7438 * Resolve path provenance more precisely and keep POSIX stream encodings for PR #7438 * Resolve qualified path classes and scope conditional imports for PR #7438 * Scope CLI stream setup to the entry point and align two encoding pairs for PR #7438 --------- Co-authored-by: pre-commit-ci[bot] <66853113+pre-commit-ci[bot]@users.noreply.github.com> Co-authored-by: danielhanchen <danielhanchen@gmail.com>
271 lines
8.8 KiB
Python
271 lines
8.8 KiB
Python
"""TOCTOU atomicity guards for the cancel path: single _CANCEL_LOCK critical sections; parallel cancel-POST vs __enter__ never drops a cancel."""
|
|
|
|
from __future__ import annotations
|
|
|
|
import ast
|
|
import random
|
|
import threading
|
|
from pathlib import Path
|
|
|
|
|
|
SOURCE_PATH = Path(__file__).resolve().parents[2] / "studio" / "backend" / "routes" / "inference.py"
|
|
_SRC = SOURCE_PATH.read_text(encoding = "utf-8")
|
|
_TREE = ast.parse(_SRC)
|
|
|
|
|
|
def _find_function(name: str) -> ast.FunctionDef | ast.AsyncFunctionDef:
|
|
for node in ast.walk(_TREE):
|
|
if isinstance(node, (ast.FunctionDef, ast.AsyncFunctionDef)) and node.name == name:
|
|
return node
|
|
raise AssertionError(f"function {name!r} not found")
|
|
|
|
|
|
def _find_class(name: str) -> ast.ClassDef:
|
|
for node in ast.walk(_TREE):
|
|
if isinstance(node, ast.ClassDef) and node.name == name:
|
|
return node
|
|
raise AssertionError(f"class {name!r} not found")
|
|
|
|
|
|
def _count_with_cancel_lock_blocks(node: ast.AST) -> int:
|
|
n = 0
|
|
for sub in ast.walk(node):
|
|
if not isinstance(sub, ast.With):
|
|
continue
|
|
for item in sub.items:
|
|
ctx = item.context_expr
|
|
if isinstance(ctx, ast.Name) and ctx.id == "_CANCEL_LOCK":
|
|
n += 1
|
|
break
|
|
return n
|
|
|
|
|
|
def test_cancel_by_cancel_id_or_stash_is_single_lock_critical_section():
|
|
fn = _find_function("_cancel_by_cancel_id_or_stash")
|
|
assert _count_with_cancel_lock_blocks(fn) == 1, (
|
|
"_cancel_by_cancel_id_or_stash must use exactly one `with "
|
|
"_CANCEL_LOCK:` block; splitting into two acquisitions reopens "
|
|
"the TOCTOU race with _TrackedCancel.__enter__"
|
|
)
|
|
src = ast.unparse(fn)
|
|
assert "_CANCEL_REGISTRY.get(cancel_id)" in src
|
|
assert "_PENDING_CANCELS[cancel_id]" in src
|
|
|
|
|
|
def test_tracked_cancel_enter_registers_and_consumes_pending_under_one_lock():
|
|
cls = _find_class("_TrackedCancel")
|
|
enter = None
|
|
for n in cls.body:
|
|
if isinstance(n, ast.FunctionDef) and n.name == "__enter__":
|
|
enter = n
|
|
break
|
|
assert enter is not None
|
|
assert _count_with_cancel_lock_blocks(enter) == 1, (
|
|
"_TrackedCancel.__enter__ must acquire _CANCEL_LOCK exactly once. "
|
|
"A second acquisition for consume-pending lets a concurrent "
|
|
"cancel POST stash after consume sees an empty map, silently "
|
|
"dropping the cancel"
|
|
)
|
|
with_block = None
|
|
for sub in ast.walk(enter):
|
|
if isinstance(sub, ast.With) and any(
|
|
isinstance(i.context_expr, ast.Name) and i.context_expr.id == "_CANCEL_LOCK"
|
|
for i in sub.items
|
|
):
|
|
with_block = sub
|
|
break
|
|
assert with_block is not None
|
|
block_src = "\n".join(ast.unparse(s) for s in with_block.body)
|
|
assert "_CANCEL_REGISTRY.setdefault" in block_src
|
|
assert "_PENDING_CANCELS.pop" in block_src, (
|
|
"__enter__ critical section must consume from _PENDING_CANCELS "
|
|
"inside the same lock, not a later re-acquisition"
|
|
)
|
|
|
|
|
|
def test_cancel_inference_uses_atomic_helper_for_cancel_id_path():
|
|
fn = _find_function("cancel_inference")
|
|
src = ast.unparse(fn)
|
|
assert "_cancel_by_cancel_id_or_stash" in src
|
|
# The pre-fix two-step idiom must be gone.
|
|
assert "_remember_pending_cancel(cancel_id)" not in src, (
|
|
"two-step _cancel_by_keys + _remember_pending_cancel produced "
|
|
"the TOCTOU race and must not return"
|
|
)
|
|
|
|
|
|
_WANTED = {
|
|
"_CANCEL_REGISTRY",
|
|
"_CANCEL_LOCK",
|
|
"_PENDING_CANCELS",
|
|
"_PENDING_CANCEL_TTL_S",
|
|
"_prune_pending",
|
|
"_remember_pending_cancel",
|
|
"_TrackedCancel",
|
|
"_cancel_by_keys",
|
|
"_cancel_by_cancel_id_or_stash",
|
|
}
|
|
|
|
|
|
def _load_registry_module():
|
|
chunks = []
|
|
for n in _TREE.body:
|
|
seg = ast.get_source_segment(_SRC, n)
|
|
if seg is None:
|
|
continue
|
|
if isinstance(n, (ast.FunctionDef, ast.ClassDef)) and n.name in _WANTED:
|
|
chunks.append(seg)
|
|
elif isinstance(n, ast.Assign):
|
|
names = [t.id for t in n.targets if isinstance(t, ast.Name)]
|
|
if any(name in _WANTED for name in names):
|
|
chunks.append(seg)
|
|
elif (
|
|
isinstance(n, ast.AnnAssign)
|
|
and isinstance(n.target, ast.Name)
|
|
and n.target.id in _WANTED
|
|
):
|
|
chunks.append(seg)
|
|
mod = {}
|
|
exec(
|
|
"import threading, time\nfrom typing import Optional\n" + "\n\n".join(chunks),
|
|
mod,
|
|
)
|
|
return mod
|
|
|
|
|
|
def test_parallel_cancel_vs_register_never_drops():
|
|
m = _load_registry_module()
|
|
trials = 500
|
|
dropped = 0
|
|
for i in range(trials):
|
|
m["_CANCEL_REGISTRY"].clear()
|
|
m["_PENDING_CANCELS"].clear()
|
|
cid = f"cid-{i}"
|
|
ev = threading.Event()
|
|
tracker = m["_TrackedCancel"](ev, cid, "thread")
|
|
start = threading.Event()
|
|
|
|
def do_cancel():
|
|
start.wait()
|
|
m["_cancel_by_cancel_id_or_stash"](cid)
|
|
|
|
def do_enter():
|
|
start.wait()
|
|
tracker.__enter__()
|
|
|
|
threads = [
|
|
threading.Thread(target = do_cancel),
|
|
threading.Thread(target = do_enter),
|
|
]
|
|
random.shuffle(threads)
|
|
for t in threads:
|
|
t.start()
|
|
start.set()
|
|
for t in threads:
|
|
t.join(timeout = 5.0)
|
|
assert not t.is_alive()
|
|
|
|
if not ev.is_set():
|
|
dropped += 1
|
|
tracker.__exit__(None, None, None)
|
|
|
|
assert dropped == 0, (
|
|
f"TOCTOU regression: {dropped}/{trials} parallel trials silently " f"dropped the cancel"
|
|
)
|
|
|
|
|
|
def test_cancel_before_register_replays_atomically():
|
|
m = _load_registry_module()
|
|
cid = "early-cid"
|
|
ev = threading.Event()
|
|
tracker = m["_TrackedCancel"](ev, cid, "thread-x")
|
|
|
|
assert m["_cancel_by_cancel_id_or_stash"](cid) == 0
|
|
assert cid in m["_PENDING_CANCELS"]
|
|
|
|
tracker.__enter__()
|
|
assert ev.is_set()
|
|
assert cid not in m["_PENDING_CANCELS"]
|
|
tracker.__exit__(None, None, None)
|
|
|
|
|
|
def test_cancel_after_register_signals_without_stash():
|
|
m = _load_registry_module()
|
|
cid = "post-cid"
|
|
ev = threading.Event()
|
|
tracker = m["_TrackedCancel"](ev, cid, "thread-y")
|
|
tracker.__enter__()
|
|
|
|
assert m["_cancel_by_cancel_id_or_stash"](cid) == 1
|
|
assert ev.is_set()
|
|
assert cid not in m["_PENDING_CANCELS"]
|
|
tracker.__exit__(None, None, None)
|
|
|
|
|
|
def test_cancel_by_keys_tolerates_empty_and_falsy_keys():
|
|
m = _load_registry_module()
|
|
m["_CANCEL_REGISTRY"].clear()
|
|
m["_PENDING_CANCELS"].clear()
|
|
assert m["_cancel_by_keys"]([]) == 0
|
|
assert m["_cancel_by_keys"](["", None, "unknown"]) == 0
|
|
# Non-stashing fallback must never leak into _PENDING_CANCELS.
|
|
assert m["_PENDING_CANCELS"] == {}
|
|
|
|
|
|
def test_cancel_by_keys_fans_out_to_all_streams_on_same_session():
|
|
# Compare mode and other flows launch concurrent streams under a
|
|
# shared session_id; a single session cancel POST must hit all of them.
|
|
m = _load_registry_module()
|
|
m["_CANCEL_REGISTRY"].clear()
|
|
m["_PENDING_CANCELS"].clear()
|
|
session = "shared-thread"
|
|
ev_a = threading.Event()
|
|
ev_b = threading.Event()
|
|
tracker_a = m["_TrackedCancel"](ev_a, "cancel-a", session, "chatcmpl-a")
|
|
tracker_b = m["_TrackedCancel"](ev_b, "cancel-b", session, "chatcmpl-b")
|
|
tracker_a.__enter__()
|
|
tracker_b.__enter__()
|
|
try:
|
|
assert m["_cancel_by_keys"]([session]) == 2
|
|
assert ev_a.is_set() and ev_b.is_set()
|
|
finally:
|
|
tracker_a.__exit__(None, None, None)
|
|
tracker_b.__exit__(None, None, None)
|
|
assert session not in m["_CANCEL_REGISTRY"]
|
|
|
|
|
|
def test_cancel_by_cancel_id_is_exclusive_to_single_run():
|
|
# cancel_id is per-run unique; cancelling run A must not touch run B
|
|
# even when both share a session_id.
|
|
m = _load_registry_module()
|
|
m["_CANCEL_REGISTRY"].clear()
|
|
m["_PENDING_CANCELS"].clear()
|
|
session = "shared-thread-2"
|
|
ev_a = threading.Event()
|
|
ev_b = threading.Event()
|
|
tracker_a = m["_TrackedCancel"](ev_a, "cancel-only-a", session, "chatcmpl-a")
|
|
tracker_b = m["_TrackedCancel"](ev_b, "cancel-only-b", session, "chatcmpl-b")
|
|
tracker_a.__enter__()
|
|
tracker_b.__enter__()
|
|
try:
|
|
assert m["_cancel_by_cancel_id_or_stash"]("cancel-only-a") == 1
|
|
assert ev_a.is_set()
|
|
assert not ev_b.is_set()
|
|
finally:
|
|
tracker_a.__exit__(None, None, None)
|
|
tracker_b.__exit__(None, None, None)
|
|
|
|
|
|
def test_tracked_cancel_exit_is_idempotent():
|
|
# Outer except BaseException + the generator's finally may both call
|
|
# __exit__ under certain race combos; must not raise.
|
|
m = _load_registry_module()
|
|
m["_CANCEL_REGISTRY"].clear()
|
|
m["_PENDING_CANCELS"].clear()
|
|
ev = threading.Event()
|
|
tracker = m["_TrackedCancel"](ev, "cid", "sess", "chatcmpl-x")
|
|
tracker.__enter__()
|
|
tracker.__exit__(None, None, None)
|
|
tracker.__exit__(None, None, None)
|
|
tracker.__exit__(None, None, None)
|
|
assert not m["_CANCEL_REGISTRY"]
|