gms.py 8.35 KB
Newer Older
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
# SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved.
# SPDX-License-Identifier: Apache-2.0

from __future__ import annotations

import asyncio
import os
import shutil
import threading
import time
from concurrent.futures import TimeoutError as FutureTimeoutError

from gpu_memory_service.client.rpc import _GMSRPCTransport
from gpu_memory_service.common.protocol.messages import (
    GetEventHistoryRequest,
    GetEventHistoryResponse,
    GetRuntimeStateRequest,
    GetRuntimeStateResponse,
)
from gpu_memory_service.common.utils import get_socket_path
from gpu_memory_service.server.rpc import GMSRPCServer

from tests.utils.managed_process import ManagedProcess

from .runtime import DYNAMO_BIN


def _socket_has_live_gms(socket_path: str) -> bool:
    if not os.path.exists(socket_path):
        return False
    try:
        with _GMSRPCTransport(socket_path) as transport:
            transport.connect()
            transport.request(
                GetRuntimeStateRequest(),
                GetRuntimeStateResponse,
            )
    except Exception:
        return False
    return True


def _prepare_socket_path_for_launch(socket_path: str) -> None:
    if not os.path.exists(socket_path):
        return
    if _socket_has_live_gms(socket_path):
        raise RuntimeError(f"GMS already active at {socket_path}")
    os.unlink(socket_path)


class GMSServerProcess(ManagedProcess):
    def __init__(self, request, device: int, tag: str = "weights"):
        self.device = device
        self.tag = tag
        self.socket_path = get_socket_path(device, tag)

        log_dir = f"{request.node.name}_gms_{tag}_{device}"
        shutil.rmtree(log_dir, ignore_errors=True)

        super().__init__(
            command=[
                "python",
                "-m",
                "gpu_memory_service",
                "--device",
                str(device),
                "--tag",
                tag,
            ],
            env={
                **os.environ,
                "PATH": f"{DYNAMO_BIN}:{os.environ.get('PATH', '')}",
                "DYN_LOG": "debug",
            },
            timeout=60,
            display_output=True,
            terminate_all_matching_process_names=False,
            log_dir=log_dir,
            display_name=f"gms_{tag}",
            health_check_funcs=[self._runtime_state_ready],
        )

    def __enter__(self):
        _prepare_socket_path_for_launch(self.socket_path)
        return super().__enter__()

    def __exit__(self, exc_type, exc_val, exc_tb):
        try:
            return super().__exit__(exc_type, exc_val, exc_tb)
        finally:
            if os.path.exists(self.socket_path) and not _socket_has_live_gms(
                self.socket_path
            ):
                os.unlink(self.socket_path)

    def _socket_has_live_gms(self) -> bool:
        return _socket_has_live_gms(self.socket_path)

    def _prepare_socket_path_for_launch(self) -> None:
        _prepare_socket_path_for_launch(self.socket_path)

    def _runtime_state_ready(self, timeout: float = 30) -> bool:
        deadline = time.monotonic() + timeout
        while time.monotonic() < deadline:
            if not os.path.exists(self.socket_path):
                time.sleep(0.1)
                continue
            try:
                self.get_runtime_state()
                return True
            except Exception:
                time.sleep(0.1)
                continue
            time.sleep(0.1)
        return False

    def get_runtime_state(self) -> GetRuntimeStateResponse:
        with _GMSRPCTransport(self.socket_path) as transport:
            transport.connect()
            return transport.request(
                GetRuntimeStateRequest(),
                GetRuntimeStateResponse,
            )

    def get_event_history(self) -> GetEventHistoryResponse:
        with _GMSRPCTransport(self.socket_path) as transport:
            transport.connect()
            return transport.request(
                GetEventHistoryRequest(),
                GetEventHistoryResponse,
            )


class ServerThread:
    def __init__(self, server, socket_path: str):
        self.server = server
        self.socket_path = socket_path
        self._loop: asyncio.AbstractEventLoop | None = None
        self._task: asyncio.Task[None] | None = None
        self._thread = threading.Thread(target=self._run, daemon=True)
        self._exception: BaseException | None = None

    def _run(self) -> None:
        loop = asyncio.new_event_loop()
        self._loop = loop
        asyncio.set_event_loop(loop)
        self._task = loop.create_task(self.server.serve())
        try:
            loop.run_until_complete(self._task)
        except asyncio.CancelledError:
            pass
        except BaseException as exc:
            self._exception = exc
        finally:
            pending = [task for task in asyncio.all_tasks(loop) if not task.done()]
            for task in pending:
                task.cancel()
            if pending:
                loop.run_until_complete(
                    asyncio.gather(*pending, return_exceptions=True)
                )
            loop.close()

    def start(self) -> None:
        self._thread.start()
        deadline = time.monotonic() + 5.0
        while True:
            if self._exception is not None:
                raise self._exception
            if self.server._server is not None and os.path.exists(self.socket_path):
                try:
                    with _GMSRPCTransport(self.socket_path) as transport:
                        transport.connect()
                        transport.request(
                            GetRuntimeStateRequest(),
                            GetRuntimeStateResponse,
                        )
                    return
                except Exception:
                    pass
            if time.monotonic() > deadline:
                raise TimeoutError(f"GMS socket did not appear at {self.socket_path}")
            time.sleep(0.01)

    def stop(self) -> None:
        if self._loop is not None:

            def cancel() -> None:
                if self.server._server is not None:
                    self.server._server.close()
                if self._task is not None:
                    self._task.cancel()

            self._loop.call_soon_threadsafe(cancel)
        self._thread.join(timeout=5)
        if self._exception is not None:
            raise self._exception
        if os.path.exists(self.socket_path):
            os.unlink(self.socket_path)

    def disconnect_rw_session(self, timeout: float = 5.0) -> None:
        if self._loop is None:
            raise RuntimeError("GMS server thread is not running")
        future = asyncio.run_coroutine_threadsafe(
            self._disconnect_rw_session(), self._loop
        )
        try:
            future.result(timeout=timeout)
        except FutureTimeoutError as exc:
            raise TimeoutError("Timed out disconnecting RW session") from exc

    async def _disconnect_rw_session(self) -> None:
        conn = self.server._gms._sessions._locking.rw_conn
        if conn is None:
            raise RuntimeError("No active RW session to disconnect")
        await self.server._gms.cleanup_connection(conn)


class ThreadedGMSServer:
    def __init__(self, device: int, tag: str = "weights"):
        self.device = device
        self.tag = tag
        self.socket_path = get_socket_path(device, tag)
        self.server = GMSRPCServer(self.socket_path, device)
        self._thread = ServerThread(self.server, self.socket_path)

    def __enter__(self) -> "ThreadedGMSServer":
        _prepare_socket_path_for_launch(self.socket_path)
        self._thread.start()
        return self

    def __exit__(self, exc_type, exc_val, exc_tb) -> None:
        self._thread.stop()

    def get_runtime_state(self) -> GetRuntimeStateResponse:
        with _GMSRPCTransport(self.socket_path) as transport:
            transport.connect()
            return transport.request(
                GetRuntimeStateRequest(),
                GetRuntimeStateResponse,
            )

    def get_event_history(self) -> GetEventHistoryResponse:
        with _GMSRPCTransport(self.socket_path) as transport:
            transport.connect()
            return transport.request(
                GetEventHistoryRequest(),
                GetEventHistoryResponse,
            )

    def disconnect_rw_session(self) -> None:
        self._thread.disconnect_rw_session()