Weifeng Sun

80 papers C 6Journal 62Unranked 12
YearRankTypeTitle / Venue / Authors
2026 J jnl
IEEE Trans. Circuits Syst. I Regul. Pap.
Yongjia Li, Feng Cheng, Jianlin Xia, Peng Su, Long Zhang, Zhongyuan Fang, Feng Lin, Encheng Zhu, Weifeng Sun
2026 J jnl
Int. J. Circuit Theory Appl.
Ziyan Zhou, Song Ding, Qi Liu, Qinsong Qian, Weifeng Sun
2025 J jnl
IEEE Trans. Circuits Syst. I Regul. Pap.
Jianlin Xia, Yongjia Li, Weiyue Qu, Zhongyuan Fang, Jin Wu, Feng Lin, Encheng Zhu, Weifeng Sun
2025 J jnl
IEEE Trans. Ind. Electron.
Ziyan Zhou, Qiang Luo, Yufan Wang, YueFei Sun, Qinsong Qian, Weifeng Sun
2025 J jnl
IEEE Trans. Ind. Electron.
Yijie Qian, Shen Xu, Xinru Wang, Lingyun Li, Weifeng Sun
2025 J jnl
IEEE Trans. Circuits Syst. I Regul. Pap.
Chenggong Wan, Yi Zhu, Yingjie Ma, Xue Li, Lixia Zheng, Jin Wu, Weifeng Sun
2025 J jnl
IEEE Trans. Circuits Syst. II Express Briefs
Chenggong Wan, Wenhao Yang, Jun Zhang, Jin Wu, Weifeng Sun, Lixia Zheng
2025 J jnl
Integr.
Lixia Zheng, Yuxiao Li, Zilu Zhao, Hehe Tian, Jin Wu, Weifeng Sun
2025 J jnl
IEEE Solid State Circuits Lett.
Yongjia Li, Jianlin Xia, Feng Cheng, Yifan Cao, Jin Wu, Encheng Zhu, Xiaofeng Sun, Dejin Wang, Long Zhang, Zhongyuan Fang, Weifeng Sun
2025 J jnl
Sci. China Inf. Sci.
Mengyao Zhao, Jie Ma, Lanlan Yang, Chuanqi Pan, Denggui Wang, Jianjun Zhou, Sheng Li, Jiaxing Wei, Long Zhang, Siyang Liu, Weifeng Sun
2024 C conf
ISCAS
Haiyang Guo, Zhongyuan Fang, Haonan Fan, Shen Xu, Xueyong Zhang, Weifeng Sun
2024 J jnl
IEEE Trans. Very Large Scale Integr. Syst.
Yifei Zheng, Boyu Li, Qianheng Dong, Yutao Ying, Deyuan Song, Jing Zhu, Weifeng Sun, Qinsong Qian, Long Zhang, Sheng Li, Denggui Wang, Jianjun Zhou
2024 J jnl
IEEE Trans. Circuits Syst. II Express Briefs
Haonan Fan, Zhongyuan Fang, Yongjia Li, Qinsong Qian, Xiaozhi Kang, Minggang Chen, Weifeng Sun
2024 C conf
ISCAS
Haonan Fan, Zhongyuan Fang, Minggang Chen, Weifeng Sun
2024 conf
BioCAS
Jiayi Yuan, Wangchen Fan, Weifeng Sun, Zhongyuan Fang
2024 C conf
ISCAS
Wangchen Fan, Qinsong Qian, Weifeng Sun, Zhongyuan Fang
2024 J jnl
J. Circuits Syst. Comput.
Lixia Zheng, Yi Zhu, Huiyong Xian, Jinwen Li, Chenggong Wan, Jin Wu, Weifeng Sun
2024 C conf
ISCAS
Jianlin Xia, Yongjia Li, Zhongyuan Fang, Jin Wu, Feng Lin, Weifeng Sun
2024 conf
ICTA
Chengwu Pan, Shipeng Chang, Jingjing Hong, Nailong He, Jie Ma, Sen Zhang, Long Zhang, Siyang Liu, Weifeng Sun
2024 C conf
ISCAS
Jiahao Liu, Wangchen Fan, Yiqing Wang, Weifeng Sun, Zhongyuan Fang
2024 J jnl
IEICE Electron. Express
Yongjia Li, Feng Cheng, Lei Zhang, Jianlin Xia, Zhongyuan Fang, Long Zhang, Weifeng Sun
2024 J jnl
IEEE Access
Wei Xi, Siyang Liu, Weifeng Sun
2024 conf
APCCAS
Haonan Fan, Jiahao Liu, Wangchen Fan, Weifeng Sun, Zhongyuan Fang
2024 conf
IEEE SENSORS
Guiqiang Zheng, Qingyin Zhong, Jie Ma, Yichen Li, Nannan Cheng, Nailong He, Sen Zhang, Yongjia Li, Long Zhang, Siyang Liu, Weifeng Sun
2024 J jnl
IEEE Trans. Circuits Syst. II Express Briefs
Lixia Zheng, Yanna Zheng, Chenggong Wan, Xuyan Li, Weifeng Sun, Jin Wu
2024 J jnl
Int. J. Circuit Theory Appl.
Lingyun Li, Shen Xu, Yujie Liu, Weifeng Sun
2024 J jnl
IEEE Trans. Circuits Syst. I Regul. Pap.
Song Ding, Zheng Zhang, Li Chen, Qinsong Qian, Weifeng Sun
2024 J jnl
Int. J. Circuit Theory Appl.
Yiqian Shan, Yuan Liang, Chenyang Li, Weifeng Sun, Zhongyuan Fang
2023 J jnl
IEEE Trans. Ind. Electron.
Siyuan Yu, Qi Zhou, Gang Shi, Tianyang Wu, Jing Zhu, Long Zhang, Weifeng Sun, Sen Zhang, Nailong He, Ye Li
2023 J jnl
IEEE Trans. Ind. Electron.
Limin Yu, Shen Xu, Chenxi Yang, Yugeng Wu, Longxing Shi, Weifeng Sun
2023 conf
APCCAS
Wangchen Fan, Zhongyuan Fang, Yongjia Li, Minggang Chen, Weifeng Sun
2023 J jnl
Microelectron. J.
Jin Wu, Minwei Hu, Xudong Wu, Yang Zuo, Chenggong Wan, Lixia Zheng, Weifeng Sun
2023 conf
ICTA
Weihao Lu, Sheng Li, Siyang Liu, Yanfeng Ma, Long Zhang, Jiaxing Wei, Weifeng Sun, Yiheng Li, Yuanyang Xia, Ke Wang, Tinggang Zhu
2022 J jnl
Microelectron. J.
Jin Wu, Jinwen Li, Xudong Wu, Chenggong Wan, Lixia Zheng, Weifeng Sun
2022 J jnl
IEEE Trans. Ind. Electron.
Sheng Li, Siyang Liu, Chi Zhang, Le Qian, Shuxuan Xin, Chen Ge, Weifeng Sun
2021 conf
ISSCC
Jing Zhu, Ding Yan, Siyuan Yu, Weifeng Sun, Gang Shi, Siyang Liu, Sen Zhang
2021 J jnl
IEEE Trans. Circuits Syst. I Regul. Pap.
Minggang Chen, Shen Xu, Linlin Huang, Weifeng Sun, Longxing Shi
2021 J jnl
Int. J. Circuit Theory Appl.
Limin Yu, Shen Xu, Huaxin Zhang, Longxing Shi, Weifeng Sun
2021 conf
ICVISP
Qi Liu, Weifeng Sun, Siyang Liu
2021 J jnl
IEEE Trans. Ind. Electron.
Jing Zhu, Siyuan Yu, Guichuang Zhu, Weifeng Sun, Yunqi Wang, Yangyang Lu, Yunwu Zhang, Sen Zhang, Yan Gu, Nailong He
2021 J jnl
IEEE Trans. Ind. Electron.
Siyang Liu, Xin Tong, Jiaxing Wei, Weifeng Sun, Zhuo Yang, Peng Ye, Yuanzheng Zhu
2021 J jnl
Circuits Syst. Signal Process.
Jin Wu, Xingsheng Sun, Mouzhao Zhou, Yue Rong, Lixia Zheng, Weifeng Sun
2021 J jnl
IEEE Access
Qi Liu, Tuanzhuang Wu, Yanfeng Ma, Siyang Liu, Xiaojuan Xia, Weifeng Sun
2021 J jnl
IEEE Trans. Ind. Electron.
Jing Zhu, Siyuan Yu, Yangyang Lu, Weifeng Sun, Chuanyi Cheng, Ding Yan, Yunwu Zhang, Shaohong Li, Long Zhang, Sen Zhang, Nailong He, Yan Gu
2020 J jnl
Microelectron. J.
Jin Wu, Shuang Chen, Kang Hu, Lixia Zheng, Weifeng Sun
2020 C conf
IECON
Minggang Chen, Shen Xu, Linlin Huang, Weifeng Sun
2020 J jnl
IEEE Trans. Ind. Electron.
Xiaoqiang Guo, Rui Zhong, Mingshu Zhang, Desheng Ding, Weifeng Sun
2020 J jnl
IEEE Trans. Circuits Syst.
Shengyou Xu, Qinsong Qian, Ran Shi, Syed Sikandar Shah, Qi Liu, Shengli Lu, Weifeng Sun
2020 J jnl
Int. J. Circuit Theory Appl.
Shen Xu, Chong Wang, Shengli Lu, Weifeng Sun
2019 J jnl
Sci. China Inf. Sci.
Jin Wu, Zhiming Qian, Yang Zhao, Xiangrong Yu, Lixia Zheng, Weifeng Sun
2019 J jnl
IEEE Trans. Circuits Syst. II Express Briefs
Jin Wu, Pengbo Zhang, Shufang Shi, Lixia Zheng, Weifeng Sun
2018 J jnl
Int. J. Circuit Theory Appl.
Chong Wang, Shen Xu, Xingpeng Kou, Shengli Lu, Weifeng Sun
2018 J jnl
IEEE Trans. Ind. Electron.
Xiaoqiang Guo, Rui Zhong, Mingshu Zhang, Desheng Ding, Weifeng Sun
2018 J jnl
IET Circuits Devices Syst.
Jin Wu, Youzhi Zhang, Rongqi Zhao, Kunpeng Zhang, Lixia Zheng, Weifeng Sun
2018 J jnl
IEEE Trans. Ind. Electron.
Jing Zhu, Yunwu Zhang, Weifeng Sun, Yangyang Lu, Longxing Shi, Yan Gu, Sen Zhang
2018 J jnl
J. Circuits Syst. Comput.
Jin Wu, Pengfei Dai, Jie Peng, Lixia Zheng, Weifeng Sun
2017 J jnl
Microelectron. J.
Jin Wu, Wenlong Zhang, Xiangrong Yu, Qi Jiang, Lixia Zheng, Weifeng Sun
2017 J jnl
J. Circuits Syst. Comput.
Lixia Zheng, Huan Hu, Ziqing Weng, Qun Yao, Jin Wu, Weifeng Sun
2017 J jnl
IEEE Trans. Circuits Syst. II Express Briefs
Jin Wu, Qi Jiang, Ke Song, Lixia Zheng, Dongchen Sun, Weifeng Sun
2016 J jnl
IET Circuits Devices Syst.
Jing Zhu, Yunwu Zhang, Weifeng Sun, Yangyang Lu, Yicheng Du, Yangbo Yi
2015 conf
ASICON
Jiaxing Wei, Jianfeng Wang, Ning Wang, Siyang Liu, Weifeng Sun
2014 J jnl
IEICE Electron. Express
Chunwei Zhang, Siyang Liu, Daying Sun, Chaohui Yu, Weifeng Sun
2014 J jnl
IET Circuits Devices Syst.
Jing Zhu, Yunwu Zhang, Weifeng Sun, Yangbo Yi
2014 J jnl
IET Circuits Devices Syst.
Jian Chen, Jing Zhu, Guodong Sun, Weifeng Sun, Weinan Dai, Zexiang Huang
2014 J jnl
Microelectron. J.
Qing Wang, Ning Chen, Shen Xu, Weifeng Sun, Longxing Shi
2013 conf
ASICON
Jing Zhu, Yunwu Zhang, Weifeng Sun, Shengli Lu
2011 conf
ASICON
Miao Yang, Weifeng Sun, Shen Xu, Shengli Lu, Longxing Shi
2011 J jnl
Microelectron. J.
Kan Jia, Weifeng Sun, Longxing Shi
2011 conf
ASICON
Daying Sun, Shen Xu, Weifeng Sun, Shengli Lu, Longxing Shi
2011 J jnl
Microelectron. J.
Qinsong Qian, Siyang Liu, Weifeng Sun, Hu Sun
2010 J jnl
Microelectron. Reliab.
Hong Wu, Qinsong Qian, Siyang Liu, Weifeng Sun, Longxing Shi
2010 J jnl
Microelectron. Reliab.
Qinsong Qian, Weifeng Sun, Jing Zhu, Longxing Shi
2009 J jnl
IEICE Trans. Electron.
Xiaojuan Xia, Liang Xie, Weifeng Sun, Longxing Shi
2009 J jnl
Microelectron. J.
Haisong Li, Qinsong Qian, Hong Wu, Weifeng Sun, Longxing Shi
2009 J jnl
IET Circuits Devices Syst.
Xiaojuan Xia, Liang Xie, Weifeng Sun, Longxing Shi
2008 J jnl
IET Circuits Devices Syst.
Weifeng Sun, Haisong Li, Yangbo Yi, Hong Wu, Longxing Shi
2008 J jnl
Microelectron. Reliab.
Hong Wu, Weifeng Sun, Yangbo Yi, Haisong Li, Longxing Shi
2006 J jnl
Microelectron. J.
Zhilin Sun, Weifeng Sun, Longxing Shi
2006 J jnl
Microelectron. Reliab.
Zhilin Sun, Weifeng Sun, Yangbo Yi, Longxing Shi
2005 J jnl
Microelectron. Reliab.
Weifeng Sun, Longxing Shi
s3-storage/s3_uploader.py
← Index s3-storage/s3_uploader.py python
# python s3_uploader.py /path/to/files repository_name --notes "Optional notes"
import hashlib
import multiprocessing
from multiprocessing import Pool
import sys
import logging
from logging.handlers import QueueHandler
from datetime import datetime
import os
import time
import json
from enum import Enum
import tempfile
import magic
from typing import Optional, Dict, Any
from dotenv import load_dotenv
from datetime import timezone

