Phase 2: 資料夾建立 & 改名同步 — Implementation Plan

For agentic workers: REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task.

Spec: design.md (§4.1, §4.2, §5, §8.6, §8.7, §9) Master Index: implementation-plan.md Depends on: Phase 1 完成(需要 OAuth token、Fernet crypto、tenant_drive_integrations 表)

Goal: PM 啟動專案時,自動在共用 Drive 帳號下建立完整資料夾結構(GuidantAI/Project/AP/CG/Control/AO/Task,6 層 — v0.2)。系統內任何 entity(含 task)改名 → best-effort 同步到 Drive。

Architecture: 引入 drive_folder_mappingsdrive_sync_jobs 兩張新表 + APScheduler 起 worker poller + channel renewer placeholder(Phase 3 才會用到)+ GoogleDriveApiClient(folder/permission ops)。把 worker 的 7 種 job_type 中的 INIT_PROJECT_FOLDERSCREATE_FOLDERRENAME_FOLDER 三種實作完成(其餘 4 種留 stub 給 Phase 3)。

v0.2 加 Task layer 後續變更(2026-04-24):原 Phase 2 已完成 5 層實作(Tasks 1–20)。後續以 Task #59 補做 6 層支援:

  • 新增 SQL migration 2026-04-24-add-task-scope-to-drive-folder-mappings.sql ALTER 既有 CHECK constraint 加入 'TASK'
  • DriveScopeType enum 補 TASK
  • ProjectTreeLoader 載 task 資料;InitProjectFoldersHandler 多一層遞迴
  • 新增 task creation hook(task add → enqueue CREATE_FOLDER scope=TASK)
  • Task rename hook(如有 task name 編輯 API) 詳見下方各 Task 的 v0.2 註記區塊與 changelog 2026-04-24-google-drive-folder-init-and-rename-sync.md 「後續變更」段。

Tech Stack: APScheduler / google-api-python-client (Drive v3) / SQLAlchemy advisory lock / dependency-injector


§1

File Structure (Phase 2)

後端新增 / 修改

scripts/sql/
├── 2026-MM-DD-google-drive-folder-mappings.sql        # 新增:drive_folder_mappings + drive_sync_jobs
└── 2026-04-24-add-task-scope-to-drive-folder-mappings.sql   # v0.2: ALTER CHECK 加入 'TASK'

infra/cloud_integration/
├── models/
│   ├── drive_folder_mapping.py                        # 新增
│   └── drive_sync_job.py                              # 新增
├── mapper/
│   ├── drive_folder_mapping_mapper.py                 # 新增
│   └── drive_sync_job_mapper.py                       # 新增
├── repository/
│   ├── drive_folder_mapping_repo_impl.py              # 新增
│   └── drive_sync_job_repo_impl.py                    # 新增
└── google_drive/
    ├── google_drive_api_client.py                     # 新增 (folders / permissions)
    └── google_drive_token_provider.py                 # 新增 (取 access token,含 refresh)

domain/cloud_integration/
├── entities/
│   ├── drive_folder_mapping_entity.py                 # 新增
│   └── drive_sync_job_entity.py                       # 新增
├── enums/
│   ├── __init__.py
│   ├── scope_type.py                                  # 新增 ROOT/PROJECT/AP/CG/CONTROL/AO
│   ├── sync_job_type.py                               # 新增 INIT_PROJECT_FOLDERS / ... 7 種
│   └── sync_job_status.py                             # 新增 PENDING/IN_PROGRESS/SUCCESS/FAILED/SKIPPED
├── repository/
│   ├── drive_folder_mapping_repository.py             # 新增
│   └── drive_sync_job_repository.py                   # 新增
└── service/
    ├── drive_folder_mapping_domain_service.py         # 新增 CRUD + 樹查詢
    ├── drive_sync_job_domain_service.py               # 新增 enqueue/claim/complete/fail
    └── google_drive_token_manager.py                  # 新增 (decrypt refresh, refresh access, return cached)

app/cloud_integration/
├── service/
│   ├── google_drive_sync_service.py                   # 新增 (對外 enqueue / list / retry API)
│   ├── handlers/
│   │   ├── __init__.py
│   │   ├── base_job_handler.py                        # 新增 ABC
│   │   ├── init_project_folders_handler.py            # 新增
│   │   ├── create_folder_handler.py                   # 新增
│   │   └── rename_folder_handler.py                   # 新增
│   └── drive_sync_worker.py                           # 新增 (poll / dispatch handler / retry)
└── dto/
    ├── drive_folder_mapping_dto.py
    └── drive_sync_job_dto.py

api/cloud_integration/
├── routes/
│   ├── google_drive_sync_route.py                     # 新增 sync trigger / sync jobs / init-folders
│   └── google_drive_integration_route.py              # 修改:status 增加 root_folder_url
└── serializers/
    └── drive_sync_job.py                              # 新增

di_containers/cloud_integration/
└── cloud_integration_containers.py                    # 修改:append new providers

core/
└── scheduler.py                                       # 新增 APScheduler 初始化
core/app_factory.py                                    # 修改:呼叫 scheduler init

main_app.py / main_socketio.py                        # 修改:start scheduler

# Project start integration
app/oscal/project/                                    # 路徑視專案實際位置
└── oscal_project_start_service.py                    # 修改:start 結束後 enqueue INIT_PROJECT_FOLDERS

# Rename hooks (各 entity 改名 service)
app/grc/service/project_service.py (or similar)       # 修改:改名後 enqueue RENAME_FOLDER
app/grc/service/assessment_plan_service.py            # 修改:同上
app/grc/service/control_group_service.py              # 修改:同上
app/grc/service/control_service.py                    # 修改:同上
app/grc/service/assessment_object_service.py          # 修改:同上

各 entity rename API 的位置需先用 grep 確認。Spec §5.4 已指明改名同步是 best-effort(失敗不阻擋系統改名)。

前端新增 / 修改

src/components/integrations/
├── GoogleDriveIntegrationCard.vue                     # 修改:root folder 連結 / 立即重建按鈕
└── SyncJobHistoryDialog.vue                           # 新增 失敗 jobs dialog

src/service/
└── CloudIntegrationService.js                         # 修改:加 listSyncJobs / retrySyncJob / initProjectFolders / triggerSync

src/config/api/api.js                                  # 修改:加 endpoint 常數

Tests

test/cloud_integration/
├── test_drive_folder_mapping_repo.py
├── test_drive_sync_job_repo.py
├── test_drive_folder_mapping_domain_service.py
├── test_drive_sync_job_domain_service.py
├── test_google_drive_token_manager.py
├── test_google_drive_api_client.py                    # mock googleapiclient
├── test_init_project_folders_handler.py
├── test_create_folder_handler.py
├── test_rename_folder_handler.py
├── test_drive_sync_worker.py
└── test_google_drive_sync_route.py

§2

Task List

Task 1: SQL Migration

Files:

  • Create: scripts/sql/<YYYY-MM-DD>-google-drive-folder-mappings.sql

v0.2 Addendum:原 migration(已執行)的 chk_drive_folder_mappings_scope 不含 'TASK'。Task #59 需新增獨立 migration scripts/sql/2026-04-24-add-task-scope-to-drive-folder-mappings.sql

-- Date: 2026-04-24
-- Purpose: drive_folder_mappings.scope_type 加入 'TASK' 以支援 6 層結構

-- 1. 替換 CHECK constraint 加入 TASK (2026-04-24)
ALTER TABLE compliance.drive_folder_mappings
    DROP CONSTRAINT IF EXISTS chk_drive_folder_mappings_scope;
ALTER TABLE compliance.drive_folder_mappings
    ADD CONSTRAINT chk_drive_folder_mappings_scope
    CHECK (scope_type IN ('ROOT','PROJECT','AP','CONTROL_GROUP','CONTROL','AO','TASK'));

不需要新增欄位 / 索引 — 既有 unique index (tenant_id, scope_type, scope_uid) 對 TASK 一樣適用(scope_uid = job_execution.uid)。

v0.3 Addendum:在 v0.2 task scope migration 之後,再新增 scripts/sql/2026-04-24-add-archive-scope-to-drive-folder-mappings.sql 把 ARCHIVE 加入 CHECK constraint:

-- Date: 2026-04-24
-- Purpose: drive_folder_mappings.scope_type 加入 'ARCHIVE'(v0.3 archive folder 機制)

-- 1. 替換 CHECK constraint 加入 ARCHIVE (2026-04-24)
ALTER TABLE compliance.drive_folder_mappings
    DROP CONSTRAINT IF EXISTS chk_drive_folder_mappings_scope;
ALTER TABLE compliance.drive_folder_mappings
    ADD CONSTRAINT chk_drive_folder_mappings_scope
    CHECK (scope_type IN ('ROOT','PROJECT','AP','CONTROL_GROUP','CONTROL','AO','TASK','ARCHIVE'));

