Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
67 changes: 32 additions & 35 deletions servers/lean_runtime.py
Original file line number Diff line number Diff line change
Expand Up @@ -339,23 +339,44 @@ def lease(
self,
project_dir: str,
*,
create: bool = True,
acquisition_timeout: float | None = None,
creation_budget: float = 0.0,
) -> Iterator[T | None]:
) -> Iterator[T]:
"""Keep a project resource alive for the complete operation."""
root = resolve_lean_project_dir(project_dir)
resource = self._acquire(
root,
create=create,
acquisition_timeout=acquisition_timeout,
creation_budget=creation_budget,
)
try:
yield resource
finally:
self._release(root, resource)

@contextmanager
def observe(self, project_dir: str) -> Iterator[tuple[T | None, str]]:
"""Borrow one atomic state snapshot without validating or refreshing it."""
root = resolve_lean_project_dir(project_dir)
with self._condition:
if self._closed:
raise RuntimeError("project resource cache is closed")
entry = self._entries.get(root)
if entry is None:
resource = None
state = "warming" if root in self._creating else "cold"
else:
entry.active += 1
resource = entry.resource
state = "warm"
try:
yield resource, state
finally:
if resource is not None:
self._release(root, resource)
with self._condition:
# A pinned entry is never evicted, replaced, or closed.
self._entries[root].active -= 1
self._condition.notify_all()

def stats(self) -> dict[str, Any]:
with self._condition:
Expand All @@ -376,16 +397,6 @@ def stats(self) -> dict[str, Any]:
"creating": sorted(str(root) for root in self._creating),
}

def state(self, project_dir: str) -> str:
"""Return ``cold``, ``warming``, or ``warm`` without creating state."""
root = resolve_lean_project_dir(project_dir)
with self._condition:
if root in self._entries:
return "warm"
if root in self._creating:
return "warming"
return "cold"

def invalidate(self, project_dir: str, resource: T) -> None:
"""Arrange to replace a failed resource after its active calls finish."""
root = resolve_lean_project_dir(project_dir)
Expand Down Expand Up @@ -430,10 +441,9 @@ def _acquire(
self,
root: Path,
*,
create: bool,
acquisition_timeout: float | None,
creation_budget: float,
) -> T | None:
) -> T:
if acquisition_timeout is not None and acquisition_timeout <= 0:
raise ProjectResourceBusyError(
"no response budget remains for a shared Lean project slot"
Expand Down Expand Up @@ -469,22 +479,14 @@ def _acquire(
)
if entry_is_stale:
assert entry is not None
self._require_creation_budget(
root,
deadline=deadline,
creation_budget=creation_budget,
)
if entry.active:
if not create:
resource = None
break
self._require_creation_budget(
root,
deadline=deadline,
creation_budget=creation_budget,
)
wait = True
else:
self._require_creation_budget(
root,
deadline=deadline,
creation_budget=creation_budget,
)
resources_to_close.append(self._entries.pop(root).resource)
self._condition.notify_all()
entry = None
Expand All @@ -499,10 +501,6 @@ def _acquire(
resource = entry.resource
break

if not wait and entry is None and not create:
resource = None
break

if not wait and root in self._creating:
wait = True

Expand Down Expand Up @@ -716,8 +714,7 @@ def dispatch(self, method: str, params: dict[str, Any]) -> Any:
return format_repl_response(pool.run(code, timeout=effective_timeout))
if method == "repl.status":
project_dir = self._string_param(params, "project_dir")
with self.repl_projects.lease(project_dir, create=False) as pool:
state = "warm" if pool is not None else self.repl_projects.state(project_dir)
with self.repl_projects.observe(project_dir) as (pool, state):
return {
"state": state,
"capacity": (
Expand Down
Loading
Loading