From bfb68392e8c6ea990d704585f418b79654cead3a Mon Sep 17 00:00:00 2001 From: badayvedat Date: Mon, 24 Aug 2026 00:33:44 +0700 Subject: [PATCH 01/27] fix: --- .github/workflows/fal-e2e-tests.yml | 2 +- projects/fal/tests/conftest.py | 2 +- projects/fal/tests/e2e/test_apps.py | 96 +++++++++---------- .../fal/tests/integration/test_stability.py | 4 +- 4 files changed, 48 insertions(+), 56 deletions(-) diff --git a/.github/workflows/fal-e2e-tests.yml b/.github/workflows/fal-e2e-tests.yml index 01d71dde9..b4931dc0e 100644 --- a/.github/workflows/fal-e2e-tests.yml +++ b/.github/workflows/fal-e2e-tests.yml @@ -74,4 +74,4 @@ jobs: FAL_GRPC_HOST: api.alpha.fal.ai FAL_REST_HOST: rest.fal.ai FAL_RUN_HOST: fal.run - run: pytest -n auto -v projects/fal/tests/e2e + run: pytest -n auto --dist loadgroup -v projects/fal/tests/e2e diff --git a/projects/fal/tests/conftest.py b/projects/fal/tests/conftest.py index 121d3d4bd..b4e6e76c9 100644 --- a/projects/fal/tests/conftest.py +++ b/projects/fal/tests/conftest.py @@ -20,7 +20,7 @@ def isolated_client(): return partial(function, machine_type="XS", keep_alive=0) -@pytest.fixture(scope="function") +@pytest.fixture(scope="module") def make_tmp_app_name() -> Callable[[str], str]: def _make_tmp_app_name(prefix: str = "test") -> str: short_id = uuid.uuid4().hex[:8] diff --git a/projects/fal/tests/e2e/test_apps.py b/projects/fal/tests/e2e/test_apps.py index a7ccf936c..414f32b94 100644 --- a/projects/fal/tests/e2e/test_apps.py +++ b/projects/fal/tests/e2e/test_apps.py @@ -541,25 +541,6 @@ async def cancel_handler(self) -> Output: return Output(result=0) -class HealthCheckApp(fal.App, keep_alive=300, max_concurrency=1, request_timeout=4): - @fal.endpoint("/") - def run(self, input: Input) -> Output: - time.sleep(input.wait_time) - return Output(result=input.lhs + input.rhs) - - @fal.endpoint( - "/health", - health_check=fal.HealthCheck( - start_period_seconds=10, - timeout_seconds=10, - failure_threshold=3, - call_regularly=True, - ), - ) - def health(self) -> Output: - return Output(result=0) - - class HealthOverrideApp(fal.App, keep_alive=300, max_concurrency=1): """fal.App declaring health at /ready with a broken default /health. @@ -711,7 +692,7 @@ def user(rest_client: Client) -> Generator[User, None, None]: yield user -@pytest.fixture() +@pytest.fixture(scope="module") def register_app( host: api.FalServerlessHost, make_tmp_app_name: Callable[[str], str], @@ -761,7 +742,7 @@ def base_app(register_app): yield app_alias, app_revision -@pytest.fixture() +@pytest.fixture(scope="module") def test_app( user: User, register_app, @@ -837,7 +818,7 @@ def test_fastapi_app( yield f"{user.username}/{app_alias}" -@pytest.fixture() +@pytest.fixture(scope="module") def test_stateful_app( user: User, register_app, @@ -847,12 +828,6 @@ def test_stateful_app( yield f"{user.username}/{app_alias}" -@pytest.fixture() -def test_pydantic_validation_error(): - with AppClient.connect(StatefulAdditionApp) as client: - yield client - - @pytest.fixture() def test_cancellable_app( user: User, @@ -863,22 +838,12 @@ def test_cancellable_app( yield f"{user.username}/{app_alias}" -@pytest.fixture() +@pytest.fixture(scope="module") def test_exception_app(): with AppClient.connect(ExceptionApp) as client: yield client -@pytest.fixture() -def test_health_check_app( - user: User, - register_app, -): - health_check_app = wrap_app(HealthCheckApp) - with register_app(health_check_app, "health-check") as (app_alias, _): - yield f"{user.username}/{app_alias}" - - @pytest.fixture() def test_sleep_app( user: User, @@ -899,7 +864,7 @@ def test_queue_blocking_app( yield f"{user.username}/{app_alias}" -@pytest.fixture() +@pytest.fixture(scope="module") def test_realtime_app( user: User, register_app, @@ -909,13 +874,14 @@ def test_realtime_app( yield f"{user.username}/{app_alias}" -def test_broken_app_failure(host: api.FalServerlessHost, user: User): +def test_broken_app_failure(): with pytest.raises(FalServerlessException) as e: wrap_app(BrokenApp) assert "Failed to generate OpenAPI" in str(e) +@pytest.mark.xdist_group(name="addition-app") def test_app_client(test_app: str): response = apps.run(test_app, arguments={"lhs": 1, "rhs": 2}) assert response["result"] == 3 @@ -1064,6 +1030,7 @@ def test_app_health_override(test_health_override_app: str): assert r.status_code == 502, r.text +@pytest.mark.xdist_group(name="addition-app") def test_ws_client(test_app: str): with apps.ws(test_app) as connection: for i in range(3): @@ -1079,6 +1046,7 @@ def test_ws_client(test_app: str): assert response["result"] == 2 + i +@pytest.mark.xdist_group(name="stateful-app") def test_app_client_path_included_in_app_id(test_stateful_app: str): response = apps.run(test_stateful_app + "/reset", arguments={}) assert response["result"] == 0 @@ -1091,6 +1059,7 @@ def test_app_client_path_included_in_app_id(test_stateful_app: str): assert response["result"] == 6 +@pytest.mark.xdist_group(name="stateful-app") def test_stateful_app_client(test_stateful_app: str): response = apps.run(test_stateful_app, arguments={}, path="/reset") assert response["result"] == 0 @@ -1108,6 +1077,7 @@ def test_stateful_app_client(test_stateful_app: str): assert response["result"] == 0 +@pytest.mark.xdist_group(name="addition-app") def test_app_cancellation(test_app: str, test_cancellable_app: str): request_handle = apps.submit( test_cancellable_app, arguments={"lhs": 1, "rhs": 2, "wait_time": 6} @@ -1147,7 +1117,7 @@ def test_app_cancellation(test_app: str, test_cancellable_app: str): assert response == {"result": 3} -def test_app_disconnect_behavior(test_app: str, test_cancellable_app: str): +def test_app_disconnect_behavior(test_cancellable_app: str): with pytest.raises(HTTPStatusError) as e: apps.run( test_cancellable_app, @@ -1278,6 +1248,7 @@ def test_app_client_async(test_sleep_app: str): @pytest.mark.xfail( reason="Temporary disabled while investigating backend issue. Ping @efiop" ) +@pytest.mark.xdist_group(name="exception-app") def test_traceback_logs(test_exception_app: AppClient, rest_client: Client): date = ( datetime.now(timezone.utc).replace(tzinfo=None) - timedelta(seconds=1) @@ -1313,11 +1284,9 @@ def test_traceback_logs(test_exception_app: AppClient, rest_client: Client): ), "Logs should contain the traceback message" -def test_app_openapi_spec_metadata( - base_app: Tuple[str, str], user: User, rest_client: Client -): - app_alias, _ = base_app - app_user_id = user.username +@pytest.mark.xdist_group(name="addition-app") +def test_app_openapi_spec_metadata(test_app: str, rest_client: Client): + app_user_id, _, app_alias = test_app.partition("/") res = app_metadata.sync_detailed( app_alias_or_id=app_alias, app_user_id=app_user_id, client=rest_client ) @@ -1350,11 +1319,13 @@ def test_app_no_serve_spec_metadata(test_fastapi_app: str, rest_client: Client): ), f"openapi should not be present in metadata {metadata}" -def test_404_response(test_app: str, request: pytest.FixtureRequest): +@pytest.mark.xdist_group(name="addition-app") +def test_404_response(test_app: str): with pytest.raises(HTTPStatusError, match="Path /.*other not found"): apps.run(test_app, path="/other", arguments={"lhs": 1, "rhs": 2}) +@pytest.mark.xdist_group(name="exception-app") def test_404_billable_units(test_exception_app: AppClient): """Test that 404 responses include x-fal-billable-units: 0 header.""" with httpx.Client(headers=_auth_headers()) as httpx_client: @@ -1511,6 +1482,7 @@ def test_app_set_delete_alias(base_app: Tuple[str, str]): assert not found, f"Found app {app_alias} in {res} after deletion" +@pytest.mark.xdist_group(name="realtime-app") def test_realtime_connection(test_realtime_app): response = apps.run(test_realtime_app, arguments={"prompt": "a cat"}) assert response["text"] == "a cat" @@ -1540,6 +1512,7 @@ def test_realtime_connection(test_realtime_app): assert batch_sizes == [4, 4, 2] +@pytest.mark.xdist_group(name="realtime-app") def test_realtime_ws_endpoint(test_realtime_app): app_id = apps._backwards_compatible_app_id(test_realtime_app) url = apps._REALTIME_URL_FORMAT.format(app_id=app_id) + "/ws" @@ -1558,6 +1531,7 @@ def test_realtime_ws_endpoint(test_realtime_app): assert messages == [{"message": "Hello world!"}] * 3 +@pytest.mark.xdist_group(name="realtime-app") def test_realtime_connection_custom_codec(test_realtime_app): with apps._connect( test_realtime_app, @@ -1569,6 +1543,7 @@ def test_realtime_connection_custom_codec(test_realtime_app): assert response["text"] == "json cat" +@pytest.mark.xdist_group(name="realtime-app") def test_realtime_server_streaming_mode(test_realtime_app): with apps._connect( test_realtime_app, path="/realtime/server-streaming" @@ -1582,6 +1557,7 @@ def test_realtime_server_streaming_mode(test_realtime_app): ] +@pytest.mark.xdist_group(name="realtime-app") def test_realtime_server_streaming_sync_mode(test_realtime_app): with apps._connect( test_realtime_app, path="/realtime/server-streaming-sync" @@ -1595,6 +1571,7 @@ def test_realtime_server_streaming_sync_mode(test_realtime_app): ] +@pytest.mark.xdist_group(name="realtime-app") def test_realtime_client_streaming_mode(test_realtime_app): with apps._connect( test_realtime_app, path="/realtime/client-streaming" @@ -1606,6 +1583,7 @@ def test_realtime_client_streaming_mode(test_realtime_app): assert response["texts"] == ["first", "second", "third"] +@pytest.mark.xdist_group(name="realtime-app") def test_realtime_bidi_mode(test_realtime_app): with apps._connect(test_realtime_app, path="/realtime/bidi") as connection: connection.send({"prompt": "one"}) @@ -1622,6 +1600,7 @@ def delete_workflow_on_exit(client: httpx.Client, workflow_url: str): client.delete(workflow_url) +@pytest.mark.xdist_group(name="addition-app") def test_workflows(test_app: str, rest_client: Client): workflow = Workflow( name="test_workflow_" + secrets.token_hex(), @@ -1671,6 +1650,7 @@ def test_workflows(test_app: str, rest_client: Client): assert data["result"] == 10 +@pytest.mark.xdist_group(name="exception-app") def test_app_exceptions(test_exception_app: AppClient): with pytest.raises(AppClientError) as app_exc: test_exception_app.app_exception({}) @@ -1703,9 +1683,12 @@ def test_app_exceptions(test_exception_app: AppClient): assert _CUDA_OOM_MESSAGE in cuda_exc.value.message -def test_pydantic_validation_billing(test_pydantic_validation_error: AppClient): - with httpx.Client(headers=_auth_headers()) as httpx_client: - url = test_pydantic_validation_error.url + "/increment" +@pytest.mark.xdist_group(name="stateful-app") +def test_pydantic_validation_billing(test_stateful_app: str): + from fal.flags import FAL_RUN_HOST + + with httpx.Client(headers=get_credentials().to_headers()) as httpx_client: + url = f"https://{FAL_RUN_HOST}/{test_stateful_app}/increment" response = httpx_client.post( url, json={"value": "this-is-not-an-integer"}, @@ -1716,6 +1699,7 @@ def test_pydantic_validation_billing(test_pydantic_validation_error: AppClient): assert response.headers.get("x-fal-billable-units") == "0" +@pytest.mark.xdist_group(name="exception-app") def test_field_exception_billing(test_exception_app: AppClient): with httpx.Client(headers=_auth_headers()) as httpx_client: url = test_exception_app.url + "/field-exception" @@ -1731,6 +1715,7 @@ def test_field_exception_billing(test_exception_app: AppClient): assert not hasattr(response.headers, "x-fal-billable-units") +@pytest.mark.xdist_group(name="exception-app") def test_field_exception_int_billable_units_formatting(test_exception_app: AppClient): """Test that int billable_units are formatted without decimal places.""" with httpx.Client(headers=_auth_headers()) as httpx_client: @@ -1745,6 +1730,7 @@ def test_field_exception_int_billable_units_formatting(test_exception_app: AppCl assert response.headers.get("x-fal-billable-units") == "42" +@pytest.mark.xdist_group(name="exception-app") def test_field_exception_float_billable_units_formatting(test_exception_app: AppClient): """Test that float billable_units are formatted with 8 decimal places.""" with httpx.Client(headers=_auth_headers()) as httpx_client: @@ -1759,6 +1745,7 @@ def test_field_exception_float_billable_units_formatting(test_exception_app: App assert response.headers.get("x-fal-billable-units") == "3.14159265" +@pytest.mark.xdist_group(name="exception-app") def test_field_exception_scientific_notation_small(test_exception_app: AppClient): """Test that small scientific notation values are properly formatted.""" with httpx.Client(headers=_auth_headers()) as httpx_client: @@ -1774,6 +1761,7 @@ def test_field_exception_scientific_notation_small(test_exception_app: AppClient assert response.headers.get("x-fal-billable-units") == "0.00001230" +@pytest.mark.xdist_group(name="exception-app") def test_field_exception_scientific_notation_large(test_exception_app: AppClient): """Test that large scientific notation values are properly formatted.""" with httpx.Client(headers=_auth_headers()) as httpx_client: @@ -1789,6 +1777,7 @@ def test_field_exception_scientific_notation_large(test_exception_app: AppClient assert response.headers.get("x-fal-billable-units") == "12300000000.00000000" +@pytest.mark.xdist_group(name="exception-app") def test_field_exception_invalid_billable_units(test_exception_app: AppClient): """Test that invalid billable_units (non-numeric string) raises an error.""" with httpx.Client(headers=_auth_headers()) as httpx_client: @@ -1804,6 +1793,7 @@ def test_field_exception_invalid_billable_units(test_exception_app: AppClient): assert response.status_code == 500 +@pytest.mark.xdist_group(name="exception-app") def test_field_exception_default_billable_units(test_exception_app: AppClient): """Test that when billable_units is not set (None), no header is included.""" with httpx.Client(headers=_auth_headers()) as httpx_client: @@ -2211,7 +2201,7 @@ def get_context(self, request: Request) -> RequestContextOutput: ) -@pytest.fixture() +@pytest.fixture(scope="module") def test_request_context_app( user: User, register_app, @@ -2222,6 +2212,7 @@ def test_request_context_app( yield f"{user.username}/{app_alias}" +@pytest.mark.xdist_group(name="request-context-app") def test_request_context_fields_populated(test_request_context_app: str): """Test that request context fields are properly populated.""" @@ -2234,6 +2225,7 @@ def test_request_context_fields_populated(test_request_context_app: str): assert result["request_id_from_context"] == result["request_id_from_header"] +@pytest.mark.xdist_group(name="request-context-app") def test_request_context_isolation_with_multiplexing(test_request_context_app: str): """Test that request context is properly isolated between concurrent requests. diff --git a/projects/fal/tests/integration/test_stability.py b/projects/fal/tests/integration/test_stability.py index f14bee918..f4b92eefa 100644 --- a/projects/fal/tests/integration/test_stability.py +++ b/projects/fal/tests/integration/test_stability.py @@ -418,7 +418,7 @@ def test_big_message(isolated_client): # try doubling that. data_length = 8 * (1024**2) - @isolated_client("virtualenv", machine_type="M") + @isolated_client("virtualenv", machine_type="M", keep_alive=10) def big_return_function(data_length): return b"0" * data_length @@ -429,7 +429,7 @@ def big_return_function(data_length): def test_futures(isolated_client): from concurrent.futures import wait - @isolated_client("virtualenv") + @isolated_client("virtualenv", keep_alive=10) def regular_function(n): return n * 2 From 1a0672e6545c1720ba54015c9954cf2deb902db2 Mon Sep 17 00:00:00 2001 From: badayvedat Date: Mon, 24 Aug 2026 02:11:19 +0700 Subject: [PATCH 02/27] test: harden timing-sensitive e2e and integration tests --- .github/workflows/fal-e2e-tests.yml | 3 +- .github/workflows/fal-integration-tests.yml | 3 +- projects/fal/tests/e2e/test_apps.py | 616 ++++++++++-------- .../tests/integration/test_environments.py | 11 +- .../fal/tests/integration/test_stability.py | 28 +- projects/fal/tests/unit/test_app.py | 41 ++ 6 files changed, 414 insertions(+), 288 deletions(-) diff --git a/.github/workflows/fal-e2e-tests.yml b/.github/workflows/fal-e2e-tests.yml index b4931dc0e..60a2538df 100644 --- a/.github/workflows/fal-e2e-tests.yml +++ b/.github/workflows/fal-e2e-tests.yml @@ -34,6 +34,7 @@ jobs: timeout-minutes: 10 strategy: fail-fast: false + max-parallel: 2 matrix: os: [ubuntu-latest] deps: ["pydantic==1.10.18", "pydantic==2.13.3"] @@ -74,4 +75,4 @@ jobs: FAL_GRPC_HOST: api.alpha.fal.ai FAL_REST_HOST: rest.fal.ai FAL_RUN_HOST: fal.run - run: pytest -n auto --dist loadgroup -v projects/fal/tests/e2e + run: pytest -n auto --dist loadgroup --timeout=120 -v projects/fal/tests/e2e diff --git a/.github/workflows/fal-integration-tests.yml b/.github/workflows/fal-integration-tests.yml index a7a12b477..554fd6547 100644 --- a/.github/workflows/fal-integration-tests.yml +++ b/.github/workflows/fal-integration-tests.yml @@ -32,6 +32,7 @@ jobs: timeout-minutes: 10 strategy: fail-fast: false + max-parallel: 4 matrix: os: [ubuntu-latest] deps: ["pydantic==1.10.18", "pydantic==2.13.3"] @@ -71,4 +72,4 @@ jobs: FAL_GRPC_HOST: api.alpha.fal.ai FAL_REST_HOST: rest.fal.ai FAL_RUN_HOST: fal.run - run: pytest -n auto -v projects/fal/tests/integration + run: pytest -n auto --dist loadgroup --timeout=120 -v projects/fal/tests/integration diff --git a/projects/fal/tests/e2e/test_apps.py b/projects/fal/tests/e2e/test_apps.py index 414f32b94..5d5418d90 100644 --- a/projects/fal/tests/e2e/test_apps.py +++ b/projects/fal/tests/e2e/test_apps.py @@ -3,6 +3,7 @@ import os import secrets import subprocess +import sys import time from contextlib import contextmanager, suppress from datetime import datetime, timedelta, timezone @@ -16,6 +17,7 @@ List, Optional, Tuple, + TypeVar, Union, ) @@ -35,7 +37,6 @@ from fal import apps from fal.api.deploy import User, _get_user from fal.app import AppClient, AppClientError, wrap_app -from fal.auth import key_credentials from fal.container import ContainerImage from fal.exceptions import ( AppException, @@ -77,25 +78,103 @@ class Output(BaseModel): result: int +class FailInput(BaseModel): + marker: str + + actual_python = active_python() +T = TypeVar("T") def _auth_headers() -> Dict[str, str]: - key_creds = key_credentials() - if not key_creds: - return {} - key_id, key_secret = key_creds - return {"Authorization": f"Key {key_id}:{key_secret}"} - - -def git_revision_short_hash() -> str: - return ( - subprocess.check_output(["git", "rev-parse", "--short", "HEAD"]) - .decode("ascii") - .strip() + return get_credentials().to_headers() + + +def _wait_until( + fetch: Callable[[], T], + predicate: Callable[[T], bool], + *, + timeout: float, + description: str, + interval: float = 0.1, +) -> T: + deadline = time.monotonic() + timeout + + while True: + value = fetch() + if predicate(value): + return value + + remaining = deadline - time.monotonic() + if remaining <= 0: + raise AssertionError(f"Timed out waiting for {description}: {value!r}") + time.sleep(min(interval, remaining)) + + +def _wait_for_request_status( + handle, + expected_status, + *, + timeout: float = 60, + logs: bool = False, +): + def fetch_status(): + status = handle.status(logs=logs) + if isinstance(status, apps.Completed) and not isinstance( + status, expected_status + ): + raise AssertionError( + f"Request completed before reaching {expected_status}: {status!r}" + ) + return status + + return _wait_until( + fetch_status, + lambda status: isinstance(status, expected_status), + timeout=timeout, + description=f"request status {expected_status}", ) +def _cancel_and_wait(handle, *, timeout: float = 30): + status = handle.status() + if isinstance(status, apps.Completed): + return status + + try: + handle.cancel() + except HTTPStatusError: + status = handle.status() + if isinstance(status, apps.Completed): + return status + raise + + return _wait_for_request_status(handle, apps.Completed, timeout=timeout) + + +def _wait_for_alias_runners( + client, + app_alias: str, + predicate, + *, + timeout: float = 45, +): + return _wait_until( + lambda: client.list_alias_runners(app_alias), + predicate, + timeout=timeout, + interval=0.5, + description=f"runner state for {app_alias}", + ) + + +GIT_REVISION_SHORT_HASH = ( + subprocess.check_output(["git", "rev-parse", "--short", "HEAD"]) + .decode("ascii") + .strip() +) + + @fal.function( keep_alive=60, machine_type=["S", "M"], @@ -115,7 +194,7 @@ def addition_app(input: Input) -> Output: @fal.function( kind="container", image=ContainerImage.from_dockerfile_str( - f"FROM python:{actual_python}-slim\n# {git_revision_short_hash()}", + f"FROM python:{actual_python}-slim\n# {GIT_REVISION_SHORT_HASH}", ), keep_alive=60, machine_type="S", @@ -134,7 +213,7 @@ def container_addition_app(input: Input) -> Output: @fal.function( kind="container", image=ContainerImage.from_dockerfile_str( - f"FROM python:{actual_python}-slim\n# {git_revision_short_hash()}", + f"FROM python:{actual_python}-slim\n# {GIT_REVISION_SHORT_HASH}", ), keep_alive=60, machine_type="S", @@ -154,7 +233,7 @@ def container_cache_enabled_app(input: Input) -> Output: @fal.function( kind="container", image=ContainerImage.from_dockerfile_str( - f"""FROM python:{actual_python}-slim\n# {git_revision_short_hash()} + f"""FROM python:{actual_python}-slim\n# {GIT_REVISION_SHORT_HASH} ARG OUTPUT="built incorrectly" ENV OUTPUT="$OUTPUT" """, @@ -458,8 +537,8 @@ class ExceptionApp(fal.App, keep_alive=300, max_concurrency=1): machine_type = "XS" @fal.endpoint("/fail") - def fail(self) -> Output: - raise Exception("this app is designed to fail!") + def fail(self, input: FailInput) -> Output: + raise Exception(f"this app is designed to fail! {input.marker}") @fal.endpoint("/app-exception") def app_exception(self) -> Output: @@ -594,7 +673,7 @@ def can_batch( other: "RTInput", current_batch_size: int = 1, ) -> bool: - return "don't batch" not in other.prompt + return "don't batch" not in self.prompt and "don't batch" not in other.prompt class RTOutput(BaseModel): @@ -643,7 +722,7 @@ def generate_rt_server_streaming_sync(self, input: RTInput) -> Iterator[RTOutput for idx in range(3): yield RTOutput(text=f"{input.prompt}:{idx}") - @fal.realtime("/realtime/client-streaming", session_timeout=0.2) + @fal.realtime("/realtime/client-streaming", session_timeout=1) async def generate_rt_client_streaming( self, inputs: AsyncIterator[RTInput] ) -> RTOutputs: @@ -669,7 +748,6 @@ def generate_rt_json(self, input: RTInput) -> RTOutput: @fal.realtime("/realtime/batched", buffering=10, max_batch_size=4) def generate_rt_batched(self, input: RTInput, *inputs: RTInput) -> RTOutputs: - time.sleep(2) # fixed cost return RTOutputs(texts=[input.prompt] + [i.prompt for i in inputs]) @@ -1083,12 +1161,7 @@ def test_app_cancellation(test_app: str, test_cancellable_app: str): test_cancellable_app, arguments={"lhs": 1, "rhs": 2, "wait_time": 6} ) - while True: - status = request_handle.status() - time.sleep(0.05) - if isinstance(status, apps.InProgress): - # The app is running - break + _wait_for_request_status(request_handle, apps.InProgress) # cancel the request request_handle.cancel() @@ -1103,12 +1176,7 @@ def test_app_cancellation(test_app: str, test_cancellable_app: str): test_app, arguments={"lhs": 1, "rhs": 2, "wait_time": 6} ) - while True: - status = request_handle.status() - time.sleep(0.05) - if isinstance(status, apps.InProgress): - # The app is running - break + _wait_for_request_status(request_handle, apps.InProgress) # cancel the request request_handle.cancel() @@ -1121,7 +1189,7 @@ def test_app_disconnect_behavior(test_cancellable_app: str): with pytest.raises(HTTPStatusError) as e: apps.run( test_cancellable_app, - arguments={"lhs": 1, "rhs": 2, "wait_time": 6}, + arguments={"lhs": 1, "rhs": 2, "wait_time": 30}, path="/well-handled", ) assert ( @@ -1141,50 +1209,41 @@ def test_app_disconnect_behavior(test_cancellable_app: str): with pytest.raises(HTTPStatusError) as e: apps.run( test_cancellable_app, - arguments={"lhs": 1, "rhs": 2, "wait_time": 6}, + arguments={"lhs": 1, "rhs": 2, "wait_time": 30}, ) assert ( e.value.response.status_code == 504 ), "Expected Gateway Timeout even though the app handled it" +@pytest.mark.timeout(120) def test_start_timeout_queue_blocking(test_queue_blocking_app: str): """ Test that start_timeout correctly times out a request waiting in queue. Scenario: - 1. Send a 10-second sleep request (occupies the only slot) - 2. While it's processing, send a second request with start_timeout=5 + 1. Send a long-running request (occupies the only slot) + 2. While it's processing, send a second request with start_timeout=2 3. The second request should return 504 because it times out waiting in queue (before processing starts) - 4. First request should complete successfully + 4. Cancel the blocking request during cleanup """ import fal_client from fal_client.client import FalClientHTTPError - # Send a long-running request that will occupy the only slot - # (max_concurrency=1, max_multiplexing=1) - # Use 10 seconds to ensure it blocks long enough for the second request to timeout - first_handle = apps.submit(test_queue_blocking_app, arguments={"wait_time": 10}) + first_handle = apps.submit(test_queue_blocking_app, arguments={"wait_time": 15}) - # Wait for the first request to start processing - while True: - status = first_handle.status() - if isinstance(status, apps.InProgress): - break - elif isinstance(status, apps.Queued): - time.sleep(0.1) - else: - raise Exception(f"Unexpected status for first request: {status}") - - # Now send a second request with a short start_timeout - # This should fail because it will timeout waiting in the queue - with pytest.raises(FalClientHTTPError) as exc_info: - fal_client.subscribe( - test_queue_blocking_app, - arguments={"wait_time": 1}, - start_timeout=5, - ) + try: + _wait_for_request_status(first_handle, apps.InProgress, timeout=60) + + with pytest.raises(FalClientHTTPError) as exc_info: + fal_client.subscribe( + test_queue_blocking_app, + arguments={"wait_time": 1}, + start_timeout=2, + ) + finally: + _cancel_and_wait(first_handle) # Should get a 504 timeout error assert ( @@ -1195,18 +1254,12 @@ def test_start_timeout_queue_blocking(test_queue_blocking_app: str): timeout_type = exc_info.value.response_headers.get("x-fal-request-timeout-type") assert timeout_type == "user", f"Expected 'user' timeout type, got {timeout_type}" - # First request should complete successfully - result = first_handle.get() - assert result == {"slept": True}, f"First request should succeed, got {result}" - -@pytest.mark.xfail( - reason="Temporary disabled while investigating backend issue. Ping @efiop" -) +@pytest.mark.timeout(120) def test_app_client_async(test_sleep_app: str): - handle = apps.submit(test_sleep_app, arguments={"wait_time": 1}) + handle = apps.submit(test_sleep_app, arguments={"wait_time": 10}) + _wait_for_request_status(handle, apps.InProgress) with pytest.raises(HTTPStatusError) as e: - # Not yet completed handle.fetch_result() assert e.value.response.status_code == 400 @@ -1225,11 +1278,12 @@ def test_app_client_async(test_sleep_app: str): elif isinstance(event, apps.Queued): assert event.position == 0 - for _ in range(10): - status = handle.status(logs=True) - assert isinstance(status, apps.Completed) - if status.logs: - break + status = _wait_until( + lambda: handle.status(logs=True), + lambda current: isinstance(current, apps.Completed) and bool(current.logs), + timeout=30, + description="completed request logs", + ) assert status.logs, "Logs missing from Completed status" assert any("sleeping..." in log["message"] for log in status.logs) @@ -1244,44 +1298,45 @@ def test_app_client_async(test_sleep_app: str): assert result == {"slept": True} -# If the logging subsystem is not working for some nodes, this test will flake @pytest.mark.xfail( reason="Temporary disabled while investigating backend issue. Ping @efiop" ) @pytest.mark.xdist_group(name="exception-app") def test_traceback_logs(test_exception_app: AppClient, rest_client: Client): + marker = f"traceback-{secrets.token_hex(8)}" date = ( - datetime.now(timezone.utc).replace(tzinfo=None) - timedelta(seconds=1) + datetime.now(timezone.utc).replace(tzinfo=None) - timedelta(minutes=5) ).isoformat() with pytest.raises(AppClientError): - test_exception_app.fail({}) + test_exception_app.fail({"marker": marker}) with httpx.Client( base_url=rest_client.base_url, headers=rest_client.get_headers(), timeout=300, ) as client: - # Give some time for logs to propagate through the logging subsystem. - for _ in range(10): - time.sleep(2) + + def fetch_matching_logs(): response = client.get( rest_client.base_url + f"/logs/?traceback=true&since={date}" ) + response.raise_for_status() + return [log for log in response.json() if marker in log["message"]] + + logs = _wait_until( + fetch_matching_logs, + bool, + timeout=45, + interval=1, + description="traceback log propagation", + ) - logs = response.json() - if len(logs) > 0: - break - - assert len(logs) > 0 for log in logs: assert log["message"].count("\n") > 1, "Logs should be multi-line" assert ( '{"traceback":' not in log["message"] ), "Logs should not be JSON-wrapped" - assert ( - "this app is designed to fail" in log["message"] - ), "Logs should contain the traceback message" @pytest.mark.xdist_group(name="addition-app") @@ -1484,6 +1539,11 @@ def test_app_set_delete_alias(base_app: Tuple[str, str]): @pytest.mark.xdist_group(name="realtime-app") def test_realtime_connection(test_realtime_app): + isolated_input = RTInput(prompt="don't batch") + batchable_input = RTInput(prompt="batchable") + assert not isolated_input.can_batch(batchable_input) + assert not batchable_input.can_batch(isolated_input) + response = apps.run(test_realtime_app, arguments={"prompt": "a cat"}) assert response["text"] == "a cat" @@ -1493,14 +1553,12 @@ def test_realtime_connection(test_realtime_app): assert response["text"] == "a cat" with apps._connect(test_realtime_app, path="/realtime/batched") as connection: - connection.send({"prompt": "keep busy"}) - time.sleep(1) + connection.send({"prompt": "don't batch"}) + assert connection.recv()["texts"] == ["don't batch"] for prompt in range(10): connection.send({"prompt": str(prompt)}) - assert connection.recv()["texts"] == ["keep busy"] - received_prompts = set() batch_sizes = [] while len(received_prompts) < 10: @@ -1508,8 +1566,9 @@ def test_realtime_connection(test_realtime_app): received_prompts.update(response["texts"]) batch_sizes.append(len(response["texts"])) - assert len(received_prompts) == 10 - assert batch_sizes == [4, 4, 2] + assert received_prompts == {str(prompt) for prompt in range(10)} + assert sum(batch_sizes) == 10 + assert all(1 <= batch_size <= 4 for batch_size in batch_sizes) @pytest.mark.xdist_group(name="realtime-app") @@ -1687,7 +1746,7 @@ def test_app_exceptions(test_exception_app: AppClient): def test_pydantic_validation_billing(test_stateful_app: str): from fal.flags import FAL_RUN_HOST - with httpx.Client(headers=get_credentials().to_headers()) as httpx_client: + with httpx.Client(headers=_auth_headers()) as httpx_client: url = f"https://{FAL_RUN_HOST}/{test_stateful_app}/increment" response = httpx_client.post( url, @@ -1809,178 +1868,167 @@ def test_field_exception_default_billable_units(test_exception_app: AppClient): assert "x-fal-billable-units" not in response.headers -def submit_and_wait_for_runner(app: str, arguments: dict = {}, *, path: str = ""): - handle = apps.submit(app, arguments=arguments, path=path) +def _active_runners(runners): + active_states = {RunnerState.RUNNING, RunnerState.IDLE} + return [runner for runner in runners if runner.state in active_states] - while True: - status = handle.status() - if isinstance(status, apps.InProgress) or isinstance(status, apps.Completed): - break - elif isinstance(status, apps.Queued): - time.sleep(0.1) - else: - raise Exception(f"Failed to start the app: {status}") +def submit_and_wait_for_runner( + app: str, arguments: Optional[dict] = None, *, path: str = "" +): + handle = apps.submit(app, arguments=arguments or {}, path=path) + status = _wait_for_request_status( + handle, + (apps.InProgress, apps.Completed), + ) + if isinstance(status, apps.Completed): + handle.fetch_result() return handle +@pytest.mark.timeout(180) def test_stop_runner(host: api.FalServerlessHost, test_sleep_app: str): - # Submit a runner and wait for it to be idle + _, _, app_alias = test_sleep_app.partition("/") submit_and_wait_for_runner(test_sleep_app, arguments={"wait_time": 1}) - original_runner_id = None with host._connection as client: - timeout = 30 - start_time = time.time() - while True: - _, _, app_alias = test_sleep_app.partition("/") - runners = client.list_alias_runners(app_alias) - assert len(runners) == 1 - - if runners[0].in_flight_requests == 0: - original_runner_id = runners[0].runner_id - break - elif time.time() - start_time > timeout: - raise Exception(f"Timeout waiting for runner to be idle: {runners[0]}") - time.sleep(1) - - # Because the runner is not requested to be stopped, it should be reused - submit_and_wait_for_runner(test_sleep_app, arguments={"wait_time": 1}) + runners = _wait_for_alias_runners( + client, + app_alias, + lambda current: len(_active_runners(current)) == 1 + and _active_runners(current)[0].in_flight_requests == 0, + ) + original_runner_id = _active_runners(runners)[0].runner_id - with host._connection as client: - _, _, app_alias = test_sleep_app.partition("/") - runners = client.list_alias_runners(app_alias) - assert len(runners) == 1 + reuse_handle = apps.submit(test_sleep_app, arguments={"wait_time": 15}) + try: + _wait_for_request_status(reuse_handle, apps.InProgress) + _wait_for_alias_runners( + client, + app_alias, + lambda current: len(_active_runners(current)) == 1 + and _active_runners(current)[0].runner_id == original_runner_id + and _active_runners(current)[0].in_flight_requests > 0, + ) + finally: + _cancel_and_wait(reuse_handle) + + _wait_for_alias_runners( + client, + app_alias, + lambda current: len(_active_runners(current)) == 1 + and _active_runners(current)[0].runner_id == original_runner_id + and _active_runners(current)[0].in_flight_requests == 0, + ) - # Request to stop the runner - with host._connection as client: - with pytest.raises(Exception) as e: + with pytest.raises(Exception) as exc_info: client.stop_runner("1234567890") + assert "not found" in str(exc_info.value).lower() + + client.stop_runner(original_runner_id) + _wait_for_alias_runners( + client, + app_alias, + lambda current: all( + runner.runner_id != original_runner_id + for runner in _active_runners(current) + ), + ) - assert "not found" in str(e).lower() - - _, _, app_alias = test_sleep_app.partition("/") - runners = client.list_alias_runners(app_alias) - assert len(runners) == 1 - - client.stop_runner(runners[0].runner_id) - - # Because the runner is requested to be stopped, - # it should not be reused and a new runner should be created - submit_and_wait_for_runner(test_sleep_app, arguments={"wait_time": 1}) - - with host._connection as client: - runners = client.list_alias_runners(app_alias) - assert original_runner_id is not None - assert any(runner.runner_id != original_runner_id for runner in runners) + submit_and_wait_for_runner(test_sleep_app, arguments={"wait_time": 1}) + _wait_for_alias_runners( + client, + app_alias, + lambda current: any( + runner.runner_id != original_runner_id + for runner in _active_runners(current) + ), + ) +@pytest.mark.timeout(180) def test_kill_runner(host: api.FalServerlessHost, test_sleep_app: str): - # Kill all the replicas of the app that is already running - handle = apps.submit(test_sleep_app, arguments={"wait_time": 10}) - - while True: - status = handle.status() - if isinstance(status, apps.InProgress): - break - elif isinstance(status, apps.Queued): - time.sleep(1) - else: - raise Exception(f"Failed to start the app: {status}") + handle = apps.submit(test_sleep_app, arguments={"wait_time": 30}) + _wait_for_request_status(handle, apps.InProgress, timeout=60) with host._connection as client: - with pytest.raises(Exception) as e: + with pytest.raises(Exception) as exc_info: client.kill_runner("1234567890") - - assert "not found" in str(e).lower() + assert "not found" in str(exc_info.value).lower() _, _, app_alias = test_sleep_app.partition("/") - runners = client.list_alias_runners(app_alias) - existing_runners = len( - [runner for runner in runners if runner.state == RunnerState.RUNNING] + runners = _wait_for_alias_runners( + client, + app_alias, + lambda current: bool(_active_runners(current)), ) - runner_id = runners[0].runner_id + runner_id = _active_runners(runners)[0].runner_id client.kill_runner(runner_id) - - timeout = 15 - start_time = time.time() - while True: - runners = client.list_alias_runners(app_alias) - running_runner_ids = { - runner.runner_id - for runner in runners - if runner.state == RunnerState.RUNNING - } - if runner_id not in running_runner_ids: - break - if time.time() - start_time > timeout: - raise AssertionError( - f"Runner {runner_id} still running after kill request: {runners}" - ) - time.sleep(0.5) - - num_runners = len(running_runner_ids) - assert num_runners <= existing_runners - 1 + _wait_for_alias_runners( + client, + app_alias, + lambda current: all( + runner.runner_id != runner_id for runner in _active_runners(current) + ), + ) +@pytest.mark.timeout(180) def test_rollout_application(host: api.FalServerlessHost, test_sleep_app: str): handle = apps.submit(test_sleep_app, arguments={"wait_time": 30}) - - while True: - status = handle.status() - if isinstance(status, apps.InProgress): - break - elif isinstance(status, apps.Queued): - time.sleep(1) - else: - raise Exception(f"Failed to start the app: {status}") + _wait_for_request_status(handle, apps.InProgress, timeout=60) with host._connection as client: _, _, app_alias = test_sleep_app.partition("/") - runners_before = client.list_alias_runners(app_alias) - assert len(runners_before) == 1 - runner_id_before = runners_before[0].runner_id + runners_before = _wait_for_alias_runners( + client, + app_alias, + lambda current: len(_active_runners(current)) == 1, + ) + runner_id_before = _active_runners(runners_before)[0].runner_id client.rollout_application(app_alias, force=True) - - time.sleep(15) - - runners_after = client.list_alias_runners(app_alias) - runner_ids_after = {r.runner_id for r in runners_after} - - assert runner_id_before not in runner_ids_after + runners_after = _wait_for_alias_runners( + client, + app_alias, + lambda current: all( + runner.runner_id != runner_id_before + for runner in _active_runners(current) + ), + timeout=60, + ) + runner_ids_after = { + runner.runner_id for runner in _active_runners(runners_after) + } client.rollout_application(app_alias, force=True) - - time.sleep(3) - - runners_final = client.list_alias_runners(app_alias) - runner_ids_final = {r.runner_id for r in runners_final} - - assert not runner_ids_after.intersection(runner_ids_final) + _wait_for_alias_runners( + client, + app_alias, + lambda current: not runner_ids_after.intersection( + runner.runner_id for runner in _active_runners(current) + ), + timeout=60, + ) +@pytest.mark.timeout(180) def test_shell_runner(host: api.FalServerlessHost, test_sleep_app: str): - handle = apps.submit(test_sleep_app, arguments={"wait_time": 30}) - - while True: - status = handle.status() - if isinstance(status, apps.InProgress): - break - elif isinstance(status, apps.Queued): - time.sleep(1) - else: - raise Exception(f"Failed to start the app: {status}") + handle = submit_and_wait_for_runner(test_sleep_app, arguments={"wait_time": 1}) + assert handle.get() == {"slept": True} with host._connection as client: _, _, app_alias = test_sleep_app.partition("/") - runners = client.list_alias_runners(app_alias) - assert len(runners) == 1 - runner_id = runners[0].runner_id + runners = _wait_for_alias_runners( + client, + app_alias, + lambda current: bool(_active_runners(current)), + ) + runner_id = _active_runners(runners)[0].runner_id proc = subprocess.Popen( - ["python", "-m", "fal", "runners", "shell", runner_id], + [sys.executable, "-m", "fal", "runners", "shell", runner_id], stdin=subprocess.PIPE, stdout=subprocess.PIPE, stderr=subprocess.PIPE, @@ -1988,7 +2036,7 @@ def test_shell_runner(host: api.FalServerlessHost, test_sleep_app: str): try: commands = b"echo 'a' > t.txt\ncat t.txt\nexit\n" - stdout, stderr = proc.communicate(input=commands, timeout=10) + stdout, stderr = proc.communicate(input=commands, timeout=30) assert b"a" in stdout, f"Expected 'a' in output, got: {stdout.decode()}" finally: if proc.poll() is None: @@ -1996,27 +2044,23 @@ def test_shell_runner(host: api.FalServerlessHost, test_sleep_app: str): proc.wait() +@pytest.mark.timeout(180) def test_exec_runner(host: api.FalServerlessHost, test_sleep_app: str): - handle = apps.submit(test_sleep_app, arguments={"wait_time": 30}) - - while True: - status = handle.status() - if isinstance(status, apps.InProgress): - break - elif isinstance(status, apps.Queued): - time.sleep(1) - else: - raise Exception(f"Failed to start the app: {status}") + handle = submit_and_wait_for_runner(test_sleep_app, arguments={"wait_time": 1}) + assert handle.get() == {"slept": True} with host._connection as client: _, _, app_alias = test_sleep_app.partition("/") - runners = client.list_alias_runners(app_alias) - assert len(runners) == 1 - runner_id = runners[0].runner_id + runners = _wait_for_alias_runners( + client, + app_alias, + lambda current: bool(_active_runners(current)), + ) + runner_id = _active_runners(runners)[0].runner_id proc = subprocess.Popen( [ - "python", + sys.executable, "-m", "fal", "runners", @@ -2032,7 +2076,7 @@ def test_exec_runner(host: api.FalServerlessHost, test_sleep_app: str): ) try: - stdout, stderr = proc.communicate(timeout=10) + stdout, stderr = proc.communicate(timeout=30) assert ( b"hello" in stdout ), f"Expected 'hello' in output, got: {stdout.decode()}" @@ -2090,15 +2134,29 @@ class AppRefOutput(BaseModel): from_external_method: str -class AppRefApp(fal.App, keep_alive=300, max_concurrency=1, max_multiplexing=3): +CONCURRENT_REQUESTS = 3 + + +class AppRefApp( + fal.App, + keep_alive=300, + max_concurrency=1, + max_multiplexing=CONCURRENT_REQUESTS, +): machine_type = "XS" + async def setup(self): + self.concurrent_requests = 0 + self.requests_ready = asyncio.Event() + @fal.endpoint("/") - def run(self, request: Request) -> AppRefOutput: + async def run(self, request: Request) -> AppRefOutput: request_id = request.headers.get("x-request-id", "") - # sleep to intentionally cause a race condition - time.sleep(3) + self.concurrent_requests += 1 + if self.concurrent_requests == CONCURRENT_REQUESTS: + self.requests_ready.set() + await asyncio.wait_for(self.requests_ready.wait(), timeout=30) return AppRefOutput( from_app=request_id, @@ -2117,14 +2175,9 @@ def test_app_ref_app( def test_app_ref_app_client(test_app_ref_app: str): - time.sleep(3) - handle_1 = apps.submit(test_app_ref_app, arguments={}) - time.sleep(1) handle_2 = apps.submit(test_app_ref_app, arguments={}) - time.sleep(1) handle_3 = apps.submit(test_app_ref_app, arguments={}) - time.sleep(1) result_1 = handle_1.get() result_2 = handle_2.get() @@ -2135,23 +2188,29 @@ def test_app_ref_app_client(test_app_ref_app: str): assert result_3["from_app"] == result_3["from_external_method"] +@pytest.mark.timeout(180) def test_runner_machine_type(host: api.FalServerlessHost, test_sleep_app: str): """Test that machine_type is populated in runner info.""" + search_start = datetime.now() - timedelta(minutes=5) submit_and_wait_for_runner(test_sleep_app, arguments={"wait_time": 1}) with host._connection as client: _, _, app_alias = test_sleep_app.partition("/") - # list_alias_runners - runners = client.list_alias_runners(app_alias) - assert len(runners) >= 1 - assert runners[0].machine_type == "XS" - - # list_runners - all_runners = client.list_runners( - start_time=datetime.now() - timedelta(seconds=60) + runners = _wait_for_alias_runners( + client, + app_alias, + lambda current: any(runner.machine_type == "XS" for runner in current), + ) + assert any(runner.machine_type == "XS" for runner in runners) + + all_runners = _wait_until( + lambda: client.list_runners(start_time=search_start), + lambda current: any(runner.alias == app_alias for runner in current), + timeout=45, + interval=0.5, + description=f"runner history for {app_alias}", ) - assert len(all_runners) >= 1 target_runner = next((r for r in all_runners if r.alias == app_alias), None) assert target_runner is not None, "Runner for test app alias not found" assert target_runner.machine_type == "XS" @@ -2164,6 +2223,10 @@ class RequestContextOutput(BaseModel): request_id_from_header: Optional[str] +class RequestContextInput(BaseModel): + synchronize: bool = False + + def _external_get_request_context() -> dict: """External function that accesses request context without request parameter.""" current_app = get_current_app() @@ -2180,17 +2243,30 @@ def _external_get_request_context() -> dict: } -class RequestContextApp(fal.App, keep_alive=300, max_concurrency=1, max_multiplexing=3): +class RequestContextApp( + fal.App, + keep_alive=300, + max_concurrency=1, + max_multiplexing=CONCURRENT_REQUESTS, +): """App to test request context fields are properly populated.""" machine_type = "XS" + async def setup(self): + self.concurrent_requests = 0 + self.requests_ready = asyncio.Event() + @fal.endpoint("/") - def get_context(self, request: Request) -> RequestContextOutput: - # Sleep to intentionally cause potential race conditions with multiplexing - time.sleep(2) + async def get_context( + self, input: RequestContextInput, request: Request + ) -> RequestContextOutput: + if input.synchronize: + self.concurrent_requests += 1 + if self.concurrent_requests == CONCURRENT_REQUESTS: + self.requests_ready.set() + await asyncio.wait_for(self.requests_ready.wait(), timeout=30) - # Get context from external function (simulates File.from_bytes usage) context_data = _external_get_request_context() return RequestContextOutput( @@ -2208,7 +2284,6 @@ def test_request_context_app( ): request_context_app = wrap_app(RequestContextApp) with register_app(request_context_app, "request-context") as (app_alias, _): - time.sleep(1) yield f"{user.username}/{app_alias}" @@ -2216,7 +2291,10 @@ def test_request_context_app( def test_request_context_fields_populated(test_request_context_app: str): """Test that request context fields are properly populated.""" - result = apps.run(test_request_context_app, arguments={}) + result = apps.run( + test_request_context_app, + arguments={"synchronize": False}, + ) assert result["request_id_from_context"] is not None assert result["endpoint_from_context"] is not None @@ -2234,10 +2312,10 @@ def test_request_context_isolation_with_multiplexing(test_request_context_app: s matches the request_id from headers for each individual request. """ - # Submit multiple concurrent requests - handle_1 = apps.submit(test_request_context_app, arguments={}) - handle_2 = apps.submit(test_request_context_app, arguments={}) - handle_3 = apps.submit(test_request_context_app, arguments={}) + arguments = {"synchronize": True} + handle_1 = apps.submit(test_request_context_app, arguments=arguments) + handle_2 = apps.submit(test_request_context_app, arguments=arguments) + handle_3 = apps.submit(test_request_context_app, arguments=arguments) # Get results result_1 = handle_1.get() diff --git a/projects/fal/tests/integration/test_environments.py b/projects/fal/tests/integration/test_environments.py index 79ca84a46..7fbc7f80b 100644 --- a/projects/fal/tests/integration/test_environments.py +++ b/projects/fal/tests/integration/test_environments.py @@ -3,7 +3,7 @@ from __future__ import annotations import uuid -from datetime import datetime, timezone +from datetime import datetime, timedelta, timezone import pytest @@ -133,17 +133,14 @@ def test_environment_created_at_timestamp( client: SyncServerlessClient, test_env_name: str ): """Test that created_at timestamp is properly set.""" - # Create environment + earliest_created_at = datetime.now(timezone.utc) - timedelta(minutes=1) env = client.environments.create(test_env_name) + latest_created_at = datetime.now(timezone.utc) + timedelta(minutes=1) try: assert env.created_at is not None - # Check that it's a datetime object assert isinstance(env.created_at, datetime) - # Check that it's recent (within last minute) - now = datetime.now(timezone.utc) - time_diff = (now - env.created_at).total_seconds() - assert -60 <= time_diff < 60, f"Created time seems wrong: {time_diff}s ago" + assert earliest_created_at <= env.created_at <= latest_created_at finally: client.environments.delete(test_env_name) diff --git a/projects/fal/tests/integration/test_stability.py b/projects/fal/tests/integration/test_stability.py index f4b92eefa..1be9f997a 100644 --- a/projects/fal/tests/integration/test_stability.py +++ b/projects/fal/tests/integration/test_stability.py @@ -16,12 +16,11 @@ PACKAGE_NAME = "fall" -def git_revision_short_hash() -> str: - return ( - subprocess.check_output(["git", "rev-parse", "--short", "HEAD"]) - .decode("ascii") - .strip() - ) +GIT_REVISION_SHORT_HASH = ( + subprocess.check_output(["git", "rev-parse", "--short", "HEAD"]) + .decode("ascii") + .strip() +) @pytest.mark.xfail(reason="Temporary mismatch due to grpc version updates. Ping @efiop") @@ -55,20 +54,24 @@ def mult(a, b): assert mult(5, 2) == 10 +@pytest.mark.timeout(180) +@pytest.mark.xdist_group(name="container-builds") def test_regular_function_in_a_container(isolated_client): - @isolated_client("container") + @isolated_client("container", keep_alive=10) def regular_function(): return 42 assert regular_function() == 42 - @isolated_client("container") + @isolated_client("container", keep_alive=10) def mult(a, b): return a * b assert mult(5, 2) == 10 +@pytest.mark.timeout(180) +@pytest.mark.xdist_group(name="container-builds") def test_container_no_venv(isolated_client): actual_python = active_python() @@ -87,6 +90,8 @@ def myfunc(): assert myfunc() == 42 +@pytest.mark.timeout(180) +@pytest.mark.xdist_group(name="container-builds") def test_container_venv(isolated_client): actual_python = active_python() @@ -108,6 +113,8 @@ def myfunc(): @pytest.mark.flaky(max_runs=3) +@pytest.mark.timeout(180) +@pytest.mark.xdist_group(name="container-builds") def test_regular_function_in_a_container_with_custom_image(isolated_client): actual_python = active_python() @@ -116,7 +123,7 @@ def test_regular_function_in_a_container_with_custom_image(isolated_client): image=ContainerImage.from_dockerfile_str( f""" FROM python:{actual_python}-slim - # {git_revision_short_hash()} + # {GIT_REVISION_SHORT_HASH} RUN env """ ), @@ -129,7 +136,7 @@ def regular_function(): @isolated_client( "container", image=ContainerImage.from_dockerfile_str( - f"FROM python:{actual_python}-slim\n# {git_revision_short_hash()}" + f"FROM python:{actual_python}-slim\n# {GIT_REVISION_SHORT_HASH}" ), ) def mult(a, b): @@ -516,6 +523,7 @@ def regular_function(n): assert out.count("computing") == 1 +@pytest.mark.timeout(120) def test_pydantic_serialization(isolated_client): from pydantic import BaseModel, Field diff --git a/projects/fal/tests/unit/test_app.py b/projects/fal/tests/unit/test_app.py index 1b032bc62..6ddf144eb 100644 --- a/projects/fal/tests/unit/test_app.py +++ b/projects/fal/tests/unit/test_app.py @@ -2439,6 +2439,47 @@ def test_openapi_websocket_realtime_metadata_and_schemas(isolate_agent_env): ) +@pytest.mark.asyncio +async def test_realtime_batches_queued_inputs(): + import asyncio + from collections import deque + + from fal.api import RouteSignature + from fal.realtime import _mirror_output + + queue = deque(InputModel(prompt=str(index)) for index in range(5)) + batches = [] + sent_messages = [] + + class RecordingWebSocket: + async def send_bytes(self, message): + sent_messages.append(message) + + async def generate_batch(_, input, *inputs): + batch = [input, *inputs] + batches.append([item.prompt for item in batch]) + if sum(len(batch) for batch in batches) == 5: + queue.append(None) + return {"results": [item.prompt for item in batch]} + + await _mirror_output( + None, + queue, + RecordingWebSocket(), + func=generate_batch, + route_signature=RouteSignature( + path="/realtime", + buffering=10, + max_batch_size=3, + ), + encode_message=lambda output: repr(output).encode(), + input_ready=asyncio.Event(), + ) + + assert batches == [["0", "1", "2"], ["3", "4"]] + assert len(sent_messages) == 2 + + def test_openapi_websocket_barebones_has_no_realtime_marker(isolate_agent_env): app = RealtimeApp() spec = app.openapi() From 5747689d1d625713e81c7b6cb74af1dd7ee72475 Mon Sep 17 00:00:00 2001 From: badayvedat Date: Tue, 25 Aug 2026 00:21:04 +0700 Subject: [PATCH 03/27] test: wait for e2e app alias readiness --- projects/fal/tests/e2e/test_apps.py | 80 ++++++++++++++++++++++++++--- 1 file changed, 72 insertions(+), 8 deletions(-) diff --git a/projects/fal/tests/e2e/test_apps.py b/projects/fal/tests/e2e/test_apps.py index 5d5418d90..564d9fd3d 100644 --- a/projects/fal/tests/e2e/test_apps.py +++ b/projects/fal/tests/e2e/test_apps.py @@ -168,6 +168,73 @@ def _wait_for_alias_runners( ) +def _wait_for_alias_revision( + client, + app_alias: str, + app_revision: str, + *, + timeout: float = 60, +): + def fetch_alias(): + return next( + (alias for alias in client.list_aliases() if alias.alias == app_alias), + None, + ) + + return _wait_until( + fetch_alias, + lambda alias: alias is not None and alias.revision == app_revision, + timeout=timeout, + interval=0.5, + description=f"alias {app_alias} to point to revision {app_revision}", + ) + + +def _wait_for_queue_alias( + app_alias: str, + queue_url: str, + *, + timeout: float = 60, +): + with httpx.Client(headers=_auth_headers()) as client: + + def queue_recognizes_alias(): + # Alias resolution precedes validation of this query parameter, whose + # invalid value makes the gateway return before enqueuing a request. + response = client.post( + queue_url, + params={"fal_max_queue_length": "readiness-probe"}, + json={}, + ) + if response.status_code == 404: + detail = response.json().get("detail", "") + missing_alias_details = { + f"Application {app_alias!r} not found", + f'Application "{app_alias}" not found', + } + if detail in missing_alias_details: + return False + + if response.status_code != 400: + response.raise_for_status() + raise AssertionError(f"Unexpected queue readiness response: {response}") + + data = response.json() + if data.get("request_id") is not None or not data.get("error", "").startswith( + "Invalid fal_max_queue_length" + ): + raise AssertionError(f"Unexpected queue readiness response: {data}") + return True + + _wait_until( + queue_recognizes_alias, + bool, + timeout=timeout, + interval=0.5, + description=f"queue gateway to recognize alias {app_alias}", + ) + + GIT_REVISION_SHORT_HASH = ( subprocess.check_output(["git", "rev-parse", "--short", "HEAD"]) .decode("ascii") @@ -801,6 +868,9 @@ def _register_app( app_revision = result.result.application_id try: + with host._connection as client: + _wait_for_alias_revision(client, app_alias, app_revision) + _wait_for_queue_alias(app_alias, result.service_urls.queue) yield app_alias, app_revision finally: with host._connection as client: @@ -1422,10 +1492,7 @@ def test_app_deploy_scale(host: api.FalServerlessHost, register_app): app_revision = result.result.application_id with host._connection as client: - res = client.list_aliases() - found = next(filter(lambda alias: alias.alias == app_alias, res), None) - assert found, f"Could not find app {app_alias} in {res}" - assert found.revision == app_revision + found = _wait_for_alias_revision(client, app_alias, app_revision) # multiplexing is revision-specific assert ( found.max_multiplexing == 3 @@ -1442,10 +1509,7 @@ def test_app_deploy_scale(host: api.FalServerlessHost, register_app): app_revision = result.result.application_id with host._connection as client: - res = client.list_aliases() - found = next(filter(lambda alias: alias.alias == app_alias, res), None) - assert found, f"Could not find app {app_alias} in {res}" - assert found.revision == app_revision + found = _wait_for_alias_revision(client, app_alias, app_revision) # when scaling, all values are updated assert found.max_multiplexing == 3 assert found.max_concurrency == 2 From 77556d4ade6a32fd4700329723d856a512a11f20 Mon Sep 17 00:00:00 2001 From: badayvedat Date: Tue, 25 Aug 2026 01:57:26 +0700 Subject: [PATCH 04/27] fix: wait for app gateway readiness --- projects/fal/tests/e2e/test_apps.py | 156 +++++++++++++++++++++------- 1 file changed, 120 insertions(+), 36 deletions(-) diff --git a/projects/fal/tests/e2e/test_apps.py b/projects/fal/tests/e2e/test_apps.py index 564d9fd3d..a57127ed0 100644 --- a/projects/fal/tests/e2e/test_apps.py +++ b/projects/fal/tests/e2e/test_apps.py @@ -190,6 +190,54 @@ def fetch_alias(): ) +def _is_alias_not_found_response(response: httpx.Response, app_alias: str) -> bool: + if response.status_code != 404: + return False + + try: + detail = response.json().get("detail", "") + except ValueError: + return False + return detail in { + f"Application {app_alias!r} not found", + f'Application "{app_alias}" not found', + } + + +def _wait_for_stable_alias( + fetch_response: Callable[[float], Optional[httpx.Response]], + app_alias: str, + *, + timeout: float, + description: str, + stable_for: float = 5, +) -> httpx.Response: + # One updated replica can succeed while another still has stale alias state. + deadline = time.monotonic() + timeout + recognized_at = None + response = None + + while True: + remaining = deadline - time.monotonic() + if remaining <= 0: + raise AssertionError(f"Timed out waiting for {description}: {response!r}") + + response = fetch_response(remaining) + if time.monotonic() >= deadline: + raise AssertionError(f"Timed out waiting for {description}: {response!r}") + if response is None or _is_alias_not_found_response(response, app_alias): + recognized_at = None + elif recognized_at is None: + recognized_at = time.monotonic() + elif time.monotonic() - recognized_at >= stable_for: + return response + + remaining = deadline - time.monotonic() + if remaining <= 0: + raise AssertionError(f"Timed out waiting for {description}: {response!r}") + time.sleep(min(0.5, remaining)) + + def _wait_for_queue_alias( app_alias: str, queue_url: str, @@ -198,43 +246,60 @@ def _wait_for_queue_alias( ): with httpx.Client(headers=_auth_headers()) as client: - def queue_recognizes_alias(): + def fetch_response(remaining: float): # Alias resolution precedes validation of this query parameter, whose # invalid value makes the gateway return before enqueuing a request. response = client.post( queue_url, params={"fal_max_queue_length": "readiness-probe"}, json={}, + timeout=min(5, remaining), ) - if response.status_code == 404: - detail = response.json().get("detail", "") - missing_alias_details = { - f"Application {app_alias!r} not found", - f'Application "{app_alias}" not found', - } - if detail in missing_alias_details: - return False + if _is_alias_not_found_response(response, app_alias): + return response if response.status_code != 400: response.raise_for_status() raise AssertionError(f"Unexpected queue readiness response: {response}") data = response.json() - if data.get("request_id") is not None or not data.get("error", "").startswith( + error = data.get("error", "") + if data.get("request_id") is not None or not error.startswith( "Invalid fal_max_queue_length" ): raise AssertionError(f"Unexpected queue readiness response: {data}") - return True + return response - _wait_until( - queue_recognizes_alias, - bool, + _wait_for_stable_alias( + fetch_response, + app_alias, timeout=timeout, - interval=0.5, description=f"queue gateway to recognize alias {app_alias}", ) +def _wait_for_run_alias( + app_alias: str, + run_url: str, + *, + timeout: float = 60, +): + with httpx.Client(headers=_auth_headers()) as client: + + def fetch_response(remaining: float): + try: + return client.get(f"{run_url}/health", timeout=min(5, remaining)) + except httpx.TimeoutException: + return None + + _wait_for_stable_alias( + fetch_response, + app_alias, + timeout=timeout, + description=f"run gateway to recognize alias {app_alias}", + ) + + GIT_REVISION_SHORT_HASH = ( subprocess.check_output(["git", "rev-parse", "--short", "HEAD"]) .decode("ascii") @@ -871,6 +936,7 @@ def _register_app( with host._connection as client: _wait_for_alias_revision(client, app_alias, app_revision) _wait_for_queue_alias(app_alias, result.service_urls.queue) + _wait_for_run_alias(app_alias, result.service_urls.run) yield app_alias, app_revision finally: with host._connection as client: @@ -957,15 +1023,6 @@ def test_health_override_app( yield f"{user.username}/{app_alias}" -@pytest.fixture() -def test_fastapi_app( - user: User, - register_app, -): - with register_app(calculator_app, "fastapi") as (app_alias, _): - yield f"{user.username}/{app_alias}" - - @pytest.fixture(scope="module") def test_stateful_app( user: User, @@ -1426,22 +1483,49 @@ def test_app_openapi_spec_metadata(test_app: str, rest_client: Client): assert key in openapi_spec, f"{key} key missing from openapi {openapi_spec}" -def test_app_no_serve_spec_metadata(test_fastapi_app: str, rest_client: Client): +def test_app_no_serve_spec_metadata( + host: api.FalServerlessHost, + user: User, + rest_client: Client, + make_tmp_app_name: Callable[[str], str], +): # We do not store the openapi spec for apps that do not use serve=True - user_id, _, app_id = test_fastapi_app.partition("/") - res = app_metadata.sync_detailed( - app_alias_or_id=app_id, app_user_id=user_id, client=rest_client + app_alias = make_tmp_app_name("fastapi") + result = host.register( + func=calculator_app.func, + options=calculator_app.options, + application_name=app_alias, + application_auth_mode="private", + deployment_strategy="recreate", ) - assert ( - res.status_code == 200 - ), f"Failed to fetch metadata for app {test_fastapi_app}" - assert res.parsed, f"Failed to parse metadata for app {test_fastapi_app}" + assert result + assert result.result - metadata = res.parsed.to_dict() - assert ( - "openapi" not in metadata - ), f"openapi should not be present in metadata {metadata}" + try: + with host._connection as client: + _wait_for_alias_revision(client, app_alias, result.result.application_id) + + res = app_metadata.sync_detailed( + app_alias_or_id=app_alias, + app_user_id=user.username, + client=rest_client, + ) + + assert ( + res.status_code == 200 + ), f"Failed to fetch metadata for app {user.username}/{app_alias}" + assert ( + res.parsed + ), f"Failed to parse metadata for app {user.username}/{app_alias}" + + metadata = res.parsed.to_dict() + assert ( + "openapi" not in metadata + ), f"openapi should not be present in metadata {metadata}" + finally: + with host._connection as client: + client.delete_alias(app_alias) @pytest.mark.xdist_group(name="addition-app") From 861e73ba1dbabf6525d71819b200e104e81e2bbe Mon Sep 17 00:00:00 2001 From: badayvedat Date: Tue, 25 Aug 2026 04:46:00 +0700 Subject: [PATCH 05/27] chore: inline max multiplexing --- projects/fal/tests/e2e/test_apps.py | 11 ++++------- 1 file changed, 4 insertions(+), 7 deletions(-) diff --git a/projects/fal/tests/e2e/test_apps.py b/projects/fal/tests/e2e/test_apps.py index a57127ed0..403f52779 100644 --- a/projects/fal/tests/e2e/test_apps.py +++ b/projects/fal/tests/e2e/test_apps.py @@ -2282,14 +2282,11 @@ class AppRefOutput(BaseModel): from_external_method: str -CONCURRENT_REQUESTS = 3 - - class AppRefApp( fal.App, keep_alive=300, max_concurrency=1, - max_multiplexing=CONCURRENT_REQUESTS, + max_multiplexing=3, ): machine_type = "XS" @@ -2302,7 +2299,7 @@ async def run(self, request: Request) -> AppRefOutput: request_id = request.headers.get("x-request-id", "") self.concurrent_requests += 1 - if self.concurrent_requests == CONCURRENT_REQUESTS: + if self.concurrent_requests == self.max_multiplexing: self.requests_ready.set() await asyncio.wait_for(self.requests_ready.wait(), timeout=30) @@ -2395,7 +2392,7 @@ class RequestContextApp( fal.App, keep_alive=300, max_concurrency=1, - max_multiplexing=CONCURRENT_REQUESTS, + max_multiplexing=3, ): """App to test request context fields are properly populated.""" @@ -2411,7 +2408,7 @@ async def get_context( ) -> RequestContextOutput: if input.synchronize: self.concurrent_requests += 1 - if self.concurrent_requests == CONCURRENT_REQUESTS: + if self.concurrent_requests == self.max_multiplexing: self.requests_ready.set() await asyncio.wait_for(self.requests_ready.wait(), timeout=30) From 5aa056cfae87d65ee4e13e5367eba7ca30f7e071 Mon Sep 17 00:00:00 2001 From: badayvedat Date: Tue, 25 Aug 2026 05:06:16 +0700 Subject: [PATCH 06/27] test: probe app readiness through queue --- projects/fal/tests/e2e/test_apps.py | 81 ++++++++++++++++++++++------- 1 file changed, 63 insertions(+), 18 deletions(-) diff --git a/projects/fal/tests/e2e/test_apps.py b/projects/fal/tests/e2e/test_apps.py index 403f52779..5c8d65862 100644 --- a/projects/fal/tests/e2e/test_apps.py +++ b/projects/fal/tests/e2e/test_apps.py @@ -278,27 +278,72 @@ def fetch_response(remaining: float): ) -def _wait_for_run_alias( - app_alias: str, - run_url: str, - *, - timeout: float = 60, -): - with httpx.Client(headers=_auth_headers()) as client: +def _wait_for_app_ready(queue_url: str, *, timeout: float = 60): + import fal_client - def fetch_response(remaining: float): - try: - return client.get(f"{run_url}/health", timeout=min(5, remaining)) - except httpx.TimeoutException: - return None + application = httpx.URL(queue_url).path.lstrip("/") + deadline = time.monotonic() + timeout + fal_client_instance = fal_client.SyncClient(default_timeout=timeout) + handle = None + terminal = False - _wait_for_stable_alias( - fetch_response, - app_alias, - timeout=timeout, - description=f"run gateway to recognize alias {app_alias}", + def remaining_time(): + remaining = deadline - time.monotonic() + if remaining <= 0: + raise AssertionError( + f"Timed out waiting for app readiness probe for {application}" + ) + return remaining + + try: + handle = fal_client_instance.submit( + application, + arguments={}, + path="health", + start_timeout=timeout, ) + with httpx.Client(headers=_auth_headers()) as client: + while True: + remaining = remaining_time() + + try: + response = client.get( + handle.status_url, + params={"logs": False}, + timeout=min(5, remaining), + ) + except httpx.TransportError: + time.sleep(min(0.1, remaining_time())) + continue + + retryable_status = response.status_code in {408, 409, 429} + retryable_ingress_error = ( + response.status_code in {502, 503, 504} + and "x-fal-request-id" not in response.headers + and "nginx" in response.text + ) + if retryable_status or retryable_ingress_error: + time.sleep(min(0.1, remaining_time())) + continue + + response.raise_for_status() + status = response.json() + terminal = status["status"] == "COMPLETED" + + remaining = remaining_time() + + if terminal: + assert status.get("error") is None, status["error"] + return + + time.sleep(min(0.1, remaining)) + finally: + if handle is not None and not terminal: + with suppress(Exception): + with httpx.Client(headers=_auth_headers()) as client: + client.put(handle.cancel_url, timeout=5).raise_for_status() + GIT_REVISION_SHORT_HASH = ( subprocess.check_output(["git", "rev-parse", "--short", "HEAD"]) @@ -936,7 +981,7 @@ def _register_app( with host._connection as client: _wait_for_alias_revision(client, app_alias, app_revision) _wait_for_queue_alias(app_alias, result.service_urls.queue) - _wait_for_run_alias(app_alias, result.service_urls.run) + _wait_for_app_ready(result.service_urls.queue) yield app_alias, app_revision finally: with host._connection as client: From 256f46410746df1c1e15303651e0956ac4ecb38e Mon Sep 17 00:00:00 2001 From: badayvedat Date: Tue, 25 Aug 2026 15:34:08 +0700 Subject: [PATCH 07/27] test: fix app readiness and multiplexing probes --- projects/fal/tests/e2e/test_apps.py | 135 +++++++++++++--------------- 1 file changed, 62 insertions(+), 73 deletions(-) diff --git a/projects/fal/tests/e2e/test_apps.py b/projects/fal/tests/e2e/test_apps.py index 5c8d65862..e56e7cc5d 100644 --- a/projects/fal/tests/e2e/test_apps.py +++ b/projects/fal/tests/e2e/test_apps.py @@ -16,6 +16,7 @@ Iterator, List, Optional, + Protocol, Tuple, TypeVar, Union, @@ -278,71 +279,42 @@ def fetch_response(remaining: float): ) -def _wait_for_app_ready(queue_url: str, *, timeout: float = 60): - import fal_client - - application = httpx.URL(queue_url).path.lstrip("/") +def _wait_for_app_ready( + app_alias: str, + run_url: str, + *, + path: str = "/health", + timeout: float = 60, +): deadline = time.monotonic() + timeout - fal_client_instance = fal_client.SyncClient(default_timeout=timeout) - handle = None - terminal = False - - def remaining_time(): - remaining = deadline - time.monotonic() - if remaining <= 0: - raise AssertionError( - f"Timed out waiting for app readiness probe for {application}" - ) - return remaining - - try: - handle = fal_client_instance.submit( - application, - arguments={}, - path="health", - start_timeout=timeout, - ) + response = None - with httpx.Client(headers=_auth_headers()) as client: - while True: - remaining = remaining_time() - - try: - response = client.get( - handle.status_url, - params={"logs": False}, - timeout=min(5, remaining), - ) - except httpx.TransportError: - time.sleep(min(0.1, remaining_time())) - continue - - retryable_status = response.status_code in {408, 409, 429} - retryable_ingress_error = ( - response.status_code in {502, 503, 504} - and "x-fal-request-id" not in response.headers - and "nginx" in response.text + with httpx.Client(headers=_auth_headers()) as client: + while True: + remaining = deadline - time.monotonic() + if remaining <= 0: + raise AssertionError( + f"Timed out waiting for app readiness for {app_alias}: {response!r}" ) - if retryable_status or retryable_ingress_error: - time.sleep(min(0.1, remaining_time())) - continue - - response.raise_for_status() - status = response.json() - terminal = status["status"] == "COMPLETED" - - remaining = remaining_time() - if terminal: - assert status.get("error") is None, status["error"] - return + try: + response = client.get(f"{run_url}{path}", timeout=remaining) + except httpx.TransportError: + time.sleep(min(0.1, max(0, deadline - time.monotonic()))) + continue + + alias_not_found = _is_alias_not_found_response(response, app_alias) + retryable_status = response.status_code in {408, 409, 429} + retryable_ingress_error = ( + response.status_code in {502, 503, 504} + and "x-fal-request-id" not in response.headers + ) + if alias_not_found or retryable_status or retryable_ingress_error: + time.sleep(min(0.5, max(0, deadline - time.monotonic()))) + continue - time.sleep(min(0.1, remaining)) - finally: - if handle is not None and not terminal: - with suppress(Exception): - with httpx.Client(headers=_auth_headers()) as client: - client.put(handle.cancel_url, timeout=5).raise_for_status() + response.raise_for_status() + return GIT_REVISION_SHORT_HASH = ( @@ -947,21 +919,25 @@ def user(rest_client: Client) -> Generator[User, None, None]: yield user +class RegisterApp(Protocol): + def __call__( + self, + app: Union[api.ServedIsolatedFunction, api.IsolatedFunction], + suffix: str = "", + readiness_path: str = "/health", + ) -> ContextManager[Tuple[str, str]]: ... + + @pytest.fixture(scope="module") def register_app( host: api.FalServerlessHost, make_tmp_app_name: Callable[[str], str], -) -> Callable[ - [ - Union[api.ServedIsolatedFunction, api.IsolatedFunction], - str, - ], - ContextManager[Tuple[str, str]], -]: +) -> RegisterApp: @contextmanager def _register_app( app: Union[api.ServedIsolatedFunction, api.IsolatedFunction], suffix: str = "", + readiness_path: str = "/health", ): app_alias = make_tmp_app_name(suffix) result = host.register( @@ -981,7 +957,11 @@ def _register_app( with host._connection as client: _wait_for_alias_revision(client, app_alias, app_revision) _wait_for_queue_alias(app_alias, result.service_urls.queue) - _wait_for_app_ready(result.service_urls.queue) + _wait_for_app_ready( + app_alias, + result.service_urls.run, + path=readiness_path, + ) yield app_alias, app_revision finally: with host._connection as client: @@ -1045,7 +1025,10 @@ def test_custom_health_path_app( user: User, register_app, ): - with register_app(custom_health_path_app, "custom-health") as (app_alias, _): + with register_app(custom_health_path_app, "custom-health", "/ready") as ( + app_alias, + _, + ): yield f"{user.username}/{app_alias}" @@ -1054,7 +1037,10 @@ def test_health_override_fn( user: User, register_app, ): - with register_app(health_override_fn, "health-override-fn") as (app_alias, _): + with register_app(health_override_fn, "health-override-fn", "/ready") as ( + app_alias, + _, + ): yield f"{user.username}/{app_alias}" @@ -1064,7 +1050,10 @@ def test_health_override_app( register_app, ): health_override_app = wrap_app(HealthOverrideApp) - with register_app(health_override_app, "health-override-app") as (app_alias, _): + with register_app(health_override_app, "health-override-app", "/ready") as ( + app_alias, + _, + ): yield f"{user.username}/{app_alias}" @@ -2331,9 +2320,9 @@ class AppRefApp( fal.App, keep_alive=300, max_concurrency=1, - max_multiplexing=3, ): machine_type = "XS" + max_multiplexing = 3 async def setup(self): self.concurrent_requests = 0 @@ -2437,11 +2426,11 @@ class RequestContextApp( fal.App, keep_alive=300, max_concurrency=1, - max_multiplexing=3, ): """App to test request context fields are properly populated.""" machine_type = "XS" + max_multiplexing = 3 async def setup(self): self.concurrent_requests = 0 From f1cfacaf8e0d8fdb9a1c2b18c9537c75792bd3d9 Mon Sep 17 00:00:00 2001 From: badayvedat Date: Tue, 25 Aug 2026 16:51:05 +0700 Subject: [PATCH 08/27] chore: no op change for probing tests --- README.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/README.md b/README.md index bafc65788..bccf2ffeb 100644 --- a/README.md +++ b/README.md @@ -2,7 +2,7 @@ [![Tests](https://img.shields.io/github/actions/workflow/status/fal-ai/fal/fal-unit-tests.yml?label=Tests)](https://github.com/fal-ai/fal/actions) # fal - +a fal is a serverless Python runtime that lets you run and scale code in the cloud with no infra management. With fal, you can build pipelines, serve ML models and scale them up to many users. You scale down to 0 when you don't use any resources. From 56196499b7f2300f92a2fa152141e87ae3305b29 Mon Sep 17 00:00:00 2001 From: badayvedat Date: Tue, 25 Aug 2026 18:30:51 +0700 Subject: [PATCH 09/27] test: extend app readiness timeout --- projects/fal/tests/e2e/test_apps.py | 18 +++++++++--------- 1 file changed, 9 insertions(+), 9 deletions(-) diff --git a/projects/fal/tests/e2e/test_apps.py b/projects/fal/tests/e2e/test_apps.py index e56e7cc5d..8afb3535f 100644 --- a/projects/fal/tests/e2e/test_apps.py +++ b/projects/fal/tests/e2e/test_apps.py @@ -51,6 +51,8 @@ from fal.toolkit.utils.endpoint import cancel_on_disconnect from fal.workflows import Workflow +pytestmark = pytest.mark.timeout(420) + @pytest.fixture(scope="module") def rest_client() -> Generator[Client, None, None]: @@ -284,7 +286,7 @@ def _wait_for_app_ready( run_url: str, *, path: str = "/health", - timeout: float = 60, + timeout: float = 180, ): deadline = time.monotonic() + timeout response = None @@ -1377,7 +1379,6 @@ def test_app_disconnect_behavior(test_cancellable_app: str): ), "Expected Gateway Timeout even though the app handled it" -@pytest.mark.timeout(120) def test_start_timeout_queue_blocking(test_queue_blocking_app: str): """ Test that start_timeout correctly times out a request waiting in queue. @@ -1416,7 +1417,6 @@ def test_start_timeout_queue_blocking(test_queue_blocking_app: str): assert timeout_type == "user", f"Expected 'user' timeout type, got {timeout_type}" -@pytest.mark.timeout(120) def test_app_client_async(test_sleep_app: str): handle = apps.submit(test_sleep_app, arguments={"wait_time": 10}) _wait_for_request_status(handle, apps.InProgress) @@ -2068,7 +2068,7 @@ def submit_and_wait_for_runner( return handle -@pytest.mark.timeout(180) +@pytest.mark.timeout(480) def test_stop_runner(host: api.FalServerlessHost, test_sleep_app: str): _, _, app_alias = test_sleep_app.partition("/") submit_and_wait_for_runner(test_sleep_app, arguments={"wait_time": 1}) @@ -2128,7 +2128,7 @@ def test_stop_runner(host: api.FalServerlessHost, test_sleep_app: str): ) -@pytest.mark.timeout(180) +@pytest.mark.timeout(480) def test_kill_runner(host: api.FalServerlessHost, test_sleep_app: str): handle = apps.submit(test_sleep_app, arguments={"wait_time": 30}) _wait_for_request_status(handle, apps.InProgress, timeout=60) @@ -2156,7 +2156,7 @@ def test_kill_runner(host: api.FalServerlessHost, test_sleep_app: str): ) -@pytest.mark.timeout(180) +@pytest.mark.timeout(480) def test_rollout_application(host: api.FalServerlessHost, test_sleep_app: str): handle = apps.submit(test_sleep_app, arguments={"wait_time": 30}) _wait_for_request_status(handle, apps.InProgress, timeout=60) @@ -2195,7 +2195,7 @@ def test_rollout_application(host: api.FalServerlessHost, test_sleep_app: str): ) -@pytest.mark.timeout(180) +@pytest.mark.timeout(480) def test_shell_runner(host: api.FalServerlessHost, test_sleep_app: str): handle = submit_and_wait_for_runner(test_sleep_app, arguments={"wait_time": 1}) assert handle.get() == {"slept": True} @@ -2226,7 +2226,7 @@ def test_shell_runner(host: api.FalServerlessHost, test_sleep_app: str): proc.wait() -@pytest.mark.timeout(180) +@pytest.mark.timeout(480) def test_exec_runner(host: api.FalServerlessHost, test_sleep_app: str): handle = submit_and_wait_for_runner(test_sleep_app, arguments={"wait_time": 1}) assert handle.get() == {"slept": True} @@ -2367,7 +2367,7 @@ def test_app_ref_app_client(test_app_ref_app: str): assert result_3["from_app"] == result_3["from_external_method"] -@pytest.mark.timeout(180) +@pytest.mark.timeout(480) def test_runner_machine_type(host: api.FalServerlessHost, test_sleep_app: str): """Test that machine_type is populated in runner info.""" search_start = datetime.now() - timedelta(minutes=5) From 949e1c3d0ca73d532cd99bc601ab9fb0907b13a0 Mon Sep 17 00:00:00 2001 From: badayvedat Date: Tue, 25 Aug 2026 19:05:17 +0700 Subject: [PATCH 10/27] test: extend remaining app startup timeouts --- projects/fal/tests/e2e/test_apps.py | 4 ++-- projects/fal/tests/integration/test_stability.py | 3 ++- 2 files changed, 4 insertions(+), 3 deletions(-) diff --git a/projects/fal/tests/e2e/test_apps.py b/projects/fal/tests/e2e/test_apps.py index 8afb3535f..51419abc5 100644 --- a/projects/fal/tests/e2e/test_apps.py +++ b/projects/fal/tests/e2e/test_apps.py @@ -1081,7 +1081,7 @@ def test_cancellable_app( @pytest.fixture(scope="module") def test_exception_app(): - with AppClient.connect(ExceptionApp) as client: + with AppClient.connect(ExceptionApp, startup_timeout=180) as client: yield client @@ -2294,7 +2294,7 @@ def test_hints_encoding(): Make sure that hints that can't be encoded in latin-1 don't crash the app https://github.com/encode/starlette/blob/a766a58d14007f07c0b5782fa78cdc370b892796/starlette/datastructures.py#L568 """ - with AppClient.connect(HintsApp) as client: + with AppClient.connect(HintsApp, startup_timeout=180) as client: with httpx.Client(headers=_auth_headers()) as httpx_client: url = client.url + "/add" resp = httpx_client.post( diff --git a/projects/fal/tests/integration/test_stability.py b/projects/fal/tests/integration/test_stability.py index 1be9f997a..6ac68122a 100644 --- a/projects/fal/tests/integration/test_stability.py +++ b/projects/fal/tests/integration/test_stability.py @@ -386,6 +386,7 @@ def regular_function(): assert regular_function() == 42 +@pytest.mark.timeout(420) def test_faulty_setup_function(isolated_client): def good_setup_function(): return 42 @@ -393,7 +394,7 @@ def good_setup_function(): def bad_setup_function(): raise ValueError() - @isolated_client("virtualenv") + @isolated_client("virtualenv", startup_timeout=180) def regular_function(result): return result * 2 From 15768fbd3313dddcbba397ec3db05de986d93e18 Mon Sep 17 00:00:00 2001 From: badayvedat Date: Tue, 25 Aug 2026 20:09:56 +0700 Subject: [PATCH 11/27] test: extend integration worker startup limits --- .github/workflows/fal-integration-tests.yml | 2 +- projects/fal/tests/conftest.py | 7 ++++++- projects/fal/tests/integration/test.py | 2 ++ projects/fal/tests/integration/test_stability.py | 11 ++++------- projects/fal/tests/integration/toolkit/test_image.py | 2 ++ 5 files changed, 15 insertions(+), 9 deletions(-) diff --git a/.github/workflows/fal-integration-tests.yml b/.github/workflows/fal-integration-tests.yml index 554fd6547..201ce7b7f 100644 --- a/.github/workflows/fal-integration-tests.yml +++ b/.github/workflows/fal-integration-tests.yml @@ -29,7 +29,7 @@ jobs: integration: name: integration (${{ matrix.os }}, py ${{ matrix.python }}, ${{ matrix.deps }}) runs-on: ${{ matrix.os }} - timeout-minutes: 10 + timeout-minutes: 15 strategy: fail-fast: false max-parallel: 4 diff --git a/projects/fal/tests/conftest.py b/projects/fal/tests/conftest.py index b4e6e76c9..640fce6b5 100644 --- a/projects/fal/tests/conftest.py +++ b/projects/fal/tests/conftest.py @@ -17,7 +17,12 @@ @pytest.fixture(scope="function") def isolated_client(): - return partial(function, machine_type="XS", keep_alive=0) + return partial( + function, + machine_type="XS", + keep_alive=0, + startup_timeout=180, + ) @pytest.fixture(scope="module") diff --git a/projects/fal/tests/integration/test.py b/projects/fal/tests/integration/test.py index e087172ee..89976ca7c 100644 --- a/projects/fal/tests/integration/test.py +++ b/projects/fal/tests/integration/test.py @@ -23,6 +23,8 @@ EXAMPLE_FILE_URL = "https://raw.githubusercontent.com/fal-ai/fal/main/projects/fal/tests/assets/cat.png" +pytestmark = pytest.mark.timeout(420) + @pytest.mark.flaky(max_runs=3) def test_isolated(isolated_client: Callable[..., Callable[..., IsolatedFunction]]): diff --git a/projects/fal/tests/integration/test_stability.py b/projects/fal/tests/integration/test_stability.py index 6ac68122a..e77e7535b 100644 --- a/projects/fal/tests/integration/test_stability.py +++ b/projects/fal/tests/integration/test_stability.py @@ -15,6 +15,8 @@ PACKAGE_NAME = "fall" +pytestmark = pytest.mark.timeout(420) + GIT_REVISION_SHORT_HASH = ( subprocess.check_output(["git", "rev-parse", "--short", "HEAD"]) @@ -54,7 +56,6 @@ def mult(a, b): assert mult(5, 2) == 10 -@pytest.mark.timeout(180) @pytest.mark.xdist_group(name="container-builds") def test_regular_function_in_a_container(isolated_client): @isolated_client("container", keep_alive=10) @@ -70,7 +71,6 @@ def mult(a, b): assert mult(5, 2) == 10 -@pytest.mark.timeout(180) @pytest.mark.xdist_group(name="container-builds") def test_container_no_venv(isolated_client): actual_python = active_python() @@ -90,7 +90,6 @@ def myfunc(): assert myfunc() == 42 -@pytest.mark.timeout(180) @pytest.mark.xdist_group(name="container-builds") def test_container_venv(isolated_client): actual_python = active_python() @@ -113,7 +112,6 @@ def myfunc(): @pytest.mark.flaky(max_runs=3) -@pytest.mark.timeout(180) @pytest.mark.xdist_group(name="container-builds") def test_regular_function_in_a_container_with_custom_image(isolated_client): actual_python = active_python() @@ -358,6 +356,7 @@ def memory_overflow_crash_on_run(): memory_overflow_crash_on_run() +@pytest.mark.timeout(600) def test_keepalive_after_agent_exit(isolated_client): # Should work (fresh) @isolated_client("virtualenv") @@ -386,7 +385,6 @@ def regular_function(): assert regular_function() == 42 -@pytest.mark.timeout(420) def test_faulty_setup_function(isolated_client): def good_setup_function(): return 42 @@ -394,7 +392,7 @@ def good_setup_function(): def bad_setup_function(): raise ValueError() - @isolated_client("virtualenv", startup_timeout=180) + @isolated_client("virtualenv") def regular_function(result): return result * 2 @@ -524,7 +522,6 @@ def regular_function(n): assert out.count("computing") == 1 -@pytest.mark.timeout(120) def test_pydantic_serialization(isolated_client): from pydantic import BaseModel, Field diff --git a/projects/fal/tests/integration/toolkit/test_image.py b/projects/fal/tests/integration/toolkit/test_image.py index 095cc6007..6975598b2 100644 --- a/projects/fal/tests/integration/toolkit/test_image.py +++ b/projects/fal/tests/integration/toolkit/test_image.py @@ -11,6 +11,8 @@ from fal.toolkit import Image +pytestmark = pytest.mark.timeout(420) + @overload def get_image(as_bytes: Literal[False] = False) -> PILImage.Image: ... From 4ac49833d4e3af52c87877bf75b611c3c4cdf263 Mon Sep 17 00:00:00 2001 From: badayvedat Date: Tue, 1 Sep 2026 03:01:04 +0700 Subject: [PATCH 12/27] test: make queue alias readiness probe non-enqueuing --- projects/fal/tests/e2e/test_apps.py | 69 +++++++++-------------------- 1 file changed, 20 insertions(+), 49 deletions(-) diff --git a/projects/fal/tests/e2e/test_apps.py b/projects/fal/tests/e2e/test_apps.py index 51419abc5..25d7df1eb 100644 --- a/projects/fal/tests/e2e/test_apps.py +++ b/projects/fal/tests/e2e/test_apps.py @@ -207,78 +207,49 @@ def _is_alias_not_found_response(response: httpx.Response, app_alias: str) -> bo } -def _wait_for_stable_alias( - fetch_response: Callable[[float], Optional[httpx.Response]], - app_alias: str, - *, - timeout: float, - description: str, - stable_for: float = 5, -) -> httpx.Response: - # One updated replica can succeed while another still has stale alias state. - deadline = time.monotonic() + timeout - recognized_at = None - response = None - - while True: - remaining = deadline - time.monotonic() - if remaining <= 0: - raise AssertionError(f"Timed out waiting for {description}: {response!r}") - - response = fetch_response(remaining) - if time.monotonic() >= deadline: - raise AssertionError(f"Timed out waiting for {description}: {response!r}") - if response is None or _is_alias_not_found_response(response, app_alias): - recognized_at = None - elif recognized_at is None: - recognized_at = time.monotonic() - elif time.monotonic() - recognized_at >= stable_for: - return response - - remaining = deadline - time.monotonic() - if remaining <= 0: - raise AssertionError(f"Timed out waiting for {description}: {response!r}") - time.sleep(min(0.5, remaining)) - - def _wait_for_queue_alias( app_alias: str, queue_url: str, *, timeout: float = 60, ): + deadline = time.monotonic() + timeout + response = None + not_found_count = 0 + with httpx.Client(headers=_auth_headers()) as client: + while True: + remaining = deadline - time.monotonic() + if remaining <= 0: + raise AssertionError( + f"Timed out waiting for queue gateway to recognize {app_alias}; " + f"404 responses={not_found_count}; last_response={response!r}" + ) - def fetch_response(remaining: float): - # Alias resolution precedes validation of this query parameter, whose - # invalid value makes the gateway return before enqueuing a request. response = client.post( queue_url, - params={"fal_max_queue_length": "readiness-probe"}, + # A zero limit always returns 429 after alias resolution and before + # enqueueing because queue lengths cannot be negative. + params={"fal_max_queue_length": "0"}, json={}, timeout=min(5, remaining), ) if _is_alias_not_found_response(response, app_alias): - return response + not_found_count += 1 + time.sleep(min(0.5, max(0, deadline - time.monotonic()))) + continue - if response.status_code != 400: + if response.status_code != 429: response.raise_for_status() raise AssertionError(f"Unexpected queue readiness response: {response}") data = response.json() error = data.get("error", "") if data.get("request_id") is not None or not error.startswith( - "Invalid fal_max_queue_length" + "Queue is too long" ): raise AssertionError(f"Unexpected queue readiness response: {data}") - return response - - _wait_for_stable_alias( - fetch_response, - app_alias, - timeout=timeout, - description=f"queue gateway to recognize alias {app_alias}", - ) + return def _wait_for_app_ready( From f32084d8773974fa0f81d353774723f9e24dcc10 Mon Sep 17 00:00:00 2001 From: badayvedat Date: Tue, 1 Sep 2026 23:07:58 +0700 Subject: [PATCH 13/27] chore: dump change to run all tests 1 --- README.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/README.md b/README.md index bccf2ffeb..6136982ec 100644 --- a/README.md +++ b/README.md @@ -2,7 +2,7 @@ [![Tests](https://img.shields.io/github/actions/workflow/status/fal-ai/fal/fal-unit-tests.yml?label=Tests)](https://github.com/fal-ai/fal/actions) # fal -a +aa fal is a serverless Python runtime that lets you run and scale code in the cloud with no infra management. With fal, you can build pipelines, serve ML models and scale them up to many users. You scale down to 0 when you don't use any resources. From 9e422186bb1bd43265638363e1de00ba866bdb4c Mon Sep 17 00:00:00 2001 From: badayvedat Date: Wed, 2 Sep 2026 00:50:50 +0700 Subject: [PATCH 14/27] test(fal): wait for alias deletion convergence --- projects/fal/tests/e2e/test_apps.py | 11 +++++++---- 1 file changed, 7 insertions(+), 4 deletions(-) diff --git a/projects/fal/tests/e2e/test_apps.py b/projects/fal/tests/e2e/test_apps.py index 25d7df1eb..e12a43906 100644 --- a/projects/fal/tests/e2e/test_apps.py +++ b/projects/fal/tests/e2e/test_apps.py @@ -1684,10 +1684,13 @@ def test_app_set_delete_alias(base_app: Tuple[str, str]): assert res == app_revision with host._connection as client: - # Get the registered values - res = client.list_aliases() - found = next(filter(lambda alias: alias.alias == app_alias, res), None) - assert not found, f"Found app {app_alias} in {res} after deletion" + _wait_until( + client.list_aliases, + lambda aliases: all(alias.alias != app_alias for alias in aliases), + timeout=30, + interval=0.5, + description=f"alias {app_alias} deletion", + ) @pytest.mark.xdist_group(name="realtime-app") From 6b5313216ab2a7ee1ba7454c88ddd6994ce866be Mon Sep 17 00:00:00 2001 From: badayvedat Date: Wed, 2 Sep 2026 01:02:09 +0700 Subject: [PATCH 15/27] ci(fal): extend e2e job timeout --- .github/workflows/fal-e2e-tests.yml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.github/workflows/fal-e2e-tests.yml b/.github/workflows/fal-e2e-tests.yml index 60a2538df..eb4eb30c4 100644 --- a/.github/workflows/fal-e2e-tests.yml +++ b/.github/workflows/fal-e2e-tests.yml @@ -31,7 +31,7 @@ jobs: e2e: name: e2e (${{ matrix.os }}, py ${{ matrix.python }}, ${{ matrix.deps }}) runs-on: ${{ matrix.os }} - timeout-minutes: 10 + timeout-minutes: 15 strategy: fail-fast: false max-parallel: 2 From e80215908c8a65b4a32da85a8f124855ac7fe4a2 Mon Sep 17 00:00:00 2001 From: badayvedat Date: Wed, 2 Sep 2026 03:22:33 +0700 Subject: [PATCH 16/27] chore: dump change to run all tests 1 --- README.md | 2 +- t.py | 20 ++++++++++++++++++++ 2 files changed, 21 insertions(+), 1 deletion(-) create mode 100644 t.py diff --git a/README.md b/README.md index 6136982ec..5e30a2bbb 100644 --- a/README.md +++ b/README.md @@ -2,7 +2,7 @@ [![Tests](https://img.shields.io/github/actions/workflow/status/fal-ai/fal/fal-unit-tests.yml?label=Tests)](https://github.com/fal-ai/fal/actions) # fal -aa +aaa fal is a serverless Python runtime that lets you run and scale code in the cloud with no infra management. With fal, you can build pipelines, serve ML models and scale them up to many users. You scale down to 0 when you don't use any resources. diff --git a/t.py b/t.py new file mode 100644 index 000000000..24b00836b --- /dev/null +++ b/t.py @@ -0,0 +1,20 @@ +import time +import uuid + +from fal.api import SyncServerlessClient + +client = SyncServerlessClient() + +for attempt in range(10): + name = f"replica-lag-repro-{uuid.uuid4().hex[:8]}" + print(f"testing {name}") + env = client.environments.create(name) + + assert name in [env.name for env in client.environments.list()] + + client.environments.delete(name) + + assert name not in [env.name for env in client.environments.list()] + + time.sleep(0.1) + print(f"Attempt {attempt + 1}/ 10 passed") From 4d2631226c35b8206cc2d19af8c24dc17038be2d Mon Sep 17 00:00:00 2001 From: badayvedat Date: Sat, 5 Sep 2026 05:37:50 +0700 Subject: [PATCH 17/27] chore: no-op random change --- README.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/README.md b/README.md index 5e30a2bbb..dbf135a55 100644 --- a/README.md +++ b/README.md @@ -2,7 +2,7 @@ [![Tests](https://img.shields.io/github/actions/workflow/status/fal-ai/fal/fal-unit-tests.yml?label=Tests)](https://github.com/fal-ai/fal/actions) # fal -aaa +aaaa fal is a serverless Python runtime that lets you run and scale code in the cloud with no infra management. With fal, you can build pipelines, serve ML models and scale them up to many users. You scale down to 0 when you don't use any resources. From aa0fc7c87cdc7ca422b929ab62d966079e599d00 Mon Sep 17 00:00:00 2001 From: badayvedat Date: Sat, 5 Sep 2026 06:02:53 +0700 Subject: [PATCH 18/27] chore: no-op random change --- README.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/README.md b/README.md index dbf135a55..8d9a92405 100644 --- a/README.md +++ b/README.md @@ -2,7 +2,7 @@ [![Tests](https://img.shields.io/github/actions/workflow/status/fal-ai/fal/fal-unit-tests.yml?label=Tests)](https://github.com/fal-ai/fal/actions) # fal -aaaa +aaaaa fal is a serverless Python runtime that lets you run and scale code in the cloud with no infra management. With fal, you can build pipelines, serve ML models and scale them up to many users. You scale down to 0 when you don't use any resources. From f20f76859841415190033dfa8cd2f219ff3c917d Mon Sep 17 00:00:00 2001 From: badayvedat Date: Sat, 5 Sep 2026 06:12:51 +0700 Subject: [PATCH 19/27] test: clean up e2e hardening changes --- README.md | 1 - projects/fal/tests/e2e/test_apps.py | 1 + t.py | 20 -------------------- 3 files changed, 1 insertion(+), 21 deletions(-) delete mode 100644 t.py diff --git a/README.md b/README.md index 8d9a92405..20faa29c9 100644 --- a/README.md +++ b/README.md @@ -2,7 +2,6 @@ [![Tests](https://img.shields.io/github/actions/workflow/status/fal-ai/fal/fal-unit-tests.yml?label=Tests)](https://github.com/fal-ai/fal/actions) # fal -aaaaa fal is a serverless Python runtime that lets you run and scale code in the cloud with no infra management. With fal, you can build pipelines, serve ML models and scale them up to many users. You scale down to 0 when you don't use any resources. diff --git a/projects/fal/tests/e2e/test_apps.py b/projects/fal/tests/e2e/test_apps.py index e12a43906..e149495cd 100644 --- a/projects/fal/tests/e2e/test_apps.py +++ b/projects/fal/tests/e2e/test_apps.py @@ -1725,6 +1725,7 @@ def test_realtime_connection(test_realtime_app): assert received_prompts == {str(prompt) for prompt in range(10)} assert sum(batch_sizes) == 10 assert all(1 <= batch_size <= 4 for batch_size in batch_sizes) + assert any(batch_size > 1 for batch_size in batch_sizes) @pytest.mark.xdist_group(name="realtime-app") diff --git a/t.py b/t.py deleted file mode 100644 index 24b00836b..000000000 --- a/t.py +++ /dev/null @@ -1,20 +0,0 @@ -import time -import uuid - -from fal.api import SyncServerlessClient - -client = SyncServerlessClient() - -for attempt in range(10): - name = f"replica-lag-repro-{uuid.uuid4().hex[:8]}" - print(f"testing {name}") - env = client.environments.create(name) - - assert name in [env.name for env in client.environments.list()] - - client.environments.delete(name) - - assert name not in [env.name for env in client.environments.list()] - - time.sleep(0.1) - print(f"Attempt {attempt + 1}/ 10 passed") From 9a6315d812c9fa2781485485cd56116e83112d6d Mon Sep 17 00:00:00 2001 From: badayvedat Date: Sun, 6 Sep 2026 17:22:10 +0700 Subject: [PATCH 20/27] test: retry queue submissions while aliases propagate --- projects/fal/tests/e2e/test_apps.py | 33 ++++- .../fal/tests/unit/test_e2e_queue_retry.py | 117 ++++++++++++++++++ 2 files changed, 149 insertions(+), 1 deletion(-) create mode 100644 projects/fal/tests/unit/test_e2e_queue_retry.py diff --git a/projects/fal/tests/e2e/test_apps.py b/projects/fal/tests/e2e/test_apps.py index e149495cd..284123e18 100644 --- a/projects/fal/tests/e2e/test_apps.py +++ b/projects/fal/tests/e2e/test_apps.py @@ -7,6 +7,7 @@ import time from contextlib import contextmanager, suppress from datetime import datetime, timedelta, timezone +from functools import partial from typing import ( AsyncIterator, Callable, @@ -198,15 +199,45 @@ def _is_alias_not_found_response(response: httpx.Response, app_alias: str) -> bo return False try: - detail = response.json().get("detail", "") + data = response.json() except ValueError: return False + if not isinstance(data, dict): + return False + detail = data.get("detail", "") + if not isinstance(detail, str): + return False return detail in { f"Application {app_alias!r} not found", f'Application "{app_alias}" not found', } +def _submit_with_alias_retry(submit, app_id: str, arguments: dict, *, path: str = ""): + # Temporary until queue gateways stop caching alias misses independently. + # Retry submission itself: a successful readiness probe can hit another pod. + normalized_id = apps._backwards_compatible_app_id(app_id) + app_alias = normalized_id.split("/")[1] if "/" in normalized_id else normalized_id + deadline = time.monotonic() + 60 + while True: + try: + return submit(app_id, arguments, path=path) + except HTTPStatusError as exc: + if not _is_alias_not_found_response(exc.response, app_alias): + raise + remaining = deadline - time.monotonic() + if remaining <= 0: + raise + time.sleep(min(0.5, remaining)) + if time.monotonic() >= deadline: + raise + + +@pytest.fixture(autouse=True) +def retry_queue_alias_submission(monkeypatch): + monkeypatch.setattr(apps, "submit", partial(_submit_with_alias_retry, apps.submit)) + + def _wait_for_queue_alias( app_alias: str, queue_url: str, diff --git a/projects/fal/tests/unit/test_e2e_queue_retry.py b/projects/fal/tests/unit/test_e2e_queue_retry.py new file mode 100644 index 000000000..3f2bc9dff --- /dev/null +++ b/projects/fal/tests/unit/test_e2e_queue_retry.py @@ -0,0 +1,117 @@ +import json +from types import SimpleNamespace + +import httpx +import pytest + +from fal import apps +from tests.e2e import test_apps as e2e_apps +from tests.e2e.test_apps import retry_queue_alias_submission # noqa: F401 + + +@pytest.fixture +def queue_client(monkeypatch): + responses = [] + requests = [] + + def handle(request): + requests.append(request) + return responses.pop(0) + + with httpx.Client(transport=httpx.MockTransport(handle)) as client: + monkeypatch.setattr(apps, "_get_http_client", lambda: client) + monkeypatch.setattr( + apps, "get_credentials", lambda: SimpleNamespace(to_headers=lambda: {}) + ) + now = [0.0] + + def sleep(seconds): + now[0] += seconds + + monkeypatch.setattr( + e2e_apps, "time", SimpleNamespace(monotonic=lambda: now[0], sleep=sleep) + ) + yield responses, requests, now + + +@pytest.mark.parametrize( + "app_id, path", + [ + ("owner/model", "/increment"), + ("owner/model/increment", ""), + ("owner-model", "increment"), + ], +) +def test_queue_run_retries_submission_only(queue_client, app_id, path): + responses, requests, now = queue_client + responses.extend( + [ + httpx.Response(404, json={"detail": "Application 'model' not found"}), + httpx.Response(404, json={"detail": 'Application "model" not found'}), + httpx.Response(200, json={"request_id": "accepted"}), + httpx.Response(200, json={"logs": []}), + httpx.Response(200, json={"result": 42}), + ] + ) + + assert apps.run(app_id, {"value": 1}, path=path) == {"result": 42} + assert [request.method for request in requests] == ["POST"] * 3 + ["GET"] * 2 + assert len({str(request.url) for request in requests[:3]}) == 1 + assert requests[0].url.path == "/owner/model/increment" + assert all(json.loads(request.content) == {"value": 1} for request in requests[:3]) + assert now[0] == 1 + + +@pytest.mark.parametrize( + "status, body", + [ + (404, {"detail": "Application 'different' not found"}), + (404, {"detail": "Not Found"}), + (404, {"detail": []}), + (404, []), + (404, "not json"), + (401, {"detail": "Application 'model' not found"}), + (429, {"error": "Queue is too long"}), + (500, {"detail": "Application 'model' not found"}), + ], +) +def test_queue_submission_preserves_other_errors(queue_client, status, body): + responses, requests, now = queue_client + response = ( + httpx.Response(status, text=body) + if isinstance(body, str) + else httpx.Response(status, json=body) + ) + responses.append(response) + with pytest.raises(httpx.HTTPStatusError) as exc: + apps.submit("owner/model", {}) + assert exc.value.response is response + assert len(requests) == 1 + assert now[0] == 0 + + +def test_queue_submission_retry_has_a_deadline(queue_client): + responses, requests, now = queue_client + responses.extend( + httpx.Response(404, json={"detail": "Application 'model' not found"}) + for _ in range(120) + ) + with pytest.raises(httpx.HTTPStatusError): + apps.submit("owner/model", {}) + assert now[0] == 60 + assert len(requests) == 120 + + +def test_queue_run_does_not_resubmit_after_acceptance(queue_client): + responses, requests, now = queue_client + responses.extend( + [ + httpx.Response(200, json={"request_id": "accepted"}), + httpx.Response(200, json={"logs": []}), + httpx.Response(404, json={"detail": "Application 'model' not found"}), + ] + ) + with pytest.raises(httpx.HTTPStatusError): + apps.run("owner/model", {}) + assert [request.method for request in requests] == ["POST", "GET", "GET"] + assert now[0] == 0 From 1ed8f0980c9a4f76dac41499be3e72ee73288618 Mon Sep 17 00:00:00 2001 From: badayvedat Date: Mon, 21 Sep 2026 16:41:33 -0700 Subject: [PATCH 21/27] chore: no-op random change --- README.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/README.md b/README.md index a16d5cf1f..aff95c849 100644 --- a/README.md +++ b/README.md @@ -1,6 +1,6 @@ [![PyPI](https://img.shields.io/pypi/v/fal.svg?logo=PyPI)](https://pypi.org/project/fal) [![Tests](https://img.shields.io/github/actions/workflow/status/fal-ai/fal/fal-unit-tests.yml?label=Tests)](https://github.com/fal-ai/fal/actions) -dummy change 1 +dummy change 2 # fal fal is a serverless Python runtime that lets you run and scale code in the cloud with no infra management. From 6e22256a8bf2706a976ae7936bc4453cc22b2130 Mon Sep 17 00:00:00 2001 From: badayvedat Date: Mon, 21 Sep 2026 17:55:25 -0700 Subject: [PATCH 22/27] chore: no-op random change --- README.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/README.md b/README.md index aff95c849..69ef6e512 100644 --- a/README.md +++ b/README.md @@ -1,6 +1,6 @@ [![PyPI](https://img.shields.io/pypi/v/fal.svg?logo=PyPI)](https://pypi.org/project/fal) [![Tests](https://img.shields.io/github/actions/workflow/status/fal-ai/fal/fal-unit-tests.yml?label=Tests)](https://github.com/fal-ai/fal/actions) -dummy change 2 +dummy change 3 # fal fal is a serverless Python runtime that lets you run and scale code in the cloud with no infra management. From 686195baa5d549d6d0dfb87a1c1cd36a071e255b Mon Sep 17 00:00:00 2001 From: badayvedat Date: Mon, 21 Sep 2026 18:20:06 -0700 Subject: [PATCH 23/27] chore: no-op random change --- README.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/README.md b/README.md index 69ef6e512..dcd53a822 100644 --- a/README.md +++ b/README.md @@ -1,6 +1,6 @@ [![PyPI](https://img.shields.io/pypi/v/fal.svg?logo=PyPI)](https://pypi.org/project/fal) [![Tests](https://img.shields.io/github/actions/workflow/status/fal-ai/fal/fal-unit-tests.yml?label=Tests)](https://github.com/fal-ai/fal/actions) -dummy change 3 +dummy change 4 # fal fal is a serverless Python runtime that lets you run and scale code in the cloud with no infra management. From db71aaa38d4d4f90ae0e86125ad9e038acec02d3 Mon Sep 17 00:00:00 2001 From: badayvedat Date: Mon, 21 Sep 2026 23:11:37 -0700 Subject: [PATCH 24/27] chore: no-op random change --- README.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/README.md b/README.md index dcd53a822..0f0b894ca 100644 --- a/README.md +++ b/README.md @@ -1,6 +1,6 @@ [![PyPI](https://img.shields.io/pypi/v/fal.svg?logo=PyPI)](https://pypi.org/project/fal) [![Tests](https://img.shields.io/github/actions/workflow/status/fal-ai/fal/fal-unit-tests.yml?label=Tests)](https://github.com/fal-ai/fal/actions) -dummy change 4 +dummy change 5 # fal fal is a serverless Python runtime that lets you run and scale code in the cloud with no infra management. From 9f3295e8fa0988ff1ab646bf54ca253f4d4b2d17 Mon Sep 17 00:00:00 2001 From: badayvedat Date: Mon, 21 Sep 2026 23:24:35 -0700 Subject: [PATCH 25/27] chore: no-op random change 6 --- README.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/README.md b/README.md index 0f0b894ca..679591a67 100644 --- a/README.md +++ b/README.md @@ -1,6 +1,6 @@ [![PyPI](https://img.shields.io/pypi/v/fal.svg?logo=PyPI)](https://pypi.org/project/fal) [![Tests](https://img.shields.io/github/actions/workflow/status/fal-ai/fal/fal-unit-tests.yml?label=Tests)](https://github.com/fal-ai/fal/actions) -dummy change 5 +dummy change 6 # fal fal is a serverless Python runtime that lets you run and scale code in the cloud with no infra management. From 53670debf1bad5d55dbdea77d518f8f32ddebd27 Mon Sep 17 00:00:00 2001 From: badayvedat Date: Tue, 22 Sep 2026 10:45:42 -0700 Subject: [PATCH 26/27] chore: no-op random change 7 --- README.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/README.md b/README.md index 679591a67..16fbcd373 100644 --- a/README.md +++ b/README.md @@ -1,6 +1,6 @@ [![PyPI](https://img.shields.io/pypi/v/fal.svg?logo=PyPI)](https://pypi.org/project/fal) [![Tests](https://img.shields.io/github/actions/workflow/status/fal-ai/fal/fal-unit-tests.yml?label=Tests)](https://github.com/fal-ai/fal/actions) -dummy change 6 +dummy change 7 # fal fal is a serverless Python runtime that lets you run and scale code in the cloud with no infra management. From 2ee4cd8507ca2b5fe6bd9ea174a7409dd21a230c Mon Sep 17 00:00:00 2001 From: badayvedat Date: Tue, 22 Sep 2026 11:52:57 -0700 Subject: [PATCH 27/27] test: decouple workflow alias retry from arithmetic inputs --- projects/fal/tests/e2e/test_apps.py | 45 ++++++++++++++++++- .../fal/tests/unit/test_e2e_queue_retry.py | 29 ++++++++++++ 2 files changed, 72 insertions(+), 2 deletions(-) diff --git a/projects/fal/tests/e2e/test_apps.py b/projects/fal/tests/e2e/test_apps.py index 4a594af30..064d87b05 100644 --- a/projects/fal/tests/e2e/test_apps.py +++ b/projects/fal/tests/e2e/test_apps.py @@ -1891,12 +1891,53 @@ def test_workflows(test_app: str, rest_client: Client): with delete_workflow_on_exit( client, rest_client.base_url + "/workflows/" + workflow_id ): - data = fal.apps.run( - "workflows/" + workflow_id, arguments={"lhs": 2, "rhs": 3} + # Replaying this arithmetic workflow is safe if a node misses the alias. + data = _retry_workflow_alias_miss( + lambda: fal.apps.run( + "workflows/" + workflow_id, arguments={"lhs": 2, "rhs": 3} + ), + app_id=test_app, ) assert data["result"] == 10 +def _retry_workflow_alias_miss(run: Callable[[], T], *, app_id: str) -> T: + # Temporary until workflow node submissions consistently resolve new aliases. + # Callers must ensure the entire operation is safe to replay. + app_alias = app_id.split("/")[1] + deadline = time.monotonic() + 60 + while True: + try: + return run() + except HTTPStatusError as exc: + try: + data = exc.response.json() + except ValueError: + data = None + if exc.response.status_code != 404 or not isinstance(data, dict): + raise + error = data.get("error") + if ( + data.get("type") != "error" + or data.get("message") != f"Error while running app {app_id!r}" + or not isinstance(error, dict) + or error.get("status") != 404 + ): + raise + body = error.get("body") + if not isinstance(body, dict) or body.get("detail") not in ( + f"Application {app_alias!r} not found", + f'Application "{app_alias}" not found', + ): + raise + remaining = deadline - time.monotonic() + if remaining <= 0: + raise + time.sleep(min(0.5, remaining)) + if time.monotonic() >= deadline: + raise + + @pytest.mark.xdist_group(name="exception-app") def test_app_exceptions(test_exception_app: AppClient): with pytest.raises(AppClientError) as app_exc: diff --git a/projects/fal/tests/unit/test_e2e_queue_retry.py b/projects/fal/tests/unit/test_e2e_queue_retry.py index 3f2bc9dff..78f7231fb 100644 --- a/projects/fal/tests/unit/test_e2e_queue_retry.py +++ b/projects/fal/tests/unit/test_e2e_queue_retry.py @@ -115,3 +115,32 @@ def test_queue_run_does_not_resubmit_after_acceptance(queue_client): apps.run("owner/model", {}) assert [request.method for request in requests] == ["POST", "GET", "GET"] assert now[0] == 0 + + +def test_workflow_retry_has_a_deadline(queue_client): + error = { + "type": "error", + "message": "Error while running app 'owner/model'", + "error": { + "status": 404, + "body": {"detail": 'Application "model" not found'}, + }, + } + responses, requests, now = queue_client + for index in range(120): + responses.extend( + [ + httpx.Response(200, json={"request_id": str(index)}), + httpx.Response(200, json={"logs": []}), + httpx.Response(404, json=error), + ] + ) + last_response = responses[-1] + with pytest.raises(httpx.HTTPStatusError) as exc: + e2e_apps._retry_workflow_alias_miss( + lambda: apps.run("workflows/owner/workflow", {"lhs": 2, "rhs": 3}), + app_id="owner/model", + ) + assert exc.value.response is last_response + assert now[0] == 60 + assert len(requests) == 360