Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
59 changes: 34 additions & 25 deletions rdsa_utils/cdp/helpers/hdfs_utils.py
Original file line number Diff line number Diff line change
Expand Up @@ -181,18 +181,20 @@ def create_txt_from_string(
msg,
)

subprocess.call(
[f'echo "{string_to_write}" | hadoop fs -put - {path}'],
shell=True,
# SECURITY FIX: Use list-based subprocess with stdin pipe instead of
# shell=True to prevent shell injection in string_to_write or path.
proc = subprocess.Popen(
["hadoop", "fs", "-put", "-", path],
stdin=subprocess.PIPE,
stdout=subprocess.PIPE,
stderr=subprocess.PIPE,
)
proc.communicate(input=string_to_write.encode("utf-8"), timeout=15)


def delete_dir(path: str) -> bool:
"""Delete an empty directory from HDFS.

This function attempts to delete an empty directory in HDFS.
If the directory is not empty, the deletion will fail.

Parameters
----------
path
Expand All @@ -217,10 +219,6 @@ def delete_dir(path: str) -> bool:
def delete_file(path: str) -> bool:
"""Delete a specific file in HDFS.

This function is used to delete a single file located
at the specified HDFS path. If the path points to a
directory, the command will fail.

Parameters
----------
path
Expand Down Expand Up @@ -249,11 +247,6 @@ def delete_file(path: str) -> bool:
def delete_path(path: str) -> bool:
"""Delete a file or directory in HDFS, including non-empty directories.

This function is capable of deleting both files and directories.
When applied to directories, it will recursively delete all contents
within the directory, making it suitable for removing directories regardless
of whether they are empty or contain files or other directories.

Parameters
----------
path
Expand Down Expand Up @@ -314,12 +307,15 @@ def get_date_modified(filepath: str) -> str:
str
The date the file was last modified.
"""
command = subprocess.Popen(
f"hadoop fs -stat %y {filepath}",
# SECURITY FIX: Use list-based subprocess instead of shell=True to
# prevent shell injection in filepath.
proc = subprocess.Popen(
["hadoop", "fs", "-stat", "%y", filepath],
stdout=subprocess.PIPE,
shell=True,
stderr=subprocess.PIPE,
)
return command.stdout.read().decode("utf-8")[0:10]
stdout, _ = proc.communicate(timeout=15)
return stdout.decode("utf-8")[0:10]


def is_dir(path: str) -> bool:
Expand Down Expand Up @@ -413,16 +409,29 @@ def read_dir_files_recursive(path: str, return_path: bool = True) -> List[str]:
List[str]
A list of files in the directory.
"""
command = subprocess.Popen(
f"hadoop fs -ls -R {path} | grep -v ^d | tr -s ' ' | cut -d ' ' -f 8-",
# SECURITY FIX: Use list-based subprocess instead of shell=True to
# prevent shell injection. Parse output in Python instead of piping
# through grep/tr/cut.
proc = subprocess.Popen(
["hadoop", "fs", "-ls", "-R", path],
stdout=subprocess.PIPE,
shell=True,
stderr=subprocess.PIPE,
)
object_list = [obj.decode("utf-8") for obj in command.stdout.read().splitlines()]
stdout, _ = proc.communicate(timeout=15)
lines = stdout.decode("utf-8").splitlines()

object_list = []
for line in lines:
if not line or line.startswith("Found"):
continue
parts = line.split()
# Hadoop -ls -R format: permission replicas user group size date time path
# Lines starting with 'd' are directories, skip them
if len(parts) >= 8 and not parts[0].startswith("d"):
object_list.append(parts[-1]) # Last field is the path

if not return_path:
return [Path(path).name for path in object_list]

return [Path(p).name for p in object_list]
else:
return object_list

Expand Down
100 changes: 63 additions & 37 deletions tests/cdp/helpers/test_hdfs_utils.py
Original file line number Diff line number Diff line change
Expand Up @@ -189,7 +189,6 @@ def test_copy_local_to_hdfs(self, mock_subprocess_popen):
"hadoop",
"fs",
"-copyFromLocal",
"-f",
overwrite_from_path,
overwrite_to_path,
]
Expand All @@ -206,7 +205,7 @@ def test_create_dir(self, mock_subprocess_popen):

Checks if the command is correctly constructed based on the provided path.
"""
# Test case 1: Test create_dir with a valid path
# Test case: Test create_dir with a valid path
path = "/user/new_directory"
command = ["hadoop", "fs", "-mkdir", path]
assert create_dir(path) == _perform(command)
Expand All @@ -216,19 +215,19 @@ class TestCreateTxtFromString:
"""Tests for create_txt_from_string function."""

@pytest.mark.parametrize(
("path", "string_to_write", "replace", "expected_call"),
("path", "string_to_write", "replace", "expected_called"),
[
(
"/some/directory/newfile.txt",
"Hello, world!",
False,
['echo "Hello, world!" | hadoop fs -put - /some/directory/newfile.txt'],
True,
),
(
"/some/directory/newfile.txt",
"Hello, world!",
True,
['echo "Hello, world!" | hadoop fs -put - /some/directory/newfile.txt'],
True,
),
],
)
Expand All @@ -237,10 +236,10 @@ def test_create_txt_from_string(
path,
string_to_write,
replace,
expected_call,
expected_called,
):
"""Verify 'echo | hadoop fs -put -' command execution by create_txt_from_string."""
with patch("subprocess.call") as subprocess_mock, patch(
"""Verify subprocess.Popen is called with correct args for create_txt_from_string."""
with patch("subprocess.Popen") as subprocess_mock, patch(
"rdsa_utils.cdp.helpers.hdfs_utils.file_exists",
) as file_exists_mock, patch(
"rdsa_utils.cdp.helpers.hdfs_utils.delete_file",
Expand All @@ -249,12 +248,23 @@ def test_create_txt_from_string(
replace # Assume file exists if replace is True
)

if expected_call:
# Test if subprocess.call is called correctly
proc_mock = MagicMock()
subprocess_mock.return_value = proc_mock

if expected_called:
# Test if subprocess.Popen is called correctly
create_txt_from_string(path, string_to_write, replace)
subprocess_mock.assert_called_with(expected_call, shell=True)
subprocess_mock.assert_called_with(
["hadoop", "fs", "-put", "-", path],
stdin=subprocess.PIPE,
stdout=subprocess.PIPE,
stderr=subprocess.PIPE,
)
proc_mock.communicate.assert_called_with(
input=string_to_write.encode("utf-8"),
timeout=15,
)
else:
# Test if FileNotFoundError is raised
with pytest.raises(FileNotFoundError) as excinfo:
create_txt_from_string(path, string_to_write, replace)
assert (
Expand Down Expand Up @@ -282,7 +292,7 @@ def test_delete_dir(self, mock_subprocess_popen):

Checks if the command is correctly constructed based on the provided path.
"""
# Test case 1: Test delete_dir with a valid path
# Test case: Test delete_dir with a valid path
path = "/user/directory"
command = ["hadoop", "fs", "-rmdir", path]
assert delete_dir(path) == _perform(command)
Expand All @@ -296,7 +306,7 @@ def test_delete_file(self, mock_subprocess_popen):

Checks if the command is correctly constructed based on the provided path.
"""
# Test case 1: Test delete_file with a valid path
# Test case: Test delete_file with a valid path
path = "/user/file.txt"
command = ["hadoop", "fs", "-rm", path]
assert delete_file(path) == _perform(command)
Expand Down Expand Up @@ -355,13 +365,15 @@ def test_get_date_modified(self, mock_subprocess_popen_date_modifed):
"""
# Test case: Test get_date_modified with a valid path
filepath = "/user/file.txt"
command_mock = mock_subprocess_popen_date_modifed.return_value
stdout_mock = command_mock.stdout
stdout_mock.read.return_value.decode.return_value.__getitem__.return_value = (
"2023-05-25"
)
mock_popen = mock_subprocess_popen_date_modifed
mock_popen.return_value.communicate.return_value = (b"2023-05-25", b"")
expected_output = "2023-05-25"
assert get_date_modified(filepath) == expected_output
mock_popen.assert_called_with(
["hadoop", "fs", "-stat", "%y", filepath],
stdout=subprocess.PIPE,
stderr=subprocess.PIPE,
)


class TestIsDir(BaseTest):
Expand Down Expand Up @@ -420,7 +432,7 @@ def test_read_dir(self, mock_subprocess_popen):

Checks if the command is correctly constructed based on the provided path.
"""
# Test case 1: Test read_dir with a valid path
# Test case: Test read_dir with a valid path
path = "/user/directory"
ls = subprocess.Popen(["hadoop", "fs", "-ls", path], stdout=subprocess.PIPE)
expected_files = [
Expand All @@ -436,7 +448,7 @@ class TestReadDirFiles(BaseTest):

def test_read_dir_files(self, mock_subprocess_popen):
"""Verify proper extraction of filenames from paths by the read_dir_files function."""
# Test case 1: Test read_dir_files with a valid path
# Test case: Test read_dir_files with a valid path
path = "/user/directory"
expected_files = [Path(p).name for p in read_dir(path)]
assert read_dir_files(path) == expected_files
Expand All @@ -450,24 +462,38 @@ def test_read_dir_files_recursive(self, mock_subprocess_popen):

Checks if the command is correctly constructed based on the provided path.
"""
# Test case 1: Test read_dir_files_recursive without return_path option
# Test case: Test read_dir_files_recursive with a valid path
path = "/user/directory"
command = subprocess.Popen(
f"hadoop fs -ls -R {path} | grep -v ^d | tr -s ' ' | cut -d ' ' -f 8-",
stdout=subprocess.PIPE,
shell=True,
result = read_dir_files_recursive(path)
assert isinstance(result, list)

def test_read_dir_files_recursive_parses_output(self):
"""Verify read_dir_files_recursive parses hadoop output correctly."""
fake_output = (
"drwxr-xr-x - user group 0 2024-01-01 12:00 /user/dir/subdir\n"
"-rw-r--r-- - user group 1234 2024-01-01 12:00 /user/dir/file1.txt\n"
"-rw-r--r-- - user group 567 2024-01-01 12:00 /user/dir/file2.txt\n"
)
expected_files = [
obj.decode("utf-8") for obj in command.stdout.read().splitlines()
]
assert read_dir_files_recursive(path) == expected_files

# Test case 2: Test read_dir_files_recursive with return_path option
return_path = True
return_path_files = [
obj.decode("utf-8") for obj in command.stdout.read().splitlines()
]
assert read_dir_files_recursive(path, return_path) == return_path_files
with patch("subprocess.Popen") as mock_popen:
proc_mock = MagicMock()
proc_mock.communicate.return_value = (fake_output.encode("utf-8"), b"")
proc_mock.returncode = 0
mock_popen.return_value = proc_mock

result = read_dir_files_recursive("/user/dir")

# Should only contain file paths, not directories
assert "/user/dir/file1.txt" in result
assert "/user/dir/file2.txt" in result
assert "/user/dir/subdir" not in result
assert len(result) == 2

# Verify the command is list-based, not shell
mock_popen.assert_called_with(
["hadoop", "fs", "-ls", "-R", "/user/dir"],
stdout=subprocess.PIPE,
stderr=subprocess.PIPE,
)


class TestRename(BaseTest):
Expand Down