第155回 AIの提案を直接反映しない:Pythonで作るドライランとPlan/Apply型CSV更新

AIに顧客情報や申請データを分類させる場合、結果をそのまま元データへ書き込むと、誤分類や入力の取り違えを後から確認しにくくなります。そこで今回は、AIや担当者が作った更新案をいったん「変更計画」に変換し、人が確認してから別のCSVへ適用する仕組みを作ります。

前回扱った増分処理と冪等性が「同じ処理を安全に再実行する設計」だとすれば、今回の主題は「提案と適用の間に確認可能な境界を置く設計」です。生成AIのAPI自体は呼び出さないため、APIキーなしで再現できます。

今回作るワークフロー

  1. 現在のデータをrecords.csvに用意する。
  2. AIまたは担当者の更新案をsuggestions.csvとして受け取る。
  3. planコマンドで入力を検証し、plan.jsonを作る。この段階では元データを変更しない。
  4. 担当者が変更内容を確認し、計画ファイルのSHA-256ハッシュを取得する。
  5. applyコマンドでハッシュと元データの一致を再確認し、更新結果を新しいCSVへ出力する。

適用時に元CSVが計画作成時と一致しない場合、または計画ファイルが承認ハッシュ取得時と一致しない場合、適用処理は失敗します。元CSVは上書きせず、出力先がすでに存在する場合も停止します。

前提と対象範囲

  • Python 3.10以上を使用します。外部ライブラリは不要です。
  • 例ではpendingapprovedrejectedの3状態だけを許可します。
  • 入力をメモリへ読み込むため、担当者がレビューできる規模のCSVを対象とします。
  • 同じ出力先へ複数のプロセスが同時に書き込まない運用を前提とします。

入力CSVを用意する

records.csvは現在の正本です。

record_id,name,status
A001,青木商店,pending
A002,北川企画,pending
A003,南製作所,approved

suggestions.csvはAIや担当者が作った提案です。理由も必須にして、レビュー時に判断根拠を確認できるようにします。

record_id,proposed_status,reason
A001,approved,必要項目を確認済み
A002,rejected,連絡先が不足
A003,approved,現状維持

この例ではA001とA002が変更対象です。A003は現在値と提案値が同じなので、計画には変更操作として追加しません。

Plan/Applyを実装する

次のコードをworkflow.pyとして保存します。CSVの重複ID、未知のID、列不足、不正な状態、空の理由を計画作成前に検出します。

from __future__ import annotations

import argparse
import csv
import hashlib
import io
import json
import os
import sys
import tempfile
from pathlib import Path
from typing import Any

ALLOWED_STATUSES = {'pending', 'approved', 'rejected'}
PLAN_VERSION = 1


class WorkflowError(Exception):
    pass


def sha256_bytes(data: bytes) -> str:
    return hashlib.sha256(data).hexdigest()


def parse_csv(
    data: bytes,
    label: str,
    required: set[str],
) -> tuple[list[str], list[dict[str, str]]]:
    try:
        text = data.decode('utf-8-sig')
    except UnicodeDecodeError as exc:
        raise WorkflowError(f'{label}: UTF-8として読み込めません') from exc

    reader = csv.DictReader(io.StringIO(text, newline=''), strict=True)
    if not reader.fieldnames:
        raise WorkflowError(f'{label}: ヘッダーがありません')

    fields = list(reader.fieldnames)
    if any(name == '' for name in fields):
        raise WorkflowError(f'{label}: 空の列名があります')
    if len(fields) != len(set(fields)):
        raise WorkflowError(f'{label}: 列名が重複しています')

    missing = sorted(required - set(fields))
    if missing:
        missing_text = ', '.join(missing)
        raise WorkflowError(f'{label}: 必須列がありません: {missing_text}')

    rows: list[dict[str, str]] = []
    for data_no, row in enumerate(reader, start=1):
        if None in row or any(value is None for value in row.values()):
            raise WorkflowError(f'{label}: データ行{data_no}の列数が不正です')
        rows.append(dict(row))
    return fields, rows