同樣不需新增欄位 / 索引 — 既有 unique index 對 ARCHIVE 適用(scope_uid = AP.uid,每 AP 一筆)。

  • -- Date: <YYYY-MM-DD>
    -- Purpose: Drive folder mappings + DB-backed sync job queue
    
    -- 1. drive_folder_mappings (<YYYY-MM-DD>)
    CREATE TABLE IF NOT EXISTS compliance.drive_folder_mappings (
        id                          SERIAL PRIMARY KEY,
        uid                         UUID NOT NULL UNIQUE DEFAULT gen_random_uuid(),
        tenant_id                   INTEGER NOT NULL,
        scope_type                  VARCHAR(20) NOT NULL,
        scope_uid                   UUID,                              -- ROOT 為 NULL
        parent_drive_folder_id      VARCHAR(100),                      -- ROOT 為 NULL
        drive_folder_id             VARCHAR(100) NOT NULL UNIQUE,
        display_name_snapshot       VARCHAR(255) NOT NULL,
        is_unlinked                 BOOLEAN NOT NULL DEFAULT FALSE,
        created_at                  TIMESTAMPTZ NOT NULL DEFAULT NOW(),
        updated_at                  TIMESTAMPTZ NOT NULL DEFAULT NOW()
    );
    
    -- 2. drive_folder_mappings constraint & index (<YYYY-MM-DD>)
    ALTER TABLE compliance.drive_folder_mappings
        ADD CONSTRAINT chk_drive_folder_mappings_scope
        CHECK (scope_type IN ('ROOT','PROJECT','AP','CONTROL_GROUP','CONTROL','AO'));
    CREATE UNIQUE INDEX uq_drive_folder_mappings_scope
        ON compliance.drive_folder_mappings(tenant_id, scope_type, scope_uid)
        WHERE scope_uid IS NOT NULL;
    CREATE UNIQUE INDEX uq_drive_folder_mappings_root
        ON compliance.drive_folder_mappings(tenant_id, scope_type)
        WHERE scope_type = 'ROOT';
    CREATE INDEX idx_drive_folder_mappings_parent
        ON compliance.drive_folder_mappings(tenant_id, parent_drive_folder_id);
    
    -- 3. drive_folder_mappings GRANT (<YYYY-MM-DD>)
    GRANT SELECT, INSERT, UPDATE, DELETE ON compliance.drive_folder_mappings TO cm_app;
    GRANT USAGE, SELECT ON SEQUENCE compliance.drive_folder_mappings_id_seq TO cm_app;
    ALTER TABLE compliance.drive_folder_mappings ENABLE ROW LEVEL SECURITY;
    CREATE POLICY drive_folder_mappings_tenant_isolation
        ON compliance.drive_folder_mappings
        USING (tenant_id = current_setting('app.tenant_id', true)::INTEGER);
    
    -- 4. drive_sync_jobs (<YYYY-MM-DD>)
    CREATE TABLE IF NOT EXISTS compliance.drive_sync_jobs (
        id                  BIGSERIAL PRIMARY KEY,
        uid                 UUID NOT NULL UNIQUE DEFAULT gen_random_uuid(),
        tenant_id           INTEGER NOT NULL,
        job_type            VARCHAR(40) NOT NULL,
        payload             JSONB NOT NULL DEFAULT '{}',
        status              VARCHAR(20) NOT NULL DEFAULT 'PENDING',
        priority            INTEGER NOT NULL DEFAULT 100,
        retry_count         INTEGER NOT NULL DEFAULT 0,
        max_retries         INTEGER NOT NULL DEFAULT 5,
        next_run_at         TIMESTAMPTZ NOT NULL DEFAULT NOW(),
        last_error          TEXT,
        started_at          TIMESTAMPTZ,
        finished_at         TIMESTAMPTZ,
        created_at          TIMESTAMPTZ NOT NULL DEFAULT NOW()
    );
    
    -- 5. drive_sync_jobs constraint & index (<YYYY-MM-DD>)
    ALTER TABLE compliance.drive_sync_jobs
        ADD CONSTRAINT chk_drive_sync_jobs_status
        CHECK (status IN ('PENDING','IN_PROGRESS','SUCCESS','FAILED','SKIPPED'));
    ALTER TABLE compliance.drive_sync_jobs
        ADD CONSTRAINT chk_drive_sync_jobs_type
        CHECK (job_type IN ('INIT_PROJECT_FOLDERS','CREATE_FOLDER','RENAME_FOLDER',
                            'PROCESS_DRIVE_CHANGES','IMPORT_DRIVE_FILE',
                            'SOFT_DELETE_EVIDENCE','RECONCILE_AO_FOLDER'));
    CREATE INDEX idx_drive_sync_jobs_pending
        ON compliance.drive_sync_jobs(tenant_id, status, next_run_at)
        WHERE status = 'PENDING';
    CREATE INDEX idx_drive_sync_jobs_status
        ON compliance.drive_sync_jobs(status, created_at DESC);
    
    -- 6. drive_sync_jobs GRANT (<YYYY-MM-DD>)
    GRANT SELECT, INSERT, UPDATE, DELETE ON compliance.drive_sync_jobs TO cm_app;
    GRANT USAGE, SELECT ON SEQUENCE compliance.drive_sync_jobs_id_seq TO cm_app;
    ALTER TABLE compliance.drive_sync_jobs ENABLE ROW LEVEL SECURITY;
    CREATE POLICY drive_sync_jobs_tenant_isolation
        ON compliance.drive_sync_jobs
        USING (tenant_id = current_setting('app.tenant_id', true)::INTEGER);
  • git add scripts/sql/<YYYY-MM-DD>-google-drive-folder-mappings.sql
    git commit -m "feat(cloud_integration): add drive_folder_mappings + drive_sync_jobs migrations"

Task 2: Enums & Models

Files:

  • Create: domain/cloud_integration/enums/__init__.py
  • Create: domain/cloud_integration/enums/scope_type.py
  • Create: domain/cloud_integration/enums/sync_job_type.py
  • Create: domain/cloud_integration/enums/sync_job_status.py
  • Create: infra/cloud_integration/models/drive_folder_mapping.py
  • Create: infra/cloud_integration/models/drive_sync_job.py
  • # domain/cloud_integration/enums/scope_type.py
    from enum import StrEnum
    
    class DriveScopeType(StrEnum):
        ROOT = "ROOT"
        PROJECT = "PROJECT"
        AP = "AP"
        CONTROL_GROUP = "CONTROL_GROUP"
        CONTROL = "CONTROL"
        AO = "AO"
        TASK = "TASK"       # v0.2: task folder = 1:1 mapped to compliance.job_executions
        ARCHIVE = "ARCHIVE" # v0.3: per-AP archive folder ("_Archive/")

    v0.2 也建議將 RECONCILE_AO_FOLDER 改名為 RECONCILE_TASK_FOLDER(同步 enum + DB CHECK + handler 名稱),參見 Phase 3 plan。Phase 2 未實作此 handler 故只需更新 enum;DB CHECK constraint 變更見 Phase 3 SQL migration(job_type CHECK)。

    # domain/cloud_integration/enums/sync_job_type.py
    from enum import StrEnum
    
    class DriveSyncJobType(StrEnum):
        INIT_PROJECT_FOLDERS = "INIT_PROJECT_FOLDERS"
        CREATE_FOLDER = "CREATE_FOLDER"
        RENAME_FOLDER = "RENAME_FOLDER"
        PROCESS_DRIVE_CHANGES = "PROCESS_DRIVE_CHANGES"
        IMPORT_DRIVE_FILE = "IMPORT_DRIVE_FILE"
        SOFT_DELETE_EVIDENCE = "SOFT_DELETE_EVIDENCE"
        RECONCILE_AO_FOLDER = "RECONCILE_AO_FOLDER"
    # domain/cloud_integration/enums/sync_job_status.py
    from enum import StrEnum
    
    class DriveSyncJobStatus(StrEnum):
        PENDING = "PENDING"
        IN_PROGRESS = "IN_PROGRESS"
        SUCCESS = "SUCCESS"
        FAILED = "FAILED"
        SKIPPED = "SKIPPED"
  • # infra/cloud_integration/models/drive_folder_mapping.py
    import uuid
    from datetime import datetime
    
    from sqlalchemy import Column, Integer, String, Boolean, DateTime, CheckConstraint, Index
    from sqlalchemy.dialects.postgresql import UUID
    from jedi_common.session.database.db import Base
    
    
    class DriveFolderMapping(Base):
        __tablename__ = "drive_folder_mappings"
        __table_args__ = (
            CheckConstraint(
                "scope_type IN ('ROOT','PROJECT','AP','CONTROL_GROUP','CONTROL','AO')",
                name="chk_drive_folder_mappings_scope",
            ),
            Index("idx_drive_folder_mappings_parent", "tenant_id", "parent_drive_folder_id"),
            {"schema": "compliance"},
        )
    
        id = Column(Integer, primary_key=True)
        uid = Column(UUID(as_uuid=True), nullable=False, unique=True, default=uuid.uuid4)
        tenant_id = Column(Integer, nullable=False)
        scope_type = Column(String(20), nullable=False)
        scope_uid = Column(UUID(as_uuid=True))
        parent_drive_folder_id = Column(String(100))
        drive_folder_id = Column(String(100), nullable=False, unique=True)
        display_name_snapshot = Column(String(255), nullable=False)
        is_unlinked = Column(Boolean, nullable=False, default=False)
        created_at = Column(DateTime(timezone=True), nullable=False, default=datetime.utcnow)
        updated_at = Column(DateTime(timezone=True), nullable=False, default=datetime.utcnow, onupdate=datetime.utcnow)
    # infra/cloud_integration/models/drive_sync_job.py
    import uuid
    from datetime import datetime
    
    from sqlalchemy import Column, Integer, BigInteger, String, Text, DateTime, CheckConstraint
    from sqlalchemy.dialects.postgresql import UUID, JSONB
    from jedi_common.session.database.db import Base
    
    
    class DriveSyncJob(Base):
        __tablename__ = "drive_sync_jobs"
        __table_args__ = (
            CheckConstraint(
                "status IN ('PENDING','IN_PROGRESS','SUCCESS','FAILED','SKIPPED')",
                name="chk_drive_sync_jobs_status",
            ),
            {"schema": "compliance"},
        )
    
        id = Column(BigInteger, primary_key=True)
        uid = Column(UUID(as_uuid=True), nullable=False, unique=True, default=uuid.uuid4)
        tenant_id = Column(Integer, nullable=False)
        job_type = Column(String(40), nullable=False)
        payload = Column(JSONB, nullable=False, default=dict)
        status = Column(String(20), nullable=False, default="PENDING")
        priority = Column(Integer, nullable=False, default=100)
        retry_count = Column(Integer, nullable=False, default=0)
        max_retries = Column(Integer, nullable=False, default=5)
        next_run_at = Column(DateTime(timezone=True), nullable=False, default=datetime.utcnow)
        last_error = Column(Text)
        started_at = Column(DateTime(timezone=True))
        finished_at = Column(DateTime(timezone=True))
        created_at = Column(DateTime(timezone=True), nullable=False, default=datetime.utcnow)
  • git add domain/cloud_integration/enums/ infra/cloud_integration/models/drive_folder_mapping.py infra/cloud_integration/models/drive_sync_job.py
    git commit -m "feat(cloud_integration): add scope/job enums + folder mapping & sync job models"

Task 3: Entities, Mappers, Repos for DriveFolderMapping

