From 8d712b70ac8e71c9fe40f85d0b8ce99f3ed67d91 Mon Sep 17 00:00:00 2001 From: Eden Zhang Date: Thu, 20 Mar 2025 17:46:46 +0000 Subject: [PATCH 01/25] Change S3Storage class to S3OutputStorage --- components/clp-py-utils/clp_py_utils/clp_config.py | 14 +++++++------- 1 file changed, 7 insertions(+), 7 deletions(-) diff --git a/components/clp-py-utils/clp_py_utils/clp_config.py b/components/clp-py-utils/clp_py_utils/clp_config.py index 1fbf5cbe63..7f2c07b8a9 100644 --- a/components/clp-py-utils/clp_py_utils/clp_config.py +++ b/components/clp-py-utils/clp_py_utils/clp_config.py @@ -379,7 +379,7 @@ def dump_to_primitive_dict(self): return d -class S3Storage(BaseModel): +class S3OutputStorage(BaseModel): type: Literal[StorageType.S3.value] = StorageType.S3.value staging_directory: pathlib.Path s3_config: S3Config @@ -407,15 +407,15 @@ class StreamFsStorage(FsStorage): directory: pathlib.Path = CLP_DEFAULT_DATA_DIRECTORY_PATH / "streams" -class ArchiveS3Storage(S3Storage): +class ArchiveS3OutputStorage(S3OutputStorage): staging_directory: pathlib.Path = CLP_DEFAULT_DATA_DIRECTORY_PATH / "staged-archives" -class StreamS3Storage(S3Storage): +class StreamS3OutputStorage(S3OutputStorage): staging_directory: pathlib.Path = CLP_DEFAULT_DATA_DIRECTORY_PATH / "staged-streams" -def _get_directory_from_storage_config(storage_config: Union[FsStorage, S3Storage]) -> pathlib.Path: +def _get_directory_from_storage_config(storage_config: Union[FsStorage, S3OutputStorage]) -> pathlib.Path: storage_type = storage_config.type if StorageType.FS == storage_type: return storage_config.directory @@ -426,7 +426,7 @@ def _get_directory_from_storage_config(storage_config: Union[FsStorage, S3Storag def _set_directory_for_storage_config( - storage_config: Union[FsStorage, S3Storage], directory + storage_config: Union[FsStorage, S3OutputStorage], directory ) -> None: storage_type = storage_config.type if StorageType.FS == storage_type: @@ -438,7 +438,7 @@ def _set_directory_for_storage_config( class ArchiveOutput(BaseModel): - storage: Union[ArchiveFsStorage, ArchiveS3Storage] = ArchiveFsStorage() + storage: Union[ArchiveFsStorage, ArchiveS3OutputStorage] = ArchiveFsStorage() target_archive_size: int = 256 * 1024 * 1024 # 256 MB target_dictionaries_size: int = 32 * 1024 * 1024 # 32 MB target_encoded_file_size: int = 256 * 1024 * 1024 # 256 MB @@ -488,7 +488,7 @@ def dump_to_primitive_dict(self): class StreamOutput(BaseModel): - storage: Union[StreamFsStorage, StreamS3Storage] = StreamFsStorage() + storage: Union[StreamFsStorage, StreamS3OutputStorage] = StreamFsStorage() target_uncompressed_size: int = 128 * 1024 * 1024 @validator("target_uncompressed_size") From 88b0e2837512b25e6c4e2bfc0a57a15fcc277259 Mon Sep 17 00:00:00 2001 From: Eden Zhang Date: Thu, 20 Mar 2025 17:49:37 +0000 Subject: [PATCH 02/25] Make S3InputStorage class --- components/clp-py-utils/clp_py_utils/clp_config.py | 9 +++++++++ 1 file changed, 9 insertions(+) diff --git a/components/clp-py-utils/clp_py_utils/clp_config.py b/components/clp-py-utils/clp_py_utils/clp_config.py index 7f2c07b8a9..686f8efcf3 100644 --- a/components/clp-py-utils/clp_py_utils/clp_config.py +++ b/components/clp-py-utils/clp_py_utils/clp_config.py @@ -379,6 +379,15 @@ def dump_to_primitive_dict(self): return d +class S3InputStorage(BaseModel): + type: Literal[StorageType.S3.value] = StorageType.S3.value + s3_config: S3Config + + def dump_to_primitive_dict(self): + d = self.dict() + return d + + class S3OutputStorage(BaseModel): type: Literal[StorageType.S3.value] = StorageType.S3.value staging_directory: pathlib.Path From 64caa6f7f5d7c0067c5300cd5839fc063d0459e4 Mon Sep 17 00:00:00 2001 From: Eden Zhang Date: Thu, 20 Mar 2025 20:08:19 +0000 Subject: [PATCH 03/25] Change `input_logs_directory` to `logs_input` --- .../clp_package_utils/general.py | 18 +++++++------ .../scripts/native/compress.py | 2 +- .../clp-py-utils/clp_py_utils/clp_config.py | 25 +++++++++++-------- .../package-template/src/etc/clp-config.yml | 4 ++- 4 files changed, 28 insertions(+), 21 deletions(-) diff --git a/components/clp-package-utils/clp_package_utils/general.py b/components/clp-package-utils/clp_package_utils/general.py index d7800ea214..bb4be0e5b4 100644 --- a/components/clp-package-utils/clp_package_utils/general.py +++ b/components/clp-package-utils/clp_package_utils/general.py @@ -20,6 +20,7 @@ REDIS_COMPONENT_NAME, REDUCER_COMPONENT_NAME, RESULTS_CACHE_COMPONENT_NAME, + StorageType, WEBUI_COMPONENT_NAME, WorkerConfig, ) @@ -216,13 +217,14 @@ def generate_container_config( docker_mounts = CLPDockerMounts(clp_home, CONTAINER_CLP_HOME) - input_logs_dir = clp_config.input_logs_directory.resolve() - container_clp_config.input_logs_directory = ( - CONTAINER_INPUT_LOGS_ROOT_DIR / input_logs_dir.relative_to(input_logs_dir.anchor) - ) - docker_mounts.input_logs_dir = DockerMount( - DockerMountType.BIND, input_logs_dir, container_clp_config.input_logs_directory, True - ) + if StorageType.FS == clp_config.logs_input.type: + input_logs_dir = clp_config.logs_input.directory.resolve() + container_clp_config.logs_input.directory = ( + CONTAINER_INPUT_LOGS_ROOT_DIR / input_logs_dir.relative_to(input_logs_dir.anchor) + ) + docker_mounts.input_logs_dir = DockerMount( + DockerMountType.BIND, input_logs_dir, container_clp_config.logs_input.directory, True + ) container_clp_config.data_directory = CONTAINER_CLP_HOME / "var" / "data" if not is_path_already_mounted( @@ -494,7 +496,7 @@ def validate_results_cache_config( def validate_worker_config(clp_config: CLPConfig): - clp_config.validate_input_logs_dir() + clp_config.validate_logs_input_config() clp_config.validate_archive_output_config() clp_config.validate_stream_output_dir() diff --git a/components/clp-package-utils/clp_package_utils/scripts/native/compress.py b/components/clp-package-utils/clp_package_utils/scripts/native/compress.py index b71907eb26..cef1f60f43 100755 --- a/components/clp-package-utils/clp_package_utils/scripts/native/compress.py +++ b/components/clp-package-utils/clp_package_utils/scripts/native/compress.py @@ -227,7 +227,7 @@ def main(argv): try: config_file_path = pathlib.Path(parsed_args.config) clp_config = load_config_file(config_file_path, default_config_file_path, clp_home) - clp_config.validate_input_logs_dir() + clp_config.validate_logs_input_config() clp_config.validate_logs_dir() except: logger.exception("Failed to load config.") diff --git a/components/clp-py-utils/clp_py_utils/clp_config.py b/components/clp-py-utils/clp_py_utils/clp_config.py index 686f8efcf3..e908eca88b 100644 --- a/components/clp-py-utils/clp_py_utils/clp_config.py +++ b/components/clp-py-utils/clp_py_utils/clp_config.py @@ -408,6 +408,10 @@ def dump_to_primitive_dict(self): return d +class InputFsStorage(FsStorage): + directory: pathlib.Path = pathlib.Path("/") + + class ArchiveFsStorage(FsStorage): directory: pathlib.Path = CLP_DEFAULT_DATA_DIRECTORY_PATH / "archives" @@ -557,7 +561,7 @@ def validate_port(cls, field): class CLPConfig(BaseModel): execution_container: Optional[str] = None - input_logs_directory: pathlib.Path = pathlib.Path("/") + logs_input: Union[InputFsStorage, S3InputStorage] = InputFsStorage() package: Package = Package() database: Database = Database() @@ -581,7 +585,8 @@ class CLPConfig(BaseModel): _os_release_file_path: pathlib.Path = PrivateAttr(default=OS_RELEASE_FILE_PATH) def make_config_paths_absolute(self, clp_home: pathlib.Path): - self.input_logs_directory = make_config_path_absolute(clp_home, self.input_logs_directory) + if StorageType.FS == self.logs_input.type: + self.logs_input.make_config_paths_absolute(clp_home) self.credentials_file_path = make_config_path_absolute(clp_home, self.credentials_file_path) self.archive_output.storage.make_config_paths_absolute(clp_home) self.stream_output.storage.make_config_paths_absolute(clp_home) @@ -589,14 +594,12 @@ def make_config_paths_absolute(self, clp_home: pathlib.Path): self.logs_directory = make_config_path_absolute(clp_home, self.logs_directory) self._os_release_file_path = make_config_path_absolute(clp_home, self._os_release_file_path) - def validate_input_logs_dir(self): - # NOTE: This can't be a pydantic validator since input_logs_dir might be a package-relative - # path that will only be resolved after pydantic validation - input_logs_dir = self.input_logs_directory - if not input_logs_dir.exists(): - raise ValueError(f"input_logs_directory '{input_logs_dir}' doesn't exist.") - if not input_logs_dir.is_dir(): - raise ValueError(f"input_logs_directory '{input_logs_dir}' is not a directory.") + def validate_logs_input_config(self): + if StorageType.FS == self.logs_input.type: + try: + validate_path_could_be_dir(self.logs_input.directory) + except ValueError as ex: + raise ValueError(f"logs_input directory is invalid: {ex}") def validate_archive_output_config(self): if ( @@ -684,10 +687,10 @@ def load_redis_credentials_from_file(self): def dump_to_primitive_dict(self): d = self.dict() + d["logs_input"] = self.logs_input.dump_to_primitive_dict() d["archive_output"] = self.archive_output.dump_to_primitive_dict() d["stream_output"] = self.stream_output.dump_to_primitive_dict() # Turn paths into primitive strings - d["input_logs_directory"] = str(self.input_logs_directory) d["credentials_file_path"] = str(self.credentials_file_path) d["data_directory"] = str(self.data_directory) d["logs_directory"] = str(self.logs_directory) diff --git a/components/package-template/src/etc/clp-config.yml b/components/package-template/src/etc/clp-config.yml index 3e86199353..71aa5f47a9 100644 --- a/components/package-template/src/etc/clp-config.yml +++ b/components/package-template/src/etc/clp-config.yml @@ -2,7 +2,9 @@ ## workers. ## - This path will be exposed inside the container, so symbolic links to files ## outside this path will be ignored. -#input_logs_directory: "/" +#logs_input: +# type: "fs" +# directory: "/" # ## File containing credentials for services #credentials_file_path: "etc/credentials.yml" From d043d94bfd3b47da85017f9a741ed74c0c170cf0 Mon Sep 17 00:00:00 2001 From: Eden Zhang Date: Fri, 21 Mar 2025 20:55:06 +0000 Subject: [PATCH 04/25] Change mount for input_logs_dir to only mount when type FS --- .../clp_package_utils/scripts/start_clp.py | 14 ++++++++------ 1 file changed, 8 insertions(+), 6 deletions(-) diff --git a/components/clp-package-utils/clp_package_utils/scripts/start_clp.py b/components/clp-package-utils/clp_package_utils/scripts/start_clp.py index 3e751c6f08..e780d20650 100755 --- a/components/clp-package-utils/clp_package_utils/scripts/start_clp.py +++ b/components/clp-package-utils/clp_package_utils/scripts/start_clp.py @@ -597,15 +597,16 @@ def generic_start_scheduler( "--mount", str(mounts.clp_home), ] # fmt: on - necessary_mounts = [ - mounts.logs_dir, - ] - if COMPRESSION_SCHEDULER_COMPONENT_NAME == component_name: + necessary_mounts = [mounts.logs_dir] + if ( + COMPRESSION_SCHEDULER_COMPONENT_NAME == component_name + and StorageType.FS == clp_config.logs_input.type + ): necessary_mounts.append(mounts.input_logs_dir) for mount in necessary_mounts: if mount: container_start_cmd.append("--mount") - container_start_cmd.append(str(mount)) + container_start_cmd.append(str(mount)) container_start_cmd.append(clp_config.execution_container) # fmt: off @@ -741,10 +742,11 @@ def generic_start_worker( mounts.clp_home, mounts.data_dir, mounts.logs_dir, - mounts.input_logs_dir, ] if worker_specific_mount: necessary_mounts.extend(worker_specific_mount) + if StorageType.FS == clp_config.logs_input.type: + necessary_mounts.append(mounts.input_logs_dir) for mount in necessary_mounts: if not mount: From f5c54b8ad8595df6c075fc7530b51c6769e4a6fe Mon Sep 17 00:00:00 2001 From: Eden Zhang Date: Mon, 24 Mar 2025 15:43:13 +0000 Subject: [PATCH 05/25] Whitespace --- .../clp-package-utils/clp_package_utils/scripts/start_clp.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/components/clp-package-utils/clp_package_utils/scripts/start_clp.py b/components/clp-package-utils/clp_package_utils/scripts/start_clp.py index e780d20650..790149c37d 100755 --- a/components/clp-package-utils/clp_package_utils/scripts/start_clp.py +++ b/components/clp-package-utils/clp_package_utils/scripts/start_clp.py @@ -606,7 +606,7 @@ def generic_start_scheduler( for mount in necessary_mounts: if mount: container_start_cmd.append("--mount") - container_start_cmd.append(str(mount)) + container_start_cmd.append(str(mount)) container_start_cmd.append(clp_config.execution_container) # fmt: off From f0d45b29808c19f9899d30217b22bd27c4389535 Mon Sep 17 00:00:00 2001 From: Eden Zhang Date: Mon, 24 Mar 2025 15:45:43 +0000 Subject: [PATCH 06/25] Shorten class names --- components/clp-py-utils/clp_py_utils/clp_config.py | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/components/clp-py-utils/clp_py_utils/clp_config.py b/components/clp-py-utils/clp_py_utils/clp_config.py index e908eca88b..e5afb19cb4 100644 --- a/components/clp-py-utils/clp_py_utils/clp_config.py +++ b/components/clp-py-utils/clp_py_utils/clp_config.py @@ -420,11 +420,11 @@ class StreamFsStorage(FsStorage): directory: pathlib.Path = CLP_DEFAULT_DATA_DIRECTORY_PATH / "streams" -class ArchiveS3OutputStorage(S3OutputStorage): +class ArchiveS3Storage(S3OutputStorage): staging_directory: pathlib.Path = CLP_DEFAULT_DATA_DIRECTORY_PATH / "staged-archives" -class StreamS3OutputStorage(S3OutputStorage): +class StreamS3Storage(S3OutputStorage): staging_directory: pathlib.Path = CLP_DEFAULT_DATA_DIRECTORY_PATH / "staged-streams" @@ -451,7 +451,7 @@ def _set_directory_for_storage_config( class ArchiveOutput(BaseModel): - storage: Union[ArchiveFsStorage, ArchiveS3OutputStorage] = ArchiveFsStorage() + storage: Union[ArchiveFsStorage, ArchiveS3Storage] = ArchiveFsStorage() target_archive_size: int = 256 * 1024 * 1024 # 256 MB target_dictionaries_size: int = 32 * 1024 * 1024 # 32 MB target_encoded_file_size: int = 256 * 1024 * 1024 # 256 MB @@ -501,7 +501,7 @@ def dump_to_primitive_dict(self): class StreamOutput(BaseModel): - storage: Union[StreamFsStorage, StreamS3OutputStorage] = StreamFsStorage() + storage: Union[StreamFsStorage, StreamS3Storage] = StreamFsStorage() target_uncompressed_size: int = 128 * 1024 * 1024 @validator("target_uncompressed_size") From c9f11b00a0bad98e60ee757e8ed2a983551ec6f9 Mon Sep 17 00:00:00 2001 From: Eden Zhang Date: Mon, 24 Mar 2025 16:01:56 +0000 Subject: [PATCH 07/25] Refactor classes --- .../clp-py-utils/clp_py_utils/clp_config.py | 18 ++++++++---------- 1 file changed, 8 insertions(+), 10 deletions(-) diff --git a/components/clp-py-utils/clp_py_utils/clp_config.py b/components/clp-py-utils/clp_py_utils/clp_config.py index e5afb19cb4..4c9c285b2e 100644 --- a/components/clp-py-utils/clp_py_utils/clp_config.py +++ b/components/clp-py-utils/clp_py_utils/clp_config.py @@ -379,7 +379,7 @@ def dump_to_primitive_dict(self): return d -class S3InputStorage(BaseModel): +class InputS3Storage(BaseModel): type: Literal[StorageType.S3.value] = StorageType.S3.value s3_config: S3Config @@ -388,10 +388,8 @@ def dump_to_primitive_dict(self): return d -class S3OutputStorage(BaseModel): - type: Literal[StorageType.S3.value] = StorageType.S3.value +class OutputS3Storage(InputS3Storage): staging_directory: pathlib.Path - s3_config: S3Config @validator("staging_directory") def validate_staging_directory(cls, field): @@ -403,7 +401,7 @@ def make_config_paths_absolute(self, clp_home: pathlib.Path): self.staging_directory = make_config_path_absolute(clp_home, self.staging_directory) def dump_to_primitive_dict(self): - d = self.dict() + d = super().dump_to_primitive_dict() d["staging_directory"] = str(d["staging_directory"]) return d @@ -420,15 +418,15 @@ class StreamFsStorage(FsStorage): directory: pathlib.Path = CLP_DEFAULT_DATA_DIRECTORY_PATH / "streams" -class ArchiveS3Storage(S3OutputStorage): +class ArchiveS3Storage(OutputS3Storage): staging_directory: pathlib.Path = CLP_DEFAULT_DATA_DIRECTORY_PATH / "staged-archives" -class StreamS3Storage(S3OutputStorage): +class StreamS3Storage(OutputS3Storage): staging_directory: pathlib.Path = CLP_DEFAULT_DATA_DIRECTORY_PATH / "staged-streams" -def _get_directory_from_storage_config(storage_config: Union[FsStorage, S3OutputStorage]) -> pathlib.Path: +def _get_directory_from_storage_config(storage_config: Union[FsStorage, OutputS3Storage]) -> pathlib.Path: storage_type = storage_config.type if StorageType.FS == storage_type: return storage_config.directory @@ -439,7 +437,7 @@ def _get_directory_from_storage_config(storage_config: Union[FsStorage, S3Output def _set_directory_for_storage_config( - storage_config: Union[FsStorage, S3OutputStorage], directory + storage_config: Union[FsStorage, OutputS3Storage], directory ) -> None: storage_type = storage_config.type if StorageType.FS == storage_type: @@ -561,7 +559,7 @@ def validate_port(cls, field): class CLPConfig(BaseModel): execution_container: Optional[str] = None - logs_input: Union[InputFsStorage, S3InputStorage] = InputFsStorage() + logs_input: Union[InputFsStorage, InputS3Storage] = InputFsStorage() package: Package = Package() database: Database = Database() From 38276d3acfbea6b96b3f9dbeebe647a3711c1e13 Mon Sep 17 00:00:00 2001 From: Eden Zhang Date: Mon, 24 Mar 2025 18:28:54 +0000 Subject: [PATCH 08/25] Change InputS3Storage to S3Storage --- components/clp-py-utils/clp_py_utils/clp_config.py | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/components/clp-py-utils/clp_py_utils/clp_config.py b/components/clp-py-utils/clp_py_utils/clp_config.py index 4c9c285b2e..4a80e031a7 100644 --- a/components/clp-py-utils/clp_py_utils/clp_config.py +++ b/components/clp-py-utils/clp_py_utils/clp_config.py @@ -379,7 +379,7 @@ def dump_to_primitive_dict(self): return d -class InputS3Storage(BaseModel): +class S3Storage(BaseModel): type: Literal[StorageType.S3.value] = StorageType.S3.value s3_config: S3Config @@ -388,7 +388,7 @@ def dump_to_primitive_dict(self): return d -class OutputS3Storage(InputS3Storage): +class OutputS3Storage(S3Storage): staging_directory: pathlib.Path @validator("staging_directory") @@ -559,7 +559,7 @@ def validate_port(cls, field): class CLPConfig(BaseModel): execution_container: Optional[str] = None - logs_input: Union[InputFsStorage, InputS3Storage] = InputFsStorage() + logs_input: Union[InputFsStorage, S3Storage] = InputFsStorage() package: Package = Package() database: Database = Database() From aa16f17fe66d29aafe711f125f2edc8d85b5dfcd Mon Sep 17 00:00:00 2001 From: Eden Zhang Date: Tue, 25 Mar 2025 21:03:10 +0000 Subject: [PATCH 09/25] Change S3 compression to use logs_input s3_config instead of URL --- .../clp_package_utils/scripts/compress.py | 92 ++++++------------- .../scripts/native/compress.py | 65 +++++-------- .../clp-py-utils/clp_py_utils/clp_config.py | 4 +- 3 files changed, 53 insertions(+), 108 deletions(-) diff --git a/components/clp-package-utils/clp_package_utils/scripts/compress.py b/components/clp-package-utils/clp_package_utils/scripts/compress.py index 2829076e5b..10c2c133d0 100755 --- a/components/clp-package-utils/clp_package_utils/scripts/compress.py +++ b/components/clp-package-utils/clp_package_utils/scripts/compress.py @@ -63,8 +63,9 @@ def _parse_aws_credentials_file(credentials_file_path: pathlib.Path, user: str) def _generate_logs_list( container_logs_list_path: pathlib.Path, parsed_args: argparse.Namespace, + clp_config: CLPConfig, ) -> None: - input_type = parsed_args.input_type + input_type = clp_config.logs_input.type if InputType.FS == input_type: host_logs_list_path = parsed_args.path_list @@ -91,23 +92,25 @@ def _generate_logs_list( elif InputType.S3 == input_type: with open(container_logs_list_path, "w") as container_logs_list_file: - container_logs_list_file.write(f"{parsed_args.url}\n") + container_logs_list_file.write(f"{parsed_args.paths[0]}\n") else: raise ValueError(f"Unsupported input type: {input_type}.") def _generate_compress_cmd( - parsed_args: argparse.Namespace, config_path: pathlib.Path, logs_list_path: pathlib.Path + parsed_args: argparse.Namespace, + config_path: pathlib.Path, + logs_list_path: pathlib.Path, + clp_config: CLPConfig, ) -> List[str]: - input_type = parsed_args.input_type + input_type = clp_config.logs_input.type # fmt: off compress_cmd = [ "python3", "-m", "clp_package_utils.scripts.native.compress", "--config", str(config_path), - input_type, ] # fmt: on if parsed_args.timestamp_key is not None: @@ -119,24 +122,6 @@ def _generate_compress_cmd( if parsed_args.no_progress_reporting is True: compress_cmd.append("--no-progress-reporting") - if InputType.FS == input_type: - pass - elif InputType.S3 == input_type: - aws_access_key_id = parsed_args.aws_access_key_id - aws_secret_access_key = parsed_args.aws_secret_access_key - if parsed_args.aws_credentials_file: - default_credentials_user = "default" - aws_access_key_id, aws_secret_access_key = _parse_aws_credentials_file( - pathlib.Path(parsed_args.aws_credentials_file), default_credentials_user - ) - if bool(aws_access_key_id) and bool(aws_secret_access_key): - compress_cmd.append("--aws-access-key-id") - compress_cmd.append(aws_access_key_id) - compress_cmd.append("--aws-secret-access-key") - compress_cmd.append(aws_secret_access_key) - else: - raise ValueError(f"Unsupported input type: {input_type}.") - compress_cmd.append("--logs-list") compress_cmd.append(str(logs_list_path)) @@ -177,26 +162,10 @@ def _validate_s3_input_args( f"Input type {InputType.S3} is only supported for the storage engine" f" {StorageEngine.CLP_S}." ) - - # Validate aws credentials were specified using only one method - aws_credential_file = parsed_args.aws_credentials_file - aws_access_key_id = parsed_args.aws_access_key_id - aws_secret_access_key = parsed_args.aws_secret_access_key - if aws_credential_file is not None: - if not pathlib.Path(aws_credential_file).exists(): - args_parser.error(f"AWS credentials file '{aws_credential_file}' doesn't exist.") - - if aws_access_key_id is not None or aws_secret_access_key is not None: - args_parser.error( - "aws_credentials_file cannot be specified together with aws_access_key_id or" - " aws_secret_access_key." - ) - - else: - if not bool(aws_access_key_id): - args_parser.error("aws_access_key_id not specified or empty") - if not bool(aws_secret_access_key): - args_parser.error("aws_secret_access_key not specified or empty") + if len(parsed_args.paths) != 1: + args_parser.error(f"Only one key prefix can be specified for input type {InputType.S3}.") + if parsed_args.path_list is not None: + args_parser.error(f"Path list file is not supported for input type {InputType.S3}.") def main(argv): @@ -212,26 +181,19 @@ def main(argv): default=str(default_config_file_path), help="CLP package configuration file.", ) - input_type_args_parser = args_parser.add_subparsers(dest="input_type") - - fs_compressor_parser = input_type_args_parser.add_parser(InputType.FS) - _add_common_arguments(fs_compressor_parser) - fs_compressor_parser.add_argument("paths", metavar="PATH", nargs="*", help="Paths to compress.") - fs_compressor_parser.add_argument( - "-f", "--path-list", dest="path_list", help="A file listing all paths to compress." + args_parser.add_argument( + "--timestamp-key", + help="The path (e.g. x.y) for the field containing the log event's timestamp.", ) - - s3_compressor_parser = input_type_args_parser.add_parser(InputType.S3) - _add_common_arguments(s3_compressor_parser) - s3_compressor_parser.add_argument("url", metavar="URL", help="URL of objects to be compressed") - s3_compressor_parser.add_argument( - "--aws-access-key-id", type=str, default=None, help="AWS access key ID." + args_parser.add_argument( + "-t", "--tags", help="A comma-separated list of tags to apply to the compressed archives." ) - s3_compressor_parser.add_argument( - "--aws-secret-access-key", type=str, default=None, help="AWS secret access key." + args_parser.add_argument( + "--no-progress-reporting", action="store_true", help="Disables progress reporting." ) - s3_compressor_parser.add_argument( - "--aws-credentials-file", type=str, default=None, help="Path to AWS credentials file." + args_parser.add_argument("paths", metavar="PATH", nargs="*", help="Paths to compress.") + args_parser.add_argument( + "-f", "--path-list", dest="path_list", help="A file listing all paths to compress." ) parsed_args = args_parser.parse_args(argv[1:]) @@ -248,7 +210,7 @@ def main(argv): logger.exception("Failed to load config.") return -1 - input_type = parsed_args.input_type + input_type = clp_config.logs_input.type if InputType.FS == input_type: _validate_fs_input_args(parsed_args, args_parser) elif InputType.S3 == input_type: @@ -263,7 +225,9 @@ def main(argv): container_clp_config, clp_config, container_name ) - necessary_mounts = [mounts.clp_home, mounts.input_logs_dir, mounts.data_dir, mounts.logs_dir] + necessary_mounts = [mounts.clp_home, mounts.data_dir, mounts.logs_dir] + if InputType.FS == input_type: + necessary_mounts.append(mounts.input_logs_dir) # Write compression logs to a file while True: @@ -276,13 +240,13 @@ def main(argv): if not container_logs_list_path.exists(): break - _generate_logs_list(container_logs_list_path, parsed_args) + _generate_logs_list(container_logs_list_path, parsed_args, clp_config) container_start_cmd = generate_container_start_cmd( container_name, necessary_mounts, clp_config.execution_container ) compress_cmd = _generate_compress_cmd( - parsed_args, generated_config_path_on_container, logs_list_path_on_container + parsed_args, generated_config_path_on_container, logs_list_path_on_container, clp_config ) cmd = container_start_cmd + compress_cmd subprocess.run(cmd, check=True) diff --git a/components/clp-package-utils/clp_package_utils/scripts/native/compress.py b/components/clp-package-utils/clp_package_utils/scripts/native/compress.py index cef1f60f43..1b229a6132 100755 --- a/components/clp-package-utils/clp_package_utils/scripts/native/compress.py +++ b/components/clp-package-utils/clp_package_utils/scripts/native/compress.py @@ -10,7 +10,7 @@ import brotli import msgpack -from clp_py_utils.clp_config import COMPRESSION_JOBS_TABLE_NAME, S3Credentials +from clp_py_utils.clp_config import CLPConfig, COMPRESSION_JOBS_TABLE_NAME, S3Credentials from clp_py_utils.pretty_size import pretty_size from clp_py_utils.s3_utils import parse_s3_url from clp_py_utils.sql_adapter import SQL_Adapter @@ -128,9 +128,9 @@ def handle_job(sql_adapter: SQL_Adapter, clp_io_config: ClpIoConfig, no_progress def _generate_clp_io_config( - logs_to_compress: List[str], parsed_args: argparse.Namespace + logs_to_compress: List[str], parsed_args: argparse.Namespace, clp_config: CLPConfig ) -> typing.Union[S3InputConfig, FsInputConfig]: - input_type = parsed_args.input_type + input_type = clp_config.logs_input.type if InputType.FS == input_type: return FsInputConfig( @@ -140,18 +140,14 @@ def _generate_clp_io_config( ) elif InputType.S3 == input_type: if len(logs_to_compress) != 1: - ValueError(f"Too many URLs: {len(logs_to_compress)} > 1") + ValueError(f"Too many prefixes: {len(logs_to_compress)} > 1") - s3_url = logs_to_compress[0] - region_code, bucket_name, key_prefix = parse_s3_url(s3_url) + s3_config = clp_config.logs_input.s3_config return S3InputConfig( - region_code=region_code, - bucket=bucket_name, - key_prefix=key_prefix, - credentials=S3Credentials( - access_key_id=parsed_args.aws_access_key_id, - secret_access_key=parsed_args.aws_secret_access_key, - ), + region_code=s3_config.region_code, + bucket=s3_config.bucket, + key_prefix=s3_config.key_prefix + logs_to_compress[0], + credentials=s3_config.credentials, timestamp_key=parsed_args.timestamp_key, ) else: @@ -175,7 +171,18 @@ def _get_logs_to_compress(logs_list_path: pathlib.Path) -> List[str]: return logs_to_compress -def _add_common_arguments(args_parser: argparse.ArgumentParser) -> None: +def main(argv): + clp_home = get_clp_home() + default_config_file_path = clp_home / CLP_DEFAULT_CONFIG_FILE_RELATIVE_PATH + args_parser = argparse.ArgumentParser(description="Compresses logs") + + # Package-level config option + args_parser.add_argument( + "--config", + "-c", + default=str(default_config_file_path), + help="CLP package configuration file.", + ) args_parser.add_argument( "-f", "--logs-list", @@ -193,34 +200,6 @@ def _add_common_arguments(args_parser: argparse.ArgumentParser) -> None: args_parser.add_argument( "-t", "--tags", help="A comma-separated list of tags to apply to the compressed archives." ) - - -def main(argv): - clp_home = get_clp_home() - default_config_file_path = clp_home / CLP_DEFAULT_CONFIG_FILE_RELATIVE_PATH - args_parser = argparse.ArgumentParser(description="Compresses logs") - - # Package-level config option - args_parser.add_argument( - "--config", - "-c", - default=str(default_config_file_path), - help="CLP package configuration file.", - ) - input_type_args_parser = args_parser.add_subparsers(dest="input_type") - - fs_compressor_parser = input_type_args_parser.add_parser(InputType.FS) - _add_common_arguments(fs_compressor_parser) - - s3_compressor_parser = input_type_args_parser.add_parser(InputType.S3) - _add_common_arguments(s3_compressor_parser) - s3_compressor_parser.add_argument( - "--aws-access-key-id", type=str, default=None, help="AWS access key ID." - ) - s3_compressor_parser.add_argument( - "--aws-secret-access-key", type=str, default=None, help="AWS secret access key." - ) - parsed_args = args_parser.parse_args(argv[1:]) # Validate and load config file @@ -238,7 +217,7 @@ def main(argv): logs_to_compress = _get_logs_to_compress(pathlib.Path(parsed_args.logs_list).resolve()) - clp_input_config = _generate_clp_io_config(logs_to_compress, parsed_args) + clp_input_config = _generate_clp_io_config(logs_to_compress, parsed_args, clp_config) clp_output_config = OutputConfig.parse_obj(clp_config.archive_output) if parsed_args.tags: tag_list = [tag.strip().lower() for tag in parsed_args.tags.split(",") if tag] diff --git a/components/clp-py-utils/clp_py_utils/clp_config.py b/components/clp-py-utils/clp_py_utils/clp_config.py index 4a80e031a7..4a4aa144d5 100644 --- a/components/clp-py-utils/clp_py_utils/clp_config.py +++ b/components/clp-py-utils/clp_py_utils/clp_config.py @@ -426,7 +426,9 @@ class StreamS3Storage(OutputS3Storage): staging_directory: pathlib.Path = CLP_DEFAULT_DATA_DIRECTORY_PATH / "staged-streams" -def _get_directory_from_storage_config(storage_config: Union[FsStorage, OutputS3Storage]) -> pathlib.Path: +def _get_directory_from_storage_config( + storage_config: Union[FsStorage, OutputS3Storage], +) -> pathlib.Path: storage_type = storage_config.type if StorageType.FS == storage_type: return storage_config.directory From 3395d782afec6c722bbc8b282235841b0dfadc74 Mon Sep 17 00:00:00 2001 From: Eden Zhang Date: Wed, 26 Mar 2025 16:19:29 +0000 Subject: [PATCH 10/25] Update error message --- .../scheduler/compress/compression_scheduler.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/components/job-orchestration/job_orchestration/scheduler/compress/compression_scheduler.py b/components/job-orchestration/job_orchestration/scheduler/compress/compression_scheduler.py index 7de797bb01..494c9fb228 100644 --- a/components/job-orchestration/job_orchestration/scheduler/compress/compression_scheduler.py +++ b/components/job-orchestration/job_orchestration/scheduler/compress/compression_scheduler.py @@ -134,13 +134,13 @@ def _process_s3_input( and adds their metadata to paths_to_compress_buffer. :param s3_input_config: :param paths_to_compress_buffer: - :raises: RuntimeError if input URL doesn't resolve to any objects. + :raises: RuntimeError if input path doesn't resolve to any objects. :raises: Propagates `s3_get_object_metadata`'s exceptions. """ object_metadata_list = s3_get_object_metadata(s3_input_config) if len(object_metadata_list) == 0: - raise RuntimeError("Input URL doesn't resolve to any object") + raise RuntimeError("Input path doesn't resolve to any object") for object_metadata in object_metadata_list: paths_to_compress_buffer.add_file(object_metadata) From 042a3924b719f466e4a438e104f7b105bc13b39c Mon Sep 17 00:00:00 2001 From: Eden Zhang Date: Wed, 26 Mar 2025 19:18:54 +0000 Subject: [PATCH 11/25] Remove parse_s3_url function --- .../scripts/native/compress.py | 1 - .../clp-py-utils/clp_py_utils/s3_utils.py | 35 ------------------- 2 files changed, 36 deletions(-) diff --git a/components/clp-package-utils/clp_package_utils/scripts/native/compress.py b/components/clp-package-utils/clp_package_utils/scripts/native/compress.py index 1b229a6132..4b3e91d6f5 100755 --- a/components/clp-package-utils/clp_package_utils/scripts/native/compress.py +++ b/components/clp-package-utils/clp_package_utils/scripts/native/compress.py @@ -12,7 +12,6 @@ import msgpack from clp_py_utils.clp_config import CLPConfig, COMPRESSION_JOBS_TABLE_NAME, S3Credentials from clp_py_utils.pretty_size import pretty_size -from clp_py_utils.s3_utils import parse_s3_url from clp_py_utils.sql_adapter import SQL_Adapter from job_orchestration.scheduler.constants import ( CompressionJobCompletionStatus, diff --git a/components/clp-py-utils/clp_py_utils/s3_utils.py b/components/clp-py-utils/clp_py_utils/s3_utils.py index 12d3755e48..0c893841cd 100644 --- a/components/clp-py-utils/clp_py_utils/s3_utils.py +++ b/components/clp-py-utils/clp_py_utils/s3_utils.py @@ -13,41 +13,6 @@ AWS_ENDPOINT = "amazonaws.com" -def parse_s3_url(s3_url: str) -> Tuple[str, str, str]: - """ - Parses the region_code, bucket, and key_prefix from the given S3 URL. - :param s3_url: A host-style URL or path-style URL. - :return: A tuple of (region_code, bucket, key_prefix). - :raise: ValueError if `s3_url` is not a valid host-style URL or path-style URL. - """ - - host_style_url_regex = re.compile( - r"https://(?P[a-z0-9.-]+)\.s3(\.(?P[a-z0-9-]+))?" - r"\.(?P[a-z0-9.-]+)/(?P[^?]+).*" - ) - match = host_style_url_regex.match(s3_url) - - if match is None: - path_style_url_regex = re.compile( - r"https://s3(\.(?P[a-z0-9-]+))?\.(?P[a-z0-9.-]+)/" - r"(?P[a-z0-9.-]+)/(?P[^?]+).*" - ) - match = path_style_url_regex.match(s3_url) - - if match is None: - raise ValueError(f"Unsupported URL format: {s3_url}") - - region_code = match.group("region_code") - bucket_name = match.group("bucket_name") - endpoint = match.group("endpoint") - key_prefix = match.group("key_prefix") - - if AWS_ENDPOINT != endpoint: - raise ValueError(f"Unsupported endpoint: {endpoint}") - - return region_code, bucket_name, key_prefix - - def generate_s3_virtual_hosted_style_url( region_code: str, bucket_name: str, object_key: str ) -> str: From 61e974cd178773b7014edc8801960e15fdbd91d0 Mon Sep 17 00:00:00 2001 From: Eden Zhang Date: Wed, 26 Mar 2025 20:21:23 +0000 Subject: [PATCH 12/25] Remove unused functions --- .../clp_package_utils/scripts/compress.py | 47 ------------------- 1 file changed, 47 deletions(-) diff --git a/components/clp-package-utils/clp_package_utils/scripts/compress.py b/components/clp-package-utils/clp_package_utils/scripts/compress.py index 10c2c133d0..0605bf657a 100755 --- a/components/clp-package-utils/clp_package_utils/scripts/compress.py +++ b/components/clp-package-utils/clp_package_utils/scripts/compress.py @@ -26,40 +26,6 @@ logger = logging.getLogger(__file__) -def _parse_aws_credentials_file(credentials_file_path: pathlib.Path, user: str) -> Tuple[str, str]: - """ - Parses the `aws_access_key_id` and `aws_secret_access_key` of `user` from the given - credentials_file_path. - :param credentials_file_path: - :param user: - :return: A tuple of (aws_access_key_id, aws_secret_access_key) - :raises: ValueError if the file doesn't exist, or doesn't contain valid aws credentials. - """ - - if not credentials_file_path.exists(): - raise ValueError(f"'{credentials_file_path}' doesn't exist.") - - config_reader = configparser.ConfigParser() - config_reader.read(credentials_file_path) - - if not config_reader.has_section(user): - raise ValueError(f"User '{user}' doesn't exist.") - - user_credentials = config_reader[user] - if "aws_session_token" in user_credentials: - raise ValueError(f"Session tokens (short-term credentials) are not supported.") - - aws_access_key_id = user_credentials.get("aws_access_key_id") - aws_secret_access_key = user_credentials.get("aws_secret_access_key") - - if aws_access_key_id is None or aws_secret_access_key is None: - raise ValueError( - "The credentials file must contain both aws_access_key_id and aws_secret_access_key." - ) - - return aws_access_key_id, aws_secret_access_key - - def _generate_logs_list( container_logs_list_path: pathlib.Path, parsed_args: argparse.Namespace, @@ -128,19 +94,6 @@ def _generate_compress_cmd( return compress_cmd -def _add_common_arguments(args_parser: argparse.ArgumentParser) -> None: - args_parser.add_argument( - "--timestamp-key", - help="The path (e.g. x.y) for the field containing the log event's timestamp.", - ) - args_parser.add_argument( - "-t", "--tags", help="A comma-separated list of tags to apply to the compressed archives." - ) - args_parser.add_argument( - "--no-progress-reporting", action="store_true", help="Disables progress reporting." - ) - - def _validate_fs_input_args( parsed_args: argparse.Namespace, args_parser: argparse.ArgumentParser, From 21869dfed4a4f9e93dd5aee6c365715264a2f53d Mon Sep 17 00:00:00 2001 From: Eden Zhang Date: Wed, 26 Mar 2025 20:23:05 +0000 Subject: [PATCH 13/25] Remove unused imports --- .../clp-package-utils/clp_package_utils/scripts/compress.py | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/components/clp-package-utils/clp_package_utils/scripts/compress.py b/components/clp-package-utils/clp_package_utils/scripts/compress.py index 0605bf657a..6210876798 100755 --- a/components/clp-package-utils/clp_package_utils/scripts/compress.py +++ b/components/clp-package-utils/clp_package_utils/scripts/compress.py @@ -1,11 +1,10 @@ import argparse -import configparser import logging import pathlib import subprocess import sys import uuid -from typing import List, Tuple +from typing import List from clp_py_utils.clp_config import CLPConfig, StorageEngine from job_orchestration.scheduler.job_config import InputType From 8b936d5116866b3632c174573871064928be7e0b Mon Sep 17 00:00:00 2001 From: Eden Zhang Date: Wed, 26 Mar 2025 20:34:54 +0000 Subject: [PATCH 14/25] Update docs --- .../guides-using-object-storage/clp-config.md | 25 ++++++++++++++ .../guides-using-object-storage/clp-usage.md | 33 ++++--------------- 2 files changed, 32 insertions(+), 26 deletions(-) diff --git a/docs/src/user-guide/guides-using-object-storage/clp-config.md b/docs/src/user-guide/guides-using-object-storage/clp-config.md index 02e3b93607..290bef4761 100644 --- a/docs/src/user-guide/guides-using-object-storage/clp-config.md +++ b/docs/src/user-guide/guides-using-object-storage/clp-config.md @@ -6,6 +6,31 @@ To use object storage with CLP, follow the steps below to configure each use cas If CLP is already running, shut it down, update its configuration, and then start it again. ::: +## Configuration for input logs + +To configure CLP to compress logs from S3, update the `logs_input` key in +`/etc/clp-config.yml` with the values in the code block below, replacing the fields in +angle brackets (`<>`) with the appropriate values: + +```yaml +logs_input: + type: "s3" + s3_config: + region_code: "" + bucket: "" + key_prefix: "" + credentials: + access_key_id: "" + secret_access_key: "" +``` +* `s3_config` configures both the S3 bucket where archives should be stored and the credentials + for accessing it. + * `` is the AWS region [code][aws-region-codes] for the bucket. + * `` is the bucket's name. + * `` is the prefix of all logs you wish to compress and should be the same as the + `` value from the [compression IAM policy][compression-iam-policy]. + * `credentials` contains the CLP IAM user's credentials. + ## Configuration for archive storage To configure CLP to store archives on S3, update the `archive_output.storage` key in diff --git a/docs/src/user-guide/guides-using-object-storage/clp-usage.md b/docs/src/user-guide/guides-using-object-storage/clp-usage.md index 6fab2db443..b09e56ce2a 100644 --- a/docs/src/user-guide/guides-using-object-storage/clp-usage.md +++ b/docs/src/user-guide/guides-using-object-storage/clp-usage.md @@ -5,40 +5,20 @@ should be able to use CLP as described in the [quick start](../quick-start-overv ## Compressing logs from S3 -To compress logs from S3, use the `s3` subcommand as follows, replacing the fields in angle brackets +To compress logs from S3, use the compress script as follows, replacing the fields in angle brackets (`<>`) with the appropriate values: ```bash sbin/compress.sh \ - s3 \ - --aws-credentials-file \ --timestamp-key \ - https://.s3..amazonaws.com/ + ``` -* `` is the path to an AWS credentials file like the following: - - ```ini - [default] - aws_access_key_id = - aws_secret_access_key = - ``` - - * CLP expects the credentials to be in the `default` section. - * `` and `` are the access key ID and secret access - key of the CLP IAM user. - * If you don't want to use a credentials file, you can specify the credentials on the command - line using the `--aws-access-key-id` and `--aws-secret-access-key` flags (note that this may - expose your credentials to other users running on the system). - -* `` is the field path of the kv-pair that contains the timestamp in each log event. -* `` is the name of the S3 bucket containing your logs. -* `` is the AWS region [code][aws-region-codes] for the S3 bucket containing your logs. -* `` is the prefix of all logs you wish to compress and must begin with the - `` value from the [compression IAM policy][compression-iam-policy]. +* `` is the prefix of all logs you wish to compress and must be relative to the prefix + configured in the [logs-input s3_config][logs-input-s3-config]. :::{note} -The `s3` subcommand only supports a single URL but will compress any logs that have the given +Compressing from S3 only supports a single prefix but will compress any logs that have the given prefix. If you wish to compress a single log file, specify the entire path to the log file. However, if that @@ -49,4 +29,5 @@ both logs to be compressed). This limitation will be addressed in a future relea [add-iam-policy]: https://docs.aws.amazon.com/IAM/latest/UserGuide/access_policies_manage-attach-detach.html#embed-inline-policy-console [aws-region-codes]: https://docs.aws.amazon.com/AmazonRDS/latest/UserGuide/Concepts.RegionsAndAvailabilityZones.html#Concepts.RegionsAndAvailabilityZones.Availability -[compression-iam-policy]: ./object-storage-config.md#configuration-for-compression \ No newline at end of file +[compression-iam-policy]: ./object-storage-config.md#configuration-for-compression +[logs-input-s3-config]: ./clp-config.md#configuration-for-input-logs \ No newline at end of file From 75426c23487e07a8f1c96d78bf460a8b445e895b Mon Sep 17 00:00:00 2001 From: Eden Zhang Date: Wed, 26 Mar 2025 20:37:10 +0000 Subject: [PATCH 15/25] Fix missing change in docs --- docs/src/user-guide/guides-using-object-storage/clp-config.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/docs/src/user-guide/guides-using-object-storage/clp-config.md b/docs/src/user-guide/guides-using-object-storage/clp-config.md index 290bef4761..4f7973d242 100644 --- a/docs/src/user-guide/guides-using-object-storage/clp-config.md +++ b/docs/src/user-guide/guides-using-object-storage/clp-config.md @@ -23,7 +23,7 @@ logs_input: access_key_id: "" secret_access_key: "" ``` -* `s3_config` configures both the S3 bucket where archives should be stored and the credentials +* `s3_config` configures both the S3 bucket where logs are to be retrieved from and the credentials for accessing it. * `` is the AWS region [code][aws-region-codes] for the bucket. * `` is the bucket's name. From 5aef10c4b9855631704911392b71ddeaea1ad685 Mon Sep 17 00:00:00 2001 From: Eden Zhang Date: Wed, 26 Mar 2025 20:38:34 +0000 Subject: [PATCH 16/25] Add link to docs --- docs/src/user-guide/guides-using-object-storage/clp-config.md | 1 + 1 file changed, 1 insertion(+) diff --git a/docs/src/user-guide/guides-using-object-storage/clp-config.md b/docs/src/user-guide/guides-using-object-storage/clp-config.md index 4f7973d242..2b85739b6c 100644 --- a/docs/src/user-guide/guides-using-object-storage/clp-config.md +++ b/docs/src/user-guide/guides-using-object-storage/clp-config.md @@ -101,3 +101,4 @@ future release. ::: [aws-region-codes]: https://docs.aws.amazon.com/AmazonRDS/latest/UserGuide/Concepts.RegionsAndAvailabilityZones.html#Concepts.RegionsAndAvailabilityZones.Availability +[compression-iam-policy]: ./object-storage-config.md#configuration-for-compression From 4cc333751d7c97a52a4f35c030242ea7ae7b59dc Mon Sep 17 00:00:00 2001 From: Eden Zhang Date: Wed, 26 Mar 2025 20:42:35 +0000 Subject: [PATCH 17/25] Remove `fs` from other docs --- docs/src/user-guide/quick-start-compression/json.md | 3 +-- docs/src/user-guide/quick-start-compression/text.md | 2 +- 2 files changed, 2 insertions(+), 3 deletions(-) diff --git a/docs/src/user-guide/quick-start-compression/json.md b/docs/src/user-guide/quick-start-compression/json.md index 6091762a89..7784a4f509 100644 --- a/docs/src/user-guide/quick-start-compression/json.md +++ b/docs/src/user-guide/quick-start-compression/json.md @@ -3,10 +3,9 @@ To compress JSON logs, from inside the package directory, run: ```bash -sbin/compress.sh fs --timestamp-key '' [ ...] +sbin/compress.sh --timestamp-key '' [ ...] ``` -* `fs` is a subcommand for compressing logs from the filesystem. * `` is the field path of the kv-pair that contains the timestamp in each log event. * E.g., if your log events look like `{"timestamp": {"iso8601": "2024-01-01 00:01:02.345", ...}}`, you should enter diff --git a/docs/src/user-guide/quick-start-compression/text.md b/docs/src/user-guide/quick-start-compression/text.md index 18179a65b1..29e798b9d8 100644 --- a/docs/src/user-guide/quick-start-compression/text.md +++ b/docs/src/user-guide/quick-start-compression/text.md @@ -3,7 +3,7 @@ To compress unstructured text logs, from inside the package directory, run: ```bash -sbin/compress.sh fs [ ...] +sbin/compress.sh [ ...] ``` `` are paths to unstructured text log files or directories containing such files. From bc7b7f707e6bcae0fa2ac68cd65cf575ba0b68bc Mon Sep 17 00:00:00 2001 From: Eden Zhang <49173122+Eden-D-Zhang@users.noreply.github.com> Date: Wed, 26 Mar 2025 18:05:49 -0400 Subject: [PATCH 18/25] Update components/clp-package-utils/clp_package_utils/scripts/native/compress.py Co-authored-by: coderabbitai[bot] <136622811+coderabbitai[bot]@users.noreply.github.com> --- .../clp_package_utils/scripts/native/compress.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/components/clp-package-utils/clp_package_utils/scripts/native/compress.py b/components/clp-package-utils/clp_package_utils/scripts/native/compress.py index 4b3e91d6f5..e085c9a3d0 100755 --- a/components/clp-package-utils/clp_package_utils/scripts/native/compress.py +++ b/components/clp-package-utils/clp_package_utils/scripts/native/compress.py @@ -139,7 +139,7 @@ def _generate_clp_io_config( ) elif InputType.S3 == input_type: if len(logs_to_compress) != 1: - ValueError(f"Too many prefixes: {len(logs_to_compress)} > 1") + raise ValueError(f"Too many prefixes: {len(logs_to_compress)} > 1") s3_config = clp_config.logs_input.s3_config return S3InputConfig( From 11097e7790e286953c5621a4142ccef51b59c993 Mon Sep 17 00:00:00 2001 From: Eden Zhang Date: Sun, 6 Apr 2025 19:55:39 +0000 Subject: [PATCH 19/25] Implement review suggestions --- .../clp_package_utils/scripts/compress.py | 9 +++------ 1 file changed, 3 insertions(+), 6 deletions(-) diff --git a/components/clp-package-utils/clp_package_utils/scripts/compress.py b/components/clp-package-utils/clp_package_utils/scripts/compress.py index 6210876798..b039dbe032 100755 --- a/components/clp-package-utils/clp_package_utils/scripts/compress.py +++ b/components/clp-package-utils/clp_package_utils/scripts/compress.py @@ -28,9 +28,8 @@ def _generate_logs_list( container_logs_list_path: pathlib.Path, parsed_args: argparse.Namespace, - clp_config: CLPConfig, + input_type: InputType, ) -> None: - input_type = clp_config.logs_input.type if InputType.FS == input_type: host_logs_list_path = parsed_args.path_list @@ -67,9 +66,7 @@ def _generate_compress_cmd( parsed_args: argparse.Namespace, config_path: pathlib.Path, logs_list_path: pathlib.Path, - clp_config: CLPConfig, ) -> List[str]: - input_type = clp_config.logs_input.type # fmt: off compress_cmd = [ @@ -192,13 +189,13 @@ def main(argv): if not container_logs_list_path.exists(): break - _generate_logs_list(container_logs_list_path, parsed_args, clp_config) + _generate_logs_list(container_logs_list_path, parsed_args, clp_config.logs_input.type) container_start_cmd = generate_container_start_cmd( container_name, necessary_mounts, clp_config.execution_container ) compress_cmd = _generate_compress_cmd( - parsed_args, generated_config_path_on_container, logs_list_path_on_container, clp_config + parsed_args, generated_config_path_on_container, logs_list_path_on_container ) cmd = container_start_cmd + compress_cmd subprocess.run(cmd, check=True) From 417b9be1f9a983a055cdb006f484a9e9a7e42a70 Mon Sep 17 00:00:00 2001 From: Eden Zhang Date: Thu, 10 Apr 2025 19:29:15 +0000 Subject: [PATCH 20/25] Simplify function arguments --- .../clp-package-utils/clp_package_utils/scripts/compress.py | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/components/clp-package-utils/clp_package_utils/scripts/compress.py b/components/clp-package-utils/clp_package_utils/scripts/compress.py index b039dbe032..32917c4678 100755 --- a/components/clp-package-utils/clp_package_utils/scripts/compress.py +++ b/components/clp-package-utils/clp_package_utils/scripts/compress.py @@ -104,9 +104,9 @@ def _validate_fs_input_args( def _validate_s3_input_args( - parsed_args: argparse.Namespace, args_parser: argparse.ArgumentParser, clp_config: CLPConfig + parsed_args: argparse.Namespace, args_parser: argparse.ArgumentParser, storage_engine: StorageEngine ) -> None: - if StorageEngine.CLP_S != clp_config.package.storage_engine: + if StorageEngine.CLP_S != storage_engine: args_parser.error( f"Input type {InputType.S3} is only supported for the storage engine" f" {StorageEngine.CLP_S}." @@ -163,7 +163,7 @@ def main(argv): if InputType.FS == input_type: _validate_fs_input_args(parsed_args, args_parser) elif InputType.S3 == input_type: - _validate_s3_input_args(parsed_args, args_parser, clp_config) + _validate_s3_input_args(parsed_args, args_parser, clp_config.package.storage_engine) else: raise ValueError(f"Unsupported input type: {input_type}.") From a2802054ae98330a15c8f6208166e2327336a553 Mon Sep 17 00:00:00 2001 From: Eden Zhang Date: Thu, 10 Apr 2025 19:38:33 +0000 Subject: [PATCH 21/25] Lint --- .../clp-package-utils/clp_package_utils/scripts/compress.py | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/components/clp-package-utils/clp_package_utils/scripts/compress.py b/components/clp-package-utils/clp_package_utils/scripts/compress.py index 32917c4678..55afacff2f 100755 --- a/components/clp-package-utils/clp_package_utils/scripts/compress.py +++ b/components/clp-package-utils/clp_package_utils/scripts/compress.py @@ -104,7 +104,9 @@ def _validate_fs_input_args( def _validate_s3_input_args( - parsed_args: argparse.Namespace, args_parser: argparse.ArgumentParser, storage_engine: StorageEngine + parsed_args: argparse.Namespace, + args_parser: argparse.ArgumentParser, + storage_engine: StorageEngine, ) -> None: if StorageEngine.CLP_S != storage_engine: args_parser.error( From a0dae73ecb796c8e0202940138f27377ae5a6ee8 Mon Sep 17 00:00:00 2001 From: Eden Zhang Date: Fri, 11 Apr 2025 14:40:55 +0000 Subject: [PATCH 22/25] Update input log directory validation --- components/clp-py-utils/clp_py_utils/clp_config.py | 11 +++++++---- 1 file changed, 7 insertions(+), 4 deletions(-) diff --git a/components/clp-py-utils/clp_py_utils/clp_config.py b/components/clp-py-utils/clp_py_utils/clp_config.py index 4a4aa144d5..74d188183c 100644 --- a/components/clp-py-utils/clp_py_utils/clp_config.py +++ b/components/clp-py-utils/clp_py_utils/clp_config.py @@ -596,10 +596,13 @@ def make_config_paths_absolute(self, clp_home: pathlib.Path): def validate_logs_input_config(self): if StorageType.FS == self.logs_input.type: - try: - validate_path_could_be_dir(self.logs_input.directory) - except ValueError as ex: - raise ValueError(f"logs_input directory is invalid: {ex}") + # NOTE: This can't be a pydantic validator since input_logs_dir might be a package-relative + # path that will only be resolved after pydantic validation + input_logs_dir = self.logs_input.directory + if not input_logs_dir.exists(): + raise ValueError(f"input_logs_directory '{input_logs_dir}' doesn't exist.") + if not input_logs_dir.is_dir(): + raise ValueError(f"input_logs_directory '{input_logs_dir}' is not a directory.") def validate_archive_output_config(self): if ( From 413cc0609795def3aea399b23e030d75154b0dea Mon Sep 17 00:00:00 2001 From: Eden Zhang Date: Fri, 11 Apr 2025 14:42:06 +0000 Subject: [PATCH 23/25] Fix comment length --- components/clp-py-utils/clp_py_utils/clp_config.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/components/clp-py-utils/clp_py_utils/clp_config.py b/components/clp-py-utils/clp_py_utils/clp_config.py index 74d188183c..f72611be83 100644 --- a/components/clp-py-utils/clp_py_utils/clp_config.py +++ b/components/clp-py-utils/clp_py_utils/clp_config.py @@ -596,8 +596,8 @@ def make_config_paths_absolute(self, clp_home: pathlib.Path): def validate_logs_input_config(self): if StorageType.FS == self.logs_input.type: - # NOTE: This can't be a pydantic validator since input_logs_dir might be a package-relative - # path that will only be resolved after pydantic validation + # NOTE: This can't be a pydantic validator since input_logs_dir might be a + # package-relative path that will only be resolved after pydantic validation input_logs_dir = self.logs_input.directory if not input_logs_dir.exists(): raise ValueError(f"input_logs_directory '{input_logs_dir}' doesn't exist.") From ebbc0b32bb6acd46f160e87cb873d390b3f1bb2f Mon Sep 17 00:00:00 2001 From: Eden Zhang <49173122+Eden-D-Zhang@users.noreply.github.com> Date: Fri, 11 Apr 2025 19:06:44 -0400 Subject: [PATCH 24/25] Apply suggestions from code review Co-authored-by: kirkrodrigues <2454684+kirkrodrigues@users.noreply.github.com> --- .../clp_package_utils/scripts/compress.py | 3 +-- .../clp_package_utils/scripts/native/compress.py | 4 ++-- .../clp_package_utils/scripts/start_clp.py | 4 ++-- components/clp-py-utils/clp_py_utils/clp_config.py | 4 ++-- .../guides-using-object-storage/clp-usage.md | 13 ++++++++----- 5 files changed, 15 insertions(+), 13 deletions(-) diff --git a/components/clp-package-utils/clp_package_utils/scripts/compress.py b/components/clp-package-utils/clp_package_utils/scripts/compress.py index 55afacff2f..c178265d2e 100755 --- a/components/clp-package-utils/clp_package_utils/scripts/compress.py +++ b/components/clp-package-utils/clp_package_utils/scripts/compress.py @@ -30,7 +30,6 @@ def _generate_logs_list( parsed_args: argparse.Namespace, input_type: InputType, ) -> None: - if InputType.FS == input_type: host_logs_list_path = parsed_args.path_list with open(container_logs_list_path, "w") as container_logs_list_file: @@ -116,7 +115,7 @@ def _validate_s3_input_args( if len(parsed_args.paths) != 1: args_parser.error(f"Only one key prefix can be specified for input type {InputType.S3}.") if parsed_args.path_list is not None: - args_parser.error(f"Path list file is not supported for input type {InputType.S3}.") + args_parser.error(f"Path list file is unsupported for input type {InputType.S3}.") def main(argv): diff --git a/components/clp-package-utils/clp_package_utils/scripts/native/compress.py b/components/clp-package-utils/clp_package_utils/scripts/native/compress.py index e085c9a3d0..1c5d7f7449 100755 --- a/components/clp-package-utils/clp_package_utils/scripts/native/compress.py +++ b/components/clp-package-utils/clp_package_utils/scripts/native/compress.py @@ -10,7 +10,7 @@ import brotli import msgpack -from clp_py_utils.clp_config import CLPConfig, COMPRESSION_JOBS_TABLE_NAME, S3Credentials +from clp_py_utils.clp_config import CLPConfig, COMPRESSION_JOBS_TABLE_NAME from clp_py_utils.pretty_size import pretty_size from clp_py_utils.sql_adapter import SQL_Adapter from job_orchestration.scheduler.constants import ( @@ -139,7 +139,7 @@ def _generate_clp_io_config( ) elif InputType.S3 == input_type: if len(logs_to_compress) != 1: - raise ValueError(f"Too many prefixes: {len(logs_to_compress)} > 1") + raise ValueError(f"Too many key prefixes: {len(logs_to_compress)} > 1") s3_config = clp_config.logs_input.s3_config return S3InputConfig( diff --git a/components/clp-package-utils/clp_package_utils/scripts/start_clp.py b/components/clp-package-utils/clp_package_utils/scripts/start_clp.py index 790149c37d..e3f8b76b1f 100755 --- a/components/clp-package-utils/clp_package_utils/scripts/start_clp.py +++ b/components/clp-package-utils/clp_package_utils/scripts/start_clp.py @@ -743,10 +743,10 @@ def generic_start_worker( mounts.data_dir, mounts.logs_dir, ] - if worker_specific_mount: - necessary_mounts.extend(worker_specific_mount) if StorageType.FS == clp_config.logs_input.type: necessary_mounts.append(mounts.input_logs_dir) + if worker_specific_mount: + necessary_mounts.extend(worker_specific_mount) for mount in necessary_mounts: if not mount: diff --git a/components/clp-py-utils/clp_py_utils/clp_config.py b/components/clp-py-utils/clp_py_utils/clp_config.py index f72611be83..e01994e573 100644 --- a/components/clp-py-utils/clp_py_utils/clp_config.py +++ b/components/clp-py-utils/clp_py_utils/clp_config.py @@ -600,9 +600,9 @@ def validate_logs_input_config(self): # package-relative path that will only be resolved after pydantic validation input_logs_dir = self.logs_input.directory if not input_logs_dir.exists(): - raise ValueError(f"input_logs_directory '{input_logs_dir}' doesn't exist.") + raise ValueError(f"logs_input.directory '{input_logs_dir}' doesn't exist.") if not input_logs_dir.is_dir(): - raise ValueError(f"input_logs_directory '{input_logs_dir}' is not a directory.") + raise ValueError(f"logs_input.directory '{input_logs_dir}' is not a directory.") def validate_archive_output_config(self): if ( diff --git a/docs/src/user-guide/guides-using-object-storage/clp-usage.md b/docs/src/user-guide/guides-using-object-storage/clp-usage.md index b09e56ce2a..c909d87b6d 100644 --- a/docs/src/user-guide/guides-using-object-storage/clp-usage.md +++ b/docs/src/user-guide/guides-using-object-storage/clp-usage.md @@ -5,8 +5,8 @@ should be able to use CLP as described in the [quick start](../quick-start-overv ## Compressing logs from S3 -To compress logs from S3, use the compress script as follows, replacing the fields in angle brackets -(`<>`) with the appropriate values: +To compress logs from S3, use the `sbin/compress.sh` script as follows, replacing the fields in +angle brackets (`<>`) with the appropriate values: ```bash sbin/compress.sh \ @@ -14,8 +14,11 @@ sbin/compress.sh \ ``` -* `` is the prefix of all logs you wish to compress and must be relative to the prefix - configured in the [logs-input s3_config][logs-input-s3-config]. +* `` is the prefix of all logs you wish to compress and must be relative to + [logs-input.s3_config.key_prefix][logs-input-s3-config]. + * E.g., if you want to compress the S3 object `/a/b/c.jsonl`, and + `logs-input.s3_config.key_prefix` is `/a/`, then you would replace `` in the command + above with `b/c.jsonl`. :::{note} Compressing from S3 only supports a single prefix but will compress any logs that have the given @@ -30,4 +33,4 @@ both logs to be compressed). This limitation will be addressed in a future relea [add-iam-policy]: https://docs.aws.amazon.com/IAM/latest/UserGuide/access_policies_manage-attach-detach.html#embed-inline-policy-console [aws-region-codes]: https://docs.aws.amazon.com/AmazonRDS/latest/UserGuide/Concepts.RegionsAndAvailabilityZones.html#Concepts.RegionsAndAvailabilityZones.Availability [compression-iam-policy]: ./object-storage-config.md#configuration-for-compression -[logs-input-s3-config]: ./clp-config.md#configuration-for-input-logs \ No newline at end of file +[logs-input-s3-config]: ./clp-config.md#configuration-for-input-logs From d8681dc5dcda55b4b6cead6f7de29946eccb2add Mon Sep 17 00:00:00 2001 From: Eden Zhang Date: Fri, 11 Apr 2025 23:30:04 +0000 Subject: [PATCH 25/25] Apply review suggestions --- .../clp_package_utils/scripts/compress.py | 4 ++-- .../clp_package_utils/scripts/native/compress.py | 4 ++-- .../scheduler/compress/compression_scheduler.py | 4 ++-- components/package-template/src/etc/clp-config.yml | 7 ++++--- .../user-guide/guides-using-object-storage/clp-usage.md | 5 +++-- 5 files changed, 13 insertions(+), 11 deletions(-) diff --git a/components/clp-package-utils/clp_package_utils/scripts/compress.py b/components/clp-package-utils/clp_package_utils/scripts/compress.py index c178265d2e..f957cbef16 100755 --- a/components/clp-package-utils/clp_package_utils/scripts/compress.py +++ b/components/clp-package-utils/clp_package_utils/scripts/compress.py @@ -26,9 +26,9 @@ def _generate_logs_list( + input_type: InputType, container_logs_list_path: pathlib.Path, parsed_args: argparse.Namespace, - input_type: InputType, ) -> None: if InputType.FS == input_type: host_logs_list_path = parsed_args.path_list @@ -190,7 +190,7 @@ def main(argv): if not container_logs_list_path.exists(): break - _generate_logs_list(container_logs_list_path, parsed_args, clp_config.logs_input.type) + _generate_logs_list(clp_config.logs_input.type, container_logs_list_path, parsed_args) container_start_cmd = generate_container_start_cmd( container_name, necessary_mounts, clp_config.execution_container diff --git a/components/clp-package-utils/clp_package_utils/scripts/native/compress.py b/components/clp-package-utils/clp_package_utils/scripts/native/compress.py index 1c5d7f7449..fc6a7df1d4 100755 --- a/components/clp-package-utils/clp_package_utils/scripts/native/compress.py +++ b/components/clp-package-utils/clp_package_utils/scripts/native/compress.py @@ -127,7 +127,7 @@ def handle_job(sql_adapter: SQL_Adapter, clp_io_config: ClpIoConfig, no_progress def _generate_clp_io_config( - logs_to_compress: List[str], parsed_args: argparse.Namespace, clp_config: CLPConfig + clp_config: CLPConfig, logs_to_compress: List[str], parsed_args: argparse.Namespace ) -> typing.Union[S3InputConfig, FsInputConfig]: input_type = clp_config.logs_input.type @@ -216,7 +216,7 @@ def main(argv): logs_to_compress = _get_logs_to_compress(pathlib.Path(parsed_args.logs_list).resolve()) - clp_input_config = _generate_clp_io_config(logs_to_compress, parsed_args, clp_config) + clp_input_config = _generate_clp_io_config(clp_config, logs_to_compress, parsed_args) clp_output_config = OutputConfig.parse_obj(clp_config.archive_output) if parsed_args.tags: tag_list = [tag.strip().lower() for tag in parsed_args.tags.split(",") if tag] diff --git a/components/job-orchestration/job_orchestration/scheduler/compress/compression_scheduler.py b/components/job-orchestration/job_orchestration/scheduler/compress/compression_scheduler.py index 494c9fb228..0a02aaaff7 100644 --- a/components/job-orchestration/job_orchestration/scheduler/compress/compression_scheduler.py +++ b/components/job-orchestration/job_orchestration/scheduler/compress/compression_scheduler.py @@ -134,13 +134,13 @@ def _process_s3_input( and adds their metadata to paths_to_compress_buffer. :param s3_input_config: :param paths_to_compress_buffer: - :raises: RuntimeError if input path doesn't resolve to any objects. + :raises: RuntimeError if input prefix doesn't resolve to any objects. :raises: Propagates `s3_get_object_metadata`'s exceptions. """ object_metadata_list = s3_get_object_metadata(s3_input_config) if len(object_metadata_list) == 0: - raise RuntimeError("Input path doesn't resolve to any object") + raise RuntimeError("Input prefix doesn't resolve to any object") for object_metadata in object_metadata_list: paths_to_compress_buffer.add_file(object_metadata) diff --git a/components/package-template/src/etc/clp-config.yml b/components/package-template/src/etc/clp-config.yml index 71aa5f47a9..8e4f177312 100644 --- a/components/package-template/src/etc/clp-config.yml +++ b/components/package-template/src/etc/clp-config.yml @@ -1,9 +1,10 @@ -## A path containing any logs you which to compress. Must be reachable by all +## Location (e.g., directory) containing any logs you wish to compress. Must be reachable by all ## workers. -## - This path will be exposed inside the container, so symbolic links to files -## outside this path will be ignored. #logs_input: # type: "fs" +# +# # NOTE: This directory will be exposed inside the container, so symbolic links to files outside +# # this directory will be ignored. # directory: "/" # ## File containing credentials for services diff --git a/docs/src/user-guide/guides-using-object-storage/clp-usage.md b/docs/src/user-guide/guides-using-object-storage/clp-usage.md index c909d87b6d..963e659121 100644 --- a/docs/src/user-guide/guides-using-object-storage/clp-usage.md +++ b/docs/src/user-guide/guides-using-object-storage/clp-usage.md @@ -24,8 +24,9 @@ sbin/compress.sh \ Compressing from S3 only supports a single prefix but will compress any logs that have the given prefix. -If you wish to compress a single log file, specify the entire path to the log file. However, if that -log file's path is a prefix of another log file's path, then both log files will be compressed +If you wish to compress a single log file, specify the entire path to the log file +(relative to `logs-input.s3_config.key_prefix`). However, if that log file's path is a +prefix of another log file's path, then both log files will be compressed (e.g., with two files "logs/syslog" and "logs/syslog.1", a prefix like "logs/syslog" will cause both logs to be compressed). This limitation will be addressed in a future release. :::