# File handling
import py7zr
import pyzipper
from magika import Magika

# S3 and database
from minio import Minio
import clickhouse_connect

load_dotenv(override=True)

class UploadResult(Enum):
    CORRECTLY = 0
    FAILED = 1
    SKIPPED = 2  # For duplicates

class FileNameFormatter(logging.Formatter):
    def format(self, record):
        record.filenameinfo = getattr(record, "filenameinfo", "unknown")
        return super().format(record)

def setup_logger(log_file, filename):
    """Set up a logger for a specific file"""
    logger = logging.getLogger(filename)
    if not logger.handlers:
        if os.getenv("SERVER_ENV") == "prod":
            logger.setLevel(logging.INFO)
        else:
            logger.setLevel(logging.DEBUG)
            
        # Create file handler
        handler = logging.FileHandler(log_file)
        formatter = FileNameFormatter(
            "%(asctime)s - %(filenameinfo)s - %(levelname)s - %(message)s"
        )
        handler.setFormatter(formatter)
        logger.addHandler(handler)
        logger.propagate = False
    return logger

def logger_thread(log_queue, log_file):
    handler = logging.FileHandler(log_file)
    formatter = FileNameFormatter(
        "%(asctime)s - %(filenameinfo)s - %(levelname)s - %(message)s"
    )
    handler.setFormatter(formatter)

    while True:
        try:
            record = log_queue.get()
            if record is None:
                break
            message = formatter.format(record)
            handler.stream.write(message + "\n")
            handler.stream.flush()
        except Exception:
            import traceback
            print("[ERR] Error in logger thread:", file=sys.stderr)
            traceback.print_exc(file=sys.stderr)