def load_records(
    data: bytes,
    label: str,
) -> tuple[list[str], list[dict[str, str]], dict[str, dict[str, str]]]:
    fields, rows = parse_csv(data, label, {'record_id', 'status'})
    index: dict[str, dict[str, str]] = {}

    for data_no, row in enumerate(rows, start=1):
        record_id = row['record_id']
        status = row['status']

        if not record_id or record_id != record_id.strip():
            raise WorkflowError(f'{label}: データ行{data_no}のrecord_idが不正です')
        if record_id in index:
            raise WorkflowError(f'{label}: record_idが重複しています: {record_id}')
        if status not in ALLOWED_STATUSES:
            raise WorkflowError(f'{label}: {record_id}のstatusが不正です: {status}')

        index[record_id] = row
    return fields, rows, index


def atomic_write_text(
    path: Path,
    text: str,
    replace: bool = False,
) -> None:
    if not path.parent.is_dir():
        raise WorkflowError(f'出力先ディレクトリがありません: {path.parent}')
    if path.exists() and not replace:
        raise WorkflowError(f'出力先がすでに存在します: {path}')

    temp_path: Path | None = None
    try:
        with tempfile.NamedTemporaryFile(
            mode='w',
            encoding='utf-8',
            newline='',
            dir=path.parent,
            prefix=f'.{path.name}.',
            suffix='.tmp',
            delete=False,
        ) as temporary:
            temp_path = Path(temporary.name)
            temporary.write(text)
            temporary.flush()
            os.fsync(temporary.fileno())

        if path.exists() and not replace:
            raise WorkflowError(f'出力先がすでに存在します: {path}')
        os.replace(temp_path, path)
    except BaseException:
        if temp_path is not None:
            try:
                temp_path.unlink(missing_ok=True)
            except OSError:
                pass
        raise


def make_plan(
    records_path: Path,
    suggestions_path: Path,
    plan_path: Path,
    replace: bool,
) -> dict[str, Any]:
    if plan_path.resolve() in {
        records_path.resolve(),
        suggestions_path.resolve(),
    }:
        raise WorkflowError('計画ファイルを入力CSVと同じ場所にはできません')

    records_bytes = records_path.read_bytes()
    suggestions_bytes = suggestions_path.read_bytes()
    _, record_rows, records = load_records(records_bytes, str(records_path))
    _, suggestions = parse_csv(
        suggestions_bytes,
        str(suggestions_path),
        {'record_id', 'proposed_status', 'reason'},
    )

    seen: set[str] = set()
    operations: list[dict[str, str]] = []
    unchanged = 0

    for data_no, suggestion in enumerate(suggestions, start=1):
        record_id = suggestion['record_id']
        proposed = suggestion['proposed_status']
        reason = suggestion['reason'].strip()

        if not record_id or record_id != record_id.strip():
            raise WorkflowError(
                f'{suggestions_path}: データ行{data_no}のrecord_idが不正です'
            )
        if record_id in seen:
            raise WorkflowError(f'提案のrecord_idが重複しています: {record_id}')
        seen.add(record_id)

        if record_id not in records:
            raise WorkflowError(f'元データに存在しないrecord_idです: {record_id}')
        if proposed not in ALLOWED_STATUSES:
            raise WorkflowError(f'{record_id}の提案状態が不正です: {proposed}')
        if not reason:
            raise WorkflowError(f'{record_id}のreasonが空です')

        current = records[record_id]['status']
        if current == proposed:
            unchanged += 1
            continue

        operations.append(
            {
                'record_id': record_id,
                'before': current,
                'after': proposed,
                'reason': reason,
            }
        )

    operations.sort(key=lambda operation: operation['record_id'])
    plan: dict[str, Any] = {
        'format_version': PLAN_VERSION,
        'records_sha256': sha256_bytes(records_bytes),
        'suggestions_sha256': sha256_bytes(suggestions_bytes),
        'record_count': len(record_rows),
        'suggestion_count': len(suggestions),
        'unchanged_count': unchanged,
        'operations': operations,
    }

    plan_text = json.dumps(
        plan,
        ensure_ascii=False,
        sort_keys=True,
        indent=2,
    ) + '\n'
    atomic_write_text(plan_path, plan_text, replace=replace)
    return plan


