test_object_store_service.py 30 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628
  1. import filecmp
  2. import os
  3. import tempfile
  4. from datetime import datetime
  5. from io import BytesIO
  6. from unittest import mock
  7. from unittest import skip
  8. import requests
  9. from cloudbridge.base import helpers as cb_helpers
  10. from cloudbridge.base.resources import BaseBucketObject
  11. from cloudbridge.interfaces.exceptions import DuplicateResourceException
  12. from cloudbridge.interfaces.provider import TestMockHelperMixin
  13. from cloudbridge.interfaces.resources import Bucket
  14. from cloudbridge.interfaces.resources import BucketObject
  15. from cloudbridge.interfaces.resources import TransferConfig
  16. from tests import helpers
  17. from tests.helpers import ProviderTestBase
  18. from tests.helpers import standard_interface_tests as sit
  19. # S3 (and Swift) require every part except the last to be >= 5 MiB. Tests use
  20. # this size so they remain valid against real cloud providers, not just moto.
  21. MIN_PART_SIZE = 5 * 1024 * 1024
  22. class CloudObjectStoreServiceTestCase(ProviderTestBase):
  23. _multiprocess_can_split_ = True
  24. @helpers.skipIfNoService(['storage._bucket_objects', 'storage.buckets'])
  25. def test_storage_services_event_pattern(self):
  26. # pylint:disable=protected-access
  27. self.assertEqual(
  28. self.provider.storage.buckets._service_event_pattern,
  29. "provider.storage.buckets",
  30. "Event pattern for {} service should be '{}', "
  31. "but found '{}'.".format("buckets",
  32. "provider.storage.buckets",
  33. self.provider.storage.buckets.
  34. _service_event_pattern))
  35. # pylint:disable=protected-access
  36. self.assertEqual(
  37. self.provider.storage._bucket_objects._service_event_pattern,
  38. "provider.storage._bucket_objects",
  39. "Event pattern for {} service should be '{}', "
  40. "but found '{}'.".format("bucket_objects",
  41. "provider.storage._bucket_objects",
  42. self.provider.storage._bucket_objects.
  43. _service_event_pattern))
  44. @helpers.skipIfNoService(['storage.buckets'])
  45. def test_crud_bucket(self):
  46. def create_bucket(name):
  47. return self.provider.storage.buckets.create(name)
  48. def cleanup_bucket(bucket):
  49. if bucket:
  50. bucket.delete()
  51. def extra_tests(bucket):
  52. # Recreating existing bucket should raise an exception
  53. with self.assertRaises(DuplicateResourceException):
  54. self.provider.storage.buckets.create(name=bucket.name)
  55. sit.check_crud(self, self.provider.storage.buckets, Bucket,
  56. "cb-crudbucket", create_bucket, cleanup_bucket,
  57. extra_test_func=extra_tests)
  58. @helpers.skipIfNoService(['storage.buckets'])
  59. def test_crud_bucket_object(self):
  60. test_bucket = None
  61. def create_bucket_obj(name):
  62. obj = test_bucket.objects.create(name)
  63. # TODO: This is wrong. We shouldn't have to have a separate
  64. # call to upload some content before being able to delete
  65. # the content. Maybe the create_object method should accept
  66. # the file content as a parameter.
  67. obj.upload("dummy content")
  68. return obj
  69. def cleanup_bucket_obj(bucket_obj):
  70. if bucket_obj:
  71. bucket_obj.delete()
  72. with cb_helpers.cleanup_action(lambda: test_bucket.delete()):
  73. name = "cb-crudbucketobj-{0}".format(helpers.get_uuid())
  74. test_bucket = self.provider.storage.buckets.create(name)
  75. sit.check_crud(self, test_bucket.objects, BucketObject,
  76. "cb-bucketobj", create_bucket_obj,
  77. cleanup_bucket_obj, skip_name_check=True)
  78. @helpers.skipIfNoService(['storage.buckets'])
  79. def test_crud_bucket_object_properties(self):
  80. # Create a new bucket, upload some contents into the bucket, and
  81. # check whether list properly detects the new content.
  82. # Delete everything afterwards.
  83. name = "cbtestbucketobjs-{0}".format(helpers.get_uuid())
  84. test_bucket = self.provider.storage.buckets.create(name)
  85. # ensure that the bucket is empty
  86. objects = test_bucket.objects.list()
  87. self.assertEqual([], objects)
  88. with cb_helpers.cleanup_action(lambda: test_bucket.delete()):
  89. obj_name_prefix = "hello"
  90. obj_name = obj_name_prefix + "_world.txt"
  91. obj = test_bucket.objects.create(obj_name)
  92. with cb_helpers.cleanup_action(lambda: obj.delete()):
  93. # TODO: This is wrong. We shouldn't have to have a separate
  94. # call to upload some content before being able to delete
  95. # the content. Maybe the create_object method should accept
  96. # the file content as a parameter.
  97. obj.upload("dummy content")
  98. objs = test_bucket.objects.list()
  99. self.assertTrue(
  100. isinstance(objs[0].size, int),
  101. "Object size property needs to be a int, not {0}".format(
  102. type(objs[0].size)))
  103. # GET an object as the size property implementation differs
  104. # for objects returned by LIST and GET.
  105. obj = test_bucket.objects.get(objs[0].id)
  106. self.assertTrue(
  107. isinstance(objs[0].size, int),
  108. "Object size property needs to be an int, not {0}".format(
  109. type(obj.size)))
  110. self.assertTrue(
  111. datetime.strptime(objs[0].last_modified[:23],
  112. "%Y-%m-%dT%H:%M:%S.%f"),
  113. "Object's last_modified field format {0} not matching."
  114. .format(objs[0].last_modified))
  115. # check iteration
  116. iter_objs = list(test_bucket.objects)
  117. self.assertListEqual(iter_objs, objs)
  118. obj_too = test_bucket.objects.get(obj_name)
  119. self.assertTrue(
  120. isinstance(obj_too, BucketObject),
  121. "Did not get object {0} of expected type.".format(obj_too))
  122. prefix_filtered_list = test_bucket.objects.list(
  123. prefix=obj_name_prefix)
  124. self.assertTrue(
  125. len(objs) == len(prefix_filtered_list) == 1,
  126. 'The number of objects returned by list function, '
  127. 'with and without a prefix, are expected to be equal, '
  128. 'but its detected otherwise.')
  129. sit.check_delete(self, test_bucket.objects, obj)
  130. @helpers.skipIfNoService(['storage.buckets'])
  131. def test_upload_download_bucket_content(self):
  132. name = "cbtestbucketobjs-{0}".format(helpers.get_uuid())
  133. test_bucket = self.provider.storage.buckets.create(name)
  134. with cb_helpers.cleanup_action(lambda: test_bucket.delete()):
  135. obj_name = "hello_upload_download.txt"
  136. obj = test_bucket.objects.create(obj_name)
  137. with cb_helpers.cleanup_action(lambda: obj.delete()):
  138. content = b"Hello World. Here's some content."
  139. # TODO: Upload and download methods accept different parameter
  140. # types. Need to make this consistent - possibly provider
  141. # multiple methods like upload_from_file, from_stream etc.
  142. obj.upload(content)
  143. target_stream = BytesIO()
  144. obj.save_content(target_stream)
  145. self.assertEqual(target_stream.getvalue(), content)
  146. target_stream2 = BytesIO()
  147. for data in obj.iter_content():
  148. target_stream2.write(data)
  149. self.assertEqual(target_stream2.getvalue(), content)
  150. @helpers.skipIfNoService(['storage.buckets'])
  151. def test_generate_url(self):
  152. name = "cbtestbucketobjs-{0}".format(helpers.get_uuid())
  153. test_bucket = self.provider.storage.buckets.create(name)
  154. with cb_helpers.cleanup_action(lambda: test_bucket.delete()):
  155. obj_name = "hello_upload_download.txt"
  156. obj = test_bucket.objects.create(obj_name)
  157. with cb_helpers.cleanup_action(lambda: obj.delete()):
  158. content = b"Hello World. Generate a url."
  159. obj.upload(content)
  160. target_stream = BytesIO()
  161. obj.save_content(target_stream)
  162. url = obj.generate_url(100)
  163. if isinstance(self.provider, TestMockHelperMixin):
  164. raise self.skipTest(
  165. "Skipping rest of test - mock providers can't"
  166. " access generated url")
  167. self.assertEqual(requests.get(url).content, content)
  168. @helpers.skipIfNoService(['storage.buckets'])
  169. def test_generate_url_write_permissions(self):
  170. name = "cbtestbucketobjs-{0}".format(helpers.get_uuid())
  171. test_bucket = self.provider.storage.buckets.create(name)
  172. with cb_helpers.cleanup_action(lambda: test_bucket.delete()):
  173. obj_name = "hello_upload_download.txt"
  174. obj = test_bucket.objects.create(obj_name)
  175. with cb_helpers.cleanup_action(lambda: obj.delete()):
  176. content = b"Hello World. Generate a url."
  177. url = obj.generate_url(100, writable=True)
  178. if isinstance(self.provider, TestMockHelperMixin):
  179. raise self.skipTest(
  180. "Skipping rest of test - mock providers can't"
  181. " access generated url")
  182. # Only Azure requires the x-ms-blob-type header to be present, but there's no harm
  183. # in sending this in for all providers.
  184. headers = {'x-ms-blob-type': 'BlockBlob'}
  185. response = requests.put(url, headers=headers, data=content)
  186. response.raise_for_status()
  187. obj = test_bucket.objects.get(obj_name)
  188. target_stream = BytesIO()
  189. obj.save_content(target_stream)
  190. self.assertEqual(target_stream.getvalue(), content)
  191. @helpers.skipIfNoService(['storage.buckets'])
  192. def test_generate_url_with_response_headers(self):
  193. name = "cbtestbucketobjs-{0}".format(helpers.get_uuid())
  194. test_bucket = self.provider.storage.buckets.create(name)
  195. with cb_helpers.cleanup_action(lambda: test_bucket.delete()):
  196. obj_name = "hello_response_headers.txt"
  197. obj = test_bucket.objects.create(obj_name)
  198. with cb_helpers.cleanup_action(lambda: obj.delete()):
  199. content = b"Hello World. Serve me with response headers."
  200. obj.upload(content)
  201. disposition = 'attachment; filename="hello.txt"'
  202. content_type = "application/octet-stream"
  203. url = obj.generate_url(100,
  204. content_disposition=disposition,
  205. content_type=content_type)
  206. if isinstance(self.provider, TestMockHelperMixin):
  207. # Presigned URLs are constructed client-side, so the
  208. # response-header overrides can be asserted on the query
  209. # string without a network round trip.
  210. self.assertIn('response-content-disposition', url)
  211. self.assertIn('response-content-type', url)
  212. raise self.skipTest(
  213. "Skipping rest of test - mock providers can't"
  214. " access generated url")
  215. response = requests.get(url)
  216. self.assertEqual(response.content, content)
  217. got_disposition = response.headers.get(
  218. 'Content-Disposition', '')
  219. self.assertTrue(got_disposition.startswith('attachment'),
  220. "Expected attachment disposition, got: "
  221. f"{got_disposition}")
  222. self.assertIn('hello.txt', got_disposition)
  223. if self.provider.PROVIDER_ID != 'openstack':
  224. # Swift's tempurl middleware cannot override Content-Type
  225. # (it only honors the filename portion of the disposition
  226. # via its `filename` query parameter).
  227. self.assertEqual(response.headers.get('Content-Type'),
  228. content_type)
  229. @helpers.skipIfNoService(['storage.buckets'])
  230. def test_generate_url_writable_ignores_response_headers(self):
  231. name = "cbtestbucketobjs-{0}".format(helpers.get_uuid())
  232. test_bucket = self.provider.storage.buckets.create(name)
  233. with cb_helpers.cleanup_action(lambda: test_bucket.delete()):
  234. obj_name = "hello_response_headers.txt"
  235. obj = test_bucket.objects.create(obj_name)
  236. with cb_helpers.cleanup_action(lambda: obj.delete()):
  237. url = obj.generate_url(
  238. 100, writable=True,
  239. content_disposition='attachment; filename="hello.txt"',
  240. content_type="application/octet-stream")
  241. # Response-header overrides only make sense for reads; a
  242. # writable URL must not carry them (across all providers:
  243. # AWS/GCP query params, Azure SAS rscd/rsct, Swift filename).
  244. self.assertNotIn('response-content-disposition', url)
  245. self.assertNotIn('rscd=', url)
  246. self.assertNotIn('filename=', url)
  247. @helpers.skipIfNoService(['storage.buckets'])
  248. def test_upload_download_bucket_content_from_file(self):
  249. name = "cbtestbucketobjs-{0}".format(helpers.get_uuid())
  250. test_bucket = self.provider.storage.buckets.create(name)
  251. with cb_helpers.cleanup_action(lambda: test_bucket.delete()):
  252. obj_name = "hello_upload_download.txt"
  253. obj = test_bucket.objects.create(obj_name)
  254. with cb_helpers.cleanup_action(lambda: obj.delete()):
  255. test_file = os.path.join(
  256. helpers.get_test_fixtures_folder(), 'logo.jpg')
  257. obj.upload_from_file(test_file)
  258. target_stream = BytesIO()
  259. obj.save_content(target_stream)
  260. with open(test_file, 'rb') as f:
  261. self.assertEqual(target_stream.getvalue(), f.read())
  262. @helpers.skipIfNoService(['storage.buckets'])
  263. def test_explicit_multipart_upload_roundtrip(self):
  264. name = "cbtest-mpu-{0}".format(helpers.get_uuid())
  265. test_bucket = self.provider.storage.buckets.create(name)
  266. with cb_helpers.cleanup_action(lambda: test_bucket.delete()):
  267. obj_name = "mpu-roundtrip.bin"
  268. obj = test_bucket.objects.create(obj_name)
  269. with cb_helpers.cleanup_action(lambda: obj.delete()):
  270. part1 = b"a" * MIN_PART_SIZE
  271. part2 = b"b" * MIN_PART_SIZE
  272. part3 = b"c" * 1024 # final part may be smaller than the min
  273. expected = part1 + part2 + part3
  274. upload = obj.create_multipart_upload()
  275. parts = [upload.upload_part(1, part1),
  276. upload.upload_part(2, part2),
  277. upload.upload_part(3, part3)]
  278. upload.complete(parts)
  279. stored = test_bucket.objects.get(obj_name)
  280. self.assertIsNotNone(
  281. stored, "Object should exist after multipart completion")
  282. self.assertEqual(stored.size, len(expected))
  283. target_stream = BytesIO()
  284. stored.save_content(target_stream)
  285. self.assertEqual(target_stream.getvalue(), expected)
  286. @helpers.skipIfNoService(['storage.buckets'])
  287. def test_multipart_upload_out_of_order_parts(self):
  288. name = "cbtest-mpu-{0}".format(helpers.get_uuid())
  289. test_bucket = self.provider.storage.buckets.create(name)
  290. with cb_helpers.cleanup_action(lambda: test_bucket.delete()):
  291. obj_name = "mpu-ooo.bin"
  292. obj = test_bucket.objects.create(obj_name)
  293. with cb_helpers.cleanup_action(lambda: obj.delete()):
  294. part1 = b"1" * MIN_PART_SIZE
  295. part2 = b"2" * MIN_PART_SIZE
  296. part3 = b"3" * 1024
  297. expected = part1 + part2 + part3
  298. upload = obj.create_multipart_upload()
  299. # Upload and collect parts out of order; complete must
  300. # assemble them in ascending part-number order regardless.
  301. p3 = upload.upload_part(3, part3)
  302. p1 = upload.upload_part(1, part1)
  303. p2 = upload.upload_part(2, part2)
  304. upload.complete([p3, p1, p2])
  305. stored = test_bucket.objects.get(obj_name)
  306. target_stream = BytesIO()
  307. stored.save_content(target_stream)
  308. self.assertEqual(target_stream.getvalue(), expected)
  309. @helpers.skipIfNoService(['storage.buckets'])
  310. def test_multipart_upload_abort(self):
  311. name = "cbtest-mpu-{0}".format(helpers.get_uuid())
  312. test_bucket = self.provider.storage.buckets.create(name)
  313. with cb_helpers.cleanup_action(lambda: test_bucket.delete()):
  314. obj_name = "mpu-abort.bin"
  315. obj = test_bucket.objects.create(obj_name)
  316. with cb_helpers.cleanup_action(lambda: obj.delete()):
  317. upload = obj.create_multipart_upload()
  318. upload.upload_part(1, b"a" * MIN_PART_SIZE)
  319. upload.abort()
  320. # Aborting must not commit any part data. Some providers
  321. # pre-materialise an empty placeholder on objects.create(), so
  322. # the target may be absent or empty afterwards -- but it must
  323. # never hold the uploaded part.
  324. stored = test_bucket.objects.get(obj_name)
  325. self.assertTrue(
  326. stored is None or stored.size == 0,
  327. "Aborted multipart upload must not commit any part data")
  328. @helpers.skipIfNoService(['storage.buckets'])
  329. def test_transparent_upload_large_stream_uses_multipart(self):
  330. name = "cbtest-mpu-{0}".format(helpers.get_uuid())
  331. test_bucket = self.provider.storage.buckets.create(name)
  332. with cb_helpers.cleanup_action(lambda: test_bucket.delete()):
  333. obj_name = "transparent.bin"
  334. obj = test_bucket.objects.create(obj_name)
  335. with cb_helpers.cleanup_action(lambda: obj.delete()):
  336. content = b"x" * (MIN_PART_SIZE * 2 + 1024)
  337. # Lower the threshold/part size so a modest stream crosses it,
  338. # and assert the multipart path is taken (each provider routes
  339. # its own way underneath) and the object round-trips exactly.
  340. with mock.patch.object(
  341. BaseBucketObject, 'CB_MULTIPART_THRESHOLD',
  342. MIN_PART_SIZE), \
  343. mock.patch.object(
  344. BaseBucketObject, 'CB_MULTIPART_PART_SIZE',
  345. MIN_PART_SIZE), \
  346. mock.patch.object(
  347. obj, '_upload_multipart',
  348. wraps=obj._upload_multipart) as spy:
  349. obj.upload(BytesIO(content))
  350. spy.assert_called_once()
  351. stored = test_bucket.objects.get(obj_name)
  352. self.assertEqual(stored.size, len(content))
  353. target_stream = BytesIO()
  354. stored.save_content(target_stream)
  355. self.assertEqual(target_stream.getvalue(), content)
  356. @helpers.skipIfNoService(['storage.buckets'])
  357. def test_small_upload_stays_single_shot(self):
  358. name = "cbtest-mpu-{0}".format(helpers.get_uuid())
  359. test_bucket = self.provider.storage.buckets.create(name)
  360. with cb_helpers.cleanup_action(lambda: test_bucket.delete()):
  361. obj = test_bucket.objects.create("small.txt")
  362. with cb_helpers.cleanup_action(lambda: obj.delete()):
  363. content = b"a small payload below the multipart threshold"
  364. # A payload below the threshold must not trigger multipart.
  365. svc = self.provider.storage._bucket_objects
  366. with mock.patch.object(
  367. svc, 'create_multipart_upload',
  368. wraps=svc.create_multipart_upload) as spy:
  369. obj.upload(content)
  370. spy.assert_not_called()
  371. target_stream = BytesIO()
  372. obj.save_content(target_stream)
  373. self.assertEqual(target_stream.getvalue(), content)
  374. @helpers.skipIfNoService(['storage.buckets'])
  375. def test_upload_from_file_uses_multipart_config(self):
  376. # AWS drives boto3's TransferManager with CloudBridge's multipart
  377. # knobs (not boto3's own defaults). This wiring is AWS-specific.
  378. if self.provider.PROVIDER_ID not in ('aws', 'mock'):
  379. self.skipTest("TransferConfig wiring is specific to AWS")
  380. from boto3.s3.transfer import TransferConfig as S3TransferConfig
  381. name = "cbtest-mpu-{0}".format(helpers.get_uuid())
  382. test_bucket = self.provider.storage.buckets.create(name)
  383. with cb_helpers.cleanup_action(lambda: test_bucket.delete()):
  384. obj = test_bucket.objects.create("config.bin")
  385. with cb_helpers.cleanup_action(lambda: obj.delete()):
  386. test_file = os.path.join(
  387. helpers.get_test_fixtures_folder(), 'logo.jpg')
  388. # pylint:disable=protected-access
  389. with mock.patch.object(
  390. BaseBucketObject, 'CB_MULTIPART_PART_SIZE',
  391. 7 * 1024 * 1024), \
  392. mock.patch.object(
  393. BaseBucketObject, 'CB_MULTIPART_MAX_CONCURRENCY', 3), \
  394. mock.patch.object(
  395. obj._obj, 'upload_file',
  396. wraps=obj._obj.upload_file) as spy:
  397. obj.upload_from_file(test_file)
  398. spy.assert_called_once()
  399. config = spy.call_args.kwargs['Config']
  400. self.assertIsInstance(config, S3TransferConfig)
  401. self.assertEqual(config.multipart_chunksize, 7 * 1024 * 1024)
  402. self.assertEqual(config.max_concurrency, 3)
  403. @helpers.skipIfNoService(['storage.buckets'])
  404. def test_per_call_upload_config_overrides_defaults(self):
  405. # A per-call TransferConfig takes precedence over the global knobs.
  406. if self.provider.PROVIDER_ID not in ('aws', 'mock'):
  407. self.skipTest("TransferConfig wiring is specific to AWS")
  408. from boto3.s3.transfer import TransferConfig as S3TransferConfig
  409. name = "cbtest-mpu-{0}".format(helpers.get_uuid())
  410. test_bucket = self.provider.storage.buckets.create(name)
  411. with cb_helpers.cleanup_action(lambda: test_bucket.delete()):
  412. obj = test_bucket.objects.create("percall.bin")
  413. with cb_helpers.cleanup_action(lambda: obj.delete()):
  414. test_file = os.path.join(
  415. helpers.get_test_fixtures_folder(), 'logo.jpg')
  416. cfg = TransferConfig(part_size=8 * 1024 * 1024,
  417. max_concurrency=2)
  418. # pylint:disable=protected-access
  419. with mock.patch.object(
  420. obj._obj, 'upload_file',
  421. wraps=obj._obj.upload_file) as spy:
  422. obj.upload_from_file(test_file, config=cfg)
  423. config = spy.call_args.kwargs['Config']
  424. self.assertIsInstance(config, S3TransferConfig)
  425. self.assertEqual(config.multipart_chunksize, 8 * 1024 * 1024)
  426. self.assertEqual(config.max_concurrency, 2)
  427. @helpers.skipIfNoService(['storage.buckets'])
  428. def test_download_to_file_roundtrip(self):
  429. name = "cbtest-dl-{0}".format(helpers.get_uuid())
  430. test_bucket = self.provider.storage.buckets.create(name)
  431. with cb_helpers.cleanup_action(lambda: test_bucket.delete()):
  432. obj = test_bucket.objects.create("roundtrip.txt")
  433. with cb_helpers.cleanup_action(lambda: obj.delete()):
  434. content = b"a small payload below the download threshold"
  435. obj.upload(content)
  436. stored = test_bucket.objects.get("roundtrip.txt")
  437. fd, path = tempfile.mkstemp()
  438. os.close(fd)
  439. with cb_helpers.cleanup_action(lambda: os.remove(path)):
  440. stored.download_to_file(path)
  441. with open(path, 'rb') as f:
  442. self.assertEqual(f.read(), content)
  443. @helpers.skipIfNoService(['storage.buckets'])
  444. def test_transparent_download_large_uses_ranged(self):
  445. name = "cbtest-dl-{0}".format(helpers.get_uuid())
  446. test_bucket = self.provider.storage.buckets.create(name)
  447. with cb_helpers.cleanup_action(lambda: test_bucket.delete()):
  448. obj = test_bucket.objects.create("ranged.bin")
  449. with cb_helpers.cleanup_action(lambda: obj.delete()):
  450. # Distinct byte pattern so out-of-order reassembly would show.
  451. content = bytes(range(256)) * ((MIN_PART_SIZE * 2) // 256)
  452. obj.upload(BytesIO(content),
  453. TransferConfig(threshold=MIN_PART_SIZE,
  454. part_size=MIN_PART_SIZE))
  455. stored = test_bucket.objects.get("ranged.bin")
  456. fd, path = tempfile.mkstemp()
  457. os.close(fd)
  458. with cb_helpers.cleanup_action(lambda: os.remove(path)):
  459. # Lower the threshold so the download crosses it (each
  460. # provider routes its own way underneath) and assert the
  461. # content round-trips exactly.
  462. stored.download_to_file(
  463. path, TransferConfig(threshold=MIN_PART_SIZE,
  464. part_size=MIN_PART_SIZE,
  465. max_concurrency=2))
  466. with open(path, 'rb') as f:
  467. self.assertEqual(f.read(), content)
  468. @helpers.skipIfNoService(['storage.buckets'])
  469. def test_download_to_file_uses_transfer_config(self):
  470. # AWS drives boto3's TransferManager with CloudBridge's transfer
  471. # knobs (not boto3's own defaults). This wiring is AWS-specific.
  472. if self.provider.PROVIDER_ID not in ('aws', 'mock'):
  473. self.skipTest("TransferConfig wiring is specific to AWS")
  474. from boto3.s3.transfer import TransferConfig as S3TransferConfig
  475. name = "cbtest-dl-{0}".format(helpers.get_uuid())
  476. test_bucket = self.provider.storage.buckets.create(name)
  477. with cb_helpers.cleanup_action(lambda: test_bucket.delete()):
  478. obj = test_bucket.objects.create("dlconfig.bin")
  479. with cb_helpers.cleanup_action(lambda: obj.delete()):
  480. obj.upload(b"some downloadable content")
  481. stored = test_bucket.objects.get("dlconfig.bin")
  482. fd, path = tempfile.mkstemp()
  483. os.close(fd)
  484. with cb_helpers.cleanup_action(lambda: os.remove(path)):
  485. cfg = TransferConfig(part_size=8 * 1024 * 1024,
  486. max_concurrency=2)
  487. # pylint:disable=protected-access
  488. with mock.patch.object(
  489. stored._obj.meta.client, 'download_file',
  490. wraps=stored._obj.meta.client
  491. .download_file) as spy:
  492. stored.download_to_file(path, config=cfg)
  493. spy.assert_called_once()
  494. config = spy.call_args.kwargs['Config']
  495. self.assertIsInstance(config, S3TransferConfig)
  496. self.assertEqual(
  497. config.multipart_chunksize, 8 * 1024 * 1024)
  498. self.assertEqual(config.max_concurrency, 2)
  499. with open(path, 'rb') as f:
  500. self.assertEqual(
  501. f.read(), b"some downloadable content")
  502. @skip("Skip unless you want to test objects bigger than 5GB")
  503. @helpers.skipIfNoService(['storage.buckets'])
  504. def test_upload_download_bucket_content_with_large_file(self):
  505. # Creates a 6 Gig file in the temp directory, then uploads it to
  506. # Swift. Once uploaded, then downloads to a new file in the temp
  507. # directory and compares the two files to see if they match.
  508. temp_dir = tempfile.gettempdir()
  509. file_name = '6GigTest.tmp'
  510. six_gig_file = os.path.join(temp_dir, file_name)
  511. with open(six_gig_file, "wb") as out:
  512. out.truncate(6 * 1024 * 1024 * 1024) # 6 Gig...
  513. with cb_helpers.cleanup_action(lambda: os.remove(six_gig_file)):
  514. download_file = "{0}/cbtestfile-{1}".format(temp_dir, file_name)
  515. bucket_name = "cbtestbucketlargeobjs-{0}".format(
  516. helpers.get_uuid())
  517. test_bucket = self.provider.storage.buckets.create(bucket_name)
  518. with cb_helpers.cleanup_action(lambda: test_bucket.delete()):
  519. test_obj = test_bucket.objects.create(file_name)
  520. with cb_helpers.cleanup_action(lambda: test_obj.delete()):
  521. file_uploaded = test_obj.upload_from_file(six_gig_file)
  522. self.assertTrue(file_uploaded, "Could not upload object?")
  523. with cb_helpers.cleanup_action(
  524. lambda: os.remove(download_file)):
  525. with open(download_file, 'wb') as f:
  526. test_obj.save_content(f)
  527. self.assertTrue(
  528. filecmp.cmp(six_gig_file, download_file),
  529. "Uploaded file != downloaded")