betterwithage Claude Opus 4.7 commited on
Commit
85bad5c
·
verified ·
1 Parent(s): 97e31a9

deploy(hf): sync szl-holdings/a11oy@5c711ffb13b8ce4e0517f03c73700eaf8bfea3bb derived COPY set

Browse files

Reusable Dockerfile-COPY-derived deploy from szl-holdings/a11oy 5c711ffb13b8ce4e0517f03c73700eaf8bfea3bb.
Files: 1179 Pruned: 0
Derived from Dockerfile COPY sources (NO hand-maintained allowlist).

Signed-off-by: SZL Holdings <noreply@szlholdings.ai>
Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>

Files changed (3) hide show
  1. gdw_proofs.py +97 -15
  2. gdw_runtime.py +29 -3
  3. routers/gdw_frontier.py +39 -0
gdw_proofs.py CHANGED
@@ -1,15 +1,30 @@
1
  """Structured theorem-input export for asynchronous Lean checking."""
2
 
 
3
  import hashlib
4
  import json
5
  import os
 
6
  import tempfile
7
  import threading
 
8
  from pathlib import Path
9
  from typing import Any, Dict
10
 
11
 
12
  _ARTIFACT_QUOTA_LOCK = threading.RLock()
 
 
 
 
 
 
 
 
 
 
 
 
13
 
14
 
15
  def canonical_json(value: Any) -> str:
@@ -20,6 +35,46 @@ def sha256_json(value: Any) -> str:
20
  return hashlib.sha256(canonical_json(value).encode("utf-8")).hexdigest()
21
 
