Files
FabledScribe/src/scribe/services/backup.py
T
bvandeusen d6c8470ab2 feat(backup): v3 backup covers rulebooks, rules, events + join tables
The v2 backup silently dropped the entire rulebook system (rulebooks, topics,
rules), the project subscription/suppression join tables, and events — so a
'full' backup wasn't. v3 adds all of them with FK re-mapping on restore, and a
_not_included field that names the still-deferred tables (ACL groups/shares,
api_keys, embeddings, transient/operational) so the gap is explicit, not silent.

restore_full_backup routes v2 and v3 through one path; v3-only sections are
guarded by data.get so a v2 payload still restores cleanly.

Tests: version/coverage constants, pure join-table row helpers, and the export
contract via a mocked session (CI has no DB; full round-trip is a manual check).

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-09 23:52:14 -04:00

892 lines
34 KiB
Python

import logging
from datetime import date, datetime, timezone
from sqlalchemy import or_, select
from scribe.models import async_session
from scribe.models.event import Event
from scribe.models.milestone import Milestone
from scribe.models.note import Note
from scribe.models.note_draft import NoteDraft
from scribe.models.note_version import NoteVersion
from scribe.models.project import Project
from scribe.models.rulebook import (
Rule,
Rulebook,
RulebookTopic,
project_rule_suppressions,
project_rulebook_subscriptions,
project_topic_suppressions,
)
from scribe.models.setting import Setting
from scribe.models.task_log import TaskLog
from scribe.models.user import User
logger = logging.getLogger(__name__)
# Backup format version. v3 (2026-06) added rulebooks/topics/rules + their
# project subscription/suppression join tables, and events — all silently
# dropped by v2. Bump when the serialized schema changes.
BACKUP_VERSION = 3
# Tables intentionally NOT in the backup, surfaced in the payload so the gap is
# explicit rather than silent. ACL (groups/shares) is a coherent follow-up;
# embeddings are derived (regenerated from note bodies); api_keys are sensitive
# credentials; the rest are transient/operational.
_NOT_INCLUDED = [
"groups", "group_memberships", "project_shares", "note_shares",
"api_keys", "embeddings", "app_logs", "notifications", "invitations",
"password_resets", "user_profiles",
]
def _dt(val: str | None) -> datetime:
return datetime.fromisoformat(val) if val else datetime.now(timezone.utc)
def _d(val: str | None) -> date | None:
return date.fromisoformat(val) if val else None
def _subscription_rows(rows) -> list[dict]:
return [{"project_id": r.project_id, "rulebook_id": r.rulebook_id} for r in rows]
def _rule_suppression_rows(rows) -> list[dict]:
return [{"project_id": r.project_id, "rule_id": r.rule_id} for r in rows]
def _topic_suppression_rows(rows) -> list[dict]:
return [{"project_id": r.project_id, "topic_id": r.topic_id} for r in rows]
# ---------------------------------------------------------------------------
# Export
# ---------------------------------------------------------------------------
async def export_full_backup() -> dict:
"""Export all data as a version-3 JSON backup."""
async with async_session() as session:
users = (await session.execute(select(User))).scalars().all()
projects = (await session.execute(select(Project))).scalars().all()
milestones = (await session.execute(select(Milestone))).scalars().all()
notes = (await session.execute(select(Note))).scalars().all()
task_logs = (await session.execute(select(TaskLog))).scalars().all()
note_drafts = (await session.execute(select(NoteDraft))).scalars().all()
note_versions = (await session.execute(
select(NoteVersion).order_by(NoteVersion.note_id, NoteVersion.id)
)).scalars().all()
settings = (await session.execute(select(Setting))).scalars().all()
rulebooks = (await session.execute(select(Rulebook))).scalars().all()
topics = (await session.execute(select(RulebookTopic))).scalars().all()
rules = (await session.execute(select(Rule))).scalars().all()
events = (await session.execute(select(Event))).scalars().all()
subscriptions = (await session.execute(
select(project_rulebook_subscriptions)
)).all()
rule_suppressions = (await session.execute(
select(project_rule_suppressions)
)).all()
topic_suppressions = (await session.execute(
select(project_topic_suppressions)
)).all()
return {
"version": BACKUP_VERSION,
"scope": "full",
"exported_at": datetime.now(timezone.utc).isoformat(),
"_security_notice": (
"This backup contains hashed passwords. "
"Store it securely and restrict access."
),
"_not_included": _NOT_INCLUDED,
"users": [
{
"id": u.id,
"username": u.username,
"email": u.email,
"password_hash": u.password_hash,
"oauth_sub": u.oauth_sub,
"role": u.role,
"session_version": u.session_version,
"created_at": u.created_at.isoformat(),
}
for u in users
],
"projects": [
{
"id": p.id,
"user_id": p.user_id,
"title": p.title,
"description": p.description,
"goal": p.goal,
"status": p.status,
"color": p.color,
"created_at": p.created_at.isoformat(),
"updated_at": p.updated_at.isoformat(),
}
for p in projects
],
"milestones": [
{
"id": m.id,
"user_id": m.user_id,
"project_id": m.project_id,
"title": m.title,
"description": m.description,
"status": m.status,
"order_index": m.order_index,
"created_at": m.created_at.isoformat(),
"updated_at": m.updated_at.isoformat(),
}
for m in milestones
],
"notes": [
{
"id": n.id,
"user_id": n.user_id,
"title": n.title,
"body": n.body,
"tags": n.tags or [],
"parent_id": n.parent_id,
"project_id": n.project_id,
"milestone_id": n.milestone_id,
"status": n.status,
"priority": n.priority,
"due_date": n.due_date.isoformat() if n.due_date else None,
"created_at": n.created_at.isoformat(),
"updated_at": n.updated_at.isoformat(),
}
for n in notes
],
"task_logs": [
{
"id": tl.id,
"user_id": tl.user_id,
"task_id": tl.task_id,
"content": tl.content,
"duration_minutes": tl.duration_minutes,
"created_at": tl.created_at.isoformat(),
"updated_at": tl.updated_at.isoformat(),
}
for tl in task_logs
],
"note_drafts": [
{
"id": nd.id,
"user_id": nd.user_id,
"note_id": nd.note_id,
"proposed_body": nd.proposed_body,
"original_body": nd.original_body,
"instruction": nd.instruction,
"scope": nd.scope,
"created_at": nd.created_at.isoformat(),
"updated_at": nd.updated_at.isoformat(),
}
for nd in note_drafts
],
"note_versions": [
{
"id": nv.id,
"user_id": nv.user_id,
"note_id": nv.note_id,
"title": nv.title,
"body": nv.body,
"tags": nv.tags or [],
"pin_kind": nv.pin_kind,
"pin_label": nv.pin_label,
"created_at": nv.created_at.isoformat(),
}
for nv in note_versions
],
"settings": [
{"user_id": s.user_id, "key": s.key, "value": s.value}
for s in settings
],
"rulebooks": [
{
"id": rb.id,
"owner_user_id": rb.owner_user_id,
"title": rb.title,
"description": rb.description,
"always_on": rb.always_on,
"created_at": rb.created_at.isoformat(),
"updated_at": rb.updated_at.isoformat(),
}
for rb in rulebooks
],
"rulebook_topics": [
{
"id": t.id,
"rulebook_id": t.rulebook_id,
"title": t.title,
"description": t.description,
"order_index": t.order_index,
"created_at": t.created_at.isoformat(),
"updated_at": t.updated_at.isoformat(),
}
for t in topics
],
"rules": [
{
"id": r.id,
"topic_id": r.topic_id,
"project_id": r.project_id,
"title": r.title,
"statement": r.statement,
"why": r.why,
"how_to_apply": r.how_to_apply,
"order_index": r.order_index,
"created_at": r.created_at.isoformat(),
"updated_at": r.updated_at.isoformat(),
}
for r in rules
],
"rulebook_subscriptions": _subscription_rows(subscriptions),
"rule_suppressions": _rule_suppression_rows(rule_suppressions),
"topic_suppressions": _topic_suppression_rows(topic_suppressions),
"events": [
{
"id": e.id,
"user_id": e.user_id,
"project_id": e.project_id,
"uid": e.uid,
"caldav_uid": e.caldav_uid,
"title": e.title,
"start_dt": e.start_dt.isoformat() if e.start_dt else None,
"duration_minutes": e.duration_minutes,
"all_day": e.all_day,
"description": e.description,
"location": e.location,
"color": e.color,
"recurrence": e.recurrence,
"reminder_minutes": e.reminder_minutes,
"created_at": e.created_at.isoformat(),
"updated_at": e.updated_at.isoformat(),
}
for e in events
],
}
async def export_user_backup(user_id: int) -> dict:
"""Export a single user's data as a version-3 JSON backup."""
async with async_session() as session:
user = await session.get(User, user_id)
projects = (await session.execute(
select(Project).where(Project.user_id == user_id)
)).scalars().all()
project_ids = [p.id for p in projects]
milestones = (await session.execute(
select(Milestone).where(Milestone.user_id == user_id)
)).scalars().all()
notes = (await session.execute(
select(Note).where(Note.user_id == user_id)
)).scalars().all()
task_logs = (await session.execute(
select(TaskLog).where(TaskLog.user_id == user_id)
)).scalars().all()
note_drafts = (await session.execute(
select(NoteDraft).where(NoteDraft.user_id == user_id)
)).scalars().all()
note_versions = (await session.execute(
select(NoteVersion).where(NoteVersion.user_id == user_id)
.order_by(NoteVersion.note_id, NoteVersion.id)
)).scalars().all()
settings = (await session.execute(
select(Setting).where(Setting.user_id == user_id)
)).scalars().all()
rulebooks = (await session.execute(
select(Rulebook).where(Rulebook.owner_user_id == user_id)
)).scalars().all()
rulebook_ids = [rb.id for rb in rulebooks]
topics = (await session.execute(
select(RulebookTopic).where(RulebookTopic.rulebook_id.in_(rulebook_ids))
)).scalars().all() if rulebook_ids else []
topic_ids = [t.id for t in topics]
# Rules owned by the user = rules in the user's topics OR project-rules
# on the user's projects.
rule_filters = []
if topic_ids:
rule_filters.append(Rule.topic_id.in_(topic_ids))
if project_ids:
rule_filters.append(Rule.project_id.in_(project_ids))
rules = (await session.execute(
select(Rule).where(or_(*rule_filters))
)).scalars().all() if rule_filters else []
events = (await session.execute(
select(Event).where(Event.user_id == user_id)
)).scalars().all()
if project_ids:
subscriptions = (await session.execute(
select(project_rulebook_subscriptions).where(
project_rulebook_subscriptions.c.project_id.in_(project_ids)
)
)).all()
rule_suppressions = (await session.execute(
select(project_rule_suppressions).where(
project_rule_suppressions.c.project_id.in_(project_ids)
)
)).all()
topic_suppressions = (await session.execute(
select(project_topic_suppressions).where(
project_topic_suppressions.c.project_id.in_(project_ids)
)
)).all()
else:
subscriptions = rule_suppressions = topic_suppressions = []
return {
"version": BACKUP_VERSION,
"scope": "user",
"exported_at": datetime.now(timezone.utc).isoformat(),
"_not_included": _NOT_INCLUDED,
"user": {
"id": user.id,
"username": user.username,
"email": user.email,
"role": user.role,
"created_at": user.created_at.isoformat(),
} if user else None,
"projects": [
{
"id": p.id,
"user_id": p.user_id,
"title": p.title,
"description": p.description,
"goal": p.goal,
"status": p.status,
"color": p.color,
"created_at": p.created_at.isoformat(),
"updated_at": p.updated_at.isoformat(),
}
for p in projects
],
"milestones": [
{
"id": m.id,
"user_id": m.user_id,
"project_id": m.project_id,
"title": m.title,
"description": m.description,
"status": m.status,
"order_index": m.order_index,
"created_at": m.created_at.isoformat(),
"updated_at": m.updated_at.isoformat(),
}
for m in milestones
],
"notes": [
{
"id": n.id,
"user_id": n.user_id,
"title": n.title,
"body": n.body,
"tags": n.tags or [],
"parent_id": n.parent_id,
"project_id": n.project_id,
"milestone_id": n.milestone_id,
"status": n.status,
"priority": n.priority,
"due_date": n.due_date.isoformat() if n.due_date else None,
"created_at": n.created_at.isoformat(),
"updated_at": n.updated_at.isoformat(),
}
for n in notes
],
"task_logs": [
{
"id": tl.id,
"user_id": tl.user_id,
"task_id": tl.task_id,
"content": tl.content,
"duration_minutes": tl.duration_minutes,
"created_at": tl.created_at.isoformat(),
"updated_at": tl.updated_at.isoformat(),
}
for tl in task_logs
],
"note_drafts": [
{
"id": nd.id,
"user_id": nd.user_id,
"note_id": nd.note_id,
"proposed_body": nd.proposed_body,
"original_body": nd.original_body,
"instruction": nd.instruction,
"scope": nd.scope,
"created_at": nd.created_at.isoformat(),
"updated_at": nd.updated_at.isoformat(),
}
for nd in note_drafts
],
"note_versions": [
{
"id": nv.id,
"user_id": nv.user_id,
"note_id": nv.note_id,
"title": nv.title,
"body": nv.body,
"tags": nv.tags or [],
"pin_kind": nv.pin_kind,
"pin_label": nv.pin_label,
"created_at": nv.created_at.isoformat(),
}
for nv in note_versions
],
"settings": [
{"user_id": s.user_id, "key": s.key, "value": s.value}
for s in settings
],
"rulebooks": [
{
"id": rb.id,
"owner_user_id": rb.owner_user_id,
"title": rb.title,
"description": rb.description,
"always_on": rb.always_on,
"created_at": rb.created_at.isoformat(),
"updated_at": rb.updated_at.isoformat(),
}
for rb in rulebooks
],
"rulebook_topics": [
{
"id": t.id,
"rulebook_id": t.rulebook_id,
"title": t.title,
"description": t.description,
"order_index": t.order_index,
"created_at": t.created_at.isoformat(),
"updated_at": t.updated_at.isoformat(),
}
for t in topics
],
"rules": [
{
"id": r.id,
"topic_id": r.topic_id,
"project_id": r.project_id,
"title": r.title,
"statement": r.statement,
"why": r.why,
"how_to_apply": r.how_to_apply,
"order_index": r.order_index,
"created_at": r.created_at.isoformat(),
"updated_at": r.updated_at.isoformat(),
}
for r in rules
],
"rulebook_subscriptions": _subscription_rows(subscriptions),
"rule_suppressions": _rule_suppression_rows(rule_suppressions),
"topic_suppressions": _topic_suppression_rows(topic_suppressions),
"events": [
{
"id": e.id,
"user_id": e.user_id,
"project_id": e.project_id,
"uid": e.uid,
"caldav_uid": e.caldav_uid,
"title": e.title,
"start_dt": e.start_dt.isoformat() if e.start_dt else None,
"duration_minutes": e.duration_minutes,
"all_day": e.all_day,
"description": e.description,
"location": e.location,
"color": e.color,
"recurrence": e.recurrence,
"reminder_minutes": e.reminder_minutes,
"created_at": e.created_at.isoformat(),
"updated_at": e.updated_at.isoformat(),
}
for e in events
],
}
# ---------------------------------------------------------------------------
# Restore
# ---------------------------------------------------------------------------
async def restore_full_backup(data: dict) -> dict:
"""Restore from backup JSON. Dispatches by version."""
version = data.get("version", 1)
if version == 1:
return await _restore_v1(data)
# v2 and v3 share one path; v3-only sections are guarded by data.get so a
# v2 payload (without them) restores cleanly.
return await _restore_v2(data)
async def _restore_v1(data: dict) -> dict:
"""Restore legacy v1 backup (original format).
Pre-pivot v1 backups included conversations + messages; those are
skipped during restore now that the chat subsystem is gone.
"""
stats = {"users": 0, "notes": 0, "settings": 0}
async with async_session() as session:
user_id_map: dict[int, int] = {}
for u_data in data.get("users", []):
old_id = u_data["id"]
user = User(
username=u_data["username"],
email=u_data.get("email"),
password_hash=u_data["password_hash"],
role=u_data.get("role", "user"),
created_at=_dt(u_data.get("created_at")),
)
session.add(user)
await session.flush()
user_id_map[old_id] = user.id
stats["users"] += 1
note_id_map: dict[int, int] = {}
for n_data in data.get("notes", []):
old_id = n_data.get("id")
mapped_user_id = user_id_map.get(n_data.get("user_id", 0))
if mapped_user_id is None:
continue
note = Note(
user_id=mapped_user_id,
title=n_data.get("title", ""),
body=n_data.get("body", ""),
tags=n_data.get("tags", []),
parent_id=None, # patched below
status=n_data.get("status"),
priority=n_data.get("priority"),
due_date=_d(n_data.get("due_date")),
created_at=_dt(n_data.get("created_at")),
updated_at=_dt(n_data.get("updated_at")),
)
session.add(note)
await session.flush()
if old_id is not None:
note_id_map[old_id] = note.id
stats["notes"] += 1
# Patch parent_id now that all notes have new IDs
for n_data in data.get("notes", []):
old_id = n_data.get("id")
old_parent = n_data.get("parent_id")
if old_id and old_parent and old_id in note_id_map and old_parent in note_id_map:
note_row = await session.get(Note, note_id_map[old_id])
if note_row:
note_row.parent_id = note_id_map[old_parent]
for s_data in data.get("settings", []):
mapped_user_id = user_id_map.get(s_data.get("user_id", 0))
if mapped_user_id is None:
continue
session.add(Setting(user_id=mapped_user_id, key=s_data["key"], value=s_data.get("value", "")))
stats["settings"] += 1
await session.commit()
logger.info("Restored v1 backup: %s", stats)
return stats
async def _restore_v2(data: dict) -> dict:
"""Restore v2/v3 backup with full FK re-mapping.
Conversations + push subscriptions in pre-pivot backups are silently
skipped — those subsystems were removed in the MCP-first pivot. v3-only
sections (rulebooks/topics/rules/join-tables/events) are guarded by
data.get so a v2 payload restores without them.
"""
stats: dict[str, int] = {
"users": 0, "projects": 0, "milestones": 0, "notes": 0,
"task_logs": 0, "note_drafts": 0, "note_versions": 0,
"settings": 0, "rulebooks": 0, "rulebook_topics": 0, "rules": 0,
"rulebook_subscriptions": 0, "rule_suppressions": 0,
"topic_suppressions": 0, "events": 0,
}
async with async_session() as session:
user_id_map: dict[int, int] = {}
project_id_map: dict[int, int] = {}
milestone_id_map: dict[int, int] = {}
note_id_map: dict[int, int] = {}
rulebook_id_map: dict[int, int] = {}
topic_id_map: dict[int, int] = {}
rule_id_map: dict[int, int] = {}
# 1. Users
for u_data in data.get("users", []):
old_id = u_data["id"]
user = User(
username=u_data["username"],
email=u_data.get("email"),
password_hash=u_data.get("password_hash"),
oauth_sub=u_data.get("oauth_sub"),
role=u_data.get("role", "user"),
session_version=u_data.get("session_version", 1),
created_at=_dt(u_data.get("created_at")),
)
session.add(user)
await session.flush()
user_id_map[old_id] = user.id
stats["users"] += 1
# 2. Projects
for p_data in data.get("projects", []):
mapped_uid = user_id_map.get(p_data.get("user_id", 0))
if mapped_uid is None:
continue
proj = Project(
user_id=mapped_uid,
title=p_data.get("title", ""),
description=p_data.get("description", ""),
goal=p_data.get("goal", ""),
status=p_data.get("status", "active"),
color=p_data.get("color"),
created_at=_dt(p_data.get("created_at")),
updated_at=_dt(p_data.get("updated_at")),
)
session.add(proj)
await session.flush()
project_id_map[p_data["id"]] = proj.id
stats["projects"] += 1
# 3. Milestones
for m_data in data.get("milestones", []):
mapped_uid = user_id_map.get(m_data.get("user_id", 0))
mapped_pid = project_id_map.get(m_data.get("project_id", 0))
if mapped_uid is None or mapped_pid is None:
continue
ms = Milestone(
user_id=mapped_uid,
project_id=mapped_pid,
title=m_data.get("title", ""),
description=m_data.get("description"),
status=m_data.get("status", "active"),
order_index=m_data.get("order_index", 0),
created_at=_dt(m_data.get("created_at")),
updated_at=_dt(m_data.get("updated_at")),
)
session.add(ms)
await session.flush()
milestone_id_map[m_data["id"]] = ms.id
stats["milestones"] += 1
# 4a. Notes — first pass (no parent_id yet)
notes_with_parents: list[tuple[int, int]] = [] # (new_note_id, old_parent_id)
for n_data in data.get("notes", []):
mapped_uid = user_id_map.get(n_data.get("user_id", 0))
if mapped_uid is None:
continue
note = Note(
user_id=mapped_uid,
title=n_data.get("title", ""),
body=n_data.get("body", ""),
tags=n_data.get("tags", []),
parent_id=None,
project_id=project_id_map.get(n_data["project_id"]) if n_data.get("project_id") else None,
milestone_id=milestone_id_map.get(n_data["milestone_id"]) if n_data.get("milestone_id") else None,
status=n_data.get("status"),
priority=n_data.get("priority"),
due_date=_d(n_data.get("due_date")),
created_at=_dt(n_data.get("created_at")),
updated_at=_dt(n_data.get("updated_at")),
)
session.add(note)
await session.flush()
note_id_map[n_data["id"]] = note.id
if n_data.get("parent_id"):
notes_with_parents.append((note.id, n_data["parent_id"]))
stats["notes"] += 1
# 4b. Patch parent_id
for new_note_id, old_parent_id in notes_with_parents:
new_parent_id = note_id_map.get(old_parent_id)
if new_parent_id:
note_row = await session.get(Note, new_note_id)
if note_row:
note_row.parent_id = new_parent_id
# 5. TaskLogs
for tl_data in data.get("task_logs", []):
mapped_uid = user_id_map.get(tl_data.get("user_id", 0))
mapped_tid = note_id_map.get(tl_data.get("task_id", 0))
if mapped_uid is None or mapped_tid is None:
continue
tl = TaskLog(
user_id=mapped_uid,
task_id=mapped_tid,
content=tl_data.get("content", ""),
duration_minutes=tl_data.get("duration_minutes"),
created_at=_dt(tl_data.get("created_at")),
updated_at=_dt(tl_data.get("updated_at")),
)
session.add(tl)
stats["task_logs"] += 1
# 6. NoteDrafts
for nd_data in data.get("note_drafts", []):
mapped_uid = user_id_map.get(nd_data.get("user_id", 0))
mapped_nid = note_id_map.get(nd_data.get("note_id", 0))
if mapped_uid is None or mapped_nid is None:
continue
nd = NoteDraft(
user_id=mapped_uid,
note_id=mapped_nid,
proposed_body=nd_data.get("proposed_body", ""),
original_body=nd_data.get("original_body", ""),
instruction=nd_data.get("instruction", ""),
scope=nd_data.get("scope", "document"),
created_at=_dt(nd_data.get("created_at")),
updated_at=_dt(nd_data.get("updated_at")),
)
session.add(nd)
stats["note_drafts"] += 1
# 7. NoteVersions
for nv_data in data.get("note_versions", []):
mapped_uid = user_id_map.get(nv_data.get("user_id", 0))
mapped_nid = note_id_map.get(nv_data.get("note_id", 0))
if mapped_uid is None or mapped_nid is None:
continue
nv = NoteVersion(
user_id=mapped_uid,
note_id=mapped_nid,
title=nv_data.get("title", ""),
body=nv_data.get("body", ""),
tags=nv_data.get("tags", []),
pin_kind=nv_data.get("pin_kind"),
pin_label=nv_data.get("pin_label"),
created_at=_dt(nv_data.get("created_at")),
)
session.add(nv)
stats["note_versions"] += 1
# 8. Settings
for s_data in data.get("settings", []):
mapped_uid = user_id_map.get(s_data.get("user_id", 0))
if mapped_uid is None:
continue
session.add(Setting(user_id=mapped_uid, key=s_data["key"], value=s_data.get("value", "")))
stats["settings"] += 1
# 9. Rulebooks (v3)
for rb_data in data.get("rulebooks", []):
mapped_uid = user_id_map.get(rb_data.get("owner_user_id", 0))
if mapped_uid is None:
continue
rb = Rulebook(
owner_user_id=mapped_uid,
title=rb_data.get("title", ""),
description=rb_data.get("description", ""),
always_on=rb_data.get("always_on", False),
created_at=_dt(rb_data.get("created_at")),
updated_at=_dt(rb_data.get("updated_at")),
)
session.add(rb)
await session.flush()
rulebook_id_map[rb_data["id"]] = rb.id
stats["rulebooks"] += 1
# 10. Topics (v3)
for t_data in data.get("rulebook_topics", []):
mapped_rbid = rulebook_id_map.get(t_data.get("rulebook_id", 0))
if mapped_rbid is None:
continue
topic = RulebookTopic(
rulebook_id=mapped_rbid,
title=t_data.get("title", ""),
description=t_data.get("description"),
order_index=t_data.get("order_index", 0),
created_at=_dt(t_data.get("created_at")),
updated_at=_dt(t_data.get("updated_at")),
)
session.add(topic)
await session.flush()
topic_id_map[t_data["id"]] = topic.id
stats["rulebook_topics"] += 1
# 11. Rules (v3) — topic-rule (topic_id) XOR project-rule (project_id)
for r_data in data.get("rules", []):
mapped_topic = topic_id_map.get(r_data["topic_id"]) if r_data.get("topic_id") else None
mapped_proj = project_id_map.get(r_data["project_id"]) if r_data.get("project_id") else None
if mapped_topic is None and mapped_proj is None:
continue # orphaned — its parent didn't restore
rule = Rule(
topic_id=mapped_topic,
project_id=mapped_proj,
title=r_data.get("title", ""),
statement=r_data.get("statement", ""),
why=r_data.get("why") or None,
how_to_apply=r_data.get("how_to_apply") or None,
order_index=r_data.get("order_index", 0),
created_at=_dt(r_data.get("created_at")),
updated_at=_dt(r_data.get("updated_at")),
)
session.add(rule)
await session.flush()
rule_id_map[r_data["id"]] = rule.id
stats["rules"] += 1
# 12. Rulebook subscriptions (v3 join table)
for sub in data.get("rulebook_subscriptions", []):
mapped_pid = project_id_map.get(sub.get("project_id", 0))
mapped_rbid = rulebook_id_map.get(sub.get("rulebook_id", 0))
if mapped_pid is None or mapped_rbid is None:
continue
await session.execute(project_rulebook_subscriptions.insert().values(
project_id=mapped_pid, rulebook_id=mapped_rbid,
))
stats["rulebook_subscriptions"] += 1
# 13. Rule suppressions (v3 join table)
for sup in data.get("rule_suppressions", []):
mapped_pid = project_id_map.get(sup.get("project_id", 0))
mapped_rid = rule_id_map.get(sup.get("rule_id", 0))
if mapped_pid is None or mapped_rid is None:
continue
await session.execute(project_rule_suppressions.insert().values(
project_id=mapped_pid, rule_id=mapped_rid,
))
stats["rule_suppressions"] += 1
# 14. Topic suppressions (v3 join table)
for sup in data.get("topic_suppressions", []):
mapped_pid = project_id_map.get(sup.get("project_id", 0))
mapped_tid = topic_id_map.get(sup.get("topic_id", 0))
if mapped_pid is None or mapped_tid is None:
continue
await session.execute(project_topic_suppressions.insert().values(
project_id=mapped_pid, topic_id=mapped_tid,
))
stats["topic_suppressions"] += 1
# 15. Events (v3)
for e_data in data.get("events", []):
mapped_uid = user_id_map.get(e_data.get("user_id", 0))
if mapped_uid is None:
continue
ev = Event(
user_id=mapped_uid,
project_id=project_id_map.get(e_data["project_id"]) if e_data.get("project_id") else None,
uid=e_data.get("uid", ""),
caldav_uid=e_data.get("caldav_uid", ""),
title=e_data.get("title", ""),
start_dt=_dt(e_data.get("start_dt")),
duration_minutes=e_data.get("duration_minutes"),
all_day=e_data.get("all_day", False),
description=e_data.get("description", ""),
location=e_data.get("location", ""),
color=e_data.get("color", ""),
recurrence=e_data.get("recurrence"),
reminder_minutes=e_data.get("reminder_minutes"),
created_at=_dt(e_data.get("created_at")),
updated_at=_dt(e_data.get("updated_at")),
)
session.add(ev)
stats["events"] += 1
await session.commit()
logger.info("Restored v2/v3 backup: %s", stats)
return stats