528 lines
17 KiB
Go
528 lines
17 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 {
|
|
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
|
|
}
|
|
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 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)
|
|
}
|