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"` } 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"), }, 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) 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) 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) }