Files
weatherreporter/internal/state/filesystem.go

556 lines
18 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 err := validatePathSegment("run id", metadata.RunID); err != nil {
return ArtifactPaths{}, err
}
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")
return ArtifactPaths{
ModuleSnapshot: s.join(s.snapshotsDir, group, validDate, "modules."+metadata.RunID+".json"),
Metadata: s.join(s.snapshotsDir, group, validDate, "metadata."+metadata.RunID+".json"),
DataPackage: s.join(s.dataPackagesDir, group, validDate, "data_package."+metadata.RunID+".yaml"),
Preflight: s.join(s.preflightDir, group, validDate, "render."+metadata.RunID+".json"),
Notification: s.join(s.notificationsDir, group, validDate, "distributor."+metadata.RunID+".json"),
RenderedReport: s.join(s.reportsDir, group, validDate, "report."+metadata.RunID+".md"),
GeneratedTextRaw: s.join(s.snapshotsDir, group, validDate, "generated_text_raw."+metadata.RunID+".json"),
GeneratedTextResult: s.join(s.snapshotsDir, group, validDate, "generated_text_result."+metadata.RunID+".json"),
GeneratedText: s.join(s.snapshotsDir, group, validDate, "generated_text."+metadata.RunID+".json"),
RenderContext: s.join(s.snapshotsDir, group, validDate, "render_context."+metadata.RunID+".json"),
}, nil
}
func (s *FilesystemStore) SaveModuleSnapshot(_ context.Context, resolved report.Resolved, snapshot module.Snapshot) (string, error) {
if err := snapshot.Validate(); err != nil {
return "", err
}
return s.saveResolvedJSON(resolved, func(paths ArtifactPaths) string {
return paths.ModuleSnapshot
}, snapshot)
}
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) {
return s.saveResolvedJSON(resolved, func(paths ArtifactPaths) string {
return paths.Preflight
}, artifact)
}
func (s *FilesystemStore) SaveDistributorNotification(_ context.Context, resolved report.Resolved, artifact DistributorNotificationArtifact) (string, error) {
if artifact.SchemaVersion == "" {
artifact.SchemaVersion = DistributorNotificationSchemaVersion
}
return s.saveResolvedJSON(resolved, func(paths ArtifactPaths) string {
return paths.Notification
}, artifact)
}
func (s *FilesystemStore) BatchDistributorNotificationPath(ref BatchDistributorNotificationRef) (string, error) {
if s == nil {
return "", fmt.Errorf("state store is required")
}
if err := validateBatchNotificationRef(ref); err != nil {
return "", err
}
localDate := ref.StartedAt.In(ref.Location).Format("2006-01-02")
return s.join(s.notificationsDir, "batches", ref.Batch, localDate, "distributor."+ref.BatchRunID+".json"), nil
}
func (s *FilesystemStore) SaveBatchDistributorNotification(_ context.Context, ref BatchDistributorNotificationRef, artifact BatchDistributorNotificationArtifact) (string, error) {
path, err := s.BatchDistributorNotificationPath(ref)
if err != nil {
return "", err
}
if artifact.SchemaVersion == "" {
artifact.SchemaVersion = BatchDistributorNotificationSchemaVersion
}
artifact.Batch = ref.Batch
artifact.BatchRunID = ref.BatchRunID
if err := fileutil.WriteJSONAtomic(path, artifact); err != nil {
return "", err
}
return path, nil
}
func (s *FilesystemStore) SaveGeneratedTextRaw(_ context.Context, resolved report.Resolved, data []byte) (string, error) {
return s.saveResolvedBytes(resolved, func(paths ArtifactPaths) string {
return paths.GeneratedTextRaw
}, data)
}
func (s *FilesystemStore) SaveGeneratedTextResult(_ context.Context, resolved report.Resolved, value any) (string, error) {
return s.saveResolvedJSON(resolved, func(paths ArtifactPaths) string {
return paths.GeneratedTextResult
}, value)
}
func (s *FilesystemStore) SaveGeneratedText(_ context.Context, resolved report.Resolved, data []byte) (string, error) {
return s.saveResolvedBytes(resolved, func(paths ArtifactPaths) string {
return paths.GeneratedText
}, data)
}
func (s *FilesystemStore) SaveRenderContext(_ context.Context, resolved report.Resolved, value any) (string, error) {
return s.saveResolvedJSON(resolved, func(paths ArtifactPaths) string {
return paths.RenderContext
}, value)
}
func (s *FilesystemStore) saveResolvedJSON(resolved report.Resolved, selectPath func(ArtifactPaths) string, value any) (string, error) {
paths, err := s.Paths(resolved)
if err != nil {
return "", err
}
path := selectPath(paths)
if err := fileutil.WriteJSONAtomic(path, value); err != nil {
return "", err
}
return path, nil
}
func (s *FilesystemStore) saveResolvedBytes(resolved report.Resolved, selectPath func(ArtifactPaths) string, data []byte) (string, error) {
paths, err := s.Paths(resolved)
if err != nil {
return "", err
}
path := selectPath(paths)
if err := fileutil.WriteFileAtomic(path, data); err != nil {
return "", err
}
return path, 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() || !isMetadataFilename(entry.Name()) {
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() || !isMetadataFilename(entry.Name()) {
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 validateBatchNotificationRef(ref BatchDistributorNotificationRef) error {
if err := validatePathSegment("batch kind", ref.Batch); err != nil {
return err
}
if err := validatePathSegment("batch run id", ref.BatchRunID); err != nil {
return err
}
if ref.StartedAt.IsZero() {
return fmt.Errorf("batch started time is required")
}
if ref.Location == nil {
return fmt.Errorf("batch location is required")
}
return nil
}
func validatePathSegment(name string, value string) error {
if strings.TrimSpace(value) == "" {
return fmt.Errorf("%s is required", name)
}
if strings.ContainsAny(value, `/\`) {
return fmt.Errorf("%s must not contain path separators", name)
}
if value == "." || value == ".." {
return fmt.Errorf("%s must be a safe path segment", name)
}
return nil
}
func isMetadataFilename(name string) bool {
if !strings.HasPrefix(name, "metadata.") || !strings.HasSuffix(name, ".json") {
return false
}
runID := strings.TrimSuffix(strings.TrimPrefix(name, "metadata."), ".json")
return strings.TrimSpace(runID) != "" && !strings.ContainsAny(runID, `/\`) && runID != "." && runID != ".."
}
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)
}