diff --git a/README.md b/README.md index 60ff81d..fa2d425 100644 --- a/README.md +++ b/README.md @@ -245,7 +245,7 @@ adf mcp --project /path/to/project │ ├─ Builder implement, then record required clause evidence │ └─ Challenger try to falsify it, before and after the build │ │ - └───────────┴─ adf_submit ──── validated, stored, reevaluated + └───────────┴─ adf_submit ──── validated and stored │ ▼ ready to merge @@ -283,15 +283,28 @@ While impact assessment is still pending, `next` and `explain` select that action before deriving repository-wide Contract health. Unrelated Result and Evidence history is not loaded for this first step. -When Contract health is required, ADF indexes Evidence and verification Results -once, then validates and hashes only records that can affect a Contract clause. -It does not repeatedly scan every Result for every clause. +When Contract health is required, ADF maintains persistent Evidence and Result +indexes under `.adf/cache/runtime/`. Unchanged tracked records are identified by +their Git blob IDs; changed and untracked records use content hashes. ADF parses +only records whose source identity changed, and it does not scan every Result +again for each Contract clause. Corrupt cache entries are rebuilt from source. +ADF writes runtime caches only when Git confirms that the cache path is ignored. + +Repository observation uses a separate cache tied to the current revision, +analysis configuration, signal catalog, and source identities. Source changes +invalidate the observation without making cache files authoritative. New Results store shared input and freshness references once at the Result level. An outcome carries its own references only when they differ from those shared values. Existing Results remain readable and are not rewritten, so their identities and downstream freshness checks remain stable. +`adf_submit` returns after the Result is stored. Its response includes +`result_id`, `already_completed`, `next_required`, and per-stage `timings_ms`. +Call `adf_next` separately when `next_required` is true. This keeps a slow next +evaluation from obscuring whether submission itself succeeded. `adf_next` also +reports per-stage timings. + Each action also carries advisory execution guidance. Impact assessment normally recommends an economy model, while challenge recommends a high-accuracy model. The listed escalation conditions tell an orchestrator when diff --git a/docs/MCP-DESIGN.md b/docs/MCP-DESIGN.md index 8e679be..aa6cb05 100644 --- a/docs/MCP-DESIGN.md +++ b/docs/MCP-DESIGN.md @@ -154,7 +154,7 @@ registered_output_refs: [] - RegistryはMCP sessionのmemoryにだけ保持する。memoryにないActionは、正本の再評価と一致するときに再構成する。再構成後に`adf_next`を呼んだ場合も、同じChangeの正本に保存済みのEvidence、Decision、Contractは提出時の出力として参照できる。 - Generated ContextをGit、Result、derived cacheへ保存しない。 - Evidence、Decision、Contractの専用Toolが保存したrefだけを`registered_output_refs`へ追加する。Evidenceには、発行済みRequirementが参照した入力digestを`input_refs`として付与する。 -- 正常な`submit`後にexact keyを消費し、再評価で返した次Actionを新しいentryとして登録する。 +- 正常な`submit`後にexact keyを消費する。次Actionは後続の`adf_next`で発行し、新しいentryとして登録する。 - submit失敗時は、修正して再試行できるようentryを残す。 - 同じkeyへ複数提出が競合した場合、Filesystem Storeのexclusive createを最終防衛線とする。 @@ -186,7 +186,7 @@ Tool名は広いMCP client互換性を優先し、ASCII英数字とunderscoreだ | `adf_execution_log` | read | 保存済みResultと実行RecordからContextサイズと計測値を集計する | | `adf_begin_execution` | write | 現在のActionに対する外部実行の開始を追記する。Agentは起動しない | | `adf_complete_execution` | write | 外部実行の成否と確定済みの利用量を追記する | -| `adf_submit` | write | 発行済みActionのResultを検証・保存し、再評価する | +| `adf_submit` | write | 発行済みActionのResultを検証・保存する | | `adf_add_evidence` | write | 発行済みEvidence ActionへEvidenceを追記する | | `adf_apply_decision` | write | Human回答を解決するDecisionを保存する | | `adf_apply_contract` | write | Decisionを反映したContractを楽観的lock付きで更新する | @@ -204,7 +204,7 @@ MCP Tool annotationはHost向けhintとして設定しますが、認可には ## 10. Tool契約 全Toolは`inputSchema`と`outputSchema`を公開し、成功時は`structuredContent`を返します。 -保存Record Schemaとは別に、MCP I/O Schemaを`schemas/mcp/v1/`へ置きます。 +保存Record Schemaとは別に、MCP I/O Schemaを`schemas/mcp/`へ置きます。`adf_next`はv1、保存完了だけを返す`adf_submit`出力はv2です。 ### 10.1 `adf_next` @@ -227,6 +227,11 @@ Output: "change_id": "change.example", "action_id": "action.example", "context_digest": "sha256:..." + }, + "timings_ms": { + "repository_load": 12, + "evaluation": 34, + "total": 46 } } ``` @@ -261,13 +266,22 @@ Output: ```json { - "schema_version": "1", + "schema_version": "2", "result_id": "result.example", "already_completed": false, - "next_response": {} + "next_required": true, + "timings_ms": { + "repository_load": 12, + "change_snapshot": 3, + "validation_and_persist": 8, + "total": 23 + } } ``` +- `adf_submit`はResultを永続化した時点で応答し、次Actionを計算しない。 +- `next_required`が`true`なら、呼出し側は`adf_next`を別に実行する。 +- `timings_ms`は処理段階ごとの実測時間をミリ秒で返す。 - exact action keyがRegistryにない提出は、正本を再評価して同じActionが現在のものであるときだけ受理する。 - 成功したResult追記後だけActionを消費する。 - Contract、Decision、Evidenceの`output_refs`は、同じentryの`registered_output_refs`に存在しなければ拒否する。再起動を跨いだActionでは、そのRecordがChangeの正本に存在することで代える。 @@ -482,8 +496,11 @@ schemas/mcp/v1/ ├── next-input.schema.json ├── next-output.schema.json ├── submit-input.schema.json -├── submit-output.schema.json +├── submit-output.schema.json # 旧出力契約 └── tool-error.schema.json + +schemas/mcp/v2/ +└── submit-output.schema.json ``` RMCP、Tokio、schema生成用crateを追加する場合も、公開Schemaの正本はRepository上のJSON @@ -507,7 +524,7 @@ Schemaとし、生成差分をtestで検査します。 - 専用Toolを経由していないContract、Decision、Evidenceの`output_refs`を拒否 - Decision、Contract、EvidenceのAction binding - submit失敗時にIssued Actionを消費しない -- 成功時にexact Actionだけを消費し、返した次Actionを登録する +- 成功時にexact Actionだけを消費し、次の`adf_next`で新しいActionを登録する ### 18.3 MCP protocol integration @@ -517,7 +534,7 @@ Schemaとし、生成差分をtestで検査します。 2. tools/list 3. `adf_next` 4. `adf_submit` -5. 再評価後の次Action +5. `adf_next`で次Actionを取得 6. lifecycle全体を`ready-to-merge`まで実行 7. stdoutへJSON-RPC以外を出さない 8. server停止時に未提出Actionを失効 diff --git a/docs/concepts.ja.md b/docs/concepts.ja.md index 20bb3e6..8fbda80 100644 --- a/docs/concepts.ja.md +++ b/docs/concepts.ja.md @@ -31,7 +31,7 @@ Contractは変更のたびに増えます。障害から学んだことはテス ## 変更の流れ -変更は`adf change init`で作ります。以降は`adf next`が「次にやること」を1件ずつ返します。エージェントはそれを実行し、MCPの`adf_submit`で結果を提出して、また次を受け取ります。作業の順番を決めるのはエージェントではなく`adf`です。 +変更は`adf change init`で作ります。以降は`adf next`が「次にやること」を1件ずつ返します。エージェントはそれを実行し、MCPの`adf_submit`で結果を保存してから、`adf_next`で次の作業を受け取ります。作業の順番を決めるのはエージェントではなく`adf`です。 最初の作業は影響評価です。変更の目的から、影響がある対象とリスクを実装前に整理します。結果は次の3種類を明示します。 @@ -54,7 +54,7 @@ Contractは変更のたびに増えます。障害から学んだことはテス │ ├─ Builder 実装し、必要な条項の証拠を記録する │ └─ Challenger 実装の前後で反証する │ │ -└───────┴─ adf_submit ─── 結果を検証して保存し、状態を進める +└───────┴─ adf_submit ─── 結果を検証して保存する │ ▼ ready-to-merge @@ -97,10 +97,14 @@ Contractまたは条項の`evidence_mode`で、検証に掛ける費用を選べ 影響評価が済んでいない間、`adf next`と`adf explain`は、リポジトリ全体のContract検証状態を計算する前に影響評価を次の作業として選びます。この最初の段階では、無関係なChangeのResultとEvidenceを読み込みません。 -Contract検証状態が必要な場合は、Evidenceと検証Resultを一度索引化し、Contract条項へ影響する記録だけを検証してハッシュを計算します。条項ごとに全Resultを繰り返し検索しません。 +Contract検証状態が必要な場合は、Evidenceと検証Resultの永続索引を`.adf/cache/runtime/`に作ります。変更のない記録はGitのBlob ID、変更中または未追跡の記録は内容ハッシュで識別します。元の記録が変わった場合だけJSONを読み直し、条項ごとに全Resultを繰り返し検索しません。索引が壊れている場合は正本から作り直します。Gitがキャッシュ先を無視すると確認できた場合だけ、索引をファイルへ保存します。 + +Repository観測も、現在のrevision、解析設定、Signal Catalog、解析対象ファイルに結び付いたキャッシュを使います。解析対象が変われば無効になるため、キャッシュを正本として扱いません。 新しいResultでは、すべての判定結果に共通する入力参照と鮮度参照をResult全体へ一度だけ保存します。判定結果ごとの参照が共通値と異なる場合だけ、その判定結果にも保存します。既存Resultは引き続き読み込めます。識別子と後続Resultの鮮度判定を保つため、既存ファイルは自動で書き換えません。 +`adf_submit`はResultの保存が完了した時点で応答します。応答にはResult ID、再送かどうか、次の`adf_next`が必要かどうか、処理段階ごとの所要時間が含まれます。次の作業の計算が遅くても、Resultが保存されたかどうかを区別できます。`adf_next`も処理段階ごとの所要時間を返します。 + 各作業には、実行環境へ向けたモデルの推奨も含まれます。影響評価には通常、軽量なモデルを推奨します。ただし、影響なしと結論付ける場合、根拠が矛盾する場合、セキュリティ、プライバシー、決済、元に戻せないデータ変更の可能性がある場合は、精度の高いモデルへの切り替えを勧めます。ADF自体はモデルを選ばず、LLMも実行しません。 実行環境がすでに把握している処理時間、モデル名、入出力Token数、ツール呼び出し数、再試行回数は、`adf_submit`で任意に記録できます。外部Runnerは`adf_begin_execution`と`adf_complete_execution`を使い、Result提出後に確定したToken数や、失敗・中断した実行も追記できます。外部実行では、キャッシュ作成Token、キャッシュ読取Token、推論Token、実行環境が報告した米ドル費用も記録できます。この実行RecordはADFの状態、Result ID、鮮度、Evidence検証には影響しません。同じResultに提出時の計測値とRunnerの完了Recordがある場合は、Runnerの値だけを集計します。 diff --git a/docs/implementation.md b/docs/implementation.md index ec51012..684c132 100644 --- a/docs/implementation.md +++ b/docs/implementation.md @@ -122,7 +122,9 @@ Explainでは、確認前を`applicability-pending`、支持された後を`not- `contract-health`は、全ChangeのResult・Evidenceと現在のRepository観測から、Contract条項ごとの実装準拠状態を再生成します。 -計算時には、条項ごとのEvidenceと、Evidenceを参照する検証Resultを一度索引化します。Schema検証と現在値のハッシュ計算は、Contract Healthへ影響するContract、Evidence、検証Result、参照先だけに限定します。Result全件を条項ごとに走査せず、無関係なResultの内容全体も検証しません。ただし、Filesystem StoreはすべてのResultファイルを列挙してJSONとして読み込むため、壊れたJSONは引き続きエラーになります。 +計算時には、条項ごとのEvidenceと、Evidenceを参照する検証Resultを索引化します。索引は`.adf/cache/runtime/`へ保存し、追跡済みRecordはGitのBlob ID、変更中または未追跡のRecordは内容ハッシュで更新を判定します。変更のないResultとEvidenceはJSONの読込み、Schema検証、ハッシュ計算を再実行しません。Result全件を条項ごとに走査せず、壊れた索引は正本から作り直します。Gitが保存先を無視すると確認できない場合は、索引をメモリ内だけで使います。 + +Repository観測も永続キャッシュを使います。現在のrevision、解析設定、Signal Catalog、解析対象の識別子がすべて一致する場合だけ再利用します。キャッシュは派生データであり、Actionの認証やRecordの正本には使いません。 Result提出時は、判定結果の`input_refs`と`freshness_refs`がResult全体の値と同じなら省略します。判定結果だけが追加または異なる参照を持つ場合は、その値を判定結果へ保存します。Kernelは判定結果に値がなければResult全体の値を使うため、既存形式と軽量形式を同じ意味で扱えます。Schema上で両項目はもともと省略可能なので、Schema versionは変更しません。既存Resultを変換するとResult IDと、それを参照する後続Resultの鮮度が変わるため、自動移行は行いません。 @@ -513,7 +515,7 @@ bindings: authority_ref: decision.repository-bindings ``` -Agentの通常利用経路はlocal stdio MCP serverです。同じRustバイナリを`mcp` subcommandで起動すると、`next`、`submit`、`explain`、`contract-health`と、発行Actionに限定されたEvidence、Decision、Contract書込みToolを利用できます。Tool契約と信頼境界は[`MCP-DESIGN.md`](MCP-DESIGN.md)、固定I/O Schemaは`schemas/mcp/v1/`にあります。既存CLIは人向け診断、CI、Release・binary管理の補助経路として残します。 +Agentの通常利用経路はlocal stdio MCP serverです。同じRustバイナリを`mcp` subcommandで起動すると、`next`、`submit`、`explain`、`contract-health`と、発行Actionに限定されたEvidence、Decision、Contract書込みToolを利用できます。Tool契約と信頼境界は[`MCP-DESIGN.md`](MCP-DESIGN.md)、固定I/O Schemaは`schemas/mcp/`にあります。既存CLIは人向け診断、CI、Release・binary管理の補助経路として残します。 ```sh adf mcp --project . @@ -521,6 +523,8 @@ adf mcp --project . MCP serverは一つのProject rootへ固定され、stdoutをJSON-RPC専用にします。Action Resultは、`adf_next`が発行した`change_id`、Action ID、Context digestの完全一致でのみ受理します。再接続後は`adf_next`を再実行します。正本を再評価して同じActionが返る場合は、接続断前に保存したEvidence、Decision、Contractを提出時の出力として参照できます。 +`adf_submit`はResultを永続化した時点で応答し、次Actionは計算しません。応答の`next_required`が`true`なら、呼出し側が`adf_next`を別に実行します。`adf_submit`と`adf_next`は処理段階ごとの`timings_ms`を返します。 + ```sh sh scripts/tests/test-rust.sh ``` @@ -764,7 +768,7 @@ sources: | `src/project_runtime.rs` | 実Projectのconfig、Release、Git観測、Storeを接続 | | `src/migration.rs` | 現行CLI Projectを診断し、Migration Draftのレビュー検証、隔離候補の生成・整合性検証・明示適用を行う | | `schemas/v1/` | 保存Recordの言語非依存Schema | -| `schemas/mcp/v1/` | Agent用MCP Toolの固定I/O Schema | +| `schemas/mcp/v1/`、`schemas/mcp/v2/` | Agent用MCP Toolの固定I/O Schema。`adf_submit`出力はv2 | | `schemas/ci/v1/` | project所有のCI policy形式。Contract Healthの停止対象を明示する | | `schemas/benchmarks/v1/` | Detector benchmark corpusとreview済み正解・閾値の固定形式 | | `schemas/catalog/v1/` | 標準Signal Domain CatalogとFramework Detection Catalogの機械可読な固定形式 | diff --git a/schemas/mcp/v1/next-output.schema.json b/schemas/mcp/v1/next-output.schema.json index 551b85d..e33810b 100644 --- a/schemas/mcp/v1/next-output.schema.json +++ b/schemas/mcp/v1/next-output.schema.json @@ -25,6 +25,13 @@ "type": "null" } ] + }, + "timings_ms": { + "type": "object", + "additionalProperties": { + "type": "integer", + "minimum": 0 + } } }, "$defs": { diff --git a/schemas/mcp/v2/submit-output.schema.json b/schemas/mcp/v2/submit-output.schema.json new file mode 100644 index 0000000..cbc91a3 --- /dev/null +++ b/schemas/mcp/v2/submit-output.schema.json @@ -0,0 +1,36 @@ +{ + "$schema": "https://json-schema.org/draft/2020-12/schema", + "$id": "adf://schemas/mcp/v2/submit-output", + "title": "adf_submit output v2", + "type": "object", + "additionalProperties": false, + "required": [ + "schema_version", + "result_id", + "already_completed", + "next_required", + "timings_ms" + ], + "properties": { + "schema_version": { + "const": "2" + }, + "result_id": { + "type": "string", + "pattern": "^result\\.[A-Za-z0-9._-]+$" + }, + "already_completed": { + "type": "boolean" + }, + "next_required": { + "const": true + }, + "timings_ms": { + "type": "object", + "additionalProperties": { + "type": "integer", + "minimum": 0 + } + } + } +} diff --git a/skill-src/adf-analyst/SKILL.md b/skill-src/adf-analyst/SKILL.md index eaa2ef6..65e0703 100644 --- a/skill-src/adf-analyst/SKILL.md +++ b/skill-src/adf-analyst/SKILL.md @@ -189,8 +189,8 @@ The gap is the finding. ## Submit and continue Call `adf_submit` with the action id, the context digest from the action, and the -payload. The control plane validates it, stores it, and returns the next action. Repeat -until it hands the work to another role. +payload. The control plane validates and stores it. After submission succeeds, call +`adf_next` separately. Repeat until it hands the work to another role. If the orchestrator already exposes execution measurements, include them in the optional `execution` object: `duration_ms`, `model`, `input_tokens`, `output_tokens`, `tool_calls`, diff --git a/skill-src/adf-builder/SKILL.md b/skill-src/adf-builder/SKILL.md index 36f9e74..69a7c48 100644 --- a/skill-src/adf-builder/SKILL.md +++ b/skill-src/adf-builder/SKILL.md @@ -82,8 +82,9 @@ Every residual risk needs someone who accepts it and a date by which it is revis ## Submit and continue Call `adf_submit` with the action id, the context digest from the action, and the -payload. The control plane validates it and returns the next action - usually a -challenge run from a context independent of yours. +payload. The control plane validates and stores it. After submission succeeds, call +`adf_next` separately to receive the next action - usually a challenge run from a +context independent of yours. If the orchestrator already knows execution time, model, token counts, tool calls, or retries, it may include them in the optional `execution` object. Do not add work solely to diff --git a/skill-src/adf-challenger/SKILL.md b/skill-src/adf-challenger/SKILL.md index 02f9f76..19c39f4 100644 --- a/skill-src/adf-challenger/SKILL.md +++ b/skill-src/adf-challenger/SKILL.md @@ -85,7 +85,8 @@ person who may decide. ## Submit and continue Call `adf_submit` with the action id, the context digest from the action, and the -payload. The control plane validates it and returns the next action. +payload. The control plane validates and stores it. After submission succeeds, call +`adf_next` separately to receive the next action. If the orchestrator already knows execution time, model, token counts, tool calls, or retries, it may include them in the optional `execution` object. Do not run extra tracing diff --git a/src/application.rs b/src/application.rs index ab2fd51..885fc7d 100644 --- a/src/application.rs +++ b/src/application.rs @@ -208,6 +208,21 @@ impl<'a, Store: ProjectStore> Application<'a, Store> { submission: &ResultSubmission, snapshot: &ProjectSnapshot, ) -> Result { + let result = self.persist_issued_with_snapshot(context, submission, snapshot)?; + let response = self.next(&submission.change_id)?; + Ok(ApplicationSubmission { result, response }) + } + + /// Validate and persist a Result without deriving the next Action. + /// + /// Long-running adapters use this as the acknowledgement boundary so a + /// successful write is not hidden behind repository-wide reevaluation. + pub(crate) fn persist_issued_with_snapshot( + &mut self, + context: &GeneratedContext, + submission: &ResultSubmission, + snapshot: &ProjectSnapshot, + ) -> Result { let result = prepare_result(context, snapshot, submission, self.schema_registry) .map_err(|error| application_error(error.to_string()))?; // The append completes before consuming the issued Context. A failed @@ -219,8 +234,7 @@ impl<'a, Store: ProjectStore> Application<'a, Store> { submission.action_id.clone(), submission.context_digest.clone(), )); - let response = self.next(&submission.change_id)?; - Ok(ApplicationSubmission { result, response }) + Ok(result) } /// Replace the Result for an Action that the current Project still issues. @@ -229,13 +243,13 @@ impl<'a, Store: ProjectStore> Application<'a, Store> { /// prove that the existing Result did not complete or supersede the Action, /// and the Store compares its ID again when writing so concurrent changes /// cannot be lost. - pub(crate) fn correct_issued_with_snapshot( + pub(crate) fn replace_issued_with_snapshot( &mut self, context: &GeneratedContext, submission: &ResultSubmission, snapshot: &ProjectSnapshot, expected_result_id: &str, - ) -> Result { + ) -> Result { let result = prepare_result(context, snapshot, submission, self.schema_registry) .map_err(|error| application_error(error.to_string()))?; self.store @@ -245,8 +259,7 @@ impl<'a, Store: ProjectStore> Application<'a, Store> { submission.action_id.clone(), submission.context_digest.clone(), )); - let response = self.next(&submission.change_id)?; - Ok(ApplicationSubmission { result, response }) + Ok(result) } /// Recompute the current decision and its trace without issuing an Action. diff --git a/src/contract_health.rs b/src/contract_health.rs index 5375c12..d803b7f 100644 --- a/src/contract_health.rs +++ b/src/contract_health.rs @@ -84,6 +84,14 @@ impl ContractHealthReport { pub fn build_contract_health_report( project: &Value, schema_registry: &SchemaRegistry, +) -> Result { + build_contract_health_report_with_digests(project, schema_registry, &BTreeMap::new()) +} + +pub(crate) fn build_contract_health_report_with_digests( + project: &Value, + schema_registry: &SchemaRegistry, + record_digests: &BTreeMap, ) -> Result { let project = project .as_object() @@ -98,7 +106,8 @@ pub fn build_contract_health_report( let verification_outcomes_by_evidence = verification_outcomes_by_evidence(results); let required_current_refs = verification_freshness_refs(results); validate_health_records(project, schema_registry, &required_current_refs)?; - let current_digests = current_artifact_digests(project, &required_current_refs)?; + let current_digests = + current_artifact_digests(project, &required_current_refs, record_digests)?; let evidence_by_clause = evidence_by_clause(evidence); let mut clauses = Vec::new(); @@ -409,13 +418,19 @@ fn is_sha256_digest(value: &str) -> bool { fn current_artifact_digests( project: &serde_json::Map, required_refs: &BTreeSet<&str>, + record_digests: &BTreeMap, ) -> Result, ContractHealthError> { let mut digests = BTreeMap::new(); for collection in ["changes", "contracts", "decisions", "results", "evidence"] { for record in record_array(project.get(collection), collection)? { let id = required_string(record, "id", "Project record")?; if required_refs.contains(id) { - digests.insert(id.to_owned(), digest_value(record)?); + let digest = record_digests + .get(id) + .cloned() + .map(Ok) + .unwrap_or_else(|| digest_value(record))?; + digests.insert(id.to_owned(), digest); } if collection == "contracts" { for clause in record["clauses"] @@ -833,9 +848,9 @@ mod tests { }); let project = project.as_object().unwrap(); - assert!(current_artifact_digests(project, &BTreeSet::new()).is_ok()); + assert!(current_artifact_digests(project, &BTreeSet::new(), &BTreeMap::new()).is_ok()); let required = BTreeSet::from(["result.unrelated"]); - assert!(current_artifact_digests(project, &required).is_err()); + assert!(current_artifact_digests(project, &required, &BTreeMap::new()).is_err()); } #[test] diff --git a/src/filesystem_project.rs b/src/filesystem_project.rs index eea74a2..c29bcf4 100644 --- a/src/filesystem_project.rs +++ b/src/filesystem_project.rs @@ -7,17 +7,21 @@ use crate::application::{ ProjectStore, ProjectStoreError, merge_contract_clause_update, validate_contract_update, validate_decision_update, }; -use crate::contract_health::{ContractHealthReport, build_contract_health_report}; +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::schema::SchemaRegistry; use regex::Regex; +use serde::{Deserialize, Serialize}; use serde_json::{Value, json}; +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::path::{Component, Path, PathBuf}; +use std::process::{Command, Stdio}; pub const DEFAULT_CONTRACT_ROOT: &str = "contracts"; pub const DEFAULT_DECISION_ROOT: &str = "decisions"; @@ -32,6 +36,23 @@ 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 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"; + +#[derive(Debug, Clone, Serialize, Deserialize)] +struct HealthIndexEntry { + source_identity: String, + record_digest: String, + projection_digest: String, + record: Value, +} + +#[derive(Debug, Default, Serialize, Deserialize)] +struct HealthIndex { + schema_version: String, + entries: BTreeMap, +} #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum DocumentFormat { @@ -259,8 +280,8 @@ impl<'a> FileProjectStore<'a> { } pub fn contract_health(&self) -> Result { - let project = self.repository_project()?; - build_contract_health_report(&project, self.schema_registry) + let (project, result_digests) = self.repository_project_for_health()?; + build_contract_health_report_with_digests(&project, self.schema_registry, &result_digests) .map_err(|error| file_error(error.to_string())) } @@ -556,31 +577,120 @@ impl<'a> FileProjectStore<'a> { } } - fn repository_project(&self) -> Result { + fn repository_project_for_health( + &self, + ) -> Result<(Value, BTreeMap), FileProjectError> { let changes = document_paths(&self.change_root)? .into_iter() .filter(|path| is_change_document(path)) .map(|path| self.load_document(&path, "change")) .collect::, _>>()?; let record_paths = recursive_json_paths(&self.change_root)?; - let results = record_paths + let result_paths = record_paths .iter() .filter(|path| parent_name(path) == Some("results")) - .map(|path| read_json(path)) - .collect::, _>>()?; - let evidence = record_paths + .cloned() + .collect::>(); + let (results, mut record_digests) = + self.load_indexed_health_records(&result_paths, RESULT_INDEX_PATH, true)?; + let evidence_paths = record_paths .iter() .filter(|path| parent_name(path) == Some("evidence")) - .map(|path| read_json(path)) - .collect::, _>>()?; - Ok(json!({ - "changes": changes, - "contracts": self.load_document_records(&self.contract_root, "contract")?, - "decisions": self.load_document_records(&self.decision_root, "decision")?, - "results": results, - "evidence": evidence, - "repository": self.repository, - })) + .cloned() + .collect::>(); + let (evidence, evidence_digests) = + self.load_indexed_health_records(&evidence_paths, EVIDENCE_INDEX_PATH, false)?; + for (record_id, digest) in evidence_digests { + if record_digests.insert(record_id.clone(), digest).is_some() { + return Err(file_error(format!("duplicate record id: {record_id}"))); + } + } + Ok(( + json!({ + "changes": changes, + "contracts": self.load_document_records(&self.contract_root, "contract")?, + "decisions": self.load_document_records(&self.decision_root, "decision")?, + "results": results, + "evidence": evidence, + "repository": self.repository, + }), + record_digests, + )) + } + + fn load_indexed_health_records( + &self, + paths: &[PathBuf], + cache_relative_path: &str, + compact_results: bool, + ) -> Result<(Vec, BTreeMap), FileProjectError> { + let cache_path = self.project_root.join(cache_relative_path); + let cache_enabled = path_is_ignored(&self.project_root, cache_relative_path); + let previous = if cache_enabled { + read_health_index(&cache_path).unwrap_or_default() + } else { + HealthIndex::default() + }; + let identities = health_record_source_identities(&self.project_root, paths)?; + let mut entries = BTreeMap::new(); + let mut results = Vec::with_capacity(paths.len()); + let mut digests = BTreeMap::new(); + for path in paths { + let relative = relative_path_string(&self.project_root, path)?; + let source = identities.get(&relative).ok_or_else(|| { + file_error(format!("Result source identity is missing: {relative}")) + })?; + let entry = match previous.entries.get(&relative) { + Some(entry) + if entry.source_identity == source.identity + && canonical_digest(&entry.record).ok().as_deref() + == Some(entry.projection_digest.as_str()) => + { + entry.clone() + } + _ => { + let bytes = source + .bytes + .clone() + .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 record_digest = canonical_digest(&original) + .map_err(|error| file_error(error.to_string()))?; + let record = if compact_results { + compact_indexed_result(original) + } else { + original + }; + HealthIndexEntry { + source_identity: source.identity.clone(), + record_digest, + projection_digest: canonical_digest(&record) + .map_err(|error| file_error(error.to_string()))?, + record, + } + } + }; + let record_id = required_string(&entry.record, "id", "Record")?.to_owned(); + if digests + .insert(record_id.clone(), entry.record_digest.clone()) + .is_some() + { + return Err(file_error(format!("duplicate record id: {record_id}"))); + } + results.push(entry.record.clone()); + entries.insert(relative, entry); + } + let current = HealthIndex { + schema_version: HEALTH_INDEX_SCHEMA_VERSION.to_owned(), + entries, + }; + if cache_enabled { + let _ = write_health_index(&cache_path, ¤t); + } + Ok((results, digests)) } fn change_path(&self, change_id: &str) -> Result { @@ -696,6 +806,207 @@ impl<'a> FileProjectStore<'a> { } } +#[derive(Debug)] +struct HealthRecordSourceIdentity { + identity: String, + bytes: Option>, +} + +fn health_record_source_identities( + root: &Path, + paths: &[PathBuf], +) -> Result, FileProjectError> { + let tracked = tracked_health_record_blobs(root).unwrap_or_default(); + let dirty = dirty_health_record_paths(root).unwrap_or_else(|_| { + paths + .iter() + .filter_map(|path| relative_path_string(root, path).ok()) + .collect() + }); + let mut identities = BTreeMap::new(); + for path in paths { + let relative = relative_path_string(root, path)?; + let (identity, bytes) = if !dirty.contains(&relative) { + match tracked.get(&relative) { + Some(blob) => (format!("git:{blob}"), None), + None => hashed_health_record_identity(path)?, + } + } else { + hashed_health_record_identity(path)? + }; + identities.insert(relative, HealthRecordSourceIdentity { identity, bytes }); + } + Ok(identities) +} + +fn tracked_health_record_blobs(root: &Path) -> Result, FileProjectError> { + let output = git_output(root, &["ls-files", "--stage", "-z", "--", ".adf/changes"])?; + let mut tracked = BTreeMap::new(); + for entry in output + .split(|byte| *byte == 0) + .filter(|entry| !entry.is_empty()) + { + let Some(tab) = entry.iter().position(|byte| *byte == b'\t') else { + continue; + }; + let metadata = String::from_utf8_lossy(&entry[..tab]); + let mut fields = metadata.split_whitespace(); + let _mode = fields.next(); + let Some(blob) = fields.next() else { + continue; + }; + if fields.next() != Some("0") { + continue; + } + let path = String::from_utf8_lossy(&entry[tab + 1..]).into_owned(); + tracked.insert(path, blob.to_owned()); + } + Ok(tracked) +} + +fn dirty_health_record_paths(root: &Path) -> Result, FileProjectError> { + let output = git_output( + root, + &[ + "status", + "--porcelain=v1", + "-z", + "--untracked-files=all", + "--", + ".adf/changes", + ], + )?; + let mut dirty = BTreeSet::new(); + for entry in output + .split(|byte| *byte == 0) + .filter(|entry| !entry.is_empty()) + { + if entry.len() > 3 && entry[2] == b' ' { + dirty.insert(String::from_utf8_lossy(&entry[3..]).into_owned()); + } + } + Ok(dirty) +} + +fn git_output(root: &Path, arguments: &[&str]) -> Result, FileProjectError> { + let output = Command::new("git") + .arg("-C") + .arg(root) + .args(arguments) + .output() + .map_err(|error| file_error(format!("cannot execute Git: {error}")))?; + if !output.status.success() { + return Err(file_error(format!( + "Git command failed ({}): {}", + arguments.join(" "), + String::from_utf8_lossy(&output.stderr).trim() + ))); + } + Ok(output.stdout) +} + +fn path_is_ignored(root: &Path, relative: &str) -> bool { + Command::new("git") + .arg("-C") + .arg(root) + .args(["check-ignore", "--quiet", "--no-index", "--", relative]) + .stdout(Stdio::null()) + .stderr(Stdio::null()) + .status() + .is_ok_and(|status| status.success()) +} + +fn hashed_health_record_identity( + path: &Path, +) -> Result<(String, Option>), FileProjectError> { + let bytes = + fs::read(path).map_err(|error| file_error(format!("{}: {error}", path.display())))?; + let identity = format!("sha256:{:x}", Sha256::digest(&bytes)); + Ok((identity, Some(bytes))) +} + +fn compact_indexed_result(mut result: Value) -> Value { + let result_inputs = result.get("input_refs").cloned(); + let result_freshness = result.get("freshness_refs").cloned(); + for outcome in result["payload"]["outcomes"] + .as_array_mut() + .into_iter() + .flatten() + { + let Some(object) = outcome.as_object_mut() else { + continue; + }; + if object.get("input_refs") == result_inputs.as_ref() { + object.remove("input_refs"); + } + if object.get("freshness_refs") == result_freshness.as_ref() { + object.remove("freshness_refs"); + } + } + result +} + +fn read_health_index(path: &Path) -> Option { + if path.symlink_metadata().ok()?.file_type().is_symlink() { + return None; + } + let index: HealthIndex = serde_json::from_slice(&fs::read(path).ok()?).ok()?; + (index.schema_version == HEALTH_INDEX_SCHEMA_VERSION).then_some(index) +} + +fn write_health_index(path: &Path, index: &HealthIndex) -> Result<(), FileProjectError> { + let parent = path + .parent() + .ok_or_else(|| file_error("Result index has no parent directory"))?; + fs::create_dir_all(parent) + .map_err(|error| file_error(format!("{}: {error}", parent.display())))?; + if path + .symlink_metadata() + .is_ok_and(|metadata| metadata.file_type().is_symlink()) + { + return Err(file_error("Result index is a symlink")); + } + let bytes = serde_json::to_vec(index).map_err(|error| file_error(error.to_string()))?; + let temporary = parent.join(format!( + ".contract-health-results-v1.tmp-{}-{}", + std::process::id(), + TEMPORARY_SEQUENCE.fetch_add(1, Ordering::Relaxed) + )); + let write_result = (|| { + let mut file = OpenOptions::new() + .write(true) + .create_new(true) + .open(&temporary) + .map_err(|error| file_error(format!("{}: {error}", temporary.display())))?; + file.write_all(&bytes) + .and_then(|()| file.sync_all()) + .map_err(|error| file_error(format!("{}: {error}", temporary.display())))?; + fs::rename(&temporary, path).map_err(|error| { + file_error(format!( + "cannot replace {} with {}: {error}", + path.display(), + temporary.display() + )) + }) + })(); + if temporary.exists() { + let _ = fs::remove_file(&temporary); + } + write_result +} + +fn relative_path_string(root: &Path, path: &Path) -> Result { + path.strip_prefix(root) + .map(|relative| { + relative + .components() + .map(|component| component.as_os_str().to_string_lossy()) + .collect::>() + .join("/") + }) + .map_err(|error| file_error(error.to_string())) +} + fn source_root(root: &Path, relative: &str) -> Result { let path = Path::new(relative); if path.is_absolute() { diff --git a/src/git_repository.rs b/src/git_repository.rs index fdbe15d..3a68cce 100644 --- a/src/git_repository.rs +++ b/src/git_repository.rs @@ -15,12 +15,17 @@ use serde_json::{Map, Value, json}; use sha2::{Digest, Sha256}; use std::collections::{BTreeMap, BTreeSet}; use std::fmt; -use std::fs; +use std::fs::{self, OpenOptions}; +use std::io::Write; use std::path::{Path, PathBuf}; -use std::process::Command; +use std::process::{Command, Stdio}; +use std::sync::atomic::{AtomicU64, Ordering}; pub const OBSERVATION_SCHEMA_VERSION: &str = "5"; const LEGACY_OBSERVATION_SCHEMA_VERSION: &str = "4"; +const OBSERVATION_CACHE_SCHEMA_VERSION: &str = "1"; +const OBSERVATION_CACHE_PATH: &str = ".adf/cache/runtime/repository-observation-v1.json"; +static CACHE_SEQUENCE: AtomicU64 = AtomicU64::new(0); pub struct GitRepositoryAdapter { root: PathBuf, @@ -295,6 +300,25 @@ impl GitRepositoryAdapter { })) } + /// Reuse a repository observation only when every relevant input has the + /// same content identity. Clean tracked files use their Git blob IDs; + /// modified and untracked files are hashed directly. + pub fn observe_cached(&self) -> Result { + self.assert_repository_state()?; + let manifest = self.read_manifest()?; + if !path_is_ignored(&self.root, OBSERVATION_CACHE_PATH) { + return self.observe(); + } + let signature = self.observation_cache_signature(&manifest)?; + let cache_path = self.root.join(OBSERVATION_CACHE_PATH); + if let Some(repository) = read_observation_cache(&cache_path, &signature) { + return Ok(repository); + } + let repository = self.observe()?; + let _ = write_observation_cache(&cache_path, &signature, &repository); + Ok(repository) + } + pub fn binding_authority_refs(&self) -> Result, GitRepositoryError> { let manifest = self.read_manifest()?; let mut authority_refs = BTreeSet::new(); @@ -406,6 +430,113 @@ impl GitRepositoryAdapter { Ok(targets) } + fn observation_cache_signature( + &self, + manifest: &Map, + ) -> Result { + let roots = analysis_roots( + manifest + .get("analysis") + .ok_or_else(|| git_error("repository analysis is missing"))?, + )?; + let mut targets = self + .analysis_targets(&roots)? + .into_iter() + .collect::>(); + for declaration in required_array(manifest, "artifacts", "repository observation")? { + let declaration = declaration + .as_object() + .ok_or_else(|| git_error("artifact declaration must be a mapping"))?; + targets.insert( + required_nonempty_string(declaration, "path", "artifact declaration")?.to_owned(), + ); + } + let targets = targets.into_iter().collect::>(); + let tracked = self.tracked_blob_ids(&targets)?; + let dirty = self.dirty_source_paths(&targets)?; + let mut inputs = Vec::with_capacity(targets.len()); + for path in targets { + let source_path = + repository_path(&self.root, &path).map_err(|error| git_error(error.to_string()))?; + let identity = if !dirty.contains(path.as_str()) { + tracked.get(path.as_str()).cloned() + } else { + None + } + .map(|blob| format!("git:{blob}")) + .map(Ok) + .unwrap_or_else(|| file_sha256(&source_path))?; + inputs.push(json!({"path": path, "identity": identity})); + } + canonical_digest(&json!({ + "schema_version": OBSERVATION_CACHE_SCHEMA_VERSION, + "revision": self.git(&["rev-parse", "HEAD"] )?, + "manifest": manifest, + "signal_catalog_digest": self.signal_registry.digest(), + "require_clean": self.require_clean, + "inputs": inputs, + })) + .map_err(|error| git_error(error.to_string())) + } + + fn tracked_blob_ids( + &self, + paths: &[String], + ) -> Result, GitRepositoryError> { + if paths.is_empty() { + return Ok(BTreeMap::new()); + } + let mut arguments = vec!["ls-files", "--stage", "-z", "--"]; + arguments.extend(paths.iter().map(String::as_str)); + let output = self.git_raw(&arguments)?; + let mut tracked = BTreeMap::new(); + for entry in output + .split(|byte| *byte == 0) + .filter(|entry| !entry.is_empty()) + { + let Some(tab) = entry.iter().position(|byte| *byte == b'\t') else { + continue; + }; + let metadata = String::from_utf8_lossy(&entry[..tab]); + let mut fields = metadata.split_whitespace(); + let _mode = fields.next(); + let Some(blob) = fields.next() else { + continue; + }; + if fields.next() != Some("0") { + continue; + } + let path = String::from_utf8_lossy(&entry[tab + 1..]).into_owned(); + tracked.insert(path, blob.to_owned()); + } + Ok(tracked) + } + + fn dirty_source_paths(&self, paths: &[String]) -> Result, GitRepositoryError> { + if paths.is_empty() { + return Ok(BTreeSet::new()); + } + let mut arguments = vec![ + "status", + "--porcelain=v1", + "-z", + "--untracked-files=all", + "--", + ]; + arguments.extend(paths.iter().map(String::as_str)); + let output = self.git_raw(&arguments)?; + let mut dirty = BTreeSet::new(); + for entry in output + .split(|byte| *byte == 0) + .filter(|entry| !entry.is_empty()) + { + if entry.len() > 3 && entry[2] == b' ' { + dirty.insert(String::from_utf8_lossy(&entry[3..]).into_owned()); + } + } + Ok(dirty) + } + fn assert_tracked(&self, relative: &str) -> Result<(), GitRepositoryError> { self.git(&["ls-files", "--error-unmatch", "--", relative]) .map(|_| ()) @@ -429,6 +560,11 @@ impl GitRepositoryAdapter { } fn git(&self, arguments: &[&str]) -> Result { + let output = self.git_raw(arguments)?; + Ok(String::from_utf8_lossy(&output).trim().to_owned()) + } + + fn git_raw(&self, arguments: &[&str]) -> Result, GitRepositoryError> { let output = Command::new("git") .arg("-C") .arg(&self.root) @@ -444,10 +580,81 @@ impl GitRepositoryAdapter { arguments.join(" ") ))); } - Ok(String::from_utf8_lossy(&output.stdout).trim().to_owned()) + Ok(output.stdout) } } +fn read_observation_cache(path: &Path, signature: &str) -> Option { + if path.symlink_metadata().ok()?.file_type().is_symlink() { + return None; + } + let cache: Value = serde_json::from_slice(&fs::read(path).ok()?).ok()?; + if cache["schema_version"].as_str() != Some(OBSERVATION_CACHE_SCHEMA_VERSION) + || cache["input_signature"].as_str() != Some(signature) + { + return None; + } + let repository = cache.get("repository")?.clone(); + let expected = cache["repository_digest"].as_str()?; + (canonical_digest(&repository).ok()?.as_str() == expected).then_some(repository) +} + +fn write_observation_cache(path: &Path, signature: &str, repository: &Value) -> Result<(), String> { + let parent = path + .parent() + .ok_or_else(|| "repository observation cache has no parent".to_owned())?; + fs::create_dir_all(parent).map_err(|error| error.to_string())?; + if path + .symlink_metadata() + .is_ok_and(|metadata| metadata.file_type().is_symlink()) + { + return Err("repository observation cache is a symlink".to_owned()); + } + let value = json!({ + "schema_version": OBSERVATION_CACHE_SCHEMA_VERSION, + "input_signature": signature, + "repository_digest": canonical_digest(repository).map_err(|error| error.to_string())?, + "repository": repository, + }); + let bytes = serde_json::to_vec(&value).map_err(|error| error.to_string())?; + let temporary = parent.join(format!( + ".repository-observation-v1.tmp-{}-{}", + std::process::id(), + CACHE_SEQUENCE.fetch_add(1, Ordering::Relaxed) + )); + let result = (|| { + let mut file = OpenOptions::new() + .write(true) + .create_new(true) + .open(&temporary) + .map_err(|error| error.to_string())?; + file.write_all(&bytes).map_err(|error| error.to_string())?; + file.sync_all().map_err(|error| error.to_string())?; + fs::rename(&temporary, path).map_err(|error| error.to_string()) + })(); + if temporary.exists() { + let _ = fs::remove_file(&temporary); + } + result +} + +fn file_sha256(path: &Path) -> Result { + let bytes = + fs::read(path).map_err(|error| git_error(format!("{}: {error}", path.display())))?; + Ok(format!("sha256:{:x}", Sha256::digest(bytes))) +} + +fn path_is_ignored(root: &Path, relative: &str) -> bool { + Command::new("git") + .arg("-C") + .arg(root) + .args(["check-ignore", "--quiet", "--no-index", "--", relative]) + .stdout(Stdio::null()) + .stderr(Stdio::null()) + .status() + .is_ok_and(|status| status.success()) +} + fn analysis_roots(value: &Value) -> Result, GitRepositoryError> { let analysis = value .as_object() diff --git a/src/mcp_server.rs b/src/mcp_server.rs index 34cd527..7bfea85 100644 --- a/src/mcp_server.rs +++ b/src/mcp_server.rs @@ -321,7 +321,7 @@ impl AgenticMcpServer { /// Validate and persist the Result for an Action issued in this MCP session. #[tool( name = "adf_submit", - description = "Validate and persist an issued Action Result, then return the reevaluated next Action.", + description = "Validate and persist an issued Action Result. Call adf_next separately after persistence succeeds.", annotations( title = "Agentic Submit", read_only_hint = false, diff --git a/src/migration.rs b/src/migration.rs index ae24b7d..365f855 100644 --- a/src/migration.rs +++ b/src/migration.rs @@ -1865,7 +1865,7 @@ fn verify_unclaimed_candidate_files( for file in files { if is_reserved_candidate_artifact(&file) || file.starts_with("migration-completions/") - || file.starts_with(".adf/cache/releases/") + || file.starts_with(".adf/cache/") { continue; } diff --git a/src/project_application.rs b/src/project_application.rs index ab3f072..7463a9b 100644 --- a/src/project_application.rs +++ b/src/project_application.rs @@ -4,7 +4,7 @@ //! Record. Git, Records, and the verified Framework Release are re-read before //! every evaluation or write. -use crate::application::{ApplicationResponse, ApplicationSubmission}; +use crate::application::ApplicationResponse; use crate::cli_output::next_response_value; use crate::context::GeneratedContext; use crate::execution_log::ExecutionLog; @@ -20,8 +20,10 @@ use serde_json::{Value, json}; use std::collections::{BTreeMap, BTreeSet}; use std::fmt; use std::path::{Path, PathBuf}; +use std::time::Instant; pub const MCP_APPLICATION_PROTOCOL_VERSION: &str = "1"; +pub const MCP_SUBMIT_PROTOCOL_VERSION: &str = "2"; /// Describe a Record-shaped JSON value as a Schema object rather than the /// boolean Schema `true` that `serde_json::Value` produces on its own. @@ -65,6 +67,7 @@ pub struct NextServiceResponse { #[schemars(schema_with = "json_object_schema")] pub next_response: Value, pub issued_action: Option, + pub timings_ms: BTreeMap, } #[derive(Debug, Clone, Serialize, JsonSchema)] @@ -72,9 +75,8 @@ pub struct SubmitServiceResponse { pub schema_version: String, pub result_id: String, pub already_completed: bool, - #[schemars(schema_with = "json_object_schema")] - pub next_response: Value, - pub issued_action: Option, + pub next_required: bool, + pub timings_ms: BTreeMap, } #[derive(Debug, Clone, Serialize, JsonSchema)] @@ -133,14 +135,19 @@ impl ProjectApplicationService { change_id: &str, require_clean: bool, ) -> Result { + let total_started = Instant::now(); + let load_started = Instant::now(); let project = self.load(require_clean)?; + let repository_load_ms = elapsed_ms(load_started); if require_clean { project .assert_tracked_inputs(change_id) .map_err(project_error)?; } + let evaluation_started = Instant::now(); let mut application = project.application().map_err(application_error)?; let response = application.next(change_id).map_err(application_error)?; + let evaluation_ms = elapsed_ms(evaluation_started); let rule_index_digest = application.rule_index_digest().to_owned(); let framework_lock_digest = application.framework_lock_digest().to_owned(); let next_response = next_response_value(change_id, &response); @@ -154,6 +161,11 @@ impl ProjectApplicationService { schema_version: MCP_APPLICATION_PROTOCOL_VERSION.to_owned(), next_response, issued_action, + timings_ms: BTreeMap::from([ + ("repository_load".to_owned(), repository_load_ms), + ("evaluation".to_owned(), evaluation_ms), + ("total".to_owned(), elapsed_ms(total_started)), + ]), }) } @@ -249,23 +261,36 @@ impl ProjectApplicationService { output_refs: Vec, execution: Option, ) -> Result { + let total_started = Instant::now(); + let mut timings_ms = BTreeMap::new(); let entry = match self.issued.get(key).cloned() { Some(entry) => entry, None => { // A Result already stored for this Action and Context means the // submission arrived twice; replay it rather than writing again. - if let Some(response) = self.replay_submission(key, &payload, &output_refs)? { - return Ok(response); + let replay_started = Instant::now(); + if let Some(result_id) = self.replay_submission(key, &payload, &output_refs)? { + timings_ms.insert("replay_check".to_owned(), elapsed_ms(replay_started)); + timings_ms.insert("total".to_owned(), elapsed_ms(total_started)); + return Ok(submit_response(result_id, true, timings_ms)); } - self.resume(key)? + timings_ms.insert("replay_check".to_owned(), elapsed_ms(replay_started)); + let resume_started = Instant::now(); + let entry = self.resume(key)?; + timings_ms.insert("action_resume".to_owned(), elapsed_ms(resume_started)); + entry } }; + let load_started = Instant::now(); let project = self.load(false)?; + timings_ms.insert("repository_load".to_owned(), elapsed_ms(load_started)); let mut application = project.application().map_err(application_error)?; + let snapshot_started = Instant::now(); let snapshot = application .snapshot(&key.change_id) .map_err(application_error)?; + timings_ms.insert("change_snapshot".to_owned(), elapsed_ms(snapshot_started)); let existing_result = snapshot .results .iter() @@ -275,10 +300,9 @@ impl ProjectApplicationService { && result["output_refs"] == json!(output_refs) { let result_id = required_record_id(result, "Result")?.to_owned(); - let response = application - .next(&key.change_id) - .map_err(application_error)?; - return Ok(self.completed_response(key, result_id, response, &application)); + self.issued.remove(key); + timings_ms.insert("total".to_owned(), elapsed_ms(total_started)); + return Ok(submit_response(result_id, true, timings_ms)); } validate_output_refs(&entry, &output_refs, &snapshot)?; assert_framework_identity(&entry, &application)?; @@ -295,10 +319,11 @@ impl ProjectApplicationService { output_refs, execution, }; - let ApplicationSubmission { result, response } = if let Some(existing) = existing_result { + let persist_started = Instant::now(); + let result = if let Some(existing) = existing_result { let expected_result_id = required_record_id(existing, "Result")?; application - .correct_issued_with_snapshot( + .replace_issued_with_snapshot( &entry.context, &submission, &snapshot, @@ -307,68 +332,23 @@ impl ProjectApplicationService { .map_err(application_error)? } else { application - .submit_issued_with_snapshot(&entry.context, &submission, &snapshot) + .persist_issued_with_snapshot(&entry.context, &submission, &snapshot) .map_err(application_error)? }; - let rule_index_digest = application.rule_index_digest().to_owned(); - let framework_lock_digest = application.framework_lock_digest().to_owned(); + timings_ms.insert( + "validation_and_persist".to_owned(), + elapsed_ms(persist_started), + ); let result_id = result["id"] .as_str() .expect("validated Result has an ID") .to_owned(); - let next_response = next_response_value(&key.change_id, &response); - let next_issued = issued_entry( - &key.change_id, - &response, - &rule_index_digest, - &framework_lock_digest, - ); drop(application); drop(project); self.issued.remove(key); - let issued_action = next_issued.map(|(next_key, next_entry)| { - self.issued.insert(next_key.clone(), next_entry); - next_key - }); - Ok(SubmitServiceResponse { - schema_version: MCP_APPLICATION_PROTOCOL_VERSION.to_owned(), - result_id, - already_completed: false, - next_response, - issued_action, - }) - } - - fn completed_response( - &mut self, - key: &IssuedActionKey, - result_id: String, - response: crate::application::ApplicationResponse, - application: &crate::application::Application< - '_, - crate::filesystem_project::FileProjectStore<'_>, - >, - ) -> SubmitServiceResponse { - let next_response = next_response_value(&key.change_id, &response); - let next_issued = issued_entry( - &key.change_id, - &response, - application.rule_index_digest(), - application.framework_lock_digest(), - ); - self.issued.remove(key); - let issued_action = next_issued.map(|(next_key, next_entry)| { - self.issued.insert(next_key.clone(), next_entry); - next_key - }); - SubmitServiceResponse { - schema_version: MCP_APPLICATION_PROTOCOL_VERSION.to_owned(), - result_id, - already_completed: true, - next_response, - issued_action, - } + timings_ms.insert("total".to_owned(), elapsed_ms(total_started)); + Ok(submit_response(result_id, false, timings_ms)) } /// Replays a submission whose Result is already stored, or `None` when this @@ -378,9 +358,9 @@ impl ProjectApplicationService { key: &IssuedActionKey, payload: &Value, output_refs: &[String], - ) -> Result, ServiceError> { + ) -> Result, ServiceError> { let project = self.load(false)?; - let mut application = project.application().map_err(application_error)?; + let application = project.application().map_err(application_error)?; let snapshot = application .snapshot(&key.change_id) .map_err(application_error)?; @@ -399,32 +379,7 @@ impl ProjectApplicationService { false, )); } - let result_id = required_record_id(result, "Result")?.to_owned(); - let response = application - .next(&key.change_id) - .map_err(application_error)?; - let rule_index_digest = application.rule_index_digest().to_owned(); - let framework_lock_digest = application.framework_lock_digest().to_owned(); - let next_response = next_response_value(&key.change_id, &response); - let next_issued = issued_entry( - &key.change_id, - &response, - &rule_index_digest, - &framework_lock_digest, - ); - drop(application); - drop(project); - let issued_action = next_issued.map(|(next_key, next_entry)| { - self.issued.insert(next_key.clone(), next_entry); - next_key - }); - Ok(Some(SubmitServiceResponse { - schema_version: MCP_APPLICATION_PROTOCOL_VERSION.to_owned(), - result_id, - already_completed: true, - next_response, - issued_action, - })) + Ok(Some(required_record_id(result, "Result")?.to_owned())) } pub fn add_evidence( @@ -675,6 +630,24 @@ impl ProjectApplicationService { } } +fn submit_response( + result_id: String, + already_completed: bool, + timings_ms: BTreeMap, +) -> SubmitServiceResponse { + SubmitServiceResponse { + schema_version: MCP_SUBMIT_PROTOCOL_VERSION.to_owned(), + result_id, + already_completed, + next_required: true, + timings_ms, + } +} + +fn elapsed_ms(started: Instant) -> u64 { + u64::try_from(started.elapsed().as_millis()).unwrap_or(u64::MAX) +} + fn validate_evidence_claims(evidence: &Value) -> Result<(), ServiceError> { let declared_clauses = string_set(&evidence["contract_clause_refs"]); let mut claimed_clauses = BTreeSet::new(); diff --git a/src/project_runtime.rs b/src/project_runtime.rs index 44b9421..0579dc5 100644 --- a/src/project_runtime.rs +++ b/src/project_runtime.rs @@ -49,7 +49,7 @@ impl LoadedProject { require_clean, signal_registry.clone(), ) - .and_then(|adapter| adapter.observe()) + .and_then(|adapter| adapter.observe_cached()) .map_err(|error| runtime_error(error.to_string()))?; Ok(Self { root, diff --git a/templates/docs/adf/README.md b/templates/docs/adf/README.md index 90ad19b..2713cb8 100644 --- a/templates/docs/adf/README.md +++ b/templates/docs/adf/README.md @@ -19,7 +19,9 @@ outside ADF. adf next ``` -Agents normally use MCP. Start `adf mcp`, then use `adf_next` and `adf_submit`. +Agents normally use MCP. Start `adf mcp`, call `adf_next`, complete the issued +action, and call `adf_submit`. A successful submission stores the Result; call +`adf_next` separately to continue. Use the Skill for the role in the issued action. diff --git a/tests/cli.rs b/tests/cli.rs index f4046ce..4d7df4a 100644 --- a/tests/cli.rs +++ b/tests/cli.rs @@ -890,7 +890,8 @@ fn migration_candidate_requires_a_signed_release_and_schema_valid_records() { .as_array() .unwrap() .iter() - .any(|issue| issue["category"] == "invalid-candidate-records") + .any(|issue| issue["category"] == "invalid-candidate-records"), + "unexpected validation report: {report}" ); fs::write(&contract_path, valid_contract).unwrap(); @@ -3200,6 +3201,32 @@ fn contract_health_policy_turns_the_report_into_an_explicit_ci_gate() { assert_eq!(passed["blocking_clause_refs"], json!([])); } +#[test] +fn derived_runtime_indexes_are_rebuilt_when_cache_files_are_corrupt() { + let project = TestProject::new(); + let first = project.run(&["contract-health", "--format", "json"]); + assert_success(&first); + + let cache_paths = [ + ".adf/cache/runtime/repository-observation-v1.json", + ".adf/cache/runtime/contract-health-results-v1.json", + ".adf/cache/runtime/contract-health-evidence-v1.json", + ]; + for relative in cache_paths { + let path = project.root.join(relative); + assert!(path.is_file(), "expected runtime cache: {relative}"); + fs::write(path, b"not json").unwrap(); + } + + let rebuilt = project.run(&["contract-health", "--format", "json"]); + assert_success(&rebuilt); + 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"); + } +} + #[test] fn contract_health_gate_rejects_invalid_or_untracked_policy() { let project = TestProject::new(); @@ -3387,6 +3414,7 @@ fn stdio_mcp_lists_typed_tools_and_persists_an_issued_result() { assert_eq!(next["result"]["isError"], false); let next = &next["result"]["structuredContent"]; validate_mcp_schema(next, "next-output.schema.json"); + assert!(next["timings_ms"]["total"].is_number()); assert_eq!(next["next_response"]["state"], "needs-analysis"); assert_eq!( next["next_response"]["next_action"]["action"], @@ -3504,8 +3532,12 @@ fn stdio_mcp_lists_typed_tools_and_persists_an_issued_result() { let submitted = mcp_receive(&mut output); assert_eq!(submitted["result"]["isError"], false); let submitted = &submitted["result"]["structuredContent"]; - validate_mcp_schema(submitted, "submit-output.schema.json"); + validate_mcp_schema_version(submitted, "v2", "submit-output.schema.json"); assert_eq!(submitted["already_completed"], false); + assert_eq!(submitted["next_required"], true); + assert!(submitted["timings_ms"]["total"].is_number()); + assert!(submitted.get("next_response").is_none()); + assert!(submitted.get("issued_action").is_none()); assert!( submitted["result_id"] .as_str() @@ -3519,6 +3551,16 @@ fn stdio_mcp_lists_typed_tools_and_persists_an_issued_result() { 1 ); + let continued = mcp_call( + &mut input, + &mut output, + 32, + "adf_next", + json!({"change_id": "change.place-order"}), + ); + assert_eq!(continued["isError"], false); + validate_mcp_schema(&continued["structuredContent"], "next-output.schema.json"); + let completed = mcp_call( &mut input, &mut output, @@ -5962,8 +6004,13 @@ fn validate_delivery_schema(value: &Value, filename: &str) { } fn validate_mcp_schema(value: &Value, filename: &str) { + validate_mcp_schema_version(value, "v1", filename); +} + +fn validate_mcp_schema_version(value: &Value, version: &str, filename: &str) { let path = PathBuf::from(env!("CARGO_MANIFEST_DIR")) - .join("schemas/mcp/v1") + .join("schemas/mcp") + .join(version) .join(filename); let schema: Value = serde_json::from_slice(&fs::read(path).unwrap()).unwrap(); validate_json_document(value, &schema).unwrap();