forked from retoor/devplacepy
Compare commits
2
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
5fe9b4f1ed | ||
|
|
ae8c354942 |
@@ -1,6 +1,7 @@
|
|||||||
# retoor <retoor@molodetz.nl>
|
# retoor <retoor@molodetz.nl>
|
||||||
|
|
||||||
import logging
|
import logging
|
||||||
|
import os
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
from typing import Annotated
|
from typing import Annotated
|
||||||
|
|
||||||
@@ -64,6 +65,7 @@ def _targets() -> list[dict]:
|
|||||||
def _dashboard(can_download: bool) -> dict:
|
def _dashboard(can_download: bool) -> dict:
|
||||||
backups = [_backup_payload(row, can_download) for row in store.list_backups()]
|
backups = [_backup_payload(row, can_download) for row in store.list_backups()]
|
||||||
schedules = store.list_schedules()
|
schedules = store.list_schedules()
|
||||||
|
threshold = int(os.environ.get("DISK_WARNING_PERCENT", "85"))
|
||||||
return {
|
return {
|
||||||
"storage": store.compute_storage_stats(),
|
"storage": store.compute_storage_stats(),
|
||||||
"backups": backups,
|
"backups": backups,
|
||||||
@@ -72,6 +74,8 @@ def _dashboard(can_download: bool) -> dict:
|
|||||||
"metrics": _metrics(backups),
|
"metrics": _metrics(backups),
|
||||||
"generated_at": store.now_iso(),
|
"generated_at": store.now_iso(),
|
||||||
"can_download_backups": can_download,
|
"can_download_backups": can_download,
|
||||||
|
"disk_warnings": store.get_disk_warnings(threshold),
|
||||||
|
"disk_warning_threshold": threshold,
|
||||||
}
|
}
|
||||||
|
|
||||||
@router.get("/backups", response_class=HTMLResponse)
|
@router.get("/backups", response_class=HTMLResponse)
|
||||||
|
|||||||
@@ -53,6 +53,13 @@ class BackupJobOut(_Out):
|
|||||||
completed_at: Optional[str] = None
|
completed_at: Optional[str] = None
|
||||||
|
|
||||||
|
|
||||||
|
class DiskWarningOut(_Out):
|
||||||
|
path: str = ""
|
||||||
|
used_pct: float = 0.0
|
||||||
|
total_gb: float = 0.0
|
||||||
|
used_gb: float = 0.0
|
||||||
|
|
||||||
|
|
||||||
class BackupScheduleOut(_Out):
|
class BackupScheduleOut(_Out):
|
||||||
uid: str = ""
|
uid: str = ""
|
||||||
name: str = ""
|
name: str = ""
|
||||||
@@ -78,3 +85,5 @@ class BackupDashboardOut(_Out):
|
|||||||
generated_at: Optional[str] = None
|
generated_at: Optional[str] = None
|
||||||
can_download_backups: bool = False
|
can_download_backups: bool = False
|
||||||
admin_section: Optional[str] = None
|
admin_section: Optional[str] = None
|
||||||
|
disk_warnings: list[DiskWarningOut] = []
|
||||||
|
disk_warning_threshold: int = 85
|
||||||
|
|||||||
@@ -34,6 +34,10 @@ class BackupService(JobService):
|
|||||||
|
|
||||||
def __init__(self):
|
def __init__(self):
|
||||||
super().__init__(name="backup", interval_seconds=15)
|
super().__init__(name="backup", interval_seconds=15)
|
||||||
|
try:
|
||||||
|
store.seed_default_schedule()
|
||||||
|
except Exception as exc:
|
||||||
|
logger.warning("failed to seed default backup schedule: %s", exc)
|
||||||
|
|
||||||
async def run_once(self) -> None:
|
async def run_once(self) -> None:
|
||||||
await super().run_once()
|
await super().run_once()
|
||||||
@@ -121,6 +125,15 @@ class BackupService(JobService):
|
|||||||
summary=f"backup {target} completed ({_human_bytes(stats['bytes_out'])})",
|
summary=f"backup {target} completed ({_human_bytes(stats['bytes_out'])})",
|
||||||
links=[audit.job(job_uid)],
|
links=[audit.job(job_uid)],
|
||||||
)
|
)
|
||||||
|
if schedule_uid:
|
||||||
|
try:
|
||||||
|
store.check_disk_and_alert()
|
||||||
|
except Exception as exc:
|
||||||
|
logger.warning("disk alert check failed: %s", exc)
|
||||||
|
try:
|
||||||
|
store.check_db_growth()
|
||||||
|
except Exception as exc:
|
||||||
|
logger.warning("db growth check failed: %s", exc)
|
||||||
return {
|
return {
|
||||||
"backup_uid": backup_uid,
|
"backup_uid": backup_uid,
|
||||||
"target": target,
|
"target": target,
|
||||||
|
|||||||
@@ -1,5 +1,7 @@
|
|||||||
# retoor <retoor@molodetz.nl>
|
# retoor <retoor@molodetz.nl>
|
||||||
|
|
||||||
|
import json
|
||||||
|
import logging
|
||||||
import shutil
|
import shutil
|
||||||
import time
|
import time
|
||||||
from datetime import datetime, timezone
|
from datetime import datetime, timezone
|
||||||
@@ -9,6 +11,8 @@ from devplacepy import config
|
|||||||
from devplacepy.database import _index, db, get_table
|
from devplacepy.database import _index, db, get_table
|
||||||
from devplacepy.utils import generate_uid
|
from devplacepy.utils import generate_uid
|
||||||
|
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
BACKUP_TARGETS: dict[str, dict] = {
|
BACKUP_TARGETS: dict[str, dict] = {
|
||||||
"database": {
|
"database": {
|
||||||
"label": "Database",
|
"label": "Database",
|
||||||
@@ -375,6 +379,58 @@ def _storage_paths() -> list[tuple[str, str, Path]]:
|
|||||||
]
|
]
|
||||||
|
|
||||||
|
|
||||||
|
def seed_default_schedule() -> None:
|
||||||
|
if "backup_schedules" not in db.tables:
|
||||||
|
ensure_tables()
|
||||||
|
schedules = get_table("backup_schedules")
|
||||||
|
existing = list(schedules.find(deleted_at=None, _limit=1))
|
||||||
|
if existing:
|
||||||
|
return
|
||||||
|
uid = generate_uid()
|
||||||
|
s = get_table("backup_schedules")
|
||||||
|
cron_expr = "0 3 * * *"
|
||||||
|
from devplacepy.services.devii.tasks.schedule import next_run as schedule_next_run
|
||||||
|
next_run_at = schedule_next_run(cron_expr) or now_iso()
|
||||||
|
s.insert(
|
||||||
|
{
|
||||||
|
"uid": uid,
|
||||||
|
"name": "Daily database backup",
|
||||||
|
"target": "database",
|
||||||
|
"kind": "cron",
|
||||||
|
"every_seconds": 0,
|
||||||
|
"cron": cron_expr,
|
||||||
|
"enabled": 1,
|
||||||
|
"keep_last": 14,
|
||||||
|
"next_run_at": next_run_at,
|
||||||
|
"last_run_at": "",
|
||||||
|
"last_job_uid": "",
|
||||||
|
"run_count": 0,
|
||||||
|
"created_by": "system",
|
||||||
|
"created_at": now_iso(),
|
||||||
|
"updated_at": now_iso(),
|
||||||
|
"deleted_at": None,
|
||||||
|
"deleted_by": None,
|
||||||
|
}
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def get_disk_warnings(threshold: int = 85) -> list[dict]:
|
||||||
|
stats = compute_storage_stats()
|
||||||
|
disk = stats.get("disk", {})
|
||||||
|
used_pct = disk.get("used_percent", 0.0)
|
||||||
|
warnings = []
|
||||||
|
if used_pct >= threshold:
|
||||||
|
warnings.append(
|
||||||
|
{
|
||||||
|
"path": stats.get("data_dir", {}).get("path", str(config.DATA_DIR)),
|
||||||
|
"used_pct": used_pct,
|
||||||
|
"total_gb": round(disk.get("total_bytes", 0) / (1024**3), 1),
|
||||||
|
"used_gb": round(disk.get("used_bytes", 0) / (1024**3), 1),
|
||||||
|
}
|
||||||
|
)
|
||||||
|
return warnings
|
||||||
|
|
||||||
|
|
||||||
def compute_storage_stats() -> dict:
|
def compute_storage_stats() -> dict:
|
||||||
now = time.monotonic()
|
now = time.monotonic()
|
||||||
if (
|
if (
|
||||||
@@ -435,3 +491,131 @@ def compute_storage_stats() -> dict:
|
|||||||
_storage_cache["data"] = data
|
_storage_cache["data"] = data
|
||||||
_storage_cache["at"] = now
|
_storage_cache["at"] = now
|
||||||
return data
|
return data
|
||||||
|
|
||||||
|
|
||||||
|
DB_GROWTH_HISTORY_FILE = None
|
||||||
|
|
||||||
|
|
||||||
|
def _db_growth_path() -> Path:
|
||||||
|
global DB_GROWTH_HISTORY_FILE
|
||||||
|
if DB_GROWTH_HISTORY_FILE is None:
|
||||||
|
DB_GROWTH_HISTORY_FILE = config.BACKUPS_DIR / "db_growth_history.json"
|
||||||
|
return DB_GROWTH_HISTORY_FILE
|
||||||
|
|
||||||
|
|
||||||
|
def _load_db_sizes() -> list[dict]:
|
||||||
|
path = _db_growth_path()
|
||||||
|
if not path.exists():
|
||||||
|
return []
|
||||||
|
try:
|
||||||
|
return json.loads(path.read_text())
|
||||||
|
except (ValueError, OSError):
|
||||||
|
return []
|
||||||
|
|
||||||
|
|
||||||
|
def _save_db_sizes(sizes: list[dict]) -> None:
|
||||||
|
path = _db_growth_path()
|
||||||
|
try:
|
||||||
|
path.parent.mkdir(parents=True, exist_ok=True)
|
||||||
|
path.write_text(json.dumps(sizes))
|
||||||
|
except OSError as exc:
|
||||||
|
logger.warning("failed to write db growth history: %s", exc)
|
||||||
|
|
||||||
|
|
||||||
|
def _database_file_size() -> int:
|
||||||
|
database_file = Path(str(config.DATABASE_URL).replace("sqlite:///", ""))
|
||||||
|
try:
|
||||||
|
return database_file.stat().st_size if database_file.exists() else 0
|
||||||
|
except OSError:
|
||||||
|
return 0
|
||||||
|
|
||||||
|
|
||||||
|
def check_db_growth() -> None:
|
||||||
|
from devplacepy.services.audit import record as audit
|
||||||
|
|
||||||
|
current = _database_file_size()
|
||||||
|
if current <= 0:
|
||||||
|
return
|
||||||
|
|
||||||
|
sizes = _load_db_sizes()
|
||||||
|
last_entry = sizes[-1] if sizes else None
|
||||||
|
now = now_iso()
|
||||||
|
|
||||||
|
sizes.append({"at": now, "size_bytes": current})
|
||||||
|
_save_db_sizes(sizes)
|
||||||
|
|
||||||
|
if last_entry and last_entry.get("size_bytes", 0) > 0:
|
||||||
|
try:
|
||||||
|
last_at = datetime.fromisoformat(last_entry["at"])
|
||||||
|
current_at = datetime.fromisoformat(now)
|
||||||
|
days = (current_at - last_at).total_seconds() / 86400
|
||||||
|
if days > 0:
|
||||||
|
growth_pct = (current - last_entry["size_bytes"]) / last_entry["size_bytes"]
|
||||||
|
daily_pct = growth_pct / days
|
||||||
|
if daily_pct > 0.10:
|
||||||
|
logger.warning(
|
||||||
|
"database growth %.1f%%/day (%.0f MB -> %.0f MB over %.1f days)",
|
||||||
|
daily_pct * 100,
|
||||||
|
last_entry["size_bytes"] / (1024 * 1024),
|
||||||
|
current / (1024 * 1024),
|
||||||
|
days,
|
||||||
|
)
|
||||||
|
try:
|
||||||
|
audit.record_system(
|
||||||
|
"backup.db_growth_high",
|
||||||
|
actor_kind="service",
|
||||||
|
origin="scheduler",
|
||||||
|
target_type="system",
|
||||||
|
metadata={
|
||||||
|
"size_bytes": current,
|
||||||
|
"previous_size_bytes": last_entry["size_bytes"],
|
||||||
|
"daily_growth_pct": round(daily_pct * 100, 1),
|
||||||
|
"span_days": round(days, 2),
|
||||||
|
},
|
||||||
|
summary=(
|
||||||
|
f"database growth {daily_pct * 100:.1f}%/day "
|
||||||
|
f"({human_bytes(last_entry['size_bytes'])} -> {human_bytes(current)})"
|
||||||
|
),
|
||||||
|
)
|
||||||
|
except Exception as exc:
|
||||||
|
logger.warning("db growth audit failed: %s", exc)
|
||||||
|
except (ValueError, TypeError) as exc:
|
||||||
|
logger.warning("failed to parse db growth dates: %s", exc)
|
||||||
|
|
||||||
|
|
||||||
|
DISK_ALERT_THRESHOLD = 85
|
||||||
|
|
||||||
|
|
||||||
|
def check_disk_and_alert() -> None:
|
||||||
|
from devplacepy.services.audit import record as audit
|
||||||
|
|
||||||
|
stats = compute_storage_stats()
|
||||||
|
disk = stats.get("disk", {})
|
||||||
|
used_pct = disk.get("used_percent", 0.0)
|
||||||
|
if used_pct < DISK_ALERT_THRESHOLD:
|
||||||
|
return
|
||||||
|
logger.warning(
|
||||||
|
"disk usage at %.1f%% (threshold %d%%)",
|
||||||
|
used_pct,
|
||||||
|
DISK_ALERT_THRESHOLD,
|
||||||
|
)
|
||||||
|
try:
|
||||||
|
audit.record_system(
|
||||||
|
"backup.disk.critical",
|
||||||
|
actor_kind="service",
|
||||||
|
origin="scheduler",
|
||||||
|
target_type="system",
|
||||||
|
metadata={
|
||||||
|
"used_percent": used_pct,
|
||||||
|
"used_bytes": disk.get("used_bytes", 0),
|
||||||
|
"total_bytes": disk.get("total_bytes", 0),
|
||||||
|
"free_bytes": disk.get("free_bytes", 0),
|
||||||
|
},
|
||||||
|
summary=(
|
||||||
|
f"disk usage at {used_pct:.1f}% "
|
||||||
|
f"({human_bytes(disk.get('used_bytes', 0))} / "
|
||||||
|
f"{human_bytes(disk.get('total_bytes', 0))})"
|
||||||
|
),
|
||||||
|
)
|
||||||
|
except Exception as exc:
|
||||||
|
logger.warning("disk alert audit failed: %s", exc)
|
||||||
|
|||||||
@@ -81,6 +81,7 @@ class ContainerService(BaseService):
|
|||||||
def __init__(self):
|
def __init__(self):
|
||||||
super().__init__(name="containers", interval_seconds=5)
|
super().__init__(name="containers", interval_seconds=5)
|
||||||
self._metric_tick = 0
|
self._metric_tick = 0
|
||||||
|
self._orphan_sweep_tick = 0
|
||||||
self._booted = False
|
self._booted = False
|
||||||
self._last_sync_at = 0.0
|
self._last_sync_at = 0.0
|
||||||
|
|
||||||
@@ -107,19 +108,6 @@ class ContainerService(BaseService):
|
|||||||
|
|
||||||
await self._periodic_sync(instances, by_uid)
|
await self._periodic_sync(instances, by_uid)
|
||||||
|
|
||||||
for uid, row in by_uid.items():
|
|
||||||
if uid not in known:
|
|
||||||
try:
|
|
||||||
await backend.rm(row.container_id, force=True)
|
|
||||||
self.log(f"reaped orphan container {row.name}")
|
|
||||||
_audit_reconcile(
|
|
||||||
{"uid": uid, "name": row.name},
|
|
||||||
"orphan_reap",
|
|
||||||
f"reconciler reaped orphan container {row.name}",
|
|
||||||
)
|
|
||||||
except Exception as exc:
|
|
||||||
self.log(f"orphan rm failed for {row.name}: {exc}")
|
|
||||||
|
|
||||||
for inst in instances:
|
for inst in instances:
|
||||||
try:
|
try:
|
||||||
await self._reconcile(backend, inst, by_uid.get(inst["uid"]))
|
await self._reconcile(backend, inst, by_uid.get(inst["uid"]))
|
||||||
@@ -128,6 +116,16 @@ class ContainerService(BaseService):
|
|||||||
|
|
||||||
await self._fire_schedules()
|
await self._fire_schedules()
|
||||||
|
|
||||||
|
try:
|
||||||
|
await self._sweep_orphaned_containers(backend, by_uid, known)
|
||||||
|
except Exception as exc:
|
||||||
|
self.log(f"orphan sweep failed: {exc}")
|
||||||
|
|
||||||
|
try:
|
||||||
|
await self._cleanup_stale_workspaces(backend)
|
||||||
|
except Exception as exc:
|
||||||
|
self.log(f"workspace cleanup failed: {exc}")
|
||||||
|
|
||||||
self._metric_tick += 1
|
self._metric_tick += 1
|
||||||
if (
|
if (
|
||||||
self._metric_tick
|
self._metric_tick
|
||||||
@@ -136,6 +134,49 @@ class ContainerService(BaseService):
|
|||||||
):
|
):
|
||||||
await self._sample_metrics(backend, by_uid)
|
await self._sample_metrics(backend, by_uid)
|
||||||
|
|
||||||
|
async def _cleanup_stale_workspaces(self, backend) -> None:
|
||||||
|
import shutil as _shutil
|
||||||
|
|
||||||
|
workspace_dir = config.CONTAINER_WORKSPACES_DIR
|
||||||
|
if not workspace_dir.is_dir():
|
||||||
|
return
|
||||||
|
instances = store.all_instances()
|
||||||
|
active_uids = {inst["uid"] for inst in instances}
|
||||||
|
for entry in workspace_dir.iterdir():
|
||||||
|
if not entry.is_dir():
|
||||||
|
continue
|
||||||
|
uid = entry.name
|
||||||
|
if uid in active_uids:
|
||||||
|
continue
|
||||||
|
try:
|
||||||
|
_shutil.rmtree(entry)
|
||||||
|
self.log(f"removed stale workspace {uid}")
|
||||||
|
except Exception as exc:
|
||||||
|
self.log(f"failed to remove workspace {uid}: {exc}")
|
||||||
|
await backend.image_prune()
|
||||||
|
|
||||||
|
async def _sweep_orphaned_containers(self, backend, by_uid: dict, known: set) -> None:
|
||||||
|
if self._orphan_sweep_tick % 60 != 0:
|
||||||
|
self._orphan_sweep_tick += 1
|
||||||
|
return
|
||||||
|
self._orphan_sweep_tick += 1
|
||||||
|
for uid, row in by_uid.items():
|
||||||
|
if uid not in known:
|
||||||
|
try:
|
||||||
|
await backend.stop(row.container_id, timeout=10)
|
||||||
|
except Exception:
|
||||||
|
pass
|
||||||
|
try:
|
||||||
|
await backend.rm(row.container_id, force=True)
|
||||||
|
self.log(f"removed orphan container {row.name}")
|
||||||
|
_audit_reconcile(
|
||||||
|
{"uid": uid, "name": row.name},
|
||||||
|
"orphan_reap",
|
||||||
|
f"sweeper removed orphan container {row.name}",
|
||||||
|
)
|
||||||
|
except Exception as exc:
|
||||||
|
self.log(f"orphan rm failed for {row.name}: {exc}")
|
||||||
|
|
||||||
async def _reconcile(self, backend, inst, ps) -> None:
|
async def _reconcile(self, backend, inst, ps) -> None:
|
||||||
uid = inst["uid"]
|
uid = inst["uid"]
|
||||||
desired = inst["desired_state"]
|
desired = inst["desired_state"]
|
||||||
|
|||||||
@@ -22,8 +22,6 @@ from .pdf import MAX_PDF_BYTES, extract_pdf_text, is_pdf
|
|||||||
|
|
||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
_pw_lock = asyncio.Lock()
|
|
||||||
|
|
||||||
RSEARCH_URL = "https://rsearch.app.molodetz.nl"
|
RSEARCH_URL = "https://rsearch.app.molodetz.nl"
|
||||||
RSEARCH_TIMEOUT_SECONDS = 45.0
|
RSEARCH_TIMEOUT_SECONDS = 45.0
|
||||||
FETCH_TIMEOUT_SECONDS = 20.0
|
FETCH_TIMEOUT_SECONDS = 20.0
|
||||||
@@ -174,7 +172,13 @@ async def _render_with_playwright(
|
|||||||
) -> tuple[str, str, int, list[tuple[str, str]]]:
|
) -> tuple[str, str, int, list[tuple[str, str]]]:
|
||||||
from playwright.async_api import async_playwright
|
from playwright.async_api import async_playwright
|
||||||
|
|
||||||
async def _render(browser) -> tuple[str, str, int, list[tuple[str, str]]]:
|
own_browser = browser is None
|
||||||
|
if own_browser:
|
||||||
|
pw = await async_playwright().__aenter__()
|
||||||
|
browser = await pw.chromium.launch(
|
||||||
|
headless=True, args=["--no-sandbox", "--disable-dev-shm-usage"]
|
||||||
|
)
|
||||||
|
try:
|
||||||
context = await browser.new_context(user_agent=USER_AGENT)
|
context = await browser.new_context(user_agent=USER_AGENT)
|
||||||
page = await context.new_page()
|
page = await context.new_page()
|
||||||
response = await page.goto(url, wait_until="load", timeout=30000)
|
response = await page.goto(url, wait_until="load", timeout=30000)
|
||||||
@@ -185,18 +189,10 @@ async def _render_with_playwright(
|
|||||||
await context.close()
|
await context.close()
|
||||||
extracted = extract_html(content, base_url=url)
|
extracted = extract_html(content, base_url=url)
|
||||||
return extracted.title, extracted.text, status, extracted.links
|
return extracted.title, extracted.text, status, extracted.links
|
||||||
|
finally:
|
||||||
if browser is None:
|
if own_browser:
|
||||||
async with async_playwright() as pw:
|
await browser.close()
|
||||||
browser = await pw.chromium.launch(
|
await pw.__aexit__(None, None, None)
|
||||||
headless=True, args=["--no-sandbox", "--disable-dev-shm-usage"]
|
|
||||||
)
|
|
||||||
try:
|
|
||||||
return await _render(browser)
|
|
||||||
finally:
|
|
||||||
await browser.close()
|
|
||||||
else:
|
|
||||||
return await _render(browser)
|
|
||||||
|
|
||||||
|
|
||||||
async def fetch_page(url: str, depth: int, browser=None) -> CrawledPage | None:
|
async def fetch_page(url: str, depth: int, browser=None) -> CrawledPage | None:
|
||||||
@@ -246,20 +242,18 @@ async def fetch_page(url: str, depth: int, browser=None) -> CrawledPage | None:
|
|||||||
except (LookupError, ValueError) as exc:
|
except (LookupError, ValueError) as exc:
|
||||||
logger.info("deepsearch decode failed for %s: %s", url, exc)
|
logger.info("deepsearch decode failed for %s: %s", url, exc)
|
||||||
if len(text) < MIN_PAGE_CHARS:
|
if len(text) < MIN_PAGE_CHARS:
|
||||||
async with _pw_lock:
|
try:
|
||||||
try:
|
r_title, r_text, r_status, r_links = await _render_with_playwright(url, browser)
|
||||||
r_title, r_text, r_status, r_links = await _render_with_playwright(url, browser)
|
if len(r_text) > len(text):
|
||||||
except Exception as exc:
|
title, text, status, source, links = (
|
||||||
logger.info("deepsearch render failed for %s: %s", url, exc)
|
r_title or title,
|
||||||
r_title, r_text, r_status, r_links = "", "", 0, []
|
r_text,
|
||||||
if len(r_text) > len(text):
|
r_status or status,
|
||||||
title, text, status, source, links = (
|
"playwright",
|
||||||
r_title or title,
|
r_links,
|
||||||
r_text,
|
)
|
||||||
r_status or status,
|
except Exception as exc:
|
||||||
"playwright",
|
logger.info("deepsearch render failed for %s: %s", url, exc)
|
||||||
r_links,
|
|
||||||
)
|
|
||||||
if len(text) < MIN_PAGE_CHARS:
|
if len(text) < MIN_PAGE_CHARS:
|
||||||
return None
|
return None
|
||||||
return CrawledPage(
|
return CrawledPage(
|
||||||
@@ -302,99 +296,102 @@ async def crawl(
|
|||||||
total = min(len(level_candidates), max_pages)
|
total = min(len(level_candidates), max_pages)
|
||||||
cancelled = False
|
cancelled = False
|
||||||
|
|
||||||
|
pw = None
|
||||||
browser = None
|
browser = None
|
||||||
async with async_playwright() as pw:
|
try:
|
||||||
try:
|
pw = await async_playwright().__aenter__()
|
||||||
browser = await pw.chromium.launch(
|
browser = await pw.chromium.launch(
|
||||||
headless=True, args=["--no-sandbox", "--disable-dev-shm-usage"]
|
headless=True, args=["--no-sandbox", "--disable-dev-shm-usage"]
|
||||||
)
|
)
|
||||||
except Exception as exc:
|
except Exception as exc:
|
||||||
logger.warning("deepsearch playwright launch failed, pages will use httpx only: %s", exc)
|
logger.warning("deepsearch playwright launch failed, pages will use httpx only: %s", exc)
|
||||||
|
|
||||||
try:
|
try:
|
||||||
for level in range(max(1, depth)):
|
for level in range(max(1, depth)):
|
||||||
if cancelled or fetched >= max_pages or not level_candidates:
|
if cancelled or fetched >= max_pages or not level_candidates:
|
||||||
|
break
|
||||||
|
next_candidates: list[dict] = []
|
||||||
|
for start in range(0, len(level_candidates), CRAWL_CONCURRENCY):
|
||||||
|
if fetched >= max_pages:
|
||||||
break
|
break
|
||||||
next_candidates: list[dict] = []
|
if await should_stop():
|
||||||
for start in range(0, len(level_candidates), CRAWL_CONCURRENCY):
|
emit({"type": "stage", "stage": "cancelled", "message": "Crawl cancelled"})
|
||||||
|
cancelled = True
|
||||||
|
break
|
||||||
|
batch = level_candidates[start : start + CRAWL_CONCURRENCY][: max_pages - fetched]
|
||||||
|
for candidate in batch:
|
||||||
|
emit(
|
||||||
|
{
|
||||||
|
"type": "progress",
|
||||||
|
"done": fetched,
|
||||||
|
"total": total,
|
||||||
|
"url": candidate["url"],
|
||||||
|
"depth": level,
|
||||||
|
"message": f"Reading {candidate['url']}",
|
||||||
|
}
|
||||||
|
)
|
||||||
|
if is_cached(candidate["url"]):
|
||||||
|
emit({"type": "page_cached", "url": candidate["url"], "reason": "seen in a prior run"})
|
||||||
|
fetch_start = time.perf_counter()
|
||||||
|
results = await asyncio.gather(
|
||||||
|
*(_resolve_candidate(candidate, level, browser) for candidate in batch),
|
||||||
|
return_exceptions=True,
|
||||||
|
)
|
||||||
|
elapsed_ms = int((time.perf_counter() - fetch_start) * 1000)
|
||||||
|
for candidate, page in zip(batch, results):
|
||||||
|
url = candidate["url"]
|
||||||
|
if isinstance(page, BaseException):
|
||||||
|
logger.info("deepsearch fetch crashed for %s: %s", url, page)
|
||||||
|
page = None
|
||||||
|
if page is None:
|
||||||
|
emit(
|
||||||
|
{
|
||||||
|
"type": "page_skipped",
|
||||||
|
"url": url,
|
||||||
|
"reason": "no readable content",
|
||||||
|
"elapsed_ms": elapsed_ms,
|
||||||
|
}
|
||||||
|
)
|
||||||
|
continue
|
||||||
if fetched >= max_pages:
|
if fetched >= max_pages:
|
||||||
break
|
break
|
||||||
if await should_stop():
|
digest = content_hash(page.text)
|
||||||
emit({"type": "stage", "stage": "cancelled", "message": "Crawl cancelled"})
|
if digest in outcome.seen_hashes:
|
||||||
cancelled = True
|
|
||||||
break
|
|
||||||
batch = level_candidates[start : start + CRAWL_CONCURRENCY][: max_pages - fetched]
|
|
||||||
for candidate in batch:
|
|
||||||
emit(
|
emit(
|
||||||
{
|
{
|
||||||
"type": "progress",
|
"type": "page_duplicate",
|
||||||
"done": fetched,
|
"url": url,
|
||||||
"total": total,
|
"reason": "duplicate content",
|
||||||
"url": candidate["url"],
|
|
||||||
"depth": level,
|
|
||||||
"message": f"Reading {candidate['url']}",
|
|
||||||
}
|
|
||||||
)
|
|
||||||
if is_cached(candidate["url"]):
|
|
||||||
emit({"type": "page_cached", "url": candidate["url"], "reason": "seen in a prior run"})
|
|
||||||
fetch_start = time.perf_counter()
|
|
||||||
results = await asyncio.gather(
|
|
||||||
*(_resolve_candidate(candidate, level, browser) for candidate in batch),
|
|
||||||
return_exceptions=True,
|
|
||||||
)
|
|
||||||
elapsed_ms = int((time.perf_counter() - fetch_start) * 1000)
|
|
||||||
for candidate, page in zip(batch, results):
|
|
||||||
url = candidate["url"]
|
|
||||||
if isinstance(page, BaseException):
|
|
||||||
logger.info("deepsearch fetch crashed for %s: %s", url, page)
|
|
||||||
page = None
|
|
||||||
if page is None:
|
|
||||||
emit(
|
|
||||||
{
|
|
||||||
"type": "page_skipped",
|
|
||||||
"url": url,
|
|
||||||
"reason": "no readable content",
|
|
||||||
"elapsed_ms": elapsed_ms,
|
|
||||||
}
|
|
||||||
)
|
|
||||||
continue
|
|
||||||
if fetched >= max_pages:
|
|
||||||
break
|
|
||||||
digest = content_hash(page.text)
|
|
||||||
if digest in outcome.seen_hashes:
|
|
||||||
emit(
|
|
||||||
{
|
|
||||||
"type": "page_duplicate",
|
|
||||||
"url": url,
|
|
||||||
"reason": "duplicate content",
|
|
||||||
"elapsed_ms": elapsed_ms,
|
|
||||||
}
|
|
||||||
)
|
|
||||||
continue
|
|
||||||
outcome.seen_hashes.add(digest)
|
|
||||||
outcome.pages.append(page)
|
|
||||||
fetched += 1
|
|
||||||
emit(
|
|
||||||
{
|
|
||||||
"type": "page_loaded",
|
|
||||||
"url": page.url,
|
|
||||||
"title": page.title,
|
|
||||||
"source": page.source,
|
|
||||||
"depth": level,
|
|
||||||
"render": page.source == "playwright",
|
|
||||||
"elapsed_ms": elapsed_ms,
|
"elapsed_ms": elapsed_ms,
|
||||||
"done": fetched,
|
|
||||||
"total": total,
|
|
||||||
}
|
}
|
||||||
)
|
)
|
||||||
if level + 1 < depth:
|
continue
|
||||||
for link in relevant_links(page.links, query, LINKS_PER_PAGE):
|
outcome.seen_hashes.add(digest)
|
||||||
if link not in seen_urls:
|
outcome.pages.append(page)
|
||||||
seen_urls.add(link)
|
fetched += 1
|
||||||
next_candidates.append({"url": link})
|
emit(
|
||||||
level_candidates = next_candidates
|
{
|
||||||
total = min(total + len(next_candidates), max_pages)
|
"type": "page_loaded",
|
||||||
finally:
|
"url": page.url,
|
||||||
if browser is not None:
|
"title": page.title,
|
||||||
await browser.close()
|
"source": page.source,
|
||||||
|
"depth": level,
|
||||||
|
"render": page.source == "playwright",
|
||||||
|
"elapsed_ms": elapsed_ms,
|
||||||
|
"done": fetched,
|
||||||
|
"total": total,
|
||||||
|
}
|
||||||
|
)
|
||||||
|
if level + 1 < depth:
|
||||||
|
for link in relevant_links(page.links, query, LINKS_PER_PAGE):
|
||||||
|
if link not in seen_urls:
|
||||||
|
seen_urls.add(link)
|
||||||
|
next_candidates.append({"url": link})
|
||||||
|
level_candidates = next_candidates
|
||||||
|
total = min(total + len(next_candidates), max_pages)
|
||||||
|
finally:
|
||||||
|
if browser is not None:
|
||||||
|
await browser.close()
|
||||||
|
if pw is not None:
|
||||||
|
await pw.__aexit__(None, None, None)
|
||||||
return outcome
|
return outcome
|
||||||
|
|||||||
@@ -3,8 +3,34 @@
|
|||||||
{{ super() }}
|
{{ super() }}
|
||||||
<link rel="stylesheet" href="{{ static_url('/static/css/services.css') }}">
|
<link rel="stylesheet" href="{{ static_url('/static/css/services.css') }}">
|
||||||
<link rel="stylesheet" href="{{ static_url('/static/css/backups.css') }}">
|
<link rel="stylesheet" href="{{ static_url('/static/css/backups.css') }}">
|
||||||
|
<style>
|
||||||
|
.disk-warning {
|
||||||
|
background: #fff3cd;
|
||||||
|
border: 1px solid #ffc107;
|
||||||
|
border-radius: 6px;
|
||||||
|
padding: 12px 16px;
|
||||||
|
margin-bottom: 16px;
|
||||||
|
display: flex;
|
||||||
|
align-items: center;
|
||||||
|
gap: 10px;
|
||||||
|
color: #856404;
|
||||||
|
}
|
||||||
|
.disk-warning .icon {
|
||||||
|
font-size: 1.4em;
|
||||||
|
flex-shrink: 0;
|
||||||
|
}
|
||||||
|
.disk-warning strong {
|
||||||
|
font-weight: 600;
|
||||||
|
}
|
||||||
|
</style>
|
||||||
{% endblock %}
|
{% endblock %}
|
||||||
{% block admin_content %}
|
{% block admin_content %}
|
||||||
|
{% for w in disk_warnings %}
|
||||||
|
<div class="disk-warning" role="alert">
|
||||||
|
<span class="icon">⚠</span>
|
||||||
|
<span>Disk usage on <strong>{{ w.path }}</strong> is <strong>{{ w.used_pct }}%</strong> ({{ w.used_gb }} GB / {{ w.total_gb }} GB) - above {{ disk_warning_threshold }}% threshold. Consider cleaning up old data or increasing disk capacity.</span>
|
||||||
|
</div>
|
||||||
|
{% endfor %}
|
||||||
<div class="admin-toolbar">
|
<div class="admin-toolbar">
|
||||||
<h2>Backups</h2>
|
<h2>Backups</h2>
|
||||||
<div class="backup-controls">
|
<div class="backup-controls">
|
||||||
|
|||||||
@@ -0,0 +1,32 @@
|
|||||||
|
2026-07-19T09:01:55 INFO logging initialised at /workspace/repo/dpc.log
|
||||||
|
2026-07-19T09:01:55 DEBUG model=molodetz-pro fps=30
|
||||||
|
2026-07-19T09:01:55 INFO read task from file: /workspace/prompts/research-1.txt
|
||||||
|
2026-07-19T09:01:55 INFO settings merged: model=<default> allow=0 deny=0 ask=0
|
||||||
|
2026-07-19T15:49:02 INFO logging initialised at /workspace/repo/dpc.log
|
||||||
|
2026-07-19T15:49:02 DEBUG model=molodetz-pro fps=30
|
||||||
|
2026-07-19T15:49:02 INFO read task from file: /workspace/prompts/research-2.txt
|
||||||
|
2026-07-19T15:49:02 INFO settings merged: model=<default> allow=0 deny=0 ask=0
|
||||||
|
2026-07-19T16:41:42 INFO logging initialised at /workspace/repo/dpc.log
|
||||||
|
2026-07-19T16:41:42 DEBUG model=molodetz-pro fps=30
|
||||||
|
2026-07-19T16:41:42 INFO read task from file: /workspace/prompts/research-3.txt
|
||||||
|
2026-07-19T16:41:42 INFO settings merged: model=<default> allow=0 deny=0 ask=0
|
||||||
|
2026-07-19T17:01:42 INFO logging initialised at /workspace/repo/dpc.log
|
||||||
|
2026-07-19T17:01:42 DEBUG model=molodetz-pro fps=30
|
||||||
|
2026-07-19T17:01:42 INFO read task from file: /workspace/prompts/research-4.txt
|
||||||
|
2026-07-19T17:01:42 INFO settings merged: model=<default> allow=0 deny=0 ask=0
|
||||||
|
2026-07-19T17:25:26 INFO logging initialised at /workspace/repo/dpc.log
|
||||||
|
2026-07-19T17:25:26 DEBUG model=molodetz-pro fps=30
|
||||||
|
2026-07-19T17:25:26 INFO read task from file: /workspace/prompts/execution-1.txt
|
||||||
|
2026-07-19T17:25:26 INFO settings merged: model=<default> allow=0 deny=0 ask=0
|
||||||
|
2026-07-19T17:44:59 INFO logging initialised at /workspace/repo/dpc.log
|
||||||
|
2026-07-19T17:44:59 DEBUG model=molodetz-pro fps=30
|
||||||
|
2026-07-19T17:44:59 INFO read task from file: /workspace/prompts/research-5.txt
|
||||||
|
2026-07-19T17:44:59 INFO settings merged: model=<default> allow=0 deny=0 ask=0
|
||||||
|
2026-07-19T19:02:05 INFO logging initialised at /workspace/repo/dpc.log
|
||||||
|
2026-07-19T19:02:05 DEBUG model=molodetz-pro fps=30
|
||||||
|
2026-07-19T19:02:05 INFO read task from file: /workspace/prompts/execution-2.txt
|
||||||
|
2026-07-19T19:02:05 INFO settings merged: model=<default> allow=0 deny=0 ask=0
|
||||||
|
2026-07-23T01:50:04 INFO logging initialised at /workspace/repo/dpc.log
|
||||||
|
2026-07-23T01:50:04 DEBUG model=molodetz-pro fps=30
|
||||||
|
2026-07-23T01:50:04 INFO read task from file: /workspace/prompts/execution-2.txt
|
||||||
|
2026-07-23T01:50:04 INFO settings merged: model=<default> allow=0 deny=0 ask=0
|
||||||
Reference in New Issue
Block a user