Skip to content

[Job] import・download・batch・source syncを専用tableへ分離してjobsを廃止する #621

Description

@hmjn023

Parent: #612

Depends on: #619, #620

背景

generic jobs にはworker taskだけでなく、承認待ちの import_request、download、batch親子進捗、source syncが混在している。寿命・一意性・再試行・取消・履歴要件が異なるため、文字列typeとJSON payloadによる共通queueから用途別のrequest/run modelへ分離する。

スコープ

  • import_requests: 受付、承認、取消、対象source
  • download_runs: 入力、実行状態、試行、結果
  • batch_runs / 必要な場合の batch_items: task種別、件数、取消、再開、親状態再計算
  • source_sync_runs: source・sync種別・revision単位の直列実行
  • typed repository、明示的mapper、API/event/worker切替
  • active rowの冪等backfill、dual-read/write、rollback期間を含む段階移行
  • 全type移行後のstring dispatch、generic job repository/table削除

非スコープ

  • CCIP vector/crop schemaの実装
  • source dump/export用LanceDB機能を残すか廃止するかの意思決定
  • 無期限の実行履歴保存

受け入れ条件

  • import_request がgeneric queueに保存されない
  • 各flowが専用の状態遷移・一意制約・retry/取消規則を持つ
  • 移行中のactive request/runを失わない
  • 再起動・ページ再読込後にDBから状態を復元できる
  • atomic claim、AI concurrency、source直列化を必要なmodelで維持する
  • 最終的に jobs へのread/writeとstring-based dispatchがなくなる
  • rollback期間後に jobs tableを削除するmigrationがある

主な参照

  • apps/server/src/infrastructure/api/routers/imports-router.ts:98
  • apps/server/src/application/services/job-dispatch-service.ts:19
  • packages/db/src/repositories/job-repository.ts:280
  • packages/db/src/repositories/job-repository.ts:385
  • packages/db/src/schema.ts:601
  • packages/db/src/schema.ts:898

評価反映(2026-07-18)

本Issueは複数の独立flowと最終table削除を含むため、tracking issueとして扱い、実装前に以下のchild issueへ分割する。

  1. import requestの専用model/API移行
  2. download runの専用model/worker移行
  3. batch run/itemの専用model移行
  4. source sync run/dirty markerの専用model移行
  5. generic jobsの参照ゼロ監査、観測期間、最終削除

各child issueは独立してbackfill、cutover、rollback、検証できる単位とし、単一PRで全flowとjobs削除を同時に行わない。

全job typeの移行matrix

実装着手前に、現存する全typeについて次の列を持つ移行matrixを本Issueへ追加する。

項目 内容
現行type DBに存在し得る正確な文字列
producer 作成・再投入・mergeする全箇所
handler/claim pool worker、AI concurrency、source直列化
reader/updater API、status照会、parent進捗、取消、retry
event/UI schema、publisher、購読画面
dedupe/merge規則 論理key、active定義、supersede
移行先 processing state、import request、download run、batch run/item、source sync run等
cutover/rollback backfill、権威切替、観測期間、削除条件

少なくとも次のfamilyを漏れなく含める。

  • processMedia
  • auto_tagging
  • extract_ccip_vector
  • bulk taggingのparent/dispatch
  • batch CCIPのparent/dispatch
  • downloadImage
  • import_request
  • sync_lancedb / sync_lancedb_full / sync_lancedb_delta

コード検索だけでなく、DB上のSELECT DISTINCT type FROM jobsとも照合し、未知のlegacy typeを削除前に分類する。

batch_itemsを必須とする条件

次のいずれかを提供するbatchでは、batch_itemsを任意ではなく必須とする。

  • batch開始時点の対象集合を固定する
  • item単位の成功・失敗理由を表示する
  • 正確な進捗を算出する
  • 取消後に未処理itemだけを再開する
  • retryや履歴をitem単位で監査する

batch_runsの集計値はitem遷移と同一transactionで更新するか、itemsから再計算可能にする。処理状態tableの現在値だけからbatch履歴を推測しない。

段階移行中のclaim authority

dual-writeは保存先を二重化するだけであり、同じ論理処理を旧jobs workerと新workerの双方がclaimしてはならない。各flowのcutoverを次のphaseで管理する。

  1. 旧jobsのみclaim authority。新tableへbackfill/dual-writeし、shadow readと差分計測を行う。
  2. producer/workerをquiesceし、active row、parent/child、retry予約、失敗をreconcileする。
  3. deployment versionまたはDB上のcutover epochを切り替え、新tableのみをclaim authorityにする。旧workerはcutover後のrowをclaimできない。
  4. rollback時もquiesceと差分reconcileを行ってから権威を戻す。

rolling deploymentで旧・新workerが混在する場合のversion gateを必須とし、feature flagだけに依存した二重実行可能な構成を作らない。

source syncのwatermark/CAS

現行dirty markerは同一source/media rowを更新し、取得したrow IDを処理後に削除する。処理中に同じrowへ新しい変更が書かれると、古いworkerの削除が新しい要求まで消す可能性がある。

  • dirty markerへ単調増加するrevision/generationを持たせる。
  • workerはclaimしたrevisionを記録し、完了時はid + claimed_revisionのCASでのみ削除または完了更新する。
  • 処理中にrevisionが進んだrowは残し、次回sync対象にする。
  • source sync runへ対象watermark、完了watermark、claim tokenを記録する。
  • source単位の直列化をプロセス内SetだけでなくDB制約・claim条件でも維持する。
  • upsert/deleteの競合順序と、削除後の再作成をテストする。

主な確認箇所:

  • packages/db/src/schema.ts:598
  • apps/server/src/application/services/backup-service.ts:1144
  • apps/server/src/application/services/backup-service.ts:1205

realtime eventの責務と配置範囲

現行RealtimeEventBusはプロセス内EventEmitterで、job replayも最大100件・60秒のメモリ保持である。したがってイベントをdurable stateや完全な進捗履歴として扱わない。

初期運用をsingle server instanceへ限定する場合は、その制約を明記し、イベントをUI更新の通知としてのみ使う。再接続、ページ再読込、イベント欠落時は各専用tableをAPIから再取得して状態を復元する。

複数server instanceを要件に含める場合は、jobs削除前にtransactional outboxまたはdurable cross-process pub/subを別Issueで導入する。既存のtyped event schema、oRPC Event Iterator、共通client transportを維持し、flowごとの独自EventSourceを追加しない。

jobs table削除gate

DROP TABLE jobs migrationは、次の全条件を満たした後にのみ作成・適用する。

  • 移行matrixの全typeについてproducer、reader、updater、worker、API、event、testの移行先が埋まっている
  • DB上の全legacy typeを分類し、active rowが0件で、非終端parent/childをreconcile済み
  • processMedia由来のmetadata・thumbnailを含め、#620側の処理移行が完了している
  • import、download、batch、source syncそれぞれのshadow差分が0件
  • 新旧どちらか一方だけがclaim authorityであることを実DB concurrency testで確認している
  • 再起動、rolling deployment、cutover途中の異常終了、rollback rehearsalが通る
  • 明示した観測期間中、jobsへの新規read/writeが0件
  • jobs repository、string dispatch、直接SQL、schema export、API参照がコード検索とtypecheckで0件
  • UIがイベント再送に依存せず、DBから状態を復元できる
  • rollback期限と、削除前に必要な履歴archive/backupの扱いが文書化されている

jobs削除migrationは上記gateを満たす実装PRとは分離し、観測期間完了後の専用PRとする。

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Projects

    No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions