feat: group private apps into andreknie-privat
This commit is contained in:
@@ -0,0 +1,124 @@
|
||||
"""Drop-folder watcher using watchdog for continuous file ingestion.
|
||||
|
||||
Monitors the inbox folder for new files and triggers the processing pipeline.
|
||||
Moves processed files to archive/ and failed files to failed/ subdirectory.
|
||||
Ignores temporary files (names starting with '.' or ending with '.tmp').
|
||||
"""
|
||||
|
||||
import logging
|
||||
import shutil
|
||||
import time
|
||||
from pathlib import Path
|
||||
|
||||
from watchdog.events import FileCreatedEvent, FileSystemEventHandler
|
||||
from watchdog.observers import Observer
|
||||
|
||||
from ingestion.config import IngestionConfig
|
||||
from ingestion.discovery import SUPPORTED_EXTENSIONS
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
def should_ignore(filename: str) -> bool:
|
||||
"""Check if a file should be ignored by the watcher.
|
||||
|
||||
Ignores files starting with '.' or ending with '.tmp'.
|
||||
"""
|
||||
return filename.startswith(".") or filename.endswith(".tmp")
|
||||
|
||||
|
||||
class InboxHandler(FileSystemEventHandler):
|
||||
"""Handle new files in the inbox directory."""
|
||||
|
||||
def __init__(self, config: IngestionConfig):
|
||||
self.config = config
|
||||
self.inbox_dir = Path(config.inbox_dir)
|
||||
self.archive_dir = self.inbox_dir / "archive"
|
||||
self.failed_dir = self.inbox_dir / "failed"
|
||||
self.archive_dir.mkdir(parents=True, exist_ok=True)
|
||||
self.failed_dir.mkdir(parents=True, exist_ok=True)
|
||||
|
||||
def on_created(self, event: FileCreatedEvent):
|
||||
"""Triggered when a new file is created in the inbox."""
|
||||
if event.is_directory:
|
||||
return
|
||||
|
||||
file_path = Path(event.src_path)
|
||||
self._process_file(file_path)
|
||||
|
||||
def _process_file(self, file_path: Path) -> None:
|
||||
"""Process a single file and move to archive or failed."""
|
||||
if should_ignore(file_path.name):
|
||||
logger.debug("Ignoring temporary file: %s", file_path.name)
|
||||
return
|
||||
|
||||
if file_path.suffix.lower() not in SUPPORTED_EXTENSIONS:
|
||||
logger.debug("Ignoring unsupported file: %s", file_path.name)
|
||||
return
|
||||
|
||||
# Small delay to ensure file is fully written
|
||||
time.sleep(0.5)
|
||||
|
||||
logger.info("Processing new file: %s", file_path.name)
|
||||
|
||||
try:
|
||||
from ingestion.pipeline import process_file
|
||||
|
||||
result = process_file(file_path=file_path, config=self.config)
|
||||
|
||||
if result.success:
|
||||
# Move to archive
|
||||
dest = self.archive_dir / file_path.name
|
||||
shutil.move(str(file_path), str(dest))
|
||||
logger.info("Archived: %s → %s", file_path.name, dest)
|
||||
else:
|
||||
# Move to failed
|
||||
dest = self.failed_dir / file_path.name
|
||||
shutil.move(str(file_path), str(dest))
|
||||
logger.warning(
|
||||
"Failed: %s → %s (error: %s)",
|
||||
file_path.name,
|
||||
dest,
|
||||
result.error,
|
||||
)
|
||||
|
||||
except Exception as e:
|
||||
logger.error("Unhandled error processing %s: %s", file_path.name, e)
|
||||
dest = self.failed_dir / file_path.name
|
||||
try:
|
||||
shutil.move(str(file_path), str(dest))
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
|
||||
def start_watcher(config: IngestionConfig) -> None:
|
||||
"""Start the drop-folder watcher.
|
||||
|
||||
Processes existing files first, then watches for new ones.
|
||||
Blocks until interrupted (Ctrl+C).
|
||||
"""
|
||||
inbox_dir = Path(config.inbox_dir)
|
||||
inbox_dir.mkdir(parents=True, exist_ok=True)
|
||||
|
||||
handler = InboxHandler(config)
|
||||
|
||||
# Process existing files on startup
|
||||
existing = sorted(inbox_dir.iterdir())
|
||||
for f in existing:
|
||||
if f.is_file() and not should_ignore(f.name):
|
||||
handler._process_file(f)
|
||||
|
||||
# Start watching
|
||||
observer = Observer()
|
||||
observer.schedule(handler, str(inbox_dir), recursive=False)
|
||||
observer.start()
|
||||
|
||||
logger.info("Watcher started on: %s", inbox_dir)
|
||||
|
||||
try:
|
||||
while True:
|
||||
time.sleep(1)
|
||||
except KeyboardInterrupt:
|
||||
observer.stop()
|
||||
|
||||
observer.join()
|
||||
Reference in New Issue
Block a user