서비스 학습
SOURCE · be/tests/test_inference_timeouts.py · 2026-09-06

기다림을 끝낸 뒤 연결 자리도 반환되는가?

timeout이라는 오류 문자열만 확인하는 것으로는 충분하지 않다. 실제로 응답을 멈춘 TCP 서버를 상대로 대기 제한, 취소 전파, 연결 종료와 후속 요청까지 확인한다.

1. 이 파일을 만든 이유

기존 테스트는 handler에서 ReadTimeout을 직접 발생시켜 “HTTP 오류를 내부 timeout으로 바꾸는가”만 확인했다. 실제 시간이 지났을 때 기다림이 끝나는지, 취소된 요청이 연결 자리를 계속 차지하지 않는지는 별도의 검증이 필요하다.

2. 검증 구조

실제 Adapter와 HTTPX2 client

각 용도별 풀을 1개 연결로 제한한다. 연결 자리가 하나라도 누수되면 다음 요청은 진행할 수 없다.

held_upstream 임시 서버

요청을 받았다는 신호만 보내고 HTTP 응답은 보류한다. release 또는 상대 연결 종료를 기다린다.

실제 TCP를 사용하지만 모델은 없다. 테스트 서버는 요청마다 Connection: close로 응답한다. 따라서 다음 요청 성공은 같은 TCP 연결의 재사용이 아니라 같은 client 풀의 자리가 반환되어 새 연결을 열 수 있다는 증거다.

3. 입력 조건과 출력

조건내부 결과추가 검증
생성·probe 각각 실제 deadline 초과InferenceError(timeout)응답 해제 없이 socket 종료, 후속 요청 성공
생성·probe 각각 Task.cancel()CancelledError 그대로 전파같은 풀의 후속 요청 성공
생성·probe 각각 연결 1개 점유 후 추가 호출InferenceError(busy)기존 요청은 유지, 추가 요청은 upstream 미도착

세 조건 × 두 client 종류로 총 여섯 경우를 확인한다. deadline 조건에서는 생성·probe의 HTTP 요청 deadline을 둘 다 0.1초로, 취소·pool 조건에서는 둘 다 5초로 준다. 테스트용 값이므로 제품의 생성 300초·probe 요청당 5초 기본값을 바꾸지 않는다. 여기서 probe는 list_models 경로이며 전체 check_ready의 health→models 순서는 별도 readiness 테스트에서 확인한다.

4. 코드 해설

PendingRequest와 held_upstream

요청마다 path, release Event, closed Event를 만든다. asyncio.Queue로 요청 도착을 테스트 본문에 전달한다. 서버는 release.wait와 reader.read(1)을 함께 기다린다. 응답을 허용하기 전에 빈 bytes가 읽히면 상대 client가 연결을 닫았다는 EOF 신호다.

finally는 내부 대기 Task를 취소·회수하고 writer를 닫는다. closed 신호는 정리 뒤에 설정하므로 테스트는 정리 완료까지 기다릴 수 있다. 이 서버는 테스트에서 사용하는 단순 HTTP 요청 전용이며 일반 HTTP 서버 구현을 대신하지 않는다.

adapter_at: 왜 read timeout을 끄는가?

httpx2.Timeout(None, pool=0.05)로 응답 읽기 제한을 없애고 pool 대기만 0.05초로 둔다. 이렇게 해야 실제 deadline 테스트가 HTTPX2의 별도 ReadTimeout 때문에 우연히 통과하지 않는다. 응답이 없을 때 실행을 끝내는 것은 Adapter의 asyncio.timeout이다.

두 client 각각 max_connections=1로 만들고 같은 Adapter에 전달한다. 나머지 변환과 오류 처리는 production 구현 그대로다.

취소와 deadline을 어떻게 구분하는가?

deadline 조건에서는 아무 취소 명령도 보내지 않고 InferenceError(timeout)을 기다린다. 취소 조건에서는 충분히 긴 deadline을 둔 뒤 Task.cancel을 호출하고 CancelledError가 그대로 나오는지 확인한다. 둘 모두 서버의 release는 설정하지 않았는데도 socket이 닫혀야 한다.

바깥 wait_for의 2초 제한은 테스트가 고장 나도 영원히 기다리지 않도록 하는 안전장치다. 이것이 발생하면 기대한 InferenceError나 CancelledError와 다르므로 테스트는 실패한다.

실제 pool timeout은 무엇을 확인하는가?

첫 요청이 서버에서 대기하는 동안 두 번째 요청을 보낸다. 연결 자리가 하나뿐이어서 두 번째는 서버에 도달하지 못한 채 PoolTimeout → busy가 된다. 첫 Task가 아직 진행 중인지, 해당 socket이 닫히지 않았는지, 서버 관찰 Queue가 비어 있는지까지 검사한다. 이후 첫 응답을 해제하고 새로운 요청도 성공해야 한다.