def load_plan(data: bytes, label: str) -> dict[str, Any]:
    try:
        parsed = json.loads(data.decode('utf-8-sig'))
    except (UnicodeDecodeError, json.JSONDecodeError) as exc:
        raise WorkflowError(f'{label}: 有効なJSONではありません') from exc

    if not isinstance(parsed, dict):
        raise WorkflowError(f'{label}: JSONオブジェクトではありません')
    if type(parsed.get('format_version')) is not int:
        raise WorkflowError(f'{label}: format_versionが不正です')
    if parsed['format_version'] != PLAN_VERSION:
        raise WorkflowError(f'{label}: 未対応のformat_versionです')

    source_hash = parsed.get('records_sha256')
    if (
        not isinstance(source_hash, str)
        or len(source_hash) != 64
        or any(character not in '0123456789abcdef' for character in source_hash.lower())
    ):
        raise WorkflowError(f'{label}: records_sha256が不正です')
    if not isinstance(parsed.get('operations'), list):
        raise WorkflowError(f'{label}: operationsが配列ではありません')
    return parsed


def apply_plan(
    records_path: Path,
    plan_path: Path,
    output_path: Path,
    approved_hash: str,
) -> dict[str, Any]:
    if output_path.resolve() in {records_path.resolve(), plan_path.resolve()}:
        raise WorkflowError('出力CSVを入力CSVまたは計画ファイルと同じ場所にはできません')

    approved = approved_hash.strip().lower()
    if (
        len(approved) != 64
        or any(character not in '0123456789abcdef' for character in approved)
    ):
        raise WorkflowError('承認ハッシュは64桁の16進数で指定してください')

    plan_bytes = plan_path.read_bytes()
    actual_plan_hash = sha256_bytes(plan_bytes)
    if actual_plan_hash != approved:
        raise WorkflowError('計画ファイルが承認時の内容と一致しません')
    plan = load_plan(plan_bytes, str(plan_path))

    records_bytes = records_path.read_bytes()
    if sha256_bytes(records_bytes) != plan['records_sha256'].lower():
        raise WorkflowError(
            '入力CSVが計画作成後に変化しています。計画を作り直してください'
        )

    fields, rows, records = load_records(records_bytes, str(records_path))
    seen: set[str] = set()

    for operation_no, operation in enumerate(plan['operations'], start=1):
        if not isinstance(operation, dict):
            raise WorkflowError(f'操作{operation_no}がオブジェクトではありません')

        required = {'record_id', 'before', 'after', 'reason'}
        if not required.issubset(operation):
            raise WorkflowError(f'操作{operation_no}の項目が不足しています')
        if not all(isinstance(operation[key], str) for key in required):
            raise WorkflowError(f'操作{operation_no}に文字列以外の値があります')

        record_id = operation['record_id']
        before = operation['before']
        after = operation['after']

        if record_id in seen:
            raise WorkflowError(f'計画内のrecord_idが重複しています: {record_id}')
        seen.add(record_id)
        if record_id not in records:
            raise WorkflowError(f'計画対象が元データにありません: {record_id}')
        if before not in ALLOWED_STATUSES or after not in ALLOWED_STATUSES:
            raise WorkflowError(f'{record_id}の状態が不正です')
        if before == after:
            raise WorkflowError(f'{record_id}は変更前後が同じです')
        if not operation['reason'].strip():
            raise WorkflowError(f'{record_id}のreasonが空です')
        if records[record_id]['status'] != before:
            raise WorkflowError(f'{record_id}の変更前状態が一致しません')

        records[record_id]['status'] = after

    buffer = io.StringIO(newline='')
    writer = csv.DictWriter(buffer, fieldnames=fields, lineterminator='\n')
    writer.writeheader()
    writer.writerows(rows)
    atomic_write_text(output_path, buffer.getvalue())

    return {
        'applied': len(plan['operations']),
        'output': str(output_path),
        'plan_sha256': actual_plan_hash,
    }


def build_parser() -> argparse.ArgumentParser:
    parser = argparse.ArgumentParser(
        description='CSVの状態変更を計画し、承認後に別ファイルへ適用します'
    )
    commands = parser.add_subparsers(dest='command', required=True)

    plan_command = commands.add_parser('plan', help='変更計画を作成します')
    plan_command.add_argument('--records', type=Path, required=True)
    plan_command.add_argument('--suggestions', type=Path, required=True)
    plan_command.add_argument('--plan', type=Path, required=True)
    plan_command.add_argument('--replace-plan', action='store_true')

    apply_command = commands.add_parser('apply', help='承認済み計画を適用します')
    apply_command.add_argument('--records', type=Path, required=True)
    apply_command.add_argument('--plan', type=Path, required=True)
    apply_command.add_argument('--output', type=Path, required=True)
    apply_command.add_argument('--approve-hash', required=True)
    return parser


