Skip to content

Instantly share code, notes, and snippets.

@jschmid1
Created August 14, 2020 17:04
Show Gist options
  • Select an option

  • Save jschmid1/09b56ce482aed479fba59521a959e8cc to your computer and use it in GitHub Desktop.

Select an option

Save jschmid1/09b56ce482aed479fba59521a959e8cc to your computer and use it in GitHub Desktop.
v5
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, namespace, version=2):
self.mgr = mgr
self.namespace = namespace
self.version = version
def save(self, data=None) -> bool:
"""
Save `blob` to store.
Example call to this function would look like:
`self.store.save("spec/host_a/config", {'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.
// This approach is overcoming the drawbacks of the non-flat hierarchy approach
// by giving each component their own namespace which allows us to update them
// individually without having to re-write the entire inventory/host blob.
"""
assert data
print(f"Saving in namespace -> {self.namespace}")
print(f"Saving data -> {data}")
# TODO: bug, existing store data is not being loaded
old = self.mgr.get_store_prefix(self.namespace, version=self.version)
if self.has_changes(data, old):
print("Changes detected. Writing to store")
return self.mgr.set_store(self.namespace, self.jsonify(data), version=self.version)
print("No changes. Not writing to the store")
return True
@staticmethod
def jsonify(data) -> Optional[str]:
try:
return json.dumps(data)
except json.JSONDecodeError:
# TODO: catch correct json errors
raise Exception("Error encoding json")
@staticmethod
def has_changes(new=None, old=None) -> bool:
"""
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.")
print(f"old -> {old}")
print(f"new -> {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.namespace, 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:
# TODO: Consider to move to __slots__ instead..
# This may give more support for attribute access (foo.bar)
# 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, mgr=None, host=None, **kwargs):
print(f"Loading kwargs -> {kwargs}")
self.mgr = mgr
self.host = host
self.version = None
self.store = self.load_store()
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, notify=False)
def load_store(self):
"""
Load the dynamically namespaced `Store`
"""
return Store(self.mgr,
self.namespace(self.host),
version=self.version)
@property
def component_name(self):
return self.__class__.__name__.lower()
def namespace(self, host):
"""
Components have a dynamic namespace based on the hostname and the
`component_name`.
TODO: Find a rule for pluralization/singularization
Currently we're overwriting methods if it's NOT pluralizable
"""
return f"inventory/{host}/{self.component_name}s"
@classmethod
def from_json(cls, data, mgr, host) -> 'Component':
print(f"Loading {cls} from json")
return cls(mgr=mgr, host=host, **{
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 __setattr__(self, key, value, notify=True):
if (not hasattr(self, key) or getattr(self, key) != value) and key in self.loadable_fields:
if notify:
self.notify_on_change()
super(Component, self).__setattr__(key, value)
def notify_on_change(self):
print("A loadable_field was changed. This triggers a save() operation")
self.save()
def save(self):
self.store.save(data=self.to_json())
def __eq__(self, other):
return (self.__class__ == other.__class__ and
self.to_json() == other.to_json)
def __hash__(self):
return hash(self.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)
def namespace(self, host):
"""
Exception: Singularize (overwrite)
"""
return f"inventory/{host}/{self.component_name}"
# 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, mgr, host):
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, mgr=mgr, host=host)
if isinstance(data, list):
return [component_class.from_json(i, mgr=mgr, host=host) for i in data]
raise Exception(f"Can't handle type {type(data)}")
@staticmethod
def update_field(key, value, component):
if isinstance(component, list):
# This makes sure that to_json exists for all objects
return [d.__setattr__(key, value) for d in component]
if isinstance(component, Component):
return component.__setattr__(key, value)
class Host:
"""
The Host class has instances of `storable_components` loaded dynamically.
This makes it easy to retrieve information about the host and its inventory.
Example:
self.daemons -> List[DaemonDescription]
TODO: more real world examples
"""
storable_components = ['networks', 'attributes', 'config', 'devices', 'daemons']
def __init__(self, hostname, mgr=None):
self.mgr = mgr
self.hostname = hostname
self.requested_version = 2 # This can come from the module config
# TODO: hostnames can change. Account for that.
self.store = Store(mgr, self.namespace(), version=self.requested_version)
def online_daemons(self):
return [x for x in self.daemons if x.running == True]
def namespace(self):
return f"inventory/{self.hostname}"
@property
def inventory_objects(self):
return [self.__getattribute__(component) for component in self.storable_components]
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 identified by their
return type:
List:
These components consists of multiple entries (i.e. Devices, Daemons, Networks etc)
Dict:
Those are flat `key: value` pairs (i.e. Attribute, Config)
This is abstracted to ComponentBuilder()
"""
component_object = ComponentBuilder.from_json(component_data, component_name, self.mgr, self.hostname)
self.__setattr__(component_name, component_object)
def __repr__(self):
printable_fields = {k: v for (k, v) in self.__dict__.items() if k in self.storable_components}
return f"{self.__class__.__name__}({printable_fields})"
class Inventory:
"""
Contains a List of Hosts that are registered to the system
"""
def __init__(self, mgr):
self.mgr = mgr
self.ident = None
self.requested_version = 2 # This can come from the module config
self.loaded_version = None
self.store = Store(mgr, 'inventory.', version=self.requested_version)
self.hosts = list()
self.load()
# debug
pp.pprint(self.__dict__)
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}>")
host = Host(hostname=host, mgr=self.mgr)
host.populate_inventory(component)
self.hosts.append(host)
assert self.loaded_version == self.requested_version
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, namespace=None, version=None):
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, namespace, data, version=None):
"""
This is the raw dump of the data into the mon_store
"""
assert data
assert namespace
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):
counter = 0
while True:
print("foo")
counter += 1
self.inventory.hosts[0].attributes.key1 = counter
self.inventory.hosts[0].config.key1 = counter
self.inventory.hosts[0].devices[0].key1 = counter
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