Low-level design
Part 2 of 3 · ArchiveSweepArchiveSweep - Streaming SHA Solution & Tests
Full runnable stdlib implementation: archive_sweep(root, workers) with scan, sample, stream hash, byte verify, ThreadPoolExecutor, and unittest coverage for duplicates, near-misses, symlinks, hard links, and empty files.
- 1Gist
- 2Maps
- 3Q&A
- 4Sandbox
Voice readout needs Web Speech Synthesis in this browser.
Question ladder
L1
What does scan_files skip?
Answer
Symlinks, non-regular files, and a second path for an inode already seen.
L2
Why is the sample keyed with the size?
Answer
The sample runs on a size bucket. The map key keeps that size with the prefix so later stages do not mix lengths.
L3
Where is the lock?
Answer
Around SweepStats updates and the error list inside the worker. The file read itself is outside the lock.
L4
What does hash_filter drop?
Answer
Files that errored and digest buckets of length one.
L5
What does byte_equal require?
Answer
Equal chunks until both reads return empty. A mismatch returns false.
L6
Which test covers a shared prefix?
Answer
test_same_prefix_different_tail writes 4KiB of P and different tails and expects no group.
L7
Which test covers hard links?
Answer
test_hardlinks_count_once expects skipped_hardlink_extras and an empty group list.
Failure modes
Truncated hash
A short read can collide. byte_equal is the check that refuses to emit the group.
Unlocked stats
Two workers incrementing scanned or errors can lose updates. The lock covers those writes.
Symlink followed
test_skips_symlinks expects the link to be counted and absent from any group.
Misconceptions
Workers must lock the whole pipeline.
They lock the shared counters. Buckets are built after the pool joins.
Empty files cannot be duplicates.
Two empty files share a size, a sample, a digest, and the bytes. They are one group.
Resolve the path before the symlink check.
The code checks is_symlink before stat and resolve, so the link is not walked as the target.
Interviewer traps
Paste a shorter stub than the tests import.
The module and the unittest file on this page are the full source.
Sort inside the worker.
Sort the finished groups in the parent so order does not depend on completion.
Design scenario
Same prompt for every reader.
Requirements
Stdlib only. Thread pool for sample and hash. Seven tests: identical, same size, same prefix, symlink, hard link, empty files, helpers.
Traffic / scale
A temp directory and a handful of workers. The test is correctness, not throughput.
Latency
Chunked reads. No whole-file buffer.
Consistency
Group membership is byte equality. Output order is sorted paths.
Availability
An OSError on one path increments errors and stays out of the groups.
Failure assumptions
- Two files share 4KiB and differ after that.
- A symlink points at a real file in the tree.
- Two names are the same inode.
Constraints
- Do not open a new full-file bytes object for the hash.
- Do not follow symlinks.
Prompt
Implement archive_sweep and the tests that prove the pipeline.
API
Which assertions does test_finds_identical_duplicates make?
Data
What fields does FileRecord hold after each stage?
Architecture
What does the lock cover, and what is built after join?
Where the pool is allowed
Prefer
Pool the independent reads, lock the shared stats
Sample and hash do not need each other's bytes. The counters do.
- Buckets are rebuilt after join.
- byte_equal runs against the first path.
- Groups are sorted in the parent.
Alternative
One thread, whole-file reads
Easier to draw and the wrong cost once the tree is large.
- No overlap of I/O.
- A big file becomes one allocation.
- The tests still pass, and the design story does not.
From the archivesweep directory, run python3 test_archive_sweep.py -v. You want identical files in one group, a shared prefix in no group, a symlink skipped, and a hard link collapsed.
Overview
Full runnable stdlib implementation: archive_sweep(root, workers) with scan, sample, stream hash, byte verify, ThreadPoolExecutor, and unittest coverage for duplicates, near-misses, symlinks, hard links, and empty files.
Run instructions
cd code/archivesweep
python3 archive_sweep.py --root /path/to/tree --workers 4
python3 test_archive_sweep.py -vFilter until the bytes match
Diagram 1. Size, sample, and digest each drop a file that is alone. A group is emitted only after the bytes match.
- 1
Scan and skip symlinks
os.walk does not follow links. Symlink directories are pruned. - 2
Collapse hard links
A seen (device, inode) increments skipped_hardlink_extras. - 3
Drop singleton sizes
group_by_size keeps only lengths that appear twice. - 4
Sample in the pool
Each worker reads 4KiB. The lock covers the stats update. - 5
Stream SHA-256
Unique digests are dropped after the pool joins. - 6
Byte verify
A mismatch is not a group. A match is sorted and emitted.
Decisions
- 1
1. Scan files skip symlinks
- next2. Collapse hard links by inode
- 2
2. Collapse hard links by inode
- next3. Group by size
- 3
3. Group by size
- next4. More than one file?
- ?
4. More than one file?
- No5. Drop the size
- Yes6. Sample prefix in a pool
- 5
5. Drop the size
- 6
6. Sample prefix in a pool
- next7. Sample collision?
- ?
7. Sample collision?
- No5. Drop the size
- Yes8. Stream SHA-256
- 8
8. Stream SHA-256
- next9. Digest collision?
- ?
9. Digest collision?
- No5. Drop the size
- Yes10. Byte verify
- 10
10. Byte verify
- next11. Emit sorted group
- 11
11. Emit sorted group
Lesson map
ArchiveSweep - Streaming SHA Solution & Tests
Full runnable stdlib implementation: archive_sweep(root, workers) with scan, sample, stream hash, byte verify, ThreadPoolExecutor, and unittest coverage for duplicates, near-misses, symlinks, hard links, and empty files.
Architecture. Architecture
Select a node to see why it exists, or an edge to see the protocol, direction, effect, and consequence.
Mermaid export
flowchart TB a["1. Scan files skip symlinks"] b["2. Collapse hard links by inode"] c["3. Group by size"] d["4. More than one file?"] z["5. Drop the size"] e["6. Sample prefix in a pool"] f["7. Sample collision?"] g["8. Stream SHA-256"] h["9. Digest collision?"] i["10. Byte verify"] j["11. Emit sorted group"] a -->|continues| b b -->|continues| c c -->|continues| d d -->|No| z d -->|Yes| e e -->|continues| f f -->|No| z f -->|Yes| g g -->|continues| h h -->|No| z h -->|Yes| i i -->|continues| j
Solution code
"""ArchiveSweep - recursive file deduplication.
Sandbox: python3 archive_sweep.py --root /tmp/demo
Tests: python3 test_archive_sweep.py
stdlib only.
"""
from __future__ import annotations
import argparse
import hashlib
import os
import stat
import threading
from concurrent.futures import ThreadPoolExecutor, as_completed
from dataclasses import dataclass, field
from pathlib import Path
from typing import Iterable, Optional
SAMPLE_BYTES = 4096
HASH_CHUNK = 1024 * 1024
@dataclass
class FileRecord:
path: str
size: int
sample: Optional[bytes] = None
digest: Optional[str] = None
error: Optional[str] = None
@dataclass
class SweepStats:
scanned: int = 0
sampled: int = 0
hashed: int = 0
verified_groups: int = 0
duplicate_files: int = 0
errors: int = 0
skipped_symlinks: int = 0
skipped_hardlink_extras: int = 0
@dataclass
class SweepResult:
groups: list[list[str]] = field(default_factory=list)
stats: SweepStats = field(default_factory=SweepStats)
errors: list[str] = field(default_factory=list)
def _is_symlink(path: Path) -> bool:
try:
return path.is_symlink()
except OSError:
return False
def scan_files(root: Path, stats: SweepStats, errors: list[str]) -> list[FileRecord]:
"""Recursive scan; skip symlinks; keep one path per inode (hard links)."""
seen_inodes: set[tuple[int, int]] = set()
out: list[FileRecord] = []
for dirpath, dirnames, filenames in os.walk(root, followlinks=False):
# prune symlink dirs
keep = []
for d in dirnames:
p = Path(dirpath) / d
if _is_symlink(p):
stats.skipped_symlinks += 1
else:
keep.append(d)
dirnames[:] = keep
for name in filenames:
p = Path(dirpath) / name
try:
if _is_symlink(p):
stats.skipped_symlinks += 1
continue
st = p.stat()
if not stat.S_ISREG(st.st_mode):
continue
key = (st.st_dev, st.st_ino)
if key in seen_inodes:
stats.skipped_hardlink_extras += 1
continue
seen_inodes.add(key)
out.append(FileRecord(path=str(p.resolve()), size=st.st_size))
stats.scanned += 1
except OSError as e:
stats.errors += 1
errors.append(f"scan:{p}:{e}")
return out
def sample_prefix(path: str, n: int = SAMPLE_BYTES) -> bytes:
with open(path, "rb") as f:
return f.read(n)
def sha256_stream(path: str, chunk: int = HASH_CHUNK) -> str:
h = hashlib.sha256()
with open(path, "rb") as f:
while True:
block = f.read(chunk)
if not block:
break
h.update(block)
return h.hexdigest()
def byte_equal(a: str, b: str, chunk: int = HASH_CHUNK) -> bool:
with open(a, "rb") as fa, open(b, "rb") as fb:
while True:
xa, xb = fa.read(chunk), fb.read(chunk)
if xa != xb:
return False
if not xa:
return True
def group_by_size(files: list[FileRecord]) -> dict[int, list[FileRecord]]:
buckets: dict[int, list[FileRecord]] = {}
for fr in files:
buckets.setdefault(fr.size, []).append(fr)
return {k: v for k, v in buckets.items() if len(v) > 1}
def sample_filter(
candidates: list[FileRecord],
stats: SweepStats,
errors: list[str],
workers: int,
) -> dict[tuple[int, bytes], list[FileRecord]]:
"""Size already equal within list; group by sample prefix."""
lock = threading.Lock()
def work(fr: FileRecord) -> None:
try:
fr.sample = sample_prefix(fr.path)
with lock:
stats.sampled += 1
except OSError as e:
fr.error = str(e)
with lock:
stats.errors += 1
errors.append(f"sample:{fr.path}:{e}")
with ThreadPoolExecutor(max_workers=max(1, workers)) as ex:
list(ex.map(work, candidates))
buckets: dict[tuple[int, bytes], list[FileRecord]] = {}
for fr in candidates:
if fr.error or fr.sample is None:
continue
buckets.setdefault((fr.size, fr.sample), []).append(fr)
return {k: v for k, v in buckets.items() if len(v) > 1}
def hash_filter(
candidates: list[FileRecord],
stats: SweepStats,
errors: list[str],
workers: int,
) -> dict[str, list[FileRecord]]:
lock = threading.Lock()
def work(fr: FileRecord) -> None:
try:
fr.digest = sha256_stream(fr.path)
with lock:
stats.hashed += 1
except OSError as e:
fr.error = str(e)
with lock:
stats.errors += 1
errors.append(f"hash:{fr.path}:{e}")
with ThreadPoolExecutor(max_workers=max(1, workers)) as ex:
list(ex.map(work, candidates))
buckets: dict[str, list[FileRecord]] = {}
for fr in candidates:
if fr.error or not fr.digest:
continue
buckets.setdefault(fr.digest, []).append(fr)
return {k: v for k, v in buckets.items() if len(v) > 1}
def verify_group(files: list[FileRecord], stats: SweepStats, errors: list[str]) -> list[str]:
"""Byte-for-byte verify against first file; deterministic path sort for group."""
ordered = sorted(files, key=lambda f: f.path)
keep = [ordered[0]]
for fr in ordered[1:]:
try:
if byte_equal(ordered[0].path, fr.path):
keep.append(fr)
else:
# hash collision survivor
errors.append(f"verify-mismatch:{fr.path}")
except OSError as e:
stats.errors += 1
errors.append(f"verify:{fr.path}:{e}")
if len(keep) > 1:
stats.verified_groups += 1
stats.duplicate_files += len(keep) - 1
return [f.path for f in keep]
return []
def archive_sweep(root: str | Path, workers: int = 4) -> SweepResult:
root_p = Path(root)
result = SweepResult()
files = scan_files(root_p, result.stats, result.errors)
size_groups = group_by_size(files)
for _size, group in size_groups.items():
sampled = sample_filter(group, result.stats, result.errors, workers)
for _key, sg in sampled.items():
hashed = hash_filter(sg, result.stats, result.errors, workers)
for _digest, hg in hashed.items():
verified = verify_group(hg, result.stats, result.errors)
if verified:
result.groups.append(verified)
# deterministic group order by first path
result.groups.sort(key=lambda g: g[0])
return result
def main(argv: Optional[list[str]] = None) -> int:
ap = argparse.ArgumentParser(description="ArchiveSweep file deduper")
ap.add_argument("--root", required=True)
ap.add_argument("--workers", type=int, default=4)
args = ap.parse_args(argv)
res = archive_sweep(args.root, workers=args.workers)
print(f"groups={len(res.groups)} scanned={res.stats.scanned} errors={res.stats.errors}")
for g in res.groups:
print("---")
for p in g:
print(p)
return 0
if __name__ == "__main__":
raise SystemExit(main())
Tests
"""Tests for ArchiveSweep."""
from __future__ import annotations
import os
import tempfile
import unittest
from pathlib import Path
from archive_sweep import archive_sweep, sample_prefix, sha256_stream, byte_equal
class ArchiveSweepTests(unittest.TestCase):
def _write(self, root: Path, rel: str, data: bytes) -> Path:
p = root / rel
p.parent.mkdir(parents=True, exist_ok=True)
p.write_bytes(data)
return p
def test_finds_identical_duplicates(self):
with tempfile.TemporaryDirectory() as td:
root = Path(td)
payload = b"hello-archive-sweep-" + (b"x" * 5000)
self._write(root, "a/one.bin", payload)
self._write(root, "b/two.bin", payload)
self._write(root, "c/unique.bin", b"different")
res = archive_sweep(root, workers=2)
self.assertEqual(len(res.groups), 1)
paths = set(res.groups[0])
self.assertEqual(len(paths), 2)
self.assertTrue(any("one.bin" in p for p in paths))
self.assertTrue(any("two.bin" in p for p in paths))
def test_same_size_different_content(self):
with tempfile.TemporaryDirectory() as td:
root = Path(td)
self._write(root, "a.bin", b"AAAA" * 1000)
self._write(root, "b.bin", b"BBBB" * 1000)
res = archive_sweep(root, workers=2)
self.assertEqual(res.groups, [])
def test_same_prefix_different_tail(self):
with tempfile.TemporaryDirectory() as td:
root = Path(td)
prefix = b"P" * 4096
self._write(root, "a.bin", prefix + b"TAIL-A")
self._write(root, "b.bin", prefix + b"TAIL-B")
res = archive_sweep(root, workers=2)
self.assertEqual(res.groups, [])
def test_skips_symlinks(self):
with tempfile.TemporaryDirectory() as td:
root = Path(td)
real = self._write(root, "real.bin", b"symlink-target-data")
link = root / "link.bin"
link.symlink_to(real)
res = archive_sweep(root, workers=2)
self.assertGreaterEqual(res.stats.skipped_symlinks, 1)
self.assertEqual(res.groups, [])
def test_hardlinks_count_once(self):
with tempfile.TemporaryDirectory() as td:
root = Path(td)
a = self._write(root, "a.bin", b"hardlink-payload")
b = root / "b.bin"
os.link(a, b)
res = archive_sweep(root, workers=2)
self.assertGreaterEqual(res.stats.skipped_hardlink_extras, 1)
self.assertEqual(res.groups, [])
def test_empty_files_group(self):
with tempfile.TemporaryDirectory() as td:
root = Path(td)
self._write(root, "e1", b"")
self._write(root, "e2", b"")
res = archive_sweep(root, workers=1)
self.assertEqual(len(res.groups), 1)
self.assertEqual(len(res.groups[0]), 2)
def test_helpers(self):
with tempfile.TemporaryDirectory() as td:
root = Path(td)
p = self._write(root, "x.bin", b"abc" * 100)
self.assertEqual(sample_prefix(str(p), 3), b"abc")
d = sha256_stream(str(p))
self.assertEqual(len(d), 64)
self.assertTrue(byte_equal(str(p), str(p)))
if __name__ == "__main__":
unittest.main()
Concepts used, learn more
Read the underlying idea on its own study page. This lesson applies it. It does not replace those pages.
- Hashing, Frequency Maps & Counting — Two Sum Family & Anagrams
- Mutexes, Condition Variables, Deadlocks & Happens-Before
- Atomics vs Locks
- Mutex vs RWLock
- When Locks Win — Contention, Fairness & Hybrid Designs
- Consistent Hashing: Rings, Virtual Nodes & Replica Placement
- Low-Level Design Under Time — Interfaces, State & Tradeoffs
Pitfalls
- Closing files inside the thread pool without
with. - Sorting paths after mutation mid-iteration.
- Assuming SHA-256 equality without
byte_equal.
Interview Q&A
Why ThreadPoolExecutor?
Answer
Sample and hash are independent reads. The pool overlaps that I/O. The parent builds buckets after join.
Why a lock around stats?
Answer
Workers share SweepStats and the error list. The file bytes are not shared.
Why check the symlink before resolve?
Answer
resolve would turn the link into the target. The skip has to happen first.
Why is workers=1 still correct?
Answer
The filters do not depend on completion order. Groups are sorted after the work finishes.
What does test_same_prefix_different_tail prove?
Answer
A 4KiB shared prefix is not a duplicate when the tail differs.
What does test_empty_files_group prove?
Answer
Two empty files are a real group. Size zero is not a reason to drop them.
Related
The series pager also walks these pages.