微信公众号:进击的大杂烩
欢迎关注我,一起学习,一起进步!
Ansible 2 API
ansible 2 API发生了很大的变化。
通过对ansible 2.4.2 的源代码(Python 环境为2.7.5)进行分析来学习如何使用ansible 2 api 并自己编写一个ansible api。
ansible 2.4.2 相对于 ansible 2.2.2 变化比较大的地方是 Inventory 类和解析 inventory 的方式。
我们分析ansible的AdHocCLI模式来了解ansible的运行过程。
入口文件分析:
入口文件:ansible 命令(通过which ansible命令来查看命令所在目录)
对源码进行了简化,只分析AdHocCLI模式相关的代码
if __name__ == '__main__':
me = os.path.basename(sys.argv[0])
target = me.split('-')
if target[0] == 'ansible':
sub = 'adhoc'
myclass = 'AdHocCLI'
mycli = getattr(__import__("ansible.cli.%s" % sub, fromlist=[myclass]), myclass)
args = [to_text(a, errors='surrogate_or_strict') for a in sys.argv]
cli = mycli(args)
cli.parse()
exit_code = cli.run()
shutil.rmtree(C.DEFAULT_LOCAL_TMP, True)
sys.exit(exit_code)
通过以上简化的代码可以知道入口文件做了以下几件事情:
确定命令工具(AdHocCLI模式使用的是 ansible 命令)
定义sub,myclass变量
导入类AdHocCLI
mycli = getattr(__import__("ansible.cli.%s" % sub, fromlist=[myclass]), myclass)等同于
from ansible.cli import adhoc
mycli = adhoc.AdHocCLI处理命令行编码格式(to_text)--用来统一编码格式(源码默认编码为utf-8)
实例化mycli类(cli = mycli(args))
通过解析器(cli.parse())来解析ansible命令行参数
运行cli(cli.run())
清理临时文件
退出命令行
AdHOCCLI 类分析
对应入口文件的mycli = getattr(__impor__("ansible.cli.%s" % sub, fromlist=[myclass]), myclass)
AdHOCCLI类--源码位置(ansible/cli/adhoc.py)
class AdHocCLI(CLI) 从CLI继承, mycli 通过 cli = AdHocCLI(args) 实例化,我们来看CLI类(源码位置ansible/cli/__init__.py)的初始化函数__init__():
注意在ansible/cli/init.py中有一个引入
from ansible import constants as C
constants中配置了ansible配置的选项和默认值--源码位置:ansible/constants.py
constants.py 也会将 ansible/conf/base.yml 中的配置一起加载到constants中
通过以下代码可以看到_init_的工作仅仅是做了一些参数的初始化
def __init__(self, args, callback=None):
"""
Base init method for all command line programs
"""
self.args = args
self.options = None
self.parser = None
self.action = None
self.callback = callback
命令行参数解析函数
对应入口文件的cli.parse(),相关代码如下:
def parse(self):
''' create an options parser for bin/ansible '''
#CLI.base_parser是CLI类的一个静态方法。
self.parser = CLI.base_parser(
usage='%prog <host-pattern> [options]',
runas_opts=True,
inventory_opts=True,
async_opts=True,
output_opts=True,
connect_opts=True,
check_opts=True,
runtask_opts=True,
vault_opts=True,
fork_opts=True,
module_opts=True,
desc="Define and run a single task 'playbook' against a set of hosts",
epilog="Some modules do not make sense in Ad-Hoc (include, meta, etc)",
)
super(AdHocCLI, self).parse()
self.validate_conflicts(runas_opts=True, vault_opts=True, fork_opts=True)
#CLI.base_parser 源码
@staticmethod
def base_parser(usage="", output_opts=False, runas_opts=False, meta_opts=False, runtask_opts=False, vault_opts=False, module_opts=False,
async_opts=False, connect_opts=False, subset_opts=False, check_opts=False, inventory_opts=False, epilog=None, fork_opts=False,
runas_prompt_opts=False, desc=None):
''' create an options parser for most ansible scripts '''
# base opts
parser = SortedOptParser(usage, version=CLI.version("%prog"), description=desc, epilog=epilog)
class SortedOptParser(optparse.OptionParser):
'''Optparser which sorts the options by opt before outputting --help'''
def format_help(self, formatter=None, epilog=None):
self.option_list.sort(key=operator.methodcaller('get_opt_string'))
return optparse.OptionParser.format_help(self, formatter=None)
通过上面的baseparser代码段可以看出ansible对于参数的解析是通过python标准库operator来实现的。baseparser 函数的具体内容可以去源码看,就是一些参数个数和说明(ansible -h)
我们继续看parse()函数, super(AdHocCLI, self).parse()。可见调用了CLI类的parse()函数:
@abstractmethod
def parse(self):
self.options, self.args = self.parser.parse_args(self.args[1:])
可以看到通过self.parser.parse_args(self.args[1:])返回命名参数和未命名参数(要记住self.parser实际上是operator.OptionParser的实例),然后根据option做了一些逻辑判断和操作。
此时self.args 就是ansible 的pattern。完成参数格式化后,通过self.validate_conflicts()又做了几个参数的合法性验证。
运行阶段--cli.run()
对应入口文件cli.run()
参数解析完成后,到了最关键的运行阶段--cli.run(),部分代码和代码执行流程如下:
def run(self):
''' create and execute the single task playbook '''
super(AdHocCLI, self).run()
pattern = to_text(self.args[0], errors='surrogate_or_strict')
# dynamically load any plugins
get_all_plugin_loaders()
loader, inventory, variable_manager = self._play_prereqs(self.options)
hosts = inventory.list_hosts(pattern)
if self.options.module_name in C.MODULE_REQUIRE_ARGS and not self.options.module_args:
err = "No argument passed to %s module" % self.options.module_name
if pattern.endswith(".yml"):
err = err + ' (did you mean to run ansible-playbook?)'
raise AnsibleOptionsError(err)
play_ds = self._play_ds(pattern, self.options.seconds, self.options.poll_interval)
play = Play().load(play_ds, variable_manager=variable_manager, loader=loader)
if self.callback:
cb = self.callback
elif self.options.one_line:
cb = 'oneline'
# Respect custom 'stdout_callback' only with enabled 'bin_ansible_callbacks'
elif C.DEFAULT_LOAD_CALLBACK_PLUGINS and C.DEFAULT_STDOUT_CALLBACK != 'default':
cb = C.DEFAULT_STDOUT_CALLBACK
else:
cb = 'minimal'
self._tqm = None
try:
self._tqm = TaskQueueManager(
inventory=inventory,
variable_manager=variable_manager,
loader=loader,
options=self.options,
passwords=passwords,
stdout_callback=cb,
run_additional_callbacks=C.DEFAULT_LOAD_CALLBACK_PLUGINS,
run_tree=run_tree,
)
result = self._tqm.run(play)
finally:
if self._tqm:
self._tqm.cleanup()
if loader:
loader.cleanup_all_tmp_files()
return result
定义通配符:pattern
加载所有模块--get_all_plugin_loaders():
作用是将模块 ansible.plugins.loader 中符合(isinstance(obj, PluginLoader))这个条件的加载器加载语句:loader, inventory, variable_manager = self._play_prereqs(self.options)的作用是:生成加载器loader, inventory(实际工作就是将source解析成inventory对象), 变量管理器实例variable_manager
匹配目标hosts--hosts = inventory.list_hosts(pattern)
判断调用模块和模块参数是否合法
生成运行对象--运行的模块,参数(seconds参数对应ansible -B参数(后台运行长时任务),poll_interval对应ansible -P 表示对后台任务的轮询的间隔时间)
根据条件确定回调函数(ansible 命令返回结果的处理函数)
创建一个任务队列去运行paly(运行对象)
返回运行结果
主要来看一下生成 inventory 对象的过程,函数_play_prereqs代码和相关解析如下:
@staticmethod
def _play_prereqs(options):
# all needs loader
loader = DataLoader()
vault_ids = options.vault_ids
default_vault_ids = C.DEFAULT_VAULT_IDENTITY_LIST
vault_ids = default_vault_ids + vault_ids
vault_secrets = CLI.setup_vault_secrets(loader,
vault_ids=vault_ids,
vault_password_files=options.vault_password_files,
ask_vault_pass=options.ask_vault_pass,
auto_prompt=False)
loader.set_vault_secrets(vault_secrets)
# create the inventory, and filter it based on the subset specified (if any)
inventory = InventoryManager(loader=loader, sources=options.inventory)
# create the variable manager, which will be shared throughout
# the code, ensuring a consistent view of global variables
variable_manager = VariableManager(loader=loader, inventory=inventory)
# load vars from cli options
variable_manager.extra_vars = load_extra_vars(loader=loader, options=options)
variable_manager.options_vars = load_options_vars(options, CLI.version_info(gitinfo=False))
return loader, inventory, variable_manager
#其中 DataLoader 是一个用来加载和解析YAML或Json格式的类
#vault 相关的是ansible对敏感文件加密解密的相关选项,并将加密->解密的对应方式通过loader.set_vault_secrets(vault_secrets)来绑定到加载器上
#ansible的inventory是通过实例化InventoryManager(loader=loader, sources=options.inventory)来生成的
#函数_play_prereqs中生成inventory对象的语句如下:
#inventory = InventoryManager(loader=loader, sources=options.inventory)
#通过InventoryManager类代码来说明具体流程:
class InventoryManager(object):
''' Creates and manages inventory '''
#初始化的时候实例化了InventoryData()类,调用了parse_sources()函数,这个函数的作用是解析inventory的源
def __init__(self, loader, sources=None):
# base objects
self._loader = loader
self._inventory = InventoryData()
# a list of host(names) to contain current inquiries to
self._restriction = None
self._subset = None
# caches
self._hosts_patterns_cache = {} # resolved full patterns
self._pattern_cache = {} # resolved individual patterns
self._inventory_plugins = [] # for generating inventory
# the inventory dirs, files, script paths or lists of hosts
if sources is None:
self._sources = []
elif isinstance(sources, string_types):
self._sources = [sources]
else:
self._sources = sources
# get to work!
self.parse_sources()
def parse_sources(self, cache=False):
''' iterate over inventory sources and parse each one to populate it'''
self._setup_inventory_plugins()
parsed = False
# allow for multiple inventory parsing
for source in self._sources:
if source:
if ',' not in source:
source = unfrackpath(source, follow=False)
parse = self.parse_source(source, cache=cache)
if parse and not parsed:
parsed = True
if parsed:
# do post processing
self._inventory.reconcile_inventory()
else:
display.warning("No inventory was parsed, only implicit localhost is available")
self._inventory_plugins = []
#真正解析_sources的函数是parse_source(),这个函数通过加载plugin的inventory插件来解析_soureces,默认插件有:['host_list', 'script', 'yaml', 'ini']。
#此处有个bug,if ',' not in source 语句中表示如果source中没有","就通过unfrackpath函数先处理source一次,处理的过程就是将source当做目录来处理。
#所以如果是单个机器作为inventory后面要加个","才能正常运行。
#测试方法:
#ansible all -i '192.168.1.x' -m shell -a date
#ansible all -i '192.168.1.x,' -m shell -a date
梳理运行流程
通过对代码的分析,根据这个流程自定义运行过程如下:
采用 ssh 的秘钥模式, 管理节点和被管理节点已经互信
常用的ansible参数为:
Options = namedtuple('Options',[
'connection', 'module_path', "remote_user",
'timeout', 'forks', 'become', 'become_method', 'become_user',
'seconds', 'poll_interval', 'check', 'diff',
])
options = Options(
connection = 'smart', #采用ssh模式,smart会在本机ssh和paramiko直接选一种方式
module_path = module_path,
remote_user = remote_user,
timeout = timeout,
forks = forks,
become = become,
become_method = become_method,
seconds = seconds,
poll_interval = poll_interval,
check = check,
deff = deff,
)
关键的inventory,从源码我们知道 inventory是通过InventoryManager类实现的:
hosts定义:
loder = DataLoader()
hosts = ','.join(['192.168.100.101:22', '192.168.100.102:22', '192.168.100.103:22'])
inventory = InventoryManager(loader=loader, sources=hosts)
variable_manager = VariableManager(loader=loader, inventory=inventory)
创建任务
tasks = [dict(action=dict(module='shell', args='ls'))]
根据任务创建运行对象
play_source = dict(
name = "Ansible AdHoc",
hosts = "all",
gather_facts = 'no',
tasks = tasks
)
play = Play().load(play_source, variable_manager=variable_manager, loader=loader)
定义callback
cb = C.DEFAULT_STDOUT_CALLBACK
启动任务
tqm = None
try:
tqm = TaskQueueManager(
inventory=inventory,
variable_manager=variable_manager,
loader=loader,
options=options,
passwords={},
stdout_callback=cb,
)
result = tqm.run(play)
finally:
if tqm is not None:
tqm.cleanup()
写成一个工具类
#coding:utf-8
import json
from collections import namedtuple
import ansible.constants as C
from ansible.parsing.dataloader import DataLoader
from ansible.vars.manager import VariableManager
from ansible.inventory.manager import InventoryManager
from ansible.playbook.play import Play
from ansible.executor.task_queue_manager import TaskQueueManager
from ansible.plugins.callback import CallbackBase
from ansible.errors import AnsibleError, AnsibleOptionsError
#结果回调类,是从CallbackBase继承的用于结果处理
class ResultCallback(CallbackBase):
"""
Callback 类,增加了result_info,用于存储返回结果。
设置async和poll的时候并没有回调到相关函数上!!!
"""
def __init__(self, result_info, display=None, options=None):
self.result_info = result_info
super(ResultCallback, self).__init__(display, options)
def v2_runner_on_ok(self, result):
host = result._host.get_name()
self.runner_on_ok(host, result._result)
def runner_on_ok(self, host, res):
self.result_info['contacted'].setdefault(host, []).append(res)
def runner_on_failed(self, host, res, ignore_errors=False):
self.result_info['dark'].setdefault(host, []).append(res)
def runner_on_skipped(self, host, item=None):
self.result_info['dark'].setdefault(host, []).append(res)
def runner_on_unreachable(self, host, res):
self.result_info['dark'].setdefault(host, []).append(res)
def runner_on_async_poll(self, host, res, jid, clock):
pass
def runner_on_async_ok(self, host, res, jid):
pass
def runner_on_async_failed(self, host, res, jid):
pass
class AdHocRunnerAPI(object):
"""
Ad-Hoc Api
"""
Options = namedtuple('Options',[
'connection', 'module_path', "remote_user",
'timeout', 'forks', 'become', 'become_method', 'become_user',
'seconds', 'poll_interval', 'check', 'diff',
])
def __init__(self,
hosts=C.DEFAULT_HOST_LIST, # Inventory Source default: etc/ansible/hosts
forks=C.DEFAULT_FORKS, # 5
timeout=C.DEFAULT_TIMEOUT, # SSH timeout = 10s
remote_user=C.DEFAULT_REMOTE_USER, # root
module_path=None, # dirs of custome modules
connection_type="smart", # ssh or paramiko
poll_interval=C.DEFAULT_POLL_INTERVAL, # 15s
seconds=0, # -B run asynchronously, failing after X seconds
become=C.DEFAULT_BECOME, # default False
become_method=C.DEFAULT_BECOME_METHOD, # privilege escalation method to use (default=sudo)
become_user=None, # run operations as this user (default=root)
check=False,
diff=False):
#配置相关
self.options = self.Options(
connection = connection_type,
module_path = module_path,
remote_user = remote_user,
timeout = timeout,
forks = forks,
become = become,
become_method = become_method,
become_user = become_user,
seconds = seconds,
poll_interval = poll_interval,
check = check,
diff = diff,
)
self.pattern = 'all'
#实例化 inventory
self.loader = DataLoader()
self.inventory = InventoryManager(loader=self.loader, sources=[hosts])
self.variable_manager = VariableManager(loader=self.loader, inventory=self.inventory)
#其他定义
self.tasks = []
self.play_source = None
self.play = None
self._tqm = None
#逻辑判断 hosts 是否合法
def check_hosts(self):
if len(self.inventory.list_hosts()) == 0:
raise AnsibleError("Inventory is empty.")
hosts = self.inventory.list_hosts(self.pattern)
if len(hosts) == 0:
raise AnsibleError("Specified --limit does not match any hosts")
#逻辑判断模块和参数合法性
def check_module(self, module, args=''):
if module in C.MODULE_REQUIRE_ARGS and not args:
raise AnsibleOptionsError("No argument passed to %s module" % module)
#清理历史运行记录
def cleanup(self):
self.tasks = []
self.play_source = None
self.paly = None
#启动任务
def run(self, tasks, pattern='all'):
"""
:param tasks: (('shell', 'ls'), ('ping', ''))
:param pattern:
:return:
"""
self.pattern = pattern
self.check_hosts()
#将返回的结果通过回调类存储到 result_info 中,便于使用和二次分析
self.result_info = {
'contacted' : {},
'dark' : {}
}
self.cb = ResultCallback(self.result_info)
for module, args in tasks:
self.check_module(module, args)
self.tasks.append(
dict(action=dict(
module=module,
args=args),
async=self.options.seconds,
poll=self.options.poll_interval
)
)
#生成运行对象
self.play_source = dict(
name='Ansible AdHoc',
hosts=self.pattern,
gather_facts='no',
tasks=self.tasks
)
self.play = Play().load(self.play_source, variable_manager=self.variable_manager, loader=self.loader)
#创建任务队列运行play
try:
self._tqm = TaskQueueManager(
inventory=self.inventory,
variable_manager=self.variable_manager,
loader=self.loader,
options=self.options,
passwords={},
stdout_callback=self.cb,
)
result = self._tqm.run(self.play)
finally:
if self._tqm:
self._tqm.cleanup()
if self.loader:
self.loader.cleanup_all_tmp_files()
self.cleanup()
return result
#测试
if __name__ == '__main__':
hosts = '192.168.100.101:22,192.168.10.102:22'
tasks = (('shell', 'sleep 10'),)
tasks = (('shell', 'date'),)
hadoc = AdHocRunnerAPI(hosts)
result_code = hadoc.run(tasks)
print result_code
print json.dumps(hadoc.result_info, indent=4)
result_code = hadoc.run(tasks)
print json.dumps(hadoc.result_info, indent=4)
参考:
http://docs.ansible.com/ansible/latest/intro.html
http://docs.ansible.com/ansible/latest/dev_guide/developing_api.html#python-api-2-0