def get_s3_key(sha256: str, is_archived: bool, original_ext: str = None) -> str:
    """
    Generate S3 key using sharding pattern with appropriate extension
    
    Args:
        sha256: The file's SHA256 hash
        is_archived: Whether the file was already archived (zip/7z)
        original_ext: The original file extension (if any)
    """
    # If file was already archived, keep its extension
    if is_archived and original_ext and original_ext.lower() in ['.zip', '.7z']:
        extension = original_ext
    else:
        # For files we archived ourselves, use .zip
        extension = '.zip'
        
    return f"{sha256[:2]}/{sha256[2:4]}/{sha256}{extension}"

def process_single_file(args):
    """Independent worker function for processing a single file"""
    filepath, file_number, total_files, config, log_file = args
    filename = os.path.basename(filepath)
    base_logger = setup_logger(log_file, filename)
    extra = {"filenameinfo": filename}
    logger = logging.LoggerAdapter(base_logger, extra)
    s3_client = None
    ch_client = None

    try:
        logger.info(f"Processing file {file_number}/{total_files}: {filepath}")

        # Create connections for this process
        s3_client = Minio(
            config["s3_endpoint"],
            access_key=config["s3_access_key"],
            secret_key=config["s3_secret_key"],
            secure=True
        )
        
        ch_client = clickhouse_connect.get_client(
            host=config['clickhouse_host'],
            port=config['clickhouse_port'],
            username=config['clickhouse_user'],
            password=config['clickhouse_password'],
            database=config['clickhouse_database'],
            verify=config['clickhouse_verify']
        )

        # Process the file and get contents
        ext = os.path.splitext(filepath.lower())[1]
        is_archived = ext in ['.zip', '.7z']

        with tempfile.TemporaryDirectory() as temp_dir:
            try:
                if filename.startswith("."):
                    logger.debug(f"Skipping hidden file: {filename}")
                    return UploadResult.SKIPPED

                elif ext == ".zip":
                    logger.debug(f"Processing ZIP file: {filepath}")
                    with pyzipper.AESZipFile(filepath) as zf:
                        zf.pwd = b"infected"
                        filename = zf.namelist()[0]
                        zf.extractall(temp_dir)
                        extracted_path = os.path.join(temp_dir, filename)
                        
                        with open(extracted_path, "rb") as f:
                            file_content = f.read()
                        with open(filepath, "rb") as f:
                            archive_content = f.read()
                            
                elif ext == ".7z":
                    logger.debug(f"Processing 7z file: {filepath}")
                    with py7zr.SevenZipFile(filepath, mode="r", password="infected") as z:
                        z.extractall(path=temp_dir)
                        # Get first file in the archive
                        for root, _, files in os.walk(temp_dir):
                            if files:
                                filename = files[0]
                                extracted_path = os.path.join(root, filename)
                                break
                                
                        with open(extracted_path, "rb") as f:
                            file_content = f.read()
                        with open(filepath, "rb") as f:
                            archive_content = f.read()
                else:
                    logger.debug(f"Processing non-archived file: {filepath}")
                    # For non-archived files, create password protected zip
                    with open(filepath, "rb") as f:
                        file_content = f.read()
                    filename = os.path.basename(filepath)
                    
                    # Create password protected zip
                    with tempfile.NamedTemporaryFile(delete=False) as temp_zip:
                        with pyzipper.AESZipFile(
                            temp_zip.name,
                            'w',
                            compression=pyzipper.ZIP_LZMA,
                            encryption=pyzipper.WZ_AES
                        ) as zf:
                            zf.pwd = b'infected'
                            zf.writestr(filename, file_content)
                        
                        with open(temp_zip.name, 'rb') as f:
                            archive_content = f.read()
                    
                    os.unlink(temp_zip.name)

                # Calculate hash
                sha256 = hashlib.sha256(file_content).hexdigest()
                # s3_key = f"{sha256[:2]}/{sha256[2:4]}/{sha256}"
                s3_key = get_s3_key(sha256, is_archived, ext if is_archived else None)
                logger.debug(f"Calculated SHA256: {sha256}")

                # Check if exists in S3
                try:
                    s3_client.stat_object(config["s3_bucket"], s3_key)
                    exists = True
                    logger.debug(f"File exists in S3: {s3_key}")
                except:
                    exists = False
                    logger.debug(f"File does not exist in S3: {s3_key}")

                # Prepare metadata
                max_timestamp = int(datetime(9999, 12, 31, 23, 59, 59, tzinfo=timezone.utc).timestamp() * 1000)
                now = datetime.now(timezone.utc)
                now_timestamp = int(now.timestamp() * 1000)
                inverted_timestamp = datetime.fromtimestamp((max_timestamp - now_timestamp)/1000, tz=timezone.utc)


                data = [[
                    sha256,                     # String
                    filename,                   # String
                    config['repository'],       # String
                    len(file_content),         # UInt64
                    magic.from_buffer(file_content),        # String
                    magic.from_buffer(file_content, mime=True),  # String
                    Magika().identify_bytes(file_content).output.ct_label,  # String
                    now,                       # upload_date
                    now,                       # first_seen
                    inverted_timestamp,        # version_date for replacing
                    config['s3_bucket'],       # String
                    s3_key,                    # String
                    config['notes'] if config['notes'] else None  # Nullable(String)
                ]]

                column_names = [
                    'sha256', 'filename', 'repository', 'file_size',
                    'filetype', 'filetype_mime', 'filetype_magika',
                    'upload_date', 'first_seen', 'version_date', 's3_bucket', 's3_key', 'notes'
                ]

                column_type_names = [
                    'String', 'String', 'LowCardinality(String)', 'UInt64',
                    'String', 'String', 'String',
                    'DateTime64(3, \'UTC\')', 'DateTime64(3, \'UTC\')', 
                    'DateTime64(3, \'UTC\')', 'String', 'String', 'Nullable(String)'
                ]
                
                if exists:
                    # Just update database
                    try:
                        ch_client.insert(
                            'samples_catalog',
                            data,
                            column_names=column_names,
                            column_type_names=column_type_names,
                            settings={'input_format_values_interpret_expressions': 0}
                        )
                        logger.info(f"Updated metadata for existing file: {sha256}")
                    except Exception as e:
                        logger.error(f"Error inserting metadata: {str(e)}")
                        import traceback
                        logger.error(f"Full traceback: {traceback.format_exc()}")
                    return UploadResult.SKIPPED
                else:
                    try:    
                        # Upload to S3 and insert metadata
                        import io
                        file_data = io.BytesIO(archive_content)
                        s3_client.put_object(
                            bucket_name=config['s3_bucket'],
                                object_name=s3_key,
                                data=file_data,
                                length=len(archive_content)
                            )
                        logger.info(f"Uploaded to S3: {s3_key}")
                    except Exception as e:
                        logger.error(f"Error uploading to S3: {str(e)}")
                        import traceback
                        logger.error(f"Full traceback: {traceback.format_exc()}")
                        return UploadResult.FAILED
                    
                    try:
                        ch_client.insert(
                            'samples_catalog',
                            data,
                            column_names=column_names,
                            column_type_names=column_type_names,
                            settings={'input_format_values_interpret_expressions': 0}
                        )
                        logger.info(f"Successfully processed file: {sha256}")
                        return UploadResult.CORRECTLY
                    except Exception as e:
                        logger.error(f"Error inserting metadata: {str(e)}")
                        import traceback
                        logger.error(f"Full traceback: {traceback.format_exc()}")
                        return UploadResult.FAILED

            except Exception as e:
                logger.error(f"Error processing file: {str(e)}")
                import traceback
                logger.error(f"Full traceback: {traceback.format_exc()}")
                return UploadResult.FAILED

    except Exception as e:
        logger.error(f"Error in worker: {str(e)}")
        return UploadResult.FAILED
    finally:
        # if s3_client:
        #     s3_client.close()
        if ch_client:
            ch_client.close()
        if base_logger and base_logger.handlers:
            for handler in base_logger.handlers:
                handler.close()
            base_logger.handlers.clear()

