Coverage for tests/test_file.py: 99%
258 statements
« prev ^ index » next coverage.py v7.16.0, created at 2026-09-23 02:10 -0700
« prev ^ index » next coverage.py v7.16.0, created at 2026-09-23 02:10 -0700
1# This file is part of lsst-resources.
2#
3# Developed for the LSST Data Management System.
4# This product includes software developed by the LSST Project
5# (https://www.lsst.org).
6# See the COPYRIGHT file at the top-level directory of this distribution
7# for details of code ownership.
8#
9# Use of this source code is governed by a 3-clause BSD-style
10# license that can be found in the LICENSE file.
12import contextlib
13import datetime
14import functools
15import os
16import pathlib
17import threading
18import unittest
19import unittest.mock
20import urllib.parse
21from collections.abc import Callable, Iterator
22from typing import Any
24from lsst.resources import ResourceInfo, ResourcePath, ResourcePathExpression
25from lsst.resources.file import FileResourcePath
26from lsst.resources.tests import GenericReadWriteTestCase, GenericTestCase
27from lsst.resources.utils import makeTestTempDir, removeTestTempDir
29TESTDIR = os.path.abspath(os.path.dirname(__file__))
32class SimpleTestCase(unittest.TestCase):
33 """Basic tests for file URIs."""
35 def test_instance(self):
36 for example in (
37 "xxx",
38 ResourcePath("xxx"),
39 pathlib.Path("xxx"),
40 urllib.parse.urlparse("file:///xxx"),
41 ):
42 self.assertIsInstance(example, ResourcePathExpression)
44 for example in ({1, 2, 3}, 42, self):
45 self.assertNotIsInstance(example, ResourcePathExpression)
48class FileTestCase(GenericTestCase, unittest.TestCase):
49 """File-specific generic test cases."""
51 scheme = "file"
52 netloc = "localhost"
54 def test_env_var(self):
55 """Test that environment variables are expanded."""
56 with unittest.mock.patch.dict(os.environ, {"MY_TEST_DIRX": "/a/b/c"}):
57 uri = ResourcePath("${MY_TEST_DIRX}/d.txt")
58 self.assertEqual(uri.path, "/a/b/c/d.txt")
59 self.assertEqual(uri.scheme, "file")
61 # This will not expand
62 uri = ResourcePath("${MY_TEST_DIRX}/d.txt", forceAbsolute=False)
63 self.assertEqual(uri.path, "${MY_TEST_DIRX}/d.txt")
64 self.assertFalse(uri.scheme)
66 def test_ospath(self):
67 """File URIs have ospath property."""
68 file = ResourcePath(self._make_uri("a/test.txt"))
69 self.assertEqual(file.ospath, "/a/test.txt")
70 self.assertEqual(file.ospath, file.path)
72 # A Schemeless URI can take unquoted files but will be quoted
73 # when it becomes a file URI.
74 something = "/a#/???.txt"
75 file = ResourcePath(something, forceAbsolute=True)
76 self.assertEqual(file.scheme, "file")
77 self.assertEqual(file.ospath, something, "From URI: {file}")
78 self.assertNotIn("???", file.path)
80 def test_path_lib(self):
81 """File URIs can be created from pathlib."""
82 file = ResourcePath(self._make_uri("a/test.txt"))
84 path_file = pathlib.Path(file.ospath)
85 from_path = ResourcePath(path_file)
86 self.assertEqual(from_path.ospath, file.ospath)
88 def test_schemeless_root(self):
89 root = ResourcePath(self._make_uri("/root"))
90 via_root = ResourcePath("b.txt", root=root)
91 self.assertEqual(via_root.ospath, "/root/b.txt")
93 def test_get_info(self):
94 now = datetime.datetime.now(tz=datetime.UTC)
95 with ResourcePath.temporary_uri(suffix=".txt") as target:
96 target.write(b"abc")
98 info = target.get_info()
99 self.assertIsInstance(info, ResourceInfo)
100 self.assertTrue(info.uri.endswith(".txt"))
101 self.assertTrue(info.is_file)
102 self.assertEqual(info.size, 3)
103 self.assertEqual(info.checksums, {})
104 self.assertEqual(info.last_modified.tzinfo, datetime.UTC)
105 self.assertGreaterEqual(info.last_modified.timestamp(), now.timestamp() - 1.0)
107 dirinfo = target.parent().get_info()
108 self.assertEqual(dirinfo.uri, str(target.parent()))
109 self.assertFalse(dirinfo.is_file)
110 self.assertEqual(dirinfo.size, 0)
111 self.assertGreaterEqual(dirinfo.last_modified.timestamp(), 0)
112 self.assertEqual(dirinfo.checksums, {})
115TEST_UMASK = 0o0333
118class FileReadWriteTestCase(GenericReadWriteTestCase, unittest.TestCase):
119 """File tests involving reading and writing of data."""
121 scheme = "file"
122 netloc = "localhost"
123 testdir = TESTDIR
124 transfer_modes = ("move", "copy", "link", "hardlink", "symlink", "relsymlink")
126 def test_transfer_identical(self):
127 """Test overwrite of identical files.
129 Only relevant for local files.
130 """
131 dir1 = self.tmpdir.join("dir1", forceDirectory=True)
132 dir1.mkdir()
133 self.assertTrue(dir1.exists())
134 dir2 = self.tmpdir.join("dir2", forceDirectory=True)
135 # A symlink can't include a trailing slash.
136 dir2_ospath = dir2.ospath
137 if dir2_ospath.endswith("/"): 137 ↛ 139line 137 didn't jump to line 139 because the condition on line 137 was always true
138 dir2_ospath = dir2_ospath[:-1]
139 os.symlink(dir1.ospath, dir2_ospath)
141 # Write a test file.
142 src_file = dir1.join("test.txt")
143 content = "0123456"
144 src_file.write(content.encode())
146 # Construct URI to destination that should be identical.
147 dest_file = dir2.join("test.txt")
148 self.assertTrue(dest_file.exists())
149 self.assertNotEqual(src_file, dest_file)
151 # Transfer it over itself.
152 dest_file.transfer_from(src_file, transfer="symlink", overwrite=True)
153 new_content = dest_file.read().decode()
154 self.assertEqual(content, new_content)
156 def test_local_temporary(self):
157 """Create temporary local file if no prefix specified."""
158 with ResourcePath.temporary_uri(suffix=".json") as tmp:
159 self.assertEqual(tmp.getExtension(), ".json", f"uri: {tmp}")
160 self.assertTrue(tmp.isabs(), f"uri: {tmp}")
161 self.assertFalse(tmp.exists(), f"uri: {tmp}")
162 tmp.write(b"abcd")
163 self.assertTrue(tmp.exists(), f"uri: {tmp}")
164 self.assertTrue(tmp.isTemporary)
165 self.assertTrue(tmp.isLocal)
167 # If we now ask for a local form of this temporary file
168 # it should still be temporary and it should not be deleted
169 # on exit.
170 with tmp.as_local() as loc:
171 self.assertEqual(tmp, loc)
172 self.assertTrue(loc.isTemporary)
173 self.assertTrue(tmp.exists())
174 self.assertFalse(tmp.exists(), f"uri: {tmp}")
176 with ResourcePath.temporary_uri(suffix=".yaml", delete=False) as tmp:
177 tmp.write(b"1234")
178 self.assertTrue(tmp.exists(), f"uri: {tmp}")
179 # If the file doesn't exist there is nothing to clean up so a failure
180 # here is not a problem.
181 self.assertTrue(tmp.exists(), f"uri: {tmp} should still exist")
183 # If removal does not work it's worth reporting that as an error.
184 tmp.remove()
186 def test_transfers_from_local(self):
187 """Extra tests for local transfers."""
188 target = self.tmpdir.join("a/target.txt")
189 with ResourcePath.temporary_uri() as tmp:
190 tmp.write(b"")
191 self.assertTrue(tmp.isTemporary)
193 # Symlink transfers for temporary resources should
194 # trigger a debug message.
195 for transfer in ("symlink", "relsymlink"):
196 with self.assertLogs("lsst.resources", level="DEBUG") as cm:
197 target.transfer_from(tmp, transfer)
198 target.remove()
199 self.assertIn("Using a symlink for a temporary", "".join(cm.output))
201 # Force the target directory to be created.
202 target.transfer_from(tmp, "move")
203 self.assertFalse(tmp.exists())
205 # Temporary file now gone so transfer should not work.
206 with self.assertRaises(FileNotFoundError):
207 target.transfer_from(tmp, "move", overwrite=True)
209 def test_write_with_restrictive_umask(self):
210 self._test_file_with_restrictive_umask(lambda target: target.write(b"123"))
212 def test_transfer_from_with_restrictive_umask(self):
213 def cb(target):
214 with ResourcePath.temporary_uri() as tmp:
215 tmp.write(b"")
216 target.transfer_from(tmp, "copy")
218 self._test_file_with_restrictive_umask(cb)
220 def test_mkdir_with_restrictive_umask(self):
221 self._test_with_restrictive_umask(lambda target: target.mkdir())
223 def test_temporary_uri_with_restrictive_umask(self):
224 with _override_umask(TEST_UMASK):
225 with ResourcePath.temporary_uri() as tmp:
226 tmp.write(b"")
227 self.assertTrue(tmp.exists())
229 def _test_file_with_restrictive_umask(self, callback):
230 def inner_cb(target):
231 callback(target)
233 # Make sure the umask was respected for the file itself
234 file_mode = os.stat(target.ospath).st_mode
235 self.assertEqual(file_mode & TEST_UMASK, 0)
237 self._test_with_restrictive_umask(inner_cb)
239 def _test_with_restrictive_umask(self, callback):
240 """Make sure that parent directories for a file can be created even if
241 the user has set a process umask that restricts the write and traverse
242 bits.
243 """
244 with _override_umask(TEST_UMASK):
245 target = self.tmpdir.join("a/b/target.txt")
246 callback(target)
247 self.assertTrue(target.exists())
249 dir_b_path = os.path.dirname(target.ospath)
250 dir_a_path = os.path.dirname(dir_b_path)
251 for dir in [dir_a_path, dir_b_path]:
252 # Make sure we only added the minimum permissions needed for it
253 # to work (owner-write and owner-traverse)
254 mode = os.stat(dir).st_mode
255 self.assertEqual(mode & TEST_UMASK, 0o0300, f"Permissions incorrect for {dir}: {mode:o}")
258class BulkOperationTestCase(unittest.TestCase):
259 """Tests for batched bulk operations on local files."""
261 def setUp(self) -> None:
262 self.tmpdir = ResourcePath(makeTestTempDir(TESTDIR), forceDirectory=True)
264 def tearDown(self) -> None:
265 removeTestTempDir(self.tmpdir.ospath)
267 def _make_files(self, prefix: str, count: int) -> list[ResourcePath]:
268 uris = [self.tmpdir.join(f"{prefix}{n}.txt") for n in range(count)]
269 for uri in uris:
270 uri.write(b"")
271 return uris
273 def test_chunk_sizes(self) -> None:
274 items = list(range(10))
276 # The batch size does not depend on how many items there are.
277 self.assertEqual([len(c) for c in FileResourcePath._chunk_work(items, 4)], [4, 4, 2])
278 self.assertEqual([len(c) for c in FileResourcePath._chunk_work(items, 1)], [1] * 10)
280 # A batch that fits in one chunk is the signal to the caller to handle
281 # it without a pool.
282 self.assertEqual([len(c) for c in FileResourcePath._chunk_work(items, 100)], [10])
284 # An empty input yields no chunks at all.
285 self.assertEqual(FileResourcePath._chunk_work([], 4), [])
287 def test_schemes_size_their_own_chunks(self) -> None:
288 # A local operation is cheap enough that a batch has to be large
289 # before threading it pays, while a scheme whose every operation is a
290 # round trip is worth overlapping immediately.
291 self.assertEqual(FileResourcePath._chunk_size, 1000)
292 self.assertEqual(ResourcePath._chunk_size, 1)
294 # A transfer is far more expensive than an existence check on the same
295 # scheme, so it batches separately.
296 self.assertLess(FileResourcePath._transfer_chunk_size, FileResourcePath._chunk_size)
298 def test_empty_bulk_operations_are_no_ops(self) -> None:
299 self.assertEqual(ResourcePath.mremove([]), {})
300 self.assertEqual(ResourcePath.mexists([]), {})
301 self.assertEqual(ResourcePath.mtransfer("copy", []), {})
303 @unittest.mock.patch.object(FileResourcePath, "_chunk_size", 1)
304 def test_removal_failure_does_not_abandon_the_rest(self) -> None:
305 uris = self._make_files("f", 20)
306 # Remove one out from under the batch so that its own removal raises.
307 uris[1].remove()
309 results = ResourcePath.mremove(uris, do_raise=False)
311 self.assertEqual(len(results), len(uris))
312 self.assertFalse(results[uris[1]].success)
313 self.assertIsInstance(results[uris[1]].exception, FileNotFoundError)
314 for uri in uris[2:]:
315 self.assertTrue(results[uri].success, f"{uri} should have been removed")
316 self.assertFalse(uri.exists())
318 @unittest.mock.patch.object(FileResourcePath, "_chunk_size", 1)
319 def test_existence_check_reports_each_uri(self) -> None:
320 present = self._make_files("p", 20)
321 absent = self.tmpdir.join("gone.txt")
323 results = ResourcePath.mexists([*present, absent])
325 self.assertEqual(len(results), len(present) + 1)
326 self.assertTrue(all(results[uri] for uri in present))
327 self.assertFalse(results[absent])
329 @unittest.mock.patch.object(FileResourcePath, "_chunk_size", 1)
330 def test_existence_check_treats_an_error_as_missing(self) -> None:
331 uris = self._make_files("e", 20)
332 failing = uris[1].ospath
333 real_exists = FileResourcePath.exists
335 def flaky(self: FileResourcePath) -> bool:
336 if self.ospath == failing:
337 raise PermissionError("cannot stat")
338 return real_exists(self)
340 with unittest.mock.patch.object(FileResourcePath, "exists", flaky):
341 results = ResourcePath.mexists(uris)
343 # The failure is reported as absent and the rest are still checked.
344 self.assertFalse(results[uris[1]])
345 self.assertTrue(all(results[uri] for uri in uris if uri != uris[1]))
347 def test_transfer_of_many_files(self) -> None:
348 sources = self._make_files("src", 50)
349 for i, src in enumerate(sources):
350 src.write(f"{i}".encode())
351 destinations = [self.tmpdir.join(f"dest{n}.txt") for n in range(len(sources))]
353 results = ResourcePath.mtransfer("copy", zip(sources, destinations, strict=True))
355 self.assertEqual(len(results), len(sources))
356 self.assertTrue(all(res.success for res in results.values()))
357 for i, dest in enumerate(destinations):
358 self.assertEqual(dest.read().decode(), str(i))
360 def test_transfer_failure_does_not_abandon_the_rest(self) -> None:
361 sources = self._make_files("s", 20)
362 destinations = [self.tmpdir.join(f"d{n}.txt") for n in range(len(sources))]
363 # An existing target fails when overwriting is not allowed.
364 destinations[1].write(b"in the way")
366 results = ResourcePath.mtransfer("copy", zip(sources, destinations, strict=True), do_raise=False)
368 self.assertEqual(len(results), len(sources))
369 self.assertFalse(results[destinations[1]].success)
370 for dest in destinations[2:]:
371 self.assertTrue(results[dest].success, f"{dest} should have been written")
372 self.assertTrue(dest.exists())
374 def test_transfer_registers_undo_actions(self) -> None:
375 sources = self._make_files("u", 20)
376 destinations = [self.tmpdir.join(f"undo{n}.txt") for n in range(len(sources))]
377 transaction = _RecordingTransaction()
379 results = ResourcePath.mtransfer(
380 "copy", zip(sources, destinations, strict=True), overwrite=True, transaction=transaction
381 )
383 self.assertTrue(all(res.success for res in results.values()))
384 # Every transfer must be undoable by the caller that supplied the
385 # transaction, no matter which worker performed it.
386 self.assertEqual(len(transaction.undone), len(sources))
388 for undo in transaction.undone:
389 undo()
390 self.assertFalse(any(dest.exists() for dest in destinations))
393class _RecordingTransaction:
394 """Transaction that collects the undo actions registered against it."""
396 def __init__(self) -> None:
397 self.undone: list[Callable[[], Any]] = []
398 self._lock = threading.Lock()
400 @contextlib.contextmanager
401 def undoWith(self, name: str, undoFunc: Callable, *args: Any, **kwargs: Any) -> Iterator[None]:
402 yield None
403 with self._lock:
404 self.undone.append(functools.partial(undoFunc, *args, **kwargs))
407@contextlib.contextmanager
408def _override_umask(temp_umask):
409 old = os.umask(temp_umask)
410 try:
411 yield
412 finally:
413 os.umask(old)
416if __name__ == "__main__":
417 unittest.main()