Wei Tu

86 papers A* 1A 2Journal 73Unranked 9
YearRankTypeTitle / Venue / Authors
2026 J jnl
Comput. Environ. Urban Syst.
Zhaoyue Cai, Siqi Chen, Yaosen Ye, Tianhong Zhao, Wei Tu, Junxian Yu, Shengao Yi, Jinzhou Cao
2026 J jnl
Geo spatial Inf. Sci.
Junxian Yu, Wei Tu, Xiana Chen, Mingxiao Li, Dongsheng Chen, Rui Cao, Qingquan Li
2026 J jnl
Geo spatial Inf. Sci.
Jiasong Zhu, Jingbiao Chen, Mingxiao Li, Wei Tu
2026 J jnl
Inf. Fusion
Xiangxu Wang, Jinzhou Cao, Tianhong Zhao, Bowen Zhang, Guanzhou Chen, Zhenhui Li, Haolin Chen, Wei Tu, Qingquan Li
2026 A* conf
AAAI
Xiangxu Wang, Tianhong Zhao, Wei Tu, Bowen Zhang, Guanzhou Chen, Jinzhou Cao
2025 conf
SIGMOD Conference Companion
Tsz Nam Chan, Bojian Zhu, Dingming Wu, Yun Peng, Leong Hou U, Wei Tu, Ruisheng Wang
2025 conf
KDD (2)
Yue Zhong, Tsz Nam Chan, Leong Hou U, Dingming Wu, Wei Tu, Ruisheng Wang, Joshua Zhexue Huang
2025 J jnl
Inf. Fusion
Jinzhou Cao, Xiangxu Wang, Guanzhou Chen, Wei Tu, Xiaole Shen, Tianhong Zhao, Jiashi Chen, Qingquan Li
2025 J jnl
CoRR
Rui Cao, Wei Tu, Dongsheng Chen, Wenyu Zhang
2025 J jnl
Comput. Environ. Urban Syst.
Rui Cao, Wei Tu, Dongsheng Chen, Wenyu Zhang
2025 J jnl
Urban Inform.
Yang Zhou, Jean-Claude Thill, Xingjian Liu, Chen Zhong, Wei Tu
2025 J jnl
Geo spatial Inf. Sci.
Yuhao Liu, Songbo Wu, Bochen Zhang, Zhen Peng, Jiayuan Zhang, Chisheng Wang, Wei Tu, Zhipeng Chen, Mi Jiang, Xiao Cheng, Jiasong Zhu, Qingquan Li
2025 J jnl
Comput. Environ. Urban Syst.
Jinzhou Cao, Xianyu Cao, Wei Tu, Xiaoliang Tan, Tong Wang, Guanzhou Chen, Xiaodong Zhang, Qingquan Li
2025 J jnl
CoRR
Shengao Yi, Xiaojiang Li, Wei Tu, Tianhong Zhao
2025 J jnl
Comput. Environ. Urban Syst.
Tianhong Zhao, Xiucheng Liang, Filip Biljecki, Wei Tu, Jinzhou Cao, Xiaojiang Li, Shengao Yi
2025 J jnl
Geo spatial Inf. Sci.
Bichen Fang, Mingxiao Li, Zhengdong Huang, Yang Yue, Wei Tu, Renzhong Guo
2025 J jnl
CoRR
Xiangxu Wang, Tianhong Zhao, Wei Tu, Bowen Zhang, Guanzhou Chen, Jinzhou Cao
2025 J jnl
CoRR
Jinzhou Cao, Xiangxu Wang, Jiashi Chen, Wei Tu, Zhenhui Li, Xindong Yang, Tianhong Zhao, Qingquan Li
2024 J jnl
IEEE Internet Things J.
Lianjie Zhou, Wei Tu, Qingquan Li, Dongjie Guan
2024 J jnl
IEEE Trans. Geosci. Remote. Sens.
Chisheng Wang, Chuanhua Zhu, Xiang-Sheng Wang, Wei Tu, Qingquan Li
2024 J jnl
Int. J. Geogr. Inf. Sci.
Wei Tu, Haoyu Ye, Ke Mai, Meng Zhou, Jincheng Jiang, Tianhong Zhao, Shengao Yi, Qingquan Li
2024 J jnl
Comput. Environ. Urban Syst.
Xiana Chen, Wei Tu, Junxian Yu, Rui Cao, Shengao Yi, Qingquan Li
2024 J jnl
Int. J. Geogr. Inf. Sci.
Wei Tu, Wei Gao, Mingxiao Li, Yao Yao, Biao He, Zhengdong Huang, Jie Zhang, Renzhong Guo
2024 J jnl
Geo spatial Inf. Sci.
Jinzhou Cao, Wei Tu, Rui Cao, Qi-Li Gao, Guanzhou Chen, Qingquan Li
2023 J jnl
Comput. Environ. Urban Syst.
Yingjing Huang, Fan Zhang, Yong Gao, Wei Tu, Fabio Duarte, Carlo Ratti, Diansheng Guo, Yu Liu
2023 J jnl
Int. J. Geogr. Inf. Sci.
Tianhong Zhao, Zhengdong Huang, Wei Tu, Filip Biljecki, Long Chen
2023 J jnl
Remote. Sens.
Runqin Deng, Meng Zhou, Yinni Huang, Wei Tu
2023 J jnl
Int. J. Appl. Earth Obs. Geoinformation
Rui Cao, Cai Liao, Qing Li, Wei Tu, Rui Zhu, Nianxue Luo, Guoping Qiu, Wenzhong Shi
2023 J jnl
Int. J. Geogr. Inf. Sci.
Zhensheng Wang, Feidong Lu, Zhaohui Liu, Wei Tu, Ke Nie, Qingyun Du, Qingquan Li, Zhiqiang Wu
2023 conf
Geoinformatics
Yunjie Zhang, Siuming Lo, Shengjun Tang, Wei Tu
2023 J jnl
Comput. Environ. Urban Syst.
Tianhong Zhao, Xiucheng Liang, Wei Tu, Zhengdong Huang, Filip Biljecki
2023 J jnl
Remote. Sens.
Xiping Yang, Zhiyuan Zhao, Chaoyang Shi, Lin Luo, Wei Tu
2022 J jnl
IEEE Internet Things J.
Lianjie Zhou, Wei Tu, Chisheng Wang, Qingquan Li
2022 ch.
Artificial Intelligence, Machine Learning, and Optimization Tools for Smart Cities
Xiaojiang Li, Yuji Yoshimura, Wei Tu, Carlo Ratti
2022 J jnl
Int. J. Appl. Earth Obs. Geoinformation
Dongsheng Chen, Wei Tu, Rui Cao, Yatao Zhang, Biao He, Chisheng Wang, Tiezhu Shi, Qingquan Li
2022 J jnl
IEEE Trans. Intell. Transp. Syst.
Jianghai Liao, Yuanhao Yue, Dejin Zhang, Wei Tu, Rui Cao, Qin Zou, Qingquan Li
2022 J jnl
Comput. Environ. Urban Syst.
Tianhong Zhao, Zhengdong Huang, Wei Tu, Biao He, Rui Cao, Jinzhou Cao, Mingxiao Li
2022 J jnl
Comput. Environ. Urban Syst.
Bowen Zhang, Chen Zhong, Qi-Li Gao, Zahratu Shabrina, Wei Tu
2022 J jnl
Geo spatial Inf. Sci.
Wei Tu, Tingting Zhu, Chen Zhong, Xiaohu Zhang, Yang Xu, Qingquan Li
2022 J jnl
IEEE Trans. Intell. Transp. Syst.
Jincheng Jiang, Wei Tu, Hui Kong, Wei Zeng, Rui Zhang, Milan Konecný
2022 J jnl
IEEE Trans. Intell. Transp. Syst.
Tianhong Zhao, Wei Tu, Zhixiang Fang, Xiaofan Wang, Zhengdong Huang, Shengwu Xiong, Meng Zheng
2022 J jnl
IEEE Trans. Intell. Transp. Syst.
Yuting Xie, Yachen Zhang, Long Chen, Hui Cheng, Wei Tu, Dongpu Cao, Qingquan Li
2022 J jnl
Remote. Sens.
Mingxiao Li, Wei Tu, Feng Lu
2022 conf
GeoSim@SIGSPATIAL
Mingxiao Li, Wei Gao, Wei Tu, Jun Yue, Zhengdong Huang, Qingquan Li
2021 J jnl
Int. J. Geogr. Inf. Sci.
Zhensheng Wang, Yang Yue, Biao He, Ke Nie, Wei Tu, Qingyun Du, Qingquan Li
2021 J jnl
IEEE Internet Things J.
Baoding Zhou, Tianjing Zheng, Jincai Huang, Yunfei Zhang, Wei Tu, Qingquan Li, Min Deng
2021 J jnl
Environ. Model. Softw.
Lianjie Zhou, Qingquan Li, Wei Tu, Chisheng Wang
2021 conf
SIGSPATIAL/GIS
Yu Gu, Wei Tu, Qingquan Li, Tianhong Zhao, Dingyi Zhao, Song Zhu, Jiasong Zhu
2021 J jnl
Int. J. Geogr. Inf. Sci.
Mingxiao Li, Song Gao, Feng Lu, Kang Liu, Hengcai Zhang, Wei Tu
2021 J jnl
IEEE Trans. Ind. Informatics
Wei Tu, Ke Mai, Yatao Zhang, Yang Xu, Jincai Huang, Min Deng, Long Chen, Qingquan Li
2021 conf
Geoinformatics
Lianjie Zhou, Wei Tu, Qingquan Li
2021 J jnl
IEEE Geosci. Remote. Sens. Lett.
Wei Tu, Yatao Zhang, Qingquan Li, Ke Mai, Jinzhou Cao
2021 J jnl
Trans. GIS
Qi-Li Gao, Yang Yue, Wei Tu, Jinzhou Cao, Qingquan Li
2021 J jnl
CoRR
Shuhui Gong, Xiaopeng Mo, Rui Cao, Yu Liu, Wei Tu, Ruibin Bai
2021 J jnl
ACM Trans. Intell. Syst. Technol.
Ling Huang, Xing-Xing Liu, Shu-Qiang Huang, Chang-Dong Wang, Wei Tu, Jia-Meng Xie, Shuai Tang, Wendi Xie
2020 J jnl
IEEE Wirel. Commun. Lett.
Baoding Zhou, Wei Tu, Ke Mai, Weixing Xue, Wei Ma, Qingquan Li
2020 J jnl
IEEE Access
Lianjie Zhou, Qingquan Li, Wei Tu
2020 J jnl
Int. J. Digit. Earth
Jizhe Xia, Sicheng Huang, Shaobiao Zhang, Xiaoming Li, Jianrong Lyu, Wenqun Xiu, Wei Tu
2020 J jnl
CoRR
Iva Bojic, Dániel Kondor, Wei Tu, Ke Mai, Paolo Santi, Carlo Ratti
2020 J jnl
IEEE Internet Things J.
Wei Tu, Tianhong Zhao, Baoding Zhou, Jincheng Jiang, Jizhe Xia, Qingquan Li
2020 J jnl
IEEE Trans. Intell. Transp. Syst.
Qiuping Li, Shaobo Zhong, Zhixiang Fang, Lin Liu, Wei Tu, Bi Yu Chen
2020 J jnl
Comput. Environ. Urban Syst.
Wei Tu, Tingting Zhu, Jizhe Xia, Yulun Zhou, Yani Lai, Jincheng Jiang, Qingquan Li
2019 J jnl
Int. J. Geogr. Inf. Sci.
Jincheng Jiang, Qingquan Li, Wei Tu, Shih-Lung Shaw, Yang Yue
2019 J jnl
Comput. Environ. Urban Syst.
Yatao Zhang, Qingquan Li, Wei Tu, Ke Mai, Yao Yao, Yiyong Chen
2019 conf
GeoAI@SIGSPATIAL
Ke Mai, Wei Tu, Qingquan Li, Tianzhong Zhao, Yatao Zhang, Haoyu Ye
2019 J jnl
Remote. Sens.
Jiasong Zhu, Siyuan Chen, Wei Tu, Ke Sun
2019 J jnl
Comput. Environ. Urban Syst.
Yang Xu, Dachi Chen, Xiaohu Zhang, Wei Tu, Yuanyang Chen, Yu Shen, Carlo Ratti
2018 J jnl
IEEE Access
Baoding Zhou, Qingquan Li, Guanxun Zhai, Qingzhou Mao, Jun Yang, Wei Tu, Weixing Xue, Long Chen
2018 J jnl
Remote. Sens.
Rui Cao, Jiasong Zhu, Wei Tu, Qingquan Li, Jinzhou Cao, Bozhi Liu, Qian Zhang, Guoping Qiu
2018 J jnl
Remote. Sens.
Wei Tu, Zhongwen Hu, Lefei Li, Jinzhou Cao, Jincheng Jiang, Qiuping Li, Qingquan Li
2018 J jnl
ISPRS Int. J. Geo Inf.
Chen Zhong, Shi Zeng, Wei Tu, Mitsuo Yoshida
2018 J jnl
ISPRS Int. J. Geo Inf.
Qiuping Li, Wei Tu, Li Zhuo
2018 J jnl
Remote. Sens.
Peng Liu, Xiaofei Chen, Zhenhong Li, Zhenguo Zhang, Jiankuan Xu, Wanpeng Feng, Chisheng Wang, Zhongwen Hu, Wei Tu, Hongzhong Li
2017 J jnl
Sensors
Baoding Zhou, Qingquan Li, Qingzhou Mao, Wei Tu
2017 J jnl
Trans. GIS
Wei Tu, Qingquan Li, Qiuping Li, Jiasong Zhu, Baoding Zhou, Bi Yu Chen
2017 J jnl
Int. J. Geogr. Inf. Sci.
Wei Tu, Jinzhou Cao, Yang Yue, Shih-Lung Shaw, Meng Zhou, Zhensheng Wang, Xiaomeng Chang, Yang Xu, Qingquan Li
2017 J jnl
ISPRS Int. J. Geo Inf.
Jincheng Jiang, Qingquan Li, Lixin Wu, Wei Tu
2015 J jnl
ISPRS Int. J. Geo Inf.
Wei Tu, Qingquan Li, Zhixiang Fang, Baoding Zhou
2015 J jnl
IEEE Trans. Intell. Transp. Syst.
Baoding Zhou, Qingquan Li, Qingzhou Mao, Wei Tu, Xing Zhang, Long Chen
2015 J jnl
IEEE Trans. Hum. Mach. Syst.
Baoding Zhou, Qingquan Li, Qingzhou Mao, Wei Tu, Xing Zhang
2015 conf
Geoinformatics
Meng Zhou, Wei Tu, Qingquan Li, Yang Yue, Xiaomeng Chang
2013 J jnl
Int. J. Geogr. Inf. Sci.
Zhixiang Fang, Wei Tu, Qingquan Li, Shih-Lung Shaw, Shunqing Chen, Bi Yu Chen
2010 conf
Geoinformatics
Wei Tu, Zhixiang Fang, Qingquan Li
2009 J jnl
IEEE Trans. Circuits Syst. Video Technol.
Wei Tu, Eckehard G. Steinbach
2007 A conf
ICME
Rajyalakshmi Mahalingam, Wei Tu, Eckehard G. Steinbach
2005 A conf
ACSAC
Benjamin Schwarz, Hao Chen, David A. Wagner, Jeremy Lin, Wei Tu, Geoff Morrison, Jacob West
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()