Files
narratio/internal/previouscache/previouscache.go

449 lines
13 KiB
Go

package previouscache
import (
"context"
"errors"
"fmt"
"path"
"path/filepath"
"sort"
"strings"
"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/manifest"
"gitea.maximumdirect.net/eric/narratio/internal/pathsafe"
)
const (
InputKindManifest = "previous_manifest"
InputKindArtifact = "previous_artifact"
InputSource = "previous_session_publish.current"
)
type Plan struct {
Records []Record
SkippedMissing []string
PreviousRunID string
}
type Record struct {
Kind string
RequirementName string
Required bool
LocalRelativePath string
LocalPath string
RemoteKey string
SHA256 string
Size int64
Generation string
S3Bucket string
}
func BuildPlan(
ctx context.Context,
cfg *config.Config,
paths artifacts.SessionPaths,
requirements []artifacts.PreviousArtifactRequirement,
store storage.ObjectStore,
) (*Plan, error) {
if len(requirements) == 0 {
return &Plan{}, nil
}
if cfg == nil || cfg.Session == nil || cfg.Pipeline == nil {
return nil, fmt.Errorf("resolved config with session/pipeline 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(cfg.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 &Plan{SkippedMissing: optionalNames}, nil
}
if store == nil {
return nil, fmt.Errorf("previous-session artifact hydration requires object store backend")
}
if cfg.Pipeline.Storage.S3 == nil {
return nil, fmt.Errorf("pipeline.storage.s3 configuration is required for previous-session artifact hydration")
}
campaign := strings.TrimSpace(cfg.Session.Campaign)
if campaign == "" {
return nil, fmt.Errorf("session campaign is required for previous-session artifact hydration")
}
bucket := strings.TrimSpace(cfg.Pipeline.Storage.S3.Bucket)
if bucket == "" {
return nil, fmt.Errorf("pipeline.storage.s3.bucket is required for previous-session artifact hydration")
}
previousSessionPrefix := artifacts.S3SessionPrefix(
cfg.Pipeline.Storage.S3.RootPrefix,
campaign,
previousSessionID,
)
result := &Plan{}
current, err := artifacts.LoadCurrentState(ctx, store, previousSessionPrefix, artifacts.CurrentStateValidation{
ExpectedSessionID: previousSessionID,
ExpectedCampaign: campaign,
ValidateRunID: true,
})
if err != nil {
var runPointerMissing *artifacts.CurrentRunPointerMissingError
var manifestMissing *artifacts.CurrentManifestMissingError
if errors.As(err, &runPointerMissing) || errors.As(err, &manifestMissing) {
if len(requiredNames) > 0 {
return nil, fmt.Errorf("required previous-session artifacts unavailable: remote %w", err)
}
result.SkippedMissing = optionalNames
return result, nil
}
return nil, fmt.Errorf("load previous-session current state: %w", err)
}
result.PreviousRunID = strings.TrimSpace(current.RunID)
currentManifestKey := current.CurrentManifestKey
previousManifest := current.Manifest
manifestRel, err := relativeToSession(paths, paths.PreviousManifestPath)
if err != nil {
return nil, err
}
manifestRecord := Record{
Kind: InputKindManifest,
LocalRelativePath: manifestRel,
LocalPath: paths.PreviousManifestPath,
RemoteKey: currentManifestKey,
S3Bucket: bucket,
}
if current.Commit != nil {
committedManifest, ok := current.Commit.Artifact(artifacts.RemoteArtifactTypeSessionManifest)
if !ok || committedManifest.DestinationKey != currentManifestKey {
return nil, fmt.Errorf("previous-session remote commit does not declare its session manifest")
}
manifestRecord.SHA256 = committedManifest.SHA256
manifestRecord.Size = committedManifest.Size
manifestRecord.Generation = committedManifest.Generation
}
result.Records = append(result.Records, manifestRecord)
for _, requirement := range orderedRequirements {
candidates := artifactRelativePathCandidates(requirement.Name, previousManifest, cfg)
if len(candidates) == 0 {
if requirement.Required {
return nil, fmt.Errorf(
"required previous-session artifact %q is unavailable in previous-session manifest/published-state",
requirement.Name,
)
}
result.SkippedMissing = append(result.SkippedMissing, requirement.Name)
continue
}
selectedRel := ""
selectedKey := ""
var selectedArtifact *artifacts.RemoteArtifact
for _, candidate := range candidates {
if current.Commit != nil {
artifact, ok := committedPublishedArtifact(current.Commit, previousSessionPrefix, candidate)
if !ok {
continue
}
selectedRel = candidate
selectedKey = artifact.DestinationKey
selectedArtifact = &artifact
break
}
remoteKey := artifacts.S3PublishedOutputKey(previousSessionPrefix, candidate)
exists, err := store.Exists(ctx, remoteKey)
if err != nil {
return nil, fmt.Errorf("check previous-session artifact object %q: %w", remoteKey, err)
}
if exists {
selectedRel = candidate
selectedKey = remoteKey
break
}
}
if selectedRel == "" {
if requirement.Required {
return nil, fmt.Errorf(
"required previous-session artifact %q object missing from published candidate keys",
requirement.Name,
)
}
result.SkippedMissing = append(result.SkippedMissing, requirement.Name)
continue
}
localPath, err := artifacts.SessionPreviousArtifactPath(paths, selectedRel)
if err != nil {
return nil, fmt.Errorf("resolve previous-session artifact path: %w", err)
}
localRel, err := relativeToSession(paths, localPath)
if err != nil {
return nil, err
}
record := Record{
Kind: InputKindArtifact,
RequirementName: requirement.Name,
Required: requirement.Required,
LocalRelativePath: localRel,
LocalPath: localPath,
RemoteKey: selectedKey,
S3Bucket: bucket,
}
if selectedArtifact != nil {
record.SHA256 = selectedArtifact.SHA256
record.Size = selectedArtifact.Size
record.Generation = selectedArtifact.Generation
}
result.Records = append(result.Records, record)
}
sort.Strings(result.SkippedMissing)
sort.Slice(result.Records, func(i, j int) bool {
if result.Records[i].LocalRelativePath != result.Records[j].LocalRelativePath {
return result.Records[i].LocalRelativePath < result.Records[j].LocalRelativePath
}
return result.Records[i].RemoteKey < result.Records[j].RemoteKey
})
return result, nil
}
func committedPublishedArtifact(commit *artifacts.RemoteCommitManifest, sessionPrefix, relativePath string) (artifacts.RemoteArtifact, bool) {
if commit == nil {
return artifacts.RemoteArtifact{}, false
}
want := artifacts.S3RunRelativeDestinationKey(artifacts.S3RunPrefix(sessionPrefix, commit.RunID), relativePath)
for _, artifact := range commit.Artifacts {
if artifact.Type == artifacts.RemoteArtifactTypePublishedOutput && artifact.DestinationKey == want {
return artifact, true
}
}
return artifacts.RemoteArtifact{}, false
}
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 artifactRelativePathCandidates(
artifactName string,
previousManifest *manifest.Manifest,
cfg *config.Config,
) []string {
candidates := []string{}
appendCandidate := func(v string) {
normalized, err := pathsafe.NormalizeRelativeDestination(v)
if err != nil {
return
}
candidates = append(candidates, normalized)
}
sourceDescriptor, err := artifactpolicy.PreviousSessionSourceDescriptorForConfiguredKey(artifactName)
sourceID := artifactpolicy.ConfiguredSourceID(artifactName)
if err == nil {
sourceID = sourceDescriptor.ConfiguredSourceID
}
if rel, ok := manifestArtifactRelativePathBySourceID(previousManifest, sourceID); ok {
appendCandidate(rel)
base := path.Base(rel)
for _, published := range manifestPublishedPaths(previousManifest) {
if path.Base(published) == base {
appendCandidate(published)
}
}
}
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 manifestArtifactRelativePathBySourceID(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 := deriveManifestRelativePath(previousManifest, out.LocalPath)
if ok {
return rel, true
}
}
}
return "", false
}
func deriveManifestRelativePath(previousManifest *manifest.Manifest, localPath string) (string, bool) {
trimmed := strings.TrimSpace(localPath)
if trimmed == "" {
return "", false
}
if !filepath.IsAbs(trimmed) {
normalized, err := pathsafe.NormalizeRelativeDestination(filepath.ToSlash(trimmed))
if err != nil {
return "", false
}
return normalized, true
}
sessionRoot, ok := manifestSessionRoot(previousManifest)
if !ok {
return "", false
}
normalized, err := pathsafe.SlashRelativeFromRoot(sessionRoot, trimmed)
if err != nil {
return "", false
}
return normalized, true
}
func manifestSessionRoot(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 manifestPublishedPaths(previousManifest *manifest.Manifest) []string {
if previousManifest == nil || len(previousManifest.Stages) == 0 {
return nil
}
sr := previousManifest.Stages["publish"]
if sr == nil || sr.Metadata == nil {
return nil
}
raw, ok := sr.Metadata["published_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 := pathsafe.NormalizeRelativeDestination(asString)
if err != nil {
continue
}
out = append(out, normalized)
}
return dedupeOrderedStrings(out)
}
func relativeToSession(paths artifacts.SessionPaths, localPath string) (string, error) {
root := filepath.Clean(paths.Root)
if strings.TrimSpace(root) == "" {
return "", fmt.Errorf("session root is required")
}
normalized, err := pathsafe.SlashRelativeFromRoot(root, localPath)
if err != nil {
return "", fmt.Errorf("resolve previous-cache relative path: %w", err)
}
return normalized, nil
}
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
}