Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
__init__.py81 linesDownload Raw Back to distributed
1from abc import abstractmethod
2from dataclasses import dataclass
3from typing import Any, Callable, List
4
5from overrides import EnforceOverrides, overrides
6from chromadb.config import Component, System
7from chromadb.types import Segment
8
9
10class SegmentDirectory(Component):
11    """A segment directory is a data interface that manages the location of segments. Concretely, this
12    means that for distributed chroma, it provides the grpc endpoint for a segment."""
13
14    @abstractmethod
15    def get_segment_endpoints(self, segment: Segment, n: int) -> List[str]:
16        """Return the segment residences for a given segment ID. Will return at most n residences.
17        Should only return less than n residences if there are less than n residences available.
18        """
19
20    @abstractmethod
21    def register_updated_segment_callback(
22        self, callback: Callable[[Segment], None]
23    ) -> None:
24        """Register a callback that will be called when a segment is updated"""
25        pass
26
27
28@dataclass
29class Member:
30    id: str
31    ip: str
32    node: str
33
34
35Memberlist = List[Member]
36
37
38class MemberlistProvider(Component, EnforceOverrides):
39    """Returns the latest memberlist and provdes a callback for when it changes. This
40    callback may be called from a different thread than the one that called. Callers should ensure
41    that they are thread-safe."""
42
43    callbacks: List[Callable[[Memberlist], Any]]
44
45    def __init__(self, system: System):
46        self.callbacks = []
47        super().__init__(system)
48
49    @abstractmethod
50    def get_memberlist(self) -> Memberlist:
51        """Returns the latest memberlist"""
52        pass
53
54    @abstractmethod
55    def set_memberlist_name(self, memberlist: str) -> None:
56        """Sets the memberlist that this provider will watch"""
57        pass
58
59    @overrides
60    def stop(self) -> None:
61        """Stops watching the memberlist"""
62        self.callbacks = []
63
64    def register_updated_memberlist_callback(
65        self, callback: Callable[[Memberlist], Any]
66    ) -> None:
67        """Registers a callback that will be called when the memberlist changes. May be called many times
68        with the same memberlist, so callers should be idempotent. May be called from a different thread.
69        """
70        self.callbacks.append(callback)
71
72    def unregister_updated_memberlist_callback(
73        self, callback: Callable[[Memberlist], Any]
74    ) -> bool:
75        """Unregisters a callback that was previously registered. Returns True if the callback was
76        successfully unregistered, False if it was not ever registered."""
77        if callback in self.callbacks:
78            self.callbacks.remove(callback)
79            return True
80        return False
81 
codekingpro/portable-devtools · Team Ai