Created
August 14, 2020 12:12
-
-
Save jschmid1/1b955d2b9feac4f64faf3dfd1a9f477d to your computer and use it in GitHub Desktop.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| import json | |
| import pprint | |
| import time | |
| import yaml | |
| from typing import * | |
| def _load_inventory_struct(): | |
| with open('inventory_schema.yaml', 'rb') as _fd: | |
| return yaml.safe_load(_fd) | |
| pp = pprint.PrettyPrinter(indent=2) | |
| store_struct = _load_inventory_struct() | |
| pp.pprint(store_struct) | |
| class Store: | |
| """ | |
| TODO: Maybe implement a database like interface for the mgr_store.. | |
| CRUD could be helpful here | |
| """ | |
| def __init__(self, mgr, prefix, version=1): | |
| self.mgr = mgr | |
| self.prefix = prefix | |
| self.version = version | |
| def save(self, ident=None, data=None, version=None): | |
| """ | |
| Save `blob` to store. | |
| `self.prefix` is usually on of ['host', 'spec.'] | |
| `ident` can be the name of a host or a spec. | |
| Example: | |
| `self.mgr.set_store("spec.host", {'foo': 'bar'})` | |
| raise Exception if an issue is encountered | |
| // There are trade-offs between having a flat vs. non-flat hierarchy. | |
| // I do however think that having a non-flat hierarchy will benefit | |
| // us in this scenario. This may reflect in the need for a lower amount | |
| // of auxiliary functions. (map host to daemon, sort by type etc) | |
| // as everything is in one place. | |
| A downside of this may be partial updates since the underlying mon_store | |
| implementation doesn't have the concept of partial writes. The whole `inventory` | |
| blob would be re-written if anything changes. | |
| """ | |
| assert ident, data | |
| old = self.mgr.get_store_prefix(prefix=self.prefix, version=version) | |
| if self.has_changes(old, data): | |
| print("Changes detected. Writing to store") | |
| return self.mgr.set_store(f"{self.prefix}.{ident}", self.jsonify(data), version=self.version) | |
| print("No changes. Not writing to the store") | |
| return True | |
| @staticmethod | |
| def jsonify(data): | |
| try: | |
| return json.dumps(data) | |
| except json.JSONDecodeError: | |
| # TODO: catch correct json errors | |
| raise Exception("Error encoding json") | |
| @staticmethod | |
| def has_changes(new, old): | |
| """ | |
| If we only write whole data blobs we can easily compute if we need to write | |
| to the store. Computing the *actual* diff is a bit more complicated. | |
| """ | |
| print("Comparing existing data with new.") | |
| return new == old | |
| def load(self) -> Generator[str, Dict[Any, Any], Any]: | |
| """ | |
| Load from the mon_store. | |
| This should only happen on object creation (ceph-mgr startup). | |
| """ | |
| for k, v in self.mgr.get_store_prefix(self.prefix, version=self.version).items(): | |
| yield k, v | |
| def migrate(self, from_v, to_v): | |
| # TODO! | |
| assert from_v > to_v | |
| # TODO: Maybe s/Component/InventoryItem/g | |
| class Component: | |
| # Base fields that all components have in common | |
| loadable_base_fields = ['key1'] | |
| # Loadable fields are supposed to hold attributes that are | |
| # good to be persisted to and loaded from the store. | |
| # These should not contain dynamic data such as `daemon.running` or similar | |
| # data that could lead to issues if getting stale. | |
| loadable_fields = [] + loadable_base_fields | |
| def __init__(self, **kwargs): | |
| print(f"Loading kwargs -> {kwargs}") | |
| for k, v in kwargs.items(): | |
| if k not in self.loadable_fields: | |
| print(f"field <{k}> is not supported in class <{self.__class__.__name__}>") | |
| continue | |
| self.__setattr__(k, v) | |
| @classmethod | |
| def from_json(cls, data) -> 'Component': | |
| print(f"Loading {cls} from json") | |
| return cls(**{ | |
| key: data.get(key, None) | |
| for key in cls.loadable_fields | |
| }) | |
| def to_json(self) -> dict: | |
| # to_json all fields that are in loadable_fields and are not None (reduces the clutter) | |
| return {field: value for (field, value) in self.__dict__.items() | |
| if all([field in self.loadable_fields, field is not None])} | |
| def __eq__(self, other): | |
| return (self.__class__ == other.__class__ and | |
| self.to_json() == other.to_json) | |
| def __repr__(self): | |
| printable_fields = {k: v for (k, v) in self.__dict__.items() if k in self.loadable_fields} | |
| return f"{self.__class__.__name__}({printable_fields})" | |
| class DaemonDescription(Component): | |
| loadable_fields = ['daemon_id', 'daemon_type', 'etc'] + Component.loadable_base_fields | |
| def __init__(self, **kwargs): | |
| super(DaemonDescription, self).__init__(**kwargs) | |
| class Network(Component): | |
| loadable_fields = ['address', 'port', 'subnet'] + Component.loadable_base_fields | |
| def __init__(self, **kwargs): | |
| super(Network, self).__init__(**kwargs) | |
| class Attribute(Component): | |
| loadable_fields = ['cpu', 'ram', 'os'] + Component.loadable_base_fields | |
| def __init__(self, **kwargs): | |
| super(Attribute, self).__init__(**kwargs) | |
| class Device(Component): | |
| loadable_fields = ['path', 'rotational', 'model'] + Component.loadable_base_fields | |
| def __init__(self, **kwargs): | |
| super(Device, self).__init__(**kwargs) | |
| class Config(Component): | |
| loadable_fields = ['foo', 'bar', 'baz'] + Component.loadable_base_fields | |
| def __init__(self, **kwargs): | |
| super(Config, self).__init__(**kwargs) | |
| # TODO: It's not only building Components(InventoryItems) but also operating | |
| # in bulk on them (to_json). The name doesn't reflect what's the purpose | |
| # of this class is right now. Rename this class | |
| class ComponentBuilder: | |
| """ | |
| This class is supposed to act as an interface between the mon_store and | |
| the object representation of it that will be used in cephadm. | |
| It translates (bidirectionally) mon_store representation (json) to python objects and vice versa. | |
| """ | |
| component_map = dict(networks=Network, | |
| daemons=DaemonDescription, | |
| attributes=Attribute, | |
| devices=Device, | |
| config=Config) | |
| @staticmethod | |
| def to_json(components: List[Component]): | |
| if isinstance(components, list): | |
| # This makes sure that to_json exists for all objects | |
| assert all([isinstance(component, Component) for component in components]) | |
| return [d.to_json() for d in components] | |
| if isinstance(components, Component): | |
| return components.to_json() | |
| @classmethod | |
| def from_json(cls, data: Any, component_name: str): | |
| component_class = cls.component_map.get(component_name) | |
| if not component_class: | |
| raise Exception(f"Do not have a base class for {component_name}") | |
| if isinstance(data, dict): | |
| return component_class.from_json(data) | |
| if isinstance(data, list): | |
| return [component_class.from_json(i) for i in data] | |
| raise Exception(f"Can't handle type {type(data)}") | |
| class Inventory: | |
| """ | |
| Contains a List of Hosts that are registered to the system | |
| """ | |
| storable_components = ['networks', 'attributes', 'config', 'devices', 'daemons'] | |
| def __init__(self, mgr): | |
| self.mgr = mgr | |
| self.ident = None | |
| self.requested_version = 2 | |
| self.loaded_version = None | |
| self.store = Store(mgr, 'inventory.', version=self.requested_version) | |
| self.load() | |
| # debug | |
| pp.pprint(self.__dict__) | |
| self.save() | |
| def save(self): | |
| """ | |
| to store | |
| """ | |
| json_data = {} | |
| for component in self.storable_components: | |
| component_obj = self.__getattribute__(component) | |
| json_data = ComponentBuilder.to_json(component_obj) | |
| print(f"Trying to save json_data {json_data} for component {component_obj}") | |
| self.store.save(ident=self.ident, data=json_data, version=self.loaded_version) | |
| def load(self): | |
| """ | |
| Load from the mon_store. | |
| This should happen on object creation. | |
| This method contains custom logic for populating attributes. | |
| """ | |
| for ident, data in self.store.load(): | |
| self.ident = ident | |
| self.loaded_version = data.pop('version') | |
| for host, component in data.items(): | |
| print(f"Loading components for host <{host}>") | |
| self.populate_inventory(component) | |
| assert self.loaded_version == self.requested_version | |
| def populate_inventory(self, component): | |
| for component_name, component_data in component.items(): | |
| print(f"Loading component -> {component_name}") | |
| """ | |
| There are different types of components which can be distinct by their | |
| return type: | |
| List: | |
| These components consists of multiple entries (Devices, Daemons, Networks etc) | |
| Dict: | |
| Those are flat `key: value` pairs (Attribute, Config) | |
| """ | |
| component_object = ComponentBuilder.from_json(component_data, component_name) | |
| self.__setattr__(component_name, component_object) | |
| class Cache: | |
| pass | |
| class Mgr: | |
| """ | |
| This is the MgrModule that implements interaction with the mgr-daemon, the mon store | |
| and more. Interesting methods are: | |
| * osdmap | |
| * osd_df | |
| * monmap | |
| * osd metadata | |
| * config store handling | |
| * ... and more | |
| """ | |
| def get_store_prefix(self, prefix=None, version=None): | |
| assert version | |
| if version == 1: | |
| return { | |
| "cephadm-dev": { | |
| "hostname": "cephadm-dev", | |
| "addr": "cephadm-dev", | |
| "labels": [], | |
| "status": "" | |
| }, | |
| "cephadm-dev-2": { | |
| "hostname": "cephadm-dev-2", | |
| "addr": "cephadm-dev-2", | |
| "labels": [], | |
| "status": "" | |
| } | |
| } | |
| if version == 2: | |
| return _load_inventory_struct() | |
| def set_store(self, data, version=None): | |
| """ | |
| This is the raw dump of the data into the mon_store | |
| """ | |
| assert version, data | |
| pass | |
| class Internal: | |
| """ | |
| This represents the respective backend orchestrator implementation (for Rook, mgr/cephadm, ceph-ansible) | |
| This is running on the ceph-mgr and implements a busy-loop which is triggered by a configurable timer | |
| called `serve`. | |
| """ | |
| def __init__(self): | |
| self.mgr = Mgr() | |
| self.inventory = Inventory(mgr=self.mgr) | |
| def serve(self): | |
| while True: | |
| print("foo") | |
| time.sleep(3) | |
| pass | |
| def example_internal_method(self, blob: Any) -> Optional[Any]: | |
| """ | |
| This is an example internal method that handles a *WRITE* request. | |
| It passes a `blob` of type Any and may return Any. | |
| (Keeping it vague here as it can return str, or any completion type object) | |
| To avoid long waiting times, hanging commands etc, the orchestrator backends | |
| return a `Future/Promise` basically saying that the task was accepted and will | |
| be executed when applicable. | |
| There are some things that are happening to reduce the feedback loop though: | |
| * Input validation | |
| ------------------ | |
| // Technically this happens during object creation in the `example_external_method` | |
| // Putting it here to *conceptually* highlight the role of this method. | |
| // Other validation require data about the state of the cluster which isn't | |
| // available in the `example_external_method` space. | |
| This can detect missing fields, bad formatting, unsupported features etc. | |
| * System validation | |
| ------------------- | |
| If you try to add something that would violate the principles of the cluster | |
| we can stop you right there. | |
| An example: | |
| // This is probably a bad example. Find a better one | |
| You're trying to add a gateway but there are no OSDs yet. The creation would | |
| fail as you can't save configuration (if stored in rados obj) or create pools. | |
| The mentioned validations should be quick and unnoticeable(speed wise) by the user. | |
| This method will eventually kick-off any according method that handles the respective | |
| case. | |
| """ | |
| return self.example_handler_method(blob) | |
| def example_handler_method(self, blob): | |
| """ | |
| After we successfully validated the input and ruled out any known | |
| complication we possibly, depending on the request, have to store something | |
| to the Store(). // Explained later | |
| For the sake of consistency we have to `raise` immediately if the `save` operation | |
| failed. | |
| """ | |
| pass | |
| class External: | |
| """ | |
| This represents the orchestrator CLI. | |
| It is agnostic of any orchestrator specifics and calls | |
| methods implemented in a `Interface`. | |
| """ | |
| def __init__(self): | |
| self.internal = Internal() | |
| def example_external_method(self, cmd="dummy command"): | |
| """ | |
| This is an example user-facing method. | |
| There are two types of commands: | |
| * READ | |
| ------ | |
| Commands like `ps`, `ls` etc. | |
| Those are not interesting for this purpose and will be omitted. | |
| * WRITE | |
| ------- | |
| Here we're trying to change the state of the cluster by | |
| deploying daemons, changing configuration, adding hosts etc. | |
| The method is responsible for the following actions: | |
| - Parse json/yaml/args | |
| - Construct objects | |
| - Call the respective method | |
| In order to stay consistent we should operate in the following order: | |
| """ | |
| print(f"Calling example method with cmd -> {cmd}") | |
| # TODO | |
| spec = object() | |
| dummy_data = dict(cmd='orch apply', spec=spec) | |
| return self.internal.example_internal_method(dummy_data) | |
| if __name__ == '__main__': | |
| module = Internal() | |
| module.serve() |
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment