codekingpro/portable-devtools
114k
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 