Repository navigation
Expand file tree
/
Copy pathanalytics.py
More file actions
executable file
·336 lines (305 loc) · 15.8 KB
/
Copy pathanalytics.py
File metadata and controls
executable file
·336 lines (305 loc) · 15.8 KB
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
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
#!/usr/bin/env python3
"""Optional deterministic read-only analytics. No installation or wallet access."""
import time
_PROCESS_STARTED = time.monotonic()
import sys
sys.dont_write_bytecode = True
from pathlib import Path
# Isolated Python (-I) deliberately omits the script directory from sys.path.
_SCRIPT_DIRECTORY = Path(__file__).resolve().parent
if str(_SCRIPT_DIRECTORY) not in sys.path:
sys.path.insert(0, str(_SCRIPT_DIRECTORY))
import argparse
import hashlib
import math
import os
import re
import select
from netstack_core import Context, RpcError, SERIALIZATION_RESERVE, StopRun, load_json, serialize_result
from netstack_output import full_fallback, summarize_advance, summarize_rfv
class ArgumentError(Exception):
pass
class Parser(argparse.ArgumentParser):
def error(self, message):
# Argparse messages can contain arbitrary user input; do not echo it.
raise ArgumentError("Invalid arguments; use --help for accepted read-only options")
def positive_seconds(text):
try:
value = float(text)
except ValueError as exc:
raise argparse.ArgumentTypeError("must be a positive number") from exc
if not math.isfinite(value) or not 0 < value <= 600:
raise argparse.ArgumentTypeError("must be positive and at most 600 seconds")
return value
def positive_days(text):
try:
value = float(text)
except ValueError as exc:
raise argparse.ArgumentTypeError("must be a positive number") from exc
if not math.isfinite(value) or value <= 0 or value > 1000000:
raise argparse.ArgumentTypeError("must be finite, positive and at most 1000000 days")
return value
def block_spec(text):
if text != "latest-2" and (len(text) > 78 or not re.fullmatch(r"0|[1-9][0-9]*", text)):
raise argparse.ArgumentTypeError("must be latest-2 or a nonnegative block number")
return text
def series_id(text):
if len(text) > 78 or not re.fullmatch(r"[1-9][0-9]*", text) or int(text) >= 2**256:
raise argparse.ArgumentTypeError("must be a positive uint256 series ID")
return int(text)
def parser():
result = Parser(
description="Read-only Robinhood mainnet analytics using one pinned block and the fixed public HTTPS RPC.",
epilog="No wallet, authentication, proxy, redirect, installation, or write RPC. JSON is emitted even for partial collections. Exit 0: requested collection completed; 2: partial/blocked; 1: invalid arguments/fatal error.")
commands = result.add_subparsers(dest="command", required=True)
for command, help_text in (("lp", "NET/USDG v2 LP holders, flows and conditional gross fees"),
("predict", "selected Predict Desk series and observed trading flows"),
("house", "selected House Vault accounts, queues and observed cash flows"),
("rfv", "reconciled Core RFV and RPC-backed Sleeve asset accounting"),
("advance", "NET Advance positions, holders, totals, capacity and parameters")):
child = commands.add_parser(command, help=help_text, description=help_text)
child.add_argument("--deadline", type=positive_seconds, default=120.0, metavar="SECONDS",
help="whole-process walltime, including a 2-second serialization reserve (default: 120; maximum: 600)")
child.add_argument("--block", type=block_spec, default="latest-2", metavar="BLOCK",
help="latest-2 or an explicit nonnegative block number (default: latest-2)")
if command == "lp":
child.add_argument("--since-days", type=positive_days, default=7.0, metavar="DAYS",
help="exact historical fee interval ending at the pinned block (default: 7)")
else:
child.set_defaults(since_days=None)
if command == "predict":
child.add_argument("--series", type=series_id, metavar="ID",
help="single-series pinned snapshot only; no logs, history or quote (default: all series and history)")
if command == "advance":
child.add_argument("action", nargs="?", choices=("provenance",),
help="offline Desk/Zap catalog provenance; no RPC calls")
child.add_argument("--view", choices=("positions", "holders", "totals", "capacity", "params"), default="totals",
help="requested pinned evidence view (default: totals); capacity/params do not enumerate positions")
if command == "rfv":
child.add_argument("--scope", choices=("core", "reports", "net-assets"), default="core",
help="Core reserves, publisher Reports composition, or adjusted net assets (default: core)")
if command in ("rfv", "advance"):
child.add_argument("--detail", choices=("summary", "full"), default="full",
help="stdout evidence detail (default: full); --output always preserves full checkpoints")
child.add_argument("--json", action="store_true", help="emit machine-readable JSON (also the default)")
child.add_argument("--output", metavar="PATH", help="atomic partial/final checkpoint outside the installed package; symlinks refused")
return result
def _failure(command, message):
return {"schema_version": 1, "command": command, "snapshot": {}, "metrics": {}, "coverage": {},
"not_proven": [], "errors": [{"kind": "input", "message": message}],
"status": "failed", "stopping_reason": "invalid_arguments"}
def _provenance(ctx):
ctx.check()
version = load_json("release-manifest.json").get("version")
files = {}
modules = ["analytics.py", "netstack_core.py"]
if ctx.command == "advance":
modules.extend(("netstack_advance.py", "netstack_output.py"))
else:
modules.extend(("netstack_reserves.py", "netstack_sleeve.py", "netstack_v4.py", "netstack_discovery.py", "netstack_methodology.py", "netstack_output.py") if ctx.command == "rfv"
else ("netstack_lp.py",) if ctx.command == "lp" else ("netstack_markets.py",))
for name in modules:
ctx.check()
try:
with (_SCRIPT_DIRECTORY / name).open("rb") as handle:
data = handle.read(1024 * 1024 + 1)
except OSError as exc:
raise RpcError("Required packaged analytics module is unavailable", kind="package") from exc
if len(data) > 1024 * 1024:
raise RpcError("Packaged analytics module exceeds size limit", kind="package")
files["scripts/" + name] = hashlib.sha256(data).hexdigest()
ctx.result["provenance"] = {"package_version": version, "runtime_sha256": files,
"meaning": "Hashes identify the local code used; they are not an authenticity or deployed-contract verification claim."}
if ctx.command == "advance":
interface = load_json("assets/analytics/advance-interface.json")
ctx.result["provenance"]["analysis"] = {
"interface": {key: interface[key] for key in
("source_id", "source_url", "source_sha256", "reviewed_on")},
"runtime_evidence": interface["runtime_evidence"],
"meaning": "Unofficial independent catalog analysis, not a NetNet publication or current runtime match. Source IDs and JSON Pointer references identify evidence, not extra observations."}
def _pending_coverage(value):
if isinstance(value, dict):
if value.get("event_coverage_complete") is False or value.get("collection_complete") is False:
return True
return any(_pending_coverage(item) for item in value.values())
if isinstance(value, list):
return any(_pending_coverage(item) for item in value)
return False
def _emit(text, hard_end=None):
"""Do not let a blocked stdout pipe outlive the process deadline."""
try:
descriptor = sys.stdout.fileno()
except (AttributeError, OSError):
sys.stdout.write(text)
sys.stdout.flush()
return True
blocking = os.get_blocking(descriptor)
try:
os.set_blocking(descriptor, False)
data = memoryview(text.encode("ascii"))
while data:
try:
count = os.write(descriptor, data)
if count <= 0:
return False
data = data[count:]
except BlockingIOError:
remaining = 0.0 if hard_end is None else max(0.0, hard_end - time.monotonic())
if not remaining or not select.select([], [descriptor], [], remaining)[1]:
return False
return True
finally:
os.set_blocking(descriptor, blocking)
def main(argv=None):
ctx = None
try:
args = parser().parse_args(argv)
if (args.command == "advance" and args.action == "provenance"
and (args.view != "totals" or args.block != "latest-2")):
raise ArgumentError("Offline provenance cannot select a live view or block")
except ArgumentError as exc:
sys.stdout.write(serialize_result(_failure(None, str(exc))))
return 1
result = _failure(args.command, "Analytics initialization failed")
exit_code = 1
try:
remaining = args.deadline - (time.monotonic() - _PROCESS_STARTED)
if remaining <= 0:
result = {"schema_version": 1, "command": args.command, "snapshot": {}, "metrics": {},
"coverage": {}, "not_proven": [], "errors": [], "status": "partial",
"stopping_reason": "deadline_exhausted",
"elapsed_seconds": time.monotonic() - _PROCESS_STARTED,
"collector_deadline_seconds": args.deadline}
exit_code = 2
else:
ctx = Context(args.command, deadline=remaining, output=args.output)
ctx._started = _PROCESS_STARTED
ctx._deadline = args.deadline
ctx._hard_end = _PROCESS_STARTED + args.deadline
ctx._collect_end = ctx._hard_end - SERIALIZATION_RESERVE
ctx._arm()
result = ctx.result
if args.output is not None:
ctx._open_output_parent()
_provenance(ctx)
offline = args.command == "advance" and args.action == "provenance"
if offline:
from netstack_advance import provenance
provenance(ctx)
else:
ctx.pin(args.block)
if args.command == "lp":
from netstack_lp import run
elif args.command == "predict":
from netstack_markets import run_predict as run
elif args.command == "rfv":
from netstack_reserves import run
elif args.command == "advance":
from netstack_advance import run
else:
from netstack_markets import run_house as run
run(ctx, args)
ctx.check()
ctx.recheck()
requested = result["coverage"].get("requested_scope") if args.command in ("rfv", "advance") else None
incomplete = (not requested.get("collection_complete", False) if requested is not None
else bool(result["errors"] or _pending_coverage(result["coverage"])))
if incomplete:
result["status"] = "partial"
result["stopping_reason"] = "incomplete_collection"
exit_code = 2
else:
result["status"] = "completed"
result["stopping_reason"] = "completed_requested_collection"
exit_code = 0
except StopRun as exc:
result["status"] = "partial"
result["stopping_reason"] = exc.reason
exit_code = 2
except RpcError as exc:
result["errors"].append({"kind": exc.kind, "message": str(exc)})
result["status"] = "partial" if ctx is not None else "failed"
result["stopping_reason"] = "rpc_or_evidence_failure"
exit_code = 2 if ctx is not None else 1
except KeyboardInterrupt:
result["status"] = "partial"
result["stopping_reason"] = "interrupted_by_SIGINT"
exit_code = 2
except Exception as exc:
# Never expose source/provider text, arbitrary paths, credentials or traceback data.
result["errors"].append({"kind": "fatal", "message": "Internal analytics failure (" + type(exc).__name__ + ")"})
result["status"] = "failed"
result["stopping_reason"] = "fatal_error"
exit_code = 1
# A useful partial collection can still confirm B when retrieval failed
# without cancellation or permission denial. Never restart after a stop.
if (ctx is not None and ctx.block is not None and ctx._stopped is None
and ctx._terminal_error is None
and result["snapshot"].get("recheck_status") == "not_performed"
and not any(isinstance(error, dict) and error.get("kind") == "permission"
for error in result["errors"])):
try:
ctx.recheck()
except StopRun as exc:
result["stopping_reason"] = exc.reason
result["status"] = "partial"
exit_code = 2
except RpcError as exc:
result["errors"].append({"kind": exc.kind, "message": str(exc)})
result["status"] = "partial"
exit_code = 2
text = None
checkpoint_status = "not_requested" if args.output is None else "not_confirmed"
try:
if ctx is not None:
ctx.begin_finalization()
if result.get("snapshot") and result["snapshot"].get("recheck_status") != "confirmed":
result["snapshot"].setdefault("recheck_status", "not_performed")
result["snapshot"]["confirmation"] = (
"invalid" if result["snapshot"]["recheck_status"] == "mismatch" else "unconfirmed")
ctx.checkpoint()
if args.output is not None:
checkpoint_status = "saved"
text = ctx._last_json
else:
text = serialize_result(result)
except RpcError as exc:
result["errors"].append({"kind": exc.kind, "message": str(exc)})
result["status"] = "partial"
result["stopping_reason"] = "checkpoint_failure"
exit_code = 2
try:
text = serialize_result(result)
except (StopRun, KeyboardInterrupt):
text = ctx.partial_json(ctx._stopped or "deadline_exhausted")
except (StopRun, KeyboardInterrupt):
exit_code = 2
# Reuse a valid aggregate body without repeating expensive serialization.
text = ctx.partial_json(ctx._stopped or "deadline_exhausted") if ctx is not None else serialize_result(result)
if args.command in ("rfv", "advance") and args.detail == "summary":
try:
# Finalization can reuse an older valid checkpoint or add an output
# error. Project exactly that document, not the mutable ctx.result.
if args.command == "advance":
text = summarize_advance(text, checkpoint_status, args.output)
else:
text = summarize_rfv(text, checkpoint_status)
except (StopRun, KeyboardInterrupt) as exc:
exit_code = 2
reason = exc.reason if isinstance(exc, StopRun) else "interrupted_by_SIGINT"
text = full_fallback(text, checkpoint_status, reason)
try:
if ctx is not None:
ctx._emitting = True
if not _emit(text, ctx._hard_end if ctx is not None else None):
return 2
if ctx is not None and ctx._stopped:
exit_code = 2
except (BrokenPipeError, OSError, StopRun, KeyboardInterrupt):
return 2
finally:
if ctx is not None:
ctx.close()
return exit_code
if __name__ == "__main__":
raise SystemExit(main())