From 46bd4b01e3364a88d94c594651f5641106538f6d Mon Sep 17 00:00:00 2001 From: Emre K <110906681+kocaemre@users.noreply.github.com> Date: Tue, 28 Jul 2026 16:56:31 +0200 Subject: [PATCH] Close workflow audit log writers Signed-off-by: Emre K <110906681+kocaemre@users.noreply.github.com> --- pkg/audit/auditor.go | 10 ++++- pkg/audit/workflow_auditor.go | 9 +++++ pkg/audit/workflow_auditor_test.go | 63 ++++++++++++++++++++++++++++++ pkg/vmcp/core/core_vmcp.go | 21 ++++++++++ 4 files changed, 102 insertions(+), 1 deletion(-) diff --git a/pkg/audit/auditor.go b/pkg/audit/auditor.go index 330fb42e60..192be2c5d2 100644 --- a/pkg/audit/auditor.go +++ b/pkg/audit/auditor.go @@ -107,7 +107,15 @@ func NewAuditorWithTransport(config *Config, transportType string) (*Auditor, er // Close closes the underlying log writer if it implements io.Closer. // This should be called when the auditor is no longer needed to properly release resources. func (a *Auditor) Close() error { - if closer, ok := a.logWriter.(io.Closer); ok { + return closeLogWriter(a.logWriter) +} + +func closeLogWriter(logWriter io.Writer) error { + if logWriter == os.Stdout || logWriter == os.Stderr { + return nil + } + + if closer, ok := logWriter.(io.Closer); ok { return closer.Close() } return nil diff --git a/pkg/audit/workflow_auditor.go b/pkg/audit/workflow_auditor.go index 5a12ae1331..e074cc0729 100644 --- a/pkg/audit/workflow_auditor.go +++ b/pkg/audit/workflow_auditor.go @@ -8,6 +8,7 @@ import ( "context" "encoding/json" "fmt" + "io" "log/slog" "time" @@ -21,6 +22,7 @@ type WorkflowAuditor struct { auditLogger *slog.Logger config *Config component string + logWriter io.Writer } // NewWorkflowAuditor creates a new workflow auditor. @@ -45,9 +47,16 @@ func NewWorkflowAuditor(config *Config) (*WorkflowAuditor, error) { auditLogger: NewAuditLogger(logWriter), config: config, component: component, + logWriter: logWriter, }, nil } +// Close closes the underlying log writer if it owns a closeable resource. +// This should be called when the workflow auditor is no longer needed. +func (w *WorkflowAuditor) Close() error { + return closeLogWriter(w.logWriter) +} + // LogWorkflowStarted logs the start of workflow execution. func (w *WorkflowAuditor) LogWorkflowStarted( ctx context.Context, diff --git a/pkg/audit/workflow_auditor_test.go b/pkg/audit/workflow_auditor_test.go index 83da6b66c7..f1800108ab 100644 --- a/pkg/audit/workflow_auditor_test.go +++ b/pkg/audit/workflow_auditor_test.go @@ -23,11 +23,25 @@ type testLogWriter struct { logs []string } +type closeTrackingWriter struct { + closed bool + closeErr error +} + func (w *testLogWriter) Write(p []byte) (n int, err error) { w.logs = append(w.logs, string(p)) return len(p), nil } +func (*closeTrackingWriter) Write(p []byte) (n int, err error) { + return len(p), nil +} + +func (w *closeTrackingWriter) Close() error { + w.closed = true + return w.closeErr +} + func (w *testLogWriter) getLastLog() string { if len(w.logs) == 0 { return "" @@ -52,6 +66,7 @@ func createTestAuditor(t *testing.T, config *Config) (*WorkflowAuditor, *testLog auditLogger: NewAuditLogger(writer), config: config, component: "vmcp-composer", + logWriter: writer, } return auditor, writer @@ -122,6 +137,54 @@ func TestNewWorkflowAuditor(t *testing.T) { } } +func TestWorkflowAuditor_Close(t *testing.T) { + t.Parallel() + + t.Run("closes retained file writer", func(t *testing.T) { + t.Parallel() + + logFilePath := t.TempDir() + "/workflow-audit.log" + auditor, err := NewWorkflowAuditor(&Config{LogFile: logFilePath}) + require.NoError(t, err) + + _, ok := auditor.logWriter.(interface{ Close() error }) + require.True(t, ok, "file-backed workflow auditor should retain a closeable writer") + + require.NoError(t, auditor.Close()) + }) + + t.Run("does not close stdout", func(t *testing.T) { + t.Parallel() + + auditor, err := NewWorkflowAuditor(&Config{}) + require.NoError(t, err) + + require.NoError(t, auditor.Close()) + assert.Same(t, os.Stdout, auditor.logWriter) + }) + + t.Run("propagates close errors", func(t *testing.T) { + t.Parallel() + + closeErr := errors.New("close failed") + writer := &closeTrackingWriter{closeErr: closeErr} + auditor := &WorkflowAuditor{logWriter: writer} + + err := auditor.Close() + + require.ErrorIs(t, err, closeErr) + assert.True(t, writer.closed) + }) +} + +func TestAuditor_CloseDoesNotCloseStdout(t *testing.T) { + t.Parallel() + + auditor := &Auditor{logWriter: os.Stdout} + + require.NoError(t, auditor.Close()) +} + func TestWorkflowAuditor_LogWorkflowStarted(t *testing.T) { t.Parallel() diff --git a/pkg/vmcp/core/core_vmcp.go b/pkg/vmcp/core/core_vmcp.go index 98c5523fca..e2781cf8b8 100644 --- a/pkg/vmcp/core/core_vmcp.go +++ b/pkg/vmcp/core/core_vmcp.go @@ -71,6 +71,10 @@ type coreVMCP struct { // by advertised tool name. workflowDefs map[string]*composer.WorkflowDefinition + // workflowAuditor owns the optional workflow audit log writer and is closed + // with the core. Nil when workflow audit logging is disabled. + workflowAuditor *audit.WorkflowAuditor + // composerFactory builds a per-call composite-tool engine bound to a routing // table, generalizing server.New's sessionComposerFactory (server.go:393). composerFactory func(sessionRT *vmcp.RoutingTable, sessionTools []vmcp.Tool) composer.Composer @@ -138,6 +142,14 @@ func New(cfg *Config) (VMCP, error) { } slog.Info("workflow audit logging enabled") } + closeWorkflowAuditor := func() { + if workflowAuditor == nil { + return + } + if err := workflowAuditor.Close(); err != nil { + slog.Warn("failed to close workflow auditor", "error", err) + } + } // The elicitation handler depends only on the domain-typed ElicitationRequester // (#5436); no mcp-go types cross this boundary (vmcp anti-pattern #5). @@ -168,6 +180,7 @@ func New(cfg *Config) (VMCP, error) { instruments, err := newWorkflowInstruments(cfg.TelemetryProvider) if err != nil { stopStore() + closeWorkflowAuditor() return nil, fmt.Errorf("failed to create workflow telemetry instruments: %w", err) } @@ -195,6 +208,7 @@ func New(cfg *Config) (VMCP, error) { workflowDefs, err := validateWorkflowDefs(validationEngine, cfg.WorkflowDefs) if err != nil { stopStore() + closeWorkflowAuditor() return nil, fmt.Errorf("workflow validation failed: %w", err) } @@ -206,6 +220,7 @@ func New(cfg *Config) (VMCP, error) { healthMonitor, healthProvider, err := buildHealthMonitor(cfg) if err != nil { stopStore() + closeWorkflowAuditor() return nil, err } @@ -217,6 +232,7 @@ func New(cfg *Config) (VMCP, error) { healthMonitor: healthMonitor, admission: admission, workflowDefs: workflowDefs, + workflowAuditor: workflowAuditor, composerFactory: composerFactory, stopStore: stopStore, }, nil @@ -543,6 +559,11 @@ func (c *coreVMCP) InvalidateCapabilityCache() { func (c *coreVMCP) Close() error { c.closeOnce.Do(func() { c.stopStore() + if c.workflowAuditor != nil { + if err := c.workflowAuditor.Close(); err != nil { + slog.Warn("failed to close workflow auditor", "error", err) + } + } if c.healthMonitor != nil { if err := c.healthMonitor.Stop(); err != nil { slog.Warn("failed to stop health monitor", "error", err)