def main(argv: list[str] | None = None) -> int:
    args = build_parser().parse_args(argv)
    try:
        if args.command == 'plan':
            plan = make_plan(
                args.records,
                args.suggestions,
                args.plan,
                replace=args.replace_plan,
            )
            result = {
                'mode': 'plan',
                'plan': str(args.plan),
                'plan_sha256': sha256_bytes(args.plan.read_bytes()),
                'changes': len(plan['operations']),
                'unchanged': plan['unchanged_count'],
            }
        else:
            result = {
                'mode': 'apply',
                **apply_plan(
                    args.records,
                    args.plan,
                    args.output,
                    args.approve_hash,
                ),
            }
    except (WorkflowError, OSError, csv.Error) as exc:
        print(f'error: {exc}', file=sys.stderr)
        return 1

    print(json.dumps(result, ensure_ascii=False, sort_keys=True))
    return 0


if __name__ == '__main__':
    raise SystemExit(main())

実行順と入出力

1. 変更計画を作る

python workflow.py plan --records records.csv --suggestions suggestions.csv --plan plan.json

標準出力には、変更件数、変更なしの件数、計画ファイルのSHA-256ハッシュがJSONで1行出力されます。この例では変更2件、変更なし1件です。元のrecords.csvは変更されません。

計画内容は次のコマンドで読みやすく表示できます。

python -m json.tool plan.json

operationsにあるrecord_idbeforeafterreasonを確認します。既存の計画を作り直す場合だけ--replace-planを指定します。置き換えると以前のハッシュは一致しなくなるため、再確認が必要です。

2. 確認した計画のハッシュを取得する

レビューを終えた後、次のコマンドで計画ファイルの現在のハッシュを表示します。

python -c "from pathlib import Path; import hashlib; print(hashlib.sha256(Path('plan.json').read_bytes()).hexdigest())"

表示された64桁の値を、確認済み計画の識別子として控えます。ハッシュを取得した後に計画ファイルを編集した場合は、内容を再確認してハッシュも取り直します。

3. 計画を適用する

次のHASHを、直前に取得した実際の値へ置き換えて実行します。

python workflow.py apply --records records.csv --plan plan.json --output updated_records.csv --approve-hash HASH

成功時は適用件数、出力先、計画ハッシュが標準出力へJSONで出ます。生成されるupdated_records.csvの論理的な内容は次のとおりです。

record_id,name,status
A001,青木商店,approved
A002,北川企画,rejected
A003,南製作所,approved

元のrecords.csvは残ります。CSVは再生成されるため、引用符、BOM、改行コードなどのバイト表現は元ファイルと異なる場合があります。

失敗時の動作

  • 利用方法の誤りはargparseが標準エラーへ案内を出し、終了コード2で終了します。
  • 入力不備、ハッシュ不一致、既存の出力先など、想定内の業務エラーは標準エラーへ理由を出し、終了コード1で停止します。
  • 全操作の検証が終わるまで出力先へ切り替えないため、途中の操作だけを反映したCSVは残りません。
  • 一時ファイルは出力先と同じディレクトリへ作り、書き込み完了後にos.replaceで出力先へ切り替えます。

たとえば計画作成後にrecords.csvを編集すると、適用時に「入力CSVが計画作成後に変化しています」と表示され、出力CSVは作成されません。改行コードやBOMだけの変更もハッシュ差として検出します。これは、論理的に同じに見える別ファイルを誤って適用しないための保守的な動作です。

最小テストを追加する

次をtest_workflow.pyとして保存します。正常適用、計画作成後の入力変更、承認後の計画変更を検査します。

import csv
import tempfile
import unittest
from pathlib import Path

from workflow import WorkflowError, apply_plan, make_plan, sha256_bytes


RECORDS = '''record_id,name,status
A001,青木商店,pending
A002,北川企画,pending
A003,南製作所,approved
'''

SUGGESTIONS = '''record_id,proposed_status,reason
A001,approved,必要項目を確認済み
A002,rejected,連絡先が不足
A003,approved,現状維持
'''


