{"metadata":{"kernelspec":{"display_name":"Python 3","language":"python","name":"python3"},"language_info":{"name":"python","version":"3.10.12","mimetype":"text/x-python","codemirror_mode":{"name":"ipython","version":3},"pygments_lexer":"ipython3","nbconvert_exporter":"python","file_extension":".py"},"kaggle":{"accelerator":"nvidiaL4","dataSources":[{"sourceId":84795,"databundleVersionId":10934030,"sourceType":"competition"}],"isInternetEnabled":false,"language":"python","sourceType":"notebook","isGpuEnabled":true},"papermill":{"default_parameters":{},"duration":22.341371,"end_time":"2024-12-11T03:22:13.479076","environment_variables":{},"exception":null,"input_path":"__notebook__.ipynb","output_path":"__notebook__.ipynb","parameters":{},"start_time":"2024-12-11T03:21:51.137705","version":"2.6.0"}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"code","source":"import io\nimport os\nimport shutil\nimport subprocess\n\nimport pandas as pd\nimport polars as pl\n\nimport kaggle_evaluation.konwinski_prize_inference_server","metadata":{"_cell_guid":"b1076dfc-b9ad-4769-8c92-a6c4dae69d19","_uuid":"8f2839f25d086af736a60e9eeb907d3b93b6e0e5","execution":{"iopub.status.busy":"2025-02-27T16:52:17.269591Z","iopub.execute_input":"2025-02-27T16:52:17.269943Z","iopub.status.idle":"2025-02-27T16:52:18.808632Z","shell.execute_reply.started":"2025-02-27T16:52:17.269916Z","shell.execute_reply":"2025-02-27T16:52:18.80756Z"},"papermill":{"duration":14.873526,"end_time":"2024-12-11T03:22:08.818755","exception":false,"start_time":"2024-12-11T03:21:53.945229","status":"completed"},"tags":[],"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"The evaluation API requires that you set up a server which will respond to inference requests. We have already defined the server; you just need write the predict function. When we evaluate your submission on the hidden test set the client defined in `konwinski_prize_gateway` will run in a different container with direct access to the hidden test set and hand off the data.\n\nYour code will always have access to the published copies of the files.","metadata":{"papermill":{"duration":0.002032,"end_time":"2024-12-11T03:22:08.823897","exception":false,"start_time":"2024-12-11T03:22:08.821865","status":"completed"},"tags":[]}},{"cell_type":"code","source":"instance_count = None\n\ndef get_number_of_instances(num_instances: int) -> None:\n    \"\"\" The very first message from the gateway will be the total number of instances to be served.\n    You don't need to edit this function.\n    \"\"\"\n    global instance_count\n    instance_count = num_instances","metadata":{"execution":{"iopub.status.busy":"2025-02-27T16:52:18.810337Z","iopub.execute_input":"2025-02-27T16:52:18.811049Z","iopub.status.idle":"2025-02-27T16:52:18.816069Z","shell.execute_reply.started":"2025-02-27T16:52:18.810985Z","shell.execute_reply":"2025-02-27T16:52:18.814831Z"},"papermill":{"duration":0.011949,"end_time":"2024-12-11T03:22:08.838279","exception":false,"start_time":"2024-12-11T03:22:08.82633","status":"completed"},"tags":[],"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"import multiprocessing\nfrom typing import Any, Callable\nimport time\nimport os\nimport signal\nfrom contextlib import contextmanager\n\nclass TimeoutError(Exception):\n    pass\n\n@contextmanager\ndef time_limit(seconds):\n    \"\"\"Context manager that raises TimeoutError if the code inside takes longer than specified seconds.\"\"\"\n    def signal_handler(signum, frame):\n        raise TimeoutError(\"Timed out!\")\n    \n    # Set the signal handler and a alarm\n    signal.signal(signal.SIGALRM, signal_handler)\n    signal.alarm(int(seconds))\n    try:\n        yield\n    finally:\n        # Cancel the alarm\n        signal.alarm(0)\n\ndef run_with_timeout(\n    func: Callable[..., Any], \n    timeout: float | int, \n    *args: Any, \n    max_poll_interval: float = 1.0,  # New parameter to control maximum polling interval\n    **kwargs: Any\n) -> Any | None:\n    \"\"\"Run a function in a separate process with a guaranteed timeout.\n    \n    This uses a modified multiprocessing approach that is container-friendly and \n    reliably terminates the function after the timeout period.\n    \n    Args:\n        func (Callable[..., Any]): The target function to execute.\n        timeout (float | int): Maximum allowed execution time in seconds.\n        *args (Any): Positional arguments to pass to `func`.\n        max_poll_interval (float): Maximum time between polling checks (default: 1.0 seconds).\n        **kwargs (Any): Keyword arguments to pass to `func`.\n    \n    Returns:\n        Any | None: \n            The result of `func(*args, **kwargs)` if it completes within `timeout` seconds; \n            Otherwise, None.\n    \"\"\"\n    t1 = time.time()\n    # Create a pipe for communication\n    parent_conn, child_conn = multiprocessing.Pipe()\n    \n    # Use start method that works well in containers\n    ctx = multiprocessing.get_context('spawn')\n    p = ctx.Process(\n        target=_process_worker,\n        args=(child_conn, func, args, kwargs)\n    )\n    \n    # Start process\n    p.start()\n    \n    # Set a timeout for receiving data\n    start_time = time.time()\n    result = None\n    received = False\n    \n    # Start with a short polling interval that gradually increases\n    poll_interval = 0.1  # Initial polling interval\n    sleep_time = 0.01    # Initial sleep time\n    \n    while time.time() - start_time < timeout and not received:\n        if parent_conn.poll(poll_interval):  # Check with adaptive timeout\n            result = parent_conn.recv()\n            received = True\n        \n        # Gradually increase polling interval for efficiency, up to max_poll_interval\n        poll_interval = min(poll_interval * 1.5, max_poll_interval)\n        \n        # Also gradually increase sleep time, up to 1/10th of poll_interval\n        sleep_time = min(sleep_time * 1.5, poll_interval / 10)\n        \n        time.sleep(sleep_time)  # Adaptive sleep to prevent CPU spinning\n    \n    # Make sure to terminate the process if running\n    if p.is_alive():\n        p.terminate()\n        p.join(1.0)  # Give it a second to clean up\n        # Force kill if still alive\n        if p.is_alive():\n            os.kill(p.pid, signal.SIGKILL)\n    \n    print(f\"TERMINATING USING MANUAL TIMEOUT AFTER: {time.time()-t1:.2f} SECONDS\")\n    return result\n\ndef _process_worker(conn, func, args, kwargs):\n    \"\"\"Worker function that runs in the separate process.\"\"\"\n    try:\n        # Additional safety: use signal alarm as a backup timeout mechanism\n        with time_limit(60*29):  # 29 minute hard limit as an extra safeguard\n            result = func(*args, **kwargs)\n            conn.send(result)\n    except Exception as e:\n        # If any exception occurs, just return None\n        pass\n    finally:\n        conn.close()","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-02-27T16:52:18.817821Z","iopub.execute_input":"2025-02-27T16:52:18.818165Z","iopub.status.idle":"2025-02-27T16:52:18.836737Z","shell.execute_reply.started":"2025-02-27T16:52:18.818139Z","shell.execute_reply":"2025-02-27T16:52:18.835502Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"import time\nfirst_prediction = True\n\nACTUAL_TIMEOUT_ERROR_LIMIT = 60*30  # 30 minutes\nOUR_TIMEOUT_LIMIT = 60*1  # 1 minute timeout\n\ndef inner_predict(random_arg_1: Any, random_arg_2: int = 2):\n    \"\"\"Included some args just for demonstration purposes\"\"\"\n    print(random_arg_1)\n    time.sleep(ACTUAL_TIMEOUT_ERROR_LIMIT*2)  # This function will run for 60 minutes if uninterrupted.\n    print(random_arg_2)\n    return \"**yawn** ... I slept so long ...\"\n    \ndef predict(problem_statement: str, repo_archive: io.BytesIO, pip_packages_archive: io.BytesIO, env_setup_cmds_templates: list[str]) -> str:\n    \"\"\" Replace this function with your inference code.\n    Args:\n        problem_statement: The text of the git issue.\n        repo_path: A BytesIO buffer path with a .tar containing the codebase that must be patched. The gateway will make this directory available immediately before this function runs.\n        pip_packages_archive: A BytesIO buffer path with a .tar containing the wheel files necessary for running unit tests.\n        env_setup_cmds_templates: Commands necessary for installing the pip_packages_archive.\n    \"\"\"\n    global first_prediction\n    if not first_prediction:\n        return None  # Skip issue.\n\n    # Unpack the codebase to be patched into a directory that won't be exported when\n    # the notebook is saved.\n    archive_path = '/tmp/repo_archive.tar'\n    with open(archive_path, 'wb') as f:\n        f.write(repo_archive.read())\n    repo_path = 'repo'\n    if os.path.exists(repo_path):\n        shutil.rmtree(repo_path)\n    shutil.unpack_archive(archive_path, extract_dir=repo_path)\n    os.remove(archive_path)\n\n    \"\"\"\n    Unpack pip_packages if you want to run unit tests on your patch.\n    Note that editing unit tests with your patch -- even to add valid tests -- can cause your submission to be flagged as a failure.\n    Most of the relevant repos use pytest for running tests. You will almost certainly need to run only a subset of the unit tests to avoid running out of inference time.\n    \"\"\"\n    pip_archive_dir = '/tmp/pip_packages_archive.tar'\n    with open(pip_archive_dir, 'wb') as f:\n        f.write(pip_packages_archive.read())\n    pip_packages_path = '/path/to/pip_packages'\n    if os.path.exists(pip_packages_path):\n        shutil.rmtree(pip_packages_path)\n    shutil.unpack_archive(pip_archive_dir, extract_dir=pip_packages_path)\n    os.remove(pip_archive_dir)\n\n    # Get env setup cmds by setting the pip_packages_path\n    env_setup_cmds = [cmd.format(pip_packages_path=pip_packages_path) for cmd in env_setup_cmds_templates]\n\n    # Run env setup for the repo\n    subprocess.run(\n        \"\\n\".join(env_setup_cmds),\n        shell=True,\n        executable=\"/bin/bash\",\n        cwd=repo_path,\n    )\n\n    first_prediction = False\n    # Instead of a valid diff, let's just submit a generic string. This will definitely fail.\n    prediction_result = run_with_timeout(inner_predict, timeout=OUR_TIMEOUT_LIMIT, random_arg_1=\"hello\", random_arg_2=7)\n    return prediction_result","metadata":{"execution":{"iopub.status.busy":"2025-02-27T16:52:18.837933Z","iopub.execute_input":"2025-02-27T16:52:18.838369Z","iopub.status.idle":"2025-02-27T16:52:18.857652Z","shell.execute_reply.started":"2025-02-27T16:52:18.838331Z","shell.execute_reply":"2025-02-27T16:52:18.856433Z"},"papermill":{"duration":0.011382,"end_time":"2024-12-11T03:22:08.852112","exception":false,"start_time":"2024-12-11T03:22:08.84073","status":"completed"},"tags":[],"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"When your notebook is run on the hidden test set, inference_server.serve must be called within 15 minutes of the notebook starting or the gateway will throw an error. If you need more than 15 minutes to load your model you can do so during the very first predict call, which does not have the usual 30 minute response deadline.","metadata":{"papermill":{"duration":0.001889,"end_time":"2024-12-11T03:22:08.856283","exception":false,"start_time":"2024-12-11T03:22:08.854394","status":"completed"},"tags":[]}},{"cell_type":"code","source":"inference_server = kaggle_evaluation.konwinski_prize_inference_server.KPrizeInferenceServer(\n    get_number_of_instances,   \n    predict\n)\n\nif os.getenv('KAGGLE_IS_COMPETITION_RERUN'):\n    inference_server.serve()\nelse:\n    inference_server.run_local_gateway(\n        data_paths=(\n            '/kaggle/input/konwinski-prize/',  # Path to the entire competition dataset\n            '/kaggle/tmp/konwinski-prize/',   # Path to a scratch directory for unpacking data.a_zip.\n        ),\n        use_concurrency=True,  # This can safely be disabled for purposes of local testing if necessary.\n    )","metadata":{"execution":{"iopub.status.busy":"2025-02-27T16:52:18.85873Z","iopub.execute_input":"2025-02-27T16:52:18.859142Z","iopub.status.idle":"2025-02-27T16:53:29.555244Z","shell.execute_reply.started":"2025-02-27T16:52:18.859102Z","shell.execute_reply":"2025-02-27T16:53:29.554103Z"},"papermill":{"duration":3.790202,"end_time":"2024-12-11T03:22:12.648591","exception":false,"start_time":"2024-12-11T03:22:08.858389","status":"completed"},"tags":[],"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null}]}