Xiaodong Chen

133 papers A* 5A 1B 2C 2Misc 2Journal 91Unranked 29
YearRankTypeTitle / Venue / Authors
2026 J jnl
Expert Syst. Appl.
Zhengliang Zhang, Yachen Wei, Xin Rao, Liyang Yu, Wen Sun, Ruixue Li, Xiaoshuai Zhang, Xiaodong Chen, Xingru Huang
2026 J jnl
CoRR
Shaopeng Chen, Chuyue Xie, Huimin Ren, Shaozong Zhang, Han Zhang, Ruobing Cheng, Zhiqiang Cao, Zehao Ju, Yu Gao, Jie Ding, Xiaodong Chen, Xuewu Jiao, Shuanglong Li, Liu Lin
2026 J jnl
Earth Sci. Informatics
Min Zhao, Qianqian He, Xiaodong Chen, Miaomiao Zhang
2026 J jnl
Graphs Comb.
Xiaodong Chen, Wenyue Zhang, Shou-Jun Xu
2025 J jnl
Discret. Appl. Math.
Xiaodong Chen, Tianhao Li, Jiayuan Zhang
2025 J jnl
IEEE Trans. Dependable Secur. Comput.
Liu Liu, Xinwen Fu, Xiaodong Chen, Jianpeng Wang, Zhongjie Ba, Feng Lin, Li Lu, Kui Ren
2025 J jnl
CoRR
Tieyuan Chen, Xiaodong Chen, Haoxing Chen, Zhenzhong Lan, Weiyao Lin, Jianguo Li
2025 A* conf
ICRA
Changhao Tian, Annan Wang, Han Fan, Thomas Wiedemann, Yifei Luo, Le Yang, Weisi Lin, Achim J. Lilienthal, Xiaodong Chen
2025 J jnl
Circuits Syst. Signal Process.
Ziyang Liu, Xiaodong Chen, Chongyu Shi, Gang Wang, Qingtang Su
2025 J jnl
Discret. Appl. Math.
Xiaodong Chen, Jiayuan Zhang, Liming Xiong, Guifu Su
2025 J jnl
CoRR
Haoyuan Wu, Haoxing Chen, Xiaodong Chen, Zhanchao Zhou, Tieyuan Chen, Yihong Zhuang, Guoshan Lu, Zenan Huang, Junbo Zhao, Lin Liu, Zhenzhong Lan, Bei Yu, Jianguo Li
2025 J jnl
Int. J. Digit. Earth
Kunpeng Shi, Hao Ding, Xiaodong Chen, Xizhi Hu, Weiping Jiang, Heping Sun
2025 J jnl
CoRR
Zhanchao Zhou, Xiaodong Chen, Haoxing Chen, Zhenzhong Lan, Jianguo Li
2025 J jnl
CoRR
Yuxuan Hu, Jing Zhang, Xiaodong Chen, Zhe Zhao, Cuiping Li, Hong Chen
2025 J jnl
IEEE Trans. Consumer Electron.
En Mou, Huiqian Wang, Xiaodong Chen, Zhangyong Li, Lisha Zhong, Shuai Xia, Guanjie Zeng, Yu Pang
2025 conf
CVPR Workshops
Alicia Li, Xiaodong Chen, Bohao Liang, Qian Bao, Wu Liu
2025 J jnl
CoRR
Xiaodong Chen, Mingming Ha, Zhenzhong Lan, Jing Zhang, Jianguo Li
2025 J jnl
BMC Medical Imaging
Ang Shi, Huaiqing Zhi, Dongze Wu, Wanda Cai, Yajin Chen, Xiaodong Chen, Chenbin Chen, Xinxin Yang, Jingwei Zheng, Hao Chen, Weiteng Zhang, Xian Shen
2025 J jnl
J. Comput. Des. Eng.
Xiaoshuai Huo, Tanghong Liu, Xiaodong Chen, Zhengwei Chen, Xinran Wang
2025 conf
ACL (1)
Xiaodong Chen, Yuxuan Hu, Xiaokang Zhang, Yanling Wang, Cuiping Li, Hong Chen, Jing Zhang
2025 J jnl
CoRR
Yuxuan Hu, Xiaodong Chen, Cuiping Li, Hong Chen, Jing Zhang
2025 A* conf
ICLR
Xiaodong Chen, Yuxuan Hu, Jing Zhang, Yanling Wang, Cuiping Li, Hong Chen
2025 J jnl
IEEE Access
Hui Liu, Xiaowan Li, Chuang Zhang, Xiaodong Chen, Weipeng Tai
2025 J jnl
CoRR
Yanling Wang, Yihan Zhao, Xiaodong Chen, Shasha Guo, Lixin Liu, Haoyang Li, Yong Xiao, Jing Zhang, Qi Li, Ke Xu
2024 conf
CBD
Yongkang Li, Xiaodong Chen, Xiaohui Wan, Jin Sun, Peng Zheng, Yunchang Wang, Zebin Wu
2024 J jnl
Int. J. Wavelets Multiresolution Inf. Process.
Shaobo Liu, Tian Xia, Xiaodong Chen, Hui Li, Guanghui Yuan, Dong Yang
2024 J jnl
Int. J. Appl. Earth Obs. Geoinformation
Zhiyuan Li, Fengxiang Jin, Jian Wang, Zhenyu Zhang, Lei Zhu, Wenxiao Sun, Xiaodong Chen
2024 J jnl
Educ. Inf. Technol.
Gang Yang, Wei Zhou, Huimin Zhou, Jiawen Li, Xiaodong Chen, Yun-Fang Tu
2024 conf
PEAI
Xiaoting Huang, Tingying Zhang, Xiaodong Chen, Ge Zhang, Xinyu Xu, Yuxin Zhang, Xinyu Chen
2024 J jnl
Frontiers Robotics AI
Saikrishna Dontu, Elgar Kanhere, Thileepan Stalin, Audelia Gumarus Dharmawan, Chidanand Hegde, Jiangtao Su, Xiaodong Chen, Shlomo Magdassi, Gim Song Soh, Pablo Valdivia y Alvarado
2024 J jnl
Axioms
Jiatong Cui, Tianhao Li, Jiayuan Zhang, Xiaodong Chen, Liming Xiong
2024 J jnl
CoRR
Xiaodong Chen, Yuxuan Hu, Jing Zhang
2024 J jnl
Int. J. Netw. Virtual Organisations
Xiaodong Chen
2024 J jnl
Remote. Sens.
Junjie Wu, Du Xiao, Bingrui Du, Yuge Liu, Qingquan Zhi, Xingchun Wang, Xiaohong Deng, Xiaodong Chen, Yi Zhao, Yue Huang
2024 J jnl
IEEE Trans. Engineering Management
Zhuang Miao, Anda Guo, Xiaodong Chen, Pengyu Zhu
2024 J jnl
Sensors
Hui Liu, Chuang Zhang, Xiaodong Chen, Weipeng Tai
2024 J jnl
BMC Medical Imaging
En Mou, Huiqian Wang, Xiaodong Chen, Zhangyong Li, Enling Cao, Yuanyuan Chen, Zhiwei Huang, Yu Pang
2024 conf
ACL (Findings)
Yuxuan Hu, Jing Zhang, Zhe Zhao, Chen Zhao, Xiaodong Chen, Cuiping Li, Hong Chen
2024 J jnl
CoRR
Xiaodong Chen, Yuxuan Hu, Jing Zhang, Xiaokang Zhang, Cuiping Li, Hong Chen
2024 A* conf
ICRA
Jana Egli, Benedek Forrai, Thomas Buchner, Jiangtao Su, Xiaodong Chen, Robert K. Katzschmann
2024 J jnl
CoRR
Jana Egli, Benedek Forrai, Thomas Buchner, Jiangtao Su, Xiaodong Chen, Robert K. Katzschmann
2024 J jnl
Comput. Biol. Medicine
Chenbin Chen, Xietao Chen, Yuanbo Hu, Bujian Pan, Qunjia Huang, Qiantong Dong, Xiangyang Xue, Xian Shen, Xiaodong Chen
2023 conf
ACM TUR-C
Kai Guo, Zhizun Qin, Xiaodong Chen, Haoran Wang, Hanjiang Luo
2023 A* conf
CCS
Liu Liu, Xinwen Fu, Xiaodong Chen, Jianpeng Wang, Zhongjie Ba, Feng Lin, Li Lu, Kui Ren
2023 J jnl
IEEE Robotics Autom. Lett.
Ke Ma, Xiaodong Chen, Jie Zhang, Zekai Xie, Jianing Wu, Jinxiu Zhang
2023 J jnl
IEEE Internet Things J.
Ning Jin, Gang Yang, Ying-Chang Liang, Songbo Fu, Xiaodong Chen
2023 J jnl
Discret. Math.
Xiaodong Chen, Guantao Chen
2023 J jnl
Inf. Sci.
Bo Liu, Weibin Li, Yanshan Xiao, Xiaodong Chen, Laiwang Liu, Changdong Liu, Kai Wang, Peng Sun
2023 J jnl
CoRR
Qifeng Lin, Rui Li, Feilong Zhang, Kazuki Kai, Ong Zong Chen, Xiaodong Chen, Hirotaka Sato
2023 J jnl
Intell. Serv. Robotics
Yanting Lan, Xiaodong Chen
2023 J jnl
Graphs Comb.
Xiaodong Chen, Jingjing Li, Fuliang Lu
2023 J jnl
Discret. Appl. Math.
Xiaodong Chen, Qinghai Liu, Xiwu Yang
2022 J jnl
Graphs Comb.
Xiaodong Chen, Xiyao Guo, Xiwu Yang
2022 J jnl
Knowl. Based Syst.
Bo Liu, Changdong Liu, Yanshan Xiao, Laiwang Liu, Weibin Li, Xiaodong Chen
2022 J jnl
Inf. Sci.
Bo Liu, Laiwang Liu, Yanshan Xiao, Changdong Liu, Xiaodong Chen, Weibin Li
2022 J jnl
Comput. Methods Programs Biomed.
Yang Bai, Dan Li, Qiongyu Duan, Xiaodong Chen
2022 J jnl
Comput. Electron. Agric.
Yifan Guo, Yanting Lan, Xiaodong Chen
2022 conf
VTC Fall
Ning Jin, Fanyi Shu, Gang Yang, Ying-Chang Liang, Xiaodong Chen
2022 J jnl
Intell. Serv. Robotics
Yanting Lan, Xiaodong Chen
2022 J jnl
Intell. Serv. Robotics
Yanting Lan, Xiaodong Chen
2022 conf
VTC Fall
Ning Jin, Yating Liao, Gang Yang, Ying-Chang Liang, Xiaodong Chen
2022 J jnl
Int. J. Comput. Assist. Radiol. Surg.
Zezhong Li, Kangming Chen, Peng Liu, Xiaodong Chen, Guoyan Zheng
2022 J jnl
ISPRS Int. J. Geo Inf.
Xiaodong Chen, Zhaoping Yang, Tian Wang, Fang Han
2022 conf
MLSys
Vijay Janapa Reddi, David Kanter, Peter Mattson, Jared Duke, Thai Nguyen, Ramesh Chukka, Kenneth Shiring, Koan-Sin Tan, Mark Charlebois, William Chou, Mostafa El-Khamy, Jungwook Hong, Tom St. John, Cindy Trinh, Michael Buch, Mark Mazumder, Relja Markovic, Thomas Atta-Fosu, Fatih Çakir, Masoud Charkhabi, Xiaodong Chen, Cheng-Ming Chiang, Dave Dexter, Terry Heo, Guenther Schmuelling, Maryam Shabani, Dylan Zika
2022 conf
ICCSIE
Zhijian Ye, Yiang Li, Zhiping Zhang, Xiaodong Chen, Dejia Liu, Weichuan Ni
2022 J jnl
Neurocomputing
Shiyuan Yang, Yi Wang, Huaiyu Cai, Xiaodong Chen
2022 J jnl
Int. J. Appl. Earth Obs. Geoinformation
Chengqian Zhang, Xiaodong Chen, Shunying Ji
2022 J jnl
Discret. Math.
Xiaodong Chen, Mingchu Li, Liming Xiong
2022 conf
ICCT
Lurui Yang, Kai Liu, Yongqing Huang, Xiaofeng Duan, Qi Wang, Xiaoxia Du, Xiaodong Chen
2021 J jnl
Inf. Sci.
Bo Liu, Xiaodong Chen, Yanshan Xiao, Weibin Li, Laiwang Liu, Changdong Liu
2021 J jnl
Remote. Sens.
Jiachen Zhang, Weisong Wen, Feng Huang, Xiaodong Chen, Li-Ta Hsu
2021 J jnl
Comput. Methods Programs Biomed.
Ying Su, Dan Li, Xiaodong Chen
2021 J jnl
Remote. Sens.
Xiongwei Tang, Rumeng Guo, Jianqiao Xu, Heping Sun, Xiaodong Chen, Jiangcun Zhou
2021 J jnl
Int. J. Digit. Earth
Miao Yu, Peng Lu, Zhiyuan Li, Zhijun Li, Qingkai Wang, Xiaowei Cao, Xiaodong Chen
2021 conf
ECOC
Xin-Tao He, Meng-Yu Li, Xiaodong Chen, Jian-Wen Dong
2020 J jnl
IEEE Trans. Multim.
Guangming Sun, Bufan Shi, Xiaodong Chen, Andrey S. Krylov, Yong Ding
2020 J jnl
CoRR
Vijay Janapa Reddi, David Kanter, Peter Mattson, Jared Duke, Thai Nguyen, Ramesh Chukka, Kenneth Shiring, Koan-Sin Tan, Mark Charlebois, William Chou, Mostafa El-Khamy, Jungwook Hong, Michael Buch, Cindy Trinh, Thomas Atta-Fosu, Fatih Çakir, Masoud Charkhabi, Xiaodong Chen, Jimmy Chiang, Dave Dexter, Woncheol Heo, Guenther Schmuelling, Maryam Shabani, Dylan Zika
2020 J jnl
IEEE Access
Ying Zhang, Weihong Yu, Xiaodong Chen, Jianhui Jiang
2020 J jnl
Graphs Comb.
Xiaodong Chen, Qing Ji, Mingda Liu
2020 J jnl
IEEE Access
Bo Ren, Jinsong Li, Yongkang Zheng, Xiaodong Chen, Yibing Zhao, Haiyang Zhang, Chao Zheng
2020 J jnl
Int. J. Distributed Sens. Networks
Chu Rouxia, Xiaodong Chen, Tao Shifang, Donghai Yang
2020 conf
OFC
Xin-Tao He, Meng-Yu Li, Hao-Yang Qiu, Xiaodong Chen, Jian-Wen Dong
2019 J jnl
IEEE Access
Xiaogang Xu, Bufan Shi, Zijin Gu, Ruizhe Deng, Xiaodong Chen, Andrey S. Krylov, Yong Ding
2019 J jnl
Int. J. Online Biomed. Eng.
Yan Ting Lan, Xiaodong Chen
2019 J jnl
Grey Syst. Theory Appl.
Xiaodong Chen, Desheng Pei, Liping Li
2019 J jnl
IEEE Access
Min Zhang, Jiantong Zhang, Ruolin Ma, Xiaodong Chen
2019 conf
ICIRA (5)
Junjie Yang, Hao Sun, Dongping Wu, Xiaodong Chen, Changhong Wang
2019 conf
ICDLT
Hongbing Xiao, Peiyuan Guo, Xiaodong Chen
2019 conf
CACRE
Zijiao Han, Xiaodong Chen, Chenqi Wang, Xin Wang, Baoshi Wang, Jiapeng Wei
2019 J jnl
IET Image Process.
Yong Ding, Yang Zhao, Xiaodong Chen, Xiaolei Zhu, Andrey S. Krylov
2019 J jnl
Big Data Cogn. Comput.
Quanchun Jiang, Olamide Timothy Tawose, Songwen Pei, Xiaodong Chen, Linhua Jiang, Jiayao Wang, Dongfang Zhao
2018 conf
EMBC
Pingao Huang, Zhenxin Chen, Menglong Fu, Hui Wang, Oluwarotimi Williams Samuel, Zhiyuan Liu, Yongmei Hu, Peng Fang, Shixiong Chen, Xiaodong Chen, Guanglin Li
2018 J jnl
IEEE Access
Yong Ding, Ruizhe Deng, Xin Xie, Xiaogang Xu, Yang Zhao, Xiaodong Chen, Andrey S. Krylov
2018 J jnl
Ars Comb.
Xiaodong Chen, Mingchu Li, Fuliang Lu
2018 J jnl
IEEE Access
Guangming Sun, Yong Ding, Ruizhe Deng, Yang Zhao, Xiaodong Chen, Andrey S. Krylov
2017 J jnl
Discret. Math.
Guantao Chen, Xiaodong Chen, Yue Zhao
2017 conf
CISP-BMEI
Jinqing Li, Xiaoqiang Di, Xingchen Liu, Xiaodong Chen
2017 J jnl
IEEE Trans. Knowl. Data Eng.
Guojie Song, Yuanhao Li, Xiaodong Chen, Xinran He, Jie Tang
2017 J jnl
Ars Comb.
Xiaodong Chen, Mingchu Li, Meijin Xu
2016 J jnl
Int. J. Online Eng.
Yan Ting Lan, Jiinying Huang, Xiaodong Chen
2016 J jnl
Ars Comb.
Xiaodong Chen, Mingchu Li, Wei Liao, Hajo Broersma
2016 J jnl
计算机科学
Xiaodong Chen, Lijuan Sun, Chong Han, Jian Guo
2015 J jnl
IEEE ACM Trans. Comput. Biol. Bioinform.
Beichen Wang, Xiaodong Chen, Hiroshi Mamitsuka, Shanfeng Zhu
2015 A conf
SDM
Xiaodong Chen, Guojie Song, Xinran He, Kunqing Xie
2015 J jnl
Graphs Comb.
Xiaodong Chen, Mingchu Li, Xin Ma
2013 C conf
ICIS
Weijie Zhao, Xiaodong Chen, Ji Cheng, Linhua Jiang
2013 J jnl
Graphs Comb.
Xiaodong Chen, Mingchu Li, Xin Ma, Xinxin Fan
2013 conf
ICCT
Chong Han, Lijuan Sun, Jian Guo, Xiaodong Chen
2013 J jnl
Secur. Commun. Networks
Xinxin Fan, Mingchu Li, Hui Zhao, Xiaodong Chen, Zhenzhou Guo, Dong Jiao, Weifeng Sun
2013 conf
ICCI*CC
Guihua Wen, Xiaodong Chen, Lijun Jiang, Haisheng Li
2012 conf
ISORC Workshops
Linhua Jiang, Chao Peng, Haibin Cai, Yixiang Chen, Xiaodong Chen
2012 J jnl
Neural Comput. Appl.
Tingting Chen, Ankur Bansal, Sheng Zhong, Xiaodong Chen
2012 J jnl
J. Appl. Math.
Jianjun Sun, Bin Huang, Xiaodong Chen, Lihong Cui
2011 J jnl
J. Graph Theory
Mingchu Li, Xiaodong Chen, Hajo Broersma
2011 conf
ACPR
Xuesong Yin, Xiaodong Chen, Xiaofang Ruan, Yarong Huang
2011 J jnl
J. Quantum Inf. Sci.
He-Zhou Wang, He-Xiang He, Jie Feng, Xiaodong Chen, Wei Lin
2010 B conf
CEC
Shaohua Liu, Yinglong Ma, Tianlu Mao, Junsheng Yu, Siyu Liu, Xiaodong Chen
2010 J jnl
J. Digit. Content Technol. its Appl.
Dongming Chen, Jing Wang, Xiaodong Chen, Xiaowei Xu
2010 Misc conf
DCAI
Ruxiu Zhong, Fei Ji, Fangjiong Chen, Shangkun Xiong, Xiaodong Chen
2009 conf
CSIE (3)
Xishuang Dong, Xiaodong Chen, Yi Guan, Zhiming Yu, Sheng Li
2009 conf
NEMS
Di Wang, Xiaodong Chen, Qun Zhang, Yu Liu, Jincheng Liu, Lijiang Hu
2009 conf
NEMS
Yu Liu, Zushun Lu, Xiaodong Chen, Di Wang, Jincheng Liu, Lijiang Hu
2009 J jnl
IET Commun.
Shuxian Chen, Marianna Setta, Xiaodong Chen, Clive Parini
2007 conf
EUSIPCO
Mahmoud Hadef, Adel Daas, Stephan Weiss, Josh Reiss, Xiaodong Chen
2006 J jnl
Afr. J. Inf. Commun. Technol.
Xiaodong Chen, Jianxin Liang, Pengcheng Li, Choo C. Chiau
2001 Misc conf
International Conference on Computational Science (2)
Angela Violi, Xiaodong Chen, Gary Lindstrom, Eric Eddings, Adel F. Sarofim
2000 B conf
DaWaK
Xiaodong Chen, Ilias Petrounias
2000 A* conf
ICDE
Xiaodong Chen, Ilias Petrounias
1999 conf
PKDD
Xiaodong Chen, Ilias Petrounias
1999
Xiaodong Chen
1998 C conf
DEXA
Xiaodong Chen, Ilias Petrounias
1998 conf
IADT
Xiaodong Chen, Ilias Petrounias, Heather Heathfield
1998 conf
PKDD
Xiaodong Chen, Ilias Petrounias
redb/ingestor.py
← Index redb/ingestor.py python
"""
REDB Ingestor - Core ingestion orchestration.

This module contains the Ingestor class which orchestrates the entire
sample processing pipeline: querying catalogs, downloading from S3,
dispatching to workers, and collecting results.

The actual implementation is split across focused modules:
- redb.queries: Database query and deduplication functions
- redb.s3_utils: S3/MinIO client and file operations
- redb.workers: File processing and worker functions
- redb.logging_utils: Logging setup and ImportResult enum

For backward compatibility, all public names from these modules are
re-exported here so that `from redb.ingestor import *` continues to work.
"""
import multiprocessing
from multiprocessing import Pool
from datetime import datetime
import os
import sys
import gc
import psutil
import time
import json
import tempfile
import warnings
from urllib3.exceptions import InsecureRequestWarning
from dotenv import load_dotenv