Files:

  • Create: domain/cloud_integration/entities/drive_folder_mapping_entity.py
  • Create: domain/cloud_integration/repository/drive_folder_mapping_repository.py
  • Create: infra/cloud_integration/mapper/drive_folder_mapping_mapper.py
  • Create: infra/cloud_integration/repository/drive_folder_mapping_repo_impl.py
  • Create: test/cloud_integration/test_drive_folder_mapping_repo.py
  • from dataclasses import dataclass
    from datetime import datetime
    from typing import Optional
    from uuid import UUID
    
    @dataclass
    class DriveFolderMappingEntity:
        tenant_id: int
        scope_type: str
        drive_folder_id: str
        display_name_snapshot: str
        id: Optional[int] = None
        uid: Optional[UUID] = None
        scope_uid: Optional[UUID] = None
        parent_drive_folder_id: Optional[str] = None
        is_unlinked: bool = False
        created_at: Optional[datetime] = None
        updated_at: Optional[datetime] = None
  • from abc import ABC, abstractmethod
    from typing import Optional, List
    from uuid import UUID
    
    from domain.cloud_integration.entities.drive_folder_mapping_entity import DriveFolderMappingEntity
    
    class IDriveFolderMappingRepo(ABC):
        @abstractmethod
        def get_root(self, tenant_id: int) -> Optional[DriveFolderMappingEntity]: ...
        @abstractmethod
        def get_by_scope(self, tenant_id: int, scope_type: str, scope_uid: UUID) -> Optional[DriveFolderMappingEntity]: ...
        @abstractmethod
        def get_by_drive_folder_id(self, drive_folder_id: str) -> Optional[DriveFolderMappingEntity]: ...
        @abstractmethod
        def list_children(self, tenant_id: int, parent_drive_folder_id: str) -> List[DriveFolderMappingEntity]: ...
        @abstractmethod
        def add(self, entity: DriveFolderMappingEntity) -> DriveFolderMappingEntity: ...
        @abstractmethod
        def update(self, entity: DriveFolderMappingEntity) -> DriveFolderMappingEntity: ...
        @abstractmethod
        def delete_by_uid(self, uid: UUID) -> bool: ...
  • git add domain/cloud_integration/entities/drive_folder_mapping_entity.py \
           domain/cloud_integration/repository/drive_folder_mapping_repository.py \
           infra/cloud_integration/mapper/drive_folder_mapping_mapper.py \
           infra/cloud_integration/repository/drive_folder_mapping_repo_impl.py \
           test/cloud_integration/test_drive_folder_mapping_repo.py
    git commit -m "feat(cloud_integration): add DriveFolderMapping entity/mapper/repo + tests"

Task 4: Entities, Mappers, Repos for DriveSyncJob

Files:

  • Create: domain/cloud_integration/entities/drive_sync_job_entity.py
  • Create: domain/cloud_integration/repository/drive_sync_job_repository.py
  • Create: infra/cloud_integration/mapper/drive_sync_job_mapper.py
  • Create: infra/cloud_integration/repository/drive_sync_job_repo_impl.py
  • Create: test/cloud_integration/test_drive_sync_job_repo.py
  • from dataclasses import dataclass, field
    from datetime import datetime
    from typing import Optional, Dict, Any
    from uuid import UUID
    
    @dataclass
    class DriveSyncJobEntity:
        tenant_id: int
        job_type: str
        payload: Dict[str, Any] = field(default_factory=dict)
        id: Optional[int] = None
        uid: Optional[UUID] = None
        status: str = "PENDING"
        priority: int = 100
        retry_count: int = 0
        max_retries: int = 5
        next_run_at: Optional[datetime] = None
        last_error: Optional[str] = None
        started_at: Optional[datetime] = None
        finished_at: Optional[datetime] = None
        created_at: Optional[datetime] = None
  • from abc import ABC, abstractmethod
    from typing import Optional, List
    from uuid import UUID
    
    from domain.cloud_integration.entities.drive_sync_job_entity import DriveSyncJobEntity
    
    class IDriveSyncJobRepo(ABC):
        @abstractmethod
        def add(self, entity: DriveSyncJobEntity) -> DriveSyncJobEntity: ...
        @abstractmethod
        def claim_next_pending(self, batch_size: int = 5) -> List[DriveSyncJobEntity]:
            """Atomically pick up to batch_size pending jobs whose next_run_at <= now,
            mark them IN_PROGRESS + started_at=now, return entities.
            Use SELECT ... FOR UPDATE SKIP LOCKED to allow parallel workers."""
        @abstractmethod
        def mark_success(self, job_id: int) -> None: ...
        @abstractmethod
        def mark_failed(self, job_id: int, error: str, retry_at: Optional[datetime]) -> None:
            """If retry_at is None → status=FAILED; else status=PENDING + next_run_at=retry_at + retry_count+=1."""
        @abstractmethod
        def mark_skipped(self, job_id: int, reason: str) -> None: ...
        @abstractmethod
        def get_by_uid(self, uid: UUID) -> Optional[DriveSyncJobEntity]: ...
        @abstractmethod
        def list_by_tenant(self, tenant_id: int, status: Optional[str], page: int, page_size: int) -> List[DriveSyncJobEntity]: ...
        @abstractmethod
        def reset_to_pending(self, uid: UUID) -> bool:
            """For manual retry of FAILED jobs from admin UI."""
    • claim_next_pending:兩個 worker 並行(用兩個 session)只各領到一筆,不重複
    • mark_failedretry_at → status=PENDING + retry_count+1
    • mark_failed 不帶 → status=FAILED
    • 過了 max_retries 自動 FAILED(這個邏輯由 worker 控,repo 不主動)
  • from datetime import datetime, timezone
    from sqlalchemy import text
    
    def claim_next_pending(self, batch_size=5):
        now = datetime.now(timezone.utc)
        sql = text("""
            UPDATE compliance.drive_sync_jobs
            SET status = 'IN_PROGRESS', started_at = :now
            WHERE id IN (
                SELECT id FROM compliance.drive_sync_jobs
                WHERE status = 'PENDING' AND next_run_at <= :now
                ORDER BY priority, next_run_at
                LIMIT :batch
                FOR UPDATE SKIP LOCKED
            )
            RETURNING id
        """)
        result = self.session.execute(sql, {"now": now, "batch": batch_size})
        ids = [r[0] for r in result]
        if not ids:
            return []
        models = self.session.query(DriveSyncJob).filter(DriveSyncJob.id.in_(ids)).all()
        return [DriveSyncJobMapper.to_entity(m) for m in models]
  • git add domain/cloud_integration/entities/drive_sync_job_entity.py \
           domain/cloud_integration/repository/drive_sync_job_repository.py \
           infra/cloud_integration/mapper/drive_sync_job_mapper.py \
           infra/cloud_integration/repository/drive_sync_job_repo_impl.py \
           test/cloud_integration/test_drive_sync_job_repo.py
    git commit -m "feat(cloud_integration): add DriveSyncJob entity/mapper/repo + tests (incl. atomic claim)"

Task 5: Domain Services for Folder Mapping & Sync Job

Files:

  • Create: domain/cloud_integration/service/drive_folder_mapping_domain_service.py
  • Create: domain/cloud_integration/service/drive_sync_job_domain_service.py
  • Create: test/cloud_integration/test_drive_folder_mapping_domain_service.py
  • Create: test/cloud_integration/test_drive_sync_job_domain_service.py
  • from datetime import datetime
    from typing import Optional, List
    from uuid import UUID
    
    from domain.cloud_integration.entities.drive_folder_mapping_entity import DriveFolderMappingEntity
    from domain.cloud_integration.repository.drive_folder_mapping_repository import IDriveFolderMappingRepo
    from domain.cloud_integration.enums.scope_type import DriveScopeType
    
    
    class DriveFolderMappingDomainService:
        def __init__(self, repo: IDriveFolderMappingRepo):
            self._repo = repo
    
        def get_root(self, tenant_id: int) -> Optional[DriveFolderMappingEntity]:
            return self._repo.get_root(tenant_id)
    
        def get_by_scope(self, tenant_id: int, scope_type: str, scope_uid: UUID) -> Optional[DriveFolderMappingEntity]:
            return self._repo.get_by_scope(tenant_id, scope_type, scope_uid)
    
        def list_children(self, tenant_id: int, parent_drive_folder_id: str) -> List[DriveFolderMappingEntity]:
            return self._repo.list_children(tenant_id, parent_drive_folder_id)
    
        def upsert(self, entity: DriveFolderMappingEntity) -> DriveFolderMappingEntity:
            if entity.scope_type == DriveScopeType.ROOT:
                existing = self._repo.get_root(entity.tenant_id)
            elif entity.scope_uid:
                existing = self._repo.get_by_scope(entity.tenant_id, entity.scope_type, entity.scope_uid)
            else:
                existing = None
            if existing:
                entity.id = existing.id
                entity.uid = existing.uid
                return self._repo.update(entity)
            return self._repo.add(entity)
    
        def mark_unlinked_by_drive_id(self, drive_folder_id: str) -> Optional[DriveFolderMappingEntity]:
            mapping = self._repo.get_by_drive_folder_id(drive_folder_id)
            if mapping is None:
                return None
            mapping.is_unlinked = True
            return self._repo.update(mapping)
  • from datetime import datetime, timedelta, timezone
    from typing import Optional, List, Dict, Any
    from uuid import UUID
    
    from domain.cloud_integration.entities.drive_sync_job_entity import DriveSyncJobEntity
    from domain.cloud_integration.repository.drive_sync_job_repository import IDriveSyncJobRepo
    
    
    RETRY_DELAYS = [60, 300, 900, 3600, 21600]   # seconds (1m, 5m, 15m, 1h, 6h)
    
    
    class DriveSyncJobDomainService:
        def __init__(self, repo: IDriveSyncJobRepo):
            self._repo = repo
    
        def enqueue(self, tenant_id: int, job_type: str, payload: Dict[str, Any], priority: int = 100, max_retries: int = 5) -> DriveSyncJobEntity:
            entity = DriveSyncJobEntity(
                tenant_id=tenant_id,
                job_type=job_type,
                payload=payload,
                priority=priority,
                max_retries=max_retries,
                next_run_at=datetime.now(timezone.utc),
            )
            return self._repo.add(entity)
    
        def claim_batch(self, batch_size: int = 5) -> List[DriveSyncJobEntity]:
            return self._repo.claim_next_pending(batch_size)
    
        def succeed(self, job_id: int) -> None:
            self._repo.mark_success(job_id)
    
        def fail_or_retry(self, job: DriveSyncJobEntity, error: str) -> None:
            if job.retry_count + 1 >= job.max_retries:
                self._repo.mark_failed(job.id, error, retry_at=None)
                return
            delay_idx = min(job.retry_count, len(RETRY_DELAYS) - 1)
            retry_at = datetime.now(timezone.utc) + timedelta(seconds=RETRY_DELAYS[delay_idx])
            self._repo.mark_failed(job.id, error, retry_at=retry_at)
    
        def skip(self, job_id: int, reason: str) -> None:
            self._repo.mark_skipped(job_id, reason)
    
        def list_for_admin(self, tenant_id: int, status: Optional[str], page: int, page_size: int) -> List[DriveSyncJobEntity]:
            return self._repo.list_by_tenant(tenant_id, status, page, page_size)
    
        def manual_retry(self, uid: UUID) -> bool:
            return self._repo.reset_to_pending(uid)
  • git add domain/cloud_integration/service/drive_folder_mapping_domain_service.py \
           domain/cloud_integration/service/drive_sync_job_domain_service.py \
           test/cloud_integration/test_drive_folder_mapping_domain_service.py \
           test/cloud_integration/test_drive_sync_job_domain_service.py
    git commit -m "feat(cloud_integration): add domain services for folder mapping + sync job"

