1004 lines
32 KiB
Go
1004 lines
32 KiB
Go
// Package app owns application orchestration and top-level use cases.
|
|
package app
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"path/filepath"
|
|
"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/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
|
|
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
|
|
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
|
|
}
|
|
debugWriter, err := state.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
|
|
}
|
|
debugWriter, err := state.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 := plannedBatchOutputPath(req.OutputDir, planned)
|
|
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 {
|
|
if outputDir == "" {
|
|
return ""
|
|
}
|
|
outputCopyName := planned.OutputCopyName
|
|
if outputCopyName == "" {
|
|
outputCopyName = planned.Resolved.Definition.BatchOutputName
|
|
}
|
|
if outputCopyName == "" {
|
|
return ""
|
|
}
|
|
return filepath.Join(outputDir, outputCopyName)
|
|
}
|
|
|
|
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.ManagedReportPath, 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 managed report %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, resolved.Definition.BatchOutputName)
|
|
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, batchOutputName string) (config.DistributorTemplateValues, error) {
|
|
values := config.DistributorTemplateValues{
|
|
LocationID: cfg.Location.ID,
|
|
ReportID: string(resolved.Definition.ID),
|
|
RunID: runID,
|
|
ArtifactGroup: resolved.Definition.ArtifactGroup,
|
|
BatchOutputName: batchOutputName,
|
|
}
|
|
if values.BatchOutputName == "" {
|
|
values.BatchOutputName = resolved.Definition.BatchOutputName
|
|
}
|
|
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)
|
|
}
|