Sitelet https://github.com/NodeOps-app/treg2/commit/46c5510db0c9677825bd06323e669c859f3d3eee
Skip to content

Commit 46c5510

Browse files
authored
fix(archive): measure write stages and skip unused normalization
* fix(archive): measure write stages and skip unused normalization * test: use UTC dates for daily counter expectations
1 parent 8c700dc commit 46c5510

6 files changed

Lines changed: 235 additions & 34 deletions

File tree

‎docs/context/architecture/archive.md‎

Lines changed: 21 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -158,6 +158,22 @@ latency cannot delay it. The separate `archive_body_stored` completion event car
158158
`call_ref`. A failed R2 upload followed by a committed DB copy is not dropped. R2-only upload
159159
failure still records a hash-only snapshot and statistics. These remain best-effort background
160160
writes: a killed process can lose completion events, but cannot withhold the calling event.
161+
The same completion event measures archive stages, including elapsed work on timeout or cancellation:
162+
163+
| Fields (milliseconds) | Scope |
164+
| --- | --- |
165+
| `compare_sem_wait_ms`, `compare_ms` | Precomparison slot wait, then pointer queries/reads and any JSON normalization. Both precede the DB write deadline. |
166+
| `record_key_wait_ms`, `record_sem_wait_ms` | Same-key lock wait, then write-slot wait, inside the DB write deadline. |
167+
| `record_db_ms` | Wall time while holding the write slot: pool checkout, body packing, SQL/commit, and integrity retries/backoff. This is not pure SQL time. |
168+
| `observe_sem_wait_ms`, `observe_ms` | Optional post-commit change-report slot wait and work, outside the DB write deadline. |
169+
170+
Unentered phases are `null`. `failure_phase` names the measured phase interrupted by an escaping
171+
exception (including cancellation); it is `null` when none was interrupted, including a queue
172+
rejection before recording starts. Internally handled comparison failures still fall back to raw
173+
hashes and use the existing counters. A post-commit observation cancellation does not mean the
174+
snapshot was lost: `storage`/`dropped` continue to describe the committed write. Timings add no
175+
DB writes or per-call events and do not change deadlines, the two archive slots, or pool sizes.
176+
161177
Stats count snapshots with recoverable bodies in DB or R2, including deduplicated versions;
162178
`kept_bytes` is logical retained response bytes, not PostgreSQL physical table size.
163179

@@ -751,8 +767,11 @@ Deleting `items[*].request_id` keeps all elements; deleting `items[*]` deletes t
751767
Comparison has no six-level reporting limit and never uses reported/truncated paths as policy.
752768

753769
`_ignored_matches` preloads at most the latest and decisive snapshot bodies before the DB write.
754-
It skips body reads for identical raw hashes. Pointer sessions close before object I/O. The pre-read
755-
uses the shared archive semaphore and a three-second budget, separate from the write deadline.
770+
It checks the key and raw hashes first, skipping new-response JSON normalization unless a differing
771+
candidate has a readable body pointer. New keys and raw-identical baselines therefore need no
772+
normalization; identical raw hashes also skip body reads. Pointer sessions close before JSON
773+
normalization or object I/O. The pre-read uses the shared archive semaphore and a three-second
774+
budget, separate from the write deadline.
756775
Only matched snapshot IDs are passed to `_store_locked`; a concurrently changed baseline falls
757776
back to raw hashes without holding a row lock across I/O or adding a reconciliation write.
758777
`ignore_body_unavailable` and `ignore_comparison_failed` retain their diagnostic names for both

‎src/treg/archive.py‎

