diff --git a/api/storage_s3.py b/api/storage_s3.py index f789998..acd198a 100644 --- a/api/storage_s3.py +++ b/api/storage_s3.py @@ -21,9 +21,12 @@ S3_REGION = os.getenv("S3_REGION", "us-east-1") router = APIRouter(prefix="/api/storage/s3", tags=["storage"]) -# Bounded scan so the dashboard stays responsive on big buckets. -SCAN_MAX_OBJECTS = int(os.getenv("S3_SCAN_MAX_OBJECTS", "50000")) -SCAN_DEADLINE_S = float(os.getenv("S3_SCAN_DEADLINE_S", "12")) +# Bounded scan so the dashboard stays responsive on big / fast-growing buckets. +SCAN_MAX_OBJECTS = int(os.getenv("S3_SCAN_MAX_OBJECTS", "80000")) +SCAN_DEADLINE_S = float(os.getenv("S3_SCAN_DEADLINE_S", "20")) +# Per top-level prefix cap so a single huge prefix (e.g. kafka/ CDC json) cannot +# starve the others — this keeps the type/size composition representative. +SCAN_PER_PREFIX = int(os.getenv("S3_SCAN_PER_PREFIX", "20000")) # In-process activity log — every S3 op routed through this API is recorded so # the dashboard can show live "what is happening on the storage" timelines. @@ -32,6 +35,25 @@ _analytics_cache: dict[str, Any] = {"ts": 0.0, "data": None} _ANALYTICS_TTL = 45.0 +def _client(): + return boto3.client( + "s3", + endpoint_url=S3_ENDPOINT, + aws_access_key_id=S3_ACCESS_KEY, + aws_secret_access_key=S3_SECRET_KEY, + region_name=S3_REGION, + config=Config(signature_version="s3v4", s3={"addressing_style": "path"}), + ) + + +def _human_size(n: float) -> str: + for unit in ("B", "KB", "MB", "GB", "TB"): + if n < 1024: + return f"{n:.0f} {unit}" if unit == "B" else f"{n:.1f} {unit}" + n /= 1024 + return f"{n:.1f} PB" + + def _track(op: str) -> None: try: _activity.append((time.time(), op)) @@ -63,25 +85,6 @@ def _size_class(n: int) -> str: return ">1 GB" -def _client(): - return boto3.client( - "s3", - endpoint_url=S3_ENDPOINT, - aws_access_key_id=S3_ACCESS_KEY, - aws_secret_access_key=S3_SECRET_KEY, - region_name=S3_REGION, - config=Config(signature_version="s3v4", s3={"addressing_style": "path"}), - ) - - -def _human_size(n: int) -> str: - for unit in ("B", "KB", "MB", "GB", "TB"): - if n < 1024: - return f"{n:.0f} {unit}" if unit == "B" else f"{n:.1f} {unit}" - n /= 1024 - return f"{n:.1f} PB" - - @router.get("/health") async def s3_health(): try: @@ -127,35 +130,87 @@ async def list_buckets(): async def list_objects( bucket: str, prefix: str = Query("", alias="prefix"), - max_keys: int = Query(200, le=500), + max_keys: int = Query(200, le=1000), + recursive: bool = Query(False), + q: str = Query(""), + ext: str = Query(""), + token: str = Query(""), ): + """Browse a bucket. + + - Folder mode (default): one level, returns sub-folders + objects. + - recursive / q / ext: flat search across the whole prefix subtree, + paginated via `token` (returned as `next_token`). + """ try: _track("list_objects") s3 = _client() - resp = s3.list_objects_v2(Bucket=bucket, Prefix=prefix, Delimiter="/", MaxKeys=max_keys) - folders = [ - {"type": "prefix", "name": p["Prefix"][len(prefix):].rstrip("/"), "prefix": p["Prefix"]} - for p in resp.get("CommonPrefixes", []) - ] - objects = [ - { - "type": "object", - "key": o["Key"], - "name": o["Key"][len(prefix):] if o["Key"].startswith(prefix) else o["Key"], - "size": o.get("Size", 0), - "size_human": _human_size(o.get("Size", 0)), - "modified": o.get("LastModified", "").isoformat() if o.get("LastModified") else None, - } - for o in resp.get("Contents", []) - if o["Key"] != prefix - ] + ql = q.lower().strip() + ext_l = ext.lower().lstrip(".").strip() + flat = bool(recursive or ql or ext_l) + folders: list[dict[str, Any]] = [] + objects: list[dict[str, Any]] = [] + cont = token or None + pages = 0 + scanned = 0 + next_token: str | None = None + PAGE_BUDGET = 40 + + while True: + kw: dict[str, Any] = {"Bucket": bucket, "Prefix": prefix, "MaxKeys": 1000} + if not flat: + kw["Delimiter"] = "/" + if cont: + kw["ContinuationToken"] = cont + r = s3.list_objects_v2(**kw) + pages += 1 + + for p in r.get("CommonPrefixes", []): + folders.append({ + "type": "prefix", + "name": p["Prefix"][len(prefix):].rstrip("/"), + "prefix": p["Prefix"], + }) + + for o in r.get("Contents", []): + key = o["Key"] + if key == prefix: + continue + scanned += 1 + if ql and ql not in key.lower(): + continue + if ext_l and _ext(key) != ext_l: + continue + objects.append({ + "type": "object", + "key": key, + "name": key[len(prefix):] if key.startswith(prefix) else key, + "ext": _ext(key), + "size": o.get("Size", 0), + "size_human": _human_size(o.get("Size", 0)), + "modified": o.get("LastModified", "").isoformat() if o.get("LastModified") else None, + }) + + cont = r.get("NextContinuationToken") + more = r.get("IsTruncated") + + if not flat: + next_token = cont if more else None + break + if len(objects) >= max_keys or not more or pages >= PAGE_BUDGET: + next_token = cont if more else None + break + return { "ok": True, "bucket": bucket, "prefix": prefix, + "flat": flat, "folders": folders, - "objects": objects, - "truncated": resp.get("IsTruncated", False), + "objects": objects[:max_keys] if flat else objects, + "scanned": scanned, + "next_token": next_token, + "truncated": bool(next_token), } except ClientError as exc: return JSONResponse({"ok": False, "error": str(exc)}, status_code=403) @@ -186,6 +241,69 @@ async def download_object(bucket: str, key: str = Query(...)): return JSONResponse({"ok": False, "error": str(exc)}, status_code=404) +# ── Analytics ──────────────────────────────────────────────────────────────── + + +def _accumulate(agg: dict[str, Any], bucket: str, o: dict[str, Any]) -> None: + sz = int(o.get("Size", 0) or 0) + key = o["Key"] + lm = o.get("LastModified") + agg["objects"] += 1 + agg["bytes"] += sz + bb = agg["by_bucket"].setdefault(bucket, [0, 0]) + bb[0] += 1 + bb[1] += sz + e = agg["by_ext"].setdefault(_ext(key), [0, 0]) + e[0] += 1 + e[1] += sz + sc = _size_class(sz) + agg["by_size"][sc] = agg["by_size"].get(sc, 0) + 1 + if lm: + d = lm.astimezone(timezone.utc).strftime("%Y-%m-%d") + dd = agg["by_day"].setdefault(d, [0, 0]) + dd[0] += 1 + dd[1] += sz + agg["recent"].append((lm.isoformat(), bucket, key, sz)) + if len(agg["recent"]) > 400: + agg["recent"] = sorted(agg["recent"], reverse=True)[:60] + rest = key[len(bucket) + 1:] if key.startswith(bucket + "/") else key + top = rest.split("/", 1)[0] if "/" in rest else "(root)" + pp = agg["by_prefix"].setdefault(f"{bucket}/{top}", [0, 0]) + pp[0] += 1 + pp[1] += sz + agg["largest"].append((sz, bucket, key, lm.isoformat() if lm else None)) + if len(agg["largest"]) > 400: + agg["largest"] = sorted(agg["largest"], reverse=True)[:60] + + +def _scan_prefix(s3, bucket: str, prefix: str, cap: int, deadline: float, agg: dict[str, Any]) -> bool: + """Scan one prefix subtree (flat). Returns True if capped/cut short.""" + tok = None + n = 0 + while True: + if time.time() > deadline or n >= cap or agg["objects"] >= SCAN_MAX_OBJECTS: + return True + kw: dict[str, Any] = {"Bucket": bucket, "MaxKeys": 1000} + if prefix: + kw["Prefix"] = prefix + if tok: + kw["ContinuationToken"] = tok + try: + r = s3.list_objects_v2(**kw) + except Exception: + return False + for o in r.get("Contents", []): + _accumulate(agg, bucket, o) + n += 1 + if n >= cap or agg["objects"] >= SCAN_MAX_OBJECTS: + return True + if r.get("IsTruncated"): + tok = r.get("NextContinuationToken") + if not tok: + return False + else: + return False + def _scan_analytics() -> dict[str, Any]: s3 = _client() @@ -198,101 +316,73 @@ def _scan_analytics() -> dict[str, Any]: } names = [b["Name"] for b in bucket_objs] - total_bytes = 0 - total_objects = 0 - by_bucket: dict[str, list[int]] = {} - by_ext: dict[str, list[int]] = {} - by_size: dict[str, int] = {} - by_day: dict[str, list[int]] = {} - by_prefix: dict[str, list[int]] = {} - largest: list[tuple[int, str, str, str | None]] = [] - recent: list[tuple[str, str, str, int]] = [] + agg: dict[str, Any] = { + "objects": 0, "bytes": 0, "by_bucket": {}, "by_ext": {}, "by_size": {}, + "by_day": {}, "by_prefix": {}, "largest": [], "recent": [], + } truncated = False - scanned = 0 for name in names: - bc = bb = 0 - token = None - while True: - if time.time() > deadline or scanned >= SCAN_MAX_OBJECTS: + if time.time() > deadline: + truncated = True + break + try: + r = s3.list_objects_v2(Bucket=name, Delimiter="/", MaxKeys=1000) + except Exception: + continue + # root-level objects + for o in r.get("Contents", []): + _accumulate(agg, name, o) + units = [p["Prefix"] for p in r.get("CommonPrefixes", [])] + agg["by_bucket"].setdefault(name, [0, 0]) + if not units: + if _scan_prefix(s3, name, "", SCAN_MAX_OBJECTS, deadline, agg): truncated = True - break - kw: dict[str, Any] = {"Bucket": name, "MaxKeys": 1000} - if token: - kw["ContinuationToken"] = token - try: - r = s3.list_objects_v2(**kw) - except Exception: - break - for o in r.get("Contents", []): - sz = int(o.get("Size", 0) or 0) - key = o["Key"] - lm = o.get("LastModified") - bc += 1 - bb += sz - scanned += 1 - total_objects += 1 - total_bytes += sz - e = by_ext.setdefault(_ext(key), [0, 0]) - e[0] += 1 - e[1] += sz - sc = _size_class(sz) - by_size[sc] = by_size.get(sc, 0) + 1 - if lm: - d = lm.astimezone(timezone.utc).strftime("%Y-%m-%d") - dd = by_day.setdefault(d, [0, 0]) - dd[0] += 1 - dd[1] += sz - recent.append((lm.isoformat(), name, key, sz)) - top = key.split("/", 1)[0] if "/" in key else "(root)" - pp = by_prefix.setdefault(f"{name}/{top}", [0, 0]) - pp[0] += 1 - pp[1] += sz - largest.append((sz, name, key, lm.isoformat() if lm else None)) - if scanned >= SCAN_MAX_OBJECTS: + else: + # breadth-first: every top-level prefix gets its own budget so a + # single huge prefix can't hide the rest of the data. + for pre in units: + if time.time() > deadline: truncated = True break - if r.get("IsTruncated") and not (time.time() > deadline or scanned >= SCAN_MAX_OBJECTS): - token = r.get("NextContinuationToken") - if not token: - break - else: - break - by_bucket[name] = [bc, bb] + if _scan_prefix(s3, name, pre, SCAN_PER_PREFIX, deadline, agg): + truncated = True + + total_bytes = agg["bytes"] + total_objects = agg["objects"] - # cumulative growth timeline growth = [] cum_b = cum_o = 0 - for d in sorted(by_day): - c, b = by_day[d] + for d in sorted(agg["by_day"]): + c, b = agg["by_day"][d] cum_o += c cum_b += b growth.append({"date": d, "objects": c, "bytes": b, "cum_objects": cum_o, "cum_bytes": cum_b}) buckets_out = sorted( - ([{"name": n, "objects": v[0], "bytes": v[1], "size_human": _human_size(v[1]), - "pct": round(100 * v[1] / total_bytes, 1) if total_bytes else 0, - "created": created_map.get(n)} for n, v in by_bucket.items()]), + [{"name": n, "objects": v[0], "bytes": v[1], "size_human": _human_size(v[1]), + "pct": round(100 * v[1] / total_bytes, 1) if total_bytes else 0, + "created": created_map.get(n)} for n, v in agg["by_bucket"].items()], key=lambda x: x["bytes"], reverse=True) types_out = sorted( - ([{"ext": k, "objects": v[0], "bytes": v[1], "size_human": _human_size(v[1])} - for k, v in by_ext.items()]), key=lambda x: x["bytes"], reverse=True)[:12] + [{"ext": k, "objects": v[0], "bytes": v[1], "size_human": _human_size(v[1])} + for k, v in agg["by_ext"].items()], key=lambda x: x["objects"], reverse=True)[:12] size_order = ["<1 KB", "1 KB–1 MB", "1–10 MB", "10–100 MB", "100 MB–1 GB", ">1 GB"] - size_out = [{"label": k, "count": by_size.get(k, 0)} for k in size_order] + size_out = [{"label": k, "count": agg["by_size"].get(k, 0)} for k in size_order] prefixes_out = sorted( - ([{"prefix": k, "objects": v[0], "bytes": v[1], "size_human": _human_size(v[1])} - for k, v in by_prefix.items()]), key=lambda x: x["bytes"], reverse=True)[:10] + [{"prefix": k, "objects": v[0], "bytes": v[1], "size_human": _human_size(v[1])} + for k, v in agg["by_prefix"].items()], key=lambda x: x["bytes"], reverse=True)[:12] - largest.sort(key=lambda x: x[0], reverse=True) + largest = sorted(agg["largest"], reverse=True)[:10] largest_out = [{"bucket": b, "key": k, "bytes": s, "size_human": _human_size(s), "modified": m} - for s, b, k, m in largest[:10]] + for s, b, k, m in largest] - recent.sort(key=lambda x: x[0], reverse=True) + recent = sorted(agg["recent"], reverse=True)[:15] recent_out = [{"modified": m, "bucket": b, "key": k, "bytes": s, "size_human": _human_size(s)} - for m, b, k, s in recent[:15]] + for m, b, k, s in recent] return { "summary": { @@ -306,7 +396,7 @@ def _scan_analytics() -> dict[str, Any]: "newest": recent_out[0]["modified"] if recent_out else None, "oldest": growth[0]["date"] if growth else None, "truncated": truncated, - "scanned": scanned, + "scanned": total_objects, }, "buckets": buckets_out, "types": types_out, diff --git a/ui/src/components/features/StorageView.tsx b/ui/src/components/features/StorageView.tsx index a8dd9a0..db097f6 100644 --- a/ui/src/components/features/StorageView.tsx +++ b/ui/src/components/features/StorageView.tsx @@ -1,13 +1,13 @@ import { useCallback, useEffect, useMemo, useState } from 'react' import { Activity, BarChart3, Boxes, ChevronRight, Clock, Database, Download, ExternalLink, - Files, Folder, Gauge, HardDrive, LayoutGrid, Loader2, RefreshCw, TrendingUp, + Files, Folder, Gauge, HardDrive, LayoutGrid, Loader2, RefreshCw, Search, TrendingUp, } from 'lucide-react' import { cn } from '../../lib/utils' import { subTabActive, subTabIdle } from '../../lib/tabActive' type Bucket = { name: string; created?: string; has_objects?: boolean } -type S3Item = { type: string; name?: string; prefix?: string; key?: string; size_human?: string; modified?: string } +type S3Item = { type: string; name?: string; prefix?: string; key?: string; ext?: string; size_human?: string; modified?: string } type GrowthPt = { date: string; objects: number; bytes: number; cum_objects: number; cum_bytes: number } type Analytics = { @@ -349,6 +349,8 @@ export function StorageView() { /* ── Browser (the original explorer, now a tab) ──────────────────────────── */ +const TYPE_CHIPS = ['json', 'jpeg', 'png', 'csv', 'log', 'parquet', 'avro', 'pdf', 'txt'] + function BucketBrowser() { const [buckets, setBuckets] = useState([]) const [bucket, setBucket] = useState(null) @@ -356,7 +358,17 @@ function BucketBrowser() { const [folders, setFolders] = useState([]) const [objects, setObjects] = useState([]) const [loading, setLoading] = useState(false) + const [loadingMore, setLoadingMore] = useState(false) const [error, setError] = useState(null) + // search / filter + const [qInput, setQInput] = useState('') + const [q, setQ] = useState('') + const [ext, setExt] = useState('') + const [recursive, setRecursive] = useState(false) + const [nextToken, setNextToken] = useState(null) + const [scanned, setScanned] = useState(0) + + const flat = recursive || !!q || !!ext const loadBuckets = useCallback(async () => { setLoading(true); setError(null) @@ -370,18 +382,31 @@ function BucketBrowser() { } catch { setError('S3 API unavailable') } finally { setLoading(false) } }, [bucket]) - const loadObjects = useCallback(async (b: string, p: string) => { - setLoading(true); setError(null) + const loadObjects = useCallback(async (b: string, p: string, opts: { q: string; ext: string; recursive: boolean }, token?: string) => { + const append = !!token + append ? setLoadingMore(true) : setLoading(true) + setError(null) try { - const r = await fetch(`/api/storage/s3/buckets/${encodeURIComponent(b)}/objects?prefix=${encodeURIComponent(p)}`) + const params = new URLSearchParams({ prefix: p, max_keys: '300' }) + if (opts.recursive) params.set('recursive', 'true') + if (opts.q) params.set('q', opts.q) + if (opts.ext) params.set('ext', opts.ext) + if (token) params.set('token', token) + const r = await fetch(`/api/storage/s3/buckets/${encodeURIComponent(b)}/objects?${params}`) const j = await r.json() if (!r.ok || !j.ok) { setError(j.error || 'List failed'); return } - setFolders(j.folders || []); setObjects(j.objects || []) - } catch { setError('Failed to list objects') } finally { setLoading(false) } + setFolders(j.flat ? [] : (j.folders || [])) + setObjects((prev) => (append ? [...prev, ...(j.objects || [])] : (j.objects || []))) + setNextToken(j.next_token || null) + setScanned((prev) => (append ? prev + (j.scanned || 0) : (j.scanned || 0))) + } catch { setError('Failed to list objects') } finally { append ? setLoadingMore(false) : setLoading(false) } }, []) useEffect(() => { loadBuckets() }, [loadBuckets]) - useEffect(() => { if (bucket) loadObjects(bucket, prefix) }, [bucket, prefix, loadObjects]) + useEffect(() => { if (bucket) loadObjects(bucket, prefix, { q, ext, recursive }) }, [bucket, prefix, q, ext, recursive, loadObjects]) + + const runSearch = () => setQ(qInput.trim()) + const clearSearch = () => { setQInput(''); setQ(''); setExt(''); setRecursive(false) } const crumbs = prefix ? prefix.split('/').filter(Boolean) : [] @@ -391,7 +416,7 @@ function BucketBrowser() {

