Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
16 commits
Select commit Hold shift + click to select a range
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
45 changes: 35 additions & 10 deletions blueprints/us-equities/mover-v3/protocol-core-draft.json

Large diffs are not rendered by default.

37 changes: 27 additions & 10 deletions blueprints/us-equities/mover-v3/study/core/count_only.py
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,7 @@
from core.fills import fill_at
from core.params import FETCH, SAMPLING, T
from core.records import MERGER_TYPES
from core.store import Store
from core.store import Store, checked_run_identity, recover_seal
from pinned.sessions_io_copy import official_price

PART1_RANGE = ("2016-01-04", "2020-12-31")
Expand Down Expand Up @@ -464,28 +464,45 @@ def dry_run_requests(cal, sessions: list, symbols: list) -> list:
return out


def dry_run(cal, sessions: list, symbols: list, transports: dict, snapshot_root, fetch_date: str,
def dry_run(cal, sessions: list, symbols: list, transports: dict, snapshot_root, utc_start: str, run_identity: dict,
clock=driver.utc_now, progress: dict | None = None) -> dict:
"""freeze_preconditions: the native dry run through core/driver.py and the transport, sealed, re-read from the
seal and counted (counts only; run.py dry-run writes the output with its run-log line). progress, when given,
receives the snapshot sha256 as soon as it is sealed, so a run that fails afterwards still logs it (review round
12, F3)."""
dry_run_requests(cal, sessions, symbols) # refuses before any fetch (review round 11, C1)
store = Store()
12, F3).

