From 85f09303250fc1cdbac76a06933ac96c2b01740f Mon Sep 17 00:00:00 2001 From: doomedraven Date: Wed, 9 Sep 2026 12:20:11 +0200 Subject: [PATCH 1/5] Update Office file detection and error logging (#3200) * Update Office file detection and error logging Refactor logic for identifying Office files and enhance error logging for password-protected files. CAB files were going to ELSE block * Fix logic for checking password-protected Office files --- lib/cuckoo/common/demux.py | 16 ++++++++++++---- 1 file changed, 12 insertions(+), 4 deletions(-) diff --git a/lib/cuckoo/common/demux.py b/lib/cuckoo/common/demux.py index ff5109d2e75..84a57fbeafb 100644 --- a/lib/cuckoo/common/demux.py +++ b/lib/cuckoo/common/demux.py @@ -123,6 +123,11 @@ "Microsoft OOXML", ] +MS_EXCLUDE = [ + "MSI Installer", + "Microsoft Cabinet archive" +] + IGNORABLE_PATTERNS = ( re.compile(br"msvcp\d+\.dll$", re.IGNORECASE), @@ -487,15 +492,18 @@ def demux_sample( magic = File(filename).get_type() or "" # --- 3. Handle Password-Protected Office Files --- - is_office = ("Microsoft" in magic or any(x in magic for x in OFFICE_TYPES)) and "MSI Installer" not in magic + is_office = ("Microsoft" in magic or any(x in magic for x in OFFICE_TYPES)) and not any(x in magic for x in MS_EXCLUDE) if is_office and use_sflock: password = options2passwd(options) - if HAS_SFLOCK and password: + if use_sflock and password: retlist = demux_office(filename, password, platform) return retlist, error_list + # elif use_sflock: + # retlist = demux_office(filename, "", platform) + # return retlist, error_list else: - log.error("Detected password protected office file, but no sflock is installed.") - return [], [{os.path.basename(filename).decode(errors='ignore'): "Detected password protected office file, but no sflock is installed"}] + log.error("Detected password protected office file, but no sflock is installed. Magic: %s, Password:%s", magic, str(password)) + return [], [{os.path.basename(filename).decode(errors='ignore'): f"Detected password protected office file, but no sflock is installed. Magic: {magic}. Password: {password}"}] # --- 4. Skip Extraction for specific types --- ignored_signatures = [ From ede0407b61b11d8540fe2f7bbe6180fb6f94eed3 Mon Sep 17 00:00:00 2001 From: doomedraven Date: Wed, 9 Sep 2026 20:56:26 +0200 Subject: [PATCH 2/5] Enable symlink support in copytree function (#3225) * Enable symlink support in copytree function * Add worker cleaning option to dist.py --- utils/dist.py | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/utils/dist.py b/utils/dist.py index 8481c981b5c..5d43f9266e8 100644 --- a/utils/dist.py +++ b/utils/dist.py @@ -363,7 +363,7 @@ def node_get_report_nfs(task_id, worker_name, main_task_id) -> bool: path_mkdir(analyses_path, mode=0o755, exist_ok=False) try: - shutil.copytree(worker_path, analyses_path, ignore=dist_ignore_patterns, ignore_dangling_symlinks=True, dirs_exist_ok=True) + shutil.copytree(worker_path, analyses_path, symlinks=True, ignore=dist_ignore_patterns, ignore_dangling_symlinks=True, dirs_exist_ok=True) except shutil.Error: log.error("Files doens't exist on worker") except Exception as e: @@ -2268,6 +2268,10 @@ def main(): if args.enable_clean: cron_cleaner(args.clean_hours) + if args.clean_workers: + cron_cleaner(args.clean_hours) + sys.exit() + if args.force_reported: with main_db.session.begin(): main_db.set_status(args.force_reported, TASK_DISTRIBUTED_COMPLETED) From bd70d465cdeeeb203e2978e14c08eb9622f268a5 Mon Sep 17 00:00:00 2001 From: doomedraven Date: Wed, 9 Sep 2026 18:57:49 +0000 Subject: [PATCH 3/5] Implement diagnostic loop for GCP Pub/Sub status logging Added a diagnostic loop to log status and queue depth for GCP Pub/Sub subscriber. --- utils/gcp_pubsub_service.py | 76 +++++++++++++++++++++++++++++++++++++ 1 file changed, 76 insertions(+) diff --git a/utils/gcp_pubsub_service.py b/utils/gcp_pubsub_service.py index 729f624ece6..583ff0dd846 100644 --- a/utils/gcp_pubsub_service.py +++ b/utils/gcp_pubsub_service.py @@ -409,9 +409,85 @@ def process_message(self, message: Any): self.processing_ids.discard(msg_id) log.info("[%s] Total processing time: %.2f seconds", correlation_id, time.time() - start_time) + def _diagnostic_loop(self): + """Periodically log status and queue depth (if monitoring is available).""" + import time + from lib.cuckoo.common.gcp import gcp_cfg + + monitoring_client = None + try: + from google.cloud import monitoring_v3 + auth_by = gcp_cfg.gcp.get("auth_by", "vm") + service_account_path = gcp_cfg.gcp.get("service_account_path") + + if auth_by == "json" and service_account_path: + if not os.path.isabs(service_account_path): + from lib.cuckoo.common.constants import CUCKOO_ROOT + service_account_path = os.path.join(CUCKOO_ROOT, service_account_path) + if os.path.exists(service_account_path): + monitoring_client = monitoring_v3.MetricServiceClient.from_service_account_json(service_account_path) + else: + monitoring_client = monitoring_v3.MetricServiceClient() + except ImportError: + log.debug("google-cloud-monitoring not installed. Install via `pip install google-cloud-monitoring` for precise queue counts.") + except Exception as e: + log.debug("Failed to initialize monitoring client: %s", e) + + # Wait briefly before first check + time.sleep(5) + + while True: + queue_size_str = "unknown (install google-cloud-monitoring)" + if monitoring_client: + try: + from google.cloud.monitoring_v3 import types + project_name = f"projects/{self.project_id}" + now = time.time() + interval = types.TimeInterval( + { + "end_time": {"seconds": int(now)}, + "start_time": {"seconds": int(now - 600)}, + } + ) + + results = monitoring_client.list_time_series( + request={ + "name": project_name, + "filter": f'metric.type = "pubsub.googleapis.com/subscription/num_undelivered_messages" AND resource.labels.subscription_id = "{self.subscription_id}"', + "interval": interval, + } + ) + + latest_val = None + for result in results: + for point in result.points: + latest_val = point.value.int64_value + break + if latest_val is not None: + break + + if latest_val is not None: + queue_size_str = str(latest_val) + else: + queue_size_str = "0" + except Exception as e: + log.debug("Error fetching queue size metric: %s", e) + queue_size_str = "error (permission or API issue)" + + with self.ids_lock: + active = len(self.processing_ids) + + log.info("[HEARTBEAT] Subscriber is healthy. Actively processing: %d Tasks. Undelivered queue size: %s.", active, queue_size_str) + time.sleep(300) + def start(self): log.info("Starting GCP Pub/Sub subscriber on %s", self.subscription_path) + # Start a background diagnostic thread so the app doesn't seem 'hung' when idle + import threading + diag_thread = threading.Thread(target=self._diagnostic_loop, daemon=True) + diag_thread.start() + from lib.cuckoo.common.gcp import gcp_cfg max_messages = 5 lease_duration = 1800 From 8d891a2f19e509673c4e902da101e66040067d11 Mon Sep 17 00:00:00 2001 From: doomedraven Date: Wed, 9 Sep 2026 19:00:39 +0000 Subject: [PATCH 4/5] Clean up blank lines in _diagnostic_loop method Removed unnecessary blank lines in the diagnostic loop method for cleaner code. --- utils/gcp_pubsub_service.py | 14 +++++++------- 1 file changed, 7 insertions(+), 7 deletions(-) diff --git a/utils/gcp_pubsub_service.py b/utils/gcp_pubsub_service.py index 583ff0dd846..bd1a53f147c 100644 --- a/utils/gcp_pubsub_service.py +++ b/utils/gcp_pubsub_service.py @@ -413,13 +413,13 @@ def _diagnostic_loop(self): """Periodically log status and queue depth (if monitoring is available).""" import time from lib.cuckoo.common.gcp import gcp_cfg - + monitoring_client = None try: from google.cloud import monitoring_v3 auth_by = gcp_cfg.gcp.get("auth_by", "vm") service_account_path = gcp_cfg.gcp.get("service_account_path") - + if auth_by == "json" and service_account_path: if not os.path.isabs(service_account_path): from lib.cuckoo.common.constants import CUCKOO_ROOT @@ -435,7 +435,7 @@ def _diagnostic_loop(self): # Wait briefly before first check time.sleep(5) - + while True: queue_size_str = "unknown (install google-cloud-monitoring)" if monitoring_client: @@ -449,7 +449,7 @@ def _diagnostic_loop(self): "start_time": {"seconds": int(now - 600)}, } ) - + results = monitoring_client.list_time_series( request={ "name": project_name, @@ -457,7 +457,7 @@ def _diagnostic_loop(self): "interval": interval, } ) - + latest_val = None for result in results: for point in result.points: @@ -465,7 +465,7 @@ def _diagnostic_loop(self): break if latest_val is not None: break - + if latest_val is not None: queue_size_str = str(latest_val) else: @@ -476,7 +476,7 @@ def _diagnostic_loop(self): with self.ids_lock: active = len(self.processing_ids) - + log.info("[HEARTBEAT] Subscriber is healthy. Actively processing: %d Tasks. Undelivered queue size: %s.", active, queue_size_str) time.sleep(300) From e29082d1b716b883928f0bf14e00d38edca7d157 Mon Sep 17 00:00:00 2001 From: doomedraven Date: Wed, 9 Sep 2026 21:13:31 +0200 Subject: [PATCH 5/5] Update views.py --- web/analysis/views.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/web/analysis/views.py b/web/analysis/views.py index 18ba39da562..9641c11fe0e 100644 --- a/web/analysis/views.py +++ b/web/analysis/views.py @@ -4510,7 +4510,7 @@ def on_demand(request, service: str, task_id: str, category: str, sha256): if not path_exists(path): extractedfile = False - if category == "static": + if category in ("static", "target.file"): path = os.path.join(ANALYSIS_BASE_PATH, "analyses", task_id, "binary") category = "target.file" elif category == "dropped":