Skip to content

Commit 2d23940

Browse files
hpliStartAgainPangjiping
authored andcommitted
feat(egress): track bounded request snapshot lifecycles
1 parent 0f0a888 commit 2d23940

3 files changed

Lines changed: 658 additions & 16 deletions

File tree

‎components/egress/mitmscripts/tls_registry.py‎

Lines changed: 144 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -39,11 +39,12 @@
3939
"snapshot_missing",
4040
"snapshot_mismatch",
4141
"receiver_unavailable",
42+
"request_registry_exhausted",
4243
]
4344

4445

4546
class RegistryError(Exception):
46-
"""A fixed activation error that never contains snapshot data."""
47+
"""A fixed registry error that never contains snapshot data."""
4748

4849

4950
@dataclass(frozen=True, slots=True)
@@ -63,13 +64,37 @@ class AdmissionResult:
6364
token: AdmissionToken | None = None
6465

6566

67+
@dataclass(frozen=True, slots=True)
68+
class PendingRequest:
69+
"""Internal metadata only; neither a finish capability nor proof of drain."""
70+
71+
serial: int
72+
connection_serial: int
73+
revision: Revision
74+
75+
76+
@dataclass(frozen=True, slots=True, eq=False)
77+
class RequestHandle:
78+
"""Exact in-process request membership; finish on every terminal path."""
79+
80+
serial: int
81+
owner: object = field(repr=False)
82+
connection: AdmissionToken = field(repr=False)
83+
snapshot: Snapshot = field(repr=False)
84+
85+
@property
86+
def revision(self) -> Revision:
87+
return self.snapshot.revision
88+
89+
6690
@dataclass(frozen=True, slots=True)
6791
class RequestAdmission:
6892
"""Request lifecycle eligibility, not authorization to inject credentials."""
6993

7094
action: Literal["allow", "deny"]
7195
reason: RequestReason
7296
snapshot: Snapshot | None = field(default=None, repr=False)
97+
handle: RequestHandle | None = field(default=None, repr=False)
7398

7499

