Dockerfile / services /file_manager.py
kkthakur's picture
Deploy Local Hybrid Engine
b336134
Raw
History Blame Contribute Delete
4.45 kB
"""
File manager — upload, CSV→Parquet conversion, download, cleanup.
"""
from __future__ import annotations
import os
import uuid
import shutil
import polars as pl
from config import DATA_DIR, UPLOAD_DIR, LAZY_THRESHOLD_BYTES, MAX_UPLOAD_BYTES
from services.session_manager import session_manager
from core.column_registry import column_registry
def _human_size(n: int) -> str:
for unit in ("B", "KB", "MB", "GB"):
if n < 1024:
return f"{n:.1f} {unit}"
n /= 1024
return f"{n:.1f} TB"
def _detect_separator(path: str) -> str:
"""Sniff the separator from the first few KB."""
with open(path, "r", errors="replace") as f:
head = f.read(8192)
for sep in (",", ";", "\t", "|"):
if sep in head:
return sep
return ","
def handle_upload(file_content: bytes, file_name: str) -> dict:
"""Save uploaded file, convert to Parquet, create session."""
# Size check
if len(file_content) > MAX_UPLOAD_BYTES:
raise ValueError(f"File too large. Max {MAX_UPLOAD_BYTES // (1024**2)} MB allowed.")
# Generate session
session_id = uuid.uuid4().hex[:12]
# Save raw upload
ext = os.path.splitext(file_name)[1].lower()
raw_path = os.path.join(UPLOAD_DIR, f"{session_id}{ext}")
with open(raw_path, "wb") as f:
f.write(file_content)
# Determine file type and read
if ext in (".csv", ".tsv", ".txt"):
sep = _detect_separator(raw_path)
lf = pl.scan_csv(raw_path, separator=sep, try_parse_dates=True)
elif ext in (".xlsx", ".xls"):
# Polars can read Excel but needs the feature flag.
# Fall back to eager read for Excel.
df = pl.read_excel(raw_path) # type: ignore[attr-defined]
parquet_path = os.path.join(DATA_DIR, f"{session_id}.parquet")
df.write_parquet(parquet_path)
os.remove(raw_path)
elif ext == ".parquet":
shutil.copy2(raw_path, os.path.join(DATA_DIR, f"{session_id}.parquet"))
os.remove(raw_path)
lf = pl.scan_parquet(os.path.join(DATA_DIR, f"{session_id}.parquet"))
else:
os.remove(raw_path)
raise ValueError(f"Unsupported file format: {ext}")
# For CSV: stream-convert to Parquet
if ext in (".csv", ".tsv", ".txt"):
parquet_path = os.path.join(DATA_DIR, f"{session_id}.parquet")
lf.sink_parquet(parquet_path)
os.remove(raw_path)
lf = pl.scan_parquet(parquet_path)
# Read metadata (no full data load)
schema = lf.collect_schema()
row_count = lf.select(pl.len()).collect().item()
columns = [{"name": name, "dtype": str(dtype)} for name, dtype in schema.items()]
# Register in session manager
session_manager.create(
session_id=session_id,
file_name=file_name,
file_size_bytes=len(file_content),
columns=columns,
row_count=row_count,
)
# Register columns for fuzzy resolution
column_registry.register(session_id, schema.names())
return {
"session_id": session_id,
"file_name": file_name,
"rows": row_count,
"columns": schema.names(),
"size_human": _human_size(len(file_content)),
}
def get_download_path(session_id: str) -> Optional[str]:
"""Return path to the Parquet file, or None if session doesn't exist."""
path = session_manager.get_filepath(session_id)
if os.path.exists(path) and session_manager.get(session_id):
return path
return None
def export_to_csv(session_id: str) -> Optional[str]:
"""Export Parquet to a temporary CSV file and return its path."""
pq_path = session_manager.get_filepath(session_id)
if not os.path.exists(pq_path):
return None
csv_path = os.path.join(UPLOAD_DIR, f"{session_id}.csv")
lf = pl.scan_parquet(pq_path)
if os.path.getsize(pq_path) > LAZY_THRESHOLD_BYTES:
lf.sink_csv(csv_path)
else:
lf.collect().write_csv(csv_path)
return csv_path
def delete_session(session_id: str) -> bool:
"""Remove session metadata, Parquet file, and column registry."""
import glob
pattern = os.path.join(DATA_DIR, f"{session_id}*.parquet")
removed = False
for pq_path in glob.glob(pattern):
try:
os.remove(pq_path)
removed = True
except Exception:
pass
session_manager.remove(session_id)
column_registry.remove(session_id)
return removed