5. 상태 변화

deadline / cancel: 연결 1개 점유 → 응답 대기 → 시간 초과 또는 취소 → client 연결 종료 → 서버 EOF 관찰 → 같은 풀의 다음 요청 성공 pool timeout: 첫 요청 연결 점유 → 두 번째 요청은 풀에서 대기 → 두 번째만 busy → 첫 요청은 유지 → 첫 응답 해제 → 다음 요청 성공

6. 실제 결과

여섯 경우 모두 통과했다. 전체 테스트는 47 passed다. 모의 예외가 아니라 실제 timer와 socket을 사용해 오류 종류뿐 아니라 자원 반환과 후속 요청을 확인했다.

검증 범위: 원격 서버가 TCP 종료를 본다는 것과 Gemma GPU kernel이 즉시 중단된다는 것은 다르다. 실제 Swift 생성 취소의 시점, browser disconnect 전달, streaming 수명은 다음 단계의 검증 대상이다.

7. 오류를 읽는 법

deadline 테스트에서 busy가 나오면 응답 대기 전에 풀에서 막힌 것이다. 취소 뒤 후속 요청이 막히면 연결 반환을 확인한다. closed Event가 끝나지 않으면 client가 socket을 닫았는지 점검한다. 외부의 2초 safety timeout 실패는 서비스 timeout 성공으로 해석하면 안 된다.

현재 전체 코드

"""Real socket tests: deadlines/cancellation must release the occupied pool slot."""

import asyncio
import json
from collections.abc import AsyncGenerator
from contextlib import asynccontextmanager
from dataclasses import dataclass, field

import httpx2
import pytest

from app.engines.base import GenerationRequest, InferenceError, Message
from app.engines.turbofieldfare import TurboFieldfareAdapter


MODEL = "gemma-4-26b-a4b-it"
REQUEST = GenerationRequest(messages=(Message(role="user", content="Synthetic"),))


@dataclass
class PendingRequest:
    path: str
    release: asyncio.Event = field(default_factory=asyncio.Event)
    closed: asyncio.Event = field(default_factory=asyncio.Event)


@asynccontextmanager
async def held_upstream() -> AsyncGenerator[
    tuple[str, asyncio.Queue[PendingRequest]], None,
]:
    pending: asyncio.Queue[PendingRequest] = asyncio.Queue()
    handlers: set[asyncio.Task] = set()

    async def handle(reader: asyncio.StreamReader, writer: asyncio.StreamWriter) -> None:
        task = asyncio.current_task()
        assert task is not None
        handlers.add(task)
        item: PendingRequest | None = None
        waits: list[asyncio.Task] = []
        try:
            head = (await reader.readuntil(b"\r\n\r\n")).decode("ascii")
            first, *headers = head.split("\r\n")
            _, path, _ = first.split()
            size = next((int(header.split(":", 1)[1]) for header in headers
                         if header.lower().startswith("content-length:")), 0)
            await reader.readexactly(size)
            item = PendingRequest(path)
            # No response until released. A peer EOF proves client cancellation
            # closed the socket, without relying on a fabricated HTTP exception.
            waits = [asyncio.create_task(item.release.wait()),
                     asyncio.create_task(reader.read(1))]
            await pending.put(item)
            done, _ = await asyncio.wait(waits, return_when=asyncio.FIRST_COMPLETED)
            if waits[1] in done:
                assert waits[1].result() == b""
                return
            if path == "/health":
                body = {"status": "ok"}
            elif path == "/v1/models":
                body = {"object": "list", "data": [{"id": MODEL, "object": "model"}]}
            else:
                assert path == "/v1/chat/completions"
                body = {
                    "id": "chatcmpl-held", "object": "chat.completion", "model": MODEL,
                    "choices": [{"index": 0,
                                 "message": {"role": "assistant", "content": "OK"},
                                 "finish_reason": "stop"}],
                    "usage": {"prompt_tokens": 1, "completion_tokens": 1,
                              "total_tokens": 2,
                              "prompt_tokens_details": {"cached_tokens": 0}},
                }
            payload = json.dumps(body).encode()
            writer.write(
                b"HTTP/1.1 200 OK\r\nContent-Type: application/json\r\n"
                b"Connection: close\r\n"
                + f"Content-Length: {len(payload)}\r\n\r\n".encode() + payload
            )
            await writer.drain()
        except (asyncio.IncompleteReadError, ConnectionError):
            pass
        finally:
            for wait in waits:
                wait.cancel()
            await asyncio.gather(*waits, return_exceptions=True)
            writer.close()
            try:
                await writer.wait_closed()
            finally:
                if item is not None:
                    item.closed.set()
                handlers.discard(task)

    server = await asyncio.start_server(handle, "127.0.0.1", 0)
    async with server:
        try:
            yield f"http://127.0.0.1:{server.sockets[0].getsockname()[1]}", pending
        finally:
            tasks = tuple(handlers)
            for task in tasks:
                task.cancel()
            await asyncio.gather(*tasks, return_exceptions=True)


