Rename publish runtime terminology to published outputs

This commit is contained in:
2026-05-23 04:42:08 +00:00
parent df2c765b7f
commit 79737edf79
32 changed files with 354 additions and 354 deletions

View File

@@ -496,13 +496,13 @@ func executeAnalyzeArtifact(
if err := requireNonEmptyFile(finalOutputPath, artifactName+" output"); err != nil {
return nil, fmt.Errorf("analyze: %w", err)
}
promotedArtifact, err := promoteRunLocalOutput(env.ArtifactStore, finalOutputPath, canonicalOutputPath, artifacts.Ref{
materializedArtifact, err := materializeRunLocalOutput(env.ArtifactStore, finalOutputPath, canonicalOutputPath, artifacts.Ref{
Kind: artifactName,
Category: "artifacts",
SessionID: sessionID,
})
if err != nil {
return nil, fmt.Errorf("analyze: promote artifact output for %q: %w", artifactName, err)
return nil, fmt.Errorf("analyze: materialize artifact output for %q: %w", artifactName, err)
}
logPaths = append(logPaths, stdoutLogPath, stderrLogPath)
@@ -529,7 +529,7 @@ func executeAnalyzeArtifact(
}
return &analyzeArtifactExecutionResult{
Output: promotedArtifact,
Output: materializedArtifact,
Logs: logPaths,
GeneratedConfigs: generatedConfigs,
Metadata: meta,

View File

@@ -277,7 +277,7 @@ func TestAnalyzeUsesRunLocalPathsAndPromotesCanonical(t *testing.T) {
t.Fatalf("outputs len = %d, want 1", len(result.Outputs))
}
if strings.Contains(result.Outputs[0].AbsolutePath, string(filepath.Separator)+"runs"+string(filepath.Separator)) {
t.Fatalf("promoted output path = %q, want canonical session path", result.Outputs[0].AbsolutePath)
t.Fatalf("materialized output path = %q, want canonical session path", result.Outputs[0].AbsolutePath)
}
}

View File

@@ -72,7 +72,7 @@ func (archiveStage) Run(ctx context.Context, env *Env, m *manifest.Manifest) (*S
}, nil
}
if err := validateArchivePrerequisites(m); err != nil {
if err := validatePublishPrerequisites(m); err != nil {
return nil, fmt.Errorf("publish: %w", err)
}
if env.ObjectStore == nil {
@@ -126,17 +126,17 @@ func (archiveStage) Run(ctx context.Context, env *Env, m *manifest.Manifest) (*S
if err != nil {
return nil, fmt.Errorf("publish: build runtime artifact catalog: %w", err)
}
promotions, skippedOptional, skippedUnselected, lockedPromotions, err := resolveArchivePromotions(
publishOutputs, skippedOptionalOutputs, skippedUnselectedOutputs, lockedOutputs, err := resolvePublishOutputs(
sessionPaths,
m,
runtimeCatalog,
env.Config.Pipeline.Publish.PromoteArtifacts,
env.Config.Pipeline.Publish.Outputs,
env.Config.Pipeline.Publish.Locks,
env.SelectedArtifactKeys,
sessionPrefix,
)
if err != nil {
return nil, fmt.Errorf("publish: resolve promotion rules: %w", err)
return nil, fmt.Errorf("publish: resolve publish output rules: %w", err)
}
runUploaded := make([]string, 0, len(runFiles))
for _, file := range runFiles {
@@ -147,18 +147,18 @@ func (archiveStage) Run(ctx context.Context, env *Env, m *manifest.Manifest) (*S
runUploaded = append(runUploaded, file.RelativePath)
}
promotedUploaded := make([]string, 0, len(promotions))
for _, promotion := range promotions {
key := artifacts.S3PromotedArtifactKey(sessionPrefix, promotion.Dest)
publishedUploaded := make([]string, 0, len(publishOutputs))
for _, promotion := range publishOutputs {
key := artifacts.S3PublishedOutputKey(sessionPrefix, promotion.Dest)
if _, err := env.ObjectStore.Upload(ctx, promotion.LocalPath, key, storage.UploadOptions{}); err != nil {
return nil, fmt.Errorf("publish: upload promoted output source %q to %q: %w", promotion.Source, key, err)
return nil, fmt.Errorf("publish: upload published output source %q to %q: %w", promotion.Source, key, err)
}
promotedUploaded = append(promotedUploaded, promotion.Dest)
publishedUploaded = append(publishedUploaded, promotion.Dest)
}
previousUploaded := make([]string, 0, len(previousFiles))
for _, file := range previousFiles {
key := artifacts.S3PromotedArtifactKey(sessionPrefix, file.RelativePath)
key := artifacts.S3PublishedOutputKey(sessionPrefix, file.RelativePath)
if _, err := env.ObjectStore.Upload(ctx, file.LocalPath, key, storage.UploadOptions{}); err != nil {
return nil, fmt.Errorf("publish: upload previous file %q to %q: %w", file.RelativePath, key, err)
}
@@ -171,11 +171,11 @@ func (archiveStage) Run(ctx context.Context, env *Env, m *manifest.Manifest) (*S
runPrefix,
sessionPrefix,
runUploaded,
promotedUploaded,
publishedUploaded,
previousUploaded,
skippedOptional,
skippedUnselected,
lockedPromotions,
skippedOptionalOutputs,
skippedUnselectedOutputs,
lockedOutputs,
currentManifestKey,
))
if err != nil {
@@ -203,29 +203,29 @@ func (archiveStage) Run(ctx context.Context, env *Env, m *manifest.Manifest) (*S
return &StageResult{
Metadata: map[string]any{
"stage": "publish",
"uploaded": true,
"s3_bucket": bucket,
"s3_run_prefix": runPrefix,
"run_files_uploaded": len(runUploaded),
"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,
"skipped_unselected_promotions": skippedUnselectedPromotionMetadata(skippedUnselected),
"locked_promotion_count": len(lockedPromotions),
"locked_promotions": lockedPromotionMetadata(lockedPromotions),
"current_manifest_key": currentManifestKey,
"current_run_id_key": currentRunPointerKey,
"current_pointer_written": true,
"audio_upload_skipped": true,
"stage": "publish",
"uploaded": true,
"s3_bucket": bucket,
"s3_run_prefix": runPrefix,
"run_files_uploaded": len(runUploaded),
"run_uploaded_paths": runUploaded,
"published_files_uploaded": len(publishedUploaded),
"published_paths": publishedUploaded,
"previous_files_uploaded": len(previousUploaded),
"previous_uploaded_paths": previousUploaded,
"skipped_optional_outputs": skippedOptionalOutputs,
"skipped_unselected_outputs": skippedUnselectedOutputMetadata(skippedUnselectedOutputs),
"locked_output_count": len(lockedOutputs),
"locked_outputs": lockedOutputMetadata(lockedOutputs),
"current_manifest_key": currentManifestKey,
"current_run_id_key": currentRunPointerKey,
"current_pointer_written": true,
"audio_upload_skipped": true,
},
}, nil
}
type archivePromotion struct {
type publishOutput struct {
Source string
Dest string
Required bool
@@ -233,7 +233,7 @@ type archivePromotion struct {
Provenance string
}
type archiveLockedPromotion struct {
type publishLockedOutput struct {
Source string
Dest string
RemoteKey string
@@ -243,7 +243,7 @@ type archiveLockedPromotion struct {
Provenance string
}
type archiveSkippedUnselectedPromotion struct {
type publishSkippedUnselectedOutput struct {
Source string
Dest string
Required bool
@@ -265,7 +265,7 @@ func archiveRunUploadDisabled(env *Env) bool {
return cfg.UploadRun != nil && !*cfg.UploadRun
}
func validateArchivePrerequisites(m *manifest.Manifest) error {
func validatePublishPrerequisites(m *manifest.Manifest) error {
if m == nil {
return fmt.Errorf("manifest is required")
}
@@ -343,32 +343,32 @@ func archiveSessionPaths(env *Env, m *manifest.Manifest) artifacts.SessionPaths
return store.SessionPathsFor(campaign, sessionID)
}
func resolveArchivePromotions(
func resolvePublishOutputs(
paths artifacts.SessionPaths,
m *manifest.Manifest,
catalog *artifacts.ArtifactCatalog,
rules []config.ArchivePromotionRule,
locks []config.ArchiveLockRule,
rules []config.PublishOutputRule,
locks []config.PublishLockRule,
selectedArtifactKeys []string,
sessionPrefix string,
) ([]archivePromotion, []string, []archiveSkippedUnselectedPromotion, []archiveLockedPromotion, error) {
out := make([]archivePromotion, 0, len(rules))
skippedOptional := make([]string, 0)
skippedUnselected := make([]archiveSkippedUnselectedPromotion, 0)
lockedPromotions := make([]archiveLockedPromotion, 0)
) ([]publishOutput, []string, []publishSkippedUnselectedOutput, []publishLockedOutput, error) {
out := make([]publishOutput, 0, len(rules))
skippedOptionalOutputs := make([]string, 0)
skippedUnselectedOutputs := make([]publishSkippedUnselectedOutput, 0)
lockedOutputs := make([]publishLockedOutput, 0)
lockSet := archiveLockSet(locks)
selectedSet := archiveSelectedArtifactSet(selectedArtifactKeys)
for _, rule := range rules {
source := strings.TrimSpace(rule.Source)
required := rule.Required == nil || *rule.Required
dest, err := resolveArchivePromotionDest(rule, catalog)
dest, err := resolvePublishOutputDest(rule, catalog)
if err != nil {
return nil, nil, nil, nil, fmt.Errorf("source %q: %w", source, err)
}
if len(selectedSet) > 0 {
if key, ok := artifacts.ConfiguredArtifactName(source); ok {
if _, selected := selectedSet[key]; !selected {
skippedUnselected = append(skippedUnselected, archiveSkippedUnselectedPromotion{
skippedUnselectedOutputs = append(skippedUnselectedOutputs, publishSkippedUnselectedOutput{
Source: source,
Dest: dest,
Required: required,
@@ -381,29 +381,29 @@ func resolveArchivePromotions(
resolved, err := artifacts.ResolveSessionArtifactWithCatalog(paths, m, source, catalog)
if err != nil {
if locked {
lockedPromotions = append(lockedPromotions, archiveLockedPromotion{
lockedOutputs = append(lockedOutputs, publishLockedOutput{
Source: source,
Dest: dest,
RemoteKey: artifacts.S3PromotedArtifactKey(sessionPrefix, dest),
RemoteKey: artifacts.S3PublishedOutputKey(sessionPrefix, dest),
Reason: strings.TrimSpace(lock.Reason),
Required: required,
})
continue
}
if errors.Is(err, artifacts.ErrSessionArtifactNotFound) && !required {
skippedOptional = append(skippedOptional, dest)
skippedOptionalOutputs = append(skippedOptionalOutputs, dest)
continue
}
if errors.Is(err, artifacts.ErrSessionArtifactNotFound) {
return nil, nil, nil, nil, fmt.Errorf("required promotion source unavailable: %q", source)
return nil, nil, nil, nil, fmt.Errorf("required output source unavailable: %q", source)
}
return nil, nil, nil, nil, fmt.Errorf("resolve source %q: %w", source, err)
}
if locked {
lockedPromotions = append(lockedPromotions, archiveLockedPromotion{
lockedOutputs = append(lockedOutputs, publishLockedOutput{
Source: source,
Dest: dest,
RemoteKey: artifacts.S3PromotedArtifactKey(sessionPrefix, dest),
RemoteKey: artifacts.S3PublishedOutputKey(sessionPrefix, dest),
Reason: strings.TrimSpace(lock.Reason),
Required: required,
LocalPath: resolved.Path,
@@ -411,7 +411,7 @@ func resolveArchivePromotions(
})
continue
}
out = append(out, archivePromotion{
out = append(out, publishOutput{
Source: source,
Dest: dest,
Required: required,
@@ -419,7 +419,7 @@ func resolveArchivePromotions(
Provenance: resolved.Provenance,
})
}
return out, skippedOptional, skippedUnselected, lockedPromotions, nil
return out, skippedOptionalOutputs, skippedUnselectedOutputs, lockedOutputs, nil
}
func archiveSelectedArtifactSet(selected []string) map[string]struct{} {
@@ -437,8 +437,8 @@ func archiveSelectedArtifactSet(selected []string) map[string]struct{} {
return out
}
func archiveLockSet(locks []config.ArchiveLockRule) map[string]config.ArchiveLockRule {
out := make(map[string]config.ArchiveLockRule, len(locks))
func archiveLockSet(locks []config.PublishLockRule) map[string]config.PublishLockRule {
out := make(map[string]config.PublishLockRule, len(locks))
for _, lock := range locks {
source := strings.TrimSpace(lock.Source)
if source == "" {
@@ -451,7 +451,7 @@ func archiveLockSet(locks []config.ArchiveLockRule) map[string]config.ArchiveLoc
return out
}
func resolveArchivePromotionDest(rule config.ArchivePromotionRule, catalog *artifacts.ArtifactCatalog) (string, error) {
func resolvePublishOutputDest(rule config.PublishOutputRule, catalog *artifacts.ArtifactCatalog) (string, error) {
dest := strings.TrimSpace(rule.Dest)
if dest == "" {
entry, ok := catalog.Lookup(strings.TrimSpace(rule.Source))
@@ -760,36 +760,36 @@ func writeCurrentRunIDPointer(runID string) (string, error) {
func archiveMetadataPreview(
bucket, runPrefix, sessionPrefix string,
runUploaded []string,
promotedUploaded []string,
publishedUploaded []string,
previousUploaded []string,
skippedOptional []string,
skippedUnselected []archiveSkippedUnselectedPromotion,
lockedPromotions []archiveLockedPromotion,
skippedOptionalOutputs []string,
skippedUnselectedOutputs []publishSkippedUnselectedOutput,
lockedOutputs []publishLockedOutput,
currentManifestKey string,
) map[string]any {
return map[string]any{
"stage": "publish",
"uploaded": true,
"s3_bucket": bucket,
"s3_run_prefix": runPrefix,
"run_files_uploaded": len(runUploaded),
"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...),
"skipped_unselected_promotions": skippedUnselectedPromotionMetadata(skippedUnselected),
"locked_promotion_count": len(lockedPromotions),
"locked_promotions": lockedPromotionMetadata(lockedPromotions),
"current_manifest_key": currentManifestKey,
"current_run_id_key": artifacts.S3CurrentRunPointerKey(sessionPrefix),
"current_pointer_written": false,
"audio_upload_skipped": true,
"stage": "publish",
"uploaded": true,
"s3_bucket": bucket,
"s3_run_prefix": runPrefix,
"run_files_uploaded": len(runUploaded),
"run_uploaded_paths": append([]string(nil), runUploaded...),
"published_files_uploaded": len(publishedUploaded),
"published_paths": append([]string(nil), publishedUploaded...),
"previous_files_uploaded": len(previousUploaded),
"previous_uploaded_paths": append([]string(nil), previousUploaded...),
"skipped_optional_outputs": append([]string(nil), skippedOptionalOutputs...),
"skipped_unselected_outputs": skippedUnselectedOutputMetadata(skippedUnselectedOutputs),
"locked_output_count": len(lockedOutputs),
"locked_outputs": lockedOutputMetadata(lockedOutputs),
"current_manifest_key": currentManifestKey,
"current_run_id_key": artifacts.S3CurrentRunPointerKey(sessionPrefix),
"current_pointer_written": false,
"audio_upload_skipped": true,
}
}
func skippedUnselectedPromotionMetadata(skipped []archiveSkippedUnselectedPromotion) []map[string]any {
func skippedUnselectedOutputMetadata(skipped []publishSkippedUnselectedOutput) []map[string]any {
out := make([]map[string]any, 0, len(skipped))
for _, item := range skipped {
out = append(out, map[string]any{
@@ -801,7 +801,7 @@ func skippedUnselectedPromotionMetadata(skipped []archiveSkippedUnselectedPromot
return out
}
func lockedPromotionMetadata(locked []archiveLockedPromotion) []map[string]any {
func lockedOutputMetadata(locked []publishLockedOutput) []map[string]any {
out := make([]map[string]any, 0, len(locked))
for _, item := range locked {
out = append(out, map[string]any{

View File

@@ -128,8 +128,8 @@ func TestArchiveUploadsRunRecordPromotionsAndCurrentPointer(t *testing.T) {
if result.Metadata["current_pointer_written"] != true {
t.Fatalf("metadata = %#v, want current_pointer_written=true", result.Metadata)
}
if result.Metadata["promoted_files_uploaded"] != 2 {
t.Fatalf("metadata promoted_files_uploaded = %#v, want 2", result.Metadata["promoted_files_uploaded"])
if result.Metadata["published_files_uploaded"] != 2 {
t.Fatalf("metadata published_files_uploaded = %#v, want 2", result.Metadata["published_files_uploaded"])
}
if result.Metadata["previous_files_uploaded"] != 0 {
t.Fatalf("metadata previous_files_uploaded = %#v, want 0", result.Metadata["previous_files_uploaded"])
@@ -179,7 +179,7 @@ func TestArchiveToleratesMissingPreviousCache(t *testing.T) {
func TestArchiveUsesCustomPromotionRules(t *testing.T) {
env, m, _ := archiveFixture(t)
env.Config.Pipeline.Publish.PromoteArtifacts = []config.ArchivePromotionRule{
env.Config.Pipeline.Publish.Outputs = []config.PublishOutputRule{
{Source: "narratio.transcript.final_trimmed", Dest: "published/trimmed.json", Required: boolPtr(true)},
{Source: "narratio.artifact.session_recap", Dest: "published/recap.md", Required: boolPtr(true)},
}
@@ -200,7 +200,7 @@ func TestArchiveUsesCustomPromotionRules(t *testing.T) {
func TestArchiveSkipsOptionalMissingPromotion(t *testing.T) {
env, m, _ := archiveFixture(t)
env.Config.Pipeline.Publish.PromoteArtifacts = []config.ArchivePromotionRule{
env.Config.Pipeline.Publish.Outputs = []config.PublishOutputRule{
{Source: "narratio.transcript.final_trimmed", Dest: "transcripts/final.trimmed.json", Required: boolPtr(true)},
{Source: "narratio.transcript.base", Dest: "transcripts/base.json", Required: boolPtr(false)},
}
@@ -209,10 +209,10 @@ func TestArchiveSkipsOptionalMissingPromotion(t *testing.T) {
if err != nil {
t.Fatalf("Run() error = %v", err)
}
got, _ := result.Metadata["skipped_optional_promotions"].([]string)
got, _ := result.Metadata["skipped_optional_outputs"].([]string)
want := []string{"transcripts/base.json"}
if !reflect.DeepEqual(got, want) {
t.Fatalf("skipped_optional_promotions = %#v, want %#v", got, want)
t.Fatalf("skipped_optional_outputs = %#v, want %#v", got, want)
}
}
@@ -236,15 +236,15 @@ func TestArchiveSkipsRequiredUnselectedConfiguredPromotion(t *testing.T) {
if _, ok := fake.Objects[m.S3SessionPrefix+"artifacts/session_recap.md"]; ok {
t.Fatalf("unexpected unselected recap promotion upload")
}
skipped := result.Metadata["skipped_unselected_promotions"].([]map[string]any)
skipped := result.Metadata["skipped_unselected_outputs"].([]map[string]any)
if len(skipped) != 1 {
t.Fatalf("skipped_unselected_promotions = %#v, want one item", skipped)
t.Fatalf("skipped_unselected_outputs = %#v, want one item", skipped)
}
if skipped[0]["source"] != "narratio.artifact.session_recap" || skipped[0]["dest"] != "artifacts/session_recap.md" || skipped[0]["required"] != true {
t.Fatalf("skipped_unselected_promotions[0] = %#v, want session recap", skipped[0])
t.Fatalf("skipped_unselected_outputs[0] = %#v, want session recap", skipped[0])
}
if result.Metadata["locked_promotion_count"] != 0 {
t.Fatalf("locked_promotion_count = %#v, want 0", result.Metadata["locked_promotion_count"])
if result.Metadata["locked_output_count"] != 0 {
t.Fatalf("locked_output_count = %#v, want 0", result.Metadata["locked_output_count"])
}
if _, ok := fake.Objects[m.S3SessionPrefix+"current/run_id.txt"]; !ok {
t.Fatalf("missing current pointer")
@@ -264,15 +264,15 @@ func TestArchiveSelectedConfiguredPromotionStillFailsWhenMissing(t *testing.T) {
}
_, err := archiveStage{}.Run(context.Background(), env, m)
if err == nil || !strings.Contains(err.Error(), `required promotion source unavailable: "narratio.artifact.session_recap"`) {
t.Fatalf("Run() error = %v, want required selected promotion failure", err)
if err == nil || !strings.Contains(err.Error(), `required output source unavailable: "narratio.artifact.session_recap"`) {
t.Fatalf("Run() error = %v, want required selected output failure", err)
}
}
func TestArchiveLockedSelectedPromotionSkipsAsLocked(t *testing.T) {
env, m, _ := archiveFixture(t)
env.SelectedArtifactKeys = []string{"session_recap"}
env.Config.Pipeline.Publish.Locks = []config.ArchiveLockRule{
env.Config.Pipeline.Publish.Locks = []config.PublishLockRule{
{Source: "narratio.artifact.session_recap", Reason: "reviewed"},
}
fake := env.ObjectStore.(*storage.FakeBackend)
@@ -284,12 +284,12 @@ func TestArchiveLockedSelectedPromotionSkipsAsLocked(t *testing.T) {
if _, ok := fake.Objects[m.S3SessionPrefix+"artifacts/session_recap.md"]; ok {
t.Fatalf("unexpected locked recap promotion upload")
}
if result.Metadata["locked_promotion_count"] != 1 {
t.Fatalf("locked_promotion_count = %#v, want 1", result.Metadata["locked_promotion_count"])
if result.Metadata["locked_output_count"] != 1 {
t.Fatalf("locked_output_count = %#v, want 1", result.Metadata["locked_output_count"])
}
skipped := result.Metadata["skipped_unselected_promotions"].([]map[string]any)
skipped := result.Metadata["skipped_unselected_outputs"].([]map[string]any)
if len(skipped) != 0 {
t.Fatalf("skipped_unselected_promotions = %#v, want empty", skipped)
t.Fatalf("skipped_unselected_outputs = %#v, want empty", skipped)
}
}
@@ -301,7 +301,7 @@ func TestArchiveLockedUnselectedConfiguredPromotionSkipsAsUnselected(t *testing.
PromptID: "dnd.player_handout",
OutputPath: "artifacts/player_handout.md",
}
env.Config.Pipeline.Publish.Locks = []config.ArchiveLockRule{
env.Config.Pipeline.Publish.Locks = []config.PublishLockRule{
{Source: "narratio.artifact.session_recap", Reason: "reviewed"},
}
@@ -309,18 +309,18 @@ func TestArchiveLockedUnselectedConfiguredPromotionSkipsAsUnselected(t *testing.
if err != nil {
t.Fatalf("Run() error = %v", err)
}
if result.Metadata["locked_promotion_count"] != 0 {
t.Fatalf("locked_promotion_count = %#v, want 0", result.Metadata["locked_promotion_count"])
if result.Metadata["locked_output_count"] != 0 {
t.Fatalf("locked_output_count = %#v, want 0", result.Metadata["locked_output_count"])
}
skipped := result.Metadata["skipped_unselected_promotions"].([]map[string]any)
skipped := result.Metadata["skipped_unselected_outputs"].([]map[string]any)
if len(skipped) != 1 || skipped[0]["source"] != "narratio.artifact.session_recap" {
t.Fatalf("skipped_unselected_promotions = %#v, want unselected recap", skipped)
t.Fatalf("skipped_unselected_outputs = %#v, want unselected recap", skipped)
}
}
func TestArchiveSkipsLockedRequiredPromotionAndCommits(t *testing.T) {
env, m, _ := archiveFixture(t)
env.Config.Pipeline.Publish.Locks = []config.ArchiveLockRule{
env.Config.Pipeline.Publish.Locks = []config.PublishLockRule{
{Source: "narratio.transcript.final_trimmed", Reason: "human reviewed"},
}
fake := env.ObjectStore.(*storage.FakeBackend)
@@ -348,15 +348,15 @@ func TestArchiveSkipsLockedRequiredPromotionAndCommits(t *testing.T) {
t.Fatalf("last upload = %#v, want current run pointer %q", fake.Uploads, currentRunIDKey)
}
if result.Metadata["promoted_files_uploaded"] != 1 {
t.Fatalf("metadata promoted_files_uploaded = %#v, want 1", result.Metadata["promoted_files_uploaded"])
if result.Metadata["published_files_uploaded"] != 1 {
t.Fatalf("metadata published_files_uploaded = %#v, want 1", result.Metadata["published_files_uploaded"])
}
if result.Metadata["locked_promotion_count"] != 1 {
t.Fatalf("metadata locked_promotion_count = %#v, want 1", result.Metadata["locked_promotion_count"])
if result.Metadata["locked_output_count"] != 1 {
t.Fatalf("metadata locked_output_count = %#v, want 1", result.Metadata["locked_output_count"])
}
locked := result.Metadata["locked_promotions"].([]map[string]any)
locked := result.Metadata["locked_outputs"].([]map[string]any)
if len(locked) != 1 {
t.Fatalf("locked_promotions = %#v, want one item", locked)
t.Fatalf("locked_outputs = %#v, want one item", locked)
}
if locked[0]["source"] != "narratio.transcript.final_trimmed" ||
locked[0]["dest"] != "transcripts/final.trimmed.json" ||
@@ -376,21 +376,21 @@ func TestArchiveSkipsLockedRequiredPromotionAndCommits(t *testing.T) {
stages := current["stages"].(map[string]any)
archive := stages["publish"].(map[string]any)
meta := archive["metadata"].(map[string]any)
if meta["locked_promotion_count"] != float64(1) {
t.Fatalf("current manifest locked_promotion_count = %#v, want 1", meta["locked_promotion_count"])
if meta["locked_output_count"] != float64(1) {
t.Fatalf("current manifest locked_output_count = %#v, want 1", meta["locked_output_count"])
}
items := meta["locked_promotions"].([]any)
items := meta["locked_outputs"].([]any)
if len(items) != 1 {
t.Fatalf("current manifest locked_promotions = %#v, want one item", items)
t.Fatalf("current manifest locked_outputs = %#v, want one item", items)
}
}
func TestArchiveLockedRequiredMissingPromotionSucceeds(t *testing.T) {
env, m, _ := archiveFixture(t)
env.Config.Pipeline.Publish.PromoteArtifacts = []config.ArchivePromotionRule{
env.Config.Pipeline.Publish.Outputs = []config.PublishOutputRule{
{Source: "narratio.transcript.base", Dest: "transcripts/base.json", Required: boolPtr(true)},
}
env.Config.Pipeline.Publish.Locks = []config.ArchiveLockRule{
env.Config.Pipeline.Publish.Locks = []config.PublishLockRule{
{Source: "narratio.transcript.base", Reason: "manual merge is locked"},
}
@@ -407,9 +407,9 @@ func TestArchiveLockedRequiredMissingPromotionSucceeds(t *testing.T) {
if _, ok := fake.Objects[m.S3SessionPrefix+"current/run_id.txt"]; !ok {
t.Fatalf("current run pointer should be written for locked missing promotion")
}
locked := result.Metadata["locked_promotions"].([]map[string]any)
locked := result.Metadata["locked_outputs"].([]map[string]any)
if len(locked) != 1 {
t.Fatalf("locked_promotions = %#v, want one item", locked)
t.Fatalf("locked_outputs = %#v, want one item", locked)
}
if locked[0]["local_path"] != "" || locked[0]["provenance"] != "" {
t.Fatalf("locked missing promotion metadata = %#v, want empty local path/provenance", locked[0])
@@ -418,7 +418,7 @@ func TestArchiveLockedRequiredMissingPromotionSucceeds(t *testing.T) {
func TestArchiveLockDoesNotOverwriteExistingPromotion(t *testing.T) {
env, m, _ := archiveFixture(t)
env.Config.Pipeline.Publish.Locks = []config.ArchiveLockRule{
env.Config.Pipeline.Publish.Locks = []config.PublishLockRule{
{Source: "narratio.transcript.final_trimmed", Reason: "already published"},
}
fake := env.ObjectStore.(*storage.FakeBackend)
@@ -442,13 +442,13 @@ func TestArchiveLockDoesNotOverwriteExistingPromotion(t *testing.T) {
func TestArchiveFailsWhenRequiredPromotionMissing(t *testing.T) {
env, m, _ := archiveFixture(t)
env.Config.Pipeline.Publish.PromoteArtifacts = []config.ArchivePromotionRule{
env.Config.Pipeline.Publish.Outputs = []config.PublishOutputRule{
{Source: "narratio.transcript.base", Dest: "transcripts/base.json", Required: boolPtr(true)},
}
_, err := archiveStage{}.Run(context.Background(), env, m)
if err == nil || !strings.Contains(err.Error(), "required promotion source unavailable") {
t.Fatalf("Run() error = %v, want required promotion source unavailable failure", err)
if err == nil || !strings.Contains(err.Error(), "required output source unavailable") {
t.Fatalf("Run() error = %v, want required output source unavailable failure", err)
}
}
@@ -485,8 +485,8 @@ func TestArchiveDoesNotWriteCurrentPointerWhenPromotionUploadFails(t *testing.T)
env.ObjectStore = &promotionFailingStore{delegate: fake, failKey: failingKey}
_, err := archiveStage{}.Run(context.Background(), env, m)
if err == nil || !strings.Contains(err.Error(), "promoted output") {
t.Fatalf("Run() error = %v, want promotion upload failure", err)
if err == nil || !strings.Contains(err.Error(), "upload published output source") {
t.Fatalf("Run() error = %v, want published output upload failure", err)
}
if _, ok := fake.Objects[m.S3SessionPrefix+"current/run_id.txt"]; ok {
t.Fatalf("unexpected current pointer write on promotion failure")
@@ -563,10 +563,10 @@ func archiveFixture(t *testing.T) (*Env, *manifest.Manifest, string) {
RootPrefix: "dnd",
},
},
Publish: &config.ArchiveConfig{
Publish: &config.PublishConfig{
Enabled: boolPtr(true),
UploadRun: boolPtr(true),
PromoteArtifacts: []config.ArchivePromotionRule{
Outputs: []config.PublishOutputRule{
{Source: "narratio.transcript.final_trimmed", Dest: "transcripts/final.trimmed.json", Required: boolPtr(true)},
{Source: "narratio.artifact.session_recap", Dest: "artifacts/session_recap.md", Required: boolPtr(true)},
},

View File

@@ -150,7 +150,7 @@ func (mergeStage) Run(ctx context.Context, env *Env, m *manifest.Manifest) (*Sta
}
}
promotedMerged, err := promoteRunLocalOutput(env.ArtifactStore, finalMergedPath, canonicalMergedPath, artifacts.Ref{
materializedMerged, err := materializeRunLocalOutput(env.ArtifactStore, finalMergedPath, canonicalMergedPath, artifacts.Ref{
Kind: "transcript_base",
Category: "transcripts",
SessionID: sessionID,
@@ -158,9 +158,9 @@ func (mergeStage) Run(ctx context.Context, env *Env, m *manifest.Manifest) (*Sta
if err != nil {
return nil, fmt.Errorf("merge: promote merged transcript: %w", err)
}
outputs := []artifacts.Ref{promotedMerged}
outputs := []artifacts.Ref{materializedMerged}
if reportEnabled {
promotedReport, err := promoteRunLocalOutput(env.ArtifactStore, finalReportPath, canonicalReportPath, artifacts.Ref{
materializedReport, err := materializeRunLocalOutput(env.ArtifactStore, finalReportPath, canonicalReportPath, artifacts.Ref{
Kind: "seriatim_report",
Category: "artifacts",
SessionID: sessionID,
@@ -168,7 +168,7 @@ func (mergeStage) Run(ctx context.Context, env *Env, m *manifest.Manifest) (*Sta
if err != nil {
return nil, fmt.Errorf("merge: promote report: %w", err)
}
outputs = append(outputs, promotedReport)
outputs = append(outputs, materializedReport)
}
coalesceGap := any(nil)

View File

@@ -315,10 +315,10 @@ func TestMergeStageUsesRunLocalPathsAndPromotesCanonical(t *testing.T) {
t.Fatalf("run output path = %q, want run-local path", req.OutputMergedTranscriptPath)
}
if len(result.Outputs) == 0 {
t.Fatalf("outputs = %#v, want promoted outputs", result.Outputs)
t.Fatalf("outputs = %#v, want materialized outputs", result.Outputs)
}
if strings.Contains(result.Outputs[0].AbsolutePath, string(filepath.Separator)+"runs"+string(filepath.Separator)) {
t.Fatalf("promoted output path = %q, want canonical session path", result.Outputs[0].AbsolutePath)
t.Fatalf("materialized output path = %q, want canonical session path", result.Outputs[0].AbsolutePath)
}
}

View File

@@ -133,7 +133,7 @@ func (normalizeStage) Run(ctx context.Context, env *Env, m *manifest.Manifest) (
}
}
promotedNormalized, err := promoteRunLocalOutput(env.ArtifactStore, finalNormalizedPath, canonicalNormalizedPath, artifacts.Ref{
materializedNormalized, err := materializeRunLocalOutput(env.ArtifactStore, finalNormalizedPath, canonicalNormalizedPath, artifacts.Ref{
Kind: "transcript_final",
Category: "transcripts",
SessionID: sessionID,
@@ -141,9 +141,9 @@ func (normalizeStage) Run(ctx context.Context, env *Env, m *manifest.Manifest) (
if err != nil {
return nil, fmt.Errorf("normalize: promote normalized transcript: %w", err)
}
outputs := []artifacts.Ref{promotedNormalized}
outputs := []artifacts.Ref{materializedNormalized}
if reportEnabled {
promotedReport, err := promoteRunLocalOutput(env.ArtifactStore, finalReportPath, canonicalReportPath, artifacts.Ref{
materializedReport, err := materializeRunLocalOutput(env.ArtifactStore, finalReportPath, canonicalReportPath, artifacts.Ref{
Kind: "seriatim_normalize_report",
Category: "artifacts",
SessionID: sessionID,
@@ -151,7 +151,7 @@ func (normalizeStage) Run(ctx context.Context, env *Env, m *manifest.Manifest) (
if err != nil {
return nil, fmt.Errorf("normalize: promote report: %w", err)
}
outputs = append(outputs, promotedReport)
outputs = append(outputs, materializedReport)
}
reportCanonicalPath := ""
if reportEnabled {

View File

@@ -231,10 +231,10 @@ func TestNormalizeStageUsesRunLocalPathsAndPromotesCanonical(t *testing.T) {
t.Fatalf("run output path = %q, want run-local path", req.OutputNormalizedPath)
}
if len(result.Outputs) == 0 {
t.Fatalf("outputs = %#v, want promoted outputs", result.Outputs)
t.Fatalf("outputs = %#v, want materialized outputs", result.Outputs)
}
if strings.Contains(result.Outputs[0].AbsolutePath, string(filepath.Separator)+"runs"+string(filepath.Separator)) {
t.Fatalf("promoted output path = %q, want canonical session path", result.Outputs[0].AbsolutePath)
t.Fatalf("materialized output path = %q, want canonical session path", result.Outputs[0].AbsolutePath)
}
}

View File

@@ -61,7 +61,7 @@ func TestStagesReturnExpectedMetadata(t *testing.T) {
RootPrefix: "dnd",
},
},
Publish: &config.ArchiveConfig{
Publish: &config.PublishConfig{
Enabled: boolPtr(true),
UploadRun: boolPtr(true),
},

View File

@@ -149,7 +149,7 @@ func (polishStage) Run(ctx context.Context, env *Env, m *manifest.Manifest) (*St
}
}
promotedProcessed, err := promoteRunLocalOutput(env.ArtifactStore, finalProcessedPath, canonicalProcessedPath, artifacts.Ref{
materializedProcessed, err := materializeRunLocalOutput(env.ArtifactStore, finalProcessedPath, canonicalProcessedPath, artifacts.Ref{
Kind: "transcript_polished",
Category: "transcripts",
SessionID: sessionID,
@@ -157,9 +157,9 @@ func (polishStage) Run(ctx context.Context, env *Env, m *manifest.Manifest) (*St
if err != nil {
return nil, fmt.Errorf("polish: promote processed transcript: %w", err)
}
outputs := []artifacts.Ref{promotedProcessed}
outputs := []artifacts.Ref{materializedProcessed}
if reportEnabled {
promotedReport, err := promoteRunLocalOutput(env.ArtifactStore, finalReportPath, canonicalReportPath, artifacts.Ref{
materializedReport, err := materializeRunLocalOutput(env.ArtifactStore, finalReportPath, canonicalReportPath, artifacts.Ref{
Kind: "audita_report",
Category: "artifacts",
SessionID: sessionID,
@@ -167,7 +167,7 @@ func (polishStage) Run(ctx context.Context, env *Env, m *manifest.Manifest) (*St
if err != nil {
return nil, fmt.Errorf("polish: promote report: %w", err)
}
outputs = append(outputs, promotedReport)
outputs = append(outputs, materializedReport)
}
var validationConcurrency any

View File

@@ -261,10 +261,10 @@ func TestPolishStageUsesRunLocalPathsAndPromotesCanonical(t *testing.T) {
t.Fatalf("run output path = %q, want run-local path", req.OutputProcessedPath)
}
if len(result.Outputs) == 0 {
t.Fatalf("outputs = %#v, want promoted outputs", result.Outputs)
t.Fatalf("outputs = %#v, want materialized outputs", result.Outputs)
}
if strings.Contains(result.Outputs[0].AbsolutePath, string(filepath.Separator)+"runs"+string(filepath.Separator)) {
t.Fatalf("promoted output path = %q, want canonical session path", result.Outputs[0].AbsolutePath)
t.Fatalf("materialized output path = %q, want canonical session path", result.Outputs[0].AbsolutePath)
}
}

View File

@@ -300,7 +300,7 @@ func seedPreviousCurrentState(
})
}
artifactKey := artifacts.S3PromotedArtifactKey(previousSessionPrefix, "artifacts/session_recap.md")
artifactKey := artifacts.S3PublishedOutputKey(previousSessionPrefix, "artifacts/session_recap.md")
if options.includeArtifactObject {
body := options.artifactBody
if body == "" {
@@ -365,9 +365,9 @@ func buildPreviousManifestForSeed(
LocalPath: artifactPath,
},
})
m.MarkStageSucceeded("archive", now, nil)
m.Stages["archive"].Metadata = map[string]any{
"promoted_paths": []string{"artifacts/session_recap.md"},
m.MarkStageSucceeded("publish", now, nil)
m.Stages["publish"].Metadata = map[string]any{
"published_paths": []string{"artifacts/session_recap.md"},
}
data, err := json.MarshalIndent(m, "", " ")
if err != nil {

View File

@@ -110,7 +110,7 @@ func runLocalPathForCanonical(layout runStageLayout, sessionPaths artifacts.Sess
return localPath, nil
}
func promoteRunLocalOutput(
func materializeRunLocalOutput(
store artifacts.Store,
srcPath, canonicalPath string,
ref artifacts.Ref,
@@ -128,11 +128,11 @@ func promoteRunLocalOutput(
return artifacts.Ref{}, fmt.Errorf("read run-local output %q: %w", srcPath, err)
}
if err := store.WriteFileAtomic(canonicalPath, data, 0o644); err != nil {
return artifacts.Ref{}, fmt.Errorf("promote output to %q: %w", canonicalPath, err)
return artifacts.Ref{}, fmt.Errorf("materialize output to %q: %w", canonicalPath, err)
}
checksum, err := store.Checksum(canonicalPath)
if err != nil {
return artifacts.Ref{}, fmt.Errorf("checksum promoted output %q: %w", canonicalPath, err)
return artifacts.Ref{}, fmt.Errorf("checksum materialized output %q: %w", canonicalPath, err)
}
ref.AbsolutePath = canonicalPath
ref.Checksum = checksum

View File

@@ -213,11 +213,11 @@ dispatch:
ref := outputRef[speaker]
runOutputPaths = append(runOutputPaths, ref.AbsolutePath)
canonicalOut := filepath.Join(paths.TranscriptsRawDir, speaker+".json")
promoted, err := promoteRunLocalOutput(env.ArtifactStore, ref.AbsolutePath, canonicalOut, ref)
materialized, err := materializeRunLocalOutput(env.ArtifactStore, ref.AbsolutePath, canonicalOut, ref)
if err != nil {
return nil, fmt.Errorf("transcribe: promote %q output: %w", speaker, err)
return nil, fmt.Errorf("transcribe: materialize %q output: %w", speaker, err)
}
outputs = append(outputs, promoted)
outputs = append(outputs, materialized)
outputPaths = append(outputPaths, canonicalOut)
orderedPerFile[speaker] = perFile[speaker]
}

View File

@@ -217,7 +217,7 @@ func TestTranscribeStageUsesRunLocalOutputAndPromotesCanonical(t *testing.T) {
t.Fatalf("outputs = %#v, want one output", result.Outputs)
}
if strings.Contains(result.Outputs[0].AbsolutePath, string(filepath.Separator)+"runs"+string(filepath.Separator)) {
t.Fatalf("promoted output path = %q, want canonical session path", result.Outputs[0].AbsolutePath)
t.Fatalf("materialized output path = %q, want canonical session path", result.Outputs[0].AbsolutePath)
}
}

View File

@@ -100,7 +100,7 @@ func (trimStage) Run(ctx context.Context, env *Env, m *manifest.Manifest) (*Stag
if err := validateProcessedTranscriptOutput(trimmedPath); err != nil {
return nil, fmt.Errorf("trim: copied trimmed transcript %q invalid: %w", trimmedPath, err)
}
promotedTrimmed, err := promoteRunLocalOutput(env.ArtifactStore, trimmedPath, canonicalTrimmedPath, artifacts.Ref{
materializedTrimmed, err := materializeRunLocalOutput(env.ArtifactStore, trimmedPath, canonicalTrimmedPath, artifacts.Ref{
Kind: "transcript_final_trimmed",
Category: "transcripts",
SessionID: sessionID,
@@ -110,7 +110,7 @@ func (trimStage) Run(ctx context.Context, env *Env, m *manifest.Manifest) (*Stag
}
metadata["trim_action"] = "copy_disabled"
return &StageResult{
Outputs: []artifacts.Ref{promotedTrimmed},
Outputs: []artifacts.Ref{materializedTrimmed},
Metadata: metadata,
}, nil
}
@@ -351,7 +351,7 @@ func (trimStage) Run(ctx context.Context, env *Env, m *manifest.Manifest) (*Stag
return nil, fmt.Errorf("trim: trimmed transcript %q invalid: %w", trimmedPath, err)
}
promotedTrimmed, err := promoteRunLocalOutput(env.ArtifactStore, trimmedPath, canonicalTrimmedPath, artifacts.Ref{
materializedTrimmed, err := materializeRunLocalOutput(env.ArtifactStore, trimmedPath, canonicalTrimmedPath, artifacts.Ref{
Kind: "transcript_final_trimmed",
Category: "transcripts",
SessionID: sessionID,
@@ -359,7 +359,7 @@ func (trimStage) Run(ctx context.Context, env *Env, m *manifest.Manifest) (*Stag
if err != nil {
return nil, fmt.Errorf("trim: promote trimmed transcript: %w", err)
}
promotedBounds, err := promoteRunLocalOutput(env.ArtifactStore, finalBoundsOutputPath, canonicalBoundsOutputPath, artifacts.Ref{
materializedBounds, err := materializeRunLocalOutput(env.ArtifactStore, finalBoundsOutputPath, canonicalBoundsOutputPath, artifacts.Ref{
Kind: "session_bounds",
Category: "artifacts",
SessionID: sessionID,
@@ -369,7 +369,7 @@ func (trimStage) Run(ctx context.Context, env *Env, m *manifest.Manifest) (*Stag
}
return &StageResult{
Outputs: []artifacts.Ref{promotedTrimmed, promotedBounds},
Outputs: []artifacts.Ref{materializedTrimmed, materializedBounds},
Logs: logPaths,
GeneratedConfigs: generatedConfigs,
Metadata: metadata,

View File

@@ -326,7 +326,7 @@ func TestTrimStageUsesRunLocalPathsAndPromotesCanonical(t *testing.T) {
}
for _, out := range result.Outputs {
if strings.Contains(out.AbsolutePath, string(filepath.Separator)+"runs"+string(filepath.Separator)) {
t.Fatalf("promoted output path = %q, want canonical session path", out.AbsolutePath)
t.Fatalf("materialized output path = %q, want canonical session path", out.AbsolutePath)
}
}
}