load_dotenv(override=True)

warnings.filterwarnings("ignore", category=InsecureRequestWarning, module="urllib3")
warnings.filterwarnings("ignore", category=UserWarning, module="elasticsearch")

# =============================================================================
# Re-exports for backward compatibility
# =============================================================================
# These imports ensure that `from redb.ingestor import X` and
# `@patch('redb.ingestor.X')` continue to work after the refactor.

from redb.logging_utils import (  # noqa: F401
    ImportResult,
    FileNameFormatter,
    setup_logger,
    logger_thread,
    setup_direct_logger,
)

from redb.queries import (  # noqa: F401
    get_supported_formats,
    get_db_catalog_connection,
    fetch_s3_objects_by_repository,
    fetch_s3_objects_by_date_range,
    fetch_analyzed_samples,
    is_in_db,
    is_in_code_db,
    is_in_db_bulk,
)

from redb.s3_utils import (  # noqa: F401
    get_minio_client,
    generate_s3_key_from_hash,
    download_s3_object,
    extract_fat_slices,
)

from redb.workers import (  # noqa: F401
    process_s3_file,
    process_file,
    _process_file_internal,
    process_zip_file,
    process_7zip_file,
    process_binary_file,
    worker,
    direct_s3_worker,
    is_binary_file,
    check_dotnet,
    check_high_swap,
    get_module_by_name,
    filter_selected_modules,
    _is_packed,
)