22
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
23
  def build_proof_payload(
24
  proposal_id: str,
25
  request_id: str,
@@ -89,11 +144,12 @@ def _export_json_artifact_unlocked(
89
  destination = owner_root / filename
90
  encoded = (json.dumps(payload, indent=2, sort_keys=True) + "\n").encode("utf-8")
91
  expected_sha256 = hashlib.sha256(encoded).hexdigest()
92
- if destination.exists():
93
- if destination.read_bytes() != encoded:
94
  raise FileExistsError(
95
  "refusing to overwrite an existing non-identical GDW artifact"
96
  )
 
97
  return {
98
  "status": "EXPORTED",
99
  "path": str(destination),
@@ -101,6 +157,7 @@ def _export_json_artifact_unlocked(
101
  "reused": True,
102
  "immutable": True,
103
  "owner_scope": owner_scope,
 
104
  }
105
  owner_limit = _bounded_artifact_limit(
106
  "GDW_OWNER_MAX_ARTIFACTS", default=10_000, maximum=100_000
@@ -116,32 +173,57 @@ def _export_json_artifact_unlocked(
116
  raise RuntimeError("per-owner artifact quota exceeded")
117
  if sum(1 for _ in root.glob("*/*.json")) >= global_limit:
118
  raise RuntimeError("global artifact quota exceeded")
119
- handle, temporary = tempfile.mkstemp(
120
- prefix=".gdw-artifact-", suffix=".tmp", dir=owner_root
 
 
 
 
 
121
  )
 
 
122
  try:
123
  with os.fdopen(handle, "wb") as stream:
124
  stream.write(encoded)
125
  stream.flush()
126
  os.fsync(stream.fileno())
127
- try:
128
- os.link(temporary, destination)
129
- except FileExistsError:
130
- if destination.read_bytes() != encoded:
131
- raise FileExistsError(
132
- "refusing to overwrite a concurrently created "
133
- "non-identical GDW artifact"
134
- )
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
135
  finally:
136
- if os.path.exists(temporary):
137
- os.unlink(temporary)
 
 
138
  return {
139
  "status": "EXPORTED",
140
  "path": str(destination),
141
  "sha256": expected_sha256,
142
- "reused": False,
143
  "immutable": True,
144
  "owner_scope": owner_scope,
 
145
  }
146
 
147
 
 
1
  """Structured theorem-input export for asynchronous Lean checking."""
2
 
3
+ import errno
4
  import hashlib
5
  import json
6
  import os
7
+ import stat
8
  import tempfile
9
  import threading
10
+ import time
11
  from pathlib import Path
12
  from typing import Any, Dict
13
 
14
 
15
  _ARTIFACT_QUOTA_LOCK = threading.RLock()
16
+ _TRANSIENT_LINK_ERRNOS = {
17
+ errno.ENOTSUP,
18
+ getattr(errno, "EOPNOTSUPP", errno.ENOTSUP),
19
+ }
20
+ _UNSUPPORTED_DIRECTORY_FSYNC_ERRNOS = {
21
+ errno.ENOSYS,
22
+ errno.EINVAL,
23
+ errno.ENOTSUP,
24
+ getattr(errno, "EOPNOTSUPP", errno.ENOTSUP),
25
+ }
26
+ _LINK_MAX_ATTEMPTS = 61
27
+ _LINK_RETRY_SECONDS = 0.5
28
 
29
 
30
  def canonical_json(value: Any) -> str:
 
35
  return hashlib.sha256(canonical_json(value).encode("utf-8")).hexdigest()
36
 
37
 
38
+ def _read_existing_regular(path: Path) -> bytes:
39
+ flags = os.O_RDONLY | getattr(os, "O_NOFOLLOW", 0)
40
+ try:
41
+ descriptor = os.open(path, flags)
42
+ except OSError as exc:
43
+ if exc.errno == getattr(errno, "ELOOP", None):
44
+ raise FileExistsError(
45
+ "existing GDW artifact is not a regular file"
46
+ ) from None
47
+ raise
48
+ try:
49
+ if not stat.S_ISREG(os.fstat(descriptor).st_mode):
50
+ raise FileExistsError(
51
+ "existing GDW artifact is not a regular file"
52
+ )
53
+ with os.fdopen(descriptor, "rb") as stream:
54
+ descriptor = -1
55
+ return stream.read()
56
+ finally:
57
+ if descriptor >= 0:
58
+ os.close(descriptor)
59
+
60
+
61
+ def _fsync_directory(path: Path) -> None:
62
+ if os.name != "posix":
63
+ return
64
+ descriptor = os.open(
65
+ path,
66
+ os.O_RDONLY | getattr(os, "O_DIRECTORY", 0),
67
+ )
68
+ try:
69
+ try:
70
+ os.fsync(descriptor)
71
+ except OSError as exc:
72
+ if exc.errno not in _UNSUPPORTED_DIRECTORY_FSYNC_ERRNOS:
73
+ raise
74
+ finally:
75
+ os.close(descriptor)
76
+
77
+
78
  def build_proof_payload(
79
  proposal_id: str,
80
  request_id: str,
 
144
  destination = owner_root / filename
145
  encoded = (json.dumps(payload, indent=2, sort_keys=True) + "\n").encode("utf-8")
146
  expected_sha256 = hashlib.sha256(encoded).hexdigest()
147
+ if os.path.lexists(destination):
148
+ if _read_existing_regular(destination) != encoded:
149
  raise FileExistsError(
150
  "refusing to overwrite an existing non-identical GDW artifact"
151
  )
152
+ _fsync_directory(owner_root)
153
  return {
154
  "status": "EXPORTED",
155
  "path": str(destination),
 
157
  "reused": True,
158
  "immutable": True,
159
  "owner_scope": owner_scope,
160
+ "publication_mode": "REUSED",
161
  }
162
  owner_limit = _bounded_artifact_limit(
163
  "GDW_OWNER_MAX_ARTIFACTS", default=10_000, maximum=100_000
 
173
  raise RuntimeError("per-owner artifact quota exceeded")
174
  if sum(1 for _ in root.glob("*/*.json")) >= global_limit:
175
  raise RuntimeError("global artifact quota exceeded")
176
+ # A unique stage is owned only by this publisher. Process death can leave
177
+ # that hidden file unreferenced, but no retry or replica may delete an
178
+ # unknown concurrent stage or expose it as the final JSON artifact.
179
+ handle, staging = tempfile.mkstemp(
180
+ prefix=".gdw-artifact-",
181
+ suffix=".tmp",
182
+ dir=owner_root,
183
  )
184
+ reused = False
185
+ publication_mode = "HARD_LINK"
186
  try:
187
  with os.fdopen(handle, "wb") as stream:
188
  stream.write(encoded)
189
  stream.flush()
190
  os.fsync(stream.fileno())
191
+ for attempt in range(_LINK_MAX_ATTEMPTS):
192
+ try:
193
+ os.link(staging, destination)
194
+ except FileExistsError:
195
+ if _read_existing_regular(destination) != encoded:
196
+ raise FileExistsError(
197
+ "refusing to overwrite a concurrently created "
198
+ "non-identical GDW artifact"
199
+ )
200
+ reused = True
201
+ publication_mode = "REUSED"
202
+ break
203
+ except OSError as exc:
204
+ if (
205
+ exc.errno not in _TRANSIENT_LINK_ERRNOS
206
+ or attempt + 1 >= _LINK_MAX_ATTEMPTS
207
+ ):
208
+ raise
209
+ publication_mode = "HARD_LINK_AFTER_FLUSH"
210
+ time.sleep(_LINK_RETRY_SECONDS)
211
+ else:
212
+ break
213
+ _fsync_directory(owner_root)
214
  finally:
215
+ try:
216
+ os.unlink(staging)
217
+ except FileNotFoundError:
218
+ pass
219
  return {
220
  "status": "EXPORTED",
221
  "path": str(destination),
222
  "sha256": expected_sha256,
223
+ "reused": reused,
224
  "immutable": True,
225
  "owner_scope": owner_scope,
226
+ "publication_mode": publication_mode,
227
  }
228
 
229
 
gdw_runtime.py CHANGED
@@ -442,8 +442,14 @@ def drain_once(
442
  "exported": exported,
443
  "failed": failed,
444
  "pending_effects": integrity["pending_effects"],
 
 
445
  "legacy_pending_proofs": integrity["pending_proofs"],
446
  "sqlite_integrity": integrity["sqlite_integrity"],
 
 
 
 
447
  "errors": errors,
448
  }
449
 
@@ -555,7 +561,28 @@ class OutboxSupervisor:
555
  lease_seconds=self.lease_seconds,
556
  worker_id=self.worker_id,
557
  )
558
- if report["failed"] or report["legacy_pending_proofs"]:
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
559
  retry_delay = min(
560
  self.retry_max_seconds,
561
  max(self.interval_seconds, retry_delay * 2),
@@ -564,8 +591,7 @@ class OutboxSupervisor:
564
  _set_drain_state(
565
  last_outcome="RETRY_SCHEDULED",
566
  last_error=(
567
- "bounded drain pass reported failures or "
568
- "unmigrated legacy proofs"
569
  ),
570
  last_report=report,
571
  )
 
442
  "exported": exported,
443
  "failed": failed,
444
  "pending_effects": integrity["pending_effects"],
445
+ "claimed_effects": integrity["claimed_effects"],
446
+ "dead_letter_effects": integrity["dead_letter_effects"],
447
  "legacy_pending_proofs": integrity["pending_proofs"],
448
  "sqlite_integrity": integrity["sqlite_integrity"],
449
+ "invalid_effect_bindings": integrity["invalid_effect_bindings"],
450
+ "invalid_exported_artifacts": integrity[
451
+ "invalid_exported_artifacts"
452
+ ],
453
  "errors": errors,
454
  }
455
 
 
561
  lease_seconds=self.lease_seconds,
562
  worker_id=self.worker_id,
563
  )
564
+ terminal_failure = (
565
+ report["dead_letter_effects"]
566
+ or report["sqlite_integrity"] != "ok"
567
+ or report["invalid_effect_bindings"]
568
+ or report["invalid_exported_artifacts"]
569
+ )
570
+ retryable_work = (
571
+ report["failed"]
572
+ or report["pending_effects"]
573
+ or report["legacy_pending_proofs"]
574
+ )
575
+ if terminal_failure:
576
+ delay = self.retry_max_seconds
577
+ _set_drain_state(
578
+ last_outcome="FAILED",
579
+ last_error=(
580
+ "bounded drain pass reported terminal "
581
+ "integrity or dead-letter failures"
582
+ ),
583
+ last_report=report,
584
+ )
585
+ elif retryable_work:
586
  retry_delay = min(
587
  self.retry_max_seconds,
588
  max(self.interval_seconds, retry_delay * 2),
 
591
  _set_drain_state(
592
  last_outcome="RETRY_SCHEDULED",
593
  last_error=(
594
+ "bounded drain pass remains non-quiescent"
 
595
  ),
596
  last_report=report,
597
  )
routers/gdw_frontier.py CHANGED
@@ -286,6 +286,17 @@ def _write_readiness(
286
  blockers.append("OUTBOX_SUPERVISOR_NOT_RUNNING")
287
  if drain.get("last_outcome") != "SUCCEEDED":
288
  blockers.append("OUTBOX_SUPERVISOR_NOT_HEALTHY")
 
 
 
 
 
 
 
 
 
 
 
289
  if drain.get("success_run_generation_id") != drain.get(
290
  "run_generation_id"
291
  ):
@@ -376,6 +387,34 @@ def _public_runtime_health(runtime: dict) -> dict:
376
  )
377
  if key in drain
378
  }
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
379
  return {
380
  "startup_state": runtime.get("startup_state"),
381
  "evidence_label": runtime.get("evidence_label"),
 
286
  blockers.append("OUTBOX_SUPERVISOR_NOT_RUNNING")
287
  if drain.get("last_outcome") != "SUCCEEDED":
288
  blockers.append("OUTBOX_SUPERVISOR_NOT_HEALTHY")
289
+ last_report = drain.get("last_report")
290
+ if isinstance(last_report, dict) and (
291
+ last_report.get("failed")
292
+ or last_report.get("pending_effects")
293
+ or last_report.get("dead_letter_effects")
294
+ or last_report.get("legacy_pending_proofs")
295
+ or last_report.get("invalid_effect_bindings")
296
+ or last_report.get("invalid_exported_artifacts")
297
+ or last_report.get("sqlite_integrity") != "ok"
298
+ ):
299
+ blockers.append("OUTBOX_SUPERVISOR_NOT_QUIESCENT")
300
  if drain.get("success_run_generation_id") != drain.get(
301
  "run_generation_id"
302
  ):
 
387
  )
388
  if key in drain
389
  }
390
+ last_report = drain.get("last_report")
391
+ if isinstance(last_report, dict):
392
+ public_report = {
393
+ key: last_report.get(key)
394
+ for key in (
395
+ "attempted",
396
+ "exported",
397
+ "failed",
398
+ "pending_effects",
399
+ "claimed_effects",
400
+ "dead_letter_effects",
401
+ "legacy_pending_proofs",
402
+ "sqlite_integrity",
403
+ "invalid_effect_bindings",
404
+ "invalid_exported_artifacts",
405
+ )
406
+ if key in last_report
407
+ }
408
+ errors = last_report.get("errors")
409
+ if isinstance(errors, list):
410
+ public_report["errors"] = [
411
+ value
412
+ for value in errors
413
+ if isinstance(value, str)
414
+ and len(value) <= 96
415
+ and re.fullmatch(r"[a-z_]+:[A-Za-z_]+", value)
416
+ ]
417
+ public_drain["last_report"] = public_report
418
  return {
419
  "startup_state": runtime.get("startup_state"),
420
  "evidence_label": runtime.get("evidence_label"),