|
|
@ -143,6 +143,7 @@ exclude_host_filters:
|
|
|
|
import hashlib
|
|
|
|
import hashlib
|
|
|
|
import json
|
|
|
|
import json
|
|
|
|
import re
|
|
|
|
import re
|
|
|
|
|
|
|
|
import uuid
|
|
|
|
|
|
|
|
|
|
|
|
try:
|
|
|
|
try:
|
|
|
|
from queue import Queue, Empty
|
|
|
|
from queue import Queue, Empty
|
|
|
@ -350,9 +351,10 @@ class InventoryModule(BaseInventoryPlugin, Constructable):
|
|
|
|
if next_link:
|
|
|
|
if next_link:
|
|
|
|
self._enqueue_get(url=next_link, api_version=self._compute_api_version, handler=self._on_vm_page_response)
|
|
|
|
self._enqueue_get(url=next_link, api_version=self._compute_api_version, handler=self._on_vm_page_response)
|
|
|
|
|
|
|
|
|
|
|
|
for h in response['value']:
|
|
|
|
if 'value' in response:
|
|
|
|
# FUTURE: add direct VM filtering by tag here (performance optimization)?
|
|
|
|
for h in response['value']:
|
|
|
|
self._hosts.append(AzureHost(h, self, vmss=vmss))
|
|
|
|
# FUTURE: add direct VM filtering by tag here (performance optimization)?
|
|
|
|
|
|
|
|
self._hosts.append(AzureHost(h, self, vmss=vmss))
|
|
|
|
|
|
|
|
|
|
|
|
def _on_vmss_page_response(self, response):
|
|
|
|
def _on_vmss_page_response(self, response):
|
|
|
|
next_link = response.get('nextLink')
|
|
|
|
next_link = response.get('nextLink')
|
|
|
@ -372,16 +374,16 @@ class InventoryModule(BaseInventoryPlugin, Constructable):
|
|
|
|
while True:
|
|
|
|
while True:
|
|
|
|
batch_requests = []
|
|
|
|
batch_requests = []
|
|
|
|
batch_item_index = 0
|
|
|
|
batch_item_index = 0
|
|
|
|
batch_response_handlers = []
|
|
|
|
batch_response_handlers = dict()
|
|
|
|
try:
|
|
|
|
try:
|
|
|
|
while batch_item_index < 500:
|
|
|
|
while batch_item_index < 100:
|
|
|
|
item = self._request_queue.get_nowait()
|
|
|
|
item = self._request_queue.get_nowait()
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
name = str(uuid.uuid4())
|
|
|
|
query_parameters = {'api-version': item.api_version}
|
|
|
|
query_parameters = {'api-version': item.api_version}
|
|
|
|
req = self._client.get(item.url, query_parameters)
|
|
|
|
req = self._client.get(item.url, query_parameters)
|
|
|
|
|
|
|
|
batch_requests.append(dict(httpMethod="GET", url=req.url, name=name))
|
|
|
|
batch_requests.append(dict(httpMethod="GET", url=req.url))
|
|
|
|
batch_response_handlers[name] = item
|
|
|
|
batch_response_handlers.append(item)
|
|
|
|
|
|
|
|
batch_item_index += 1
|
|
|
|
batch_item_index += 1
|
|
|
|
except Empty:
|
|
|
|
except Empty:
|
|
|
|
pass
|
|
|
|
pass
|
|
|
@ -391,15 +393,23 @@ class InventoryModule(BaseInventoryPlugin, Constructable):
|
|
|
|
|
|
|
|
|
|
|
|
batch_resp = self._send_batch(batch_requests)
|
|
|
|
batch_resp = self._send_batch(batch_requests)
|
|
|
|
|
|
|
|
|
|
|
|
for idx, r in enumerate(batch_resp['responses']):
|
|
|
|
key_name = None
|
|
|
|
|
|
|
|
if 'responses' in batch_resp:
|
|
|
|
|
|
|
|
key_name = 'responses'
|
|
|
|
|
|
|
|
elif 'value' in batch_resp:
|
|
|
|
|
|
|
|
key_name = 'value'
|
|
|
|
|
|
|
|
else:
|
|
|
|
|
|
|
|
raise AnsibleError("didn't find expected key responses/value in batch response")
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
for idx, r in enumerate(batch_resp[key_name]):
|
|
|
|
status_code = r.get('httpStatusCode')
|
|
|
|
status_code = r.get('httpStatusCode')
|
|
|
|
|
|
|
|
returned_name = r['name']
|
|
|
|
|
|
|
|
result = batch_response_handlers[returned_name]
|
|
|
|
if status_code != 200:
|
|
|
|
if status_code != 200:
|
|
|
|
# FUTURE: error-tolerant operation mode (eg, permissions)
|
|
|
|
# FUTURE: error-tolerant operation mode (eg, permissions)
|
|
|
|
raise AnsibleError("a batched request failed with status code {0}, url {1}".format(status_code, batch_requests[idx].get('url')))
|
|
|
|
raise AnsibleError("a batched request failed with status code {0}, url {1}".format(status_code, result.url))
|
|
|
|
|
|
|
|
|
|
|
|
item = batch_response_handlers[idx]
|
|
|
|
|
|
|
|
# FUTURE: store/handle errors from individual handlers
|
|
|
|
# FUTURE: store/handle errors from individual handlers
|
|
|
|
item.handler(r['content'], **item.handler_args)
|
|
|
|
result.handler(r['content'], **result.handler_args)
|
|
|
|
|
|
|
|
|
|
|
|
def _send_batch(self, batched_requests):
|
|
|
|
def _send_batch(self, batched_requests):
|
|
|
|
url = '/batch'
|
|
|
|
url = '/batch'
|
|
|
@ -409,8 +419,11 @@ class InventoryModule(BaseInventoryPlugin, Constructable):
|
|
|
|
|
|
|
|
|
|
|
|
body_content = self._serializer.body(body_obj, 'object')
|
|
|
|
body_content = self._serializer.body(body_obj, 'object')
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
header = {'x-ms-client-request-id': str(uuid.uuid4())}
|
|
|
|
|
|
|
|
header.update(self._default_header_parameters)
|
|
|
|
|
|
|
|
|
|
|
|
request = self._client.post(url, query_parameters)
|
|
|
|
request = self._client.post(url, query_parameters)
|
|
|
|
initial_response = self._client.send(request, self._default_header_parameters, body_content)
|
|
|
|
initial_response = self._client.send(request, header, body_content)
|
|
|
|
|
|
|
|
|
|
|
|
# FUTURE: configurable timeout?
|
|
|
|
# FUTURE: configurable timeout?
|
|
|
|
poller = ARMPolling(timeout=2)
|
|
|
|
poller = ARMPolling(timeout=2)
|
|
|
@ -499,7 +512,7 @@ class AzureHost(object):
|
|
|
|
pip_id = ipc['properties'].get('publicIPAddress', {}).get('id')
|
|
|
|
pip_id = ipc['properties'].get('publicIPAddress', {}).get('id')
|
|
|
|
if pip_id:
|
|
|
|
if pip_id:
|
|
|
|
pip = nic.public_ips[pip_id]
|
|
|
|
pip = nic.public_ips[pip_id]
|
|
|
|
new_hostvars['public_ipv4_addresses'].append(pip._pip_model['properties']['ipAddress'])
|
|
|
|
new_hostvars['public_ipv4_addresses'].append(pip._pip_model['properties'].get('ipAddress', None))
|
|
|
|
pip_fqdn = pip._pip_model['properties'].get('dnsSettings', {}).get('fqdn')
|
|
|
|
pip_fqdn = pip._pip_model['properties'].get('dnsSettings', {}).get('fqdn')
|
|
|
|
if pip_fqdn:
|
|
|
|
if pip_fqdn:
|
|
|
|
new_hostvars['public_dns_hostnames'].append(pip_fqdn)
|
|
|
|
new_hostvars['public_dns_hostnames'].append(pip_fqdn)
|
|
|
@ -514,8 +527,9 @@ class AzureHost(object):
|
|
|
|
for s in vm_instanceview_model.get('statuses', []) if self._powerstate_regex.match(s.get('code', ''))), 'unknown')
|
|
|
|
for s in vm_instanceview_model.get('statuses', []) if self._powerstate_regex.match(s.get('code', ''))), 'unknown')
|
|
|
|
|
|
|
|
|
|
|
|
def _on_nic_response(self, nic_model, is_primary=False):
|
|
|
|
def _on_nic_response(self, nic_model, is_primary=False):
|
|
|
|
nic = AzureNic(nic_model=nic_model, inventory_client=self._inventory_client, is_primary=is_primary)
|
|
|
|
if nic_model.get('type') == 'Microsoft.Network/networkInterfaces':
|
|
|
|
self.nics.append(nic)
|
|
|
|
nic = AzureNic(nic_model=nic_model, inventory_client=self._inventory_client, is_primary=is_primary)
|
|
|
|
|
|
|
|
self.nics.append(nic)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
class AzureNic(object):
|
|
|
|
class AzureNic(object):
|
|
|
|