// Package app owns application orchestration and top-level use cases. package app import ( "context" "fmt" "os" "path/filepath" "strings" "time" distributoradapter "gitea.maximumdirect.net/eric/weatherreporter/internal/adapters/distributor" "gitea.maximumdirect.net/eric/weatherreporter/internal/briefing" "gitea.maximumdirect.net/eric/weatherreporter/internal/collect" "gitea.maximumdirect.net/eric/weatherreporter/internal/config" "gitea.maximumdirect.net/eric/weatherreporter/internal/facts" "gitea.maximumdirect.net/eric/weatherreporter/internal/fileutil" "gitea.maximumdirect.net/eric/weatherreporter/internal/forecast" "gitea.maximumdirect.net/eric/weatherreporter/internal/module" "gitea.maximumdirect.net/eric/weatherreporter/internal/promptdebug" "gitea.maximumdirect.net/eric/weatherreporter/internal/promptexec" "gitea.maximumdirect.net/eric/weatherreporter/internal/promptinput" "gitea.maximumdirect.net/eric/weatherreporter/internal/report" "gitea.maximumdirect.net/eric/weatherreporter/internal/state" "gitea.maximumdirect.net/eric/weatherreporter/internal/timeutil" "gitea.maximumdirect.net/eric/weatherreporter/internal/weatherdata" ) type ReportKind string const ( ReportDaily ReportKind = ReportKind(report.CommandNameDaily) ReportToday ReportKind = ReportKind(report.CommandNameToday) ReportTomorrow ReportKind = ReportKind(report.CommandNameTomorrow) ReportHourly ReportKind = ReportKind(report.CommandNameHourly) ) type BatchKind string const ( BatchMorning BatchKind = BatchKind(report.BatchNameMorning) BatchEvening BatchKind = BatchKind(report.BatchNameEvening) ) type GenerateRequest struct { Config config.Config Report ReportKind WorkingDir string OutputPath string LLMDebugDir string Now time.Time Date time.Time Collector Collector Notifier Notifier Executor promptexec.Executor Store state.Store } type BatchRequest struct { Config config.Config Batch BatchKind Now time.Time WorkingDir string OutputDir string LLMDebugDir string Collector Collector Executor promptexec.Executor Store state.Store Notifier Notifier } type FetchBundleRequest struct { Config config.Config OutputPath string } type ModuleSnapshotRequest struct { Config config.Config Resolved report.Resolved } type ReportFacts struct { Collected facts.CollectedFacts Derived facts.DerivedFacts } type ReportResult struct { ModuleSnapshot module.Snapshot ModuleSnapshotPath string DataPackage promptinput.Package DataPackagePath string PreparationPath string ExecutionPath string LLMDebugPath string ReportPath string OutputPath string NotificationPath string Metadata state.Metadata MetadataPath string GeneratedTextRawPath string GeneratedTextPath string RenderContextPath string Notification *NotificationResult } type BatchResult struct { Batch BatchKind `json:"batch"` StartedAt time.Time `json:"startedAt"` FinishedAt time.Time `json:"finishedAt"` Total int `json:"total"` Succeeded int `json:"succeeded"` Failed int `json:"failed"` Notification *BatchNotificationResult `json:"notification,omitempty"` Reports []BatchReportResult `json:"reports"` } type BatchNotificationResult struct { Status string `json:"status"` Reason string `json:"reason,omitempty"` RunID string `json:"runId,omitempty"` PipelineID string `json:"pipelineId,omitempty"` BundleID string `json:"bundleId,omitempty"` IdempotencyKey string `json:"idempotencyKey,omitempty"` Path string `json:"path,omitempty"` IncludedReports []BatchNotificationReport `json:"includedReports,omitempty"` Error string `json:"error,omitempty"` } type BatchNotificationReport struct { ReportID report.ID `json:"reportId"` RunID string `json:"runId"` SourcePath string `json:"sourcePath"` BundlePaths []string `json:"bundlePaths"` } type BatchReportResult struct { ReportID report.ID `json:"reportId"` ReportName string `json:"reportName"` PromptID string `json:"promptId"` RunID string `json:"runId"` Status string `json:"status"` Error string `json:"error,omitempty"` NotificationStatus string `json:"notificationStatus,omitempty"` NotificationRunID string `json:"notificationRunId,omitempty"` NotificationPipelineID string `json:"notificationPipelineId,omitempty"` NotificationError string `json:"notificationError,omitempty"` NotificationPath string `json:"notificationPath,omitempty"` GeneratedAt time.Time `json:"generatedAt"` ValidPeriod timeutil.Period `json:"validPeriod"` DataPackagePath string `json:"dataPackagePath,omitempty"` PreparationPath string `json:"preparationPath,omitempty"` ExecutionPath string `json:"executionPath,omitempty"` LLMDebugPath string `json:"llmDebugPath,omitempty"` ReportPath string `json:"reportPath,omitempty"` OutputPath string `json:"outputPath,omitempty"` MetadataPath string `json:"metadataPath,omitempty"` } type BatchError struct { Result *BatchResult } func (e BatchError) Error() string { if e.Result == nil { return "batch failed" } if batchNotificationFailed(e.Result) && batchReportFailures(e.Result) == 0 { if e.Result.Notification.Error != "" { return fmt.Sprintf("batch %s notification failed: %s", e.Result.Batch, e.Result.Notification.Error) } return fmt.Sprintf("batch %s notification failed", e.Result.Batch) } return fmt.Sprintf("batch %s failed: %d of %d reports failed", e.Result.Batch, e.Result.Failed, e.Result.Total) } func batchNotificationFailed(result *BatchResult) bool { return result != nil && result.Notification != nil && result.Notification.Status == "failed" } func batchReportFailures(result *BatchResult) int { if result == nil { return 0 } failures := 0 for _, item := range result.Reports { if item.Status == "failed" { failures++ } } return failures } type Collector interface { Run(context.Context, collect.Request) (*collect.Result, error) } type defaultCollector struct{} func (defaultCollector) Run(ctx context.Context, req collect.Request) (*collect.Result, error) { return collect.Run(ctx, req) } type Notifier interface { Notify(context.Context, NotificationRequest) (*NotificationResult, error) } type NotificationRequest struct { ReportID report.ID RunID string PipelineID string BundleID string IdempotencyKey string ReportPath string BundlePaths []string CreatedAt time.Time } type NotificationResult struct { BundleID string IdempotencyKey string RunID string Status string UploadStatus string StatusError string PipelineID string AcceptedAt time.Time StartedAt *time.Time FinishedAt *time.Time Report []byte Error string } type NotificationError struct { Request NotificationRequest Err error } func (e *NotificationError) Error() string { if e == nil || e.Err == nil { return "notification failed" } return e.Err.Error() } func (e *NotificationError) Unwrap() error { if e == nil { return nil } return e.Err } func Generate(ctx context.Context, req GenerateRequest) error { _, err := GenerateDetailed(ctx, req) return err } func GenerateDetailed(ctx context.Context, req GenerateRequest) (*ReportResult, error) { now := req.Now if now.IsZero() { now = time.Now() } resolved, err := ResolveGenerate(req, now) if err != nil { return nil, err } outputPath, err := resolveReportOutputPath(req.WorkingDir, req.OutputPath, resolved) if err != nil { return nil, err } req.OutputPath = outputPath debugWriter, err := promptdebug.NewPromptDebugWriter(req.LLMDebugDir) if err != nil { return nil, promptexec.NewError(promptexec.InvalidConfiguration, "initialize prompt debug", err) } inspection, err := InspectPromptExecution(ctx, PromptInspectionRequest{ Resolved: resolved, Executor: req.Executor, Promptkit: req.Config.Promptkit, }) if err != nil { return nil, err } collection, err := collectWeather(ctx, req.Config, req.Collector) if err != nil { return nil, err } return generatePromptReport(ctx, promptReportRequest{ GenerateRequest: req, Resolved: resolved, Collection: *collection, Inspection: inspection, DebugWriter: debugWriter, }) } func RunBatch(ctx context.Context, req BatchRequest) error { result, err := RunBatchDetailed(ctx, req) if err != nil { return err } if result.Failed > 0 { return BatchError{Result: result} } return nil } func RunBatchDetailed(ctx context.Context, req BatchRequest) (*BatchResult, error) { now := req.Now if now.IsZero() { now = time.Now() } if _, err := report.BatchForCommandName(string(req.Batch)); err != nil { return nil, err } outputDir, err := resolveOutputDir(req.WorkingDir, req.OutputDir) if err != nil { return nil, err } req.OutputDir = outputDir debugWriter, err := promptdebug.NewPromptDebugWriter(req.LLMDebugDir) if err != nil { return nil, promptexec.NewError(promptexec.InvalidConfiguration, "initialize prompt debug", err) } candidates, err := batchInspectionCandidates(req, now) if err != nil { return nil, err } inspections, err := InspectPromptExecutions(ctx, PromptExecutionsInspectionRequest{ Resolved: candidates, Executor: req.Executor, Promptkit: req.Config.Promptkit, }) if err != nil { return nil, err } collection, err := collectWeather(ctx, req.Config, req.Collector) if err != nil { return nil, err } plannedReports, err := planBatchRun(req, now, *collection) if err != nil { return nil, err } if req.Batch == BatchEvening || req.Batch == BatchMorning { store := req.Store if store == nil { defaultStore, err := defaultStore(req.Config) if err != nil { return nil, err } store = defaultStore } startedAt := now result := &BatchResult{Batch: req.Batch, StartedAt: startedAt} for _, planned := range plannedReports { resolved := planned.Resolved item := batchReportResult(planned) outputPath, err := plannedBatchOutputPath(req.OutputDir, planned) if err != nil { return nil, err } reportResult, err := generatePromptReport(ctx, promptReportRequest{ GenerateRequest: GenerateRequest{ Config: req.Config, OutputPath: outputPath, Notifier: req.Notifier, Executor: req.Executor, Store: store, }, Resolved: resolved, Collection: *collection, Inspection: inspections[resolved.Definition.ID], DebugWriter: debugWriter, noNotify: true, }) if reportResult != nil { copyBatchReportPaths(&item, reportResult) } if err != nil { item.Status = "failed" item.Error = err.Error() result.Failed++ } else { item.Status = "succeeded" result.Succeeded++ } result.Reports = append(result.Reports, item) } result.Total = len(result.Reports) batchNotification, err := notifyBatch(ctx, req.Config, req.Batch, batchRunID(startedAt, req.Batch), startedAt, result, plannedReports, store, req.Notifier) if batchNotification != nil { result.Notification = batchNotification } if err != nil { result.Failed++ } result.FinishedAt = time.Now() return result, nil } return nil, fmt.Errorf("run is not implemented") } func copyBatchReportPaths(item *BatchReportResult, result *ReportResult) { item.DataPackagePath = result.DataPackagePath item.PreparationPath = result.PreparationPath item.ExecutionPath = result.ExecutionPath item.LLMDebugPath = result.LLMDebugPath item.ReportPath = result.ReportPath item.OutputPath = result.OutputPath item.MetadataPath = result.MetadataPath item.NotificationPath = result.NotificationPath if result.Notification != nil { item.NotificationStatus = result.Notification.Status item.NotificationRunID = result.Notification.RunID item.NotificationPipelineID = result.Notification.PipelineID } } func batchInspectionCandidates(req BatchRequest, now time.Time) ([]report.Resolved, error) { location, err := timeutil.LoadLocation(req.Config.WeatherAPI.Timezone) if err != nil { return nil, err } registry, err := reportRegistry(req.Config) if err != nil { return nil, err } ids := []report.ID{report.Tomorrow, report.Daily} if req.Batch == BatchMorning { ids = []report.ID{report.Today, report.Tomorrow, report.Daily} } date := timeutil.LocalDate(now, location).AddDate(0, 0, 2) candidates := make([]report.Resolved, 0, len(ids)) for _, id := range ids { resolveReq := report.ResolveRequest{Now: now, Location: location} if id == report.Daily { resolveReq.Date = date } resolved, err := registry.Resolve(id, resolveReq) if err != nil { return nil, err } candidates = append(candidates, resolved) } return candidates, nil } func batchReportResult(planned plannedBatchReport) BatchReportResult { resolved := planned.Resolved metadata := resolved.Metadata() return BatchReportResult{ ReportID: resolved.Definition.ID, ReportName: resolved.Definition.Name, PromptID: resolved.Definition.PromptID, RunID: metadata.RunID, GeneratedAt: metadata.GeneratedAt, ValidPeriod: metadata.ValidPeriod, } } func plannedBatchOutputPath(outputDir string, planned plannedBatchReport) (string, error) { outputName, err := planned.Resolved.OutputName() if err != nil { return "", err } return validateOutputPath(filepath.Join(outputDir, outputName)) } func resolveReportOutputPath(workingDir, override string, resolved report.Resolved) (string, error) { outputName, err := resolved.OutputName() if err != nil { return "", err } return resolveOutputPath(workingDir, override, outputName) } func resolveOutputDir(workingDir, override string) (string, error) { workingDir, err := validateWorkingDir(workingDir) if err != nil { return "", err } if override == "" { return workingDir, nil } if strings.TrimSpace(override) == "" { return "", fmt.Errorf("output directory is required") } directory := override if !filepath.IsAbs(directory) { directory = filepath.Join(workingDir, directory) } directory = filepath.Clean(directory) if info, err := os.Stat(directory); err == nil && !info.IsDir() { return "", fmt.Errorf("output directory %q is not a directory", directory) } else if err != nil && !os.IsNotExist(err) { return "", fmt.Errorf("inspect output directory %q: %w", directory, err) } return directory, nil } func resolveOutputPath(workingDir, override, defaultName string) (string, error) { workingDir, err := validateWorkingDir(workingDir) if err != nil { return "", err } path := override if path == "" { path = defaultName } if strings.TrimSpace(path) == "" { return "", fmt.Errorf("final output path is required") } if !filepath.IsAbs(path) { path = filepath.Join(workingDir, path) } return validateOutputPath(path) } func validateWorkingDir(workingDir string) (string, error) { if strings.TrimSpace(workingDir) == "" { return "", fmt.Errorf("working directory is required") } if !filepath.IsAbs(workingDir) { return "", fmt.Errorf("working directory %q must be absolute", workingDir) } return filepath.Clean(workingDir), nil } func validateOutputPath(path string) (string, error) { if strings.TrimSpace(path) == "" { return "", fmt.Errorf("final output path is required") } path = filepath.Clean(path) if !filepath.IsAbs(path) { return "", fmt.Errorf("final output path %q must be absolute", path) } if filepath.Dir(path) == path { return "", fmt.Errorf("final output path %q must not be a filesystem root", path) } if info, err := os.Stat(path); err == nil && info.IsDir() { return "", fmt.Errorf("final output path %q is a directory", path) } else if err != nil && !os.IsNotExist(err) { return "", fmt.Errorf("inspect final output path %q: %w", path, err) } return path, nil } func ResolveGenerate(req GenerateRequest, now time.Time) (report.Resolved, error) { location, err := timeutil.LoadLocation(req.Config.WeatherAPI.Timezone) if err != nil { return report.Resolved{}, err } id, err := report.IDForCommandName(string(req.Report)) if err != nil { return report.Resolved{}, err } registry, err := reportRegistry(req.Config) if err != nil { return report.Resolved{}, err } return registry.Resolve(id, report.ResolveRequest{ Now: now, Location: location, Date: req.Date, }) } func reportRegistry(cfg config.Config) (report.Registry, error) { overrides, err := cfg.ReportModuleOverrides() if err != nil { return report.Registry{}, err } registry, err := report.DefaultRegistry().WithModuleOverrides(overrides) if err != nil { return report.Registry{}, err } return registry, nil } func FetchBundle(ctx context.Context, req FetchBundleRequest) (*weatherdata.Bundle, error) { result, err := collectWeather(ctx, req.Config, nil) if err != nil { return nil, err } return result.Bundle, nil } func collectWeather(ctx context.Context, cfg config.Config, collector Collector) (*collect.Result, error) { if collector == nil { collector = defaultCollector{} } result, err := collector.Run(ctx, collect.Request{Config: cfg}) if err != nil { return nil, err } if result == nil { return nil, fmt.Errorf("collect weather bundle: collector returned nil result") } if result.Bundle == nil { return nil, fmt.Errorf("collect weather bundle: collector returned nil bundle") } return result, nil } func FetchAndSaveBundle(ctx context.Context, req FetchBundleRequest) (*weatherdata.Bundle, error) { if req.OutputPath == "" { return nil, fmt.Errorf("output path is required") } bundle, err := FetchBundle(ctx, req) if err != nil { return nil, err } if err := fileutil.WriteJSONAtomic(req.OutputPath, bundle); err != nil { return nil, fmt.Errorf("save bundle: %w", err) } return bundle, nil } type finalizeRenderedReportRequest struct { Config config.Config Store state.Store Resolved report.Resolved Metadata state.Metadata MetadataPath string ExecutionArtifact *state.PromptExecutionArtifact ManagedReportPath string OutputPath string Notifier Notifier GenerationErr error noNotify bool } type finalizeRenderedReportResult struct { OutputPath string NotificationPath string Metadata state.Metadata MetadataPath string Notification *NotificationResult } func finalizeRenderedReport(ctx context.Context, req finalizeRenderedReportRequest) (finalizeRenderedReportResult, error) { if req.Store == nil { return finalizeRenderedReportResult{}, fmt.Errorf("state store is required") } if req.ManagedReportPath == "" { return finalizeRenderedReportResult{}, fmt.Errorf("managed report path is required for report %q", req.Resolved.Definition.ID) } if req.ExecutionArtifact == nil { return finalizeRenderedReportResult{}, fmt.Errorf("prompt execution artifact is required for report %q", req.Resolved.Definition.ID) } result := finalizeRenderedReportResult{Metadata: req.Metadata, MetadataPath: req.MetadataPath} if req.OutputPath != "" && req.GenerationErr == nil { if req.OutputPath != req.ManagedReportPath { if err := fileutil.CopyFileAtomic(req.ManagedReportPath, req.OutputPath); err != nil { return result, err } result.OutputPath = req.OutputPath if err := persistReachedPromptPath(ctx, req.Store, req.Resolved, req.ExecutionArtifact, func(paths *state.PromptExecutionPaths) { paths.OutputPath = req.OutputPath }); err != nil { return result, err } } else { result.OutputPath = req.OutputPath } } metadata := req.Metadata metadata.RenderedReportPath = req.ManagedReportPath metadataPath, err := req.Store.SaveMetadata(ctx, metadata) if err != nil { return result, err } result.Metadata = metadata result.MetadataPath = metadataPath if req.GenerationErr != nil { return result, req.GenerationErr } if req.noNotify { return result, nil } notification, notificationPath, err := notifyReport(ctx, req.Config, req.Resolved, req.OutputPath, metadata, req.Notifier, req.Store) if notificationPath != "" { result.NotificationPath = notificationPath result.Notification = notification metadata.NotificationPath = notificationPath result.Metadata = metadata if saveErr := persistReachedPromptPath(ctx, req.Store, req.Resolved, req.ExecutionArtifact, func(paths *state.PromptExecutionPaths) { paths.NotificationPath = notificationPath }); saveErr != nil { return result, saveErr } metadataPath, saveErr := req.Store.SaveMetadata(ctx, metadata) if saveErr != nil { return result, saveErr } result.Metadata = metadata result.MetadataPath = metadataPath } result.Notification = notification if err != nil { return result, err } return result, nil } func notifyReport(ctx context.Context, cfg config.Config, resolved report.Resolved, reportPath string, metadata state.Metadata, notifier Notifier, store state.Store) (*NotificationResult, string, error) { notifier, enabled := reportNotifier(cfg, notifier) if !enabled { return nil, "", nil } notificationRequest, err := buildNotificationRequest(cfg, resolved, reportPath, metadata) if err != nil { notificationPath, saveErr := saveNotificationArtifact(ctx, store, resolved, cfg, metadata, NotificationRequest{}, nil, err) if saveErr != nil { return nil, "", saveErr } return nil, notificationPath, err } result, err := notifier.Notify(ctx, notificationRequest) notificationPath, saveErr := saveNotificationArtifact(ctx, store, resolved, cfg, metadata, notificationRequest, result, err) if saveErr != nil { return nil, "", saveErr } if err != nil { return result, notificationPath, &NotificationError{ Request: notificationRequest, Err: fmt.Errorf("notify report %q run %q from output %q: %w", resolved.Definition.ID, metadata.RunID, reportPath, err), } } return result, notificationPath, nil } func reportNotifier(cfg config.Config, notifier Notifier) (Notifier, bool) { if !cfg.Notify.Distributor.Enabled { return noopNotifier{}, false } if notifier != nil { return notifier, true } return distributorNotifier{ client: distributoradapter.New(cfg.Notify.Distributor), }, true } func buildNotificationRequest(cfg config.Config, resolved report.Resolved, reportPath string, metadata state.Metadata) (NotificationRequest, error) { values, err := distributorTemplateValuesForReport(cfg, resolved, metadata.RunID, filepath.Base(reportPath)) if err != nil { return NotificationRequest{}, err } bundleID, err := config.RenderDistributorBundleID(cfg.Notify.Distributor.BundleIDTemplate, values) if err != nil { return NotificationRequest{}, err } values.BundleID = bundleID pipelineID, err := config.RenderDistributorPipelineID(cfg.Notify.Distributor.PipelineIDTemplate, values) if err != nil { return NotificationRequest{}, err } idempotencyKey, err := config.RenderDistributorIdempotencyKey(cfg.Notify.Distributor.IdempotencyKeyTemplate, values) if err != nil { return NotificationRequest{}, err } bundlePaths, err := renderDistributorReportBundlePaths(cfg, resolved, metadata.RunID, reportPath, values) if err != nil { return NotificationRequest{}, err } return NotificationRequest{ ReportID: resolved.Definition.ID, RunID: metadata.RunID, PipelineID: pipelineID, BundleID: bundleID, IdempotencyKey: idempotencyKey, ReportPath: reportPath, BundlePaths: bundlePaths, CreatedAt: metadata.GeneratedAt, }, nil } func distributorTemplateValuesForReport(cfg config.Config, resolved report.Resolved, runID string, outputName string) (config.DistributorTemplateValues, error) { values := config.DistributorTemplateValues{ LocationID: cfg.Location.ID, ReportID: string(resolved.Definition.ID), RunID: runID, ArtifactGroup: resolved.Definition.ArtifactGroup, BatchOutputName: outputName, } if values.BatchOutputName == "" { var err error values.BatchOutputName, err = resolved.OutputName() if err != nil { return config.DistributorTemplateValues{}, err } } if err := addDistributorValidPeriodValues(&values, resolved.ValidPeriod, cfg.WeatherAPI.Timezone); err != nil { return config.DistributorTemplateValues{}, err } return values, nil } func renderDistributorReportBundlePaths(cfg config.Config, resolved report.Resolved, runID string, sourcePath string, values config.DistributorTemplateValues) ([]string, error) { templates, name, err := distributorPathTemplatesForReport(cfg, resolved.Definition) if err != nil { return nil, distributorReportPathError(resolved.Definition.ID, runID, sourcePath, err) } paths, err := config.RenderDistributorReportPaths(name, templates, values) if err != nil { return nil, distributorReportPathError(resolved.Definition.ID, runID, sourcePath, err) } return paths, nil } func distributorPathTemplatesForReport(cfg config.Config, definition report.Definition) ([]string, string, error) { overrides, err := cfg.ReportDistributorPathOverrides() if err != nil { return nil, "", err } if templates, ok := overrides[definition.ID]; ok { return append([]string(nil), templates...), fmt.Sprintf("reports.%s.distributor.path_templates", definition.ID), nil } if len(definition.DistributorPathTemplates) > 0 { return append([]string(nil), definition.DistributorPathTemplates...), fmt.Sprintf("report.%s.distributor_path_templates", definition.ID), nil } return nil, "", fmt.Errorf("no distributor path templates configured") } func distributorReportPathError(id report.ID, runID string, sourcePath string, err error) error { if sourcePath != "" { return fmt.Errorf("report %q run %q source path %q: %w", id, runID, sourcePath, err) } return fmt.Errorf("report %q run %q: %w", id, runID, err) } func addDistributorValidPeriodValues(values *config.DistributorTemplateValues, period timeutil.Period, timezone string) error { location, err := timeutil.LoadLocation(timezone) if err != nil { return err } start := period.Start.In(location) end := period.End.In(location) values.ValidStartDate = start.Format(timeutil.DateLayout) values.ValidEndDate = end.Format(timeutil.DateLayout) values.ValidStartTime = start.Format("1504") values.ValidEndTime = end.Format("1504") values.ValidStartStamp = start.Format("2006-01-02T1504") values.ValidEndStamp = end.Format("2006-01-02T1504") return nil } func saveNotificationArtifact(ctx context.Context, store state.Store, resolved report.Resolved, cfg config.Config, metadata state.Metadata, req NotificationRequest, result *NotificationResult, notifyErr error) (string, error) { if store == nil { return "", fmt.Errorf("state store is required") } artifact := state.DistributorNotificationArtifact{ SchemaVersion: state.DistributorNotificationSchemaVersion, RunID: metadata.RunID, ReportID: resolved.Definition.ID, AttemptedAt: time.Now(), Endpoint: cfg.Notify.Distributor.Endpoint, PipelineID: req.PipelineID, BundleID: req.BundleID, IdempotencyKey: req.IdempotencyKey, SourcePath: req.ReportPath, BundlePaths: append([]string(nil), req.BundlePaths...), BundleCreated: req.CreatedAt, Status: "attempted", } if result != nil { artifact.Status = result.Status artifact.Upload = &state.DistributorUploadResult{ RunID: result.RunID, Status: result.UploadStatus, } if result.PipelineID != "" || !result.AcceptedAt.IsZero() || result.StartedAt != nil || result.FinishedAt != nil || len(result.Report) > 0 || result.Error != "" { artifact.RunStatus = &state.DistributorRunStatus{ RunID: result.RunID, PipelineID: result.PipelineID, Status: result.Status, AcceptedAt: result.AcceptedAt, StartedAt: result.StartedAt, FinishedAt: result.FinishedAt, Report: append([]byte(nil), result.Report...), Error: result.Error, } } artifact.StatusError = result.StatusError } if notifyErr != nil { artifact.Status = "failed" artifact.Error = notifyErr.Error() } if artifact.Status == "" { artifact.Status = "unknown" } return store.SaveDistributorNotification(ctx, resolved, artifact) } type noopNotifier struct{} func (noopNotifier) Notify(context.Context, NotificationRequest) (*NotificationResult, error) { return nil, nil } type distributorNotifier struct { client *distributoradapter.Client } func (n distributorNotifier) Notify(ctx context.Context, req NotificationRequest) (*NotificationResult, error) { result, err := n.client.Upload(ctx, distributoradapter.UploadRequest{ PipelineID: req.PipelineID, BundleID: req.BundleID, IdempotencyKey: req.IdempotencyKey, Files: distributorUploadFiles(req.ReportPath, req.BundlePaths), CreatedAt: req.CreatedAt, }) notification := &NotificationResult{ PipelineID: req.PipelineID, BundleID: req.BundleID, IdempotencyKey: req.IdempotencyKey, RunID: result.RunID, Status: result.Status, UploadStatus: result.UploadStatus, StatusError: result.StatusError, } if result.RunStatus != nil { if result.RunStatus.PipelineID != "" { notification.PipelineID = result.RunStatus.PipelineID } notification.AcceptedAt = result.RunStatus.AcceptedAt notification.StartedAt = result.RunStatus.StartedAt notification.FinishedAt = result.RunStatus.FinishedAt notification.Report = append([]byte(nil), result.RunStatus.Report...) notification.Error = result.RunStatus.Error } if err != nil { return notification, err } return notification, nil } func (n distributorNotifier) NotifyBatch(ctx context.Context, req batchNotificationRequest) (*NotificationResult, error) { result, err := n.client.Upload(ctx, batchDistributorUploadRequest(req)) notification := notificationResultFromUpload(req.PipelineID, req.BundleID, req.IdempotencyKey, result) if err != nil { return notification, err } return notification, nil } func notificationResultFromUpload(pipelineID string, bundleID string, idempotencyKey string, result distributoradapter.UploadResult) *NotificationResult { notification := &NotificationResult{ PipelineID: pipelineID, BundleID: bundleID, IdempotencyKey: idempotencyKey, RunID: result.RunID, Status: result.Status, UploadStatus: result.UploadStatus, StatusError: result.StatusError, } if result.RunStatus != nil { if result.RunStatus.PipelineID != "" { notification.PipelineID = result.RunStatus.PipelineID } notification.AcceptedAt = result.RunStatus.AcceptedAt notification.StartedAt = result.RunStatus.StartedAt notification.FinishedAt = result.RunStatus.FinishedAt notification.Report = append([]byte(nil), result.RunStatus.Report...) notification.Error = result.RunStatus.Error } return notification } func distributorUploadFiles(sourcePath string, bundlePaths []string) []distributoradapter.UploadFile { files := make([]distributoradapter.UploadFile, 0, len(bundlePaths)) for _, bundlePath := range bundlePaths { files = append(files, distributoradapter.UploadFile{ SourcePath: sourcePath, BundlePath: bundlePath, }) } return files } func BuildModuleSnapshot(req ModuleSnapshotRequest, bundle *weatherdata.Bundle) (module.Snapshot, error) { reportFacts, err := BuildReportFacts(req, bundle) if err != nil { return module.Snapshot{}, err } return BuildModuleSnapshotFromFacts(req, reportFacts) } func BuildReportFacts(req ModuleSnapshotRequest, bundle *weatherdata.Bundle) (ReportFacts, error) { collected := facts.BuildCollected(bundle) derived, err := buildDerivedFacts(req.Config, req.Resolved, collected) if err != nil { return ReportFacts{}, err } return ReportFacts{ Collected: collected, Derived: derived, }, nil } func BuildModuleSnapshotFromFacts(req ModuleSnapshotRequest, reportFacts ReportFacts) (module.Snapshot, error) { if !req.Resolved.ValidPeriod.IsValid() { return module.Snapshot{}, fmt.Errorf("resolved valid period is required") } registry, err := briefing.DefaultModuleRegistry() if err != nil { return module.Snapshot{}, err } moduleContext := briefing.ModuleContext{ Resolved: req.Resolved, Collected: reportFacts.Collected, Derived: reportFacts.Derived, Units: req.Config.WeatherAPI.Units, Timezone: req.Config.WeatherAPI.Timezone, Location: briefingLocation(req.Config), } var outputs []module.Output for _, item := range req.Resolved.Definition.Modules { output, err := registry.BuildModule(moduleContext, item) if err != nil { return module.Snapshot{}, err } if output == nil { continue } outputs = append(outputs, *output) } return module.NewSnapshot(outputs) } func briefingBuildContext(cfg config.Config, resolved report.Resolved, collected facts.CollectedFacts) briefing.BuildContext { return briefing.BuildContext{ Resolved: resolved, Bundle: collected.Bundle(), Units: cfg.WeatherAPI.Units, Timezone: cfg.WeatherAPI.Timezone, Location: briefingLocation(cfg), } } func promptMetadata(metadata state.Metadata) promptinput.Metadata { return promptinput.Metadata{ RunID: metadata.RunID, ReportID: metadata.ReportID, Variant: metadata.Variant, PromptID: metadata.PromptID, GeneratedAt: metadata.GeneratedAt, Timezone: metadata.Timezone, ValidPeriod: metadata.ValidPeriod, SourceWarnings: metadata.SourceWarnings, } } func buildDerivedFacts(cfg config.Config, resolved report.Resolved, collected facts.CollectedFacts) (facts.DerivedFacts, error) { dayparts := make([]forecast.DaypartDefinition, 0, len(cfg.Dayparts)) for _, daypart := range cfg.Dayparts { dayparts = append(dayparts, forecast.DaypartDefinition{ Name: daypart.Name, Start: daypart.Start, End: daypart.End, }) } return facts.BuildDerived(facts.BuildDerivedRequest{ Resolved: resolved, Timezone: cfg.WeatherAPI.Timezone, Dayparts: dayparts, Collected: collected, }) } func briefingLocation(cfg config.Config) *briefing.LocationContext { location := briefing.LocationContext{ ID: cfg.Location.ID, Name: cfg.Location.Name, Region: cfg.Location.Region, Timezone: cfg.WeatherAPI.Timezone, } if location.ID == "" && location.Name == "" && location.Region == "" && location.Timezone == "" { return nil } return &location } func defaultStore(cfg config.Config) (*state.FilesystemStore, error) { return state.NewFilesystemStore(cfg.Workspace) } func generatedReportError(resolved report.Resolved, runID string, operation string, err error) error { if err == nil { return nil } return fmt.Errorf("generate report %q run %q: %s: %w", resolved.Definition.ID, runID, operation, err) }