You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
 
 

312 lines
13 KiB

  1. #!/usr/bin/env python
  2. import base64
  3. import datetime
  4. import json
  5. import logging
  6. import os
  7. import pathlib
  8. import shutil
  9. import time
  10. from typing import Optional
  11. import urllib.parse
  12. import re
  13. import certifi
  14. import click
  15. import minio
  16. import requests
  17. import urllib3
  18. from progress import Progress
  19. logging.basicConfig(level=logging.INFO)
  20. BACKFEED_DELIM = "\n"
  21. SKIP_BACKFEED = "SKIPBF"
  22. SKIP_DISPATCHER = "SKIPDISPATCH://"
  23. # TODO: Add rsync support
  24. # TODO: Add rsync+ssh support
  25. # TODO: Add webdav support.
  26. # TODO: Fix the "ctrl-c handling" logic so it actually cleans up in the s3 bucket.
  27. def retry_failures(fn, msg, *args, **kwargs):
  28. tries = 0
  29. while True:
  30. try:
  31. return fn(*args, **kwargs)
  32. except Exception:
  33. logging.exception(msg)
  34. delay = min(2 ** tries, 64)
  35. tries = tries + 1
  36. logging.info(f"Sleeping {delay} seconds...")
  37. time.sleep(delay)
  38. @click.group()
  39. def sender():
  40. pass
  41. def watch_pass(input_directory: pathlib.Path, work_directory: pathlib.Path, ia_collection: str, ia_item_title: str,
  42. ia_item_prefix: str, ia_item_date: str, project: str, dispatcher: str, delete: bool, backfeed_key: str):
  43. logging.info("Checking for new items...")
  44. for original_directory in input_directory.iterdir():
  45. if original_directory.is_dir():
  46. original_name = original_directory.name
  47. new_directory = work_directory.joinpath(original_name)
  48. try:
  49. original_directory.rename(new_directory)
  50. except FileNotFoundError:
  51. logging.warning(f"Unable to move item {original_directory}")
  52. continue
  53. single_impl(new_directory, ia_collection, ia_item_title, ia_item_prefix, ia_item_date, project,
  54. dispatcher, delete, backfeed_key)
  55. return True
  56. return False
  57. @sender.command()
  58. @click.option('--input-directory', envvar='UPLOAD_QUEUE_DIR', default="/data/upload-queue",
  59. type=click.Path(exists=True))
  60. @click.option('--work-directory', envvar='UPLOADER_WORKING_DIR', default="/data/uploader-work",
  61. type=click.Path(exists=True))
  62. @click.option('--ia-collection', envvar='IA_COLLECTION', required=True)
  63. @click.option('--ia-item-title', envvar='IA_ITEM_TITLE', required=True)
  64. @click.option('--ia-item-prefix', envvar='IA_ITEM_PREFIX', required=True)
  65. @click.option('--ia-item-date', envvar='IA_ITEM_DATE', required=False)
  66. @click.option('--project', envvar='PROJECT', required=True)
  67. @click.option('--dispatcher', envvar='DISPATCHER', required=True)
  68. @click.option('--delete/--no-delete', envvar='DELETE', default=False)
  69. @click.option('--backfeed-key', envvar='BACKFEED_KEY', required=True)
  70. def watch(input_directory: pathlib.Path, work_directory: pathlib.Path, ia_collection: str, ia_item_title: str,
  71. ia_item_prefix: str, ia_item_date: str, project: str, dispatcher: str, delete: bool, backfeed_key: str):
  72. if not isinstance(input_directory, pathlib.Path):
  73. input_directory = pathlib.Path(input_directory)
  74. if not isinstance(work_directory, pathlib.Path):
  75. work_directory = pathlib.Path(work_directory)
  76. while True:
  77. if not watch_pass(input_directory, work_directory, ia_collection, ia_item_title, ia_item_prefix, ia_item_date,
  78. project, dispatcher, delete, backfeed_key):
  79. logging.info("No item found, sleeping...")
  80. time.sleep(10)
  81. @sender.command()
  82. @click.option('--item-directory', type=click.Path(exists=True), required=True)
  83. @click.option('--ia-collection', envvar='IA_COLLECTION', required=True)
  84. @click.option('--ia-item-title', envvar='IA_ITEM_TITLE', required=True)
  85. @click.option('--ia-item-prefix', envvar='IA_ITEM_PREFIX', required=True)
  86. @click.option('--ia-item-date', envvar='IA_ITEM_DATE', required=False)
  87. @click.option('--project', envvar='PROJECT', required=True)
  88. @click.option('--dispatcher', envvar='DISPATCHER', required=True)
  89. @click.option('--delete/--no-delete', envvar='DELETE', default=False)
  90. @click.option('--backfeed-key', envvar='BACKFEED_KEY', required=True)
  91. def single(item_directory: pathlib.Path, ia_collection: str, ia_item_title: str, ia_item_prefix: str,
  92. ia_item_date: Optional[str], project: str, dispatcher: str, delete: bool, backfeed_key: str):
  93. single_impl(item_directory, ia_collection, ia_item_title, ia_item_prefix, ia_item_date, project, dispatcher, delete,
  94. backfeed_key)
  95. def single_impl(item_directory: pathlib.Path, ia_collection: str, ia_item_title: str, ia_item_prefix: str,
  96. ia_item_date: Optional[str], project: str, dispatcher: str, delete: bool, backfeed_key: str):
  97. if not isinstance(item_directory, pathlib.Path):
  98. item_directory = pathlib.Path(item_directory)
  99. logging.info(f"Processing item {item_directory}...")
  100. if ia_item_date is None:
  101. s = item_directory.name.split("_")
  102. if len(s) > 0:
  103. ds = s[0]
  104. try:
  105. d = datetime.datetime.strptime(ds, "%Y%m%d%H%M%S")
  106. ia_item_date = d.strftime("%Y-%m")
  107. except ValueError:
  108. pass
  109. meta_json_loc = item_directory.joinpath('__upload_meta.json')
  110. if meta_json_loc.exists():
  111. raise Exception("META JSON EXISTS WTF")
  112. meta_json = {
  113. "IA_COLLECTION": ia_collection,
  114. "IA_ITEM_TITLE": f"{ia_item_title} {item_directory.name}",
  115. "IA_ITEM_DATE": ia_item_date,
  116. "IA_ITEM_NAME": f"{ia_item_prefix}{item_directory.name}",
  117. "PROJECT": project,
  118. }
  119. with open(meta_json_loc, 'w') as f:
  120. f.write(json.dumps(meta_json))
  121. logging.info("Wrote metadata json.")
  122. total_size = 0
  123. files = list(item_directory.glob("**/*"))
  124. for item in files:
  125. total_size = total_size + os.path.getsize(item)
  126. logging.info(f"Item size is {total_size} bytes across {len(files)} files.")
  127. meta_json["SIZE_HINT"] = str(total_size)
  128. def assign_target():
  129. if dispatcher.startswith(SKIP_DISPATCHER):
  130. target = dispatcher[len(SKIP_DISPATCHER):]
  131. logging.info(f"Skipping dispatcher, using {target} as target.")
  132. return target
  133. logging.info("Attempting to assign target...")
  134. r = requests.get(f"{dispatcher}/offload_target", params=meta_json, timeout=60)
  135. if r.status_code == 200:
  136. data = r.json()
  137. return data["url"]
  138. else:
  139. raise Exception(f"Invalid status code {r.status_code}: {r.text}")
  140. url = retry_failures(assign_target, "Failed to fetch target")
  141. logging.info(f"Assigned target {url}")
  142. parsed_url = urllib.parse.urlparse(url)
  143. parsed_qs = urllib.parse.parse_qs(parsed_url.query)
  144. def get_q(key, default):
  145. return parsed_qs.get(key, [str(default)])[0]
  146. bf_item = None
  147. if parsed_url.scheme == "minio+http" or parsed_url.scheme == "minio+https":
  148. secure = (parsed_url.scheme == "minio+https")
  149. ep = parsed_url.hostname
  150. if parsed_url.port is not None:
  151. ep = f"{ep}:{parsed_url.port}"
  152. client = None
  153. def create_client():
  154. logging.info("Connecting to minio...")
  155. cert_check = True
  156. timeout = datetime.timedelta(seconds=int(get_q("timeout", 60))).seconds
  157. total_timeout = datetime.timedelta(seconds=int(get_q("total_timeout", timeout*2))).seconds
  158. hclient = urllib3.PoolManager(
  159. timeout=urllib3.util.Timeout(connect=timeout, read=timeout, total=total_timeout),
  160. maxsize=10,
  161. cert_reqs='CERT_REQUIRED' if cert_check else 'CERT_NONE',
  162. ca_certs=os.environ.get('SSL_CERT_FILE') or certifi.where(),
  163. retries=urllib3.Retry(
  164. total=5,
  165. backoff_factor=0.2,
  166. backoff_max=30,
  167. status_forcelist=[500, 502, 503, 504]
  168. )
  169. )
  170. return minio.Minio(endpoint=ep, access_key=parsed_url.username, secret_key=parsed_url.password,
  171. secure=secure, http_client=hclient)
  172. client = retry_failures(create_client, "Failed to connect to minio")
  173. bucket_name = item_directory.name.replace("_", "-")
  174. def make_bucket():
  175. logging.info("Attempting to make bucket...")
  176. if client.bucket_exists(bucket_name=bucket_name):
  177. # If bucket already exists a previous attempt was aborted.
  178. logging.warning("Bucket already exists!")
  179. return
  180. client.make_bucket(bucket_name=bucket_name)
  181. retry_failures(make_bucket, "Failed to make bucket")
  182. logging.info("Starting uploads...")
  183. for file in files:
  184. rel_file = file.relative_to(item_directory)
  185. def upload_file():
  186. logging.info(f"Uploading file {rel_file}...")
  187. client.fput_object(bucket_name=bucket_name, object_name=str(rel_file), file_path=file, num_parallel_uploads=8)
  188. retry_failures(upload_file, f"Failed to upload {rel_file}")
  189. item_data = {"url": url, "item_name": item_directory.name, "bucket_name": bucket_name}
  190. bf_item_part = base64.urlsafe_b64encode(str(json.dumps(item_data)).encode("UTF-8")).decode("UTF-8")
  191. bf_item = f"{project}:{parsed_url.hostname}:{bf_item_part}"
  192. elif parsed_url.scheme == "s3+cred":
  193. secure = True
  194. ep = parsed_url.hostname
  195. if parsed_url.port is not None:
  196. ep = f"{ep}:{parsed_url.port}"
  197. if not parsed_url.path:
  198. raise Exception("Invalid s3+cred url, no bucket specified. vhost style s3 urls are not supported.")
  199. bucket_name = parsed_url.path[1:].strip("/")
  200. if "/" in bucket_name:
  201. raise Exception("Invalid s3+cred url, path does not support sub-directories - only a main bucket in the url is supported.")
  202. client = None
  203. # ensure the cred name is reasonable
  204. if not re.match("^[a-zA-Z0-9_.-]+$", parsed_url.username):
  205. raise Exception(f"Invalid cred name in url: {parsed_url.username}")
  206. creds_file = f'creds/{parsed_url.username}.json'
  207. if not os.path.exists(creds_file):
  208. raise Exception(f"Unable to find credentials file {creds_file}")
  209. with open(creds_file) as f:
  210. creds = json.load(f)
  211. if "access_key" not in creds or "secret_key" not in creds:
  212. raise Exception(f"Unable to find access_key or secret_key in credentials file {creds_file}")
  213. def create_client():
  214. logging.info("Connecting to minio...")
  215. cert_check = True
  216. timeout = datetime.timedelta(seconds=int(get_q("timeout", 60))).seconds
  217. total_timeout = datetime.timedelta(seconds=int(get_q("total_timeout", timeout * 2))).seconds
  218. hclient = urllib3.PoolManager(
  219. timeout=urllib3.util.Timeout(connect=timeout, read=timeout, total=total_timeout),
  220. maxsize=10,
  221. cert_reqs='CERT_REQUIRED' if cert_check else 'CERT_NONE',
  222. ca_certs=os.environ.get('SSL_CERT_FILE') or certifi.where(),
  223. retries=urllib3.Retry(
  224. total=5,
  225. backoff_factor=0.2,
  226. backoff_max=30,
  227. status_forcelist=[500, 502, 503, 504]
  228. )
  229. )
  230. return minio.Minio(endpoint=ep, access_key=creds['access_key'], secret_key=creds['secret_key'],
  231. secure=secure, http_client=hclient)
  232. client = retry_failures(create_client, "Failed to connect to minio")
  233. folder_name = item_directory.name.replace("_", "-")
  234. logging.info("Starting uploads...")
  235. for file in files:
  236. rel_file = file.relative_to(item_directory)
  237. def upload_file():
  238. logging.info(f"Uploading file {rel_file}...")
  239. client.fput_object(bucket_name = bucket_name, object_name=f"{folder_name}/{str(rel_file)}", file_path=file, num_parallel_uploads=8)
  240. retry_failures(upload_file, f"Failed to upload {rel_file}")
  241. item_data = {"url": url, "item_name": item_directory.name, "folder_name": folder_name}
  242. bf_item_part = base64.urlsafe_b64encode(str(json.dumps(item_data)).encode("UTF-8")).decode("UTF-8")
  243. bf_item = f"{project}:{parsed_url.hostname}:{bf_item_part}"
  244. else:
  245. raise Exception("Unable to upload, don't understand url: {url}")
  246. if bf_item is None:
  247. raise Exception("Unable to create backfeed item")
  248. if backfeed_key == SKIP_BACKFEED:
  249. logging.warning(f"Skipping backfeed! Would have submitted: {bf_item}")
  250. else:
  251. def submit_item():
  252. u = f"https://legacy-api.arpa.li/backfeed/legacy/{backfeed_key}"
  253. logging.info(f"Attempting to submit bf item {bf_item} to {u}...")
  254. resp = requests.post(u, params={"skipbloom": "1", "delimiter": BACKFEED_DELIM},
  255. data=f"{bf_item}{BACKFEED_DELIM}".encode("UTF-8"), timeout=60)
  256. if resp.status_code != 200:
  257. raise Exception(f"Failed to submit to backfeed {resp.status_code}: {resp.text}")
  258. retry_failures(submit_item, "Failed to submit to backfeed")
  259. logging.info("Backfeed submit complete!")
  260. if delete:
  261. logging.info("Removing item...")
  262. shutil.rmtree(item_directory)
  263. logging.info("Upload complete!")
  264. if __name__ == '__main__':
  265. sender()