# Re-export settings for patches like @patch('redb.ingestor.settings')
from redb import settings  # noqa: F401

# Re-export hashlib and Magika for patches like @patch('redb.ingestor.hashlib')
import hashlib  # noqa: F401
try:
    from magika import Magika  # noqa: F401
except ImportError:
    pass

# Re-export py7zr for patches like @patch('redb.ingestor.py7zr')
try:
    import py7zr  # noqa: F401
except ImportError:
    pass


class Ingestor:
    def __init__(
        self,
        path=None,
        decompile=False,
        yara_scan=False,
        with_yara=False,
        repository="",
        index_prefix="",
        selected_modules=None,
        s3_mode=False,
        s3_notes=None,
        magika_filter=None,
        s3_solo=False,
        s3_solo_hash=None,
        s3_solo_key=None,
        dry_run=False,
        force=False,
        job_id=None,
        start_date=None,
        end_date=None,
        analyzed=False,
        decompile_modules=None,
        rerun=False,
    ):

        # Set multiprocessing start method as early as possible
        try:
            multiprocessing.set_start_method('spawn', force=True)
        except RuntimeError:
            current_method = multiprocessing.get_start_method()
            if current_method != 'spawn':
                print(f"[WARNING] Multiprocessing start method is {current_method}, not 'spawn'. This may cause issues.")

        self.path = path
        self.decompile = decompile
        self.yara_scan = yara_scan
        self.with_yara = with_yara
        self.repository = repository
        self.index_prefix = index_prefix
        self.selected_modules = selected_modules
        self.s3_mode = s3_mode
        self.s3_notes = s3_notes
        self.magika_filter = magika_filter
        self.s3_solo = s3_solo
        self.s3_solo_hash = s3_solo_hash
        self.s3_solo_key = s3_solo_key
        self.dry_run = dry_run
        # --rerun implies force at the worker level: the query already selects
        # only already-disassembled samples, so the per-file is_in_code_db
        # dedup check must be skipped or every sample gets skipped.
        self.force = force or rerun
        self.job_id = job_id
        self.start_date = start_date
        self.end_date = end_date
        self.analyzed = analyzed
        self.decompile_modules = decompile_modules or {"all"}
        self.rerun = rerun

        self.manager = multiprocessing.Manager()
        self.file_type_stats = self.manager.dict()
        self.total_results = self.manager.dict({result: 0 for result in ImportResult})

        self.today = datetime.today().strftime("%Y%m%dT%H%M%S")
        log_base_path = os.getenv("LOG_FILE_PATH", "/app/logs/")
        if not log_base_path.endswith("/"):
            log_base_path += "/"
        index_suffix = self.index_prefix.upper() if self.index_prefix else "DEFAULT"
        self.log_file = log_base_path + f"{self.today}-{self.repository}-{index_suffix}.txt"

        with open(self.log_file, "a") as f:
            f.write(f"CMD: {' '.join(sys.argv)}\n")
            f.write(f"=== Ingestor started at {datetime.now()} ===\n")


    def restart_worker_pool(self):
        """Restart the worker pool to help address memory issues"""
        if hasattr(self, 'pool') and self.pool:
            try:
                print("[INFO] Restarting worker pool to address memory fragmentation")
                self.pool.close()
                self.pool.join()
                self.pool = None
            except:
                pass

        # Force garbage collection
        gc.collect(2)


    def _process_files_streaming(self, s3_files, temp_dir, total_files, parallel_proc, decompile=None):
        """Process files in a streaming fashion using direct process management."""
        import queue

        # Use the instance's decompile flag if not provided
        if decompile is None:
            decompile = self.decompile

        # Use local tracking for statistics
        completed_count = 0
        skipped_count = 0
        failed_count = 0
        correctly_processed = 0
        partially_processed = 0
        filetype_stats = {}

        # Create a result queue for workers to return their results
        result_queue = multiprocessing.Queue()

        # Create a process ID tracking dict
        active_processes = {}  # {proc_id: (process, start_time, s3_key)}

        mode_str = "decompile" if decompile else "analysis"
        print(f"[INFO] Starting streaming processing of {len(s3_files)} files with {parallel_proc} workers ({mode_str} mode)")

        # Process files
        file_index = 0
        # Use different timeouts based on mode
        if decompile:
            worker_timeout = int(os.getenv("DECOMPILE_WORKER_TIMEOUT", "2700"))
        else:
            worker_timeout = int(os.getenv("REDB_TIMEOUT", "1200"))

        # Main processing loop
        while file_index < len(s3_files) or active_processes:
            # Start new processes if we have capacity and files to process
            while len(active_processes) < parallel_proc and file_index < len(s3_files):
                s3_bucket, s3_key, first_seen = s3_files[file_index]
                file_number = file_index + 1

                # Create and start a new process
                p = multiprocessing.Process(
                    target=direct_s3_worker,
                    args=(
                        s3_bucket,
                        s3_key,
                        temp_dir,
                        decompile,  # Pass the actual decompile flag
                        self.index_prefix,
                        self.log_file,
                        file_number,
                        total_files,
                        self.selected_modules,
                        result_queue,
                        self.dry_run,
                        self.yara_scan,
                        self.with_yara,
                        self.force,
                        self.decompile_modules,
                        first_seen,
                    )
                )
                p.start()

                # Track the process
                active_processes[p.pid] = (p, time.time(), s3_key, file_number)
                file_index += 1

                # Small delay to avoid overloading
                time.sleep(0.05)

            # Check for completed processes
            try:
                # Poll the result queue with a timeout
                while True:
                    try:
                        result = result_queue.get(block=True, timeout=1)

                        # Process result
                        s3_key = result.get('s3_key')
                        status = result.get('status', 'FAILED')
                        filetype = result.get('filetype')
                        file_number = result.get('file_number')
                        worker_pid = result.get('worker_pid')

                        # Remove from active processes if present
                        if worker_pid in active_processes:
                            del active_processes[worker_pid]

                        # Update statistics
                        if status == 'CORRECTLY':
                            correctly_processed += 1
                        elif status == 'PARTIALLY':
                            partially_processed += 1
                        elif status == 'SKIPPED':
                            skipped_count += 1
                        else:  # Any other status is treated as failure
                            failed_count += 1

                        if filetype:
                            filetype_stats[filetype] = filetype_stats.get(filetype, 0) + 1

                        # Update progress
                        completed_count += 1
                        if completed_count % 50 == 0 or completed_count == 1:
                            print(f"[INFO] Completed {completed_count}/{total_files} files. Last: {s3_key}")
                            print(f"[INFO] Progress - OK: {correctly_processed}, Partial: {partially_processed}, Failed: {failed_count}, Skipped: {skipped_count}, Active: {len(active_processes)}")

                    except queue.Empty:
                        # No results in the queue, break and check for timeouts
                        break

                # Check for timed-out processes
                current_time = time.time()
                timed_out_pids = []

                for pid, (proc, start_time, s3_key, file_number) in active_processes.items():
                    runtime = current_time - start_time

                    # Check if process has exceeded timeout
                    if runtime > worker_timeout:
                        print(f"[WARNING] Process {pid} processing {s3_key} exceeded timeout ({runtime:.0f}s > {worker_timeout}s)")

                        # Terminate the process
                        try:
                            proc.terminate()
                            time.sleep(0.1)  # Give it a moment to terminate
                            if proc.is_alive():
                                # If still alive, force kill
                                proc.kill()
                        except:
                            pass

                        # Clean up any child processes
                        try:
                            parent = psutil.Process(pid)
                            for child in parent.children(recursive=True):
                                try:
                                    child.kill()
                                except:
                                    pass
                        except:
                            pass

                        # Mark as failed
                        failed_count += 1
                        completed_count += 1
                        timed_out_pids.append(pid)

                    # Check if process has terminated without returning a result
                    elif not proc.is_alive():
                        print(f"[WARNING] Process {pid} processing {s3_key} terminated without result")

                        # Mark as failed
                        failed_count += 1
                        completed_count += 1
                        timed_out_pids.append(pid)

                # Remove timed-out processes from tracking
                for pid in timed_out_pids:
                    if pid in active_processes:
                        del active_processes[pid]

                # Sleep briefly to avoid hogging CPU
                time.sleep(0.1)

            except Exception as e:
                print(f"[ERROR] Exception in main processing loop: {e}")
                time.sleep(1)  # Sleep to avoid tight loop on error

        # Set final results in the shared dictionaries
        self.total_results[ImportResult.CORRECTLY] = correctly_processed
        self.total_results[ImportResult.PARTIALLY] = partially_processed
        self.total_results[ImportResult.FAILED] = failed_count
        self.total_results[ImportResult.SKIPPED] = skipped_count

        for filetype, count in filetype_stats.items():
            self.file_type_stats[filetype] = count

        print(f"[INFO] Streaming processing completed. Processed {completed_count}/{total_files} files.")

        # Final cleanup
        killed = self._kill_all_python_processes()
        if killed > 0:
            print(f"[INFO] Killed {killed} lingering processes during final cleanup")


    def _kill_all_python_processes(self):
        """Kill all python worker processes."""
        killed_count = 0
        for proc in psutil.process_iter(['pid', 'name', 'cmdline']):
            try:
                if proc.info['name'] == 'python' and proc.info['cmdline']:
                    # Check if it's one of our processes
                    is_worker = False
                    for cmd in proc.info['cmdline']:
                        if 'deploy/redb/venv312bin/python' in cmd and 'multiprocessing' in cmd:
                            is_worker = True
                            break

                    if is_worker:
                        try:
                            proc.kill()
                            killed_count += 1
                        except Exception as e:
                            print(f"[ERROR] Failed to kill process {proc.info['pid']}: {e}")
            except (psutil.NoSuchProcess, psutil.AccessDenied, psutil.ZombieProcess):
                pass

        if killed_count > 0:
            print(f"[INFO] Killed {killed_count} python processes during cleanup")

        return killed_count


    def ingest(self):
        general_start_time = time.time()
        BATCH_SIZE = int(os.getenv("BATCH_SIZE", 1000))
        pool = None

        try:
            # Handle S3-solo mode
            if self.s3_solo:
                # Create a temporary directory for S3 downloads
                with tempfile.TemporaryDirectory() as temp_dir:
                    s3_bucket = os.getenv('S3_BUCKET')
                    if not s3_bucket:
                        print("[ERROR] S3_BUCKET environment variable is required for S3-solo mode")
                        return

                    # Use the provided S3 key directly (supports both sharded and private paths)
                    s3_key = self.s3_solo_key
                    print(f"[INFO] S3-solo mode: processing {s3_bucket}/{s3_key}")
                    print(f"[INFO] Extracted hash: {self.s3_solo_hash}")

                    # Fetch first_seen from catalog_samples
                    first_seen = None
                    try:
                        client = get_db_catalog_connection()
                        result = client.query(
                            "SELECT first_seen FROM catalog_samples WHERE sha256 = %(hash)s LIMIT 1",
                            parameters={"hash": self.s3_solo_hash}
                        )
                        if result.result_rows:
                            first_seen = result.result_rows[0][0]
                            print(f"[INFO] first_seen from catalog: {first_seen}")
                        else:
                            print(f"[INFO] No catalog entry found, first_seen will default to epoch zero")
                        client.close()
                    except Exception as e:
                        print(f"[WARNING] Could not fetch first_seen from catalog: {e}")

                    # Process single S3 file directly
                    result = process_s3_file(
                        s3_bucket,
                        s3_key,
                        temp_dir,
                        self.decompile,
                        self.index_prefix,
                        self.log_file,
                        1,  # file_number
                        1,  # total_files
                        selected_modules=self.selected_modules,
                        dry_run=self.dry_run,
                        yara_scan=self.yara_scan,
                        with_yara=self.with_yara,
                        force=self.force,
                        decompile_modules=self.decompile_modules,
                        first_seen=first_seen,
                    )

                    if result:
                        print(f"[INFO] S3-solo processing completed successfully")
                    else:
                        print(f"[ERROR] S3-solo processing failed")

                    return

            # Handle S3 bulk mode
            elif self.s3_mode:
                # Create a temporary directory for S3 downloads
                with tempfile.TemporaryDirectory() as temp_dir:
                    # Query catalog for S3 objects based on mode
                    if self.start_date and self.end_date:
                        # Date-based query: join catalog_samples with repository_upload_sessions
                        print(f"[INFO] Querying samples by date range: {self.start_date} to {self.end_date}")
                        # Pass repository as None if it's a default placeholder (not a real repo name)
                        repo_filter = self.repository if self.repository not in ("date-range", "analyzed") else None
                        catalog_entries = fetch_s3_objects_by_date_range(
                            index_prefix=self.index_prefix,
                            decompile=self.decompile,
                            start_date=self.start_date,
                            end_date=self.end_date,
                            repository=repo_filter,
                            notes=self.s3_notes,
                            magika_filter=self.magika_filter,
                            yara_scan=self.yara_scan,
                            force=self.force,
                            analyzed=self.analyzed
                        )
                        date_info = f" from {self.start_date} to {self.end_date}"
                    elif self.analyzed:
                        # Analyzed mode (standalone, no date filter): query basic_properties for already-analyzed samples
                        print(f"[INFO] Querying already-analyzed samples from {self.index_prefix}_basic_properties")
                        catalog_entries = fetch_analyzed_samples(
                            index_prefix=self.index_prefix,
                            decompile=self.decompile,
                            magika_filter=self.magika_filter,
                            yara_scan=self.yara_scan,
                            force=self.force,
                            rerun=self.rerun
                        )
                        date_info = ""
                    else:
                        # Repository-based query: direct query to repository_upload_sessions
                        catalog_entries = fetch_s3_objects_by_repository(
                            self.repository,
                            self.index_prefix,
                            self.decompile,
                            self.s3_notes,
                            self.magika_filter,
                            self.yara_scan,
                            self.force
                        )
                        date_info = ""

                    if not catalog_entries:
                        print(f"[INFO] No files found for repository: {self.repository}{date_info}" +
                              (f" with notes: {self.s3_notes}" if self.s3_notes else ""))
                        return

                    # Extract S3 bucket, keys, and first_seen
                    s3_files = []
                    for entry in catalog_entries:
                        s3_bucket = entry.get('s3_bucket')
                        s3_key = entry.get('s3_key')
                        first_seen = entry.get('first_seen')
                        if s3_bucket and s3_key:
                            s3_files.append((s3_bucket, s3_key, first_seen))

                    total_files = len(s3_files)
                    num_cores = multiprocessing.cpu_count()
                    parallel_proc = num_cores - 1

                    if self.decompile:
                        if self.magika_filter == 'apk':
                            parallel_proc = num_cores - 1
                            print(f"[INFO] Using decompile mode (APK) with {parallel_proc} parallel processes")
                        else:
                            parallel_proc = max(1, int((num_cores - 1) / 2))
                            print(f"[INFO] Using decompile mode with {parallel_proc} parallel processes")
                    else:
                        print(f"[INFO] Using analysis mode with {parallel_proc} parallel processes")

                    # Use streaming approach for both decompile and non-decompile cases
                    self._process_files_streaming(s3_files, temp_dir, total_files, parallel_proc, self.decompile)
            else:
                # Original file/directory processing logic
                if os.path.isfile(self.path):
                    # Check if it's a text file containing paths
                    if self.path.endswith('.txt'):
                        try:
                            with open(self.path, 'r') as f:
                                files = [line.strip() for line in f
                                    if line.strip() and not os.path.basename(line.strip()).startswith('.')]
                        except Exception as e:
                            print(f"[ERR] Failed to read file list from {self.path}: {str(e)}")
                            return
                    else:
                        files = [self.path]
                elif os.path.isdir(self.path):
                    files = [
                        os.path.join(root, file)
                        for root, dirs, files_list in os.walk(self.path)
                        for file in files_list
                        if not file.startswith('.')
                    ]
                else:
                    print(f"[ERR] Invalid path: {self.path}")
                    return

                total_files = len(files)
                num_cores = multiprocessing.cpu_count()
                parallel_proc = num_cores - 1
                if self.decompile:
                    if self.magika_filter == 'apk':
                        parallel_proc = num_cores - 1
                    else:
                        parallel_proc = int(num_cores/2)
                print(f"Number of CPU cores: {num_cores}")
                print(f"Number of parallel processes: {parallel_proc}")

                # Process files in batches
                for i in range(0, len(files), BATCH_SIZE):
                    batch_files = files[i:i + BATCH_SIZE]
                    batch_start = i
                    print(f"\nProcessing batch {i//BATCH_SIZE + 1}/{(len(files) + BATCH_SIZE - 1)//BATCH_SIZE}")
                    pool = None
                    try:
                        pool = Pool(processes=parallel_proc)
                        batch_results = pool.map(
                            worker,
                            [
                                (
                                    f,
                                    self.decompile,
                                    self.index_prefix,
                                    self.log_file,
                                    batch_start + idx + 1,
                                    total_files,
                                    self.selected_modules,
                                    self.dry_run,
                                    self.yara_scan,
                                    self.with_yara,
                                    self.force,
                                    self.decompile_modules,
                                )
                                for idx, f in enumerate(batch_files)
                            ],
                        )

                        # Update statistics for this batch
                        for result, filetype in batch_results:
                            if result is not None:
                                self.total_results[result] += 1
                            if filetype:
                                self.file_type_stats[filetype] = (
                                    self.file_type_stats.get(filetype, 0) + 1
                                )

                    finally:
                        # Properly close the pool after each batch
                        if pool:
                            try:
                                pool.close()
                                pool.join()
                                pool = None
                                gc.collect()
                            except Exception as e:
                                print(f"[ERROR] Error cleaning up pool: {e}")
                                try:
                                    pool.terminate()
                                    pool.join()
                                except:
                                    pass
                                pool = None
                                gc.collect()

                        if check_high_swap():
                            print("[INFO] High swap detected, restarting pool")
                            self.restart_worker_pool()

                            gc.collect(2)

                            # Create a fresh pool
                            self.pool = Pool(processes=parallel_proc)

                        # Export strings after each batch
            general_end_time = time.time()
            general_elapsed_time = general_end_time - general_start_time
            general_elapsed_time_pretty = time.strftime("%H:%M:%S", time.gmtime(general_elapsed_time))

            summary = (
                f"\n\nIngestion finished for {self.path}."
                f"\nTime required: {general_elapsed_time_pretty}"
                f"\nResults:"
                f"\n- Total analyzed: {sum(self.total_results.values())}"
                f"\n- Correctly imported: {self.total_results[ImportResult.CORRECTLY]}"
                f"\n- Partially imported: {self.total_results[ImportResult.PARTIALLY]}"
                f"\n- Failed: {self.total_results[ImportResult.FAILED]}"
                f"\n- Skipped: {self.total_results[ImportResult.SKIPPED]}"
                f"\n\nFiletype stats:\n{json.dumps(dict(self.file_type_stats))}\n"
            )

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

            print("Ingestion completed. Check the log file for details.")
            print(summary)

        finally:
            try:
                self._kill_all_python_processes()
            except Exception as e:
                print(f"[ERROR] Error in final cleanup: {e}")

            try:
                gc.collect(2)
            except Exception as e:
                print(f"[ERROR] Final garbage collection error: {e}")