|
|
@@ -76,6 +76,7 @@ class RazClient(object):
|
|
|
self.raz_token = raz_token
|
|
|
self.username = username
|
|
|
self.service = service
|
|
|
+
|
|
|
if self.service == 'adls':
|
|
|
self.service_params = {
|
|
|
'endpoint_prefix': 'adls',
|
|
|
@@ -88,6 +89,7 @@ class RazClient(object):
|
|
|
'service_name': 's3',
|
|
|
'serviceType': 's3'
|
|
|
}
|
|
|
+
|
|
|
self.service_name = service_name
|
|
|
self.cluster_name = cluster_name
|
|
|
self.requestid = str(uuid.uuid4())
|
|
|
@@ -100,51 +102,34 @@ class RazClient(object):
|
|
|
params = params if params is not None else {}
|
|
|
headers = headers if headers is not None else {}
|
|
|
|
|
|
- allparams = [raz_signer.StringListStringMapProto(key=key, value=[val]) for key, val in url_params.items()]
|
|
|
- allparams.extend([raz_signer.StringListStringMapProto(key=key, value=[val]) for key, val in params.items()])
|
|
|
- headers = [raz_signer.StringStringMapProto(key=key, value=val) for key, val in headers.items()]
|
|
|
endpoint = "%s://%s" % (path.scheme, path.netloc)
|
|
|
resource_path = path.path.lstrip("/")
|
|
|
|
|
|
- LOG.debug(
|
|
|
- "Preparing sign request with http_method: {%s}, headers: {%s}, parameters: {%s}, endpoint: {%s}, resource_path: {%s}" %
|
|
|
- (method, headers, allparams, endpoint, resource_path)
|
|
|
- )
|
|
|
- raz_req = raz_signer.SignRequestProto(
|
|
|
- endpoint_prefix=self.service_params['endpoint_prefix'],
|
|
|
- service_name=self.service_params['service_name'],
|
|
|
- endpoint=endpoint,
|
|
|
- http_method=method,
|
|
|
- headers=headers,
|
|
|
- parameters=allparams,
|
|
|
- resource_path=resource_path,
|
|
|
- time_offset=0
|
|
|
- )
|
|
|
- raz_req_serialized = raz_req.SerializeToString()
|
|
|
- signed_request = base64.b64encode(raz_req_serialized)
|
|
|
-
|
|
|
request_data = {
|
|
|
"requestId": self.requestid,
|
|
|
"serviceType": self.service_params['serviceType'],
|
|
|
"serviceName": self.service_name,
|
|
|
"user": self.username,
|
|
|
"userGroups": [],
|
|
|
- "accessTime": "",
|
|
|
"clientIpAddress": "",
|
|
|
"clientType": "",
|
|
|
"clusterName": self.cluster_name,
|
|
|
"clusterType": "",
|
|
|
"sessionId": "",
|
|
|
- "context": {
|
|
|
- "S3_SIGN_REQUEST": signed_request
|
|
|
- }
|
|
|
+ "accessTime": "",
|
|
|
+ "context": {}
|
|
|
}
|
|
|
- headers = {"Content-Type":"application/json", "Accept-Encoding":"gzip,deflate"}
|
|
|
- raz_url = "%s/api/authz/s3/access?delegation=%s" % (self.raz_url, self.raz_token)
|
|
|
- LOG.debug('Raz url: %s' % raz_url)
|
|
|
+ request_headers = {"Content-Type": "application/json"}
|
|
|
+ raz_url = "%s/api/authz/%s/access?delegation=%s" % (self.raz_url, self.service, self.raz_token)
|
|
|
|
|
|
- LOG.debug("Sending access check headers: {%s} request_data: {%s}" % (headers, request_data))
|
|
|
- raz_req = requests.post(raz_url, headers=headers, json=request_data, verify=False)
|
|
|
+ if self.service == 'adls':
|
|
|
+ self._make_adls_request(request_data, path, resource_path)
|
|
|
+ elif self.service == 's3':
|
|
|
+ self._make_s3_request(request_data, request_headers, method, params, headers, url_params, endpoint, resource_path)
|
|
|
+
|
|
|
+ LOG.debug('Raz url: %s' % raz_url)
|
|
|
+ LOG.debug("Sending access check headers: {%s} request_data: {%s}" % (request_headers, request_data))
|
|
|
+ raz_req = requests.post(raz_url, headers=request_headers, json=request_data, verify=False)
|
|
|
|
|
|
signed_response_result = None
|
|
|
signed_response = None
|
|
|
@@ -164,21 +149,67 @@ class RazClient(object):
|
|
|
if result == "ALLOWED":
|
|
|
LOG.debug('Received allowed response %s' % raz_req.json())
|
|
|
signed_response_data = raz_req.json()["operResult"]["additionalInfo"]
|
|
|
+
|
|
|
if self.service == 'adls':
|
|
|
LOG.debug("Received SAS %s" % signed_response_data["ADLS_DSAS"])
|
|
|
return {'token': signed_response_data["ADLS_DSAS"]}
|
|
|
else:
|
|
|
signed_response_result = signed_response_data["S3_SIGN_RESPONSE"]
|
|
|
|
|
|
- if signed_response_result:
|
|
|
+ if signed_response_result is not None:
|
|
|
raz_response_proto = raz_signer.SignResponseProto()
|
|
|
signed_response = raz_response_proto.FromString(base64.b64decode(signed_response_result))
|
|
|
LOG.debug("Received signed Response %s" % signed_response)
|
|
|
|
|
|
# Signed headers "only"
|
|
|
- if signed_response:
|
|
|
+ if signed_response is not None:
|
|
|
return dict([(i.key, i.value) for i in signed_response.signer_generated_headers])
|
|
|
|
|
|
+ def _make_adls_request(self, request_data, path, resource_path):
|
|
|
+ storage_account = path.netloc.split('.')[0]
|
|
|
+ container, relative_path = resource_path.split('/', 1)
|
|
|
+
|
|
|
+ request_data.update({
|
|
|
+ "clientType": "adls",
|
|
|
+ "operation": {
|
|
|
+ "resource": {
|
|
|
+ "storageaccount": storage_account,
|
|
|
+ "container": container,
|
|
|
+ "relativepath": relative_path,
|
|
|
+ },
|
|
|
+ "resourceOwner": "",
|
|
|
+ "action": "read",
|
|
|
+ "accessTypes":["read"]
|
|
|
+ }
|
|
|
+ })
|
|
|
+
|
|
|
+ def _make_s3_request(self, request_data, request_headers, method, params, headers, url_params, endpoint, resource_path):
|
|
|
+
|
|
|
+ allparams = [raz_signer.StringListStringMapProto(key=key, value=[val]) for key, val in url_params.items()]
|
|
|
+ allparams.extend([raz_signer.StringListStringMapProto(key=key, value=[val]) for key, val in params.items()])
|
|
|
+ headers = [raz_signer.StringStringMapProto(key=key, value=val) for key, val in headers.items()]
|
|
|
+
|
|
|
+ LOG.debug(
|
|
|
+ "Preparing sign request with http_method: {%s}, headers: {%s}, parameters: {%s}, endpoint: {%s}, resource_path: {%s}" %
|
|
|
+ (method, headers, allparams, endpoint, resource_path)
|
|
|
+ )
|
|
|
+ raz_req = raz_signer.SignRequestProto(
|
|
|
+ endpoint_prefix=self.service_params['endpoint_prefix'],
|
|
|
+ service_name=self.service_params['service_name'],
|
|
|
+ endpoint=endpoint,
|
|
|
+ http_method=method,
|
|
|
+ headers=headers,
|
|
|
+ parameters=allparams,
|
|
|
+ resource_path=resource_path,
|
|
|
+ time_offset=0
|
|
|
+ )
|
|
|
+ raz_req_serialized = raz_req.SerializeToString()
|
|
|
+ signed_request = base64.b64encode(raz_req_serialized)
|
|
|
+
|
|
|
+ request_headers["Accept-Encoding"] = {"gzip,deflate"}
|
|
|
+ request_data["context"] = {
|
|
|
+ "S3_SIGN_REQUEST": signed_request
|
|
|
+ }
|
|
|
|
|
|
def get_raz_client(raz_url, username, auth='kerberos', service='s3', service_name='cm_s3', cluster_name='myCluster'):
|
|
|
if not username:
|