Compare commits

..
Author SHA1 Message Date
Typosaurus 54a7805a90 ticket #85 attempt 1 2026-07-23 01:31:16 +00:00
Typosaurus 1d1efc0dfc ticket #85 attempt 1 2026-07-23 01:19:19 +00:00
12 changed files with 115 additions and 333 deletions
File diff suppressed because one or more lines are too long
+21
View File
@@ -30,6 +30,26 @@ def migrate_bug_tables_to_issue_tables() -> None:
logger.info("Dropped table %s after migration", source_name) logger.info("Dropped table %s after migration", source_name)
def _add_message_unique_index() -> None:
if "messages" not in db.tables:
return
with db:
db.query(
"""
DELETE FROM messages WHERE id NOT IN (
SELECT MIN(id) FROM messages GROUP BY sender_uid, receiver_uid, content
)
"""
)
_index(
db,
"messages",
"idx_messages_unique_sender_receiver_content",
["sender_uid", "receiver_uid", "content"],
unique=True,
)
def init_db(): def init_db():
tables = db.tables tables = db.tables
_index(db, "users", "idx_users_username", ["username"]) _index(db, "users", "idx_users_username", ["username"])
@@ -131,6 +151,7 @@ def init_db():
"idx_messages_conversation_rev", "idx_messages_conversation_rev",
["receiver_uid", "sender_uid"], ["receiver_uid", "sender_uid"],
) )
_add_message_unique_index()
_index(db, "notifications", "idx_notifications_user", ["user_uid"]) _index(db, "notifications", "idx_notifications_user", ["user_uid"])
_index(db, "notifications", "idx_notifications_user_read", ["user_uid", "read"]) _index(db, "notifications", "idx_notifications_user_read", ["user_uid", "read"])
_index(db, "push_registration", "idx_push_registration_user", ["user_uid"]) _index(db, "push_registration", "idx_push_registration_user", ["user_uid"])
-4
View File
@@ -1,7 +1,6 @@
# 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
@@ -65,7 +64,6 @@ 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,
@@ -74,8 +72,6 @@ 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)
-9
View File
@@ -53,13 +53,6 @@ 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 = ""
@@ -85,5 +78,3 @@ 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
-13
View File
@@ -34,10 +34,6 @@ 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()
@@ -125,15 +121,6 @@ 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,
-184
View File
@@ -1,7 +1,5 @@
# 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
@@ -11,8 +9,6 @@ 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",
@@ -379,58 +375,6 @@ 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 (
@@ -491,131 +435,3 @@ 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)
+13 -54
View File
@@ -81,7 +81,6 @@ 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
@@ -108,6 +107,19 @@ 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"]))
@@ -116,16 +128,6 @@ 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
@@ -134,49 +136,6 @@ 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"]
+1 -1
View File
@@ -88,7 +88,7 @@ def gateway_complete(
{"role": "system", "content": system}, {"role": "system", "content": system},
{"role": "user", "content": text}, {"role": "user", "content": text},
], ],
"temperature": 0.1, "temperature": 0.0,
} }
headers = { headers = {
"Content-Type": "application/json", "Content-Type": "application/json",
+59 -9
View File
@@ -1,6 +1,9 @@
# retoor <retoor@molodetz.nl> # retoor <retoor@molodetz.nl>
import hashlib
import logging import logging
import time
from collections import OrderedDict
from datetime import datetime, timezone from datetime import datetime, timezone
from typing import Any, Optional from typing import Any, Optional
@@ -15,12 +18,16 @@ from devplacepy.utils import (
track_action, track_action,
) )
from devplacepy.services.audit import record as audit from devplacepy.services.audit import record as audit
from sqlalchemy.exc import IntegrityError
from devplacepy.services.correction import schedule_correction from devplacepy.services.correction import schedule_correction
from devplacepy.services.ai_modifier import schedule_modification from devplacepy.services.ai_modifier import schedule_modification
logger = logging.getLogger("messaging.persist") logger = logging.getLogger("messaging.persist")
MAX_CONTENT_LENGTH = 2000 MAX_CONTENT_LENGTH = 2000
DEDUP_WINDOW_SECONDS = 3
_content_cache: dict[str, tuple[float, str]] = OrderedDict()
def _slim_attachment(attachment: dict[str, Any]) -> dict[str, Any]: def _slim_attachment(attachment: dict[str, Any]) -> dict[str, Any]:
@@ -81,19 +88,60 @@ def persist_message(
sender_uid = sender["uid"] sender_uid = sender["uid"]
sender_username = sender.get("username", "") sender_username = sender.get("username", "")
content_hash = hashlib.sha256(
f"{sender_uid}:{receiver_uid}:{content}".encode()
).hexdigest()[:16]
now = time.time()
last_seen, cached_uid = _content_cache.get(content_hash, (0.0, None))
if now - last_seen < DEDUP_WINDOW_SECONDS and cached_uid is not None:
logger.debug(
"Dedup hit for message hash %s (original uid %s)", content_hash, cached_uid
)
cached = get_table("messages").find_one(uid=cached_uid)
if cached:
return {
"uid": cached["uid"],
"sender_uid": cached["sender_uid"],
"receiver_uid": cached["receiver_uid"],
"content": cached["content"],
"read": cached.get("read", False),
"created_at": cached["created_at"],
}
messages_table = get_table("messages") messages_table = get_table("messages")
msg_uid = generate_uid() msg_uid = generate_uid()
created_at = datetime.now(timezone.utc).isoformat() created_at = datetime.now(timezone.utc).isoformat()
messages_table.insert(
{ try:
"uid": msg_uid, messages_table.insert(
"sender_uid": sender_uid, {
"receiver_uid": receiver_uid, "uid": msg_uid,
"content": content, "sender_uid": sender_uid,
"read": False, "receiver_uid": receiver_uid,
"created_at": created_at, "content": content,
"read": False,
"created_at": created_at,
}
)
except IntegrityError:
existing = messages_table.find_one(
sender_uid=sender_uid, receiver_uid=receiver_uid, content=content
)
if not existing:
raise
logger.debug(
"Dedup via unique constraint for message (uid %s)", existing["uid"]
)
_content_cache[content_hash] = (time.time(), existing["uid"])
return {
"uid": existing["uid"],
"sender_uid": existing["sender_uid"],
"receiver_uid": existing["receiver_uid"],
"content": existing["content"],
"read": existing.get("read", False),
"created_at": existing["created_at"],
} }
)
link_attachments(attachment_uids, "message", msg_uid) link_attachments(attachment_uids, "message", msg_uid)
schedule_correction(sender, "messages", msg_uid, request) schedule_correction(sender, "messages", msg_uid, request)
@@ -114,6 +162,8 @@ def persist_message(
) )
track_action(sender_uid, "message") track_action(sender_uid, "message")
_content_cache[content_hash] = (time.time(), msg_uid)
logger.info( logger.info(
"Message %s sent from %s to %s via %s", "Message %s sent from %s to %s via %s",
msg_uid, msg_uid,
-26
View File
@@ -3,34 +3,8 @@
{{ 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">&#x26A0;</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">
-32
View File
@@ -1,32 +0,0 @@
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
+20
View File
@@ -166,3 +166,23 @@ def test_send_attachment_only_empty_content_succeeds(seeded_db):
refresh_snapshot() refresh_snapshot()
row = get_table("messages").find_one(uid=msg["uid"]) row = get_table("messages").find_one(uid=msg["uid"])
assert row["content"] == "" assert row["content"] == ""
def test_duplicate_message_returns_same_uid(seeded_db):
s, _ = _member()
receiver = _db_user("bob_test")["uid"]
content = _unique("dupmsg")
first = s.post(
f"{BASE_URL}/messages/send",
headers=JSON_audit_log,
data={"content": content, "receiver_uid": receiver},
).json()["data"]
second = s.post(
f"{BASE_URL}/messages/send",
headers=JSON_audit_log,
data={"content": content, "receiver_uid": receiver},
).json()["data"]
assert first["uid"] == second["uid"], "duplicate messages should return the same uid"