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

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. 

11 

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 

23 

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 

28 

29TESTDIR = os.path.abspath(os.path.dirname(__file__)) 

30 

31 

32class SimpleTestCase(unittest.TestCase): 

33 """Basic tests for file URIs.""" 

34 

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) 

43 

44 for example in ({1, 2, 3}, 42, self): 

45 self.assertNotIsInstance(example, ResourcePathExpression) 

46 

47 

48class FileTestCase(GenericTestCase, unittest.TestCase): 

49 """File-specific generic test cases.""" 

50 

51 scheme = "file" 

52 netloc = "localhost" 

53 

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") 

60 

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) 

65 

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) 

71 

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) 

79 

80 def test_path_lib(self): 

81 """File URIs can be created from pathlib.""" 

82 file = ResourcePath(self._make_uri("a/test.txt")) 

83 

84 path_file = pathlib.Path(file.ospath) 

85 from_path = ResourcePath(path_file) 

86 self.assertEqual(from_path.ospath, file.ospath) 

87 

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") 

92 

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") 

97 

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) 

106 

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, {}) 

113 

114 

115TEST_UMASK = 0o0333 

116 

117 

118class FileReadWriteTestCase(GenericReadWriteTestCase, unittest.TestCase): 

119 """File tests involving reading and writing of data.""" 

120 

121 scheme = "file" 

122 netloc = "localhost" 

123 testdir = TESTDIR 

124 transfer_modes = ("move", "copy", "link", "hardlink", "symlink", "relsymlink") 

125 

126 def test_transfer_identical(self): 

127 """Test overwrite of identical files. 

128 

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) 

140 

141 # Write a test file. 

142 src_file = dir1.join("test.txt") 

143 content = "0123456" 

144 src_file.write(content.encode()) 

145 

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) 

150 

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) 

155 

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) 

166 

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}") 

175 

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") 

182 

183 # If removal does not work it's worth reporting that as an error. 

184 tmp.remove() 

185 

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) 

192 

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)) 

200 

201 # Force the target directory to be created. 

202 target.transfer_from(tmp, "move") 

203 self.assertFalse(tmp.exists()) 

204 

205 # Temporary file now gone so transfer should not work. 

206 with self.assertRaises(FileNotFoundError): 

207 target.transfer_from(tmp, "move", overwrite=True) 

208 

209 def test_write_with_restrictive_umask(self): 

210 self._test_file_with_restrictive_umask(lambda target: target.write(b"123")) 

211 

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") 

217 

218 self._test_file_with_restrictive_umask(cb) 

219 

220 def test_mkdir_with_restrictive_umask(self): 

221 self._test_with_restrictive_umask(lambda target: target.mkdir()) 

222 

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()) 

228 

229 def _test_file_with_restrictive_umask(self, callback): 

230 def inner_cb(target): 

231 callback(target) 

232 

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) 

236 

237 self._test_with_restrictive_umask(inner_cb) 

238 

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()) 

248 

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}") 

256 

257 

258class BulkOperationTestCase(unittest.TestCase): 

259 """Tests for batched bulk operations on local files.""" 

260 

261 def setUp(self) -> None: 

262 self.tmpdir = ResourcePath(makeTestTempDir(TESTDIR), forceDirectory=True) 

263 

264 def tearDown(self) -> None: 

265 removeTestTempDir(self.tmpdir.ospath) 

266 

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 

272 

273 def test_chunk_sizes(self) -> None: 

274 items = list(range(10)) 

275 

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) 

279 

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]) 

283 

284 # An empty input yields no chunks at all. 

285 self.assertEqual(FileResourcePath._chunk_work([], 4), []) 

286 

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) 

293 

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) 

297 

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", []), {}) 

302 

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() 

308 

309 results = ResourcePath.mremove(uris, do_raise=False) 

310 

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()) 

317 

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") 

322 

323 results = ResourcePath.mexists([*present, absent]) 

324 

325 self.assertEqual(len(results), len(present) + 1) 

326 self.assertTrue(all(results[uri] for uri in present)) 

327 self.assertFalse(results[absent]) 

328 

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 

334 

335 def flaky(self: FileResourcePath) -> bool: 

336 if self.ospath == failing: 

337 raise PermissionError("cannot stat") 

338 return real_exists(self) 

339 

340 with unittest.mock.patch.object(FileResourcePath, "exists", flaky): 

341 results = ResourcePath.mexists(uris) 

342 

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])) 

346 

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))] 

352 

353 results = ResourcePath.mtransfer("copy", zip(sources, destinations, strict=True)) 

354 

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)) 

359 

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") 

365 

366 results = ResourcePath.mtransfer("copy", zip(sources, destinations, strict=True), do_raise=False) 

367 

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()) 

373 

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() 

378 

379 results = ResourcePath.mtransfer( 

380 "copy", zip(sources, destinations, strict=True), overwrite=True, transaction=transaction 

381 ) 

382 

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)) 

387 

388 for undo in transaction.undone: 

389 undo() 

390 self.assertFalse(any(dest.exists() for dest in destinations)) 

391 

392 

393class _RecordingTransaction: 

394 """Transaction that collects the undo actions registered against it.""" 

395 

396 def __init__(self) -> None: 

397 self.undone: list[Callable[[], Any]] = [] 

398 self._lock = threading.Lock() 

399 

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)) 

405 

406 

407@contextlib.contextmanager 

408def _override_umask(temp_umask): 

409 old = os.umask(temp_umask) 

410 try: 

411 yield 

412 finally: 

413 os.umask(old) 

414 

415 

416if __name__ == "__main__": 

417 unittest.main()