package state import ( "context" "encoding/json" "fmt" "os" "path/filepath" "sort" "strings" "time" "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"` } 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"` BriefingPath string `json:"briefingPath"` 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, } { 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) 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.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(ctx context.Context, resolved report.Resolved) (*PriorSnapshot, error) { return s.FindPriorSnapshot(ctx, resolved) } 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, err := reportGroup(resolved.Definition.ID) if err != nil { return nil, err } 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 !compatiblePriorReport(group, metadata.ReportID, resolved.Definition.ID) { 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], BriefingPath: candidates[0].BriefingPath, }, 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") } var pkg promptinput.Package if err := readJSON(path, &pkg); err != nil { return promptinput.Package{}, err } return pkg, 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, BriefingPath: metadata.BriefingPath, 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 compatiblePriorReport(group string, prior report.ID, current report.ID) bool { switch group { case "daily": return prior == report.DailyToday || prior == report.DailyTomorrow case "three-day": return prior == report.ThreeDay && current == report.ThreeDay case "weekend": return prior == report.Weekend && current == report.Weekend default: return false } } 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") } 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) } func (s *FilesystemStore) LoadBriefing(_ context.Context, path string) (briefing.Package, error) { if path == "" { return briefing.Package{}, fmt.Errorf("briefing path is required") } var pkg briefing.Package if err := readJSON(path, &pkg); err != nil { return briefing.Package{}, err } return pkg, nil }