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"` Preparation string `json:"preparation,omitempty"` Execution string `json:"execution,omitempty"` Notification string `json:"notification,omitempty"` RenderedReport string `json:"renderedReport,omitempty"` GeneratedTextRaw string `json:"generatedTextRaw,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"), Preparation: s.join(s.preflightDir, group, validDate, "prompt_preparation."+metadata.RunID+".json"), Execution: s.join(s.snapshotsDir, group, validDate, "prompt_execution."+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"), 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(ctx context.Context, resolved report.Resolved, pkg promptinput.Package) (string, error) { data, err := promptinput.MarshalYAML(pkg) if err != nil { return "", err } return s.SaveDataPackageBytes(ctx, resolved, data) } func (s *FilesystemStore) SaveDataPackageBytes(_ context.Context, resolved report.Resolved, data []byte) (string, error) { paths, err := s.Paths(resolved) if err != nil { return "", err } if err := fileutil.WriteFileAtomic(paths.DataPackage, data); err != nil { return "", err } return paths.DataPackage, nil } func (s *FilesystemStore) SavePromptPreparation(_ context.Context, resolved report.Resolved, artifact PromptPreparationArtifact) (string, error) { if artifact.SchemaVersion == "" { artifact.SchemaVersion = PromptPreparationSchemaVersion } if err := artifact.Validate(); err != nil { return "", err } return s.saveResolvedJSON(resolved, func(paths ArtifactPaths) string { return paths.Preparation }, artifact) } func (s *FilesystemStore) SavePromptExecution(_ context.Context, resolved report.Resolved, artifact PromptExecutionArtifact) (string, error) { if artifact.SchemaVersion == "" { artifact.SchemaVersion = PromptExecutionSchemaVersion } if err := artifact.Validate(); err != nil { return "", err } return s.saveResolvedJSON(resolved, func(paths ArtifactPaths) string { return paths.Execution }, 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) 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.SchemaVersion != MetadataSchemaVersion { return "", fmt.Errorf("new metadata must use schema version %q", MetadataSchemaVersion) } if err := metadata.Validate(); err != nil { return "", err } if err := s.validateManagedPath("metadata path", metadata.MetadataPath); err != nil { return "", err } 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 { 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) LoadPromptPreparation(_ context.Context, path string) (PromptPreparationArtifact, error) { if path == "" { return PromptPreparationArtifact{}, fmt.Errorf("prompt preparation path is required") } var artifact PromptPreparationArtifact if err := readJSON(path, &artifact); err != nil { return PromptPreparationArtifact{}, err } if err := artifact.Validate(); err != nil { return PromptPreparationArtifact{}, err } return artifact, nil } func (s *FilesystemStore) LoadPromptExecution(_ context.Context, path string) (PromptExecutionArtifact, error) { if path == "" { return PromptExecutionArtifact{}, fmt.Errorf("prompt execution path is required") } var artifact PromptExecutionArtifact if err := readJSON(path, &artifact); err != nil { return PromptExecutionArtifact{}, err } if err := artifact.Validate(); err != nil { return PromptExecutionArtifact{}, err } return artifact, nil } 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 } metadata.MetadataPath = path 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 } return []string{filepath.Dir(paths.Metadata)}, nil } func (s *FilesystemStore) join(parts ...string) string { all := append([]string{s.root}, parts...) return filepath.Join(all...) } func (s *FilesystemStore) validateManagedPath(name, path string) error { if s == nil { return fmt.Errorf("state store is required") } root, err := filepath.Abs(s.root) if err != nil { return fmt.Errorf("resolve workspace root: %w", err) } target, err := filepath.Abs(path) if err != nil { return fmt.Errorf("resolve %s: %w", name, err) } relative, err := filepath.Rel(root, target) if err != nil { return fmt.Errorf("resolve %s relative to workspace root: %w", name, err) } if relative == "." || relative == ".." || strings.HasPrefix(relative, ".."+string(filepath.Separator)) { return fmt.Errorf("%s must stay within workspace root", name) } return nil } 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 { return resolved.Definition.ComparisonStrategy == report.CompareSameValidDate && sameValidDate(metadata, resolved) }