12 Commits

50 changed files with 3068 additions and 199 deletions

View File

@@ -30,6 +30,7 @@ For config semantics, see [docs/config.md](./config.md). For operator lifecycle
- `--config <path>`: optional explicit `pipeline.yml` path.
- `--session <path>`: optional explicit `session.yml` path.
- `--session-id <value>`: session template variable value.
- `--previous-session-id <value>`: previous-session template variable value.
- `--force`: force stage execution.
- `--artifacts <names>`: analyze artifact keys to execute (repeatable or comma-separated).
@@ -38,6 +39,7 @@ For config semantics, see [docs/config.md](./config.md). For operator lifecycle
- `--config <path>`
- `--session <path>`
- `--session-id <value>`
- `--previous-session-id <value>`
- `--force`
### `resume`
@@ -45,6 +47,7 @@ For config semantics, see [docs/config.md](./config.md). For operator lifecycle
- `--config <path>`
- `--session <path>`
- `--session-id <value>`
- `--previous-session-id <value>`
- `--force`
- `--artifacts <names>`: analyze artifact keys to execute (repeatable or comma-separated).
@@ -53,6 +56,7 @@ For config semantics, see [docs/config.md](./config.md). For operator lifecycle
- `--config <path>`
- `--session <path>`
- `--session-id <value>`
- `--previous-session-id <value>`
- `--force`
- `--artifacts <names>`: analyze artifact keys to execute (repeatable or comma-separated).
- positional `<stage>`: required stage name.
@@ -74,6 +78,7 @@ Valid stage names:
- `--config <path>`
- `--session <path>`
- `--session-id <value>`
- `--previous-session-id <value>`
- `--dry-run`: plan restore actions without writing local files.
- `--force`: overwrite local conflicting files with remote archive files.
- `--include-audio`: include durable archived `audio/**` files in restore scope.
@@ -92,7 +97,7 @@ Purpose:
Syntax:
```bash
narratio run [--config <pipeline.yml>] [--session <session.yml>] [--session-id <id>] [--force] [--artifacts <name[,name...]>]
narratio run [--config <pipeline.yml>] [--session <session.yml>] [--session-id <id>] [--previous-session-id <id>] [--force] [--artifacts <name[,name...]>]
```
Success output:
@@ -112,7 +117,7 @@ Purpose:
Syntax:
```bash
narratio plan [--config <pipeline.yml>] [--session <session.yml>] [--session-id <id>] [--force]
narratio plan [--config <pipeline.yml>] [--session <session.yml>] [--session-id <id>] [--previous-session-id <id>] [--force]
```
Success output includes:
@@ -132,7 +137,7 @@ Purpose:
Syntax:
```bash
narratio resume [--config <pipeline.yml>] [--session <session.yml>] [--session-id <id>] [--force] [--artifacts <name[,name...]>]
narratio resume [--config <pipeline.yml>] [--session <session.yml>] [--session-id <id>] [--previous-session-id <id>] [--force] [--artifacts <name[,name...]>]
```
Success output:
@@ -172,7 +177,7 @@ Purpose:
Syntax:
```bash
narratio run-stage [--config <pipeline.yml>] [--session <session.yml>] [--session-id <id>] [--force] [--artifacts <name[,name...]>] <stage>
narratio run-stage [--config <pipeline.yml>] [--session <session.yml>] [--session-id <id>] [--previous-session-id <id>] [--force] [--artifacts <name[,name...]>] <stage>
```
Success output:
@@ -191,12 +196,12 @@ Common failure cases:
### `restore`
Purpose:
- Restore durable session state (`manifest.json`, `transcripts/**`, `artifacts/**`, and optional `audio/**`) from the committed remote archive current state.
- Restore durable session state (`manifest.json`, `transcripts/**`, `artifacts/**`, `previous/**`, and optional `audio/**`) from the committed remote archive current state.
Syntax:
```bash
narratio restore [--config <pipeline.yml>] [--session <session.yml>] [--session-id <id>] [--dry-run] [--force] [--include-audio]
narratio restore [--config <pipeline.yml>] [--session <session.yml>] [--session-id <id>] [--previous-session-id <id>] [--dry-run] [--force] [--include-audio]
```
Success output (dry-run):
@@ -260,6 +265,12 @@ narratio restore --session-id 2026-04-04
narratio run-stage --session-id 2026-04-04 --force analyze
```
Rehydrate canonical previous-session inputs after artifact-input changes:
```bash
narratio run-stage --session-id 2026-04-04 --force prepare
```
## Diagnostic / Recovery Commands
Inspect stage status:

View File

@@ -48,9 +48,13 @@ Template behavior:
- supported placeholders:
- `{{session_id}}`
- `{{ session_id }}`
- `{{previous_session_id}}`
- `{{ previous_session_id }}`
- `--session-id <value>` supplies the placeholder value.
- `--previous-session-id <value>` supplies the previous-session placeholder value.
- unresolved placeholders fail load.
- if rendered `session_id` mismatches `--session-id`, load fails.
- if rendered `previous_session_id` mismatches `--previous-session-id`, load fails.
## 4. Minimal pipeline config
@@ -83,6 +87,23 @@ Usage:
narratio run --config /path/to/pipeline.yml --session ./session.yml --session-id 2026-05-03
```
Previous-session-enabled variant:
```yaml
session_id: "{{ session_id }}"
previous_session_id: "{{ previous_session_id }}"
campaign: sample-campaign
inputs:
audio_dir: ./audio
speakers_file: ./examples/speakers.yml
autocorrect_file: ./examples/autocorrect.yml
glossary_file: ./examples/glossary.yml
```
```bash
narratio run --config /path/to/pipeline.yml --session ./session.yml --session-id 2026-05-03 --previous-session-id 2026-04-26
```
## 6. Production-oriented config
```yaml
@@ -127,6 +148,9 @@ scriptorium:
transcript:
source: narratio.transcript.trimmed
required: true
previous_recap:
source: narratio.previous_session.artifact.session_recap
required: false
```
Operational notes:
@@ -243,13 +267,14 @@ Scriptorium artifact-key and dependency rules:
Allowed `pipeline.scriptorium.artifacts.<name>.inputs.<key>.source` values:
- `previous_session_artifact`
- `narratio.previous_session.artifact.<configured_artifact_key>`
- `narratio.transcript.merged`
- `narratio.transcript.polished`
- `narratio.transcript.full`
- `narratio.transcript.trimmed`
- `narratio.bounds.session`
- `narratio.artifact.<configured_artifact_key>`
- `previous_session_artifact` (legacy path-based source; uses `inputs.<key>.path`)
`pipeline.archive.promote_artifacts[].source` values:
@@ -272,13 +297,14 @@ Archive promotion destination rules:
Restore-related implications:
- restore remote identity requires archive S3 identity to resolve (`pipeline.storage.s3.bucket` and session prefix derivation inputs).
- restore scope considers only committed current state and durable paths (`manifest.json`, `transcripts/**`, `artifacts/**`, optional `audio/**`).
- restore scope considers committed current state and durable paths (`manifest.json`, `transcripts/**`, `artifacts/**`, `previous/**`, optional `audio/**`).
## 8. Full session reference
| Path | Type | Required | Default |
| --- | --- | --- | --- |
| `session.session_id` | string | Yes | none |
| `session.previous_session_id` | string | No | empty |
| `session.campaign` | string | Yes | none |
| `session.date` | string | No | empty |
| `session.title` | string | No | empty |
@@ -297,6 +323,11 @@ Audio-source rule:
- `audio_s3.prefix`
- `audio_s3` cannot be combined with local audio fields.
Previous-session rule:
- if `session.previous_session_id` is set, it must not equal `session.session_id`.
- canonical previous-session sources (`narratio.previous_session.artifact.<name>`) are hydrated during `prepare` from archive current state when required by enabled configured artifacts.
## 9. Secrets
Narratio supports filesystem-based secret injection via `pipeline.secrets.env_dir`.

View File

@@ -1,41 +1,35 @@
# Internal: Artifacts
## Purpose
Define Narratio's artifact identity and resolution model for built-in transcript/bounds artifacts and runtime-configured analyze artifacts.
Define Narratio artifact identity, catalog, and source-resolution behavior for:
- built-in session artifacts;
- configured analyze artifacts;
- canonical previous-session artifact sources.
## Inputs and outputs
Inputs:
- artifact sources from config/runtime (`pipeline.scriptorium.artifacts.*.inputs.*.source`)
- session paths and optional session manifest stage outputs
- runtime artifact catalog state for configured artifact sources
- configured input sources (`pipeline.scriptorium.artifacts.*.inputs.*.source`);
- session paths and manifest inputs/outputs;
- runtime catalog state.
Outputs:
- resolved local artifact path and provenance (`ResolvedSessionArtifact`)
- runtime catalog entries for planned/executable/available artifacts
- validation errors for unsupported, missing, or invalid artifact sources
- resolved artifact path + provenance (`ResolvedSessionArtifact`);
- runtime catalog entries for built-ins and configured artifacts;
- requirement sets for canonical previous-session inputs.
## Boundaries
Owns:
- built-in artifact registry and content validation rules
- runtime artifact catalog for configured artifact source IDs
- source resolution behavior for built-in and configured artifact sources
- built-in source registry and validation;
- configured artifact catalog identity (`narratio.artifact.<name>`);
- canonical previous-session source parsing and resolution;
- previous-session requirement collection (`CollectPreviousArtifactRequirements`).
Does not own:
- artifact generation (stages produce files)
- manifest transition policy
- archive promotion behavior
## Config fields used
- `pipeline.scriptorium.artifacts.<name>.enabled`
- `pipeline.scriptorium.artifacts.<name>.output_path`
- `pipeline.scriptorium.artifacts.<name>.inputs.<key>.source`
## External adapters used
- none
## State and manifest behavior
Built-in registry entries:
- prepare-stage remote hydration;
- stage success/skip transitions;
- archive upload orchestration.
## Built-in IDs
| Artifact ID | Canonical file | Producer stage | Output kind |
| --- | --- | --- | --- |
| `narratio.transcript.merged` | `transcripts/merged.json` | `merge` | `transcript_merged` |
@@ -44,44 +38,58 @@ Built-in registry entries:
| `narratio.transcript.trimmed` | `transcripts/trimmed.json` | `trim` | `transcript_trimmed` |
| `narratio.bounds.session` | `artifacts/session_bounds.json` | `trim` | `session_bounds` |
Runtime catalog entries include built-ins and configured `narratio.artifact.<name>` sources.
## Source families
- built-in: `narratio.transcript.*`, `narratio.bounds.session`
- configured artifact: `narratio.artifact.<artifact_key>`
- canonical previous-session artifact: `narratio.previous_session.artifact.<artifact_key>`
Catalog states:
- `planned`: source is registered and known for this run
- `executable`: configured artifact is selected for analyze execution
- `available`: artifact has a usable file path (generated this run or reused from disk)
## Runtime catalog model
Catalog entries track:
- `planned`: source is registered for this run;
- `executable`: configured artifact is selected for analyze execution;
- `available`: usable local file exists (generated this run or reused from disk).
Resolution behavior:
- built-in sources resolve via manifest producer outputs first, then canonical fallback path
- configured `narratio.artifact.<name>` sources resolve through runtime catalog availability
- configured source lookup requires catalog context
Configured artifact provenance values:
Configured artifact provenance values include:
- `generated.current_analyze_run`
- `filesystem.disabled_artifact_output`
Content validation:
- transcript built-ins: JSON with top-level `segments` array
- bounds built-in: valid JSON
- configured artifacts: non-empty text file
Previous-session canonical provenance values include:
- `manifest.inputs.previous_cache`
- `current_session.previous_cache`
## Skip and resume behavior
- resolver and catalog have no direct skip/resume decisions
- stage/runner skip-resume behavior consumes catalog/resolver results
## Resolution behavior
- Built-ins resolve via manifest producer outputs first, then canonical fallback paths.
- Configured `narratio.artifact.<name>` sources resolve through catalog availability.
- Canonical previous-session sources resolve to current-session `previous/` cache candidates derived from configured artifact canonical output paths.
- Previous-session canonical resolution prefers manifest-recorded input paths when present, then filesystem fallback under `previous/artifacts/**`.
## Previous-session requirement scanning
`CollectPreviousArtifactRequirements`:
- scans enabled configured artifacts only;
- includes canonical previous-session sources only;
- deduplicates by artifact key;
- merges required/optional references (`required` wins);
- records deterministic sorted source locations for diagnostics.
## Validation behavior
- transcript built-ins: JSON with top-level `segments` array;
- bounds built-in: valid JSON;
- configured and previous-session artifact files: non-empty text content.
## Failure behavior
- unsupported source -> source validation error
- known source unavailable -> `ErrSessionArtifactNotFound`
- configured source without catalog -> resolution error
- resolved file with invalid content -> validation error
- unsupported source or malformed canonical previous source: validation/resolution error;
- known source unavailable: `ErrSessionArtifactNotFound`;
- configured/previous canonical source without catalog: error;
- resolved invalid file content: validation error.
## Tests to inspect before changing
- `internal/artifacts/artifact_resolver_test.go`
- `internal/artifacts/catalog_test.go`
- `internal/artifacts/previous_requirements_test.go`
- `internal/stage/prepare_previous_test.go`
- `internal/stage/analyze_test.go`
- `internal/config/scriptorium_test.go`
## Architectural invariants
- built-in IDs are static and registry-backed
- configured artifact IDs are runtime-derived (`narratio.artifact.<name>`) and catalog-backed
- built-in/source resolution remains deterministic and validation-gated
- Built-in source IDs are static.
- Configured and previous-session source IDs are artifact-key based and validation-gated.
- Resolution behavior remains deterministic and manifest-aware.

View File

@@ -1,11 +1,11 @@
# Internal: Command Restore
## Purpose
Define the implemented `narratio restore` command contract: committed remote-state discovery, deterministic planning, safe file installation, conflict policy, and restore reporting.
Define the implemented `narratio restore` contract: committed remote-state discovery, deterministic plan classification, safe file install semantics, and restore reporting.
## Inputs and outputs
Inputs:
- CLI flags: `--config`, `--session`, `--session-id`, `--dry-run`, `--force`, `--include-audio`.
- CLI flags: `--config`, `--session`, `--session-id`, `--previous-session-id`, `--dry-run`, `--force`, `--include-audio`.
- Resolved/validated `pipeline.yml` and `session.yml`.
- Configured remote object store.
- Remote committed current-state markers (`current/run_id.txt`, `current/manifest.json`).
@@ -54,6 +54,21 @@ Does not own:
- existing local manifest is preserved if restored manifest validation/install fails.
- Non-dry-run report persists summary/action status metadata in `reports/restore-latest.json`.
Restore path scope:
- includes:
- `manifest.json`
- `transcripts/**`
- `artifacts/**`
- `previous/**`
- `audio/**` only when `--include-audio` is set
- excludes:
- `runs/**`
- `logs/**`
- `reports/**`
- `config/**`
- `inputs/**`
- remote `current/**` pointer files as local restore targets
## Skip and resume behavior
- Restore does not participate in stage skip/resume decisions.
- Restore provides durable local state so subsequent stage commands can resume or rerun based on restored manifest state.
@@ -80,7 +95,7 @@ Does not own:
- `current/run_id.txt` is the remote commit marker; restore must not infer committed state from incidental files.
- Local path mapping is traversal-safe and constrained to session root.
- Restore scope is deterministic and path-classified:
- include `manifest.json`, `transcripts/**`, `artifacts/**`
- include `manifest.json`, `transcripts/**`, `artifacts/**`, `previous/**`
- include `audio/**` only with `--include-audio`
- exclude `runs/**`, `logs/**`, `reports/**`, `config/**`, `inputs/**`
- Command remains standalone; no implicit `run --restore` behavior.

View File

@@ -3,31 +3,36 @@
## Purpose
Execute selected configured Scriptorium artifacts in deterministic dependency order and promote successful outputs to canonical session artifact paths.
## Inputs and Outputs
## Inputs and outputs
Inputs:
- configured artifact definitions from `pipeline.scriptorium.artifacts`
- selected artifact filter from runtime (`--artifacts`) when provided
- resolved artifact input sources declared per artifact (`inputs.*.source`)
- optional previous-session file inputs (`previous_session_artifact`)
- configured artifact definitions from `pipeline.scriptorium.artifacts`;
- selected artifact filter (`--artifacts`) when provided;
- resolved artifact sources from resolver/catalog.
Source types used by analyze:
- built-ins: `narratio.transcript.*`, `narratio.bounds.session`;
- configured artifacts: `narratio.artifact.<artifact_key>`;
- canonical previous-session artifacts: `narratio.previous_session.artifact.<artifact_key>`;
- legacy path-based previous-session source: `previous_session_artifact` (uses `inputs.*.path`).
Outputs:
- one promoted output file per executed configured artifact at that artifact's configured `output_path`
- stage metadata containing generated artifact entries and reused disabled-artifact entries
- promoted configured artifact files at each configured `output_path`;
- stage metadata (`generated_artifacts`, `reused_artifacts`, selected/order info).
## Boundaries
Owns:
- runtime artifact catalog construction for analyze execution
- selected-artifact planning and dependency ordering
- per-artifact input resolution, var resolution, timeout/render-debug resolution
- Scriptorium run/render invocation for each selected artifact
- run-local output generation and canonical promotion
- runtime artifact catalog construction;
- selected-artifact planning and dependency ordering;
- per-input resolution and required/optional handling;
- Scriptorium render/run invocation;
- run-local output generation and canonical promotion.
Does not own:
- transcript generation/processing stages
- archive promotion policy
- per-artifact resume semantics
- prepare-time previous-session hydration;
- object-store access for previous-session sources;
- archive promotion policy.
## Config Fields Used
## Config fields used
- `session.session_id`
- `session.campaign`
- `pipeline.workspace.root`
@@ -36,49 +41,41 @@ Does not own:
- `pipeline.scriptorium.timeout`
- `pipeline.scriptorium.render_debug`
- `pipeline.scriptorium.artifacts.<name>.*`
- `enabled`
- `depends_on`
- `prompt_id`
- `profile_id`
- `timeout`
- `output_path`
- `render_debug`
- `inputs`
- `vars`
## External Adapters Used
## External adapters used
- Scriptorium adapter:
- optional `RenderArtifact` (render debug)
- `RunArtifact` (artifact generation)
- optional `RenderArtifact` when render-debug is enabled;
- `RunArtifact` for artifact generation.
## State and Manifest Behavior
- If `pipeline.scriptorium` is absent, stage returns success metadata with `skipped=true`.
- If no artifacts are configured, stage returns success metadata with `skipped=true`.
- If zero artifacts are executable after `enabled` + `--artifacts` filtering, stage returns success metadata with `skipped=true`.
- Builds runtime catalog with built-ins and configured artifacts.
- Non-executable configured artifacts are marked available only when their configured output file exists and is valid on disk.
- Executes selected configured artifacts in topological order with deterministic tie-breaking.
- For each generated artifact, records metadata fields including `name`, `source_id`, `output_kind`, `path`, `prompt_id`, `profile_id`, and `provenance`.
- Reused disabled artifacts are recorded separately in `reused_artifacts` with provenance `filesystem.disabled_artifact_output`.
## State and manifest behavior
- If Scriptorium config is absent, or no artifacts are executable after filtering, analyze returns success metadata with `skipped=true`.
- Builds runtime catalog with built-ins and configured `narratio.artifact.<name>` entries.
- Non-executable configured artifacts may still be marked available from existing canonical output files.
- Resolves canonical previous-session sources from local prepared `previous/` cache:
- prefers manifest-backed previous input paths when present;
- may fall back to current-session `previous/` filesystem paths.
- Analyze does not call object storage for canonical previous-session source resolution.
- Required canonical previous-session input missing:
- fails with guidance to run `narratio run-stage --force prepare`.
- Optional missing sources are omitted from adapter input paths.
## Skip and Resume Behavior
## Skip and resume behavior
- Runner-level skip applies when analyze is already `succeeded` and `--force` is not set.
- Analyze remains stage-scoped for resume/skip; there is no per-artifact resume state.
- `--artifacts` filters which configured artifacts are executable when analyze runs; it does not imply `--force`.
- Analyze is stage-scoped for resume; no per-artifact manifest resume state.
- `--artifacts` filters executable artifacts but does not imply force rerun.
## Failure Behavior
- Fails on invalid dependency ordering, unavailable required configured inputs, invalid built-in input prerequisites, render/run adapter failures, validation-failed adapter results, or missing/empty outputs.
- Required configured dependency missing from catalog availability fails clearly before invocation.
- Optional missing inputs are omitted.
## Failure behavior
- Fails on dependency-order violations, missing required inputs, resolver validation failures, adapter errors, and missing/empty generated outputs.
- Required unavailable configured artifact source (`narratio.artifact.<name>`) fails before invocation.
- Required canonical previous-session source fails with prepare-rerun guidance.
## Tests to Inspect Before Changing
## Tests to inspect before changing
- `internal/stage/analyze_test.go`
- `internal/artifacts/catalog_test.go`
- `internal/artifacts/artifact_resolver_test.go`
- `internal/adapters/scriptorium/subprocess_test.go`
- `internal/app/restore_workflow_test.go`
## Architectural Invariants
- Configured artifacts are identified by `narratio.artifact.<name>` source IDs.
- Artifact-to-artifact references rely on explicit `depends_on` declarations validated in config.
- Generated analyze outputs are treated uniformly as Scriptorium artifacts.
- Successful outputs must exist and be non-empty before promotion.
## Architectural invariants
- Canonical previous-session behavior is local-cache only during analyze.
- Generated outputs are validated and promoted before stage success is recorded.
- Resolver/catalog decisions stay deterministic and validation-gated.

