MB-X Bilibili Pipeline
6 days ago 643d038b717c97958e7e9dac25fb67d56645f12a
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
"""Cooperative authoritative reader for DIRECT_STABLE_PATH_SET_V1.
 
The stable paths are intentionally readable by ordinary filesystem tools, but
only this double-read protocol may announce them as the current generation.
It never returns a partial payload: lock presence yields RETRY and any identity
or exact-set mismatch yields STOP_NOT_CURRENT.
"""
 
from __future__ import annotations
 
import os
from pathlib import Path
from typing import Any
 
from . import publisher as base
 
 
READER_SCHEMA = "SHARED_CONTENT_PUBLISHER_AUTHORITATIVE_READER_V1"
 
 
def _early_paths(config: dict[str, Any], config_path: Path) -> tuple[Path, Path, Path, Path]:
    if config.get("schema_version") != base.DIRECT_CONFIG_SCHEMA:
        raise base.PublisherError("READER_CONFIG_SCHEMA", str(config.get("schema_version")))
    commit = config.get("commit")
    roots = config.get("roots")
    expected = config.get("expected")
    if not isinstance(commit, dict) or commit.get("strategy") != base.DIRECT_COMMIT_STRATEGY:
        raise base.PublisherError("READER_COMMIT_STRATEGY", str(commit))
    if not isinstance(roots, dict) or not isinstance(expected, dict):
        raise base.PublisherError("READER_CONFIG_TYPE", "roots/expected")
    operation_root = Path(os.path.abspath(base._need_string(roots.get("operation_root"), "roots.operation_root")))
    base._physical_chain(operation_root, final_kind="dir")
    if base._norm(operation_root) != base._norm(base._need_string(roots.get("resolved_root"), "roots.resolved_root")):
        raise base.PublisherError("READER_ROOT_RESOLUTION", str(operation_root))
    if os.path.commonpath([base._norm(operation_root), base._norm(config_path)]) != base._norm(operation_root):
        raise base.PublisherError("READER_CONFIG_OUTSIDE_ROOT", str(config_path))
    lock_path = base._inside(operation_root, base._safe_relative(roots.get("lock_path"), "roots.lock_path"))
    case_index = base._inside(operation_root, base._safe_relative(roots.get("case_current_index_path"), "roots.case_current_index_path"))
    candidate = expected.get("candidate")
    if not isinstance(candidate, dict):
        raise base.PublisherError("READER_CONFIG_TYPE", "expected.candidate")
    rows = base._parse_artifact_rows(candidate.get("rows"), "expected.candidate.rows")
    manifest_rel = base._safe_relative(candidate.get("manifest_self_formal_relative_path"), "expected.candidate.manifest_self_formal_relative_path")
    manifest_matches = [row for row in rows if row.formal_relative_path.casefold() == manifest_rel.casefold()]
    if len(manifest_matches) != 1:
        raise base.PublisherError("READER_MANIFEST_ROW", str(len(manifest_matches)))
    current_manifest = base._inside(operation_root, manifest_matches[0].formal_relative_path)
    return operation_root, lock_path, case_index, current_manifest
 
 
def _identity_or_stop(path: Path, label: str) -> tuple[int, str]:
    try:
        return base._identity(path)
    except base.PublisherError:
        raise
    except Exception as exc:
        raise base.PublisherError("READER_IDENTITY", f"{label}:{type(exc).__name__}:{exc}", exit_code=20) from exc
 
 
def read_direct_current(config_path: Path, *, include_payload: bool = False) -> dict[str, Any]:
    """Return authoritative identities, or a retry/stop status with no payload."""
 
    config_path = Path(os.path.abspath(config_path))
    config, config_data = base._load_config(config_path)
    operation_root, lock_path, case_index, current_manifest = _early_paths(config, config_path)
    common = {
        "schema_version": READER_SCHEMA,
        "commit_strategy": base.DIRECT_COMMIT_STRATEGY,
        "config_bytes": len(config_data),
        "config_sha256": base._sha_bytes(config_data),
        "resolved_root": os.path.abspath(operation_root),
    }
    if base._exists(lock_path):
        return {**common, "status": "RETRY_PUBLISH_IN_PROGRESS", "exit_code": 75, "payload_returned": False, "error_code": "PUBLISH_LOCK_PRESENT"}
    try:
        before_case = _identity_or_stop(case_index, "case_current_index")
        before_manifest = _identity_or_stop(current_manifest, "current_manifest")
        # Full strict validation is intentionally between the two identity reads.
        # It validates every declared stable target, current-manifest coverage,
        # link/evidence closure, exact scope sets, and the complete receipt chain.
        from .direct_stable_path import (
            _destination_exact_set,
            _validate_direct_config,
            _verify_candidate_destinations,
        )
 
        plan = _validate_direct_config(config, config_path, config_data)
        if plan.lifecycle != "REPLAY":
            raise base.PublisherError("READER_NOT_COMMITTED", plan.lifecycle, exit_code=20)
        _verify_candidate_destinations(plan, "READER_TARGET_SET")
        payloads: dict[str, bytes] = {}
        for row in plan.candidate_rows:
            data = base._read_bytes(plan.destination_paths[row.member_id])
            if (len(data), base._sha_bytes(data)) != (row.bytes, row.sha256):
                raise base.PublisherError("READER_PAYLOAD_IDENTITY", row.member_id, exit_code=20)
            payloads[row.formal_relative_path] = data
        after_case = _identity_or_stop(case_index, "case_current_index")
        after_manifest = _identity_or_stop(current_manifest, "current_manifest")
        if base._exists(lock_path):
            return {**common, "status": "RETRY_PUBLISH_IN_PROGRESS", "exit_code": 75, "payload_returned": False, "error_code": "PUBLISH_LOCK_PRESENT_AFTER_READ"}
        if before_case != after_case or before_manifest != after_manifest:
            return {**common, "status": "STOP_NOT_CURRENT", "exit_code": 20, "payload_returned": False, "error_code": "CURRENT_IDENTITY_DRIFT"}
        return {
            **common,
            "status": "CURRENT_READ_COMPLETE",
            "exit_code": 0,
            "payload_returned": include_payload,
            "current_announced": True,
            "case_current_index_bytes": after_case[0],
            "case_current_index_sha256": after_case[1],
            "current_manifest_bytes": after_manifest[0],
            "current_manifest_sha256": after_manifest[1],
            "destination_count": len(plan.candidate_rows),
            "destination_set_sha256": _destination_exact_set(plan.candidate_rows),
            **({"payloads": payloads} if include_payload else {}),
        }
    except base.PublisherError as exc:
        if base._exists(lock_path):
            return {**common, "status": "RETRY_PUBLISH_IN_PROGRESS", "exit_code": 75, "payload_returned": False, "error_code": "PUBLISH_LOCK_PRESENT_DURING_READ"}
        return {**common, "status": "STOP_NOT_CURRENT", "exit_code": 20, "payload_returned": False, "error_code": exc.code, "detail": exc.detail}
 
 
def read_direct_current_snapshot(config_path: Path) -> dict[str, Any]:
    """Python API returning the full frozen byte snapshot only on success."""
 
    return read_direct_current(config_path, include_payload=True)