@asynccontextmanager
async def adapter_at(url: str, *, deadline: float = 5) -> AsyncGenerator[
    TurboFieldfareAdapter, None,
]:
    # Disable HTTP read timeout: only the adapter's real total deadline can end
    # a held response. One connection makes a leaked pool slot observable.
    options = {
        "base_url": url, "trust_env": False,
        "timeout": httpx2.Timeout(None, pool=0.05),
        "limits": httpx2.Limits(max_connections=1, max_keepalive_connections=1),
    }
    async with (
        httpx2.AsyncClient(**options) as generation_client,
        httpx2.AsyncClient(**options) as probe_client,
    ):
        yield TurboFieldfareAdapter(
            generation_client, probe_client=probe_client, model=MODEL,
            probe_timeout=deadline, generation_timeout=deadline,
        )


async def invoke(adapter: TurboFieldfareAdapter, kind: str):
    if kind == "generation":
        return await adapter.generate(REQUEST)
    return await adapter.list_models()


async def assert_next_request_works(adapter, pending, kind: str) -> None:
    task = asyncio.create_task(invoke(adapter, kind))
    try:
        item = await asyncio.wait_for(pending.get(), timeout=2)
        item.release.set()
        result = await asyncio.wait_for(task, timeout=2)
        if kind == "generation":
            assert result.text == "OK"
        else:
            assert result[0].id == MODEL
    finally:
        task.cancel()
        await asyncio.gather(task, return_exceptions=True)


@pytest.mark.parametrize("kind", ["generation", "probe"])
@pytest.mark.parametrize("termination", ["deadline", "cancel"])
def test_real_deadline_or_cancellation_closes_connection_and_releases_slot(
    kind: str, termination: str,
) -> None:
    async def scenario() -> None:
        async with held_upstream() as (url, pending):
            async with adapter_at(url, deadline=0.1 if termination == "deadline" else 5) as adapter:
                task = asyncio.create_task(invoke(adapter, kind))
                try:
                    item = await asyncio.wait_for(pending.get(), timeout=2)
                    assert item.path == ("/v1/chat/completions" if kind == "generation" else "/v1/models")
                    if termination == "cancel":
                        task.cancel()
                        with pytest.raises(asyncio.CancelledError):
                            await asyncio.wait_for(task, timeout=2)
                    else:
                        with pytest.raises(InferenceError) as caught:
                            await asyncio.wait_for(task, timeout=2)
                        assert caught.value.code == "timeout"
                        assert caught.value.upstream_status is None
                    await asyncio.wait_for(item.closed.wait(), timeout=2)
                    assert not item.release.is_set()
                    # Same client, same one-slot pool, a new TCP connection.
                    await assert_next_request_works(adapter, pending, kind)
                finally:
                    task.cancel()
                    await asyncio.gather(task, return_exceptions=True)

    asyncio.run(scenario())


@pytest.mark.parametrize("kind", ["generation", "probe"])
def test_real_pool_timeout_is_busy_and_does_not_cancel_active_request(kind: str) -> None:
    async def scenario() -> None:
        async with held_upstream() as (url, pending):
            async with adapter_at(url) as adapter:
                first = asyncio.create_task(invoke(adapter, kind))
                try:
                    item = await asyncio.wait_for(pending.get(), timeout=2)
                    with pytest.raises(InferenceError) as caught:
                        await asyncio.wait_for(invoke(adapter, kind), timeout=2)
                    assert caught.value.code == "busy"
                    assert caught.value.upstream_status is None
                    assert not first.done()
                    assert not item.closed.is_set()
                    assert pending.empty()  # The rejected request never reached upstream.
                    item.release.set()
                    await asyncio.wait_for(first, timeout=2)
                    await assert_next_request_works(adapter, pending, kind)
                finally:
                    first.cancel()
                    await asyncio.gather(first, return_exceptions=True)

    asyncio.run(scenario())

실제 테스트 소스와 동일한 코드다. 실행 명령: cd be.venv/bin/python -m pytest -q tests/test_inference_timeouts.py.

8. 다음 파일과의 연결

lifespan 종료 검증에서 다른 경계의 보장을 확인한다. 전체 구조는 Backend 개요, client 조립은 main.py, 오류 번역은 Adapter 구현과 연결된다.

← 개요Backend 검증 목록다음 →lifespan 종료 검증