diff --git a/.github/workflows/release.yml b/.github/workflows/release.yml index de7009c..71a9e45 100644 --- a/.github/workflows/release.yml +++ b/.github/workflows/release.yml @@ -3,6 +3,10 @@ name: Release candidate on: workflow_dispatch: inputs: + release_id: + description: Unique Framework Release ID (published as framework-) + required: true + type: string source_id: description: Logical Release source ID required: true @@ -31,6 +35,7 @@ jobs: ADF_RELEASE_SOURCE_ID: ${{ inputs.source_id }} ADF_RELEASE_SIGNER_KEY_ID: ${{ inputs.signer_key_id }} ADF_RELEASE_OUTPUT_DIR: ${{ github.workspace }}/dist/framework + ADF_RELEASE_ID: ${{ inputs.release_id }} steps: - name: Check out reviewed source diff --git a/README.md b/README.md index ad2b19f..3536996 100644 --- a/README.md +++ b/README.md @@ -417,6 +417,68 @@ cargo build --locked --bin adf-claude-runner The signed binary release currently continues to publish only `adf`; runner distribution is a later compatibility milestone. +## Reducing stored Record size + +ADF can store identical input and freshness reference maps once within each Result. +Each file remains self-contained JSON. Reading it restores the exact logical Record, +including explicit versus omitted fields, before validation and digest calculation. +Record IDs, evidence, explanations and the pinned Framework identity are preserved. +Small Records remain plain JSON when sharing would cost more space. + +Existing projects keep their current write format until explicitly migrated. Before +migration, stop agents, MCP sessions and CI writers for that working directory and +update every reader/writer to a version supporting `adf-record-refmaps-v1`. Older +CLIs reject the new project configuration on most paths, but some old execution +commands bypass that check; configuration alone cannot stop a running old writer. + +```sh +adf project storage inspect --format json +adf project storage migrate --to adaptive-refmaps-v1 --dry-run +adf project storage migrate --to adaptive-refmaps-v1 +adf project storage verify +adf project storage export --record result. --format json +``` + +These commands accept `--project ` and the existing offline `--release ` +option. Inspect, dry-run, verify and export leave Records, configuration and derived +indexes unchanged; they acquire a small maintenance lock under `.adf/cache/locks/`. +Reports describe Result/Evidence JSON bytes only, excluding execution events, +non-JSON evidence, caches, Git history and other working directories. A skipped-file +list identifies non-JSON files and interrupted temporary files. + +Migration validates Records against the pinned signed release and checks lossless +roundtrips before writing. It enables project config version 2, atomically replaces +individual files, and keeps a small Git-ignored recovery journal under +`.adf/local/storage-migrations/`. A local `.gitignore` is created there if needed; +existing ignore rules are never overwritten. Configuration YAML may be reformatted. +Normal operations cooperate with the migration lock. Do not edit files, switch Git +revisions, or remove lock files while migration is running. + +Repeat the same migration command after interruption. If a Record changed since an +interrupted migration, the command stops instead of overwriting it. Inspect that +conflict before continuing. After verifying the changed Records, move the recovery +journal aside and run dry-run again to review a new plan; do not rewrite its hashes +to conceal the conflict. To restore compatibility with older CLIs, expand all +Records before downgrading the configuration: + +```sh +adf project storage migrate --to plain-json-v1 --dry-run +adf project storage migrate --to plain-json-v1 +``` + +Rollback restores equivalent JSON, including untracked Records, without keeping +second copies of their contents. Original whitespace is not retained. Free space is +checked before either direction. Completed replacement files are synced before +publication; an interruption during a temporary write can leave a temporary file, +which is reported and must be reviewed before removal. A restored Record may be up +to 256 MiB; a sharing envelope may be up to 64 MiB. Excessive expansion is rejected +before copying shared maps. JSON nesting is subject to the parser's depth limit. + +No commits, Git history rewrites, cache deletion or automatic legacy Kit migration +are performed. Existing indexes rebuild from restored Records when their physical +source changes; their own reference duplication is a separate capacity cost. A Git +revision or source change still triggers the normal freshness rules. + ## Commands | Command | What it does | diff --git a/docs/concepts.ja.md b/docs/concepts.ja.md index a9e4368..6ac5c9b 100644 --- a/docs/concepts.ja.md +++ b/docs/concepts.ja.md @@ -252,6 +252,35 @@ Feature Contractは、上位Contractの写しでも上書きでもありませ 現在有効な規範は`docs/`に留めず、`contracts/`へ上げます。`docs/adf/`は案内と索引です。 +## 記録の保存容量を減らす + +作業結果の中で同じ入力・鮮度確認用の一覧が繰り返される場合、一件の記録内で一覧を共有できます。読込み時に元と同じJSONの内容へ戻してから検証するため、記録ID、根拠、説明文、項目の省略・空の区別を保持します。一覧を共有しても小さくならない記録は通常のJSONを使います。 + +既存プロジェクトは、バイナリを更新しただけでは保存形式を変更しません。まず、その作業ディレクトリへ書き込むエージェント、MCP、CIを止め、新形式に対応するバイナリへ揃えてください。旧CLIの一部経路は設定の版を検査しないため、設定変更だけで旧プロセスを止められるわけではありません。 + +```sh +adf project storage inspect +adf project storage migrate --to adaptive-refmaps-v1 --dry-run +adf project storage migrate --to adaptive-refmaps-v1 +adf project storage verify +adf project storage export --record result. +``` + +対象は作業結果と証拠のJSONです。結果に表示する容量は、この対象のバイト数です。実行イベント、JSON以外の証拠、キャッシュ、Git履歴、別作業ディレクトリは含みません。対象外のファイルと停止時の一時ファイルは別途一覧に表示します。参照系の操作も小さな移行ロックを取得しますが、正規記録・設定・索引を変更しません。 + +移行は設定をversion 2へ切り替え、一件ずつ検証して置き換えます。本文の複製は残さず、再開に必要なハッシュなどをGit管理外の`.adf/local/storage-migrations/`へ保存します。必要なら、この専用ディレクトリに自身を無視する`.gitignore`を新規作成します。既存の無視設定は上書きしません。設定YAMLの整形は変わる場合があります。 + +停止した場合は同じ移行コマンドで再開できます。停止後に記録が変わっていれば、上書きせず止まります。内容を確認して新しい計画に切り替える場合は、再開用の`current.json`を別名で保管し、dry-runからやり直します。通常JSONへ戻す場合は、次を実行します。 + +```sh +adf project storage migrate --to plain-json-v1 --dry-run +adf project storage migrate --to plain-json-v1 +``` + +復元には追加の空き容量が必要です。旧CLIへ戻すのは、通常JSONへの展開と設定の復元が完了してからです。元の空白・改行までは復元しません。途中の一時書込みで停止した場合は、一時ファイルが残ることがあるため、報告された内容を確認してから整理します。移行中の手動編集、Git操作、ロックファイルの削除は避けてください。 + +新形式でもGit revisionやコード変更による証拠の失効は維持します。Git履歴やキャッシュを削除して容量を小さく見せる操作は行いません。索引の重複削減と、旧Kitを使うプロジェクトの移行は別の改善です。 + ## 制御基盤・エージェント・人の分担 | 担当 | 役割 | diff --git a/docs/publishing.md b/docs/publishing.md index 90f837f..2584abd 100644 --- a/docs/publishing.md +++ b/docs/publishing.md @@ -75,6 +75,17 @@ linked into them, and uploads the whole set as workflow artifacts. Nothing is published. Stop here and look at what it produced. +Set `release_id` to a new identifier, for example `adf-2026-09-11-storage`. +The publication tag will be `framework-adf-2026-09-11-storage`. Never reuse an +existing identifier for updated binaries. Set `source_id` and `signer_key_id` +to the configured repository values so existing projects retain their trust pins. +Local candidate builds accept the same identifier through `ADF_RELEASE_ID`; +the default `adf-dev` is retained for regression fixtures. +Signed release names identify archives; compatibility still requires exact +protocol versions and Rule/Schema digests. The development lock retains its +fixed `adf-dev` identity. Changing a signed release name does not relax signature +or pinned archive verification. + ### 2. Publish it Run the **Publish a Framework Release** workflow with the candidate run's ID diff --git a/schemas/storage/project-config-v2.schema.json b/schemas/storage/project-config-v2.schema.json new file mode 100644 index 0000000..c672ec4 --- /dev/null +++ b/schemas/storage/project-config-v2.schema.json @@ -0,0 +1,17 @@ +{ + "$schema": "https://json-schema.org/draft/2020-12/schema", + "$id": "adf://schemas/storage/project-config-v2", + "type": "object", + "required": ["schema_version", "record_storage", "project_sources", "repository_observation"], + "properties": { + "schema_version": {"const": "2"}, + "record_storage": {"const": "adaptive-refmaps-v1"}, + "project_sources": { + "type": "object", "required": ["contracts", "decisions"], + "properties": {"contracts": {"type": "string"}, "decisions": {"type": "string"}}, + "additionalProperties": false + }, + "repository_observation": {"type": "string"} + }, + "additionalProperties": false +} diff --git a/schemas/storage/record-refmaps-v1.schema.json b/schemas/storage/record-refmaps-v1.schema.json new file mode 100644 index 0000000..1cb436f --- /dev/null +++ b/schemas/storage/record-refmaps-v1.schema.json @@ -0,0 +1,24 @@ +{ + "$schema": "https://json-schema.org/draft/2020-12/schema", + "$id": "adf://schemas/storage/record-refmaps-v1", + "type": "object", + "required": ["storage_format", "record_kind", "record_digest", "record", "reference_maps", "map_bindings"], + "properties": { + "storage_format": {"const": "adf-record-refmaps-v1"}, + "record_kind": {"enum": ["result", "evidence"]}, + "record_digest": {"type": "string", "pattern": "^sha256:[0-9a-f]{64}$"}, + "record": {"type": "object"}, + "reference_maps": { + "type": "object", "minProperties": 1, + "propertyNames": {"pattern": "^sha256:[0-9a-f]{64}$"}, + "additionalProperties": {"type": "object", "additionalProperties": {"type": "string"}} + }, + "map_bindings": { + "type": "object", "minProperties": 1, + "propertyNames": {"pattern": "^/(input_refs|freshness_refs|payload/outcomes/(0|[1-9][0-9]*)/(input_refs|freshness_refs))$"}, + "additionalProperties": {"type": "string", "pattern": "^sha256:[0-9a-f]{64}$"} + } + }, + "additionalProperties": false, + "$comment": "The codec also verifies map hashes, allowed paths by record kind, missing/unused maps, no overwrites, resource limits and restored record_digest. Validate the restored Record with the pinned logical Schema. This storage Schema does not change the Framework identity." +} diff --git a/scripts/release-ci.sh b/scripts/release-ci.sh index 05462af..acb1713 100755 --- a/scripts/release-ci.sh +++ b/scripts/release-ci.sh @@ -12,6 +12,13 @@ BASE_LOCK=$KIT_ROOT/testdata/fixtures/db-sqs/framework-lock.yaml OUTPUT_DIR=${ADF_RELEASE_OUTPUT_DIR:-"$KIT_ROOT/dist/framework"} SOURCE_ID=${ADF_RELEASE_SOURCE_ID:-remote:official} SIGNER_KEY_ID=${ADF_RELEASE_SIGNER_KEY_ID:-framework.release.prototype} +RELEASE_ID=${ADF_RELEASE_ID:-adf-dev} +case "$RELEASE_ID" in + ''|*[!A-Za-z0-9._-]*) + echo "ADF_RELEASE_ID contains unsupported characters" >&2 + exit 2 + ;; +esac PUBLIC_KEY=${ADF_RELEASE_SIGNING_PUBLIC_KEY_HEX:?ADF_RELEASE_SIGNING_PUBLIC_KEY_HEX is required} : "${ADF_RELEASE_SIGNING_KEY_HEX:?ADF_RELEASE_SIGNING_KEY_HEX is required}" @@ -34,6 +41,19 @@ cleanup() { trap cleanup EXIT HUP INT TERM mkdir -p "$WORK_ROOT/source/schemas" "$WORK_ROOT/first" "$WORK_ROOT/second" +python3 - "$BASE_LOCK" "$WORK_ROOT/framework-lock.yaml" "$RELEASE_ID" <<'PY' +import re +import sys +from pathlib import Path + +source, output, release_id = sys.argv[1:] +text, count = re.subn(r'^framework_release:.*$', 'framework_release: "' + release_id + '"', + Path(source).read_text(), flags=re.MULTILINE) +if count != 1: + raise SystemExit('Expected exactly one Framework Release ID in the base lock') +Path(output).write_text(text) +PY +BASE_LOCK=$WORK_ROOT/framework-lock.yaml cp "$SOURCE_RULES" "$WORK_ROOT/source/rules.yaml" cp "$SOURCE_FRAMEWORK_CATALOG" "$WORK_ROOT/source/framework-catalog.yaml" cp -R "$SOURCE_SCHEMAS" "$WORK_ROOT/source/schemas/v1" diff --git a/scripts/tests/test-release-ci.sh b/scripts/tests/test-release-ci.sh index 5be8192..c18e0ba 100755 --- a/scripts/tests/test-release-ci.sh +++ b/scripts/tests/test-release-ci.sh @@ -32,6 +32,23 @@ test -s "$OUTPUT/candidate-framework.lock" test -s "$OUTPUT/distribution-trust.json" test -s "$OUTPUT/publish-receipt.json" +# A new CLI release must not collide with the development release tag. +VERSIONED_OUTPUT=$TEST_ROOT/versioned +ADF_RELEASE_ID=adf-storage-test run_release_ci "$VERSIONED_OUTPUT" "$PUBLIC_KEY" +python3 - "$VERSIONED_OUTPUT" <<'PY' +import json +import sys +from pathlib import Path +root = Path(sys.argv[1]) +assert json.loads((root / 'publish-receipt.json').read_text())['release_id'] == 'adf-storage-test' +assert 'adf-storage-test' in (root / 'candidate-framework.lock').read_text() +PY +if ADF_RELEASE_ID='../invalid' run_release_ci "$TEST_ROOT/invalid" "$PUBLIC_KEY" >/dev/null 2>&1; then + echo "Release CI accepted an unsafe release ID" >&2 + exit 1 +fi +test ! -e "$TEST_ROOT/invalid/framework-release.tar" + # A rerun must not replace an already reviewed candidate. if run_release_ci "$OUTPUT" "$PUBLIC_KEY" >/dev/null 2>&1; then echo "Release CI unexpectedly overwrote existing outputs" >&2 diff --git a/skill-src/adf-analyst/SKILL.md b/skill-src/adf-analyst/SKILL.md index 4384cf9..2a918df 100644 --- a/skill-src/adf-analyst/SKILL.md +++ b/skill-src/adf-analyst/SKILL.md @@ -207,3 +207,11 @@ measurements remain unknown. If submission is rejected as stale, the inputs moved under you. Call `adf_next` again and redo the work against the fresh action - do not retry the old payload. + +## Reading stored Records + +CLI and MCP return logical Records even when files use shared reference maps. To +inspect a file-backed Result or Evidence as ordinary JSON, run +`adf project storage export --record `. Do not hand-edit storage envelopes, +reference tables, or digests. Storage migration is an explicit maintenance task; +continue using the issued Context and normal submission tools for development. diff --git a/skill-src/adf-builder/SKILL.md b/skill-src/adf-builder/SKILL.md index 709541c..28c9479 100644 --- a/skill-src/adf-builder/SKILL.md +++ b/skill-src/adf-builder/SKILL.md @@ -103,3 +103,11 @@ collect these metrics. ADF records Context size without an extra model invocatio If submission is rejected as stale, the inputs moved under you. Call `adf_next` again and work from the fresh action rather than retrying the old payload. + +## Reading stored Records + +CLI and MCP return logical Records even when files use shared reference maps. To +inspect a file-backed Result or Evidence as ordinary JSON, run +`adf project storage export --record `. Do not hand-edit storage envelopes, +reference tables, or digests. Storage migration is an explicit maintenance task; +continue using the issued Context and normal submission tools for development. diff --git a/skill-src/adf-challenger/SKILL.md b/skill-src/adf-challenger/SKILL.md index 19c39f4..704ae2e 100644 --- a/skill-src/adf-challenger/SKILL.md +++ b/skill-src/adf-challenger/SKILL.md @@ -94,3 +94,11 @@ or another model call solely to collect them. If submission is rejected as stale, the work moved under you. Call `adf_next` again and challenge the fresh state rather than retrying the old payload. + +## Reading stored Records + +CLI and MCP return logical Records even when files use shared reference maps. To +inspect a file-backed Result or Evidence as ordinary JSON, run +`adf project storage export --record `. Do not hand-edit storage envelopes, +reference tables, or digests. Storage migration is an explicit maintenance task; +continue using the issued Context and normal submission tools for development. diff --git a/src/delivery.rs b/src/delivery.rs index 24e5a23..77c634c 100644 --- a/src/delivery.rs +++ b/src/delivery.rs @@ -289,6 +289,8 @@ pub fn switch_framework_lock( project_root: &Path, candidate_lock_path: &Path, ) -> Result { + let _storage_guard = + crate::storage_io::StorageGuard::shared(project_root).map_err(delivery_error)?; switch_framework_lock_for(project_root, candidate_lock_path, TrustUse::NewActivation) } @@ -355,6 +357,8 @@ pub fn rollback_framework_lock( project_root: &Path, backup_lock_path: &Path, ) -> Result { + let _storage_guard = + crate::storage_io::StorageGuard::shared(project_root).map_err(delivery_error)?; let project_root = canonical_project_root(project_root)?; let backups_root = project_root.join(".adf/cache/framework-lock-backups"); let backup_path = absolute_from(&project_root, backup_lock_path) diff --git a/src/execution_record.rs b/src/execution_record.rs index 6fc71b7..a36b3d6 100644 --- a/src/execution_record.rs +++ b/src/execution_record.rs @@ -146,6 +146,7 @@ impl ExecutionEvent { #[derive(Debug, Clone)] pub struct ExecutionEventStore { + _storage_guard: crate::storage_io::StorageGuard, project_root: PathBuf, } @@ -155,7 +156,14 @@ impl ExecutionEventStore { .as_ref() .canonicalize() .map_err(|error| format!("cannot resolve project root: {error}"))?; - Ok(Self { project_root }) + let storage_guard = crate::storage_io::StorageGuard::shared(&project_root)?; + if project_root.join(".adf/config.yaml").exists() { + crate::project_config::load_project_config(&project_root).map_err(|e| e.to_string())?; + } + Ok(Self { + project_root, + _storage_guard: storage_guard, + }) } pub fn begin( @@ -291,8 +299,9 @@ impl ExecutionEventStore { paths.sort(); for path in paths { let bytes = read_regular_file(&path)?; - let result: Value = serde_json::from_slice(&bytes) - .map_err(|error| format!("{}: {error}", path.display()))?; + let result = + crate::record_storage::decode(&bytes, crate::record_storage::RecordKind::Result) + .map_err(|error| format!("{}: {error}", path.display()))?; if result["id"].as_str() == Some(result_id) { if result["action_id"].as_str() == Some(started.action_id.as_str()) && result["context_digest"].as_str() == Some(started.context_digest.as_str()) diff --git a/src/filesystem_project.rs b/src/filesystem_project.rs index c29bcf4..4864d8b 100644 --- a/src/filesystem_project.rs +++ b/src/filesystem_project.rs @@ -11,7 +11,9 @@ use crate::canonical_digest; use crate::contract_health::{ContractHealthReport, build_contract_health_report_with_digests}; use crate::kernel::ProjectSnapshot; use crate::project::build_project_snapshot; +use crate::record_storage::{self, RecordKind, StoragePolicy}; use crate::schema::SchemaRegistry; +use crate::storage_io::{self, StorageGuard}; use regex::Regex; use serde::{Deserialize, Serialize}; use serde_json::{Value, json}; @@ -19,7 +21,7 @@ use sha2::{Digest, Sha256}; use std::collections::{BTreeMap, BTreeSet}; use std::fmt; use std::fs::{self, File, OpenOptions}; -use std::io::{Read, Write}; +use std::io::Write; use std::path::{Component, Path, PathBuf}; use std::process::{Command, Stdio}; @@ -36,7 +38,7 @@ const DISALLOWED_SOURCE_ROOTS: [&str; 5] = [ ]; pub const FILESYSTEM_PROJECT_PROTOCOL_VERSION: &str = "3"; static TEMPORARY_SEQUENCE: AtomicU64 = AtomicU64::new(0); -const HEALTH_INDEX_SCHEMA_VERSION: &str = "1"; +const HEALTH_INDEX_SCHEMA_VERSION: &str = "2"; const RESULT_INDEX_PATH: &str = ".adf/cache/runtime/contract-health-results-v1.json"; const EVIDENCE_INDEX_PATH: &str = ".adf/cache/runtime/contract-health-evidence-v1.json"; @@ -75,6 +77,8 @@ struct InitializationTarget { } pub struct FileProjectStore<'a> { + _storage_guard: StorageGuard, + storage_policy: StoragePolicy, project_root: PathBuf, contract_root: PathBuf, decision_root: PathBuf, @@ -118,10 +122,20 @@ impl<'a> FileProjectStore<'a> { root.display() ))); } + let storage_guard = StorageGuard::shared(&root).map_err(file_error)?; + let storage_policy = if root.join(".adf/config.yaml").exists() { + crate::project_config::load_project_config(&root) + .map_err(|e| file_error(e.to_string()))? + .record_storage + } else { + StoragePolicy::Plain + }; let contract_root = source_root(&root, contract_root)?; let decision_root = source_root(&root, decision_root)?; let change_root = root.join(".adf").join("changes"); Ok(Self { + _storage_guard: storage_guard, + storage_policy, project_root: root, contract_root, decision_root, @@ -352,6 +366,7 @@ impl<'a> FileProjectStore<'a> { .join(change_id) .join("results") .join(result_filename(result)?); + let _correction = storage_io::correction_lock(&self.project_root).map_err(file_error)?; let existing = read_json(&path)?; if existing["id"].as_str() != Some(expected_result_id) { return Err(file_error("Result changed before correction")); @@ -655,8 +670,15 @@ impl<'a> FileProjectStore<'a> { .map(Ok) .unwrap_or_else(|| fs::read(path)) .map_err(|error| file_error(format!("{}: {error}", path.display())))?; - let original: Value = serde_json::from_slice(&bytes) - .map_err(|error| file_error(format!("{}: {error}", path.display())))?; + let original = record_storage::decode( + &bytes, + if compact_results { + RecordKind::Result + } else { + RecordKind::Evidence + }, + ) + .map_err(|error| file_error(format!("{}: {error}", path.display())))?; let record_digest = canonical_digest(&original) .map_err(|error| file_error(error.to_string()))?; let record = if compact_results { @@ -729,26 +751,31 @@ impl<'a> FileProjectStore<'a> { value: &Value, format: FileFormat, ) -> Result<(), FileProjectError> { - let serialized = serialize_record(value, format, None)?; - if let Some(parent) = path.parent() { - fs::create_dir_all(parent).map_err(|error| { - file_error(format!("cannot create {}: {error}", parent.display())) - })?; + storage_io::reject_symlinks(&self.project_root, path).map_err(file_error)?; + let serialized = self.serialize_stored(value, format, None)?; + storage_io::atomic_write(path, serialized.as_bytes(), true).map_err(file_error) + } + + fn serialize_stored( + &self, + value: &Value, + format: FileFormat, + existing: Option<&str>, + ) -> Result { + if matches!(format, FileFormat::Json) && self.storage_policy == StoragePolicy::Adaptive { + let kind = if value["id"] + .as_str() + .is_some_and(|id| id.starts_with("result.")) + { + RecordKind::Result + } else { + RecordKind::Evidence + }; + let bytes = + record_storage::encode(value, kind, self.storage_policy).map_err(file_error)?; + return String::from_utf8(bytes).map_err(|e| file_error(e.to_string())); } - let mut file = OpenOptions::new() - .write(true) - .create_new(true) - .open(path) - .map_err(|error| { - if error.kind() == std::io::ErrorKind::AlreadyExists { - file_error(format!("record already exists: {}", path.display())) - } else { - file_error(format!("cannot create {}: {error}", path.display())) - } - })?; - file.write_all(serialized.as_bytes()) - .and_then(|()| file.sync_all()) - .map_err(|error| file_error(format!("cannot write {}: {error}", path.display()))) + serialize_record(value, format, existing) } fn write_atomic( @@ -770,7 +797,12 @@ impl<'a> FileProjectStore<'a> { } else { None }; - let serialized = serialize_record(value, format, existing.as_deref())?; + let serialized = self.serialize_stored(value, format, existing.as_deref())?; + storage_io::reject_symlinks(&self.project_root, path).map_err(file_error)?; + if matches!(format, FileFormat::Json) { + return storage_io::atomic_write(path, serialized.as_bytes(), false) + .map_err(file_error); + } let temporary = temporary_path(path)?; let write_result = (|| { let mut file = OpenOptions::new() @@ -1202,11 +1234,13 @@ fn assert_unique_ids(values: &[Value]) -> Result<(), FileProjectError> { } fn read_json(path: &Path) -> Result { - let mut bytes = Vec::new(); - File::open(path) - .and_then(|mut file| file.read_to_end(&mut bytes)) - .map_err(|error| file_error(format!("{}: {error}", path.display())))?; - let value: Value = serde_json::from_slice(&bytes) + let bytes = storage_io::read_record(path).map_err(file_error)?; + let kind = match parent_name(path) { + Some("results") => RecordKind::Result, + Some("evidence") => RecordKind::Evidence, + _ => return Err(file_error("unsupported stored Record location")), + }; + let value = record_storage::decode(&bytes, kind) .map_err(|error| file_error(format!("{}: {error}", path.display())))?; require_object(value, path) } diff --git a/src/framework_lock.rs b/src/framework_lock.rs index 09a1d90..c199f99 100644 --- a/src/framework_lock.rs +++ b/src/framework_lock.rs @@ -66,8 +66,27 @@ pub fn validate_framework_lock( rule_index: &RuleIndex, schema_registry: &SchemaRegistry, ) -> Result { - let expected = build_framework_lock(rule_source, rule_index, schema_registry)?; + let mut expected = build_framework_lock(rule_source, rule_index, schema_registry)?; let comparable = comparable_core_lock(lock_source)?; + if lock_source.get("schema_version").and_then(Value::as_str) + == Some(SIGNED_FRAMEWORK_LOCK_SCHEMA_VERSION) + { + let release_id = lock_source + .get("framework_release") + .and_then(Value::as_str) + .filter(|id| { + !id.is_empty() + && *id != "." + && *id != ".." + && id + .bytes() + .all(|b| b.is_ascii_alphanumeric() || b"._-".contains(&b)) + }) + .ok_or_else(|| FrameworkLockError::new("invalid signed Framework Release ID"))?; + // Signed delivery checks this label against the pinned archive. Runtime + // compatibility is determined by the exact protocol and content fields. + expected["framework_release"] = Value::String(release_id.to_owned()); + } let differences = differences(&expected, &comparable, ""); if !differences.is_empty() { return Err(FrameworkLockError::new(format!( @@ -208,6 +227,37 @@ impl std::error::Error for FrameworkLockError {} mod tests { use super::*; + #[test] + fn signed_release_labels_do_not_relax_runtime_compatibility() { + let root = std::path::Path::new(env!("CARGO_MANIFEST_DIR")); + let registry = SchemaRegistry::load(root.join("schemas/v1")).unwrap(); + let rules: Value = serde_yaml::from_str( + &std::fs::read_to_string(root.join("testdata/fixtures/db-sqs/rules.yaml")).unwrap(), + ) + .unwrap(); + let index = crate::rules::compile_rule_index(&rules, ®istry).unwrap(); + let mut lock = build_framework_lock(&rules, &index, ®istry).unwrap(); + lock["framework_release"] = json!("adf-storage-test"); + assert!(validate_framework_lock(&lock, &rules, &index, ®istry).is_err()); + lock["schema_version"] = json!("2"); + lock["release_artifact"] = json!({ + "artifact_digest": format!("sha256:{}", "a".repeat(64)), + "signer_key_id": "test.key", + "source_id": "remote:test" + }); + assert!(validate_framework_lock(&lock, &rules, &index, ®istry).is_ok()); + for field in ["protocols", "schema_bundle", "rule_set"] { + let mut invalid = lock.clone(); + invalid[field] = json!({}); + assert!(validate_framework_lock(&invalid, &rules, &index, ®istry).is_err()); + } + for label in ["", ".", "..", "../escape", "a/b"] { + let mut invalid = lock.clone(); + invalid["framework_release"] = json!(label); + assert!(validate_framework_lock(&invalid, &rules, &index, ®istry).is_err()); + } + } + #[test] fn reports_nested_missing_and_unexpected_fields_in_order() { assert_eq!( diff --git a/src/lib.rs b/src/lib.rs index aa7fa72..3eb87a8 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -44,6 +44,7 @@ pub mod project_config; pub mod project_runtime; pub mod project_setup; mod python_detection; +pub mod record_storage; pub mod release_publisher; pub mod remote_delivery; mod ruby_detection; @@ -54,6 +55,8 @@ pub mod schema; mod script_detection; pub mod signal_catalog; mod source_detection; +pub mod storage_io; +pub mod storage_migration; pub mod submission; mod swift_detection; diff --git a/src/main.rs b/src/main.rs index 082a67f..740b551 100644 --- a/src/main.rs +++ b/src/main.rs @@ -124,6 +124,21 @@ fn main() -> ExitCode { } }; } + if command == "project" && arguments.get(2).map(String::as_str) == Some("storage") { + return match adf::storage_migration::run_cli(&arguments[3..]) { + Ok(value) => { + println!( + "{}", + serde_json::to_string_pretty(&value).expect("JSON response") + ); + ExitCode::SUCCESS + } + Err(error) => { + eprintln!("{error}"); + ExitCode::FAILURE + } + }; + } if command == "project" { let options = match parse_project_management_command(&arguments[2..]) { Ok(options) => options, @@ -2385,6 +2400,7 @@ fn usage_text() -> &'static str { "Agentic Development Framework\n\ \n\ usage:\n\ + adf project storage [--project ] [--to ] [--dry-run] [--record ] [--release ] [--format json]\n\ adf project init [--project ] [--candidate-dir ] [--analysis-root ]...\n\ adf project observe [--project ] [--analysis-root ]... [--format ] [--output ]\n\ adf project validate-bindings [--project ] [--draft ] [--format ] [--require-clean]\n\ diff --git a/src/migration.rs b/src/migration.rs index 365f855..9b117ea 100644 --- a/src/migration.rs +++ b/src/migration.rs @@ -1320,6 +1320,15 @@ pub fn apply_migration_candidate( ))); } + let _storage_guard = crate::storage_io::StorageGuard::shared(&root).map_err(migration_error)?; + if root.join(".adf/config.yaml").exists() { + let config = load_project_config(&root).map_err(|e| migration_error(e.to_string()))?; + if config.record_storage != crate::record_storage::StoragePolicy::Plain { + return Err(migration_error( + "expand Record storage to plain-json-v1 before applying a legacy migration candidate", + )); + } + } let manifest = load_candidate_manifest(&candidate_root)?; let draft = load_candidate_draft(&candidate_root)?; let manifest_value = serde_json::to_value(&manifest).map_err(|error| { diff --git a/src/project_config.rs b/src/project_config.rs index eecba63..73ab81d 100644 --- a/src/project_config.rs +++ b/src/project_config.rs @@ -10,6 +10,7 @@ pub const PROJECT_CONFIG_SCHEMA_VERSION: &str = "1"; #[derive(Debug, Clone, PartialEq, Eq)] pub struct ProjectConfig { + pub record_storage: crate::record_storage::StoragePolicy, pub contract_root: String, pub decision_root: String, pub repository_observation: String, @@ -24,21 +25,36 @@ pub fn load_project_config(root: &Path) -> Result ( + vec![ + "schema_version", + "project_sources", + "repository_observation", + ], + crate::record_storage::StoragePolicy::Plain, + ), + Some("2") + if object.get("record_storage").and_then(Value::as_str) + == Some("adaptive-refmaps-v1") => + { + ( + vec![ + "schema_version", + "project_sources", + "repository_observation", + "record_storage", + ], + crate::record_storage::StoragePolicy::Adaptive, + ) + } + _ => { + return Err(config_error( + "unsupported project config schema or record_storage policy", + )); + } + }; + assert_exact_fields(object.keys().map(String::as_str), &fields, "project config")?; let sources = object["project_sources"] .as_object() .ok_or_else(|| config_error("project_sources must be a mapping"))?; @@ -57,6 +73,7 @@ pub fn load_project_config(root: &Path) -> Result with letters, digits, '.', '_', or '-'", )); } + let _storage_guard = + crate::storage_io::StorageGuard::shared(project_root).map_err(setup_error)?; let root = project_root .canonicalize() .map_err(|error| setup_error(format!("cannot resolve project root: {error}")))?; @@ -450,6 +457,8 @@ pub fn promote_observation_draft( project_root: &Path, draft: &Value, ) -> Result { + let _storage_guard = + crate::storage_io::StorageGuard::shared(project_root).map_err(setup_error)?; let root = project_root .canonicalize() .map_err(|error| setup_error(format!("cannot resolve project root: {error}")))?; @@ -933,25 +942,25 @@ fn atomic_replace_if_digest( } #[cfg(unix)] -fn sync_directory(path: &Path) -> Result<(), ProjectSetupError> { +pub(crate) fn sync_directory(path: &Path) -> Result<(), ProjectSetupError> { fs::File::open(path) .and_then(|directory| directory.sync_all()) .map_err(|error| setup_error(format!("{}: {error}", path.display()))) } #[cfg(windows)] -fn sync_directory(_path: &Path) -> Result<(), ProjectSetupError> { +pub(crate) fn sync_directory(_path: &Path) -> Result<(), ProjectSetupError> { Ok(()) } #[cfg(not(windows))] -fn replace_file(source: &Path, target: &Path) -> Result<(), ProjectSetupError> { +pub(crate) fn replace_file(source: &Path, target: &Path) -> Result<(), ProjectSetupError> { fs::rename(source, target) .map_err(|error| setup_error(format!("{}: {error}", target.display()))) } #[cfg(windows)] -fn replace_file(source: &Path, target: &Path) -> Result<(), ProjectSetupError> { +pub(crate) fn replace_file(source: &Path, target: &Path) -> Result<(), ProjectSetupError> { use std::os::windows::ffi::OsStrExt; const MOVEFILE_REPLACE_EXISTING: u32 = 0x1; diff --git a/src/record_storage.rs b/src/record_storage.rs new file mode 100644 index 0000000..22c82fa --- /dev/null +++ b/src/record_storage.rs @@ -0,0 +1,433 @@ +//! Lossless, self-contained storage of repeated Record reference maps. +//! This module never changes logical Records or their identity rules. + +use crate::{canonical_digest, canonical_json}; +use serde::Deserialize; +use serde::de::{self, MapAccess, SeqAccess, Visitor}; +use serde_json::{Map, Value, json}; +use std::collections::{BTreeMap, BTreeSet}; +use std::fmt; + +pub const STORAGE_FORMAT: &str = "adf-record-refmaps-v1"; +pub const MAX_STORED_BYTES: usize = 64 * 1024 * 1024; +pub const MAX_EXPANDED_BYTES: usize = 256 * 1024 * 1024; + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum RecordKind { + Result, + Evidence, +} + +impl RecordKind { + pub fn as_str(self) -> &'static str { + match self { + Self::Result => "result", + Self::Evidence => "evidence", + } + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)] +pub enum StoragePolicy { + #[default] + Plain, + Adaptive, +} + +impl StoragePolicy { + pub fn as_str(self) -> &'static str { + match self { + Self::Plain => "plain-json-v1", + Self::Adaptive => "adaptive-refmaps-v1", + } + } +} + +fn canonical(value: &Value) -> Result { + canonical_json(value).map_err(|error| error.to_string()) +} + +fn digest(value: &Value) -> Result { + canonical_digest(value).map_err(|error| error.to_string()) +} + +fn slots(record: &Value, kind: RecordKind) -> Vec { + let mut paths = Vec::new(); + for key in ["input_refs", "freshness_refs"] { + if (kind == RecordKind::Result || key == "input_refs") && record.get(key).is_some() { + paths.push(format!("/{key}")); + } + } + if kind == RecordKind::Result + && let Some(outcomes) = record + .pointer("/payload/outcomes") + .and_then(Value::as_array) + { + for (i, outcome) in outcomes.iter().enumerate() { + for key in ["input_refs", "freshness_refs"] { + if outcome.get(key).is_some() { + paths.push(format!("/payload/outcomes/{i}/{key}")); + } + } + } + } + paths +} + +fn location(path: &str, kind: RecordKind) -> Result<(&str, &str), String> { + if path == "/input_refs" || (kind == RecordKind::Result && path == "/freshness_refs") { + return Ok(("", &path[1..])); + } + if kind == RecordKind::Result + && let Some(rest) = path.strip_prefix("/payload/outcomes/") + && let Some((index, field)) = rest.split_once('/') + && matches!(field, "input_refs" | "freshness_refs") + && index.parse::().is_ok_and(|i| i.to_string() == index) + { + return Ok((&path[..path.len() - field.len() - 1], field)); + } + Err(format!("invalid reference map location: {path}")) +} + +fn reference_map(value: &Value) -> Result<(), String> { + if value + .as_object() + .is_some_and(|map| map.values().all(Value::is_string)) + { + Ok(()) + } else { + Err("reference map must be a string-to-string object".to_owned()) + } +} + +/// Caller validates the logical Record with the pinned Schema before writing. +pub fn encode(record: &Value, kind: RecordKind, policy: StoragePolicy) -> Result, String> { + let plain = canonical(record)? + "\n"; + if plain.len() > MAX_EXPANDED_BYTES { + return Err("expanded Record exceeds storage limit".to_owned()); + } + if policy == StoragePolicy::Plain { + return Ok(plain.into_bytes()); + } + let mut counts = BTreeMap::::new(); + let mut values = BTreeMap::::new(); + let mut bindings = Map::new(); + for path in slots(record, kind) { + let value = record.pointer(&path).expect("slot exists"); + reference_map(value)?; + let hash = digest(value)?; + if let Some(previous) = values.get(&hash) + && previous != value + { + return Err("reference map hash collision".to_owned()); + } + values.entry(hash.clone()).or_insert_with(|| value.clone()); + *counts.entry(hash.clone()).or_default() += 1; + bindings.insert(path, Value::String(hash)); + } + bindings.retain(|_, hash| counts[hash.as_str().expect("hash string")] > 1); + if bindings.is_empty() { + return Ok(plain.into_bytes()); + } + values.retain(|hash, _| counts[hash] > 1); + let mut body = record.clone(); + for path in bindings.keys() { + let (parent, field) = location(path, kind)?; + body.pointer_mut(parent) + .and_then(Value::as_object_mut) + .ok_or("reference map parent is not an object")? + .remove(field); + } + let envelope = json!({ + "storage_format": STORAGE_FORMAT, + "record_kind": kind.as_str(), + "record_digest": digest(record)?, + "record": body, + "reference_maps": values, + "map_bindings": bindings, + }); + let packed = canonical(&envelope)? + "\n"; + if packed.len() < plain.len() && packed.len() <= MAX_STORED_BYTES { + Ok(packed.into_bytes()) + } else { + Ok(plain.into_bytes()) + } +} + +/// Decode before Schema validation, hashing, cache projection or identity checks. +pub fn decode(bytes: &[u8], kind: RecordKind) -> Result { + decode_with_limit(bytes, kind, MAX_EXPANDED_BYTES) +} + +fn decode_with_limit(bytes: &[u8], kind: RecordKind, limit: usize) -> Result { + if bytes.len() > MAX_EXPANDED_BYTES { + return Err("Record exceeds storage limit".to_owned()); + } + let envelope = parse_strict(bytes)?; + let object = envelope.as_object().ok_or("Record must be an object")?; + if !object.contains_key("storage_format") { + if bytes.len() > limit { + return Err("expanded Record exceeds storage limit".to_owned()); + } + return Ok(envelope); + } + if bytes.len() > MAX_STORED_BYTES { + return Err("stored Record exceeds storage limit".to_owned()); + } + let expected = [ + "map_bindings", + "record", + "record_digest", + "record_kind", + "reference_maps", + "storage_format", + ]; + if object.keys().map(String::as_str).collect::>() != expected.into_iter().collect() + || envelope["storage_format"] != STORAGE_FORMAT + || envelope["record_kind"] != kind.as_str() + { + return Err("unsupported or invalid Record storage envelope".to_owned()); + } + let maps = envelope["reference_maps"] + .as_object() + .ok_or("reference_maps must be an object")?; + let bindings = envelope["map_bindings"] + .as_object() + .ok_or("map_bindings must be an object")?; + if maps.is_empty() || bindings.is_empty() { + return Err("storage envelope must contain reference maps and bindings".to_owned()); + } + let mut sizes = BTreeMap::new(); + for (hash, value) in maps { + reference_map(value)?; + if hash != &digest(value)? { + return Err("reference map digest mismatch".to_owned()); + } + sizes.insert(hash.as_str(), canonical(value)?.len()); + } + let mut body = envelope["record"].clone(); + if !body.is_object() { + return Err("stored record must be an object".to_owned()); + } + // Budget all insertions before cloning any shared map into its destinations. + let mut expanded = canonical(&body)?.len(); + let mut used = BTreeSet::new(); + for (path, hash) in bindings { + let hash = hash + .as_str() + .ok_or("reference map binding must be a string")?; + let size = sizes.get(hash).ok_or("missing reference map")?; + used.insert(hash); + let (parent, field) = location(path, kind)?; + let parent = body + .pointer(parent) + .and_then(Value::as_object) + .ok_or("reference map parent is missing")?; + if parent.contains_key(field) { + return Err("reference map binding would overwrite an inline field".to_owned()); + } + expanded = expanded + .checked_add(*size) + .and_then(|n| n.checked_add(field.len() + 4)) + .filter(|n| *n <= limit) + .ok_or("expanded Record exceeds storage limit")?; + } + if used.len() != maps.len() { + return Err("unused reference map".to_owned()); + } + for (path, hash) in bindings { + let (parent, field) = location(path, kind)?; + body.pointer_mut(parent) + .and_then(Value::as_object_mut) + .expect("validated location") + .insert( + field.to_owned(), + maps[hash.as_str().expect("validated hash")].clone(), + ); + } + if envelope["record_digest"].as_str() != Some(digest(&body)?.as_str()) { + return Err("restored Record digest mismatch".to_owned()); + } + Ok(body) +} + +/// Duplicate object keys cannot silently change a reference table or binding. +pub fn parse_strict(bytes: &[u8]) -> Result { + let mut deserializer = serde_json::Deserializer::from_slice(bytes); + let value = StrictValue::deserialize(&mut deserializer).map_err(|e| e.to_string())?; + deserializer.end().map_err(|e| e.to_string())?; + Ok(value.0) +} + +struct StrictValue(Value); + +impl<'de> Deserialize<'de> for StrictValue { + fn deserialize>(deserializer: D) -> Result { + struct StrictVisitor; + impl<'de> Visitor<'de> for StrictVisitor { + type Value = StrictValue; + fn expecting(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + f.write_str("JSON without duplicate keys") + } + fn visit_bool(self, v: bool) -> Result { + Ok(StrictValue(v.into())) + } + fn visit_i64(self, v: i64) -> Result { + Ok(StrictValue(v.into())) + } + fn visit_u64(self, v: u64) -> Result { + Ok(StrictValue(v.into())) + } + fn visit_f64(self, v: f64) -> Result { + serde_json::Number::from_f64(v) + .map(|n| StrictValue(Value::Number(n))) + .ok_or_else(|| E::custom("invalid number")) + } + fn visit_str(self, v: &str) -> Result { + Ok(StrictValue(v.into())) + } + fn visit_unit(self) -> Result { + Ok(StrictValue(Value::Null)) + } + fn visit_seq>(self, mut seq: A) -> Result { + let mut values = Vec::new(); + while let Some(StrictValue(value)) = seq.next_element()? { + values.push(value); + } + Ok(StrictValue(Value::Array(values))) + } + fn visit_map>(self, mut map: A) -> Result { + let mut values = Map::new(); + while let Some((key, StrictValue(value))) = + map.next_entry::()? + { + if values.insert(key, value).is_some() { + return Err(de::Error::custom("duplicate JSON key")); + } + } + Ok(StrictValue(Value::Object(values))) + } + } + deserializer.deserialize_any(StrictVisitor) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + fn example() -> Value { + let refs: Map = (0..30) + .map(|i| { + ( + format!("code.日本語.{i}"), + json!(format!("sha256:{}", "a".repeat(64))), + ) + }) + .collect(); + json!({"id":"result.example", "input_refs":refs, "freshness_refs":refs, + "execution":{"context_bytes":30}, "payload":{"outcomes":[ + {"input_refs": refs}, {"freshness_refs":{}, "summary":"keep empty and missing distinct"}], + "custom":{"input_refs":refs}}}) + } + + #[test] + fn preserves_exact_logical_record_and_opaque_payload() { + let record = example(); + let packed = encode(&record, RecordKind::Result, StoragePolicy::Adaptive).unwrap(); + let envelope = parse_strict(&packed).unwrap(); + assert_eq!(envelope["storage_format"], STORAGE_FORMAT); + let schema: Value = serde_json::from_str(include_str!( + "../schemas/storage/record-refmaps-v1.schema.json" + )) + .unwrap(); + crate::schema::validate_json_document(&envelope, &schema).unwrap(); + assert_eq!( + envelope["record"]["payload"]["custom"], + record["payload"]["custom"] + ); + let restored = decode(&packed, RecordKind::Result).unwrap(); + assert_eq!(restored, record); + assert_eq!( + encode(&restored, RecordKind::Result, StoragePolicy::Adaptive).unwrap(), + packed + ); + assert!(packed.len() < canonical(&record).unwrap().len()); + assert!(decode_with_limit(&packed, RecordKind::Result, 100).is_err()); + } + + #[test] + fn rejects_corrupt_ambiguous_and_unknown_envelopes() { + let original = + parse_strict(&encode(&example(), RecordKind::Result, StoragePolicy::Adaptive).unwrap()) + .unwrap(); + for pointer in ["/storage_format", "/record_kind", "/record_digest"] { + let mut value = original.clone(); + *value.pointer_mut(pointer).unwrap() = json!("bad"); + assert!(decode(&serde_json::to_vec(&value).unwrap(), RecordKind::Result).is_err()); + } + let hash = original["reference_maps"] + .as_object() + .unwrap() + .keys() + .next() + .unwrap(); + for path in [ + "/payload/custom/input_refs", + "/payload/outcomes/00/input_refs", + "/payload/outcomes/99/input_refs", + ] { + let mut value = original.clone(); + value["map_bindings"][path] = json!(hash); + assert!(decode(&serde_json::to_vec(&value).unwrap(), RecordKind::Result).is_err()); + } + let mut overwrite = original.clone(); + overwrite["record"]["input_refs"] = json!({}); + assert!(decode(&serde_json::to_vec(&overwrite).unwrap(), RecordKind::Result).is_err()); + let mut missing = original.clone(); + missing["reference_maps"] = json!({}); + assert!(decode(&serde_json::to_vec(&missing).unwrap(), RecordKind::Result).is_err()); + let mut tampered = original; + tampered["record"]["execution"]["context_bytes"] = json!(31); + assert!(decode(&serde_json::to_vec(&tampered).unwrap(), RecordKind::Result).is_err()); + assert!(parse_strict(br#"{"a":1,"a":2}"#).is_err()); + } + + #[test] + fn python_storage_fixture_preserves_schema_and_identity() { + let logical = + parse_strict(include_bytes!("../testdata/storage/v1/result.logical.json")).unwrap(); + let stored = include_bytes!("../testdata/storage/v1/result.stored.json"); + assert_eq!(decode(stored, RecordKind::Result).unwrap(), logical); + assert_eq!( + encode(&logical, RecordKind::Result, StoragePolicy::Adaptive).unwrap(), + stored + ); + let registry = crate::schema::SchemaRegistry::load( + std::path::Path::new(env!("CARGO_MANIFEST_DIR")).join("schemas/v1"), + ) + .unwrap(); + registry.validate("result", &logical).unwrap(); + let mut body = logical.clone(); + body.as_object_mut().unwrap().remove("id"); + body.as_object_mut().unwrap().remove("execution"); + assert_eq!( + logical["id"], + format!("result.{}", &digest(&body).unwrap()[7..27]) + ); + } + + #[test] + fn small_records_stay_plain() { + let record = json!({"input_refs":{"code.a":"sha256:ABC"}}); + let bytes = encode(&record, RecordKind::Evidence, StoragePolicy::Adaptive).unwrap(); + assert_eq!(decode(&bytes, RecordKind::Evidence).unwrap(), record); + assert!( + !parse_strict(&bytes) + .unwrap() + .as_object() + .unwrap() + .contains_key("storage_format") + ); + } +} diff --git a/src/release_publisher.rs b/src/release_publisher.rs index 8d1c2a4..2e582e1 100644 --- a/src/release_publisher.rs +++ b/src/release_publisher.rs @@ -140,8 +140,18 @@ pub fn publish_release( SchemaRegistry::load(&schemas_root).map_err(|error| publish_error(error.to_string()))?; let rule_index = compile_rule_index(&rule_source, &schema_registry) .map_err(|error| publish_error(error.to_string()))?; - validate_framework_lock(&base_lock, &rule_source, &rule_index, &schema_registry) - .map_err(|error| publish_error(error.to_string()))?; + let mut runtime_base_lock = base_lock.clone(); + // The publisher assigns the signed archive's release label. Keep the + // unsigned development lock check strict for every technical field. + runtime_base_lock["framework_release"] = + Value::String(crate::framework_lock::FRAMEWORK_RELEASE.to_owned()); + validate_framework_lock( + &runtime_base_lock, + &rule_source, + &rule_index, + &schema_registry, + ) + .map_err(|error| publish_error(error.to_string()))?; let inventory = files .iter() diff --git a/src/storage_io.rs b/src/storage_io.rs new file mode 100644 index 0000000..a51affa --- /dev/null +++ b/src/storage_io.rs @@ -0,0 +1,149 @@ +//! Cooperative maintenance locking and atomic Record publication. + +use std::fs::{self, File, OpenOptions}; +use std::io::{Read, Write}; +use std::path::Path; +use std::sync::Arc; +use std::sync::atomic::{AtomicU64, Ordering}; + +static SEQUENCE: AtomicU64 = AtomicU64::new(0); + +#[derive(Debug, Clone)] +pub struct StorageGuard { + _file: Arc, +} + +impl StorageGuard { + pub fn shared(root: &Path) -> Result { + Self::acquire(root, false) + } + + pub fn exclusive(root: &Path) -> Result { + Self::acquire(root, true) + } + + fn acquire(root: &Path, exclusive: bool) -> Result { + let root = root.canonicalize().map_err(|e| e.to_string())?; + let relative = ".adf/cache/locks/record-storage.lock"; + let path = root.join(relative); + reject_symlinks(&root, &path)?; + fs::create_dir_all(path.parent().ok_or("missing lock parent")?) + .map_err(|e| e.to_string())?; + let file = OpenOptions::new() + .read(true) + .write(true) + .create(true) + .truncate(false) + .open(&path) + .map_err(|e| e.to_string())?; + let locked = if exclusive { + file.try_lock() + } else { + file.try_lock_shared() + }; + locked.map_err(|e| { + format!("Record storage is busy; finish active ADF operations and retry: {e}") + })?; + Ok(Self { + _file: Arc::new(file), + }) + } +} + +pub fn reject_symlinks(root: &Path, path: &Path) -> Result<(), String> { + let relative = path + .strip_prefix(root) + .map_err(|_| "path escapes project")?; + let mut current = root.to_path_buf(); + for part in relative.components() { + if !matches!(part, std::path::Component::Normal(_)) { + return Err("invalid project path".to_owned()); + } + current.push(part); + if fs::symlink_metadata(¤t).is_ok_and(|m| m.file_type().is_symlink()) { + return Err(format!( + "symlinked project path is not allowed for Record storage: {}", + current.display() + )); + } + } + Ok(()) +} + +pub fn read_record(path: &Path) -> Result, String> { + if !fs::symlink_metadata(path) + .map_err(|e| format!("{}: {e}", path.display()))? + .is_file() + { + return Err(format!("Record must be a regular file: {}", path.display())); + } + let mut bytes = Vec::new(); + File::open(path) + .map_err(|e| e.to_string())? + .take(crate::record_storage::MAX_EXPANDED_BYTES as u64 + 1) + .read_to_end(&mut bytes) + .map_err(|e| e.to_string())?; + if bytes.len() > crate::record_storage::MAX_EXPANDED_BYTES { + return Err("Record exceeds storage limit".to_owned()); + } + Ok(bytes) +} + +/// Publish a fully synced sibling, preserving the previous file on failure. +/// A cooperative caller holds the relevant lock for compare-and-replace. +pub fn atomic_write(path: &Path, bytes: &[u8], create_new: bool) -> Result<(), String> { + let parent = path.parent().ok_or("missing file parent")?; + fs::create_dir_all(parent).map_err(|e| e.to_string())?; + let temporary = parent.join(format!( + ".adf-storage-{}-{}.tmp", + std::process::id(), + SEQUENCE.fetch_add(1, Ordering::Relaxed) + )); + let result = (|| { + let mut file = OpenOptions::new() + .write(true) + .create_new(true) + .open(&temporary) + .map_err(|e| e.to_string())?; + if let Ok(metadata) = fs::metadata(path) { + file.set_permissions(metadata.permissions()) + .map_err(|e| e.to_string())?; + } + file.write_all(bytes) + .and_then(|()| file.sync_all()) + .map_err(|e| e.to_string())?; + drop(file); + if read_record(&temporary)? != bytes { + return Err("temporary Record verification failed".to_owned()); + } + if create_new { + fs::hard_link(&temporary, path).map_err(|e| { + format!( + "record already exists or cannot be published: {}: {e}", + path.display() + ) + })?; + fs::remove_file(&temporary).map_err(|e| e.to_string())?; + } else { + crate::project_setup::replace_file(&temporary, path).map_err(|e| e.to_string())?; + } + crate::project_setup::sync_directory(parent).map_err(|e| e.to_string()) + })(); + let _ = fs::remove_file(&temporary); + result +} + +pub fn correction_lock(root: &Path) -> Result { + let path = root.join(".adf/cache/locks/result-corrections.lock"); + reject_symlinks(root, &path)?; + fs::create_dir_all(path.parent().ok_or("missing lock parent")?).map_err(|e| e.to_string())?; + let file = OpenOptions::new() + .read(true) + .write(true) + .create(true) + .truncate(false) + .open(path) + .map_err(|e| e.to_string())?; + file.lock().map_err(|e| e.to_string())?; + Ok(file) +} diff --git a/src/storage_migration.rs b/src/storage_migration.rs new file mode 100644 index 0000000..7a457ae --- /dev/null +++ b/src/storage_migration.rs @@ -0,0 +1,615 @@ +//! Explicit, resumable changes of physical Record representation. +//! No Git index/history changes, implicit upgrades, or logical Record edits. + +use crate::canonical_digest; +use crate::delivery::{read_framework_lock, resolve_verified_release}; +use crate::project_config::load_project_config; +use crate::record_storage::{self, RecordKind, StoragePolicy}; +use crate::schema::SchemaRegistry; +use crate::storage_io::{self, StorageGuard}; +use serde::{Deserialize, Serialize}; +use serde_json::{Value, json}; +use sha2::{Digest, Sha256}; +use std::collections::{BTreeMap, BTreeSet}; +use std::fs; +use std::path::{Path, PathBuf}; +use std::process::Command; + +const JOURNAL: &str = ".adf/local/storage-migrations/current.json"; + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +struct Entry { + path: String, + kind: String, + id: String, + before_hash: String, + after_hash: String, + logical_digest: String, + before_bytes: usize, + after_bytes: usize, +} + +#[derive(Debug, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +struct Journal { + schema_version: String, + target: String, + complete: bool, + entries: Vec, +} + +fn hash(bytes: &[u8]) -> String { + format!("sha256:{:x}", Sha256::digest(bytes)) +} + +fn kind(name: &str) -> Result { + match name { + "result" => Ok(RecordKind::Result), + "evidence" => Ok(RecordKind::Evidence), + _ => Err("unsupported Record kind".to_owned()), + } +} + +fn regular_entries(root: &Path, directory: &Path) -> Result, String> { + storage_io::reject_symlinks(root, directory)?; + if !directory.exists() { + return Ok(Vec::new()); + } + let mut paths = fs::read_dir(directory) + .map_err(|e| e.to_string())? + .map(|entry| entry.map(|e| e.path()).map_err(|e| e.to_string())) + .collect::, _>>()?; + paths.sort(); + Ok(paths) +} + +fn inventory(root: &Path) -> Result, String> { + let mut records = Vec::new(); + for change in regular_entries(root, &root.join(".adf/changes"))? { + storage_io::reject_symlinks(root, &change)?; + if !change.is_dir() { + continue; + } + for (name, record_kind) in [ + ("results", RecordKind::Result), + ("evidence", RecordKind::Evidence), + ] { + for path in regular_entries(root, &change.join(name))? { + storage_io::reject_symlinks(root, &path)?; + if path.extension().is_some_and(|s| s == "json") { + if !path.is_file() { + return Err(format!("not a regular Record: {}", path.display())); + } + records.push((path, record_kind)); + } else if path.is_dir() { + return Err(format!( + "nested Record directory is not supported: {}", + path.display() + )); + } + } + } + } + Ok(records) +} + +fn skipped_files(root: &Path) -> Result, String> { + let mut skipped = Vec::new(); + for change in regular_entries(root, &root.join(".adf/changes"))? { + if !change.is_dir() { + continue; + } + for name in ["results", "evidence"] { + for path in regular_entries(root, &change.join(name))? { + if path.extension().is_none_or(|s| s != "json") { + skipped.push(json!({"path":path.strip_prefix(root).map_err(|e| e.to_string())?.to_str().ok_or("non-UTF8 skipped file path")?, + "bytes":fs::symlink_metadata(&path).map_err(|e| e.to_string())?.len(), + "reason": if path.file_name().and_then(|n| n.to_str()).is_some_and(|n| n.starts_with(".adf-storage-")) { "temporary-file-review-required" } else { "not-a-json-record" }})); + } + } + } + } + Ok(skipped) +} + +fn converted( + raw: &[u8], + record_kind: RecordKind, + target: StoragePolicy, +) -> Result<(Value, Vec), String> { + let value = record_storage::decode(raw, record_kind)?; + let encoded = record_storage::encode(&value, record_kind, target)?; + let bytes = if target == StoragePolicy::Adaptive && encoded.len() >= raw.len() { + raw.to_vec() + } else { + encoded + }; + if record_storage::decode(&bytes, record_kind)? != value { + return Err("Record storage roundtrip mismatch".to_owned()); + } + Ok((value, bytes)) +} + +fn scan( + root: &Path, + registry: &SchemaRegistry, + target: StoragePolicy, +) -> Result, String> { + let mut entries = Vec::new(); + let mut ids = BTreeSet::new(); + for (path, record_kind) in inventory(root)? { + let raw = storage_io::read_record(&path)?; + let (value, bytes) = + converted(&raw, record_kind, target).map_err(|e| format!("{}: {e}", path.display()))?; + registry + .validate(record_kind.as_str(), &value) + .map_err(|e| format!("{}: {e}", path.display()))?; + let id = value["id"] + .as_str() + .ok_or("Record id is missing")? + .to_owned(); + if !ids.insert(id.clone()) { + return Err(format!("duplicate Record id: {id}")); + } + let change_id = path + .parent() + .and_then(Path::parent) + .and_then(Path::file_name) + .and_then(|v| v.to_str()) + .ok_or("invalid Record path")?; + if value["change_id"].as_str() != Some(change_id) { + return Err(format!( + "Record change_id disagrees with path: {}", + path.display() + )); + } + entries.push(Entry { + path: path + .strip_prefix(root) + .map_err(|e| e.to_string())? + .to_str() + .ok_or("non-UTF8 Record path")? + .to_owned(), + kind: record_kind.as_str().to_owned(), + id, + before_hash: hash(&raw), + after_hash: hash(&bytes), + logical_digest: canonical_digest(&value).map_err(|e| e.to_string())?, + before_bytes: raw.len(), + after_bytes: bytes.len(), + }); + } + Ok(entries) +} + +fn ignored(root: &Path, relative: &str) -> bool { + Command::new("git") + .arg("-C") + .arg(root) + .args(["check-ignore", "-q", "--", relative]) + .status() + .is_ok_and(|status| status.success()) +} + +fn journal_path(root: &Path) -> Result { + let path = root.join(JOURNAL); + storage_io::reject_symlinks(root, &path)?; + Ok(path) +} + +fn validate_pending(root: &Path, entries: &[Entry]) -> Result<(), String> { + let path = journal_path(root)?; + if !path.exists() { + return Ok(()); + } + let previous: Journal = serde_json::from_value(record_storage::parse_strict( + &storage_io::read_record(&path)?, + )?) + .map_err(|e| e.to_string())?; + if previous.schema_version != "1" { + return Err("unsupported storage migration journal".to_owned()); + } + if previous.complete { + return Ok(()); + } + let current: BTreeMap<_, _> = entries + .iter() + .map(|entry| (entry.path.as_str(), entry)) + .collect(); + if current.len() != previous.entries.len() { + return Err( + "Records changed during interrupted storage migration; inspect before resuming" + .to_owned(), + ); + } + for prior in previous.entries { + let entry = current + .get(prior.path.as_str()) + .ok_or("Record paths changed during interrupted storage migration")?; + if (entry.before_hash != prior.before_hash && entry.before_hash != prior.after_hash) + || entry.logical_digest != prior.logical_digest + || entry.id != prior.id + || entry.kind != prior.kind + { + return Err(format!( + "Record changed during interrupted storage migration: {}", + prior.path + )); + } + } + Ok(()) +} + +fn save_journal(root: &Path, journal: &Journal) -> Result<(), String> { + if !ignored(root, JOURNAL) { + let ignore = root.join(".adf/local/storage-migrations/.gitignore"); + storage_io::reject_symlinks(root, &ignore)?; + if ignore.exists() { + return Err( + "storage migration journal must be Git-ignored; review its existing .gitignore" + .to_owned(), + ); + } + storage_io::atomic_write(&ignore, b"*\n", true)?; + if !ignored(root, JOURNAL) { + return Err("storage migration journal must be Git-ignored".to_owned()); + } + } + storage_io::atomic_write( + &journal_path(root)?, + &serde_json::to_vec(journal).map_err(|e| e.to_string())?, + false, + ) +} + +fn switch_config(root: &Path, target: StoragePolicy) -> Result<(), String> { + let path = root.join(".adf/config.yaml"); + storage_io::reject_symlinks(root, &path)?; + let old = fs::read(&path).map_err(|e| e.to_string())?; + let mut value: Value = serde_yaml::from_slice(&old).map_err(|e| e.to_string())?; + let config = value.as_object_mut().ok_or("invalid project config")?; + config.insert( + "schema_version".to_owned(), + json!(if target == StoragePolicy::Adaptive { + "2" + } else { + "1" + }), + ); + if target == StoragePolicy::Adaptive { + config.insert("record_storage".to_owned(), json!(target.as_str())); + } else { + config.remove("record_storage"); + } + let bytes = serde_yaml::to_string(&value).map_err(|e| e.to_string())?; + if fs::read(&path).map_err(|e| e.to_string())? != old { + return Err("project config changed before migration".to_owned()); + } + storage_io::atomic_write(&path, bytes.as_bytes(), false) +} + +fn available_bytes(root: &Path) -> Result { + let output = Command::new("df") + .arg("-Pk") + .arg(root) + .output() + .map_err(|e| e.to_string())?; + if !output.status.success() { + return Err("cannot determine free space before migration".to_owned()); + } + String::from_utf8_lossy(&output.stdout) + .lines() + .last() + .and_then(|line| line.split_whitespace().nth(3)) + .and_then(|v| v.parse::().ok()) + .and_then(|v| v.checked_mul(1024)) + .ok_or_else(|| "cannot parse free space before migration".to_owned()) +} + +fn apply( + root: &Path, + target: StoragePolicy, + entries: Vec, + stop_after: Option, +) -> Result<(), String> { + validate_pending(root, &entries)?; + let mut journal = Journal { + schema_version: "1".to_owned(), + target: target.as_str().to_owned(), + complete: false, + entries, + }; + let journal_bytes = serde_json::to_vec(&journal) + .map_err(|e| e.to_string())? + .len() as u64; + let growth: u64 = journal + .entries + .iter() + .map(|e| e.after_bytes.saturating_sub(e.before_bytes) as u64) + .sum(); + let largest = journal + .entries + .iter() + .map(|e| e.after_bytes as u64) + .max() + .unwrap_or(0); + if available_bytes(root)? < growth + largest + journal_bytes * 2 + 65536 { + return Err("insufficient free space for storage migration".to_owned()); + } + save_journal(root, &journal)?; + // Keep v2 enabled throughout either direction, including interrupted rollback. + switch_config(root, StoragePolicy::Adaptive)?; + if stop_after == Some(0) { + return Err("simulated storage migration interruption".to_owned()); + } + for (i, entry) in journal.entries.iter().enumerate() { + let path = root.join(&entry.path); + storage_io::reject_symlinks(root, &path)?; + let raw = storage_io::read_record(&path)?; + if hash(&raw) != entry.before_hash { + return Err(format!("Record changed before replacement: {}", entry.path)); + } + let (value, bytes) = converted(&raw, kind(&entry.kind)?, target)?; + if hash(&bytes) != entry.after_hash + || canonical_digest(&value).map_err(|e| e.to_string())? != entry.logical_digest + { + return Err("migration plan changed before replacement".to_owned()); + } + if storage_io::read_record(&path)? != raw { + return Err(format!("Record changed before replacement: {}", entry.path)); + } + if raw != bytes { + storage_io::atomic_write(&path, &bytes, false)?; + } + if stop_after == Some(i + 1) { + return Err("simulated storage migration interruption".to_owned()); + } + } + // Verify all bytes after replacement before lowering the reader requirement. + for entry in &journal.entries { + let raw = storage_io::read_record(&root.join(&entry.path))?; + if hash(&raw) != entry.after_hash { + return Err("Record changed after migration".to_owned()); + } + let restored = record_storage::decode(&raw, kind(&entry.kind)?)?; + if canonical_digest(&restored).map_err(|e| e.to_string())? != entry.logical_digest { + return Err("Record verification failed after migration".to_owned()); + } + } + if inventory(root)?.len() != journal.entries.len() { + return Err("Record inventory changed during migration".to_owned()); + } + switch_config(root, target)?; + journal.complete = true; + save_journal(root, &journal) +} + +/// CLI-only explicit maintenance interface; ordinary project operations never migrate. +pub fn run_cli(arguments: &[String]) -> Result { + let operation = arguments + .first() + .map(String::as_str) + .ok_or("expected storage inspect, verify, export or migrate")?; + if !matches!(operation, "inspect" | "verify" | "export" | "migrate") { + return Err("unsupported storage operation".to_owned()); + } + let mut options = BTreeMap::new(); + let mut dry_run = false; + let mut i = 1; + while i < arguments.len() { + let flag = &arguments[i]; + if flag == "--dry-run" { + if dry_run { + return Err("duplicate --dry-run".to_owned()); + } + dry_run = true; + i += 1; + continue; + } + if !matches!( + flag.as_str(), + "--project" | "--to" | "--release" | "--record" | "--format" + ) { + return Err(format!("unsupported storage argument: {flag}")); + } + let value = arguments + .get(i + 1) + .filter(|v| !v.starts_with("--")) + .ok_or_else(|| format!("missing value for {flag}"))?; + if options.insert(flag.as_str(), value.as_str()).is_some() { + return Err(format!("duplicate argument: {flag}")); + } + i += 2; + } + if options.get("--format").is_some_and(|v| *v != "json") { + return Err("storage commands support --format json".to_owned()); + } + if operation != "migrate" && (dry_run || options.contains_key("--to")) { + return Err("--to and --dry-run require storage migrate".to_owned()); + } + if operation != "export" && options.contains_key("--record") { + return Err("--record requires storage export".to_owned()); + } + let root = Path::new(options.get("--project").copied().unwrap_or(".")) + .canonicalize() + .map_err(|e| e.to_string())?; + if !root.join(".adf/config.yaml").is_file() { + return Err("project is not initialized".to_owned()); + } + let _guard = if operation == "migrate" && !dry_run { + StorageGuard::exclusive(&root)? + } else { + StorageGuard::shared(&root)? + }; + let config = load_project_config(&root).map_err(|e| e.to_string())?; + let lock = read_framework_lock(&root.join(".adf/framework.lock")).map_err(|e| e.to_string())?; + let release = resolve_verified_release(&root, &lock, options.get("--release").map(Path::new)) + .map_err(|e| e.to_string())?; + let target = match options + .get("--to") + .copied() + .unwrap_or("adaptive-refmaps-v1") + { + "plain-json-v1" => StoragePolicy::Plain, + "adaptive-refmaps-v1" => StoragePolicy::Adaptive, + _ => return Err("unsupported target storage policy".to_owned()), + }; + if operation == "migrate" && !options.contains_key("--to") { + return Err("storage migrate requires --to".to_owned()); + } + let entries = scan(&root, &release.schema_registry, target)?; + validate_pending(&root, &entries)?; + if operation == "export" { + let id = options + .get("--record") + .ok_or("storage export requires --record")?; + let entry = entries + .iter() + .find(|entry| entry.id == *id) + .ok_or("Record does not exist")?; + return record_storage::decode( + &storage_io::read_record(&root.join(&entry.path))?, + kind(&entry.kind)?, + ); + } + let before: usize = entries.iter().map(|e| e.before_bytes).sum(); + let after: usize = entries.iter().map(|e| e.after_bytes).sum(); + let changed = entries + .iter() + .filter(|e| e.before_hash != e.after_hash) + .count(); + let skipped = skipped_files(&root)?; + let growth: u64 = entries + .iter() + .map(|e| e.after_bytes.saturating_sub(e.before_bytes) as u64) + .sum(); + let temporary_record_bytes = entries.iter().map(|e| e.after_bytes).max().unwrap_or(0); + let journal_estimate = serde_json::to_vec(&entries) + .map_err(|e| e.to_string())? + .len() as u64 + + 1024; + let required_free_bytes = growth + temporary_record_bytes as u64 + journal_estimate * 2 + 65536; + let mut report = json!({"schema_version":"1", "operation":operation, "dry_run":dry_run, + "current_policy":config.record_storage.as_str(), "target_policy":target.as_str(), + "records":entries.len(), "changed_records":changed, "current_bytes":before, "proposed_bytes":after, + "saved_bytes": before as i64 - after as i64, "logical_records_verified":true, "applied":false, "skipped_files":skipped, + "temporary_record_bytes":temporary_record_bytes, "required_free_bytes":required_free_bytes}); + if operation == "migrate" && !dry_run { + apply(&root, target, entries, None)?; + report["applied"] = json!(true); + } + Ok(report) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn interrupted_migration_can_resume_or_reverse_without_original_copies() { + let root = + std::env::temp_dir().join(format!("adf-storage-migration-{}", std::process::id())); + let _ = fs::remove_dir_all(&root); + fs::create_dir_all(root.join(".adf/changes/change.example/results")).unwrap(); + let root = root.canonicalize().unwrap(); + assert!( + Command::new("git") + .args(["init", "-q"]) + .arg(&root) + .status() + .unwrap() + .success() + ); + fs::write(root.join(".gitignore"), ".adf/local/\n.adf/cache/\n").unwrap(); + fs::write(root.join(".adf/config.yaml"), "schema_version: '1'\nproject_sources:\n contracts: contracts\n decisions: decisions\nrepository_observation: .adf/repository-observation.yaml\n").unwrap(); + let refs: BTreeMap<_, _> = (0..20) + .map(|i| (format!("code.{i}"), format!("sha256:{}", "a".repeat(64)))) + .collect(); + let mut entries = Vec::new(); + for i in 0..2 { + let value = json!({"id":format!("result.{i}"), "change_id":"change.example", "input_refs":refs, "freshness_refs":refs}); + let raw = serde_json::to_vec_pretty(&value).unwrap(); + let (_, after) = converted(&raw, RecordKind::Result, StoragePolicy::Adaptive).unwrap(); + let path = format!(".adf/changes/change.example/results/{i}.json"); + fs::write(root.join(&path), &raw).unwrap(); + entries.push(Entry { + path, + kind: "result".to_owned(), + id: format!("result.{i}"), + before_hash: hash(&raw), + after_hash: hash(&after), + logical_digest: canonical_digest(&value).unwrap(), + before_bytes: raw.len(), + after_bytes: after.len(), + }); + } + let guard = StorageGuard::exclusive(&root).unwrap(); + assert!(StorageGuard::shared(&root).is_err()); + assert!(apply(&root, StoragePolicy::Adaptive, entries.clone(), Some(1)).is_err()); + assert_eq!( + load_project_config(&root).unwrap().record_storage, + StoragePolicy::Adaptive + ); + let mut reverse = entries.clone(); + for entry in &mut reverse { + let raw = storage_io::read_record(&root.join(&entry.path)).unwrap(); + let (_, bytes) = converted(&raw, RecordKind::Result, StoragePolicy::Plain).unwrap(); + entry.before_hash = hash(&raw); + entry.before_bytes = raw.len(); + entry.after_hash = hash(&bytes); + entry.after_bytes = bytes.len(); + } + validate_pending(&root, &reverse).unwrap(); + let mut altered = reverse.clone(); + altered[0].before_hash = "sha256:external-edit".to_owned(); + assert!(validate_pending(&root, &altered).is_err()); + let mut resumed = reverse.clone(); + for entry in &mut resumed { + let raw = storage_io::read_record(&root.join(&entry.path)).unwrap(); + let (_, bytes) = converted(&raw, RecordKind::Result, StoragePolicy::Adaptive).unwrap(); + entry.after_hash = hash(&bytes); + entry.after_bytes = bytes.len(); + } + apply(&root, StoragePolicy::Adaptive, resumed, None).unwrap(); + for entry in &mut reverse { + let raw = storage_io::read_record(&root.join(&entry.path)).unwrap(); + let (_, bytes) = converted(&raw, RecordKind::Result, StoragePolicy::Plain).unwrap(); + entry.before_hash = hash(&raw); + entry.before_bytes = raw.len(); + entry.after_hash = hash(&bytes); + entry.after_bytes = bytes.len(); + } + apply(&root, StoragePolicy::Plain, reverse, Some(0)).unwrap_err(); + let mut rollback = entries.clone(); + for entry in &mut rollback { + let raw = storage_io::read_record(&root.join(&entry.path)).unwrap(); + let (_, bytes) = converted(&raw, RecordKind::Result, StoragePolicy::Plain).unwrap(); + entry.before_hash = hash(&raw); + entry.before_bytes = raw.len(); + entry.after_hash = hash(&bytes); + entry.after_bytes = bytes.len(); + } + apply(&root, StoragePolicy::Plain, rollback, None).unwrap(); + assert_eq!( + load_project_config(&root).unwrap().record_storage, + StoragePolicy::Plain + ); + for entry in entries { + let raw = storage_io::read_record(&root.join(entry.path)).unwrap(); + assert!( + record_storage::parse_strict(&raw) + .unwrap() + .get("storage_format") + .is_none() + ); + assert_eq!( + canonical_digest(&record_storage::decode(&raw, RecordKind::Result).unwrap()) + .unwrap(), + entry.logical_digest + ); + } + drop(guard); + fs::remove_dir_all(root).unwrap(); + } +} diff --git a/testdata/storage/v1/result.logical.json b/testdata/storage/v1/result.logical.json new file mode 100644 index 0000000..f7575d2 --- /dev/null +++ b/testdata/storage/v1/result.logical.json @@ -0,0 +1 @@ +{"action_id":"action.example","change_id":"change.example","context_digest":"sha256:0000000000000000000000000000000000000000000000000000000000000000","execution":{"context_bytes":128},"freshness_refs":{"code.example.0":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","code.example.1":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","code.example.10":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","code.example.11":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","code.example.12":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","code.example.13":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","code.example.14":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","code.example.15":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","code.example.2":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","code.example.3":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","code.example.4":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","code.example.5":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","code.example.6":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","code.example.7":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","code.example.8":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","code.example.9":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"},"id":"result.387e1487ec27151dc41b","input_refs":{"code.example.0":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","code.example.1":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","code.example.10":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","code.example.11":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","code.example.12":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","code.example.13":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","code.example.14":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","code.example.15":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","code.example.2":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","code.example.3":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","code.example.4":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","code.example.5":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","code.example.6":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","code.example.7":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","code.example.8":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","code.example.9":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"},"output_refs":[],"payload":{"summary":"Preserve the recorded implementation and input references."},"result_schema":"result.build","role":"Builder","schema_version":"1"} diff --git a/testdata/storage/v1/result.stored.json b/testdata/storage/v1/result.stored.json new file mode 100644 index 0000000..504b2d7 --- /dev/null +++ b/testdata/storage/v1/result.stored.json @@ -0,0 +1 @@ +{"map_bindings":{"/freshness_refs":"sha256:456f4a72704603e6806b05c942264067ca4feb6ef37a5325b8fff4d2aeeb2359","/input_refs":"sha256:456f4a72704603e6806b05c942264067ca4feb6ef37a5325b8fff4d2aeeb2359"},"record":{"action_id":"action.example","change_id":"change.example","context_digest":"sha256:0000000000000000000000000000000000000000000000000000000000000000","execution":{"context_bytes":128},"id":"result.387e1487ec27151dc41b","output_refs":[],"payload":{"summary":"Preserve the recorded implementation and input references."},"result_schema":"result.build","role":"Builder","schema_version":"1"},"record_digest":"sha256:af6693bddb2b05e4090a08d6d3a8aac410a8c993245cc4156c5a7fddf7355eaa","record_kind":"result","reference_maps":{"sha256:456f4a72704603e6806b05c942264067ca4feb6ef37a5325b8fff4d2aeeb2359":{"code.example.0":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","code.example.1":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","code.example.10":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","code.example.11":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","code.example.12":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","code.example.13":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","code.example.14":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","code.example.15":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","code.example.2":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","code.example.3":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","code.example.4":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","code.example.5":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","code.example.6":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","code.example.7":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","code.example.8":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","code.example.9":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"}},"storage_format":"adf-record-refmaps-v1"} diff --git a/tests/cli.rs b/tests/cli.rs index c6e208a..95a7e68 100644 --- a/tests/cli.rs +++ b/tests/cli.rs @@ -3360,7 +3360,14 @@ fn derived_runtime_indexes_are_rebuilt_when_cache_files_are_corrupt() { for relative in cache_paths { let value: Value = serde_json::from_slice(&fs::read(project.root.join(relative)).unwrap()).unwrap(); - assert_eq!(value["schema_version"], "1"); + assert_eq!( + value["schema_version"], + if relative.contains("contract-health-") { + "2" + } else { + "1" + } + ); } } @@ -3973,6 +3980,226 @@ fn mcp_call( mcp_receive(output)["result"].clone() } +#[test] +fn storage_migration_preserves_restart_health_execution_and_correction() { + use adf::project_application::ProjectApplicationService; + use adf::record_storage::{self, RecordKind, StoragePolicy}; + let project = TestProject::new(); + let mut service = ProjectApplicationService::new(&project.root, None).unwrap(); + let issued = service.next("change.place-order", false).unwrap(); + let key = issued.issued_action.unwrap(); + let begun = service + .begin_execution( + &key, + serde_json::from_value(json!({ + "runner":{"provider":"test", "surface":"fixture"} + })) + .unwrap(), + ) + .unwrap(); + let submission = risk_signal_submission( + &serde_json::to_value(&key).unwrap(), + &issued.next_response["context"]["payload"], + ); + service + .submit(&key, submission["payload"].clone(), vec![], None) + .unwrap(); + let results = project.root.join(".adf/changes/change.place-order/results"); + let path = fs::read_dir(&results) + .unwrap() + .next() + .unwrap() + .unwrap() + .path(); + let mut original: Value = serde_json::from_slice(&fs::read(&path).unwrap()).unwrap(); + // Historical Records explicitly repeated inherited outcome maps. + for field in ["input_refs", "freshness_refs"] { + let inherited = original[field].clone(); + for outcome in original["payload"]["outcomes"].as_array_mut().unwrap() { + outcome + .as_object_mut() + .unwrap() + .entry(field) + .or_insert_with(|| inherited.clone()); + } + } + let mut identity = original.clone(); + identity.as_object_mut().unwrap().remove("execution"); + identity.as_object_mut().unwrap().remove("id"); + original["id"] = json!(format!( + "result.{}", + &canonical_digest(&identity).unwrap()[7..27] + )); + fs::write(&path, serde_json::to_vec_pretty(&original).unwrap()).unwrap(); + let record_id = original["id"].as_str().unwrap(); + drop(service); + let next_before = project.run(&["next", "change.place-order", "--format", "json"]); + let health_before = project.run(&["contract-health", "--format", "json"]); + assert_success(&next_before); + assert_success(&health_before); + let before_bytes = fs::read(&path).unwrap(); + let before_config = fs::read(project.root.join(".adf/config.yaml")).unwrap(); + let head = git_output(&project.root, &["rev-parse", "HEAD"]); + let dry = project.run(&[ + "project", + "storage", + "migrate", + "--to", + "adaptive-refmaps-v1", + "--dry-run", + ]); + assert_success(&dry); + assert_eq!(fs::read(&path).unwrap(), before_bytes); + assert_eq!( + fs::read(project.root.join(".adf/config.yaml")).unwrap(), + before_config + ); + assert!( + !project + .root + .join(".adf/local/storage-migrations/current.json") + .exists() + ); + let migrated = project.run(&[ + "project", + "storage", + "migrate", + "--to", + "adaptive-refmaps-v1", + ]); + assert_success(&migrated); + let packed = fs::read(&path).unwrap(); + assert_eq!( + record_storage::parse_strict(&packed).unwrap()["storage_format"], + record_storage::STORAGE_FORMAT + ); + assert_eq!( + record_storage::decode(&packed, RecordKind::Result).unwrap(), + original + ); + assert_eq!(git_output(&project.root, &["rev-parse", "HEAD"]), head); + assert_eq!( + read_yaml(&project.root.join(".adf/config.yaml"))["schema_version"], + "2" + ); + for (command, before) in [("next", &next_before), ("contract-health", &health_before)] { + let args = if command == "next" { + vec![command, "change.place-order", "--format", "json"] + } else { + vec![command, "--format", "json"] + }; + let after = project.run(&args); + assert_success(&after); + assert_eq!( + serde_json::from_slice::(&before.stdout).unwrap(), + serde_json::from_slice::(&after.stdout).unwrap() + ); + } + fs::remove_dir_all(project.root.join(".adf/cache/runtime")).unwrap(); + let cold = project.run(&["contract-health", "--format", "json"]); + assert_success(&cold); + assert_eq!( + serde_json::from_slice::(&health_before.stdout).unwrap(), + serde_json::from_slice::(&cold.stdout).unwrap() + ); + let exported = project.run(&["project", "storage", "export", "--record", record_id]); + assert_success(&exported); + assert_eq!( + serde_json::from_slice::(&exported.stdout).unwrap(), + original + ); + let completed = project.run(&[ + "execution", + "complete", + &begun.execution_id, + "--change", + "change.place-order", + "--status", + "succeeded", + "--result", + record_id, + "--format", + "json", + ]); + assert_success(&completed); + let (mut child, mut input, mut output) = start_mcp_server(&project.root); + let restarted = mcp_call( + &mut input, + &mut output, + 2, + "adf_next", + json!({"change_id":"change.place-order"}), + ); + assert_eq!(restarted["isError"], false, "{restarted}"); + drop(input); + assert!(child.wait().unwrap().success()); + let source = project.root.join("src/place_order.py"); + let mut text = fs::read_to_string(&source).unwrap(); + text.push_str("\n# changed after verification\n"); + fs::write(&source, text).unwrap(); + let stale_packed = project.run(&["next", "change.place-order", "--format", "json"]); + assert_success(&stale_packed); + assert_ne!( + serde_json::from_slice::(&stale_packed.stdout).unwrap()["context"]["digest"], + serde_json::from_slice::(&next_before.stdout).unwrap()["context"]["digest"] + ); + let restored = project.run(&["project", "storage", "migrate", "--to", "plain-json-v1"]); + assert_success(&restored); + assert_eq!( + read_yaml(&project.root.join(".adf/config.yaml"))["schema_version"], + "1" + ); + assert_eq!( + serde_json::from_slice::(&fs::read(&path).unwrap()).unwrap(), + original + ); + let stale_plain = project.run(&["next", "change.place-order", "--format", "json"]); + assert_success(&stale_plain); + assert_eq!( + serde_json::from_slice::(&stale_packed.stdout).unwrap(), + serde_json::from_slice::(&stale_plain.stdout).unwrap() + ); + assert_success(&project.run(&[ + "project", + "storage", + "migrate", + "--to", + "adaptive-refmaps-v1", + ])); + let schemas = + adf::schema::SchemaRegistry::load(project.release_root.join("schemas/v1")).unwrap(); + let mut store = + adf::filesystem_project::FileProjectStore::open(&project.root, json!({}), &schemas) + .unwrap(); + assert!(store.append_result(&original).is_err()); + let mut corrected = original.clone(); + corrected["id"] = json!("result.corrected"); + store.replace_result(&corrected, record_id).unwrap(); + assert!(store.replace_result(&original, record_id).is_err()); + assert_eq!( + record_storage::decode(&fs::read(&path).unwrap(), RecordKind::Result).unwrap(), + corrected + ); + let exclusive = project.run(&["project", "storage", "migrate", "--to", "plain-json-v1"]); + assert!(!exclusive.status.success()); + assert!(String::from_utf8_lossy(&exclusive.stderr).contains("storage is busy")); + drop(store); + let mut additional = original.clone(); + additional["id"] = json!("result.additional"); + additional["action_id"] = json!("action.additional"); + let mut store = + adf::filesystem_project::FileProjectStore::open(&project.root, json!({}), &schemas) + .unwrap(); + store.append_result(&additional).unwrap(); + drop(store); + let report = project.run(&["project", "storage", "verify"]); + assert_success(&report); + assert_eq!( + record_storage::encode(&corrected, RecordKind::Result, StoragePolicy::Adaptive).unwrap(), + fs::read(&path).unwrap() + ); +} + fn risk_signal_submission(key: &Value, context: &Value) -> Value { let reviewed_candidates = context["signal_candidates"] .as_array()