+5
-10
@@ -14,22 +14,19 @@ acts:
|
||||
- name: "DB Dump and Backup"
|
||||
steps:
|
||||
- name: "find database container"
|
||||
actions:
|
||||
find_container: "get_service_container_name"
|
||||
action: "get_service_container_name"
|
||||
context:
|
||||
service_name: "gitea-db"
|
||||
|
||||
- name: "dump database"
|
||||
actions:
|
||||
dump_database: "dump_container_pg_db"
|
||||
action: "dump_container_pg_db"
|
||||
context:
|
||||
dump_path: "/tmp/db_dump.sql"
|
||||
db_user: "gitea"
|
||||
database: "gitea"
|
||||
|
||||
- name: "restic database backup"
|
||||
actions:
|
||||
dump_database: "backup_data_to_restic_repo"
|
||||
action: "backup_data_to_restic_repo"
|
||||
context:
|
||||
source_path: "/tmp/db_dump.sql"
|
||||
tags: ["gitea", "db-dump"]
|
||||
@@ -37,8 +34,7 @@ acts:
|
||||
- name: "Backup gitea data"
|
||||
steps:
|
||||
- name: "backup gitea mount points data"
|
||||
actions:
|
||||
backup_data_to_restic_repo: "backup_data_to_restic_repo"
|
||||
action: "backup_data_to_restic_repo"
|
||||
context:
|
||||
source_path: "/home/esilva/Trash/docker-compose/resources/gitea/gitea_data/"
|
||||
tags: ["gitea", "data"]
|
||||
@@ -46,8 +42,7 @@ acts:
|
||||
- name: "Backup gitea PG data"
|
||||
steps:
|
||||
- name: "backup gitea database mount points data"
|
||||
actions:
|
||||
backup_data_to_restic_repo: "backup_data_to_restic_repo"
|
||||
action: "backup_data_to_restic_repo"
|
||||
context:
|
||||
source_path: "/home/esilva/Trash/docker-compose/resources/gitea/postgres_data/"
|
||||
tags: ["gitea", "db-data"]
|
||||
|
||||
+100
-72
@@ -1,13 +1,13 @@
|
||||
from __future__ import annotations
|
||||
from re import A
|
||||
|
||||
from dataclasses import dataclass
|
||||
import yaml
|
||||
from abc import ABC, abstractmethod
|
||||
from pydantic import BaseModel
|
||||
|
||||
from playbook.models import ActModel, PlaybookModel, timed_run
|
||||
from playbook.models import ActModel, PlaybookModel, StepModel, timed_run
|
||||
from playbook.logging_models import Status, StepLogModel
|
||||
from playbook.action_registry import ActionFn
|
||||
|
||||
from playbook.action_registry import ActionFn, ActionRegistry
|
||||
|
||||
|
||||
class StepIF(ABC):
|
||||
@@ -21,18 +21,6 @@ class StepIF(ABC):
|
||||
""" Return a mapping of action roles (pre, play, post) to function names. """
|
||||
pass
|
||||
|
||||
|
||||
# class PlaybookError(Exception):
|
||||
# """ Base exception for Playbook parsing and execution failures. """
|
||||
# pass
|
||||
|
||||
|
||||
# @dataclass
|
||||
# class StepEntry(object):
|
||||
# name: str
|
||||
# step: StepIF
|
||||
|
||||
|
||||
# class Play(object):
|
||||
# def __init__(self, name: str, registries: dict[str, ActionRegistry]):
|
||||
# self.name: str = name
|
||||
@@ -130,57 +118,101 @@ class StepIF(ABC):
|
||||
# playLog.msg = "success"
|
||||
#
|
||||
# return playLog
|
||||
#
|
||||
# def _play_steps(self, steps: list[CustomStep]) -> StepLogModel:
|
||||
# playLog: StepLogModel = StepLogModel(
|
||||
# step_name=self.name,
|
||||
# status=Status.GOOD,
|
||||
# duration_sec=0,
|
||||
# msg="",
|
||||
# error=[],
|
||||
# substeps=[],
|
||||
# )
|
||||
#
|
||||
# prev_ctx: dict[str, object] = {}
|
||||
# for step in steps:
|
||||
# log: StepLogModel = timed_run(step.run, self._context | prev_ctx)
|
||||
# prev_ctx = log.pipe_ctx
|
||||
#
|
||||
# playLog.substeps.append(log)
|
||||
#
|
||||
# if log.failed:
|
||||
# playLog.error.append({
|
||||
# "status": "failed",
|
||||
# "output": f"failed in step: {step.name}"
|
||||
# })
|
||||
# playLog.status = Status.BAD
|
||||
# break
|
||||
#
|
||||
# if playLog.status == Status.GOOD:
|
||||
# playLog.msg = "success"
|
||||
#
|
||||
# return playLog
|
||||
#
|
||||
# def _get_function_from_registry(self, fn_name) -> ActionFn:
|
||||
# for reg in self._registries.values():
|
||||
# try:
|
||||
# return reg.get(fn_name, "*")
|
||||
# except ValueError:
|
||||
# continue
|
||||
#
|
||||
# raise ValueError(f"no function with name '{fn_name}' was found")
|
||||
#
|
||||
# def _build_custom_step(self, name, actions, ctx) -> CustomStep:
|
||||
# functionlist: list[ActionFn] = []
|
||||
# for f in actions:
|
||||
# try:
|
||||
# fn = self._get_function_from_registry(actions[f])
|
||||
# functionlist.append(fn)
|
||||
# except KeyError as e:
|
||||
# raise PlaybookError(
|
||||
# f"Step '{name}' is missing required action field: {e}") from e
|
||||
#
|
||||
# return CustomStep(name, functionlist, ctx)
|
||||
|
||||
|
||||
class PlaybookYAMLParser():
|
||||
@staticmethod
|
||||
def _parse_registries(reg_names: list[str]) -> list[ActionRegistry]:
|
||||
registries: list[ActionRegistry] = []
|
||||
for reg in reg_names:
|
||||
registries.extend(ActionRegistry.load_registries_from_file(reg))
|
||||
return registries
|
||||
|
||||
@staticmethod
|
||||
def _parse_steps(steps: list[dict[str, object]],
|
||||
registries: list[ActionRegistry]) -> list[StepModel]:
|
||||
step_list: list[StepModel] = []
|
||||
for step in steps:
|
||||
if not isinstance(step['name'], str):
|
||||
raise ValueError(
|
||||
f"name must of type `str` not `{type(step["name"])}`")
|
||||
if not isinstance(step['action'], str):
|
||||
raise ValueError(
|
||||
f"action must of type `str` not `{type(step["action"])}`")
|
||||
if not isinstance(step['context'], dict):
|
||||
raise ValueError(
|
||||
f"context must of type `dict` not `{type(step["context"])}`")
|
||||
|
||||
name: str = step['name']
|
||||
action: str = step['action']
|
||||
context: dict[str, object] = step['context']
|
||||
|
||||
found = False
|
||||
fn: ActionFn
|
||||
for reg in registries:
|
||||
try:
|
||||
fn = reg.get(action)
|
||||
found = True
|
||||
except:
|
||||
continue
|
||||
|
||||
if not found:
|
||||
raise ValueError(
|
||||
f"Unable to find function `{action}` in any registry")
|
||||
|
||||
step_list.append(StepModel(name=name, action=fn, context=context))
|
||||
|
||||
return step_list
|
||||
|
||||
@staticmethod
|
||||
def _parse_acts(acts_data: list[dict[str, object]],
|
||||
registries: list[ActionRegistry]) -> list[ActModel]:
|
||||
acts: list[ActModel] = []
|
||||
for act in acts_data:
|
||||
if not isinstance(act["name"], str):
|
||||
raise ValueError(
|
||||
f"name must of type `str` not `{type(act["name"])}`")
|
||||
|
||||
if not isinstance(act["steps"], list):
|
||||
raise ValueError(
|
||||
f"steps must of type `list` not `{type(act["steps"])}`")
|
||||
|
||||
name: str = act["name"]
|
||||
steps: list[dict[str, object]] = act["steps"]
|
||||
|
||||
step_list = PlaybookYAMLParser._parse_steps(steps, registries)
|
||||
modeled_act: ActModel = ActModel(name=name, steps=step_list)
|
||||
acts.append(modeled_act)
|
||||
|
||||
return acts
|
||||
|
||||
@staticmethod
|
||||
def from_yaml(fp: str) -> Playbook:
|
||||
with open(fp, "rb") as source_file:
|
||||
data = yaml.safe_load(source_file)
|
||||
|
||||
reg_yaml = data.get("registries", [])
|
||||
registries = PlaybookYAMLParser._parse_registries(reg_yaml)
|
||||
|
||||
acts_yaml = data.get("acts", {})
|
||||
acts: list[ActModel] = PlaybookYAMLParser._parse_acts(
|
||||
acts_yaml, registries)
|
||||
|
||||
global_context = data.get("global_context", {})
|
||||
|
||||
log_dir: str = data.get("log_dir", "./")
|
||||
playbook_name: str = data.get("playbook_name")
|
||||
|
||||
new_playbook: PlaybookModel = PlaybookModel(
|
||||
playbook_name=playbook_name,
|
||||
acts=acts,
|
||||
log_dir=log_dir,
|
||||
registries=registries,
|
||||
global_context=global_context
|
||||
)
|
||||
|
||||
return Playbook(new_playbook)
|
||||
pass
|
||||
|
||||
|
||||
class Playbook():
|
||||
@@ -195,11 +227,7 @@ class Playbook():
|
||||
|
||||
@classmethod
|
||||
def from_yaml(cls, fp: str) -> Playbook:
|
||||
try:
|
||||
new_playboook = PlaybookModel.from_yaml_file(fp)
|
||||
return Playbook(new_playboook)
|
||||
except Exception as e:
|
||||
raise e
|
||||
return PlaybookYAMLParser.from_yaml(fp)
|
||||
|
||||
def view_playbook(self) -> str:
|
||||
return self.model.model_dump_json()
|
||||
|
||||
@@ -4,6 +4,7 @@ import importlib.util
|
||||
from playbook.logging_models import StepLogModel
|
||||
from typing import Callable
|
||||
from pydantic import BaseModel, PrivateAttr
|
||||
from dataclasses import field
|
||||
|
||||
|
||||
ActionFn = Callable[[dict[str, object], str], StepLogModel]
|
||||
@@ -17,21 +18,21 @@ class Action(BaseModel):
|
||||
|
||||
class ActionRegistry(BaseModel):
|
||||
name: str
|
||||
_actions: dict[str, Action] = PrivateAttr(default_factory=dict)
|
||||
actions: dict[str, Action] = field(default_factory=dict)
|
||||
|
||||
def register(self, name: str, version: str):
|
||||
"""Decorator to register python functions with a name and version."""
|
||||
def decorator(function: ActionFn):
|
||||
self._actions[name] = Action(name=name, fn=function, ver=version)
|
||||
self.actions[name] = Action(name=name, fn=function, ver=version)
|
||||
return function
|
||||
return decorator
|
||||
|
||||
def get(self, name: str, ver: str = '*') -> ActionFn:
|
||||
if name not in self._actions:
|
||||
if name not in self.actions:
|
||||
raise ValueError(
|
||||
f"Action '{name}' is not registered in registry '{self.name}'")
|
||||
|
||||
return self._actions[name].fn
|
||||
return self.actions[name].fn
|
||||
|
||||
@staticmethod
|
||||
def load_registries_from_file(file_path: str) -> list[ActionRegistry]:
|
||||
|
||||
@@ -15,14 +15,17 @@ class LogTiming(BaseModel):
|
||||
duration_sec: float = 0.0
|
||||
|
||||
|
||||
class StepLogModel(BaseModel):
|
||||
class BaseLog(BaseModel):
|
||||
name: str
|
||||
timing: LogTiming = field(default=LogTiming())
|
||||
status: Status
|
||||
log_file_path: str = "" # If set we have logged this particular step to a file
|
||||
|
||||
|
||||
class ActLog(BaseLog):
|
||||
error: str = ""
|
||||
msg: str = ""
|
||||
pipe_ctx: dict[str, object] = field(default_factory=dict)
|
||||
log_file_path: str = "" # If set we have logged this particular step to a file
|
||||
logs: list[StepLogModel]
|
||||
|
||||
@property
|
||||
def failed(self) -> bool:
|
||||
@@ -30,9 +33,49 @@ class StepLogModel(BaseModel):
|
||||
return len(self.error) > 0 or self.status == Status.BAD
|
||||
|
||||
@classmethod
|
||||
def ok(cls, name: str, msg: str = "success", pipe_ctx: dict[str, object] = {}) -> StepLogModel:
|
||||
return cls(name=name, status=Status.GOOD, msg=msg, pipe_ctx=pipe_ctx)
|
||||
def ok(cls, name: str, logs: list[StepLogModel], msg: str = "success") -> ActLog:
|
||||
return cls(name=name, status=Status.GOOD, msg=msg, logs=logs)
|
||||
|
||||
@classmethod
|
||||
def fail(cls, name: str, err: str, logs: list[StepLogModel]) -> ActLog:
|
||||
return cls(name=name, status=Status.BAD, error=err, logs=logs)
|
||||
|
||||
|
||||
class PlaybookLog(BaseLog):
|
||||
errors: list[str] = []
|
||||
summary: str = ""
|
||||
act_logs: list[ActLog]
|
||||
|
||||
@property
|
||||
def failed(self) -> bool:
|
||||
return len(self.errors) > 0 or self.status == Status.BAD
|
||||
|
||||
@classmethod
|
||||
def ok(cls, name: str, logs: list[ActLog], msg: str = "success") -> PlaybookLog:
|
||||
return cls(name=name, status=Status.GOOD, summary=msg, act_logs=logs)
|
||||
|
||||
@classmethod
|
||||
def fail(cls, name: str, logs: list[ActLog], err: list[str]) -> PlaybookLog:
|
||||
return cls(name=name, status=Status.BAD, errors=err, act_logs=logs)
|
||||
|
||||
|
||||
class StepLogModel(BaseLog):
|
||||
error: str = ""
|
||||
msg: str = ""
|
||||
substeps_logs: list[StepLogModel] = []
|
||||
pipe_ctx: dict[str, object] = field(default_factory=dict)
|
||||
|
||||
@property
|
||||
def failed(self) -> bool:
|
||||
""" Checks if the step failed """
|
||||
return len(self.error) > 0 or self.status == Status.BAD
|
||||
|
||||
@classmethod
|
||||
def ok(cls, name: str, msg: str = "success", pipe_ctx: dict[str, object] = {},
|
||||
substeps: list[StepLogModel] = []) -> StepLogModel:
|
||||
return cls(name=name, status=Status.GOOD, msg=msg, pipe_ctx=pipe_ctx,
|
||||
substeps_logs=substeps)
|
||||
|
||||
@ classmethod
|
||||
def fail(cls, name: str, err: str) -> StepLogModel:
|
||||
return cls(name=name, status=Status.BAD, error=err)
|
||||
|
||||
+70
-71
@@ -5,62 +5,46 @@ import time
|
||||
import functools
|
||||
from typing import Callable
|
||||
from datetime import datetime, timezone
|
||||
from pydantic import BaseModel, field_validator
|
||||
from pydantic import BaseModel, ValidationInfo, field_serializer, field_validator, model_validator
|
||||
|
||||
from playbook.action_registry import ActionRegistry, ActionFn
|
||||
from playbook.logging_models import StepLogModel
|
||||
from playbook.logging_models import PlaybookLog, StepLogModel, ActLog, Status
|
||||
|
||||
|
||||
CtxType = dict[str, object]
|
||||
|
||||
|
||||
class StepModel(BaseModel):
|
||||
name: str
|
||||
_actions: list[ActionFn]
|
||||
_context: dict[str, object]
|
||||
action: ActionFn
|
||||
context: dict[str, object]
|
||||
|
||||
def get_action_names(self) -> list[str]:
|
||||
return [action.__name__ for action in self._actions]
|
||||
@field_serializer("action")
|
||||
def serialize_action_fn(self, action_fn: ActionFn, _info) -> str:
|
||||
return getattr(action_fn, "__name__", str(action_fn))
|
||||
|
||||
# NOTE: refactor this method to be more readable
|
||||
def run(self, ctx) -> StepLogModel:
|
||||
""" Run the step. """
|
||||
return StepLogModel.ok(self.name, msg="success")
|
||||
# substeps: list[StepLogModel] = []
|
||||
# status: Status = Status.GOOD
|
||||
#
|
||||
# # NOTE:
|
||||
# # Use this to move context from the previous step to the next step
|
||||
# # This is wrong, we are moving away from having a list of steps here to a single
|
||||
# # step for each operation
|
||||
# prev_step_ctx: dict[str, object] = {}
|
||||
#
|
||||
# for op in self.__actions:
|
||||
# try:
|
||||
# log: StepLogModel = timed_run(
|
||||
# op, ctx | prev_step_ctx | self.__context, self.name)
|
||||
# except Exception as e:
|
||||
# log: StepLogModel = StepLogModel.fail(
|
||||
# self.name, [{"status": "failed", "output": str(e)}])
|
||||
#
|
||||
# status = Status.BAD if log.failed else Status.GOOD
|
||||
# prev_step_ctx = log.pipe_ctx
|
||||
#
|
||||
# substeps.append(log)
|
||||
#
|
||||
# if status == Status.BAD:
|
||||
# break
|
||||
#
|
||||
# errors = []
|
||||
# msg: str = ""
|
||||
# if status == Status.BAD:
|
||||
# # Note, use this to set the standar error output/formatting
|
||||
# # errors = [{"status": "failed", "output": ""}]
|
||||
# errors = []
|
||||
# else:
|
||||
# msg = "success"
|
||||
#
|
||||
# return StepLogModel(
|
||||
# step_name=self.name, status=status,
|
||||
# msg=msg, error=errors, substeps=substeps, pipe_ctx=prev_step_ctx
|
||||
# )
|
||||
substeps: list[StepLogModel] = []
|
||||
previous_step_ctx: dict[str, object] = {}
|
||||
for action in self.actions:
|
||||
try:
|
||||
local_ctx: CtxType = ctx | previous_step_ctx | self.context
|
||||
log: StepLogModel = timed_run(
|
||||
action, local_ctx, self.name
|
||||
)
|
||||
except Exception as e:
|
||||
log: StepLogModel = StepLogModel.fail(
|
||||
self.name, str(e)
|
||||
)
|
||||
|
||||
previous_step_ctx = log.pipe_ctx
|
||||
|
||||
substeps.append(log)
|
||||
if log.status == Status.BAD:
|
||||
return StepLogModel.fail(self.name, err=f"failed on step: {log.name}")
|
||||
|
||||
return StepLogModel.ok(self.name, msg="success", substeps=substeps)
|
||||
|
||||
|
||||
class ActModel(BaseModel):
|
||||
@@ -68,8 +52,15 @@ class ActModel(BaseModel):
|
||||
steps: list[StepModel]
|
||||
|
||||
# NOTE: this shouldnt return a steplogmodel but an actlogmodel or something like that
|
||||
def run(self, ctx) -> StepLogModel:
|
||||
return StepLogModel.ok(self.name, msg="success")
|
||||
def run(self, ctx) -> ActLog:
|
||||
step_logs: list[StepLogModel] = []
|
||||
for step in self.steps:
|
||||
stepLog: StepLogModel = step.run(ctx)
|
||||
step_logs.append(stepLog)
|
||||
if stepLog.failed:
|
||||
break
|
||||
|
||||
return ActLog.ok(self.name, msg="success", logs=step_logs)
|
||||
|
||||
|
||||
class PlaybookModel(BaseModel):
|
||||
@@ -79,36 +70,43 @@ class PlaybookModel(BaseModel):
|
||||
global_context: dict[str, str]
|
||||
acts: list[ActModel]
|
||||
|
||||
@field_validator("registries", mode="before")
|
||||
@model_validator(mode="before")
|
||||
@classmethod
|
||||
def parse_registries(cls, v: object) -> list[ActionRegistry]:
|
||||
result = []
|
||||
if isinstance(v, list):
|
||||
for item in v:
|
||||
# def parse_registries(cls, v: object) -> list[ActionRegistry]:
|
||||
def load_registries(cls, data: object, info: ValidationInfo) -> object:
|
||||
if not isinstance(data, dict):
|
||||
return data
|
||||
|
||||
raw_registries = data.get("registries", [])
|
||||
loaded_registries: list[ActionRegistry] = []
|
||||
|
||||
for item in raw_registries:
|
||||
if isinstance(item, str):
|
||||
try:
|
||||
registries: list[ActionRegistry] = \
|
||||
ActionRegistry.load_registries_from_file(item)
|
||||
result.extend(registries)
|
||||
except Exception as e:
|
||||
raise e
|
||||
else:
|
||||
raise ValueError(f"Invalid registry type: {type(item)}")
|
||||
return result
|
||||
loaded_registries.extend(ActionRegistry.load_registries_from_file(item))
|
||||
|
||||
data["registries"] = loaded_registries
|
||||
|
||||
if info.context is not None:
|
||||
info.context["registries"] = loaded_registries
|
||||
|
||||
return data
|
||||
|
||||
@classmethod
|
||||
def from_yaml_file(cls, fp: str):
|
||||
with open(fp, "rb") as f:
|
||||
data = yaml.safe_load(f)
|
||||
return cls(**data)
|
||||
return cls.model_validate(data, context={})
|
||||
# return cls(**data)
|
||||
|
||||
def _run(self) -> StepLogModel:
|
||||
def _run(self) -> PlaybookLog:
|
||||
act_logs: list[ActLog] = []
|
||||
for act in self.acts:
|
||||
act.run(self.global_context)
|
||||
log: ActLog = act.run(self.global_context)
|
||||
act_logs.append(log)
|
||||
|
||||
return StepLogModel.ok(name="somasjd", msg="ajshdsajhd")
|
||||
return PlaybookLog.ok(name="somasjd", msg="ajshdsajhd", logs=act_logs)
|
||||
|
||||
def run(self) -> StepLogModel:
|
||||
def run(self) -> PlaybookLog:
|
||||
return timed_run(self._run)
|
||||
|
||||
|
||||
@@ -128,7 +126,7 @@ class ContextChecker:
|
||||
else:
|
||||
missing_keys = [key for key in args if key not in ctx]
|
||||
missing_keys_msg = f"missing context keys: {missing_keys}"
|
||||
return StepLog.fail(name, [{"status": "failed", "output": missing_keys_msg}])
|
||||
return StepLogModel.fail(name, err=missing_keys_msg)
|
||||
return wrapper
|
||||
return decorator
|
||||
|
||||
@@ -142,17 +140,18 @@ def human_readable_date(time: int) -> str:
|
||||
return dt.strftime(f"%Y-%m-%dT%H:%M:%S.{nanos:09d}%:z")
|
||||
|
||||
|
||||
def timed_run(op: Callable[..., StepLogModel], *args: object, **kwargs: object) -> StepLogModel:
|
||||
# def timed_run(op: Callable[..., StepLogModel], *args: object, **kwargs: object) -> StepLogModel:
|
||||
def timed_run(op, *args: object, **kwargs: object):
|
||||
start_date: int = time.time_ns()
|
||||
start_time: int = time.perf_counter_ns()
|
||||
|
||||
try:
|
||||
log: StepLogModel = op(*args, **kwargs)
|
||||
log = op(*args, **kwargs)
|
||||
finally:
|
||||
delta = time.perf_counter_ns() - start_time
|
||||
|
||||
# NOTE: we could simplyfy this by having a marshalling method
|
||||
if isinstance(log, StepLogModel):
|
||||
if isinstance(log, StepLogModel | PlaybookLog | ActLog):
|
||||
end_date: int = time.time_ns()
|
||||
log.timing.duration_sec = delta / 1_000_000_000.0
|
||||
log.timing.start_date = human_readable_date(start_date)
|
||||
|
||||
Reference in New Issue
Block a user