Buckets

{buckets.map((b) => ( - + + + {TYPE_CHIPS.map((t) => ( + + ))} + {(q || ext || recursive) && ( + + )} +
+ + {bucket && !flat && ( )} + {flat && ( +

+ {loading ? 'Searching…' : <>Found {objects.length} object{objects.length === 1 ? '' : 's'} + {q && <> matching “{q}”}{ext && <> of type .{ext}} + {prefix && <> under {prefix}} · scanned {fmtNum(scanned)}{nextToken ? '+' : ''} keys} +

+ )} {loading &&

Loading…

} {error &&

{error}

}
- @@ -431,12 +492,14 @@ function BucketBrowser() { {f.name}/ + ))} {objects.map((o) => ( - + +
NameSizeModified + {flat ? 'Key' : 'Name'}TypeSizeModified
folder
{o.name || o.key}{flat ? o.key : (o.name || o.key)}{o.ext || '—'} {o.size_human} {o.modified?.slice(0, 19) || '—'} @@ -450,7 +513,17 @@ function BucketBrowser() { ))}
- {!loading && folders.length === 0 && objects.length === 0 && bucket &&

This prefix is empty.

} + {!loading && folders.length === 0 && objects.length === 0 && bucket && ( +

{flat ? 'No matching objects found.' : 'This prefix is empty.'}

+ )} + {nextToken && bucket && ( +
+ +
+ )}