class WorkflowTest(unittest.TestCase):
    def prepare(self, root: Path) -> tuple[Path, Path, Path]:
        records = root / 'records.csv'
        suggestions = root / 'suggestions.csv'
        plan = root / 'plan.json'
        records.write_text(RECORDS, encoding='utf-8')
        suggestions.write_text(SUGGESTIONS, encoding='utf-8')
        return records, suggestions, plan

    def test_apply_changes_only_planned_rows(self) -> None:
        with tempfile.TemporaryDirectory() as directory:
            root = Path(directory)
            records, suggestions, plan = self.prepare(root)
            output = root / 'updated.csv'

            plan_data = make_plan(records, suggestions, plan, replace=False)
            approval = sha256_bytes(plan.read_bytes())
            result = apply_plan(records, plan, output, approval)

            self.assertEqual(len(plan_data['operations']), 2)
            self.assertEqual(result['applied'], 2)
            self.assertEqual(
                sha256_bytes(records.read_bytes()),
                plan_data['records_sha256'],
            )

            with output.open(encoding='utf-8', newline='') as file:
                rows = list(csv.DictReader(file))
            statuses = {row['record_id']: row['status'] for row in rows}
            self.assertEqual(
                statuses,
                {
                    'A001': 'approved',
                    'A002': 'rejected',
                    'A003': 'approved',
                },
            )

    def test_rejects_changed_input(self) -> None:
        with tempfile.TemporaryDirectory() as directory:
            root = Path(directory)
            records, suggestions, plan = self.prepare(root)
            output = root / 'updated.csv'

            make_plan(records, suggestions, plan, replace=False)
            approval = sha256_bytes(plan.read_bytes())
            changed = RECORDS.replace(
                'A001,青木商店,pending',
                'A001,青木商店,approved',
            )
            records.write_text(changed, encoding='utf-8')

            with self.assertRaisesRegex(
                WorkflowError,
                '入力CSVが計画作成後に変化',
            ):
                apply_plan(records, plan, output, approval)
            self.assertFalse(output.exists())

    def test_rejects_changed_plan(self) -> None:
        with tempfile.TemporaryDirectory() as directory:
            root = Path(directory)
            records, suggestions, plan = self.prepare(root)
            output = root / 'updated.csv'

            make_plan(records, suggestions, plan, replace=False)
            approval = sha256_bytes(plan.read_bytes())
            plan.write_text(
                plan.read_text(encoding='utf-8') + '\n',
                encoding='utf-8',
            )

            with self.assertRaisesRegex(
                WorkflowError,
                '計画ファイルが承認時の内容と一致しません',
            ):
                apply_plan(records, plan, output, approval)
            self.assertFalse(output.exists())


if __name__ == '__main__':
    unittest.main()

テストは次のコマンドで実行します。

python -m unittest -v test_workflow

3つのテストがokになれば、サンプルの正常系、元データ変更の拒否、計画ファイル変更の拒否を確認できます。

この実装で保証しないこと

SHA-256ハッシュはファイル内容の一致確認には使えますが、誰が承認したかを証明するものではありません。ハッシュを渡すだけで承認者の認証や電子署名になるわけでもありません。チーム運用では、計画ファイルとハッシュをアクセス制御されたレビューやCIの承認手順へ渡し、提案者、確認者、適用日時を別途記録します。

また、この例は同時実行ロック、大容量CSVのストリーミング、状態遷移ごとの権限制御を実装していません。複数プロセスが同じ出力先を扱う場合は排他制御を追加し、大容量データでは変更対象をデータベースや不変スナップショットへ移す設計を検討してください。表計算ソフトで開くCSVに外部由来の文字列を含める場合は、数式として解釈される文字列への対策も別途必要です。

導入時のチェックリスト

  • AIの出力を正本へ直接書き込まず、提案ファイルとして分離する。
  • 許可する値、必須列、重複ID、未知のIDを機械的に検証する。
  • 変更前後の値と理由を計画ファイルで確認する。
  • 確認後に取得した計画ファイルのハッシュを適用時に照合する。
  • 元データを上書きせず、出力結果を比較してから次工程へ渡す。

最初は実際の業務CSVを少量だけ複製し、許可状態と必須列を自社のルールへ置き換えて試してください。AIの判断精度とは別に、誤った提案を適用前に止められる工程を持つことが重要です。