Lines changed: 35 additions & 24 deletions
Original file line numberDiff line numberDiff line change
@@ -536,32 +536,39 @@ async def _store(
536536
ignore_paths = cache.get("ignore_paths", []) if isinstance(cache, dict) else []
537537
ignored_matches = set()
538538
if plan.storage is not None and origin in ("caller", "refresh"):
539-
async with _get_sem():
540-
ignored_matches = await _ignored_matches(kh, body, ignore_paths)
539+
async with observation.wait(_get_sem(), "compare_sem_wait"):
540+
with observation.measure("compare"):
541+
ignored_matches = await _ignored_matches(kh, body, ignore_paths)
541542

542543
# Same-key waiters must queue before taking a scarce database-write slot. Otherwise four
543544
# duplicate recordings can occupy the whole semaphore while only one touches the database.
544-
async with asyncio.timeout(_TERMINAL_DB_S if origin == "async_terminal" else _STORE_TIMEOUT_S), _get_key_lock(kh):
545-
async with _get_sem():
546-
# Postgres row locking handles other processes. A retry also covers the narrow
547-
# first-key race and multi-process SQLite, where SELECT FOR UPDATE is ignored.
548-
for attempt in range(4):
549-
try:
550-
change = await _store_locked(
551-
method=method, endpoint_id=endpoint_id, provider=provider, url=url,
552-
caller_body=caller_body, headers=headers, status_code=status_code,
553-
media_type=media_type, body=body, origin=origin,
554-
key_hash=kh, body_hash=ch, plan=plan, ignored_matches=ignored_matches)
555-
stored, reason = plan.storage, plan.reason
556-
break
557-
except IntegrityError:
558-
if attempt == 3:
559-
raise
560-
await asyncio.sleep(0.01 * (attempt + 1))
545+
async with (
546+
asyncio.timeout(_TERMINAL_DB_S if origin == "async_terminal" else _STORE_TIMEOUT_S),
547+
observation.wait(_get_key_lock(kh), "record_key_wait"),
548+
):
549+
async with observation.wait(_get_sem(), "record_sem_wait"):
550+
# Wall time includes pool checkout, packing, SQL/commit and integrity retries.
551+
with observation.measure("record_db"):
552+
# Postgres row locking handles other processes. A retry also covers the narrow
553+
# first-key race and multi-process SQLite, where SELECT FOR UPDATE is ignored.
554+
for attempt in range(4):
555+
try:
556+
change = await _store_locked(
557+
method=method, endpoint_id=endpoint_id, provider=provider, url=url,
558+
caller_body=caller_body, headers=headers, status_code=status_code,
559+
media_type=media_type, body=body, origin=origin,
560+
key_hash=kh, body_hash=ch, plan=plan, ignored_matches=ignored_matches)
561+
stored, reason = plan.storage, plan.reason
562+
break
563+
except IntegrityError:
564+
if attempt == 3:
565+
raise
566+
await asyncio.sleep(0.01 * (attempt + 1))
561567
if change is not None and get_settings().archive_change_observation_enabled:
562568
previous_id, masked_by_ignore = change
563-
async with _get_sem():
564-
await _observe_change(previous_id, body, endpoint_id, provider, masked_by_ignore)
569+
async with observation.wait(_get_sem(), "observe_sem_wait"):
570+
with observation.measure("observe"):
571+
await _observe_change(previous_id, body, endpoint_id, provider, masked_by_ignore)
565572
except asyncio.CancelledError:
566573
reason = "cancelled"
567574
raise
@@ -666,9 +673,6 @@ async def _ignored_matches(key_hash: str, body: bytes, paths: list[str]) -> set[
666673
matches = set()
667674
try:
668675
async with asyncio.timeout(_CHANGE_TIMEOUT_S):
669-
new_hash = await _change_compute(_normalized_hash, body, paths)
670-
if new_hash is None:
671-
return matches
672676
async with background_session_maker() as s:
673677
key = (await s.execute(select(ArchiveKey).where(
674678
ArchiveKey.key_hash == key_hash))).scalar_one_or_none()
@@ -685,6 +689,13 @@ async def _ignored_matches(key_hash: str, body: bytes, paths: list[str]) -> set[
685689
matches.update(row.id for row in rows if row.content_hash == raw_hash)
686690
pointers = [(row.id, await archive_bodies.pointer(s, row, "observation"))
687691
for row in rows if row.content_hash != raw_hash and _has_change_body(row)]
692+
if not pointers:
693+
return matches
694+
# New keys and raw-identical baselines need no JSON parsing or serialization.
695+
# Close the pointer session before any off-thread work or object I/O.
696+
new_hash = await _change_compute(_normalized_hash, body, paths)
697+
if new_hash is None:
698+
return matches
688699
for snapshot_id, pointer in pointers:
689700
previous = await archive_bodies.read(pointer, "observation")
690701
if previous is None:

‎src/treg/archive_bodies.py‎

Lines changed: 29 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
"""Archive body I/O and upload scheduling, separate from archive indexing and TTL learning."""
22
import asyncio
33
from collections import Counter, OrderedDict
4+
from contextlib import asynccontextmanager, contextmanager
45
from dataclasses import dataclass
56
import logging
67
import re
@@ -62,6 +63,31 @@ def __init__(self, *, call_ref=None, emit=None):
6263
self.finished = False
6364
self.props = {"archive_body_upload_status": "not_requested", "archive_body_upload_ms": 0.0,
6465
"archive_body_queue_wait_ms": 0.0}
66+
self.timings = dict.fromkeys(("compare_sem_wait", "compare", "record_key_wait",
67+
"record_sem_wait", "record_db", "observe_sem_wait", "observe"))
68+
self.failure_phase = None
69+
70+
@contextmanager
71+
def measure(self, phase):
72+
"""Include interrupted work; an unentered phase stays unknown rather than zero."""
73+
started = time.monotonic()
74+
try:
75+
yield
76+
except BaseException:
77+
self.failure_phase = phase
78+
raise
79+
finally:
80+
self.timings[phase] = (time.monotonic() - started) * 1000
81+
82+
@asynccontextmanager
83+
async def wait(self, gate, phase):
84+
# Measure acquisition only, then retain the original lock/semaphore lifetime.
85+
with self.measure(phase):
86+
await gate.acquire()
87+
try:
88+
yield
89+
finally:
90+
gate.release()
6591

6692
def finish(self, *, storage=None, reason=None):
6793
if self.finished:
@@ -73,6 +99,9 @@ def finish(self, *, storage=None, reason=None):
7399
"upload_ms": round(self.props["archive_body_upload_ms"], 3),
74100
"queue_wait_ms": round(self.props["archive_body_queue_wait_ms"], 3),
75101
"dropped": storage is None, "drop_reason": reason or "none"}
102+
data.update({phase + "_ms": round(ms, 3) if ms is not None else None
103+
for phase, ms in self.timings.items()})
104+
data["failure_phase"] = self.failure_phase
76105
if self.emit is not None:
77106
self.emit(data)
78107
elif self.call_ref:

‎tests/test_archive_r2.py‎

Lines changed: 138 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -89,6 +89,52 @@ async def test_upload_precedes_pointer_and_call_does_not_wait(clients, r2, monke
8989
assert len(props) == 1 and props[0]['upload_status'] == 'uploaded'
9090
assert props[0]['upload_ms'] >= 0
9191
assert props[0]['dropped'] is False
92+
assert props[0]['failure_phase'] is None
93+
for phase in ('compare_sem_wait', 'compare', 'record_key_wait', 'record_sem_wait', 'record_db'):
94+
assert props[0][phase + '_ms'] >= 0
95+
96+
97+
@pytest.mark.parametrize('phase', ['record_key_wait', 'record_sem_wait', 'record_db'])
98+
async def test_call_reports_where_archive_db_deadline_expires(clients, r2, monkeypatch, phase):
99+
# The upstream answer and R2 object succeed even when archive indexing cannot finish.
100+
monkeypatch.setattr(archive, '_STORE_TIMEOUT_S', .05)
101+
blocked = asyncio.Lock()
102+
await blocked.acquire()
103+
if phase == 'record_key_wait':
104+
monkeypatch.setattr(archive, '_get_key_lock', lambda key: blocked)
105+
elif phase == 'record_sem_wait':
106+
comparison_sem = asyncio.Semaphore(2)
107+
calls = 0
108+
def get_sem():
109+
nonlocal calls
110+
calls += 1
111+
return comparison_sem if calls == 1 else blocked
112+
monkeypatch.setattr(archive, '_get_sem', get_sem)
113+
else:
114+
async def blocked_store(**kwargs):
115+
await asyncio.Event().wait()
116+
monkeypatch.setattr(archive, '_store_locked', blocked_store)
117+
reports = []
118+
monkeypatch.setattr(service.analytics, 'capture',
119+
lambda who, name, props, **kw: reports.append(props)
120+
if name == 'archive_body_stored' else None)
121+
try:
122+
response = await clients.get(URL)
123+
assert response.status_code == 200 and response.content == RAW
124+
await archive.drain()
125+
finally:
126+
blocked.release()
127+
assert r2.objects[archive.content_hash(RAW)] == RAW
128+
assert await snapshots() == []
129+
assert len(reports) == 1
130+
report = reports[0]
131+
assert report['drop_reason'] == 'record_timeout' and report['dropped'] is True
132+
assert report['failure_phase'] == phase
133+
assert report[phase + '_ms'] >= 40
134+
if phase != 'record_db':
135+
assert report['record_db_ms'] is None
136+
if phase == 'record_key_wait':
137+
assert report['record_sem_wait_ms'] is None
92138

93139

94140
@pytest.mark.parametrize('mode,expected_rows', [('both', 1), ('r2', 1)])
@@ -642,6 +688,65 @@ async def blocked(**kwargs):
642688
assert any('record_timeout' in record.message for record in caplog.records)
643689

644690

691+
@pytest.mark.parametrize('phase', [
692+
'compare_sem_wait', 'compare', 'record_key_wait', 'record_sem_wait', 'record_db',
693+
'observe_sem_wait', 'observe',
694+
])
695+
async def test_cancelled_record_reports_elapsed_phase(clients, r2, monkeypatch, phase):
696+
entered = asyncio.Event()
697+
class BlockedGate(asyncio.Semaphore):
698+
async def acquire(self):
699+
entered.set()
700+
return await super().acquire()
701+
blocked_gate = BlockedGate(0)
702+
async def blocked_work(*args, **kwargs):
703+
entered.set()
704+
await asyncio.Event().wait()
705+
sem = asyncio.Semaphore(2)
706+
key_lock = asyncio.Lock()
707+
sem_calls = 0
708+
def get_sem():
709+
nonlocal sem_calls
710+
sem_calls += 1
711+
blocked_call = {'compare_sem_wait': 1, 'record_sem_wait': 2, 'observe_sem_wait': 3}.get(phase)
712+
return blocked_gate if sem_calls == blocked_call else sem
713+
monkeypatch.setattr(archive, '_get_sem', get_sem)
714+
monkeypatch.setattr(archive, '_get_key_lock',
715+
lambda key: blocked_gate if phase == 'record_key_wait' else key_lock)
716+
if phase == 'compare':
717+
monkeypatch.setattr(archive, '_ignored_matches', blocked_work)
718+
elif phase == 'record_db':
719+
monkeypatch.setattr(archive, '_store_locked', blocked_work)
720+
elif phase == 'observe':
721+
monkeypatch.setattr(archive, '_observe_change', blocked_work)
722+
# A prior version ensures the post-commit observation phases are reached.
723+
if phase.startswith('observe'):
724+
await archive._store(method='GET', endpoint_id=EP, provider='tikhub', url=URL,
725+
caller_body=b'', headers={}, status_code=200,
726+
media_type='application/json', body=b'{"before":1}',
727+
plan=archive_bodies.WritePlan('db'))
728+
sem_calls = 0
729+
reports = []
730+
task = asyncio.create_task(archive._store(
731+
method='GET', endpoint_id=EP, provider='tikhub', url=URL, caller_body=b'',
732+
headers={}, status_code=200, media_type='application/json', body=RAW,
733+
plan=archive_bodies.WritePlan('db'),
734+
observation=archive_bodies.StorageReport(emit=reports.append)))
735+
try:
736+
await asyncio.wait_for(entered.wait(), 2)
737+
await asyncio.sleep(.01)
738+
finally:
739+
task.cancel()
740+
with pytest.raises(asyncio.CancelledError):
741+
await task
742+
assert len(reports) == 1
743+
assert reports[0]['drop_reason'] == 'cancelled'
744+
assert reports[0]['failure_phase'] == phase
745+
assert reports[0][phase + '_ms'] > 0
746+
assert reports[0]['dropped'] is (not phase.startswith('observe'))
747+
assert sem._value == 2 and not key_lock.locked() and blocked_gate._value == 0
748+
749+
645750
async def test_cancel_before_start_reports_once(clients, r2, monkeypatch):
646751
reports = []
647752
archive.record(method='GET', endpoint_id=EP, provider='tikhub', url=URL,
@@ -652,6 +757,8 @@ async def test_cancel_before_start_reports_once(clients, r2, monkeypatch):
652757
task.cancel()
653758
await archive.drain()
654759
assert len(reports) == 1 and reports[0]['drop_reason'] == 'cancelled'
760+
assert reports[0]['failure_phase'] is None
761+
assert reports[0]['record_key_wait_ms'] is None and reports[0]['record_db_ms'] is None
655762

656763

657764
@pytest.mark.parametrize('error,reason', [
@@ -690,6 +797,37 @@ async def test_same_body_refetch_skips_duplicate_put(clients, r2, monkeypatch):
690797
assert all(p['storage'] == 'both' and not p['dropped'] for p in reports)
691798

692799

800+
@pytest.mark.parametrize('write_mode', ['both', 'r2'])
801+
async def test_new_and_identical_calls_skip_unused_normalization(clients, r2, monkeypatch, write_mode):
802+
monkeypatch.setattr(get_settings(), 'archive_body_write', write_mode)
803+
if write_mode == 'r2':
804+
monkeypatch.setattr(get_settings(), 'archive_body_read_lookup', 'r2-first')
805+
normalized = []
806+
real_normalize = archive._normalized_hash
807+
def normalize(body, paths):
808+
normalized.append(body)
809+
return real_normalize(body, paths)
810+
monkeypatch.setattr(archive, '_normalized_hash', normalize)
811+
for _ in range(2):
812+
response = await clients.get(URL, headers={'Cache-Control': 'no-cache'})
813+
assert response.status_code == 200 and response.content == RAW
814+
await archive.drain()
815+
assert normalized == []
816+
817+
# Default JSON equality still applies when the raw bytes actually differ.
818+
reformatted = b' ' + RAW + b'\n'
819+
monkeypatch.setattr(service, 'relay', _fake_relay(200, reformatted))
820+
assert (await clients.get(URL, headers={'Cache-Control': 'no-cache'})).content == reformatted
821+
await archive.drain()
822+
assert RAW in normalized and reformatted in normalized
823+
async with db.session_maker() as s:
824+
key = (await s.execute(select(ArchiveKey))).scalar_one()
825+
assert key.stable_seen == 2 and key.change_seen == 0
826+
rows = await snapshots()
827+
assert len(rows) == 3
828+
assert [row.content_hash for row in rows] == [archive.content_hash(raw) for raw in (RAW, RAW, reformatted)]
829+
830+
693831
async def test_concurrent_same_body_calls_share_one_put(clients, r2):
694832
r2.gate = asyncio.Event()
695833
try:

‎tests/test_daily_spend_counter.py‎

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,7 @@
55
semantics the journal view (`spent_today_from_ledger`) has always had: settled today plus still
66
held from today, reset at the UTC day boundary, with a hold opened yesterday belonging to yesterday.
77
"""
8-
from datetime import date, timedelta
8+
from datetime import timedelta
99

1010
import pytest
1111
from httpx import ASGITransport, AsyncClient
@@ -15,6 +15,7 @@
1515
from treg.domain import money as ledger
1616
from treg.infra.db import reset_db, session_maker
1717
from treg.models import Hold, Org
18+
from treg.timeutil import utcnow_naive
1819
from conftest import make_upstream, verified_signup
1920

2021
EP = "acme.thing.get"
@@ -86,7 +87,7 @@ async def test_counter_resets_on_a_new_utc_day(c: AsyncClient):
8687
assert await _both(org_id) == (spent, spent)
8788

8889
# Move the counter to "yesterday" - as the day rolling over would leave it - and it reads 0.
89-
yesterday = date.today() - timedelta(days=1)
90+
yesterday = utcnow_naive().date() - timedelta(days=1)
9091
async with session_maker() as db:
9192
await db.execute(update(Org).where(Org.id == org_id).values(spent_today_day=yesterday))
9293
await db.commit()

0 commit comments

Comments
 (0)