Files
narratio/internal/stage/publish.go

945 lines
30 KiB
Go

package stage
import (
"context"
"crypto/sha256"
"encoding/hex"
"encoding/json"
"errors"
"fmt"
"io"
"os"
"path/filepath"
"sort"
"strings"
"time"
"gitea.maximumdirect.net/eric/narratio/internal/adapters/storage"
"gitea.maximumdirect.net/eric/narratio/internal/artifactpolicy"
"gitea.maximumdirect.net/eric/narratio/internal/artifacts"
"gitea.maximumdirect.net/eric/narratio/internal/config"
"gitea.maximumdirect.net/eric/narratio/internal/fileops"
"gitea.maximumdirect.net/eric/narratio/internal/manifest"
"gitea.maximumdirect.net/eric/narratio/internal/pathsafe"
)
type publishStage struct{}
type publishUploadFile struct {
RelativePath string
SourceRelativePath string
Checksum string
}
var publishPrerequisiteStages = []string{
"prepare",
"transcribe",
"merge",
"polish",
"normalize",
"trim",
"render",
"analyze",
}
func (publishStage) Name() string { return "publish" }
func (publishStage) Declares() IODecl {
return IODecl{
Inputs: []artifacts.Ref{
{Kind: "manifest", Category: "input", RelativePath: "manifest.json"},
},
}
}
func (publishStage) Run(ctx context.Context, env *Env, m *manifest.Manifest) (*StageResult, error) {
if env == nil || env.Config == nil || env.Config.Pipeline == nil || env.Config.Session == nil {
return nil, fmt.Errorf("publish: resolved config must include pipeline and session")
}
if publishDisabled(env) {
return &StageResult{
Metadata: map[string]any{
"stage": "publish",
"skipped": true,
"publish_enabled": false,
"audio_upload_skipped": true,
"current_pointer_written": false,
},
}, nil
}
if publishRunUploadDisabled(env) {
return &StageResult{
Metadata: map[string]any{
"stage": "publish",
"skipped": true,
"upload_run_enabled": false,
"audio_upload_skipped": true,
"current_pointer_written": false,
},
}, nil
}
if err := validatePublishPrerequisites(m); err != nil {
return nil, fmt.Errorf("publish: %w", err)
}
if env.ObjectStore == nil {
return nil, fmt.Errorf("publish: remote object store backend is required when publish run upload is enabled")
}
runRoot, err := resolvePublishRunRoot(env, m)
if err != nil {
return nil, fmt.Errorf("publish: resolve run root: %w", err)
}
runRootInfo, err := os.Lstat(runRoot)
if err != nil {
return nil, fmt.Errorf("publish: run root %q: %w", runRoot, err)
}
if runRootInfo.Mode()&os.ModeSymlink != 0 || !runRootInfo.IsDir() {
return nil, fmt.Errorf("publish: run root %q is not a directory", runRoot)
}
runPrefix, err := artifacts.ResolvePublishRunPrefix(env.Config, m)
if err != nil {
return nil, fmt.Errorf("publish: resolve s3 run prefix: %w", err)
}
sessionPrefix, err := artifacts.ResolvePublishSessionPrefix(env.Config, m)
if err != nil {
return nil, fmt.Errorf("publish: resolve s3 session prefix: %w", err)
}
bucket := artifacts.ResolvePublishBucket(env.Config, m)
if bucket == "" {
return nil, fmt.Errorf("publish: resolve s3 bucket: bucket is required")
}
runID := strings.TrimSpace(m.RunID)
if runID == "" {
return nil, fmt.Errorf("publish: run id is required")
}
runManifestFile, err := fileops.OpenConfinedRegularFile(runRoot, "manifest.json")
if err != nil {
return nil, fmt.Errorf("publish: open run manifest: %w", err)
}
defer runManifestFile.Close()
runManifest, err := (&manifest.LocalStore{}).LoadRunReader(ctx, runManifestFile)
if err != nil {
return nil, fmt.Errorf("publish: load run manifest: %w", err)
}
if _, err := runManifestFile.Seek(0, io.SeekStart); err != nil {
return nil, fmt.Errorf("publish: rewind run manifest: %w", err)
}
runFiles, err := collectPublishRunFiles(runRoot, runManifest)
if err != nil {
return nil, fmt.Errorf("publish: collect run files: %w", err)
}
runSources, err := openPublishRunFiles(runRoot, runFiles, runManifestFile)
if err != nil {
return nil, fmt.Errorf("publish: verify run files: %w", err)
}
defer closePublishSources(runSources)
uploader, ok := env.ObjectStore.(storage.ReaderUploader)
if !ok {
return nil, fmt.Errorf("publish: object store does not support verified source uploads")
}
sessionPaths := publishSessionPaths(env, m)
previousFiles, previousManifestFile, err := collectPublishPreviousFiles(ctx, sessionPaths.PreviousDir)
if err != nil {
return nil, fmt.Errorf("publish: collect previous files: %w", err)
}
if previousManifestFile != nil {
defer previousManifestFile.Close()
}
previousSources, err := openPublishRunFiles(sessionPaths.PreviousDir, previousFiles, previousManifestFile)
if err != nil {
return nil, fmt.Errorf("publish: verify previous files: %w", err)
}
defer closePublishSources(previousSources)
runtimeCatalog, err := buildPublishRuntimeArtifactCatalog(
sessionPaths,
m,
env.Config.Pipeline.Scriptorium,
env.Config.Pipeline.Notarius,
)
if err != nil {
return nil, fmt.Errorf("publish: build runtime artifact catalog: %w", err)
}
publishOutputs, skippedOptionalOutputs, skippedUnselectedOutputs, lockedOutputs, err := resolvePublishOutputs(
sessionPaths,
m,
runtimeCatalog,
env.Config.Pipeline.Publish.Outputs,
env.Config.Pipeline.Publish.Locks,
env.SelectedArtifactKeys,
sessionPrefix,
)
if err != nil {
return nil, fmt.Errorf("publish: resolve publish output rules: %w", err)
}
runUploaded := make([]string, 0, len(runFiles))
for index, file := range runFiles {
key := artifacts.S3RunRelativeDestinationKey(runPrefix, file.RelativePath)
if _, err := uploader.UploadReader(ctx, runSources[index], key, storage.UploadOptions{}); err != nil {
return nil, fmt.Errorf("publish: upload run file %q to %q: %w", file.RelativePath, key, err)
}
runUploaded = append(runUploaded, file.RelativePath)
}
publishedUploaded := make([]string, 0, len(publishOutputs))
for _, publishedOutput := range publishOutputs {
key := artifacts.S3PublishedOutputKey(sessionPrefix, publishedOutput.Dest)
if _, err := env.ObjectStore.Upload(ctx, publishedOutput.LocalPath, key, storage.UploadOptions{}); err != nil {
return nil, fmt.Errorf("publish: upload published output source %q to %q: %w", publishedOutput.Source, key, err)
}
publishedUploaded = append(publishedUploaded, publishedOutput.Dest)
}
previousUploaded := make([]string, 0, len(previousFiles))
for index, file := range previousFiles {
key := artifacts.S3PublishedOutputKey(sessionPrefix, file.RelativePath)
if _, err := uploader.UploadReader(ctx, previousSources[index], key, storage.UploadOptions{}); err != nil {
return nil, fmt.Errorf("publish: upload previous file %q to %q: %w", file.RelativePath, key, err)
}
previousUploaded = append(previousUploaded, file.RelativePath)
}
currentManifestKey, currentRunPointerKey := artifacts.ResolveCurrentStateKeys(sessionPrefix)
manifestTempPath, err := writeCurrentManifestSnapshot(m, publishMetadataPreview(
bucket,
runPrefix,
sessionPrefix,
runUploaded,
publishedUploaded,
previousUploaded,
skippedOptionalOutputs,
skippedUnselectedOutputs,
lockedOutputs,
currentManifestKey,
))
if err != nil {
return nil, fmt.Errorf("publish: build current manifest snapshot: %w", err)
}
defer func() { _ = os.Remove(manifestTempPath) }()
if _, err := env.ObjectStore.Upload(ctx, manifestTempPath, currentManifestKey, storage.UploadOptions{
ContentType: "application/json",
}); err != nil {
return nil, fmt.Errorf("publish: upload current manifest to %q: %w", currentManifestKey, err)
}
runIDTempPath, err := writeCurrentRunIDPointer(runID)
if err != nil {
return nil, fmt.Errorf("publish: build current run id pointer: %w", err)
}
defer func() { _ = os.Remove(runIDTempPath) }()
if _, err := env.ObjectStore.Upload(ctx, runIDTempPath, currentRunPointerKey, storage.UploadOptions{
ContentType: "text/plain; charset=utf-8",
}); err != nil {
return nil, fmt.Errorf("publish: upload current run pointer to %q: %w", currentRunPointerKey, err)
}
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,
"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 publishOutput struct {
Source string
Dest string
Required bool
LocalPath string
Provenance string
}
type publishLockedOutput struct {
Source string
Dest string
RemoteKey string
Reason string
Required bool
LocalPath string
Provenance string
}
type publishSkippedUnselectedOutput struct {
Source string
Dest string
Required bool
}
func publishDisabled(env *Env) bool {
cfg := env.Config.Pipeline.Publish
if cfg == nil {
return true
}
return cfg.Enabled != nil && !*cfg.Enabled
}
func publishRunUploadDisabled(env *Env) bool {
cfg := env.Config.Pipeline.Publish
if cfg == nil {
return true
}
return cfg.UploadRun != nil && !*cfg.UploadRun
}
func validatePublishPrerequisites(m *manifest.Manifest) error {
if m == nil {
return fmt.Errorf("manifest is required")
}
for _, stageName := range publishPrerequisiteStages {
sr := m.Stages[stageName]
if sr == nil {
return fmt.Errorf("prerequisite stage %q has not succeeded", stageName)
}
if sr.Status != manifest.StatusSucceeded {
return fmt.Errorf("prerequisite stage %q status is %q (want %q)", stageName, sr.Status, manifest.StatusSucceeded)
}
}
return nil
}
func resolvePublishRunRoot(env *Env, m *manifest.Manifest) (string, error) {
sessionID := strings.TrimSpace(env.Config.Session.SessionID)
if sessionID == "" && m != nil {
sessionID = strings.TrimSpace(m.SessionID)
}
campaign := strings.TrimSpace(env.Config.Session.Campaign)
if campaign == "" && m != nil {
campaign = strings.TrimSpace(m.Campaign)
}
runID := ""
if m != nil {
runID = strings.TrimSpace(m.RunID)
}
if sessionID == "" || campaign == "" {
return "", fmt.Errorf("campaign and session id are required")
}
if runID == "" {
return "", fmt.Errorf("run id is required")
}
canonical := filepath.Clean(artifacts.SessionRunRootForCampaign(env.Config.Pipeline.Workspace.Root, campaign, sessionID, runID))
canonicalExists, err := directoryExists(canonical)
if err != nil {
return "", fmt.Errorf("check canonical run root %q: %w", canonical, err)
}
if !canonicalExists {
return "", fmt.Errorf("run root not found for campaign %q session %q run %q at canonical path %q", campaign, sessionID, runID, canonical)
}
return canonical, nil
}
func resolvePublishSessionRoot(env *Env, m *manifest.Manifest) (string, error) {
sessionID := strings.TrimSpace(env.Config.Session.SessionID)
if sessionID == "" && m != nil {
sessionID = strings.TrimSpace(m.SessionID)
}
campaign := strings.TrimSpace(env.Config.Session.Campaign)
if campaign == "" && m != nil {
campaign = strings.TrimSpace(m.Campaign)
}
if sessionID == "" {
return "", fmt.Errorf("session id is required")
}
if campaign == "" {
return "", fmt.Errorf("campaign is required")
}
return filepath.Clean(artifacts.SessionWorkDirForCampaign(env.Config.Pipeline.Workspace.Root, campaign, sessionID)), nil
}
func publishSessionPaths(env *Env, m *manifest.Manifest) artifacts.SessionPaths {
sessionID := strings.TrimSpace(env.Config.Session.SessionID)
if sessionID == "" && m != nil {
sessionID = strings.TrimSpace(m.SessionID)
}
campaign := strings.TrimSpace(env.Config.Session.Campaign)
if campaign == "" && m != nil {
campaign = strings.TrimSpace(m.Campaign)
}
store := artifacts.NewLocalStore(env.Config.Pipeline.Workspace.Root)
return store.SessionPathsFor(campaign, sessionID)
}
func resolvePublishOutputs(
paths artifacts.SessionPaths,
m *manifest.Manifest,
catalog *artifacts.ArtifactCatalog,
rules []config.PublishOutputRule,
locks []config.PublishLockRule,
selectedArtifactKeys []string,
sessionPrefix string,
) ([]publishOutput, []string, []publishSkippedUnselectedOutput, []publishLockedOutput, error) {
out := make([]publishOutput, 0, len(rules))
skippedOptionalOutputs := make([]string, 0)
skippedUnselectedOutputs := make([]publishSkippedUnselectedOutput, 0)
lockedOutputs := make([]publishLockedOutput, 0)
lockSet := publishLockSet(locks)
selectedSet := publishSelectedArtifactSet(selectedArtifactKeys)
configuredOutputs := configuredOutputPathMapFromCatalog(catalog)
extractionOutputs := extractionOutputSetFromCatalog(catalog)
for _, rule := range rules {
source := strings.TrimSpace(rule.Source)
required := rule.Required == nil || *rule.Required
dest, err := resolvePublishOutputDest(rule, configuredOutputs, extractionOutputs)
if err != nil {
return nil, nil, nil, nil, fmt.Errorf("source %q: %w", source, err)
}
if len(selectedSet) > 0 {
if key, ok := artifactpolicy.ParseConfiguredSource(source); ok {
if _, selected := selectedSet[key]; !selected {
skippedUnselectedOutputs = append(skippedUnselectedOutputs, publishSkippedUnselectedOutput{
Source: source,
Dest: dest,
Required: required,
})
continue
}
}
}
lock, locked := lockSet[source]
resolved, err := artifacts.ResolveSessionArtifactWithCatalog(paths, m, source, catalog)
if err != nil {
if locked {
lockedOutputs = append(lockedOutputs, publishLockedOutput{
Source: source,
Dest: dest,
RemoteKey: artifacts.S3PublishedOutputKey(sessionPrefix, dest),
Reason: strings.TrimSpace(lock.Reason),
Required: required,
})
continue
}
if errors.Is(err, artifacts.ErrSessionArtifactNotFound) && !required {
skippedOptionalOutputs = append(skippedOptionalOutputs, dest)
continue
}
if errors.Is(err, artifacts.ErrSessionArtifactNotFound) {
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 {
lockedOutputs = append(lockedOutputs, publishLockedOutput{
Source: source,
Dest: dest,
RemoteKey: artifacts.S3PublishedOutputKey(sessionPrefix, dest),
Reason: strings.TrimSpace(lock.Reason),
Required: required,
LocalPath: resolved.Path,
Provenance: resolved.Provenance,
})
continue
}
out = append(out, publishOutput{
Source: source,
Dest: dest,
Required: required,
LocalPath: resolved.Path,
Provenance: resolved.Provenance,
})
}
return out, skippedOptionalOutputs, skippedUnselectedOutputs, lockedOutputs, nil
}
func publishSelectedArtifactSet(selected []string) map[string]struct{} {
if len(selected) == 0 {
return nil
}
out := make(map[string]struct{}, len(selected))
for _, key := range selected {
trimmed := strings.TrimSpace(key)
if trimmed == "" {
continue
}
out[trimmed] = struct{}{}
}
return out
}
func publishLockSet(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 == "" {
continue
}
lock.Source = source
lock.Reason = strings.TrimSpace(lock.Reason)
out[source] = lock
}
return out
}
func resolvePublishOutputDest(rule config.PublishOutputRule, configured map[string]string, extractions map[string]struct{}) (string, error) {
return artifactpolicy.ResolvePublishedDestinationWithExtractions(rule.Source, rule.Dest, configured, extractions)
}
func configuredOutputPathMapFromCatalog(catalog *artifacts.ArtifactCatalog) map[string]string {
out := map[string]string{}
if catalog == nil {
return out
}
for _, entry := range catalog.ListConfigured() {
if strings.TrimSpace(entry.ConfiguredKey) == "" {
continue
}
out[entry.ConfiguredKey] = strings.TrimSpace(entry.CanonicalRelPath)
}
return out
}
func extractionOutputSetFromCatalog(catalog *artifacts.ArtifactCatalog) map[string]struct{} {
out := map[string]struct{}{}
if catalog == nil {
return out
}
for _, entry := range catalog.ListExtraction() {
if strings.TrimSpace(entry.ExtractionKey) != "" {
out[entry.ExtractionKey] = struct{}{}
}
}
return out
}
func buildPublishRuntimeArtifactCatalog(
paths artifacts.SessionPaths,
m *manifest.Manifest,
scriptoriumCfg *config.ScriptoriumConfig,
notariusCfg *config.NotariusConfig,
) (*artifacts.ArtifactCatalog, error) {
catalog := artifacts.NewArtifactCatalog()
if err := catalog.RegisterBuiltIns(); err != nil {
return nil, err
}
extractionDefinitions := artifacts.ExtractionDefinitionsFromConfig(notariusCfg)
if err := catalog.RegisterExtractionArtifacts(extractionDefinitions); err != nil {
return nil, err
}
if notariusCfg != nil && notariusCfg.Enabled {
catalog.HydrateExtractionArtifacts(paths, m, extractionDefinitions)
}
if scriptoriumCfg == nil {
return catalog, nil
}
configured := map[string]artifacts.ConfiguredArtifactDefinition{}
for key, artifactCfg := range scriptoriumCfg.Artifacts {
configured[key] = artifacts.ConfiguredArtifactDefinition{
Enabled: artifactCfg.Enabled,
OutputPath: artifactCfg.OutputPath,
}
}
if err := catalog.RegisterConfiguredArtifacts(configured, nil); err != nil {
return nil, err
}
for _, entry := range catalog.ListConfigured() {
if strings.TrimSpace(entry.CanonicalRelPath) == "" {
continue
}
localPath, err := resolveConfiguredArtifactLocalPath(paths, entry.CanonicalRelPath)
if err != nil {
continue
}
info, statErr := os.Stat(localPath)
if statErr != nil {
if os.IsNotExist(statErr) {
continue
}
return nil, fmt.Errorf("stat configured artifact %q: %w", entry.SourceID, statErr)
}
if info.IsDir() {
continue
}
if err := catalog.MarkAvailableFromDisk(entry.SourceID, localPath); err != nil {
return nil, err
}
}
return catalog, nil
}
func resolveConfiguredArtifactLocalPath(paths artifacts.SessionPaths, configured string) (string, error) {
outputPath := strings.TrimSpace(configured)
if outputPath == "" {
return "", fmt.Errorf("configured artifact output path is required")
}
if filepath.IsAbs(outputPath) {
return filepath.Clean(outputPath), nil
}
rel := filepath.Clean(outputPath)
if rel == "." || rel == "" {
return "", fmt.Errorf("relative output path is required")
}
if rel == ".." || strings.HasPrefix(rel, ".."+string(filepath.Separator)) {
return "", fmt.Errorf("relative output path escapes session root: %q", configured)
}
return filepath.Join(paths.Root, rel), nil
}
func collectPublishRunFiles(runRoot string, runManifest *manifest.RunManifest) ([]publishUploadFile, error) {
if runManifest == nil {
return nil, fmt.Errorf("run manifest is required")
}
files := map[string]publishUploadFile{
"manifest.json": {RelativePath: "manifest.json", SourceRelativePath: "manifest.json"},
}
add := func(path, checksum string) error {
if strings.TrimSpace(path) == "" {
return nil
}
rel, err := pathsafe.SlashRelativeFromRoot(runRoot, path)
if err != nil {
return nil
}
candidate := publishUploadFile{RelativePath: rel, SourceRelativePath: rel, Checksum: strings.TrimSpace(checksum)}
if prior, exists := files[rel]; exists && prior.Checksum != "" && candidate.Checksum != "" && prior.Checksum != candidate.Checksum {
return fmt.Errorf("source %q has conflicting declared checksums", rel)
}
if prior, exists := files[rel]; exists && prior.Checksum != "" {
candidate.Checksum = prior.Checksum
}
files[rel] = candidate
return nil
}
for _, record := range runManifest.Stages {
if record == nil {
continue
}
for _, output := range record.Outputs {
if err := add(output.ArchivePath, output.Checksum); err != nil {
return nil, err
}
if output.ArchivePath == "" {
if err := add(output.LocalPath, output.Checksum); err != nil {
return nil, err
}
}
}
for _, path := range append(append([]string(nil), record.Logs...), record.GeneratedConfigs...) {
if err := add(path, ""); err != nil {
return nil, err
}
}
}
out := make([]publishUploadFile, 0, len(files))
for _, file := range files {
out = append(out, file)
}
sort.Slice(out, func(i, j int) bool { return out[i].RelativePath < out[j].RelativePath })
return out, nil
}
func collectPublishPreviousFiles(ctx context.Context, previousDir string) ([]publishUploadFile, *os.File, error) {
previousDir = filepath.Clean(strings.TrimSpace(previousDir))
if previousDir == "" || previousDir == "." {
return nil, nil, fmt.Errorf("previous directory is required")
}
info, err := os.Lstat(previousDir)
if errors.Is(err, os.ErrNotExist) {
return nil, nil, nil
}
if err != nil {
return nil, nil, fmt.Errorf("inspect previous directory: %w", err)
}
if info.Mode()&os.ModeSymlink != 0 || !info.IsDir() {
return nil, nil, fmt.Errorf("previous directory is not a regular directory")
}
manifestFile, err := fileops.OpenConfinedRegularFile(previousDir, "manifest.json")
if err != nil {
if errors.Is(err, os.ErrNotExist) {
return nil, nil, nil
}
return nil, nil, fmt.Errorf("open previous manifest: %w", err)
}
previousManifest, err := (&manifest.LocalStore{}).LoadReader(ctx, manifestFile)
if err != nil {
_ = manifestFile.Close()
return nil, nil, fmt.Errorf("load previous manifest: %w", err)
}
if _, err := manifestFile.Seek(0, io.SeekStart); err != nil {
_ = manifestFile.Close()
return nil, nil, fmt.Errorf("rewind previous manifest: %w", err)
}
files := map[string]publishUploadFile{
"manifest.json": {RelativePath: filepath.ToSlash(filepath.Join(config.PathPreviousDirSegment, "manifest.json")), SourceRelativePath: "manifest.json"},
}
add := func(path, checksum string) error {
if strings.TrimSpace(path) == "" {
return nil
}
rel, err := pathsafe.SlashRelativeFromRoot(previousDir, path)
if err != nil {
return nil
}
archivePath := filepath.ToSlash(filepath.Join(config.PathPreviousDirSegment, rel))
candidate := publishUploadFile{RelativePath: archivePath, SourceRelativePath: rel, Checksum: strings.TrimSpace(checksum)}
if prior, exists := files[rel]; exists && prior.Checksum != "" && candidate.Checksum != "" && prior.Checksum != candidate.Checksum {
return fmt.Errorf("source %q has conflicting declared checksums", rel)
}
if prior, exists := files[rel]; exists && prior.Checksum != "" {
candidate.Checksum = prior.Checksum
}
files[rel] = candidate
return nil
}
for _, artifact := range previousManifest.Artifacts {
if err := add(artifact.ArchivePath, artifact.Checksum); err != nil {
_ = manifestFile.Close()
return nil, nil, err
}
if artifact.ArchivePath == "" {
if err := add(artifact.LocalPath, artifact.Checksum); err != nil {
_ = manifestFile.Close()
return nil, nil, err
}
}
}
for _, record := range previousManifest.Stages {
if record == nil {
continue
}
for _, artifact := range record.Outputs {
if err := add(artifact.ArchivePath, artifact.Checksum); err != nil {
_ = manifestFile.Close()
return nil, nil, err
}
if artifact.ArchivePath == "" {
if err := add(artifact.LocalPath, artifact.Checksum); err != nil {
_ = manifestFile.Close()
return nil, nil, err
}
}
}
for _, path := range append(append([]string(nil), record.Logs...), record.GeneratedConfigs...) {
if err := add(path, ""); err != nil {
_ = manifestFile.Close()
return nil, nil, err
}
}
}
out := make([]publishUploadFile, 0, len(files))
for _, file := range files {
out = append(out, file)
}
sort.Slice(out, func(i, j int) bool { return out[i].RelativePath < out[j].RelativePath })
return out, manifestFile, nil
}
func openPublishRunFiles(runRoot string, files []publishUploadFile, manifestFile *os.File) ([]*os.File, error) {
sources := make([]*os.File, 0, len(files))
for _, candidate := range files {
if manifestFile == nil {
closePublishSources(sources)
return nil, fmt.Errorf("archive manifest is required")
}
file := manifestFile
if candidate.SourceRelativePath != "manifest.json" {
var err error
file, err = fileops.OpenConfinedRegularFile(runRoot, candidate.SourceRelativePath)
if err != nil {
closePublishSources(sources)
return nil, err
}
}
if err := verifyPublishSourceChecksum(file, candidate.RelativePath, candidate.Checksum); err != nil {
if file != manifestFile {
_ = file.Close()
}
closePublishSources(sources)
return nil, err
}
sources = append(sources, file)
}
return sources, nil
}
func verifyPublishSourceChecksum(file *os.File, relativePath, expected string) error {
if strings.TrimSpace(expected) == "" {
return nil
}
hash := sha256.New()
if _, err := io.Copy(hash, file); err != nil {
return fmt.Errorf("read source %q for checksum: %w", relativePath, err)
}
if _, err := file.Seek(0, io.SeekStart); err != nil {
return fmt.Errorf("rewind source %q after checksum: %w", relativePath, err)
}
if actual := hex.EncodeToString(hash.Sum(nil)); !strings.EqualFold(actual, expected) {
return fmt.Errorf("source %q does not match its declared checksum", relativePath)
}
return nil
}
func closePublishSources(sources []*os.File) {
for _, source := range sources {
_ = source.Close()
}
}
func directoryExists(path string) (bool, error) {
info, err := os.Stat(path)
if err == nil {
return info.IsDir(), nil
}
if os.IsNotExist(err) {
return false, nil
}
return false, err
}
func writeCurrentManifestSnapshot(m *manifest.Manifest, publishMetadata map[string]any) (string, error) {
if m == nil {
return "", fmt.Errorf("manifest is required")
}
clone := *m
clone.Stages = make(map[string]*manifest.StageRecord, len(m.Stages))
for name, sr := range m.Stages {
if sr == nil {
continue
}
stageCopy := *sr
if sr.Outputs != nil {
stageCopy.Outputs = append([]manifest.ArtifactRecord(nil), sr.Outputs...)
}
if sr.Logs != nil {
stageCopy.Logs = append([]string(nil), sr.Logs...)
}
if sr.GeneratedConfigs != nil {
stageCopy.GeneratedConfigs = append([]string(nil), sr.GeneratedConfigs...)
}
if sr.Metadata != nil {
metaCopy := make(map[string]any, len(sr.Metadata))
for k, v := range sr.Metadata {
metaCopy[k] = v
}
stageCopy.Metadata = metaCopy
}
clone.Stages[name] = &stageCopy
}
now := time.Now().UTC()
clone.MarkStageSucceeded("publish", now, nil)
if sr := clone.Stages["publish"]; sr != nil {
sr.Metadata = publishMetadata
}
data, err := json.MarshalIndent(&clone, "", " ")
if err != nil {
return "", fmt.Errorf("marshal manifest: %w", err)
}
data = append(data, '\n')
tmp, err := os.CreateTemp("", "narratio-current-manifest-*.json")
if err != nil {
return "", fmt.Errorf("create temp file: %w", err)
}
path := tmp.Name()
if _, err := tmp.Write(data); err != nil {
_ = tmp.Close()
return "", fmt.Errorf("write temp manifest: %w", err)
}
if err := tmp.Close(); err != nil {
return "", fmt.Errorf("close temp manifest: %w", err)
}
return path, nil
}
func writeCurrentRunIDPointer(runID string) (string, error) {
tmp, err := os.CreateTemp("", "narratio-current-run-id-*.txt")
if err != nil {
return "", fmt.Errorf("create temp file: %w", err)
}
path := tmp.Name()
if _, err := tmp.WriteString(runID + "\n"); err != nil {
_ = tmp.Close()
return "", fmt.Errorf("write temp run id pointer: %w", err)
}
if err := tmp.Close(); err != nil {
return "", fmt.Errorf("close temp run id pointer: %w", err)
}
return path, nil
}
func publishMetadataPreview(
bucket, runPrefix, sessionPrefix string,
runUploaded []string,
publishedUploaded []string,
previousUploaded []string,
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...),
"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 skippedUnselectedOutputMetadata(skipped []publishSkippedUnselectedOutput) []map[string]any {
out := make([]map[string]any, 0, len(skipped))
for _, item := range skipped {
out = append(out, map[string]any{
"source": item.Source,
"dest": item.Dest,
"required": item.Required,
})
}
return out
}
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{
"source": item.Source,
"dest": item.Dest,
"remote_key": item.RemoteKey,
"reason": item.Reason,
"required": item.Required,
"local_path": item.LocalPath,
"provenance": item.Provenance,
})
}
return out
}