Task 6: GoogleDriveTokenManager

Files:

  • Create: domain/cloud_integration/service/google_drive_token_manager.py
  • Create: test/cloud_integration/test_google_drive_token_manager.py

注意:這個 service 同時依賴 domain (TenantDriveIntegrationDomainService) 與 infra (GoogleOAuthClient + TokenCryptoService)。為了保持 DDD 純度可以放 app/,但因為它表達的核心邏輯是「取一個 valid token」這個業務規則,放 domain/ 也合理(infra 是 abstraction 注入的)。

  • from datetime import datetime, timedelta, timezone
    from typing import Tuple
    
    from jedi_common.handler.exception import UnauthorizedError
    
    from common.code.grc_error_code import GrcErrorCode
    from domain.cloud_integration.entities.tenant_drive_integration_entity import TenantDriveIntegrationEntity
    from domain.cloud_integration.service.tenant_drive_integration_domain_service import TenantDriveIntegrationDomainService
    from domain.cloud_integration.service.token_crypto_service import TokenCryptoService
    from infra.cloud_integration.google_drive.google_oauth_client import GoogleOAuthClient
    
    
    TOKEN_EXPIRY_BUFFER_SECONDS = 60
    
    
    class DriveTokenRevokedError(Exception):
        """Raised when refresh token is invalid_grant."""
    
    
    class GoogleDriveTokenManager:
        def __init__(
            self,
            tenant_drive_integration_domain_service: TenantDriveIntegrationDomainService,
            google_oauth_client: GoogleOAuthClient,
            token_crypto_service: TokenCryptoService,
        ):
            self._domain = tenant_drive_integration_domain_service
            self._oauth = google_oauth_client
            self._crypto = token_crypto_service
    
        def get_valid_access_token(self, tenant_id: int) -> str:
            entity = self._domain.get_required(tenant_id)
            if entity.status != "CONNECTED":
                raise UnauthorizedError(GrcErrorCode.GRC_DRIVE_TOKEN_REVOKED)
    
            now = datetime.now(timezone.utc)
            if (
                entity.access_token_cache
                and entity.access_token_expires_at
                and entity.access_token_expires_at > now + timedelta(seconds=TOKEN_EXPIRY_BUFFER_SECONDS)
            ):
                return self._crypto.decrypt(entity.access_token_cache)
    
            # need refresh
            refresh_token = self._crypto.decrypt(entity.refresh_token_encrypted)
            try:
                tokens = self._oauth.refresh_access_token(refresh_token)
            except RuntimeError as e:
                if "invalid_grant" in str(e):
                    entity.status = "REVOKED"
                    entity.last_sync_error = str(e)
                    self._domain.update(entity)
                    raise DriveTokenRevokedError(str(e)) from e
                raise
    
            access_token = tokens["access_token"]
            expires_in = tokens.get("expires_in", 3600)
            entity.access_token_cache = self._crypto.encrypt(access_token)
            entity.access_token_expires_at = now + timedelta(seconds=expires_in - TOKEN_EXPIRY_BUFFER_SECONDS)
            self._domain.update(entity)
            return access_token
    • cache hit → 直接回,不 refresh
    • cache 過期 → call refresh,更新 entity
    • refresh invalid_grant → 標 REVOKED,raise DriveTokenRevokedError
    • status 不是 CONNECTED → raise UnauthorizedError
  • git add domain/cloud_integration/service/google_drive_token_manager.py test/cloud_integration/test_google_drive_token_manager.py
    git commit -m "feat(cloud_integration): add GoogleDriveTokenManager (cache + refresh + revoke detect)"

Task 7: GoogleDriveApiClient — Folder & Permission ops

Files:

  • Create: infra/cloud_integration/google_drive/google_drive_api_client.py
  • Create: test/cloud_integration/test_google_drive_api_client.py

v0.3 Addendum (move_to_archive helper):為支援 archive 機制,需在 GoogleDriveApiClient 加新 method:

def move_to_archive(
    self,
    tenant_id: int,
    file_id: str,
    archive_folder_id: str,
    current_parent_id: str,
) -> Dict[str, Any]:
    """把 file/folder 從 current_parent_id 搬到 archive_folder_id。
    用 Drive API files.update?addParents=&removeParents=(單筆 atomic move)。
    對 file 與 folder 都適用(folder 整個搬,內含檔案不動)。
    """
    return self._service(tenant_id).files().update(
        fileId=file_id,
        addParents=archive_folder_id,
        removeParents=current_parent_id,
        fields="id,name,parents",
    ).execute()

對應測試:mock googleapiclient,驗證呼叫帶正確 query params。

Caller(DriveSyncOrchestrationService)負責先查到 current_parent_id(可從 evidence 對應 task 的 task folder ID,或從 mapping 取);對單一檔案 archive,current_parent_id 為該檔案目前所在的 task folder ID(由 webhook 變更或 evidence 上的 metadata 推得;若取不到,可以 drive.files.get(fileId, fields="parents") 拿)。

  • # infra/cloud_integration/google_drive/google_drive_api_client.py
    from typing import Optional, Dict, Any, List
    
    from google.oauth2.credentials import Credentials
    from googleapiclient.discovery import build
    
    
    GOOGLE_FOLDER_MIME = "application/vnd.google-apps.folder"
    
    
    class GoogleDriveApiClient:
        """Thin wrapper over googleapiclient. Constructor takes a callable that returns
        a fresh access token (so we delegate refresh logic to TokenManager)."""
    
        def __init__(self, token_provider):
            """
            token_provider: callable(tenant_id) -> str
            """
            self._token_provider = token_provider
    
        def _service(self, tenant_id: int):
            access_token = self._token_provider(tenant_id)
            creds = Credentials(token=access_token)
            return build("drive", "v3", credentials=creds, cache_discovery=False)
    
        # ── Folder ops ──────────────────────────────
        def create_folder(self, tenant_id: int, name: str, parent_id: Optional[str]) -> Dict[str, Any]:
            body = {"name": name, "mimeType": GOOGLE_FOLDER_MIME}
            if parent_id:
                body["parents"] = [parent_id]
            return self._service(tenant_id).files().create(body=body, fields="id,name,parents,webViewLink").execute()
    
        def rename(self, tenant_id: int, file_id: str, new_name: str) -> Dict[str, Any]:
            return self._service(tenant_id).files().update(
                fileId=file_id, body={"name": new_name}, fields="id,name"
            ).execute()
    
        def get_folder(self, tenant_id: int, file_id: str) -> Optional[Dict[str, Any]]:
            try:
                return self._service(tenant_id).files().get(
                    fileId=file_id, fields="id,name,parents,trashed,webViewLink"
                ).execute()
            except Exception as e:
                if "404" in str(e) or "notFound" in str(e):
                    return None
                raise
    
        def list_children(self, tenant_id: int, parent_id: str) -> List[Dict[str, Any]]:
            q = f"'{parent_id}' in parents and trashed=false"
            page_token = None
            out = []
            svc = self._service(tenant_id)
            while True:
                resp = svc.files().list(
                    q=q, fields="nextPageToken, files(id,name,mimeType,parents,modifiedTime,size,webViewLink)",
                    pageSize=200, pageToken=page_token,
                ).execute()
                out.extend(resp.get("files", []))
                page_token = resp.get("nextPageToken")
                if not page_token:
                    break
            return out
    
        # ── Permission ops ──────────────────────────
        def share_anyone_writer(self, tenant_id: int, file_id: str) -> Dict[str, Any]:
            return self._service(tenant_id).permissions().create(
                fileId=file_id,
                body={"type": "anyone", "role": "writer", "allowFileDiscovery": False},
                fields="id",
            ).execute()
  • git add infra/cloud_integration/google_drive/google_drive_api_client.py test/cloud_integration/test_google_drive_api_client.py
    git commit -m "feat(cloud_integration): add GoogleDriveApiClient (folders/permissions/list_children)"

Task 8: Job Handler 基底類別

Files:

  • Create: app/cloud_integration/service/handlers/__init__.py
  • Create: app/cloud_integration/service/handlers/base_job_handler.py
  • from abc import ABC, abstractmethod
    from domain.cloud_integration.entities.drive_sync_job_entity import DriveSyncJobEntity
    
    
    class BaseJobHandler(ABC):
        @property
        @abstractmethod
        def job_type(self) -> str: ...
    
        @abstractmethod
        def handle(self, job: DriveSyncJobEntity) -> None:
            """Execute the job. Raise on failure (worker handles retry)."""
  • git add app/cloud_integration/service/handlers/
    git commit -m "feat(cloud_integration): add BaseJobHandler abstract class"

