summaryrefslogtreecommitdiff
path: root/docker
diff options
context:
space:
mode:
authorStephen Newey <github@s-n.me>2015-08-12 18:12:56 +0100
committerStephen Newey <github@s-n.me>2015-08-12 18:12:56 +0100
commit1c1d7eee5a7a583635892a2ba753a67c82a80325 (patch)
tree130220c40987db2fd1a849a796e72df3830a67d9 /docker
parent2febf104a02e5d7c17ba321cf179848f4c9a7945 (diff)
parentc697f0c26c0c06dd4626f4cc77c3a5e0e70494ee (diff)
downloaddocker-py-1c1d7eee5a7a583635892a2ba753a67c82a80325.tar.gz
Merge branch 'master' into exec_create_user
Diffstat (limited to 'docker')
-rw-r--r--docker/auth/__init__.py1
-rw-r--r--docker/auth/auth.py38
-rw-r--r--docker/client.py356
-rw-r--r--docker/clientbase.py277
-rw-r--r--docker/constants.py4
-rw-r--r--docker/errors.py4
-rw-r--r--docker/utils/__init__.py2
-rw-r--r--docker/utils/types.py3
-rw-r--r--docker/utils/utils.py79
-rw-r--r--docker/version.py2
10 files changed, 430 insertions, 336 deletions
diff --git a/docker/auth/__init__.py b/docker/auth/__init__.py
index d068b7f..6fc83f8 100644
--- a/docker/auth/__init__.py
+++ b/docker/auth/__init__.py
@@ -1,4 +1,5 @@
from .auth import (
+ INDEX_NAME,
INDEX_URL,
encode_header,
load_config,
diff --git a/docker/auth/auth.py b/docker/auth/auth.py
index 1c29615..4af741e 100644
--- a/docker/auth/auth.py
+++ b/docker/auth/auth.py
@@ -16,38 +16,34 @@ import base64
import fileinput
import json
import os
+import warnings
import six
-from ..utils import utils
+from .. import constants
from .. import errors
-INDEX_URL = 'https://index.docker.io/v1/'
+INDEX_NAME = 'index.docker.io'
+INDEX_URL = 'https://{0}/v1/'.format(INDEX_NAME)
DOCKER_CONFIG_FILENAME = os.path.join('.docker', 'config.json')
LEGACY_DOCKER_CONFIG_FILENAME = '.dockercfg'
-def expand_registry_url(hostname, insecure=False):
- if hostname.startswith('http:') or hostname.startswith('https:'):
- return hostname
- if utils.ping_registry('https://' + hostname):
- return 'https://' + hostname
- elif insecure:
- return 'http://' + hostname
- else:
- raise errors.DockerException(
- "HTTPS endpoint unresponsive and insecure mode isn't enabled."
+def resolve_repository_name(repo_name, insecure=False):
+ if insecure:
+ warnings.warn(
+ constants.INSECURE_REGISTRY_DEPRECATION_WARNING.format(
+ 'resolve_repository_name()'
+ ), DeprecationWarning
)
-
-def resolve_repository_name(repo_name, insecure=False):
if '://' in repo_name:
raise errors.InvalidRepository(
'Repository name cannot contain a scheme ({0})'.format(repo_name))
parts = repo_name.split('/', 1)
if '.' not in parts[0] and ':' not in parts[0] and parts[0] != 'localhost':
# This is a docker index repo (ex: foo/bar or ubuntu)
- return INDEX_URL, repo_name
+ return INDEX_NAME, repo_name
if len(parts) < 2:
raise errors.InvalidRepository(
'Invalid repository name ({0})'.format(repo_name))
@@ -57,7 +53,7 @@ def resolve_repository_name(repo_name, insecure=False):
'Invalid repository name, try "{0}" instead'.format(parts[1])
)
- return expand_registry_url(parts[0], insecure), parts[1]
+ return parts[0], parts[1]
def resolve_authconfig(authconfig, registry=None):
@@ -68,7 +64,7 @@ def resolve_authconfig(authconfig, registry=None):
Returns None if no match was found.
"""
# Default to the public index server
- registry = convert_to_hostname(registry) if registry else INDEX_URL
+ registry = convert_to_hostname(registry) if registry else INDEX_NAME
if registry in authconfig:
return authconfig[registry]
@@ -102,12 +98,6 @@ def encode_header(auth):
return base64.b64encode(auth_json)
-def encode_full_header(auth):
- """ Returns the given auth block encoded for the X-Registry-Config header.
- """
- return encode_header({'configs': auth})
-
-
def parse_auth(entries):
"""
Parses authentication entries
@@ -185,7 +175,7 @@ def load_config(config_path=None):
'Invalid or empty configuration file!')
username, password = decode_auth(data[0])
- conf[INDEX_URL] = {
+ conf[INDEX_NAME] = {
'username': username,
'password': password,
'email': data[1],
diff --git a/docker/client.py b/docker/client.py
index 74c33a2..e4a6735 100644
--- a/docker/client.py
+++ b/docker/client.py
@@ -12,237 +12,23 @@
# See the License for the specific language governing permissions and
# limitations under the License.
-import json
import os
import re
import shlex
-import struct
import warnings
from datetime import datetime
-import requests
-import requests.exceptions
import six
-import websocket
-
+from . import clientbase
from . import constants
from . import errors
from .auth import auth
-from .unixconn import unixconn
-from .ssladapter import ssladapter
from .utils import utils, check_resource
-from .tls import TLSConfig
-
-
-class Client(requests.Session):
- def __init__(self, base_url=None, version=None,
- timeout=constants.DEFAULT_TIMEOUT_SECONDS, tls=False):
- super(Client, self).__init__()
-
- if tls and not base_url.startswith('https://'):
- raise errors.TLSParameterError(
- 'If using TLS, the base_url argument must begin with '
- '"https://".')
-
- self.base_url = base_url
- self.timeout = timeout
-
- self._auth_configs = auth.load_config()
-
- base_url = utils.parse_host(base_url)
- if base_url.startswith('http+unix://'):
- unix_socket_adapter = unixconn.UnixAdapter(base_url, timeout)
- self.mount('http+docker://', unix_socket_adapter)
- self.base_url = 'http+docker://localunixsocket'
- else:
- # Use SSLAdapter for the ability to specify SSL version
- if isinstance(tls, TLSConfig):
- tls.configure_client(self)
- elif tls:
- self.mount('https://', ssladapter.SSLAdapter())
- self.base_url = base_url
-
- # version detection needs to be after unix adapter mounting
- if version is None:
- self._version = constants.DEFAULT_DOCKER_API_VERSION
- elif isinstance(version, six.string_types):
- if version.lower() == 'auto':
- self._version = self._retrieve_server_version()
- else:
- self._version = version
- else:
- raise errors.DockerException(
- 'Version parameter must be a string or None. Found {0}'.format(
- type(version).__name__
- )
- )
-
- def _retrieve_server_version(self):
- try:
- return self.version(api_version=False)["ApiVersion"]
- except KeyError:
- raise errors.DockerException(
- 'Invalid response from docker daemon: key "ApiVersion"'
- ' is missing.'
- )
- except Exception as e:
- raise errors.DockerException(
- 'Error while fetching server API version: {0}'.format(e)
- )
-
- def _set_request_timeout(self, kwargs):
- """Prepare the kwargs for an HTTP request by inserting the timeout
- parameter, if not already present."""
- kwargs.setdefault('timeout', self.timeout)
- return kwargs
+from .constants import INSECURE_REGISTRY_DEPRECATION_WARNING
- def _post(self, url, **kwargs):
- return self.post(url, **self._set_request_timeout(kwargs))
-
- def _get(self, url, **kwargs):
- return self.get(url, **self._set_request_timeout(kwargs))
-
- def _delete(self, url, **kwargs):
- return self.delete(url, **self._set_request_timeout(kwargs))
-
- def _url(self, path, versioned_api=True):
- if versioned_api:
- return '{0}/v{1}{2}'.format(self.base_url, self._version, path)
- else:
- return '{0}{1}'.format(self.base_url, path)
-
- def _raise_for_status(self, response, explanation=None):
- """Raises stored :class:`APIError`, if one occurred."""
- try:
- response.raise_for_status()
- except requests.exceptions.HTTPError as e:
- raise errors.APIError(e, response, explanation=explanation)
-
- def _result(self, response, json=False, binary=False):
- assert not (json and binary)
- self._raise_for_status(response)
-
- if json:
- return response.json()
- if binary:
- return response.content
- return response.text
-
- def _post_json(self, url, data, **kwargs):
- # Go <1.1 can't unserialize null to a string
- # so we do this disgusting thing here.
- data2 = {}
- if data is not None:
- for k, v in six.iteritems(data):
- if v is not None:
- data2[k] = v
-
- if 'headers' not in kwargs:
- kwargs['headers'] = {}
- kwargs['headers']['Content-Type'] = 'application/json'
- return self._post(url, data=json.dumps(data2), **kwargs)
-
- def _attach_params(self, override=None):
- return override or {
- 'stdout': 1,
- 'stderr': 1,
- 'stream': 1
- }
-
- @check_resource
- def _attach_websocket(self, container, params=None):
- url = self._url("/containers/{0}/attach/ws".format(container))
- req = requests.Request("POST", url, params=self._attach_params(params))
- full_url = req.prepare().url
- full_url = full_url.replace("http://", "ws://", 1)
- full_url = full_url.replace("https://", "wss://", 1)
- return self._create_websocket_connection(full_url)
-
- def _create_websocket_connection(self, url):
- return websocket.create_connection(url)
-
- def _get_raw_response_socket(self, response):
- self._raise_for_status(response)
- if six.PY3:
- sock = response.raw._fp.fp.raw
- else:
- sock = response.raw._fp.fp._sock
- try:
- # Keep a reference to the response to stop it being garbage
- # collected. If the response is garbage collected, it will
- # close TLS sockets.
- sock._response = response
- except AttributeError:
- # UNIX sockets can't have attributes set on them, but that's
- # fine because we won't be doing TLS over them
- pass
-
- return sock
-
- def _stream_helper(self, response, decode=False):
- """Generator for data coming from a chunked-encoded HTTP response."""
- if response.raw._fp.chunked:
- reader = response.raw
- while not reader.closed:
- # this read call will block until we get a chunk
- data = reader.read(1)
- if not data:
- break
- if reader._fp.chunk_left:
- data += reader.read(reader._fp.chunk_left)
- if decode:
- if six.PY3:
- data = data.decode('utf-8')
- data = json.loads(data)
- yield data
- else:
- # Response isn't chunked, meaning we probably
- # encountered an error immediately
- yield self._result(response)
-
- def _multiplexed_buffer_helper(self, response):
- """A generator of multiplexed data blocks read from a buffered
- response."""
- buf = self._result(response, binary=True)
- walker = 0
- while True:
- if len(buf[walker:]) < 8:
- break
- _, length = struct.unpack_from('>BxxxL', buf[walker:])
- start = walker + constants.STREAM_HEADER_SIZE_BYTES
- end = start + length
- walker = end
- yield buf[start:end]
-
- def _multiplexed_response_stream_helper(self, response):
- """A generator of multiplexed data blocks coming from a response
- stream."""
-
- # Disable timeout on the underlying socket to prevent
- # Read timed out(s) for long running processes
- socket = self._get_raw_response_socket(response)
- if six.PY3:
- socket._sock.settimeout(None)
- else:
- socket.settimeout(None)
-
- while True:
- header = response.raw.read(constants.STREAM_HEADER_SIZE_BYTES)
- if not header:
- break
- _, length = struct.unpack('>BxxxL', header)
- if not length:
- break
- data = response.raw.read(length)
- if not data:
- break
- yield data
-
- @property
- def api_version(self):
- return self._version
+class Client(clientbase.ClientBase):
@check_resource
def attach(self, container, stdout=True, stderr=True,
stream=False, logs=False):
@@ -255,28 +41,7 @@ class Client(requests.Session):
u = self._url("/containers/{0}/attach".format(container))
response = self._post(u, params=params, stream=stream)
- # Stream multi-plexing was only introduced in API v1.6. Anything before
- # that needs old-style streaming.
- if utils.compare_version('1.6', self._version) < 0:
- def stream_result():
- self._raise_for_status(response)
- for line in response.iter_lines(chunk_size=1,
- decode_unicode=True):
- # filter out keep-alive new lines
- if line:
- yield line
-
- return stream_result() if stream else \
- self._result(response, binary=True)
-
- sep = bytes() if six.PY3 else str()
-
- if stream:
- return self._multiplexed_response_stream_helper(response)
- else:
- return sep.join(
- [x for x in self._multiplexed_buffer_helper(response)]
- )
+ return self._get_result(container, stream, response)
@check_resource
def attach_socket(self, container, params=None, ws=False):
@@ -317,7 +82,7 @@ class Client(requests.Session):
elif fileobj is not None:
context = utils.mkbuildcontext(fileobj)
elif path.startswith(('http://', 'https://',
- 'git://', 'github.com/')):
+ 'git://', 'github.com/', 'git@')):
remote = path
elif not os.path.isdir(path):
raise TypeError("You must specify a directory to build in path")
@@ -375,9 +140,14 @@ class Client(requests.Session):
if self._auth_configs:
if headers is None:
headers = {}
- headers['X-Registry-Config'] = auth.encode_full_header(
- self._auth_configs
- )
+ if utils.compare_version('1.19', self._version) >= 0:
+ headers['X-Registry-Config'] = auth.encode_header(
+ self._auth_configs
+ )
+ else:
+ headers['X-Registry-Config'] = auth.encode_header({
+ 'configs': self._auth_configs
+ })
response = self._post(
u,
@@ -450,11 +220,11 @@ class Client(requests.Session):
def create_container(self, image, command=None, hostname=None, user=None,
detach=False, stdin_open=False, tty=False,
- mem_limit=0, ports=None, environment=None, dns=None,
- volumes=None, volumes_from=None,
+ mem_limit=None, ports=None, environment=None,
+ dns=None, volumes=None, volumes_from=None,
network_disabled=False, name=None, entrypoint=None,
cpu_shares=None, working_dir=None, domainname=None,
- memswap_limit=0, cpuset=None, host_config=None,
+ memswap_limit=None, cpuset=None, host_config=None,
mac_address=None, labels=None, volume_driver=None):
if isinstance(volumes, six.string_types):
@@ -503,23 +273,12 @@ class Client(requests.Session):
'filters': filters
}
- return self._stream_helper(self.get(self._url('/events'),
- params=params, stream=True),
- decode=decode)
-
- @check_resource
- def execute(self, container, cmd, detach=False, stdout=True, stderr=True,
- stream=False, tty=False):
- warnings.warn(
- 'Client.execute is being deprecated. Please use exec_create & '
- 'exec_start instead', DeprecationWarning
+ return self._stream_helper(
+ self.get(self._url('/events'), params=params, stream=True),
+ decode=decode
)
- create_res = self.exec_create(
- container, cmd, stdout, stderr, tty
- )
-
- return self.exec_start(create_res, detach, tty, stream)
+ @check_resource
def exec_create(self, container, cmd, stdout=True, stderr=True, tty=False,
privileged=False, user=''):
if utils.compare_version('1.15', self._version) < 0:
@@ -582,17 +341,7 @@ class Client(requests.Session):
res = self._post_json(self._url('/exec/{0}/start'.format(exec_id)),
data=data, stream=stream)
- self._raise_for_status(res)
- if stream:
- return self._multiplexed_response_stream_helper(res)
- elif six.PY3:
- return bytes().join(
- [x for x in self._multiplexed_buffer_helper(res)]
- )
- else:
- return str().join(
- [x for x in self._multiplexed_buffer_helper(res)]
- )
+ return self._get_result_tty(stream, res, tty)
@check_resource
def export(self, container):
@@ -740,7 +489,9 @@ class Client(requests.Session):
@check_resource
def inspect_image(self, image):
return self._result(
- self._get(self._url("/images/{0}/json".format(image))),
+ self._get(
+ self._url("/images/{0}/json".format(image.replace('/', '%2F')))
+ ),
True
)
@@ -760,6 +511,12 @@ class Client(requests.Session):
def login(self, username, password=None, email=None, registry=None,
reauth=False, insecure_registry=False, dockercfg_path=None):
+ if insecure_registry:
+ warnings.warn(
+ INSECURE_REGISTRY_DEPRECATION_WARNING.format('login()'),
+ DeprecationWarning
+ )
+
# If we don't have any auth data so far, try reloading the config file
# one more time in case anything showed up in there.
# If dockercfg_path is passed check to see if the config file exists,
@@ -805,16 +562,7 @@ class Client(requests.Session):
params['tail'] = tail
url = self._url("/containers/{0}/logs".format(container))
res = self._get(url, params=params, stream=stream)
- if stream:
- return self._multiplexed_response_stream_helper(res)
- elif six.PY3:
- return bytes().join(
- [x for x in self._multiplexed_buffer_helper(res)]
- )
- else:
- return str().join(
- [x for x in self._multiplexed_buffer_helper(res)]
- )
+ return self._get_result(container, stream, res)
return self.attach(
container,
stdout=stdout,
@@ -854,11 +602,15 @@ class Client(requests.Session):
def pull(self, repository, tag=None, stream=False,
insecure_registry=False, auth_config=None):
+ if insecure_registry:
+ warnings.warn(
+ INSECURE_REGISTRY_DEPRECATION_WARNING.format('pull()'),
+ DeprecationWarning
+ )
+
if not tag:
repository, tag = utils.parse_repository_tag(repository)
- registry, repo_name = auth.resolve_repository_name(
- repository, insecure=insecure_registry
- )
+ registry, repo_name = auth.resolve_repository_name(repository)
if repo_name.count(":") == 1:
repository, tag = repository.rsplit(":", 1)
@@ -901,11 +653,15 @@ class Client(requests.Session):
def push(self, repository, tag=None, stream=False,
insecure_registry=False):
+ if insecure_registry:
+ warnings.warn(
+ INSECURE_REGISTRY_DEPRECATION_WARNING.format('push()'),
+ DeprecationWarning
+ )
+
if not tag:
repository, tag = utils.parse_repository_tag(repository)
- registry, repo_name = auth.resolve_repository_name(
- repository, insecure=insecure_registry
- )
+ registry, repo_name = auth.resolve_repository_name(repository)
u = self._url("/images/{0}/push".format(repository))
params = {
'tag': tag
@@ -981,7 +737,7 @@ class Client(requests.Session):
@check_resource
def start(self, container, binds=None, port_bindings=None, lxc_conf=None,
- publish_all_ports=False, links=None, privileged=False,
+ publish_all_ports=None, links=None, privileged=None,
dns=None, dns_search=None, volumes_from=None, network_mode=None,
restart_policy=None, cap_add=None, cap_drop=None, devices=None,
extra_hosts=None, read_only=None, pid_mode=None, ipc_mode=None,
@@ -1023,7 +779,7 @@ class Client(requests.Session):
'ulimits is only supported for API version >= 1.18'
)
- start_config = utils.create_host_config(
+ start_config_kwargs = dict(
binds=binds, port_bindings=port_bindings, lxc_conf=lxc_conf,
publish_all_ports=publish_all_ports, links=links, dns=dns,
privileged=privileged, dns_search=dns_search, cap_add=cap_add,
@@ -1032,16 +788,18 @@ class Client(requests.Session):
extra_hosts=extra_hosts, read_only=read_only, pid_mode=pid_mode,
ipc_mode=ipc_mode, security_opt=security_opt, ulimits=ulimits
)
+ start_config = None
+
+ if any(v is not None for v in start_config_kwargs.values()):
+ if utils.compare_version('1.15', self._version) > 0:
+ warnings.warn(
+ 'Passing host config parameters in start() is deprecated. '
+ 'Please use host_config in create_container instead!',
+ DeprecationWarning
+ )
+ start_config = utils.create_host_config(**start_config_kwargs)
url = self._url("/containers/{0}/start".format(container))
- if not start_config:
- start_config = None
- elif utils.compare_version('1.15', self._version) > 0:
- warnings.warn(
- 'Passing host config parameters in start() is deprecated. '
- 'Please use host_config in create_container instead!',
- DeprecationWarning
- )
res = self._post_json(url, data=start_config)
self._raise_for_status(res)
diff --git a/docker/clientbase.py b/docker/clientbase.py
new file mode 100644
index 0000000..ce52ffa
--- /dev/null
+++ b/docker/clientbase.py
@@ -0,0 +1,277 @@
+import json
+import struct
+
+import requests
+import requests.exceptions
+import six
+import websocket
+
+
+from . import constants
+from . import errors
+from .auth import auth
+from .unixconn import unixconn
+from .ssladapter import ssladapter
+from .utils import utils, check_resource
+from .tls import TLSConfig
+
+
+class ClientBase(requests.Session):
+ def __init__(self, base_url=None, version=None,
+ timeout=constants.DEFAULT_TIMEOUT_SECONDS, tls=False):
+ super(ClientBase, self).__init__()
+
+ if tls and not base_url.startswith('https://'):
+ raise errors.TLSParameterError(
+ 'If using TLS, the base_url argument must begin with '
+ '"https://".')
+
+ self.base_url = base_url
+ self.timeout = timeout
+
+ self._auth_configs = auth.load_config()
+
+ base_url = utils.parse_host(base_url)
+ if base_url.startswith('http+unix://'):
+ self._custom_adapter = unixconn.UnixAdapter(base_url, timeout)
+ self.mount('http+docker://', self._custom_adapter)
+ self.base_url = 'http+docker://localunixsocket'
+ else:
+ # Use SSLAdapter for the ability to specify SSL version
+ if isinstance(tls, TLSConfig):
+ tls.configure_client(self)
+ elif tls:
+ self._custom_adapter = ssladapter.SSLAdapter()
+ self.mount('https://', self._custom_adapter)
+ self.base_url = base_url
+
+ # version detection needs to be after unix adapter mounting
+ if version is None:
+ self._version = constants.DEFAULT_DOCKER_API_VERSION
+ elif isinstance(version, six.string_types):
+ if version.lower() == 'auto':
+ self._version = self._retrieve_server_version()
+ else:
+ self._version = version
+ else:
+ raise errors.DockerException(
+ 'Version parameter must be a string or None. Found {0}'.format(
+ type(version).__name__
+ )
+ )
+
+ def _retrieve_server_version(self):
+ try:
+ return self.version(api_version=False)["ApiVersion"]
+ except KeyError:
+ raise errors.DockerException(
+ 'Invalid response from docker daemon: key "ApiVersion"'
+ ' is missing.'
+ )
+ except Exception as e:
+ raise errors.DockerException(
+ 'Error while fetching server API version: {0}'.format(e)
+ )
+
+ def _set_request_timeout(self, kwargs):
+ """Prepare the kwargs for an HTTP request by inserting the timeout
+ parameter, if not already present."""
+ kwargs.setdefault('timeout', self.timeout)
+ return kwargs
+
+ def _post(self, url, **kwargs):
+ return self.post(url, **self._set_request_timeout(kwargs))
+
+ def _get(self, url, **kwargs):
+ return self.get(url, **self._set_request_timeout(kwargs))
+
+ def _delete(self, url, **kwargs):
+ return self.delete(url, **self._set_request_timeout(kwargs))
+
+ def _url(self, path, versioned_api=True):
+ if versioned_api:
+ return '{0}/v{1}{2}'.format(self.base_url, self._version, path)
+ else:
+ return '{0}{1}'.format(self.base_url, path)
+
+ def _raise_for_status(self, response, explanation=None):
+ """Raises stored :class:`APIError`, if one occurred."""
+ try:
+ response.raise_for_status()
+ except requests.exceptions.HTTPError as e:
+ if e.response.status_code == 404:
+ raise errors.NotFound(e, response, explanation=explanation)
+ raise errors.APIError(e, response, explanation=explanation)
+
+ def _result(self, response, json=False, binary=False):
+ assert not (json and binary)
+ self._raise_for_status(response)
+
+ if json:
+ return response.json()
+ if binary:
+ return response.content
+ return response.text
+
+ def _post_json(self, url, data, **kwargs):
+ # Go <1.1 can't unserialize null to a string
+ # so we do this disgusting thing here.
+ data2 = {}
+ if data is not None:
+ for k, v in six.iteritems(data):
+ if v is not None:
+ data2[k] = v
+
+ if 'headers' not in kwargs:
+ kwargs['headers'] = {}
+ kwargs['headers']['Content-Type'] = 'application/json'
+ return self._post(url, data=json.dumps(data2), **kwargs)
+
+ def _attach_params(self, override=None):
+ return override or {
+ 'stdout': 1,
+ 'stderr': 1,
+ 'stream': 1
+ }
+
+ @check_resource
+ def _attach_websocket(self, container, params=None):
+ url = self._url("/containers/{0}/attach/ws".format(container))
+ req = requests.Request("POST", url, params=self._attach_params(params))
+ full_url = req.prepare().url
+ full_url = full_url.replace("http://", "ws://", 1)
+ full_url = full_url.replace("https://", "wss://", 1)
+ return self._create_websocket_connection(full_url)
+
+ def _create_websocket_connection(self, url):
+ return websocket.create_connection(url)
+
+ def _get_raw_response_socket(self, response):
+ self._raise_for_status(response)
+ if six.PY3:
+ sock = response.raw._fp.fp.raw
+ else:
+ sock = response.raw._fp.fp._sock
+ try:
+ # Keep a reference to the response to stop it being garbage
+ # collected. If the response is garbage collected, it will
+ # close TLS sockets.
+ sock._response = response
+ except AttributeError:
+ # UNIX sockets can't have attributes set on them, but that's
+ # fine because we won't be doing TLS over them
+ pass
+
+ return sock
+
+ def _stream_helper(self, response, decode=False):
+ """Generator for data coming from a chunked-encoded HTTP response."""
+ if response.raw._fp.chunked:
+ reader = response.raw
+ while not reader.closed:
+ # this read call will block until we get a chunk
+ data = reader.read(1)
+ if not data:
+ break
+ if reader._fp.chunk_left:
+ data += reader.read(reader._fp.chunk_left)
+ if decode:
+ if six.PY3:
+ data = data.decode('utf-8')
+ data = json.loads(data)
+ yield data
+ else:
+ # Response isn't chunked, meaning we probably
+ # encountered an error immediately
+ yield self._result(response)
+
+ def _multiplexed_buffer_helper(self, response):
+ """A generator of multiplexed data blocks read from a buffered
+ response."""
+ buf = self._result(response, binary=True)
+ walker = 0
+ while True:
+ if len(buf[walker:]) < 8:
+ break
+ _, length = struct.unpack_from('>BxxxL', buf[walker:])
+ start = walker + constants.STREAM_HEADER_SIZE_BYTES
+ end = start + length
+ walker = end
+ yield buf[start:end]
+
+ def _multiplexed_response_stream_helper(self, response):
+ """A generator of multiplexed data blocks coming from a response
+ stream."""
+
+ # Disable timeout on the underlying socket to prevent
+ # Read timed out(s) for long running processes
+ socket = self._get_raw_response_socket(response)
+ if six.PY3:
+ socket._sock.settimeout(None)
+ else:
+ socket.settimeout(None)
+
+ while True:
+ header = response.raw.read(constants.STREAM_HEADER_SIZE_BYTES)
+ if not header:
+ break
+ _, length = struct.unpack('>BxxxL', header)
+ if not length:
+ continue
+ data = response.raw.read(length)
+ if not data:
+ break
+ yield data
+
+ def _stream_raw_result_old(self, response):
+ ''' Stream raw output for API versions below 1.6 '''
+ self._raise_for_status(response)
+ for line in response.iter_lines(chunk_size=1,
+ decode_unicode=True):
+ # filter out keep-alive new lines
+ if line:
+ yield line
+
+ def _stream_raw_result(self, response):
+ ''' Stream result for TTY-enabled container above API 1.6 '''
+ self._raise_for_status(response)
+ for out in response.iter_content(chunk_size=1, decode_unicode=True):
+ yield out
+
+ def _get_result(self, container, stream, res):
+ cont = self.inspect_container(container)
+ return self._get_result_tty(stream, res, cont['Config']['Tty'])
+
+ def _get_result_tty(self, stream, res, is_tty):
+ # Stream multi-plexing was only introduced in API v1.6. Anything
+ # before that needs old-style streaming.
+ if utils.compare_version('1.6', self._version) < 0:
+ return self._stream_raw_result_old(res)
+
+ # We should also use raw streaming (without keep-alives)
+ # if we're dealing with a tty-enabled container.
+ if is_tty:
+ return self._stream_raw_result(res) if stream else \
+ self._result(res, binary=True)
+
+ self._raise_for_status(res)
+ sep = six.binary_type()
+ if stream:
+ return self._multiplexed_response_stream_helper(res)
+ else:
+ return sep.join(
+ [x for x in self._multiplexed_buffer_helper(res)]
+ )
+
+ def get_adapter(self, url):
+ try:
+ return super(ClientBase, self).get_adapter(url)
+ except requests.exceptions.InvalidSchema as e:
+ if self._custom_adapter:
+ return self._custom_adapter
+ else:
+ raise e
+
+ @property
+ def api_version(self):
+ return self._version
diff --git a/docker/constants.py b/docker/constants.py
index f99f192..10a2fee 100644
--- a/docker/constants.py
+++ b/docker/constants.py
@@ -4,3 +4,7 @@ STREAM_HEADER_SIZE_BYTES = 8
CONTAINER_LIMITS_KEYS = [
'memory', 'memswap', 'cpushares', 'cpusetcpus'
]
+
+INSECURE_REGISTRY_DEPRECATION_WARNING = \
+ 'The `insecure_registry` argument to {} ' \
+ 'is deprecated and non-functional. Please remove it.'
diff --git a/docker/errors.py b/docker/errors.py
index d15e332..066406a 100644
--- a/docker/errors.py
+++ b/docker/errors.py
@@ -53,6 +53,10 @@ class DockerException(Exception):
pass
+class NotFound(APIError):
+ pass
+
+
class InvalidVersion(DockerException):
pass
diff --git a/docker/utils/__init__.py b/docker/utils/__init__.py
index 81cc8a6..6189ed8 100644
--- a/docker/utils/__init__.py
+++ b/docker/utils/__init__.py
@@ -2,7 +2,7 @@ from .utils import (
compare_version, convert_port_bindings, convert_volume_binds,
mkbuildcontext, tar, parse_repository_tag, parse_host,
kwargs_from_env, convert_filters, create_host_config,
- create_container_config, parse_bytes, ping_registry
+ create_container_config, parse_bytes, ping_registry, parse_env_file
) # flake8: noqa
from .types import Ulimit, LogConfig # flake8: noqa
diff --git a/docker/utils/types.py b/docker/utils/types.py
index d742fd0..ca67467 100644
--- a/docker/utils/types.py
+++ b/docker/utils/types.py
@@ -5,9 +5,10 @@ class LogConfigTypesEnum(object):
_values = (
'json-file',
'syslog',
+ 'journald',
'none'
)
- JSON, SYSLOG, NONE = _values
+ JSON, SYSLOG, JOURNALD, NONE = _values
class DictType(dict):
diff --git a/docker/utils/utils.py b/docker/utils/utils.py
index e4e665f..d979c96 100644
--- a/docker/utils/utils.py
+++ b/docker/utils/utils.py
@@ -19,6 +19,7 @@ import json
import shlex
import tarfile
import tempfile
+import warnings
from distutils.version import StrictVersion
from fnmatch import fnmatch
from datetime import datetime
@@ -120,6 +121,11 @@ def compare_version(v1, v2):
def ping_registry(url):
+ warnings.warn(
+ 'The `ping_registry` method is deprecated and will be removed.',
+ DeprecationWarning
+ )
+
return ping(url + '/v2/', [401]) or ping(url + '/v1/_ping')
@@ -333,9 +339,9 @@ def convert_filters(filters):
return json.dumps(result)
-def datetime_to_timestamp(dt=datetime.now()):
- """Convert a datetime in local timezone to a unix timestamp"""
- delta = dt - datetime.fromtimestamp(0)
+def datetime_to_timestamp(dt):
+ """Convert a UTC datetime to a Unix timestamp"""
+ delta = dt - datetime.utcfromtimestamp(0)
return delta.seconds + delta.days * 24 * 3600
@@ -383,10 +389,21 @@ def create_host_config(
dns=None, dns_search=None, volumes_from=None, network_mode=None,
restart_policy=None, cap_add=None, cap_drop=None, devices=None,
extra_hosts=None, read_only=None, pid_mode=None, ipc_mode=None,
- security_opt=None, ulimits=None, log_config=None
+ security_opt=None, ulimits=None, log_config=None, mem_limit=None,
+ memswap_limit=None
):
host_config = {}
+ if mem_limit is not None:
+ if isinstance(mem_limit, six.string_types):
+ mem_limit = parse_bytes(mem_limit)
+ host_config['Memory'] = mem_limit
+
+ if memswap_limit is not None:
+ if isinstance(memswap_limit, six.string_types):
+ memswap_limit = parse_bytes(memswap_limit)
+ host_config['MemorySwap'] = memswap_limit
+
if pid_mode not in (None, 'host'):
raise errors.DockerException(
'Invalid value for pid param: {0}'.format(pid_mode)
@@ -411,6 +428,8 @@ def create_host_config(
if network_mode:
host_config['NetworkMode'] = network_mode
+ elif network_mode is None:
+ host_config['NetworkMode'] = 'default'
if restart_policy:
host_config['RestartPolicy'] = restart_policy
@@ -501,16 +520,42 @@ def create_host_config(
return host_config
+def parse_env_file(env_file):
+ """
+ Reads a line-separated environment file.
+ The format of each line should be "key=value".
+ """
+ environment = {}
+
+ with open(env_file, 'r') as f:
+ for line in f:
+
+ if line[0] == '#':
+ continue
+
+ parse_line = line.strip().split('=')
+ if len(parse_line) == 2:
+ k, v = parse_line
+ environment[k] = v
+ else:
+ raise errors.DockerException(
+ 'Invalid line in environment file {0}:\n{1}'.format(
+ env_file, line))
+
+ return environment
+
+
def create_container_config(
version, image, command, hostname=None, user=None, detach=False,
- stdin_open=False, tty=False, mem_limit=0, ports=None, environment=None,
+ stdin_open=False, tty=False, mem_limit=None, ports=None, environment=None,
dns=None, volumes=None, volumes_from=None, network_disabled=False,
entrypoint=None, cpu_shares=None, working_dir=None, domainname=None,
- memswap_limit=0, cpuset=None, host_config=None, mac_address=None,
+ memswap_limit=None, cpuset=None, host_config=None, mac_address=None,
labels=None, volume_driver=None
):
if isinstance(command, six.string_types):
command = shlex.split(str(command))
+
if isinstance(environment, dict):
environment = [
six.text_type('{0}={1}').format(k, v)
@@ -522,10 +567,24 @@ def create_container_config(
'labels were only introduced in API version 1.18'
)
- if volume_driver is not None and compare_version('1.19', version) < 0:
- raise errors.InvalidVersion(
- 'Volume drivers were only introduced in API version 1.19'
- )
+ if compare_version('1.19', version) < 0:
+ if volume_driver is not None:
+ raise errors.InvalidVersion(
+ 'Volume drivers were only introduced in API version 1.19'
+ )
+ mem_limit = mem_limit if mem_limit is not None else 0
+ memswap_limit = memswap_limit if memswap_limit is not None else 0
+ else:
+ if mem_limit is not None:
+ raise errors.InvalidVersion(
+ 'mem_limit has been moved to host_config in API version 1.19'
+ )
+
+ if memswap_limit is not None:
+ raise errors.InvalidVersion(
+ 'memswap_limit has been moved to host_config in API '
+ 'version 1.19'
+ )
if isinstance(labels, list):
labels = dict((lbl, six.text_type('')) for lbl in labels)
diff --git a/docker/version.py b/docker/version.py
index 88859a6..d0aad76 100644
--- a/docker/version.py
+++ b/docker/version.py
@@ -1,2 +1,2 @@
-version = "1.3.0-dev"
+version = "1.4.0-dev"
version_info = tuple([int(d) for d in version.split("-")[0].split(".")])