View File

@@ -1,17 +1,19 @@
# Stage: archive
## Purpose
Publish run records and promoted session artifacts to object storage, then atomically advance the remote current pointer.
Publish durable run/session state to object storage, then atomically advance remote current state.
## Inputs and Outputs
Inputs:
- session manifest and prerequisite stage records
- run root contents under `runs/{run_id}/`
- promotion rules with artifact `source` IDs and archive `dest` paths (`archive.promote_artifacts`)
- session-level `previous/**` cache files when present
Outputs:
- uploaded run files under `{session_prefix}/runs/{run_id}/...`
- uploaded promoted artifacts under `{session_prefix}/...`
- uploaded session previous-cache files under `{session_prefix}/previous/...` when present
- `{session_prefix}/current/manifest.json`
- `{session_prefix}/current/run_id.txt` written last
@@ -21,6 +23,7 @@ Owns:
- Prerequisite stage success enforcement
- Run file collection and upload (excluding `audio/`)
- Promotion rule resolution and upload
- Session previous-cache file collection/upload
- Commit pointer publish order
Does not own:
@@ -43,8 +46,10 @@ Does not own:
## State and Manifest Behavior
- Requires `prepare`, `transcribe`, `merge`, `polish`, `normalize`, `trim`, and `analyze` status `succeeded`.
- Resolves bucket/prefix from manifest identity first, then config fallback.
- Uploads session `previous/**` files as durable session state when the local `previous/` directory exists.
- Writes metadata including:
- upload counts/paths
- `previous_files_uploaded` and `previous_uploaded_paths`
- `current_manifest_key`
- `current_run_id_key`
- `current_pointer_written`
@@ -64,5 +69,6 @@ Does not own:
## Architectural Invariants
- Run upload excludes `audio/` subtree.
- Session `previous/**` is archiveable durable input/provenance state, not run-local output.
- `current/manifest.json` uploads before `current/run_id.txt`.
- `current/run_id.txt` is the remote publish commit marker.

View File

@@ -1,40 +1,49 @@
# Stage: prepare
## Purpose
Materialize all required session inputs into canonical local workspace paths and record input provenance in the session manifest.
Materialize canonical current-session input state and provenance before downstream stages run.
## Inputs and Outputs
Prepare owns:
- local input file materialization (`inputs/**`);
- audio input materialization (`audio/**`);
- previous-session cache hydration (`previous/**`) for canonical previous-session artifact sources.
## Inputs and outputs
Inputs:
- `session.yml` (resolved session config)
- `pipeline.resolved.yml` (materialized from resolved pipeline config)
- `speakers.yml`
- `autocorrect.yml`
- `glossary.yml`
- resolved config/session (`pipeline.yml`, `session.yml`);
- session-local input files (`speakers`, `autocorrect`, `glossary`);
- audio source:
- local (`session.inputs.audio_dir` or `session.inputs.audio_files`), or
- S3 (`session.inputs.audio_s3.prefix`)
- local: `session.inputs.audio_dir` or `session.inputs.audio_files`;
- S3: `session.inputs.audio_s3.prefix`;
- configured enabled Scriptorium artifact inputs (for previous-session requirement scanning);
- remote previous-session current archive state when previous hydration is required.
Outputs:
- `inputs/session.yml`
- `inputs/pipeline.resolved.yml`
- `inputs/speakers.yml`
- `inputs/autocorrect.yml`
- `inputs/glossary.yml`
- `audio/*.flac` in session workdir
- `manifest.Inputs` records with checksums and source metadata
- `inputs/session.yml`;
- `inputs/pipeline.resolved.yml`;
- `inputs/speakers.yml`;
- `inputs/autocorrect.yml`;
- `inputs/glossary.yml`;
- `audio/*.flac` in canonical session `audio/`;
- optional `previous/manifest.json`;
- optional `previous/artifacts/**`;
- deterministic `manifest.Inputs` records with checksums and provenance metadata.
## Boundaries
Owns:
- Input path resolution and validation
- Local copy/materialization of configs and audio files
- S3 audio download to run-scoped spool, then copy into work audio dir
- input path resolution and materialization;
- S3 audio list/download/copy flow;
- previous-session artifact requirement collection from enabled configured artifacts;
- previous cache lifecycle when requirements exist (clear and rehydrate managed `previous/` state).
Does not own:
- Transcript generation/processing
- Archive publish behavior
- transcript or artifact generation;
- analyze-stage source resolution;
- archive commit behavior.
## Config Fields Used
## Config fields used
- `session.session_id`
- `session.previous_session_id`
- `session.campaign`
- `session.inputs.speakers_file`
- `session.inputs.autocorrect_file`
@@ -46,29 +55,56 @@ Does not own:
- `pipeline.spool.root`
- `pipeline.storage.s3.bucket`
- `pipeline.storage.s3.root_prefix`
- `pipeline.scriptorium.artifacts.<name>.enabled`
- `pipeline.scriptorium.artifacts.<name>.inputs.<key>.source`
- `pipeline.scriptorium.artifacts.<name>.inputs.<key>.required`
## External Adapters Used
- Object storage backend (`env.ObjectStore`) for S3 audio list/download when `audio_s3` is configured.
## External adapters used
- `storage.ObjectStore` for:
- S3 audio listing/downloads;
- previous-session current pointer/manifest/artifact object checks and downloads.
## State and Manifest Behavior
## State and manifest behavior
- Ensures workspace layout exists.
- Writes resolved config and input files to canonical `inputs/` paths.
- Records all prepared inputs into `manifest.Inputs` (sorted deterministically by kind/path).
- For S3 audio, records `S3Bucket`, `S3Key`, `S3Size`, `S3ETag`, and `SpoolPath` in each audio input record.
- Materializes canonical input files and audio files.
- Scans enabled configured artifact inputs for canonical sources:
- `narratio.previous_session.artifact.<artifact_key>`
- If one or more canonical previous-session requirements exist:
- clears managed `previous/` state;
- hydrates required/optional previous artifacts from the configured previous sessions committed archive current state;
- writes `previous/manifest.json` and hydrated `previous/artifacts/**`;
- records hydrated previous inputs in `manifest.Inputs` with source `previous_session_archive.current`.
- If no canonical previous-session requirements exist, prepare does not manage `previous/`.
- `manifest.Inputs` is sorted deterministically by `(kind, path)`.
## Skip and Resume Behavior
- Runner-level skip applies when stage already `succeeded` and `--force` is not set.
- Stage itself is deterministic/idempotent for unchanged inputs (`copyFileIfChanged`, `writeBytesIfChanged`).
## Required and optional previous-session behavior
- `previous_session_id` unset:
- if any referenced previous artifact is required: fail;
- if all referenced previous artifacts are optional: continue and omit them.
- Previous session archive current pointer or manifest missing:
- if any referenced previous artifact is required: fail;
- if all referenced previous artifacts are optional: continue and omit missing ones.
- Missing required previous artifact object: fail.
- Missing optional previous artifact object: omit.
- Downloaded previous artifacts must validate as non-empty files.
## Failure Behavior
- Fails on missing required files, invalid audio source combinations, no discoverable `.flac` files, duplicate audio basenames, missing object store for S3 mode, or S3 list/download failures.
## Skip and resume behavior
- Runner-level skip remains authoritative:
- if `prepare` already succeeded and run is not forced, `prepare` does not run and no hydration/download occurs.
- If `prepare` runs (including with `--force`), it owns managed `previous/` state for canonical previous-session inputs.
## Tests to Inspect Before Changing
## Failure behavior
- Fails on missing required input files, invalid audio-source combinations, empty/duplicate audio inputs, missing object store for S3 modes, and remote access/download/validation errors.
- For required canonical previous-session inputs, analyze-time missing-input guidance is to rerun:
- `narratio run-stage --force prepare`
## Tests to inspect before changing
- `internal/stage/prepare_test.go`
- `internal/app/session_cli_test.go`
- `internal/config/load_validate_test.go`
- `internal/stage/prepare_previous_test.go`
- `internal/artifacts/previous_requirements_test.go`
- `internal/app/runner_test.go`
## Architectural Invariants
## Architectural invariants
- `audio_dir`/`audio_files` and `audio_s3` are mutually exclusive.
- Audio files must be `.flac`.
- Canonical `inputs/*` and `audio/*` paths are the durable source for downstream stages.
- Storage keys are computed by callers using archive/path helpers; storage adapter receives explicit keys.
- `prepare` is the only stage that hydrates canonical previous-session cache state.

View File

@@ -17,7 +17,8 @@ Outputs:
## Boundaries
Owns:
- Session-level path layout (`inputs/`, `audio/`, `transcripts/`, `artifacts/`, `reports/`, `logs/`, `config/`, `current/`, `runs/`)
- Session-level path layout (`inputs/`, `audio/`, `transcripts/`, `artifacts/`, `reports/`, `logs/`, `config/`, `current/`, `runs/`, `previous/`)
- `previous/manifest.json` and `previous/artifacts/**` are reserved for prepared previous-session state
- Run-local stage sandbox layout under `runs/{run_id}/{stage}/`
- Session lock acquisition/release (`.lock`)
@@ -65,4 +66,5 @@ None directly in this subsystem. Stages may use object storage adapters and then
- Session root is campaign-aware: `{workspace.root}/work/{campaign}/{session_id}`.
- Run roots are always nested: `runs/{run_id}` under the session root.
- Run-local output promotion must end in canonical session paths.
- `previous/**` is session-durable state and must not be treated as run-local output scratch state.
- Cleanup only targets run-scoped directories and must never delete configured root directories.

View File

@@ -48,7 +48,7 @@ Restore source-of-truth:
- remote current manifest: `current/manifest.json`
Restore default scope:
- includes `manifest.json`, `transcripts/**`, `artifacts/**`
- includes `manifest.json`, `transcripts/**`, `artifacts/**`, `previous/**`
- includes `audio/**` only with `--include-audio`
- excludes `runs/**`, `logs/**`, `reports/**`, `config/**`, `inputs/**`, and `current/**` (except remote `current/manifest.json` as source)
@@ -67,6 +67,7 @@ Canonical session directories:
- `audio/`
- `transcripts/`
- `artifacts/`
- `previous/`
- `reports/`
- `logs/`
- `config/`
@@ -99,6 +100,12 @@ Configured artifact source reuse:
- accepted on `run`, `resume`, and `run-stage analyze`.
- filters analyze execution only; does not force stage rerun.
Canonical previous-session input behavior:
- canonical sources use `narratio.previous_session.artifact.<artifact_key>`.
- these inputs are hydrated by `prepare`, not `analyze`.
- if analyze fails due to missing canonical previous cache, rerun:
- `narratio run-stage --session-id <id> --force prepare`
## Remote archive layout and publish contract
When archive is enabled and run upload is enabled, archive publishes under:

View File

@@ -2,7 +2,7 @@
## Status
Planned.
Completed.
This roadmap describes the implementation strategy for first-class previous-session artifact support in Narratio. It belongs under `docs/roadmap/previous.md` until the feature is implemented. After implementation, current behavior should be documented in the appropriate user-facing and internal documentation files, and this roadmap should be removed or marked complete according to the documentation policy.

View File

@@ -133,9 +133,7 @@ scriptorium:
source: narratio.transcript.trimmed
required: true
previous_recap:
source: previous_session_artifact
artifact: session_recap
path: ""
source: narratio.previous_session.artifact.session_recap
required: false
vars:
session_id: true

View File

@@ -81,8 +81,7 @@ scriptorium:
source: narratio.transcript.trimmed
required: true
previous_recap:
source: previous_session_artifact
artifact: session_recap
source: narratio.previous_session.artifact.session_recap
required: false
vars:
session_id: true

View File