A snapshot already sealed under <snapshot_root>/dry-run (a run that failed or was killed after its seal) is
adopted without fetching only through core.store.recover_seal (review round 18, R5, and its repair): its seal
record verifies the ledger and pages (G-H), and its binding names the running study tree, protocol and runtime
lock (run_identity; H1, so a tree whose fetch plumbing never ran cannot take another tree's seal), the sessions,
the symbols and the request plan. utc_start, the fetch start, is bound beside that identity and read back from
the seal into the output (fetch_utc_start), never matched (M1): a retry on a later UTC day adopts the seal."""
reqs = dry_run_requests(cal, sessions, symbols) # refuses before any fetch (review round 11, C1)
if not isinstance(utc_start, str) or not utc_start:
raise ValueError("the dry run needs its fetch start (utc_start)")
root = Path(snapshot_root) / "dry-run"
driver.stage_fetch(lambda st: dry_run_requests(cal, sessions, symbols), transports, store, fetch_date,
clock=clock)
sha = store.write(root)
seal_identity = {**checked_run_identity(run_identity), "sessions": list(sessions), "symbols": list(symbols),
"requests": sorted(plan.record(r) for r in reqs)}
recovered = recover_seal(root, seal_identity, "dry-run", extra=("utc_start",))
if recovered is None:
store = Store()
driver.stage_fetch(lambda st: reqs, transports, store, utc_start[:10], clock=clock)
sha = store.write(root, binding={"identity": seal_identity, "utc_start": utc_start}, seal_record=True)
else:
sha = recovered[1]
if progress is not None:
progress["dry-run"] = sha
return dry_run_output(snapshot_root, sha, sessions, symbols)


def dry_run_output(snapshot_root, sha: str, sessions: list, symbols: list) -> dict:
"""The dry-run output from its sealed snapshot alone (review round 15, N02 / R14-open-2: run.py dry-run
recomputes it to adopt an output that a hard kill left with no run-log line)."""
recomputes it to adopt an output that a hard kill left with no run-log line). fetch_utc_start is the fetch start
the seal's binding holds (review round 18 repair, M1), provenance metadata rather than market data; None for a
seal without one."""
sealed = Store.read(Path(snapshot_root) / "dry-run", sha)
return {"kind": "mover_v3_dry_run_output", "snapshot_sha256": sha, "sessions": list(sessions),
return {"kind": "mover_v3_dry_run_output", "snapshot_sha256": sha,
"fetch_utc_start": (sealed.binding or {}).get("utc_start"), "sessions": list(sessions),
"symbols_count": len(symbols), "counts": dry_run_counts(sealed),
"incomplete_by_kind": sealed.incomplete_by_kind(), "measured": dry_run_measured(sealed)}

Expand Down
109 changes: 103 additions & 6 deletions blueprints/us-equities/mover-v3/study/core/guards.py
Original file line number Diff line number Diff line change
Expand Up @@ -17,10 +17,16 @@
time comes only from a commit signed by the pinned GitHub web-flow key (verified_merge, M3); a transport deviation
governs only with the hashed output of a logged run.py transport-check run (check_study_tree, H1); void deviations
are typed (voids, F4); and every locked package's version and dist-info RECORD digest is checked (check_runtime, L3).

Review round 18, second repair: the running tree is HEAD's only when no index flag hides a tracked file and every
tracked file's own bytes are its blob (running_tree, R2-2), and every git command this module runs ignores replace
refs (git_command, follow-up F2), so each guard reads the repository's own commits, trees and blobs.
"""
from __future__ import annotations

import json
import os
import stat
import subprocess
import sys
from pathlib import Path
Expand All @@ -43,28 +49,45 @@ class Refused(Exception):
pass


def git_command(repo, *args) -> list:
"""The argument list of every git command this module runs (git, git_ok, committed_bytes, verified_merge and the
running-tree checks). Review round 18, second repair, follow-up F2: --no-replace-objects (git(1)) makes git ignore
replace refs (git-replace(1)), with which a local refs/replace entry puts another commit, tree or blob in place of
the one a guard names. So a guard reads the repository's own objects: the tree a run names, the protocol and logs
at a commit, the tree diff behind a transport deviation, the commits whose times and signatures it uses."""
return ["git", "--no-replace-objects", "-C", str(repo), *args]


def git(repo, *args) -> str:
return subprocess.check_output(["git", "-C", str(repo), *args], text=True).strip()
return subprocess.check_output(git_command(repo, *args), text=True).strip()


def git_ok(repo, *args) -> bool:
return subprocess.run(["git", "-C", str(repo), *args], stdout=subprocess.DEVNULL,
return subprocess.run(git_command(repo, *args), stdout=subprocess.DEVNULL,
stderr=subprocess.DEVNULL).returncode == 0


def committed_bytes(repo, ref: str, path: str):
"""The bytes of path at ref, or None if the path does not exist there."""
if not git_ok(repo, "cat-file", "-e", f"{ref}:{path}"):
return None
return subprocess.check_output(["git", "-C", str(repo), "show", f"{ref}:{path}"])
return subprocess.check_output(git_command(repo, "show", f"{ref}:{path}"))


# ---------------------------------------------------------------- the running tree

def running_tree(repo, study_path: str = STUDY_PATH) -> str:
"""The git tree hash of the study directory at HEAD; refuses if the working copy of it differs from HEAD, or
if it holds an ignored file other than a __pycache__ entry (review round 9, F9: an ignored file is invisible to
the tree hash but could still be imported)."""
the tree hash but could still be imported).

Review round 18, second repair (R2-2): git status alone does not show that the working copy is HEAD's. It reports
no change to a tracked file whose index entry is marked assume-unchanged or skip-worktree (git-update-index(1)),
nor one that a clean filter maps back to the committed bytes, nor files that match a tree which a replace ref puts
in place of HEAD's (git-replace(1)). So an index entry with either flag is refused, and every file of the tree
returned must be a regular file whose own bytes hash to its blob: the tree is read without replace refs
(git_command), and the files without the index and without filters. Both checks run here, once, when a run builds
its context: a file changed later in the same run is not seen."""
dirty = git(repo, "status", "--porcelain", "--untracked-files=all", "--", study_path)
if dirty:
raise Refused("the study tree has uncommitted changes; a run executes a committed tree only")
Expand All @@ -73,7 +96,81 @@ def running_tree(repo, study_path: str = STUDY_PATH) -> str:
planted = [p for p in ignored if "__pycache__" not in Path(p).parts]
if planted:
raise Refused(f"the study tree holds ignored files that are not bytecode caches: {planted[:5]}")
return git(repo, "rev-parse", f"HEAD:{study_path}")
masked = _masked_index_entries(repo, study_path)
if masked:
# the independent verification of ace56485 (C7): with core.ignoreStat=true every file git checks out is marked
# assume-unchanged, so a bit need not hide a change; the refusal stays, and the message says how to clear it
raise Refused(f"the index marks study files assume-unchanged or skip-worktree, so git status does not check "
f"them and the run cannot verify them: {masked[:5]}; clear the bits (git update-index "
"--no-assume-unchanged or --no-skip-worktree on those files) and, if core.ignoreStat is true, "
"unset it, or git marks every file it checks out assume-unchanged again")
tree = git(repo, "rev-parse", f"HEAD:{study_path}")
unverified = _unverified_files(repo, tree, study_path)
if unverified:
raise Refused(f"the study tree's files are not all regular files that hold the committed bytes, although git "
f"status reports no change: {unverified[:5]}; a run executes a committed tree only")
return tree


def _masked_index_entries(repo, study_path: str) -> list:
"""Every index entry under study_path that is marked assume-unchanged or skip-worktree, as '<path> (<flags>)'.
git-ls-files(1) -v prints the tag of an assume-unchanged entry in lowercase and that of a skip-worktree entry as S
(s with both flags)."""
out = subprocess.check_output(git_command(repo, "ls-files", "-v", "-z", "--", study_path))
masked = []
for row in out.split(b"\0"):
if row:
tag = row[:1].decode("ascii", "replace")
flags = [name for name, marked in (("assume-unchanged", tag.islower()), ("skip-worktree", tag in "Ss"))
if marked]
if flags:
masked.append(f"{os.fsdecode(row[2:])} ({' and '.join(flags)})")
return masked


def _unverified_files(repo, tree: str, study_path: str) -> list:
"""The files of `tree` (the study tree at HEAD) that the working copy does not hold as committed, each with its
problem. An entry must be a regular-file blob: a symbolic link or submodule entry is refused, because the tree hash
covers a link's text and not the file it names. Its working file must be a regular file, and its bytes must hash to
the entry's blob. One git hash-object --no-filters --stdin-paths call hashes the files themselves, so the result
depends on no index entry (its flags or its cached stat data) and on no clean filter. The tree is listed without
replace refs (git_command), so the blobs are its own and not those of a tree that a replace ref puts in its
place."""
out = subprocess.check_output(git_command(repo, "ls-tree", "-r", "-z", tree))
problems, paths, blobs = [], [], []
for row in out.split(b"\0"):
if not row:
continue
meta, name = row.split(b"\t", 1)
mode, kind, blob = meta.decode("ascii").split()
path = f"{study_path.rstrip('/')}/{os.fsdecode(name)}"
if kind != "blob" or mode not in ("100644", "100755"):
problems.append(f"{path} is a symbolic link or submodule entry (mode {mode})")
elif any(ord(c) < 32 for c in path):
problems.append(f"{path!r} has a control character in its name") # --stdin-paths reads a path per line
else:
try:
file_mode = os.lstat(Path(repo) / path).st_mode # the entry itself: a link is not followed
except OSError:
problems.append(f"{path} is missing")
else:
if stat.S_ISREG(file_mode):
paths.append(path)
blobs.append(blob)
else:
problems.append(f"{path} is not a regular file")
if paths:
base = Path(repo).resolve()
hashed = subprocess.run(git_command(repo, "hash-object", "--no-filters", "--stdin-paths"),
input=b"".join(os.fsencode(str(base / p)) + b"\n" for p in paths),
capture_output=True)
got = hashed.stdout.decode("ascii", "replace").split()
if hashed.returncode != 0 or len(got) != len(paths):
problems.append("git hash-object could not read every tracked file: "
+ hashed.stderr.decode("utf-8", "replace").strip()[:200])
else:
problems.extend(f"{p} differs from HEAD's blob" for p, want, have in zip(paths, blobs, got) if have != want)
return problems


def check_bytecode_isolated() -> None:
Expand Down Expand Up @@ -117,7 +214,7 @@ def verified_merge(repo, commit: str, key: dict | None = None) -> bool:
email = git(repo, "show", "-s", "--format=%ce", commit)
if email == key["committer_email"]:
env = {"GNUPGHOME": _gnupg_home(key), "PATH": __import__("os").environ.get("PATH", "/usr/bin:/bin")}
res = subprocess.run(["git", "-C", str(repo), "-c", "gpg.program=gpg", "verify-commit", "--raw", commit],
res = subprocess.run(git_command(repo, "-c", "gpg.program=gpg", "verify-commit", "--raw", commit),
capture_output=True, text=True, env=env)
for line in res.stderr.splitlines():
parts = line.split()
Expand Down
8 changes: 8 additions & 0 deletions blueprints/us-equities/mover-v3/study/core/runner.py
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,14 @@ def count_only_dependencies(ctx: dict, protocol: dict, rate: dict, allowance: di
"rate_limit": dict(rate), "pipeline_allowance": dict(allowance) if allowance is not None else None}


def seal_run_identity(ctx: dict) -> dict:
"""The running study tree, protocol sha256 and runtime-lock sha256 that a recoverable seal binds, as
core.holdout.collect's identity does (review round 18 repair, H1): the native dry run's snapshot and each
transport-check live sample (core.store.recover_seal)."""
return {"study_tree": ctx["tree"], "protocol_sha256": ctx["protocol_sha256"],
"runtime_lock_sha256": ctx["runtime_lock_sha256"]}


def adopt_uncited_output(out_path, expected: dict, ignore=("code_revision",)) -> str:
"""Review round 15, N02 (R14-open-2): a pre-freeze output that a hard kill left with no run-log line (the process
died between the output's rename and its line) is adopted, never fetched again: its content must equal the
Expand Down
Loading
Loading