From 233b1d3101ff5ca274ba584b63e40433c4867d5a Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Wed, 18 Mar 2026 19:59:35 +0000 Subject: [PATCH] fix(ci): aggregate all mock calls in langfuse e2e test to fix race condition The _verify_langfuse_call helper only inspected the last mock call (mock_post.call_args), but the Langfuse SDK may split trace-create and generation-create events across separate HTTP flush cycles. This caused an IndexError when the last call's batch contained only one event type. Fix: iterate over mock_post.call_args_list to collect batch items from ALL calls. Also add a safety assertion after filtering by trace_id and mark all langfuse e2e tests with @pytest.mark.flaky(retries=3) as an extra safety net for any residual timing issues. Co-authored-by: Ishaan Jaff --- .../test_langfuse_e2e_test.py | 50 ++++++++++++++++--- 1 file changed, 42 insertions(+), 8 deletions(-) diff --git a/tests/logging_callback_tests/test_langfuse_e2e_test.py b/tests/logging_callback_tests/test_langfuse_e2e_test.py index 369988ed405..f7affaec295 100644 --- a/tests/logging_callback_tests/test_langfuse_e2e_test.py +++ b/tests/logging_callback_tests/test_langfuse_e2e_test.py @@ -64,6 +64,13 @@ def assert_langfuse_request_matches_expected( "actual_request_body after filtering", json.dumps(actual_request_body, indent=4) ) + assert len(actual_request_body["batch"]) >= 2, ( + f"Expected at least 2 batch items (trace-create + generation-create) " + f"after filtering by trace_id={trace_id}, " + f"but got {len(actual_request_body['batch'])}. " + f"Items: {json.dumps(actual_request_body['batch'], indent=2)}" + ) + # Replace dynamic values in actual request body for item in actual_request_body["batch"]: @@ -150,19 +157,36 @@ class TestLangfuseLogging: """Helper method to verify Langfuse API calls""" await asyncio.sleep(3) - # Verify the call + # Verify at least one call was made assert mock_post.call_count >= 1 - url = mock_post.call_args[0][0] - request_body = mock_post.call_args[1].get("content") - # Parse the JSON string into a dict for assertions - actual_request_body = json.loads(request_body) + # Aggregate batch items from ALL calls — the Langfuse SDK may split + # trace-create and generation-create across separate HTTP flushes. + langfuse_url = "https://us.cloud.langfuse.com/api/public/ingestion" + all_batch_items: list = [] + metadata: Optional[dict] = None + for call in mock_post.call_args_list: + url = call[0][0] + if url != langfuse_url: + continue + request_body = call[1].get("content") + if request_body: + body = json.loads(request_body) + all_batch_items.extend(body.get("batch", [])) + if metadata is None: + metadata = body.get("metadata") - print("\nMocked Request Details:") - print(f"URL: {url}") + assert len(all_batch_items) > 0, "No Langfuse ingestion calls found" + assert metadata is not None, "No metadata found in Langfuse calls" + + actual_request_body = { + "batch": all_batch_items, + "metadata": metadata, + } + + print("\nMocked Request Details (aggregated from all calls):") print(f"Request Body: {json.dumps(actual_request_body, indent=4)}") - assert url == "https://us.cloud.langfuse.com/api/public/ingestion" assert_langfuse_request_matches_expected( actual_request_body, expected_file_name, @@ -170,6 +194,7 @@ class TestLangfuseLogging: ) @pytest.mark.asyncio + @pytest.mark.flaky(retries=3, delay=1) async def test_langfuse_logging_completion(self, mock_setup): """Test Langfuse logging for chat completion""" setup = mock_setup @@ -185,6 +210,7 @@ class TestLangfuseLogging: ) @pytest.mark.asyncio + @pytest.mark.flaky(retries=3, delay=1) async def test_langfuse_logging_completion_with_tags(self, mock_setup): """Test Langfuse logging for chat completion with tags""" setup = mock_setup @@ -203,6 +229,7 @@ class TestLangfuseLogging: ) @pytest.mark.asyncio + @pytest.mark.flaky(retries=3, delay=1) async def test_langfuse_logging_completion_with_tags_stream(self, mock_setup): """Test Langfuse logging for chat completion with tags""" setup = mock_setup @@ -223,6 +250,7 @@ class TestLangfuseLogging: ) @pytest.mark.asyncio + @pytest.mark.flaky(retries=3, delay=1) async def test_langfuse_logging_completion_with_langfuse_metadata(self, mock_setup): """Test Langfuse logging for chat completion with metadata for langfuse""" setup = mock_setup @@ -252,6 +280,7 @@ class TestLangfuseLogging: ) @pytest.mark.asyncio + @pytest.mark.flaky(retries=3, delay=1) async def test_langfuse_logging_with_non_serializable_metadata(self, mock_setup): """Test Langfuse logging with metadata that requires preparation (Pydantic models, sets, etc)""" from pydantic import BaseModel @@ -358,6 +387,7 @@ class TestLangfuseLogging: ) @pytest.mark.asyncio + @pytest.mark.flaky(retries=3, delay=1) async def test_langfuse_logging_completion_with_malformed_llm_response( self, mock_setup ): @@ -387,6 +417,7 @@ class TestLangfuseLogging: ) @pytest.mark.asyncio + @pytest.mark.flaky(retries=3, delay=1) async def test_langfuse_logging_completion_with_bedrock_llm_response( self, mock_setup ): @@ -418,6 +449,7 @@ class TestLangfuseLogging: setup["mock_post"], "completion_with_bedrock_call.json", setup["trace_id"] ) @pytest.mark.asyncio + @pytest.mark.flaky(retries=3, delay=1) async def test_langfuse_logging_completion_with_vertex_llm_response( self, mock_setup ): @@ -449,6 +481,7 @@ class TestLangfuseLogging: ) @pytest.mark.asyncio + @pytest.mark.flaky(retries=3, delay=1) async def test_langfuse_logging_vllm_embedding(self, mock_setup): """ Test that the request sent to the vllm embedding endpoint is correct. @@ -500,6 +533,7 @@ class TestLangfuseLogging: ) @pytest.mark.asyncio + @pytest.mark.flaky(retries=3, delay=1) async def test_langfuse_logging_with_router(self, mock_setup): """Test Langfuse logging with router""" litellm._turn_on_debug()