@@ -22,10 +22,12 @@ func Plan(ctx context.Context, args []string, out io.Writer) error {
var pipelinePath string
var sessionPath string
var sessionID string
var previousSessionID string
var force bool
fs.StringVar(&pipelinePath, "config", "", "path to pipeline.yml (optional; defaults searched)")
fs.StringVar(&sessionPath, "session", "", "path to session.yml")
fs.StringVar(&sessionID, "session-id", "", "session identifier for session.yml templates")
fs.StringVar(&previousSessionID, "previous-session-id", "", "previous session identifier for session.yml templates")
fs.BoolVar(&force, "force", false, "force stage execution (reserved for future behavior)")
if err := fs.Parse(args); err != nil {
@@ -44,7 +46,8 @@ func Plan(ctx context.Context, args []string, out io.Writer) error {
}
cfg, err := config.LoadWithSessionOptions(resolvedPipelinePath, resolvedSessionPath, config.SessionLoadOptions{
SessionID: sessionID,
SessionID: sessionID,
PreviousSessionID: previousSessionID,
})
if err != nil {
return fmt.Errorf("plan: %w", err)

View File

@@ -82,6 +82,7 @@ func TestPostArchiveCleanupWorkdirOnly(t *testing.T) {
assertExists(t, cfg.Pipeline.Workspace.Root)
assertExists(t, seed.otherRunDir)
assertExists(t, seed.previousCachePath)
assertMissing(t, seed.runWorkDir)
assertExists(t, seed.spoolAudioDir)
}
@@ -98,6 +99,7 @@ func TestPostArchiveCleanupBothPolicies(t *testing.T) {
assertMissing(t, seed.spoolAudioDir)
assertMissing(t, seed.runWorkDir)
assertExists(t, seed.otherRunDir)
assertExists(t, seed.previousCachePath)
}
func TestPostArchiveCleanupNotRunWhenArchiveFails(t *testing.T) {
@@ -256,11 +258,12 @@ func TestPostArchiveCleanupNotRunWhenCurrentPointerUploadFails(t *testing.T) {
}
type cleanupSeed struct {
runWorkDir string
otherRunDir string
spoolAudioDir string
localSourceAudio string
sessionPrefix string
runWorkDir string
otherRunDir string
spoolAudioDir string
localSourceAudio string
previousCachePath string
sessionPrefix string
}
func cleanupFixtureConfig(t *testing.T) (*config.Config, cleanupSeed) {
@@ -274,11 +277,18 @@ func cleanupFixtureConfig(t *testing.T) (*config.Config, cleanupSeed) {
runWorkDir := artifacts.SessionRunRootForCampaign(cfg.Pipeline.Workspace.Root, cfg.Session.Campaign, cfg.Session.SessionID, runID)
otherRunDir := artifacts.SessionRunRootForCampaign(cfg.Pipeline.Workspace.Root, cfg.Session.Campaign, cfg.Session.SessionID, "20260516T010204Z-5e6f7a8b")
spoolAudioDir := artifacts.SessionSpoolAudioDir(cfg.Pipeline.Spool.Root, cfg.Session.Campaign, cfg.Session.SessionID, runID)
previousCachePath := artifacts.SessionPreviousArtifactPathForCampaign(
cfg.Pipeline.Workspace.Root,
cfg.Session.Campaign,
cfg.Session.SessionID,
"session_recap.md",
)
mustWriteFile(t, filepath.Join(runWorkDir, "manifest.json"), "{}\n")
mustWriteFile(t, filepath.Join(runWorkDir, "logs", "stage.log"), "log\n")
mustWriteFile(t, filepath.Join(otherRunDir, "logs", "stage.log"), "other\n")
mustWriteFile(t, filepath.Join(spoolAudioDir, "speaker.flac"), "flac\n")
mustWriteFile(t, previousCachePath, "# previous recap\n")
localSourceAudio := filepath.Join(filepath.Dir(cfg.SessionPath), "audio", "alice.flac")
mustWriteFile(t, localSourceAudio, "source\n")
@@ -301,11 +311,12 @@ func cleanupFixtureConfig(t *testing.T) (*config.Config, cleanupSeed) {
}
return cfg, cleanupSeed{
runWorkDir: runWorkDir,
otherRunDir: otherRunDir,
spoolAudioDir: spoolAudioDir,
localSourceAudio: localSourceAudio,
sessionPrefix: seed.S3SessionPrefix,
runWorkDir: runWorkDir,
otherRunDir: otherRunDir,
spoolAudioDir: spoolAudioDir,
localSourceAudio: localSourceAudio,
previousCachePath: previousCachePath,
sessionPrefix: seed.S3SessionPrefix,
}
}

View File

@@ -28,17 +28,19 @@ func Restore(ctx context.Context, args []string, out io.Writer) error {
var pipelinePath string
var sessionPath string
var sessionID string
var previousSessionID string
var dryRun bool
var force bool
var includeAudio bool
fs.StringVar(&pipelinePath, "config", "", "path to pipeline.yml (optional; defaults searched)")
fs.StringVar(&sessionPath, "session", "", "path to session.yml")
fs.StringVar(&sessionID, "session-id", "", "session identifier for session.yml templates")
fs.StringVar(&previousSessionID, "previous-session-id", "", "previous session identifier for session.yml templates")
fs.BoolVar(&dryRun, "dry-run", false, "plan restore actions without writing local files")
fs.BoolVar(&force, "force", false, "overwrite local conflicts with remote state")
fs.BoolVar(&includeAudio, "include-audio", false, "include archived session-level audio objects")
fs.Usage = func() {
_, _ = fmt.Fprintln(out, "Usage: narratio restore [--config <path>] [--session <path>] [--session-id <value>] [--dry-run] [--force] [--include-audio]")
_, _ = fmt.Fprintln(out, "Usage: narratio restore [--config <path>] [--session <path>] [--session-id <value>] [--previous-session-id <value>] [--dry-run] [--force] [--include-audio]")
_, _ = fmt.Fprintln(out)
_, _ = fmt.Fprintln(out, "Flags:")
fs.PrintDefaults()
@@ -63,7 +65,8 @@ func Restore(ctx context.Context, args []string, out io.Writer) error {
}
cfg, err := config.LoadWithSessionOptions(resolvedPipelinePath, resolvedSessionPath, config.SessionLoadOptions{
SessionID: sessionID,
SessionID: sessionID,
PreviousSessionID: previousSessionID,
})
if err != nil {
return fmt.Errorf("restore: %w", err)

View File

@@ -87,6 +87,33 @@ func TestExecuteRestoreIncludeAudioRestoresAudio(t *testing.T) {
}
}
func TestExecuteRestoreRestoresPreviousCacheWhenPresent(t *testing.T) {
workspaceRoot := t.TempDir()
pipelinePath, sessionPath := writeValidConfigFiles(t, workspaceRoot)
fake := &storage.FakeBackend{}
cfg, sessionPrefix, _, _ := seedRestoreCommittedState(t, fake, pipelinePath, sessionPath)
seedRestoreObject(fake, sessionPrefix+"previous/manifest.json", []byte(`{"session_id":"2026-04-26"}`))
seedRestoreObject(fake, sessionPrefix+"previous/artifacts/session_recap.md", []byte("# previous recap\n"))
restoreWithStoreAndRealPhases(t, fake)
var stdout bytes.Buffer
var stderr bytes.Buffer
code := Execute([]string{"restore", "--config", pipelinePath, "--session", sessionPath}, &stdout, &stderr)
if code != 0 {
t.Fatalf("exit code = %d, want 0; stderr=%q", code, stderr.String())
}
sessionRoot := artifacts.SessionWorkDirForCampaign(workspaceRoot, cfg.Session.Campaign, cfg.Session.SessionID)
mustReadEquals(t, filepath.Join(sessionRoot, "previous", "manifest.json"), `{"session_id":"2026-04-26"}`)
mustReadEquals(t, filepath.Join(sessionRoot, "previous", "artifacts", "session_recap.md"), "# previous recap\n")
report := mustReadRestoreReport(t, filepath.Join(sessionRoot, "reports", "restore-latest.json"))
if report.Execution.Downloaded != 3 {
t.Fatalf("report execution.downloaded = %d, want 3", report.Execution.Downloaded)
}
}
func TestExecuteRestoreConflictWithoutForceDoesNotOverwrite(t *testing.T) {
workspaceRoot := t.TempDir()
pipelinePath, sessionPath := writeValidConfigFiles(t, workspaceRoot)
@@ -145,6 +172,28 @@ func TestExecuteRestoreForceOverwritesDifferingFile(t *testing.T) {
}
}
func TestExecuteRestoreForceOverwritesDifferingPreviousCacheFile(t *testing.T) {
workspaceRoot := t.TempDir()
pipelinePath, sessionPath := writeValidConfigFiles(t, workspaceRoot)
fake := &storage.FakeBackend{}
cfg, sessionPrefix, _, _ := seedRestoreCommittedState(t, fake, pipelinePath, sessionPath)
seedRestoreObject(fake, sessionPrefix+"previous/artifacts/session_recap.md", []byte("# remote previous recap\n"))
sessionRoot := artifacts.SessionWorkDirForCampaign(workspaceRoot, cfg.Session.Campaign, cfg.Session.SessionID)
mustWriteTestFile(t, filepath.Join(sessionRoot, "previous", "artifacts", "session_recap.md"), "# local previous recap\n")
restoreWithStoreAndRealPhases(t, fake)
var stdout bytes.Buffer
var stderr bytes.Buffer
code := Execute([]string{"restore", "--config", pipelinePath, "--session", sessionPath, "--force"}, &stdout, &stderr)
if code != 0 {
t.Fatalf("exit code = %d, want 0; stderr=%q", code, stderr.String())
}
mustReadEquals(t, filepath.Join(sessionRoot, "previous", "artifacts", "session_recap.md"), "# remote previous recap\n")
}
func TestExecuteRestoreLockConflictFailsAndWritesNothing(t *testing.T) {
workspaceRoot := t.TempDir()
pipelinePath, sessionPath := writeValidConfigFiles(t, workspaceRoot)

View File

@@ -194,6 +194,9 @@ func restoreLocalRelativePathForKey(sessionPrefix, currentManifestKey, key strin
if cleanRel == config.PathArtifactsDirSegment || strings.HasPrefix(cleanRel, config.PathArtifactsDirSegment+"/") {
return cleanRel, true, nil
}
if cleanRel == config.PathPreviousDirSegment || strings.HasPrefix(cleanRel, config.PathPreviousDirSegment+"/") {
return cleanRel, true, nil
}
if includeAudio && (cleanRel == config.PathAudioDirSegment || strings.HasPrefix(cleanRel, config.PathAudioDirSegment+"/")) {
return cleanRel, true, nil
}

View File

@@ -58,6 +58,27 @@ func TestRestorePlanIncludeAudio(t *testing.T) {
}
}
func TestRestorePlanIncludesPreviousCacheByDefault(t *testing.T) {
cfg := restorePlanConfig(t)
current := restorePlanCurrentState(t, cfg)
store := &storage.FakeBackend{}
seedRestoreObject(store, current.CurrentManifestKey, []byte(`{"session_id":"2026-05-03"}`))
seedRestoreObject(store, current.SessionPrefix+"previous/manifest.json", []byte(`{"session_id":"2026-04-26"}`))
seedRestoreObject(store, current.SessionPrefix+"previous/artifacts/session_recap.md", []byte("# previous recap\n"))
plan, err := buildRestorePlan(context.Background(), cfg, current, store, RestorePlanOptions{})
if err != nil {
t.Fatalf("buildRestorePlan() error = %v", err)
}
got := actionRelPaths(plan.Actions)
want := []string{"manifest.json", "previous/artifacts/session_recap.md", "previous/manifest.json"}
if !reflect.DeepEqual(got, want) {
t.Fatalf("action local paths = %#v, want %#v", got, want)
}
}
func TestRestorePlanClassifiesSameAndConflict(t *testing.T) {
cfg := restorePlanConfig(t)
current := restorePlanCurrentState(t, cfg)

View File

@@ -155,6 +155,116 @@ func TestRestoreThenRunStageForceAnalyzeUsesRestoredDurableState(t *testing.T) {
}
}
func TestRestoreThenAnalyzeUsesRestoredPreviousCacheWithoutObjectStore(t *testing.T) {
workspaceRoot := t.TempDir()
pipelinePath, sessionPath := writeValidConfigFiles(t, workspaceRoot)
appendRestoreWorkflowScriptoriumConfig(t, pipelinePath, `
scriptorium:
binary: scriptorium
artifacts:
session_recap:
enabled: true
prompt_id: dnd.session_recap
output_path: artifacts/session_recap.md
inputs:
transcript:
source: narratio.transcript.trimmed
required: true
previous_recap:
source: narratio.previous_session.artifact.session_recap
required: true
`)
fakeStore := &storage.FakeBackend{}
cfg, sessionPrefix, manifestKey, runIDKey := seedRestoreCommittedState(t, fakeStore, pipelinePath, sessionPath)
seedRestoreObject(fakeStore, runIDKey, []byte("20260519T010203Z-a1b2c3d4\n"))
seedRestoreObject(fakeStore, manifestKey, restoreWorkflowManifestJSON(t, cfg.Session.SessionID, cfg.Session.Campaign))
seedRestoreObject(fakeStore, sessionPrefix+"transcripts/trimmed.json", []byte(`{"segments":[]}`+"\n"))
seedRestoreObject(fakeStore, sessionPrefix+"previous/manifest.json", []byte(`{"session_id":"2026-04-26"}`))
seedRestoreObject(fakeStore, sessionPrefix+"previous/artifacts/session_recap.md", []byte("# previous recap\n"))
restoreWithStoreAndRealPhases(t, fakeStore)
var stdout bytes.Buffer
var stderr bytes.Buffer
restoreCode := Execute(
[]string{
"restore",
"--config", pipelinePath,
"--session", sessionPath,
"--session-id", cfg.Session.SessionID,
},
&stdout,
&stderr,
)
if restoreCode != 0 {
t.Fatalf("restore exit code = %d, want 0; stderr=%q", restoreCode, stderr.String())
}
if stderr.Len() != 0 {
t.Fatalf("restore stderr = %q, want empty", stderr.String())
}
sessionRoot := artifacts.SessionWorkDirForCampaign(workspaceRoot, cfg.Session.Campaign, cfg.Session.SessionID)
mustReadEquals(t, filepath.Join(sessionRoot, "transcripts", "trimmed.json"), `{"segments":[]}`+"\n")
mustReadEquals(t, filepath.Join(sessionRoot, "previous", "manifest.json"), `{"session_id":"2026-04-26"}`)
mustReadEquals(t, filepath.Join(sessionRoot, "previous", "artifacts", "session_recap.md"), "# previous recap\n")
scriptoriumFake := &scriptorium.FakeRunner{}
origExecuteStagesFn := executeStagesFn
origObjectStoreFn := newObjectStoreFromConfigFn
objectStoreConstructed := false
t.Cleanup(func() {
executeStagesFn = origExecuteStagesFn
newObjectStoreFromConfigFn = origObjectStoreFn
})
executeStagesFn = func(ctx context.Context, cfg *config.Config, stages []stage.Stage, opts RunOptions) (*RunSummary, error) {
if opts.Env == nil {
opts.Env = &Env{}
}
opts.Env.Scriptorium = scriptoriumFake
return executeStages(ctx, cfg, stages, opts)
}
newObjectStoreFromConfigFn = func(context.Context, *config.Config) (storage.ObjectStore, error) {
objectStoreConstructed = true
return nil, context.Canceled
}
stdout.Reset()
stderr.Reset()
runStageCode := Execute(
[]string{
"run-stage",
"--config", pipelinePath,
"--session", sessionPath,
"--session-id", cfg.Session.SessionID,
"--force",
"--artifacts", "session_recap",
"analyze",
},
&stdout,
&stderr,
)
if runStageCode != 0 {
t.Fatalf("run-stage exit code = %d, want 0; stderr=%q", runStageCode, stderr.String())
}
if stderr.Len() != 0 {
t.Fatalf("run-stage stderr = %q, want empty", stderr.String())
}
if objectStoreConstructed {
t.Fatal("analyze run-stage should not construct object store for previous-session input resolution")
}
if len(scriptoriumFake.RunRequests) != 1 {
t.Fatalf("scriptorium run requests = %d, want 1", len(scriptoriumFake.RunRequests))
}
req := scriptoriumFake.RunRequests[0]
if got := req.InputPaths["transcript"]; got != filepath.Join(sessionRoot, "transcripts", "trimmed.json") {
t.Fatalf("transcript input = %q, want trimmed transcript path", got)
}
if got := req.InputPaths["previous_recap"]; got != filepath.Join(sessionRoot, "previous", "artifacts", "session_recap.md") {
t.Fatalf("previous_recap input = %q, want restored previous cache path", got)
}
}
func restoreWorkflowManifestJSON(t *testing.T, sessionID, campaign string) []byte {
t.Helper()
store := &manifest.LocalStore{}
@@ -176,3 +286,15 @@ func restoreWorkflowManifestJSON(t *testing.T, sessionID, campaign string) []byt
}
return data
}
func appendRestoreWorkflowScriptoriumConfig(t *testing.T, pipelinePath, extra string) {
t.Helper()
f, err := os.OpenFile(pipelinePath, os.O_APPEND|os.O_WRONLY, 0)
if err != nil {
t.Fatalf("open pipeline config for append: %v", err)
}
defer f.Close()
if _, err := f.WriteString(extra); err != nil {
t.Fatalf("append pipeline config: %v", err)
}
}

View File

@@ -19,11 +19,13 @@ func Resume(ctx context.Context, args []string, out io.Writer) error {
var pipelinePath string
var sessionPath string
var sessionID string
var previousSessionID string
var force bool
var selectedArtifacts artifactSelectionFlag
fs.StringVar(&pipelinePath, "config", "", "path to pipeline.yml (optional; defaults searched)")
fs.StringVar(&sessionPath, "session", "", "path to session.yml")
fs.StringVar(&sessionID, "session-id", "", "session identifier for session.yml templates")
fs.StringVar(&previousSessionID, "previous-session-id", "", "previous session identifier for session.yml templates")
fs.BoolVar(&force, "force", false, "force stage execution")
fs.Var(&selectedArtifacts, "artifacts", "artifact names to execute during analyze (comma-separated or repeatable)")
@@ -43,7 +45,8 @@ func Resume(ctx context.Context, args []string, out io.Writer) error {
}
cfg, err := config.LoadWithSessionOptions(resolvedPipelinePath, resolvedSessionPath, config.SessionLoadOptions{
SessionID: sessionID,
SessionID: sessionID,
PreviousSessionID: previousSessionID,
})
if err != nil {
return fmt.Errorf("resume: %w", err)

View File

@@ -17,11 +17,13 @@ func Run(ctx context.Context, args []string, out io.Writer) error {
var pipelinePath string
var sessionPath string
var sessionID string
var previousSessionID string
var force bool
var selectedArtifacts artifactSelectionFlag
fs.StringVar(&pipelinePath, "config", "", "path to pipeline.yml (optional; defaults searched)")
fs.StringVar(&sessionPath, "session", "", "path to session.yml")
fs.StringVar(&sessionID, "session-id", "", "session identifier for session.yml templates")
fs.StringVar(&previousSessionID, "previous-session-id", "", "previous session identifier for session.yml templates")
fs.BoolVar(&force, "force", false, "force stage execution (reserved for future behavior)")
fs.Var(&selectedArtifacts, "artifacts", "artifact names to execute during analyze (comma-separated or repeatable)")
@@ -41,7 +43,8 @@ func Run(ctx context.Context, args []string, out io.Writer) error {
}
cfg, err := config.LoadWithSessionOptions(resolvedPipelinePath, resolvedSessionPath, config.SessionLoadOptions{
SessionID: sessionID,
SessionID: sessionID,
PreviousSessionID: previousSessionID,
})
if err != nil {
return fmt.Errorf("run: %w", err)

View File

@@ -17,11 +17,13 @@ func RunStage(ctx context.Context, args []string, out io.Writer) error {
var pipelinePath string
var sessionPath string
var sessionID string
var previousSessionID string
var force bool
var selectedArtifacts artifactSelectionFlag
fs.StringVar(&pipelinePath, "config", "", "path to pipeline.yml (optional; defaults searched)")
fs.StringVar(&sessionPath, "session", "", "path to session.yml")
fs.StringVar(&sessionID, "session-id", "", "session identifier for session.yml templates")
fs.StringVar(&previousSessionID, "previous-session-id", "", "previous session identifier for session.yml templates")
fs.BoolVar(&force, "force", false, "force stage execution (reserved for future behavior)")
fs.Var(&selectedArtifacts, "artifacts", "artifact names to execute during analyze (comma-separated or repeatable)")
@@ -54,7 +56,8 @@ func RunStage(ctx context.Context, args []string, out io.Writer) error {
}
cfg, err := config.LoadWithSessionOptions(resolvedPipelinePath, resolvedSessionPath, config.SessionLoadOptions{
SessionID: sessionID,
SessionID: sessionID,
PreviousSessionID: previousSessionID,
})
if err != nil {
return fmt.Errorf("run-stage: %w", err)

View File

@@ -555,6 +555,12 @@ func needsObjectStoreForRun(cfg *config.Config, stages []stage.Stage) bool {
if cfg.Session.Inputs.AudioS3 != nil && stageRequested("prepare") {
return true
}
if stageRequested("prepare") {
requirements := artifacts.CollectPreviousArtifactRequirements(configuredScriptoriumArtifacts(cfg))
if len(requirements) > 0 && strings.TrimSpace(cfg.Session.PreviousSessionID) != "" {
return true
}
}
if !stageRequested("archive") {
return false
}
@@ -569,3 +575,10 @@ func needsObjectStoreForRun(cfg *config.Config, stages []stage.Stage) bool {
}
return true
}
func configuredScriptoriumArtifacts(cfg *config.Config) map[string]config.ScriptoriumArtifactConfig {
if cfg == nil || cfg.Pipeline == nil || cfg.Pipeline.Scriptorium == nil {
return nil
}
return cfg.Pipeline.Scriptorium.Artifacts
}

View File

@@ -199,6 +199,55 @@ func TestExecuteStagesAnalyzeOutputsPersistAsScriptoriumArtifacts(t *testing.T)
}
}
func TestNeedsObjectStoreForRunPrepareWithPreviousRequirements(t *testing.T) {
tests := []struct {
name string
previousSessionID string
want bool
}{
{
name: "previous session configured",
previousSessionID: "2026-05-10",
want: true,
},
{
name: "previous session missing",
previousSessionID: "",
want: false,
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
cfg := &config.Config{
Pipeline: &config.PipelineConfig{
Scriptorium: &config.ScriptoriumConfig{
Artifacts: map[string]config.ScriptoriumArtifactConfig{
"session_recap": {
Enabled: true,
Inputs: map[string]config.ScriptoriumInputConfig{
"previous_recap": {
Source: "narratio.previous_session.artifact.session_recap",
Required: true,
},
},
},
},
},
},
Session: &config.SessionConfig{
PreviousSessionID: tt.previousSessionID,
},
}
got := needsObjectStoreForRun(cfg, []stage.Stage{countingStage{name: "prepare", runs: new(int)}})
if got != tt.want {
t.Fatalf("needsObjectStoreForRun() = %v, want %v", got, tt.want)
}
})
}
}
func TestExecuteStagesArchiveFailsWhenRequiredRecapPromotionMissingForSelectedArtifacts(t *testing.T) {
cfg := testConfig(t)
cfg.Pipeline.Storage.S3 = &config.StorageS3Config{

View File

@@ -9,11 +9,12 @@ import (
"testing"
)
func TestPlanUsesDiscoveredSessionTemplateWithSessionID(t *testing.T) {
func TestPlanUsesDiscoveredSessionTemplateWithSessionIDs(t *testing.T) {
workspaceRoot := t.TempDir()
pipelinePath, sessionPath := writeValidConfigFiles(t, workspaceRoot)
sessionTemplate := `session_id: "{{ session_id }}"
previous_session_id: "{{ previous_session_id }}"
campaign: sample-campaign
inputs:
audio_dir: ./audio
@@ -36,7 +37,11 @@ inputs:
t.Cleanup(func() { _ = os.Chdir(originalWD) })
var out bytes.Buffer
if err := Plan(context.Background(), []string{"--config", pipelinePath, "--session-id", "2026-04-04"}, &out); err != nil {
if err := Plan(context.Background(), []string{
"--config", pipelinePath,
"--session-id", "2026-04-04",
"--previous-session-id", "2026-03-28",
}, &out); err != nil {
t.Fatalf("Plan() error = %v", err)
}
if !strings.Contains(out.String(), "narratio plan: workdir prepared") {
@@ -58,6 +63,38 @@ func TestPlanFailsWhenSessionIDMismatchesConcreteSession(t *testing.T) {
}
}
func TestPlanFailsWhenPreviousSessionIDMismatchesConcreteSession(t *testing.T) {
workspaceRoot := t.TempDir()
pipelinePath, sessionPath := writeValidConfigFiles(t, workspaceRoot)
sessionYAML := `session_id: 2026-05-03
previous_session_id: 2026-04-26
campaign: sample-campaign
inputs:
audio_dir: ./audio
speakers_file: ./speakers.yml
autocorrect_file: ./autocorrect.yml
glossary_file: ./glossary.yml
`
if err := os.WriteFile(sessionPath, []byte(sessionYAML), 0o644); err != nil {
t.Fatalf("write session.yml: %v", err)
}
var out bytes.Buffer
err := Plan(context.Background(), []string{
"--config", pipelinePath,
"--session", sessionPath,
"--session-id", "2026-05-03",
"--previous-session-id", "2026-04-25",
}, &out)
if err == nil {
t.Fatal("expected error, got nil")
}
if !strings.Contains(err.Error(), "previous_session_id mismatch") {
t.Fatalf("error = %q, want mismatch context", err.Error())
}
}
func TestRunStageAcceptsSessionIDFlagAndParsesStageName(t *testing.T) {
workspaceRoot := t.TempDir()
pipelinePath, sessionPath := writeValidConfigFiles(t, workspaceRoot)

View File

@@ -18,11 +18,15 @@ const (
ArtifactTranscriptFull = "narratio.transcript.full"
ArtifactTranscriptTrimmed = "narratio.transcript.trimmed"
ArtifactBoundsSession = "narratio.bounds.session"
ArtifactProvenancePreviousCacheManifestInput = "manifest.inputs.previous_cache"
ArtifactProvenancePreviousCacheFilesystem = "current_session.previous_cache"
)
// ErrSessionArtifactNotFound is returned when no readable artifact exists for a known ID.
var ErrSessionArtifactNotFound = errors.New("session artifact not found")
var configuredArtifactSourceRE = regexp.MustCompile(`^narratio\.artifact\.[a-z][a-z0-9_]*$`)
var previousSessionArtifactSourceRE = regexp.MustCompile(`^narratio\.previous_session\.artifact\.([a-z][a-z0-9_]*)$`)
type artifactContentKind string
@@ -118,6 +122,21 @@ func IsConfiguredArtifactSource(source string) bool {
return configuredArtifactSourceRE.MatchString(strings.TrimSpace(source))
}
// IsPreviousSessionArtifactSource returns true when source is narratio.previous_session.artifact.<name>.
func IsPreviousSessionArtifactSource(source string) bool {
_, ok := PreviousSessionArtifactName(source)
return ok
}
// PreviousSessionArtifactName extracts <name> from narratio.previous_session.artifact.<name>.
func PreviousSessionArtifactName(source string) (string, bool) {
matches := previousSessionArtifactSourceRE.FindStringSubmatch(strings.TrimSpace(source))
if len(matches) != 2 {
return "", false
}
return matches[1], true
}
// ResolveSessionArtifact resolves a symbolic source to a readable local session artifact path.
// Resolution order is manifest producer outputs first, then canonical session path fallback.
func ResolveSessionArtifact(paths SessionPaths, m *manifest.Manifest, source string) (ResolvedSessionArtifact, error) {
@@ -170,6 +189,9 @@ func ResolveSessionArtifact(paths SessionPaths, m *manifest.Manifest, source str
// configured narratio.artifact.<name> sources through runtime catalog availability.
func ResolveSessionArtifactWithCatalog(paths SessionPaths, m *manifest.Manifest, source string, catalog *ArtifactCatalog) (ResolvedSessionArtifact, error) {
normalized := strings.TrimSpace(source)
if IsPreviousSessionArtifactSource(normalized) {
return ResolvePreviousSessionArtifactWithCatalog(paths, m, normalized, catalog)
}
if !IsConfiguredArtifactSource(normalized) {
return ResolveSessionArtifact(paths, m, normalized)
}
@@ -195,6 +217,72 @@ func ResolveSessionArtifactWithCatalog(paths SessionPaths, m *manifest.Manifest,
}, nil
}
// ResolvePreviousSessionArtifactWithCatalog resolves one canonical previous-session source id
// to the prepared current-session previous-cache path.
func ResolvePreviousSessionArtifactWithCatalog(
paths SessionPaths,
m *manifest.Manifest,
source string,
catalog *ArtifactCatalog,
) (ResolvedSessionArtifact, error) {
artifactName, ok := PreviousSessionArtifactName(source)
if !ok {
return ResolvedSessionArtifact{}, fmt.Errorf("unsupported previous-session artifact source %q", source)
}
if catalog == nil {
return ResolvedSessionArtifact{}, fmt.Errorf("previous-session artifact source %q requires runtime artifact catalog", source)
}
configuredSourceID := ConfiguredArtifactSourceID(artifactName)
entry, ok := catalog.Lookup(configuredSourceID)
if !ok {
return ResolvedSessionArtifact{}, fmt.Errorf("unsupported previous-session artifact source %q", source)
}
candidates := previousSessionCacheCandidatePaths(paths, entry.CanonicalRelPath)
if len(candidates) == 0 {
return ResolvedSessionArtifact{}, &SessionArtifactNotFoundError{ArtifactID: source}
}
manifestInputPaths := manifestInputPathSet(paths, m)
fallback := ""
for _, candidate := range candidates {
exists, isDir, statErr := pathExists(candidate)
if statErr != nil {
return ResolvedSessionArtifact{}, fmt.Errorf("stat %q: %w", candidate, statErr)
}
if !exists || isDir {
continue
}
if err := validateResolvedContent(candidate, contentText); err != nil {
return ResolvedSessionArtifact{}, fmt.Errorf("validate %q: %w", source, err)
}
if _, ok := manifestInputPaths[candidate]; ok {
return ResolvedSessionArtifact{
ID: source,
Path: candidate,
ProducerStage: "prepare",
OutputKind: "previous_session_artifact",
Provenance: ArtifactProvenancePreviousCacheManifestInput,
}, nil
}
if fallback == "" {
fallback = candidate
}
}
if fallback != "" {
return ResolvedSessionArtifact{
ID: source,
Path: fallback,
ProducerStage: "prepare",
OutputKind: "previous_session_artifact",
Provenance: ArtifactProvenancePreviousCacheFilesystem,
}, nil
}
return ResolvedSessionArtifact{}, &SessionArtifactNotFoundError{ArtifactID: source}
}
func manifestArtifactCandidates(paths SessionPaths, m *manifest.Manifest, spec artifactSpec) []ResolvedSessionArtifact {
if m == nil || len(m.Stages) == 0 || spec.ProducerStage == "" || spec.OutputKind == "" {
return nil
@@ -243,6 +331,50 @@ func dedupeResolvedArtifacts(values []ResolvedSessionArtifact) []ResolvedSession
return out
}
func previousSessionCacheCandidatePaths(paths SessionPaths, canonicalRelPath string) []string {
trimmed := strings.TrimSpace(canonicalRelPath)
if trimmed == "" {
return nil
}
normalized := filepath.ToSlash(filepath.Clean(filepath.FromSlash(trimmed)))
if normalized == "." || normalized == "" || normalized == ".." || strings.HasPrefix(normalized, "../") || strings.HasPrefix(normalized, "/") {
return nil
}
relCandidates := []string{normalized}
const artifactsPrefix = "artifacts/"
if strings.HasPrefix(normalized, artifactsPrefix) && len(normalized) > len(artifactsPrefix) {
relCandidates = append(relCandidates, strings.TrimPrefix(normalized, artifactsPrefix))
}
out := make([]string, 0, len(relCandidates))
seen := map[string]struct{}{}
for _, rel := range relCandidates {
abs := filepath.Clean(SessionPreviousArtifactPath(paths, rel))
if _, ok := seen[abs]; ok {
continue
}
seen[abs] = struct{}{}
out = append(out, abs)
}
return out
}
func manifestInputPathSet(paths SessionPaths, m *manifest.Manifest) map[string]struct{} {
if m == nil || len(m.Inputs) == 0 {
return nil
}
out := make(map[string]struct{}, len(m.Inputs))
for _, in := range m.Inputs {
resolved := filepath.Clean(ResolveSessionLocalPathForRead(paths, in.Path))
if strings.TrimSpace(resolved) == "" {
continue
}
out[resolved] = struct{}{}
}
return out
}
func pathExists(path string) (exists bool, isDir bool, err error) {
info, err := os.Stat(path)
if err == nil {

View File

@@ -45,6 +45,58 @@ func TestNormalizeSessionArtifactSource(t *testing.T) {
}
}
func TestPreviousSessionArtifactSourceHelpers(t *testing.T) {
tests := []struct {
name string
source string
wantName string
wantMatch bool
}{
{
name: "valid",
source: "narratio.previous_session.artifact.session_recap",
wantName: "session_recap",
wantMatch: true,
},
{
name: "valid with surrounding whitespace",
source: " narratio.previous_session.artifact.quest_log ",
wantName: "quest_log",
wantMatch: true,
},
{
name: "missing name",
source: "narratio.previous_session.artifact.",
wantMatch: false,
},
{
name: "invalid key characters",
source: "narratio.previous_session.artifact.session-recap",
wantMatch: false,
},
{
name: "wrong prefix",
source: "narratio.previous.artifact.session_recap",
wantMatch: false,
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
name, ok := PreviousSessionArtifactName(tt.source)
if ok != tt.wantMatch {
t.Fatalf("PreviousSessionArtifactName(%q) ok = %t, want %t", tt.source, ok, tt.wantMatch)
}
if name != tt.wantName {
t.Fatalf("PreviousSessionArtifactName(%q) name = %q, want %q", tt.source, name, tt.wantName)
}
if got := IsPreviousSessionArtifactSource(tt.source); got != tt.wantMatch {
t.Fatalf("IsPreviousSessionArtifactSource(%q) = %t, want %t", tt.source, got, tt.wantMatch)
}
})
}
}
func TestResolveSessionArtifactPrefersManifestOutput(t *testing.T) {
workspace := t.TempDir()
paths := buildSessionPaths(workspace, "campaign", "session")
@@ -266,3 +318,120 @@ func TestResolveSessionArtifactWithCatalogUnsupportedConfiguredSourceFails(t *te
t.Fatalf("error = %q, want unsupported artifact source", err.Error())
}
}
func TestResolvePreviousSessionArtifactWithCatalogPrefersManifestInputRecord(t *testing.T) {
workspace := t.TempDir()
paths := buildSessionPaths(workspace, "campaign", "session")
manifestBackedPath := SessionPreviousArtifactPath(paths, "artifacts/session_recap.md")
fallbackPath := SessionPreviousArtifactPath(paths, "session_recap.md")
if err := os.MkdirAll(filepath.Dir(manifestBackedPath), 0o755); err != nil {
t.Fatalf("MkdirAll() error = %v", err)
}
if err := os.MkdirAll(filepath.Dir(fallbackPath), 0o755); err != nil {
t.Fatalf("MkdirAll() error = %v", err)
}
if err := os.WriteFile(manifestBackedPath, []byte("recap from manifest input\n"), 0o644); err != nil {
t.Fatalf("WriteFile() error = %v", err)
}
if err := os.WriteFile(fallbackPath, []byte("recap fallback\n"), 0o644); err != nil {
t.Fatalf("WriteFile() error = %v", err)
}
catalog := NewArtifactCatalog()
if err := catalog.RegisterConfiguredArtifacts(
map[string]ConfiguredArtifactDefinition{
"session_recap": {Enabled: true, OutputPath: "artifacts/session_recap.md"},
},
nil,
); err != nil {
t.Fatalf("RegisterConfiguredArtifacts() error = %v", err)
}
m := manifest.New("session", time.Now().UTC())
m.Inputs = []manifest.InputRecord{
{Kind: "previous_artifact", Path: manifestBackedPath},
}
resolved, err := ResolvePreviousSessionArtifactWithCatalog(
paths,
m,
"narratio.previous_session.artifact.session_recap",
catalog,
)
if err != nil {
t.Fatalf("ResolvePreviousSessionArtifactWithCatalog() error = %v", err)
}
if resolved.Path != manifestBackedPath {
t.Fatalf("resolved path = %q, want %q", resolved.Path, manifestBackedPath)
}
if resolved.Provenance != ArtifactProvenancePreviousCacheManifestInput {
t.Fatalf("provenance = %q, want %q", resolved.Provenance, ArtifactProvenancePreviousCacheManifestInput)
}
}
func TestResolvePreviousSessionArtifactWithCatalogFallsBackToPreparedCachePath(t *testing.T) {
workspace := t.TempDir()
paths := buildSessionPaths(workspace, "campaign", "session")
fallbackPath := SessionPreviousArtifactPath(paths, "session_recap.md")
if err := os.MkdirAll(filepath.Dir(fallbackPath), 0o755); err != nil {
t.Fatalf("MkdirAll() error = %v", err)
}
if err := os.WriteFile(fallbackPath, []byte("recap fallback\n"), 0o644); err != nil {
t.Fatalf("WriteFile() error = %v", err)
}
catalog := NewArtifactCatalog()
if err := catalog.RegisterConfiguredArtifacts(
map[string]ConfiguredArtifactDefinition{
"session_recap": {Enabled: true, OutputPath: "artifacts/session_recap.md"},
},
nil,
); err != nil {
t.Fatalf("RegisterConfiguredArtifacts() error = %v", err)
}
resolved, err := ResolvePreviousSessionArtifactWithCatalog(
paths,
nil,
"narratio.previous_session.artifact.session_recap",
catalog,
)
if err != nil {
t.Fatalf("ResolvePreviousSessionArtifactWithCatalog() error = %v", err)
}
if resolved.Path != fallbackPath {
t.Fatalf("resolved path = %q, want %q", resolved.Path, fallbackPath)
}
if resolved.Provenance != ArtifactProvenancePreviousCacheFilesystem {
t.Fatalf("provenance = %q, want %q", resolved.Provenance, ArtifactProvenancePreviousCacheFilesystem)
}
}
func TestResolvePreviousSessionArtifactWithCatalogMissingReturnsTypedError(t *testing.T) {
workspace := t.TempDir()
paths := buildSessionPaths(workspace, "campaign", "session")
catalog := NewArtifactCatalog()
if err := catalog.RegisterConfiguredArtifacts(
map[string]ConfiguredArtifactDefinition{
"session_recap": {Enabled: true, OutputPath: "artifacts/session_recap.md"},
},
nil,
); err != nil {
t.Fatalf("RegisterConfiguredArtifacts() error = %v", err)
}
_, err := ResolvePreviousSessionArtifactWithCatalog(
paths,
nil,
"narratio.previous_session.artifact.session_recap",
catalog,
)
if err == nil {
t.Fatal("expected error, got nil")
}
if !errors.Is(err, ErrSessionArtifactNotFound) {
t.Fatalf("errors.Is(err, ErrSessionArtifactNotFound)=false; err=%v", err)
}
}

View File

@@ -73,6 +73,8 @@ func (s *LocalStore) ensureLayout(paths SessionPaths) (SessionPaths, error) {
paths.LogsDir,
paths.CurrentDir,
paths.RunsDir,
paths.PreviousDir,
paths.PreviousArtifactsDir,
}
for _, dir := range dirs {

View File

@@ -27,10 +27,15 @@ func TestEnsureLayoutCreatesExpectedDirectories(t *testing.T) {
checkDirExists(t, paths.LogsDir)
checkDirExists(t, paths.CurrentDir)
checkDirExists(t, paths.RunsDir)
checkDirExists(t, paths.PreviousDir)
checkDirExists(t, paths.PreviousArtifactsDir)
if filepath.Base(paths.ManifestPath) != "manifest.json" {
t.Fatalf("ManifestPath = %q, want basename manifest.json", paths.ManifestPath)
}
if filepath.Base(paths.PreviousManifestPath) != "manifest.json" {
t.Fatalf("PreviousManifestPath = %q, want basename manifest.json", paths.PreviousManifestPath)
}
if filepath.Base(paths.LockPath) != ".lock" {
t.Fatalf("LockPath = %q, want basename .lock", paths.LockPath)
}

View File

@@ -23,6 +23,9 @@ type SessionPaths struct {
LogsDir string
CurrentDir string
RunsDir string
PreviousDir string
PreviousManifestPath string
PreviousArtifactsDir string
ManifestPath string
LockPath string
}
@@ -42,6 +45,29 @@ func SessionRunsDirForCampaign(rootDir, campaign, sessionID string) string {
return filepath.Join(SessionWorkDirForCampaign(rootDir, campaign, sessionID), config.PathRunsDirSegment)
}
// SessionPreviousDirForCampaign returns the canonical previous-session state directory.
func SessionPreviousDirForCampaign(rootDir, campaign, sessionID string) string {
return filepath.Join(SessionWorkDirForCampaign(rootDir, campaign, sessionID), config.PathPreviousDirSegment)
}
// SessionPreviousManifestPathForCampaign returns the canonical previous-session manifest cache path.
func SessionPreviousManifestPathForCampaign(rootDir, campaign, sessionID string) string {
return filepath.Join(SessionPreviousDirForCampaign(rootDir, campaign, sessionID), config.PathManifestFile)
}
// SessionPreviousArtifactsDirForCampaign returns the canonical previous-session artifact cache directory.
func SessionPreviousArtifactsDirForCampaign(rootDir, campaign, sessionID string) string {
return filepath.Join(SessionPreviousDirForCampaign(rootDir, campaign, sessionID), config.PathArtifactsDirSegment)
}
// SessionPreviousArtifactPathForCampaign returns a path under previous/artifacts for one artifact.
func SessionPreviousArtifactPathForCampaign(rootDir, campaign, sessionID, artifactRelativePath string) string {
return filepath.Join(
SessionPreviousArtifactsDirForCampaign(rootDir, campaign, sessionID),
filepath.FromSlash(artifactRelativePath),
)
}
// SessionRunRootForCampaign returns the canonical run root under runs/{run_id}.
func SessionRunRootForCampaign(rootDir, campaign, sessionID, runID string) string {
return filepath.Join(SessionRunsDirForCampaign(rootDir, campaign, sessionID), runID)
@@ -62,6 +88,26 @@ func SessionSpoolAudioDir(spoolRoot, campaign, sessionID, runID string) string {
return filepath.Join(spoolRoot, campaign, sessionID, runID, config.PathAudioDirSegment)
}
// SessionPreviousDir returns the previous-session state directory for already-resolved session paths.
func SessionPreviousDir(paths SessionPaths) string {
return paths.PreviousDir
}
// SessionPreviousManifestPath returns the previous-session manifest path for already-resolved session paths.
func SessionPreviousManifestPath(paths SessionPaths) string {
return paths.PreviousManifestPath
}
// SessionPreviousArtifactsDir returns the previous-session artifact directory for already-resolved session paths.
func SessionPreviousArtifactsDir(paths SessionPaths) string {
return paths.PreviousArtifactsDir
}
// SessionPreviousArtifactPath returns a path under previous/artifacts for already-resolved session paths.
func SessionPreviousArtifactPath(paths SessionPaths, artifactRelativePath string) string {
return filepath.Join(paths.PreviousArtifactsDir, filepath.FromSlash(artifactRelativePath))
}
func buildSessionPaths(workspaceRoot, campaign, sessionID string) SessionPaths {
root := SessionWorkDirForCampaign(workspaceRoot, campaign, sessionID)
return buildSessionPathsFromRoot(workspaceRoot, campaign, sessionID, root)
@@ -84,6 +130,9 @@ func buildSessionPathsFromRoot(workspaceRoot, campaign, sessionID, root string)
LogsDir: filepath.Join(root, config.PathLogsDirSegment),
CurrentDir: filepath.Join(root, config.PathCurrentDirSegment),
RunsDir: filepath.Join(root, config.PathRunsDirSegment),
PreviousDir: filepath.Join(root, config.PathPreviousDirSegment),
PreviousManifestPath: filepath.Join(root, config.PathPreviousDirSegment, config.PathManifestFile),
PreviousArtifactsDir: filepath.Join(root, config.PathPreviousDirSegment, config.PathArtifactsDirSegment),
ManifestPath: filepath.Join(root, config.PathManifestFile),
LockPath: filepath.Join(root, config.PathLockFile),
}

View File

@@ -49,6 +49,52 @@ func TestSessionRunManifestPathForCampaign(t *testing.T) {
}
}
func TestSessionPreviousPathsForCampaign(t *testing.T) {
root := "/tmp/workspace"
previousDir := SessionPreviousDirForCampaign(root, "forsaken", "2026-04-19")
wantPreviousDir := filepath.Join(root, "work", "forsaken", "2026-04-19", "previous")
if previousDir != wantPreviousDir {
t.Fatalf("SessionPreviousDirForCampaign() = %q, want %q", previousDir, wantPreviousDir)
}
manifestPath := SessionPreviousManifestPathForCampaign(root, "forsaken", "2026-04-19")
wantManifestPath := filepath.Join(root, "work", "forsaken", "2026-04-19", "previous", "manifest.json")
if manifestPath != wantManifestPath {
t.Fatalf("SessionPreviousManifestPathForCampaign() = %q, want %q", manifestPath, wantManifestPath)
}
artifactsDir := SessionPreviousArtifactsDirForCampaign(root, "forsaken", "2026-04-19")
wantArtifactsDir := filepath.Join(root, "work", "forsaken", "2026-04-19", "previous", "artifacts")
if artifactsDir != wantArtifactsDir {
t.Fatalf("SessionPreviousArtifactsDirForCampaign() = %q, want %q", artifactsDir, wantArtifactsDir)
}
artifactPath := SessionPreviousArtifactPathForCampaign(root, "forsaken", "2026-04-19", "session_recap.md")
wantArtifactPath := filepath.Join(root, "work", "forsaken", "2026-04-19", "previous", "artifacts", "session_recap.md")
if artifactPath != wantArtifactPath {
t.Fatalf("SessionPreviousArtifactPathForCampaign() = %q, want %q", artifactPath, wantArtifactPath)
}
}
func TestSessionPreviousPathsFromSessionPaths(t *testing.T) {
paths := buildSessionPaths("/tmp/workspace", "forsaken", "2026-04-19")
if got := SessionPreviousDir(paths); got != paths.PreviousDir {
t.Fatalf("SessionPreviousDir() = %q, want %q", got, paths.PreviousDir)
}
if got := SessionPreviousManifestPath(paths); got != paths.PreviousManifestPath {
t.Fatalf("SessionPreviousManifestPath() = %q, want %q", got, paths.PreviousManifestPath)
}
if got := SessionPreviousArtifactsDir(paths); got != paths.PreviousArtifactsDir {
t.Fatalf("SessionPreviousArtifactsDir() = %q, want %q", got, paths.PreviousArtifactsDir)
}
got := SessionPreviousArtifactPath(paths, "quest_log.json")
want := filepath.Join(paths.PreviousArtifactsDir, "quest_log.json")
if got != want {
t.Fatalf("SessionPreviousArtifactPath() = %q, want %q", got, want)
}
}
func TestSessionSpoolAudioDir(t *testing.T) {
root := "/var/spool/narratio"
got := SessionSpoolAudioDir(root, "forsaken", "2026-04-19", "20260515T031522Z-a1b2c3d4")

View File

@@ -0,0 +1,112 @@
package artifacts
import (
"fmt"
"sort"
"strings"
"gitea.maximumdirect.net/eric/narratio/internal/config"
)
// PreviousArtifactRequirement describes one previous-session artifact dependency.
type PreviousArtifactRequirement struct {
Name string
Required bool
Sources []string
}
// CollectPreviousArtifactRequirements scans enabled Scriptorium artifacts and returns
// deduplicated previous-session artifact requirements in deterministic order.
func CollectPreviousArtifactRequirements(
artifactsCfg map[string]config.ScriptoriumArtifactConfig,
) []PreviousArtifactRequirement {
if len(artifactsCfg) == 0 {
return nil
}
artifactNames := sortedScriptoriumArtifactNames(artifactsCfg)
byName := map[string]PreviousArtifactRequirement{}
for _, artifactName := range artifactNames {
artifactCfg := artifactsCfg[artifactName]
if !artifactCfg.Enabled {
continue
}
inputNames := sortedScriptoriumInputKeys(artifactCfg.Inputs)
for _, inputName := range inputNames {
inputCfg := artifactCfg.Inputs[inputName]
previousName, ok := PreviousSessionArtifactName(inputCfg.Source)
if !ok {
continue
}
location := fmt.Sprintf(
"pipeline.scriptorium.artifacts.%s.inputs.%s.source",
artifactName,
inputName,
)
requirement := byName[previousName]
requirement.Name = previousName
requirement.Required = requirement.Required || inputCfg.Required
requirement.Sources = append(requirement.Sources, location)
byName[previousName] = requirement
}
}
if len(byName) == 0 {
return nil
}
requirements := make([]PreviousArtifactRequirement, 0, len(byName))
for _, requirement := range byName {
requirement.Sources = dedupeAndSortStrings(requirement.Sources)
requirements = append(requirements, requirement)
}
sort.Slice(requirements, func(i, j int) bool {
return requirements[i].Name < requirements[j].Name
})
return requirements
}
func sortedScriptoriumArtifactNames(artifactsCfg map[string]config.ScriptoriumArtifactConfig) []string {
names := make([]string, 0, len(artifactsCfg))
for name := range artifactsCfg {
names = append(names, name)
}
sort.Strings(names)
return names
}
func sortedScriptoriumInputKeys(inputs map[string]config.ScriptoriumInputConfig) []string {
if len(inputs) == 0 {
return nil
}
names := make([]string, 0, len(inputs))
for name := range inputs {
names = append(names, name)
}
sort.Strings(names)
return names
}
func dedupeAndSortStrings(values []string) []string {
if len(values) == 0 {
return nil
}
seen := map[string]struct{}{}
out := make([]string, 0, len(values))
for _, value := range values {
trimmed := strings.TrimSpace(value)
if trimmed == "" {
continue
}
if _, ok := seen[trimmed]; ok {
continue
}
seen[trimmed] = struct{}{}
out = append(out, trimmed)
}
sort.Strings(out)
return out
}

View File

@@ -0,0 +1,190 @@
package artifacts
import (
"reflect"
"testing"
"gitea.maximumdirect.net/eric/narratio/internal/config"
)
func TestCollectPreviousArtifactRequirements(t *testing.T) {
tests := []struct {
name string
artifactsCfg map[string]config.ScriptoriumArtifactConfig
want []PreviousArtifactRequirement
}{
{
name: "no artifacts",
artifactsCfg: nil,
want: nil,
},
{
name: "no previous inputs",
artifactsCfg: map[string]config.ScriptoriumArtifactConfig{
"session_recap": {
Enabled: true,
Inputs: map[string]config.ScriptoriumInputConfig{
"transcript": {Source: "narratio.transcript.trimmed", Required: true},
},
},
},
want: nil,
},
{
name: "one optional previous input",
artifactsCfg: map[string]config.ScriptoriumArtifactConfig{
"session_recap": {
Enabled: true,
Inputs: map[string]config.ScriptoriumInputConfig{
"previous_recap": {Source: "narratio.previous_session.artifact.session_recap", Required: false},
},
},
},
want: []PreviousArtifactRequirement{
{
Name: "session_recap",
Required: false,
Sources: []string{"pipeline.scriptorium.artifacts.session_recap.inputs.previous_recap.source"},
},
},
},
{
name: "one required previous input",
artifactsCfg: map[string]config.ScriptoriumArtifactConfig{
"quest_log": {
Enabled: true,
Inputs: map[string]config.ScriptoriumInputConfig{
"previous_quest_log": {Source: "narratio.previous_session.artifact.quest_log", Required: true},
},
},
},
want: []PreviousArtifactRequirement{
{
Name: "quest_log",
Required: true,
Sources: []string{"pipeline.scriptorium.artifacts.quest_log.inputs.previous_quest_log.source"},
},
},
},
{
name: "duplicate references are deduped",
artifactsCfg: map[string]config.ScriptoriumArtifactConfig{
"a": {
Enabled: true,
Inputs: map[string]config.ScriptoriumInputConfig{
"x": {Source: "narratio.previous_session.artifact.session_recap", Required: false},
},
},
"b": {
Enabled: true,
Inputs: map[string]config.ScriptoriumInputConfig{
"y": {Source: "narratio.previous_session.artifact.session_recap", Required: false},
},
},
},
want: []PreviousArtifactRequirement{
{
Name: "session_recap",
Required: false,
Sources: []string{
"pipeline.scriptorium.artifacts.a.inputs.x.source",
"pipeline.scriptorium.artifacts.b.inputs.y.source",
},
},
},
},
{
name: "required plus optional reference becomes required",
artifactsCfg: map[string]config.ScriptoriumArtifactConfig{
"a": {
Enabled: true,
Inputs: map[string]config.ScriptoriumInputConfig{
"x": {Source: "narratio.previous_session.artifact.session_recap", Required: false},
},
},
"b": {
Enabled: true,
Inputs: map[string]config.ScriptoriumInputConfig{
"y": {Source: "narratio.previous_session.artifact.session_recap", Required: true},
},
},
},
want: []PreviousArtifactRequirement{
{
Name: "session_recap",
Required: true,
Sources: []string{
"pipeline.scriptorium.artifacts.a.inputs.x.source",
"pipeline.scriptorium.artifacts.b.inputs.y.source",
},
},
},
},
{
name: "disabled artifact references are ignored",
artifactsCfg: map[string]config.ScriptoriumArtifactConfig{
"disabled_artifact": {
Enabled: false,
Inputs: map[string]config.ScriptoriumInputConfig{
"x": {Source: "narratio.previous_session.artifact.session_recap", Required: true},
},
},
"enabled_artifact": {
Enabled: true,
Inputs: map[string]config.ScriptoriumInputConfig{
"y": {Source: "narratio.previous_session.artifact.quest_log", Required: false},
},
},
},
want: []PreviousArtifactRequirement{
{
Name: "quest_log",
Required: false,
Sources: []string{"pipeline.scriptorium.artifacts.enabled_artifact.inputs.y.source"},
},
},
},
{
name: "deterministic ordering",
artifactsCfg: map[string]config.ScriptoriumArtifactConfig{
"zz": {
Enabled: true,
Inputs: map[string]config.ScriptoriumInputConfig{
"b_input": {Source: "narratio.previous_session.artifact.quest_log", Required: false},
"a_input": {Source: "narratio.previous_session.artifact.session_recap", Required: false},
},
},
"aa": {
Enabled: true,
Inputs: map[string]config.ScriptoriumInputConfig{
"c_input": {Source: "narratio.previous_session.artifact.session_recap", Required: false},
},
},
},
want: []PreviousArtifactRequirement{
{
Name: "quest_log",
Required: false,
Sources: []string{"pipeline.scriptorium.artifacts.zz.inputs.b_input.source"},
},
{
Name: "session_recap",
Required: false,
Sources: []string{
"pipeline.scriptorium.artifacts.aa.inputs.c_input.source",
"pipeline.scriptorium.artifacts.zz.inputs.a_input.source",
},
},
},
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
got := CollectPreviousArtifactRequirements(tt.artifactsCfg)
if !reflect.DeepEqual(got, tt.want) {
t.Fatalf("CollectPreviousArtifactRequirements() = %#v, want %#v", got, tt.want)
}
})
}
}

View File

@@ -27,11 +27,12 @@ type PipelineConfig struct {
// SessionConfig contains per-session inputs and metadata.
type SessionConfig struct {
SessionID string `yaml:"session_id"`
Campaign string `yaml:"campaign"`
Date string `yaml:"date"`
Title string `yaml:"title"`
Inputs SessionInputsConfig `yaml:"inputs"`
SessionID string `yaml:"session_id"`
PreviousSessionID string `yaml:"previous_session_id"`
Campaign string `yaml:"campaign"`
Date string `yaml:"date"`
Title string `yaml:"title"`
Inputs SessionInputsConfig `yaml:"inputs"`
}
// WorkspaceConfig configures local workspace behavior.

View File

@@ -56,6 +56,7 @@ const (
PathLogsDirSegment = "logs"
PathCurrentDirSegment = "current"
PathRunsDirSegment = "runs"
PathPreviousDirSegment = "previous"
PathManifestFile = "manifest.json"
PathLockFile = ".lock"
PathTranscriptMerged = "transcripts/merged.json"

View File

@@ -6,6 +6,7 @@ import (
"os"
"path/filepath"
"regexp"
"sort"
"strings"
"gopkg.in/yaml.v3"
@@ -28,7 +29,8 @@ func LoadSession(path string) (*SessionConfig, error) {
// SessionLoadOptions configures session template rendering behavior.
type SessionLoadOptions struct {
SessionID string
SessionID string
PreviousSessionID string
}
// LoadSessionWithOptions loads session configuration from a YAML file with
@@ -56,6 +58,16 @@ func LoadSessionWithOptions(path string, opts SessionLoadOptions) (*SessionConfi
strings.TrimSpace(cfg.SessionID),
)
}
if strings.TrimSpace(opts.PreviousSessionID) != "" &&
strings.TrimSpace(cfg.PreviousSessionID) != "" &&
strings.TrimSpace(cfg.PreviousSessionID) != strings.TrimSpace(opts.PreviousSessionID) {
return nil, fmt.Errorf(
"load session config: session file %q: previous_session_id mismatch: --previous-session-id %q does not match rendered previous_session_id %q",
path,
strings.TrimSpace(opts.PreviousSessionID),
strings.TrimSpace(cfg.PreviousSessionID),
)
}
return &cfg, nil
}
@@ -114,24 +126,36 @@ var sessionTemplatePattern = regexp.MustCompile(`\{\{\s*([a-zA-Z_][a-zA-Z0-9_]*)
func renderSessionTemplate(content string, opts SessionLoadOptions) (string, error) {
sessionID := strings.TrimSpace(opts.SessionID)
previousSessionID := strings.TrimSpace(opts.PreviousSessionID)
rendered := content
if sessionID != "" {
rendered = strings.ReplaceAll(rendered, "{{session_id}}", sessionID)
rendered = strings.ReplaceAll(rendered, "{{ session_id }}", sessionID)
rendered = replaceTemplateVariable(rendered, "session_id", sessionID)
}
if previousSessionID != "" {
rendered = replaceTemplateVariable(rendered, "previous_session_id", previousSessionID)
}
unresolved := sessionTemplatePattern.FindAllStringSubmatch(rendered, -1)
if len(unresolved) > 0 {
seenVars := map[string]struct{}{}
vars := make([]string, 0, len(unresolved))
for _, m := range unresolved {
if len(m) > 1 {
vars = append(vars, m[1])
name := m[1]
if _, ok := seenVars[name]; ok {
continue
}
seenVars[name] = struct{}{}
vars = append(vars, name)
}
}
sort.Strings(vars)
if len(vars) > 0 {
hints := unresolvedTemplateHints(vars)
return "", fmt.Errorf(
"session file template rendering failed: unresolved template variable(s): %s; pass --session-id when using {{ session_id }}",
"session file template rendering failed: unresolved template variable(s): %s%s",
strings.Join(vars, ", "),
hints,
)
}
return "", fmt.Errorf("session file template rendering failed: unresolved template placeholders remain")
@@ -140,6 +164,35 @@ func renderSessionTemplate(content string, opts SessionLoadOptions) (string, err
return rendered, nil
}
func replaceTemplateVariable(content, name, value string) string {
rendered := strings.ReplaceAll(content, "{{"+name+"}}", value)
rendered = strings.ReplaceAll(rendered, "{{ "+name+" }}", value)
return rendered
}
func unresolvedTemplateHints(vars []string) string {
seen := map[string]struct{}{}
flags := make([]string, 0, 2)
for _, name := range vars {
switch name {
case "session_id":
if _, ok := seen["--session-id"]; !ok {
seen["--session-id"] = struct{}{}
flags = append(flags, "--session-id")
}
case "previous_session_id":
if _, ok := seen["--previous-session-id"]; !ok {
seen["--previous-session-id"] = struct{}{}
flags = append(flags, "--previous-session-id")
}
}
}
if len(flags) == 0 {
return ""
}
return "; pass " + strings.Join(flags, " and ") + " when using those template variable(s)"
}
func shortName(path, fallback string) string {
base := filepath.Base(path)
if base == "." || base == string(filepath.Separator) {

View File

@@ -187,6 +187,43 @@ inputs:
`,
wantValidate: "session config \"session.yml\" invalid: session.session_id is required",
},
{
name: "valid previous_session_id passes",
pipelineYAML: `workspace:
root: /tmp/narratio
whisperx:
transcribe_url: https://transcription.ai.rakestrawhome.com/transcribe
seriatim:
binary: seriatim
`,
sessionYAML: `session_id: 2026-05-03
previous_session_id: 2026-04-26
inputs:
audio_dir: ./audio
speakers_file: ./speakers.yml
autocorrect_file: ./autocorrect.yml
glossary_file: ./glossary.yml
`,
},
{
name: "previous_session_id equal to session_id fails",
pipelineYAML: `workspace:
root: /tmp/narratio
whisperx:
transcribe_url: https://transcription.ai.rakestrawhome.com/transcribe
seriatim:
binary: seriatim
`,
sessionYAML: `session_id: 2026-05-03
previous_session_id: 2026-05-03
inputs:
audio_dir: ./audio
speakers_file: ./speakers.yml
autocorrect_file: ./autocorrect.yml
glossary_file: ./glossary.yml
`,
wantValidate: "session config \"session.yml\" invalid: session.previous_session_id must not equal session.session_id",
},
{
name: "missing transcribe_url fails",
pipelineYAML: `workspace:

View File

@@ -116,6 +116,75 @@ func TestScriptoriumLoadAndValidate(t *testing.T) {
output_kind: session_recap
`,
},
{
name: "canonical previous-session source is accepted",
scriptoriumYAML: `scriptorium:
binary: scriptorium
artifacts:
session_recap:
enabled: true
prompt_id: dnd.session_recap
output_path: artifacts/session_recap.md
inputs:
transcript:
source: narratio.transcript.polished
required: true
previous_recap:
source: narratio.previous_session.artifact.session_recap
required: false
vars:
session_id: true
output_kind: session_recap
`,
},
{
name: "canonical previous-session source missing artifact key fails validation",
scriptoriumYAML: `scriptorium:
binary: scriptorium
artifacts:
session_recap:
enabled: true
prompt_id: dnd.session_recap
output_path: artifacts/session_recap.md
inputs:
previous_recap:
source: narratio.previous_session.artifact.
required: false
`,
wantValidateErr: `pipeline.scriptorium.artifacts.session_recap.inputs.previous_recap.source "narratio.previous_session.artifact." must reference configured artifact key matching ^[a-z][a-z0-9_]*$`,
},
{
name: "canonical previous-session source invalid artifact key fails validation",
scriptoriumYAML: `scriptorium:
binary: scriptorium
artifacts:
session_recap:
enabled: true
prompt_id: dnd.session_recap
output_path: artifacts/session_recap.md
inputs:
previous_recap:
source: narratio.previous_session.artifact.session-recap
required: false
`,
wantValidateErr: `pipeline.scriptorium.artifacts.session_recap.inputs.previous_recap.source "narratio.previous_session.artifact.session-recap" must reference configured artifact key matching ^[a-z][a-z0-9_]*$`,
},
{
name: "canonical previous-session source unknown artifact fails validation",
scriptoriumYAML: `scriptorium:
binary: scriptorium
artifacts:
session_recap:
enabled: true
prompt_id: dnd.session_recap
output_path: artifacts/session_recap.md
inputs:
previous_recap:
source: narratio.previous_session.artifact.quest_log
required: false
`,
wantValidateErr: `pipeline.scriptorium.artifacts.session_recap.inputs.previous_recap.source "narratio.previous_session.artifact.quest_log" references unknown artifact "quest_log"`,
},
{
name: "canonical artifact source is accepted",
scriptoriumYAML: `scriptorium:

View File

@@ -55,6 +55,34 @@ inputs:
}
}
func TestLoadSessionWithOptionsRendersPreviousSessionPlaceholder(t *testing.T) {
dir := t.TempDir()
sessionPath := filepath.Join(dir, "session.yml")
sessionYAML := `session_id: "{{ session_id }}"
previous_session_id: "{{ previous_session_id }}"
campaign: sample-campaign
inputs:
audio_dir: ./audio
speakers_file: ./speakers.yml
autocorrect_file: ./autocorrect.yml
glossary_file: ./glossary.yml
`
if err := os.WriteFile(sessionPath, []byte(sessionYAML), 0o644); err != nil {
t.Fatalf("write session.yml: %v", err)
}
cfg, err := LoadSessionWithOptions(sessionPath, SessionLoadOptions{
SessionID: "2026-04-04",
PreviousSessionID: "2026-03-28",
})
if err != nil {
t.Fatalf("LoadSessionWithOptions() error = %v", err)
}
if cfg.PreviousSessionID != "2026-03-28" {
t.Fatalf("PreviousSessionID = %q, want 2026-03-28", cfg.PreviousSessionID)
}
}
func TestLoadSessionWithOptionsUnresolvedPlaceholderFails(t *testing.T) {
dir := t.TempDir()
sessionPath := filepath.Join(dir, "session.yml")
@@ -82,6 +110,37 @@ inputs:
}
}
func TestLoadSessionWithOptionsUnresolvedPreviousSessionPlaceholderFails(t *testing.T) {
dir := t.TempDir()
sessionPath := filepath.Join(dir, "session.yml")
sessionYAML := `session_id: 2026-05-03
previous_session_id: "{{ previous_session_id }}"
campaign: sample-campaign
inputs:
audio_dir: ./audio
speakers_file: ./speakers.yml
autocorrect_file: ./autocorrect.yml
glossary_file: ./glossary.yml
`
if err := os.WriteFile(sessionPath, []byte(sessionYAML), 0o644); err != nil {
t.Fatalf("write session.yml: %v", err)
}
_, err := LoadSessionWithOptions(sessionPath, SessionLoadOptions{})
if err == nil {
t.Fatal("expected error, got nil")
}
if !strings.Contains(err.Error(), "unresolved template variable") {
t.Fatalf("error = %q, want unresolved-variable context", err.Error())
}
if !strings.Contains(err.Error(), "previous_session_id") {
t.Fatalf("error = %q, want previous_session_id variable", err.Error())
}
if !strings.Contains(err.Error(), "--previous-session-id") {
t.Fatalf("error = %q, want previous-session-id guidance", err.Error())
}
}
func TestLoadSessionWithOptionsMismatchFails(t *testing.T) {
dir := t.TempDir()
sessionPath := filepath.Join(dir, "session.yml")
@@ -106,6 +165,34 @@ inputs:
}
}
func TestLoadSessionWithOptionsPreviousSessionMismatchFails(t *testing.T) {
dir := t.TempDir()
sessionPath := filepath.Join(dir, "session.yml")
sessionYAML := `session_id: 2026-05-03
previous_session_id: 2026-04-26
campaign: sample-campaign
inputs:
audio_dir: ./audio
speakers_file: ./speakers.yml
autocorrect_file: ./autocorrect.yml
glossary_file: ./glossary.yml
`
if err := os.WriteFile(sessionPath, []byte(sessionYAML), 0o644); err != nil {
t.Fatalf("write session.yml: %v", err)
}
_, err := LoadSessionWithOptions(sessionPath, SessionLoadOptions{
SessionID: "2026-05-03",
PreviousSessionID: "2026-04-25",
})
if err == nil {
t.Fatal("expected error, got nil")
}
if !strings.Contains(err.Error(), "previous_session_id mismatch") {
t.Fatalf("error = %q, want mismatch context", err.Error())
}
}
func TestLoadSessionWithOptionsUnknownFieldStillRejectedAfterRendering(t *testing.T) {
dir := t.TempDir()
sessionPath := filepath.Join(dir, "session.yml")

View File

@@ -487,8 +487,14 @@ func validateScriptorium(cfg *ScriptoriumConfig) error {
}
func validateSession(cfg *SessionConfig) error {
if strings.TrimSpace(cfg.SessionID) == "" {
return fmt.Errorf("session.session_id is required")
if err := validateSessionIdentifier("session.session_id", cfg.SessionID, true); err != nil {
return err
}
if err := validateSessionIdentifier("session.previous_session_id", cfg.PreviousSessionID, false); err != nil {
return err
}
if strings.TrimSpace(cfg.PreviousSessionID) != "" && strings.TrimSpace(cfg.PreviousSessionID) == strings.TrimSpace(cfg.SessionID) {
return fmt.Errorf("session.previous_session_id must not equal session.session_id")
}
if strings.TrimSpace(cfg.Campaign) == "" {
return fmt.Errorf("session.campaign is required")
@@ -525,6 +531,16 @@ func validateSession(cfg *SessionConfig) error {
return nil
}
func validateSessionIdentifier(fieldName, value string, required bool) error {
if strings.TrimSpace(value) == "" {
if required {
return fmt.Errorf("%s is required", fieldName)
}
return nil
}
return nil
}
func validateCrossConfig(pipeline *PipelineConfig, session *SessionConfig) error {
if pipeline == nil || session == nil {
return nil
@@ -563,12 +579,38 @@ var windowsAbsPathRE = regexp.MustCompile(`^[A-Za-z]:[\\/].*`)
var envVarNameRE = regexp.MustCompile(`^[A-Za-z_][A-Za-z0-9_]*$`)
var scriptoriumArtifactKeyRE = regexp.MustCompile(`^[a-z][a-z0-9_]*$`)
var narratioArtifactSourceRE = regexp.MustCompile(`^narratio\.artifact\.([a-z][a-z0-9_]*)$`)
var narratioPreviousSessionArtifactSourceRE = regexp.MustCompile(`^narratio\.previous_session\.artifact\.([a-z][a-z0-9_]*)$`)
func validateScriptoriumInputSource(artifactName, inputName, source string, configuredArtifacts map[string]struct{}) (string, error) {
if isStaticSupportedScriptoriumInputSource(source) {
trimmedSource := strings.TrimSpace(source)
if isStaticSupportedScriptoriumInputSource(trimmedSource) {
return "", nil
}
matches := narratioArtifactSourceRE.FindStringSubmatch(source)
if strings.HasPrefix(trimmedSource, "narratio.previous_session.artifact") {
matches := narratioPreviousSessionArtifactSourceRE.FindStringSubmatch(trimmedSource)
if len(matches) != 2 {
return "", fmt.Errorf(
"pipeline.scriptorium.artifacts.%s.inputs.%s.source %q must reference configured artifact key matching ^[a-z][a-z0-9_]*$",
artifactName,
inputName,
source,
)
}
referenced := matches[1]
if _, ok := configuredArtifacts[referenced]; !ok {
return "", fmt.Errorf(
"pipeline.scriptorium.artifacts.%s.inputs.%s.source %q references unknown artifact %q",
artifactName,
inputName,
source,
referenced,
)
}
return "", nil
}
matches := narratioArtifactSourceRE.FindStringSubmatch(trimmedSource)
if len(matches) != 2 {
return "", fmt.Errorf(
"pipeline.scriptorium.artifacts.%s.inputs.%s.source %q is unsupported",
@@ -591,7 +633,7 @@ func validateScriptoriumInputSource(artifactName, inputName, source string, conf
}
func isStaticSupportedScriptoriumInputSource(source string) bool {
switch strings.TrimSpace(source) {
switch source {
case "previous_session_artifact":
return true
case "narratio.transcript.merged":

View File

@@ -631,6 +631,23 @@ func resolveScriptoriumInput(
runtimeCatalog *artifacts.ArtifactCatalog,
) (string, bool, *artifacts.ResolvedSessionArtifact, error) {
source := strings.TrimSpace(inputCfg.Source)
if artifacts.IsPreviousSessionArtifactSource(source) {
resolved, err := artifacts.ResolvePreviousSessionArtifactWithCatalog(paths, m, source, runtimeCatalog)
if err == nil {
copy := resolved
return resolved.Path, true, &copy, nil
}
if errors.Is(err, artifacts.ErrSessionArtifactNotFound) {
if inputCfg.Required {
return "", false, nil, fmt.Errorf(
"required previous-session input source %q is unavailable; run narratio run-stage --force prepare",
source,
)
}
return "", false, nil, nil
}
return "", false, nil, err
}
switch source {
case "previous_session_artifact":
if strings.TrimSpace(inputCfg.Path) == "" {

View File

@@ -10,6 +10,7 @@ import (
"time"
"gitea.maximumdirect.net/eric/narratio/internal/adapters/scriptorium"
"gitea.maximumdirect.net/eric/narratio/internal/adapters/storage"
"gitea.maximumdirect.net/eric/narratio/internal/artifacts"
"gitea.maximumdirect.net/eric/narratio/internal/config"
"gitea.maximumdirect.net/eric/narratio/internal/manifest"
@@ -717,6 +718,146 @@ func TestAnalyzeOmitsOptionalMissingConfiguredArtifactInput(t *testing.T) {
}
}
func TestAnalyzeRequiredPreviousSessionArtifactInputGuidesPrepareForce(t *testing.T) {
env, m, _ := setupAnalyzeEnv(t)
sessionRecap := env.Config.Pipeline.Scriptorium.Artifacts["session_recap"]
sessionRecap.Inputs["previous_recap"] = config.ScriptoriumInputConfig{
Source: "narratio.previous_session.artifact.session_recap",
Required: true,
}
env.Config.Pipeline.Scriptorium.Artifacts["session_recap"] = sessionRecap
_, err := (analyzeStage{}).Run(context.Background(), env, m)
if err == nil {
t.Fatal("expected error, got nil")
}
if !strings.Contains(err.Error(), "run narratio run-stage --force prepare") {
t.Fatalf("error = %q, want guidance to run force prepare", err.Error())
}
}
func TestAnalyzeResolvesCanonicalPreviousSessionArtifactFromManifestInput(t *testing.T) {
env, m, fake := setupAnalyzeEnv(t)
paths := sessionPathsForEnv(env, m.SessionID)
writeAnalyzeFile(t, filepath.Join(paths.TranscriptsDir, "trimmed.json"), `{"segments":[]}`)
previousPath := artifacts.SessionPreviousArtifactPath(paths, "artifacts/session_recap.md")
writeAnalyzeFile(t, previousPath, "previous recap\n")
m.Inputs = append(m.Inputs, manifest.InputRecord{
Kind: "previous_artifact",
Path: previousPath,
})
sessionRecap := env.Config.Pipeline.Scriptorium.Artifacts["session_recap"]
sessionRecap.Inputs["previous_recap"] = config.ScriptoriumInputConfig{
Source: "narratio.previous_session.artifact.session_recap",
Required: true,
}
env.Config.Pipeline.Scriptorium.Artifacts["session_recap"] = sessionRecap
_, err := (analyzeStage{}).Run(context.Background(), env, m)
if err != nil {
t.Fatalf("Run() error = %v", err)
}
if len(fake.RunRequests) != 1 {
t.Fatalf("run requests = %d, want 1", len(fake.RunRequests))
}
if got := fake.RunRequests[0].InputPaths["previous_recap"]; got != previousPath {
t.Fatalf("previous_recap input = %q, want %q", got, previousPath)
}
}
func TestAnalyzeOmitsOptionalMissingCanonicalPreviousSessionArtifact(t *testing.T) {
env, m, fake := setupAnalyzeEnv(t)
paths := sessionPathsForEnv(env, m.SessionID)
writeAnalyzeFile(t, filepath.Join(paths.TranscriptsDir, "trimmed.json"), `{"segments":[]}`)
sessionRecap := env.Config.Pipeline.Scriptorium.Artifacts["session_recap"]
sessionRecap.Inputs["previous_recap"] = config.ScriptoriumInputConfig{
Source: "narratio.previous_session.artifact.session_recap",
Required: false,
}
env.Config.Pipeline.Scriptorium.Artifacts["session_recap"] = sessionRecap
_, err := (analyzeStage{}).Run(context.Background(), env, m)
if err != nil {
t.Fatalf("Run() error = %v", err)
}
if len(fake.RunRequests) != 1 {
t.Fatalf("run requests = %d, want 1", len(fake.RunRequests))
}
if _, exists := fake.RunRequests[0].InputPaths["previous_recap"]; exists {
t.Fatalf("optional canonical previous_recap should be omitted when unavailable")
}
}
func TestAnalyzeRenderDebugWithCanonicalPreviousSessionInput(t *testing.T) {
env, m, fake := setupAnalyzeEnv(t)
paths := sessionPathsForEnv(env, m.SessionID)
writeAnalyzeFile(t, filepath.Join(paths.TranscriptsDir, "trimmed.json"), `{"segments":[]}`)
env.Config.Pipeline.Scriptorium.RenderDebug = true
previousPath := artifacts.SessionPreviousArtifactPath(paths, "artifacts/session_recap.md")
writeAnalyzeFile(t, previousPath, "previous recap\n")
m.Inputs = append(m.Inputs, manifest.InputRecord{
Kind: "previous_artifact",
Path: previousPath,
})
sessionRecap := env.Config.Pipeline.Scriptorium.Artifacts["session_recap"]
sessionRecap.Inputs["previous_recap"] = config.ScriptoriumInputConfig{
Source: "narratio.previous_session.artifact.session_recap",
Required: true,
}
env.Config.Pipeline.Scriptorium.Artifacts["session_recap"] = sessionRecap
_, err := (analyzeStage{}).Run(context.Background(), env, m)
if err != nil {
t.Fatalf("Run() error = %v", err)
}
if len(fake.RenderRequests) != 1 || len(fake.RunRequests) != 1 {
t.Fatalf("render/run requests = %d/%d, want 1/1", len(fake.RenderRequests), len(fake.RunRequests))
}
if got := fake.RenderRequests[0].InputPaths["previous_recap"]; got != previousPath {
t.Fatalf("render previous_recap input = %q, want %q", got, previousPath)
}
if got := fake.RunRequests[0].InputPaths["previous_recap"]; got != previousPath {
t.Fatalf("run previous_recap input = %q, want %q", got, previousPath)
}
}
func TestAnalyzeDoesNotCallObjectStoreForCanonicalPreviousSessionInput(t *testing.T) {
env, m, _ := setupAnalyzeEnv(t)
paths := sessionPathsForEnv(env, m.SessionID)
writeAnalyzeFile(t, filepath.Join(paths.TranscriptsDir, "trimmed.json"), `{"segments":[]}`)
previousPath := artifacts.SessionPreviousArtifactPath(paths, "artifacts/session_recap.md")
writeAnalyzeFile(t, previousPath, "previous recap\n")
m.Inputs = append(m.Inputs, manifest.InputRecord{
Kind: "previous_artifact",
Path: previousPath,
})
sessionRecap := env.Config.Pipeline.Scriptorium.Artifacts["session_recap"]
sessionRecap.Inputs["previous_recap"] = config.ScriptoriumInputConfig{
Source: "narratio.previous_session.artifact.session_recap",
Required: true,
}
env.Config.Pipeline.Scriptorium.Artifacts["session_recap"] = sessionRecap
tracker := &analyzeObjectStoreTracker{}
env.ObjectStore = tracker
_, err := (analyzeStage{}).Run(context.Background(), env, m)
if err != nil {
t.Fatalf("Run() error = %v", err)
}
if tracker.called {
t.Fatal("analyze should not call object store for previous-session input resolution")
}
}
func TestAnalyzeFailsWhenOutputPathMissing(t *testing.T) {
env, m, fake := setupAnalyzeEnv(t)
paths := sessionPathsForEnv(env, m.SessionID)
@@ -1145,3 +1286,27 @@ func mustArtifactEntryList(t *testing.T, metadata map[string]any, key string) []
}
return out
}
type analyzeObjectStoreTracker struct {
called bool
}
func (s *analyzeObjectStoreTracker) List(context.Context, string) ([]storage.ObjectInfo, error) {
s.called = true
return nil, errors.New("unexpected object store list call")
}
func (s *analyzeObjectStoreTracker) Download(context.Context, string, string) error {
s.called = true
return errors.New("unexpected object store download call")
}
func (s *analyzeObjectStoreTracker) Upload(context.Context, string, string, storage.UploadOptions) (storage.ObjectInfo, error) {
s.called = true
return storage.ObjectInfo{}, errors.New("unexpected object store upload call")
}
func (s *analyzeObjectStoreTracker) Exists(context.Context, string) (bool, error) {
s.called = true
return false, errors.New("unexpected object store exists call")
}

View File

@@ -118,6 +118,10 @@ func (archiveStage) Run(ctx context.Context, env *Env, m *manifest.Manifest) (*S
return nil, fmt.Errorf("archive: collect run files: %w", err)
}
sessionPaths := archiveSessionPaths(env, m)
previousFiles, err := collectArchivePreviousFiles(sessionPaths.PreviousDir)
if err != nil {
return nil, fmt.Errorf("archive: collect previous files: %w", err)
}
runtimeCatalog, err := buildArchiveRuntimeArtifactCatalog(sessionPaths, env.Config.Pipeline.Scriptorium)
if err != nil {
return nil, fmt.Errorf("archive: build runtime artifact catalog: %w", err)
@@ -144,6 +148,15 @@ func (archiveStage) Run(ctx context.Context, env *Env, m *manifest.Manifest) (*S
promotedUploaded = append(promotedUploaded, promotion.Dest)
}
previousUploaded := make([]string, 0, len(previousFiles))
for _, file := range previousFiles {
key := artifacts.S3PromotedArtifactKey(sessionPrefix, file.RelativePath)
if _, err := env.ObjectStore.Upload(ctx, file.LocalPath, key, storage.UploadOptions{}); err != nil {
return nil, fmt.Errorf("archive: upload previous file %q to %q: %w", file.RelativePath, key, err)
}
previousUploaded = append(previousUploaded, file.RelativePath)
}
currentManifestKey, currentRunPointerKey := artifacts.ResolveArchiveCurrentStateKeys(sessionPrefix)
manifestTempPath, err := writeCurrentManifestSnapshot(m, archiveMetadataPreview(
bucket,
@@ -151,6 +164,7 @@ func (archiveStage) Run(ctx context.Context, env *Env, m *manifest.Manifest) (*S
sessionPrefix,
runUploaded,
promotedUploaded,
previousUploaded,
skippedOptional,
currentManifestKey,
))
@@ -187,6 +201,8 @@ func (archiveStage) Run(ctx context.Context, env *Env, m *manifest.Manifest) (*S
"run_uploaded_paths": runUploaded,
"promoted_files_uploaded": len(promotedUploaded),
"promoted_paths": promotedUploaded,
"previous_files_uploaded": len(previousUploaded),
"previous_uploaded_paths": previousUploaded,
"skipped_optional_promotions": skippedOptional,
"current_manifest_key": currentManifestKey,
"current_run_id_key": currentRunPointerKey,
@@ -497,6 +513,51 @@ func collectArchiveRunFiles(runRoot, manifestPath string) ([]archiveUploadFile,
return files, nil
}
func collectArchivePreviousFiles(previousDir string) ([]archiveUploadFile, error) {
previousDir = filepath.Clean(strings.TrimSpace(previousDir))
if previousDir == "" {
return nil, fmt.Errorf("previous directory is required")
}
info, err := os.Stat(previousDir)
if err != nil {
if os.IsNotExist(err) {
return nil, nil
}
return nil, fmt.Errorf("stat %q: %w", previousDir, err)
}
if !info.IsDir() {
return nil, fmt.Errorf("previous path %q is not a directory", previousDir)
}
files := make([]archiveUploadFile, 0, 16)
err = filepath.WalkDir(previousDir, func(path string, d fs.DirEntry, walkErr error) error {
if walkErr != nil {
return walkErr
}
if d.IsDir() {
return nil
}
rel, err := filepath.Rel(previousDir, path)
if err != nil {
return fmt.Errorf("relative path from %q to %q: %w", previousDir, path, err)
}
rel = filepath.ToSlash(rel)
files = append(files, archiveUploadFile{
RelativePath: filepath.ToSlash(filepath.Join(config.PathPreviousDirSegment, rel)),
LocalPath: path,
})
return nil
})
if err != nil {
return nil, fmt.Errorf("walk %q: %w", previousDir, err)
}
sort.Slice(files, func(i, j int) bool {
return files[i].RelativePath < files[j].RelativePath
})
return files, nil
}
func resolveArchiveRunManifestSource(runRoot string) (string, error) {
path := filepath.Join(filepath.Clean(runRoot), "manifest.json")
info, err := os.Stat(path)
@@ -600,6 +661,7 @@ func archiveMetadataPreview(
bucket, runPrefix, sessionPrefix string,
runUploaded []string,
promotedUploaded []string,
previousUploaded []string,
skippedOptional []string,
currentManifestKey string,
) map[string]any {
@@ -612,6 +674,8 @@ func archiveMetadataPreview(
"run_uploaded_paths": append([]string(nil), runUploaded...),
"promoted_files_uploaded": len(promotedUploaded),
"promoted_paths": append([]string(nil), promotedUploaded...),
"previous_files_uploaded": len(previousUploaded),
"previous_uploaded_paths": append([]string(nil), previousUploaded...),
"skipped_optional_promotions": append([]string(nil), skippedOptional...),
"current_manifest_key": currentManifestKey,
"current_run_id_key": artifacts.S3CurrentRunPointerKey(sessionPrefix),

View File

@@ -130,6 +130,50 @@ func TestArchiveUploadsRunRecordPromotionsAndCurrentPointer(t *testing.T) {
if result.Metadata["promoted_files_uploaded"] != 2 {
t.Fatalf("metadata promoted_files_uploaded = %#v, want 2", result.Metadata["promoted_files_uploaded"])
}
if result.Metadata["previous_files_uploaded"] != 0 {
t.Fatalf("metadata previous_files_uploaded = %#v, want 0", result.Metadata["previous_files_uploaded"])
}
}
func TestArchiveUploadsPreviousCacheWhenPresent(t *testing.T) {
env, m, _ := archiveFixture(t)
fake := env.ObjectStore.(*storage.FakeBackend)
sessionRoot := artifacts.SessionWorkDirForCampaign(
env.Config.Pipeline.Workspace.Root,
env.Config.Session.Campaign,
env.Config.Session.SessionID,
)
writeStageTestFile(t, filepath.Join(sessionRoot, "previous", "manifest.json"), "{\"session_id\":\"2026-04-12\"}\n")
writeStageTestFile(t, filepath.Join(sessionRoot, "previous", "artifacts", "session_recap.md"), "# previous recap\n")
result, err := archiveStage{}.Run(context.Background(), env, m)
if err != nil {
t.Fatalf("Run() error = %v", err)
}
previousManifestKey := m.S3SessionPrefix + "previous/manifest.json"
previousRecapKey := m.S3SessionPrefix + "previous/artifacts/session_recap.md"
if _, ok := fake.Objects[previousManifestKey]; !ok {
t.Fatalf("missing archived previous manifest key %q", previousManifestKey)
}
if _, ok := fake.Objects[previousRecapKey]; !ok {
t.Fatalf("missing archived previous artifact key %q", previousRecapKey)
}
if result.Metadata["previous_files_uploaded"] != 2 {
t.Fatalf("metadata previous_files_uploaded = %#v, want 2", result.Metadata["previous_files_uploaded"])
}
}
func TestArchiveToleratesMissingPreviousCache(t *testing.T) {
env, m, _ := archiveFixture(t)
result, err := archiveStage{}.Run(context.Background(), env, m)
if err != nil {
t.Fatalf("Run() error = %v", err)
}
if result.Metadata["previous_files_uploaded"] != 0 {
t.Fatalf("metadata previous_files_uploaded = %#v, want 0", result.Metadata["previous_files_uploaded"])
}
}
func TestArchiveUsesCustomPromotionRules(t *testing.T) {

View File

@@ -145,6 +145,20 @@ func (prepareStage) Run(ctx context.Context, env *Env, m *manifest.Manifest) (*S
}
}
previousRequirements := collectPreparePreviousRequirements(env.Config)
var previousHydration *previousSessionHydrationResult
if len(previousRequirements) > 0 {
if err := clearManagedPreviousState(paths); err != nil {
return nil, fmt.Errorf("prepare: clear previous-session cache: %w", err)
}
hydration, err := hydratePreviousSessionArtifacts(ctx, env, paths, previousRequirements)
if err != nil {
return nil, fmt.Errorf("prepare: hydrate previous-session artifacts: %w", err)
}
previousHydration = hydration
inputs = append(inputs, hydration.Inputs...)
}
sort.Slice(inputs, func(i, j int) bool {
if inputs[i].Kind != inputs[j].Kind {
return inputs[i].Kind < inputs[j].Kind
@@ -153,13 +167,27 @@ func (prepareStage) Run(ctx context.Context, env *Env, m *manifest.Manifest) (*S
})
m.Inputs = inputs
metadata := map[string]any{
"prepared": true,
"stage": "prepare",
"inputs_count": len(inputs),
"audio_files_resolved": countAudioInputs(inputs),
}
if len(previousRequirements) > 0 {
metadata["previous_requirements_count"] = len(previousRequirements)
if previousHydration != nil {
metadata["previous_artifacts_hydrated"] = append([]string(nil), previousHydration.Hydrated...)
metadata["previous_artifacts_hydrated_count"] = len(previousHydration.Hydrated)
metadata["previous_artifacts_missing_optional"] = append([]string(nil), previousHydration.SkippedMissing...)
metadata["previous_artifacts_missing_optional_count"] = len(previousHydration.SkippedMissing)
if strings.TrimSpace(previousHydration.PreviousRunID) != "" {
metadata["previous_session_run_id"] = strings.TrimSpace(previousHydration.PreviousRunID)
}
}
}
return &StageResult{
Metadata: map[string]any{
"prepared": true,
"stage": "prepare",
"inputs_count": len(inputs),
"audio_files_resolved": countAudioInputs(inputs),
},
Metadata: metadata,
}, nil
}
@@ -354,6 +382,35 @@ func countAudioInputs(inputs []manifest.InputRecord) int {
return count
}
func collectPreparePreviousRequirements(cfg *config.Config) []artifacts.PreviousArtifactRequirement {
if cfg == nil || cfg.Pipeline == nil || cfg.Pipeline.Scriptorium == nil {
return nil
}
return artifacts.CollectPreviousArtifactRequirements(cfg.Pipeline.Scriptorium.Artifacts)
}
func clearManagedPreviousState(paths artifacts.SessionPaths) error {
previousDir := filepath.Clean(paths.PreviousDir)
sessionRoot := filepath.Clean(paths.Root)
if strings.TrimSpace(previousDir) == "" || strings.TrimSpace(sessionRoot) == "" {
return fmt.Errorf("previous/session root paths are required")
}
if previousDir == sessionRoot {
return fmt.Errorf("refusing to clear session root as previous cache: %q", previousDir)
}
prefix := sessionRoot + string(filepath.Separator)
if !strings.HasPrefix(previousDir, prefix) {
return fmt.Errorf("refusing to clear path outside session root: %q", previousDir)
}
if filepath.Base(previousDir) != config.PathPreviousDirSegment {
return fmt.Errorf("refusing to clear non-previous path %q", previousDir)
}
if err := os.RemoveAll(previousDir); err != nil {
return err
}
return os.MkdirAll(previousDir, 0o755)
}
func pathsWorkDirForManifest(env *Env, m *manifest.Manifest, sessionID string) string {
if env == nil || env.Config == nil || env.Config.Pipeline == nil || env.Config.Session == nil {
return ""

View File

@@ -0,0 +1,475 @@
package stage
import (
"context"
"fmt"
"os"
"path"
"path/filepath"
"sort"
"strings"
"gitea.maximumdirect.net/eric/narratio/internal/artifacts"
"gitea.maximumdirect.net/eric/narratio/internal/config"
"gitea.maximumdirect.net/eric/narratio/internal/manifest"
)
const (
preparePreviousInputKindManifest = "previous_manifest"
preparePreviousInputKindArtifact = "previous_artifact"
preparePreviousInputSource = "previous_session_archive.current"
)
// previousSessionHydrationResult captures prepare-time previous-session cache materialization.
type previousSessionHydrationResult struct {
Inputs []manifest.InputRecord
Hydrated []string
SkippedMissing []string
PreviousRunID string
}
func hydratePreviousSessionArtifacts(
ctx context.Context,
env *Env,
paths artifacts.SessionPaths,
requirements []artifacts.PreviousArtifactRequirement,
) (*previousSessionHydrationResult, error) {
if len(requirements) == 0 {
return &previousSessionHydrationResult{}, nil
}
if env == nil || env.Config == nil || env.Config.Session == nil || env.Config.Pipeline == nil {
return nil, fmt.Errorf("resolved config with session/pipeline is required")
}
if env.ArtifactStore == nil {
return nil, fmt.Errorf("artifact store is required")
}
orderedRequirements := append([]artifacts.PreviousArtifactRequirement(nil), requirements...)
sort.Slice(orderedRequirements, func(i, j int) bool {
return orderedRequirements[i].Name < orderedRequirements[j].Name
})
requiredNames := requiredPreviousArtifactNames(orderedRequirements)
optionalNames := optionalPreviousArtifactNames(orderedRequirements)
previousSessionID := strings.TrimSpace(env.Config.Session.PreviousSessionID)
if previousSessionID == "" {
if len(requiredNames) > 0 {
return nil, fmt.Errorf(
"previous_session_id is required for required previous-session artifacts: %s",
strings.Join(requiredNames, ", "),
)
}
return &previousSessionHydrationResult{SkippedMissing: optionalNames}, nil
}
if env.ObjectStore == nil {
return nil, fmt.Errorf("previous-session artifact hydration requires object store backend")
}
if env.Config.Pipeline.Storage.S3 == nil {
return nil, fmt.Errorf("pipeline.storage.s3 configuration is required for previous-session artifact hydration")
}
campaign := strings.TrimSpace(env.Config.Session.Campaign)
if campaign == "" {
return nil, fmt.Errorf("session campaign is required for previous-session artifact hydration")
}
bucket := strings.TrimSpace(env.Config.Pipeline.Storage.S3.Bucket)
if bucket == "" {
return nil, fmt.Errorf("pipeline.storage.s3.bucket is required for previous-session artifact hydration")
}
previousSessionPrefix := artifacts.S3SessionPrefix(
env.Config.Pipeline.Storage.S3.RootPrefix,
campaign,
previousSessionID,
)
currentManifestKey, currentRunIDKey := artifacts.ResolveArchiveCurrentStateKeys(previousSessionPrefix)
result := &previousSessionHydrationResult{}
runPointerExists, err := env.ObjectStore.Exists(ctx, currentRunIDKey)
if err != nil {
return nil, fmt.Errorf("check previous-session current run pointer %q: %w", currentRunIDKey, err)
}
if !runPointerExists {
if len(requiredNames) > 0 {
return nil, fmt.Errorf("required previous-session artifacts unavailable: remote current run pointer missing: %q", currentRunIDKey)
}
result.SkippedMissing = optionalNames
return result, nil
}
runIDTemp, err := downloadObjectToTempStage(ctx, env.ObjectStore, currentRunIDKey, "narratio-prepare-previous-run-id-*.txt")
if err != nil {
return nil, fmt.Errorf("download previous-session current run pointer %q: %w", currentRunIDKey, err)
}
defer func() { _ = os.Remove(runIDTemp) }()
runIDBytes, err := os.ReadFile(runIDTemp)
if err != nil {
return nil, fmt.Errorf("read previous-session current run pointer %q: %w", currentRunIDKey, err)
}
previousRunID := strings.TrimSpace(string(runIDBytes))
if previousRunID == "" {
return nil, fmt.Errorf("previous-session current run pointer %q is empty", currentRunIDKey)
}
result.PreviousRunID = previousRunID
manifestExists, err := env.ObjectStore.Exists(ctx, currentManifestKey)
if err != nil {
return nil, fmt.Errorf("check previous-session current manifest %q: %w", currentManifestKey, err)
}
if !manifestExists {
if len(requiredNames) > 0 {
return nil, fmt.Errorf("required previous-session artifacts unavailable: remote current manifest missing: %q", currentManifestKey)
}
result.SkippedMissing = optionalNames
return result, nil
}
if err := os.MkdirAll(filepath.Dir(paths.PreviousManifestPath), 0o755); err != nil {
return nil, fmt.Errorf("create previous manifest directory: %w", err)
}
if err := env.ObjectStore.Download(ctx, currentManifestKey, paths.PreviousManifestPath); err != nil {
return nil, fmt.Errorf("download previous-session current manifest %q: %w", currentManifestKey, err)
}
manifestStore := &manifest.LocalStore{}
previousManifest, err := manifestStore.Load(ctx, paths.PreviousManifestPath)
if err != nil {
return nil, fmt.Errorf("decode downloaded previous-session manifest %q: %w", currentManifestKey, err)
}
if strings.TrimSpace(previousManifest.SessionID) != previousSessionID {
return nil, fmt.Errorf(
"previous-session manifest session_id %q does not match configured previous_session_id %q",
strings.TrimSpace(previousManifest.SessionID),
previousSessionID,
)
}
if strings.TrimSpace(previousManifest.Campaign) != campaign {
return nil, fmt.Errorf(
"previous-session manifest campaign %q does not match current campaign %q",
strings.TrimSpace(previousManifest.Campaign),
campaign,
)
}
if strings.TrimSpace(previousManifest.RunID) == "" {
return nil, fmt.Errorf("previous-session manifest run_id is required")
}
if strings.TrimSpace(previousManifest.RunID) != previousRunID {
return nil, fmt.Errorf(
"previous-session current run pointer %q references run %q but current manifest run_id is %q",
currentRunIDKey,
previousRunID,
strings.TrimSpace(previousManifest.RunID),
)
}
manifestChecksum, err := env.ArtifactStore.Checksum(paths.PreviousManifestPath)
if err != nil {
return nil, fmt.Errorf("checksum downloaded previous-session manifest: %w", err)
}
result.Inputs = append(result.Inputs, manifest.InputRecord{
Kind: preparePreviousInputKindManifest,
Path: paths.PreviousManifestPath,
Checksum: manifestChecksum,
Source: preparePreviousInputSource,
S3Bucket: bucket,
S3Key: currentManifestKey,
})
for _, requirement := range orderedRequirements {
candidates := previousArtifactRelativePathCandidates(requirement.Name, previousManifest, env.Config)
if len(candidates) == 0 {
if requirement.Required {
return nil, fmt.Errorf(
"required previous-session artifact %q is unavailable in previous-session manifest/archive",
requirement.Name,
)
}
result.SkippedMissing = append(result.SkippedMissing, requirement.Name)
continue
}
selectedRel := ""
selectedKey := ""
for _, candidate := range candidates {
remoteKey := artifacts.S3PromotedArtifactKey(previousSessionPrefix, candidate)
exists, err := env.ObjectStore.Exists(ctx, remoteKey)
if err != nil {
return nil, fmt.Errorf("check previous-session artifact object %q: %w", remoteKey, err)
}
if !exists {
continue
}
selectedRel = candidate
selectedKey = remoteKey
break
}
if selectedRel == "" {
if requirement.Required {
return nil, fmt.Errorf(
"required previous-session artifact %q object missing from archive candidate keys",
requirement.Name,
)
}
result.SkippedMissing = append(result.SkippedMissing, requirement.Name)
continue
}
localPath := artifacts.SessionPreviousArtifactPath(paths, selectedRel)
if err := os.MkdirAll(filepath.Dir(localPath), 0o755); err != nil {
return nil, fmt.Errorf("create previous-session artifact directory for %q: %w", localPath, err)
}
if err := env.ObjectStore.Download(ctx, selectedKey, localPath); err != nil {
return nil, fmt.Errorf("download previous-session artifact %q from %q: %w", requirement.Name, selectedKey, err)
}
if err := requireNonEmptyFile(localPath, "previous-session artifact "+requirement.Name); err != nil {
return nil, err
}
checksum, err := env.ArtifactStore.Checksum(localPath)
if err != nil {
return nil, fmt.Errorf("checksum previous-session artifact %q: %w", requirement.Name, err)
}
result.Inputs = append(result.Inputs, manifest.InputRecord{
Kind: preparePreviousInputKindArtifact,
Path: localPath,
Checksum: checksum,
Source: preparePreviousInputSource,
S3Bucket: bucket,
S3Key: selectedKey,
})
result.Hydrated = append(result.Hydrated, requirement.Name)
}
sort.Strings(result.Hydrated)
sort.Strings(result.SkippedMissing)
sort.Slice(result.Inputs, func(i, j int) bool {
if result.Inputs[i].Kind != result.Inputs[j].Kind {
return result.Inputs[i].Kind < result.Inputs[j].Kind
}
return result.Inputs[i].Path < result.Inputs[j].Path
})
return result, nil
}
func requiredPreviousArtifactNames(requirements []artifacts.PreviousArtifactRequirement) []string {
names := make([]string, 0, len(requirements))
for _, requirement := range requirements {
if requirement.Required {
names = append(names, strings.TrimSpace(requirement.Name))
}
}
sort.Strings(names)
return names
}
func optionalPreviousArtifactNames(requirements []artifacts.PreviousArtifactRequirement) []string {
names := make([]string, 0, len(requirements))
for _, requirement := range requirements {
if requirement.Required {
continue
}
names = append(names, strings.TrimSpace(requirement.Name))
}
sort.Strings(names)
return names
}
func previousArtifactRelativePathCandidates(
artifactName string,
previousManifest *manifest.Manifest,
cfg *config.Config,
) []string {
candidates := []string{}
appendCandidate := func(v string) {
normalized, err := normalizeArchiveRelativePath(v)
if err != nil {
return
}
candidates = append(candidates, normalized)
}
sourceID := artifacts.ConfiguredArtifactSourceID(artifactName)
if rel, ok := previousManifestArtifactRelativePathBySourceID(previousManifest, sourceID); ok {
appendCandidate(rel)
base := path.Base(rel)
for _, promoted := range previousManifestPromotedPaths(previousManifest) {
if path.Base(promoted) == base {
appendCandidate(promoted)
}
}
}
if cfg != nil && cfg.Pipeline != nil && cfg.Pipeline.Scriptorium != nil {
if artifactCfg, ok := cfg.Pipeline.Scriptorium.Artifacts[artifactName]; ok {
appendCandidate(artifactCfg.OutputPath)
}
}
return dedupeOrderedStrings(candidates)
}
func previousManifestArtifactRelativePathBySourceID(previousManifest *manifest.Manifest, sourceID string) (string, bool) {
if previousManifest == nil || len(previousManifest.Stages) == 0 {
return "", false
}
sourceID = strings.TrimSpace(sourceID)
if sourceID == "" {
return "", false
}
stageNames := make([]string, 0, len(previousManifest.Stages))
if _, ok := previousManifest.Stages["analyze"]; ok {
stageNames = append(stageNames, "analyze")
}
for stageName := range previousManifest.Stages {
if stageName == "analyze" {
continue
}
stageNames = append(stageNames, stageName)
}
start := 0
if len(stageNames) > 0 && stageNames[0] == "analyze" {
start = 1
}
sort.Strings(stageNames[start:])
for _, stageName := range stageNames {
sr := previousManifest.Stages[stageName]
if sr == nil {
continue
}
for _, out := range sr.Outputs {
if strings.TrimSpace(out.SourceID) != sourceID {
continue
}
rel, ok := derivePreviousManifestRelativePath(previousManifest, out.LocalPath)
if ok {
return rel, true
}
}
}
return "", false
}
func derivePreviousManifestRelativePath(previousManifest *manifest.Manifest, localPath string) (string, bool) {
trimmed := strings.TrimSpace(localPath)
if trimmed == "" {
return "", false
}
if !filepath.IsAbs(trimmed) {
normalized, err := normalizeArchiveRelativePath(filepath.ToSlash(trimmed))
if err != nil {
return "", false
}
return normalized, true
}
sessionRoot, ok := previousManifestSessionRoot(previousManifest)
if !ok {
return "", false
}
rel, err := filepath.Rel(sessionRoot, trimmed)
if err != nil {
return "", false
}
normalized, err := normalizeArchiveRelativePath(filepath.ToSlash(rel))
if err != nil {
return "", false
}
return normalized, true
}
func previousManifestSessionRoot(previousManifest *manifest.Manifest) (string, bool) {
if previousManifest == nil {
return "", false
}
runRoot := filepath.Clean(strings.TrimSpace(previousManifest.LocalWorkDir))
runID := strings.TrimSpace(previousManifest.RunID)
if runRoot == "" || runID == "" {
return "", false
}
if filepath.Base(runRoot) != runID {
return "", false
}
runsDir := filepath.Dir(runRoot)
if filepath.Base(runsDir) != config.PathRunsDirSegment {
return "", false
}
return filepath.Dir(runsDir), true
}
func previousManifestPromotedPaths(previousManifest *manifest.Manifest) []string {
if previousManifest == nil || len(previousManifest.Stages) == 0 {
return nil
}
sr := previousManifest.Stages["archive"]
if sr == nil || sr.Metadata == nil {
return nil
}
raw, ok := sr.Metadata["promoted_paths"]
if !ok {
return nil
}
values, ok := raw.([]any)
if !ok {
return nil
}
out := make([]string, 0, len(values))
for _, value := range values {
asString, ok := value.(string)
if !ok {
continue
}
normalized, err := normalizeArchiveRelativePath(asString)
if err != nil {
continue
}
out = append(out, normalized)
}
return dedupeOrderedStrings(out)
}
func dedupeOrderedStrings(values []string) []string {
if len(values) == 0 {
return nil
}
seen := map[string]struct{}{}
out := make([]string, 0, len(values))
for _, value := range values {
trimmed := strings.TrimSpace(value)
if trimmed == "" {
continue
}
if _, ok := seen[trimmed]; ok {
continue
}
seen[trimmed] = struct{}{}
out = append(out, trimmed)
}
return out
}
func downloadObjectToTempStage(
ctx context.Context,
store interface {
Download(context.Context, string, string) error
},
key, pattern string,
) (string, error) {
tmp, err := os.CreateTemp("", pattern)
if err != nil {
return "", fmt.Errorf("create temp file: %w", err)
}
path := tmp.Name()
if err := tmp.Close(); err != nil {
_ = os.Remove(path)
return "", fmt.Errorf("close temp file: %w", err)
}
if err := store.Download(ctx, key, path); err != nil {
_ = os.Remove(path)
return "", err
}
return path, nil
}

View File

@@ -0,0 +1,415 @@
package stage
import (
"context"
"encoding/json"
"os"
"path/filepath"
"strings"
"testing"
"time"
"gitea.maximumdirect.net/eric/narratio/internal/adapters/storage"
"gitea.maximumdirect.net/eric/narratio/internal/artifacts"
"gitea.maximumdirect.net/eric/narratio/internal/config"
"gitea.maximumdirect.net/eric/narratio/internal/manifest"
)
func TestHydratePreviousSessionArtifactsDownloadsManifestAndRequiredArtifact(t *testing.T) {
env, sessionPaths, fake := previousHydrationFixture(t)
requirements := []artifacts.PreviousArtifactRequirement{
{Name: "session_recap", Required: true},
}
seedPreviousCurrentState(t, env, fake, previousStateSeedOptions{
includeRunPointerObject: true,
includeArtifactObject: true,
artifactBody: "# previous recap\n",
})
result, err := hydratePreviousSessionArtifacts(context.Background(), env, sessionPaths, requirements)
if err != nil {
t.Fatalf("hydratePreviousSessionArtifacts() error = %v", err)
}
if result == nil {
t.Fatal("result is nil")
}
if len(result.Inputs) != 2 {
t.Fatalf("inputs len = %d, want 2", len(result.Inputs))
}
if !containsString(result.Hydrated, "session_recap") {
t.Fatalf("hydrated = %#v, want session_recap", result.Hydrated)
}
if len(result.SkippedMissing) != 0 {
t.Fatalf("skipped missing = %#v, want none", result.SkippedMissing)
}
if _, err := os.Stat(sessionPaths.PreviousManifestPath); err != nil {
t.Fatalf("previous manifest missing: %v", err)
}
recapPath := artifacts.SessionPreviousArtifactPath(sessionPaths, "artifacts/session_recap.md")
if _, err := os.Stat(recapPath); err != nil {
t.Fatalf("previous artifact missing: %v", err)
}
manifestInput := findInputByKind(result.Inputs, preparePreviousInputKindManifest)
if manifestInput == nil {
t.Fatalf("missing input kind %q", preparePreviousInputKindManifest)
}
if manifestInput.Source != preparePreviousInputSource {
t.Fatalf("manifest input source = %q, want %q", manifestInput.Source, preparePreviousInputSource)
}
artifactInput := findInputByKind(result.Inputs, preparePreviousInputKindArtifact)
if artifactInput == nil {
t.Fatalf("missing input kind %q", preparePreviousInputKindArtifact)
}
if artifactInput.Source != preparePreviousInputSource {
t.Fatalf("artifact input source = %q, want %q", artifactInput.Source, preparePreviousInputSource)
}
if artifactInput.Path != recapPath {
t.Fatalf("artifact input path = %q, want %q", artifactInput.Path, recapPath)
}
}
func TestHydratePreviousSessionArtifactsSkipsMissingOptionalArtifact(t *testing.T) {
env, sessionPaths, fake := previousHydrationFixture(t)
requirements := []artifacts.PreviousArtifactRequirement{
{Name: "session_recap", Required: false},
}
seedPreviousCurrentState(t, env, fake, previousStateSeedOptions{
includeRunPointerObject: true,
includeManifestObject: true,
includeArtifactObject: false,
})
result, err := hydratePreviousSessionArtifacts(context.Background(), env, sessionPaths, requirements)
if err != nil {
t.Fatalf("hydratePreviousSessionArtifacts() error = %v", err)
}
if result == nil {
t.Fatal("result is nil")
}
if !containsString(result.SkippedMissing, "session_recap") {
t.Fatalf("skipped missing = %#v, want session_recap", result.SkippedMissing)
}
if len(result.Hydrated) != 0 {
t.Fatalf("hydrated = %#v, want none", result.Hydrated)
}
manifestInput := findInputByKind(result.Inputs, preparePreviousInputKindManifest)
if manifestInput == nil {
t.Fatalf("missing input kind %q", preparePreviousInputKindManifest)
}
if findInputByKind(result.Inputs, preparePreviousInputKindArtifact) != nil {
t.Fatalf("unexpected %q input for missing optional artifact", preparePreviousInputKindArtifact)
}
}
func TestHydratePreviousSessionArtifactsFailsMissingRequiredArtifact(t *testing.T) {
env, sessionPaths, fake := previousHydrationFixture(t)
requirements := []artifacts.PreviousArtifactRequirement{
{Name: "session_recap", Required: true},
}
seedPreviousCurrentState(t, env, fake, previousStateSeedOptions{
includeRunPointerObject: true,
includeArtifactObject: false,
})
_, err := hydratePreviousSessionArtifacts(context.Background(), env, sessionPaths, requirements)
if err == nil || !strings.Contains(err.Error(), "required previous-session artifact") {
t.Fatalf("error = %v, want required artifact failure", err)
}
}
func TestHydratePreviousSessionArtifactsFailsRequiredWhenPreviousSessionIDUnset(t *testing.T) {
env, sessionPaths, _ := previousHydrationFixture(t)
env.Config.Session.PreviousSessionID = ""
requirements := []artifacts.PreviousArtifactRequirement{
{Name: "session_recap", Required: true},
}
_, err := hydratePreviousSessionArtifacts(context.Background(), env, sessionPaths, requirements)
if err == nil || !strings.Contains(err.Error(), "previous_session_id is required") {
t.Fatalf("error = %v, want previous_session_id required failure", err)
}
}
func TestHydratePreviousSessionArtifactsOptionalWithNoPreviousSessionID(t *testing.T) {
env, sessionPaths, _ := previousHydrationFixture(t)
env.Config.Session.PreviousSessionID = ""
requirements := []artifacts.PreviousArtifactRequirement{
{Name: "session_recap", Required: false},
}
result, err := hydratePreviousSessionArtifacts(context.Background(), env, sessionPaths, requirements)
if err != nil {
t.Fatalf("hydratePreviousSessionArtifacts() error = %v", err)
}
if result == nil {
t.Fatal("result is nil")
}
if len(result.Inputs) != 0 {
t.Fatalf("inputs len = %d, want 0", len(result.Inputs))
}
if !containsString(result.SkippedMissing, "session_recap") {
t.Fatalf("skipped missing = %#v, want session_recap", result.SkippedMissing)
}
if _, err := os.Stat(sessionPaths.PreviousManifestPath); !os.IsNotExist(err) {
t.Fatalf("previous manifest should not be created, stat err = %v", err)
}
}
func TestHydratePreviousSessionArtifactsRespectsCurrentCommitMarker(t *testing.T) {
env, sessionPaths, fake := previousHydrationFixture(t)
requirements := []artifacts.PreviousArtifactRequirement{
{Name: "session_recap", Required: true},
}
seedPreviousCurrentState(t, env, fake, previousStateSeedOptions{
includeRunPointerObject: false,
includeArtifactObject: true,
artifactBody: "# previous recap\n",
})
_, err := hydratePreviousSessionArtifacts(context.Background(), env, sessionPaths, requirements)
if err == nil || !strings.Contains(err.Error(), "remote current run pointer missing") {
t.Fatalf("error = %v, want current run pointer missing failure", err)
}
}
func TestHydratePreviousSessionArtifactsDoesNotUseLocalPreviousWorkspaceState(t *testing.T) {
env, sessionPaths, fake := previousHydrationFixture(t)
requirements := []artifacts.PreviousArtifactRequirement{
{Name: "session_recap", Required: true},
}
seed := seedPreviousCurrentState(t, env, fake, previousStateSeedOptions{
includeRunPointerObject: true,
includeManifestObject: true,
includeArtifactObject: false,
})
// Write a local previous-session workspace file that should be ignored.
writeFile(t, filepath.Join(seed.PreviousSessionRoot, "artifacts", "session_recap.md"), "# local stale recap\n")
_, err := hydratePreviousSessionArtifacts(context.Background(), env, sessionPaths, requirements)
if err == nil || !strings.Contains(err.Error(), "object missing from archive") {
t.Fatalf("error = %v, want remote-object-missing failure", err)
}
}
func TestHydratePreviousSessionArtifactsUsesExplicitStorageKeys(t *testing.T) {
env, sessionPaths, fake := previousHydrationFixture(t)
requirements := []artifacts.PreviousArtifactRequirement{
{Name: "session_recap", Required: true},
}
seed := seedPreviousCurrentState(t, env, fake, previousStateSeedOptions{
includeRunPointerObject: true,
includeArtifactObject: true,
artifactBody: "# previous recap\n",
})
capture := &preparePreviousCaptureStore{delegate: fake}
env.ObjectStore = capture
_, err := hydratePreviousSessionArtifacts(context.Background(), env, sessionPaths, requirements)
if err != nil {
t.Fatalf("hydratePreviousSessionArtifacts() error = %v", err)
}
if !containsString(capture.existsKeys, seed.RunPointerKey) {
t.Fatalf("exists keys = %#v, want %q", capture.existsKeys, seed.RunPointerKey)
}
if !containsString(capture.existsKeys, seed.ManifestKey) {
t.Fatalf("exists keys = %#v, want %q", capture.existsKeys, seed.ManifestKey)
}
if !containsString(capture.existsKeys, seed.ArtifactKey) {
t.Fatalf("exists keys = %#v, want %q", capture.existsKeys, seed.ArtifactKey)
}
if !containsString(capture.downloadKeys, seed.RunPointerKey) {
t.Fatalf("download keys = %#v, want %q", capture.downloadKeys, seed.RunPointerKey)
}
if !containsString(capture.downloadKeys, seed.ManifestKey) {
t.Fatalf("download keys = %#v, want %q", capture.downloadKeys, seed.ManifestKey)
}
if !containsString(capture.downloadKeys, seed.ArtifactKey) {
t.Fatalf("download keys = %#v, want %q", capture.downloadKeys, seed.ArtifactKey)
}
}
type previousStateSeedResult struct {
PreviousSessionPrefix string
PreviousSessionRoot string
ManifestKey string
RunPointerKey string
ArtifactKey string
}
type previousStateSeedOptions struct {
includeRunPointerObject bool
includeManifestObject bool
includeArtifactObject bool
artifactBody string
}
func seedPreviousCurrentState(
t *testing.T,
env *Env,
fake *storage.FakeBackend,
options previousStateSeedOptions,
) previousStateSeedResult {
t.Helper()
if !options.includeRunPointerObject && !options.includeManifestObject && !options.includeArtifactObject {
// Keep default behavior deterministic when caller omits explicit flags.
options.includeRunPointerObject = true
options.includeManifestObject = true
}
if options.includeManifestObject == false && options.includeArtifactObject {
options.includeManifestObject = true
}
previousSessionID := strings.TrimSpace(env.Config.Session.PreviousSessionID)
campaign := strings.TrimSpace(env.Config.Session.Campaign)
rootPrefix := strings.TrimSpace(env.Config.Pipeline.Storage.S3.RootPrefix)
previousSessionPrefix := artifacts.S3SessionPrefix(rootPrefix, campaign, previousSessionID)
manifestKey, runPointerKey := artifacts.ResolveArchiveCurrentStateKeys(previousSessionPrefix)
previousRunID := "20260510T010203Z-a1b2c3d4"
previousSessionRoot := filepath.Join(t.TempDir(), "work", campaign, previousSessionID)
artifactLocalPath := filepath.Join(previousSessionRoot, "artifacts", "session_recap.md")
previousManifest := buildPreviousManifestForSeed(
t,
previousSessionID,
campaign,
previousRunID,
filepath.Join(previousSessionRoot, "runs", previousRunID),
artifactLocalPath,
)
if options.includeRunPointerObject || (!options.includeRunPointerObject && !options.includeManifestObject && !options.includeArtifactObject) {
fake.SeedObject(storage.FakeObject{
Key: runPointerKey,
Data: []byte(previousRunID + "\n"),
})
}
if options.includeManifestObject || options.includeArtifactObject {
fake.SeedObject(storage.FakeObject{
Key: manifestKey,
Data: previousManifest,
})
}
artifactKey := artifacts.S3PromotedArtifactKey(previousSessionPrefix, "artifacts/session_recap.md")
if options.includeArtifactObject {
body := options.artifactBody
if body == "" {
body = "# previous recap\n"
}
fake.SeedObject(storage.FakeObject{
Key: artifactKey,
Data: []byte(body),
})
}
return previousStateSeedResult{
PreviousSessionPrefix: previousSessionPrefix,
PreviousSessionRoot: previousSessionRoot,
ManifestKey: manifestKey,
RunPointerKey: runPointerKey,
ArtifactKey: artifactKey,
}
}
func previousHydrationFixture(t *testing.T) (*Env, artifacts.SessionPaths, *storage.FakeBackend) {
t.Helper()
env, m := setupPrepareEnv(t)
env.Config.Session.Campaign = "forsaken"
env.Config.Session.PreviousSessionID = "2026-05-10"
env.Config.Pipeline.Storage.S3 = &config.StorageS3Config{
Bucket: "my-dnd-archive",
RootPrefix: "dnd",
}
env.Config.Pipeline.Scriptorium = &config.ScriptoriumConfig{
Artifacts: map[string]config.ScriptoriumArtifactConfig{
"session_recap": {
Enabled: true,
OutputPath: "artifacts/session_recap.md",
},
},
}
fake := &storage.FakeBackend{}
env.ObjectStore = fake
paths, err := ensureLayoutForEnv(env, m.SessionID)
if err != nil {
t.Fatalf("ensureLayoutForEnv() error = %v", err)
}
return env, paths, fake
}
func buildPreviousManifestForSeed(
t *testing.T,
sessionID, campaign, runID, runRoot, artifactPath string,
) []byte {
t.Helper()
now := time.Date(2026, 5, 19, 22, 0, 0, 0, time.UTC)
m := manifest.New(sessionID, now)
m.Campaign = campaign
m.RunID = runID
m.LocalWorkDir = runRoot
m.MarkStageSucceeded("analyze", now, []manifest.ArtifactRecord{
{
Kind: "scriptorium_artifact",
SourceID: artifacts.ConfiguredArtifactSourceID("session_recap"),
LocalPath: artifactPath,
},
})
m.MarkStageSucceeded("archive", now, nil)
m.Stages["archive"].Metadata = map[string]any{
"promoted_paths": []string{"artifacts/session_recap.md"},
}
data, err := json.MarshalIndent(m, "", " ")
if err != nil {
t.Fatalf("marshal manifest: %v", err)
}
return append(data, '\n')
}
type preparePreviousCaptureStore struct {
delegate storage.ObjectStore
existsKeys []string
downloadKeys []string
}
func (s *preparePreviousCaptureStore) List(ctx context.Context, prefix string) ([]storage.ObjectInfo, error) {
return s.delegate.List(ctx, prefix)
}
func (s *preparePreviousCaptureStore) Download(ctx context.Context, key, localPath string) error {
s.downloadKeys = append(s.downloadKeys, key)
return s.delegate.Download(ctx, key, localPath)
}
func (s *preparePreviousCaptureStore) Upload(ctx context.Context, localPath, key string, opts storage.UploadOptions) (storage.ObjectInfo, error) {
return s.delegate.Upload(ctx, localPath, key, opts)
}
func (s *preparePreviousCaptureStore) Exists(ctx context.Context, key string) (bool, error) {
s.existsKeys = append(s.existsKeys, key)
return s.delegate.Exists(ctx, key)
}
func findInputByKind(inputs []manifest.InputRecord, kind string) *manifest.InputRecord {
for i := range inputs {
if inputs[i].Kind == kind {
return &inputs[i]
}
}
return nil
}
func containsString(values []string, target string) bool {
for _, value := range values {
if value == target {
return true
}
}
return false
}

View File

@@ -284,6 +284,188 @@ func TestPrepareStageAudioSourceConflictFails(t *testing.T) {
}
}
func TestPrepareStageWithoutPreviousRequirementsDoesNotTouchPreviousState(t *testing.T) {
env, m := setupPrepareEnv(t)
root := filepath.Dir(env.Config.SessionPath)
writeFile(t, filepath.Join(root, "audio", "a.flac"), "a")
env.Config.Session.Inputs.AudioFiles = []string{"./audio/a.flac"}
paths := sessionPathsForEnv(env, m.SessionID)
stalePath := filepath.Join(paths.PreviousDir, "stale.txt")
writeFile(t, stalePath, "stale")
_, err := (prepareStage{}).Run(context.Background(), env, m)
if err != nil {
t.Fatalf("prepare.Run() error = %v", err)
}
if _, err := os.Stat(stalePath); err != nil {
t.Fatalf("expected previous stale file to remain untouched: %v", err)
}
}
func TestPrepareStageOptionalPreviousArtifactWithoutPreviousSessionIDSucceeds(t *testing.T) {
env, m := setupPrepareEnv(t)
root := filepath.Dir(env.Config.SessionPath)
writeFile(t, filepath.Join(root, "audio", "a.flac"), "a")
env.Config.Session.Inputs.AudioFiles = []string{"./audio/a.flac"}
configurePreparePreviousArtifactSource(env, false)
env.Config.Session.PreviousSessionID = ""
paths := sessionPathsForEnv(env, m.SessionID)
writeFile(t, filepath.Join(paths.PreviousArtifactsDir, "stale.txt"), "stale")
result, err := (prepareStage{}).Run(context.Background(), env, m)
if err != nil {
t.Fatalf("prepare.Run() error = %v", err)
}
if result.Metadata["previous_requirements_count"] != 1 {
t.Fatalf("metadata previous_requirements_count = %#v, want 1", result.Metadata["previous_requirements_count"])
}
if result.Metadata["previous_artifacts_hydrated_count"] != 0 {
t.Fatalf("metadata previous_artifacts_hydrated_count = %#v, want 0", result.Metadata["previous_artifacts_hydrated_count"])
}
if result.Metadata["previous_artifacts_missing_optional_count"] != 1 {
t.Fatalf("metadata previous_artifacts_missing_optional_count = %#v, want 1", result.Metadata["previous_artifacts_missing_optional_count"])
}
if _, err := os.Stat(filepath.Join(paths.PreviousArtifactsDir, "stale.txt")); !os.IsNotExist(err) {
t.Fatalf("expected stale previous state to be cleared, stat err = %v", err)
}
for _, in := range m.Inputs {
if in.Kind == preparePreviousInputKindManifest || in.Kind == preparePreviousInputKindArtifact {
t.Fatalf("unexpected previous input record: %#v", in)
}
}
}
func TestPrepareStageRequiredPreviousArtifactWithoutPreviousSessionIDFails(t *testing.T) {
env, m := setupPrepareEnv(t)
root := filepath.Dir(env.Config.SessionPath)
writeFile(t, filepath.Join(root, "audio", "a.flac"), "a")
env.Config.Session.Inputs.AudioFiles = []string{"./audio/a.flac"}
configurePreparePreviousArtifactSource(env, true)
env.Config.Session.PreviousSessionID = ""
_, err := (prepareStage{}).Run(context.Background(), env, m)
if err == nil || !strings.Contains(err.Error(), "previous_session_id is required") {
t.Fatalf("error = %v, want previous_session_id required failure", err)
}
}
func TestPrepareStageHydratesRequiredPreviousArtifactAndRecordsInputs(t *testing.T) {
env, m := setupPrepareEnv(t)
root := filepath.Dir(env.Config.SessionPath)
writeFile(t, filepath.Join(root, "audio", "a.flac"), "a")
env.Config.Session.Inputs.AudioFiles = []string{"./audio/a.flac"}
env.Config.Session.Campaign = "forsaken"
env.Config.Session.PreviousSessionID = "2026-05-10"
env.Config.Pipeline.Storage.S3 = &config.StorageS3Config{Bucket: "my-dnd-archive", RootPrefix: "dnd"}
configurePreparePreviousArtifactSource(env, true)
fake := &storage.FakeBackend{}
env.ObjectStore = fake
seedPreviousCurrentState(t, env, fake, previousStateSeedOptions{
includeRunPointerObject: true,
includeArtifactObject: true,
artifactBody: "# prior recap\n",
})
paths := sessionPathsForEnv(env, m.SessionID)
result, err := (prepareStage{}).Run(context.Background(), env, m)
if err != nil {
t.Fatalf("prepare.Run() error = %v", err)
}
if _, err := os.Stat(paths.PreviousManifestPath); err != nil {
t.Fatalf("expected previous manifest: %v", err)
}
recapPath := artifacts.SessionPreviousArtifactPath(paths, "artifacts/session_recap.md")
if _, err := os.Stat(recapPath); err != nil {
t.Fatalf("expected previous artifact: %v", err)
}
var hasPreviousManifest, hasPreviousArtifact bool
for _, in := range m.Inputs {
if in.Kind == preparePreviousInputKindManifest {
hasPreviousManifest = true
}
if in.Kind == preparePreviousInputKindArtifact {
hasPreviousArtifact = true
}
}
if !hasPreviousManifest || !hasPreviousArtifact {
t.Fatalf("manifest inputs missing previous provenance: %#v", m.Inputs)
}
if result.Metadata["previous_artifacts_hydrated_count"] != 1 {
t.Fatalf("metadata previous_artifacts_hydrated_count = %#v, want 1", result.Metadata["previous_artifacts_hydrated_count"])
}
}
func TestPrepareStageRerunOverwritesPreviousCache(t *testing.T) {
env, m := setupPrepareEnv(t)
root := filepath.Dir(env.Config.SessionPath)
writeFile(t, filepath.Join(root, "audio", "a.flac"), "a")
env.Config.Session.Inputs.AudioFiles = []string{"./audio/a.flac"}
env.Config.Session.Campaign = "forsaken"
env.Config.Session.PreviousSessionID = "2026-05-10"
env.Config.Pipeline.Storage.S3 = &config.StorageS3Config{Bucket: "my-dnd-archive", RootPrefix: "dnd"}
configurePreparePreviousArtifactSource(env, true)
fake := &storage.FakeBackend{}
env.ObjectStore = fake
seed := seedPreviousCurrentState(t, env, fake, previousStateSeedOptions{
includeRunPointerObject: true,
includeArtifactObject: true,
artifactBody: "# old recap\n",
})
paths := sessionPathsForEnv(env, m.SessionID)
if _, err := (prepareStage{}).Run(context.Background(), env, m); err != nil {
t.Fatalf("first prepare run error = %v", err)
}
recapPath := artifacts.SessionPreviousArtifactPath(paths, "artifacts/session_recap.md")
firstBytes, err := os.ReadFile(recapPath)
if err != nil {
t.Fatalf("read first hydrated artifact: %v", err)
}
if strings.TrimSpace(string(firstBytes)) != "# old recap" {
t.Fatalf("first hydrated content = %q, want %q", strings.TrimSpace(string(firstBytes)), "# old recap")
}
writeFile(t, filepath.Join(paths.PreviousArtifactsDir, "stale.txt"), "stale")
fake.SeedObject(storage.FakeObject{Key: seed.ArtifactKey, Data: []byte("# new recap\n")})
if _, err := (prepareStage{}).Run(context.Background(), env, m); err != nil {
t.Fatalf("second prepare run error = %v", err)
}
secondBytes, err := os.ReadFile(recapPath)
if err != nil {
t.Fatalf("read second hydrated artifact: %v", err)
}
if strings.TrimSpace(string(secondBytes)) != "# new recap" {
t.Fatalf("second hydrated content = %q, want %q", strings.TrimSpace(string(secondBytes)), "# new recap")
}
if _, err := os.Stat(filepath.Join(paths.PreviousArtifactsDir, "stale.txt")); !os.IsNotExist(err) {
t.Fatalf("expected stale previous cache file to be removed, stat err = %v", err)
}
}
func configurePreparePreviousArtifactSource(env *Env, required bool) {
env.Config.Pipeline.Scriptorium = &config.ScriptoriumConfig{
Artifacts: map[string]config.ScriptoriumArtifactConfig{
"session_recap": {
Enabled: true,
OutputPath: "artifacts/session_recap.md",
Inputs: map[string]config.ScriptoriumInputConfig{
"previous_recap": {
Source: "narratio.previous_session.artifact.session_recap",
Required: required,
},
},
},
},
}
}
func setupPrepareEnv(t *testing.T) (*Env, *manifest.Manifest) {
t.Helper()
workspace := t.TempDir()

View File

@@ -7,6 +7,7 @@ import (
"strings"
"gitea.maximumdirect.net/eric/narratio/internal/artifacts"
"gitea.maximumdirect.net/eric/narratio/internal/config"
"gitea.maximumdirect.net/eric/narratio/internal/manifest"
)
@@ -99,6 +100,9 @@ func runLocalPathForCanonical(layout runStageLayout, sessionPaths artifacts.Sess
if rel == "." || rel == ".." || strings.HasPrefix(rel, ".."+string(filepath.Separator)) {
return "", fmt.Errorf("canonical path %q is outside session root %q", cleanCanonical, sessionPaths.Root)
}
if rel == config.PathPreviousDirSegment || strings.HasPrefix(rel, config.PathPreviousDirSegment+string(filepath.Separator)) {
return cleanCanonical, nil
}
localPath := filepath.Join(layout.OutputsDir, rel)
if err := os.MkdirAll(filepath.Dir(localPath), 0o755); err != nil {
return "", fmt.Errorf("create run-local output parent for %q: %w", localPath, err)

View File

@@ -33,3 +33,24 @@ func TestRunLocalPathForCanonicalCreatesParentDirectories(t *testing.T) {
t.Fatalf("expected run-local parent directory to exist: %v", err)
}
}
func TestRunLocalPathForCanonicalKeepsPreviousStateSessionDurable(t *testing.T) {
root := t.TempDir()
sessionRoot := filepath.Join(root, "work", "dilfs", "2026-05-17")
layout := runStageLayout{
Enabled: true,
OutputsDir: filepath.Join(sessionRoot, "runs", "run-1", "prepare", "outputs"),
}
if err := os.MkdirAll(layout.OutputsDir, 0o755); err != nil {
t.Fatalf("mkdir outputs dir: %v", err)
}
canonical := filepath.Join(sessionRoot, "previous", "artifacts", "session_recap.md")
got, err := runLocalPathForCanonical(layout, artifacts.SessionPaths{Root: sessionRoot}, canonical)
if err != nil {
t.Fatalf("runLocalPathForCanonical() error = %v", err)
}
if got != canonical {
t.Fatalf("runLocalPathForCanonical() = %q, want canonical %q", got, canonical)
}
}