diff --git a/docs/_build/html/_modules/gen3/auth.html b/docs/_build/html/_modules/gen3/auth.html
index bac0cdd77..7f88fdc8c 100644
--- a/docs/_build/html/_modules/gen3/auth.html
+++ b/docs/_build/html/_modules/gen3/auth.html
@@ -58,35 +58,35 @@
def decode_token(token_str):
- """
- jq -r '.api_key' < ~/.gen3/qa-covid19.planx-pla.net.json | awk -F . '{ print $2 }' | base64 --decode | jq -r .
- """
- tokenParts = token_str.split(".")
+ """
+ jq -r '.api_key' < ~/.gen3/qa-covid19.planx-pla.net.json | awk -F . '{ print $2 }' | base64 --decode | jq -r .
+ """
+ tokenParts = token_str.split(".")
if len(tokenParts) < 3:
- raise Exception("Invalid JWT. Could not split into parts.")
- padding = "===="
+ raise Exception("Invalid JWT. Could not split into parts.")
+ padding = "===="
infoStr = tokenParts[1] + padding[0 : len(tokenParts[1]) % 4]
jsonStr = base64.urlsafe_b64decode(infoStr)
return json.loads(jsonStr)
def endpoint_from_token(token_str):
- """
- Extract the endpoint from a JWT issue ("iss" property)
- """
+ """
+ Extract the endpoint from a JWT issue ("iss" property)
+ """
info = decode_token(token_str)
- urlparts = urlparse(info["iss"])
- endpoint = urlparts.scheme + "://" + urlparts.hostname
+ urlparts = urlparse(info["iss"])
+ endpoint = urlparts.scheme + "://" + urlparts.hostname
if urlparts.port:
- endpoint += ":" + str(urlparts.port)
+ endpoint += ":" + str(urlparts.port)
return remove_trailing_whitespace_and_slashes_in_url(endpoint)
def _handle_access_token_response(resp, token_key):
- """
+ """
Shared helper for both get_access_token_with_key and get_access_token_from_wts
- """
- err_msg = "Failed to get an access token from {}:\n{}"
+ """
+ err_msg = "Failed to get an access token from {}:\n{}"
if resp.status_code != 200:
raise Gen3AuthError(err_msg.format(resp.url, resp.text))
try:
@@ -99,64 +99,64 @@ Source code for gen3.auth
def get_access_token_with_key(api_key):
- """
+ """
Try to fetch an access token given the api key
- """
- endpoint = endpoint_from_token(api_key["api_key"])
+ """
+ endpoint = endpoint_from_token(api_key["api_key"])
# attempt to get a token from Fence
- auth_url = "{}/user/credentials/cdis/access_token".format(endpoint)
+ auth_url = "{}/user/credentials/cdis/access_token".format(endpoint)
resp = requests.post(auth_url, json=api_key)
- token_key = "access_token"
+ token_key = "access_token"
return _handle_access_token_response(resp, token_key)
def get_access_token_with_client_credentials(endpoint, client_credentials, scopes):
- """
+ """
Try to get an access token from Fence using client credentials
Args:
endpoint (str): URL of the Gen3 instance to get an access token for
client_credentials ((str, str) tuple): (client ID, client secret) tuple
scopes (str): space-delimited list of scopes to request
- """
+ """
if not endpoint:
- raise ValueError("'endpoint' must be specified when using client credentials")
- url = f"{endpoint}/user/oauth2/token?grant_type=client_credentials&scope={scopes}"
+ raise ValueError("'endpoint' must be specified when using client credentials")
+ url = f"{endpoint}/user/oauth2/token?grant_type=client_credentials&scope={scopes}"
resp = requests.post(url, auth=client_credentials)
- return _handle_access_token_response(resp, "access_token")
+ return _handle_access_token_response(resp, "access_token")
-def get_wts_endpoint(namespace=os.getenv("NAMESPACE", "default")):
- return "http://workspace-token-service.{}.svc.cluster.local".format(namespace)
+def get_wts_endpoint(namespace=os.getenv("NAMESPACE", "default")):
+ return "http://workspace-token-service.{}.svc.cluster.local".format(namespace)
-def get_wts_idps(namespace=os.getenv("NAMESPACE", "default"), external_wts_host=None):
+def get_wts_idps(namespace=os.getenv("NAMESPACE", "default"), external_wts_host=None):
wts_url = None
if external_wts_host == None:
wts_url = get_wts_endpoint(namespace)
else:
wts_url = external_wts_host
- url = wts_url.rstrip("/") + "/external_oidc/"
+ url = wts_url.rstrip("/") + "/external_oidc/"
resp = requests.get(url)
raise_for_status_and_print_error(resp)
return resp.json()
def get_token_cache_file_name(key):
- """Compute the path to the access-token cache file"""
- cache_folder = "{}/.cache/gen3/".format(os.path.expanduser("~"))
+ """Compute the path to the access-token cache file"""
+ cache_folder = "{}/.cache/gen3/".format(os.path.expanduser("~"))
os.makedirs(cache_folder, exist_ok=True)
- cache_prefix = cache_folder + "token_cache_"
+ cache_prefix = cache_folder + "token_cache_"
s = hashlib.sha256()
- s.update(key.encode("utf-8"))
+ s.update(key.encode("utf-8"))
return cache_prefix + s.hexdigest()
[docs]
class Gen3Auth(AuthBase):
-
"""Gen3 auth helper class for use with requests auth.
+
"""Gen3 auth helper class for use with requests auth.
Implements requests.auth.AuthBase in order to support JWT authentication.
Generates access tokens from the provided refresh token file or string.
@@ -164,17 +164,17 @@
Source code for gen3.auth
Args:
refresh_file (str, opt): The file containing the downloaded JSON web token. Optional if working in a Gen3 Workspace.
- Defaults to (env["GEN3_API_KEY"] || "credentials") if refresh_token and idp not set.
+ Defaults to (env["GEN3_API_KEY"] || "credentials") if refresh_token and idp not set.
Includes ~/.gen3/ in search path if value does not include /.
- Interprets "idp://wts/<idp>" as an idp.
- Interprets "accesstoken:///<token>" as an access token
+ Interprets "idp://wts/<idp>" as an idp.
+ Interprets "accesstoken:///<token>" as an access token
refresh_token (str, opt): The JSON web token. Optional if working in a Gen3 Workspace.
idp (str, opt): If working in a Gen3 Workspace, the IDP to use can be specified -
- "local" indicates the local environment fence idp
+ "local" indicates the local environment fence idp
client_credentials (tuple, opt): The (client_id, client_secret) credentials for an OIDC client
- that has the 'client_credentials' grant, allowing it to obtain access tokens.
+ that has the 'client_credentials' grant, allowing it to obtain access tokens.
client_scopes (str, opt): Space-separated list of scopes requested for access tokens obtained from client
- credentials. Default: "user data openid"
+ credentials. Default: "user data openid"
access_token (str, opt): provide an access token to override the use of any
API key/refresh token. This is intended for cases where you may want to
pass a token that was issued to a particular OIDC client (rather than acting on
@@ -189,30 +189,30 @@ Source code for gen3.auth
or use ~/.gen3/crdc.json:
- >>> auth = Gen3Auth(refresh_file="crdc")
+ >>> auth = Gen3Auth(refresh_file="crdc")
or use some arbitrary file:
- >>> auth = Gen3Auth(refresh_file="./key.json")
+ >>> auth = Gen3Auth(refresh_file="./key.json")
or set the GEN3_API_KEY environment variable rather
than pass the refresh_file argument to the Gen3Auth
constructor.
- If working with an OIDC client that has the 'client_credentials' grant, allowing it to obtain
+ If working with an OIDC client that has the 'client_credentials' grant, allowing it to obtain
access tokens, provide the client ID and secret:
Note: client secrets should never be hardcoded!
>>> auth = Gen3Auth(
- endpoint="https://datacommons.example",
- client_credentials=("client ID", os.environ["GEN3_OIDC_CLIENT_CREDS_SECRET"])
+ endpoint="https://datacommons.example",
+ client_credentials=("client ID", os.environ["GEN3_OIDC_CLIENT_CREDS_SECRET"])
)
If working in a Gen3 Workspace, initialize as follows:
>>> auth = Gen3Auth()
- """
+ """
def __init__(
self,
@@ -224,40 +224,40 @@ Source code for gen3.auth
client_scopes=None,
access_token=None,
):
- logging.debug("Initializing auth..")
+ logging.debug("Initializing auth..")
self.endpoint = remove_trailing_whitespace_and_slashes_in_url(endpoint)
- # note - `_refresh_token` is not actually a JWT refresh token - it's a
- # gen3 api key with a token as the "api_key" property
+ # note - `_refresh_token` is not actually a JWT refresh token - it's a
+ # gen3 api key with a token as the "api_key" property
self._refresh_token = refresh_token
self._access_token = access_token
self._access_token_info = None
- self._wts_idp = idp or "local"
- self._wts_namespace = os.environ.get("NAMESPACE", "default")
+ self._wts_idp = idp or "local"
+ self._wts_namespace = os.environ.get("NAMESPACE", "default")
self._use_wts = False
self._external_wts_host = None
self._refresh_file = refresh_file
self._client_credentials = client_credentials
if self._client_credentials:
- self._client_scopes = client_scopes or "user data openid"
+ self._client_scopes = client_scopes or "user data openid"
elif client_scopes:
raise ValueError(
- "'client_scopes' cannot be specified without 'client_credentials'"
+ "'client_scopes' cannot be specified without 'client_credentials'"
)
if refresh_file and refresh_token:
raise ValueError(
- "Only one of 'refresh_file' and 'refresh_token' can be specified."
+ "Only one of 'refresh_file' and 'refresh_token' can be specified."
)
if endpoint and idp:
- raise ValueError("Only one of 'endpoint' and 'idp' can be specified.")
+ raise ValueError("Only one of 'endpoint' and 'idp' can be specified.")
if not refresh_file and not refresh_token and not idp:
- refresh_file = os.getenv("GEN3_API_KEY", "credentials")
+ refresh_file = os.getenv("GEN3_API_KEY", "credentials")
if refresh_file and not idp:
- idp_prefix = "idp://wts/"
- access_token_prefix = "accesstoken:///"
+ idp_prefix = "idp://wts/"
+ access_token_prefix = "accesstoken:///"
if refresh_file[0 : len(idp_prefix)] == idp_prefix:
idp = refresh_file[len(idp_prefix) :]
refresh_file = None
@@ -267,22 +267,22 @@ Source code for gen3.auth
refresh_file = None
elif (
not os.path.isfile(refresh_file)
- and "/" not in refresh_file
- and "\\" not in refresh_file
+ and "/" not in refresh_file
+ and "\\" not in refresh_file
):
- refresh_file = "{}/.gen3/{}".format(
- os.path.expanduser("~"), refresh_file
+ refresh_file = "{}/.gen3/{}".format(
+ os.path.expanduser("~"), refresh_file
)
- if not os.path.isfile(refresh_file) and refresh_file[-5:] != ".json":
- refresh_file += ".json"
+ if not os.path.isfile(refresh_file) and refresh_file[-5:] != ".json":
+ refresh_file += ".json"
if not os.path.isfile(refresh_file):
- logging.warning("Unable to find refresh_file")
+ logging.warning("Unable to find refresh_file")
refresh_file = None
if self._client_credentials:
if not endpoint:
raise ValueError(
- "'endpoint' must be specified when '_client_credentials' is specified"
+ "'endpoint' must be specified when '_client_credentials' is specified"
)
self._access_token = get_access_token_with_client_credentials(
endpoint, self._client_credentials, self._client_scopes
@@ -292,7 +292,7 @@ Source code for gen3.auth
# at this point - refresh_file either exists or is None
if not refresh_file and not refresh_token:
# check if this is a Gen3 workspace environment
- # most production environments are in the "default" namespace
+ # most production environments are in the "default" namespace
# attempt to get a token from the workspace-token-service
self._use_wts = True
# hate calling a method from the constructor, but avoids copying code
@@ -304,34 +304,34 @@ Source code for gen3.auth
self._refresh_token = json.loads(file_data)
except Exception as e:
raise ValueError(
- "Couldn't load your refresh token file: {}\n{}".format(
+ "Couldn't load your refresh token file: {}\n{}".format(
refresh_file, str(e)
)
)
- assert "api_key" in self._refresh_token
+ assert "api_key" in self._refresh_token
# if both endpoint and refresh file are provided, compare endpoint with iss in refresh file
# May need to use network wts endpoint
if idp or (
endpoint
and (
- not endpoint.rstrip("/")
- == endpoint_from_token(self._refresh_token["api_key"])
+ not endpoint.rstrip("/")
+ == endpoint_from_token(self._refresh_token["api_key"])
)
):
try:
logging.debug(
- "Switch to using WTS and set external WTS host url.."
+ "Switch to using WTS and set external WTS host url.."
)
self._use_wts = True
self._external_wts_host = (
- endpoint_from_token(self._refresh_token["api_key"])
- + "/wts/"
+ endpoint_from_token(self._refresh_token["api_key"])
+ + "/wts/"
)
self.get_access_token()
except Gen3AuthError as g:
logging.warning(
- "Could not obtain access token from WTS service."
+ "Could not obtain access token from WTS service."
)
raise g
@@ -339,20 +339,20 @@ Source code for gen3.auth
if self._access_token:
self.endpoint = endpoint_from_token(self._access_token)
else:
- self.endpoint = endpoint_from_token(self._refresh_token["api_key"])
+ self.endpoint = endpoint_from_token(self._refresh_token["api_key"])
@property
def _token_info(self):
- """
+ """
Wrapper to fix intermittent errors when the token is being refreshed
and `_access_token_info` == None
- """
+ """
if not self._access_token_info:
self.refresh_access_token()
return self._access_token_info
def __call__(self, request):
- """Adds authorization header to the request
+ """Adds authorization header to the request
This gets called by the python.requests package on outbound requests
so that authentication can be added.
@@ -360,13 +360,13 @@ Source code for gen3.auth
Args:
request (object): The incoming request object
- """
- request.headers["Authorization"] = self._get_auth_value()
- request.register_hook("response", self._handle_401)
+ """
+ request.headers["Authorization"] = self._get_auth_value()
+ request.register_hook("response", self._handle_401)
return request
def _handle_401(self, response, **kwargs):
- """Handles failed requests when authorization failed.
+ """Handles failed requests when authorization failed.
This gets called after a failed request when an HTTP 401 error
occurs. This then tries to refresh the access token in the event
@@ -375,7 +375,7 @@ Source code for gen3.auth
Args:
request (object): The failed request object
- """
+ """
if not response.status_code == 401 and not response.status_code == 403:
return response
@@ -386,7 +386,7 @@ Source code for gen3.auth
# copy the request to resend
newreq = response.request.copy()
- newreq.headers["Authorization"] = self._get_auth_value()
+ newreq.headers["Authorization"] = self._get_auth_value()
_response = response.connection.send(newreq, **kwargs)
_response.history.append(response)
@@ -397,7 +397,7 @@ Source code for gen3.auth
[docs]
def refresh_access_token(self, endpoint=None):
-
"""Get a new access token"""
+
"""Get a new access token"""
if self._use_wts:
self._access_token = self.get_access_token_from_wts(endpoint)
elif self._client_credentials:
@@ -408,8 +408,8 @@
Source code for gen3.auth
self._access_token = get_access_token_with_key(self._refresh_token)
else:
logging.warning(
- f"Unable to refresh access token. "
- f"Authorized API calls will stop working when this token expires."
+ f"Unable to refresh access token. "
+ f"Authorized API calls will stop working when this token expires."
)
self._access_token_info = decode_token(self._access_token)
@@ -418,14 +418,14 @@ Source code for gen3.auth
if self._use_wts:
cache_file = get_token_cache_file_name(self._wts_idp)
elif self._refresh_file:
- cache_file = get_token_cache_file_name(self._refresh_token["api_key"])
+ cache_file = get_token_cache_file_name(self._refresh_token["api_key"])
if cache_file:
try:
self._write_to_file(cache_file, self._access_token)
except Exception as e:
logging.warning(
- f"Unable to write access token to cache file. Exceeded number of retries. Details: {e}"
+ f"Unable to write access token to cache file. Exceeded number of retries. Details: {e}"
)
return self._access_token
@@ -438,37 +438,37 @@ Source code for gen3.auth
# write a temp file, then rename - to avoid
# simultaneous writes to same file race condition
temp = cache_file + (
- ".tmp_eraseme_%d_%d" % (random.randrange(100000), time.time())
+ ".tmp_eraseme_%d_%d" % (random.randrange(100000), time.time())
)
try:
- with open(temp, "w") as f:
+ with open(temp, "w") as f:
f.write(content)
os.rename(temp, cache_file)
return True
except Exception as e:
- logging.warning("failed to write token cache file: " + cache_file)
+ logging.warning("failed to write token cache file: " + cache_file)
logging.warning(str(e))
raise e
[docs]
def get_access_token(self):
-
"""Get the access token - auto refresh if within 5 minutes of expiration"""
+
"""Get the access token - auto refresh if within 5 minutes of expiration"""
if not self._access_token:
if self._use_wts == True:
cache_file = get_token_cache_file_name(self._wts_idp)
else:
if self._refresh_token:
cache_file = get_token_cache_file_name(
-
self._refresh_token["api_key"]
+
self._refresh_token["api_key"]
)
if cache_file and os.path.isfile(cache_file):
-
try: # don't freak out on invalid cache
+
try: # don't freak out on invalid cache
with open(cache_file) as f:
self._access_token = f.read()
self._access_token_info = decode_token(self._access_token)
except Exception as e:
-
logging.warning("ignoring invalid token cache: " + cache_file)
+
logging.warning("ignoring invalid token cache: " + cache_file)
self._access_token = None
self._access_token_info = None
logging.warning(str(e))
@@ -476,108 +476,108 @@
Source code for gen3.auth
need_new_token = (
not self._access_token
or not self._access_token_info
- or time.time() + 300 > self._access_token_info["exp"]
+ or time.time() + 300 > self._access_token_info["exp"]
)
if need_new_token:
return self.refresh_access_token(
- self.endpoint if hasattr(self, "endpoint") else None
+ self.endpoint if hasattr(self, "endpoint") else None
)
# use cache
return self._access_token
def _get_auth_value(self):
-
"""Returns the Authorization header value for the request
+
"""Returns the Authorization header value for the request
This gets called when added the Authorization header to the request.
This fetches the access token from the refresh token if the access token is missing.
-
"""
-
return "bearer " + self.get_access_token()
+
"""
+
return "bearer " + self.get_access_token()
[docs]
def curl(self, path, request=None, data=None):
-
"""
+
"""
Curl the given endpoint - ex: gen3 curl /user/user. Return requests.Response
Args:
path (str): path under the commons to curl (/user/user, /index/index, /authz/mapping, ...)
request (str in GET|POST|PUT|DELETE): default to GET if data is not set, else default to POST
-
data (str): json string or "@filename" of a json file
-
"""
+
data (str): json string or "@filename" of a json file
+
"""
if not request:
-
request = "GET"
+
request = "GET"
if data:
-
request = "POST"
+
request = "POST"
json_data = data
output = None
-
if data and data[0] == "@":
+
if data and data[0] == "@":
with open(data[1:]) as f:
json_data = f.read()
-
if request == "GET":
-
output = requests.get(self.endpoint + "/" + path, auth=self)
-
elif request == "POST":
+
if request == "GET":
+
output = requests.get(self.endpoint + "/" + path, auth=self)
+
elif request == "POST":
output = requests.post(
-
self.endpoint + "/" + path, json=json_data, auth=self
+
self.endpoint + "/" + path, json=json_data, auth=self
)
-
elif request == "PUT":
-
output = requests.put(self.endpoint + "/" + path, json=json_data, auth=self)
-
elif request == "DELETE":
-
output = requests.delete(self.endpoint + "/" + path, auth=self)
+
elif request == "PUT":
+
output = requests.put(self.endpoint + "/" + path, json=json_data, auth=self)
+
elif request == "DELETE":
+
output = requests.delete(self.endpoint + "/" + path, auth=self)
else:
-
raise Exception("Invalid request type: " + request)
+
raise Exception("Invalid request type: " + request)
return output
[docs]
def get_access_token_from_wts(self, endpoint=None):
-
"""
+
"""
Try to fetch an access token for the given idp from the wts
-
in the given namespace. If idp is not set, then default to "local"
-
"""
+
in the given namespace. If idp is not set, then default to "local"
+
"""
# attempt to get a token from the workspace-token-service
-
logging.debug("getting access token from wts..")
-
auth_url = get_wts_endpoint(self._wts_namespace) + "/token/"
+
logging.debug("getting access token from wts..")
+
auth_url = get_wts_endpoint(self._wts_namespace) + "/token/"
-
# If non "local" idp value exists, append to auth url
+
# If non "local" idp value exists, append to auth url
# If user specified endpoint value, then first attempt to determine idp value.
-
if self.endpoint or (self._wts_idp and self._wts_idp != "local"):
+
if self.endpoint or (self._wts_idp and self._wts_idp != "local"):
# If user supplied endpoint value and not idp, figure out the idp value
if self.endpoint:
logging.debug(
-
"First try to use the local WTS to figure out idp name for the supplied endpoint.."
+
"First try to use the local WTS to figure out idp name for the supplied endpoint.."
)
try:
provider_List = get_wts_idps(self._wts_namespace)
matchProviders = list(
filter(
-
lambda provider: provider["base_url"] == endpoint,
-
provider_List["providers"],
+
lambda provider: provider["base_url"] == endpoint,
+
provider_List["providers"],
)
)
if len(matchProviders) == 1:
-
logging.debug("Found matching idp from local WTS.")
-
self._wts_idp = matchProviders[0]["idp"]
+
logging.debug("Found matching idp from local WTS.")
+
self._wts_idp = matchProviders[0]["idp"]
elif len(matchProviders) > 1:
raise ValueError(
-
"Multiple idps matched with endpoint value provided."
+
"Multiple idps matched with endpoint value provided."
)
else:
-
logging.debug("Could not find matching idp from local WTS.")
+
logging.debug("Could not find matching idp from local WTS.")
except Exception as e:
logging.debug(
-
"Exception occured when making network call to local WTS."
+
"Exception occured when making network call to local WTS."
)
if not self._external_wts_host:
raise e
else:
-
logging.debug("Since external WTS host exists, continuing on..")
+
logging.debug("Since external WTS host exists, continuing on..")
pass
-
if self._wts_idp and self._wts_idp != "local":
-
auth_url += "?idp={}".format(self._wts_idp)
+
if self._wts_idp and self._wts_idp != "local":
+
auth_url += "?idp={}".format(self._wts_idp)
# If endpoint value exists, only get WTS token if idp value has been successfully determined
# Otherwise skip to querying external WTS
@@ -585,97 +585,97 @@
Source code for gen3.auth
if (
not self._external_wts_host
or not self.endpoint
- or (self.endpoint and self._wts_idp != "local")
+ or (self.endpoint and self._wts_idp != "local")
):
try:
- logging.debug("Try to get access token from local WTS..")
- logging.debug(f"{auth_url=}")
+ logging.debug("Try to get access token from local WTS..")
+ logging.debug(f"{auth_url=}")
resp = requests.get(auth_url)
if (resp and resp.status_code == 200) or (not self._external_wts_host):
- return _handle_access_token_response(resp, "token")
+ return _handle_access_token_response(resp, "token")
except Exception as e:
if not self._external_wts_host:
raise e
else:
# Try to obtain token from external wts
- logging.debug("Could get obtain token from Local WTS.")
+ logging.debug("Could get obtain token from Local WTS.")
pass
# local workspace wts call failed, try using a network call
# First get access token with WTS host
- logging.debug("Trying to get access token from external WTS Host..")
+ logging.debug("Trying to get access token from external WTS Host..")
wts_token = get_access_token_with_key(self._refresh_token)
- auth_url = self._external_wts_host + "token/"
+ auth_url = self._external_wts_host + "token/"
provider_List = get_wts_idps(self._wts_namespace, self._external_wts_host)
# if user already supplied idp, use that
- if self._wts_idp and self._wts_idp != "local":
+ if self._wts_idp and self._wts_idp != "local":
matchProviders = list(
filter(
- lambda provider: provider["idp"] == self._wts_idp,
- provider_List["providers"],
+ lambda provider: provider["idp"] == self._wts_idp,
+ provider_List["providers"],
)
)
elif endpoint:
matchProviders = list(
filter(
- lambda provider: provider["base_url"] == endpoint,
- provider_List["providers"],
+ lambda provider: provider["base_url"] == endpoint,
+ provider_List["providers"],
)
)
else:
raise Exception(
- "Unable to generate matching identity providers (no IdP or endpoint provided)"
+ "Unable to generate matching identity providers (no IdP or endpoint provided)"
)
if len(matchProviders) == 1:
- self._wts_idp = matchProviders[0]["idp"]
- logging.debug("Succesfully determined idp value: {}".format(self._wts_idp))
+ self._wts_idp = matchProviders[0]["idp"]
+ logging.debug("Succesfully determined idp value: {}".format(self._wts_idp))
else:
- idp_list = "\n "
+ idp_list = "\n "
if len(matchProviders) > 1:
for idp in matchProviders:
idp_list = (
idp_list
- + "idp name: "
- + idp["idp"]
- + " url: "
- + idp["base_url"]
- + "\n "
+ + "idp name: "
+ + idp["idp"]
+ + " url: "
+ + idp["base_url"]
+ + "\n "
)
raise ValueError(
- "Multiple idps matched with endpoint value provided."
+ "Multiple idps matched with endpoint value provided."
+ idp_list
- + "Query /wts/external_oidc/ for more information."
+ + "Query /wts/external_oidc/ for more information."
)
else:
- for idp in provider_List["providers"]:
+ for idp in provider_List["providers"]:
idp_list = (
idp_list
- + "idp name: "
- + idp["idp"]
- + " Endpoint url: "
- + idp["base_url"]
- + "\n "
+ + "idp name: "
+ + idp["idp"]
+ + " Endpoint url: "
+ + idp["base_url"]
+ + "\n "
)
raise ValueError(
- "No idp matched with the endpoint or idp value provided.\n"
- + "Please make sure your endpoint or idp value matches exactly with the output below.\n"
- + "i.e. check trailing '/' character for the endpoint url\n"
- + "Available Idps:"
+ "No idp matched with the endpoint or idp value provided.\n"
+ + "Please make sure your endpoint or idp value matches exactly with the output below.\n"
+ + "i.e. check trailing '/' character for the endpoint url\n"
+ + "Available Idps:"
+ idp_list
- + "Query /wts/external_oidc/ for more information."
+ + "Query /wts/external_oidc/ for more information."
)
- logging.debug("Finally getting access token..")
- auth_url += "?idp={}".format(self._wts_idp)
- header = {"Authorization": "Bearer " + wts_token}
+ logging.debug("Finally getting access token..")
+ auth_url += "?idp={}".format(self._wts_idp)
+ header = {"Authorization": "Bearer " + wts_token}
resp = requests.get(auth_url, headers=header)
- err_msg = "Please make sure the target commons is connected on your profile page and that connection has not expired."
+ err_msg = "Please make sure the target commons is connected on your profile page and that connection has not expired."
if resp.status_code != 200:
logging.warning(err_msg)
- return _handle_access_token_response(resp, "token")
+
return _handle_access_token_response(resp, "token")
diff --git a/docs/_build/html/_modules/gen3/file.html b/docs/_build/html/_modules/gen3/file.html
index 53916adc8..1cd1cbfa6 100644
--- a/docs/_build/html/_modules/gen3/file.html
+++ b/docs/_build/html/_modules/gen3/file.html
@@ -59,7 +59,7 @@ Source code for gen3.file
[docs]
class Gen3File:
-
"""For interacting with Gen3 file management features.
+
"""For interacting with Gen3 file management features.
A class for interacting with the Gen3 file download services.
Supports getting presigned urls right now.
@@ -71,10 +71,10 @@
Source code for gen3.file
This generates the Gen3File class pointed at the sandbox commons while
using the credentials.json downloaded from the commons profile page.
- >>> auth = Gen3Auth(refresh_file="credentials.json")
+ >>> auth = Gen3Auth(refresh_file="credentials.json")
... file = Gen3File(auth)
- """
+ """
def __init__(self, endpoint=None, auth_provider=None):
# auth_provider legacy interface required endpoint as 1st arg
@@ -85,7 +85,7 @@ Source code for gen3.file
[docs]
def get_presigned_url(self, guid, protocol=None):
-
"""Generates a presigned URL for a file.
+
"""Generates a presigned URL for a file.
Retrieves a presigned url for a file giving access to a file for a limited time.
@@ -97,10 +97,10 @@
Source code for gen3.file
>>> Gen3File.get_presigned_url(query)
- """
- api_url = "{}/user/data/download/{}".format(self._endpoint, guid)
+ """
+ api_url = "{}/user/data/download/{}".format(self._endpoint, guid)
if protocol:
- api_url += "?protocol={}".format(protocol)
+ api_url += "?protocol={}".format(protocol)
resp = requests.get(api_url, auth=self._auth_provider)
raise_for_status_and_print_error(resp)
@@ -113,7 +113,7 @@ Source code for gen3.file
[docs]
def delete_file(self, guid):
-
"""
+
"""
This method is DEPRECATED. Use delete_file_locations() instead.
Delete all locations of a stored data file and remove its record from indexd
@@ -121,9 +121,9 @@
Source code for gen3.file
guid (str): provide a UUID for file id to delete
Returns:
text: requests.delete text result
- """
- print("This method is DEPRECATED. Use delete_file_locations() instead.")
- api_url = "{}/user/data/{}".format(self._endpoint, guid)
+ """
+ print("This method is DEPRECATED. Use delete_file_locations() instead.")
+ api_url = "{}/user/data/{}".format(self._endpoint, guid)
output = requests.delete(api_url, auth=self._auth_provider).text
return output
@@ -132,15 +132,15 @@
Source code for gen3.file
[docs]
def delete_file_locations(self, guid):
-
"""
+
"""
Delete all locations of a stored data file and remove its record from indexd
Args:
guid (str): provide a UUID for file id to delete
Returns:
requests.Response : requests.delete result
-
"""
-
api_url = "{}/user/data/{}".format(self._endpoint, guid)
+
"""
+
api_url = "{}/user/data/{}".format(self._endpoint, guid)
output = requests.delete(api_url, auth=self._auth_provider)
return output
@@ -151,37 +151,37 @@ Source code for gen3.file
def upload_file(
self, file_name, authz=None, protocol=None, expires_in=None, bucket=None
):
- """
+ """
Get a presigned url for a file to upload
Args:
file_name (str): file_name to use for upload
authz (list): authorization scope for the file as list of paths, optional.
- protocol (str): Storage protocol to use for upload: "s3", "az".
- If this isn't set, the default will be "s3"
+ protocol (str): Storage protocol to use for upload: "s3", "az".
+ If this isn't set, the default will be "s3"
expires_in (int): Amount in seconds that the signed url will expire from datetime.utcnow().
Be sure to use a positive integer.
This value will also be treated as <= MAX_PRESIGNED_URL_TTL in the fence configuration.
- bucket (str): Bucket to upload to. The bucket must be configured in the Fence instance's
+ bucket (str): Bucket to upload to. The bucket must be configured in the Fence instance's
`ALLOWED_DATA_UPLOAD_BUCKETS` setting. If not specified, Fence defaults to the
`DATA_UPLOAD_BUCKET` setting.
Returns:
Document: json representation for the file upload
- """
- api_url = f"{self._endpoint}/user/data/upload"
+ """
+ api_url = f"{self._endpoint}/user/data/upload"
body = {}
if protocol:
- body["protocol"] = protocol
+ body["protocol"] = protocol
if authz:
- body["authz"] = authz
+ body["authz"] = authz
if expires_in:
- body["expires_in"] = expires_in
+ body["expires_in"] = expires_in
if file_name:
- body["file_name"] = file_name
+ body["file_name"] = file_name
if bucket:
- body["bucket"] = bucket
+ body["bucket"] = bucket
- headers = {"Content-Type": "application/json"}
+ headers = {"Content-Type": "application/json"}
resp = requests.post(
api_url, auth=self._auth_provider, json=body, headers=headers
)
@@ -195,13 +195,13 @@ Source code for gen3.file
def _ensure_dirpath_exists(path: Path) -> Path:
- """Utility to create a directory if missing.
+ """Utility to create a directory if missing.
Returns the path so that the call can be inlined in another call
Args:
path (Path): path to create
Returns
path of created directory
- """
+ """
assert path
out_path: Path = path
@@ -213,58 +213,58 @@ Source code for gen3.file
[docs]
def download_single(self, object_id, path):
-
"""
+
"""
Download a single file using its GUID.
Args:
-
object_id (str): The file's unique ID
+
object_id (str): The file's unique ID
path (str): Path to store the downloaded file at
-
"""
+
"""
try:
url = self.get_presigned_url(object_id)
except Exception as e:
-
logging.critical(f"Unable to get a presigned URL for download: {e}")
+
logging.critical(f"Unable to get a presigned URL for download: {e}")
return False
-
response = requests.get(url["url"], stream=True)
+
response = requests.get(url["url"], stream=True)
if response.status_code != 200:
-
logging.error(f"Response code: {response.status_code}")
+
logging.error(f"Response code: {response.status_code}")
if response.status_code >= 500:
for _ in range(MAX_RETRIES):
-
logging.info("Retrying now...")
+
logging.info("Retrying now...")
# NOTE could be updated with exponential backoff
time.sleep(1)
-
response = requests.get(url["url"], stream=True)
+
response = requests.get(url["url"], stream=True)
if response.status == 200:
break
if response.status != 200:
-
logging.critical("Response status not 200, try again later")
+
logging.critical("Response status not 200, try again later")
return False
else:
return False
response.raise_for_status()
-
total_size_in_bytes = int(response.headers.get("content-length"))
+
total_size_in_bytes = int(response.headers.get("content-length"))
total_downloaded = 0
index = Gen3Index(self._auth_provider)
record = index.get_record(object_id)
-
filename = record["file_name"]
+
filename = record["file_name"]
out_path = Gen3File._ensure_dirpath_exists(Path(path))
-
with open(os.path.join(out_path, filename), "wb") as f:
+
with open(os.path.join(out_path, filename), "wb") as f:
for data in response.iter_content(4096):
total_downloaded += len(data)
f.write(data)
if total_size_in_bytes == total_downloaded:
-
logging.info(f"File {filename} downloaded successfully")
+
logging.info(f"File {filename} downloaded successfully")
else:
-
logging.error(f"File {filename} not downloaded successfully")
+
logging.error(f"File {filename} not downloaded successfully")
return False
return True
@@ -275,32 +275,32 @@ Source code for gen3.file
def upload_file_to_guid(
self, guid, file_name, protocol=None, expires_in=None, bucket=None
):
- """
+ """
Get a presigned url for a file to upload to the specified existing GUID
Args:
file_name (str): file_name to use for upload
- protocol (str): Storage protocol to use for upload: "s3", "az".
- If this isn't set, the default will be "s3"
+ protocol (str): Storage protocol to use for upload: "s3", "az".
+ If this isn't set, the default will be "s3"
expires_in (int): Amount in seconds that the signed url will expire from datetime.utcnow().
Be sure to use a positive integer.
This value will also be treated as <= MAX_PRESIGNED_URL_TTL in the fence configuration.
- bucket (str): Bucket to upload to. The bucket must be configured in the Fence instance's
+ bucket (str): Bucket to upload to. The bucket must be configured in the Fence instance's
`ALLOWED_DATA_UPLOAD_BUCKETS` setting. If not specified, Fence defaults to the
`DATA_UPLOAD_BUCKET` setting.
Returns:
Document: json representation for the file upload
- """
- url = f"{self._endpoint}/user/data/upload/{guid}"
+ """
+ url = f"{self._endpoint}/user/data/upload/{guid}"
params = {}
if protocol:
- params["protocol"] = protocol
+ params["protocol"] = protocol
if expires_in:
- params["expires_in"] = expires_in
+ params["expires_in"] = expires_in
if file_name:
- params["file_name"] = file_name
+ params["file_name"] = file_name
if bucket:
- params["bucket"] = bucket
+ params["bucket"] = bucket
url_parts = list(urlparse(url))
query = dict(parse_qsl(url_parts[4]))
diff --git a/docs/_build/html/_modules/gen3/index.html b/docs/_build/html/_modules/gen3/index.html
index b393f29c4..694009f6e 100644
--- a/docs/_build/html/_modules/gen3/index.html
+++ b/docs/_build/html/_modules/gen3/index.html
@@ -50,7 +50,7 @@ Source code for gen3.index
[docs]
class Gen3Index:
-
"""
+
"""
A class for interacting with the Gen3 Index services.
@@ -62,26 +62,26 @@
Source code for gen3.index
This generates the Gen3Index class pointed at the sandbox commons while
using the credentials.json downloaded from the commons profile page.
- >>> auth = Gen3Auth(refresh_file="credentials.json")
+ >>> auth = Gen3Auth(refresh_file="credentials.json")
... index = Gen3Index(auth)
- """
+ """
- def __init__(self, endpoint=None, auth_provider=None, service_location="index"):
+ def __init__(self, endpoint=None, auth_provider=None, service_location="index"):
# legacy interface required endpoint as 1st arg
if endpoint and isinstance(endpoint, Gen3Auth):
auth_provider = endpoint
endpoint = None
if auth_provider and isinstance(auth_provider, Gen3Auth):
endpoint = auth_provider.endpoint
- endpoint = endpoint.strip("/")
+ endpoint = endpoint.strip("/")
# if running locally, indexd is deployed by itself without a location relative
# to the commons
- if "http://localhost" in endpoint:
- service_location = ""
+ if "http://localhost" in endpoint:
+ service_location = ""
if not endpoint.endswith(service_location):
- endpoint += "/" + service_location
+ endpoint += "/" + service_location
self.endpoint = endpoint
self.client = client.IndexClient(endpoint, auth=auth_provider)
@@ -90,29 +90,29 @@ Source code for gen3.index
[docs]
def is_healthy(self):
-
"""
+
"""
Return if indexd is healthy or not
-
"""
+
"""
try:
-
response = self.client._get("_status")
+
response = self.client._get("_status")
response.raise_for_status()
except Exception:
return False
-
return response.text == "Healthy"
+ return response.text == "Healthy"
[docs]
@backoff.on_exception(backoff.expo, Exception, **DEFAULT_BACKOFF_SETTINGS)
def get_version(self):
-
"""
+
"""
Return the version of indexd
-
"""
-
response = self.client._get("_version")
+
"""
+
response = self.client._get("_version")
raise_for_status_and_print_error(response)
return response.json()
@@ -121,12 +121,12 @@ Source code for gen3.index
[docs]
@backoff.on_exception(backoff.expo, Exception, **DEFAULT_BACKOFF_SETTINGS)
def get_stats(self):
- """
+ """
Return basic info about the records in indexd
- """
- response = self.client._get("_stats")
+ """
+ response = self.client._get("_stats")
raise_for_status_and_print_error(response)
return response.json()
@@ -135,31 +135,31 @@ Source code for gen3.index
[docs]
@backoff.on_exception(backoff.expo, Exception, **DEFAULT_BACKOFF_SETTINGS)
def get_all_records(self, limit=None, paginate=False):
- """
+ """
Get a list of all records
- """
+ """
all_records = []
- url = "index/"
+ url = "index/"
if limit:
- url += f"?limit={limit}"
+ url += f"?limit={limit}"
response = self.client._get(url)
raise_for_status_and_print_error(response)
- records = response.json().get("records")
+ records = response.json().get("records")
all_records.extend(records)
if paginate and records:
previous_did = None
- start_did = records[-1].get("did")
+ start_did = records[-1].get("did")
while start_did != previous_did:
previous_did = start_did
- params = {"start": f"{start_did}"}
+ params = {"start": f"{start_did}"}
url_parts = list(urllib.parse.urlparse(url))
query = dict(urllib.parse.parse_qsl(url_parts[4]))
query.update(params)
@@ -170,11 +170,11 @@ Source code for gen3.index
response = self.client._get(url)
raise_for_status_and_print_error(response)
- records = response.json().get("records")
+ records = response.json().get("records")
all_records.extend(records)
if records:
- start_did = response.json().get("records")[-1].get("did")
+ start_did = response.json().get("records")[-1].get("did")
return all_records
@@ -183,33 +183,33 @@ Source code for gen3.index
[docs]
@backoff.on_exception(backoff.expo, Exception, **DEFAULT_BACKOFF_SETTINGS)
def get_records_on_page(self, limit=None, page=None):
- """
+ """
Get a list of all records given the page and page size limit
- """
+ """
params = {}
- url = "index/"
+ url = "index/"
if limit is not None:
- params["limit"] = limit
+ params["limit"] = limit
if page is not None:
- params["page"] = page
+ params["page"] = page
query = urllib.parse.urlencode(params)
- response = self.client._get(url + "?" + query)
+ response = self.client._get(url + "?" + query)
raise_for_status_and_print_error(response)
- return response.json().get("records")
+ return response.json().get("records")
[docs]
@backoff.on_exception(backoff.expo, Exception, **DEFAULT_BACKOFF_SETTINGS)
async def async_get_record(self, guid=None, _ssl=None):
-
"""
+
"""
Asynchronous function to request a record from indexd.
Args:
@@ -217,8 +217,8 @@
Source code for gen3.index
Returns:
dict: indexd record
- """
- url = f"{self.client.url}/index/{guid}"
+ """
+ url = f"{self.client.url}/index/{guid}"
async with aiohttp.ClientSession() as session:
async with session.get(url, ssl=_ssl) as response:
raise_for_status_and_print_error(response)
@@ -231,7 +231,7 @@ Source code for gen3.index
[docs]
@backoff.on_exception(backoff.expo, Exception, **DEFAULT_BACKOFF_SETTINGS)
async def async_get_records_on_page(self, limit=None, page=None, _ssl=None):
- """
+ """
Asynchronous function to request a page from indexd.
Args:
@@ -239,33 +239,33 @@ Source code for gen3.index
Returns:
List[dict]: List of indexd records from the page
- """
+ """
all_records = []
params = {}
if limit is not None:
- params["limit"] = limit
+ params["limit"] = limit
if page is not None:
- params["page"] = page
+ params["page"] = page
query = urllib.parse.urlencode(params)
- url = f"{self.client.url}/index" + "?" + query
+ url = f"{self.client.url}/index" + "?" + query
async with aiohttp.ClientSession() as session:
async with session.get(url, ssl=_ssl) as response:
response = await response.json()
- return response.get("records")
+ return response.get("records")
[docs]
@backoff.on_exception(backoff.expo, Exception, **DEFAULT_BACKOFF_SETTINGS)
async def async_get_records_from_checksum(
-
self, checksum, checksum_type="md5", _ssl=None
+
self, checksum, checksum_type="md5", _ssl=None
):
-
"""
+
"""
Asynchronous function to request records from indexd matching checksum.
Args:
@@ -274,27 +274,27 @@
Source code for gen3.index
Returns:
List[dict]: List of indexd records
- """
+ """
all_records = []
params = {}
- params["hash"] = f"{checksum_type}:{checksum}"
+ params["hash"] = f"{checksum_type}:{checksum}"
query = urllib.parse.urlencode(params)
- url = f"{self.client.url}/index" + "?" + query
+ url = f"{self.client.url}/index" + "?" + query
async with aiohttp.ClientSession() as session:
async with session.get(url, ssl=_ssl) as response:
response = await response.json()
- return response.get("records")
+
return response.get("records")
[docs]
@backoff.on_exception(backoff.expo, Exception, **DEFAULT_BACKOFF_SETTINGS)
def get(self, guid, dist_resolution=True):
-
"""
+
"""
Get the metadata associated with the given id, alias, or
distributed identifier
@@ -305,7 +305,7 @@
Source code for gen3.index
dist_resolution: boolean
- *optional* Specify if we want distributed dist_resolution or not
- """
+ """
rec = self.client.global_get(guid, dist_resolution)
if not rec:
@@ -318,7 +318,7 @@ Source code for gen3.index
[docs]
@backoff.on_exception(backoff.expo, Exception, **DEFAULT_BACKOFF_SETTINGS)
def get_urls(self, size=None, hashes=None, guids=None):
- """
+ """
Get a list of urls that match query params
@@ -330,11 +330,11 @@ Source code for gen3.index
guids: list
- list of ids
- """
+ """
if guids:
- guids = ",".join(guids)
- p = {"size": size, "hash": hashes, "ids": guids}
- urls = self.client._get("urls", params=p).json()
+ guids = ",".join(guids)
+ p = {"size": size, "hash": hashes, "ids": guids}
+ urls = self.client._get("urls", params=p).json()
return [url for _, url in urls.items()]
@@ -342,11 +342,11 @@ Source code for gen3.index
[docs]
@backoff.on_exception(backoff.expo, Exception, **DEFAULT_BACKOFF_SETTINGS)
def get_record(self, guid):
- """
+ """
Get the metadata associated with a given id
- """
+ """
rec = self.client.get(guid)
if not rec:
@@ -359,11 +359,11 @@ Source code for gen3.index
[docs]
@backoff.on_exception(backoff.expo, Exception, **DEFAULT_BACKOFF_SETTINGS)
def get_record_doc(self, guid):
- """
+ """
Get the metadata associated with a given id
- """
+ """
return self.client.get(guid)
@@ -371,16 +371,16 @@ Source code for gen3.index
[docs]
@backoff.on_exception(backoff.expo, Exception, **DEFAULT_BACKOFF_SETTINGS)
def get_with_params(self, params=None):
- """
+ """
Return a document object corresponding to the supplied parameters, such
- as ``{'hashes': {'md5': '...'}, 'size': '...', 'metadata': {'file_state': '...'}}``.
+ as ``{'hashes': {'md5': '...'}, 'size': '...', 'metadata': {'file_state': '...'}}``.
- need to include all the hashes in the request
- index client like signpost or indexd will need to handle the
- query param `'hash': 'hash_type:hash'`
+ query param `'hash': 'hash_type:hash'`
- """
+ """
rec = self.client.get_with_params(params)
if not rec:
@@ -393,12 +393,12 @@ Source code for gen3.index
[docs]
@backoff.on_exception(backoff.expo, Exception, **DEFAULT_BACKOFF_SETTINGS)
async def async_get_with_params(self, params, _ssl=None):
- """
+ """
Return a document object corresponding to the supplied parameter
- need to include all the hashes in the request
- - need to handle the query param `'hash': 'hash_type:hash'`
+ - need to handle the query param `'hash': 'hash_type:hash'`
Args:
params (dict): params to search with
@@ -407,9 +407,9 @@ Source code for gen3.index
Returns:
Document: json representation of an entry in indexd
- """
+ """
query_params = urllib.parse.urlencode(params)
- url = f"{self.client.url}/index/?{query_params}"
+ url = f"{self.client.url}/index/?{query_params}"
async with aiohttp.ClientSession() as session:
async with session.get(url, ssl=_ssl) as response:
await response.raise_for_status()
@@ -422,7 +422,7 @@ Source code for gen3.index
[docs]
@backoff.on_exception(backoff.expo, Exception, **DEFAULT_BACKOFF_SETTINGS)
def get_latest_version(self, guid, has_version=False):
- """
+ """
Get the metadata of the latest index record version associated
with the given id
@@ -433,7 +433,7 @@ Source code for gen3.index
has_version: boolean
- *optional* exclude entries without a version
- """
+ """
rec = self.client.get_latest_version(guid, has_version)
if not rec:
@@ -446,7 +446,7 @@ Source code for gen3.index
[docs]
@backoff.on_exception(backoff.expo, Exception, **DEFAULT_BACKOFF_SETTINGS)
def get_versions(self, guid):
- """
+ """
Get the metadata of index record version associated with the
given id
@@ -455,8 +455,8 @@ Source code for gen3.index
guid: string
- record id
- """
- response = self.client._get(f"/index/{guid}/versions")
+ """
+ response = self.client._get(f"/index/{guid}/versions")
raise_for_status_and_print_error(response)
versions = response.json()
@@ -485,13 +485,13 @@ Source code for gen3.index
content_created_date=None,
content_updated_date=None,
):
- """
+ """
Create a new record and add it to the index
Args:
hashes (dict): {hash type: hash value,}
- eg ``hashes={'md5': ab167e49d25b488939b1ede42752458b'}``
+ eg ``hashes={'md5': ab167e49d25b488939b1ede42752458b'}``
size (int): file size metadata associated with a given uuid
did (str): provide a UUID for the new indexd to be made
urls (list): list of URLs where you can download the UUID
@@ -508,26 +508,26 @@ Source code for gen3.index
Returns:
Document: json representation of an entry in indexd
- """
+ """
if urls is None:
urls = []
json = {
- "urls": urls,
- "hashes": hashes,
- "size": size,
- "file_name": file_name,
- "metadata": metadata,
- "urls_metadata": urls_metadata,
- "baseid": baseid,
- "acl": acl,
- "authz": authz,
- "version": version,
- "description": description,
- "content_created_date": content_created_date,
- "content_updated_date": content_updated_date,
+ "urls": urls,
+ "hashes": hashes,
+ "size": size,
+ "file_name": file_name,
+ "metadata": metadata,
+ "urls_metadata": urls_metadata,
+ "baseid": baseid,
+ "acl": acl,
+ "authz": authz,
+ "version": version,
+ "description": description,
+ "content_created_date": content_created_date,
+ "content_updated_date": content_updated_date,
}
if did:
- json["did"] = did
+ json["did"] = did
rec = self.client.create(**json)
return rec.to_json()
@@ -554,12 +554,12 @@ Source code for gen3.index
content_created_date=None,
content_updated_date=None,
):
- """
+ """
Asynchronous function to create a record in indexd.
Args:
hashes (dict): {hash type: hash value,}
- eg ``hashes={'md5': ab167e49d25b488939b1ede42752458b'}``
+ eg ``hashes={'md5': ab167e49d25b488939b1ede42752458b'}``
size (int): file size metadata associated with a given uuid
did (str): provide a UUID for the new indexd to be made
urls (list): list of URLs where you can download the UUID
@@ -576,45 +576,45 @@ Source code for gen3.index
Returns:
Document: json representation of an entry in indexd
- """
+ """
async with aiohttp.ClientSession() as session:
if urls is None:
urls = []
json = {
- "form": "object",
- "hashes": hashes,
- "size": size,
- "urls": urls or [],
+ "form": "object",
+ "hashes": hashes,
+ "size": size,
+ "urls": urls or [],
}
if did:
- json["did"] = did
+ json["did"] = did
if file_name:
- json["file_name"] = file_name
+ json["file_name"] = file_name
if metadata:
- json["metadata"] = metadata
+ json["metadata"] = metadata
if baseid:
- json["baseid"] = baseid
+ json["baseid"] = baseid
if acl:
- json["acl"] = acl
+ json["acl"] = acl
if urls_metadata:
- json["urls_metadata"] = urls_metadata
+ json["urls_metadata"] = urls_metadata
if version:
- json["version"] = version
+ json["version"] = version
if authz:
- json["authz"] = authz
+ json["authz"] = authz
if description:
- json["description"] = description
+ json["description"] = description
if content_created_date:
- json["content_created_date"] = content_created_date
+ json["content_created_date"] = content_created_date
if content_updated_date:
- json["content_updated_date"] = content_updated_date
+ json["content_updated_date"] = content_updated_date
# aiohttp only allows basic auth with their built in auth, so we
# need to manually add JWT auth header
- headers = {"Authorization": self.client.auth._get_auth_value()}
+ headers = {"Authorization": self.client.auth._get_auth_value()}
async with session.post(
- f"{self.client.url}/index/",
+ f"{self.client.url}/index/",
json=json,
headers=headers,
ssl=_ssl,
@@ -629,29 +629,29 @@ Source code for gen3.index
[docs]
@backoff.on_exception(backoff.expo, Exception, **DEFAULT_BACKOFF_SETTINGS)
def create_blank(self, uploader, file_name=None):
- """
+ """
Create a blank record
Args:
json - json in the format:
{
- 'uploader': type(string)
- 'file_name': type(string) (optional*)
+ 'uploader': type(string)
+ 'file_name': type(string) (optional*)
}
- """
- json = {"uploader": uploader, "file_name": file_name}
+ """
+ json = {"uploader": uploader, "file_name": file_name}
response = self.client._post(
- "index/blank",
- headers={"content-type": "application/json"},
+ "index/blank",
+ headers={"content-type": "application/json"},
auth=self.client.auth,
data=client.json_dumps(json),
)
raise_for_status_and_print_error(response)
rec = response.json()
- return self.get_record(rec["did"])
+ return self.get_record(rec["did"])
@@ -674,7 +674,7 @@
Source code for gen3.index
content_created_date=None,
content_updated_date=None,
):
- """
+ """
Add new version for the document associated to the provided uuid
@@ -686,7 +686,7 @@ Source code for gen3.index
Args:
guid: (string): record id
hashes (dict): {hash type: hash value,}
- eg ``hashes={'md5': ab167e49d25b488939b1ede42752458b'}``
+ eg ``hashes={'md5': ab167e49d25b488939b1ede42752458b'}``
size (int): file size metadata associated with a given uuid
did (str): provide a UUID for the new indexd to be made
urls (list): list of URLs where you can download the UUID
@@ -706,38 +706,38 @@ Source code for gen3.index
sufficient. Note: it is a good idea to add a version
number
- """
+ """
if urls is None:
urls = []
json = {
- "urls": urls,
- "form": "object",
- "hashes": hashes,
- "size": size,
- "file_name": file_name,
- "metadata": metadata,
- "urls_metadata": urls_metadata,
- "acl": acl,
- "authz": authz,
- "version": version,
- "description": description,
- "content_created_date": content_created_date,
- "content_updated_date": content_updated_date,
+ "urls": urls,
+ "form": "object",
+ "hashes": hashes,
+ "size": size,
+ "file_name": file_name,
+ "metadata": metadata,
+ "urls_metadata": urls_metadata,
+ "acl": acl,
+ "authz": authz,
+ "version": version,
+ "description": description,
+ "content_created_date": content_created_date,
+ "content_updated_date": content_updated_date,
}
if did:
- json["did"] = did
+ json["did"] = did
response = self.client._post(
- "index",
+ "index",
guid,
- headers={"content-type": "application/json"},
+ headers={"content-type": "application/json"},
data=client.json_dumps(json),
auth=self.client.auth,
)
raise_for_status_and_print_error(response)
rec = response.json()
- if rec and "did" in rec:
- return self.get_record(rec["did"])
+ if rec and "did" in rec:
+ return self.get_record(rec["did"])
return None
@@ -745,7 +745,7 @@ Source code for gen3.index
[docs]
@backoff.on_exception(backoff.expo, Exception, **DEFAULT_BACKOFF_SETTINGS)
def get_records(self, dids):
- """
+ """
Get a list of documents given a list of dids
@@ -756,10 +756,10 @@ Source code for gen3.index
Returns:
list: json representing index records
- """
+ """
try:
response = self.client._post(
- "bulk/documents", json=dids, auth=self.client.auth
+ "bulk/documents", json=dids, auth=self.client.auth
)
except requests.HTTPError as exception:
if exception.response.status_code == 404:
@@ -776,7 +776,7 @@ Source code for gen3.index
[docs]
@backoff.on_exception(backoff.expo, Exception, **DEFAULT_BACKOFF_SETTINGS)
def update_blank(self, guid, rev, hashes, size, urls=None, authz=None):
- """
+ """
Update only hashes and size for a blank index
@@ -784,21 +784,21 @@ Source code for gen3.index
guid (string): record id
rev (string): data revision - simple consistency mechanism
hashes (dict): {hash type: hash value,}
- eg ``hashes={'md5': ab167e49d25b488939b1ede42752458b'}``
+ eg ``hashes={'md5': ab167e49d25b488939b1ede42752458b'}``
size (int): file size metadata associated with a given uuid
- """
- params = {"rev": rev}
- json = {"hashes": hashes, "size": size}
+ """
+ params = {"rev": rev}
+ json = {"hashes": hashes, "size": size}
if urls:
- json["urls"] = urls
+ json["urls"] = urls
if authz:
- json["authz"] = authz
+ json["authz"] = authz
response = self.client._put(
- "index/blank",
+ "index/blank",
guid,
- headers={"content-type": "application/json"},
+ headers={"content-type": "application/json"},
params=params,
auth=self.client.auth,
data=client.json_dumps(json),
@@ -806,7 +806,7 @@ Source code for gen3.index
raise_for_status_and_print_error(response)
rec = response.json()
- return self.get_record(rec["did"])
+ return self.get_record(rec["did"])
@@ -826,7 +826,7 @@
Source code for gen3.index
content_created_date=None,
content_updated_date=None,
):
- """
+ """
Update an existing entry in the index
@@ -837,27 +837,27 @@ Source code for gen3.index
- index record information that needs to be updated.
- can not update size or hash, use new version for that
- """
+ """
updatable_attrs = {
- "file_name": file_name,
- "urls": urls,
- "version": version,
- "metadata": metadata,
- "acl": acl,
- "authz": authz,
- "urls_metadata": urls_metadata,
- "description": description,
- "content_created_date": content_created_date,
- "content_updated_date": content_updated_date,
+ "file_name": file_name,
+ "urls": urls,
+ "version": version,
+ "metadata": metadata,
+ "acl": acl,
+ "authz": authz,
+ "urls_metadata": urls_metadata,
+ "description": description,
+ "content_created_date": content_created_date,
+ "content_updated_date": content_updated_date,
}
rec = self.client.get(guid)
if not rec:
raise ValueError(
- f"No indexd record found for GUID '{guid}' at '{self.endpoint}'"
+ f"No indexd record found for GUID '{guid}' at '{self.endpoint}'"
)
for k, v in updatable_attrs.items():
if v is not None:
- exec(f"rec.{k} = v")
+ exec(f"rec.{k} = v")
rec.patch()
return rec.to_json()
@@ -881,7 +881,7 @@ Source code for gen3.index
content_updated_date=None,
**kwargs,
):
- """
+ """
Asynchronous function to update a record in indexd.
Args:
@@ -890,47 +890,47 @@ Source code for gen3.index
body: json/dictionary format
- index record information that needs to be updated.
- can not update size or hash, use new version for that
- """
+ """
async with aiohttp.ClientSession() as session:
updatable_attrs = {
- "file_name": file_name,
- "urls": urls,
- "version": version,
- "metadata": metadata,
- "acl": acl,
- "authz": authz,
- "urls_metadata": urls_metadata,
- "description": description,
- "content_created_date": content_created_date,
- "content_updated_date": content_updated_date,
+ "file_name": file_name,
+ "urls": urls,
+ "version": version,
+ "metadata": metadata,
+ "acl": acl,
+ "authz": authz,
+ "urls_metadata": urls_metadata,
+ "description": description,
+ "content_created_date": content_created_date,
+ "content_updated_date": content_updated_date,
}
record = await self.async_get_record(guid)
- revision = record.get("rev")
+ revision = record.get("rev")
for key, value in updatable_attrs.items():
if value is not None:
record[key] = value
- del record["created_date"]
- del record["rev"]
- del record["updated_date"]
- del record["version"]
- del record["uploader"]
- del record["form"]
- del record["urls_metadata"]
- del record["baseid"]
- del record["size"]
- del record["hashes"]
- del record["did"]
+ del record["created_date"]
+ del record["rev"]
+ del record["updated_date"]
+ del record["version"]
+ del record["uploader"]
+ del record["form"]
+ del record["urls_metadata"]
+ del record["baseid"]
+ del record["size"]
+ del record["hashes"]
+ del record["did"]
- logging.info(f"PUT-ing record: {record}")
+ logging.info(f"PUT-ing record: {record}")
# aiohttp only allows basic auth with their built in auth, so we
# need to manually add JWT auth header
- headers = {"Authorization": self.client.auth._get_auth_value()}
+ headers = {"Authorization": self.client.auth._get_auth_value()}
async with session.put(
- f"{self.client.url}/index/{guid}?rev={revision}",
+ f"{self.client.url}/index/{guid}?rev={revision}",
json=record,
headers=headers,
ssl=_ssl,
@@ -947,7 +947,7 @@ Source code for gen3.index
[docs]
@backoff.on_exception(backoff.expo, Exception, **DEFAULT_BACKOFF_SETTINGS)
def delete_record(self, guid):
- """
+ """
Delete an entry from the index
@@ -957,7 +957,7 @@ Source code for gen3.index
Returns: Nothing
- """
+ """
rec = self.client.get(guid)
if rec:
rec.delete()
@@ -970,7 +970,7 @@ Source code for gen3.index
[docs]
@backoff.on_exception(backoff.expo, Exception, **DEFAULT_BACKOFF_SETTINGS)
def query_urls(self, pattern):
- """
+ """
Query all record URLs for given pattern
@@ -979,8 +979,8 @@ Source code for gen3.index
Returns:
List[records]: indexd records with urls matching pattern
- """
- response = self.client._get(f"/_query/urls/q?include={pattern}")
+ """
+ response = self.client._get(f"/_query/urls/q?include={pattern}")
raise_for_status_and_print_error(response)
return response.json()
@@ -989,7 +989,7 @@ Source code for gen3.index
[docs]
@backoff.on_exception(backoff.expo, Exception, **DEFAULT_BACKOFF_SETTINGS)
async def async_query_urls(self, pattern, _ssl=None):
- """
+ """
Asynchronous function to query urls from indexd.
Args:
@@ -997,10 +997,10 @@ Source code for gen3.index
Returns:
List[records]: indexd records with urls matching pattern
- """
- url = f"{self.client.url}/_query/urls/q?include={pattern}"
+ """
+ url = f"{self.client.url}/_query/urls/q?include={pattern}"
async with aiohttp.ClientSession() as session:
- logging.debug(f"request: {url}")
+ logging.debug(f"request: {url}")
async with session.get(url, ssl=_ssl) as response:
raise_for_status_and_print_error(response)
response = await response.json()
@@ -1014,44 +1014,44 @@ Source code for gen3.index
[docs]
@backoff.on_exception(backoff.expo, Exception, **DEFAULT_BACKOFF_SETTINGS)
def get_valid_guids(self, count=None):
- """
+ """
Get a list of valid GUIDs without indexing
Args:
count (int): number of GUIDs to request
Returns:
List[str]: list of valid indexd GUIDs
- """
- url = "/guid/mint"
+ """
+ url = "/guid/mint"
if count:
- url += f"?count={count}"
+ url += f"?count={count}"
response = self.client._get(url)
response.raise_for_status()
- return response.json().get("guids", [])
+ return response.json().get("guids", [])
[docs]
@backoff.on_exception(backoff.expo, Exception, **DEFAULT_BACKOFF_SETTINGS)
def get_guids_prefix(self):
-
"""
+
"""
Get the prefix for GUIDs if there is one
Returns:
str: prefix for this instance
-
"""
-
response = self.client._get("/guid/prefix")
+
"""
+
response = self.client._get("/guid/prefix")
response.raise_for_status()
-
return response.json().get("prefix")
+ return response.json().get("prefix")
def _print_func_name(function):
- return "{}.{}".format(function.__module__, function.__name__)
+ return "{}.{}".format(function.__module__, function.__name__)
def _print_kwargs(kwargs):
- return ", ".join("{}={}".format(k, repr(v)) for k, v in list(kwargs.items()))
+ return ", ".join("{}={}".format(k, repr(v)) for k, v in list(kwargs.items()))
diff --git a/docs/_build/html/_modules/gen3/jobs.html b/docs/_build/html/_modules/gen3/jobs.html
index ed3448caa..a81656f82 100644
--- a/docs/_build/html/_modules/gen3/jobs.html
+++ b/docs/_build/html/_modules/gen3/jobs.html
@@ -31,9 +31,9 @@
Source code for gen3.jobs
-"""
-Contains class for interacting with Gen3's Job Dispatching Service(s).
-"""
+"""
+Contains class for interacting with Gen3's Job Dispatching Service(s).
+"""
import aiohttp
import asyncio
import backoff
@@ -51,12 +51,12 @@ Source code for gen3.jobs
raise_for_status_and_print_error,
)
-# sower's "action" mapping to the relevant job
-INGEST_METADATA_JOB = "ingest-metadata-manifest"
-DBGAP_METADATA_JOB = "get-dbgap-metadata"
-INDEX_MANIFEST_JOB = "index-object-manifest"
-DOWNLOAD_MANIFEST_JOB = "download-indexd-manifest"
-MERGE_MANIFEST_JOB = "merge-manifests"
+# sower's "action" mapping to the relevant job
+INGEST_METADATA_JOB = "ingest-metadata-manifest"
+DBGAP_METADATA_JOB = "get-dbgap-metadata"
+INDEX_MANIFEST_JOB = "index-object-manifest"
+DOWNLOAD_MANIFEST_JOB = "download-indexd-manifest"
+MERGE_MANIFEST_JOB = "merge-manifests"
logging = get_logger(__name__)
@@ -64,19 +64,19 @@ Source code for gen3.jobs
[docs]
class Gen3Jobs:
-
"""
-
A class for interacting with the Gen3's Job Dispatching Service(s).
+
"""
+
A class for interacting with the Gen3's Job Dispatching Service(s).
Examples:
This generates the Gen3Jobs class pointed at the sandbox commons while
using the credentials.json downloaded from the commons profile page.
-
>>> auth = Gen3Auth(refresh_file="credentials.json")
+
>>> auth = Gen3Auth(refresh_file="credentials.json")
... jobs = Gen3Jobs(auth)
-
"""
+
"""
-
def __init__(self, endpoint=None, auth_provider=None, service_location="job"):
-
"""
+
def __init__(self, endpoint=None, auth_provider=None, service_location="job"):
+
"""
Initialization for instance of the class to setup basic endpoint info.
Args:
@@ -84,25 +84,25 @@
Source code for gen3.jobs
token, required for admin endpoints
service_location (str, optional): deployment location relative to the
endpoint provided
- """
+ """
# auth_provider legacy interface required endpoint as 1st arg
auth_provider = auth_provider or endpoint
- endpoint = auth_provider.endpoint.strip("/")
+ endpoint = auth_provider.endpoint.strip("/")
# if running locally, mds is deployed by itself without a location relative
# to the commons
- if "http://localhost" in endpoint:
- service_location = ""
+ if "http://localhost" in endpoint:
+ service_location = ""
if not endpoint.endswith(service_location):
- endpoint += "/" + service_location
+ endpoint += "/" + service_location
- self.endpoint = endpoint.rstrip("/")
+ self.endpoint = endpoint.rstrip("/")
self._auth_provider = auth_provider
[docs]
async def async_run_job_and_wait(self, job_name, job_input, _ssl=None, **kwargs):
-
"""
+
"""
Asynchronous function to create a job, wait for output, and return. Will
sleep in a linear delay until the job is done, starting with 1 second.
@@ -113,71 +113,71 @@
Source code for gen3.jobs
Returns:
Dict: Response from the endpoint
- """
+ """
job_create_response = await self.async_create_job(job_name, job_input)
- status = {"status": "Running"}
+ status = {"status": "Running"}
sleep_time = 3
- while status.get("status") == "Running":
- logging.info(f"job still running, waiting for {sleep_time} seconds...")
+ while status.get("status") == "Running":
+ logging.info(f"job still running, waiting for {sleep_time} seconds...")
time.sleep(sleep_time)
sleep_time *= 1.5
- status = await self.async_get_status(job_create_response.get("uid"))
- logging.info(f"{status}")
+ status = await self.async_get_status(job_create_response.get("uid"))
+ logging.info(f"{status}")
- logging.info(f"Job is finished!")
+ logging.info(f"Job is finished!")
- if status.get("status") != "Completed":
- raise Exception(f"Job status not complete: {status.get('status')}.")
+ if status.get("status") != "Completed":
+ raise Exception(f"Job status not complete: {status.get('status')}.")
- response = await self.async_get_output(job_create_response.get("uid"))
+ response = await self.async_get_output(job_create_response.get("uid"))
return response
[docs]
def is_healthy(self):
-
"""
+
"""
Return if is healthy or not
Returns:
bool: True if healthy
-
"""
+
"""
try:
response = requests.get(
-
self.endpoint + "/_status", auth=self._auth_provider
+
self.endpoint + "/_status", auth=self._auth_provider
)
raise_for_status_and_print_error(response)
except Exception as exc:
logging.error(exc)
return False
-
return response.text == "Healthy"
+
return response.text == "Healthy"
[docs]
@backoff.on_exception(backoff.expo, Exception, **DEFAULT_BACKOFF_SETTINGS)
def get_version(self):
-
"""
+
"""
Return the version
Returns:
str: the version
-
"""
-
response = requests.get(self.endpoint + "/_version", auth=self._auth_provider)
+
"""
+
response = requests.get(self.endpoint + "/_version", auth=self._auth_provider)
raise_for_status_and_print_error(response)
-
return response.json().get("version")
+ return response.json().get("version")
[docs]
@backoff.on_exception(backoff.expo, Exception, **DEFAULT_BACKOFF_SETTINGS)
def list_jobs(self):
-
"""
+
"""
List all jobs
-
"""
-
response = requests.get(self.endpoint + "/list", auth=self._auth_provider)
+
"""
+
response = requests.get(self.endpoint + "/list", auth=self._auth_provider)
raise_for_status_and_print_error(response)
return response.json()
@@ -186,7 +186,7 @@
Source code for gen3.jobs
[docs]
@backoff.on_exception(backoff.expo, Exception, **DEFAULT_BACKOFF_SETTINGS)
def create_job(self, job_name, job_input):
- """
+ """
Create a job with given name and input
Args:
@@ -195,10 +195,10 @@ Source code for gen3.jobs
Returns:
Dict: Response from the endpoint
- """
- data = {"action": job_name, "input": job_input}
+ """
+ data = {"action": job_name, "input": job_input}
response = requests.post(
- self.endpoint + "/dispatch", json=data, auth=self._auth_provider
+ self.endpoint + "/dispatch", json=data, auth=self._auth_provider
)
raise_for_status_and_print_error(response)
return response.json()
@@ -207,14 +207,14 @@ Source code for gen3.jobs
@backoff.on_exception(backoff.expo, Exception, **DEFAULT_BACKOFF_SETTINGS)
async def async_create_job(self, job_name, job_input, _ssl=None, **kwargs):
async with aiohttp.ClientSession() as session:
- url = self.endpoint + f"/dispatch"
+ url = self.endpoint + f"/dispatch"
url_with_params = append_query_params(url, **kwargs)
- data = json.dumps({"action": job_name, "input": job_input})
+ data = json.dumps({"action": job_name, "input": job_input})
# aiohttp only allows basic auth with their built in auth, so we
# need to manually add JWT auth header
- headers = {"Authorization": self._auth_provider._get_auth_value()}
+ headers = {"Authorization": self._auth_provider._get_auth_value()}
async with session.post(
url_with_params, data=data, headers=headers, ssl=_ssl
@@ -227,11 +227,11 @@ Source code for gen3.jobs
[docs]
@backoff.on_exception(backoff.expo, Exception, **DEFAULT_BACKOFF_SETTINGS)
def get_status(self, job_id):
- """
+ """
Get the status of a previously created job
- """
+ """
response = requests.get(
- self.endpoint + f"/status?UID={job_id}", auth=self._auth_provider
+ self.endpoint + f"/status?UID={job_id}", auth=self._auth_provider
)
raise_for_status_and_print_error(response)
return response.json()
@@ -240,12 +240,12 @@ Source code for gen3.jobs
@backoff.on_exception(backoff.expo, Exception, **DEFAULT_BACKOFF_SETTINGS)
async def async_get_status(self, job_id, _ssl=None, **kwargs):
async with aiohttp.ClientSession() as session:
- url = self.endpoint + f"/status?UID={job_id}"
+ url = self.endpoint + f"/status?UID={job_id}"
url_with_params = append_query_params(url, **kwargs)
# aiohttp only allows basic auth with their built in auth, so we
# need to manually add JWT auth header
- headers = {"Authorization": self._auth_provider._get_auth_value()}
+ headers = {"Authorization": self._auth_provider._get_auth_value()}
async with session.get(
url_with_params, headers=headers, ssl=_ssl
@@ -258,11 +258,11 @@ Source code for gen3.jobs
[docs]
@backoff.on_exception(backoff.expo, Exception, **DEFAULT_BACKOFF_SETTINGS)
def get_output(self, job_id):
- """
+ """
Get the output of a previously completed job
- """
+ """
response = requests.get(
- self.endpoint + f"/output?UID={job_id}", auth=self._auth_provider
+ self.endpoint + f"/output?UID={job_id}", auth=self._auth_provider
)
raise_for_status_and_print_error(response)
return response.json()
@@ -271,12 +271,12 @@ Source code for gen3.jobs
@backoff.on_exception(backoff.expo, Exception, **DEFAULT_BACKOFF_SETTINGS)
async def async_get_output(self, job_id, _ssl=None, **kwargs):
async with aiohttp.ClientSession() as session:
- url = self.endpoint + f"/output?UID={job_id}"
+ url = self.endpoint + f"/output?UID={job_id}"
url_with_params = append_query_params(url, **kwargs)
# aiohttp only allows basic auth with their built in auth, so we
# need to manually add JWT auth header
- headers = {"Authorization": self._auth_provider._get_auth_value()}
+ headers = {"Authorization": self._auth_provider._get_auth_value()}
async with session.get(
url_with_params, headers=headers, ssl=_ssl
diff --git a/docs/_build/html/_modules/gen3/metadata.html b/docs/_build/html/_modules/gen3/metadata.html
index cb785a927..7e1dc353c 100644
--- a/docs/_build/html/_modules/gen3/metadata.html
+++ b/docs/_build/html/_modules/gen3/metadata.html
@@ -31,9 +31,9 @@
Source code for gen3.metadata
-"""
-Contains class for interacting with Gen3's Metadata Service.
-"""
+"""
+Contains class for interacting with Gen3's Metadata Service.
+"""
import aiohttp
import backoff
from datetime import datetime
@@ -66,24 +66,24 @@ Source code for gen3.metadata
logging = get_logger(__name__)
-PACKAGE_CONTENTS_STANDARD_KEY = "package_contents"
+PACKAGE_CONTENTS_STANDARD_KEY = "package_contents"
PACKAGE_CONTENTS_SCHEMA = {
- "type": "array",
- "items": {
- "type": "object",
- "properties": {
- "file_name": {
- "type": "string",
+ "type": "array",
+ "items": {
+ "type": "object",
+ "properties": {
+ "file_name": {
+ "type": "string",
},
- "size": {
- "type": "integer",
+ "size": {
+ "type": "integer",
},
- "hashes": {
- "type": "object",
+ "hashes": {
+ "type": "object",
},
},
- "required": ["file_name"],
- "additionalProperties": True,
+ "required": ["file_name"],
+ "additionalProperties": True,
},
}
@@ -91,29 +91,29 @@ Source code for gen3.metadata
[docs]
class Gen3Metadata:
-
"""
+
"""
A class for interacting with the Gen3 Metadata services.
Examples:
This generates the Gen3Metadata class pointed at the sandbox commons while
using the credentials.json downloaded from the commons profile page.
-
>>> auth = Gen3Auth(refresh_file="credentials.json")
+
>>> auth = Gen3Auth(refresh_file="credentials.json")
... metadata = Gen3Metadata(auth)
Attributes:
endpoint (str): public endpoint for reading/querying metadata - only necessary if auth_provider not provided
auth_provider (Gen3Auth): auth manager
-
"""
+
"""
def __init__(
self,
endpoint=None,
auth_provider=None,
-
service_location="mds",
-
admin_endpoint_suffix="-admin",
+
service_location="mds",
+
admin_endpoint_suffix="-admin",
):
-
"""
+
"""
Initialization for instance of the class to setup basic endpoint info.
Args:
@@ -122,59 +122,59 @@
Source code for gen3.metadata
token, required for admin endpoints
service_location (str, optional): deployment location relative to the
endpoint provided
- """
+ """
# legacy interface required endpoint as 1st arg
if endpoint and isinstance(endpoint, Gen3Auth):
auth_provider = endpoint
endpoint = None
if auth_provider and isinstance(auth_provider, Gen3Auth):
endpoint = auth_provider.endpoint
- endpoint = endpoint.strip("/")
+ endpoint = endpoint.strip("/")
# if running locally, mds is deployed by itself without a location relative
# to the commons
- if "http://localhost" in endpoint:
- service_location = ""
- admin_endpoint_suffix = ""
+ if "http://localhost" in endpoint:
+ service_location = ""
+ admin_endpoint_suffix = ""
if not endpoint.endswith(service_location):
- endpoint += "/" + service_location
+ endpoint += "/" + service_location
- self.endpoint = endpoint.rstrip("/")
- self.admin_endpoint = endpoint.rstrip("/") + admin_endpoint_suffix
+ self.endpoint = endpoint.rstrip("/")
+ self.admin_endpoint = endpoint.rstrip("/") + admin_endpoint_suffix
self._auth_provider = auth_provider
+ return response.json().get("status") == "OK"
@@ -183,14 +183,14 @@
Source code for gen3.metadata
[docs]
@backoff.on_exception(backoff.expo, Exception, **DEFAULT_BACKOFF_SETTINGS)
def get_index_key_paths(self):
- """
+ """
List all the metadata key paths indexed in the database.
Returns:
List: list of metadata key paths
- """
+ """
response = requests.get(
- self.admin_endpoint + "/metadata_index", auth=self._auth_provider
+ self.admin_endpoint + "/metadata_index", auth=self._auth_provider
)
response.raise_for_status()
return response.json()
@@ -200,14 +200,14 @@
Source code for gen3.metadata
[docs]
@backoff.on_exception(backoff.expo, Exception, **DEFAULT_BACKOFF_SETTINGS)
def create_index_key_path(self, path):
- """
+ """
Create a metadata key path indexed in the database.
Args:
path (str): metadata key path
- """
+ """
response = requests.post(
- self.admin_endpoint + f"/metadata_index/{path}", auth=self._auth_provider
+ self.admin_endpoint + f"/metadata_index/{path}", auth=self._auth_provider
)
response.raise_for_status()
return response.json()
@@ -217,14 +217,14 @@
Source code for gen3.metadata
[docs]
@backoff.on_exception(backoff.expo, Exception, **DEFAULT_BACKOFF_SETTINGS)
def delete_index_key_path(self, path):
- """
+ """
List all the metadata key paths indexed in the database.
Args:
path (str): metadata key path
- """
+ """
response = requests.delete(
- self.admin_endpoint + f"/metadata_index/{path}", auth=self._auth_provider
+ self.admin_endpoint + f"/metadata_index/{path}", auth=self._auth_provider
)
response.raise_for_status()
return response
@@ -242,38 +242,38 @@
Source code for gen3.metadata
use_agg_mds=False,
**kwargs,
):
- """
+ """
Query the metadata given a query.
Query format is based off the logic used in the service:
- '''
+ '''
Without filters, this will return all data. Add filters as query strings like this:
GET /metadata?a=1&b=2
This will match all records that have metadata containing all of:
- {"a": 1, "b": 2}
+ {"a": 1, "b": 2}
The values are always treated as strings for filtering. Nesting is supported:
GET /metadata?a.b.c=3
Matching records containing:
- {"a": {"b": {"c": 3}}}
+ {"a": {"b": {"c": 3}}}
Providing the same key with more than one value filters records whose value of the given key matches any of the given values. But values of different keys must all match. For example:
GET /metadata?a.b.c=3&a.b.c=33&a.b.d=4
Matches these:
- {"a": {"b": {"c": 3, "d": 4}}}
- {"a": {"b": {"c": 33, "d": 4}}}
- {"a": {"b": {"c": "3", "d": 4, "e": 5}}}
- But won't match these:
+ {"a": {"b": {"c": 3, "d": 4}}}
+ {"a": {"b": {"c": 33, "d": 4}}}
+ {"a": {"b": {"c": "3", "d": 4, "e": 5}}}
+ But won't match these:
- {"a": {"b": {"c": 3}}}
- {"a": {"b": {"c": 3, "d": 5}}}
- {"a": {"b": {"d": 5}}}
- {"a": {"b": {"c": "333", "d": 4}}}
- '''
+ {"a": {"b": {"c": 3}}}
+ {"a": {"b": {"c": 3, "d": 5}}}
+ {"a": {"b": {"d": 5}}}
+ {"a": {"b": {"c": "333", "d": 4}}}
+ '''
Args:
query (str): mds query as defined by the metadata api
@@ -286,14 +286,14 @@ Source code for gen3.metadata
OR if return_full_metadata=True
Dict{guid: {metadata}}: Dictionary with GUIDs as keys and associated
metadata JSON blobs as values
- """
+ """
- url = self.endpoint + f"/metadata?{query}"
+ url = self.endpoint + f"/metadata?{query}"
url_with_params = append_query_params(
url, data=return_full_metadata, limit=limit, offset=offset, **kwargs
)
- logging.debug(f"hitting: {url_with_params}")
+ logging.debug(f"hitting: {url_with_params}")
response = requests.get(url_with_params, auth=self._auth_provider)
response.raise_for_status()
@@ -304,7 +304,7 @@ Source code for gen3.metadata
[docs]
@backoff.on_exception(backoff.expo, Exception, **BACKOFF_NO_LOG_IF_NOT_RETRIED)
async def async_get(self, guid, _ssl=None, **kwargs):
- """
+ """
Asynchronous function to get metadata
Args:
@@ -313,12 +313,12 @@ Source code for gen3.metadata
Returns:
Dict: metadata for given guid
- """
+ """
async with aiohttp.ClientSession() as session:
- url = self.endpoint + f"/metadata/{guid}"
+ url = self.endpoint + f"/metadata/{guid}"
url_with_params = append_query_params(url, **kwargs)
- logging.debug(f"hitting: {url_with_params}")
+ logging.debug(f"hitting: {url_with_params}")
async with session.get(url_with_params, ssl=_ssl) as response:
response.raise_for_status()
@@ -331,16 +331,16 @@ Source code for gen3.metadata
[docs]
@backoff.on_exception(backoff.expo, Exception, **BACKOFF_NO_LOG_IF_NOT_RETRIED)
def get(self, guid, **kwargs):
- """
+ """
Get the metadata associated with the guid
Args:
guid (str): guid to use
Returns:
Dict: metadata for given guid
- """
- url = self.endpoint + f"/metadata/{guid}"
+ """
+ url = self.endpoint + f"/metadata/{guid}"
url_with_params = append_query_params(url, **kwargs)
- logging.debug(f"hitting: {url_with_params}")
+ logging.debug(f"hitting: {url_with_params}")
response = requests.get(url_with_params, auth=self._auth_provider)
response.raise_for_status()
@@ -352,30 +352,30 @@ Source code for gen3.metadata
[docs]
@backoff.on_exception(backoff.expo, Exception, **DEFAULT_BACKOFF_SETTINGS)
def batch_create(self, metadata_list, overwrite=True, **kwargs):
- """
+ """
Create the list of metadata associated with the list of guids
Args:
- metadata_list (List[Dict{"guid": "", "data": {}}]): list of metadata
- objects in a specific format. Expects a dict with "guid" and "data"
- fields where "data" is another JSON blob to add to the mds
+ metadata_list (List[Dict{"guid": "", "data": {}}]): list of metadata
+ objects in a specific format. Expects a dict with "guid" and "data"
+ fields where "data" is another JSON blob to add to the mds
overwrite (bool, optional): whether or not to overwrite existing data
- """
- url = self.admin_endpoint + f"/metadata"
+ """
+ url = self.admin_endpoint + f"/metadata"
if len(metadata_list) > 1 and (
- "guid" not in metadata_list[0] and "data" not in metadata_list[0]
+ "guid" not in metadata_list[0] and "data" not in metadata_list[0]
):
logging.warning(
- "it looks like your metadata list for bulk create is malformed. "
- "the expected format is a list of dicts that have 2 keys: 'guid' "
- "and 'data', where 'guid' is a string and 'data' is another dict. "
- f"The first element doesn't match that pattern: {metadata_list[0]}"
+ "it looks like your metadata list for bulk create is malformed. "
+ "the expected format is a list of dicts that have 2 keys: 'guid' "
+ "and 'data', where 'guid' is a string and 'data' is another dict. "
+ f"The first element doesn't match that pattern: {metadata_list[0]}"
)
url_with_params = append_query_params(url, overwrite=overwrite, **kwargs)
- logging.debug(f"hitting: {url_with_params}")
- logging.debug(f"data: {metadata_list}")
+ logging.debug(f"hitting: {url_with_params}")
+ logging.debug(f"data: {metadata_list}")
response = requests.post(
url_with_params, json=metadata_list, auth=self._auth_provider
)
@@ -388,7 +388,7 @@ Source code for gen3.metadata
[docs]
@backoff.on_exception(backoff.expo, Exception, **BACKOFF_NO_LOG_IF_NOT_RETRIED)
def create(self, guid, metadata, aliases=None, overwrite=False, **kwargs):
- """
+ """
Create the metadata associated with the guid
Args:
@@ -396,14 +396,14 @@ Source code for gen3.metadata
metadata (Dict): dictionary representing what will end up a JSON blob
attached to the provided GUID as metadata
overwrite (bool, optional): whether or not to overwrite existing data
- """
+ """
aliases = aliases or []
- url = self.admin_endpoint + f"/metadata/{guid}"
+ url = self.admin_endpoint + f"/metadata/{guid}"
url_with_params = append_query_params(url, overwrite=overwrite, **kwargs)
- logging.debug(f"hitting: {url_with_params}")
- logging.debug(f"data: {metadata}")
+ logging.debug(f"hitting: {url_with_params}")
+ logging.debug(f"data: {metadata}")
response = requests.post(
url_with_params, json=metadata, auth=self._auth_provider
)
@@ -414,10 +414,10 @@ Source code for gen3.metadata
self.create_aliases(guid=guid, aliases=aliases, merge=overwrite)
except Exception:
logging.error(
- "Error while attempting to create aliases: "
- f"'{aliases}' to GUID: '{guid}' with merge={overwrite}. "
- "GUID metadata record was created successfully and "
- "will NOT be deleted."
+ "Error while attempting to create aliases: "
+ f"'{aliases}' to GUID: '{guid}' with merge={overwrite}. "
+ "GUID metadata record was created successfully and "
+ "will NOT be deleted."
)
return response.json()
@@ -435,7 +435,7 @@ Source code for gen3.metadata
_ssl=None,
**kwargs,
):
- """
+ """
Asynchronous function to create metadata
Args:
@@ -444,19 +444,19 @@ Source code for gen3.metadata
attached to the provided GUID as metadata
overwrite (bool, optional): whether or not to overwrite existing data
_ssl (None, optional): whether or not to use ssl
- """
+ """
aliases = aliases or []
async with aiohttp.ClientSession() as session:
- url = self.admin_endpoint + f"/metadata/{guid}"
+ url = self.admin_endpoint + f"/metadata/{guid}"
url_with_params = append_query_params(url, overwrite=overwrite, **kwargs)
# aiohttp only allows basic auth with their built in auth, so we
# need to manually add JWT auth header
- headers = {"Authorization": self._auth_provider._get_auth_value()}
+ headers = {"Authorization": self._auth_provider._get_auth_value()}
- logging.debug(f"hitting: {url_with_params}")
- logging.debug(f"data: {metadata}")
+ logging.debug(f"hitting: {url_with_params}")
+ logging.debug(f"data: {metadata}")
async with session.post(
url_with_params, json=metadata, headers=headers, ssl=_ssl
) as response:
@@ -464,17 +464,17 @@ Source code for gen3.metadata
response = await response.json()
if aliases:
- logging.info(f"creating aliases: {aliases}")
+ logging.info(f"creating aliases: {aliases}")
try:
await self.async_create_aliases(
guid=guid, aliases=aliases, _ssl=_ssl
)
except Exception:
logging.error(
- "Error while attempting to create aliases: "
- f"'{aliases}' to GUID: '{guid}'. "
- "GUID metadata record was created successfully and "
- "will NOT be deleted."
+ "Error while attempting to create aliases: "
+ f"'{aliases}' to GUID: '{guid}'. "
+ "GUID metadata record was created successfully and "
+ "will NOT be deleted."
)
return response
@@ -484,21 +484,21 @@ Source code for gen3.metadata
[docs]
@backoff.on_exception(backoff.expo, Exception, **DEFAULT_BACKOFF_SETTINGS)
def update(self, guid, metadata, aliases=None, merge=False, **kwargs):
- """
+ """
Update the metadata associated with the guid
Args:
guid (str): guid to use
metadata (Dict): dictionary representing what will end up a JSON blob
attached to the provided GUID as metadata
- """
+ """
aliases = aliases or []
- url = self.admin_endpoint + f"/metadata/{guid}"
+ url = self.admin_endpoint + f"/metadata/{guid}"
url_with_params = append_query_params(url, **kwargs)
- logging.debug(f"hitting: {url_with_params}")
- logging.debug(f"data: {metadata}")
+ logging.debug(f"hitting: {url_with_params}")
+ logging.debug(f"data: {metadata}")
response = requests.put(
url_with_params, json=metadata, auth=self._auth_provider
)
@@ -509,10 +509,10 @@ Source code for gen3.metadata
self.update_aliases(guid=guid, aliases=aliases, merge=merge)
except Exception:
logging.error(
- "Error while attempting to update aliases: "
- f"'{aliases}' to GUID: '{guid}'. "
- "GUID metadata record was created successfully and "
- "will NOT be deleted."
+ "Error while attempting to update aliases: "
+ f"'{aliases}' to GUID: '{guid}'. "
+ "GUID metadata record was created successfully and "
+ "will NOT be deleted."
)
return response.json()
@@ -524,7 +524,7 @@ Source code for gen3.metadata
async def async_update(
self, guid, metadata, aliases=None, merge=False, _ssl=None, **kwargs
):
- """
+ """
Asynchronous function to update metadata
Args:
@@ -536,16 +536,16 @@ Source code for gen3.metadata
with existing values
_ssl (None, optional): whether or not to use ssl
**kwargs: Description
- """
+ """
aliases = aliases or []
async with aiohttp.ClientSession() as session:
- url = self.admin_endpoint + f"/metadata/{guid}"
+ url = self.admin_endpoint + f"/metadata/{guid}"
url_with_params = append_query_params(url, merge=merge, **kwargs)
# aiohttp only allows basic auth with their built in auth, so we
# need to manually add JWT auth header
- headers = {"Authorization": self._auth_provider._get_auth_value()}
+ headers = {"Authorization": self._auth_provider._get_auth_value()}
async with session.put(
url_with_params, json=metadata, headers=headers, ssl=_ssl
@@ -560,10 +560,10 @@ Source code for gen3.metadata
)
except Exception:
logging.error(
- "Error while attempting to update aliases: "
- f"'{aliases}' to GUID: '{guid}' with merge={merge}. "
- "GUID metadata record was created successfully and "
- "will NOT be deleted."
+ "Error while attempting to update aliases: "
+ f"'{aliases}' to GUID: '{guid}' with merge={merge}. "
+ "GUID metadata record was created successfully and "
+ "will NOT be deleted."
)
return response
@@ -573,16 +573,16 @@ Source code for gen3.metadata
[docs]
@backoff.on_exception(backoff.expo, Exception, **DEFAULT_BACKOFF_SETTINGS)
def delete(self, guid, **kwargs):
- """
+ """
Delete the metadata associated with the guid
Args:
guid (str): guid to use
- """
- url = self.admin_endpoint + f"/metadata/{guid}"
+ """
+ url = self.admin_endpoint + f"/metadata/{guid}"
url_with_params = append_query_params(url, **kwargs)
- logging.debug(f"hitting: {url_with_params}")
+ logging.debug(f"hitting: {url_with_params}")
response = requests.delete(url_with_params, auth=self._auth_provider)
response.raise_for_status()
@@ -597,7 +597,7 @@ Source code for gen3.metadata
[docs]
@backoff.on_exception(backoff.expo, Exception, **BACKOFF_NO_LOG_IF_NOT_RETRIED)
def get_aliases(self, guid, **kwargs):
- """
+ """
Get Aliases for the given guid
Args:
@@ -606,11 +606,11 @@ Source code for gen3.metadata
Returns:
requests.Response: response from the request to get aliases
- """
- url = self.endpoint + f"/metadata/{guid}/aliases"
+ """
+ url = self.endpoint + f"/metadata/{guid}/aliases"
url_with_params = append_query_params(url, **kwargs)
- logging.debug(f"hitting: {url_with_params}")
+ logging.debug(f"hitting: {url_with_params}")
response = requests.get(url_with_params, auth=self._auth_provider)
response.raise_for_status()
@@ -621,7 +621,7 @@ Source code for gen3.metadata
[docs]
@backoff.on_exception(backoff.expo, Exception, **BACKOFF_NO_LOG_IF_NOT_RETRIED)
async def async_get_aliases(self, guid, _ssl=None, **kwargs):
- """
+ """
Asyncronously get Aliases for the given guid
Args:
@@ -631,16 +631,16 @@ Source code for gen3.metadata
Returns:
requests.Response: response from the request to get aliases
- """
+ """
async with aiohttp.ClientSession() as session:
- url = self.endpoint + f"/metadata/{guid}/aliases"
+ url = self.endpoint + f"/metadata/{guid}/aliases"
url_with_params = append_query_params(url, **kwargs)
# aiohttp only allows basic auth with their built in auth, so we
# need to manually add JWT auth header
- headers = {"Authorization": self._auth_provider._get_auth_value()}
+ headers = {"Authorization": self._auth_provider._get_auth_value()}
- logging.debug(f"hitting: {url_with_params}")
+ logging.debug(f"hitting: {url_with_params}")
async with session.get(
url_with_params, headers=headers, ssl=_ssl
) as response:
@@ -651,7 +651,7 @@ Source code for gen3.metadata
@backoff.on_exception(backoff.expo, Exception, **BACKOFF_NO_LOG_IF_NOT_RETRIED)
def delete_alias(self, guid, alias, **kwargs):
- """
+ """
Delete single Alias for the given guid
Args:
@@ -660,11 +660,11 @@ Source code for gen3.metadata
Returns:
requests.Response: response from the request to delete aliases
- """
- url = self.admin_endpoint + f"/metadata/{guid}/aliases/{alias}"
+ """
+ url = self.admin_endpoint + f"/metadata/{guid}/aliases/{alias}"
url_with_params = append_query_params(url, **kwargs)
- logging.debug(f"hitting: {url_with_params}")
+ logging.debug(f"hitting: {url_with_params}")
response = requests.delete(url_with_params, auth=self._auth_provider)
response.raise_for_status()
@@ -672,7 +672,7 @@ Source code for gen3.metadata
@backoff.on_exception(backoff.expo, Exception, **BACKOFF_NO_LOG_IF_NOT_RETRIED)
async def async_delete_alias(self, guid, alias, _ssl=None, **kwargs):
- """
+ """
Asyncronously delete single Aliases for the given guid
Args:
@@ -682,16 +682,16 @@ Source code for gen3.metadata
Returns:
requests.Response: response from the request to delete aliases
- """
+ """
async with aiohttp.ClientSession() as session:
- url = self.admin_endpoint + f"/metadata/{guid}/aliases/{alias}"
+ url = self.admin_endpoint + f"/metadata/{guid}/aliases/{alias}"
url_with_params = append_query_params(url, **kwargs)
# aiohttp only allows basic auth with their built in auth, so we
# need to manually add JWT auth header
- headers = {"Authorization": self._auth_provider._get_auth_value()}
+ headers = {"Authorization": self._auth_provider._get_auth_value()}
- logging.debug(f"hitting: {url_with_params}")
+ logging.debug(f"hitting: {url_with_params}")
async with session.delete(
url_with_params, headers=headers, ssl=_ssl
) as response:
@@ -703,7 +703,7 @@ Source code for gen3.metadata
[docs]
@backoff.on_exception(backoff.expo, Exception, **BACKOFF_NO_LOG_IF_NOT_RETRIED)
def create_aliases(self, guid, aliases, **kwargs):
- """
+ """
Create Aliases for the given guid
Args:
@@ -713,14 +713,14 @@ Source code for gen3.metadata
Returns:
requests.Response: response from the request to create aliases
- """
- url = self.admin_endpoint + f"/metadata/{guid}/aliases"
+ """
+ url = self.admin_endpoint + f"/metadata/{guid}/aliases"
url_with_params = append_query_params(url, **kwargs)
- data = {"aliases": aliases}
+ data = {"aliases": aliases}
- logging.debug(f"hitting: {url_with_params}")
- logging.debug(f"data: {data}")
+ logging.debug(f"hitting: {url_with_params}")
+ logging.debug(f"data: {data}")
response = requests.post(url_with_params, json=data, auth=self._auth_provider)
response.raise_for_status()
@@ -731,7 +731,7 @@ Source code for gen3.metadata
[docs]
@backoff.on_exception(backoff.expo, Exception, **BACKOFF_NO_LOG_IF_NOT_RETRIED)
async def async_create_aliases(self, guid, aliases, _ssl=None, **kwargs):
- """
+ """
Asyncronously create Aliases for the given guid
Args:
@@ -742,19 +742,19 @@ Source code for gen3.metadata
Returns:
requests.Response: response from the request to create aliases
- """
+ """
async with aiohttp.ClientSession() as session:
- url = self.admin_endpoint + f"/metadata/{guid}/aliases"
+ url = self.admin_endpoint + f"/metadata/{guid}/aliases"
url_with_params = append_query_params(url, **kwargs)
# aiohttp only allows basic auth with their built in auth, so we
# need to manually add JWT auth header
- headers = {"Authorization": self._auth_provider._get_auth_value()}
+ headers = {"Authorization": self._auth_provider._get_auth_value()}
- data = {"aliases": aliases}
+ data = {"aliases": aliases}
- logging.debug(f"hitting: {url_with_params}")
- logging.debug(f"data: {data}")
+ logging.debug(f"hitting: {url_with_params}")
+ logging.debug(f"data: {data}")
async with session.post(
url_with_params, json=data, headers=headers, ssl=_ssl
) as response:
@@ -766,7 +766,7 @@ Source code for gen3.metadata
[docs]
@backoff.on_exception(backoff.expo, Exception, **BACKOFF_NO_LOG_IF_NOT_RETRIED)
def update_aliases(self, guid, aliases, merge=False, **kwargs):
- """
+ """
Update Aliases for the given guid
Args:
@@ -777,14 +777,14 @@ Source code for gen3.metadata
Returns:
requests.Response: response from the request to update aliases
- """
- url = self.admin_endpoint + f"/metadata/{guid}/aliases"
+ """
+ url = self.admin_endpoint + f"/metadata/{guid}/aliases"
url_with_params = append_query_params(url, merge=merge, **kwargs)
- data = {"aliases": aliases}
+ data = {"aliases": aliases}
- logging.debug(f"hitting: {url_with_params}")
- logging.debug(f"data: {data}")
+ logging.debug(f"hitting: {url_with_params}")
+ logging.debug(f"data: {data}")
response = requests.put(url_with_params, json=data, auth=self._auth_provider)
response.raise_for_status()
@@ -797,7 +797,7 @@ Source code for gen3.metadata
async def async_update_aliases(
self, guid, aliases, merge=False, _ssl=None, **kwargs
):
- """
+ """
Asyncronously update Aliases for the given guid
Args:
@@ -809,19 +809,19 @@ Source code for gen3.metadata
Returns:
requests.Response: response from the request to update aliases
- """
+ """
async with aiohttp.ClientSession() as session:
- url = self.admin_endpoint + f"/metadata/{guid}/aliases"
+ url = self.admin_endpoint + f"/metadata/{guid}/aliases"
url_with_params = append_query_params(url, merge=merge, **kwargs)
# aiohttp only allows basic auth with their built in auth, so we
# need to manually add JWT auth header
- headers = {"Authorization": self._auth_provider._get_auth_value()}
+ headers = {"Authorization": self._auth_provider._get_auth_value()}
- data = {"aliases": aliases}
+ data = {"aliases": aliases}
- logging.debug(f"hitting: {url_with_params}")
- logging.debug(f"data: {data}")
+ logging.debug(f"hitting: {url_with_params}")
+ logging.debug(f"data: {data}")
async with session.put(
url_with_params, json=data, headers=headers, ssl=_ssl
) as response:
@@ -834,7 +834,7 @@ Source code for gen3.metadata
[docs]
@backoff.on_exception(backoff.expo, Exception, **BACKOFF_NO_LOG_IF_NOT_RETRIED)
def delete_aliases(self, guid, **kwargs):
- """
+ """
Delete all Aliases for the given guid
Args:
@@ -843,11 +843,11 @@ Source code for gen3.metadata
Returns:
requests.Response: response from the request to delete aliases
- """
- url = self.admin_endpoint + f"/metadata/{guid}/aliases"
+ """
+ url = self.admin_endpoint + f"/metadata/{guid}/aliases"
url_with_params = append_query_params(url, **kwargs)
- logging.debug(f"hitting: {url_with_params}")
+ logging.debug(f"hitting: {url_with_params}")
response = requests.delete(url_with_params, auth=self._auth_provider)
response.raise_for_status()
@@ -858,7 +858,7 @@ Source code for gen3.metadata
[docs]
@backoff.on_exception(backoff.expo, Exception, **BACKOFF_NO_LOG_IF_NOT_RETRIED)
async def async_delete_aliases(self, guid, _ssl=None, **kwargs):
- """
+ """
Asyncronously delete all Aliases for the given guid
Args:
@@ -868,16 +868,16 @@ Source code for gen3.metadata
Returns:
requests.Response: response from the request to delete aliases
- """
+ """
async with aiohttp.ClientSession() as session:
- url = self.admin_endpoint + f"/metadata/{guid}/aliases"
+ url = self.admin_endpoint + f"/metadata/{guid}/aliases"
url_with_params = append_query_params(url, **kwargs)
# aiohttp only allows basic auth with their built in auth, so we
# need to manually add JWT auth header
- headers = {"Authorization": self._auth_provider._get_auth_value()}
+ headers = {"Authorization": self._auth_provider._get_auth_value()}
- logging.debug(f"hitting: {url_with_params}")
+ logging.debug(f"hitting: {url_with_params}")
async with session.delete(
url_with_params, headers=headers, ssl=_ssl
) as response:
@@ -890,7 +890,7 @@ Source code for gen3.metadata
[docs]
@backoff.on_exception(backoff.expo, Exception, **BACKOFF_NO_LOG_IF_NOT_RETRIED)
def delete_alias(self, guid, alias, **kwargs):
- """
+ """
Delete single Alias for the given guid
Args:
@@ -900,11 +900,11 @@ Source code for gen3.metadata
Returns:
requests.Response: response from the request to delete aliases
- """
- url = self.admin_endpoint + f"/metadata/{guid}/aliases/{alias}"
+ """
+ url = self.admin_endpoint + f"/metadata/{guid}/aliases/{alias}"
url_with_params = append_query_params(url, **kwargs)
- logging.debug(f"hitting: {url_with_params}")
+ logging.debug(f"hitting: {url_with_params}")
response = requests.delete(url_with_params, auth=self._auth_provider)
response.raise_for_status()
@@ -915,7 +915,7 @@ Source code for gen3.metadata
[docs]
@backoff.on_exception(backoff.expo, Exception, **BACKOFF_NO_LOG_IF_NOT_RETRIED)
async def async_delete_alias(self, guid, alias, _ssl=None, **kwargs):
- """
+ """
Asyncronously delete single Aliases for the given guid
Args:
@@ -926,16 +926,16 @@ Source code for gen3.metadata
Returns:
requests.Response: response from the request to delete aliases
- """
+ """
async with aiohttp.ClientSession() as session:
- url = self.admin_endpoint + f"/metadata/{guid}/aliases/{alias}"
+ url = self.admin_endpoint + f"/metadata/{guid}/aliases/{alias}"
url_with_params = append_query_params(url, **kwargs)
# aiohttp only allows basic auth with their built in auth, so we
# need to manually add JWT auth header
- headers = {"Authorization": self._auth_provider._get_auth_value()}
+ headers = {"Authorization": self._auth_provider._get_auth_value()}
- logging.debug(f"hitting: {url_with_params}")
+ logging.debug(f"hitting: {url_with_params}")
async with session.delete(
url_with_params, headers=headers, ssl=_ssl
) as response:
@@ -947,11 +947,11 @@ Source code for gen3.metadata
def _prepare_metadata(
self, metadata, indexd_doc, force_metadata_columns_even_if_empty
):
- """
+ """
Validate and generate the provided metadata for submission to the metadata
service.
- If the record is of type "package", also prepare package metadata.
+ If the record is of type "package", also prepare package metadata.
Args:
metadata (dict): metadata provided by the submitter
@@ -960,13 +960,13 @@ Source code for gen3.metadata
Returns:
dict: metadata ready to be submitted to the metadata service
- """
+ """
def _extract_non_indexd_metadata(metadata):
- """
- Get the "additional metadata": metadata that was provided but is
+ """
+ Get the "additional metadata": metadata that was provided but is
not stored in indexd, so should be stored in the metadata service.
- """
+ """
return {
k: v
for k, v in metadata.items()
@@ -987,14 +987,14 @@ Source code for gen3.metadata
valid = True
# validate package columns
- record_type = to_submit.pop(RECORD_TYPE_STANDARD_KEY, "").strip().lower()
+ record_type = to_submit.pop(RECORD_TYPE_STANDARD_KEY, "").strip().lower()
package_contents = to_submit.pop(PACKAGE_CONTENTS_STANDARD_KEY, None)
- if record_type == "package":
+ if record_type == "package":
if package_contents:
package_contents = json.loads(package_contents)
if not _verify_schema(package_contents, PACKAGE_CONTENTS_SCHEMA):
logging.error(
- f"ERROR: {package_contents} is not in package contents format"
+ f"ERROR: {package_contents} is not in package contents format"
)
valid = False
# generate package metadata
@@ -1009,19 +1009,19 @@ Source code for gen3.metadata
to_submit.update(package_metadata)
elif package_contents:
logging.error(
- f"ERROR: tried to set '{PACKAGE_CONTENTS_STANDARD_KEY}' for a non-package row. Ignoring '{PACKAGE_CONTENTS_STANDARD_KEY}'. Set '{RECORD_TYPE_STANDARD_KEY}' to 'package' to create packages."
+ f"ERROR: tried to set '{PACKAGE_CONTENTS_STANDARD_KEY}' for a non-package row. Ignoring '{PACKAGE_CONTENTS_STANDARD_KEY}'. Set '{RECORD_TYPE_STANDARD_KEY}' to 'package' to create packages."
)
valid = False
if not valid:
- raise Exception(f"Metadata is not valid: {metadata}")
+ raise Exception(f"Metadata is not valid: {metadata}")
if not force_metadata_columns_even_if_empty:
- # remove any empty columns if we're not being forced to include them
+ # remove any empty columns if we're not being forced to include them
to_submit = {
key: value
for key, value in to_submit.items()
- if value is not None and value != ""
+ if value is not None and value != ""
}
return to_submit
@@ -1029,18 +1029,18 @@ Source code for gen3.metadata
def _get_package_metadata(
self, submitted_metadata, file_name, file_size, hashes, urls, contents
):
- """
+ """
The MDS Objects API currently expects files that have not been
uploaded yet. For files we only needs to index, not upload, create
object records manually by generating the expected object fields.
TODO: update the MDS objects API to not create upload URLs if the
relevant data is provided.
- """
+ """
def _get_filename_from_urls(submitted_metadata, urls):
- file_name = ""
+ file_name = ""
if not urls:
- logging.warning(f"No URLs provided for: {submitted_metadata}")
+ logging.warning(f"No URLs provided for: {submitted_metadata}")
for url in urls:
_file_name = os.path.basename(url)
if not file_name:
@@ -1048,7 +1048,7 @@ Source code for gen3.metadata
else:
if file_name != _file_name:
logging.warning(
- f"Received multiple URLs with different file names; will use the first URL (file name '{file_name}'): {submitted_metadata}"
+ f"Received multiple URLs with different file names; will use the first URL (file name '{file_name}'): {submitted_metadata}"
)
return file_name
@@ -1058,17 +1058,17 @@ Source code for gen3.metadata
now = str(datetime.utcnow())
metadata = {
- "type": "package",
- "package": {
- "version": "0.1",
- "file_name": file_name,
- "created_time": now,
- "updated_time": now,
- "size": file_size,
- "hashes": hashes,
- "contents": contents or None,
+ "type": "package",
+ "package": {
+ "version": "0.1",
+ "file_name": file_name,
+ "created_time": now,
+ "updated_time": now,
+ "size": file_size,
+ "hashes": hashes,
+ "contents": contents or None,
},
- "_upload_status": "uploaded",
+ "_upload_status": "uploaded",
}
return metadata
diff --git a/docs/_build/html/_modules/gen3/object.html b/docs/_build/html/_modules/gen3/object.html
index 19cc6697d..ac26f96c0 100644
--- a/docs/_build/html/_modules/gen3/object.html
+++ b/docs/_build/html/_modules/gen3/object.html
@@ -42,7 +42,7 @@ Source code for gen3.object
[docs]
class Gen3Object:
-
"""For interacting with Gen3 object level features.
+
"""For interacting with Gen3 object level features.
A class for interacting with the Gen3 object services.
Currently allows creating and deleting of an object from the Gen3 System.
@@ -54,36 +54,36 @@
Source code for gen3.object
This generates the Gen3Object class pointed at the sandbox commons while
using the credentials.json downloaded from the commons profile page.
- >>> auth = Gen3Auth(refresh_file="credentials.json")
+ >>> auth = Gen3Auth(refresh_file="credentials.json")
... object = Gen3Object(auth)
- """
+ """
def __init__(self, auth_provider=None):
self._auth_provider = auth_provider
- self.service_endpoint = "/mds"
+ self.service_endpoint = "/mds"
def create_object(self, file_name, authz, metadata=None, aliases=None):
url = (
- self._auth_provider.endpoint.rstrip("/")
+ self._auth_provider.endpoint.rstrip("/")
+ self.service_endpoint
- + "/objects"
+ + "/objects"
)
body = {
- "file_name": file_name,
- "authz": authz,
- "metadata": metadata,
- "aliases": aliases,
+ "file_name": file_name,
+ "authz": authz,
+ "metadata": metadata,
+ "aliases": aliases,
}
response = requests.post(url, json=body, auth=self._auth_provider)
raise_for_status_and_print_error(response)
data = response.json()
- return data["guid"], data["upload_url"]
+ return data["guid"], data["upload_url"]
[docs]
def delete_object(self, guid, delete_file_locations=False):
-
"""
+
"""
Delete the object from indexd, metadata service and optionally all storage locations
Args:
@@ -91,12 +91,12 @@
Source code for gen3.object
`delete_file_locations` -- if True, removes the object from existing bucket location(s) through fence
Returns:
Nothing
- """
- delete_param = "?delete_file_locations" if delete_file_locations else ""
+ """
+ delete_param = "?delete_file_locations" if delete_file_locations else ""
url = (
- self._auth_provider.endpoint.rstrip("/")
+ self._auth_provider.endpoint.rstrip("/")
+ self.service_endpoint
- + "/objects/"
+ + "/objects/"
+ guid
+ delete_param
)
diff --git a/docs/_build/html/_modules/gen3/query.html b/docs/_build/html/_modules/gen3/query.html
index af2587093..295750ef2 100644
--- a/docs/_build/html/_modules/gen3/query.html
+++ b/docs/_build/html/_modules/gen3/query.html
@@ -39,7 +39,7 @@ Source code for gen3.query
[docs]
class Gen3Query:
-
"""
+
"""
Query ElasticSearch data from a Gen3 system.
Args:
@@ -49,9 +49,9 @@
Source code for gen3.query
This generates the Gen3Query class pointed at the sandbox commons while
using the credentials.json downloaded from the commons profile page.
- >>> auth = Gen3Auth(endpoint, refresh_file="credentials.json")
+ >>> auth = Gen3Auth(endpoint, refresh_file="credentials.json")
... query = Gen3Query(auth)
- """
+ """
def __init__(self, auth_provider):
self._auth_provider = auth_provider
@@ -70,7 +70,7 @@ Source code for gen3.query
accessibility=None,
verbose=True,
):
- """
+ """
Execute a query against a Data Commons.
Args:
@@ -81,23 +81,23 @@ Source code for gen3.query
filters: (object, optional): { field: sort method } object. Will filter data with ALL fields EQUAL to the provided respective value. If more complex filters are needed, use the `filter_object` parameter instead.
filter_object (object, optional): Filter to apply. For syntax details, see https://github.com/uc-cdis/guppy/blob/master/doc/queries.md#filter.
sort_object (object, optional): { field: sort method } object.
- accessibility (list, optional): One of ["accessible" (default), "unaccessible", "all"]. Only valid when querying a data type in "regular" tier access mode.
+ accessibility (list, optional): One of ["accessible" (default), "unaccessible", "all"]. Only valid when querying a data type in "regular" tier access mode.
Returns:
- Object: {"data": {<data_type>: [<record>, <record>, ...]}}
+ Object: {"data": {<data_type>: [<record>, <record>, ...]}}
Examples:
>>> Gen3Query.query(
- data_type="subject",
+ data_type="subject",
first=50,
fields=[
- "vital_status",
- "submitter_id",
+ "vital_status",
+ "submitter_id",
],
- filters={"vital_status": "Alive"},
- sort_object={"submitter_id": "asc"},
+ filters={"vital_status": "Alive"},
+ sort_object={"submitter_id": "asc"},
)
- """
+ """
if not first:
first = 10
if not offset:
@@ -105,14 +105,14 @@ Source code for gen3.query
if not sort_object:
sort_object = {}
if not accessibility:
- accessibility = "accessible"
+ accessibility = "accessible"
if filters and filter_object:
raise Exception(
- "Only one of `filters` and `filter_object` can be used at a time."
+ "Only one of `filters` and `filter_object` can be used at a time."
)
if filters:
filter_object = {
- "AND": [{"=": {field: val}} for field, val in filters.items()]
+ "AND": [{"=": {field: val}} for field, val in filters.items()]
}
if first + offset > 10000: # ElasticSearch limitation
@@ -126,13 +126,13 @@ Source code for gen3.query
first=first,
offset=offset,
)
- return {"data": {data_type: data}}
+ return {"data": {data_type: data}}
- # convert sort_object to graphql: [ { field_name: "sort_method" } ]
- sorts = [f'{{{field}: "{val}"}}' for field, val in sort_object.items()]
- sort_string = f'[{", ".join(sorts)}]'
+ # convert sort_object to graphql: [ { field_name: "sort_method" } ]
+ sorts = [f'{{{field}: "{val}"}}' for field, val in sort_object.items()]
+ sort_string = f'[{", ".join(sorts)}]'
- query_string = f"""query($filter: JSON) {{
+ query_string = f"""query($filter: JSON) {{
{data_type}(
first: {first},
offset: {offset},
@@ -140,17 +140,17 @@ Source code for gen3.query
accessibility: {accessibility},
filter: $filter
) {{
- {" ".join(fields)}
+ {" ".join(fields)}
}}
- }}"""
- variables = {"filter": filter_object}
+ }}"""
+ variables = {"filter": filter_object}
return self.graphql_query(query_string=query_string, variables=variables)
[docs]
def graphql_query(self, query_string, variables=None):
-
"""
+
"""
Execute a GraphQL query against a Data Commons.
Args:
@@ -158,29 +158,29 @@
Source code for gen3.query
variables (:obj:`object`, optional): Dictionary of variables to pass with the query.
Returns:
- Object: {"data": {<data_type>: [<record>, <record>, ...]}}
+ Object: {"data": {<data_type>: [<record>, <record>, ...]}}
Examples:
- >>> query_string = "{ my_index { my_field } }"
+ >>> query_string = "{ my_index { my_field } }"
... Gen3Query.graphql_query(query_string)
- """
- url = f"{self._auth_provider.endpoint}/guppy/graphql"
+ """
+ url = f"{self._auth_provider.endpoint}/guppy/graphql"
response = requests.post(
url,
- json={"query": query_string, "variables": variables},
+ json={"query": query_string, "variables": variables},
auth=self._auth_provider,
)
try:
raise_for_status_and_print_error(response)
except Exception:
print(
- f"Unable to query.\nQuery: {query_string}\nVariables: {variables}\n{response.text}"
+ f"Unable to query.\nQuery: {query_string}\nVariables: {variables}\n{response.text}"
)
raise
try:
return response.json()
except Exception:
- print(f"Did not receive JSON: {response.text}")
+ print(f"Did not receive JSON: {response.text}")
raise
@@ -196,7 +196,7 @@
Source code for gen3.query
first=None,
offset=None,
):
- """
+ """
Execute a raw data download against a Data Commons.
Args:
@@ -204,7 +204,7 @@ Source code for gen3.query
fields (list): List of fields to return.
filter_object (object, optional): Filter to apply. For syntax details, see https://github.com/uc-cdis/guppy/blob/master/doc/queries.md#filter.
sort_fields (list, optional): List of { field: sort method } objects.
- accessibility (list, optional): One of ["accessible" (default), "unaccessible", "all"]. Only valid when downloading from a data type in "regular" tier access mode.
+ accessibility (list, optional): One of ["accessible" (default), "unaccessible", "all"]. Only valid when downloading from a data type in "regular" tier access mode.
first (int, optional): Number of rows to return (default: all rows).
offset (int, optional): Starting position (default: 0).
@@ -213,29 +213,29 @@ Source code for gen3.query
Examples:
>>> Gen3Query.raw_data_download(
- data_type="subject",
+ data_type="subject",
fields=[
- "vital_status",
- "submitter_id",
- "project_id"
+ "vital_status",
+ "submitter_id",
+ "project_id"
],
- filter_object={"=": {"project_id": "my_program-my_project"}},
- sort_fields=[{"submitter_id": "asc"}],
- accessibility="accessible"
+ filter_object={"=": {"project_id": "my_program-my_project"}},
+ sort_fields=[{"submitter_id": "asc"}],
+ accessibility="accessible"
)
- """
+ """
if not accessibility:
- accessibility = "accessible"
+ accessibility = "accessible"
if not offset:
offset = 0
- body = {"type": data_type, "fields": fields, "accessibility": accessibility}
+ body = {"type": data_type, "fields": fields, "accessibility": accessibility}
if filter_object:
- body["filter"] = filter_object
+ body["filter"] = filter_object
if sort_fields:
- body["sort"] = sort_fields
+ body["sort"] = sort_fields
- url = f"{self._auth_provider.endpoint}/guppy/download"
+ url = f"{self._auth_provider.endpoint}/guppy/download"
response = requests.post(
url,
json=body,
@@ -244,12 +244,12 @@ Source code for gen3.query
try:
raise_for_status_and_print_error(response)
except Exception:
- print(f"Unable to download.\nBody: {body}\n{response.text}")
+ print(f"Unable to download.\nBody: {body}\n{response.text}")
raise
try:
data = response.json()
except Exception:
- print(f"Did not receive JSON: {response.text}")
+ print(f"Did not receive JSON: {response.text}")
raise
if offset:
diff --git a/docs/_build/html/_modules/gen3/submission.html b/docs/_build/html/_modules/gen3/submission.html
index 04223a5fb..b82fd483e 100644
--- a/docs/_build/html/_modules/gen3/submission.html
+++ b/docs/_build/html/_modules/gen3/submission.html
@@ -58,7 +58,7 @@ Source code for gen3.submission
[docs]
class Gen3Submission:
-
"""Submit/Export/Query data from a Gen3 Submission system.
+
"""Submit/Export/Query data from a Gen3 Submission system.
A class for interacting with the Gen3 submission services.
Supports submitting and exporting from Sheepdog.
@@ -71,10 +71,10 @@
Source code for gen3.submission
This generates the Gen3Submission class pointed at the sandbox commons while
using the credentials.json downloaded from the commons profile page.
- >>> auth = Gen3Auth(refresh_file="credentials.json")
+ >>> auth = Gen3Auth(refresh_file="credentials.json")
... sub = Gen3Submission(auth)
- """
+ """
def __init__(self, endpoint=None, auth_provider=None):
# auth_provider legacy interface required endpoint as 1st arg
@@ -82,18 +82,18 @@ Source code for gen3.submission
self._endpoint = self._auth_provider.endpoint
def __export_file(self, filename, output):
- """Writes an API response to a file."""
- with open(filename, "w") as outfile:
+ """Writes an API response to a file."""
+ with open(filename, "w") as outfile:
outfile.write(output)
- print("\nOutput written to file: " + filename)
+ print("\nOutput written to file: " + filename)
### Program functions