From ad322543d089c7684ea0aebe49a8f7a82b73f86b Mon Sep 17 00:00:00 2001 From: 1445043649 <> Date: Wed, 13 May 2026 09:07:44 +0800 Subject: [PATCH] =?UTF-8?q?feat(persistence):=20backfill=5Fworkspace=5Fid.?= =?UTF-8?q?py=20=E2=80=94=20script=20skeleton?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit argparse with --dry-run, async orchestrator that loops the 4 business tables, and three step stubs (Step 1: users -> workspace creation; Step 2: UPDATE tables from users.default_workspace_id; Step 3: orphan rows -> legacy_workspace anchor at the LOCK'd nil UUID). Step bodies are filled in T5.5-T5.7 alongside their tests. Co-Authored-By: Claude Opus 4.7 (1M context) --- backend/scripts/backfill_workspace_id.py | 154 +++++++++++++++++++++++ 1 file changed, 154 insertions(+) create mode 100644 backend/scripts/backfill_workspace_id.py diff --git a/backend/scripts/backfill_workspace_id.py b/backend/scripts/backfill_workspace_id.py new file mode 100644 index 00000000..54a8866c --- /dev/null +++ b/backend/scripts/backfill_workspace_id.py @@ -0,0 +1,154 @@ +"""Backfill ``workspace_id`` on PR5 business tables. + +Three-step backfill (each idempotent — re-running picks up where a crash +left off because every step's WHERE clause filters already-processed rows): + + 1. For each user without ``default_workspace_id``: create a personal + workspace + ``owner`` membership + write back the user's + ``default_workspace_id``. + 2. ``UPDATE`` each of ``threads_meta`` / ``runs`` / ``feedback`` / + ``run_events`` setting ``workspace_id`` from the row's owner's + ``users.default_workspace_id``. Only touches rows where + ``workspace_id IS NULL`` and ``user_id IS NOT NULL``. + 3. Any rows still with ``workspace_id IS NULL`` (truly orphan — they had + ``user_id = NULL`` to begin with) are assigned the *legacy* workspace + UUID ``00000000-0000-0000-0000-000000000000``. The script creates + that workspace on demand, owned by the platform admin. + +Usage:: + + PYTHONPATH=. python scripts/backfill_workspace_id.py [--dry-run] + +T5.4 only ships the skeleton — the three step bodies are filled in by +T5.5 / T5.6 / T5.7 along with their per-step tests. +""" + +from __future__ import annotations + +import argparse +import asyncio +import logging +from typing import Any + +from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker + +logger = logging.getLogger(__name__) + +# The legacy workspace anchor. Stage 0 LOCK'd UUID — chosen as the standard +# nil UUID so SQL log scans can spot it instantly. +LEGACY_WORKSPACE_ID = "00000000-0000-0000-0000-000000000000" +LEGACY_WORKSPACE_SLUG = "legacy" +LEGACY_WORKSPACE_NAME = "Legacy Workspace" + +# The four business tables that gained ``workspace_id`` in alembic 0002. +_BUSINESS_TABLES: tuple[str, ...] = ("threads_meta", "runs", "feedback", "run_events") + + +async def _step1_create_workspaces_for_users( + session_factory: async_sessionmaker[AsyncSession], + *, + dry_run: bool, +) -> int: + """Create one workspace + owner membership for each user missing default_workspace_id. + + Filled in by T5.5. + """ + _ = session_factory # silenced until T5.5 + _ = dry_run + return 0 + + +async def _step2_update_table_from_users( + session_factory: async_sessionmaker[AsyncSession], + table: str, + *, + dry_run: bool, +) -> int: + """UPDATE *table* setting workspace_id from owner's users.default_workspace_id. + + Filled in by T5.6. + """ + _ = session_factory + _ = table + _ = dry_run + return 0 + + +async def _step3_assign_legacy_workspace( + session_factory: async_sessionmaker[AsyncSession], + table: str, + *, + dry_run: bool, +) -> int: + """Assign LEGACY_WORKSPACE_ID to *table* rows still missing workspace_id. + + Filled in by T5.7 (also responsible for ensuring the legacy workspace row exists). + """ + _ = session_factory + _ = table + _ = dry_run + return 0 + + +async def backfill( + session_factory: async_sessionmaker[AsyncSession], + *, + dry_run: bool = False, +) -> dict[str, Any]: + """Run all three backfill steps; return a per-step row-count report. + + Order matters: Step 1 must populate ``users.default_workspace_id`` + before Step 2 can correlate business rows back through ``users``. + """ + report: dict[str, Any] = {"dry_run": dry_run} + + report["users_workspaces_created"] = await _step1_create_workspaces_for_users(session_factory, dry_run=dry_run) + for table in _BUSINESS_TABLES: + report[f"{table}_from_users"] = await _step2_update_table_from_users(session_factory, table, dry_run=dry_run) + for table in _BUSINESS_TABLES: + report[f"{table}_legacy"] = await _step3_assign_legacy_workspace(session_factory, table, dry_run=dry_run) + + return report + + +def _build_session_factory_from_config() -> async_sessionmaker[AsyncSession]: + """Build an async session factory from the active config.yaml. + + Avoids importing on module load so unit tests can stub + ``session_factory`` directly without booting the full config pipeline. + """ + from deerflow.config import get_app_config + from deerflow.persistence.engine import get_session_factory, init_engine_from_config + + asyncio.run(init_engine_from_config(get_app_config().database)) + sf = get_session_factory() + if sf is None: + raise RuntimeError( + "database.backend=memory: nothing to backfill. Switch config.yaml to sqlite/postgres first.", + ) + return sf + + +def main() -> None: + parser = argparse.ArgumentParser(description="Backfill workspace_id on Stage 0 business tables (idempotent).") + parser.add_argument( + "--dry-run", + action="store_true", + help="Print the rows each step would touch without writing.", + ) + args = parser.parse_args() + + logging.basicConfig(level=logging.INFO, format="%(levelname)s: %(message)s") + + sf = _build_session_factory_from_config() + report = asyncio.run(backfill(sf, dry_run=args.dry_run)) + + logger.info("Backfill report (dry_run=%s):", args.dry_run) + for key, value in report.items(): + if key == "dry_run": + continue + logger.info(" %s: %s", key, value) + + +if __name__ == "__main__": + main()