Files
weatherreporter/internal/state/filesystem.go

504 lines
16 KiB
Go

package state
import (
"context"
"encoding/json"
"fmt"
"os"
"path/filepath"
"sort"
"strings"
"time"
"gitea.maximumdirect.net/eric/weatherreporter/internal/config"
"gitea.maximumdirect.net/eric/weatherreporter/internal/fileutil"
"gitea.maximumdirect.net/eric/weatherreporter/internal/module"
"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
notificationsDir string
}
type ArtifactPaths struct {
ModuleSnapshot string `json:"moduleSnapshot"`
Metadata string `json:"metadata"`
DataPackage string `json:"dataPackage"`
Preflight string `json:"preflight"`
Notification string `json:"notification,omitempty"`
RenderedReport string `json:"renderedReport,omitempty"`
GeneratedTextRaw string `json:"generatedTextRaw,omitempty"`
GeneratedTextResult string `json:"generatedTextResult,omitempty"`
GeneratedText string `json:"generatedText,omitempty"`
RenderContext string `json:"renderContext,omitempty"`
}
type ReportRecord struct {
RunID string `json:"runId"`
ReportID report.ID `json:"reportId"`
Variant string `json:"variant,omitempty"`
PromptID string `json:"promptId"`
GeneratedAt string `json:"generatedAt"`
ValidStart string `json:"validStart"`
ValidEnd string `json:"validEnd"`
MetadataPath string `json:"metadataPath"`
ReportPath string `json:"reportPath,omitempty"`
Warnings int `json:"warnings"`
metadata Metadata
}
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,
"notifications_dir": cfg.NotificationsDir,
} {
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),
notificationsDir: filepath.Clean(cfg.NotificationsDir),
}, 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 := resolved.Definition.ArtifactGroup
if group == "" {
return ArtifactPaths{}, fmt.Errorf("report %q has no artifact group", resolved.Definition.ID)
}
validDate := resolved.ValidPeriod.Start.Format("2006-01-02")
filenameBase := metadata.RunID
return ArtifactPaths{
ModuleSnapshot: s.join(s.snapshotsDir, group, validDate, filenameBase+".modules.json"),
Metadata: s.join(s.snapshotsDir, group, validDate, filenameBase+".metadata.json"),
DataPackage: s.join(s.dataPackagesDir, group, validDate, filenameBase+".data_package.yaml"),
Preflight: s.join(s.preflightDir, group, validDate, filenameBase+".render.json"),
Notification: s.join(s.notificationsDir, group, validDate, filenameBase+".distributor.json"),
RenderedReport: s.join(s.reportsDir, group, filenameBase+".md"),
GeneratedTextRaw: s.join(s.snapshotsDir, group, validDate, filenameBase+".generated_text.raw.json"),
GeneratedTextResult: s.join(s.snapshotsDir, group, validDate, filenameBase+".generated_text.run.json"),
GeneratedText: s.join(s.snapshotsDir, group, validDate, filenameBase+".generated_text.json"),
RenderContext: s.join(s.snapshotsDir, group, validDate, filenameBase+".render_context.json"),
}, nil
}
func (s *FilesystemStore) SaveModuleSnapshot(_ context.Context, resolved report.Resolved, snapshot module.Snapshot) (string, error) {
paths, err := s.Paths(resolved)
if err != nil {
return "", err
}
if err := snapshot.Validate(); err != nil {
return "", err
}
if err := fileutil.WriteJSONAtomic(paths.ModuleSnapshot, snapshot); err != nil {
return "", err
}
return paths.ModuleSnapshot, 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.Save(paths.DataPackage, pkg); err != nil {
return "", err
}
return paths.DataPackage, nil
}
func (s *FilesystemStore) SavePreflight(_ context.Context, resolved report.Resolved, artifact PreflightArtifact) (string, error) {
paths, err := s.Paths(resolved)
if err != nil {
return "", err
}
if err := fileutil.WriteJSONAtomic(paths.Preflight, artifact); err != nil {
return "", err
}
return paths.Preflight, nil
}
func (s *FilesystemStore) SaveDistributorNotification(_ context.Context, resolved report.Resolved, artifact DistributorNotificationArtifact) (string, error) {
paths, err := s.Paths(resolved)
if err != nil {
return "", err
}
if artifact.SchemaVersion == "" {
artifact.SchemaVersion = DistributorNotificationSchemaVersion
}
if err := fileutil.WriteJSONAtomic(paths.Notification, artifact); err != nil {
return "", err
}
return paths.Notification, nil
}
func (s *FilesystemStore) SaveGeneratedTextRaw(_ context.Context, resolved report.Resolved, data []byte) (string, error) {
paths, err := s.Paths(resolved)
if err != nil {
return "", err
}
if err := fileutil.WriteFileAtomic(paths.GeneratedTextRaw, data); err != nil {
return "", err
}
return paths.GeneratedTextRaw, nil
}
func (s *FilesystemStore) SaveGeneratedTextResult(_ context.Context, resolved report.Resolved, value any) (string, error) {
paths, err := s.Paths(resolved)
if err != nil {
return "", err
}
if err := fileutil.WriteJSONAtomic(paths.GeneratedTextResult, value); err != nil {
return "", err
}
return paths.GeneratedTextResult, nil
}
func (s *FilesystemStore) SaveGeneratedText(_ context.Context, resolved report.Resolved, data []byte) (string, error) {
paths, err := s.Paths(resolved)
if err != nil {
return "", err
}
if err := fileutil.WriteFileAtomic(paths.GeneratedText, data); err != nil {
return "", err
}
return paths.GeneratedText, nil
}
func (s *FilesystemStore) SaveRenderContext(_ context.Context, resolved report.Resolved, value any) (string, error) {
paths, err := s.Paths(resolved)
if err != nil {
return "", err
}
if err := fileutil.WriteJSONAtomic(paths.RenderContext, value); err != nil {
return "", err
}
return paths.RenderContext, nil
}
func (s *FilesystemStore) PrepareRenderedReport(_ context.Context, resolved report.Resolved) (string, error) {
paths, err := s.Paths(resolved)
if err != nil {
return "", err
}
if err := os.MkdirAll(filepath.Dir(paths.RenderedReport), 0o755); err != nil {
return "", fmt.Errorf("create rendered report directory %q: %w", filepath.Dir(paths.RenderedReport), err)
}
return paths.RenderedReport, nil
}
func (s *FilesystemStore) SaveMetadata(_ context.Context, metadata Metadata) (string, error) {
if metadata.RunID == "" {
return "", fmt.Errorf("metadata run id is required")
}
if metadata.ModuleSnapshotPath == "" {
return "", fmt.Errorf("metadata module snapshot 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")
}
if metadata.MetadataPath == "" {
return "", fmt.Errorf("metadata path is required")
}
if err := fileutil.WriteJSONAtomic(metadata.MetadataPath, metadata); err != nil {
return "", err
}
return metadata.MetadataPath, nil
}
func (s *FilesystemStore) FindPriorSnapshot(_ context.Context, resolved report.Resolved) (*PriorSnapshot, error) {
if resolved.Definition.ComparisonStrategy != report.CompareSameValidDate && resolved.Definition.ComparisonStrategy != report.CompareWeekendWindow {
return nil, nil
}
group := resolved.Definition.ArtifactGroup
if group == "" {
return nil, fmt.Errorf("report %q has no artifact group", resolved.Definition.ID)
}
dirs, err := s.metadataDirectories(resolved, group)
if err != nil {
return nil, err
}
var candidates []Metadata
for _, dir := range dirs {
entries, err := os.ReadDir(dir)
if err != nil {
if os.IsNotExist(err) {
continue
}
return nil, fmt.Errorf("read snapshot metadata directory %q: %w", dir, err)
}
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 !resolved.Definition.CompatibleWithPrior(metadata.ReportID) {
continue
}
if !comparablePeriod(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],
ModuleSnapshotPath: candidates[0].ModuleSnapshotPath,
}, nil
}
func (s *FilesystemStore) ListReports(_ context.Context, limit int) ([]ReportRecord, error) {
if s == nil {
return nil, fmt.Errorf("state store is required")
}
root := s.join(s.snapshotsDir)
if _, err := os.Stat(root); err != nil {
if os.IsNotExist(err) {
return nil, nil
}
return nil, fmt.Errorf("inspect %q: %w", root, err)
}
var records []ReportRecord
err := filepath.WalkDir(root, func(path string, entry os.DirEntry, err error) error {
if err != nil {
return fmt.Errorf("inspect %q: %w", path, err)
}
if entry.IsDir() || !strings.HasSuffix(entry.Name(), ".metadata.json") {
return nil
}
record, err := s.reportRecord(path)
if err != nil {
return err
}
records = append(records, record)
return nil
})
if err != nil {
if os.IsNotExist(err) {
return nil, nil
}
return nil, err
}
sort.Slice(records, func(i, j int) bool {
return records[i].metadata.GeneratedAt.After(records[j].metadata.GeneratedAt)
})
if limit > 0 && len(records) > limit {
records = records[:limit]
}
return records, nil
}
func (s *FilesystemStore) LoadMetadataByRunID(ctx context.Context, runID string) (Metadata, string, error) {
if strings.TrimSpace(runID) == "" {
return Metadata{}, "", fmt.Errorf("run id is required")
}
records, err := s.ListReports(ctx, 0)
if err != nil {
return Metadata{}, "", err
}
for _, record := range records {
if record.RunID == runID {
return record.metadata, record.MetadataPath, nil
}
}
return Metadata{}, "", fmt.Errorf("metadata for run id %q was not found", runID)
}
func (s *FilesystemStore) LoadDataPackage(_ context.Context, path string) (promptinput.Package, error) {
if path == "" {
return promptinput.Package{}, fmt.Errorf("data package path is required")
}
data, err := os.ReadFile(path)
if err != nil {
return promptinput.Package{}, fmt.Errorf("read %q: %w", path, err)
}
pkg, err := promptinput.LoadYAML(data)
if err != nil {
return promptinput.Package{}, err
}
return pkg, nil
}
func (s *FilesystemStore) LoadModuleSnapshot(_ context.Context, path string) (module.Snapshot, error) {
if path == "" {
return module.Snapshot{}, fmt.Errorf("module snapshot path is required")
}
var snapshot module.Snapshot
if err := readJSON(path, &snapshot); err != nil {
return module.Snapshot{}, err
}
if err := snapshot.Validate(); err != nil {
return module.Snapshot{}, err
}
return snapshot, nil
}
func (s *FilesystemStore) LoadGeneratedText(_ context.Context, path string) ([]byte, error) {
if path == "" {
return nil, fmt.Errorf("generated text path is required")
}
data, err := os.ReadFile(path)
if err != nil {
return nil, fmt.Errorf("read %q: %w", path, err)
}
return data, nil
}
func (s *FilesystemStore) LoadGeneratedTextResult(_ context.Context, path string, target any) error {
if path == "" {
return fmt.Errorf("generated text result path is required")
}
if target == nil {
return fmt.Errorf("generated text result target is required")
}
return readJSON(path, target)
}
func (s *FilesystemStore) LoadRenderContext(_ context.Context, path string, target any) error {
if path == "" {
return fmt.Errorf("render context path is required")
}
if target == nil {
return fmt.Errorf("render context target is required")
}
return readJSON(path, target)
}
func (s *FilesystemStore) reportRecord(path string) (ReportRecord, error) {
var metadata Metadata
if err := readJSON(path, &metadata); err != nil {
return ReportRecord{}, err
}
return ReportRecord{
RunID: metadata.RunID,
ReportID: metadata.ReportID,
Variant: metadata.Variant,
PromptID: metadata.PromptID,
GeneratedAt: metadata.GeneratedAt.Format(time.RFC3339Nano),
ValidStart: metadata.ValidPeriod.Start.Format(time.RFC3339Nano),
ValidEnd: metadata.ValidPeriod.End.Format(time.RFC3339Nano),
MetadataPath: path,
ReportPath: metadata.RenderedReportPath,
Warnings: len(metadata.SourceWarnings),
metadata: metadata,
}, nil
}
func (s *FilesystemStore) metadataDirectories(resolved report.Resolved, group string) ([]string, error) {
paths, err := s.Paths(resolved)
if err != nil {
return nil, err
}
if resolved.Definition.ComparisonStrategy != report.CompareWeekendWindow {
return []string{filepath.Dir(paths.Metadata)}, nil
}
root := s.join(s.snapshotsDir, group)
entries, err := os.ReadDir(root)
if err != nil {
if os.IsNotExist(err) {
return nil, nil
}
return nil, fmt.Errorf("read snapshot group directory %q: %w", root, err)
}
var dirs []string
for _, entry := range entries {
if entry.IsDir() {
dirs = append(dirs, filepath.Join(root, entry.Name()))
}
}
return dirs, 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 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 sameValidDate(metadata Metadata, resolved report.Resolved) bool {
return metadata.ValidPeriod.Start.Format("2006-01-02") == resolved.ValidPeriod.Start.Format("2006-01-02")
}
func comparablePeriod(metadata Metadata, resolved report.Resolved) bool {
switch resolved.Definition.ComparisonStrategy {
case report.CompareSameValidDate:
return sameValidDate(metadata, resolved)
case report.CompareWeekendWindow:
return sameWeekendWindow(metadata, resolved)
default:
return false
}
}
func sameWeekendWindow(metadata Metadata, resolved report.Resolved) bool {
return metadata.ValidPeriod.End.Equal(resolved.ValidPeriod.End) && !metadata.ValidPeriod.Start.After(resolved.ValidPeriod.Start)
}