class S3Uploader:
    def __init__(self, path: str, repository: str, notes: Optional[str] = None):
        """
        Initialize S3 uploader
        
        Args:
            path: Path to file or directory to process
            repository: Repository name
            notes: Optional notes to add to all files
        """
        self.path = path
        self.repository = repository
        self.notes = notes
        
        # Initialize multiprocessing components
        self.manager = multiprocessing.Manager()
        # self.log_queue = self.manager.Queue()
        self.total_results = self.manager.dict({result: 0 for result in UploadResult})
        
        # Set up logging
        self.today = datetime.today().strftime("%Y%m%dT%H%M%S")
        # self.log_file = os.getenv("LOG_FILE_PATH", "./logs/") + f"{self.today}-{self.repository}-upload.txt"
        log_base_path = os.getenv("LOG_FILE_PATH")
        if not log_base_path:
            log_base_path = os.path.join(os.getcwd(), "logs")  # Default to ./logs directory
        self.log_file = os.path.join(log_base_path, f"{self.today}-{self.repository}-upload.txt")
        os.makedirs(os.path.dirname(self.log_file), exist_ok=True)

        # Initialize S3 client
        self.s3_client = Minio(
            os.getenv("S3_ENDPOINT"),
            access_key=os.getenv("S3_ACCESS_KEY"),
            secret_key=os.getenv("S3_SECRET_KEY"),
            secure=True  # Set to False if not using HTTPS
        )
        self.bucket_name = os.getenv("S3_BUCKET")

        # Initialize ClickHouse client
        self.ch_client = clickhouse_connect.get_client(
            host=os.getenv('CLICKHOUSE_HOST'),
            port=os.getenv('CLICKHOUSE_PORT'),
            username=os.getenv('CLICKHOUSE_USER'),
            password=os.getenv('CLICKHOUSE_PASSWORD'),
            database=os.getenv('CLICKHOUSE_DATABASE'),
            verify=os.getenv('CLICKHOUSE_ENFORCE_SSL', 'False')
        )
        
        # Ensure database table exists
        self._create_table()

    def _create_table(self):
        """Create the file catalog table if it doesn't exist"""
        create_table_query = """
            CREATE TABLE IF NOT EXISTS samples_catalog (
                sha256 String,
                filename String,
                repository LowCardinality(String),
                file_size UInt64,
                filetype String,
                filetype_mime String,
                filetype_magika String,
                upload_date DateTime64(3, 'UTC'),
                first_seen DateTime64(3, 'UTC'),
                version_date DateTime64(3, 'UTC'),  -- This will be inverted timestamp
                s3_bucket String,
                s3_key String,
                notes Nullable(String),
                PRIMARY KEY (sha256, repository)
            ) ENGINE = ReplacingMergeTree(version_date)
            ORDER BY (sha256, repository)
        """
        self.ch_client.command(create_table_query)

    # # def upload(self):
    # #     """Main upload function with multiprocessing support"""
    # #     start_time = time.time()

    # #     config = {
    # #         "s3_endpoint": os.getenv("S3_ENDPOINT"),
    # #         "s3_access_key": os.getenv("S3_ACCESS_KEY"),
    # #         "s3_secret_key": os.getenv("S3_SECRET_KEY"),
    # #         "s3_bucket": os.getenv("S3_BUCKET"),
    # #         "clickhouse_host": os.getenv("CLICKHOUSE_HOST"),
    # #         "clickhouse_port": os.getenv("CLICKHOUSE_PORT"),
    # #         "clickhouse_user": os.getenv("CLICKHOUSE_USER"),
    # #         "clickhouse_password": os.getenv("CLICKHOUSE_PASSWORD"),
    # #         "clickhouse_database": os.getenv("CLICKHOUSE_DATABASE"),
    # #         "clickhouse_verify": os.getenv("CLICKHOUSE_ENFORCE_SSL", "False"),
    # #         "repository": self.repository,
    # #         "notes": self.notes
    # #     }

    # #     try:
    # #         # Get list of files
    # #         if os.path.isfile(self.path):
    # #             files = [self.path] if not self.path.endswith('.txt') else [
    # #                 line.strip() for line in open(self.path) 
    # #                 if line.strip() and not os.path.basename(line.strip()).startswith('.')
    # #             ]
    # #         elif os.path.isdir(self.path):
    # #             files = [
    # #                 os.path.join(root, file)
    # #                 for root, _, files in os.walk(self.path)
    # #                 for file in files
    # #                 if not file.startswith('.')
    # #             ]
    # #         else:
    # #             print(f"[ERR] Invalid path: {self.path}")
    # #             return

    # #         total_files = len(files)
    # #         print(f"Found {total_files} files to process")

    # #         with Pool(processes=max(1, multiprocessing.cpu_count() - 1)) as pool:
    # #             results = pool.map(
    # #                 process_single_file,
    # #                 [(f, i + 1, total_files, config, self.log_file) for i, f in enumerate(files)]
    # #             )

    # #         # Update statistics
    # #         for result in results:
    # #             if result is not None:
    # #                 self.total_results[result] += 1

    # #         # Generate summary
    # #         end_time = time.time()
    # #         elapsed_time = end_time - start_time
    # #         elapsed_time_pretty = time.strftime("%H:%M:%S", time.gmtime(elapsed_time))

    # #         summary = (
    # #             f"\n\nUpload finished for {self.path}"
    # #             f"\nTime required: {elapsed_time_pretty}"
    # #             f"\nResults:"
    # #             f"\n- Total processed: {sum(self.total_results.values())}"
    # #             f"\n- Successfully uploaded: {self.total_results[UploadResult.CORRECTLY]}"
    # #             f"\n- Skipped (already exists): {self.total_results[UploadResult.SKIPPED]}"
    # #             f"\n- Failed: {self.total_results[UploadResult.FAILED]}\n"
    # #         )

    # #         with open(self.log_file, "a") as f:
    # #             f.write(summary)

    # #         print(summary)

    # #     except Exception as e:
    # #         print(f"[ERR] Error in upload process: {str(e)}")

    # def upload(self):
    #     """Main upload function with multiprocessing support"""
    #     start_time = time.time()
    #     BATCH_SIZE = 1000  # Process files in batches of 1000

    #     config = {
    #         "s3_endpoint": os.getenv("S3_ENDPOINT"),
    #         "s3_access_key": os.getenv("S3_ACCESS_KEY"),
    #         "s3_secret_key": os.getenv("S3_SECRET_KEY"),
    #         "s3_bucket": os.getenv("S3_BUCKET"),
    #         "clickhouse_host": os.getenv("CLICKHOUSE_HOST"),
    #         "clickhouse_port": os.getenv("CLICKHOUSE_PORT"),
    #         "clickhouse_user": os.getenv("CLICKHOUSE_USER"),
    #         "clickhouse_password": os.getenv("CLICKHOUSE_PASSWORD"),
    #         "clickhouse_database": os.getenv("CLICKHOUSE_DATABASE"),
    #         "clickhouse_verify": os.getenv("CLICKHOUSE_ENFORCE_SSL", "False"),
    #         "repository": self.repository,
    #         "notes": self.notes
    #     }

    #     try:
    #         # Get list of files
    #         if os.path.isfile(self.path):
    #             files = [self.path] if not self.path.endswith('.txt') else [
    #                 line.strip() for line in open(self.path) 
    #                 if line.strip() and not os.path.basename(line.strip()).startswith('.')
    #             ]
    #         elif os.path.isdir(self.path):
    #             files = [
    #                 os.path.join(root, file)
    #                 for root, _, files in os.walk(self.path)
    #                 for file in files
    #                 if not file.startswith('.')
    #             ]
    #         else:
    #             print(f"[ERR] Invalid path: {self.path}")
    #             return

    #         total_files = len(files)
    #         print(f"Found {total_files} files to process")

    #         # Process files in batches
    #         for i in range(0, total_files, BATCH_SIZE):
    #             batch_files = files[i:i + BATCH_SIZE]
    #             batch_start = i + 1
    #             batch_end = min(i + BATCH_SIZE, total_files)
    #             print(f"\nProcessing batch {batch_start}-{batch_end} of {total_files}")
                
    #             with Pool(processes=max(1, multiprocessing.cpu_count() - 1)) as pool:
    #                 batch_results = pool.map(
    #                     process_single_file,
    #                     [(f, j + 1, total_files, config, self.log_file) 
    #                         for j, f in enumerate(batch_files, start=i)]
    #                 )
                    
    #                 # Update statistics for this batch
    #                 for result in batch_results:
    #                     if result is not None:
    #                         self.total_results[result] += 1

    #                 # Print intermediate summary
    #                 print(f"Batch {batch_start}-{batch_end} completed:")
    #                 print(f"- Successfully uploaded: {sum(1 for r in batch_results if r == UploadResult.CORRECTLY)}")
    #                 print(f"- Skipped (exists): {sum(1 for r in batch_results if r == UploadResult.SKIPPED)}")
    #                 print(f"- Failed: {sum(1 for r in batch_results if r == UploadResult.FAILED)}")

    #         # Generate final summary
    #         end_time = time.time()
    #         elapsed_time = end_time - start_time
    #         elapsed_time_pretty = time.strftime("%H:%M:%S", time.gmtime(elapsed_time))

    #         summary = (
    #             f"\n\nUpload finished for {self.path}"
    #             f"\nTime required: {elapsed_time_pretty}"
    #             f"\nResults:"
    #             f"\n- Total processed: {sum(self.total_results.values())}"
    #             f"\n- Successfully uploaded: {self.total_results[UploadResult.CORRECTLY]}"
    #             f"\n- Skipped (already exists): {self.total_results[UploadResult.SKIPPED]}"
    #             f"\n- Failed: {self.total_results[UploadResult.FAILED]}\n"
    #         )

    #         with open(self.log_file, "a") as f:
    #             f.write(summary)

    #         print(summary)

    #     except Exception as e:
    #         print(f"[ERR] Error in upload process: {str(e)}")

    def upload(self):
        """Main upload function with multiprocessing support"""
        start_time = time.time()
        BATCH_SIZE = 1000  # Process files in batches of 1000

        config = {
            "s3_endpoint": os.getenv("S3_ENDPOINT"),
            "s3_access_key": os.getenv("S3_ACCESS_KEY"),
            "s3_secret_key": os.getenv("S3_SECRET_KEY"),
            "s3_bucket": os.getenv("S3_BUCKET"),
            "clickhouse_host": os.getenv("CLICKHOUSE_HOST"),
            "clickhouse_port": os.getenv("CLICKHOUSE_PORT"),
            "clickhouse_user": os.getenv("CLICKHOUSE_USER"),
            "clickhouse_password": os.getenv("CLICKHOUSE_PASSWORD"),
            "clickhouse_database": os.getenv("CLICKHOUSE_DATABASE"),
            "clickhouse_verify": os.getenv("CLICKHOUSE_ENFORCE_SSL", "False"),
            "repository": self.repository,
            "notes": self.notes
        }

        try:
            # Get all existing filenames for this repository
            query = f"SELECT filename FROM samples_catalog WHERE repository = '{self.repository}'"
            # existing_files = set()
            # for row in self.ch_client.query(query).result_rows:
            #     filename = os.path.splitext(row[0])[0]  # Strip extension
            #     existing_files.add(filename)
            existing_files = set(row[0] for row in self.ch_client.query(query).result_rows)
            print(f"Found {len(existing_files)} existing files in repository")

            # Get list of files to process
            if os.path.isfile(self.path):
                files = [self.path] if not self.path.endswith('.txt') else [
                    line.strip() for line in open(self.path) 
                    if line.strip() and not os.path.basename(line.strip()).startswith('.')
                ]
            elif os.path.isdir(self.path):
                files = [
                    os.path.join(root, file)
                    for root, _, files in os.walk(self.path)
                    for file in files
                    if not file.startswith('.')
                ]
            else:
                print(f"[ERR] Invalid path: {self.path}")
                return

            # Filter out files that are already in the repository
            files_to_process = []
            for f in files:
                basename = os.path.splitext(os.path.basename(f))[0]
                if basename not in existing_files:
                    files_to_process.append(f)

            total_files = len(files_to_process)
            skipped_files = len(files) - total_files
            print(f"Found {total_files} new files to process (skipped {skipped_files} existing files)")

            if total_files == 0:
                print("No new files to process")
                return

            # Process files in batches
            for i in range(0, total_files, BATCH_SIZE):
                batch_files = files_to_process[i:i + BATCH_SIZE]
                batch_start = i + 1
                batch_end = min(i + BATCH_SIZE, total_files)
                print(f"\nProcessing batch {batch_start}-{batch_end} of {total_files}")
                
                with Pool(processes=max(1, multiprocessing.cpu_count() - 1)) as pool:
                    batch_results = pool.map(
                        process_single_file,
                        [(f, j + 1, total_files, config, self.log_file) 
                        for j, f in enumerate(batch_files, start=i)]
                    )
                    
                    # Update statistics for this batch
                    for result in batch_results:
                        if result is not None:
                            self.total_results[result] += 1

                    # Print intermediate summary
                    print(f"Batch {batch_start}-{batch_end} completed:")
                    print(f"- Successfully uploaded: {sum(1 for r in batch_results if r == UploadResult.CORRECTLY)}")
                    print(f"- Skipped (exists): {sum(1 for r in batch_results if r == UploadResult.SKIPPED)}")
                    print(f"- Failed: {sum(1 for r in batch_results if r == UploadResult.FAILED)}")

            # Generate final summary
            end_time = time.time()
            elapsed_time = end_time - start_time
            elapsed_time_pretty = time.strftime("%H:%M:%S", time.gmtime(elapsed_time))

            summary = (
                f"\n\nUpload finished for {self.path}"
                f"\nTime required: {elapsed_time_pretty}"
                f"\nResults:"
                f"\n- Files already in repository: {skipped_files}"
                f"\n- New files processed: {total_files}"
                f"\n- Successfully uploaded: {self.total_results[UploadResult.CORRECTLY]}"
                f"\n- Skipped (already exists): {self.total_results[UploadResult.SKIPPED]}"
                f"\n- Failed: {self.total_results[UploadResult.FAILED]}\n"
            )

            with open(self.log_file, "a") as f:
                f.write(summary)

            print(summary)

        except Exception as e:
            print(f"[ERR] Error in upload process: {str(e)}")
            import traceback
            print(traceback.format_exc())

if __name__ == "__main__":
    import argparse

    parser = argparse.ArgumentParser(description='Upload malware samples to S3 and catalog them.')
    parser.add_argument('path', help='Path to file or directory to process')
    parser.add_argument('repository', help='Repository name')
    parser.add_argument('--notes', help='Optional notes to add to all files', default=None)
    
    args = parser.parse_args()
    
    uploader = S3Uploader(args.path, args.repository, args.notes)
    uploader.upload()