Spaces:
Running
Running
deploy(hf): sync szl-holdings/a11oy@8b76c9d882b6d6a0e106aaa597cab2c125bbd676 derived COPY set
Browse filesReusable Dockerfile-COPY-derived deploy from szl-holdings/a11oy 8b76c9d882b6d6a0e106aaa597cab2c125bbd676.
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>
- gdw_runtime.py +17 -1
- gdw_workspace.py +348 -83
- routers/gdw_frontier.py +31 -0
gdw_runtime.py
CHANGED
|
@@ -347,10 +347,17 @@ def drain_once(
|
|
| 347 |
exported = 0
|
| 348 |
failed = 0
|
| 349 |
errors = []
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 350 |
identities = (
|
| 351 |
[(store.namespace, store.owner_id)]
|
| 352 |
if workspace is not None
|
| 353 |
-
else store.
|
| 354 |
)
|
| 355 |
|
| 356 |
for namespace, owner_id in identities:
|
|
@@ -436,6 +443,14 @@ def drain_once(
|
|
| 436 |
failed += 1
|
| 437 |
errors.append(f"{row['kind']}:{type(exc).__name__}")
|
| 438 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 439 |
integrity = store.integrity(global_scope=True)
|
| 440 |
return {
|
| 441 |
"attempted": exported + failed,
|
|
@@ -450,6 +465,7 @@ def drain_once(
|
|
| 450 |
"invalid_exported_artifacts": integrity[
|
| 451 |
"invalid_exported_artifacts"
|
| 452 |
],
|
|
|
|
| 453 |
"errors": errors,
|
| 454 |
}
|
| 455 |
|
|
|
|
| 347 |
exported = 0
|
| 348 |
failed = 0
|
| 349 |
errors = []
|
| 350 |
+
garbage_collected = {
|
| 351 |
+
"sessions_tombstoned": 0,
|
| 352 |
+
"requests_tombstoned": 0,
|
| 353 |
+
"effects_compacted": 0,
|
| 354 |
+
"proofs_compacted": 0,
|
| 355 |
+
"tombstones_purged": 0,
|
| 356 |
+
}
|
| 357 |
identities = (
|
| 358 |
[(store.namespace, store.owner_id)]
|
| 359 |
if workspace is not None
|
| 360 |
+
else store.lifecycle_identities()
|
| 361 |
)
|
| 362 |
|
| 363 |
for namespace, owner_id in identities:
|
|
|
|
| 443 |
failed += 1
|
| 444 |
errors.append(f"{row['kind']}:{type(exc).__name__}")
|
| 445 |
|
| 446 |
+
collected = store.collect_garbage(
|
| 447 |
+
limit=bounded,
|
| 448 |
+
namespace=namespace,
|
| 449 |
+
owner_id=owner_id,
|
| 450 |
+
)
|
| 451 |
+
for key in garbage_collected:
|
| 452 |
+
garbage_collected[key] += collected[key]
|
| 453 |
+
|
| 454 |
integrity = store.integrity(global_scope=True)
|
| 455 |
return {
|
| 456 |
"attempted": exported + failed,
|
|
|
|
| 465 |
"invalid_exported_artifacts": integrity[
|
| 466 |
"invalid_exported_artifacts"
|
| 467 |
],
|
| 468 |
+
"garbage_collected": garbage_collected,
|
| 469 |
"errors": errors,
|
| 470 |
}
|
| 471 |
|
gdw_workspace.py
CHANGED
|
@@ -665,9 +665,11 @@ class GDWWorkspace:
|
|
| 665 |
"v2 proof effect differs from persisted request"
|
| 666 |
)
|
| 667 |
rebound_response = dict(response)
|
|
|
|
| 668 |
rebound_response["request_digest"] = str(
|
| 669 |
request_row["request_digest"]
|
| 670 |
)
|
|
|
|
| 671 |
rebound_response["database_generation_id"] = generation
|
| 672 |
rebound_response["principal"] = principal
|
| 673 |
rebound_response["proposal_id"] = hashlib.sha256(
|
|
@@ -1026,7 +1028,8 @@ class GDWWorkspace:
|
|
| 1026 |
ns, owner = self._identity(namespace, owner_id)
|
| 1027 |
row = connection.execute(
|
| 1028 |
"""
|
| 1029 |
-
SELECT
|
|
|
|
| 1030 |
FROM requests
|
| 1031 |
WHERE namespace = ? AND owner_id = ? AND request_id = ?
|
| 1032 |
""",
|
|
@@ -1036,7 +1039,54 @@ class GDWWorkspace:
|
|
| 1036 |
return None
|
| 1037 |
if row["lifecycle"] != "ACTIVE" or row["response_json"] is None:
|
| 1038 |
raise GDWLifecycleError("idempotency record is outside its replay window")
|
| 1039 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 1040 |
|
| 1041 |
def session_state(
|
| 1042 |
self,
|
|
@@ -1049,8 +1099,8 @@ class GDWWorkspace:
|
|
| 1049 |
ns, owner = self._identity(namespace, owner_id)
|
| 1050 |
row = connection.execute(
|
| 1051 |
"""
|
| 1052 |
-
SELECT
|
| 1053 |
-
expires_at
|
| 1054 |
FROM session_state
|
| 1055 |
WHERE namespace = ? AND owner_id = ? AND session_id = ?
|
| 1056 |
""",
|
|
@@ -1058,20 +1108,36 @@ class GDWWorkspace:
|
|
| 1058 |
).fetchone()
|
| 1059 |
if row is None or row["lifecycle"] != "ACTIVE" or row["state_json"] is None:
|
| 1060 |
return None
|
| 1061 |
-
|
| 1062 |
-
""
|
| 1063 |
-
|
| 1064 |
-
|
| 1065 |
-
|
| 1066 |
-
(
|
| 1067 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 1068 |
return {
|
| 1069 |
"namespace": ns,
|
| 1070 |
"owner_id": owner,
|
| 1071 |
"session_id": session_id,
|
| 1072 |
"database_generation_id": self.database_generation_id,
|
| 1073 |
"step": int(row["step"]),
|
| 1074 |
-
"state":
|
| 1075 |
"state_hash": row["state_hash"],
|
| 1076 |
"updated_at": row["updated_at"],
|
| 1077 |
"expires_at": row["expires_at"],
|
|
@@ -2265,9 +2331,127 @@ class GDWWorkspace:
|
|
| 2265 |
finally:
|
| 2266 |
connection.close()
|
| 2267 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 2268 |
@staticmethod
|
| 2269 |
-
def
|
|
|
|
|
|
|
|
|
|
|
|
|
| 2270 |
timestamp = _text_time()
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 2271 |
identities = connection.execute(
|
| 2272 |
"""
|
| 2273 |
SELECT namespace, owner_id FROM session_state
|
|
@@ -2279,73 +2463,10 @@ class GDWWorkspace:
|
|
| 2279 |
).fetchall()
|
| 2280 |
connection.execute("DELETE FROM usage")
|
| 2281 |
for identity in identities:
|
| 2282 |
-
|
| 2283 |
-
|
| 2284 |
-
|
| 2285 |
-
|
| 2286 |
-
SELECT COUNT(*) FROM session_state
|
| 2287 |
-
WHERE namespace = ? AND owner_id = ? AND lifecycle = 'ACTIVE'
|
| 2288 |
-
""",
|
| 2289 |
-
(ns, owner),
|
| 2290 |
-
).fetchone()[0]
|
| 2291 |
-
)
|
| 2292 |
-
requests = int(
|
| 2293 |
-
connection.execute(
|
| 2294 |
-
"""
|
| 2295 |
-
SELECT COUNT(*) FROM requests
|
| 2296 |
-
WHERE namespace = ? AND owner_id = ? AND lifecycle = 'ACTIVE'
|
| 2297 |
-
""",
|
| 2298 |
-
(ns, owner),
|
| 2299 |
-
).fetchone()[0]
|
| 2300 |
-
)
|
| 2301 |
-
pending = int(
|
| 2302 |
-
connection.execute(
|
| 2303 |
-
"""
|
| 2304 |
-
SELECT
|
| 2305 |
-
(SELECT COUNT(*) FROM effect_outbox
|
| 2306 |
-
WHERE namespace = ? AND owner_id = ?
|
| 2307 |
-
AND status IN ('PENDING', 'CLAIMED')) +
|
| 2308 |
-
(SELECT COUNT(*) FROM proof_outbox
|
| 2309 |
-
WHERE namespace = ? AND owner_id = ? AND status = 'PENDING')
|
| 2310 |
-
""",
|
| 2311 |
-
(ns, owner, ns, owner),
|
| 2312 |
-
).fetchone()[0]
|
| 2313 |
-
)
|
| 2314 |
-
stored = int(
|
| 2315 |
-
connection.execute(
|
| 2316 |
-
"""
|
| 2317 |
-
SELECT
|
| 2318 |
-
COALESCE((SELECT SUM(LENGTH(CAST(state_json AS BLOB)))
|
| 2319 |
-
FROM session_state
|
| 2320 |
-
WHERE namespace = ? AND owner_id = ?), 0) +
|
| 2321 |
-
COALESCE((SELECT SUM(LENGTH(CAST(response_json AS BLOB)))
|
| 2322 |
-
FROM requests
|
| 2323 |
-
WHERE namespace = ? AND owner_id = ?), 0) +
|
| 2324 |
-
COALESCE((SELECT SUM(LENGTH(CAST(receipt_json AS BLOB)))
|
| 2325 |
-
FROM receipts
|
| 2326 |
-
WHERE namespace = ? AND owner_id = ?), 0) +
|
| 2327 |
-
COALESCE((SELECT SUM(
|
| 2328 |
-
LENGTH(CAST(payload_json AS BLOB)) +
|
| 2329 |
-
COALESCE(LENGTH(CAST(artifact_json AS BLOB)), 0))
|
| 2330 |
-
FROM proof_outbox
|
| 2331 |
-
WHERE namespace = ? AND owner_id = ?), 0) +
|
| 2332 |
-
COALESCE((SELECT SUM(
|
| 2333 |
-
LENGTH(CAST(payload_json AS BLOB)) +
|
| 2334 |
-
COALESCE(LENGTH(CAST(artifact_json AS BLOB)), 0))
|
| 2335 |
-
FROM effect_outbox
|
| 2336 |
-
WHERE namespace = ? AND owner_id = ?), 0)
|
| 2337 |
-
""",
|
| 2338 |
-
(ns, owner, ns, owner, ns, owner, ns, owner, ns, owner),
|
| 2339 |
-
).fetchone()[0]
|
| 2340 |
-
)
|
| 2341 |
-
connection.execute(
|
| 2342 |
-
"""
|
| 2343 |
-
INSERT INTO usage(
|
| 2344 |
-
namespace, owner_id, active_sessions, active_requests,
|
| 2345 |
-
pending_effects, stored_bytes, updated_at
|
| 2346 |
-
) VALUES (?, ?, ?, ?, ?, ?, ?)
|
| 2347 |
-
""",
|
| 2348 |
-
(ns, owner, sessions, requests, pending, stored, timestamp),
|
| 2349 |
)
|
| 2350 |
|
| 2351 |
def reconcile_usage(self) -> Dict[str, int]:
|
|
@@ -2547,7 +2668,7 @@ class GDWWorkspace:
|
|
| 2547 |
(ns, owner, purge_before),
|
| 2548 |
).rowcount
|
| 2549 |
result["tombstones_purged"] = purged
|
| 2550 |
-
self.
|
| 2551 |
return result
|
| 2552 |
|
| 2553 |
def integrity(
|
|
@@ -2556,9 +2677,12 @@ class GDWWorkspace:
|
|
| 2556 |
namespace: Optional[str] = None,
|
| 2557 |
owner_id: Optional[str] = None,
|
| 2558 |
global_scope: bool = False,
|
|
|
|
| 2559 |
) -> Dict[str, Any]:
|
| 2560 |
ns, owner = self._identity(namespace, owner_id)
|
| 2561 |
-
|
|
|
|
|
|
|
| 2562 |
try:
|
| 2563 |
check = connection.execute("PRAGMA integrity_check").fetchone()[0]
|
| 2564 |
predicate = "" if global_scope else " WHERE namespace = ? AND owner_id = ?"
|
|
@@ -2634,6 +2758,144 @@ class GDWWorkspace:
|
|
| 2634 |
params,
|
| 2635 |
).fetchone()[0]
|
| 2636 |
)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 2637 |
effect_rows = connection.execute(
|
| 2638 |
"""
|
| 2639 |
SELECT namespace, owner_id, idempotency_key,
|
|
@@ -2678,6 +2940,7 @@ class GDWWorkspace:
|
|
| 2678 |
"ok": (
|
| 2679 |
check == "ok"
|
| 2680 |
and orphan_receipts == 0
|
|
|
|
| 2681 |
and invalid_effect_bindings == 0
|
| 2682 |
and invalid_exported_artifacts == 0
|
| 2683 |
),
|
|
@@ -2691,6 +2954,7 @@ class GDWWorkspace:
|
|
| 2691 |
"dead_letter_effects": dead_letter_effects,
|
| 2692 |
"invalid_effect_bindings": invalid_effect_bindings,
|
| 2693 |
"invalid_exported_artifacts": invalid_exported_artifacts,
|
|
|
|
| 2694 |
"counts": counts,
|
| 2695 |
"scope": "global" if global_scope else "owner",
|
| 2696 |
"journal_mode": str(
|
|
@@ -2705,4 +2969,5 @@ class GDWWorkspace:
|
|
| 2705 |
result["owner_id"] = owner
|
| 2706 |
return result
|
| 2707 |
finally:
|
| 2708 |
-
|
|
|
|
|
|
| 665 |
"v2 proof effect differs from persisted request"
|
| 666 |
)
|
| 667 |
rebound_response = dict(response)
|
| 668 |
+
rebound_response["request_id"] = request_id
|
| 669 |
rebound_response["request_digest"] = str(
|
| 670 |
request_row["request_digest"]
|
| 671 |
)
|
| 672 |
+
rebound_response["session_id"] = str(request_row["session_id"])
|
| 673 |
rebound_response["database_generation_id"] = generation
|
| 674 |
rebound_response["principal"] = principal
|
| 675 |
rebound_response["proposal_id"] = hashlib.sha256(
|
|
|
|
| 1028 |
ns, owner = self._identity(namespace, owner_id)
|
| 1029 |
row = connection.execute(
|
| 1030 |
"""
|
| 1031 |
+
SELECT namespace, owner_id, request_id, request_digest, session_id,
|
| 1032 |
+
response_json, response_hash, lifecycle
|
| 1033 |
FROM requests
|
| 1034 |
WHERE namespace = ? AND owner_id = ? AND request_id = ?
|
| 1035 |
""",
|
|
|
|
| 1039 |
return None
|
| 1040 |
if row["lifecycle"] != "ACTIVE" or row["response_json"] is None:
|
| 1041 |
raise GDWLifecycleError("idempotency record is outside its replay window")
|
| 1042 |
+
try:
|
| 1043 |
+
response = json.loads(row["response_json"])
|
| 1044 |
+
except (TypeError, json.JSONDecodeError) as exc:
|
| 1045 |
+
raise GDWConfigurationError(
|
| 1046 |
+
"idempotency response is not canonical JSON"
|
| 1047 |
+
) from exc
|
| 1048 |
+
if not isinstance(response, dict):
|
| 1049 |
+
raise GDWConfigurationError("idempotency response is not an object")
|
| 1050 |
+
observed_hash = hashlib.sha256(
|
| 1051 |
+
_json_text(response).encode("utf-8")
|
| 1052 |
+
).hexdigest()
|
| 1053 |
+
expected_identity = {
|
| 1054 |
+
"request_id": row["request_id"],
|
| 1055 |
+
"request_digest": row["request_digest"],
|
| 1056 |
+
"session_id": row["session_id"],
|
| 1057 |
+
"database_generation_id": self.database_generation_id,
|
| 1058 |
+
}
|
| 1059 |
+
if observed_hash != row["response_hash"] or any(
|
| 1060 |
+
response.get(field) != expected
|
| 1061 |
+
for field, expected in expected_identity.items()
|
| 1062 |
+
):
|
| 1063 |
+
raise GDWConfigurationError(
|
| 1064 |
+
"idempotency response digest or identity is invalid"
|
| 1065 |
+
)
|
| 1066 |
+
principal = response.get("principal")
|
| 1067 |
+
if not isinstance(principal, dict) or (
|
| 1068 |
+
principal.get("namespace") != row["namespace"]
|
| 1069 |
+
or principal.get("owner_id") != row["owner_id"]
|
| 1070 |
+
):
|
| 1071 |
+
raise GDWConfigurationError(
|
| 1072 |
+
"idempotency response principal identity is invalid"
|
| 1073 |
+
)
|
| 1074 |
+
receipt_hash = response.get("receipt_hash")
|
| 1075 |
+
if receipt_hash:
|
| 1076 |
+
receipt_anchor = self._receipt_anchor(
|
| 1077 |
+
connection,
|
| 1078 |
+
row["namespace"],
|
| 1079 |
+
row["owner_id"],
|
| 1080 |
+
row["request_id"],
|
| 1081 |
+
)
|
| 1082 |
+
if (
|
| 1083 |
+
receipt_anchor["receipt_hash"] != receipt_hash
|
| 1084 |
+
or receipt_anchor["receipt"].get("session_id") != row["session_id"]
|
| 1085 |
+
):
|
| 1086 |
+
raise GDWConfigurationError(
|
| 1087 |
+
"idempotency response receipt binding is invalid"
|
| 1088 |
+
)
|
| 1089 |
+
return row["request_digest"], response
|
| 1090 |
|
| 1091 |
def session_state(
|
| 1092 |
self,
|
|
|
|
| 1099 |
ns, owner = self._identity(namespace, owner_id)
|
| 1100 |
row = connection.execute(
|
| 1101 |
"""
|
| 1102 |
+
SELECT namespace, owner_id, session_id, step, state_json, state_hash,
|
| 1103 |
+
updated_at, lifecycle, expires_at
|
| 1104 |
FROM session_state
|
| 1105 |
WHERE namespace = ? AND owner_id = ? AND session_id = ?
|
| 1106 |
""",
|
|
|
|
| 1108 |
).fetchone()
|
| 1109 |
if row is None or row["lifecycle"] != "ACTIVE" or row["state_json"] is None:
|
| 1110 |
return None
|
| 1111 |
+
try:
|
| 1112 |
+
state = json.loads(row["state_json"])
|
| 1113 |
+
except (TypeError, json.JSONDecodeError) as exc:
|
| 1114 |
+
raise GDWConfigurationError("session state is not canonical JSON") from exc
|
| 1115 |
+
if not isinstance(state, dict):
|
| 1116 |
+
raise GDWConfigurationError("session state is not an object")
|
| 1117 |
+
observed_hash = hashlib.sha256(
|
| 1118 |
+
_json_text(state).encode("utf-8")
|
| 1119 |
+
).hexdigest()
|
| 1120 |
+
expected_identity = {
|
| 1121 |
+
"namespace": row["namespace"],
|
| 1122 |
+
"owner_id": row["owner_id"],
|
| 1123 |
+
"session_id": row["session_id"],
|
| 1124 |
+
"step": int(row["step"]),
|
| 1125 |
+
"database_generation_id": self.database_generation_id,
|
| 1126 |
+
}
|
| 1127 |
+
if observed_hash != row["state_hash"] or any(
|
| 1128 |
+
state.get(field) != expected
|
| 1129 |
+
for field, expected in expected_identity.items()
|
| 1130 |
+
):
|
| 1131 |
+
raise GDWConfigurationError(
|
| 1132 |
+
"session state digest or identity is invalid"
|
| 1133 |
+
)
|
| 1134 |
return {
|
| 1135 |
"namespace": ns,
|
| 1136 |
"owner_id": owner,
|
| 1137 |
"session_id": session_id,
|
| 1138 |
"database_generation_id": self.database_generation_id,
|
| 1139 |
"step": int(row["step"]),
|
| 1140 |
+
"state": state,
|
| 1141 |
"state_hash": row["state_hash"],
|
| 1142 |
"updated_at": row["updated_at"],
|
| 1143 |
"expires_at": row["expires_at"],
|
|
|
|
| 2331 |
finally:
|
| 2332 |
connection.close()
|
| 2333 |
|
| 2334 |
+
def lifecycle_identities(self) -> list[Tuple[str, str]]:
|
| 2335 |
+
"""Return principals whose retained state may need supervised cleanup."""
|
| 2336 |
+
|
| 2337 |
+
connection = self._connect()
|
| 2338 |
+
try:
|
| 2339 |
+
rows = connection.execute(
|
| 2340 |
+
"""
|
| 2341 |
+
SELECT namespace, owner_id FROM session_state
|
| 2342 |
+
UNION SELECT namespace, owner_id FROM requests
|
| 2343 |
+
UNION SELECT namespace, owner_id FROM receipts
|
| 2344 |
+
UNION SELECT namespace, owner_id FROM proof_outbox
|
| 2345 |
+
UNION SELECT namespace, owner_id FROM effect_outbox
|
| 2346 |
+
ORDER BY namespace, owner_id
|
| 2347 |
+
"""
|
| 2348 |
+
).fetchall()
|
| 2349 |
+
return [(row["namespace"], row["owner_id"]) for row in rows]
|
| 2350 |
+
finally:
|
| 2351 |
+
connection.close()
|
| 2352 |
+
|
| 2353 |
@staticmethod
|
| 2354 |
+
def _reconcile_identity_usage(
|
| 2355 |
+
connection: sqlite3.Connection,
|
| 2356 |
+
namespace: str,
|
| 2357 |
+
owner_id: str,
|
| 2358 |
+
) -> None:
|
| 2359 |
timestamp = _text_time()
|
| 2360 |
+
sessions = int(
|
| 2361 |
+
connection.execute(
|
| 2362 |
+
"""
|
| 2363 |
+
SELECT COUNT(*) FROM session_state
|
| 2364 |
+
WHERE namespace = ? AND owner_id = ? AND lifecycle = 'ACTIVE'
|
| 2365 |
+
""",
|
| 2366 |
+
(namespace, owner_id),
|
| 2367 |
+
).fetchone()[0]
|
| 2368 |
+
)
|
| 2369 |
+
requests = int(
|
| 2370 |
+
connection.execute(
|
| 2371 |
+
"""
|
| 2372 |
+
SELECT COUNT(*) FROM requests
|
| 2373 |
+
WHERE namespace = ? AND owner_id = ? AND lifecycle = 'ACTIVE'
|
| 2374 |
+
""",
|
| 2375 |
+
(namespace, owner_id),
|
| 2376 |
+
).fetchone()[0]
|
| 2377 |
+
)
|
| 2378 |
+
pending = int(
|
| 2379 |
+
connection.execute(
|
| 2380 |
+
"""
|
| 2381 |
+
SELECT
|
| 2382 |
+
(SELECT COUNT(*) FROM effect_outbox
|
| 2383 |
+
WHERE namespace = ? AND owner_id = ?
|
| 2384 |
+
AND status IN ('PENDING', 'CLAIMED')) +
|
| 2385 |
+
(SELECT COUNT(*) FROM proof_outbox
|
| 2386 |
+
WHERE namespace = ? AND owner_id = ? AND status = 'PENDING')
|
| 2387 |
+
""",
|
| 2388 |
+
(namespace, owner_id, namespace, owner_id),
|
| 2389 |
+
).fetchone()[0]
|
| 2390 |
+
)
|
| 2391 |
+
stored = int(
|
| 2392 |
+
connection.execute(
|
| 2393 |
+
"""
|
| 2394 |
+
SELECT
|
| 2395 |
+
COALESCE((SELECT SUM(LENGTH(CAST(state_json AS BLOB)))
|
| 2396 |
+
FROM session_state
|
| 2397 |
+
WHERE namespace = ? AND owner_id = ?), 0) +
|
| 2398 |
+
COALESCE((SELECT SUM(LENGTH(CAST(response_json AS BLOB)))
|
| 2399 |
+
FROM requests
|
| 2400 |
+
WHERE namespace = ? AND owner_id = ?), 0) +
|
| 2401 |
+
COALESCE((SELECT SUM(LENGTH(CAST(receipt_json AS BLOB)))
|
| 2402 |
+
FROM receipts
|
| 2403 |
+
WHERE namespace = ? AND owner_id = ?), 0) +
|
| 2404 |
+
COALESCE((SELECT SUM(
|
| 2405 |
+
LENGTH(CAST(payload_json AS BLOB)) +
|
| 2406 |
+
COALESCE(LENGTH(CAST(artifact_json AS BLOB)), 0))
|
| 2407 |
+
FROM proof_outbox
|
| 2408 |
+
WHERE namespace = ? AND owner_id = ?), 0) +
|
| 2409 |
+
COALESCE((SELECT SUM(
|
| 2410 |
+
LENGTH(CAST(payload_json AS BLOB)) +
|
| 2411 |
+
COALESCE(LENGTH(CAST(artifact_json AS BLOB)), 0))
|
| 2412 |
+
FROM effect_outbox
|
| 2413 |
+
WHERE namespace = ? AND owner_id = ?), 0)
|
| 2414 |
+
""",
|
| 2415 |
+
(
|
| 2416 |
+
namespace,
|
| 2417 |
+
owner_id,
|
| 2418 |
+
namespace,
|
| 2419 |
+
owner_id,
|
| 2420 |
+
namespace,
|
| 2421 |
+
owner_id,
|
| 2422 |
+
namespace,
|
| 2423 |
+
owner_id,
|
| 2424 |
+
namespace,
|
| 2425 |
+
owner_id,
|
| 2426 |
+
),
|
| 2427 |
+
).fetchone()[0]
|
| 2428 |
+
)
|
| 2429 |
+
connection.execute(
|
| 2430 |
+
"""
|
| 2431 |
+
INSERT INTO usage(
|
| 2432 |
+
namespace, owner_id, active_sessions, active_requests,
|
| 2433 |
+
pending_effects, stored_bytes, updated_at
|
| 2434 |
+
) VALUES (?, ?, ?, ?, ?, ?, ?)
|
| 2435 |
+
ON CONFLICT(namespace, owner_id) DO UPDATE SET
|
| 2436 |
+
active_sessions = excluded.active_sessions,
|
| 2437 |
+
active_requests = excluded.active_requests,
|
| 2438 |
+
pending_effects = excluded.pending_effects,
|
| 2439 |
+
stored_bytes = excluded.stored_bytes,
|
| 2440 |
+
updated_at = excluded.updated_at
|
| 2441 |
+
""",
|
| 2442 |
+
(
|
| 2443 |
+
namespace,
|
| 2444 |
+
owner_id,
|
| 2445 |
+
sessions,
|
| 2446 |
+
requests,
|
| 2447 |
+
pending,
|
| 2448 |
+
stored,
|
| 2449 |
+
timestamp,
|
| 2450 |
+
),
|
| 2451 |
+
)
|
| 2452 |
+
|
| 2453 |
+
@staticmethod
|
| 2454 |
+
def _reconcile_usage(connection: sqlite3.Connection) -> None:
|
| 2455 |
identities = connection.execute(
|
| 2456 |
"""
|
| 2457 |
SELECT namespace, owner_id FROM session_state
|
|
|
|
| 2463 |
).fetchall()
|
| 2464 |
connection.execute("DELETE FROM usage")
|
| 2465 |
for identity in identities:
|
| 2466 |
+
GDWWorkspace._reconcile_identity_usage(
|
| 2467 |
+
connection,
|
| 2468 |
+
identity["namespace"],
|
| 2469 |
+
identity["owner_id"],
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 2470 |
)
|
| 2471 |
|
| 2472 |
def reconcile_usage(self) -> Dict[str, int]:
|
|
|
|
| 2668 |
(ns, owner, purge_before),
|
| 2669 |
).rowcount
|
| 2670 |
result["tombstones_purged"] = purged
|
| 2671 |
+
self._reconcile_identity_usage(connection, ns, owner)
|
| 2672 |
return result
|
| 2673 |
|
| 2674 |
def integrity(
|
|
|
|
| 2677 |
namespace: Optional[str] = None,
|
| 2678 |
owner_id: Optional[str] = None,
|
| 2679 |
global_scope: bool = False,
|
| 2680 |
+
connection: Optional[sqlite3.Connection] = None,
|
| 2681 |
) -> Dict[str, Any]:
|
| 2682 |
ns, owner = self._identity(namespace, owner_id)
|
| 2683 |
+
owns_connection = connection is None
|
| 2684 |
+
if connection is None:
|
| 2685 |
+
connection = self._connect()
|
| 2686 |
try:
|
| 2687 |
check = connection.execute("PRAGMA integrity_check").fetchone()[0]
|
| 2688 |
predicate = "" if global_scope else " WHERE namespace = ? AND owner_id = ?"
|
|
|
|
| 2758 |
params,
|
| 2759 |
).fetchone()[0]
|
| 2760 |
)
|
| 2761 |
+
digest_violations = {
|
| 2762 |
+
"invalid_state_digests": 0,
|
| 2763 |
+
"invalid_request_digests": 0,
|
| 2764 |
+
"invalid_receipt_digests": 0,
|
| 2765 |
+
"invalid_proof_digests": 0,
|
| 2766 |
+
}
|
| 2767 |
+
scoped_suffix = (
|
| 2768 |
+
""
|
| 2769 |
+
if global_scope
|
| 2770 |
+
else " WHERE namespace = ? AND owner_id = ?"
|
| 2771 |
+
)
|
| 2772 |
+
for row in connection.execute(
|
| 2773 |
+
"SELECT namespace, owner_id, session_id, step, state_json, "
|
| 2774 |
+
"state_hash FROM session_state"
|
| 2775 |
+
+ scoped_suffix,
|
| 2776 |
+
params,
|
| 2777 |
+
):
|
| 2778 |
+
if row["state_json"] is None:
|
| 2779 |
+
continue
|
| 2780 |
+
try:
|
| 2781 |
+
state = json.loads(row["state_json"])
|
| 2782 |
+
observed = hashlib.sha256(
|
| 2783 |
+
_json_text(state).encode("utf-8")
|
| 2784 |
+
).hexdigest()
|
| 2785 |
+
expected_identity = {
|
| 2786 |
+
"namespace": row["namespace"],
|
| 2787 |
+
"owner_id": row["owner_id"],
|
| 2788 |
+
"session_id": row["session_id"],
|
| 2789 |
+
"step": int(row["step"]),
|
| 2790 |
+
"database_generation_id": self.database_generation_id,
|
| 2791 |
+
}
|
| 2792 |
+
if (
|
| 2793 |
+
not isinstance(state, dict)
|
| 2794 |
+
or observed != row["state_hash"]
|
| 2795 |
+
or any(
|
| 2796 |
+
state.get(field) != expected
|
| 2797 |
+
for field, expected in expected_identity.items()
|
| 2798 |
+
)
|
| 2799 |
+
):
|
| 2800 |
+
raise ValueError("state digest mismatch")
|
| 2801 |
+
except (TypeError, ValueError, json.JSONDecodeError):
|
| 2802 |
+
digest_violations["invalid_state_digests"] += 1
|
| 2803 |
+
for row in connection.execute(
|
| 2804 |
+
"SELECT namespace, owner_id, request_id, request_digest, "
|
| 2805 |
+
"session_id, response_json, response_hash FROM requests"
|
| 2806 |
+
+ scoped_suffix,
|
| 2807 |
+
params,
|
| 2808 |
+
):
|
| 2809 |
+
if row["response_json"] is None:
|
| 2810 |
+
continue
|
| 2811 |
+
try:
|
| 2812 |
+
response = json.loads(row["response_json"])
|
| 2813 |
+
observed = hashlib.sha256(
|
| 2814 |
+
_json_text(response).encode("utf-8")
|
| 2815 |
+
).hexdigest()
|
| 2816 |
+
principal = (
|
| 2817 |
+
response.get("principal")
|
| 2818 |
+
if isinstance(response, dict)
|
| 2819 |
+
else None
|
| 2820 |
+
)
|
| 2821 |
+
expected_identity = {
|
| 2822 |
+
"request_id": row["request_id"],
|
| 2823 |
+
"request_digest": row["request_digest"],
|
| 2824 |
+
"session_id": row["session_id"],
|
| 2825 |
+
"database_generation_id": self.database_generation_id,
|
| 2826 |
+
}
|
| 2827 |
+
if (
|
| 2828 |
+
not isinstance(response, dict)
|
| 2829 |
+
or observed != row["response_hash"]
|
| 2830 |
+
or any(
|
| 2831 |
+
response.get(field) != expected
|
| 2832 |
+
for field, expected in expected_identity.items()
|
| 2833 |
+
)
|
| 2834 |
+
or not isinstance(principal, dict)
|
| 2835 |
+
or principal.get("namespace") != row["namespace"]
|
| 2836 |
+
or principal.get("owner_id") != row["owner_id"]
|
| 2837 |
+
):
|
| 2838 |
+
raise ValueError("request digest mismatch")
|
| 2839 |
+
except (TypeError, ValueError, json.JSONDecodeError):
|
| 2840 |
+
digest_violations["invalid_request_digests"] += 1
|
| 2841 |
+
for row in connection.execute(
|
| 2842 |
+
"SELECT namespace, owner_id, request_id, session_id, step, "
|
| 2843 |
+
"receipt_json, receipt_hash FROM receipts"
|
| 2844 |
+
+ scoped_suffix,
|
| 2845 |
+
params,
|
| 2846 |
+
):
|
| 2847 |
+
if row["receipt_json"] is None:
|
| 2848 |
+
continue
|
| 2849 |
+
try:
|
| 2850 |
+
receipt = json.loads(row["receipt_json"])
|
| 2851 |
+
if not isinstance(receipt, dict):
|
| 2852 |
+
raise ValueError("receipt must be an object")
|
| 2853 |
+
claimed = str(receipt.pop("receipt_hash", ""))
|
| 2854 |
+
observed = hashlib.sha256(
|
| 2855 |
+
_json_text(receipt).encode("utf-8")
|
| 2856 |
+
).hexdigest()
|
| 2857 |
+
expected_identity = {
|
| 2858 |
+
"namespace": row["namespace"],
|
| 2859 |
+
"owner_id": row["owner_id"],
|
| 2860 |
+
"request_id": row["request_id"],
|
| 2861 |
+
"session_id": row["session_id"],
|
| 2862 |
+
"step": int(row["step"]),
|
| 2863 |
+
"database_generation_id": self.database_generation_id,
|
| 2864 |
+
}
|
| 2865 |
+
if (
|
| 2866 |
+
claimed != row["receipt_hash"]
|
| 2867 |
+
or observed != row["receipt_hash"]
|
| 2868 |
+
or any(
|
| 2869 |
+
receipt.get(field) != expected
|
| 2870 |
+
for field, expected in expected_identity.items()
|
| 2871 |
+
)
|
| 2872 |
+
):
|
| 2873 |
+
raise ValueError("receipt digest mismatch")
|
| 2874 |
+
except (TypeError, ValueError, json.JSONDecodeError):
|
| 2875 |
+
digest_violations["invalid_receipt_digests"] += 1
|
| 2876 |
+
for row in connection.execute(
|
| 2877 |
+
"SELECT payload_json, payload_sha256 FROM proof_outbox"
|
| 2878 |
+
+ scoped_suffix,
|
| 2879 |
+
params,
|
| 2880 |
+
):
|
| 2881 |
+
if row["payload_json"] is None:
|
| 2882 |
+
continue
|
| 2883 |
+
try:
|
| 2884 |
+
payload = json.loads(row["payload_json"])
|
| 2885 |
+
if not isinstance(payload, dict):
|
| 2886 |
+
raise ValueError("proof payload must be an object")
|
| 2887 |
+
observed = self._effect_payload_digest(
|
| 2888 |
+
"proof_export", payload
|
| 2889 |
+
)
|
| 2890 |
+
if observed != row["payload_sha256"]:
|
| 2891 |
+
raise ValueError("proof digest mismatch")
|
| 2892 |
+
except (
|
| 2893 |
+
GDWConfigurationError,
|
| 2894 |
+
TypeError,
|
| 2895 |
+
ValueError,
|
| 2896 |
+
json.JSONDecodeError,
|
| 2897 |
+
):
|
| 2898 |
+
digest_violations["invalid_proof_digests"] += 1
|
| 2899 |
effect_rows = connection.execute(
|
| 2900 |
"""
|
| 2901 |
SELECT namespace, owner_id, idempotency_key,
|
|
|
|
| 2940 |
"ok": (
|
| 2941 |
check == "ok"
|
| 2942 |
and orphan_receipts == 0
|
| 2943 |
+
and not any(digest_violations.values())
|
| 2944 |
and invalid_effect_bindings == 0
|
| 2945 |
and invalid_exported_artifacts == 0
|
| 2946 |
),
|
|
|
|
| 2954 |
"dead_letter_effects": dead_letter_effects,
|
| 2955 |
"invalid_effect_bindings": invalid_effect_bindings,
|
| 2956 |
"invalid_exported_artifacts": invalid_exported_artifacts,
|
| 2957 |
+
**digest_violations,
|
| 2958 |
"counts": counts,
|
| 2959 |
"scope": "global" if global_scope else "owner",
|
| 2960 |
"journal_mode": str(
|
|
|
|
| 2969 |
result["owner_id"] = owner
|
| 2970 |
return result
|
| 2971 |
finally:
|
| 2972 |
+
if owns_connection:
|
| 2973 |
+
connection.close()
|
routers/gdw_frontier.py
CHANGED
|
@@ -1024,6 +1024,37 @@ def register(app, ns: str = "a11oy"):
|
|
| 1024 |
authorised_state_hash,
|
| 1025 |
)
|
| 1026 |
with workspace.transaction() as connection:
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 1027 |
previous = workspace.session_state(connection, payload.session_id)
|
| 1028 |
if previous is None:
|
| 1029 |
before_step = 0
|
|
|
|
| 1024 |
authorised_state_hash,
|
| 1025 |
)
|
| 1026 |
with workspace.transaction() as connection:
|
| 1027 |
+
cached = workspace.cached_request(connection, request_id)
|
| 1028 |
+
if cached is not None:
|
| 1029 |
+
cached_digest, cached_response = cached
|
| 1030 |
+
if cached_digest != request_digest:
|
| 1031 |
+
raise HTTPException(
|
| 1032 |
+
status_code=409,
|
| 1033 |
+
detail="X-Request-Id was already used with different content",
|
| 1034 |
+
)
|
| 1035 |
+
current_bundle = _policy_bundle_sha256()
|
| 1036 |
+
cached_bundle = (
|
| 1037 |
+
cached_response.get("audit", {})
|
| 1038 |
+
.get("governance", {})
|
| 1039 |
+
.get("colang", {})
|
| 1040 |
+
.get("bundle_sha256")
|
| 1041 |
+
)
|
| 1042 |
+
if not current_bundle or cached_bundle != current_bundle:
|
| 1043 |
+
raise HTTPException(
|
| 1044 |
+
status_code=409,
|
| 1045 |
+
detail="policy snapshot changed; replay refused",
|
| 1046 |
+
)
|
| 1047 |
+
cached_response["replayed"] = True
|
| 1048 |
+
selected_mode = cached_response["scheduler_mode"]
|
| 1049 |
+
decision = cached_response["decision"]
|
| 1050 |
+
receipt_hash = cached_response.get("receipt_hash") or ""
|
| 1051 |
+
_TELEMETRY.observe(
|
| 1052 |
+
(time.perf_counter() - started) * 1000.0,
|
| 1053 |
+
decision,
|
| 1054 |
+
selected_mode,
|
| 1055 |
+
False,
|
| 1056 |
+
)
|
| 1057 |
+
return cached_response
|
| 1058 |
previous = workspace.session_state(connection, payload.session_id)
|
| 1059 |
if previous is None:
|
| 1060 |
before_step = 0
|