For agentic workers: REQUIRED SUB-SKILL: Use
superpowers:subagent-driven-development(recommended) orsuperpowers:executing-plansto 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_mappings 與 drive_sync_jobs 兩張新表 + APScheduler 起 worker poller + channel renewer placeholder(Phase 3 才會用到)+ GoogleDriveApiClient(folder/permission ops)。把 worker 的 7 種 job_type 中的 INIT_PROJECT_FOLDERS、CREATE_FOLDER、RENAME_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.sqlALTER 既有 CHECK constraint 加入'TASK'DriveScopeTypeenum 補TASKProjectTreeLoader載 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
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 常數
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
Files:
scripts/sql/<YYYY-MM-DD>-google-drive-folder-mappings.sqlv0.2 Addendum:原 migration(已執行)的
chk_drive_folder_mappings_scope不含'TASK'。Task #59 需新增獨立 migrationscripts/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"Files:
domain/cloud_integration/enums/__init__.pydomain/cloud_integration/enums/scope_type.pydomain/cloud_integration/enums/sync_job_type.pydomain/cloud_integration/enums/sync_job_status.pyinfra/cloud_integration/models/drive_folder_mapping.pyinfra/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"DriveFolderMappingFiles:
domain/cloud_integration/entities/drive_folder_mapping_entity.pydomain/cloud_integration/repository/drive_folder_mapping_repository.pyinfra/cloud_integration/mapper/drive_folder_mapping_mapper.pyinfra/cloud_integration/repository/drive_folder_mapping_repo_impl.pytest/cloud_integration/test_drive_folder_mapping_repo.pyfrom 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] = Nonefrom 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"DriveSyncJobFiles:
domain/cloud_integration/entities/drive_sync_job_entity.pydomain/cloud_integration/repository/drive_sync_job_repository.pyinfra/cloud_integration/mapper/drive_sync_job_mapper.pyinfra/cloud_integration/repository/drive_sync_job_repo_impl.pytest/cloud_integration/test_drive_sync_job_repo.pyfrom 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] = Nonefrom 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_failed 帶 retry_at → status=PENDING + retry_count+1mark_failed 不帶 → status=FAILEDmax_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)"Files:
domain/cloud_integration/service/drive_folder_mapping_domain_service.pydomain/cloud_integration/service/drive_sync_job_domain_service.pytest/cloud_integration/test_drive_folder_mapping_domain_service.pytest/cloud_integration/test_drive_sync_job_domain_service.pyfrom 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"Files:
domain/cloud_integration/service/google_drive_token_manager.pytest/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_tokeninvalid_grant → 標 REVOKED,raise DriveTokenRevokedErrorgit 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)"Files:
infra/cloud_integration/google_drive/google_drive_api_client.pytest/cloud_integration/test_google_drive_api_client.pyv0.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)"Files:
app/cloud_integration/service/handlers/__init__.pyapp/cloud_integration/service/handlers/base_job_handler.pyfrom 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"Files:
app/cloud_integration/service/handlers/create_folder_handler.pytest/cloud_integration/test_create_folder_handler.pyfrom 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)"Files:
app/cloud_integration/service/handlers/rename_folder_handler.pytest/cloud_integration/test_rename_folder_handler.pyfrom 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"Files:
app/cloud_integration/service/handlers/init_project_folders_handler.pytest/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 層時要:
- 在
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']}", )- AP 名稱改用
f"第{ap['round_no']}輪 - {ap['name']}"格式ProjectTreeLoader.load()多 query task — 每個 AO 拉compliance.job_executions找該 AO 在該 AP 週期下的所有 task。實作時需 grep 確認 task / job_execution 對應 query 介面(候選:ext_workflow_execution_domain_service內既有方法、或在JobExecutionDomainService補list_by_ao_uid_in_ap(ao_uid, ap_uid))- 確認 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_idNote:
project_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 補一個。
create_folder 呼叫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"Files:
app/cloud_integration/service/drive_sync_worker.pytest/cloud_integration/test_drive_sync_worker.pyimport 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)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)"Files:
core/scheduler.pycore/extensions.pycore/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 = Nonefrom 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"Files:
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"OscalProjectStartServiceFiles:
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 resultSpec §9.1:enqueue 在 start transaction commit 後跑、自帶獨立
@transaction、失敗只 log 不 raise。 不要用no_autoflush或在原 transaction 內塞額外操作 — 那會把兩個邏輯耦合。
drive_sync_jobs 表有一筆 INIT_PROJECT_FOLDERSdrive_folder_mappings 表有對應 mappingsgit add app/oscal/project/ di_containers/
git commit -m "feat(cloud_integration): enqueue INIT_PROJECT_FOLDERS on project start (best-effort)"Files:
對每個 entity 改名 service 重複:
v0.2 Addendum (revised):
- Task rename hook — 已確認 endpoint:
PUT /grc/project/<pid>/job/<job_uid>→JobService.update_job(inapp/grc/service/job_service.py,line ~123)。
- 比照 PROJECT / AP 的 pattern:service 層偵測 name 變動 → return
(dto, name_changed)→ route (api/grc/routes/job_route.pyProjectJobDetailResource.put) 呼叫try_enqueue_rename_folder('TASK', task_uid, name_for_drive)name_for_drive = f"[T-{str(task_uid)[:8]}] {new_name}"- CG / Control / AO 不需要 rename hook — 這三層名稱由 OSCAL profile import 決定,對 user 為 immutable(系統內無 rename API endpoint),故 spec / plan 明確不為它們補 hook
- 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 rename →
app/grc/service/project_service.py(Phase 2 已實作)- AP rename →
app/oscal/...(Phase 2 已實作)- Task rename →
app/grc/service/job_service.pyJobService.update_job(line ~123);endpointPUT /grc/project/<pid>/job/<job_uid>(api/grc/routes/job_route.pyProjectJobDetailResource.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"[T-xxxxxxxx] prefix)Files:
api/cloud_integration/routes/google_drive_sync_route.pyapi/cloud_integration/serializers/drive_sync_job.pyapi/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-foldersroute 直接 enqueueINIT_PROJECT_FOLDERSjob。
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"Files:
src/components/integrations/GoogleDriveIntegrationCard.vuesrc/components/integrations/SyncJobHistoryDialog.vuesrc/service/CloudIntegrationService.jssrc/config/api/api.jsINTEGRATION_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_id 組出(Phase 1 已實作)triggerSync)SyncJobHistoryDialoggit add src/
git commit -m "feat(cloud-integration): add sync-jobs history dialog + manual sync button"GuidantAI/
└── <Project Name>/
└── 第N輪 - <AP Name>/
└── [CG-XX] <Group>/
└── [CTRL-XX] <Control>/
└── <AO Title>/
└── [T-xxxxxxxx] <Task Name>/Files:
core/scheduler.pydef _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"Files:
app/grc/service/job_service.py — JobService.create_job (line ~43)api/grc/routes/job_route.py — ProjectJobCreateResource.post (line ~129)v0.2 為了讓「AO init 完成後再加的 task」能自動補出 Drive folder,需在 task 建立流程加 hook。 已確認 endpoint:
POST /grc/project/<pid>/assessment-object/<ao_uid>/jobs→JobService.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)Files:
app/grc/service/job_service.py — JobService.delete_job (line ~91)app/cloud_integration/service/drive_sync_orchestration_service.py已確認 endpoint:
DELETE /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):
_Archive/(呼叫 Drive API),不再「保持原位」# 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)_Archive/ 內is_unlinked=TRUEFiles:
docs/changelog/<YYYY-MM-DD>-google-drive-folder-init-and-rename-sync.md完成 → 進 Phase 3。