test(azure_sentinel): pin batch_size as a per-request bound under concurrent events (#40320)

* test(azure_sentinel): pin batch_size as a per-request bound under concurrent events

Adds a regression test to the mapped Azure Sentinel test file for the concurrency scenario from LIT-6920: 40 records logged concurrently at batch_size=5 while each ingestion request is still in flight. Asserts no request carries more than batch_size records, every record arrives exactly once in order, and the queue is empty afterwards. Runs for both the standard log queue and the audit log queue.

The test fails on the tree before #39880 (whole shared queue serialized per threshold send, then cleared) and passes on current staging. It is independent of the size-split coverage that #39880 added for LIT-5899.

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* test(azure_sentinel): gate the first send on events so later records provably arrive while it is in flight

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>

---------

Co-authored-by: yucheng <yucheng@berri.ai>
Co-authored-by: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
devin-ai-integration[bot] 2026-09-08 16:55:40 -07:00 committed by GitHub
parent 54dc1d7644
commit 075655c7ee
No known key found for this signature in database
GPG key ID: B5690EEEBB952194

View file

@ -867,6 +867,45 @@ async def test_azure_sentinel_concurrent_threshold_sends_collapse_into_one_attem
assert getattr(logger, queue_attr) == []
@pytest.mark.asyncio
@pytest.mark.parametrize("queue_attr, send_method, build_payloads", QUEUE_CASES)
async def test_azure_sentinel_batch_size_bounds_every_request_under_concurrent_events(
queue_attr, send_method, build_payloads
):
"""Lowering batch_size is the documented way to stay under the ingestion cap, so no request may
carry more than batch_size records even when events keep landing while a send is on the wire,
and every one of those records still has to arrive exactly once."""
logger = _build_logger(batch_size=5)
records = build_payloads(40)
attempts = []
first_send_started = asyncio.Event()
release_first_send = asyncio.Event()
async def _on_ingest(data):
attempts.append([record["id"] for record in json.loads(data.decode("utf-8"))])
if len(attempts) == 1:
first_send_started.set()
await release_first_send.wait()
return _accepted()
_install_ingestion(logger, _on_ingest)
sends = [asyncio.create_task(_log(logger, queue_attr, record)) for record in records]
await asyncio.wait_for(first_send_started.wait(), timeout=10)
assert attempts == [[record["id"] for record in records[:5]]]
assert getattr(logger, queue_attr) == records[5:]
release_first_send.set()
await asyncio.wait_for(asyncio.gather(*sends), timeout=10)
await logger.flush_queue()
assert max(len(attempt) for attempt in attempts) <= 5
assert [record_id for attempt in attempts for record_id in attempt] == [record["id"] for record in records]
assert getattr(logger, queue_attr) == []
@pytest.mark.asyncio
@pytest.mark.parametrize("queue_attr, send_method, build_payloads", QUEUE_CASES)
async def test_azure_sentinel_requeues_a_cancelled_send(