75100
class BoundConnectionRegistry:
@@ -79,12 +104,21 @@ class BoundConnectionRegistry:
79104
connections before acknowledging a host-removal or generation transition.
80105
"""
81106

82-
def __init__(self, *, capacity: int, receiver: Receiver | None = None) -> None:
107+
def __init__(
108+
self, *, capacity: int, receiver: Receiver | None = None,
109+
request_capacity: int | None = None,
110+
) -> None:
83111
if type(capacity) is not int or capacity <= 0:
84112
raise ValueError("positive TLS registry capacity required")
85113
if receiver is not None and type(receiver) is not Receiver:
86114
raise TypeError("TLS registry receiver required")
115+
if request_capacity is None:
116+
request_capacity = capacity
117+
if type(request_capacity) is not int or request_capacity <= 0:
118+
raise ValueError("positive request registry capacity required")
87119
self._capacity = capacity
120+
# This compatibility default is not a production HTTP/2 sizing policy.
121+
self._request_capacity = request_capacity
88122
self._receiver = receiver
89123
self._lock = threading.Lock()
90124
self._owner = object()
@@ -94,6 +128,18 @@ def __init__(self, *, capacity: int, receiver: Receiver | None = None) -> None:
94128
self._entries: dict[int, AdmissionToken] = {}
95129
self._request_fenced: set[int] = set()
96130
self._next_serial = 0
131+
self._requests: dict[int, RequestHandle] = {}
132+
self._requests_by_connection: dict[int, set[int]] = {}
133+
self._next_request_serial = 0
134+
135+
def _owns_connection(self, token: AdmissionToken | None) -> bool:
136+
"""Check exact membership with the Registry lock already held."""
137+
return (
138+
type(token) is AdmissionToken
139+
and token.owner is self._owner
140+
and type(token.serial) is int
141+
and self._entries.get(token.serial) is token
142+
)
97143

98144
def activate(self, snapshot: Snapshot) -> tuple[AdmissionToken, ...]:
99145
"""Publish a confirmed snapshot and return newly uncovered memberships.
@@ -185,24 +231,22 @@ def admit(
185231
return AdmissionResult("decrypt", "binding_host", token)
186232

187233
def acquire_request(self, token: AdmissionToken | None) -> RequestAdmission:
188-
"""Pin one coherent snapshot while connection eligibility is stable.
234+
"""Pin and register one snapshot while connection eligibility is stable.
189235
190236
Lock order is Registry -> Receiver. The receiver must never call back
191237
into this registry while holding its state lock. A successful request
192238
is admitted when Receiver.acquire pins its snapshot, even if commit or
193239
close happens before this method returns. Its caller must retain that
194240
same snapshot through binding checks, injection and response redaction.
241+
The caller must finish its handle in every completion/cancellation/error
242+
path. Terminal connection release also removes its request records, but
243+
does not revoke external handles or cancel work still using them.
195244
196245
Independent Receiver/Registry publications may temporarily deny requests;
197246
this primitive is not their joint commit or a transport drain owner.
198247
"""
199248
with self._lock:
200-
if (
201-
type(token) is not AdmissionToken
202-
or token.owner is not self._owner
203-
or type(token.serial) is not int
204-
or self._entries.get(token.serial) is not token
205-
):
249+
if not self._owns_connection(token):
206250
return RequestAdmission("deny", "invalid_token")
207251
if self._closed:
208252
return RequestAdmission("deny", "registry_closed")
@@ -216,6 +260,8 @@ def acquire_request(self, token: AdmissionToken | None) -> RequestAdmission:
216260
current.control_generation, current.subject_generation
217261
):
218262
return RequestAdmission("deny", "invalid_token")
263+
if len(self._requests) >= self._request_capacity:
264+
return RequestAdmission("deny", "request_registry_exhausted")
219265
try:
220266
snapshot = self._receiver.acquire()
221267
except Exception: # noqa: BLE001 - never expose receiver error contents
@@ -224,19 +270,103 @@ def acquire_request(self, token: AdmissionToken | None) -> RequestAdmission:
224270
return RequestAdmission("deny", "snapshot_missing")
225271
if snapshot.revision != current:
226272
return RequestAdmission("deny", "snapshot_mismatch")
227-
return RequestAdmission("allow", "admitted", snapshot)
273+
serial = self._next_request_serial + 1
274+
handle = RequestHandle(serial, self._owner, token, snapshot)
275+
result = RequestAdmission("allow", "admitted", snapshot, handle)
276+
members = self._requests_by_connection.get(token.serial)
277+
if members is None:
278+
members = set()
279+
self._next_request_serial = serial
280+
registered = True
281+
try:
282+
self._requests[serial] = handle
283+
self._requests_by_connection[token.serial] = members
284+
members.add(serial)
285+
except Exception: # noqa: BLE001 - roll back without exposing payloads
286+
self._requests.pop(serial, None)
287+
members.discard(serial)
288+
if not members:
289+
self._requests_by_connection.pop(token.serial, None)
290+
registered = False
291+
if not registered:
292+
# Raise outside the handler: no secret-bearing exception context.
293+
raise RegistryError("request registration failed")
294+
return result
295+
296+
def finish_request(self, handle: RequestHandle | None) -> bool:
297+
"""Idempotently remove an exact live request, leaving its connection open."""
298+
with self._lock:
299+
if (
300+
type(handle) is not RequestHandle
301+
or handle.owner is not self._owner
302+
or type(handle.serial) is not int
303+
or self._requests.get(handle.serial) is not handle
304+
):
305+
return False
306+
del self._requests[handle.serial]
307+
connection_serial = handle.connection.serial
308+
members = self._requests_by_connection[connection_serial]
309+
members.remove(handle.serial)
310+
if not members:
311+
del self._requests_by_connection[connection_serial]
312+
return True
228313

229314
def release(self, token: AdmissionToken | None) -> bool:
230-
"""Idempotently remove an exact admission; serials are never reused."""
231-
if type(token) is not AdmissionToken or token.owner is not self._owner:
232-
return False
315+
"""Remove an exact terminal connection and all of its request records.
316+
317+
Only call after the transport owner confirms terminal state, not to start
318+
drain. This clears accounting, not external references or running work.
319+
A request admitted first may return its handle after terminal release.
320+
"""
233321
with self._lock:
234-
if self._entries.get(token.serial) != token:
322+
if not self._owns_connection(token):
235323
return False
324+
for serial in self._requests_by_connection.pop(token.serial, ()):
325+
del self._requests[serial]
236326
del self._entries[token.serial]
237327
self._request_fenced.discard(token.serial)
238328
return True
239329

330+
def pending_requests(
331+
self, *, connection: AdmissionToken | None = None,
332+
revision: Revision | None = None, after_serial: int = 0, limit: int = 128,
333+
) -> tuple[PendingRequest, ...]:
334+
"""Return a bounded metadata page, without snapshots or finish handles.
335+
336+
Each page is consistent under the lock; pages are not a frozen view.
337+
Completions can disappear and new admissions can appear between pages.
338+
The last returned serial is the next cursor. Empty is not proof of
339+
transport drain or permission to acknowledge a public mutation.
340+
"""
341+
if (
342+
type(after_serial) is not int or after_serial < 0
343+
or type(limit) is not int or not 1 <= limit <= 128
344+
or revision is not None and type(revision) is not Revision
345+
):
346+
raise ValueError("invalid pending request query")
347+
with self._lock:
348+
if connection is not None and not self._owns_connection(connection):
349+
return ()
350+
result = []
351+
# Dict insertion order is serial order; serials are never reused.
352+
for serial, handle in self._requests.items():
353+
if (
354+
serial <= after_serial
355+
or connection is not None and handle.connection is not connection
356+
or revision is not None and handle.revision != revision
357+
):
358+
continue
359+
result.append(PendingRequest(serial, handle.connection.serial, handle.revision))
360+
if len(result) == limit:
361+
break
362+
return tuple(result)
363+
364+
@property
365+
def request_count(self) -> int:
366+
"""Registry-held records only, not all externally retained snapshots."""
367+
with self._lock:
368+
return len(self._requests)
369+
240370
@property
241371
def count(self) -> int:
242372
with self._lock:

0 commit comments

Comments
 (0)