#-------------------------------------------------------------------------------
# Copyright (c) 2012-2013 by European Organization for Nuclear Research (CERN)
# Author: Justin Salmon <jsalmon@cern.ch>
#-------------------------------------------------------------------------------
# XRootD is free software: you can redistribute it and/or modify
# it under the terms of the GNU Lesser General Public License as published by
# the Free Software Foundation, either version 3 of the License, or
# (at your option) any later version.
#
# XRootD is distributed in the hope that it will be useful,
# but WITHOUT ANY WARRANTY; without even the implied warranty of
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
# GNU General Public License for more details.
#
# You should have received a copy of the GNU Lesser General Public License
# along with XRootD. If not, see <http:#www.gnu.org/licenses/>.
#-------------------------------------------------------------------------------
from __future__ import absolute_import, division, print_function
try:
from urllib.parse import urlparse
except ImportError:
from urlparse import urlparse
from XRootD.client.url import URL
[docs]
class XRootDError(RuntimeError):
"""Base exception raised from unsuccessful :class:`XRootDStatus` objects."""
def __init__(self, status):
self.status = status
RuntimeError.__init__(self, str(status))
[docs]
class XRootDNotFoundError(XRootDError):
"""The requested file or resource was not found."""
[docs]
class XRootDAuthorizationError(XRootDError):
"""Authentication or authorization failed."""
[docs]
class XRootDTimeoutError(XRootDError):
"""The request timed out or expired."""
[docs]
class XRootDChecksumError(XRootDError):
"""The request failed checksum validation."""
[docs]
class XRootDOperationError(XRootDError):
"""Generic unsuccessful XRootD operation."""
class Struct(object):
"""Convert a dict into an object by adding each dict entry to __dict__"""
def __init__(self, entries):
self.__dict__.update(**entries)
def __repr__(self):
return '<%s>' % str(', '.join('%s: %s' % (k, repr(v))
for (k, v) in self.__dict__.items()))
[docs]
class LocationInfo(Struct):
"""Path location information (a list of discovered file locations).
:param locations: (List of :mod:`XRootD.client.responses.Location` objects)
List of discovered locations
This object is iterable::
>>> status, locations = filesystem.locate('/tmp', OpenFlags.REFRESH)
>>> print locations
<XRootD.client.responses.LocationInfo object at 0x288b9f0>
>>> for location in locations:
... print location.address
...
[::127.0.0.1]:1094
"""
def __init__(self, locations):
super(LocationInfo, self).__init__({'locations':
[Location(l) for l in locations]})
def __iter__(self):
return iter(self.locations)
[docs]
class Location(Struct):
"""Information about a single location.
:var address: The address of this location
:var type: The type of this location, one of
:mod:`XRootD.client.flags.LocationType`
:var accesstype: The allowed access type of this location, one of
:mod:`XRootD.client.flags.AccessType`
:var is_manager: Is the location a manager
:var is_server: Is the location a server
"""
def __init__(self, location):
super(Location, self).__init__(location)
[docs]
class XRootDStatus(Struct):
"""Status of a request. Returned with all requests.
:var message: Message describing the status of this request
:var ok: The request was successful
:var error: Error making request
:var fatal: Fatal error making request
:var status: Status of the request
:var code: Error type, or additional hints on what to do
:var shellcode: Status code that may be returned to the shell
:var errno: Errno, if any
"""
#----------------------------------------------------------------------------
# Additional info for the stOK status
#----------------------------------------------------------------------------
suDone = 0
suContinue = 1
suRetry = 2
suPartial = 3
suAlreadyDone = 4
suNotStarted = 5
#----------------------------------------------------------------------------
# Generic errors
#----------------------------------------------------------------------------
errNone = 0; # No error
errRetry = 1; # Try again for whatever reason
errUnknown = 2; # Unknown error
errInvalidOp = 3; # The operation cannot be performed in the
# given circumstances
errFcntl = 4; # failed manipulate file descriptor
errPoll = 5; # error while polling descriptors
errConfig = 6; # System misconfigured
errInternal = 7; # Internal error
errUnknownCommand = 8;
errInvalidArgs = 9;
errInProgress = 10;
errUninitialized = 11;
errOSError = 12;
errNotSupported = 13;
errDataError = 14; # data is corrupted
errNotImplemented = 15; # Operation is not implemented
errNoMoreReplicas = 16; # No more replicas to try
errPipelineError = 17; # Backward-compatible spelling
errPipelineFailed = 17; # Pipeline failed and operation couldn't be executed
#----------------------------------------------------------------------------
# Socket related errors
#----------------------------------------------------------------------------
errInvalidAddr = 101;
errSocketError = 102;
errSocketTimeout = 103;
errSocketDisconnected = 104;
errPollerError = 105;
errSocketOptError = 106;
errStreamDisconnect = 107;
errConnectionError = 108;
errInvalidSession = 109;
errTlsError = 110;
#----------------------------------------------------------------------------
# Post Master related errors
#----------------------------------------------------------------------------
errInvalidMessage = 201;
errHandShakeFailed = 202;
errLoginFailed = 203;
errAuthFailed = 204;
errQueryNotSupported = 205;
errOperationExpired = 206;
errOperationInterrupted = 207;
errThresholdExceeded = 208;
#----------------------------------------------------------------------------
# XRootD related errors
#----------------------------------------------------------------------------
errNoMoreFreeSIDs = 301;
errInvalidRedirectURL = 302;
errInvalidResponse = 303;
errNotFound = 304;
errCheckSumError = 305;
errRedirectLimit = 306;
errCorruptedHeader = 307;
errErrorResponse = 400;
errRedirect = 401;
errLocalError = 402;
errResponseNegative = 500; # Query response was negative
_ERROR_NAMES = dict(
(value, name) for name, value in locals().copy().items()
if name.startswith('err')
)
# XRootD protocol error numbers carried by errErrorResponse statuses.
_SERVER_NOT_AUTHORIZED = 3010
_SERVER_NOT_FOUND = 3011
_SERVER_CHECKSUM_ERROR = 3019
_SERVER_AUTH_FAILED = 3030
_SERVER_REQUEST_TIMED_OUT = 3034
_SERVER_TIMER_EXPIRED = 3035
def __init__(self, status):
super(XRootDStatus, self).__init__(status)
def __str__(self):
return self.message
@property
def error_name(self):
"""Symbolic name for the status code, when known."""
return self._ERROR_NAMES.get(getattr(self, 'code', None))
[docs]
def exception(self):
"""Return a Python exception representing this status, or ``None`` if OK."""
if self.ok:
return None
code = getattr(self, 'code', None)
errno = getattr(self, 'errno', None)
server_error = code == self.errErrorResponse
if code == self.errNotFound or (
server_error and errno == self._SERVER_NOT_FOUND):
return XRootDNotFoundError(self)
if code in (self.errAuthFailed, self.errLoginFailed) or (
server_error and errno in (
self._SERVER_NOT_AUTHORIZED, self._SERVER_AUTH_FAILED)):
return XRootDAuthorizationError(self)
if code in (self.errSocketTimeout, self.errOperationExpired) or (
server_error and errno in (
self._SERVER_REQUEST_TIMED_OUT, self._SERVER_TIMER_EXPIRED)):
return XRootDTimeoutError(self)
if code == self.errCheckSumError or (
server_error and errno == self._SERVER_CHECKSUM_ERROR):
return XRootDChecksumError(self)
return XRootDOperationError(self)
[docs]
def raise_on_error(self):
"""Raise a mapped Python exception if this status is not OK."""
error = self.exception()
if error:
raise error
return self
[docs]
def raise_on_error(status):
"""Raise a mapped Python exception if ``status`` is not OK.
:param status: :class:`XRootDStatus` or raw status dictionary
:returns: the normalized :class:`XRootDStatus`
"""
if not isinstance(status, XRootDStatus):
status = XRootDStatus(status)
return status.raise_on_error()
class TapeEndpoint(Struct):
"""WLCG Tape REST API endpoint selected through discovery.
:var uri: Base URI of the selected Tape REST API endpoint
:var version: Endpoint API version
:var sitename: Storage site name advertised by discovery
"""
def __init__(self, endpoint):
super(TapeEndpoint, self).__init__(endpoint)
class TapeArchiveInfo(Struct):
"""Archive locality information for a file.
:var url: Original URL requested by the caller
:var path: Storage path sent to the Tape REST API
:var locality: Locality reported by the Tape REST API
:var error: Per-path error reported by the service, if any
"""
def __init__(self, info):
payload = {'locality': None, 'error': None}
payload.update(info)
super(TapeArchiveInfo, self).__init__(payload)
class TapeStageResponse(Struct):
"""Response returned after submitting a Tape REST stage request.
:var requestId: Server-side stage request identifier
"""
def __init__(self, response):
super(TapeStageResponse, self).__init__(response)
@property
def request_id(self):
return self.requestId
class TapeStageFileStatus(Struct):
"""Status of one file in a Tape REST stage request.
:var path: Storage path submitted in the stage request
:var onDisk: Whether the file is available on disk, if reported
:var state: Stage processing state, if reported
:var error: File-specific error, if reported
:var startedAt: Stage start timestamp, if reported
:var finishedAt: Terminal-state timestamp, if reported
"""
def __init__(self, status):
super(TapeStageFileStatus, self).__init__(status)
@property
def on_disk(self):
if hasattr(self, 'onDisk'):
return self.onDisk
return getattr(self, 'state', '') == 'COMPLETED'
class TapeStageStatus(Struct):
"""Status of a Tape REST stage request.
:var id: Server-side stage request identifier
:var createdAt: Request creation timestamp, if reported
:var startedAt: Request start timestamp, if reported
:var completedAt: Request completion timestamp, if reported
:var files: Per-file stage statuses
"""
def __init__(self, status):
status = dict(status)
status['files'] = [TapeStageFileStatus(f)
for f in status.get('files', [])]
super(TapeStageStatus, self).__init__(status)
@staticmethod
def _normalize_path(path):
parsed = urlparse(path)
if parsed.scheme and parsed.netloc:
path = parsed.path or '/'
if path.startswith('//'):
path = path[1:]
return path
def file_status(self, path):
path = self._normalize_path(path)
for status in self.files:
if status.path == path:
return status
return None
def is_on_disk(self, path):
status = self.file_status(path)
return bool(status and status.on_disk)
[docs]
class ProtocolInfo(Struct):
"""Protocol information for a server.
:var version: The version of the protocol this server is speaking
:var hostinfo: Informational flags for this host. An `ORed` combination of
:mod:`XRootD.client.flags.HostTypes`
"""
def __init__(self, info):
super(ProtocolInfo, self).__init__(info)
[docs]
class StatInfo(Struct):
"""Status information for files and directories.
:var id: This file's unique identifier
:var size: The file size (in bytes)
:var flags: Informational flags. An `ORed` combination of
:mod:`XRootD.client.flags.StatInfoFlags`
:var mtime: Modification time (in seconds since epoch)
:var modtime: Deprecated alias for ``mtime``
:var modtimestr: Deprecated modification time (as readable string)
:var ctime: Change time (in seconds since epoch)
:var atime: Access time (in seconds since epoch)
:var mode: File mode
:var modeoctstr: File mode as a readable permissions string
:var owner: File owner
:var group: File group
:var extended: Whether extended stat information is available
:var haschecksum: Whether checksum information is available
:var checksum: File checksum
"""
def __init__(self, info):
super(StatInfo, self).__init__(info)
[docs]
class StatInfoVFS(Struct):
"""Status information for Virtual File Systems.
:var nodes_rw: Number of nodes that can provide read/write space
:var free_rw: Size of the largest contiguous area of free r/w
space (in MB)
:var utilization_rw: Percentage of the partition utilization represented
by ``free_rw``
:var nodes_staging: Number of nodes that can provide staging space
:var free_staging: Size of the largest contiguous area of free staging
space (in MB)
:var utilization_staging: Percentage of the partition utilization represented
by ``free_staging``
"""
def __init__(self, info):
super(StatInfoVFS, self).__init__(info)
[docs]
class DirectoryList(Struct):
"""Directory listing.
This object is iterable::
>>> status, dirlist = filesystem.dirlist('/tmp', DirListFlags.STAT)
>>> print dirlist
<XRootD.client.responses.DirectoryList object at 0x288b9f0>
>>> print 'Entries:', dirlist.size
Entries: 2
>>> for item in dirlist:
... print item.name, item.statinfo.size
...
spam 1024
eggs 2048
:var size: The size of this listing (number of entries)
:var parent: The name of the parent directory of this directory
:var dirlist: (List of :mod:`XRootD.client.responses.ListEntry` objects) -
The list of directory entries
"""
def __init__(self, dirlist):
dirlist.update({'dirlist': [ListEntry(e) for e in dirlist['dirlist']]})
super(DirectoryList, self).__init__(dirlist)
def __iter__(self):
return iter(self.dirlist)
[docs]
class ListEntry(Struct):
"""An entry in a directory listing.
:var name: The name of the file/directory
:var hostaddr: The address of the host on which this file/directory lives
:var statinfo: (Instance of :mod:`XRootD.client.responses.StatInfo`) -
Status information about this file/directory. You must pass
`DirListFlags.STAT` with the call to
:mod:`XRootD.client.FileSystem.dirlist()` to retrieve status
information.
"""
def __init__(self, entry):
if entry['statinfo']: entry.update({'statinfo': StatInfo(entry['statinfo'])})
super(ListEntry, self).__init__(entry)
[docs]
class ChunkInfo(Struct):
"""Describes a data chunk for a vector read.
:var offset: The offset in the file from which this chunk came
:var length: The length of this chunk
:var buffer: The actual chunk data
"""
def __init__(self, info):
super(ChunkInfo, self).__init__(info)
[docs]
class VectorReadInfo(Struct):
"""Vector read response object.
Returned by :mod:`XRootD.client.File.vector_read()`.
This object is iterable::
>>> f.open('root://localhost/tmp/spam')
>>> status, chunks = file.vector_read([(0, 10), (10, 10)])
>>> print chunks
<XRootD.client.responses.VectorReadInfo object at 0x288b9f0>
>>> print chunks.size
20
>>> for chunk in chunks:
... print chunk.offset, chunk.length
...
0 10
10 10
:var size: Total size of all chunks
:var chunks: (List of :mod:`XRootD.client.responses.ChunkInfo` objects) -
The list of chunks that were read
"""
def __init__(self, info):
info.update({'chunks': [ChunkInfo(c) for c in info['chunks']]})
super(VectorReadInfo, self).__init__(info)
def __iter__(self):
return iter(self.chunks)
[docs]
class HostList(Struct):
"""A list of hosts that were involved in the request.
This object is iterable::
>>> print hostlist
<XRootD.client.responses.HostList object at 0x288b9f0>
>>> for host in hostlist:
... print host.url
...
root://localhost
:var hosts: (List of :mod:`XRootD.client.responses.HostInfo` objects) -
The list of hosts
"""
def __init__(self, hostlist):
super(HostList, self).__init__({'hosts': [HostInfo(h) for h in hostlist]})
def __iter__(self):
return iter(self.hosts)
[docs]
class HostInfo(Struct):
"""Information about a single host.
:var url: URL of the host, instance of :mod:`XRootD.client.URL`
:var protocol: Version of the protocol the host is speaking
:var flags: Host type, an `ORed` combination of
:mod:`XRootD.client.flags.HostTypes`
:var load_balancer: Was the host used as a load balancer
"""
def __init__(self, info):
super(HostInfo, self).__init__(info)