Task 9: CreateFolderHandler

Files:

  • Create: app/cloud_integration/service/handlers/create_folder_handler.py
  • Create: test/cloud_integration/test_create_folder_handler.py
  • from uuid import UUID
    
    from jedi_common.session.database.db import transaction
    
    from app.cloud_integration.service.handlers.base_job_handler import BaseJobHandler
    from domain.cloud_integration.entities.drive_folder_mapping_entity import DriveFolderMappingEntity
    from domain.cloud_integration.enums.sync_job_type import DriveSyncJobType
    from domain.cloud_integration.service.drive_folder_mapping_domain_service import DriveFolderMappingDomainService
    from infra.cloud_integration.google_drive.google_drive_api_client import GoogleDriveApiClient
    
    
    class CreateFolderHandler(BaseJobHandler):
        """Payload: { scope_type, scope_uid (nullable for ROOT), parent_drive_folder_id (nullable for ROOT), name }"""
    
        def __init__(
            self,
            folder_mapping_domain_service: DriveFolderMappingDomainService,
            drive_api_client: GoogleDriveApiClient,
        ):
            self._folder_domain = folder_mapping_domain_service
            self._drive = drive_api_client
    
        @property
        def job_type(self) -> str:
            return DriveSyncJobType.CREATE_FOLDER
    
        @transaction
        def handle(self, job):
            payload = job.payload
            scope_type = payload["scope_type"]
            scope_uid = UUID(payload["scope_uid"]) if payload.get("scope_uid") else None
            parent_drive_folder_id = payload.get("parent_drive_folder_id")
            name = payload["name"]
    
            drive_folder = self._drive.create_folder(job.tenant_id, name=name, parent_id=parent_drive_folder_id)
            drive_folder_id = drive_folder["id"]
    
            self._drive.share_anyone_writer(job.tenant_id, drive_folder_id)
    
            mapping = DriveFolderMappingEntity(
                tenant_id=job.tenant_id,
                scope_type=scope_type,
                scope_uid=scope_uid,
                parent_drive_folder_id=parent_drive_folder_id,
                drive_folder_id=drive_folder_id,
                display_name_snapshot=name,
                is_unlinked=False,
            )
            self._folder_domain.upsert(mapping)
  • git add app/cloud_integration/service/handlers/create_folder_handler.py test/cloud_integration/test_create_folder_handler.py
    git commit -m "feat(cloud_integration): add CreateFolderHandler (single folder creation + permission)"

Task 10: RenameFolderHandler

Files:

  • Create: app/cloud_integration/service/handlers/rename_folder_handler.py
  • Create: test/cloud_integration/test_rename_folder_handler.py
  • from uuid import UUID
    
    from jedi_common.handler.exception import NotFound
    from jedi_common.session.database.db import transaction
    
    from common.code.grc_error_code import GrcErrorCode
    from app.cloud_integration.service.handlers.base_job_handler import BaseJobHandler
    from domain.cloud_integration.enums.sync_job_type import DriveSyncJobType
    from domain.cloud_integration.service.drive_folder_mapping_domain_service import DriveFolderMappingDomainService
    from infra.cloud_integration.google_drive.google_drive_api_client import GoogleDriveApiClient
    
    
    class RenameFolderHandler(BaseJobHandler):
        """Payload: { scope_type, scope_uid, new_name }"""
    
        def __init__(
            self,
            folder_mapping_domain_service: DriveFolderMappingDomainService,
            drive_api_client: GoogleDriveApiClient,
        ):
            self._folder_domain = folder_mapping_domain_service
            self._drive = drive_api_client
    
        @property
        def job_type(self) -> str:
            return DriveSyncJobType.RENAME_FOLDER
    
        @transaction
        def handle(self, job):
            scope_type = job.payload["scope_type"]
            scope_uid = UUID(job.payload["scope_uid"])
            new_name = job.payload["new_name"]
    
            mapping = self._folder_domain.get_by_scope(job.tenant_id, scope_type, scope_uid)
            if mapping is None or mapping.is_unlinked:
                raise NotFound(GrcErrorCode.GRC_DRIVE_FOLDER_MAPPING_NOT_FOUND)
    
            self._drive.rename(job.tenant_id, mapping.drive_folder_id, new_name)
            mapping.display_name_snapshot = new_name
            self._folder_domain.upsert(mapping)
  • git add app/cloud_integration/service/handlers/rename_folder_handler.py test/cloud_integration/test_rename_folder_handler.py
    git commit -m "feat(cloud_integration): add RenameFolderHandler"

Task 11: InitProjectFoldersHandler

Files:

  • Create: app/cloud_integration/service/handlers/init_project_folders_handler.py
  • Create: test/cloud_integration/test_init_project_folders_handler.py

此 handler 是整個 phase 最複雜的:要拉 project / AP / control_group / control / AO / task 樹狀資料,逐層建資料夾(v0.2: 6 層)。為了避免一次 transaction 太大太久,建議每建立一個資料夾就 commit 並寫對應 mapping。

