refactor for separation of concern for better readabiliy and maintainability

This commit is contained in:
Robin Olsen
2026-03-03 11:41:07 +01:00
parent e8d9f9d1b6
commit 763f079063
11 changed files with 530 additions and 471 deletions

39
pkg/task/lifecsycle.go Normal file
View File

@@ -0,0 +1,39 @@
package task
import (
"log/slog"
"github.com/LazyBachelor/LazyPM/internal/models"
)
type RunLifecycle struct {
collector *taskRunCollector
config Config
details models.TaskDetails
app *App
logger *slog.Logger
metricsStore MetricsStore
}
func NewRunLifecycle(app *App, config Config, details models.TaskDetails, iType InterfaceType, logger *slog.Logger) *RunLifecycle {
collector := newTaskRunCollector(details.Title, iType, logger)
config = config.WithActionLogger(func(action string) {
collector.recordUserAction(action)
})
var store MetricsStore
if config.StatisticsStoragePath != "" {
store = NewFileMetricsStore(config.StatisticsStoragePath, logger)
}
return &RunLifecycle{
collector: collector,
config: config,
details: details,
app: app,
logger: logger,
metricsStore: store,
}
}

View File

@@ -1,11 +1,8 @@
package task
import (
"encoding/json"
"fmt"
"log/slog"
"os"
"path/filepath"
"regexp"
"sort"
"strings"
@@ -35,32 +32,51 @@ func newTaskRunCollector(taskName string, interfaceType InterfaceType, logger *s
}
}
func (c *taskRunCollector) log(level string, message string) {
action := normalizeAction(message)
result := ""
func (c *taskRunCollector) appendLog(entry models.TaskLogEntry) {
entry.Timestamp = time.Now()
c.run.Logs = append(c.run.Logs, entry)
if c.logger == nil {
return
}
attrs := []any{
"task", c.run.TaskName,
"interface", c.run.InterfaceType,
"action", entry.Action,
"result", entry.Result,
}
switch entry.Level {
case "error":
c.logger.Error(entry.Message, attrs...)
case "warn":
c.logger.Warn(entry.Message, attrs...)
default:
c.logger.Info(entry.Message, attrs...)
}
}
func (c *taskRunCollector) log(level, message string) {
result := "ok"
switch level {
case "error":
result = "failed"
case "warn":
result = "warning"
default:
result = "ok"
}
c.mu.Lock()
defer c.mu.Unlock()
c.run.Logs = append(c.run.Logs, models.TaskLogEntry{
Timestamp: time.Now(),
Level: level,
Message: message,
Source: "system",
Action: action,
Result: result,
c.appendLog(models.TaskLogEntry{
Level: level,
Message: message,
Source: "system",
Action: normalizeAction(message),
Result: result,
})
attrs := []any{"action", action, "result", result, "task", c.run.TaskName, "interface", c.run.InterfaceType}
c.logWithLogger(level, message, attrs...)
}
func (c *taskRunCollector) recordUserAction(raw string) {
@@ -69,42 +85,14 @@ func (c *taskRunCollector) recordUserAction(raw string) {
c.mu.Lock()
defer c.mu.Unlock()
c.run.Logs = append(c.run.Logs, models.TaskLogEntry{
Timestamp: time.Now(),
Level: "user_action",
Message: raw,
Source: source,
Action: normalizeAction(actionText),
Target: target,
Result: result,
c.appendLog(models.TaskLogEntry{
Level: "user_action",
Message: raw,
Source: source,
Action: normalizeAction(actionText),
Target: target,
Result: result,
})
c.logWithLogger("info", "user action recorded",
"task", c.run.TaskName,
"interface", c.run.InterfaceType,
"source", source,
"action", normalizeAction(actionText),
"target", target,
"result", result,
)
}
func (c *taskRunCollector) setCompleted(completed bool) {
c.mu.Lock()
defer c.mu.Unlock()
c.run.Completed = completed
}
func (c *taskRunCollector) setError(err error) {
if err == nil {
return
}
c.mu.Lock()
defer c.mu.Unlock()
c.run.Error = err.Error()
}
func (c *taskRunCollector) recordQuestionnaire(completed bool, userQuit bool, answers map[string]any) {
@@ -113,6 +101,7 @@ func (c *taskRunCollector) recordQuestionnaire(completed bool, userQuit bool, an
c.run.QuestionnaireCompleted = completed
c.run.QuestionnaireUserQuit = userQuit
if len(answers) > 0 {
c.run.QuestionnaireAnswers = answers
}
@@ -124,128 +113,123 @@ func (c *taskRunCollector) recordQuestionnaire(completed bool, userQuit bool, an
result = "incomplete"
}
c.run.Logs = append(c.run.Logs, models.TaskLogEntry{
Timestamp: time.Now(),
Level: "questionnaire",
Message: "questionnaire finished",
Source: "system",
Action: "questionnaire_finish",
Result: result,
c.appendLog(models.TaskLogEntry{
Level: "questionnaire",
Message: "questionnaire finished",
Source: "system",
Action: "questionnaire_finish",
Result: result,
})
if len(answers) > 0 {
keys := make([]string, 0, len(answers))
for key := range answers {
keys = append(keys, key)
}
sort.Strings(keys)
for _, key := range keys {
value := answers[key]
if value == nil {
continue
}
valueText := fmt.Sprintf("%v", value)
c.run.Logs = append(c.run.Logs, models.TaskLogEntry{
Timestamp: time.Now(),
Level: "questionnaire",
Message: fmt.Sprintf("questionnaire answer: %s=%s", key, valueText),
Source: "system",
Action: "questionnaire_answer",
Target: key,
Result: valueText,
})
}
if len(answers) == 0 {
return
}
c.logWithLogger("info", "questionnaire recorded",
"task", c.run.TaskName,
"completed", completed,
"user_quit", userQuit,
"answers_count", len(answers),
)
keys := make([]string, 0, len(answers))
for k := range answers {
keys = append(keys, k)
}
sort.Strings(keys)
for _, k := range keys {
v := answers[k]
if v == nil {
continue
}
valueText := fmt.Sprintf("%v", v)
c.appendLog(models.TaskLogEntry{
Level: "questionnaire",
Message: fmt.Sprintf("questionnaire answer: %s=%s", k, valueText),
Source: "system",
Action: "questionnaire_answer",
Target: k,
Result: valueText,
})
}
}
func (c *taskRunCollector) recordValidation(feedback ValidationFeedback) {
c.mu.Lock()
defer c.mu.Unlock()
c.run.ValidationAttempts++
attempt := c.run.ValidationAttempts
result := "failed"
if feedback.Success {
result = "passed"
c.run.ValidationSuccesses++
} else {
c.run.ValidationFailures++
}
c.run.Logs = append(c.run.Logs, models.TaskLogEntry{
Timestamp: time.Now(),
Level: "validation",
Message: fmt.Sprintf("validation attempt %d success=%t", attempt, feedback.Success),
Source: "system",
Action: "validate_attempt",
Result: validationResult(feedback.Success),
Attempt: attempt,
c.appendLog(models.TaskLogEntry{
Level: "validation",
Message: fmt.Sprintf("validation attempt %d", attempt),
Source: "system",
Action: "validate_attempt",
Result: result,
Attempt: attempt,
})
c.logWithLogger("info", "validation attempt",
"task", c.run.TaskName,
"interface", c.run.InterfaceType,
"attempt", attempt,
"result", validationResult(feedback.Success),
)
for _, check := range feedback.Checks {
checkResult := "failed"
if check.Valid {
checkResult = "passed"
c.run.ValidationChecksPassed++
c.run.Logs = append(c.run.Logs, models.TaskLogEntry{
Timestamp: time.Now(),
Level: "validation",
Message: fmt.Sprintf("validation check passed: %s", check.Message),
Source: "system",
Action: "validate_check",
Target: check.Message,
Result: "passed",
Attempt: attempt,
})
c.logWithLogger("info", "validation check",
"task", c.run.TaskName,
"attempt", attempt,
"result", "passed",
"target", check.Message,
)
continue
} else {
c.run.ValidationChecksFailed++
}
c.run.ValidationChecksFailed++
c.run.Logs = append(c.run.Logs, models.TaskLogEntry{
Timestamp: time.Now(),
Level: "validation",
Message: fmt.Sprintf("validation check failed: %s", check.Message),
Source: "system",
Action: "validate_check",
Target: check.Message,
Result: "failed",
Attempt: attempt,
c.appendLog(models.TaskLogEntry{
Level: "validation",
Message: fmt.Sprintf("validation check: %s", check.Message),
Source: "system",
Action: "validate_check",
Target: check.Message,
Result: checkResult,
Attempt: attempt,
})
c.logWithLogger("warn", "validation check",
"task", c.run.TaskName,
"attempt", attempt,
"result", "failed",
"target", check.Message,
)
}
c.run.LastValidationMessage = feedback.Message
c.mu.Unlock()
}
func validationResult(success bool) string {
if success {
return "passed"
func (c *taskRunCollector) setCompleted(completed bool) {
c.mu.Lock()
c.run.Completed = completed
c.mu.Unlock()
}
func (c *taskRunCollector) setError(err error) {
if err == nil {
return
}
return "failed"
c.mu.Lock()
c.run.Error = err.Error()
c.mu.Unlock()
}
func normalizeUserAction(raw string) (source string, actionText string, target string, result string) {
func (c *taskRunCollector) finalize() models.TaskRunMetrics {
c.mu.Lock()
defer c.mu.Unlock()
if c.run.EndedAt.IsZero() {
c.run.EndedAt = time.Now()
}
c.run.DurationMs =
c.run.EndedAt.Sub(c.run.StartedAt).Milliseconds()
final := c.run
final.Logs = append([]models.TaskLogEntry(nil), c.run.Logs...)
return final
}
func normalizeUserAction(raw string) (source, actionText, target, result string) {
trimmed := strings.TrimSpace(raw)
if trimmed == "" {
return "unknown", "unknown_action", "", "unknown"
@@ -256,57 +240,37 @@ func normalizeUserAction(raw string) (source string, actionText string, target s
if source == "" {
source = "unknown"
}
actionText = strings.TrimSpace(event.Action)
if actionText == "" {
actionText = "unknown_action"
}
actionText = stripSourcePrefix(actionText, source)
actionText = stripSourcePrefix(
strings.TrimSpace(event.Action),
source,
)
target = strings.TrimSpace(event.Target)
result = strings.TrimSpace(event.Result)
if result == "" {
result = inferActionResult(actionText)
}
return source, actionText, target, result
return
}
result = "ok"
lower := strings.ToLower(trimmed)
if strings.HasPrefix(lower, "web request:") {
rest := strings.TrimSpace(trimmed[len("web request:"):])
parts := strings.Fields(rest)
if len(parts) >= 2 {
return "web", "request", strings.ToUpper(parts[0]) + " " + parts[1], "ok"
}
return "web", "request", rest, "ok"
}
if strings.HasPrefix(lower, "repl command:") {
command := strings.TrimSpace(trimmed[len("repl command:"):])
return "repl", "run_command", command, "ok"
cmd := strings.TrimSpace(trimmed[len("repl command:"):])
return "repl", "run_command", cmd, "ok"
}
words := strings.Fields(trimmed)
if len(words) > 1 {
head := strings.ToLower(words[0])
if head == "tui" || head == "web" || head == "repl" {
source = head
actionText = strings.Join(words[1:], " ")
} else {
source = "unknown"
actionText = trimmed
}
} else {
source = "unknown"
actionText = trimmed
}
result = inferActionResult(actionText)
return source, actionText, "", result
return "unknown", trimmed, "", inferActionResult(trimmed)
}
func stripSourcePrefix(actionText string, source string) string {
func stripSourcePrefix(actionText, source string) string {
actionText = strings.TrimSpace(actionText)
if actionText == "" || source == "" {
return actionText
@@ -327,22 +291,21 @@ func stripSourcePrefix(actionText string, source string) string {
func inferActionResult(actionText string) string {
lower := strings.ToLower(actionText)
if strings.Contains(lower, "failed") {
switch {
case strings.Contains(lower, "failed"):
return "failed"
}
if strings.Contains(lower, "canceled") {
case strings.Contains(lower, "canceled"):
return "canceled"
}
if strings.Contains(lower, "started") {
case strings.Contains(lower, "started"):
return "started"
}
if strings.Contains(lower, "submitted") {
case strings.Contains(lower, "submitted"):
return "submitted"
}
if strings.Contains(lower, "requested") {
case strings.Contains(lower, "requested"):
return "requested"
default:
return "ok"
}
return "ok"
}
func normalizeAction(input string) string {
@@ -350,153 +313,15 @@ func normalizeAction(input string) string {
if lower == "" {
return "unknown_action"
}
clean := nonWord.ReplaceAllString(lower, "_")
clean = strings.Trim(clean, "_")
clean := strings.Trim(
nonWord.ReplaceAllString(lower, "_"),
"_",
)
if clean == "" {
return "unknown_action"
}
return clean
}
func (c *taskRunCollector) finalize() models.TaskRunMetrics {
c.mu.Lock()
defer c.mu.Unlock()
if c.run.EndedAt.IsZero() {
c.run.EndedAt = time.Now()
}
c.run.DurationMs = c.run.EndedAt.Sub(c.run.StartedAt).Milliseconds()
finalLogs := make([]models.TaskLogEntry, len(c.run.Logs))
copy(finalLogs, c.run.Logs)
final := c.run
final.Logs = finalLogs
return final
}
func appendTaskMetrics(path string, taskName string, run models.TaskRunMetrics, logger *slog.Logger) error {
if path == "" {
return nil
}
dir := filepath.Dir(path)
if dir != "." {
if err := os.MkdirAll(dir, 0o755); err != nil {
return fmt.Errorf("failed to create metrics directory: %w", err)
}
}
metrics := models.TaskMetricsFile{
TaskName: taskName,
Runs: []models.TaskRunMetrics{},
}
if bytes, err := os.ReadFile(path); err == nil {
if len(bytes) > 0 {
if err := json.Unmarshal(bytes, &metrics); err != nil {
return fmt.Errorf("failed to parse metrics file %q: %w", path, err)
}
}
} else if !os.IsNotExist(err) {
return fmt.Errorf("failed to read metrics file %q: %w", path, err)
}
if metrics.TaskName == "" {
metrics.TaskName = taskName
}
run.RunID = len(metrics.Runs) + 1
metrics.Runs = append(metrics.Runs, run)
metrics.Summary = buildTaskStatsSummary(metrics.Runs)
metrics.UpdatedAt = time.Now()
data, err := json.MarshalIndent(metrics, "", " ")
if err != nil {
return fmt.Errorf("failed to encode task metrics: %w", err)
}
if err := os.WriteFile(path, data, 0o644); err != nil {
return fmt.Errorf("failed to write metrics file %q: %w", path, err)
}
if logger != nil {
logger.Info("task metrics persisted",
"path", path,
"task", taskName,
"run_id", run.RunID,
"total_runs", metrics.Summary.TotalRuns,
"completed_runs", metrics.Summary.CompletedRuns,
)
}
return nil
}
func (c *taskRunCollector) logWithLogger(level string, message string, attrs ...any) {
if c.logger == nil {
return
}
switch level {
case "error":
c.logger.Error(message, attrs...)
case "warn":
c.logger.Warn(message, attrs...)
default:
c.logger.Info(message, attrs...)
}
}
func buildTaskStatsSummary(runs []models.TaskRunMetrics) models.TaskStatsSummary {
summary := models.TaskStatsSummary{}
if len(runs) == 0 {
return summary
}
summary.TotalRuns = len(runs)
summary.FirstRunStartedAt = runs[0].StartedAt
for _, run := range runs {
summary.TotalDurationMs += run.DurationMs
summary.ValidationAttempts += run.ValidationAttempts
summary.ValidationSuccesses += run.ValidationSuccesses
summary.ValidationFailures += run.ValidationFailures
summary.ValidationChecksPassed += run.ValidationChecksPassed
summary.ValidationChecksFailed += run.ValidationChecksFailed
if run.Completed {
summary.CompletedRuns++
} else {
summary.IncompleteRuns++
}
if run.QuestionnaireCompleted {
summary.QuestionnairesCompleted++
}
if run.QuestionnaireUserQuit {
summary.QuestionnairesAbandoned++
}
if run.StartedAt.Before(summary.FirstRunStartedAt) {
summary.FirstRunStartedAt = run.StartedAt
}
if run.StartedAt.After(summary.LastRunStartedAt) {
summary.LastRunStartedAt = run.StartedAt
}
if run.EndedAt.After(summary.LastRunEndedAt) {
summary.LastRunEndedAt = run.EndedAt
summary.LastInterfaceType = run.InterfaceType
}
for _, log := range run.Logs {
if log.Level == "user_action" {
summary.TotalUserActions++
}
}
}
summary.AverageDurationMs = summary.TotalDurationMs / int64(summary.TotalRuns)
return summary
}

87
pkg/task/metrics_store.go Normal file
View File

@@ -0,0 +1,87 @@
package task
import (
"context"
"encoding/json"
"fmt"
"log/slog"
"os"
"path/filepath"
"time"
"github.com/LazyBachelor/LazyPM/internal/models"
)
type MetricsStore interface {
Append(ctx context.Context, taskName string, run models.TaskRunMetrics) error
}
type FileMetricsStore struct {
path string
logger *slog.Logger
}
func NewFileMetricsStore(path string, logger *slog.Logger) *FileMetricsStore {
return &FileMetricsStore{
path: path,
logger: logger,
}
}
func (s *FileMetricsStore) Append(ctx context.Context, taskName string, run models.TaskRunMetrics) error {
if s.path == "" {
return nil
}
dir := filepath.Dir(s.path)
if dir != "." {
if err := os.MkdirAll(dir, 0o755); err != nil {
return fmt.Errorf("create metrics directory: %w", err)
}
}
metrics := models.TaskMetricsFile{
TaskName: taskName,
Runs: []models.TaskRunMetrics{},
}
if bytes, err := os.ReadFile(s.path); err == nil {
if len(bytes) > 0 {
if err := json.Unmarshal(bytes, &metrics); err != nil {
return fmt.Errorf("parse metrics file: %w", err)
}
}
} else if !os.IsNotExist(err) {
return fmt.Errorf("read metrics file: %w", err)
}
if metrics.TaskName == "" {
metrics.TaskName = taskName
}
run.RunID = len(metrics.Runs) + 1
metrics.Runs = append(metrics.Runs, run)
metrics.Summary = buildTaskStatsSummary(metrics.Runs)
metrics.UpdatedAt = time.Now()
data, err := json.MarshalIndent(metrics, "", " ")
if err != nil {
return fmt.Errorf("encode metrics: %w", err)
}
if err := os.WriteFile(s.path, data, 0o644); err != nil {
return fmt.Errorf("write metrics file: %w", err)
}
if s.logger != nil {
s.logger.Info(
"task metrics persisted",
"path", s.path,
"task", taskName,
"run_id", run.RunID,
)
}
return nil
}

View File

@@ -0,0 +1,72 @@
package task
import (
"github.com/LazyBachelor/LazyPM/internal/models"
)
func buildTaskStatsSummary(runs []models.TaskRunMetrics) models.TaskStatsSummary {
var summary models.TaskStatsSummary
if len(runs) == 0 {
return summary
}
summary.TotalRuns = len(runs)
for i, run := range runs {
// First run handling
if i == 0 || run.StartedAt.Before(summary.FirstRunStartedAt) {
summary.FirstRunStartedAt = run.StartedAt
}
if run.StartedAt.After(summary.LastRunStartedAt) {
summary.LastRunStartedAt = run.StartedAt
}
if run.EndedAt.After(summary.LastRunEndedAt) {
summary.LastRunEndedAt = run.EndedAt
summary.LastInterfaceType = run.InterfaceType
}
// Duration
summary.TotalDurationMs += run.DurationMs
// Completion
if run.Completed {
summary.CompletedRuns++
} else {
summary.IncompleteRuns++
}
// Validation
summary.ValidationAttempts += run.ValidationAttempts
summary.ValidationSuccesses += run.ValidationSuccesses
summary.ValidationFailures += run.ValidationFailures
summary.ValidationChecksPassed += run.ValidationChecksPassed
summary.ValidationChecksFailed += run.ValidationChecksFailed
// Questionnaire
if run.QuestionnaireCompleted {
summary.QuestionnairesCompleted++
}
if run.QuestionnaireUserQuit {
summary.QuestionnairesAbandoned++
}
// User actions
for _, log := range run.Logs {
if log.Level == "user_action" {
summary.TotalUserActions++
}
}
}
if summary.TotalRuns > 0 {
summary.AverageDurationMs =
summary.TotalDurationMs / int64(summary.TotalRuns)
}
return summary
}

View File

@@ -4,7 +4,6 @@ import (
"context"
"fmt"
"log/slog"
"time"
"github.com/LazyBachelor/LazyPM/internal/models"
tea "github.com/charmbracelet/bubbletea"
@@ -12,12 +11,9 @@ import (
type App = models.App
type Config = models.Config
type Tasker = models.Tasker
type Interface = models.Interface
type InterfaceType = models.InterfaceType
type ValidatedInterface = models.ValidatedInterface
type ValidationFeedback = models.ValidationFeedback
@@ -27,179 +23,174 @@ type QuestionnaireKeysProvider interface {
var ErrUserQuit = models.ErrUserQuit
// RunTask orchestrates the complete task execution flow:
// 1. Setup the task
// 2. Show task intro screen
// 3. Run the interface
// 4. Start validation loop in background
// 5. Show questionnaire when done
func RunTask(ctx context.Context, app *App, t Tasker, i Interface, iType InterfaceType) (runErr error) {
config := t.Config()
details := t.Details()
type TaskRunner struct {
app *App
logger *slog.Logger
}
func NewTaskRunner(app *App) *TaskRunner {
var logger *slog.Logger
if app != nil {
logger = app.Logger
}
return &TaskRunner{
app: app,
logger: logger,
}
}
collector := newTaskRunCollector(details.Title, iType, logger)
collector.log("info", "task run started")
func (r *TaskRunner) Run(ctx context.Context, t Tasker, i Interface, iType InterfaceType) (runErr error) {
config = config.WithActionLogger(func(action string) {
collector.recordUserAction(action)
})
config := t.Config()
details := t.Details()
lifecycle := NewRunLifecycle(r.app, config, details, iType, r.logger)
defer func() {
if runErr != nil {
collector.setError(runErr)
collector.log("error", runErr.Error())
}
run := collector.finalize()
if err := appendTaskMetrics(config.StatisticsStoragePath, details.Title, run, logger); err != nil {
if runErr == nil {
runErr = fmt.Errorf("failed to persist task metrics: %w", err)
return
}
if logger != nil {
logger.Warn("failed to persist task metrics", "error", err, "task", details.Title)
}
}
if app != nil && app.Stats != nil {
if err := app.Stats.RecordTaskRun(ctx, run); err != nil {
if runErr == nil {
runErr = fmt.Errorf("failed to update global statistics: %w", err)
return
}
if logger != nil {
logger.Warn("failed to update global statistics", "error", err, "task", details.Title)
}
}
}
runErr = lifecycle.Finish(ctx, runErr)
}()
doneChan := make(chan bool, 1)
quitChan := make(chan bool, 1)
collector := lifecycle.collector
config = lifecycle.config
collector.log("info", "task run started")
// Setup
if err := t.Setup(ctx); err != nil {
return fmt.Errorf("failed to setup task: %w", err)
}
// Intro screen
if err := runIntro(details); err != nil {
return err
}
// Validation
feedbackChan := make(chan ValidationFeedback, 10)
quitChan := make(chan bool, 1)
if validated, ok := i.(ValidatedInterface); ok {
validated.SetChannels(feedbackChan, quitChan)
}
// Setup task
collector.log("info", "setting up task")
if err := t.Setup(ctx); err != nil {
return fmt.Errorf("failed to setup task: %w", err)
}
collector.log("info", "task setup completed")
engine := &ValidationEngine{task: t}
doneChan, stopChan := engine.Start(ctx, func(feedback ValidationFeedback) {
collector.recordValidation(feedback)
// Show task intro
collector.log("info", "showing task intro")
detailsScreen := NewTaskModel(details)
model, err := tea.NewProgram(detailsScreen, tea.WithAltScreen()).Run()
if err != nil {
return err
}
if m, ok := model.(interface{ GetUserQuit() bool }); ok && m.GetUserQuit() {
collector.log("info", "user quit during task intro")
return ErrUserQuit
}
collector.log("info", "task intro completed")
if feedback.Success {
feedback.Message = "Task completed successfully!"
} else if feedback.Message == "" {
feedback.Message = "Task not completed!"
}
// Start validation loop
collector.log("info", "starting validation loop")
go startValidationLoop(ctx, t, feedbackChan, doneChan, quitChan, collector.recordValidation)
select {
case feedbackChan <- feedback:
default:
}
})
// Run interface
collector.log("info", "starting task interface")
interfaceDone := make(chan error, 1)
interfaceErr := make(chan error, 1)
go func() {
interfaceDone <- i.Run(ctx, config)
interfaceErr <- i.Run(ctx, config)
}()
select {
case <-doneChan:
close(stopChan)
close(quitChan)
collector.setCompleted(true)
collector.log("info", "task validation completed")
if err := <-interfaceDone; err != nil {
if logger != nil {
logger.Warn("interface error after task completion", "error", err, "task", details.Title)
}
collector.log("warn", fmt.Sprintf("interface error after task completion: %v", err))
}
fmt.Println("Task completed successfully!")
case err := <-interfaceDone:
case err := <-interfaceErr:
close(stopChan)
close(quitChan)
if err != nil {
collector.log("error", fmt.Sprintf("task interface failed: %v", err))
return fmt.Errorf("failed to start task interface: %w", err)
return fmt.Errorf("task interface failed: %w", err)
}
collector.log("info", "task interface exited before completion")
fmt.Println("Task incomplete - you exited early")
}
// Show questionnaire
collector.log("info", "showing post-task questionnaire")
questions := t.Questions(iType)
questionnaireKeys := []string{}
if provider, ok := t.(QuestionnaireKeysProvider); ok {
questionnaireKeys = provider.QuestionnaireKeys(iType)
// Questionnaire
if err := runQuestionnaire(t, iType, collector); err != nil {
return err
}
questionare := NewQuestionnaireModel(questions, questionnaireKeys)
model, err = tea.NewProgram(questionare, tea.WithAltScreen()).Run()
collector.log("info", "task run finished")
return nil
}
func (r *RunLifecycle) Finish(ctx context.Context, runErr error) error {
if runErr != nil {
r.collector.setError(runErr)
r.collector.log("error", runErr.Error())
}
run := r.collector.finalize()
if r.metricsStore != nil {
if err := r.metricsStore.Append(ctx, r.details.Title, run); err != nil && runErr == nil {
return fmt.Errorf("persist metrics: %w", err)
}
}
if r.app != nil && r.app.Stats != nil {
if err := r.app.Stats.RecordTaskRun(ctx, run); err != nil && runErr == nil {
return fmt.Errorf("update global stats: %w", err)
}
}
return runErr
}
func runIntro(details models.TaskDetails) error {
model, err := tea.NewProgram(NewTaskModel(details), tea.WithAltScreen()).Run()
if err != nil {
return err
}
questionnaireCompleted := false
questionnaireAnswers := map[string]any(nil)
if m, ok := model.(interface {
GetCompleted() bool
GetAnswers() map[string]any
}); ok {
questionnaireCompleted = m.GetCompleted()
questionnaireAnswers = m.GetAnswers()
if m, ok := model.(interface{ GetUserQuit() bool }); ok {
if m.GetUserQuit() {
return ErrUserQuit
}
}
if m, ok := model.(interface{ GetUserQuit() bool }); ok && m.GetUserQuit() {
collector.recordQuestionnaire(questionnaireCompleted, true, questionnaireAnswers)
collector.log("info", "user quit during questionnaire")
return ErrUserQuit
}
collector.recordQuestionnaire(questionnaireCompleted, false, questionnaireAnswers)
collector.log("info", "task run finished")
return nil
}
func startValidationLoop(ctx context.Context, t Tasker, feedbackChan chan ValidationFeedback, doneChan chan bool, quitChan chan bool, onFeedback func(ValidationFeedback)) {
ticker := time.NewTicker(1 * time.Second)
defer ticker.Stop()
func runQuestionnaire(t Tasker, iType InterfaceType, collector *taskRunCollector) error {
questions := t.Questions(iType)
for {
select {
case <-ticker.C:
feedback := t.Validate(ctx)
if onFeedback != nil {
onFeedback(feedback)
}
if feedback.Success {
feedback.Message = "Task completed successfully!"
feedbackChan <- feedback
time.Sleep(4 * time.Second)
doneChan <- true
return
}
feedback.Message = "Task not completed!"
feedbackChan <- feedback
case <-quitChan:
return
case <-ctx.Done():
return
}
keys := []string{}
if provider, ok := t.(QuestionnaireKeysProvider); ok {
keys = provider.QuestionnaireKeys(iType)
}
model, err := tea.NewProgram(NewQuestionnaireModel(questions, keys), tea.WithAltScreen()).Run()
if err != nil {
return err
}
var completed bool
var answers map[string]any
if m, ok := model.(interface {
GetCompleted() bool
GetAnswers() map[string]any
}); ok {
completed = m.GetCompleted()
answers = m.GetAnswers()
}
userQuit := false
if m, ok := model.(interface{ GetUserQuit() bool }); ok {
userQuit = m.GetUserQuit()
}
collector.recordQuestionnaire(completed, userQuit, answers)
if userQuit {
return ErrUserQuit
}
return nil
}

45
pkg/task/validation.go Normal file
View File

@@ -0,0 +1,45 @@
package task
import (
"context"
"time"
)
type ValidationEngine struct {
task Tasker
}
func (v *ValidationEngine) Start(ctx context.Context, onFeedback func(ValidationFeedback)) (done <-chan struct{}, stop chan<- struct{}) {
doneChan := make(chan struct{}, 1)
stopChan := make(chan struct{}, 1)
go func() {
ticker := time.NewTicker(time.Second)
defer ticker.Stop()
for {
select {
case <-ticker.C:
feedback := v.task.Validate(ctx)
if onFeedback != nil {
onFeedback(feedback)
}
if feedback.Success {
time.Sleep(3 * time.Second)
doneChan <- struct{}{}
return
}
case <-stopChan:
return
case <-ctx.Done():
return
}
}
}()
return doneChan, stopChan
}