mirror of
https://github.com/jkingsman/Remote-Terminal-for-MeshCore.git
synced 2026-08-07 09:13:04 +02:00
Add packet cleanup
This commit is contained in:
@@ -522,3 +522,14 @@ class RawPacketRepository:
|
||||
(int(time.time()), packet_id),
|
||||
)
|
||||
await db.conn.commit()
|
||||
|
||||
@staticmethod
|
||||
async def prune_old_undecrypted(max_age_days: int) -> int:
|
||||
"""Delete undecrypted packets older than max_age_days. Returns count deleted."""
|
||||
cutoff = int(time.time()) - (max_age_days * 86400)
|
||||
cursor = await db.conn.execute(
|
||||
"DELETE FROM raw_packets WHERE decrypted = 0 AND timestamp < ?",
|
||||
(cutoff,),
|
||||
)
|
||||
await db.conn.commit()
|
||||
return cursor.rowcount
|
||||
|
||||
@@ -1,6 +1,9 @@
|
||||
import os
|
||||
|
||||
from fastapi import APIRouter
|
||||
from pydantic import BaseModel
|
||||
|
||||
from app.config import settings
|
||||
from app.radio import radio_manager
|
||||
|
||||
|
||||
@@ -11,13 +14,23 @@ class HealthResponse(BaseModel):
|
||||
status: str
|
||||
radio_connected: bool
|
||||
serial_port: str | None
|
||||
database_size_mb: float
|
||||
|
||||
|
||||
@router.get("/health", response_model=HealthResponse)
|
||||
async def healthcheck() -> HealthResponse:
|
||||
"""Check if the API is running and if the radio is connected."""
|
||||
# Get database file size in MB
|
||||
db_size_mb = 0.0
|
||||
try:
|
||||
db_size_bytes = os.path.getsize(settings.database_path)
|
||||
db_size_mb = round(db_size_bytes / (1024 * 1024), 2)
|
||||
except OSError:
|
||||
pass
|
||||
|
||||
return HealthResponse(
|
||||
status="ok" if radio_manager.is_connected else "degraded",
|
||||
radio_connected=radio_manager.is_connected,
|
||||
serial_port=radio_manager.port,
|
||||
database_size_mb=db_size_mb,
|
||||
)
|
||||
|
||||
+165
-1
@@ -4,7 +4,8 @@ from hashlib import sha256
|
||||
from fastapi import APIRouter, BackgroundTasks
|
||||
from pydantic import BaseModel, Field
|
||||
|
||||
from app.decoder import try_decrypt_packet_with_channel_key
|
||||
from app.database import db
|
||||
from app.decoder import extract_payload, try_decrypt_packet_with_channel_key
|
||||
from app.packet_processor import create_message_from_decrypted
|
||||
from app.repository import RawPacketRepository
|
||||
|
||||
@@ -174,3 +175,166 @@ async def decrypt_historical_packets(
|
||||
async def get_decrypt_progress() -> DecryptProgress | None:
|
||||
"""Get the current progress of historical decryption."""
|
||||
return _decrypt_progress
|
||||
|
||||
|
||||
class MaintenanceRequest(BaseModel):
|
||||
prune_undecrypted_days: int = Field(
|
||||
ge=1,
|
||||
description="Delete undecrypted packets older than this many days"
|
||||
)
|
||||
|
||||
|
||||
class MaintenanceResult(BaseModel):
|
||||
packets_deleted: int
|
||||
vacuumed: bool
|
||||
|
||||
|
||||
@router.post("/maintenance", response_model=MaintenanceResult)
|
||||
async def run_maintenance(request: MaintenanceRequest) -> MaintenanceResult:
|
||||
"""
|
||||
Clean up old undecrypted packets and reclaim disk space.
|
||||
|
||||
- Deletes undecrypted packets older than the specified number of days
|
||||
- Runs VACUUM to reclaim disk space
|
||||
"""
|
||||
logger.info("Running maintenance: pruning packets older than %d days", request.prune_undecrypted_days)
|
||||
|
||||
# Prune old undecrypted packets
|
||||
deleted = await RawPacketRepository.prune_old_undecrypted(request.prune_undecrypted_days)
|
||||
logger.info("Deleted %d old undecrypted packets", deleted)
|
||||
|
||||
# Run VACUUM to reclaim space
|
||||
await db.conn.execute("VACUUM")
|
||||
logger.info("Database vacuumed")
|
||||
|
||||
return MaintenanceResult(packets_deleted=deleted, vacuumed=True)
|
||||
|
||||
|
||||
class DedupProgress(BaseModel):
|
||||
total: int
|
||||
processed: int
|
||||
duplicates_removed: int
|
||||
in_progress: bool
|
||||
|
||||
|
||||
class DedupResult(BaseModel):
|
||||
started: bool
|
||||
total_packets: int
|
||||
message: str
|
||||
|
||||
|
||||
# Global state for tracking dedup progress
|
||||
_dedup_progress: DedupProgress | None = None
|
||||
|
||||
|
||||
async def _run_payload_dedup() -> None:
|
||||
"""Background task to remove duplicate-payload packets."""
|
||||
global _dedup_progress
|
||||
|
||||
# Get all undecrypted packets
|
||||
packets = await RawPacketRepository.get_all_undecrypted()
|
||||
total = len(packets)
|
||||
|
||||
_dedup_progress = DedupProgress(
|
||||
total=total, processed=0, duplicates_removed=0, in_progress=True
|
||||
)
|
||||
|
||||
logger.info("Starting payload deduplication of %d packets", total)
|
||||
|
||||
# Group packets by payload hash
|
||||
payload_groups: dict[str, list[int]] = {} # hash -> list of packet IDs
|
||||
|
||||
for packet_id, packet_data, _timestamp in packets:
|
||||
payload = extract_payload(packet_data)
|
||||
if payload is None:
|
||||
continue
|
||||
|
||||
payload_hash = sha256(payload).hexdigest()
|
||||
|
||||
if payload_hash not in payload_groups:
|
||||
payload_groups[payload_hash] = []
|
||||
payload_groups[payload_hash].append(packet_id)
|
||||
|
||||
_dedup_progress = DedupProgress(
|
||||
total=total,
|
||||
processed=_dedup_progress.processed + 1,
|
||||
duplicates_removed=_dedup_progress.duplicates_removed,
|
||||
in_progress=True,
|
||||
)
|
||||
|
||||
# Delete duplicates (keep the first/oldest packet in each group)
|
||||
duplicates_removed = 0
|
||||
for payload_hash, packet_ids in payload_groups.items():
|
||||
if len(packet_ids) > 1:
|
||||
# Keep the first one, delete the rest
|
||||
ids_to_delete = packet_ids[1:]
|
||||
for packet_id in ids_to_delete:
|
||||
await db.conn.execute(
|
||||
"DELETE FROM raw_packets WHERE id = ?", (packet_id,)
|
||||
)
|
||||
duplicates_removed += 1
|
||||
|
||||
_dedup_progress = DedupProgress(
|
||||
total=total,
|
||||
processed=_dedup_progress.processed,
|
||||
duplicates_removed=duplicates_removed,
|
||||
in_progress=True,
|
||||
)
|
||||
|
||||
await db.conn.commit()
|
||||
|
||||
# Run VACUUM to reclaim space
|
||||
await db.conn.execute("VACUUM")
|
||||
|
||||
_dedup_progress = DedupProgress(
|
||||
total=total,
|
||||
processed=total,
|
||||
duplicates_removed=duplicates_removed,
|
||||
in_progress=False,
|
||||
)
|
||||
|
||||
logger.info("Payload deduplication complete: removed %d duplicates", duplicates_removed)
|
||||
|
||||
|
||||
@router.post("/dedup", response_model=DedupResult)
|
||||
async def deduplicate_packets(background_tasks: BackgroundTasks) -> DedupResult:
|
||||
"""
|
||||
Remove packets with duplicate payloads (keeps one copy of each unique payload).
|
||||
|
||||
This operation runs in the background and may take a long time for large databases.
|
||||
Use GET /packets/dedup/progress to check status.
|
||||
|
||||
Note: This only affects undecrypted packets. Packets that arrive through different
|
||||
repeater paths have different full data but the same payload - this removes those duplicates.
|
||||
"""
|
||||
global _dedup_progress
|
||||
|
||||
# Check if dedup is already in progress
|
||||
if _dedup_progress and _dedup_progress.in_progress:
|
||||
return DedupResult(
|
||||
started=False,
|
||||
total_packets=_dedup_progress.total,
|
||||
message=f"Deduplication already in progress: {_dedup_progress.processed}/{_dedup_progress.total}",
|
||||
)
|
||||
|
||||
# Get count of undecrypted packets
|
||||
count = await RawPacketRepository.get_undecrypted_count()
|
||||
if count == 0:
|
||||
return DedupResult(
|
||||
started=False, total_packets=0, message="No undecrypted packets to process"
|
||||
)
|
||||
|
||||
# Start background dedup
|
||||
background_tasks.add_task(_run_payload_dedup)
|
||||
|
||||
return DedupResult(
|
||||
started=True,
|
||||
total_packets=count,
|
||||
message=f"Started deduplication of {count} packets in background. This may take a while.",
|
||||
)
|
||||
|
||||
|
||||
@router.get("/dedup/progress", response_model=DedupProgress | None)
|
||||
async def get_dedup_progress() -> DedupProgress | None:
|
||||
"""Get the current progress of payload deduplication."""
|
||||
return _dedup_progress
|
||||
|
||||
@@ -1,9 +1,11 @@
|
||||
"""WebSocket router for real-time updates."""
|
||||
|
||||
import logging
|
||||
import os
|
||||
|
||||
from fastapi import APIRouter, WebSocket, WebSocketDisconnect
|
||||
|
||||
from app.config import settings
|
||||
from app.radio import radio_manager
|
||||
from app.repository import ChannelRepository, ContactRepository
|
||||
from app.websocket import ws_manager
|
||||
@@ -20,9 +22,18 @@ async def websocket_endpoint(websocket: WebSocket) -> None:
|
||||
# Send initial state
|
||||
try:
|
||||
# Health status
|
||||
db_size_mb = 0.0
|
||||
try:
|
||||
db_size_bytes = os.path.getsize(settings.database_path)
|
||||
db_size_mb = round(db_size_bytes / (1024 * 1024), 2)
|
||||
except OSError:
|
||||
pass
|
||||
|
||||
health_data = {
|
||||
"status": "ok" if radio_manager.is_connected else "degraded",
|
||||
"radio_connected": radio_manager.is_connected,
|
||||
"serial_port": radio_manager.port,
|
||||
"database_size_mb": db_size_mb,
|
||||
}
|
||||
await ws_manager.send_personal(websocket, "health", health_data)
|
||||
|
||||
|
||||
@@ -3,10 +3,13 @@
|
||||
import asyncio
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
from typing import Any
|
||||
|
||||
from fastapi import WebSocket
|
||||
|
||||
from app.config import settings
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
@@ -85,8 +88,17 @@ def broadcast_error(message: str, details: str | None = None) -> None:
|
||||
|
||||
def broadcast_health(radio_connected: bool, serial_port: str | None = None) -> None:
|
||||
"""Broadcast health status change to all connected clients."""
|
||||
# Get database file size in MB
|
||||
db_size_mb = 0.0
|
||||
try:
|
||||
db_size_bytes = os.path.getsize(settings.database_path)
|
||||
db_size_mb = round(db_size_bytes / (1024 * 1024), 2)
|
||||
except OSError:
|
||||
pass
|
||||
|
||||
asyncio.create_task(ws_manager.broadcast("health", {
|
||||
"status": "ok" if radio_connected else "degraded",
|
||||
"radio_connected": radio_connected,
|
||||
"serial_port": serial_port,
|
||||
"database_size_mb": db_size_mb,
|
||||
}))
|
||||
|
||||
Reference in New Issue
Block a user