diff --git a/go.mod b/go.mod index f03c868..4f51de8 100644 --- a/go.mod +++ b/go.mod @@ -31,6 +31,8 @@ require ( github.com/xeipuuv/gojsonreference v0.0.0-20180127040603-bd5ef7bd5415 // indirect github.com/yosida95/uritemplate/v3 v3.0.2 // indirect github.com/yusufpapurcu/wmi v1.2.4 // indirect + go.uber.org/multierr v1.11.0 // indirect + go.uber.org/zap v1.27.0 // indirect golang.org/x/sys v0.25.0 // indirect gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c // indirect ) diff --git a/go.sum b/go.sum index 30fba1d..a22cc85 100644 --- a/go.sum +++ b/go.sum @@ -70,6 +70,10 @@ github.com/yosida95/uritemplate/v3 v3.0.2 h1:Ed3Oyj9yrmi9087+NczuL5BwkIc4wvTb5zI github.com/yosida95/uritemplate/v3 v3.0.2/go.mod h1:ILOh0sOhIJR3+L/8afwt/kE++YT040gmv5BQTMR2HP4= github.com/yusufpapurcu/wmi v1.2.4 h1:zFUKzehAFReQwLys1b/iSMl+JQGSCSjtVqQn9bBrPo0= github.com/yusufpapurcu/wmi v1.2.4/go.mod h1:SBZ9tNy3G9/m5Oi98Zks0QjeHVDvuK0qfxQmPyzfmi0= +go.uber.org/multierr v1.11.0 h1:blXXJkSxSSfBVBlC76pxqeO+LN3aDfLQo+309xJstO0= +go.uber.org/multierr v1.11.0/go.mod h1:20+QtiLqy0Nd6FdQB9TLXag12DsQkrbs3htMFfDN80Y= +go.uber.org/zap v1.27.0 h1:aJMhYGrd5QSmlpLMr2MftRKl7t8J8PTZPA732ud/XR8= +go.uber.org/zap v1.27.0/go.mod h1:GB2qFLM7cTU87MWRP2mPIjqfIDnGu+VIO4V/SdhGo2E= golang.org/x/sys v0.0.0-20190916202348-b4ddaad3f8a3/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20201204225414-ed752295db88/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20220811171246-fbc7d0a398ab/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= diff --git a/internal/commands/run.go b/internal/commands/run.go index 02fa28d..f54d383 100644 --- a/internal/commands/run.go +++ b/internal/commands/run.go @@ -20,6 +20,7 @@ type RunCommand struct { workflowFile string prompt bool mcpServerURI string + logger common.FileLogger } // NewRunCommand creates a new run command @@ -43,6 +44,7 @@ func NewRunCommand() *cobra.Command { agentsFile = args[0] workflowFile = args[1] } + filelogger, _ := common.NewFileLogger("") runCmd := &RunCommand{ BaseCommand: NewBaseCommand(options), @@ -50,6 +52,7 @@ func NewRunCommand() *cobra.Command { workflowFile: workflowFile, prompt: prompt, mcpServerURI: mcpServerURI, + logger: *filelogger, } return runCmd.Run() @@ -304,4 +307,16 @@ func (c *RunCommand) runWorkflow(workflow common.YAMLDocument, agents []common.Y // logWorkflowRun logs the workflow run func (c *RunCommand) logWorkflowRun(logger *common.Logger, workflowID, workflowName, prompt, output string, modelsUsed []string, status string, startTime, endTime time.Time, durationMs int) { c.Console().Ok(fmt.Sprintf("Workflow %s completed with status: %s", workflowID, status)) + + c.logger.LogWorkflowRun( + workflowID, + workflowName, + prompt, + output, + modelsUsed, + status, + &startTime, + &endTime, + int64(durationMs), + ) } diff --git a/internal/common/file_logger.go b/internal/common/file_logger.go new file mode 100644 index 0000000..ff2550e --- /dev/null +++ b/internal/common/file_logger.go @@ -0,0 +1,381 @@ +// SPDX-License-Identifier: Apache-2.0 +// Copyright © 2025 IBM + +package common + +import ( + "crypto/rand" + "encoding/hex" + "encoding/json" + "fmt" + "os" + "path/filepath" + "strings" + "time" + + "go.uber.org/zap" + "go.uber.org/zap/zapcore" +) + +var ( + // DefaultLogDir is the default directory for log files + DefaultLogDir string +) + +func init() { + homeDir, err := os.UserHomeDir() + if err == nil { + // Check if home directory is writable + if _, err := os.Stat(homeDir); err == nil { + info, err := os.Stat(homeDir) + if err == nil && info.Mode().Perm()&(1<<(uint(7))) != 0 { + DefaultLogDir = filepath.Join(homeDir, ".maestro", "logs") + } else { + DefaultLogDir = "./logs" + } + } else { + DefaultLogDir = "./logs" + } + } else { + DefaultLogDir = "./logs" + } +} + +// generateUUID generates a random UUID-like string +func generateUUID() string { + b := make([]byte, 16) + _, err := rand.Read(b) + if err != nil { + // If we can't generate random bytes, use timestamp as fallback + return fmt.Sprintf("%x", time.Now().UnixNano()) + } + return hex.EncodeToString(b) +} + +// FileLogger handles logging of workflow and agent activities to files +type FileLogger struct { + LogDir string + loggers map[string]*zap.Logger // Map of workflow IDs to loggers +} + +// NewFileLogger creates a new FileLogger instance +func NewFileLogger(logDir string) (*FileLogger, error) { + dir := logDir + if dir == "" { + dir = DefaultLogDir + } + + // Create log directory if it doesn't exist + if err := os.MkdirAll(dir, 0755); err != nil { + return nil, fmt.Errorf("failed to create log directory: %w", err) + } + + return &FileLogger{ + LogDir: dir, + loggers: make(map[string]*zap.Logger), + }, nil +} + +// createLogger creates a new zap logger for a specific workflow +func (l *FileLogger) createLogger(workflowID string) (*zap.Logger, error) { + logPath := filepath.Join(l.LogDir, fmt.Sprintf("maestro_run_%s.jsonl", workflowID)) + + // Create encoder config for JSON format + encoderConfig := zapcore.EncoderConfig{ + TimeKey: "timestamp", + LevelKey: zapcore.OmitKey, // Omit log level as it's not in original format + NameKey: zapcore.OmitKey, + CallerKey: zapcore.OmitKey, + FunctionKey: zapcore.OmitKey, + MessageKey: zapcore.OmitKey, // We'll use custom fields instead of message + StacktraceKey: zapcore.OmitKey, + LineEnding: zapcore.DefaultLineEnding, + EncodeLevel: zapcore.LowercaseLevelEncoder, + EncodeTime: zapcore.ISO8601TimeEncoder, + EncodeDuration: zapcore.MillisDurationEncoder, + EncodeCaller: zapcore.ShortCallerEncoder, + } + + // Create file for logging + file, err := os.OpenFile(logPath, os.O_APPEND|os.O_CREATE|os.O_WRONLY, 0644) + if err != nil { + return nil, fmt.Errorf("failed to open log file: %w", err) + } + + // Create core with JSON encoder and file writer + core := zapcore.NewCore( + zapcore.NewJSONEncoder(encoderConfig), + zapcore.AddSync(file), + zap.InfoLevel, + ) + + // Create logger + return zap.New(core), nil +} + +// getLogger gets or creates a logger for the specified workflow +func (l *FileLogger) getLogger(workflowID string) (*zap.Logger, error) { + if logger, ok := l.loggers[workflowID]; ok { + return logger, nil + } + + logger, err := l.createLogger(workflowID) + if err != nil { + return nil, err + } + + l.loggers[workflowID] = logger + return logger, nil +} + +// Close closes all loggers and releases resources +func (l *FileLogger) Close() { + for _, logger := range l.loggers { + // Sync ensures all buffered logs are written + _ = logger.Sync() + } + l.loggers = make(map[string]*zap.Logger) +} + +// GenerateWorkflowID generates a unique workflow ID +func (l *FileLogger) GenerateWorkflowID() string { + return generateUUID() +} + +// writeJSONLine writes a JSON line to the specified log file +// Kept for backward compatibility with tests +func (l *FileLogger) writeJSONLine(logPath string, data interface{}) error { + // For backward compatibility with tests, use the direct file approach + jsonData, err := json.Marshal(data) + if err != nil { + return fmt.Errorf("failed to marshal JSON: %w", err) + } + + f, err := os.OpenFile(logPath, os.O_APPEND|os.O_CREATE|os.O_WRONLY, 0644) + if err != nil { + return fmt.Errorf("failed to open log file: %w", err) + } + defer f.Close() + + if _, err := f.Write(jsonData); err != nil { + return fmt.Errorf("failed to write to log file: %w", err) + } + if _, err := f.WriteString("\n"); err != nil { + return fmt.Errorf("failed to write newline to log file: %w", err) + } + + return nil +} + +// writeJSONLineWithZap writes a JSON line to the specified log file using zap +// This is an internal method used by the new implementation +func (l *FileLogger) writeJSONLineWithZap(logPath string, data interface{}) error { + // Extract the workflow ID from the log path + base := filepath.Base(logPath) + // Expected format: maestro_run_{workflowID}.jsonl + workflowID := "" + prefix := "maestro_run_" + suffix := ".jsonl" + + if len(base) > len(prefix) && strings.HasPrefix(base, prefix) && strings.HasSuffix(base, suffix) { + workflowID = base[len(prefix) : len(base)-len(suffix)] + } else { + // If we can't extract the workflow ID, create a temporary logger + encoderConfig := zapcore.EncoderConfig{ + TimeKey: "timestamp", + LevelKey: zapcore.OmitKey, + NameKey: zapcore.OmitKey, + CallerKey: zapcore.OmitKey, + FunctionKey: zapcore.OmitKey, + MessageKey: zapcore.OmitKey, + StacktraceKey: zapcore.OmitKey, + LineEnding: zapcore.DefaultLineEnding, + EncodeLevel: zapcore.LowercaseLevelEncoder, + EncodeTime: zapcore.ISO8601TimeEncoder, + EncodeDuration: zapcore.MillisDurationEncoder, + EncodeCaller: zapcore.ShortCallerEncoder, + } + + file, err := os.OpenFile(logPath, os.O_APPEND|os.O_CREATE|os.O_WRONLY, 0644) + if err != nil { + return fmt.Errorf("failed to open log file: %w", err) + } + defer file.Close() + + core := zapcore.NewCore( + zapcore.NewJSONEncoder(encoderConfig), + zapcore.AddSync(file), + zap.InfoLevel, + ) + + logger := zap.New(core) + defer logger.Sync() + + // Convert data to zap fields + jsonData, err := json.Marshal(data) + if err != nil { + return fmt.Errorf("failed to marshal JSON: %w", err) + } + + var fields map[string]interface{} + if err := json.Unmarshal(jsonData, &fields); err != nil { + return fmt.Errorf("failed to unmarshal JSON: %w", err) + } + + zapFields := make([]zap.Field, 0, len(fields)) + for k, v := range fields { + zapFields = append(zapFields, zap.Any(k, v)) + } + + logger.Info("", zapFields...) + return nil + } + + // Get or create a logger for this workflow + logger, err := l.getLogger(workflowID) + if err != nil { + return err + } + + // Convert data to zap fields + jsonData, err := json.Marshal(data) + if err != nil { + return fmt.Errorf("failed to marshal JSON: %w", err) + } + + var fields map[string]interface{} + if err := json.Unmarshal(jsonData, &fields); err != nil { + return fmt.Errorf("failed to unmarshal JSON: %w", err) + } + + zapFields := make([]zap.Field, 0, len(fields)) + for k, v := range fields { + zapFields = append(zapFields, zap.Any(k, v)) + } + + logger.Info("", zapFields...) + return nil +} + +// TokenUsage represents token usage information +type TokenUsage struct { + PromptTokens int `json:"prompt_tokens,omitempty"` + CompletionTokens int `json:"completion_tokens,omitempty"` + TotalTokens int `json:"total_tokens,omitempty"` +} + +// AgentResponseLog represents a log entry for an agent response +type AgentResponseLog struct { + LogType string `json:"log_type"` + Timestamp string `json:"timestamp"` + WorkflowID string `json:"workflow_id"` + StepIndex int `json:"step_index"` + AgentName string `json:"agent_name"` + Model string `json:"model"` + Input string `json:"input"` + Response string `json:"response"` + ToolUsed string `json:"tool_used,omitempty"` + StartTime string `json:"start_time,omitempty"` + EndTime string `json:"end_time,omitempty"` + DurationMS int64 `json:"duration_ms,omitempty"` + TokenUsage *TokenUsage `json:"token_usage,omitempty"` +} + +// LogAgentResponse logs an agent response +func (l *FileLogger) LogAgentResponse( + workflowID string, + stepIndex int, + agentName string, + model string, + inputText string, + responseText string, + toolUsed string, + startTime *time.Time, + endTime *time.Time, + durationMS int64, + tokenUsage *TokenUsage, +) error { + logPath := filepath.Join(l.LogDir, fmt.Sprintf("maestro_run_%s.jsonl", workflowID)) + + var startTimeStr, endTimeStr string + if startTime != nil { + startTimeStr = startTime.UTC().Format(time.RFC3339Nano) + } + if endTime != nil { + endTimeStr = endTime.UTC().Format(time.RFC3339Nano) + } + + data := AgentResponseLog{ + LogType: "agent_response", + Timestamp: time.Now().UTC().Format(time.RFC3339Nano), + WorkflowID: workflowID, + StepIndex: stepIndex, + AgentName: agentName, + Model: model, + Input: inputText, + Response: responseText, + ToolUsed: toolUsed, + StartTime: startTimeStr, + EndTime: endTimeStr, + DurationMS: durationMS, + TokenUsage: tokenUsage, + } + + return l.writeJSONLine(logPath, data) +} + +// WorkflowRunLog represents a log entry for a workflow run +type WorkflowRunLog struct { + LogType string `json:"log_type"` + Timestamp string `json:"timestamp"` + WorkflowID string `json:"workflow_id"` + WorkflowName string `json:"workflow_name"` + Status string `json:"status"` + Prompt string `json:"prompt"` + Output string `json:"output"` + ModelsUsed []string `json:"models_used"` + StartTime string `json:"start_time,omitempty"` + EndTime string `json:"end_time,omitempty"` + DurationMS int64 `json:"duration_ms,omitempty"` +} + +// LogWorkflowRun logs a workflow run +func (l *FileLogger) LogWorkflowRun( + workflowID string, + workflowName string, + prompt string, + output string, + modelsUsed []string, + status string, + startTime *time.Time, + endTime *time.Time, + durationMS int64, +) error { + logPath := filepath.Join(l.LogDir, fmt.Sprintf("maestro_run_%s.jsonl", workflowID)) + + var startTimeStr, endTimeStr string + if startTime != nil { + startTimeStr = startTime.UTC().Format(time.RFC3339Nano) + } + if endTime != nil { + endTimeStr = endTime.UTC().Format(time.RFC3339Nano) + } + + data := WorkflowRunLog{ + LogType: "workflow_summary", + Timestamp: time.Now().UTC().Format(time.RFC3339Nano), + WorkflowID: workflowID, + WorkflowName: workflowName, + Status: status, + Prompt: prompt, + Output: output, + ModelsUsed: modelsUsed, + StartTime: startTimeStr, + EndTime: endTimeStr, + DurationMS: durationMS, + } + + return l.writeJSONLine(logPath, data) +} + +// Made with Bob diff --git a/internal/common/file_logger_test.go b/internal/common/file_logger_test.go new file mode 100644 index 0000000..d493ada --- /dev/null +++ b/internal/common/file_logger_test.go @@ -0,0 +1,568 @@ +// SPDX-License-Identifier: Apache-2.0 +// Copyright © 2025 IBM + +package common + +import ( + "encoding/json" + "fmt" + "os" + "path/filepath" + "regexp" + "testing" + "time" +) + +// readJSONLine reads the first line from a JSONL file and unmarshals it into the provided interface +func readJSONLine(filePath string, v interface{}) error { + content, err := os.ReadFile(filePath) + if err != nil { + return fmt.Errorf("failed to read file: %w", err) + } + + // Find the first newline character + data := content + for i, b := range content { + if b == '\n' { + data = content[:i] + break + } + } + + // Unmarshal the JSON data + if err := json.Unmarshal(data, v); err != nil { + return fmt.Errorf("failed to unmarshal JSON: %w", err) + } + + return nil +} + +func TestNewFileLogger(t *testing.T) { + // Create a temporary directory for testing + tempDir, err := os.MkdirTemp("", "file_logger_test") + if err != nil { + t.Fatalf("Failed to create temp directory: %v", err) + } + defer os.RemoveAll(tempDir) + + // Test with custom log directory + t.Run("CustomLogDir", func(t *testing.T) { + logger, err := NewFileLogger(tempDir) + if err != nil { + t.Fatalf("NewFileLogger failed: %v", err) + } + if logger.LogDir != tempDir { + t.Errorf("Expected LogDir to be %s, got %s", tempDir, logger.LogDir) + } + // Check that loggers map is initialized + if logger.loggers == nil { + t.Error("Expected loggers map to be initialized") + } + }) + + // Test with default log directory + t.Run("DefaultLogDir", func(t *testing.T) { + // Save original DefaultLogDir + origDefaultLogDir := DefaultLogDir + defer func() { DefaultLogDir = origDefaultLogDir }() + + // Set DefaultLogDir to our temp directory for testing + DefaultLogDir = tempDir + + logger, err := NewFileLogger("") + if err != nil { + t.Fatalf("NewFileLogger failed: %v", err) + } + if logger.LogDir != tempDir { + t.Errorf("Expected LogDir to be %s, got %s", tempDir, logger.LogDir) + } + }) + + // Test error handling when directory creation fails + t.Run("DirectoryCreationFails", func(t *testing.T) { + // Create a file with the same name as our intended directory + // This will cause MkdirAll to fail + filePath := filepath.Join(tempDir, "cannot_create_dir") + file, err := os.Create(filePath) + if err != nil { + t.Fatalf("Failed to create test file: %v", err) + } + file.Close() + + // Try to create a logger with a directory path that conflicts with the file + _, err = NewFileLogger(filePath) + if err == nil { + t.Error("Expected an error when directory creation fails, got nil") + } + }) +} + +func TestGenerateWorkflowID(t *testing.T) { + // Create a logger for testing + tempDir, err := os.MkdirTemp("", "file_logger_test") + if err != nil { + t.Fatalf("Failed to create temp directory: %v", err) + } + defer os.RemoveAll(tempDir) + + logger, err := NewFileLogger(tempDir) + if err != nil { + t.Fatalf("NewFileLogger failed: %v", err) + } + + // Generate multiple IDs and check they're unique + ids := make(map[string]bool) + for i := 0; i < 10; i++ { + id := logger.GenerateWorkflowID() + + // Check format (should be a hex string of length 32) + matched, err := regexp.MatchString("^[0-9a-f]{32}$", id) + if err != nil { + t.Fatalf("Regex match failed: %v", err) + } + if !matched { + t.Errorf("Generated ID %s doesn't match expected format", id) + } + + // Check uniqueness + if ids[id] { + t.Errorf("Generated duplicate ID: %s", id) + } + + ids[id] = true + } +} + +func TestWriteJSONLine(t *testing.T) { + // Create a temporary directory for testing + tempDir, err := os.MkdirTemp("", "file_logger_test") + if err != nil { + t.Fatalf("Failed to create temp directory: %v", err) + } + defer os.RemoveAll(tempDir) + + logger, err := NewFileLogger(tempDir) + if err != nil { + t.Fatalf("NewFileLogger failed: %v", err) + } + + // Test successful write + t.Run("SuccessfulWrite", func(t *testing.T) { + logPath := filepath.Join(tempDir, "test_log.jsonl") + testData := map[string]string{"key": "value"} + + err := logger.writeJSONLine(logPath, testData) + if err != nil { + t.Fatalf("writeJSONLine failed: %v", err) + } + + // Read the file and verify content + content, err := os.ReadFile(logPath) + if err != nil { + t.Fatalf("Failed to read log file: %v", err) + } + + expectedJSON, _ := json.Marshal(testData) + expectedContent := string(expectedJSON) + "\n" + if string(content) != expectedContent { + t.Errorf("Expected content %q, got %q", expectedContent, string(content)) + } + }) + + // Test error handling for JSON marshaling + t.Run("JSONMarshalError", func(t *testing.T) { + logPath := filepath.Join(tempDir, "test_log_marshal_error.jsonl") + + // Create a struct with a channel which cannot be marshaled to JSON + type UnmarshalableStruct struct { + Ch chan int + } + unmarshalable := UnmarshalableStruct{Ch: make(chan int)} + + err := logger.writeJSONLine(logPath, unmarshalable) + if err == nil { + t.Error("Expected an error when marshaling invalid JSON, got nil") + } + }) + + // Test error handling for file operations + t.Run("FileOperationError", func(t *testing.T) { + // Create a directory with the same name as our intended file + // This will cause OpenFile to fail + dirPath := filepath.Join(tempDir, "dir_not_file") + err := os.Mkdir(dirPath, 0755) + if err != nil { + t.Fatalf("Failed to create test directory: %v", err) + } + + err = logger.writeJSONLine(dirPath, map[string]string{"key": "value"}) + if err == nil { + t.Error("Expected an error when file operation fails, got nil") + } + }) +} + +func TestLogAgentResponse(t *testing.T) { + // Create a temporary directory for testing + tempDir, err := os.MkdirTemp("", "file_logger_test") + if err != nil { + t.Fatalf("Failed to create temp directory: %v", err) + } + defer os.RemoveAll(tempDir) + + logger, err := NewFileLogger(tempDir) + if err != nil { + t.Fatalf("NewFileLogger failed: %v", err) + } + + workflowID := "test-workflow-id" + stepIndex := 1 + agentName := "test-agent" + model := "test-model" + inputText := "test input" + responseText := "test response" + toolUsed := "test-tool" + + // Test with all parameters provided + t.Run("AllParametersProvided", func(t *testing.T) { + startTime := time.Now().Add(-1 * time.Minute) + endTime := time.Now() + durationMS := int64(endTime.Sub(startTime).Milliseconds()) + tokenUsage := &TokenUsage{ + PromptTokens: 10, + CompletionTokens: 20, + TotalTokens: 30, + } + + err := logger.LogAgentResponse( + workflowID, + stepIndex, + agentName, + model, + inputText, + responseText, + toolUsed, + &startTime, + &endTime, + durationMS, + tokenUsage, + ) + if err != nil { + t.Fatalf("LogAgentResponse failed: %v", err) + } + + // Verify log file exists + logPath := filepath.Join(logger.LogDir, "maestro_run_"+workflowID+".jsonl") + if _, err := os.Stat(logPath); os.IsNotExist(err) { + t.Fatalf("Log file was not created: %v", err) + } + + // Read and parse the log file + var logEntry AgentResponseLog + err = readJSONLine(logPath, &logEntry) + if err != nil { + t.Fatalf("Failed to parse log entry: %v", err) + } + + // Verify log entry fields + if logEntry.LogType != "agent_response" { + t.Errorf("Expected LogType to be 'agent_response', got %s", logEntry.LogType) + } + if logEntry.WorkflowID != workflowID { + t.Errorf("Expected WorkflowID to be %s, got %s", workflowID, logEntry.WorkflowID) + } + if logEntry.StepIndex != stepIndex { + t.Errorf("Expected StepIndex to be %d, got %d", stepIndex, logEntry.StepIndex) + } + if logEntry.AgentName != agentName { + t.Errorf("Expected AgentName to be %s, got %s", agentName, logEntry.AgentName) + } + if logEntry.Model != model { + t.Errorf("Expected Model to be %s, got %s", model, logEntry.Model) + } + if logEntry.Input != inputText { + t.Errorf("Expected Input to be %s, got %s", inputText, logEntry.Input) + } + if logEntry.Response != responseText { + t.Errorf("Expected Response to be %s, got %s", responseText, logEntry.Response) + } + if logEntry.ToolUsed != toolUsed { + t.Errorf("Expected ToolUsed to be %s, got %s", toolUsed, logEntry.ToolUsed) + } + if logEntry.DurationMS != durationMS { + t.Errorf("Expected DurationMS to be %d, got %d", durationMS, logEntry.DurationMS) + } + if logEntry.TokenUsage.PromptTokens != tokenUsage.PromptTokens { + t.Errorf("Expected TokenUsage.PromptTokens to be %d, got %d", tokenUsage.PromptTokens, logEntry.TokenUsage.PromptTokens) + } + }) + + // Test with nil optional parameters + t.Run("NilOptionalParameters", func(t *testing.T) { + // Remove previous log file if it exists + logPath := filepath.Join(logger.LogDir, "maestro_run_"+workflowID+".jsonl") + os.Remove(logPath) + + err := logger.LogAgentResponse( + workflowID, + stepIndex, + agentName, + model, + inputText, + responseText, + toolUsed, + nil, // startTime + nil, // endTime + 0, // durationMS + nil, // tokenUsage + ) + if err != nil { + t.Fatalf("LogAgentResponse failed: %v", err) + } + + // Verify log file exists + if _, err := os.Stat(logPath); os.IsNotExist(err) { + t.Fatalf("Log file was not created: %v", err) + } + + // Read and parse the log file + content, err := os.ReadFile(logPath) + if err != nil { + t.Fatalf("Failed to read log file: %v", err) + } + + var logEntry AgentResponseLog + err = json.Unmarshal(content, &logEntry) + if err != nil { + t.Fatalf("Failed to parse log entry: %v", err) + } + + // Verify log entry fields + if logEntry.StartTime != "" { + t.Errorf("Expected StartTime to be empty, got %s", logEntry.StartTime) + } + if logEntry.EndTime != "" { + t.Errorf("Expected EndTime to be empty, got %s", logEntry.EndTime) + } + if logEntry.TokenUsage != nil { + t.Errorf("Expected TokenUsage to be nil, got %+v", logEntry.TokenUsage) + } + }) +} + +func TestLogWorkflowRun(t *testing.T) { + // Create a temporary directory for testing + tempDir, err := os.MkdirTemp("", "file_logger_test") + if err != nil { + t.Fatalf("Failed to create temp directory: %v", err) + } + defer os.RemoveAll(tempDir) + + logger, err := NewFileLogger(tempDir) + if err != nil { + t.Fatalf("NewFileLogger failed: %v", err) + } + + workflowID := "test-workflow-id" + workflowName := "test-workflow" + prompt := "test prompt" + output := "test output" + modelsUsed := []string{"model1", "model2"} + status := "completed" + + // Test with all parameters provided + t.Run("AllParametersProvided", func(t *testing.T) { + startTime := time.Now().Add(-1 * time.Minute) + endTime := time.Now() + durationMS := int64(endTime.Sub(startTime).Milliseconds()) + + err := logger.LogWorkflowRun( + workflowID, + workflowName, + prompt, + output, + modelsUsed, + status, + &startTime, + &endTime, + durationMS, + ) + if err != nil { + t.Fatalf("LogWorkflowRun failed: %v", err) + } + + // Verify log file exists + logPath := filepath.Join(logger.LogDir, "maestro_run_"+workflowID+".jsonl") + if _, err := os.Stat(logPath); os.IsNotExist(err) { + t.Fatalf("Log file was not created: %v", err) + } + + // Read and parse the log file + var logEntry WorkflowRunLog + err = readJSONLine(logPath, &logEntry) + if err != nil { + t.Fatalf("Failed to parse log entry: %v", err) + } + + // Verify log entry fields + if logEntry.LogType != "workflow_summary" { + t.Errorf("Expected LogType to be 'workflow_summary', got %s", logEntry.LogType) + } + if logEntry.WorkflowID != workflowID { + t.Errorf("Expected WorkflowID to be %s, got %s", workflowID, logEntry.WorkflowID) + } + if logEntry.WorkflowName != workflowName { + t.Errorf("Expected WorkflowName to be %s, got %s", workflowName, logEntry.WorkflowName) + } + if logEntry.Status != status { + t.Errorf("Expected Status to be %s, got %s", status, logEntry.Status) + } + if logEntry.Prompt != prompt { + t.Errorf("Expected Prompt to be %s, got %s", prompt, logEntry.Prompt) + } + if logEntry.Output != output { + t.Errorf("Expected Output to be %s, got %s", output, logEntry.Output) + } + if len(logEntry.ModelsUsed) != len(modelsUsed) { + t.Errorf("Expected ModelsUsed length to be %d, got %d", len(modelsUsed), len(logEntry.ModelsUsed)) + } + if logEntry.DurationMS != durationMS { + t.Errorf("Expected DurationMS to be %d, got %d", durationMS, logEntry.DurationMS) + } + }) + + // Test with nil optional parameters + t.Run("NilOptionalParameters", func(t *testing.T) { + // Remove previous log file if it exists + logPath := filepath.Join(logger.LogDir, "maestro_run_"+workflowID+".jsonl") + os.Remove(logPath) + + err := logger.LogWorkflowRun( + workflowID, + workflowName, + prompt, + output, + modelsUsed, + status, + nil, // startTime + nil, // endTime + 0, // durationMS + ) + if err != nil { + t.Fatalf("LogWorkflowRun failed: %v", err) + } + + // Verify log file exists + if _, err := os.Stat(logPath); os.IsNotExist(err) { + t.Fatalf("Log file was not created: %v", err) + } + + // Read and parse the log file + content, err := os.ReadFile(logPath) + if err != nil { + t.Fatalf("Failed to read log file: %v", err) + } + + var logEntry WorkflowRunLog + err = json.Unmarshal(content, &logEntry) + if err != nil { + t.Fatalf("Failed to parse log entry: %v", err) + } + + // Verify log entry fields + if logEntry.StartTime != "" { + t.Errorf("Expected StartTime to be empty, got %s", logEntry.StartTime) + } + if logEntry.EndTime != "" { + t.Errorf("Expected EndTime to be empty, got %s", logEntry.EndTime) + } + }) +} + +// Made with Bob + +func TestGetLogger(t *testing.T) { + // Create a temporary directory for testing + tempDir, err := os.MkdirTemp("", "file_logger_test") + if err != nil { + t.Fatalf("Failed to create temp directory: %v", err) + } + defer os.RemoveAll(tempDir) + + logger, err := NewFileLogger(tempDir) + if err != nil { + t.Fatalf("NewFileLogger failed: %v", err) + } + defer logger.Close() // Clean up resources + + // Test getting a new logger + t.Run("GetNewLogger", func(t *testing.T) { + workflowID := "test-workflow" + zapLogger, err := logger.getLogger(workflowID) + if err != nil { + t.Fatalf("getLogger failed: %v", err) + } + if zapLogger == nil { + t.Error("Expected non-nil logger") + } + + // Check that logger was cached + if cachedLogger, ok := logger.loggers[workflowID]; !ok || cachedLogger != zapLogger { + t.Error("Logger was not properly cached") + } + }) + + // Test getting an existing logger + t.Run("GetExistingLogger", func(t *testing.T) { + workflowID := "test-workflow-2" + firstLogger, err := logger.getLogger(workflowID) + if err != nil { + t.Fatalf("First getLogger failed: %v", err) + } + + secondLogger, err := logger.getLogger(workflowID) + if err != nil { + t.Fatalf("Second getLogger failed: %v", err) + } + + if firstLogger != secondLogger { + t.Error("Expected same logger instance to be returned") + } + }) +} + +func TestClose(t *testing.T) { + // Create a temporary directory for testing + tempDir, err := os.MkdirTemp("", "file_logger_test") + if err != nil { + t.Fatalf("Failed to create temp directory: %v", err) + } + defer os.RemoveAll(tempDir) + + logger, err := NewFileLogger(tempDir) + if err != nil { + t.Fatalf("NewFileLogger failed: %v", err) + } + + // Create some loggers + workflowIDs := []string{"workflow1", "workflow2", "workflow3"} + for _, id := range workflowIDs { + _, err := logger.getLogger(id) + if err != nil { + t.Fatalf("Failed to get logger for %s: %v", id, err) + } + } + + // Verify loggers exist + if len(logger.loggers) != len(workflowIDs) { + t.Errorf("Expected %d loggers, got %d", len(workflowIDs), len(logger.loggers)) + } + + // Close loggers + logger.Close() + + // Verify loggers map is empty + if len(logger.loggers) != 0 { + t.Errorf("Expected empty loggers map after Close, got %d entries", len(logger.loggers)) + } +}