v0.2 Addendum:以下 sample code 仍為 5 層版本(與已實作 code 一致)。Task #59 補做 6 層時要:

  1. for ao in ctrl["aos"] 內補一層 for task in ao["tasks"] 迴圈:
    for ao in ctrl["aos"]:
        ao_mapping = self._ensure_folder(
            tenant_id, DriveScopeType.AO, ao["uid"], ctrl_mapping.drive_folder_id,
            ao["title"],   # AO 不再加 [code] prefix(v0.2)
        )
        for task in ao["tasks"]:
            self._ensure_folder(
                tenant_id, DriveScopeType.TASK, task["uid"], ao_mapping.drive_folder_id,
                f"[T-{task['uid'][:8]}] {task['name']}",
            )
  2. AP 名稱改用 f"第{ap['round_no']}輪 - {ap['name']}" 格式
  3. ProjectTreeLoader.load() 多 query task — 每個 AO 拉 compliance.job_executions 找該 AO 在該 AP 週期下的所有 task。實作時需 grep 確認 task / job_execution 對應 query 介面(候選:ext_workflow_execution_domain_service 內既有方法、或在 JobExecutionDomainServicelist_by_ao_uid_in_ap(ao_uid, ap_uid)
  4. 確認 task entity 的 name 欄位實際命名(候選 name / display_name / description)— 若都沒有,先 fallback 為 [T-{task_id}] 純 prefix

v0.3 Addendum (Archive folder per AP):每建一個 AP folder 後,緊接著建一個 _Archive/ 子資料夾並登錄為 scope_type=ARCHIVE, scope_uid=ap_uid。在 AP 迴圈內 ap_mapping = self._ensure_folder(... AP ...) 之後加:

# v0.3: 每個 AP 下建 _Archive/ folder(用於系統刪除時 archive 檔案 / task folder)
self._ensure_folder(
    tenant_id,
    DriveScopeType.ARCHIVE,
    ap["uid"],                          # scope_uid = AP.uid(每 AP 一筆)
    ap_mapping.drive_folder_id,         # parent 為 AP folder
    "_Archive",                         # 固定名稱(底線開頭使排序置頂)
)

後續 archive 操作(try_archive_drive_file / try_archive_task_folder)取此 mapping 拿 archive folder ID。Idempotency:_ensure_folder 已有「mapping 存在且 not unlinked → skip」邏輯,重跑 INIT 不會建重。

Permission:archive folder 也走 share_anyone_writer(同其他 folder),讓 user 可手動操作(拖回 / 下載歸檔證據)。

  • from typing import List
    from uuid import UUID
    
    from jedi_common.session.database.db import transaction
    
    from app.cloud_integration.service.handlers.base_job_handler import BaseJobHandler
    from domain.cloud_integration.entities.drive_folder_mapping_entity import DriveFolderMappingEntity
    from domain.cloud_integration.enums.scope_type import DriveScopeType
    from domain.cloud_integration.enums.sync_job_type import DriveSyncJobType
    from domain.cloud_integration.service.drive_folder_mapping_domain_service import DriveFolderMappingDomainService
    from domain.cloud_integration.service.tenant_drive_integration_domain_service import TenantDriveIntegrationDomainService
    from infra.cloud_integration.google_drive.google_drive_api_client import GoogleDriveApiClient
    
    
    ROOT_FOLDER_NAME = "GuidantAI"
    
    
    class InitProjectFoldersHandler(BaseJobHandler):
        """Payload: { project_uid }
        Creates Project / AP / CG / Control / AO folders for the given project under tenant root."""
    
        def __init__(
            self,
            folder_mapping_domain_service: DriveFolderMappingDomainService,
            tenant_drive_integration_domain_service: TenantDriveIntegrationDomainService,
            drive_api_client: GoogleDriveApiClient,
            project_tree_loader,   # injected: provides project metadata + nested AP/CG/Control/AO
        ):
            self._folder_domain = folder_mapping_domain_service
            self._tenant_domain = tenant_drive_integration_domain_service
            self._drive = drive_api_client
            self._tree_loader = project_tree_loader
    
        @property
        def job_type(self) -> str:
            return DriveSyncJobType.INIT_PROJECT_FOLDERS
    
        def handle(self, job):
            project_uid = UUID(job.payload["project_uid"])
            tenant_id = job.tenant_id
    
            # 1. 確保根目錄存在
            root = self._ensure_root(tenant_id)
    
            # 2. 拉 project 樹(從 jedi_oscal / jedi_project 等)
            tree = self._tree_loader.load(project_uid)
            # tree 結構:{
            #     'name': 'Project Name',
            #     'aps': [{'uid', 'name', 'control_groups': [{'uid','code','name','controls':[{...,'aos':[...]}]}]}]
            # }
    
            # 3. 逐層建立
            self._ensure_folder(tenant_id, DriveScopeType.PROJECT, project_uid, root.drive_folder_id, tree["name"])
            for ap in tree["aps"]:
                ap_parent_id = self._get_drive_id(tenant_id, DriveScopeType.PROJECT, project_uid)
                ap_mapping = self._ensure_folder(tenant_id, DriveScopeType.AP, ap["uid"], ap_parent_id, ap["name"])
    
                for cg in ap["control_groups"]:
                    cg_mapping = self._ensure_folder(
                        tenant_id, DriveScopeType.CONTROL_GROUP, cg["uid"], ap_mapping.drive_folder_id,
                        f"[{cg['code']}] {cg['name']}",
                    )
                    for ctrl in cg["controls"]:
                        ctrl_mapping = self._ensure_folder(
                            tenant_id, DriveScopeType.CONTROL, ctrl["uid"], cg_mapping.drive_folder_id,
                            f"[{ctrl['code']}] {ctrl['name']}",
                        )
                        for ao in ctrl["aos"]:
                            self._ensure_folder(
                                tenant_id, DriveScopeType.AO, ao["uid"], ctrl_mapping.drive_folder_id,
                                f"[{ao['code']}] {ao['title']}",
                            )
    
        # ── helpers ────────────────────────────────────────
        @transaction
        def _ensure_root(self, tenant_id: int) -> DriveFolderMappingEntity:
            existing = self._folder_domain.get_root(tenant_id)
            if existing and not existing.is_unlinked:
                return existing
            drive = self._drive.create_folder(tenant_id, ROOT_FOLDER_NAME, parent_id=None)
            self._drive.share_anyone_writer(tenant_id, drive["id"])
            mapping = DriveFolderMappingEntity(
                tenant_id=tenant_id,
                scope_type=DriveScopeType.ROOT,
                drive_folder_id=drive["id"],
                display_name_snapshot=ROOT_FOLDER_NAME,
            )
            saved = self._folder_domain.upsert(mapping)
            # update tenant_drive_integrations.root_folder_id
            tenant_int = self._tenant_domain.get_required(tenant_id)
            tenant_int.root_folder_id = drive["id"]
            self._tenant_domain.update(tenant_int)
            return saved
    
        @transaction
        def _ensure_folder(self, tenant_id, scope_type, scope_uid, parent_drive_id, name) -> DriveFolderMappingEntity:
            if isinstance(scope_uid, str):
                scope_uid = UUID(scope_uid)
            existing = self._folder_domain.get_by_scope(tenant_id, scope_type, scope_uid)
            if existing and not existing.is_unlinked:
                return existing
            drive = self._drive.create_folder(tenant_id, name, parent_id=parent_drive_id)
            self._drive.share_anyone_writer(tenant_id, drive["id"])
            mapping = DriveFolderMappingEntity(
                tenant_id=tenant_id,
                scope_type=scope_type,
                scope_uid=scope_uid,
                parent_drive_folder_id=parent_drive_id,
                drive_folder_id=drive["id"],
                display_name_snapshot=name,
            )
            return self._folder_domain.upsert(mapping)
    
        def _get_drive_id(self, tenant_id, scope_type, scope_uid) -> str:
            if isinstance(scope_uid, str):
                scope_uid = UUID(scope_uid)
            mapping = self._folder_domain.get_by_scope(tenant_id, scope_type, scope_uid)
            return mapping.drive_folder_id

    Noteproject_tree_loader 是一個小 service,負責從 jedi_oscal / jedi_project 載入該 project 的完整 AP/CG/Control/AO 樹。建議獨立 class(見 Step 2)。

  • from uuid import UUID
    
    # 注意:這個 loader 透過既有 domain services 取資料,禁止直接 import infra ORM
    from domain.grc.service.grc_project_domain_service import GrcProjectDomainService    # 路徑視實際
    # ... 其他 domain services
    
    class ProjectTreeLoader:
        def __init__(
            self,
            grc_project_domain_service,
            assessment_plan_domain_service,
            control_group_domain_service,
            control_domain_service,
            assessment_object_domain_service,
        ):
            self._project = grc_project_domain_service
            self._ap = assessment_plan_domain_service
            self._cg = control_group_domain_service
            self._ctrl = control_domain_service
            self._ao = assessment_object_domain_service
    
        def load(self, project_uid: UUID) -> dict:
            project = self._project.get_by_uid(project_uid)
            aps = self._ap.list_by_project(project.id)
            tree = {"name": project.name, "aps": []}
            for ap in aps:
                cgs = self._cg.list_by_ap(ap.id)
                ap_node = {"uid": str(ap.uid), "name": ap.name, "control_groups": []}
                for cg in cgs:
                    ctrls = self._ctrl.list_by_control_group(cg.id)
                    cg_node = {"uid": str(cg.uid), "code": cg.code, "name": cg.name, "controls": []}
                    for ctrl in ctrls:
                        aos = self._ao.list_by_control(ctrl.id)
                        ctrl_node = {"uid": str(ctrl.uid), "code": ctrl.code, "name": ctrl.name, "aos": []}
                        for ao in aos:
                            ctrl_node["aos"].append({"uid": str(ao.uid), "code": ao.code, "title": ao.title})
                        cg_node["controls"].append(ctrl_node)
                    ap_node["control_groups"].append(cg_node)
                tree["aps"].append(ap_node)
            return tree

    實作前請用 grep 找出實際 domain services 名稱與 method(可能跟 jedi-oscal / jedi-project 內 already-existing service 重疊)。如果找不到 list method,需要在對應 domain service 補一個。

    • 一個 project / 1 AP / 1 CG / 1 Control / 1 AO → 5 次 create_folder 呼叫
    • 已有 mapping → 跳過 create
    • root 不存在 → 先建 root + update tenant integration
  • git add app/cloud_integration/service/handlers/init_project_folders_handler.py \
           app/cloud_integration/service/project_tree_loader.py \
           test/cloud_integration/test_init_project_folders_handler.py
    git commit -m "feat(cloud_integration): add InitProjectFoldersHandler + ProjectTreeLoader"

Task 12: DriveSyncWorker

Files:

  • Create: app/cloud_integration/service/drive_sync_worker.py
  • Create: test/cloud_integration/test_drive_sync_worker.py
  • import logging
    from typing import Dict, List
    
    from app.cloud_integration.service.handlers.base_job_handler import BaseJobHandler
    from domain.cloud_integration.entities.drive_sync_job_entity import DriveSyncJobEntity
    from domain.cloud_integration.service.drive_sync_job_domain_service import DriveSyncJobDomainService
    
    logger = logging.getLogger(__name__)
    
    
    class DriveSyncWorker:
        def __init__(
            self,
            drive_sync_job_domain_service: DriveSyncJobDomainService,
            handlers: List[BaseJobHandler],
        ):
            self._domain = drive_sync_job_domain_service
            self._handlers: Dict[str, BaseJobHandler] = {h.job_type: h for h in handlers}
    
        def run_once(self, batch_size: int = 5) -> int:
            """Process up to batch_size pending jobs. Returns how many were processed."""
            jobs = self._domain.claim_batch(batch_size)
            for job in jobs:
                handler = self._handlers.get(job.job_type)
                if handler is None:
                    logger.warning("No handler for job_type=%s, skipping", job.job_type)
                    self._domain.skip(job.id, f"No handler for {job.job_type}")
                    continue
                try:
                    handler.handle(job)
                    self._domain.succeed(job.id)
                    logger.info("Job %s (%s) succeeded", job.id, job.job_type)
                except Exception as e:
                    logger.exception("Job %s (%s) failed", job.id, job.job_type)
                    self._domain.fail_or_retry(job, str(e))
            return len(jobs)
    • 無 pending → 回 0,handler 不被呼叫
    • 1 pending + handler 成功 → succeed 被呼叫
    • 1 pending + handler 拋 → fail_or_retry 被呼叫
    • 1 pending + 沒 handler → skip 被呼叫
  • git add app/cloud_integration/service/drive_sync_worker.py test/cloud_integration/test_drive_sync_worker.py
    git commit -m "feat(cloud_integration): add DriveSyncWorker (poll + dispatch + retry)"

Task 13: APScheduler 整合

Files:

  • Create: core/scheduler.py
  • Modify: core/extensions.py
  • Modify: core/app_factory.py (and/or main_app.py / main_socketio.py)
  • ⚠️ eventlet 相容性:專案 main_socketio.py 使用 eventlet(已 monkey-patch)。 APScheduler 的 BackgroundScheduler + ThreadPoolExecutor 與 eventlet 同 process 共存時,threading 可能被 eventlet 偷換成 greenlet,行為不確定。 建議

    • 若部署用 main_socketio.py (eventlet) → 改用 apscheduler.executors.pool.ProcessPoolExecutor 或更乾淨地用 apscheduler.schedulers.gevent.GeventScheduler(如裝 gevent)。最簡單:在 main_socketio.py 啟動 scheduler 前先 monkey.patch_all() 已 done,BackgroundScheduler 在多數情境仍可工作但會有 warning。
    • 若部署用 main_app.py(無 eventlet) → BackgroundScheduler 直接可用。
    • 實作時先用 BackgroundScheduler,部署到 staging 後測試 WebSocket + scheduler 同時運作是否穩定;不穩定再切換 executor。
    # core/scheduler.py
    import logging
    from typing import Optional
    
    from apscheduler.schedulers.background import BackgroundScheduler
    from apscheduler.executors.pool import ThreadPoolExecutor
    
    logger = logging.getLogger(__name__)
    
    
    _scheduler: Optional[BackgroundScheduler] = None
    
    
    def init_scheduler(app, worker_pool_size: int = 4) -> BackgroundScheduler:
        global _scheduler
        if _scheduler is not None:
            return _scheduler
        executors = {"default": ThreadPoolExecutor(max_workers=worker_pool_size)}
        _scheduler = BackgroundScheduler(executors=executors, timezone="UTC")
    
        # 5 秒一次跑 worker(會在後續 task 加 channel renewer)
        from di_containers.containers import Containers
    
        def _run_worker_job():
            # 取 worker singleton from DI
            worker = Containers.cloud_integration_container.drive_sync_worker()
            try:
                with app.app_context():
                    worker.run_once(batch_size=10)
            except Exception:
                logger.exception("DriveSyncWorker tick failed")
    
        _scheduler.add_job(_run_worker_job, "interval", seconds=5, id="drive_sync_worker", max_instances=1)
        _scheduler.start()
        logger.info("APScheduler started with DriveSyncWorker (5s interval)")
        return _scheduler
    
    
    def shutdown_scheduler():
        global _scheduler
        if _scheduler is not None:
            _scheduler.shutdown(wait=False)
            _scheduler = None
  • from core.scheduler import init_scheduler
    init_scheduler(app, worker_pool_size=app.config.get("DRIVE_SYNC_WORKER_POOL_SIZE", 4))

    並在 atexit 註冊 shutdown:

    import atexit
    from core.scheduler import shutdown_scheduler
    atexit.register(shutdown_scheduler)
  • git add core/scheduler.py core/app_factory.py main_app.py main_socketio.py
    git commit -m "feat(cloud_integration): integrate APScheduler running DriveSyncWorker every 5s"

Task 14: DI Container 補齊新元件

Files:

  • Modify: di_containers/cloud_integration/cloud_integration_containers.py
  • # 補在 Phase 1 已有的 container 末尾
    from infra.cloud_integration.repository.drive_folder_mapping_repo_impl import DriveFolderMappingRepoImpl
    from infra.cloud_integration.repository.drive_sync_job_repo_impl import DriveSyncJobRepoImpl
    from infra.cloud_integration.google_drive.google_drive_api_client import GoogleDriveApiClient
    
    from domain.cloud_integration.service.drive_folder_mapping_domain_service import DriveFolderMappingDomainService
    from domain.cloud_integration.service.drive_sync_job_domain_service import DriveSyncJobDomainService
    from domain.cloud_integration.service.google_drive_token_manager import GoogleDriveTokenManager
    
    from app.cloud_integration.service.handlers.create_folder_handler import CreateFolderHandler
    from app.cloud_integration.service.handlers.rename_folder_handler import RenameFolderHandler
    from app.cloud_integration.service.handlers.init_project_folders_handler import InitProjectFoldersHandler
    from app.cloud_integration.service.project_tree_loader import ProjectTreeLoader
    from app.cloud_integration.service.drive_sync_worker import DriveSyncWorker
    
    # repos
    drive_folder_mapping_repo = providers.Factory(DriveFolderMappingRepoImpl, session=providers.Resource(get_session))
    drive_sync_job_repo = providers.Factory(DriveSyncJobRepoImpl, session=providers.Resource(get_session))
    
    # domain services
    drive_folder_mapping_domain_service = providers.Factory(
        DriveFolderMappingDomainService, repo=drive_folder_mapping_repo,
    )
    drive_sync_job_domain_service = providers.Factory(
        DriveSyncJobDomainService, repo=drive_sync_job_repo,
    )
    
    # token manager
    google_drive_token_manager = providers.Singleton(
        GoogleDriveTokenManager,
        tenant_drive_integration_domain_service=tenant_drive_integration_domain_service,
        google_oauth_client=google_oauth_client,
        token_crypto_service=token_crypto_service,
    )
    
    # drive api client (token_provider 從 token manager 取)
    drive_api_client = providers.Singleton(
        GoogleDriveApiClient,
        token_provider=providers.Callable(
            lambda mgr: mgr.get_valid_access_token,
            mgr=google_drive_token_manager,
        ),
    )
    
    # project tree loader (注入既有 grc / oscal 等 domain services)
    # ── Cross-container DI 範例 ────────────────────────────────────────
    # 在外層 di_containers/containers.py 把對應的 container 注入:
    #     cloud_integration_container = providers.Container(
    #         CloudIntegrationContainer,
    #         config=config,
    #         auth_container=auth_container,
    #         grc_container=grc_container,             # ← 新增依賴
    #         oscal_container=oscal_container,         # ← 新增依賴
    #     )
    # 然後在 CloudIntegrationContainer 宣告:
    #     grc_container = providers.DependenciesContainer()
    #     oscal_container = providers.DependenciesContainer()
    # 引用時:
    project_tree_loader = providers.Factory(
        ProjectTreeLoader,
        grc_project_domain_service=grc_container.grc_project_domain_service,
        assessment_plan_domain_service=oscal_container.assessment_plan_domain_service,
        control_group_domain_service=grc_container.control_group_domain_service,
        control_domain_service=grc_container.control_domain_service,
        assessment_object_domain_service=grc_container.assessment_object_domain_service,
    )
    
    # > ⚠️ 上面的 attribute 名稱要對齊現有 grc_container / oscal_container 實際 provider 名。
    # > 用 `grep -n "domain_service" di_containers/grc/grc_containers.py` 確認名稱。
    # > 若對應 list method(list_by_project / list_by_ap / list_by_control_group / list_by_control)
    # > 在 domain service 內不存在,需要先在該 service + repo 補一個簡單的 list method
    # > (這是 ProjectTreeLoader 必要的前置子任務,列在 Task 11 Step 2 的 grep 步驟之後)。
    
    # handlers
    create_folder_handler = providers.Factory(
        CreateFolderHandler,
        folder_mapping_domain_service=drive_folder_mapping_domain_service,
        drive_api_client=drive_api_client,
    )
    rename_folder_handler = providers.Factory(
        RenameFolderHandler,
        folder_mapping_domain_service=drive_folder_mapping_domain_service,
        drive_api_client=drive_api_client,
    )
    init_project_folders_handler = providers.Factory(
        InitProjectFoldersHandler,
        folder_mapping_domain_service=drive_folder_mapping_domain_service,
        tenant_drive_integration_domain_service=tenant_drive_integration_domain_service,
        drive_api_client=drive_api_client,
        project_tree_loader=project_tree_loader,
    )
    
    # worker
    drive_sync_worker = providers.Singleton(
        DriveSyncWorker,
        drive_sync_job_domain_service=drive_sync_job_domain_service,
        handlers=providers.List(
            init_project_folders_handler,
            create_folder_handler,
            rename_folder_handler,
        ),
    )
  • git add di_containers/cloud_integration/cloud_integration_containers.py di_containers/containers.py
    git commit -m "feat(cloud_integration): wire folder/sync_job/handler/worker/api_client into DI"

Task 15: Integration with OscalProjectStartService

Files:

  • Modify: app/oscal/project/oscal_project_start_service.py (或 grep 找出實際路徑)
  • grep -rn "class OscalProjectStartService\|def start" app/oscal/ app/grc/ app/project/
  • # 在 service class 內加 helper(自帶獨立 transaction):
    @transaction
    def _enqueue_drive_folder_init(self, tenant_id: int, project_uid) -> None:
        tenant_int = self._tenant_drive_integration_domain_service.get_or_none(tenant_id)
        if not tenant_int or tenant_int.status != "CONNECTED":
            return
        self._drive_sync_job_domain_service.enqueue(
            tenant_id=tenant_id,
            job_type="INIT_PROJECT_FOLDERS",
            payload={"project_uid": str(project_uid)},
            priority=50,
        )
    
    # 在 start() 結尾呼叫:
    result = ...原本的 21 張 table 寫入流程結果...
    try:
        self._enqueue_drive_folder_init(tenant_id, result.project_uid)
    except Exception as e:
        logger.warning("Failed to enqueue Drive folder init for project %s: %s",
                       result.project_uid, e)
    return result

    Spec §9.1:enqueue 在 start transaction commit 後跑、自帶獨立 @transaction、失敗只 log 不 raise。 不要用 no_autoflush 或在原 transaction 內塞額外操作 — 那會把兩個邏輯耦合。

    • 連 Drive
    • 啟動一個 project
    • drive_sync_jobs 表有一筆 INIT_PROJECT_FOLDERS
    • 等 5-10 秒,看 Drive 上出現完整資料夾結構
    • drive_folder_mappings 表有對應 mappings
  • git add app/oscal/project/ di_containers/
    git commit -m "feat(cloud_integration): enqueue INIT_PROJECT_FOLDERS on project start (best-effort)"

Task 16: Rename Hooks 整合

Files:

  • Modify: project / AP / task 各自的 update name service(CG / Control / AO 無 rename API,不需要 hook

對每個 entity 改名 service 重複:

v0.2 Addendum (revised)

  1. Task rename hook — 已確認 endpoint:PUT /grc/project/<pid>/job/<job_uid>JobService.update_job(in app/grc/service/job_service.py,line ~123)。
    • 比照 PROJECT / AP 的 pattern:service 層偵測 name 變動 → return (dto, name_changed) → route (api/grc/routes/job_route.py ProjectJobDetailResource.put) 呼叫 try_enqueue_rename_folder('TASK', task_uid, name_for_drive)
    • name_for_drive = f"[T-{str(task_uid)[:8]}] {new_name}"
  2. CG / Control / AO 不需要 rename hook — 這三層名稱由 OSCAL profile import 決定,對 user 為 immutable(系統內無 rename API endpoint),故 spec / plan 明確不為它們補 hook
  3. AP rename 時 name 組合要改用 f"第{round_no}輪 - {new_name}",否則 sync 過去會失去前綴
  • grep -rn "def.*update.*name\|def.*rename" app/grc/ app/oscal/ app/project/ app/associations/

    已確認位置(v0.2 補充):

    • Project renameapp/grc/service/project_service.py (Phase 2 已實作)
    • AP renameapp/oscal/...(Phase 2 已實作)
    • Task renameapp/grc/service/job_service.py JobService.update_job (line ~123);endpoint PUT /grc/project/<pid>/job/<job_uid> (api/grc/routes/job_route.py ProjectJobDetailResource.put)
    • CG / Control / AO → 無 rename API(OSCAL profile 為 source of truth),不需 hook
    • 若該 entity 沒有專門的 rename method(只有泛用 update),改在 update method 內偵測 name 欄位是否變動,變動才 enqueue RENAME_FOLDER
  • # Best-effort sync to Drive
    try:
        tenant_int = self._tenant_drive_integration_domain_service.get_or_none(tenant_id)
        if tenant_int and tenant_int.status == "CONNECTED":
            name_for_drive = (
                new_name if scope_type in ("PROJECT", "AP")
                else f"[{code}] {new_name}"
            )
            self._drive_sync_job_domain_service.enqueue(
                tenant_id=tenant_id,
                job_type="RENAME_FOLDER",
                payload={
                    "scope_type": scope_type,   # "PROJECT" | "AP" | "CONTROL_GROUP" | "CONTROL" | "AO"
                    "scope_uid": str(entity.uid),
                    "new_name": name_for_drive,
                },
            )
    except Exception:
        logger.warning("Failed to enqueue RENAME_FOLDER for %s/%s", scope_type, entity.uid)
  • git commit -m "feat(cloud_integration): enqueue RENAME_FOLDER on <entity> rename"
    • 改一個 project 名稱 → 等 worker 跑 → Drive 上對應 project 資料夾改名
    • 改一個 task 名稱 → 等 worker 跑 → Drive 上對應 task 資料夾改名(含 [T-xxxxxxxx] prefix)

Task 17: API — Sync Trigger / Sync Jobs / Manual Init

Files:

  • Create: api/cloud_integration/routes/google_drive_sync_route.py
  • Create: api/cloud_integration/serializers/drive_sync_job.py
  • Modify: api/cloud_integration/__init__.py (register new routes)
  • from marshmallow import Schema, fields
    
    class DriveSyncJobResponse(Schema):
        uid = fields.UUID()
        job_type = fields.String()
        status = fields.String()
        retry_count = fields.Integer()
        max_retries = fields.Integer()
        next_run_at = fields.DateTime(allow_none=True)
        last_error = fields.String(allow_none=True)
        created_at = fields.DateTime()
        finished_at = fields.DateTime(allow_none=True)
        payload = fields.Dict()
  • # api/cloud_integration/routes/google_drive_sync_route.py
    # 4 個 Resource:
    # - POST /api/integrations/google-drive/sync             — manual trigger PROCESS_DRIVE_CHANGES (Phase 3 才用,先 stub return 501 or 直接 enqueue)
    # - GET  /api/integrations/google-drive/sync-jobs        — list (admin)
    # - POST /api/integrations/google-drive/sync-jobs/<uid>/retry  — manual retry (admin)
    # - POST /api/integrations/google-drive/projects/<project_uid>/init-folders  — re-init (admin)
    
    # 每個 Resource 都要 jwt_required + admin check
    # 細節參考 Phase 1 Task 13 pattern

    詳細 code 略;模式同 Phase 1 Task 13。注意:sync 觸發 (POST /sync) 在 Phase 2 可以直接 enqueue('INIT_PROJECT_FOLDERS', ...) 或 noop(Phase 3 改為 enqueue PROCESS_DRIVE_CHANGES)。init-folders route 直接 enqueue INIT_PROJECT_FOLDERS job。

  • git add api/cloud_integration/routes/google_drive_sync_route.py api/cloud_integration/serializers/drive_sync_job.py api/cloud_integration/__init__.py
    git commit -m "feat(cloud_integration): add sync-jobs admin API + manual init-folders endpoint"

Task 18: 前端 — 整合設定頁加 root folder 連結 + 失敗 jobs dialog

Files:

  • Modify: src/components/integrations/GoogleDriveIntegrationCard.vue
  • Create: src/components/integrations/SyncJobHistoryDialog.vue
  • Modify: src/service/CloudIntegrationService.js
  • Modify: src/config/api/api.js
  • INTEGRATION_GDRIVE_SYNC: '/integrations/google-drive/sync',
    INTEGRATION_GDRIVE_SYNC_JOBS: '/integrations/google-drive/sync-jobs',
  • async listSyncJobs(params) { return this.get(API.INTEGRATION_GDRIVE_SYNC_JOBS, params) }
    async retrySyncJob(uid) { return this.post(`${API.INTEGRATION_GDRIVE_SYNC_JOBS}/${uid}/retry`) }
    async triggerSync() { return this.post(API.INTEGRATION_GDRIVE_SYNC) }
    • root folder 連結已從 status root_folder_id 組出(Phase 1 已實作)
    • 加「立即同步」按鈕(call triggerSync
    • 加「查看失敗紀錄」按鈕 → 開 SyncJobHistoryDialog
  • git add src/
    git commit -m "feat(cloud-integration): add sync-jobs history dialog + manual sync button"

Task 19: Playwright 驗收(資料夾自動建立 + 改名同步)

  • GuidantAI/
      └── <Project Name>/
          └── 第N輪 - <AP Name>/
              └── [CG-XX] <Group>/
                  └── [CTRL-XX] <Control>/
                      └── <AO Title>/
                          └── [T-xxxxxxxx] <Task Name>/

Task 20: Channel Renewer Placeholder(為 Phase 3 預留)

Files:

  • Modify: core/scheduler.py
  • def _renew_channels_job():
        # Phase 3 才實作 actual renewal
        logger.debug("WebhookChannelRenewer tick (placeholder, no-op until Phase 3)")
    
    _scheduler.add_job(
        _renew_channels_job, "interval", hours=24,
        id="webhook_channel_renewer", max_instances=1,
    )
  • git add core/scheduler.py
    git commit -m "chore(cloud_integration): add webhook channel renewer placeholder job"

Task 20.5 (v0.2 Addendum): Task Creation Hook

Files:

  • Modify: app/grc/service/job_service.pyJobService.create_job (line ~43)
  • Modify: api/grc/routes/job_route.pyProjectJobCreateResource.post (line ~129)

v0.2 為了讓「AO init 完成後再加的 task」能自動補出 Drive folder,需在 task 建立流程加 hook。 已確認 endpointPOST /grc/project/<pid>/assessment-object/<ao_uid>/jobsJobService.create_job 觸發點:BPMN editor 新增 Task 節點後 save / planning 頁面新增 task。

  • # app/cloud_integration/service/drive_sync_orchestration_service.py
    def try_enqueue_create_task_folder(self, tenant_id, ao_uid, task_uid, task_name):
        try:
            tenant_int = self._tenant_drive_integration_domain_service.get_or_none(tenant_id)
            if not tenant_int or tenant_int.status != "CONNECTED":
                return
            ao_mapping = self._drive_folder_mapping_domain_service.get_by_scope(
                tenant_id=tenant_id, scope_type="AO", scope_uid=ao_uid,
            )
            if ao_mapping is None or ao_mapping.is_unlinked:
                return  # AO mapping 還沒建(project init 失敗)→ 跳過
            self._drive_sync_job_domain_service.enqueue(
                tenant_id=tenant_id,
                job_type="CREATE_FOLDER",
                payload={
                    "scope_type": "TASK",
                    "scope_uid": str(task_uid),
                    "parent_drive_folder_id": ao_mapping.drive_folder_id,
                    "name": f"[T-{str(task_uid)[:8]}] {task_name or ''}".strip(),
                },
            )
        except Exception:
            logger.warning("Failed to enqueue CREATE_FOLDER for task %s", task_uid)

Task 20.6 (v0.2 Addendum): Task Deletion Hook

Files:

  • Modify: app/grc/service/job_service.pyJobService.delete_job (line ~91)
  • Modify: app/cloud_integration/service/drive_sync_orchestration_service.py

已確認 endpointDELETE /grc/project/<pid>/job/<job_uid>JobService.delete_job 觸發點:BPMN editor 移除 Task 節點後 save / planning 頁面刪除 task。

v0.3 Addendum (重要:取代 v0.2 設計):v0.2 設計為「不刪 Drive folder + 標 mapping unlinked」,v0.3 改為「把 task folder move 到 AP 的 _Archive/ + 標 mapping unlinked」。實際實作以 v0.3 為準(請勿再做 v0.2 版本)。

v0.3 邏輯改在 Phase 3 Task 14 旁邊新增的「Archive APIs」段落定義 — 這裡 Task 20.6 標題保留,但 step 內容改為呼叫新的 try_archive_task_folder(詳見 Phase 3 plan「Archive APIs」段)。

設計決策(v0.3)

  • 把 task folder move 到 _Archive/(呼叫 Drive API),不再「保持原位」
  • 直接 DB write 標 mapping unlinked + 觸發 best-effort Drive move:mapping unlinked 走原本一筆 SQL UPDATE;Drive move 透過 orchestration 同步呼叫 Drive API(best-effort,失敗只 log,不 enqueue retry job — archive 失敗的 evidence 仍可在 Drive 原 task folder 由 user 手動處理)
  • 後續 webhook 收到該 folder 內變更時走 ancestor=ARCHIVE → SKIP(v0.3 §8.2 step 0)
  • # app/cloud_integration/service/drive_sync_orchestration_service.py
    def try_archive_task_folder(self, tenant_id, task_uid):
        """Best-effort: 把 task folder 搬到對應 AP 的 _Archive/ 並標 mapping unlinked。
        Drive 失敗只 log,不阻擋系統內軟刪。"""
        try:
            tenant_int = self._tenant_drive_integration_domain_service.get_or_none(tenant_id)
            if not tenant_int or tenant_int.status != "CONNECTED":
                return
            task_mapping = self._drive_folder_mapping_domain_service.get_by_scope(
                tenant_id=tenant_id, scope_type="TASK", scope_uid=task_uid,
            )
            if task_mapping is None:
                return
            # 反推 AP uid(透過 task → job_execution → AP)
            ap_uid = self._resolve_ap_uid_for_task(task_uid)   # helper, 可注入 JobExecutionDomainService 等
            archive_mapping = self._ensure_archive_folder(tenant_id, ap_uid)  # 不存在則重建
            # 把 task folder move 進 archive
            self._drive_api_client.move_to_archive(
                tenant_id=tenant_id,
                file_id=task_mapping.drive_folder_id,
                archive_folder_id=archive_mapping.drive_folder_id,
                current_parent_id=task_mapping.parent_drive_folder_id,
            )
            # 標 task mapping unlinked
            self._drive_folder_mapping_domain_service.mark_unlinked_by_scope(
                tenant_id=tenant_id, scope_type="TASK", scope_uid=task_uid,
            )
        except Exception:
            logger.warning("Failed to archive task folder for task %s", task_uid)
    • 系統內刪一個有 evidence 的 task → 確認 task folder 出現在 Drive 上對應 AP 的 _Archive/
    • 確認 task mapping is_unlinked=TRUE
    • 確認 webhook 對該 archived folder 內變更走 SKIP(不重新 import)

Task 21: Changelog

Files:

  • Create: docs/changelog/<YYYY-MM-DD>-google-drive-folder-init-and-rename-sync.md

§3

Phase 2 驗收 Checklist

完成 → 進 Phase 3。