Add filesystem state metadata baseline
This commit is contained in:
285
internal/state/filesystem.go
Normal file
285
internal/state/filesystem.go
Normal file
@@ -0,0 +1,285 @@
|
||||
package state
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"sort"
|
||||
"strings"
|
||||
|
||||
"gitea.maximumdirect.net/eric/weatherreporter/internal/adapters/scriptorium"
|
||||
"gitea.maximumdirect.net/eric/weatherreporter/internal/briefing"
|
||||
"gitea.maximumdirect.net/eric/weatherreporter/internal/config"
|
||||
"gitea.maximumdirect.net/eric/weatherreporter/internal/promptinput"
|
||||
"gitea.maximumdirect.net/eric/weatherreporter/internal/report"
|
||||
)
|
||||
|
||||
type FilesystemStore struct {
|
||||
root string
|
||||
snapshotsDir string
|
||||
reportsDir string
|
||||
dataPackagesDir string
|
||||
preflightDir string
|
||||
}
|
||||
|
||||
type ArtifactPaths struct {
|
||||
Briefing string `json:"briefing"`
|
||||
Metadata string `json:"metadata"`
|
||||
DataPackage string `json:"dataPackage"`
|
||||
Preflight string `json:"preflight"`
|
||||
RenderedReport string `json:"renderedReport,omitempty"`
|
||||
}
|
||||
|
||||
func NewFilesystemStore(cfg config.WorkspaceConfig) (*FilesystemStore, error) {
|
||||
if cfg.Root == "" {
|
||||
return nil, fmt.Errorf("workspace root is required")
|
||||
}
|
||||
for name, value := range map[string]string{
|
||||
"snapshots_dir": cfg.SnapshotsDir,
|
||||
"reports_dir": cfg.ReportsDir,
|
||||
"data_packages_dir": cfg.DataPackagesDir,
|
||||
"preflight_dir": cfg.PreflightDir,
|
||||
} {
|
||||
if err := validateRelativeDir(name, value); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
return &FilesystemStore{
|
||||
root: filepath.Clean(cfg.Root),
|
||||
snapshotsDir: filepath.Clean(cfg.SnapshotsDir),
|
||||
reportsDir: filepath.Clean(cfg.ReportsDir),
|
||||
dataPackagesDir: filepath.Clean(cfg.DataPackagesDir),
|
||||
preflightDir: filepath.Clean(cfg.PreflightDir),
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (s *FilesystemStore) Paths(resolved report.Resolved) (ArtifactPaths, error) {
|
||||
if s == nil {
|
||||
return ArtifactPaths{}, fmt.Errorf("state store is required")
|
||||
}
|
||||
metadata := resolved.Metadata()
|
||||
if metadata.RunID == "" {
|
||||
return ArtifactPaths{}, fmt.Errorf("run id is required")
|
||||
}
|
||||
group, err := reportGroup(resolved.Definition.ID)
|
||||
if err != nil {
|
||||
return ArtifactPaths{}, err
|
||||
}
|
||||
validDate := resolved.ValidPeriod.Start.Format("2006-01-02")
|
||||
filenameBase := metadata.RunID
|
||||
return ArtifactPaths{
|
||||
Briefing: s.join(s.snapshotsDir, group, validDate, filenameBase+".briefing.json"),
|
||||
Metadata: s.join(s.snapshotsDir, group, validDate, filenameBase+".metadata.json"),
|
||||
DataPackage: s.join(s.dataPackagesDir, group, validDate, filenameBase+".data_package.json"),
|
||||
Preflight: s.join(s.preflightDir, group, validDate, filenameBase+".render.json"),
|
||||
RenderedReport: s.join(s.reportsDir, group, filenameBase+".md"),
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (s *FilesystemStore) SaveBriefing(_ context.Context, resolved report.Resolved, pkg briefing.Package) (string, error) {
|
||||
paths, err := s.Paths(resolved)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
if err := writeJSONAtomic(paths.Briefing, pkg); err != nil {
|
||||
return "", err
|
||||
}
|
||||
return paths.Briefing, nil
|
||||
}
|
||||
|
||||
func (s *FilesystemStore) SaveDataPackage(_ context.Context, resolved report.Resolved, pkg promptinput.Package) (string, error) {
|
||||
paths, err := s.Paths(resolved)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
if err := promptinput.Validate(pkg); err != nil {
|
||||
return "", err
|
||||
}
|
||||
if err := writeJSONAtomic(paths.DataPackage, pkg); err != nil {
|
||||
return "", err
|
||||
}
|
||||
return paths.DataPackage, nil
|
||||
}
|
||||
|
||||
func (s *FilesystemStore) SavePreflight(_ context.Context, resolved report.Resolved, result *scriptorium.RenderResult) (string, error) {
|
||||
if result == nil {
|
||||
return "", fmt.Errorf("render result is required")
|
||||
}
|
||||
paths, err := s.Paths(resolved)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
if err := writeJSONAtomic(paths.Preflight, result); err != nil {
|
||||
return "", err
|
||||
}
|
||||
return paths.Preflight, nil
|
||||
}
|
||||
|
||||
func (s *FilesystemStore) SaveMetadata(_ context.Context, metadata Metadata) (string, error) {
|
||||
if metadata.RunID == "" {
|
||||
return "", fmt.Errorf("metadata run id is required")
|
||||
}
|
||||
if metadata.BriefingPath == "" {
|
||||
return "", fmt.Errorf("metadata briefing path is required")
|
||||
}
|
||||
if metadata.DataPackagePath == "" {
|
||||
return "", fmt.Errorf("metadata data package path is required")
|
||||
}
|
||||
if metadata.PreflightPath == "" {
|
||||
return "", fmt.Errorf("metadata preflight path is required")
|
||||
}
|
||||
path := metadataPathFromStored(metadata)
|
||||
if path == "" {
|
||||
return "", fmt.Errorf("metadata path cannot be resolved")
|
||||
}
|
||||
if err := writeJSONAtomic(path, metadata); err != nil {
|
||||
return "", err
|
||||
}
|
||||
return path, nil
|
||||
}
|
||||
|
||||
func (s *FilesystemStore) FindPriorDailySnapshot(_ context.Context, resolved report.Resolved) (*PriorSnapshot, error) {
|
||||
if resolved.Definition.ComparisonStrategy != report.CompareSameValidDate {
|
||||
return nil, nil
|
||||
}
|
||||
group, err := reportGroup(resolved.Definition.ID)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if group != "daily" {
|
||||
return nil, nil
|
||||
}
|
||||
paths, err := s.Paths(resolved)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
dir := filepath.Dir(paths.Metadata)
|
||||
entries, err := os.ReadDir(dir)
|
||||
if err != nil {
|
||||
if os.IsNotExist(err) {
|
||||
return nil, nil
|
||||
}
|
||||
return nil, fmt.Errorf("read snapshot metadata directory %q: %w", dir, err)
|
||||
}
|
||||
|
||||
var candidates []Metadata
|
||||
for _, entry := range entries {
|
||||
if entry.IsDir() || !strings.HasSuffix(entry.Name(), ".metadata.json") {
|
||||
continue
|
||||
}
|
||||
path := filepath.Join(dir, entry.Name())
|
||||
var metadata Metadata
|
||||
if err := readJSON(path, &metadata); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if metadata.RunID == resolved.Metadata().RunID {
|
||||
continue
|
||||
}
|
||||
if metadata.ReportID != report.DailyToday && metadata.ReportID != report.DailyTomorrow {
|
||||
continue
|
||||
}
|
||||
if !sameValidDate(metadata, resolved) {
|
||||
continue
|
||||
}
|
||||
if !metadata.GeneratedAt.Before(resolved.GeneratedAt) {
|
||||
continue
|
||||
}
|
||||
candidates = append(candidates, metadata)
|
||||
}
|
||||
if len(candidates) == 0 {
|
||||
return nil, nil
|
||||
}
|
||||
sort.Slice(candidates, func(i, j int) bool {
|
||||
return candidates[i].GeneratedAt.After(candidates[j].GeneratedAt)
|
||||
})
|
||||
return &PriorSnapshot{
|
||||
Metadata: candidates[0],
|
||||
BriefingPath: candidates[0].BriefingPath,
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (s *FilesystemStore) join(parts ...string) string {
|
||||
all := append([]string{s.root}, parts...)
|
||||
return filepath.Join(all...)
|
||||
}
|
||||
|
||||
func validateRelativeDir(name string, value string) error {
|
||||
if value == "" {
|
||||
return fmt.Errorf("%s is required", name)
|
||||
}
|
||||
if filepath.IsAbs(value) {
|
||||
return fmt.Errorf("%s must be relative to workspace root", name)
|
||||
}
|
||||
cleaned := filepath.Clean(value)
|
||||
if cleaned == "." || cleaned == ".." || strings.HasPrefix(cleaned, ".."+string(filepath.Separator)) {
|
||||
return fmt.Errorf("%s must stay within workspace root", name)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func reportGroup(id report.ID) (string, error) {
|
||||
switch id {
|
||||
case report.DailyToday, report.DailyTomorrow:
|
||||
return "daily", nil
|
||||
case report.ThreeDay:
|
||||
return "three-day", nil
|
||||
case report.Weekend:
|
||||
return "weekend", nil
|
||||
case report.Storm:
|
||||
return "storm", nil
|
||||
default:
|
||||
return "", fmt.Errorf("unknown report %q", id)
|
||||
}
|
||||
}
|
||||
|
||||
func writeJSONAtomic(path string, value any) error {
|
||||
data, err := json.MarshalIndent(value, "", " ")
|
||||
if err != nil {
|
||||
return fmt.Errorf("marshal %q: %w", path, err)
|
||||
}
|
||||
if err := os.MkdirAll(filepath.Dir(path), 0o755); err != nil {
|
||||
return fmt.Errorf("create directory %q: %w", filepath.Dir(path), err)
|
||||
}
|
||||
tmp, err := os.CreateTemp(filepath.Dir(path), "."+filepath.Base(path)+".*.tmp")
|
||||
if err != nil {
|
||||
return fmt.Errorf("create temporary file for %q: %w", path, err)
|
||||
}
|
||||
tmpName := tmp.Name()
|
||||
defer os.Remove(tmpName)
|
||||
if _, err := tmp.Write(data); err != nil {
|
||||
tmp.Close()
|
||||
return fmt.Errorf("write temporary file for %q: %w", path, err)
|
||||
}
|
||||
if err := tmp.Close(); err != nil {
|
||||
return fmt.Errorf("close temporary file for %q: %w", path, err)
|
||||
}
|
||||
if err := os.Rename(tmpName, path); err != nil {
|
||||
return fmt.Errorf("save %q: %w", path, err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func readJSON(path string, target any) error {
|
||||
data, err := os.ReadFile(path)
|
||||
if err != nil {
|
||||
return fmt.Errorf("read %q: %w", path, err)
|
||||
}
|
||||
if err := json.Unmarshal(data, target); err != nil {
|
||||
return fmt.Errorf("decode %q: %w", path, err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func metadataPathFromStored(metadata Metadata) string {
|
||||
if metadata.BriefingPath == "" {
|
||||
return ""
|
||||
}
|
||||
filename := metadata.RunID + ".metadata.json"
|
||||
return filepath.Join(filepath.Dir(metadata.BriefingPath), filename)
|
||||
}
|
||||
|
||||
func sameValidDate(metadata Metadata, resolved report.Resolved) bool {
|
||||
return metadata.ValidPeriod.Start.Format("2006-01-02") == resolved.ValidPeriod.Start.Format("2006-01-02")
|
||||
}
|
||||
Reference in New Issue
Block a user