apps/whatsapp_control_api/app/db.py
text
from __future__ import annotations
import json
import sqlite3
from contextlib import contextmanager
from pathlib import Path
from .core.config import get_settings
def init_db() -> None:
settings = get_settings()
settings.database_path.parent.mkdir(parents=True, exist_ok=True)
settings.task_logs_dir.mkdir(parents=True, exist_ok=True)
settings.task_artifacts_dir.mkdir(parents=True, exist_ok=True)
with get_connection() as conn:
conn.executescript(
"""
PRAGMA journal_mode=WAL;
CREATE TABLE IF NOT EXISTS users (
id TEXT PRIMARY KEY,
whatsapp_number TEXT NOT NULL UNIQUE,
display_name TEXT,
role TEXT NOT NULL,
status TEXT NOT NULL,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL,
last_seen_at TEXT
);
CREATE TABLE IF NOT EXISTS sessions (
id TEXT PRIMARY KEY,
user_id TEXT NOT NULL,
channel TEXT NOT NULL,
channel_chat_id TEXT NOT NULL,
state TEXT NOT NULL,
last_message_at TEXT NOT NULL,
context_summary TEXT,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS messages (
id TEXT PRIMARY KEY,
session_id TEXT,
user_id TEXT,
direction TEXT NOT NULL,
message_type TEXT NOT NULL,
wa_message_id TEXT UNIQUE,
reply_to_wa_id TEXT,
body_text TEXT,
payload_json TEXT NOT NULL,
status TEXT NOT NULL,
created_at TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS tasks (
id TEXT PRIMARY KEY,
user_id TEXT NOT NULL,
session_id TEXT,
source_message_id TEXT,
project_name TEXT,
channel TEXT NOT NULL,
prompt TEXT NOT NULL,
status TEXT NOT NULL,
current_step TEXT,
progress_percent INTEGER NOT NULL DEFAULT 0,
hermes_task_ref TEXT,
result_summary TEXT,
error_summary TEXT,
created_at TEXT NOT NULL,
started_at TEXT,
finished_at TEXT,
updated_at TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS task_logs (
id INTEGER PRIMARY KEY AUTOINCREMENT,
task_id TEXT NOT NULL,
level TEXT NOT NULL,
message TEXT NOT NULL,
created_at TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS files (
id TEXT PRIMARY KEY,
task_id TEXT NOT NULL,
kind TEXT NOT NULL,
storage_path TEXT,
public_url TEXT,
mime_type TEXT,
file_size_bytes INTEGER,
checksum_sha256 TEXT,
created_at TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS approvals (
id TEXT PRIMARY KEY,
task_id TEXT NOT NULL,
requested_by_user_id TEXT NOT NULL,
approved_by_user_id TEXT,
approval_type TEXT NOT NULL,
reason TEXT NOT NULL,
risk_summary TEXT NOT NULL,
status TEXT NOT NULL,
approval_code TEXT NOT NULL,
nonce_hash TEXT NOT NULL,
requested_at TEXT NOT NULL,
responded_at TEXT,
expires_at TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS config_changes (
id TEXT PRIMARY KEY,
user_id TEXT NOT NULL,
scope TEXT NOT NULL,
target_name TEXT,
status TEXT NOT NULL,
request_json TEXT NOT NULL,
request_summary TEXT NOT NULL,
reason TEXT NOT NULL,
requested_at TEXT NOT NULL,
responded_at TEXT,
applied_at TEXT,
error_summary TEXT
);
CREATE TABLE IF NOT EXISTS webhook_events (
id TEXT PRIMARY KEY,
provider TEXT NOT NULL,
event_type TEXT NOT NULL,
delivery_key TEXT NOT NULL UNIQUE,
signature_valid INTEGER NOT NULL,
payload_json TEXT NOT NULL,
received_at TEXT NOT NULL,
processed_at TEXT
);
CREATE TABLE IF NOT EXISTS audit_logs (
id INTEGER PRIMARY KEY AUTOINCREMENT,
actor_type TEXT NOT NULL,
actor_ref TEXT NOT NULL,
action TEXT NOT NULL,
target_type TEXT NOT NULL,
target_ref TEXT NOT NULL,
metadata_json TEXT NOT NULL,
created_at TEXT NOT NULL
);
"""
)
conn.commit()
@contextmanager
def get_connection() -> sqlite3.Connection:
settings = get_settings()
connection = sqlite3.connect(settings.database_path, check_same_thread=False)
connection.row_factory = sqlite3.Row
try:
yield connection
finally:
connection.close()
def dumps_json(value: dict) -> str:
return json.dumps(value, ensure_ascii=True, separators=(",", ":"))