暂无图片
暂无图片
暂无图片
暂无图片
暂无图片

Ansible 2 Api 源码分析及实现

进击的大杂烩 2017-12-21
1144

微信公众号:进击的大杂烩
欢迎关注我,一起学习,一起进步!

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)

通过以上简化的代码可以知道入口文件做了以下几件事情:

  1. 确定命令工具(AdHocCLI模式使用的是 ansible 命令)

  2. 定义sub,myclass变量

  3. 导入类AdHocCLI
    mycli = getattr(__import__("ansible.cli.%s" % sub, fromlist=[myclass]), myclass)等同于
    from ansible.cli import adhoc
    mycli = adhoc.AdHocCLI

  4. 处理命令行编码格式(to_text)--用来统一编码格式(源码默认编码为utf-8)

  5. 实例化mycli类(cli = mycli(args))

  6. 通过解析器(cli.parse())来解析ansible命令行参数

  7. 运行cli(cli.run())

  8. 清理临时文件

  9. 退出命令行

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

  1. 定义通配符:pattern

  2. 加载所有模块--get_all_plugin_loaders():
    作用是将模块 ansible.plugins.loader 中符合(isinstance(obj, PluginLoader))这个条件的加载器加载

  3. 语句:loader, inventory, variable_manager = self._play_prereqs(self.options)的作用是:生成加载器loader, inventory(实际工作就是将source解析成inventory对象), 变量管理器实例variable_manager

  4. 匹配目标hosts--hosts = inventory.list_hosts(pattern)

  5. 判断调用模块和模块参数是否合法

  6. 生成运行对象--运行的模块,参数(seconds参数对应ansible -B参数(后台运行长时任务),poll_interval对应ansible -P 表示对后台任务的轮询的间隔时间)

  7. 根据条件确定回调函数(ansible 命令返回结果的处理函数)

  8. 创建一个任务队列去运行paly(运行对象)

  9. 返回运行结果

主要来看一下生成 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



文章转载自进击的大杂烩,如果涉嫌侵权,请发送邮件至:contact@modb.pro进行举报,并提供相关证据,一经查实,墨天轮将立刻删除相关内容。

评论